Commit history detached from cancellation on EVERY path, not just failure
Found by using it: the agent planted 18 beans, the client's stream dropped, and those 18 plantings ended up in the garden with no change set behind them — real changes with no way to undo them, which is the one guarantee this whole design rests on. The earlier fix detached the FAILURE path and reasoned that a cancelled context is why fn failed. That missed the commoner case. An agent turn whose client disconnects can still COMPLETE — the model finishes, the tools have already written — and then the success path committed with a dead context, the write failed, WithChangeSet returned an error, and the work was orphaned. commitScope now detaches from cancellation itself, so every caller gets it. That is the right home for the rule: by the time a commit runs, the data it describes has already been written, so cancelling it cannot undo anything — it can only lose the record of what happened. There is no path on which that is the behaviour anyone wants. Also added an SSE keep-alive, since the dropped connection is worth not having in the first place: a model thinking between tool calls sends nothing for a while, and an idle proxy will cut a quiet stream. A comment frame every 20s keeps it open, and SSE ignores comment frames so the client is unaffected. That introduced a second goroutine writing the response, so the event stream now serializes writes behind a mutex — step events come from the agent's run goroutine while the ticker writes from its own, and two goroutines writing a ResponseWriter concurrently corrupts frames long before it crashes anything. Verified under -race. Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]> Claude-Session: https://claude.ai/code/session_01H3zbym8Doka2d7D48maSgZ
This commit is contained in:
+64
-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"
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user