Build image / build-and-push (push) Successful in 5s
Gadfly caught me making the same mistake one level up. This PR claimed the history write is "detached from cancellation, always" — and record()'s auto-scope path still called store.WriteChangeSet directly with the caller's context, so every plain REST mutation kept the orphan-history window the PR was written to close. That is the path virtually every change takes; the agent is the exception. It goes through commitScope now, which is where the rule lives, so no caller can be the one that forgets. Tested the same way as the agent path: cancel right after the mutation returns, assert the change set exists and reverts. Also: stopBeat is deferred so a panic in the run can't leak the ticker goroutine (and is idempotent, since the handler stops it explicitly on the normal path); the keep-alive interval is a named constant rather than a literal duplicated between comment and call site; and commitScope's doc states the rule and why it lives there, instead of narrating the incident that produced it — that belongs in a commit message, which is where it now is. Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]> Claude-Session: https://claude.ai/code/session_01H3zbym8Doka2d7D48maSgZ
254 lines
8.4 KiB
Go
254 lines
8.4 KiB
Go
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})
|
||
}
|
||
|
||
// 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) *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}
|
||
}
|
||
|
||
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")
|
||
}
|
||
}
|
||
}()
|
||
// 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."
|
||
}
|
||
}
|