package api import ( "context" "encoding/json" "errors" "fmt" "io" "log/slog" "net/http" "sync" "time" mdagent "gitea.stevedudenhoeffer.com/steve/majordomo/agent" "gitea.stevedudenhoeffer.com/steve/majordomo/llm" "github.com/gin-gonic/gin" "gitea.stevedudenhoeffer.com/steve/pansy/internal/agent" "gitea.stevedudenhoeffer.com/steve/pansy/internal/domain" ) // The garden assistant's chat surface (#56). // // Streaming, because a turn that clears a bed and replants it makes a dozen tool // calls over tens of seconds. Without streaming that is a long silence followed // 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"` Message string `json:"message" binding:"required"` } // chatEvent is one server-sent event. Exactly one field is set. type chatEvent struct { // Step reports a completed model round trip: which tools it called. Step *stepEvent `json:"step,omitempty"` // Done carries the finished turn. Done *agent.Turn `json:"done,omitempty"` // Error is a turn that failed, in words meant for a person. Error string `json:"error,omitempty"` // Warning rides alongside Done: the turn worked, but something adjacent to it // didn't, and saying nothing would be the quieter lie. Warning string `json:"warning,omitempty"` } type stepEvent struct { Index int `json:"index"` Tools []string `json:"tools"` } func (h *handlers) agentChat(c *gin.Context) { var req chatRequest if err := c.ShouldBindJSON(&req); err != nil { writeAPIError(c, http.StatusBadRequest, "INVALID_INPUT", "a gardenId and a message are required") return } actor := mustActor(c) history, err := h.svc.AgentHistory(c.Request.Context(), actor.ID, req.GardenID) if err != nil { writeServiceError(c, err) return } 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. send(chatEvent{Error: chatErrorMessage(err)}) return } // The turn itself succeeded — the garden really did change — so Done goes out // regardless. But if the transcript couldn't be saved, say so: a clean "done" // followed by a conversation that has forgotten the exchange after a reload is // exactly the kind of quiet inconsistency that makes a tool feel unreliable. // // Detached from the request context, because the commonest reason this fails // is the client having gone away — and the exchange is worth keeping either // way, since the change set it produced certainly is. if _, err := h.svc.RecordAgentExchange(context.WithoutCancel(c.Request.Context()), actor.ID, req.GardenID, req.Message, turn.Reply, turn.ChangeSetID); err != nil { slog.Error("api: record agent exchange", "error", err, "garden", req.GardenID) send(chatEvent{Done: turn, Warning: "I couldn't save this exchange, so it won't be here after a reload. Anything I changed is still on the canvas, and in History."}) return } send(chatEvent{Done: turn}) } // sseWriteTimeout bounds ONE write to the stream, not the stream itself. // // It is refreshed per frame, which is the only shape that satisfies both ends: // the server's absolute WriteTimeout would cut a long turn (#78), while removing // the deadline entirely would let a client that stops reading block a write // forever once the socket buffer fills — pinning the run goroutine and this // stream's mutex with it, and taking the keep-alive down too since it needs the // same lock. Generous, because it is a backstop against a stuck peer and not a // pacing mechanism. // // A var, not a const, ONLY so the test can shrink it to prove the deadline is // refreshed per frame rather than set once — a set-once 30s deadline would pass // a test whose whole run is under a second. Production never reassigns it. var sseWriteTimeout = 30 * time.Second // 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 rc *http.ResponseController 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. // // Taking the write deadline off the server's absolute WriteTimeout and onto a // per-write one is what makes a turn longer than 30s possible at all (#78). // WriteTimeout is an ABSOLUTE deadline from when the request header was read, // not an idle timeout, so a streaming response is cut mid-turn however recently // it wrote. Without this the 4-minute runTimeout is unreachable and the // keep-alive below tops out at one tick — pacing a connection that is destroyed // underneath it. // // That failure is INVISIBLE from in here: writes past the deadline return // err == nil and their bytes are dropped, so there is nothing to detect on the // write path. Only the client sees it, as a truncated stream it reports as a // dropped connection. Hence a deadline set up front and refreshed per frame, // rather than anything checked after the fact. 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") s := &eventStream{c: c, rc: http.NewResponseController(c.Writer)} // Probe once here rather than reporting per frame: a writer that can't take // deadlines will fail identically on every write, and the operator needs to // hear it once. If this fails the stream still works — it is just back to // being cut at WriteTimeout, which is worth saying out loud. if err := s.rc.SetWriteDeadline(time.Now().Add(sseWriteTimeout)); err != nil { slog.Error("api: SSE write deadlines unavailable; long turns will be truncated at the server WriteTimeout", "error", err) } c.Writer.Flush() return s } 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() // Refresh for THIS write, so the stream as a whole is unbounded but no single // write is. Error deliberately unchecked: openEventStream already reported // whether deadlines work at all, and this call can only fail the same way, so // checking here would log once per frame to say the same thing. _ = s.rc.SetWriteDeadline(time.Now().Add(sseWriteTimeout)) _, _ = 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") } } }() // 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 } } // getAgentHistory returns the actor's thread for a garden. func (h *handlers) getAgentHistory(c *gin.Context) { gardenID, ok := parseIDParam(c, "id") if !ok { return } msgs, err := h.svc.AgentHistory(c.Request.Context(), mustActor(c).ID, gardenID) if err != nil { writeServiceError(c, err) return } c.JSON(http.StatusOK, gin.H{"messages": msgs}) } // deleteAgentHistory is the "start over" escape hatch. func (h *handlers) deleteAgentHistory(c *gin.Context) { gardenID, ok := parseIDParam(c, "id") if !ok { return } if err := h.svc.ClearAgentHistory(c.Request.Context(), mustActor(c).ID, gardenID); err != nil { writeServiceError(c, err) return } c.Status(http.StatusNoContent) } // replayHistory turns stored text into model messages. // // Only the text is replayed — no stored tool calls. Continuity needs what was // said and what came back; replaying a tool call would be replaying a decision // made against a garden that has since moved on. func replayHistory(msgs []domain.AgentMessage) []llm.Message { out := make([]llm.Message, 0, len(msgs)) for _, m := range msgs { if m.Role == domain.AgentRoleUser { out = append(out, llm.UserText(m.Body)) } else { out = append(out, llm.AssistantText(m.Body)) } } return out } func toolNames(s mdagent.Step) []string { names := make([]string, 0, len(s.Results)) for _, r := range s.Results { names = append(names, r.Name) } return names } // chatErrorMessage turns a failure into something worth reading. Permission // errors in particular get named: "the agent broke" and "you can only view this // garden" want very different reactions. func chatErrorMessage(err error) string { switch { case errors.Is(err, domain.ErrForbidden): return "You can only view this garden, so I can't change anything in it." case errors.Is(err, domain.ErrNotFound): return "I can't find that garden." case errors.Is(err, domain.ErrInvalidInput): return "I didn't get a message to work from." case errors.Is(err, io.EOF), errors.Is(err, context.Canceled): return "The connection dropped partway through. Anything I'd already changed is on the canvas, and in History." case errors.Is(err, context.DeadlineExceeded): return "That took too long and I stopped. Anything I'd already changed is on the canvas, and in History." default: slog.Error("api: agent run failed", "error", err) return "Something went wrong talking to the model. Anything I'd already changed is on the canvas, and in History." } }