Files
multica/server/internal/daemon/client_test.go
Bohan Jiang 30318b79bc MUL-5426: fix(daemon): retire sessions whose history the provider refuses to replay (#6083)
* fix(daemon): retire sessions whose history the provider refuses to replay

A run killed mid-reply (machine shutdown, force-quit, SIGKILL) can leave an
empty assistant message in the agent CLI's transcript. Every later resume
replays it, the provider rejects the request, and the (agent, issue) pair is
bricked with no self-healing and no user-facing recovery.

Multica already has the mechanism for this — poisoned-session classification —
but its detector paired "400" with "invalid_request_error", which is the
Anthropic wire shape. The same defect reported by any other provider carried
neither token, so it classified as agent_error.unknown: resume-safe by
omission. GetLastTaskSession kept handing back the dead session on every
follow-up, manual Rerun resolved it through the same predicate, and the
in-turn fresh-session retry never fired because ResumeRejected is false here
(nothing rejected the resume — the transcript loaded and the provider refused
to replay it).

Add taskfailure.UnresumableHistory, which recognises the defect by what the
provider says is wrong — some content is empty, and here is which message in
the history — rather than by status code or provider name. Both signals are
required, so a tool reporting "field must not be empty" does not match.

Wire it into the four places that decide whether a session survives:

- classifyPoisonedError, so the task is written as api_invalid_request
- shouldRetryWithFreshSession, so the turn recovers on all 17 backends
  instead of the subset whose adapter learned to detect it; the tools == 0
  gate is unchanged, so a run that already used a tool is never re-run
- ResumeUnsafeFailure, covering the manual-Rerun path
- both resume queries, as defense-in-depth for hosts whose daemon predates
  this (self-host daemons upgrade on their own cadence)

Fixes #6066. Also covers the daemon half of #5760.

Co-authored-by: multica-agent <github@multica.ai>

* fix(session): close the Chat and fresh-retry paths that resurrect a poisoned session

Review found the previous commit stopped short in two places, both of which
put the dead transcript back in play.

Chat never consulted the guarded query. The claim handler reads
chat_session.session_id first and only falls back to GetLastChatTaskSession
when it is empty, so a poisoned pointer there bypasses every filter that query
applies. The fail path merely declined to OVERWRITE the pointer, leaving it in
place. It now clears it in the same transaction, matched on session and
runtime so a concurrent turn's newer pointer survives. The promote guard moves
to ResumeUnsafeFailure as well — the reason-only check passed an un-upgraded
daemon's agent_error.unknown row and re-pinned what the clear had just removed.

GetLastChatTaskSession also kept the row-level filter the issue query dropped
in GH #5975: it discarded the newest poisoned row and fell back to an older
completed row carrying the same dead session. It now judges each session by
its latest terminal state, matching GetLastTaskSession.

A recovered turn could not retire anything. A terminal report carried one
session_id, and an empty one meant both "nothing to report" and "forget the
old session", so a fresh-session retry that SUCCEEDED left the id it retried
away from selectable — through an older completed row on the issue, or through
the chat pointer. agent_task_queue.retired_session_id records the abandonment
itself, reported on every terminal path including completed, and both resume
lookups exclude it. This is the contract gap the previous PR deferred; the
fresh-retry path now runs on all backends, so deferring it is not safe.

Also narrows what the cross-backend test claims: it pins the shared decision,
not that all 17 adapters surface the error into Result.Error (#5760 is the
counter-example), and says so.

Co-authored-by: multica-agent <github@multica.ai>

* test(session): require pgx.ErrNoRows in the resume-exclusion assertions

The `if err == nil && prior.SessionID.Valid` form these tests shared is
false-green: any real fault — undefined column, syntax error, dead connection
— makes err non-nil, so the condition is false and the test passes. Run
against a database missing this branch's new column, the exclusion tests
reported PASS on a SQLSTATE 42703, meaning they could not have caught a broken
query.

requireSessionExcluded demands pgx.ErrNoRows specifically and fails loudly on
anything else, so a green run now means the filter worked rather than the
query never ran.

Applied to all nine sites, not just the four this branch added: the other five
guard the same GetLastTaskSession exclusion behaviour that this branch
changes, so leaving them false-green would leave the change under-tested. All
nine pass on a correctly migrated database.

Co-authored-by: multica-agent <github@multica.ai>

---------

Co-authored-by: Bohan-J <bohan@devv.ai>
Co-authored-by: multica-agent <github@multica.ai>
2026-07-29 15:54:51 +08:00

473 lines
15 KiB
Go

package daemon
import (
"context"
"encoding/json"
"errors"
"net/http"
"net/http/httptest"
"runtime"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/multica-ai/multica/server/pkg/protocol"
)
func TestClient_IdentityHeaders_PostJSON(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if got := r.Header.Get("X-Client-Platform"); got != "daemon" {
t.Errorf("expected X-Client-Platform daemon, got %q", got)
}
if got := r.Header.Get("X-Client-Version"); got != "9.9.9" {
t.Errorf("expected X-Client-Version 9.9.9, got %q", got)
}
if got := r.Header.Get("X-Client-OS"); got != normalizeGOOS(runtime.GOOS) {
t.Errorf("expected X-Client-OS %q, got %q", normalizeGOOS(runtime.GOOS), got)
}
if got := r.Header.Get("Authorization"); got != "Bearer tok" {
t.Errorf("expected Authorization Bearer tok, got %q", got)
}
capabilities := make(map[string]bool)
for _, capability := range strings.Split(r.Header.Get("X-Client-Capabilities"), ",") {
capabilities[strings.TrimSpace(capability)] = true
}
for _, want := range []string{
protocol.DaemonCapabilitySkillBundlesV1,
protocol.DaemonCapabilityCoalescedCommentsV1,
} {
if !capabilities[want] {
t.Errorf("X-Client-Capabilities missing %q: %v", want, capabilities)
}
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]string{"ok": "1"})
}))
defer srv.Close()
c := NewClient(srv.URL)
c.SetToken("tok")
c.SetVersion("9.9.9")
if err := c.postJSON(context.Background(), "/api/daemon/test", map[string]any{}, nil); err != nil {
t.Fatalf("postJSON: %v", err)
}
}
func TestClient_IdentityHeaders_GetJSON(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if got := r.Header.Get("X-Client-Platform"); got != "daemon" {
t.Errorf("expected X-Client-Platform daemon, got %q", got)
}
if got := r.Header.Get("X-Client-Version"); got != "1.2.3" {
t.Errorf("expected X-Client-Version 1.2.3, got %q", got)
}
if got := r.Header.Get("X-Client-OS"); got == "" {
t.Errorf("expected X-Client-OS to be set")
}
w.Header().Set("Content-Type", "application/json")
w.Write([]byte(`{}`))
}))
defer srv.Close()
c := NewClient(srv.URL)
c.SetToken("tok")
c.SetVersion("1.2.3")
var out map[string]any
if err := c.getJSON(context.Background(), "/api/daemon/test", &out); err != nil {
t.Fatalf("getJSON: %v", err)
}
}
func TestClient_VersionOmittedWhenUnset(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if got := r.Header.Get("X-Client-Platform"); got != "daemon" {
t.Errorf("expected X-Client-Platform daemon, got %q", got)
}
// SetVersion not called → header must be omitted (not "").
if vals := r.Header.Values("X-Client-Version"); len(vals) != 0 {
t.Errorf("expected X-Client-Version absent, got %v", vals)
}
w.WriteHeader(http.StatusNoContent)
}))
defer srv.Close()
c := NewClient(srv.URL)
if err := c.postJSON(context.Background(), "/api/daemon/test", nil, nil); err != nil {
t.Fatalf("postJSON: %v", err)
}
}
func TestClient_ListWorkspacesUsesDaemonEndpointAndETag(t *testing.T) {
var calls atomic.Int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/api/daemon/workspaces" {
t.Errorf("path = %q, want /api/daemon/workspaces", r.URL.Path)
http.NotFound(w, r)
return
}
call := calls.Add(1)
if call == 1 {
if got := r.Header.Get("If-None-Match"); got != "" {
t.Errorf("first If-None-Match = %q, want empty", got)
}
w.Header().Set("ETag", `W/"workspace-v1"`)
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`[{"id":"ws-1","name":"One"}]`))
return
}
if got := r.Header.Get("If-None-Match"); got != `W/"workspace-v1"` {
t.Errorf("If-None-Match = %q, want cached ETag", got)
}
w.WriteHeader(http.StatusNotModified)
}))
defer srv.Close()
c := NewClient(srv.URL)
first, err := c.ListWorkspaces(context.Background())
if err != nil {
t.Fatalf("first ListWorkspaces: %v", err)
}
second, err := c.ListWorkspaces(context.Background())
if err != nil {
t.Fatalf("second ListWorkspaces: %v", err)
}
if len(first) != 1 || len(second) != 1 || second[0] != first[0] {
t.Fatalf("cached workspaces mismatch: first=%+v second=%+v", first, second)
}
}
func TestClient_ListWorkspacesFallsBackToLegacyEndpointOnce(t *testing.T) {
var daemonCalls atomic.Int32
var legacyCalls atomic.Int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/api/daemon/workspaces":
daemonCalls.Add(1)
http.NotFound(w, r)
case "/api/workspaces":
legacyCalls.Add(1)
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`[{"id":"ws-legacy","name":"Legacy"}]`))
default:
http.NotFound(w, r)
}
}))
defer srv.Close()
c := NewClient(srv.URL)
for i := 0; i < 2; i++ {
workspaces, err := c.ListWorkspaces(context.Background())
if err != nil {
t.Fatalf("ListWorkspaces call %d: %v", i+1, err)
}
if len(workspaces) != 1 || workspaces[0].ID != "ws-legacy" {
t.Fatalf("workspaces = %+v, want legacy response", workspaces)
}
}
if got := daemonCalls.Load(); got != 1 {
t.Fatalf("daemon endpoint calls = %d, want 1", got)
}
if got := legacyCalls.Load(); got != 2 {
t.Fatalf("legacy endpoint calls = %d, want 2", got)
}
}
// noSleepRetry replaces retrySleep with an immediate no-op so tests don't
// actually wait the 4s/8s/16s/... backoffs. Returns a restore func.
func noSleepRetry(t *testing.T) func() {
t.Helper()
prev := retrySleep
retrySleep = func(ctx context.Context, _ time.Duration) error {
if err := ctx.Err(); err != nil {
return err
}
return nil
}
return func() { retrySleep = prev }
}
func TestIsTransientError(t *testing.T) {
cases := []struct {
name string
err error
want bool
}{
{"nil is not transient", nil, false},
{"5xx is transient", &requestError{StatusCode: http.StatusBadGateway}, true},
{"503 is transient", &requestError{StatusCode: http.StatusServiceUnavailable}, true},
{"408 is transient", &requestError{StatusCode: http.StatusRequestTimeout}, true},
{"429 is transient", &requestError{StatusCode: http.StatusTooManyRequests}, true},
{"400 is permanent", &requestError{StatusCode: http.StatusBadRequest}, false},
{"401 is permanent", &requestError{StatusCode: http.StatusUnauthorized}, false},
{"404 is permanent", &requestError{StatusCode: http.StatusNotFound}, false},
{"transport-level error is transient", errors.New("connection reset by peer"), true},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
if got := isTransientError(tc.err); got != tc.want {
t.Fatalf("isTransientError(%v) = %v, want %v", tc.err, got, tc.want)
}
})
}
}
func TestIsIssueGCBatchUnsupported(t *testing.T) {
tests := []struct {
name string
err error
want bool
}{
{
name: "old server unmatched route",
err: &requestError{StatusCode: http.StatusNotFound, Body: "404 page not found"},
want: true,
},
{
name: "workspace access denied",
err: &requestError{StatusCode: http.StatusNotFound, Body: `{"error":"not found"}`},
want: false,
},
{
name: "transient server error",
err: &requestError{StatusCode: http.StatusInternalServerError, Body: "failure"},
want: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := isIssueGCBatchUnsupported(tt.err); got != tt.want {
t.Fatalf("isIssueGCBatchUnsupported() = %v, want %v", got, tt.want)
}
})
}
}
func TestPostJSONWithRetry_TransientThenSuccess(t *testing.T) {
defer noSleepRetry(t)()
var calls atomic.Int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
n := calls.Add(1)
if n < 3 {
w.WriteHeader(http.StatusBadGateway)
return
}
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
c := NewClient(srv.URL)
schedule := []time.Duration{time.Nanosecond, time.Nanosecond, time.Nanosecond}
if err := c.postJSONWithRetry(context.Background(), "/x", map[string]any{}, nil, schedule); err != nil {
t.Fatalf("postJSONWithRetry: %v", err)
}
if got := calls.Load(); got != 3 {
t.Fatalf("expected 3 attempts (2 transient + 1 success), got %d", got)
}
}
// TestFailTask_RetriesOnTransient5xxThenSucceeds pins the callback half of
// MUL-5305 Must-fix 1: FailTask's terminal transaction is now the sole
// persistence point for the withheld session and continuity-gap flag, so if the
// server returns a transient 5xx (the terminal tx rolled back), the daemon MUST
// retry until it lands — a 400 would make it bail immediately
// (TestPostJSONWithRetry_PermanentBailsImmediately) and drop the gap forever.
func TestFailTask_RetriesOnTransient5xxThenSucceeds(t *testing.T) {
defer noSleepRetry(t)()
var calls atomic.Int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
if calls.Add(1) < 3 {
w.WriteHeader(http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
c := NewClient(srv.URL)
if err := c.FailTask(context.Background(), "task-1", "boom", "", "", "timeout", true, ""); err != nil {
t.Fatalf("FailTask: %v", err)
}
if got := calls.Load(); got != 3 {
t.Fatalf("expected 3 attempts (2 transient 5xx + 1 success), got %d", got)
}
}
func TestPostJSONWithRetry_TransientExhausts(t *testing.T) {
defer noSleepRetry(t)()
var calls atomic.Int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
calls.Add(1)
w.WriteHeader(http.StatusBadGateway)
}))
defer srv.Close()
c := NewClient(srv.URL)
schedule := []time.Duration{time.Nanosecond, time.Nanosecond}
err := c.postJSONWithRetry(context.Background(), "/x", map[string]any{}, nil, schedule)
if err == nil {
t.Fatal("expected error after schedule exhausted, got nil")
}
if !isTransientError(err) {
t.Fatalf("expected transient error, got %v", err)
}
if got := calls.Load(); got != int32(len(schedule)+1) {
t.Fatalf("expected %d attempts (initial + %d retries), got %d", len(schedule)+1, len(schedule), got)
}
}
func TestPostJSONWithRetry_PermanentBailsImmediately(t *testing.T) {
defer noSleepRetry(t)()
var calls atomic.Int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
calls.Add(1)
w.WriteHeader(http.StatusBadRequest)
}))
defer srv.Close()
c := NewClient(srv.URL)
schedule := []time.Duration{time.Nanosecond, time.Nanosecond, time.Nanosecond}
err := c.postJSONWithRetry(context.Background(), "/x", map[string]any{}, nil, schedule)
if err == nil {
t.Fatal("expected error, got nil")
}
if got := calls.Load(); got != 1 {
t.Fatalf("expected exactly 1 attempt on permanent error, got %d", got)
}
}
func TestPostJSONWithRetry_CtxCancelStopsRetries(t *testing.T) {
// Use the real sleeper here so we can observe a cancel preempting it.
var calls atomic.Int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
calls.Add(1)
w.WriteHeader(http.StatusBadGateway)
}))
defer srv.Close()
ctx, cancel := context.WithCancel(context.Background())
go func() {
// Cancel quickly so the first sleep is aborted long before its 1s.
time.Sleep(50 * time.Millisecond)
cancel()
}()
c := NewClient(srv.URL)
schedule := []time.Duration{time.Second, time.Second, time.Second}
start := time.Now()
err := c.postJSONWithRetry(ctx, "/x", map[string]any{}, nil, schedule)
elapsed := time.Since(start)
if err == nil {
t.Fatal("expected error after ctx cancel, got nil")
}
if elapsed > 750*time.Millisecond {
t.Fatalf("expected ctx cancel to short-circuit retry, took %s", elapsed)
}
if got := calls.Load(); got != 1 {
t.Fatalf("expected exactly 1 attempt before cancel, got %d", got)
}
}
func TestDefaultTerminalRetrySchedule_MatchesAgreedPlan(t *testing.T) {
// MUL-2780 settled on a 5-step exponential backoff (4s, 8s, 16s, 32s, 64s).
// Pin it so a future "tidy this up" refactor can't silently flatten or
// shorten the recovery window without explicit discussion.
want := []time.Duration{4 * time.Second, 8 * time.Second, 16 * time.Second, 32 * time.Second, 64 * time.Second}
if len(defaultTerminalRetrySchedule) != len(want) {
t.Fatalf("schedule length: got %d, want %d", len(defaultTerminalRetrySchedule), len(want))
}
for i, d := range want {
if defaultTerminalRetrySchedule[i] != d {
t.Errorf("schedule[%d]: got %s, want %s", i, defaultTerminalRetrySchedule[i], d)
}
}
}
func TestNormalizeGOOS(t *testing.T) {
cases := map[string]string{
"darwin": "macos",
"windows": "windows",
"linux": "linux",
"freebsd": "freebsd",
}
for in, want := range cases {
if got := normalizeGOOS(in); got != want {
t.Errorf("normalizeGOOS(%q) = %q, want %q", in, got, want)
}
}
}
// TestTerminalReportsCarryRetiredSessionID pins the daemon half of the
// retire-session contract (GH #6066). Before it, a terminal report could only
// say "here is a session" or say nothing — and saying nothing was how a
// recovered turn silently left the poisoned id selectable. The completed path
// matters most: that is exactly the case where a fresh-session retry SUCCEEDED
// and the abandoned transcript would otherwise survive on an older row.
func TestTerminalReportsCarryRetiredSessionID(t *testing.T) {
for _, tc := range []struct {
name string
endpoint string
call func(*Client) error
}{
{
name: "complete",
endpoint: "/api/daemon/tasks/task-1/complete",
call: func(c *Client) error {
return c.CompleteTask(context.Background(), "task-1", "done", "", "", "/tmp/wd", false, "POISONED-S")
},
},
{
name: "fail",
endpoint: "/api/daemon/tasks/task-1/fail",
call: func(c *Client) error {
return c.FailTask(context.Background(), "task-1", "boom", "", "/tmp/wd", "api_invalid_request", false, "POISONED-S")
},
},
} {
t.Run(tc.name, func(t *testing.T) {
var body map[string]any
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != tc.endpoint {
t.Errorf("unexpected path %q", r.URL.Path)
}
_ = json.NewDecoder(r.Body).Decode(&body)
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
if err := tc.call(NewClient(srv.URL)); err != nil {
t.Fatalf("terminal report: %v", err)
}
if got, _ := body["retired_session_id"].(string); got != "POISONED-S" {
t.Fatalf("retired_session_id = %v, want POISONED-S (body: %v)", body["retired_session_id"], body)
}
})
}
}
// TestTerminalReportsOmitEmptyRetiredSessionID keeps the common case off the
// wire: nearly every run retires nothing, and an empty field would be
// indistinguishable from "retire the empty session".
func TestTerminalReportsOmitEmptyRetiredSessionID(t *testing.T) {
var body map[string]any
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
_ = json.NewDecoder(r.Body).Decode(&body)
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
if err := NewClient(srv.URL).CompleteTask(context.Background(), "task-1", "done", "", "sess-1", "/tmp/wd", false, ""); err != nil {
t.Fatalf("CompleteTask: %v", err)
}
if _, present := body["retired_session_id"]; present {
t.Fatalf("retired_session_id must be omitted when nothing was retired, got %v", body)
}
}