From 5aaecd2a53c93dce41150107bec4281ede742e94 Mon Sep 17 00:00:00 2001 From: kami Date: Sat, 1 Aug 2026 14:13:58 +0400 Subject: [PATCH 1/5] store: give ecosystem traces their own table Traces were written as facts. A single Praxis action wrote several of them, so machine-rate rows crowded out the bounded fact readers that humans and evaluation consume. The habit profile window of 2000 facts and the memeval snapshot both filled with call records instead of what Maven learned about the owner. Traces now go to ecosystem_traces, with correlation, causation, duration and HTTP status as columns, pruned to the most recent 5000. The new reader is exposed over IPC and rendered as the Calls card on the ecosystem page, so it is a table someone actually looks at. Found in review of #84. --- cmd/mavend/ecosystem_trace_test.go | 177 ++++++++++++++++++----------- cmd/mavweb/ecosystem.go | 17 ++- cmd/mavweb/ecosystem.html | 15 ++- cmd/mavweb/main.go | 2 +- internal/ipc/api.go | 21 ++++ internal/ipc/client.go | 9 ++ internal/ipc/server.go | 26 +++++ internal/ipc/unimplemented.go | 3 + internal/ipc/wire.go | 1 + internal/store/ecotraces.go | 107 +++++++++++++++++ internal/store/migrations.go | 24 ++++ 11 files changed, 330 insertions(+), 72 deletions(-) create mode 100644 internal/store/ecotraces.go diff --git a/cmd/mavend/ecosystem_trace_test.go b/cmd/mavend/ecosystem_trace_test.go index ee59d19..86d3b1c 100644 --- a/cmd/mavend/ecosystem_trace_test.go +++ b/cmd/mavend/ecosystem_trace_test.go @@ -2,7 +2,6 @@ package main import ( "context" - "encoding/json" "strings" "testing" @@ -11,34 +10,12 @@ import ( // Versioning, authentication and tracing of ecosystem calls (Vikunja #273). -func ecoTraces(t *testing.T, h *reactiveHandler) []map[string]any { +func findTrace(t *testing.T, h *reactiveHandler, service, op string) *store.EcosystemTrace { t.Helper() - facts, err := h.dataStore.RecentFacts(context.Background(), 100) - if err != nil { - t.Fatalf("read facts: %v", err) - } - var out []map[string]any - for _, f := range facts { - if f.Source != "ecosystem:trace" { - continue - } - i := strings.Index(f.Value, "{") - if i < 0 { - t.Fatalf("trace fact carries no detail object: %q", f.Value) - } - var d map[string]any - if err := json.Unmarshal([]byte(f.Value[i:]), &d); err != nil { - t.Fatalf("decode trace %q: %v", f.Value, err) - } - out = append(out, d) - } - return out -} - -func findTrace(traces []map[string]any, service, op string) map[string]any { - for _, d := range traces { - if d["service"] == service && d["operation"] == op { - return d + for _, tr := range traces(t, h) { + if tr.Service == service && tr.Operation == op { + found := tr + return &found } } return nil @@ -58,7 +35,9 @@ func TestEcosystemHeaders_VersionRequesterAndAuth(t *testing.T) { if err != nil { t.Fatalf("resolve: %v", err) } - if _, err := h.ecosystem.praxis.ListAttention(ctx, 5); err != nil { + // A bare client call carries whatever the caller assigned. Entry points + // assign the ID, the header layer only reads it, so mirror an action here. + if _, err := h.ecosystem.praxis.ListAttention(withCorrelationID(ctx, newCorrelationID()), 5); err != nil { t.Fatalf("attention: %v", err) } @@ -167,31 +146,69 @@ func TestEcosystemTrace_SuccessfulActionTracesEveryHop(t *testing.T) { t.Fatalf("setup: expected success, got %q", reply) } - traces := ecoTraces(t, h) + var chain string for _, want := range [][2]string{{"nexus", "resolve"}, {"hexis", "capabilities"}, {"hexis", "execute"}} { - d := findTrace(traces, want[0], want[1]) + d := findTrace(t, h, want[0], want[1]) if d == nil { - t.Fatalf("missing trace for %s %s, got %+v", want[0], want[1], traces) + t.Fatalf("missing trace for %s %s, got %+v", want[0], want[1], traces(t, h)) } - if d["status"] != traceOK { - t.Errorf("%s %s status = %v, want ok", want[0], want[1], d["status"]) + if d.Status != traceOK { + t.Errorf("%s %s status = %v, want ok", want[0], want[1], d.Status) } - if _, ok := d["duration_ms"]; !ok { - t.Errorf("%s %s trace has no timing", want[0], want[1]) - } - if d["correlation_id"] == nil || d["correlation_id"] == "" { + if d.CorrelationID == "" { t.Errorf("%s %s trace has no correlation id", want[0], want[1]) } + if want[1] != "execute" { + if chain == "" { + chain = d.CorrelationID + } else if d.CorrelationID != chain { + t.Errorf("%s %s left the correlation chain: %s != %s", want[0], want[1], d.CorrelationID, chain) + } + } } - exec := findTrace(traces, "hexis", "execute") - if exec["causation_id"] == nil || exec["causation_id"] == "" { + exec := findTrace(t, h, "hexis", "execute") + if exec.CausationID == "" { t.Error("execute trace must carry the causation id of the turn that caused it") } - if exec["correlation_id"] == exec["causation_id"] { + if exec.CorrelationID == exec.CausationID { t.Error("execute correlation and causation must be distinguishable") } } +// TestEcosystemTrace_OneCorrelationIDPerPraxisAction: a digest calls attention +// once and surface once per item. All of it is one turn, so the far side must +// see one ID and not N+1 unrelated ones. +func TestEcosystemTrace_OneCorrelationIDPerPraxisAction(t *testing.T) { + ctx := context.Background() + praxis := newFakePraxis(t, fixturePraxisAttentionItems( + map[string]any{"id": "item_1", "title": "disk almost full", "importance": 3.0}, + map[string]any{"id": "item_2", "title": "backup is stale", "importance": 2.0}, + )) + h := ecoHandler(t, nil, praxis, nil) + + if reply := h.handlePraxisAct(ctx, praxisActDec("list_attention")); !strings.Contains(reply, "disk almost full") { + t.Fatalf("setup: expected the digest, got %q", reply) + } + + reqs := praxis.Requests() + if len(reqs) < 3 { + t.Fatalf("expected attention plus one surface per item, got %d requests", len(reqs)) + } + first := reqs[0].Header.Get("X-Correlation-ID") + if first == "" { + t.Fatal("every ecosystem request must carry a correlation id") + } + for _, r := range reqs { + if got := r.Header.Get("X-Correlation-ID"); got != first { + t.Fatalf("%s %s carried %q, want the action's id %q", r.Method, r.Path, got, first) + } + } + tr := findTrace(t, h, "praxis", "list_attention") + if tr == nil || tr.CorrelationID != first { + t.Fatalf("the trace must carry the id that was actually sent, got %+v", tr) + } +} + // TestEcosystemTrace_FailuresAreTracedToo: the whole point of the change — // a failed hop is exactly the one worth having recorded. func TestEcosystemTrace_FailuresAreTracedToo(t *testing.T) { @@ -203,18 +220,39 @@ func TestEcosystemTrace_FailuresAreTracedToo(t *testing.T) { _ = h.handleHexisAct(ctx, actDec("muzick indexer")) - d := findTrace(ecoTraces(t, h), "nexus", "resolve") + d := findTrace(t, h, "nexus", "resolve") if d == nil { t.Fatal("a failed resolve must still be traced") } - if d["status"] != traceFailed { - t.Errorf("status = %v, want failed", d["status"]) + if d.Status != traceRefused { + t.Errorf("status = %v, want refused: the far side answered", d.Status) } - if d["class"] != "unauthorized" { - t.Errorf("class = %v, want unauthorized", d["class"]) + if d.Fields["class"] != "unauthorized" { + t.Errorf("class = %v, want unauthorized", d.Fields["class"]) } - if d["http_status"] != float64(401) { - t.Errorf("http_status = %v, want 401", d["http_status"]) + if d.HTTPStatus != 401 { + t.Errorf("http_status = %v, want 401", d.HTTPStatus) + } +} + +// TestEcosystemTrace_UnreachableIsNotRefused: never got an answer and answered +// with a refusal are different failures, and the trace must say which. +func TestEcosystemTrace_UnreachableIsNotRefused(t *testing.T) { + ctx := context.Background() + h := ecoHandler(t, nil, nil, nil) + h.ecosystem.nexus = newNexusClient("http://127.0.0.1:1") + + _ = h.handleHexisAct(ctx, actDec("muzick indexer")) + + d := findTrace(t, h, "nexus", "resolve") + if d == nil { + t.Fatal("an unreachable resolve must still be traced") + } + if d.Status != traceFailed { + t.Errorf("status = %v, want failed", d.Status) + } + if d.Fields["class"] != "unreachable" { + t.Errorf("class = %v, want unreachable", d.Fields["class"]) } } @@ -227,28 +265,23 @@ func TestEcosystemTrace_RedactsTheUtterance(t *testing.T) { _ = h.handleHexisAct(ctx, actDec("перезапусти кофемашину")) - facts, err := h.dataStore.RecentFacts(ctx, 100) - if err != nil { - t.Fatalf("read facts: %v", err) - } - var traced []store.Fact - for _, f := range facts { - if f.Source == "ecosystem:trace" { - traced = append(traced, f) - } - if strings.Contains(f.Value, "кофемашину") && f.Source == "ecosystem:trace" { - t.Fatalf("trace leaked the utterance: %q", f.Value) - } - } - if len(traced) == 0 { + recorded := traces(t, h) + if len(recorded) == 0 { t.Fatal("expected a not_found resolve trace") } - d := findTrace(ecoTraces(t, h), "nexus", "resolve") - if d["status"] != traceNotFound { - t.Errorf("status = %v, want not_found", d["status"]) + for _, tr := range recorded { + for k, v := range tr.Fields { + if s, ok := v.(string); ok && strings.Contains(s, "кофемашину") { + t.Fatalf("trace leaked the utterance in %s: %q", k, s) + } + } } - if d["subject"] != redactSubject("перезапусти кофемашину") { - t.Errorf("subject = %v, want a redacted length", d["subject"]) + d := findTrace(t, h, "nexus", "resolve") + if d.Status != traceNotFound { + t.Errorf("status = %v, want not_found", d.Status) + } + if d.Fields["subject"] != redactSubject("перезапусти кофемашину") { + t.Errorf("subject = %v, want a redacted length", d.Fields["subject"]) } } @@ -263,7 +296,7 @@ func TestEcosystemTrace_AmbiguityAndConfirmationAreRecorded(t *testing.T) { hexis := newFakeHexis(t, restartCaps(), fixtureHexisExecuted("exec_1", "succeeded")) h := ecoHandler(t, ambig, nil, hexis) _ = h.handleHexisAct(ctx, actDec("muzick")) - if d := findTrace(ecoTraces(t, h), "nexus", "resolve"); d == nil || d["status"] != traceAmbig { + if d := findTrace(t, h, "nexus", "resolve"); d == nil || d.Status != traceAmbig { t.Fatalf("ambiguous resolve must be traced as such, got %+v", d) } @@ -271,7 +304,13 @@ func TestEcosystemTrace_AmbiguityAndConfirmationAreRecorded(t *testing.T) { mutating := fixtureHexisCapabilities(map[string]any{"id": "cap_restart", "name": "restart", "read_only": false}) h2 := ecoHandler(t, nexus, nil, newFakeHexis(t, mutating, fixtureHexisExecuted("exec_1", "succeeded"))) _ = h2.handleHexisAct(ctx, actDec("restart")) - if d := findTrace(ecoTraces(t, h2), "hexis", "confirmation"); d == nil || d["status"] != "pending" { + d := findTrace(t, h2, "hexis", "confirmation") + if d == nil || d.Status != tracePending { t.Fatalf("a parked confirmation must be traced, got %+v", d) } + // The confirmation hop is measured from the top of the action, not from + // the instant it is recorded, which was always zero. + if d.DurationMs == 0 { + t.Error("the confirmation trace must report the time the action took to get there") + } } diff --git a/cmd/mavweb/ecosystem.go b/cmd/mavweb/ecosystem.go index 7198853..77a329c 100644 --- a/cmd/mavweb/ecosystem.go +++ b/cmd/mavweb/ecosystem.go @@ -8,6 +8,8 @@ import ( "net/http" "sync" "time" + + "github.com/kami/maven/internal/ipc" ) // The three sibling services Maven coordinates are headless JSON APIs (no web UI @@ -84,9 +86,13 @@ type ecoData struct { Nexus ecoPanel[ecoEntity] Praxis ecoPanel[ecoItem] Hexis ecoPanel[ecoCap] + Calls ecoPanel[ipc.EcosystemTrace] } -func handleEcosystem(w http.ResponseWriter, r *http.Request, urls ecoURLs) { +// handleEcosystem renders the three sibling panels plus Maven's own log of the +// calls she made to them. The call log comes from core, not from the siblings: +// it is what Maven saw, including the hops that never got an answer. +func handleEcosystem(w http.ResponseWriter, r *http.Request, urls ecoURLs, core ipc.CoreAPI) { ctx := r.Context() var d ecoData var wg sync.WaitGroup @@ -102,6 +108,15 @@ func handleEcosystem(w http.ResponseWriter, r *http.Request, urls ecoURLs) { go func() { defer wg.Done(); d.Hexis.Err = getEco(ctx, urls.hexis, "/api/v1/capabilities", &d.Hexis.Rows) }() wg.Wait() + if core == nil { + d.Calls.Err = "not configured" + } else if rows, err := core.RecentEcosystemTraces(ctx, 50); err != nil { + log.Printf("ecosystem traces: %v", err) + d.Calls.Err = "core read failed" + } else { + d.Calls.Rows = rows + } + w.Header().Set("Content-Type", "text/html; charset=utf-8") if err := ecosystemTmpl.Execute(w, d); err != nil { log.Printf("ecosystem render: %v", err) diff --git a/cmd/mavweb/ecosystem.html b/cmd/mavweb/ecosystem.html index 1bce15f..719173e 100644 --- a/cmd/mavweb/ecosystem.html +++ b/cmd/mavweb/ecosystem.html @@ -41,11 +41,24 @@ {{end}} +
+
+

Calls what Maven asked them

+
+{{with .Calls}} +{{if .Err}}
calls — {{.Err}}
+{{else if not .Rows}}
no ecosystem calls yet.
+{{else}}
+{{range .Rows}}{{end}} +
whenserviceoperationstatusmshttpcorrelation
{{ago .Ts}}{{.Service}}{{.Operation}}{{if eq .Status "ok"}}ok{{else}}{{.Status}}{{end}}{{.DurationMs}}{{if .HTTPStatus}}{{.HTTPStatus}}{{else}}—{{end}}{{.CorrelationID}}
{{end}} +{{end}} +
+ {{template "shellBottom"}}