Commit history detached from cancellation on every path, not just failure #73
+74
-10
@@ -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"
|
||||
@@ -24,6 +26,11 @@ import (
|
||||
// by everything at once, which reads as a hang — and the whole design rests on
|
||||
// watching the canvas change as it happens.
|
||||
|
||||
// keepAliveInterval is how often a quiet stream emits a comment frame. Well
|
||||
// under the 30–60s idle timeout typical of reverse proxies, which is the thing
|
||||
// it exists to stay ahead of.
|
||||
const keepAliveInterval = 20 * time.Second
|
||||
|
||||
// chatRequest is the body of POST /agent/chat.
|
||||
type chatRequest struct {
|
||||
GardenID int64 `json:"gardenId" binding:"required"`
|
||||
@@ -62,13 +69,21 @@ 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. Deferred so a panic in the run
|
||||
// can't leak the ticker goroutine; stopping it twice is harmless.
|
||||
stopBeat := stream.keepAlive(keepAliveInterval)
|
||||
defer stopBeat()
|
||||
|
gitea-actions
commented
🟡 stopBeat is called manually instead of via defer, so a panic in h.agent.Run skips cleanup (mitigated only by request-context cancellation) error-handling · flagged by 2 models
🪰 Gadfly · advisory 🟡 **stopBeat is called manually instead of via defer, so a panic in h.agent.Run skips cleanup (mitigated only by request-context cancellation)**
_error-handling · flagged by 2 models_
- **`internal/api/agent.go:79` — `stopBeat` is called manually instead of via `defer`.** `stopBeat := stream.keepAlive(...)` is started at line 73, then `h.agent.Run(...)` runs (lines 74–78), then `stopBeat()` is called as a plain statement at line 79 (before the error check). If `Run` (or anything between the two lines) panics, `stopBeat()` is never reached. The keep-alive goroutine does have a self-cleanup escape hatch — it also selects on `s.c.Request.Context().Done()` (line 159), and a panic…
<sub>🪰 Gadfly · advisory</sub>
|
||||
|
||||
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 +108,74 @@ 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()
|
||||
}()
|
||||
// Idempotent: the handler stops it explicitly when the run returns and again
|
||||
// via defer, so a panic can't leak the goroutine.
|
||||
var once sync.Once
|
||||
return func() {
|
||||
once.Do(func() { close(done) })
|
||||
<-stopped
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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,19 @@ 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, and that belongs here rather
|
||||
// than at each call site so no caller can be the one that forgets. By the time a
|
||||
// commit runs, the data changes it describes have already been written — 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, which is the one
|
||||
// thing the whole change-set design exists to prevent.
|
||||
|
gitea-actions
commented
🟡 commitScope doc embeds a production incident post-mortem that belongs in the commit message, not godoc maintainability · flagged by 1 model
🪰 Gadfly · advisory 🟡 **commitScope doc embeds a production incident post-mortem that belongs in the commit message, not godoc**
_maintainability · flagged by 1 model_
- **`internal/service/revisions.go:159` — doc comment carries a production post-mortem that will rot.** The `commitScope` comment now embeds a historical narrative ("Found in production with 18 plantings and no change set behind them") alongside the durable rule. The rule ("the write is detached from cancellation, always") is worth keeping; the incident report belongs in the commit message / PR description, not the godoc, where it reads as rationale that future readers can't act on and will be t…
<sub>🪰 Gadfly · advisory</sub>
|
||||
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,
|
||||
@@ -199,9 +203,13 @@ func (s *Service) record(ctx context.Context, gardenID, actorID int64, summary s
|
||||
sc.append(revs)
|
||||
return
|
||||
}
|
||||
if _, err := s.store.WriteChangeSet(ctx, &domain.ChangeSet{
|
||||
GardenID: gardenID, ActorID: actorID, Source: domain.SourceUI, Summary: summary,
|
||||
}, revs); err != nil {
|
||||
// Auto-scope: one operation, its own change set. Written through the same
|
||||
// detached path as everything else — a REST client that hangs up right after
|
||||
// its PATCH landed must not leave that change without history, and this is
|
||||
// the path virtually every mutation takes.
|
||||
auto := &changeScope{gardenID: gardenID, actorID: actorID, source: domain.SourceUI, summary: summary}
|
||||
|
gitea-actions
commented
🟠 Auto-scope record() path still commits history with the caller's (possibly dead) context — same orphan-history bug class the PR claims to close universally, left unfixed for every non-agent mutation correctness, error-handling · flagged by 1 model
🪰 Gadfly · advisory 🟠 **Auto-scope record() path still commits history with the caller's (possibly dead) context — same orphan-history bug class the PR claims to close universally, left unfixed for every non-agent mutation**
_correctness, error-handling · flagged by 1 model_
- **`internal/service/revisions.go:210` — auto-scope `record` path still commits history with the caller's (possibly dead) context, the same orphan-history bug class the PR claims to close universally, left unfixed for every non-agent mutation.** The PR's thesis (the new `commitScope` comment, revisions.go:154–157: "The write is DETACHED FROM CANCELLATION, always… it can only lose the record of what happened and leave real changes with no way to undo them") applies equally to the auto-scope bran…
<sub>🪰 Gadfly · advisory</sub>
|
||||
auto.append(revs)
|
||||
if _, err := s.commitScope(ctx, auto, nil); err != nil {
|
||||
slog.Error("service: record change set", "error", err, "garden", gardenID, "summary", summary)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -805,3 +805,91 @@ 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, "[email protected]")
|
||||
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))
|
||||
}
|
||||
}
|
||||
|
||||
// TestAutoScopedMutationRecordsEvenIfTheCallerWentAway — the same rule for the
|
||||
// path virtually every mutation takes.
|
||||
//
|
||||
// A plain REST PATCH auto-scopes into its own change set. If the client hangs up
|
||||
// between the row landing and the change set being written, that change is
|
||||
// orphaned exactly as an agent turn's was — and this path is used far more.
|
||||
func TestAutoScopedMutationRecordsEvenIfTheCallerWentAway(t *testing.T) {
|
||||
s := newTestService(t, openConfig())
|
||||
owner := seedUser(t, s, "[email protected]")
|
||||
g := seedGarden(t, s, owner)
|
||||
bed := seedBed(t, s, owner, g.ID)
|
||||
before := len(history(t, s, owner, g.ID))
|
||||
|
||||
// The row write and the history write share this context; cancelling after
|
||||
// the mutation returns is the client hanging up mid-request.
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
if _, err := s.UpdateObject(ctx, owner, bed.ID, ObjectPatch{XCM: f64Ptr(300)}, bed.Version); err != nil {
|
||||
t.Fatalf("UpdateObject: %v", err)
|
||||
}
|
||||
cancel()
|
||||
|
||||
after := history(t, s, owner, g.ID)
|
||||
if len(after) != before+1 {
|
||||
t.Fatalf("recorded %d change sets, want 1 — the move is otherwise un-undoable", len(after)-before)
|
||||
}
|
||||
if _, conflicts, err := s.RevertChangeSet(context.Background(), owner, after[0].ID, domain.SourceUI); err != nil || len(conflicts) != 0 {
|
||||
t.Fatalf("the recorded move should be undoable: err=%v conflicts=%+v", err, conflicts)
|
||||
}
|
||||
back, _ := s.store.GetObject(context.Background(), bed.ID)
|
||||
if back.XCM != bed.XCM {
|
||||
t.Errorf("undo left x at %v, want %v", back.XCM, bed.XCM)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user
🟡 Magic number keep-alive interval should be a named constant
maintainability · flagged by 2 models
internal/api/agent.go:73— Magic number heartbeat interval. The keep-alive duration20 * time.Secondis hardcoded as a literal at the call site, while every other timeout/duration in this package is a named constant (oidcDiscoveryTimeout,oidcExchangeTimeout, etc.). This breaks the project's existing pattern, buries the tuning knob, and makes it non-discoverable for anyone adjusting proxy timeouts later. It should be a package-level constant such as `sseKeepAliveInterval = 20 * tim…🪰 Gadfly · advisory