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}}whenserviceoperationstatusmshttpcorrelation +{{range .Rows}}{{ago .Ts}}{{.Service}}{{.Operation}}{{if eq .Status "ok"}}ok{{else}}{{.Status}}{{end}}{{.DurationMs}}{{if .HTTPStatus}}{{.HTTPStatus}}{{else}}—{{end}}{{.CorrelationID}}{{end}} +{{end}} +{{end}} + + {{template "shellBottom"}}