Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 23 additions & 18 deletions bt-daemon/src/sink/braintrust.rs
Original file line number Diff line number Diff line change
Expand Up @@ -270,6 +270,13 @@ impl BraintrustSink {
Ok(client)
}

fn span_origin(&self) -> SpanOrigin {
SpanOrigin::new()
.name(format!("braintrust.plugin.{}", self.source))
.version(self.version.clone())
.instrumentation("braintrust-plugin")
}

fn ensure_handle(&mut self, client: &BraintrustClient, row: &SpanRow) -> anyhow::Result<()> {
if self.open.contains_key(&row.span_id) {
return Ok(());
Expand All @@ -287,12 +294,7 @@ impl BraintrustSink {
.span_id(row.span_id.clone())
.row_id(row.span_id.clone())
.parent_info(parent)
.span_origin(
SpanOrigin::new()
.name(format!("braintrust.plugin.{}", self.source))
.version(self.version.clone())
.instrumentation("braintrust-plugin"),
);
.span_origin(self.span_origin());
if let Some(org_name) = &creds.org_name {
builder = builder.org_name(org_name.clone());
}
Expand All @@ -306,7 +308,7 @@ impl BraintrustSink {
fn update_open(&mut self, client: &BraintrustClient, row: &SpanRow) -> anyhow::Result<()> {
self.ensure_handle(client, row)?;
let handle = self.open.get(&row.span_id).expect("just inserted");
handle.log(build_log(row, &self.daemon_version)?);
handle.log(build_log(row, &self.daemon_version, self.span_origin())?);
if let Some(end) = row.end_ms {
handle.end_with_time(ms_to_secs(end));
// SpanHandle retains the complete accumulated input/output. Once a
Expand All @@ -329,7 +331,7 @@ impl BraintrustSink {
creds.token.clone(),
creds.org_id.clone(),
&components,
build_log(row, &self.daemon_version)?,
build_log(row, &self.daemon_version, self.span_origin())?,
)
.map_err(|error| anyhow::anyhow!("braintrust span merge failed: {error}"))
}
Expand Down Expand Up @@ -627,20 +629,22 @@ fn ms_to_secs(ms: i64) -> f64 {
ms as f64 / 1000.0
}

fn build_log(row: &SpanRow, daemon_version: &str) -> anyhow::Result<SpanLog> {
fn build_log(row: &SpanRow, daemon_version: &str, origin: SpanOrigin) -> anyhow::Result<SpanLog> {
// Repeat the plugin origin so stateless merges cannot use SDK defaults.
let mut builder = SpanLog::builder().span_origin(origin);

// The span's display name is carried on the log event, not the builder.
// An empty name means "unchanged" (many merge ops use `..Default::default()`
// and don't rename the span) — omitting `.name()` avoids overwriting the
// already-set name with an empty string on merge.
let mut lb = SpanLog::builder();
if !row.name.is_empty() {
lb = lb.name(row.name.clone());
builder = builder.name(row.name.clone());
}
if let Some(input) = &row.input {
lb = lb.input(input.clone());
builder = builder.input(input.clone());
}
if let Some(output) = &row.output {
lb = lb.output(output.clone());
builder = builder.output(output.clone());
}
let mut metadata = row
.metadata
Expand All @@ -655,7 +659,7 @@ fn build_log(row: &SpanRow, daemon_version: &str) -> anyhow::Result<SpanLog> {
"bt_daemon_version".into(),
Value::String(daemon_version.to_string()),
);
lb = lb.metadata(metadata);
builder = builder.metadata(metadata);
let mut metrics = row
.metrics
.as_ref()
Expand All @@ -673,17 +677,18 @@ fn build_log(row: &SpanRow, daemon_version: &str) -> anyhow::Result<SpanLog> {
metrics.insert("end".into(), ms_to_secs(end));
}
if !metrics.is_empty() {
lb = lb.metrics(metrics);
builder = builder.metrics(metrics);
}
if let Some(err) = &row.error {
lb = lb.error(Value::String(err.clone()));
builder = builder.error(Value::String(err.clone()));
}
if let Some(tags) = &row.tags {
if !tags.is_empty() {
lb = lb.tags(tags.clone());
builder = builder.tags(tags.clone());
}
}
lb.build()
builder
.build()
.map_err(|e| anyhow::anyhow!("span log build failed: {e}"))
}

Expand Down
116 changes: 74 additions & 42 deletions bt-daemon/tests/braintrust_sink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -410,49 +410,81 @@ async fn late_merge_updates_a_completed_span_without_an_open_handle() {
app_url: Some(base.clone()),
version: "test".into(),
});
let mut sink = factory.create("sess-late", "codex", None).unwrap();
sink.configure(&session_config(&base));

sink.emit(&[SpanOp::Insert(row(
"finished",
"trace-root",
&["turn-parent"],
"original name",
SpanType::Tool,
1,
Some(2),
))])
.await
.unwrap();
sink.flush().await.unwrap();
let mut late = SpanRow {
span_id: "finished".into(),
root_span_id: "trace-root".into(),
parent_span_ids: vec!["turn-parent".into()],
output: Some(json!({"status":"late"})),
..Default::default()
};
late.name.clear();
sink.emit(&[SpanOp::Merge(late)]).await.unwrap();
sink.flush().await.unwrap();
// The same cached SDK client serves every source. Provenance must stay
// session-specific, including when the plugin version falls back to bt's.
for (source, plugin_version) in [
("codex", Some("1.2.3")),
("claude", Some("2.3.4")),
("grok", Some("3.4.5")),
("opencode", Some("4.5.6")),
("pi", Some("5.6.7")),
("antigravity", None),
] {
let session_id = format!("sess-late-{source}");
let span_id = format!("finished-{source}");
let mut sink = factory.create(&session_id, source, plugin_version).unwrap();
sink.configure(&session_config(&base));

sink.emit(&[SpanOp::Insert(row(
&span_id,
"trace-root",
&["turn-parent"],
"original name",
SpanType::Tool,
1,
Some(2),
))])
.await
.unwrap();
sink.flush().await.unwrap();
let late = SpanRow {
span_id: span_id.clone(),
root_span_id: "trace-root".into(),
parent_span_ids: vec!["turn-parent".into()],
output: Some(json!({"status":"late"})),
..Default::default()
};
sink.emit(&[SpanOp::Merge(late.clone())]).await.unwrap();
sink.flush().await.unwrap();

// Recovery also starts without an open handle for an existing span.
drop(sink);
let mut recovered = factory.create(&session_id, source, plugin_version).unwrap();
recovered.configure(&session_config(&base));
recovered
.emit(&[SpanOp::Merge(SpanRow {
output: Some(json!({"status":"recovered"})),
..late
})])
.await
.unwrap();
recovered.flush().await.unwrap();

let bodies = logs3_bodies(&server).await;
assert!(
bodies.contains("original name"),
"initial row absent: {bodies}"
);
let rows = logs3_rows(&server).await;
let late = rows
.iter()
.find(|row| row.pointer("/output/status") == Some(&json!("late")))
.unwrap_or_else(|| panic!("late merge absent: {bodies}"));
assert_eq!(late["_is_merge"], true);
assert_eq!(late["root_span_id"], "trace-root");
assert_eq!(
late["span_parents"],
json!(["turn-parent"]),
"stateless merge did not repeat the child parent identity: {bodies}"
);
let rows: Vec<_> = logs3_rows(&server)
.await
.into_iter()
.filter(|row| row["span_id"] == span_id)
.collect();
assert!(rows
.iter()
.any(|row| row.pointer("/span_attributes/name") == Some(&json!("original name"))));
for status in ["late", "recovered"] {
let update = rows
.iter()
.find(|row| row.pointer("/output/status") == Some(&json!(status)))
.unwrap_or_else(|| panic!("{source}: {status} merge absent: {rows:?}"));
assert_eq!(update["_is_merge"], true);
assert_eq!(update["root_span_id"], "trace-root");
assert_eq!(update["span_parents"], json!(["turn-parent"]));
}
for row in rows {
let origin = &row["context"]["span_origin"];
assert_eq!(origin["name"], format!("braintrust.plugin.{source}"));
assert_eq!(origin["version"], plugin_version.unwrap_or("test"));
assert_eq!(origin["instrumentation"]["name"], "braintrust-plugin");
assert_eq!(row["metadata"]["bt_daemon_version"], "test");
}
}
}

#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
Expand Down
6 changes: 4 additions & 2 deletions src/plugins/pi/content/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,9 @@ export default function braintrustPiExtension(pi: ExtensionAPI): void {
requestTimeoutMs: UI_STATUS_TIMEOUT_MS,
});

const remember = async (ctx: ExtensionContext): Promise<ReturnType<typeof sessionDescriptor>> => {
const updateSession = async (
ctx: ExtensionContext,
): Promise<ReturnType<typeof sessionDescriptor>> => {
lastContext = ctx;
const descriptor = sessionDescriptor(ctx);
if (sessionId !== descriptor.sessionId) uiGeneration += 1;
Expand Down Expand Up @@ -132,7 +134,7 @@ export default function braintrustPiExtension(pi: ExtensionAPI): void {
ctx?: ExtensionContext,
updateUi = false,
): Promise<void> => {
const descriptor = ctx ? await remember(ctx) : undefined;
const descriptor = ctx ? await updateSession(ctx) : undefined;
if (!sessionId) return;
await client.log({
source: "pi",
Expand Down
Loading