From f42d3584cb8b6cb23a27d4e82c256403838192c2 Mon Sep 17 00:00:00 2001 From: Antoine Lecompte <38678863+nutgood@users.noreply.github.com> Date: Tue, 15 Sep 2026 09:30:48 -0400 Subject: [PATCH] feat(columbo): seed client-side restic telemetry into investigations (#686) --- docs/columbo.md | 29 ++++++-- packages/columbo/internal/agent/agent.go | 32 +++++++++ packages/columbo/internal/agent/tools.go | 17 +++++ packages/columbo/internal/agent/tools_test.go | 28 ++++++++ packages/columbo/internal/o11y/client.go | 69 +++++++++++++++++++ packages/columbo/internal/o11y/client_test.go | 49 +++++++++++++ 6 files changed, 219 insertions(+), 5 deletions(-) diff --git a/docs/columbo.md b/docs/columbo.md index deb60d4a..543ef41c 100644 --- a/docs/columbo.md +++ b/docs/columbo.md @@ -58,11 +58,30 @@ blob_type:="snapshots"` marks a completed backup, steady `data` saves a backup in progress, `locks` writes any operation including read-only ones; path regexes remain the fallback for older entries in retention). The catalog is maintained by hand in `investigateSystemPrompt`; update it when a -service adds or renames per-user telemetry. On top of that, each investigation opens with a free scoped -lookup of which of those names actually carry data for THIS account in the -last 30 days (`/api/v1/label/__name__/values` + `extra_label`), so an -account with no backup traffic is recognized in turn one instead of after a -string of empty queries. +service adds or renames per-user telemetry. + +The catalog also covers **client-side telemetry**: the user's own backup +client (yucca-sdk's orchestration-api) ships structured logs home, and +yucca-api records them as `_msg:"[telemetry] "` with the payload +flattened into `data.*`. This is the only view of what happened on the +user's machine — `[telemetry] Backup finished` carries `data.lastBackupStatus`, +`data.version` (client version) and, on failure, restic's verbatim stderr in +`data.error.message`. Those failures are frequently invisible server-side +because the request never arrived (DNS, TLS, local permissions, restic's +stuck-request timeout), which is exactly when an investigation would +otherwise conclude "nothing found". Telemetry is opt-in, so its absence means +the user declined it or runs an old client, not that no backups ran. + +On top of that, each investigation opens with two free scoped lookups (no +tool budget): which metric names actually carry data for THIS account in the +last 30 days (`/api/v1/label/__name__/values` + `extra_label`), and a digest +of what its client reported home over the same window — one line per distinct +(event, status, client version) with a count, the last occurrence and one +example error, newest first. So an account with no backup traffic is +recognized in turn one instead of after a string of empty queries, and a +client failing on its own side usually tells its whole story before the model +spends a single tool call. Both are fixed queries owned by the harness; the +digest goes through the same `QueryLogs` scoping as everything else. ## Trust model diff --git a/packages/columbo/internal/agent/agent.go b/packages/columbo/internal/agent/agent.go index 79650ef2..fd22426f 100644 --- a/packages/columbo/internal/agent/agent.go +++ b/packages/columbo/internal/agent/agent.go @@ -129,6 +129,11 @@ Michael operation semantics (restic REST protocol): - op:="save_blob" = a blob write; what it MEANS depends on blob_type: data = backup content uploading; index = index flush; snapshots = a backup COMPLETED (the snapshot record is written last); keys = repository key setup. locks is the exception — restic writes a lock at the start of EVERY operation, including read-only ones (restore, check), so lock writes prove activity, not backups. - get_blob = blob read (restores, checks); check_blob = existence probe; delete_blob = cleanup (locks after every operation; data/index during prune); list_blobs = listing; save_config = repository initialization (happens once, before the first backup); create_repository = repository creation. So: op:="save_blob" blob_type:="snapshots" = completed backups; op:="save_blob" blob_type:="data" = backup traffic; status:>=400 on michael = failing restic requests. +- client telemetry: the user's own backup client (the Immich integration / standalone container) ships structured logs home, which yucca-api records as _msg:"[telemetry] " with the payload flattened into data.* fields. This is the ONLY view of what happened on the user's MACHINE — including restic's own stderr — so when the server side looks clean, look here before concluding nothing is wrong. + - "[telemetry] Backup finished" is the key event: data.lastBackupStatus (complete/warn/failed), data.repositoryId, data.version (the client's version), and on failure data.error.message — the verbatim restic retry/error log as the client saw it. Typical contents: DNS/TLS failures, "request timeout" (restic's stuck-request timeout), local permission and cache errors. None of these are visible server-side, because the request never arrived. + - Lifecycle events trace the rest of a run: "Running backup", "Finished backup to primary backend", "Finished prune on primary backend", "Creating"/"Created Immich database backup", "Running"/"Finished repository prune|import|snapshot restore", "Connected FUTO Backups backend", "Configured Immich integration", "Unhandled request error" (a 5xx inside the client itself). + - Telemetry is opt-in: no [telemetry] logs means the user declined it or runs an old client, NOT that nothing happened. + - The digest below is prefetched for free. Drill in with query_logs for the full text, e.g. _msg:"[telemetry] Backup finished" data.lastBackupStatus:="failed" — data.error.message is often multi-KB, so prefer a narrow limit or a "| fields _time, data.error.message" pipe. Platform health (fleet-wide, not user-specific): the query_health tool runs fixed named probes over platform telemetry — michael error rates and latency, storage-backend health, Ceph/RGW health, pool capacity. When the user's telemetry shows server-side errors (5xx, timeouts), check whether a platform incident overlaps their error window; a healthy platform during that window is itself evidence. One or two probes over the incident window usually suffice — do not audit the whole platform. @@ -190,6 +195,7 @@ func (r *Runner) run(ctx context.Context, userID, userMessage string) (Outcome, } userMessage += "\n\n" + availableMetricsLine(box.availableMetrics(ctx)) + userMessage += "\n\n" + clientTelemetryBlock(box.clientTelemetry(ctx)) zerolog.Ctx(ctx).Info(). Str("audit", "investigation_start"). @@ -228,6 +234,32 @@ func availableMetricsLine(names []string) string { return "Metrics with data for this account (last 30d): " + strings.Join(names, ", ") } +func clientTelemetryBlock(events []o11y.ClientEvent) string { + if events == nil { + return "Client telemetry for this account: (lookup unavailable — query the logs to find out)" + } + if len(events) == 0 { + return "Client telemetry for this account (last 30d): none — this account's backup client has never reported in. " + + "Telemetry is opt-in, so this means the user declined it or runs a client too old to send it, NOT that no backups ran." + } + var b strings.Builder + b.WriteString("What this account's own backup client reported home (last 30d, newest first, count after ×):") + for _, e := range events { + b.WriteString("\n- " + e.Last + " " + e.Event) + if e.Status != "" { + b.WriteString(" status=" + e.Status) + } + if e.Version != "" { + b.WriteString(" client=" + e.Version) + } + b.WriteString(" ×" + e.Count) + if e.Error != "" { + b.WriteString("\n example error: " + e.Error) + } + } + return b.String() +} + type tokenTally struct { mu sync.Mutex prompt int diff --git a/packages/columbo/internal/agent/tools.go b/packages/columbo/internal/agent/tools.go index e4307f1b..2a7b9a72 100644 --- a/packages/columbo/internal/agent/tools.go +++ b/packages/columbo/internal/agent/tools.go @@ -118,6 +118,23 @@ func (t *toolbox) availableMetrics(ctx context.Context) []string { return names } +// clientTelemetry is the free (no tool budget) prefetch of what this user's +// own backup client reported home; nil means the lookup failed and the model +// is told to query instead. +func (t *toolbox) clientTelemetry(ctx context.Context) []o11y.ClientEvent { + ctx, cancel := context.WithTimeout(ctx, 10*time.Second) + defer cancel() + events, err := t.o11y.ClientTelemetry(ctx, maxLookback) + if err != nil { + zerolog.Ctx(ctx).Warn().Err(err).Msg("client-telemetry prefetch failed") + return nil + } + if events == nil { + events = []o11y.ClientEvent{} + } + return events +} + func (t *toolbox) callsMade() int { t.mu.Lock() defer t.mu.Unlock() diff --git a/packages/columbo/internal/agent/tools_test.go b/packages/columbo/internal/agent/tools_test.go index db4b9096..957f4a8f 100644 --- a/packages/columbo/internal/agent/tools_test.go +++ b/packages/columbo/internal/agent/tools_test.go @@ -6,6 +6,8 @@ import ( "testing" "time" + "columbo/internal/o11y" + "github.com/cloudwego/eino/components/tool" ) @@ -160,3 +162,29 @@ func TestTruncateNote(t *testing.T) { t.Fatalf("short note mangled: %q", got) } } + +func TestClientTelemetryBlock(t *testing.T) { + if got := clientTelemetryBlock(nil); !strings.Contains(got, "lookup unavailable") { + t.Fatalf("nil case = %q", got) + } + empty := clientTelemetryBlock([]o11y.ClientEvent{}) + if !strings.Contains(empty, "never reported in") || !strings.Contains(empty, "opt-in") { + t.Fatalf("empty case must not be read as 'no backups ran': %q", empty) + } + got := clientTelemetryBlock([]o11y.ClientEvent{ + {Event: "Backup finished", Status: "failed", Version: "0.40.1", Count: "77", Last: "2026-09-11T10:16:49Z", Error: "Unknown site 'local'"}, + {Event: "Running backup", Version: "0.43.0", Count: "2", Last: "2026-09-15T08:44:08Z"}, + }) + for _, want := range []string{ + "2026-09-11T10:16:49Z Backup finished status=failed client=0.40.1 ×77", + "example error: Unknown site 'local'", + "2026-09-15T08:44:08Z Running backup client=0.43.0 ×2", + } { + if !strings.Contains(got, want) { + t.Fatalf("block missing %q:\n%s", want, got) + } + } + if strings.Contains(got, "status= ") { + t.Fatalf("events without a status must omit the field: %q", got) + } +} diff --git a/packages/columbo/internal/o11y/client.go b/packages/columbo/internal/o11y/client.go index bcb85679..36dc862c 100644 --- a/packages/columbo/internal/o11y/client.go +++ b/packages/columbo/internal/o11y/client.go @@ -191,3 +191,72 @@ func truncate(s string, n int) string { } return s[:n] + "…" } + +// ClientEvent is one row of the client-telemetry digest: a distinct +// (event, status, version) combination the user's own backup client reported, +// with how often it happened, when it last did, and one example error. +type ClientEvent struct { + Event string + Status string + Version string + Count string + Last string + Error string +} + +// The client's error strings carry whole restic stack traces; the digest +// keeps only enough to recognise the failure, and the agent can pull the +// full text with a normal log query. +const maxClientErrorChars = 400 + +const clientTelemetryQuery = `_msg:"[telemetry]" | stats by (_msg, data.lastBackupStatus, data.version) ` + + `count() c, max(_time) last, row_any(data.error.message) err | sort by (last desc) | limit 100` + +// ClientTelemetry digests what this user's own backup client reported home. +// The client ships structured logs to yucca-api, which records them under +// `[telemetry] ` with the payload flattened into `data.*` — the only +// view anyone has of what happened on the user's machine, restic's own stderr +// included. The query is fixed here rather than composed by the model, and +// still goes through QueryLogs so it inherits the same per-user scoping as +// everything else. +func (c *Client) ClientTelemetry(ctx context.Context, lookback time.Duration) ([]ClientEvent, error) { + start := time.Now().Add(-lookback).Format(time.RFC3339) + end := time.Now().Format(time.RFC3339) + body, err := c.QueryLogs(ctx, clientTelemetryQuery, start, end, 1000) + if err != nil { + return nil, err + } + + events := []ClientEvent{} + for _, line := range strings.Split(body, "\n") { + if strings.TrimSpace(line) == "" { + continue + } + var row map[string]string + if err := json.Unmarshal([]byte(line), &row); err != nil { + return nil, err + } + events = append(events, ClientEvent{ + Event: strings.TrimPrefix(row["_msg"], "[telemetry] "), + Status: row["data.lastBackupStatus"], + Version: row["data.version"], + Count: row["c"], + Last: row["last"], + Error: parseRowAnyError(row["err"]), + }) + } + return events, nil +} + +// parseRowAnyError unwraps LogsQL's row_any output, which arrives as a JSON +// object of the selected fields and is `{}` for rows that had no error. +func parseRowAnyError(raw string) string { + if raw == "" { + return "" + } + var fields map[string]string + if err := json.Unmarshal([]byte(raw), &fields); err != nil { + return "" + } + return truncate(strings.Join(strings.Fields(fields["data.error.message"]), " "), maxClientErrorChars) +} diff --git a/packages/columbo/internal/o11y/client_test.go b/packages/columbo/internal/o11y/client_test.go index 458224b4..9159f50e 100644 --- a/packages/columbo/internal/o11y/client_test.go +++ b/packages/columbo/internal/o11y/client_test.go @@ -179,3 +179,52 @@ func formValue(t *testing.T, encoded []byte, key string) string { } return values.Get(key) } + +func TestClientTelemetryParsesDigestAndStaysScoped(t *testing.T) { + srv, captured, form := recordingServer(t, http.StatusOK, strings.Join([]string{ + `{"_msg":"[telemetry] Backup finished","data.lastBackupStatus":"failed","data.version":"0.40.1","c":"77","last":"2026-09-11T10:16:49Z","err":"{\"data.error.message\":\"Unknown site\\n 'local'\"}"}`, + `{"_msg":"[telemetry] Running backup","data.version":"0.43.0","c":"2","last":"2026-09-15T08:44:08Z","err":"{}"}`, + ``, + }, "\n")) + client := NewClient(srv.URL, srv.URL, "user-1") + + events, err := client.ClientTelemetry(context.Background(), 30*24*time.Hour) + if err != nil { + t.Fatal(err) + } + if captured.URL.Path != "/select/logsql/query" { + t.Fatalf("unexpected path %q", captured.URL.Path) + } + query := formValue(t, *form, "query") + if !strings.HasPrefix(query, `(user:="user-1" or customerId:="user-1") and (_msg:"[telemetry]")`) { + t.Fatalf("digest query was not scoped: %q", query) + } + if len(events) != 2 { + t.Fatalf("events = %+v", events) + } + first := events[0] + if first.Event != "Backup finished" || first.Status != "failed" || first.Version != "0.40.1" || first.Count != "77" { + t.Fatalf("first = %+v", first) + } + if first.Error != "Unknown site 'local'" { + t.Fatalf("error = %q, want the whitespace-collapsed restic message", first.Error) + } + if events[1].Error != "" { + t.Fatalf("row_any's empty object should yield no error, got %q", events[1].Error) + } +} + +func TestClientTelemetryTruncatesLongErrors(t *testing.T) { + stack := strings.Repeat("a", maxClientErrorChars*2) + srv, _, _ := recordingServer(t, http.StatusOK, + `{"_msg":"[telemetry] Backup finished","err":"{\"data.error.message\":\"`+stack+`\"}"}`) + client := NewClient(srv.URL, srv.URL, "user-1") + + events, err := client.ClientTelemetry(context.Background(), time.Hour) + if err != nil { + t.Fatal(err) + } + if got := events[0].Error; got != stack[:maxClientErrorChars]+"…" { + t.Fatalf("error was not truncated to %d chars: len=%d", maxClientErrorChars, len(got)) + } +}