diff --git a/internal/api/agent.go b/internal/api/agent.go index 2b19313..622d8c1 100644 --- a/internal/api/agent.go +++ b/internal/api/agent.go @@ -8,6 +8,8 @@ import ( "io" "log/slog" "net/http" + "sync" + "time" mdagent "gitea.stevedudenhoeffer.com/steve/majordomo/agent" "gitea.stevedudenhoeffer.com/steve/majordomo/llm" @@ -62,13 +64,19 @@ func (h *handlers) agentChat(c *gin.Context) { return } - send := openEventStream(c) + stream := openEventStream(c) + send := stream.send + // A model thinking hard between tool calls sends nothing for a while, and an + // idle proxy will cut a quiet connection. A comment frame every 20s keeps it + // open; SSE ignores comments, so this costs the client nothing. + stopBeat := stream.keepAlive(20 * time.Second) turn, err := h.agent.Run(c.Request.Context(), actor.ID, req.GardenID, req.Message, replayHistory(history), func(s mdagent.Step) { send(chatEvent{Step: &stepEvent{Index: s.Index, Tools: toolNames(s)}}) }) + stopBeat() if err != nil { // The stream is already open, so an error is an event rather than a // status code — the client has committed to reading a stream by now. @@ -93,25 +101,71 @@ func (h *handlers) agentChat(c *gin.Context) { send(chatEvent{Done: turn}) } -// openEventStream puts the response into SSE mode and returns a sender. +// eventStream serializes writes to one SSE response. +// +// The mutex is load-bearing, not decoration: step events are sent from the +// agent's run goroutine while the keep-alive ticker writes from its own, and two +// goroutines writing a ResponseWriter concurrently is a data race that corrupts +// frames long before it crashes anything. +type eventStream struct { + c *gin.Context + mu sync.Mutex +} + +// openEventStream puts the response into SSE mode. // // Headers go out before the first write and the stream is flushed immediately, // so a proxy holding the response until it looks complete can't reintroduce // exactly the silence streaming exists to remove. -func openEventStream(c *gin.Context) func(chatEvent) { +func openEventStream(c *gin.Context) *eventStream { c.Header("Content-Type", "text/event-stream") c.Header("Cache-Control", "no-cache") c.Header("X-Accel-Buffering", "no") c.Writer.Flush() + return &eventStream{c: c} +} - return func(ev chatEvent) { - b, err := json.Marshal(ev) - if err != nil { - slog.Error("api: encode chat event", "error", err) - return +func (s *eventStream) send(ev chatEvent) { + b, err := json.Marshal(ev) + if err != nil { + slog.Error("api: encode chat event", "error", err) + return + } + s.write(fmt.Sprintf("data: %s\n\n", b)) +} + +func (s *eventStream) write(frame string) { + s.mu.Lock() + defer s.mu.Unlock() + _, _ = io.WriteString(s.c.Writer, frame) + s.c.Writer.Flush() +} + +// keepAlive writes an SSE comment frame on an interval until the returned +// function is called, so a long silence while the model thinks doesn't look like +// a dead connection to whatever sits in between. SSE ignores comment frames, so +// this costs the client nothing. +func (s *eventStream) keepAlive(every time.Duration) func() { + done := make(chan struct{}) + stopped := make(chan struct{}) + go func() { + defer close(stopped) + t := time.NewTicker(every) + defer t.Stop() + for { + select { + case <-done: + return + case <-s.c.Request.Context().Done(): + return + case <-t.C: + s.write(": keep-alive\n\n") + } } - _, _ = fmt.Fprintf(c.Writer, "data: %s\n\n", b) - c.Writer.Flush() + }() + return func() { + close(done) + <-stopped } } diff --git a/internal/service/revisions.go b/internal/service/revisions.go index a0ae39b..273f469 100644 --- a/internal/service/revisions.go +++ b/internal/service/revisions.go @@ -130,11 +130,7 @@ func (s *Service) WithChangeSet(ctx context.Context, actorID, gardenID int64, op // which is exactly the situation undo exists for. Record what happened, // mark the summary, and still report the failure. sc.summary = partialSummary(sc.summary) - // Detached from cancellation on purpose. The commonest reason fn failed is - // that ctx was cancelled or timed out — and using that same dead context - // to record what committed would fail too, losing the history for changes - // that really happened. That is precisely the case this path exists for. - if _, cerr := s.commitScope(context.WithoutCancel(ctx), sc, nil); cerr != nil { + if _, cerr := s.commitScope(ctx, sc, nil); cerr != nil { slog.Error("service: partial turn could not be recorded", "error", cerr, "garden", gardenID) } return nil, err @@ -154,11 +150,23 @@ func partialSummary(summary string) string { // commitScope writes a scope's buffered revisions as one change set. revertsID is // set only by RevertChangeSet. A scope with no revisions writes nothing — an // operation that changed nothing doesn't belong in history. +// +// The write is DETACHED FROM CANCELLATION, always. By the time it runs, the data +// changes it describes have already committed — so cancelling it cannot undo +// anything, it can only lose the record of what happened and leave real changes +// with no way to undo them. +// +// This was originally only done on the failure path, on the reasoning that a +// cancelled context is why fn failed. That missed the commoner case: an agent +// turn whose client disconnects can still COMPLETE, and then the success path +// commits with a dead context and loses the history anyway. Found in production +// with 18 plantings and no change set behind them. func (s *Service) commitScope(ctx context.Context, sc *changeScope, revertsID *int64) (*domain.ChangeSet, error) { revs := sc.taken() if len(revs) == 0 { return nil, nil } + ctx = context.WithoutCancel(ctx) return s.store.WriteChangeSet(ctx, &domain.ChangeSet{ GardenID: sc.gardenID, ActorID: sc.actorID, @@ -312,10 +320,9 @@ func (s *Service) RevertChangeSet(ctx context.Context, actorID, changeSetID int6 changes, conflict, err := s.applyInverse(ctx, r, applied) if err != nil { // Record what actually landed before surfacing the failure, so the - // partial revert is visible and undoable rather than orphaned. Detached - // from cancellation for the same reason as WithChangeSet's path: a - // timed-out context can't be used to write the record of what it did. - if _, cerr := s.commitScope(context.WithoutCancel(ctx), sc, &target.ID); cerr != nil { + // partial revert is visible and undoable rather than orphaned. + // (commitScope detaches from cancellation itself.) + if _, cerr := s.commitScope(ctx, sc, &target.ID); cerr != nil { slog.Error("service: partial revert could not be recorded", "error", cerr, "changeSet", changeSetID) } return nil, nil, err diff --git a/internal/service/revisions_test.go b/internal/service/revisions_test.go index fecd22f..15234f9 100644 --- a/internal/service/revisions_test.go +++ b/internal/service/revisions_test.go @@ -805,3 +805,57 @@ func countsTotal(counts []domain.ChangeCount) int { } return n } + +// TestSucceededTurnRecordsEvenIfTheCallerWentAway is the production bug. +// +// The failure path was detached from cancellation; the SUCCESS path was not. An +// agent turn whose client disconnects can still COMPLETE — and then the success +// path committed with a dead context, the write failed, and real changes were +// left with no change set and no way to undo them. Found live with 18 plantings +// behind no history at all. +// +// Nothing about a commit needs the caller to still be there: by the time it +// runs, the data it describes has already been written. +func TestSucceededTurnRecordsEvenIfTheCallerWentAway(t *testing.T) { + s := newTestService(t, openConfig()) + owner := seedUser(t, s, "a@example.com") + g := seedGarden(t, s, owner) + bed := seedBed(t, s, owner, g.ID) + plant := seedOwnPlant(t, s, owner, 15) + + ctx, cancel := context.WithCancel(context.Background()) + before := len(history(t, s, owner, g.ID)) + + // fn does its work and SUCCEEDS, but the caller goes away before it returns. + cs, err := s.WithChangeSet(ctx, owner, g.ID, ChangeSetOptions{ + Source: domain.SourceAgent, Summary: "plant beans in the second bed", + }, func(ctx context.Context) error { + if _, err := s.FillNamedRegion(ctx, owner, bed.ID, "all", plant.ID, nil); err != nil { + return err + } + cancel() // the client disconnects, mid-turn, after the work landed + return nil + }) + if err != nil { + t.Fatalf("WithChangeSet: %v", err) + } + if cs == nil { + t.Fatal("the turn changed things but produced no change set") + } + + after := history(t, s, owner, g.ID) + if len(after) != before+1 { + t.Fatalf("recorded %d change sets, want 1", len(after)-before) + } + if after[0].Summary != "plant beans in the second bed" { + t.Errorf("summary = %q; a completed turn shouldn't be marked partial", after[0].Summary) + } + // And it's undoable, which is the entire point. + if _, conflicts, err := s.RevertChangeSet(context.Background(), owner, cs.ID, domain.SourceUI); err != nil || len(conflicts) != 0 { + t.Fatalf("the recorded turn should be undoable: err=%v conflicts=%+v", err, conflicts) + } + active, _ := s.store.ListActivePlantingsForObject(context.Background(), bed.ID) + if len(active) != 0 { + t.Errorf("%d plantings survived the undo", len(active)) + } +}