Commit history detached from cancellation on every path, not just failure (#73)
Build image / build-and-push (push) Successful in 7s
Build image / build-and-push (push) Successful in 7s
Co-authored-by: Steve Dudenhoeffer <[email protected]>
This commit was merged in pull request #73.
This commit is contained in:
+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()
|
||||
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user