mirror of
https://github.com/alphaXiv/OpenResearch.git
synced 2026-10-02 09:44:51 +08:00
Fix telemetry opt-out, partial accounting, and identity retention
This commit is contained in:
@@ -210,6 +210,7 @@ pub async fn run_antigravity() -> Result<()> {
|
||||
crate::store::Store::open()?.record_native_invocation(
|
||||
&crate::local::harness::antigravity::invocation_sample_id(conversation, step),
|
||||
&identity,
|
||||
Some(conversation),
|
||||
)?;
|
||||
println!("{{}}");
|
||||
return Ok(());
|
||||
|
||||
@@ -7048,9 +7048,13 @@ impl TurnCtx {
|
||||
for part in parts {
|
||||
if part.get("type").and_then(Value::as_str) == Some("tool_use") {
|
||||
if let Some(call_id) = part.get("id").and_then(Value::as_str) {
|
||||
if let Err(error) = Store::open()
|
||||
.and_then(|store| store.record_native_invocation(call_id, &identity))
|
||||
{
|
||||
if let Err(error) = Store::open().and_then(|store| {
|
||||
store.record_native_invocation(
|
||||
call_id,
|
||||
&identity,
|
||||
Some(&self.session_id),
|
||||
)
|
||||
}) {
|
||||
eprintln!("orx up: could not capture native tool identity: {error}");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -665,6 +665,8 @@ impl Store {
|
||||
// Best-effort migrations for pre-existing dbs; re-runs fail with
|
||||
// "duplicate column name", which is exactly the no-op we want.
|
||||
for ddl in [
|
||||
"ALTER TABLE native_invocation_identities ADD COLUMN session_id TEXT",
|
||||
"ALTER TABLE native_invocation_identities ADD COLUMN created_at INTEGER NOT NULL DEFAULT 0",
|
||||
"ALTER TABLE chat_usage_samples ADD COLUMN complete INTEGER NOT NULL DEFAULT 0",
|
||||
"ALTER TABLE runs ADD COLUMN commit_sha TEXT",
|
||||
"ALTER TABLE runs ADD COLUMN result_markdown TEXT",
|
||||
@@ -706,6 +708,11 @@ impl Store {
|
||||
] {
|
||||
let _ = conn.execute(ddl, []);
|
||||
}
|
||||
conn.execute_batch("CREATE TRIGGER IF NOT EXISTS delete_chat_invocation_identities AFTER DELETE ON chat_sessions BEGIN DELETE FROM native_invocation_identities WHERE session_id = OLD.id; END;")?;
|
||||
conn.execute(
|
||||
"DELETE FROM native_invocation_identities WHERE session_id IS NULL AND created_at < ?1",
|
||||
[now_ms() - 7 * 24 * 60 * 60 * 1000],
|
||||
)?;
|
||||
// Legacy tool failures cannot identify the missing dependency, so require one fresh check.
|
||||
conn.execute(
|
||||
"DELETE FROM ssh_host_tests
|
||||
|
||||
+219
-17
@@ -107,9 +107,12 @@ impl Store {
|
||||
&self,
|
||||
call_id: &str,
|
||||
identity: &InvocationIdentity,
|
||||
session_id: Option<&str>,
|
||||
) -> Result<()> {
|
||||
identity.validate()?;
|
||||
self.conn.execute("INSERT INTO native_invocation_identities (harness, call_id, identity_json) VALUES (?1, ?2, ?3) ON CONFLICT(harness, call_id) DO NOTHING", params![identity.harness, call_id, serde_json::to_string(identity)?])?;
|
||||
let tx = self.begin_immediate()?;
|
||||
let owner: Option<String> = self.conn.query_row("SELECT id FROM chat_sessions WHERE harness = ?1 AND (id = ?2 OR native_session_id = ?2)", params![identity.harness, session_id], |row| row.get(0)).optional()?;
|
||||
self.conn.execute("INSERT INTO native_invocation_identities (harness, call_id, identity_json, session_id, created_at) VALUES (?1, ?2, ?3, ?4, ?5) ON CONFLICT(harness, call_id) DO NOTHING", params![identity.harness, call_id, serde_json::to_string(identity)?, owner, now_ms()])?;
|
||||
if self
|
||||
.native_invocation_identity(&identity.harness, call_id)?
|
||||
.as_ref()
|
||||
@@ -117,6 +120,7 @@ impl Store {
|
||||
{
|
||||
return Err(anyhow!("Native tool identity changed after capture"));
|
||||
}
|
||||
tx.commit()?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -363,7 +367,7 @@ impl Store {
|
||||
for (execution_id, harness, report_id, suppressed) in rows {
|
||||
let mut samples = self.conn.prepare("SELECT s.model, s.provider, s.usage_json, s.complete, n.identity_json FROM chat_usage_samples s LEFT JOIN native_invocation_identities n ON n.harness = s.harness AND n.call_id = s.sample_id WHERE s.execution_id = ?1 ORDER BY s.sample_id")?;
|
||||
let mut grouped = std::collections::BTreeMap::<
|
||||
(Option<String>, Option<String>),
|
||||
(Option<String>, Option<String>, [bool; 5]),
|
||||
Vec<(TokenUsage, bool)>,
|
||||
>::new();
|
||||
for row in samples.query_map([&execution_id], |row| {
|
||||
@@ -382,17 +386,26 @@ impl Store {
|
||||
let model =
|
||||
model.or_else(|| identity.as_ref().map(|identity| identity.model.clone()));
|
||||
let provider = provider.or_else(|| identity.and_then(|identity| identity.provider));
|
||||
let usage: TokenUsage = serde_json::from_str(&json)?;
|
||||
let measured = [
|
||||
usage.input_tokens,
|
||||
usage.output_tokens,
|
||||
usage.cache_read_tokens,
|
||||
usage.cache_write_tokens,
|
||||
usage.reasoning_tokens,
|
||||
]
|
||||
.map(|counter| counter.is_some());
|
||||
grouped
|
||||
.entry((model, provider))
|
||||
.entry((model, provider, measured))
|
||||
.or_default()
|
||||
.push((serde_json::from_str(&json)?, complete));
|
||||
.push((usage, complete));
|
||||
}
|
||||
if grouped.is_empty() {
|
||||
grouped.insert((None, None), Vec::new());
|
||||
grouped.insert((None, None, [false; 5]), Vec::new());
|
||||
}
|
||||
let mut reports = Vec::new();
|
||||
if !suppressed {
|
||||
for ((model, provider), samples) in grouped {
|
||||
for ((model, provider, _), samples) in grouped {
|
||||
let sum = |field: fn(&TokenUsage) -> Option<u64>| -> Option<u64> {
|
||||
if samples.is_empty() {
|
||||
return None;
|
||||
@@ -450,18 +463,28 @@ impl Store {
|
||||
}
|
||||
}
|
||||
}
|
||||
let tx = self.begin()?;
|
||||
let applied = tx.execute("UPDATE chat_usage_executions SET outcome = ?2 WHERE execution_id = ?1 AND outcome IS NULL", params![execution_id, outcome])?;
|
||||
if applied == 1 {
|
||||
for (id, payload) in reports {
|
||||
self.stage_telemetry(&id, &payload)?;
|
||||
}
|
||||
}
|
||||
tx.commit()?;
|
||||
self.finalize_usage_execution(&execution_id, outcome, reports)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn finalize_usage_execution(
|
||||
&self,
|
||||
execution_id: &str,
|
||||
outcome: &str,
|
||||
reports: Vec<(String, serde_json::Value)>,
|
||||
) -> Result<()> {
|
||||
let tx = self.begin_immediate()?;
|
||||
let suppressed: Option<bool> = tx.query_row("UPDATE chat_usage_executions SET outcome = ?2 WHERE execution_id = ?1 AND outcome IS NULL RETURNING suppressed", params![execution_id, outcome], |row| row.get(0)).optional()?;
|
||||
if suppressed == Some(false) {
|
||||
for (id, payload) in reports {
|
||||
self.stage_telemetry(&id, &payload)?;
|
||||
}
|
||||
}
|
||||
tx.commit()?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) fn stage_telemetry(
|
||||
&self,
|
||||
event_id: &str,
|
||||
@@ -526,6 +549,185 @@ impl Store {
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn prepared_chat_report_cannot_escape_opt_out_or_duplicate_finalization() {
|
||||
let dir = std::env::temp_dir().join(format!("orx-opt-out-race-{}", uuid::Uuid::new_v4()));
|
||||
let store = Store::open_at(dir.clone()).unwrap();
|
||||
store
|
||||
.begin_usage_execution("execution", "turn", "claude-code")
|
||||
.unwrap();
|
||||
store
|
||||
.conn
|
||||
.execute("UPDATE chat_usage_executions SET suppressed = 0", [])
|
||||
.unwrap();
|
||||
let prepared = vec![("prepared".into(), serde_json::json!({"usage": 12}))];
|
||||
store.purge_pending_telemetry().unwrap();
|
||||
store
|
||||
.finalize_usage_execution("execution", "done", prepared.clone())
|
||||
.unwrap();
|
||||
assert!(store.pending_telemetry().unwrap().is_empty());
|
||||
store
|
||||
.conn
|
||||
.execute("UPDATE chat_usage_executions SET suppressed = 0", [])
|
||||
.unwrap();
|
||||
store
|
||||
.finalize_usage_execution("execution", "done", prepared)
|
||||
.unwrap();
|
||||
assert!(store.pending_telemetry().unwrap().is_empty());
|
||||
let outcome: String = store
|
||||
.conn
|
||||
.query_row("SELECT outcome FROM chat_usage_executions", [], |row| {
|
||||
row.get(0)
|
||||
})
|
||||
.unwrap();
|
||||
assert_eq!(outcome, "done");
|
||||
drop(store);
|
||||
std::fs::remove_dir_all(dir).unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn complementary_partial_samples_finalize_without_losing_measured_counters() {
|
||||
let dir = std::env::temp_dir().join(format!("orx-partial-usage-{}", uuid::Uuid::new_v4()));
|
||||
let store = Store::open_at(dir.clone()).unwrap();
|
||||
store
|
||||
.begin_usage_execution("execution", "turn", "claude-code")
|
||||
.unwrap();
|
||||
store
|
||||
.conn
|
||||
.execute("UPDATE chat_usage_executions SET suppressed = 0", [])
|
||||
.unwrap();
|
||||
for (id, usage) in [
|
||||
(
|
||||
"input",
|
||||
TokenUsage {
|
||||
input_tokens: Some(10),
|
||||
cache_read_tokens: Some(8),
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
(
|
||||
"cache",
|
||||
TokenUsage {
|
||||
cache_read_tokens: Some(9),
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
(
|
||||
"output",
|
||||
TokenUsage {
|
||||
output_tokens: Some(3),
|
||||
reasoning_tokens: Some(2),
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
(
|
||||
"reasoning",
|
||||
TokenUsage {
|
||||
reasoning_tokens: Some(4),
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
] {
|
||||
store
|
||||
.record_usage_sample(
|
||||
"execution",
|
||||
id,
|
||||
"claude-code",
|
||||
Some("claude-opus-5-5"),
|
||||
None,
|
||||
&usage,
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
store.finalize_turn_usage("turn", "failed").unwrap();
|
||||
store.finalize_turn_usage("turn", "failed").unwrap();
|
||||
let outcome: String = store
|
||||
.conn
|
||||
.query_row("SELECT outcome FROM chat_usage_executions", [], |row| {
|
||||
row.get(0)
|
||||
})
|
||||
.unwrap();
|
||||
assert_eq!(outcome, "failed");
|
||||
let counters: Vec<String> = store
|
||||
.conn
|
||||
.prepare("SELECT usage_json FROM chat_usage_samples")
|
||||
.unwrap()
|
||||
.query_map([], |row| row.get(0))
|
||||
.unwrap()
|
||||
.collect::<std::result::Result<_, _>>()
|
||||
.unwrap();
|
||||
let cache: u64 = counters
|
||||
.iter()
|
||||
.map(|json| {
|
||||
serde_json::from_str::<TokenUsage>(json)
|
||||
.unwrap()
|
||||
.cache_read_tokens
|
||||
.unwrap_or(0)
|
||||
})
|
||||
.sum();
|
||||
assert_eq!(cache, 17);
|
||||
drop(store);
|
||||
std::fs::remove_dir_all(dir).unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn invocation_identities_follow_session_deletion_and_unowned_retention() {
|
||||
let dir =
|
||||
std::env::temp_dir().join(format!("orx-identity-retention-{}", uuid::Uuid::new_v4()));
|
||||
let store = Store::open_at(dir.clone()).unwrap();
|
||||
for id in ["session", "project-session"] {
|
||||
store.conn.execute("INSERT INTO chat_sessions (id, project_id, harness, created_at, updated_at) VALUES (?1, 'project', 'claude-code', 1, 1)", [id]).unwrap();
|
||||
}
|
||||
let identity = InvocationIdentity {
|
||||
harness: "claude-code".into(),
|
||||
model: "claude-opus-5-5".into(),
|
||||
provider: None,
|
||||
};
|
||||
store
|
||||
.record_native_invocation("session-call", &identity, Some("session"))
|
||||
.unwrap();
|
||||
store
|
||||
.record_native_invocation("project-call", &identity, Some("project-session"))
|
||||
.unwrap();
|
||||
store
|
||||
.record_native_invocation("legacy-call", &identity, None)
|
||||
.unwrap();
|
||||
store
|
||||
.record_native_invocation("recent-unowned", &identity, None)
|
||||
.unwrap();
|
||||
store.conn.execute("UPDATE native_invocation_identities SET created_at = 0 WHERE call_id = 'legacy-call'", []).unwrap();
|
||||
store.delete_chat_session("session").unwrap();
|
||||
assert_eq!(
|
||||
store
|
||||
.native_invocation_identity("claude-code", "session-call")
|
||||
.unwrap(),
|
||||
None
|
||||
);
|
||||
store.delete_local_project("project").unwrap();
|
||||
assert_eq!(
|
||||
store
|
||||
.native_invocation_identity("claude-code", "project-call")
|
||||
.unwrap(),
|
||||
None
|
||||
);
|
||||
drop(store);
|
||||
let store = Store::open_at(dir.clone()).unwrap();
|
||||
assert_eq!(
|
||||
store
|
||||
.native_invocation_identity("claude-code", "legacy-call")
|
||||
.unwrap(),
|
||||
None
|
||||
);
|
||||
assert_eq!(
|
||||
store
|
||||
.native_invocation_identity("claude-code", "recent-unowned")
|
||||
.unwrap(),
|
||||
Some(identity)
|
||||
);
|
||||
drop(store);
|
||||
std::fs::remove_dir_all(dir).unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pending_reports_commit_and_survive_restart() {
|
||||
let dir = std::env::temp_dir().join(format!("orx-telemetry-{}", uuid::Uuid::new_v4()));
|
||||
@@ -896,13 +1098,13 @@ mod tests {
|
||||
..parent.clone()
|
||||
};
|
||||
store
|
||||
.record_native_invocation("parent-call", &parent)
|
||||
.record_native_invocation("parent-call", &parent, None)
|
||||
.unwrap();
|
||||
store
|
||||
.record_native_invocation("child-call", &child)
|
||||
.record_native_invocation("child-call", &child, None)
|
||||
.unwrap();
|
||||
assert!(store
|
||||
.record_native_invocation("parent-call", &child)
|
||||
.record_native_invocation("parent-call", &child, None)
|
||||
.is_err());
|
||||
assert_eq!(
|
||||
store
|
||||
|
||||
Reference in New Issue
Block a user