Thinking
Let's start by understanding the task. Issue #59: "Retry logic" β we need retry logic in the agent loop in case the LLM endpoint fails. The comments show actual failure examples: 429 (rate limit / model cooldown) and 502 (engine not healthy).
First, let's explore the project structure to understand the codebase.
Tool
π§ list_files
{"path":".","recursive":true}
Arguments
{"path":".","recursive":true}
Result
/project/README.md
/project/cmd
/project/cmd/zoo
/project/cmd/zoo/main.go
/project/docker
/project/docker/Dockerfile
/project/docker/group
/project/docker/passwd
/project/docker/shadow
/project/docker/uid1000-sudo
/project/go.mod
/project/go.sum
/project/internal
/project/internal/agentrun
/project/internal/agentrun/context.go
/project/internal/agentrun/docker.go
/project/internal/agentrun/docker_integration_test.go
/project/internal/agentrun/forgejo.go
/project/internal/agentrun/loop.go
/project/internal/agentrun/roster.go
/project/internal/agentrun/roster_test.go
/project/internal/agentrun/run.go
/project/internal/agentrun/sandboxgit.go
/project/internal/agentrun/sandboxgit_test.go
/project/internal/agentrun/seed_test.go
/project/internal/agentrun/system.md
/project/internal/config
/project/internal/config/config.go
/project/internal/config/config_test.go
/project/internal/forgejo
/project/internal/forgejo/client.go
/project/internal/forgejo/client_test.go
/project/internal/forgejo/event.go
/project/internal/forgejo/poll.go
/project/internal/forgejo/poll_test.go
/project/internal/forgejo/repo_filter.go
/project/internal/forgejo/repo_filter_test.go
/project/internal/forgejo/watcher.go
/project/internal/forgejo/watcher_test.go
/project/internal/forgejo/webhook.go
/project/internal/forgejo/webhook_test.go
/project/internal/livelog
/project/internal/livelog/hub.go
/project/internal/livelog/hub_test.go
/project/internal/llm
/project/internal/llm/client.go
/project/internal/llm/streamer.go
/project/internal/llm/streamer_test.go
/project/internal/llm/types.go
/project/internal/scheduler
/project/internal/scheduler/scheduler.go
/project/internal/scheduler/scheduler_test.go
/project/internal/store
/project/internal/store/schema.sql
/project/internal/store/store.go
/project/internal/store/store_test.go
/project/internal/tools
/project/internal/tools/bash.go
/project/internal/tools/context.go
/project/internal/tools/edit-file.go
/project/internal/tools/forgejo.go
/project/internal/tools/grep_search.go
/project/internal/tools/list_files.go
/project/internal/tools/move_file.go
/project/internal/tools/read_file.go
/project/internal/tools/registry.go
/project/internal/tools/remove_file.go
/project/internal/tools/shell.go
/project/internal/tools/tools_test.go
/project/internal/tools/write_file.go
/project/internal/web
/project/internal/web/templates.go
/project/internal/web/web.go
/project/internal/web/web_test.go
/project/zoo.hcl.example
Tool
π§ bash
{"command":"cd /project && git log --oneline -5 && git status"}
Arguments
{"command":"cd /project && git log --oneline -5 && git status"}
Result
98298f6 Merge pull request 'Fix excessive vertical space in job view blocks (#53)' (#55) from issue-53-greg into main
b75b64d Fix excessive vertical space in job view blocks
10c3dce Better Docker/git inter.
315f993 Add docker image build.
5717258 Merge pull request 'forgejo: repository allowlist with wildcard patterns' (#40) from issue-39-greg into main
On branch issue-59-greg
nothing to commit, working tree clean
Thinking
Let's take a look at the main files: the LLM client, the agent loop, and the config.
Tool
π§ read_file
{"path":"internal/llm/client.go","start_line":1,"end_line":400}
Arguments
{"path":"internal/llm/client.go","start_line":1,"end_line":400}
Result
{"path":"internal/llm/client.go","content":"package llm\n\nimport (\n\t\"bytes\"\n\t\"context\"\n\t\"encoding/json\"\n\t\"net/http\"\n\n\t\"github.com/abrander/zoo/internal/config\"\n)\n\ntype Client struct {\n\tendpoint string\n\ttoken string\n\tmodel string\n\n\thttpClient *http.Client\n}\n\n// NewClient builds a Client from a configured llm block.\nfunc NewClient(cfg config.LLM) *Client {\n\treturn \u0026Client{\n\t\tendpoint: cfg.OpenAI,\n\t\ttoken: cfg.Token,\n\t\tmodel: cfg.Model,\n\t\thttpClient: http.DefaultClient,\n\t}\n}\n\nfunc (c *Client) StreamChatCompletion(ctx context.Context, req *ChatCompletionRequest) (*Streamer, error) {\n\treq2 := *req\n\n\tif req2.Model == \"\" {\n\t\treq2.Model = c.model\n\t}\n\n\treq2.Stream = true\n\treq2.Seed = 5\n\treq2.Temperature = 0.6\n\treq2.TopP = 0.95\n\n\tbody, err := json.Marshal(req2)\n\tif err != nil {\n\t\treturn nil, err\n\t}\n\n\t// Bound to ctx so a stalled connection or silent server (headers or\n\t// stream body never arriving) is unblocked by the run's deadline\n\t// instead of hanging the goroutine forever β http.DefaultClient has\n\t// no timeout of its own.\n\tr, err := http.NewRequestWithContext(ctx, http.MethodPost, c.endpoint+\"/chat/completions\", bytes.NewBuffer(body))\n\tif err != nil {\n\t\treturn nil, err\n\t}\n\n\tr.Header.Set(\"Authorization\", \"Bearer \"+c.token)\n\tr.Header.Set(\"Content-Type\", \"application/json\")\n\n\tresp, err := c.httpClient.Do(r)\n\tif err != nil {\n\t\treturn nil, err\n\t}\n\n\treturn NewStreamer(resp)\n}","start_line":1,"end_line":65,"total_lines":65,"truncated":false}
Tool
π§ read_file
{"path":"internal/agentrun/loop.go","start_line":1,"end_line":400}
Arguments
{"path":"internal/agentrun/loop.go","start_line":1,"end_line":400}
Result
{"path":"internal/agentrun/loop.go","content":"package agentrun\n\nimport (\n\t\"context\"\n\t\"fmt\"\n\t\"io\"\n\n\t\"github.com/abrander/zoo/internal/llm\"\n\t\"github.com/abrander/zoo/internal/tools\"\n)\n\n// Hooks lets a caller observe a runLoop invocation live, as the model\n// streams and tool calls execute, in addition to the final []llm.Message\n// it returns. Any of these may be nil.\ntype Hooks struct {\n\t// OnReasoningDelta and OnContentDelta fire with just the newly\n\t// streamed text for the current turn, not the accumulated total.\n\tOnReasoningDelta func(delta string)\n\tOnContentDelta func(delta string)\n\n\t// OnTurnEnd fires once per completed streamer round-trip, after the\n\t// model's message for that turn is fully received and before any of\n\t// its tool calls run.\n\tOnTurnEnd func()\n\n\t// OnTool fires once per tool call, after it has run.\n\tOnTool func(name, arguments, result string, toolErr bool)\n}\n\n// runLoop is a headless port of ../a's App.generate(): send messages +\n// tool defs, get a completion, run any tool_calls and append their\n// results, repeat until a plain finish or ctx is done.\nfunc runLoop(ctx context.Context, client *llm.Client, toolsCtx tools.Context, messages []llm.Message, hooks Hooks) ([]llm.Message, error) {\n\tfor {\n\t\tif err := ctx.Err(); err != nil {\n\t\t\treturn messages, err\n\t\t}\n\n\t\tstreamer, err := client.StreamChatCompletion(ctx, \u0026llm.ChatCompletionRequest{\n\t\t\tMessages: messages,\n\t\t\tStream: true,\n\t\t\tTools: tools.All(),\n\t\t})\n\t\tif err != nil {\n\t\t\treturn messages, fmt.Errorf(\"chat completion: %w\", err)\n\t\t}\n\n\t\tvar completion *llm.ChatCompletion\n\n\t\tvar prevContent, prevReasoning string\n\n\t\tfor {\n\t\t\tc, err := streamer.Get()\n\t\t\tif err == io.EOF {\n\t\t\t\tbreak\n\t\t\t}\n\n\t\t\tif err != nil {\n\t\t\t\treturn messages, fmt.Errorf(\"stream completion: %w\", err)\n\t\t\t}\n\n\t\t\tcompletion = c\n\n\t\t\tif len(c.Choices) \u003e 0 {\n\t\t\t\tmsg := c.Choices[0].Message\n\n\t\t\t\tif hooks.OnReasoningDelta != nil \u0026\u0026 len(msg.ReasoningContent) \u003e len(prevReasoning) {\n\t\t\t\t\thooks.OnReasoningDelta(msg.ReasoningContent[len(prevReasoning):])\n\t\t\t\t}\n\t\t\t\tprevReasoning = msg.ReasoningContent\n\n\t\t\t\tif hooks.OnContentDelta != nil \u0026\u0026 len(msg.Content) \u003e len(prevContent) {\n\t\t\t\t\thooks.OnContentDelta(msg.Content[len(prevContent):])\n\t\t\t\t}\n\t\t\t\tprevContent = msg.Content\n\t\t\t}\n\t\t}\n\n\t\tif hooks.OnTurnEnd != nil {\n\t\t\thooks.OnTurnEnd()\n\t\t}\n\n\t\tif completion == nil || len(completion.Choices) == 0 {\n\t\t\treturn messages, fmt.Errorf(\"model returned an empty completion\")\n\t\t}\n\n\t\tchoice := completion.Choices[0]\n\t\tmessages = append(messages, choice.Message)\n\n\t\tif choice.FinishReason != \"tool_calls\" || len(choice.Message.ToolCalls) == 0 {\n\t\t\treturn messages, nil\n\t\t}\n\n\t\tfor _, tc := range choice.Message.ToolCalls {\n\t\t\tif err := ctx.Err(); err != nil {\n\t\t\t\treturn messages, err\n\t\t\t}\n\n\t\t\tresult, err := tools.Call(toolsCtx, tc.FunctionCall)\n\n\t\t\ttoolErr := err != nil\n\t\t\tif toolErr {\n\t\t\t\tresult = fmt.Sprintf(\"ERROR CALLING TOOL: %s\", err.Error())\n\t\t\t}\n\n\t\t\tif hooks.OnTool != nil {\n\t\t\t\thooks.OnTool(tc.FunctionCall.Name, tc.FunctionCall.Arguments, result, toolErr)\n\t\t\t}\n\n\t\t\tmessages = append(messages, llm.Message{\n\t\t\t\tRole: \"tool\",\n\t\t\t\tContent: result,\n\t\t\t\tToolCallID: tc.Id,\n\t\t\t\tName: tc.FunctionCall.Name,\n\t\t\t})\n\t\t}\n\t}\n}","start_line":1,"end_line":118,"total_lines":118,"truncated":false}
Thinking
Next, let's look at the streamer to understand how errors surface, and check the config to see what's available.
Tool
π§ read_file
{"path":"internal/llm/streamer.go","start_line":1,"end_line":200}
Arguments
{"path":"internal/llm/streamer.go","start_line":1,"end_line":200}
Result
{"path":"internal/llm/streamer.go","content":"package llm\n\nimport (\n\t\"bufio\"\n\t\"bytes\"\n\t\"encoding/json\"\n\t\"fmt\"\n\t\"io\"\n\t\"net/http\"\n)\n\ntype Streamer struct {\n\tcompletion ChatCompletion\n\tresp *http.Response\n\tscanner *bufio.Scanner\n\teof bool\n}\n\nfunc NewStreamer(resp *http.Response) (*Streamer, error) {\n\tif resp.StatusCode != http.StatusOK {\n\t\tbody, _ := io.ReadAll(resp.Body)\n\t\tresp.Body.Close()\n\t\treturn nil, fmt.Errorf(\"status: %d: %s\", resp.StatusCode, body)\n\t}\n\n\treturn \u0026Streamer{\n\t\tresp: resp,\n\t\tscanner: bufio.NewScanner(resp.Body),\n\t}, nil\n}\n\n// Get returns the accumulated completion after merging in the next SSE\n// chunk. Callers should keep calling Get until it returns io.EOF.\nfunc (s *Streamer) Get() (*ChatCompletion, error) {\n\tif s.eof {\n\t\treturn nil, io.EOF\n\t}\n\n\tprefix := []byte(\"data: \")\n\n\tvar buf []byte\n\n\tfor s.scanner.Scan() {\n\t\tif err := s.scanner.Err(); err != nil {\n\t\t\treturn nil, err\n\t\t}\n\n\t\tbuf = bytes.TrimSpace(s.scanner.Bytes())\n\n\t\tif len(buf) == 0 {\n\t\t\tcontinue\n\t\t}\n\n\t\tif bytes.HasPrefix(buf, prefix) {\n\t\t\tbreak\n\t\t}\n\t}\n\n\tbuf = bytes.TrimPrefix(buf, prefix)\n\n\tif bytes.Equal(buf, []byte(\"[DONE]\")) {\n\t\ts.resp.Body.Close()\n\t\ts.eof = true\n\t\treturn nil, io.EOF\n\t}\n\n\tvar completion ChatCompletion\n\n\t// Unknown fields are tolerated on purpose: servers (e.g. newer\n\t// llama.cpp builds) may add fields we don't know about yet, and\n\t// failing on them would break compatibility.\n\tif err := json.Unmarshal(buf, \u0026completion); err != nil {\n\t\treturn nil, fmt.Errorf(\"decode completion chunk: %w: %s\", err, buf)\n\t}\n\n\ts.completion.Merge(\u0026completion)\n\n\treturn \u0026s.completion, nil\n}","start_line":1,"end_line":79,"total_lines":79,"truncated":false}
Tool
π§ read_file
{"path":"internal/config/config.go","start_line":1,"end_line":250}
Arguments
{"path":"internal/config/config.go","start_line":1,"end_line":250}
Result
{"path":"internal/config/config.go","content":"// Package config loads and validates zoo's HCL configuration file.\npackage config\n\nimport (\n\t\"fmt\"\n\t\"strings\"\n\n\t\"github.com/hashicorp/hcl/v2/hclsimple\"\n)\n\n// Known event kinds. issue:assigned and pr:review are resolved\n// dynamically (agent name must match the Forgejo assignee's username,\n// or the pull request author's username, respectively) so they never\n// carry an `agent` attribute; the rest map statically to one\n// configured agent.\nconst (\n\tEventIssueNew = \"issue:new\"\n\tEventIssueComment = \"issue:comment\"\n\tEventIssueAssigned = \"issue:assigned\"\n\tEventPRNew = \"pr:new\"\n\tEventPRReview = \"pr:review\"\n)\n\nvar staticEventKinds = map[string]bool{\n\tEventIssueNew: true,\n\tEventIssueComment: true,\n\tEventPRNew: true,\n}\n\ntype Config struct {\n\tLLMs []LLM `hcl:\"llm,block\"`\n\tForgejo Forgejo `hcl:\"forgejo,block\"`\n\tEnvironment Environment `hcl:\"environment,block\"`\n\tAgents []Agent `hcl:\"agent,block\"`\n\tEvents []Event `hcl:\"event,block\"`\n\tWeb *Web `hcl:\"web,block\"`\n}\n\n// Web configures the dashboard's optional bearer-token gate. Leave the\n// block out of zoo.hcl entirely to run without one (fine on localhost;\n// put a real gate or a proxy in front for anything else).\ntype Web struct {\n\tToken string `hcl:\"token,optional\"`\n}\n\ntype LLM struct {\n\tName string `hcl:\"name,label\"`\n\tOpenAI string `hcl:\"openai\"`\n\tToken string `hcl:\"token\"`\n\tModel string `hcl:\"model\"`\n}\n\ntype Forgejo struct {\n\tURL string `hcl:\"url\"`\n\tToken string `hcl:\"token\"`\n\tWebhookSecret string `hcl:\"webhook_secret,optional\"`\n\n\t// Repos is the allowlist of repository patterns to watch, e.g.\n\t// [\"acme/*\", \"acme/widgets\"]. Patterns are \"owner/repo\" pairs with\n\t// glob wildcards; \"*\" watches everything on the instance. An empty\n\t// list keeps the historical behavior of watching every repository\n\t// the token can see.\n\tRepos []string `hcl:\"repos,optional\"`\n}\n\ntype Environment struct {\n\tDockerImage string `hcl:\"docker_image\"`\n\tMaxLive int `hcl:\"max_live_agents\"`\n}\n\ntype Agent struct {\n\tName string `hcl:\"name,label\"`\n\tLLM string `hcl:\"llm\"`\n\tToken string `hcl:\"token,optional\"`\n}\n\ntype Event struct {\n\tKind string `hcl:\"name,label\"`\n\tAgent string `hcl:\"agent,optional\"`\n\tInstructions string `hcl:\"instructions,optional\"`\n}\n\n// Load reads and validates the config file at path.\nfunc Load(path string) (*Config, error) {\n\tvar cfg Config\n\n\tif err := hclsimple.DecodeFile(path, nil, \u0026cfg); err != nil {\n\t\treturn nil, fmt.Errorf(\"parse config: %w\", err)\n\t}\n\n\tif err := cfg.Validate(); err != nil {\n\t\treturn nil, fmt.Errorf(\"invalid config: %w\", err)\n\t}\n\n\treturn \u0026cfg, nil\n}\n\n// Validate checks that the config is internally consistent: every\n// reference between blocks resolves, and required values are set.\nfunc (c *Config) Validate() error {\n\tllmNames := make(map[string]bool, len(c.LLMs))\n\tfor _, l := range c.LLMs {\n\t\tif l.OpenAI == \"\" || l.Token == \"\" || l.Model == \"\" {\n\t\t\treturn fmt.Errorf(\"llm %q: openai, token, and model are required\", l.Name)\n\t\t}\n\t\tllmNames[l.Name] = true\n\t}\n\n\tif c.Forgejo.URL == \"\" || c.Forgejo.Token == \"\" {\n\t\treturn fmt.Errorf(\"forgejo: url and token are required\")\n\t}\n\n\tfor _, p := range c.Forgejo.Repos {\n\t\tif err := validRepoPattern(p); err != nil {\n\t\t\treturn fmt.Errorf(\"forgejo: %w\", err)\n\t\t}\n\t}\n\n\tif c.Environment.MaxLive \u003c 1 {\n\t\treturn fmt.Errorf(\"environment: max_live_agents must be \u003e= 1, got %d\", c.Environment.MaxLive)\n\t}\n\n\tif c.Environment.DockerImage == \"\" {\n\t\treturn fmt.Errorf(\"environment: docker_image is required\")\n\t}\n\n\tagentNames := make(map[string]bool, len(c.Agents))\n\tfor _, a := range c.Agents {\n\t\tif !llmNames[a.LLM] {\n\t\t\treturn fmt.Errorf(\"agent %q: references undeclared llm %q\", a.Name, a.LLM)\n\t\t}\n\t\tagentNames[a.Name] = true\n\t}\n\n\tseenEventKinds := make(map[string]bool, len(c.Events))\n\tfor _, e := range c.Events {\n\t\tif seenEventKinds[e.Kind] {\n\t\t\treturn fmt.Errorf(\"event %q: declared more than once\", e.Kind)\n\t\t}\n\t\tseenEventKinds[e.Kind] = true\n\n\t\tswitch {\n\t\tcase e.Kind == EventIssueAssigned:\n\t\t\tif e.Agent != \"\" {\n\t\t\t\treturn fmt.Errorf(\"event %q: must not set agent; the agent whose name matches the Forgejo assignee's username is triggered dynamically\", e.Kind)\n\t\t\t}\n\n\t\tcase e.Kind == EventPRReview:\n\t\t\tif e.Agent != \"\" {\n\t\t\t\treturn fmt.Errorf(\"event %q: must not set agent; the agent whose name matches the pull request author's username is triggered dynamically\", e.Kind)\n\t\t\t}\n\n\t\tcase staticEventKinds[e.Kind]:\n\t\t\tif e.Agent == \"\" {\n\t\t\t\treturn fmt.Errorf(\"event %q: agent is required\", e.Kind)\n\t\t\t}\n\t\t\tif !agentNames[e.Agent] {\n\t\t\t\treturn fmt.Errorf(\"event %q: references undeclared agent %q\", e.Kind, e.Agent)\n\t\t\t}\n\n\t\tdefault:\n\t\t\treturn fmt.Errorf(\"event %q: unknown event kind\", e.Kind)\n\t\t}\n\t}\n\n\treturn nil\n}\n\n// validRepoPattern checks the shape of a forgejo `repos` entry. It must\n// be a non-empty \"owner/repo\" pair (or a bare owner name, or \"*\") β at\n// most one slash, none at either end. Glob syntax within the segments\n// is validated later by forgejo.NewRepoFilter.\nfunc validRepoPattern(p string) error {\n\tif p == \"\" {\n\t\treturn fmt.Errorf(\"repos: pattern must not be empty\")\n\t}\n\n\tif strings.Count(p, \"/\") \u003e 1 {\n\t\treturn fmt.Errorf(\"repos: pattern %q must be \\\"owner/repo\\\" (at most one \\\"/\\\")\", p)\n\t}\n\n\tif strings.HasPrefix(p, \"/\") || strings.HasSuffix(p, \"/\") {\n\t\treturn fmt.Errorf(\"repos: pattern %q must not start or end with \\\"/\\\"\", p)\n\t}\n\n\treturn nil\n}\n\n// AgentConfig holds an agent's full configuration including its\n// Forgejo token (if configured).\ntype AgentConfig struct {\n\tName string\n\tLLM string\n\tToken string\n}\n\n// AgentByName returns the agent config with the given name, if declared.\nfunc (c *Config) AgentByName(name string) (AgentConfig, bool) {\n\tfor _, a := range c.Agents {\n\t\tif a.Name == name {\n\t\t\treturn AgentConfig{\n\t\t\t\tName: a.Name,\n\t\t\t\tLLM: a.LLM,\n\t\t\t\tToken: a.Token,\n\t\t\t}, true\n\t\t}\n\t}\n\n\treturn AgentConfig{}, false\n}\n\n// LLMByName returns the llm config with the given name, if declared.\nfunc (c *Config) LLMByName(name string) (LLM, bool) {\n\tfor _, l := range c.LLMs {\n\t\tif l.Name == name {\n\t\t\treturn l, true\n\t\t}\n\t}\n\n\treturn LLM{}, false\n}\n\n// EventAgent returns the statically configured agent name for the given\n// event kind, if one is declared and set.\nfunc (c *Config) EventAgent(kind string) (string, bool) {\n\tfor _, e := range c.Events {\n\t\tif e.Kind == kind {\n\t\t\treturn e.Agent, e.Agent != \"\"\n\t\t}\n\t}\n\n\treturn \"\", false\n}\n\n// EventInstructions returns the configured freeform instructions for the\n// given event kind, e.g. \"Please review this pull request and respond\n// appropriately.\" Returns \"\" if the event isn't declared or has none.\nfunc (c *Config) EventInstructions(kind string) string {\n\tfor _, e := range c.Events {\n\t\tif e.Kind == kind {\n\t\t\treturn e.Instructions\n\t\t}\n\t}\n\n\treturn \"\"\n}","start_line":1,"end_line":246,"total_lines":246,"truncated":false}
Thinking
Next, let's look at run.go, types.go, and streamer_test.go to understand the overall flow and how errors propagate.
Tool
π§ read_file
{"path":"internal/agentrun/run.go","start_line":1,"end_line":300}
Arguments
{"path":"internal/agentrun/run.go","start_line":1,"end_line":300}
Result
{"path":"internal/agentrun/run.go","content":"// Package agentrun creates the throwaway per-job environment an agent\n// runs in: a Docker container with the agent's git working tree cloned\n// inside it (authenticated with the agent's own Forgejo token), the\n// triggering event written to /event, and the tool-calling loop\n// (internal/llm + internal/tools) driven against it.\npackage agentrun\n\nimport (\n\t\"context\"\n\t_ \"embed\"\n\t\"encoding/json\"\n\t\"fmt\"\n\t\"log/slog\"\n\t\"os\"\n\t\"path/filepath\"\n\t\"strings\"\n\t\"sync\"\n\t\"time\"\n\n\t\"github.com/abrander/zoo/internal/config\"\n\t\"github.com/abrander/zoo/internal/forgejo\"\n\t\"github.com/abrander/zoo/internal/livelog\"\n\t\"github.com/abrander/zoo/internal/llm\"\n\t\"github.com/abrander/zoo/internal/store\"\n)\n\n//go:embed system.md\nvar defaultSystemPrompt string\n\n// DefaultTimeout bounds a single agent run's wall-clock time if the\n// caller doesn't override it.\nconst DefaultTimeout = 120 * time.Minute\n\ntype Runner struct {\n\tdocker *dockerRuntime\n\tforgejo *forgejo.Client\n\tstore *store.Store\n\thub *livelog.Hub\n\tcfg *config.Config\n\tlogger *slog.Logger\n\ttimeout time.Duration\n\tkeepOnFailure bool\n\n\tagentClientsMu sync.Mutex\n\tagentClients map[string]*forgejo.Client\n}\n\nfunc NewRunner(cfg *config.Config, fg *forgejo.Client, st *store.Store, hub *livelog.Hub, logger *slog.Logger, timeout time.Duration, keepOnFailure bool) (*Runner, error) {\n\tdocker, err := newDockerRuntime()\n\tif err != nil {\n\t\treturn nil, err\n\t}\n\n\tif timeout \u003c= 0 {\n\t\ttimeout = DefaultTimeout\n\t}\n\n\treturn \u0026Runner{\n\t\tdocker: docker,\n\t\tforgejo: fg,\n\t\tstore: st,\n\t\thub: hub,\n\t\tcfg: cfg,\n\t\tlogger: logger,\n\t\ttimeout: timeout,\n\t\tkeepOnFailure: keepOnFailure,\n\t\tagentClients: make(map[string]*forgejo.Client),\n\t}, nil\n}\n\n// forgejoAs returns a Forgejo client that authenticates as the given\n// agent (using the agent's own token from config). This lets each agent\n// act as themselves on Forgejo without needing a global token with sudo\n// privileges. Clients are built once per agent and cached, since\n// constructing one costs an extra API round trip.\n//\n// If the agent has no token configured, falls back to the shared zoo\n// identity so existing deployments without per-agent tokens still work.\nfunc (r *Runner) forgejoAs(agentName, token string) *forgejo.Client {\n\tr.agentClientsMu.Lock()\n\tdefer r.agentClientsMu.Unlock()\n\n\tif c, ok := r.agentClients[agentName]; ok {\n\t\treturn c\n\t}\n\n\tvar c *forgejo.Client\n\tif token != \"\" {\n\t\tc = r.forgejo.As(token)\n\t} else {\n\t\t// Fallback: use shared identity. Optionally log a warning\n\t\t// if we ever want to enforce per-agent tokens.\n\t\tc = r.forgejo\n\t}\n\n\tr.agentClients[agentName] = c\n\n\treturn c\n}\n\n// Run implements scheduler.Runner.\nfunc (r *Runner) Run(ctx context.Context, jobID string, agent config.AgentConfig, llmCfg config.LLM, dockerImage string, ev forgejo.Event) error {\n\tctx, cancel := context.WithTimeout(ctx, r.timeout)\n\tdefer cancel()\n\n\tlogger := r.logger.With(\"job\", jobID, \"agent\", agent.Name)\n\n\trepoInfo, err := r.forgejo.RepositoryInfo(ev.Owner, ev.Repo)\n\tif err != nil {\n\t\treturn fmt.Errorf(\"look up repository: %w\", err)\n\t}\n\n\tworkDir, err := os.MkdirTemp(\"\", \"zoo-run-*\")\n\tif err != nil {\n\t\treturn fmt.Errorf(\"create work dir: %w\", err)\n\t}\n\n\tsucceeded := false\n\n\tdefer func() {\n\t\tif succeeded || !r.keepOnFailure {\n\t\t\tos.RemoveAll(workDir)\n\t\t} else {\n\t\t\tlogger.Warn(\"keeping work dir after failure\", \"dir\", workDir)\n\t\t}\n\t}()\n\n\t// The container bind-mounts projectDir as /project and does the\n\t// initial clone into it, so the (empty) directory must exist on the\n\t// host before the container is created β otherwise Docker would\n\t// create it itself, root-owned.\n\tprojectDir := filepath.Join(workDir, \"project\")\n\n\tif err := os.MkdirAll(projectDir, 0o755); err != nil {\n\t\treturn fmt.Errorf(\"create project dir: %w\", err)\n\t}\n\n\t// A pr:review run works on the PR's own head branch, so the agent's\n\t// commits push straight to the PR. Every other event kind branches\n\t// off the default branch as usual.\n\tvar review *forgejo.ReviewDetail\n\tbranch := fmt.Sprintf(\"issue-%d-%s\", ev.Index, agent.Name)\n\n\tif ev.Kind == forgejo.EventPRReview {\n\t\t// Always fetch the current head ref, not just when the event\n\t\t// lacks one (the polling path doesn't carry it): the webhook's\n\t\t// copy could be stale if the PR's head branch was renamed since\n\t\t// the review, and the push target depends on it.\n\t\theadRef := ev.HeadRef\n\n\t\tif prInfo, err := r.forgejo.PullRequestInfo(ev.Owner, ev.Repo, ev.Index); err != nil {\n\t\t\tlogger.Warn(\"fetch pull request head failed; falling back to the event's head ref\", \"error\", err)\n\t\t} else if prInfo.HeadRef != \"\" {\n\t\t\theadRef = prInfo.HeadRef\n\t\t}\n\n\t\tif headRef == \"\" {\n\t\t\treturn fmt.Errorf(\"pr:review event has no pull request head branch to check out\")\n\t\t}\n\n\t\tbranch = headRef\n\n\t\t// Fetch the full review (verdict, body, inline comments) so the\n\t\t// agent sees all the feedback, not just the triggering event. A\n\t\t// failure degrades to no review detail rather than failing the\n\t\t// run: the agent can still do its job, just without the inline\n\t\t// comments.\n\t\treview, err = r.forgejo.ReviewDetail(ev.Owner, ev.Repo, ev.Index, ev.ReviewID)\n\t\tif err != nil {\n\t\t\tlogger.Warn(\"fetch review detail failed; agent will not see inline review comments\", \"error\", err)\n\t\t\treview = nil\n\t\t}\n\t}\n\n\troster := buildRoster(r.forgejo, r.cfg.Agents, logger)\n\tgitName, gitEmail := gitIdentity(agent.Name, roster)\n\n\t// The credential the sandbox's git uses for remote operations: the\n\t// agent's own Forgejo token when configured, so its git activity is\n\t// attributed to its own account, falling back to the shared zoo\n\t// identity for deployments without per-agent tokens (mirroring\n\t// forgejoAs).\n\tgitUser, gitToken := \"zoo\", r.forgejo.Token()\n\n\tif agent.Token != \"\" {\n\t\tgitUser, gitToken = agent.Name, agent.Token\n\t}\n\n\teventPath := filepath.Join(workDir, \"event.json\")\n\tif err := os.WriteFile(eventPath, ev.Raw, 0o644); err != nil {\n\t\treturn fmt.Errorf(\"write event file: %w\", err)\n\t}\n\n\tcontainerID, err := r.docker.createContainer(ctx, dockerImage, []string{\n\t\tprojectDir + \":/project\",\n\t\teventPath + \":/event:ro\",\n\t}, fmt.Sprintf(\"zoo-issue-%d-%s\", ev.Index, agent.Name))\n\tif err != nil {\n\t\treturn fmt.Errorf(\"start container: %w\", err)\n\t}\n\n\tdefer func() {\n\t\tcleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 30*time.Second)\n\t\tdefer cleanupCancel()\n\t\tif err := r.docker.remove(cleanupCtx, containerID); err != nil {\n\t\t\tlogger.Warn(\"failed to remove container\", \"container\", containerID, \"error\", err)\n\t\t}\n\t}()\n\n\t// Git must simply work inside the sandbox: safe.directory, commit\n\t// identity, and the remote credential all go into the container's\n\t// system gitconfig (see configureSandboxGit).\n\tif err := configureSandboxGit(ctx, r.docker, containerID, repoInfo.CloneURL, gitUser, gitToken, gitName, gitEmail); err != nil {\n\t\treturn fmt.Errorf(\"configure git in container: %w\", err)\n\t}\n\n\t// The initial clone happens inside the sandbox, so the working tree\n\t// is owned by the container's user and git never runs on the host.\n\tif ev.Kind == forgejo.EventPRReview {\n\t\tif err := clonePRHead(ctx, r.docker, containerID, repoInfo.CloneURL, repoInfo.DefaultBranch, branch, ev.Index); err != nil {\n\t\t\treturn fmt.Errorf(\"prepare git working tree: %w\", err)\n\t\t}\n\t} else {\n\t\tif err := cloneAndBranch(ctx, r.docker, containerID, repoInfo.CloneURL, repoInfo.DefaultBranch, branch); err != nil {\n\t\t\treturn fmt.Errorf(\"prepare git working tree: %w\", err)\n\t\t}\n\t}\n\n\tlogAppend := func(stream, line string) {\n\t\tif err := r.store.AppendLog(context.Background(), jobID, stream, line); err != nil {\n\t\t\tlogger.Warn(\"failed to append log\", \"error\", err)\n\t\t}\n\t}\n\n\trunCtx := \u0026runContext{\n\t\tdocker: r.docker,\n\t\tcontainerID: containerID,\n\t\tforgejo: \u0026runForgejoActions{\n\t\t\tclient: r.forgejoAs(agent.Name, agent.Token),\n\t\t\towner: ev.Owner,\n\t\t\trepo: ev.Repo,\n\t\t\tindex: ev.Index,\n\t\t\tlogger: logger,\n\t\t},\n\t}\n\n\tllmClient := llm.NewClient(llmCfg)\n\n\tsystemPrompt := defaultSystemPrompt + identitySection(agent.Name, roster)\n\n\tinstructions := r.cfg.EventInstructions(ev.Kind)\n\n\t// Fetch the full comment thread so the agent sees everything that's\n\t// been said on the issue/PR, not just the triggering event (which\n\t// only carries the latest comment, if any). A failure degrades to\n\t// no comments rather than failing the run: the agent can still do\n\t// its job, just without prior context.\n\tcomments, err := r.forgejo.ListIssueComments(ev.Owner, ev.Repo, ev.Index)\n\tif err != nil {\n\t\tlogger.Warn(\"fetch issue comments failed; agent will not see prior comments\", \"error\", err)\n\t\tcomments = nil\n\t}\n\n\tmessages := []llm.Message{\n\t\t{Role: \"system\", Content: systemPrompt},\n\t\t{Role: \"user\", Content: seedMessage(ev, branch, repoInfo.DefaultBranch, instructions, comments, review)},\n\t}\n\n\thooks := r.streamHooks(jobID, logAppend)\n\n\t_, err = runLoop(ctx, llmClient, runCtx, messages, hooks)\n\tif err != nil {\n\t\treturn fmt.Errorf(\"agent loop: %w\", err)\n\t}\n\n\tsucceeded = true\n\n\treturn nil\n}\n\n// streamHooks builds the Hooks a single Run passes to runLoop: every\n// delta is published live to the hub for connected dashboard viewers,\n// and once a reasoning/content block or tool call is complete, it's\n// persisted to the store as one row and the hub's replay buffer for\n// jobID is checkpointed β so a viewer connecting from this point on\n// sees it via the persisted history instead of a live replay, and is\n// never shown it twice.\nfunc (r *Runner) streamHooks(jobID string, logAppend func(stream, line string)) Hooks {\n\tvar reasoningBuf, contentBuf strings.Builder\n\n\treasoningOpen, contentOpen := false, false\n\n\treturn Hooks{\n\t\tOnReasoningDelta: func(delta string) {\n\t\t\tif !reasoningOpen {\n\t\t\t\treasoningOpen = true\n\t\t\t\treasoningBuf.Reset()\n\t\t\t\tr.hub.Publish(jobID, livelog.Event{Type: livelog.ReasoningStart})\n\t\t\t}\n","start_line":1,"end_line":300,"total_lines":437,"truncated":true}
Tool
π§ read_file
{"path":"internal/llm/types.go","start_line":1,"end_line":200}
Arguments
{"path":"internal/llm/types.go","start_line":1,"end_line":200}
Result
{"path":"internal/llm/types.go","content":"// Package llm is an OpenAI-compatible chat completions client with\n// streaming and tool-call support, ported from ../a's llamacpp package.\npackage llm\n\nimport \"fmt\"\n\ntype Message struct {\n\tContent string `json:\"content\"`\n\tRole string `json:\"role\"`\n\tReasoningContent string `json:\"reasoning_content,omitempty\"`\n\tToolCalls []ToolCall `json:\"tool_calls,omitempty\"`\n\tToolCallID string `json:\"tool_call_id,omitempty\"`\n\tName string `json:\"name,omitempty\"`\n}\n\nfunc (m *Message) String() string {\n\tstr := fmt.Sprintf(\"Role: %s, Content: %s ReasoningContent: %s\", m.Role, m.Content, m.ReasoningContent)\n\n\tfor _, call := range m.ToolCalls {\n\t\tstr += fmt.Sprintf(\"\\n ToolCall: %s\", call.String())\n\t}\n\n\treturn str\n}\n\n// Merge folds a streamed delta into the accumulated message.\nfunc (m *Message) Merge(next *Message) {\n\tm.Content += next.Content\n\n\tif next.Role != \"\" {\n\t\tm.Role = next.Role\n\t}\n\n\tm.ReasoningContent += next.ReasoningContent\n\n\tminLength := 0\n\tfor _, choice := range next.ToolCalls {\n\t\tif choice.Index+1 \u003e minLength {\n\t\t\tminLength = choice.Index + 1\n\t\t}\n\t}\n\n\tif len(m.ToolCalls) \u003c minLength {\n\t\tnewToolCalls := make([]ToolCall, minLength)\n\t\tcopy(newToolCalls, m.ToolCalls)\n\t\tm.ToolCalls = newToolCalls\n\t}\n\n\tfor _, nextChoice := range next.ToolCalls {\n\t\tm.ToolCalls[nextChoice.Index].Merge(\u0026nextChoice)\n\t}\n}\n\ntype FunctionCall struct {\n\tName string `json:\"name\"`\n\tArguments string `json:\"arguments\"`\n}\n\ntype ToolCall struct {\n\tType string `json:\"type\"`\n\tId string `json:\"id\"`\n\tIndex int `json:\"index\"`\n\tName string `json:\"name\"`\n\tFunctionCall FunctionCall `json:\"function\"`\n}\n\nfunc (t *ToolCall) String() string {\n\treturn fmt.Sprintf(\"Type: %s, Id: %s, Index: %d, Name: %s, FunctionCall: {Name: %s, Arguments: %s}\",\n\t\tt.Type, t.Id, t.Index, t.Name, t.FunctionCall.Name, t.FunctionCall.Arguments)\n}\n\nfunc (t *ToolCall) Merge(next *ToolCall) {\n\tif next.Type != \"\" {\n\t\tt.Type = next.Type\n\t}\n\n\tif next.Id != \"\" {\n\t\tt.Id = next.Id\n\t}\n\n\tif next.Name != \"\" {\n\t\tt.Name = next.Name\n\t}\n\n\tt.Index = next.Index\n\n\tif next.FunctionCall.Name != \"\" {\n\t\tt.FunctionCall.Name = next.FunctionCall.Name\n\t}\n\n\tt.FunctionCall.Arguments += next.FunctionCall.Arguments\n}\n\ntype Property struct {\n\tType string `json:\"type\"`\n\tDescription string `json:\"description,omitempty\"`\n\tEnum []string `json:\"enum,omitempty\"`\n\tItems *Property `json:\"items,omitempty\"`\n}\n\ntype Parameters struct {\n\tType string `json:\"type\"` // object\n\tProperties map[string]Property `json:\"properties,omitempty\"`\n\tRequired []string `json:\"required,omitempty\"`\n}\n\ntype Function struct {\n\tName string `json:\"name\"`\n\tDescription string `json:\"description,omitempty\"`\n\tParameters Parameters `json:\"parameters,omitempty\"`\n\tStrict bool `json:\"strict\"`\n}\n\ntype Tool struct {\n\tType string `json:\"type\"`\n\tFunction Function `json:\"function\"`\n}\n\nfunc NewTool(name, description string) Tool {\n\treturn Tool{\n\t\tType: \"function\",\n\t\tFunction: Function{\n\t\t\tName: name,\n\t\t\tDescription: description,\n\t\t\tStrict: true,\n\t\t},\n\t}\n}\n\nfunc (t *Tool) AddProperty(name string, property Property, required bool) *Tool {\n\tif t.Function.Parameters.Properties == nil {\n\t\tt.Function.Parameters.Properties = make(map[string]Property)\n\t} else if _, exists := t.Function.Parameters.Properties[name]; exists {\n\t\tpanic(\"property already exists: \" + name)\n\t}\n\n\tt.Function.Parameters.Properties[name] = property\n\n\tif required {\n\t\tt.Function.Parameters.Required = append(t.Function.Parameters.Required, name)\n\t}\n\n\treturn t\n}\n\nfunc (t *Tool) AddStringProperty(name, description string, required bool) *Tool {\n\treturn t.AddProperty(name, Property{Type: \"string\", Description: description}, required)\n}\n\nfunc (t *Tool) AddBooleanProperty(name, description string, required bool) *Tool {\n\treturn t.AddProperty(name, Property{Type: \"boolean\", Description: description}, required)\n}\n\nfunc (t *Tool) AddIntegerProperty(name, description string, required bool) *Tool {\n\treturn t.AddProperty(name, Property{Type: \"integer\", Description: description}, required)\n}\n\nfunc (t *Tool) AddEnumProperty(name, description string, values []string, required bool) *Tool {\n\treturn t.AddProperty(name, Property{Type: \"string\", Description: description, Enum: values}, required)\n}\n\nfunc (t *Tool) AddArrayProperty(name, itemType, description string, required bool) *Tool {\n\treturn t.AddProperty(name, Property{\n\t\tType: \"array\",\n\t\tDescription: description,\n\t\tItems: \u0026Property{Type: itemType},\n\t}, required)\n}\n\nfunc (t *Tool) AddStringArrayProperty(name, description string, required bool) *Tool {\n\treturn t.AddArrayProperty(name, \"string\", description, required)\n}\n\ntype ChatCompletionRequest struct {\n\tModel string `json:\"model\"`\n\tMessages []Message `json:\"messages,omitempty\"`\n\tStream bool `json:\"stream\"`\n\tTools []Tool `json:\"tools,omitempty\"`\n\tSeed int `json:\"seed,omitempty\"`\n\tTemperature float64 `json:\"temperature,omitempty\"`\n\tTopP float64 `json:\"top_p,omitempty\"`\n}\n\ntype Choice struct {\n\tDelta Message `json:\"delta\"`\n\tFinishReason string `json:\"finish_reason,omitempty\"`\n\tIndex int `json:\"index\"`\n\tMessage Message `json:\"message\"`\n\tRole string `json:\"role\"`\n}\n\nfunc (c *Choice) String() string {\n\treturn fmt.Sprintf(\"Index: %d, Role: %s, FinishReason: %v, Delta: %s, Message: %s\",\n\t\tc.Index, c.Role, c.FinishReason, c.Delta.String(), c.Message.String())\n}\n\nfunc (c *Choice) Merge(next *Choice) {\n\tc.Message.Merge(\u0026next.Delta)\n\n\tif next.FinishReason != \"\" {","start_line":1,"end_line":200,"total_lines":267,"truncated":true}
Thinking
Now I understand the structure. Let me also check the remaining parts of run.go and the streamer tests to understand the test conventions.
Key observation points:
1. `runLoop` in `internal/agentrun/loop.go` calls `client.StreamChatCompletion(ctx, ...)` and returns an error if it fails. The error propagates as "agent loop: chat completion: ..." β exactly as seen in the issue's comments.
2. The LLM client (`internal/llm/client.go`) does a single HTTP POST. No retries.
3. The streamer reads SSE chunks. If the stream fails mid-way (e.g., connection drops), `streamer.Get()` returns an error and the run fails.
Issue's request: "We need some kind of retry logic in the agent loop in case the LLM endpoint is down."
Errors observed in the comments:
- 429: rate limit / model cooldown ("Try again in 152.0 seconds")
- 502: engine not healthy
So the retry logic should:
- Retry on transient errors: 429, 500, 502, 503, 504, and connection-level errors (EOF, connection reset, timeouts).
- Do not retry on 4xx (400, 401, 403, 404) β these are permanent client errors.
- Use exponential backoff with jitter.
- Respect the `Retry-After` header if present (litellm's 429 says "Try again in 152.0 seconds" β that's in the body, not the header. But standard HTTP has Retry-After).
- Bound the total retry time by the run's context deadline.
- Retry mid-stream failures too? This is trickier: if the stream fails partway through, we've already consumed some deltas and published them to the live log. Retrying the whole turn would re-stream content that's already been published. Hmm. But for the agent loop's correctness, a mid-stream failure means the completion is incomplete, and retrying the entire chat completion request (with the same messages) is safe in the sense that no state has been committed to the messages list yet β the assistant message is only appended after the stream completes. So retrying the request from scratch is semantically safe with respect to the message history. The only side effect is that partial deltas have been published to the dashboard's live log. That's a cosmetic issue.
Let me think about where to put the retry:
Option A: Inside `llm.Client.StreamChatCompletion` β retry the HTTP request there. But mid-stream failures happen in `streamer.Get()`, which is after StreamChatCompletion has returned. So Option A only covers request-level failures (status codes, connection errors).
Option B: Inside `runLoop` β wrap the entire "make request + consume stream" in a retry loop. This covers both request-level and mid-stream failures. The issue says "retry logic in the agent loop," which points to the loop.
The cleanest architecture:
1. Add a retry helper in the `llm` package (or a small `retry` utility) that wraps a function with backoff.
2. In `runLoop`, wrap the request + stream consumption in a retry function that:
- Classifies errors: retryable (429, 5xx, connection errors, mid-stream EOF) vs. permanent (4xx).
- Exponential backoff with jitter, capped.
- Respects Retry-After header if present.
- Bounded by ctx.
Wait, but there's a subtlety with mid-stream retries and the hooks: OnReasoningDelta/OnContentDelta have already fired for the partial content. If we retry, the same content will be re-published. In the dashboard's live log, that would look like duplicated text. But correctness-wise it's fine. Given the alternative (the whole run failing), it's a reasonable trade-off. We could log a warning when retrying.
Actually, let me reconsider. Where does the error from the stream surface? In `runLoop`:
```go
streamer, err := client.StreamChatCompletion(ctx, &llm.ChatCompletionRequest{...})
if err != nil {
return messages, fmt.Errorf("chat completion: %w", err)
}
...
for {
c, err := streamer.Get()
if err == io.EOF {
break
}
if err != nil {
return messages, fmt.Errorf("stream completion: %w", err)
}
...
}
```
The issue's error messages: "agent loop: chat completion: status: 429: ..." and "agent loop: chat completion: status: 502: engine not healthy". These come from `NewStreamer` (non-200 status). So the primary failures are request-level.
Design decision: implement retry at the `runLoop` level, covering both request and stream errors. This matches the issue's ask ("in the agent loop") and covers the most ground.
Now, let me think about error classification. The current errors are `fmt.Errorf("status: %d: %s", ...)` β not typed. To classify properly, I should introduce a typed error in the llm package:
```go
// APIError represents an HTTP error response from the LLM endpoint.
type APIError struct {
StatusCode int
Status string
Body string
}
func (e *APIError) Error() string {
return fmt.Sprintf("status: %d: %s", e.StatusCode, e.Body)
}
```
And a classification function:
```go
// IsRetryable reports whether err is a transient failure worth
// retrying: rate limits (429), server errors (5xx), and transport
// failures (connection reset, mid-stream EOF, timeouts).
func IsRetryable(err error) bool
```
For transport errors, `http.Client.Do` returns errors like `*url.Error` wrapping `*net.OpError` (connection refused, reset). Those are generally retryable. But context deadline exceeded / context canceled must not be retried β the caller's deadline has passed, and retrying is pointless.
Mid-stream errors from `streamer.Get()`: scanner errors (connection reset mid-stream, EOF before [DONE]). `scanner.Scan()` returns false on EOF, and `scanner.Err()` is nil on clean EOF. Wait, let me look at Get():
```go
for s.scanner.Scan() {
if err := s.scanner.Err(); err != nil {
return nil, err
}
...
if bytes.HasPrefix(buf, prefix) {
break
}
}
buf = bytes.TrimPrefix(buf, prefix)
if bytes.Equal(buf, []byte("[DONE]")) { ... }
```
Hmm, if the scanner hits a clean EOF without seeing "data: ", the loop exits with buf still empty (from the previous iteration). Then `bytes.TrimPrefix(buf, prefix)` β empty, not [DONE], and `json.Unmarshal("")` β error "decode completion chunk: unexpected end of JSON input". So a mid-stream clean disconnect surfaces as a decode error. If the connection is reset, scanner.Err() is non-nil and that's returned.
So mid-stream failures surface as either a scanner error or a decode error. Both are retryable (transient network issues).
Hmm, but wait β should I retry mid-stream? Let me think about what's safe:
- The messages slice hasn't been modified for this turn yet (the assistant message is only appended after the stream completes).
- Tool calls from this turn haven't been executed yet.
- So re-sending the same request is safe.
The only side effect: partial deltas have been published to the dashboard's live log and possibly persisted via OnTurnEnd⦠no wait, OnTurnEnd only fires after the stream completes. The hooks OnReasoningDelta/OnContentDelta publish to the hub (live) and buffer for persistence. On a retry, the buffer would get more appended⦠let me check how the hooks buffer.
Looking at streamHooks in run.go β I've only seen the first part. Let me read the rest of run.go to understand the buffering.
Actually, let me read the rest of run.go.
Let me read the rest of `run.go` and the existing tests to understand conventions.
Tool
π§ read_file
{"path":"internal/agentrun/run.go","start_line":300,"end_line":437}
Arguments
{"path":"internal/agentrun/run.go","start_line":300,"end_line":437}
Result
{"path":"internal/agentrun/run.go","content":"\n\t\t\treasoningBuf.WriteString(delta)\n\t\t\tr.hub.Publish(jobID, livelog.Event{Type: livelog.ReasoningDelta, Text: delta})\n\t\t},\n\t\tOnContentDelta: func(delta string) {\n\t\t\tif !contentOpen {\n\t\t\t\tcontentOpen = true\n\t\t\t\tcontentBuf.Reset()\n\t\t\t\tr.hub.Publish(jobID, livelog.Event{Type: livelog.ContentStart})\n\t\t\t}\n\n\t\t\tcontentBuf.WriteString(delta)\n\t\t\tr.hub.Publish(jobID, livelog.Event{Type: livelog.ContentDelta, Text: delta})\n\t\t},\n\t\tOnTurnEnd: func() {\n\t\t\tif reasoningOpen {\n\t\t\t\tr.hub.Publish(jobID, livelog.Event{Type: livelog.ReasoningEnd})\n\t\t\t\tlogAppend(\"reasoning\", reasoningBuf.String())\n\t\t\t\tr.hub.Checkpoint(jobID)\n\t\t\t\treasoningOpen = false\n\t\t\t}\n\n\t\t\tif contentOpen {\n\t\t\t\tr.hub.Publish(jobID, livelog.Event{Type: livelog.ContentEnd})\n\t\t\t\tlogAppend(\"content\", contentBuf.String())\n\t\t\t\tr.hub.Checkpoint(jobID)\n\t\t\t\tcontentOpen = false\n\t\t\t}\n\t\t},\n\t\tOnTool: func(name, arguments, result string, toolErr bool) {\n\t\t\tr.hub.Publish(jobID, livelog.Event{\n\t\t\t\tType: livelog.Tool,\n\t\t\t\tName: name,\n\t\t\t\tArguments: arguments,\n\t\t\t\tResult: result,\n\t\t\t\tError: toolErr,\n\t\t\t})\n\n\t\t\tline, err := json.Marshal(store.ToolLogEntry{Name: name, Arguments: arguments, Result: result, Error: toolErr})\n\t\t\tif err != nil {\n\t\t\t\tr.logger.Warn(\"failed to marshal tool log entry\", \"job\", jobID, \"error\", err)\n\t\t\t} else {\n\t\t\t\tlogAppend(\"tool\", string(line))\n\t\t\t}\n\n\t\t\tr.hub.Checkpoint(jobID)\n\t\t},\n\t}\n}\n\nfunc seedMessage(ev forgejo.Event, branch, defaultBranch, instructions string, comments []forgejo.IssueComment, review *forgejo.ReviewDetail) string {\n\traw, _ := json.MarshalIndent(json.RawMessage(ev.Raw), \"\", \" \")\n\n\tvar instructionsSection string\n\tif instructions != \"\" {\n\t\tinstructionsSection = fmt.Sprintf(\"Instructions for this event, from zoo.hcl:\\n%s\\n\\n\", instructions)\n\t}\n\n\t// A pr:review run works on the PR's own head branch, not a fresh\n\t// branch off the default branch.\n\tbranchLine := fmt.Sprintf(\"Your working branch is %q, checked out from the default branch %q.\\n\\n\", branch, defaultBranch)\n\tif ev.Kind == forgejo.EventPRReview {\n\t\tbranchLine = fmt.Sprintf(\"Your working branch is %q, the pull request's head branch β commits you push here update the pull request directly.\\n\\n\", branch)\n\t}\n\n\tvar reviewSection string\n\tif review != nil {\n\t\treviewSection = renderReviewSection(review)\n\t}\n\n\tvar commentsSection string\n\tif len(comments) \u003e 0 {\n\t\tvar b strings.Builder\n\t\tfmt.Fprintf(\u0026b, \"Comments (%d):\\n\\n\", len(comments))\n\n\t\tfor i, c := range comments {\n\t\t\tfmt.Fprintf(\u0026b, \"%d. %s (%s):\\n%s\\n\\n\", i+1, c.Author, c.Created.Format(time.RFC3339), c.Body)\n\t\t}\n\n\t\tcommentsSection = b.String()\n\t}\n\n\treturn fmt.Sprintf(\n\t\t\"You were triggered by a %q event on %s/%s.\\n\\n\"+\n\t\t\t\"%s%s\"+\n\t\t\t\"%sTitle: %s\\n\\nBody:\\n%s\\n\\n%sFull event payload:\\n```json\\n%s\\n```\",\n\t\tev.Kind, ev.Owner, ev.Repo, instructionsSection, branchLine, reviewSection, ev.Title, ev.Body, commentsSection, raw)\n}\n\n// renderReviewSection renders the submitted review as a briefing\n// section: the verdict, the review body, and each inline comment with\n// its location and id (the id lets the agent refer to a specific\n// comment in its reply).\nfunc renderReviewSection(review *forgejo.ReviewDetail) string {\n\tvar b strings.Builder\n\n\tfmt.Fprintf(\u0026b, \"Review submitted by %s:\\n\", review.Reviewer)\n\tfmt.Fprintf(\u0026b, \" Verdict: %s\\n\", humanReviewState(review.State))\n\n\tif review.Body != \"\" {\n\t\tb.WriteString(\" Review body:\\n\")\n\t\tfor _, line := range strings.Split(review.Body, \"\\n\") {\n\t\t\tfmt.Fprintf(\u0026b, \" %s\\n\", line)\n\t\t}\n\t}\n\n\tif len(review.Comments) \u003e 0 {\n\t\tfmt.Fprintf(\u0026b, \" Inline comments (%d):\\n\", len(review.Comments))\n\n\t\tfor i, c := range review.Comments {\n\t\t\twhere := c.Path\n\t\t\tif c.Line \u003e 0 {\n\t\t\t\twhere = fmt.Sprintf(\"%s:%d\", c.Path, c.Line)\n\t\t\t}\n\n\t\t\tfmt.Fprintf(\u0026b, \" %d. [%s] %s (comment id %d):\\n\", i+1, where, c.Author, c.ID)\n\t\t\tfor _, line := range strings.Split(c.Body, \"\\n\") {\n\t\t\t\tfmt.Fprintf(\u0026b, \" %s\\n\", line)\n\t\t\t}\n\t\t}\n\t}\n\n\treturn b.String() + \"\\n\"\n}\n\n// humanReviewState maps a Forgejo review state to a readable verdict.\nfunc humanReviewState(state string) string {\n\tswitch state {\n\tcase \"APPROVED\":\n\t\treturn \"approved\"\n\tcase \"REQUEST_CHANGES\":\n\t\treturn \"changes requested\"\n\tcase \"COMMENT\":\n\t\treturn \"comment\"\n\tdefault:\n\t\treturn state\n\t}\n}","start_line":300,"end_line":437,"total_lines":437,"truncated":false}
Tool
π§ read_file
{"path":"internal/llm/streamer_test.go","start_line":1,"end_line":200}
Arguments
{"path":"internal/llm/streamer_test.go","start_line":1,"end_line":200}
Result
{"path":"internal/llm/streamer_test.go","content":"package llm\n\nimport (\n\t\"encoding/json\"\n\t\"io\"\n\t\"net/http\"\n\t\"strings\"\n\t\"testing\"\n)\n\nfunc sseChunk(t *testing.T, delta Message, finish string) string {\n\tt.Helper()\n\n\tbody, err := json.Marshal(ChatCompletion{\n\t\tChoices: []Choice{{Delta: delta, FinishReason: finish}},\n\t})\n\tif err != nil {\n\t\tt.Fatal(err)\n\t}\n\n\treturn \"data: \" + string(body) + \"\\n\\n\"\n}\n\nfunc newTestResponse(body string) *http.Response {\n\treturn \u0026http.Response{\n\t\tStatusCode: http.StatusOK,\n\t\tBody: io.NopCloser(strings.NewReader(body)),\n\t}\n}\n\nfunc TestStreamerMergesChunksAndStops(t *testing.T) {\n\tsse := sseChunk(t, Message{Content: \"Hel\", Role: \"assistant\"}, \"\") +\n\t\tsseChunk(t, Message{Content: \"lo\"}, \"stop\") +\n\t\t\"data: [DONE]\\n\\n\"\n\n\ts, err := NewStreamer(newTestResponse(sse))\n\tif err != nil {\n\t\tt.Fatal(err)\n\t}\n\n\tvar last *ChatCompletion\n\tfor {\n\t\tc, err := s.Get()\n\t\tif err == io.EOF {\n\t\t\tbreak\n\t\t}\n\t\tif err != nil {\n\t\t\tt.Fatal(err)\n\t\t}\n\t\tlast = c\n\t}\n\n\tif last == nil {\n\t\tt.Fatal(\"expected at least one completion\")\n\t}\n\tif len(last.Choices) != 1 {\n\t\tt.Fatalf(\"expected 1 choice, got %d\", len(last.Choices))\n\t}\n\tif last.Choices[0].Message.Content != \"Hello\" {\n\t\tt.Fatalf(\"expected merged content %q, got %q\", \"Hello\", last.Choices[0].Message.Content)\n\t}\n\tif last.Choices[0].FinishReason != \"stop\" {\n\t\tt.Fatalf(\"expected finish_reason stop, got %q\", last.Choices[0].FinishReason)\n\t}\n\tif last.Choices[0].Message.Role != \"assistant\" {\n\t\tt.Fatalf(\"expected role to stick from first chunk, got %q\", last.Choices[0].Message.Role)\n\t}\n}\n\nfunc TestStreamerNonOKStatus(t *testing.T) {\n\tresp := newTestResponse(\"boom\")\n\tresp.StatusCode = http.StatusInternalServerError\n\n\t_, err := NewStreamer(resp)\n\tif err == nil {\n\t\tt.Fatal(\"expected error for non-200 response\")\n\t}\n}\n\nfunc TestStreamerToleratesUnknownFields(t *testing.T) {\n\t// Servers may add fields we don't know about; decoding must not fail.\n\tchunk := `data: {\"id\":\"x\",\"model\":\"m\",\"choices\":[{\"index\":0,\"delta\":{\"role\":\"assistant\",\"content\":\"hi\",\"new_field\":\"boom\"},\"finish_reason\":null,\"extra\":\"y\"}],\"timings\":{\"prompt_n\":1,\"brand_new_timing\":42},\"surprise\":true}\n\n`\n\n\ts, err := NewStreamer(newTestResponse(chunk))\n\tif err != nil {\n\t\tt.Fatal(err)\n\t}\n\n\tc, err := s.Get()\n\tif err != nil {\n\t\tt.Fatalf(\"expected unknown fields to be tolerated, got error: %v\", err)\n\t}\n\tif c.Choices[0].Message.Content != \"hi\" {\n\t\tt.Fatalf(\"expected content %q, got %q\", \"hi\", c.Choices[0].Message.Content)\n\t}\n}\n\nfunc TestToolCallMerge(t *testing.T) {\n\tsse := sseChunk(t, Message{ToolCalls: []ToolCall{{Index: 0, Id: \"call1\", FunctionCall: FunctionCall{Name: \"bash\", Arguments: `{\"command\":`}}}}, \"\") +\n\t\tsseChunk(t, Message{ToolCalls: []ToolCall{{Index: 0, FunctionCall: FunctionCall{Arguments: `\"ls\"}`}}}}, \"tool_calls\") +\n\t\t\"data: [DONE]\\n\\n\"\n\n\ts, err := NewStreamer(newTestResponse(sse))\n\tif err != nil {\n\t\tt.Fatal(err)\n\t}\n\n\tvar last *ChatCompletion\n\tfor {\n\t\tc, err := s.Get()\n\t\tif err == io.EOF {\n\t\t\tbreak\n\t\t}\n\t\tif err != nil {\n\t\t\tt.Fatal(err)\n\t\t}\n\t\tlast = c\n\t}\n\n\tcalls := last.Choices[0].Message.ToolCalls\n\tif len(calls) != 1 {\n\t\tt.Fatalf(\"expected 1 tool call, got %d\", len(calls))\n\t}\n\tif calls[0].Id != \"call1\" || calls[0].FunctionCall.Name != \"bash\" {\n\t\tt.Fatalf(\"unexpected tool call: %+v\", calls[0])\n\t}\n\tif calls[0].FunctionCall.Arguments != `{\"command\":\"ls\"}` {\n\t\tt.Fatalf(\"expected merged arguments, got %q\", calls[0].FunctionCall.Arguments)\n\t}\n}","start_line":1,"end_line":132,"total_lines":132,"truncated":false}
Thinking
Now I have the full picture. Let's design the solution.
## Design
**Goal**: Retry logic in the agent loop when the LLM endpoint fails. Observed failures:
- `429` (rate limit / model cooldown, "try again in 152.0 seconds")
- `502` (engine not healthy)
- Plus other transient failures: connection reset, mid-stream disconnect, 500/503/504.
**Architecture**:
1. **`internal/llm/client.go`**: Introduce a typed `APIError` (status code + body) so callers can classify errors. Keep the same error text format (`status: %d: %s`) for backward compatibility with logs.
2. **`internal/llm/retry.go`** (new):
- `IsRetryable(err error) bool` β classification:
- `*APIError`: 429, 500, 502, 503, 504 β retryable. Other 4xx (400, 401, 403, 404, 413...) β not retryable (client errors won't be fixed by retrying the same request).
- Transport errors from `http.Client.Do` (e.g. `*url.Error` wrapping connection errors) β retryable, except when the context is done.
- Mid-stream errors (scanner errors, decode errors from truncated streams) β retryable.
- `context.Canceled` / `context.DeadlineExceeded` β not retryable.
- `Do(ctx, opts, fn func() error) error` β a generic retry executor with exponential backoff + jitter, respecting `Retry-After`.
Hmm, wait. Let me think about where the retry should live. The issue says "in the agent loop." Two candidate layers:
a) Inside `llm.Client` β retry the HTTP request. But mid-stream failures happen outside of it.
b) Inside `agentrun.runLoop` β retry the whole turn (request + stream consumption). This covers both, and the issue explicitly says "in the agent loop."
Going with (b): wrap the request + stream consumption inside `runLoop` in a retry. The `llm` package provides the typed error + classification + backoff helpers.
3. **`internal/agentrun/loop.go`**: Restructure the top of the loop:
```go
for {
if err := ctx.Err(); err != nil {
return messages, err
}
completion, err := retryTurn(ctx, func() (*llm.ChatCompletion, error) {
streamer, err := client.StreamChatCompletion(ctx, &llm.ChatCompletionRequest{
Messages: messages,
Stream: true,
Tools: tools.All(),
})
if err != nil {
return nil, err
}
// consume stream
...
return completion, nil
})
if err != nil {
return messages, fmt.Errorf("chat completion: %w", err)
}
...
}
```
Wait, but the hooks: OnReasoningDelta/OnContentDelta fire during streaming. If the stream fails and we retry, the deltas will fire again. That's fine for the dashboard (live log will show repeated text). OnTurnEnd only fires on success. The buffered persistence (reasoningBuf/contentBuf) accumulates across the retry β the buffer gets Reset() at the start of a new block (when `!reasoningOpen`), so after a mid-stream failure, `reasoningOpen` is still true, and on retry the buffer keeps appending the re-streamed content β the persisted log will have duplicate content. Hmm.
Options:
- Reset the buffer on retry. We can add a hook `OnTurnRetry` that lets the runner reset its state. Or more simply: have the retry closure reset the state.
Actually, let me think about this more carefully. In `streamHooks`:
- `OnReasoningDelta`: if !reasoningOpen β open it, Reset the buffer, publish ReasoningStart. Then append + publish delta.
- On mid-stream failure: reasoningOpen=true, buffer has partial content.
- On retry: new stream starts, first reasoning delta β reasoningOpen is still true β no Reset β buffer = partial + full new stream β duplicates in the persisted log.
To handle this cleanly, add a `Hooks.OnRetry` callback that fires before each retry attempt (with the attempt number and error). The runner's hook implementation resets the open blocks (reasoningOpen/contentOpen = false, buffers reset) and logs a "retrying" line to the job log. This keeps the persisted log clean: the partial content from the failed attempt is discarded, and the retried attempt's content is persisted as a fresh block.
Hmm, but wait: should we persist the partial content of the failed attempt? No β it's incomplete and misleading. Discarding is correct.
Alternatively, we could not reset and just let it duplicate... no, clean is better. Let's add `OnRetry func(attempt int, err error)` to Hooks.
Actually wait β let me reconsider the scope. Is mid-stream retry overkill? The observed failures in the issue are request-level (status 429/502 from NewStreamer). But the issue title is generic: "if the LLM endpoint is failing." A robust implementation should handle mid-stream disconnects too, because those are common with LLM endpoints (long streams over flaky proxies). And the architecture is the same: wrap the whole turn.
But there's a subtle danger with mid-stream retry: the model might be non-deterministic (temperature 0.6, seed 5 β seed is set, so it's fairly deterministic, but the server may not honor it). The retried turn may produce different content. That's fine β it's a fresh attempt at the same turn.
Another consideration: **idempotency of tool calls**. Tool calls only execute after a turn completes successfully, so no double-execution risk. Good.
**Backoff policy**:
- Initial delay: ~1 second? For 429 with "try again in 152 seconds", exponential backoff from 1s: 1, 2, 4, 8, 16, 32, 64, 128, 256... with full jitter. It'd take about 9 attempts to reach 152s, totaling ~8 minutes. The run timeout is 120 minutes, so that's fine.
- Cap the per-attempt delay at, say, 60 seconds (or 120 seconds?). Let's cap at 60 seconds. Then for a 152-second cooldown: attempts at 0, 1, 3, 7, 15, 31, 60 (cap), 60, 60... cumulative: 1, 4, 11, 26, 57, 117, 177, 237... So after about 8 attempts (~3-4 minutes) we'd get past 152 seconds. Good.
- Honor `Retry-After` header when present (429/503). litellm's message "try again in 152.0 seconds" is in the body, not the header β we can't reliably parse that, but Retry-After header support is the standard mechanism, and many gateways set it.
- Jitter: full jitter (random between 0 and min(cap, base*2^attempt))? Or equal jitter? Full jitter is the AWS-recommended approach. But full jitter can give very short delays (near 0), which is fine. Actually, for rate-limit cooldowns, a very short first retry is wasteful but harmless. Let's use "equal jitter": delay/2 + rand*delay/2, ensuring at least half the computed delay. Hmm, more simply: delay = min(cap, base * 2^attempt); sleep = delay/2 + rand(delay/2). This guarantees a minimum of delay/2 and avoids thundering herd.
Actually, let's keep it simple and standard: exponential backoff with full jitter is the canonical recommendation (AWS architecture blog). `sleep = rand(0, min(cap, base * 2^n))`. But with full jitter, the expected total time to cover 152s is longer. With equal jitter, the expected value is 3/4 of the exponential. Either works. Let's go with full jitter but with a base of 2s: expected cumulative after n attempts β sum of base*2^(i-1)/2 = base*(2^n - 1)/2... with base 2s: after 8 attempts, expected cumulative β 2*(256-1)/2 = 255s. Covers 152s comfortably.
Hmm, let's not overthink. I'll implement:
```go
type RetryPolicy struct {
MaxAttempts int // total attempts including the first
BaseDelay time.Duration // delay before the 2nd attempt
MaxDelay time.Duration // cap on any single delay
// Now and Sleep for testability
}
```
And `retry.Do(ctx, policy, fn func() error) error`.
For testability, the sleep function should be injectable. In tests we can set the policy's BaseDelay/MaxDelay to tiny values, or inject a sleeper. Let's inject a `sleep func(ctx, time.Duration) error` field with a default.
**Where to put the retry helper?** Options:
- `internal/llm/retry.go` β next to the client, since it's LLM-specific (classifying LLM API errors).
- `internal/retry` β a generic package.
The classification is LLM-specific (APIError, mid-stream errors), so let's put both the policy and the classifier in `internal/llm`. The `Do` function itself is generic, but co-locating it is fine and keeps the change surgical.
**Classification details**:
```go
// APIError is a non-2xx response from the LLM endpoint.
type APIError struct {
Status int
Body string
}
func (e *APIError) Error() string {
return fmt.Sprintf("status: %d: %s", e.Status, e.Body)
}
```
Wait, the current format is `fmt.Errorf("status: %d: %s", resp.StatusCode, body)`. Keep that exact format so existing log messages don't change.
`IsRetryable(err error) bool`:
- If err is `*APIError`: return 429 || (500..599).
- If err is `context.Canceled` or `context.DeadlineExceeded` (via errors.Is): return false.
- Otherwise (transport errors, scanner errors, decode errors from truncated streams): return true.
Hmm, is "otherwise β retryable" too broad? What about a persistent 400 that comes back as a body error? No β 400 comes back as an APIError from NewStreamer (non-2xx status). So any error that isn't an APIError is a transport/decode problem, which is transient by nature (connection reset, EOF, timeout). Decode errors on a *valid* chunk would be a bug, but... a decode error on a truncated chunk (mid-stream disconnect) is transient. If the server consistently sends garbage, we'd retry MaxAttempts times and then fail β acceptable, bounded.
Actually, wait, one more: `http.Client.Do` with a context that has a deadline β if the request times out, the error is `*url.Error` wrapping `context.DeadlineExceeded`. We check `errors.Is(err, context.DeadlineExceeded)` first β not retryable. Good.
**Retry-After**: `APIError` should capture the `Retry-After` header (seconds or HTTP date). In `NewStreamer`, when the status is non-2xx, read the header. Then the retry loop can query it. Let's add a method:
```go
func (e *APIError) RetryAfter() (time.Duration, bool)
```
Parse `Retry-After` as integer seconds, or as an HTTP date (http.ParseTime). Store the raw header string in the struct and parse on demand.
**Restructuring runLoop**:
```go
func runLoop(ctx context.Context, client *llm.Client, toolsCtx tools.Context, messages []llm.Message, hooks Hooks) ([]llm.Message, error) {
policy := llm.DefaultRetryPolicy()
for {
if err := ctx.Err(); err != nil {
return messages, err
}
completion, err := llm.Do(ctx, policy, func() (*llm.ChatCompletion, error) {
...
})
...
}
}
```
Hmm, but `llm.Do` needs to call hooks on retry. Let's define it in the llm package:
```go
// Do runs fn with retries according to the policy, waiting between
// attempts with exponential backoff and jitter. fn returns an error
// that decides whether to retry: IsRetryable errors are retried,
// others are returned immediately. The context bounds the whole
// sequence of attempts.
func Do[T any](ctx context.Context, p Policy, fn func() (T, error)) (T, error)
```
Go 1.18+ generics β let me check the go.mod for the Go version.
Actually, do we need generics? The turn function returns (*llm.ChatCompletion, error). Generics keep it clean. Let me check go.mod.
Also, where does the retry hook fire? The retry executor doesn't know about hooks. We can pass an `OnRetry func(attempt int, delay time.Duration, err error)` option to Do, or handle it in runLoop by wrapping fn... no, the retry decision happens inside Do. Let's give Policy an optional `OnRetry` callback:
```go
type Policy struct {
MaxAttempts int
BaseDelay time.Duration
MaxDelay time.Duration
// OnRetry, if set, is called before each retry wait with the
// 1-based index of the failed attempt and the error that caused
// it.
OnRetry func(attempt int, err error)
}
```
Hmm, wait. Let me reconsider: should the retry live in the llm package as a generic `Do`, or should it be a method on Client?
Alternative cleaner design: **`client.StreamChatCompletion` already returns (*Streamer, error)**. We could add a method to the client: `client.ChatCompletion(ctx, req) (*ChatCompletion, error)` that does the request + full stream consumption + retry internally. Then runLoop just calls it. This keeps the retry entirely inside the llm package and makes runLoop simpler. But the hooks (deltas) are needed during streaming, so the stream consumption must be visible to runLoop...
Option: `client.ChatCompletion(ctx, req, onDelta func(msg *Message)) (*ChatCompletion, error)` β the client drives the stream and calls onDelta for each merged message. runLoop passes a closure that invokes the hooks. This is clean!
But wait β the delta logic in runLoop tracks prevContent/prevReasoning and only fires hooks for *new* suffixes. The streamer's Get() returns the accumulated completion each time. So onDelta would receive the accumulated message, and runLoop's closure computes the suffix. That works.
Hmm, but this changes the Streamer's public surface... actually no, Streamer stays as is. Let me sketch:
```go
// in llm/client.go
// ChatCompletion runs a streaming chat completion to completion,
// retrying transient failures (rate limits, server errors, dropped
// connections) with exponential backoff. onDelta, if non-nil, is
// called after each streamed chunk with the accumulated message so
// far.
func (c *Client) ChatCompletion(ctx context.Context, req *ChatCompletionRequest, onDelta func(*Message)) (*Message, error)
```
Wait, what should we return? The completion's choice message (with FinishReason?). runLoop needs `choice.FinishReason` and `choice.Message`. So return `*Choice`? Or `*ChatCompletion`? Let's return `*Choice` β the first choice. Hmm, runLoop currently checks `completion == nil || len(completion.Choices) == 0` β "model returned an empty completion".
Let's keep it minimal and not over-abstract. Decision:
**Put the retry in `agentrun/loop.go`** (the issue says "in the agent loop"), and use a small retry helper in the `llm` package for classification + backoff. The turn closure in runLoop does request + stream + hooks, and the retry wraps it. This keeps the llm package's API surface small (just APIError + IsRetryable + Do), and the loop change is localized.
Hmm, but actually... let me reconsider once more. There's a real advantage to putting request+retry inside the llm client: the client is the natural owner of "how to talk to the LLM endpoint," and the retry policy is an endpoint concern. And the mid-stream retry: the client can own stream consumption via an onDelta callback. Then runLoop becomes:
```go
choice, err := client.ChatCompletion(ctx, &llm.ChatCompletionRequest{...}, func(msg *llm.Message) {
// fire delta hooks
})
```
And all the retry/backoff logic lives in one place in the llm package, testable with a fake HTTP server. runLoop stays focused on the conversation state machine.
I think this is architecturally cleaner. But it moves the stream-consumption loop from runLoop to the client. The delta hook logic (prevContent/prevReasoning tracking) moves into a closure in runLoop. That's fine.
Wait, one issue: the current runLoop consumes the stream and fires hooks per chunk. If the client owns consumption, the onDelta callback receives the accumulated `*Message` per chunk. runLoop's closure:
```go
var prevContent, prevReasoning string
choice, err := client.ChatCompletion(ctx, req, func(msg *llm.Message) {
if hooks.OnReasoningDelta != nil && len(msg.ReasoningContent) > len(prevReasoning) {
hooks.OnReasoningDelta(msg.ReasoningContent[len(prevReasoning):])
}
prevReasoning = msg.ReasoningContent
if hooks.OnContentDelta != nil && len(msg.Content) > len(prevContent) {
hooks.OnContentDelta(msg.Content[len(prevContent):])
}
prevContent = msg.Content
})
```
And on retry, we need to reset prevContent/prevReasoning (the new stream restarts from empty). This is where the OnRetry hook comes in: runLoop's retry callback resets the prev* state, and the runner's OnRetry hook resets the buffers.
Hmm, OK, wait, I keep going back and forth. Let me settle it with a clear principle:
**Principle: the retry policy belongs with the thing being retried, and the thing being retried is "one model turn" (request + stream). The agent loop is the owner of the turn.**
Both placements are defensible. The issue literally says "retry logic in the agent loop." Let's put the retry mechanism in the `llm` package (as a reusable `Do` + `IsRetryable` + `Policy`), and have the **agent loop** use it to wrap the turn. This matches the issue's wording, keeps the llm package's API stable (no new client methods), and makes runLoop the single place where "a turn can fail transiently" is handled.
Final design:
### `internal/llm/errors.go` (or in client.go/streamer.go)
- `APIError{Status int, Body string, retryAfter string}` with `Error()`, `Unwrap`? No. And `RetryAfter() (time.Duration, bool)`.
- `NewStreamer` returns `*APIError` for non-2xx.
### `internal/llm/retry.go`
```go
// IsRetryable reports whether err is a transient failure worth
// retrying: HTTP 429 and 5xx responses, and transport-level failures
// (connection errors, mid-stream disconnects). Context cancellation
// and deadline errors are never retried.
func IsRetryable(err error) bool
// Policy configures retry.Do.
type Policy struct {
MaxAttempts int
BaseDelay time.Duration
MaxDelay time.Duration
// OnRetry, if non-nil, is called before each retry wait with the
// 1-based index of the attempt that just failed and its error.
OnRetry func(attempt int, err error)
// sleep is for tests; defaults to time.Sleep bound to ctx.
sleep func(ctx context.Context, d time.Duration) error
}
func DefaultPolicy() Policy // MaxAttempts: 8? BaseDelay: 2s? MaxDelay: 60s?
// Do runs fn, retrying while IsRetryable(err) and attempts remain.
func Do[T any](ctx context.Context, p Policy, fn func() (T, error)) (T, error)
```
Wait, one design question: should `Do` take the policy with OnRetry, or should Do take a separate observer? Policy with a callback field is a bit odd (policy describes *when* to retry, callbacks are *side effects*). But it's practical and keeps the signature small. Alternatively:
```go
func Do(ctx context.Context, p Policy, onRetry func(attempt int, err error), fn func() error) error
```
Hmm. Let's keep OnRetry on Policy β no wait. Let me think about who calls Do. Only runLoop. And runLoop needs to reset state on retry. Let's just make it a simple function signature:
Actually, the simplest, most robust approach that avoids API bikeshedding:
```go
// in llm/retry.go
type Retry struct {
MaxAttempts int
BaseDelay time.Duration
MaxDelay time.Duration
// BeforeRetry, if set, is called with the failed attempt number
// (1-based) and its error before the backoff wait.
BeforeRetry func(attempt int, err error)
// Sleep, if set, is used instead of sleeping bound to ctx (tests).
Sleep func(ctx context.Context, d time.Duration) error
}
func DefaultRetry() Retry { ... }
// Do runs fn with retries: while fn returns an error that IsRetryable
// reports as transient and attempts remain, it waits with exponential
// backoff (capped, jittered, honoring Retry-After) and tries again.
func Do[T any](ctx context.Context, r Retry, fn func() (T, error)) (T, error) {
...
}
```
Let me write the Do implementation:
```go
func Do[T any](ctx context.Context, r Retry, fn func() (T, error)) (T, error) {
if r.MaxAttempts < 1 {
r.MaxAttempts = 1
}
if r.BaseDelay <= 0 {
r.BaseDelay = 2 * time.Second
}
if r.MaxDelay <= 0 {
r.MaxDelay = 60 * time.Second
}
sleep := r.Sleep
if sleep == nil {
sleep = func(ctx context.Context, d time.Duration) error {
t := time.NewTimer(d)
defer t.Stop()
select {
case <-ctx.Done():
return ctx.Err()
case <-t.C:
return nil
}
}
}
var zero T
var lastErr error
for attempt := 1; ; attempt++ {
if err := ctx.Err(); err != nil {
return zero, err
}
result, err := fn()
if err == nil {
return result, nil
}
lastErr = err
if attempt >= r.MaxAttempts || !IsRetryable(err) {
return zero, err
}
if r.BeforeRetry != nil {
r.BeforeRetry(attempt, err)
}
delay := backoffDelay(r.BaseDelay, r.MaxDelay, attempt)
if ra, ok := retryAfter(err); ok && ra > delay {
delay = ra
}
if err := sleep(ctx, delay); err != nil {
return zero, err
}
}
}
```
Wait, subtle point: if sleep returns ctx.Err() (deadline exceeded during the wait), we return that. Good.
`backoffDelay`: exponential with jitter:
```go
// backoffDelay returns the wait before the given (1-based) failed
// attempt's retry: base * 2^(attempt-1) capped at max, with equal
// jitter (half fixed, half random) so concurrent runs don't retry in
// lockstep.
func backoffDelay(base, max time.Duration, attempt int) time.Duration {
d := base << (attempt - 1) // careful about overflow for large attempt
if d < 0 || d > max {
d = max
}
return d/2 + time.Duration(rand.Int63n(int64(d/2)+1))
}
```
Hmm, `base << (attempt-1)` β with base 2s = 2e9 ns, attempt up to ~8 β 2e9 * 128 = 2.56e11, fine. But if someone sets MaxAttempts high, the shift could overflow int64. Cap the shift: if attempt-1 >= 32, just use max. Let me write it safely:
```go
func backoffDelay(base, max time.Duration, attempt int) time.Duration {
d := base
for i := 1; i < attempt && d < max; i++ {
d *= 2
if d < 0 { // overflow
d = max
break
}
}
if d > max {
d = max
}
jitter := d / 2
if jitter > 0 {
d = jitter + time.Duration(rand.Int63n(int64(jitter)+1))
}
return d
}
```
Hmm wait, this loop: attempt=1 β no doubling (i starts at 1, i < 1 is false) β d = base. So the delay before retry #1 (after the first failure) = base. attempt=2 β d = base*2. Good.
Equal jitter: d/2 + rand(0, d/2]. Expected = 3/4 d. Good.
`rand.Int63n` β use `math/rand/v2`? Let me check the Go version in go.mod. If Go 1.20+, math/rand/v2 isn't available (it's 1.22+). Let me check.
`retryAfter(err)`:
```go
func retryAfter(err error) (time.Duration, bool) {
var apiErr *APIError
if errors.As(err, &apiErr) {
return apiErr.RetryAfter()
}
return 0, false
}
```
And `APIError.RetryAfter()`:
```go
// RetryAfter parses the Retry-After header (delta-seconds or HTTP
// date) if the server sent one.
func (e *APIError) RetryAfter() (time.Duration, bool) {
if e.retryAfter == "" {
return 0, false
}
if secs, err := strconv.Atoi(e.retryAfter); err == nil {
if secs < 0 {
return 0, false
}
return time.Duration(secs) * time.Second, true
}
if t, err := http.ParseTime(e.retryAfter); err == nil {
if d := time.Until(t); d > 0 {
return d, true
}
}
return 0, false
}
```
Hmm, `time.Until` inside a method is impure (depends on now) β fine for this use, but it makes testing a bit awkward. Acceptable. Actually, for testability, maybe parse to an absolute time and let the caller compute... Let's keep it simple: `time.Until`.
### `IsRetryable`
```go
// IsRetryable reports whether err is a transient failure worth
// retrying.
//
// HTTP responses are retried only for rate limits (429) and server
// errors (5xx): client errors (400, 401, ...) will not change on
// retry. Anything else β connection failures, timeouts, a stream that
// dropped mid-response β is treated as transient. A canceled or
// deadline-exceeded context is never retried.
func IsRetryable(err error) bool {
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
return false
}
var apiErr *APIError
if errors.As(err, &apiErr) {
return apiErr.Status == http.StatusTooManyRequests || apiErr.Status >= 500
}
return true
}
```
Wait β one more case: what about a 408 Request Timeout? That's technically a client-side timeout, but it's transient (the server gave up waiting; retrying is usually fine). OpenAI's guidance: retry on 408, 409, 429, 5xx. Let's include 408.
Also, what about errors that are neither APIError nor context errors but are actually permanent? E.g. a bug in our JSON marshaling (json.Marshal error in StreamChatCompletion) β that's not retryable, but we'd retry it MaxAttempts times with backoff before failing. Wasteful but bounded and rare. Hmm. Actually, json.Marshal on our own request struct essentially never fails. Fine.
But wait β there's a real concern: "model returned an empty completion" β that's created in runLoop *after* the stream completes, outside the retry closure. So it's not retried. Should it be? An empty completion (200 OK but no choices) is a server bug... could be transient. Hmm, let's keep the scope tight: retry transport/HTTP failures. Empty completion stays a hard error. Actually... let me reconsider β where do we define the boundary of the retry closure?
If the closure = request + stream consumption (returning the final *ChatCompletion), then "empty completion" can be inside or outside. Let's put it inside: if the stream ends cleanly with no choices, return an error "model returned an empty completion" β IsRetryable β true (not an APIError, not a context error) β retried. Is that desirable? A server that consistently returns empty completions would burn all attempts, then fail with the same message. A server that occasionally does it (buggy proxy returning an empty 200) benefits. I think retrying is the right call β it's a weird response, and retrying is cheap relative to a failed run.
Hmm, wait, actually there's a subtle issue: if the stream completes with a valid completion but no choices... the current code errors. OK, let's include it in the retry. Fine.
### Changes to runLoop
```go
func runLoop(ctx context.Context, client *llm.Client, toolsCtx tools.Context, messages []llm.Message, hooks Hooks) ([]llm.Message, error) {
retry := llm.DefaultRetry()
retry.BeforeRetry = func(attempt int, err error) {
if hooks.OnRetry != nil {
hooks.OnRetry(attempt, err)
}
}
for {
if err := ctx.Err(); err != nil {
return messages, err
}
var (
completion *llm.ChatCompletion
prevContent string
prevReason string
)
err := llm.Do(ctx, retry, func() error {
prevContent, prevReason = "", ""
streamer, err := client.StreamChatCompletion(ctx, &llm.ChatCompletionRequest{
Messages: messages,
Stream: true,
Tools: tools.All(),
})
if err != nil {
return fmt.Errorf("chat completion: %w", err)
}
completion = nil
for {
c, err := streamer.Get()
if err == io.EOF {
return nil
}
if err != nil {
return fmt.Errorf("stream completion: %w", err)
}
completion = c
if len(c.Choices) > 0 {
msg := c.Choices[0].Message
if hooks.OnReasoningDelta != nil && len(msg.ReasoningContent) > len(prevReason) {
hooks.OnReasoningDelta(msg.ReasoningContent[len(prevReason):])
}
prevReason = msg.ReasoningContent
if hooks.OnContentDelta != nil && len(msg.Content) > len(prevContent) {
hooks.OnContentDelta(msg.Content[len(prevContent):])
}
prevContent = msg.Content
}
}
})
if err != nil {
return messages, err
}
if hooks.OnTurnEnd != nil {
hooks.OnTurnEnd()
}
if completion == nil || len(completion.Choices) == 0 {
return messages, fmt.Errorf("model returned an empty completion")
}
...
}
}
```
Wait, problem: `prevContent`/`prevReason` are captured by the closure and reset at the start of each attempt β good, that handles the delta tracking across retries.
But hmm, the "empty completion" check is now *outside* the retry. If the stream ends with no choices, Do returns nil error (the stream completed "successfully"), and then we error out without retrying. Inconsistent with what I said earlier. Options:
1. Move the empty check inside the closure β retried.
2. Leave it outside β not retried.
I said earlier I'd prefer retrying. Let's move it inside:
```go
err := llm.Do(ctx, retry, func() error {
...
for { ... }
if completion == nil || len(completion.Choices) == 0 {
return fmt.Errorf("model returned an empty completion")
}
return nil
})
```
Hmm wait, but the error wrapping: currently the outer error is "chat completion: %w" for request errors and "stream completion: %w" for stream errors. run.go wraps with "agent loop: %w". The issue's log shows "agent loop: chat completion: status: 429: ...". Let's keep those prefixes inside the closure so the final error text is unchanged.
Now, on the final error after exhausting retries β should we annotate it with the attempt count? E.g. "after 8 attempts: ..."? Useful for debugging. We could wrap in Do: if all attempts fail with a retryable error, return fmt.Errorf("after %d attempts: %w", attempts, lastErr)? Hmm, that changes the error text. I think adding attempt info is genuinely useful for ops ("we retried and still failed" vs "failed immediately"). But it also complicates the text. Let me keep Do returning the last error as-is (simple, predictable), and log the retries via BeforeRetry (visible in the job log via the OnRetry hook + slog). Hmm, but runLoop has no logger... the hooks.OnRetry callback β the runner logs it.
Wait, actually, should the final error from Do include the attempt count? Let me think about what abrander sees: "zoo agent run failed: agent loop: chat completion: status: 502: engine not healthy". If we retried 8 times over ~3 minutes, that message alone doesn't tell you that. With "after 8 attempts: chat completion: status: 502: ..." it does. Let's add it β but only when at least one retry happened (attempt > 1), to keep the single-attempt text identical to today.
Hmm, where? In Do:
```go
if attempt >= r.MaxAttempts {
if attempt > 1 {
return zero, fmt.Errorf("after %d attempts: %w", attempt, err)
}
return zero, err
}
```
Wait, but that wraps the error β errors.As/Is still work through %w. Fine.
Hmm, actually, hold on. Let me reconsider whether to include this. It's a nice touch for observability. OK, include it.
### Hooks.OnRetry
Add to the Hooks struct:
```go
// OnRetry fires before the loop waits to retry a failed turn: attempt
// is the 1-based index of the attempt that just failed, err its
// error. The partial content of the failed attempt has been streamed
// already (OnReasoningDelta/OnContentDelta fired for it); implementers
// should reset any per-turn accumulation so the retried attempt is
// recorded fresh.
OnRetry func(attempt int, err error)
```
And in run.go's streamHooks:
```go
OnRetry: func(attempt int, err error) {
// Discard the partial content of the failed attempt so the retry
// is recorded as a fresh block.
reasoningBuf.Reset()
contentBuf.Reset()
reasoningOpen = false
contentOpen = false
logAppend("system", fmt.Sprintf("llm endpoint failed (attempt %d): %s β retrying", attempt, err))
r.hub.Publish(jobID, livelog.Event{...})
},
```
Hmm wait, what log streams exist? "reasoning", "content", "tool"... let me check the store's schema and web to see what streams are expected. Let me grep for the stream names.
Also, should we publish a livelog event? Let me check the livelog event types. Let me look at livelog/hub.go and the web templates to see what's rendered.
Actually, let me be careful about scope. The core ask is retry logic. The OnRetry hook + log line is a small, high-value addition (visibility into retries in the job log). Let me check what log streams the web UI supports.
### Default retry parameters
What's reasonable?
- MaxAttempts: 8? Total wait with base 2s, cap 60s, equal jitter: worst-case total β 2+4+8+16+32+60+60 = 182s + jitter... wait, attempts: after attempt 1 fails, wait d1 (base 2s), attempt 2, wait d2 (4s), attempt 3, d3 (8s), attempt 4, d4 (16s), attempt 5, d5 (32s), attempt 6, d6 (60s cap), attempt 7, d7 (60s cap), attempt 8 β fail. Total worst-case wait β 2+4+8+16+32+60+60 = 182s β 3 minutes. With equal jitter, expected β 0.75 * that β 2.25 minutes.
For the litellm case "try again in 152 seconds": if the Retry-After header is set to 152, we'd wait 152s on the first retry β success on attempt 2. If not set (body only), the backoff reaches cumulative 152s around attempt 6-7 (cumulative: 2, 6, 14, 30, 62, 122, 182 β attempt 7 starts at ~122s+... hmm, attempt 7 is at cumulative 122s, attempt 8 at 182s. So the 152s cooldown is covered by attempt 7-8. Cutting it close but OK.
Should we make MaxAttempts higher, like 10? With cap 60s: attempts 8, 9, 10 add 60s each β total worst-case 362s β 6 minutes. The run timeout is 120 minutes, so even 20 attempts (worst-case ~15 minutes) is fine. But each attempt also costs request time.
Let me think about what's most reasonable for a long-running agent job: LLM endpoints with model cooldowns (like litellm's "try again in 152 seconds") can take minutes. A run that gives up after 3 minutes when the model would've been available in 4 is a bad trade. Let's go with MaxAttempts 10, BaseDelay 2s, MaxDelay 60s. Worst-case total wait β 2+4+8+16+32+60*5 = 362s β 6 minutes. Expected ~4.5 minutes. Reasonable.
Hmm, should these be configurable in zoo.hcl (llm block)? E.g. `retry_max_attempts`? That would be nice, but expands the config surface. The issue asks for "some kind of retry logic." Let's keep the defaults in code and not add config knobs β surgical. Hmm... but config knobs would let abrander tune per-deployment. The llm block has name/openai/token/model. Adding optional fields is low-risk:
```hcl
llm "qwen" {
openai = "..."
token = "..."
model = "..."
}
```
I'm leaning toward not adding config now β keep it simple, defaults are reasonable, can be added later if needed. Actually... let me reconsider. The retry duration matters operationally (how long a job can stall on a dead endpoint before failing). With MaxAttempts 10, worst-case ~6 minutes of stalling β fine for a 120-minute job. No config needed. Decision: code defaults, no config.
### Now, the OnRetry hook and log streams
Let me check the store schema + web to see the log streams.
Also β important β let me check `go.mod` for the Go version (generics require 1.18+; math/rand/v2 requires 1.22+).
And check how runLoop is called elsewhere (tests? seed_test.go?). Let me grep for runLoop and StreamChatCompletion usages.
Now let me check the go.mod, how `runLoop`/`StreamChatCompletion` are used elsewhere, and the log stream conventions.
Tool
π§ bash
{"command":"cd /project && cat go.mod && go version"}
Arguments
{"command":"cd /project && cat go.mod && go version"}
Result
module github.com/abrander/zoo
go 1.26.4
require (
codeberg.org/mvdkleijn/forgejo-sdk/forgejo/v2 v2.2.0
github.com/docker/docker v28.5.2+incompatible
github.com/google/uuid v1.6.0
github.com/hashicorp/hcl/v2 v2.24.0
modernc.org/sqlite v1.56.0
)
require (
github.com/42wim/httpsig v1.2.3 // indirect
github.com/Microsoft/go-winio v0.6.2 // indirect
github.com/agext/levenshtein v1.2.1 // indirect
github.com/apparentlymart/go-textseg/v15 v15.0.0 // indirect
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/containerd/errdefs v1.0.0 // indirect
github.com/containerd/errdefs/pkg v0.3.0 // indirect
github.com/containerd/log v0.1.0 // indirect
github.com/davidmz/go-pageant v1.0.2 // indirect
github.com/distribution/reference v0.6.0 // indirect
github.com/docker/go-connections v0.8.1 // indirect
github.com/docker/go-units v0.5.0 // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
github.com/felixge/httpsnoop v1.1.0 // indirect
github.com/go-fed/httpsig v1.1.0 // indirect
github.com/go-logr/logr v1.4.4 // indirect
github.com/go-logr/stdr v1.2.2 // indirect
github.com/google/go-cmp v0.7.0 // indirect
github.com/hashicorp/go-version v1.7.0 // indirect
github.com/mattn/go-isatty v0.0.24 // indirect
github.com/mitchellh/go-wordwrap v1.0.1 // indirect
github.com/moby/docker-image-spec v1.3.1 // indirect
github.com/moby/sys/atomicwriter v0.1.0 // indirect
github.com/moby/term v0.5.2 // indirect
github.com/morikuni/aec v1.1.0 // indirect
github.com/ncruces/go-strftime v1.0.0 // indirect
github.com/opencontainers/go-digest v1.0.0 // indirect
github.com/opencontainers/image-spec v1.1.1 // indirect
github.com/pkg/errors v0.9.1 // indirect
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
github.com/zclconf/go-cty v1.16.3 // indirect
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.70.0 // indirect
go.opentelemetry.io/otel v1.45.0 // indirect
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.45.0 // indirect
go.opentelemetry.io/otel/metric v1.45.0 // indirect
go.opentelemetry.io/otel/trace v1.45.0 // indirect
golang.org/x/crypto v0.39.0 // indirect
golang.org/x/mod v0.37.0 // indirect
golang.org/x/sync v0.22.0 // indirect
golang.org/x/sys v0.47.0 // indirect
golang.org/x/text v0.40.0 // indirect
golang.org/x/time v0.15.0 // indirect
golang.org/x/tools v0.47.0 // indirect
gotest.tools/v3 v3.5.2 // indirect
modernc.org/libc v1.74.4 // indirect
modernc.org/mathutil v1.7.1 // indirect
modernc.org/memory v1.11.0 // indirect
)
go version go1.26.7 linux/amd64
Tool
π§ grep_search
{"query":"runLoop|StreamChatCompletion|NewClient","context_after":2}
Arguments
{"query":"runLoop|StreamChatCompletion|NewClient","context_after":2}
Result
/project/cmd/zoo/main.go:64: fgClient, err := forgejo.NewClient(cfg.Forgejo)
/project/cmd/zoo/main.go-65- if err != nil {
/project/cmd/zoo/main.go-66- return fmt.Errorf("forgejo client: %w", err)
--
/project/internal/web/web_test.go:210: fg, err := forgejo.NewClient(config.Forgejo{URL: api.URL, Token: "test"})
/project/internal/web/web_test.go-211- if err != nil {
/project/internal/web/web_test.go-212- t.Fatal(err)
--
/project/internal/llm/client.go:20:// NewClient builds a Client from a configured llm block.
/project/internal/llm/client.go:21:func NewClient(cfg config.LLM) *Client {
/project/internal/llm/client.go-22- return &Client{
/project/internal/llm/client.go-23- endpoint: cfg.OpenAI,
--
/project/internal/llm/client.go:30:func (c *Client) StreamChatCompletion(ctx context.Context, req *ChatCompletionRequest) (*Streamer, error) {
/project/internal/llm/client.go-31- req2 := *req
/project/internal/llm/client.go-32-
--
/project/internal/forgejo/client.go:24:func NewClient(cfg config.Forgejo) (*Client, error) {
/project/internal/forgejo/client.go:25: c, err := sdk.NewClient(cfg.URL, sdk.SetToken(cfg.Token))
/project/internal/forgejo/client.go-26- if err != nil {
/project/internal/forgejo/client.go-27- return nil, fmt.Errorf("forgejo client: %w", err)
--
/project/internal/forgejo/client.go:46: client, _ := sdk.NewClient(c.baseURL, sdk.SetToken(token))
/project/internal/forgejo/client.go-47- return &Client{sdk: client, baseURL: c.baseURL, token: token}
/project/internal/forgejo/client.go-48-}
--
/project/internal/forgejo/client.go:61: sudoClient, err := sdk.NewClient(c.baseURL, sdk.SetToken(c.token), sdk.SetSudo(username))
/project/internal/forgejo/client.go-62- if err != nil {
/project/internal/forgejo/client.go-63- return nil, fmt.Errorf("forgejo client sudo %q: %w", username, err)
--
/project/internal/forgejo/poll_test.go:38: client, err := NewClient(config.Forgejo{URL: serverURL, Token: "test"})
/project/internal/forgejo/poll_test.go-39- if err != nil {
/project/internal/forgejo/poll_test.go-40- t.Fatalf("new client: %v", err)
--
/project/internal/forgejo/client_test.go:29:// /api/v1/version probe (made by NewClient) plus the routes registered
/project/internal/forgejo/client_test.go-30-// on the returned mux.
/project/internal/forgejo/client_test.go-31-func newTestServer(t *testing.T) (*httptest.Server, *http.ServeMux) {
--
/project/internal/forgejo/client_test.go:76: client, err := NewClient(config.Forgejo{URL: server.URL, Token: "test"})
/project/internal/forgejo/client_test.go-77- if err != nil {
/project/internal/forgejo/client_test.go-78- t.Fatalf("new client: %v", err)
--
/project/internal/forgejo/client_test.go:110: client, err := NewClient(config.Forgejo{URL: server.URL, Token: "test"})
/project/internal/forgejo/client_test.go-111- if err != nil {
/project/internal/forgejo/client_test.go-112- t.Fatalf("new client: %v", err)
--
/project/internal/forgejo/client_test.go:147: client, err := NewClient(config.Forgejo{URL: server.URL, Token: "test"})
/project/internal/forgejo/client_test.go-148- if err != nil {
/project/internal/forgejo/client_test.go-149- t.Fatalf("new client: %v", err)
--
/project/internal/forgejo/client_test.go:188: client, err := NewClient(config.Forgejo{URL: server.URL, Token: "test"})
/project/internal/forgejo/client_test.go-189- if err != nil {
/project/internal/forgejo/client_test.go-190- t.Fatalf("new client: %v", err)
--
/project/internal/forgejo/client_test.go:226: client, err := NewClient(config.Forgejo{URL: server.URL, Token: "test"})
/project/internal/forgejo/client_test.go-227- if err != nil {
/project/internal/forgejo/client_test.go-228- t.Fatalf("new client: %v", err)
--
/project/internal/forgejo/client_test.go:246: client, err := NewClient(config.Forgejo{URL: server.URL, Token: "test"})
/project/internal/forgejo/client_test.go-247- if err != nil {
/project/internal/forgejo/client_test.go-248- t.Fatalf("new client: %v", err)
--
/project/internal/agentrun/loop.go:12:// Hooks lets a caller observe a runLoop invocation live, as the model
/project/internal/agentrun/loop.go-13-// streams and tool calls execute, in addition to the final []llm.Message
/project/internal/agentrun/loop.go-14-// it returns. Any of these may be nil.
--
/project/internal/agentrun/loop.go:30:// runLoop is a headless port of ../a's App.generate(): send messages +
/project/internal/agentrun/loop.go-31-// tool defs, get a completion, run any tool_calls and append their
/project/internal/agentrun/loop.go-32-// results, repeat until a plain finish or ctx is done.
/project/internal/agentrun/loop.go:33:func runLoop(ctx context.Context, client *llm.Client, toolsCtx tools.Context, messages []llm.Message, hooks Hooks) ([]llm.Message, error) {
/project/internal/agentrun/loop.go-34- for {
/project/internal/agentrun/loop.go-35- if err := ctx.Err(); err != nil {
--
/project/internal/agentrun/loop.go:39: streamer, err := client.StreamChatCompletion(ctx, &llm.ChatCompletionRequest{
/project/internal/agentrun/loop.go-40- Messages: messages,
/project/internal/agentrun/loop.go-41- Stream: true,
--
/project/internal/agentrun/docker.go:26: cli, err := client.NewClientWithOpts(client.FromEnv, client.WithAPIVersionNegotiation())
/project/internal/agentrun/docker.go-27- if err != nil {
/project/internal/agentrun/docker.go-28- return nil, fmt.Errorf("docker client: %w", err)
--
/project/internal/agentrun/run.go:247: llmClient := llm.NewClient(llmCfg)
/project/internal/agentrun/run.go-248-
/project/internal/agentrun/run.go-249- systemPrompt := defaultSystemPrompt + identitySection(agent.Name, roster)
--
/project/internal/agentrun/run.go:271: _, err = runLoop(ctx, llmClient, runCtx, messages, hooks)
/project/internal/agentrun/run.go-272- if err != nil {
/project/internal/agentrun/run.go-273- return fmt.Errorf("agent loop: %w", err)
--
/project/internal/agentrun/run.go:281:// streamHooks builds the Hooks a single Run passes to runLoop: every
/project/internal/agentrun/run.go-282-// delta is published live to the hub for connected dashboard viewers,
/project/internal/agentrun/run.go-283-// and once a reasoning/content block or tool call is complete, it's
Tool
π§ grep_search
{"query":"AppendLog|stream ==|\"reasoning\"|\"content\"|\"tool\"","glob":"internal/web/*.go","context_after":2}
Arguments
{"query":"AppendLog|stream ==|\"reasoning\"|\"content\"|\"tool\"","glob":"internal/web/*.go","context_after":2}
Result
No matches found.
Thinking
Go 1.26 β generics are no problem, math/rand/v2 is also no problem.
Let's check the store's schema, look at the log stream, and verify the log stream names in the web template.
Tool
π§ grep_search
{"query":"stream","glob":"internal/store/*.go","context_after":3,"max_results":60}
Arguments
{"query":"stream","glob":"internal/store/*.go","context_after":3,"max_results":60}
Result
No matches found.
Tool
π§ grep_search
{"query":"reasoning|content|tool","glob":"internal/web/templates.go","context_after":2,"max_results":60}
Arguments
{"query":"reasoning|content|tool","glob":"internal/web/templates.go","context_after":2,"max_results":60}
Result
No matches found.
Thinking
Hmm, no matches? That's weird. Let me try searching a different way.
Tool
π§ grep_search
{"query":"AppendLog","context_after":5}
Arguments
{"query":"AppendLog","context_after":5}
Result
/project/internal/web/web_test.go:60: if err := st.AppendLog(ctx, "job-1", "stdout", "hello"); err != nil {
/project/internal/web/web_test.go-61- t.Fatal(err)
/project/internal/web/web_test.go-62- }
/project/internal/web/web_test.go-63-
/project/internal/web/web_test.go-64- rr := httptest.NewRecorder()
/project/internal/web/web_test.go-65- s.Handler().ServeHTTP(rr, httptest.NewRequest(http.MethodGet, "/jobs", nil))
--
/project/internal/web/web_test.go:122: if err := st.AppendLog(ctx, "job-1", "content", "\nNow the remaining tool callers:\n\n\n"); err != nil {
/project/internal/web/web_test.go-123- t.Fatal(err)
/project/internal/web/web_test.go-124- }
/project/internal/web/web_test.go-125- // A block that is nothing but whitespace should not render at all.
/project/internal/web/web_test.go:126: if err := st.AppendLog(ctx, "job-1", "content", " \n\t "); err != nil {
/project/internal/web/web_test.go-127- t.Fatal(err)
/project/internal/web/web_test.go-128- }
/project/internal/web/web_test.go:129: if err := st.AppendLog(ctx, "job-1", "reasoning", " \ninner\nlines\nkept\n "); err != nil {
/project/internal/web/web_test.go-130- t.Fatal(err)
/project/internal/web/web_test.go-131- }
/project/internal/web/web_test.go-132-
/project/internal/web/web_test.go-133- r := httptest.NewRecorder()
/project/internal/web/web_test.go-134- s.Handler().ServeHTTP(r, httptest.NewRequest(http.MethodGet, "/jobs/job-1", nil))
--
/project/internal/store/store.go:325:// AppendLog appends one log line for a job, assigning it the next
/project/internal/store/store.go-326-// sequence number for that job.
/project/internal/store/store.go:327:func (s *Store) AppendLog(ctx context.Context, jobID, stream, line string) error {
/project/internal/store/store.go-328- _, err := s.db.ExecContext(ctx,
/project/internal/store/store.go-329- `INSERT INTO job_logs (job_id, seq, stream, line, ts)
/project/internal/store/store.go-330- VALUES (?, COALESCE((SELECT MAX(seq) + 1 FROM job_logs WHERE job_id = ?), 0), ?, ?, ?)`,
/project/internal/store/store.go-331- jobID, jobID, stream, line, time.Now().UTC())
/project/internal/store/store.go-332- if err != nil {
--
/project/internal/store/store_test.go:146: if err := s.AppendLog(ctx, "job-1", "stdout", line); err != nil {
/project/internal/store/store_test.go-147- t.Fatal(err)
/project/internal/store/store_test.go-148- }
/project/internal/store/store_test.go-149- }
/project/internal/store/store_test.go-150-
/project/internal/store/store_test.go-151- lines, err := s.TailLogs(ctx, "job-1", -1)
--
/project/internal/agentrun/run.go:230: if err := r.store.AppendLog(context.Background(), jobID, stream, line); err != nil {
/project/internal/agentrun/run.go-231- logger.Warn("failed to append log", "error", err)
/project/internal/agentrun/run.go-232- }
/project/internal/agentrun/run.go-233- }
/project/internal/agentrun/run.go-234-
/project/internal/agentrun/run.go-235- runCtx := &runContext{
Tool
π§ grep_search
{"query":"ReasoningDelta","context_after":3}
Arguments
{"query":"ReasoningDelta","context_after":3}
Result
/project/internal/livelog/hub.go:14: ReasoningDelta Type = "reasoning_delta"
/project/internal/livelog/hub.go-15- ReasoningEnd Type = "reasoning_end"
/project/internal/livelog/hub.go-16- ContentStart Type = "content_start"
/project/internal/livelog/hub.go-17- ContentDelta Type = "content_delta"
--
/project/internal/agentrun/loop.go:16: // OnReasoningDelta and OnContentDelta fire with just the newly
/project/internal/agentrun/loop.go-17- // streamed text for the current turn, not the accumulated total.
/project/internal/agentrun/loop.go:18: OnReasoningDelta func(delta string)
/project/internal/agentrun/loop.go-19- OnContentDelta func(delta string)
/project/internal/agentrun/loop.go-20-
/project/internal/agentrun/loop.go-21- // OnTurnEnd fires once per completed streamer round-trip, after the
--
/project/internal/agentrun/loop.go:67: if hooks.OnReasoningDelta != nil && len(msg.ReasoningContent) > len(prevReasoning) {
/project/internal/agentrun/loop.go:68: hooks.OnReasoningDelta(msg.ReasoningContent[len(prevReasoning):])
/project/internal/agentrun/loop.go-69- }
/project/internal/agentrun/loop.go-70- prevReasoning = msg.ReasoningContent
/project/internal/agentrun/loop.go-71-
--
/project/internal/agentrun/run.go:294: OnReasoningDelta: func(delta string) {
/project/internal/agentrun/run.go-295- if !reasoningOpen {
/project/internal/agentrun/run.go-296- reasoningOpen = true
/project/internal/agentrun/run.go-297- reasoningBuf.Reset()
--
/project/internal/agentrun/run.go:302: r.hub.Publish(jobID, livelog.Event{Type: livelog.ReasoningDelta, Text: delta})
/project/internal/agentrun/run.go-303- },
/project/internal/agentrun/run.go-304- OnContentDelta: func(delta string) {
/project/internal/agentrun/run.go-305- if !contentOpen {
Thinking
Let me look at the web template and understand how the log stream gets rendered (to decide which stream name to use for the retry notice).
Tool
π§ grep_search
{"query":"Stream|stream","glob":"internal/web/web.go","context_after":4,"max_results":80}
Arguments
{"query":"Stream|stream","glob":"internal/web/web.go","context_after":4,"max_results":80}
Result
No matches found.
Thinking
Weird β the grep with glob `internal/web/web.go` specified didn't find anything, but the earlier grep did find web.go... no wait, what was found was web_test.go. Let me just read the file directly.
Tool
π§ bash
{"command":"cd /project && grep -n \"stream\" internal/web/web.go | head -40 && echo ---- && grep -n \"stream\\|Stream\" internal/web/templates.go | head -40"}
Arguments
{"command":"cd /project && grep -n \"stream\" internal/web/web.go | head -40 && echo ---- && grep -n \"stream\\|Stream\" internal/web/templates.go | head -40"}
Result
188:// toolBlock is the parsed form of a stream="tool" store.LogLine, for the
200:// maps to exactly one block β grouping/streaming happens upstream, when
273:// handleJobEvents streams job jobID's live output as Server-Sent
295: http.Error(w, "streaming unsupported", http.StatusInternalServerError)
300: w.Header().Set("Content-Type", "text/event-stream")
366: // the auto location.reload() a job's SSE stream triggers when it
368: // streaming vanishing once the page reloads.
----
Thinking
Let me look around lines 188β280 of web.go to check how the log lines are rendered (block types), and determine whether the "system" stream will be displayed.
Tool
π§ read_file
{"path":"internal/web/web.go","start_line":150,"end_line":280}
Arguments
{"path":"internal/web/web.go","start_line":150,"end_line":280}
Result
{"path":"internal/web/web.go","content":"\t}\n\n\ts.render(w, \"jobs\", rows)\n}\n\n// avatarFor returns the Forgejo avatar URL of the agent named username,\n// or \"\" if it can't be resolved (no Forgejo client configured, unknown\n// user, API error). The dashboard must never fail to render because of a\n// missing avatar, so every failure mode degrades to no image. Results are\n// cached per username for avatarCacheTTL so a page refresh doesn't turn\n// into one GetUserInfo call per unique agent.\nfunc (s *Server) avatarFor(username string) string {\n\tif s.fg == nil || username == \"\" {\n\t\treturn \"\"\n\t}\n\n\ts.avatarMu.Lock()\n\n\tif e, ok := s.avatarCache[username]; ok \u0026\u0026 time.Since(e.fetchedAt) \u003c avatarCacheTTL {\n\t\ts.avatarMu.Unlock()\n\n\t\treturn e.url\n\t}\n\n\ts.avatarMu.Unlock()\n\n\tprofile, err := s.fg.AgentProfile(username)\n\tif err != nil {\n\t\treturn \"\"\n\t}\n\n\ts.avatarMu.Lock()\n\ts.avatarCache[username] = avatarCacheEntry{url: profile.AvatarURL, fetchedAt: time.Now()}\n\ts.avatarMu.Unlock()\n\n\treturn profile.AvatarURL\n}\n\n// toolBlock is the parsed form of a stream=\"tool\" store.LogLine, for the\n// template to render as a single collapsed detail.\ntype toolBlock struct {\n\tName string\n\tArguments string\n\tResult string\n\tError bool\n}\n\n// block is one self-contained, already-complete unit of job output: a\n// finished reasoning or assistant-message block, a finished tool call,\n// or a misc system note. Unlike the old flat log view, one store.LogLine\n// maps to exactly one block β grouping/streaming happens upstream, when\n// agentrun persists the row.\ntype block struct {\n\tKind string // \"reasoning\" | \"content\" | \"tool\" | \"system\"\n\tText string\n\tTool *toolBlock\n}\n\nfunc buildBlocks(logs []store.LogLine) []block {\n\tblocks := make([]block, 0, len(logs))\n\n\tfor _, l := range logs {\n\t\tswitch l.Stream {\n\t\tcase \"reasoning\", \"content\":\n\t\t\t// Model output routinely starts/ends with newlines. The block\n\t\t\t// body renders with white-space: pre-wrap, so those would show\n\t\t\t// up as visible blank lines inflating the block's height. Trim\n\t\t\t// them for display (internal newlines are kept) and drop\n\t\t\t// blocks that are nothing but whitespace.\n\t\t\ttext := strings.TrimSpace(l.Line)\n\t\t\tif text == \"\" {\n\t\t\t\tcontinue\n\t\t\t}\n\t\t\tblocks = append(blocks, block{Kind: l.Stream, Text: text})\n\n\t\tcase \"tool\":\n\t\t\tvar entry store.ToolLogEntry\n\t\t\tif err := json.Unmarshal([]byte(l.Line), \u0026entry); err != nil {\n\t\t\t\tblocks = append(blocks, block{Kind: \"system\", Text: l.Line})\n\n\t\t\t\tcontinue\n\t\t\t}\n\n\t\t\tblocks = append(blocks, block{Kind: \"tool\", Tool: \u0026toolBlock{\n\t\t\t\tName: entry.Name,\n\t\t\t\tArguments: entry.Arguments,\n\t\t\t\tResult: entry.Result,\n\t\t\t\tError: entry.Error,\n\t\t\t}})\n\n\t\tdefault:\n\t\t\tblocks = append(blocks, block{Kind: \"system\", Text: l.Line})\n\t\t}\n\t}\n\n\treturn blocks\n}\n\nfunc (s *Server) handleJobDetail(w http.ResponseWriter, r *http.Request) {\n\tid := r.PathValue(\"id\")\n\n\tjob, err := s.store.GetJob(r.Context(), id)\n\tif err != nil {\n\t\thttp.Error(w, \"job not found\", http.StatusNotFound)\n\n\t\treturn\n\t}\n\n\tlogs, err := s.store.TailLogs(r.Context(), id, -1)\n\tif err != nil {\n\t\thttp.Error(w, err.Error(), http.StatusInternalServerError)\n\n\t\treturn\n\t}\n\n\ts.render(w, \"job_detail\", struct {\n\t\tJob store.Job\n\t\tBlocks []block\n\t\tLive bool\n\t\tAvatarURL string\n\t}{job, buildBlocks(logs), job.Status == store.JobPending || job.Status == store.JobRunning, s.avatarFor(job.Agent)})\n}\n\n// handleJobEvents streams job jobID's live output as Server-Sent\n// Events. It only ever carries events for the current, not-yet-persisted\n// block (see livelog.Hub.Checkpoint) β everything already written to\n// the store is rendered once, statically, by handleJobDetail. If the\n// job isn't live, the response just ends immediately and the client\n// falls back to the static page.\nfunc (s *Server) handleJobEvents(w http.ResponseWriter, r *http.Request) {\n\tid := r.PathValue(\"id\")","start_line":150,"end_line":280,"total_lines":374,"truncated":true}
Thinking
The web UI already supports a "system" block kind: any log line whose stream is not "reasoning"/"content"/"tool" renders as a `system` block. So we can log retry notifications to a stream like "system" (or "retry") and it'll be displayed as a system note in the job view.
Now, do we also want live SSE for retries? Let's check the livelog hub's event types and the SSE handler to see if we can add a "retry" event type. Let's look at livelog/hub.go and the template to see how events are handled on the client side.
Tool
π§ read_file
{"path":"internal/livelog/hub.go","start_line":1,"end_line":120}
Arguments
{"path":"internal/livelog/hub.go","start_line":1,"end_line":120}
Result
{"path":"internal/livelog/hub.go","content":"// Package livelog is an in-memory pub/sub hub that lets agentrun publish\n// fine-grained, not-yet-persisted events (reasoning/content deltas, tool\n// calls) for a running job, and lets the web dashboard subscribe to them\n// live over SSE. It carries only the tail of a run that hasn't made it\n// into SQLite yet β see Hub.Checkpoint.\npackage livelog\n\nimport \"sync\"\n\ntype Type string\n\nconst (\n\tReasoningStart Type = \"reasoning_start\"\n\tReasoningDelta Type = \"reasoning_delta\"\n\tReasoningEnd Type = \"reasoning_end\"\n\tContentStart Type = \"content_start\"\n\tContentDelta Type = \"content_delta\"\n\tContentEnd Type = \"content_end\"\n\tTool Type = \"tool\"\n\tStatus Type = \"status\"\n)\n\ntype Event struct {\n\tType Type `json:\"type\"`\n\tText string `json:\"text,omitempty\"`\n\tName string `json:\"name,omitempty\"`\n\tArguments string `json:\"arguments,omitempty\"`\n\tResult string `json:\"result,omitempty\"`\n\tError bool `json:\"error,omitempty\"`\n\tStatus string `json:\"status,omitempty\"`\n}\n\n// subChanBuffer bounds how many events a slow subscriber can lag behind\n// by before new events are dropped for it. Generous for a single-user\n// local dashboard; a dropped event just means a subscriber's browser\n// tab misses a chunk and catches up on the next one, never a hang.\nconst subChanBuffer = 256\n\n// maxBufferedEvents caps the per-job replay buffer as a safety net\n// against unbounded growth if a caller forgets to Checkpoint. Ordinary\n// operation clears the buffer far below this on every persisted block.\nconst maxBufferedEvents = 10000\n\ntype stream struct {\n\tmu sync.Mutex\n\tbuf []Event\n\tsubs map[chan Event]struct{}\n}\n\n// Hub fans out per-job Events to any number of live subscribers.\ntype Hub struct {\n\tmu sync.Mutex\n\tjobs map[string]*stream\n}\n\nfunc NewHub() *Hub {\n\treturn \u0026Hub{jobs: make(map[string]*stream)}\n}\n\nfunc (h *Hub) get(jobID string) *stream {\n\th.mu.Lock()\n\tdefer h.mu.Unlock()\n\n\ts, ok := h.jobs[jobID]\n\tif !ok {\n\t\ts = \u0026stream{subs: make(map[chan Event]struct{})}\n\t\th.jobs[jobID] = s\n\t}\n\n\treturn s\n}\n\n// Publish appends ev to jobID's replay buffer and fans it out to every\n// current subscriber. Never blocks: a subscriber whose channel is full\n// simply misses this event.\nfunc (h *Hub) Publish(jobID string, ev Event) {\n\ts := h.get(jobID)\n\n\ts.mu.Lock()\n\tdefer s.mu.Unlock()\n\n\ts.buf = append(s.buf, ev)\n\tif len(s.buf) \u003e maxBufferedEvents {\n\t\ts.buf = s.buf[len(s.buf)-maxBufferedEvents:]\n\t}\n\n\tfor ch := range s.subs {\n\t\tselect {\n\t\tcase ch \u003c- ev:\n\t\tdefault:\n\t\t}\n\t}\n}\n\n// Checkpoint clears jobID's replay buffer. Call this right after a\n// block derived from buffered events has been durably written to\n// storage, so a subscriber that connects afterward β and therefore\n// already sees that block via the persisted history β is never replayed\n// it a second time over the live channel.\nfunc (h *Hub) Checkpoint(jobID string) {\n\ts := h.get(jobID)\n\n\ts.mu.Lock()\n\tdefer s.mu.Unlock()\n\n\ts.buf = nil\n}\n\n// Subscribe returns the events currently buffered for jobID (i.e. not\n// yet checkpointed) plus a channel of everything published from this\n// point on. The returned cancel func must be called when the caller is\n// done reading, to unregister and close the channel.\nfunc (h *Hub) Subscribe(jobID string) (buffered []Event, ch \u003c-chan Event, cancel func()) {\n\ts := h.get(jobID)\n\n\ts.mu.Lock()\n\tdefer s.mu.Unlock()\n\n\tbuffered = append([]Event(nil), s.buf...)\n","start_line":1,"end_line":120,"total_lines":164,"truncated":true}
Tool
π§ read_file
{"path":"internal/web/web.go","start_line":280,"end_line":374}
Arguments
{"path":"internal/web/web.go","start_line":280,"end_line":374}
Result
{"path":"internal/web/web.go","content":"\tid := r.PathValue(\"id\")\n\n\tjob, err := s.store.GetJob(r.Context(), id)\n\tif err != nil {\n\t\thttp.Error(w, \"job not found\", http.StatusNotFound)\n\n\t\treturn\n\t}\n\n\tif job.Status != store.JobPending \u0026\u0026 job.Status != store.JobRunning {\n\t\treturn\n\t}\n\n\tflusher, ok := w.(http.Flusher)\n\tif !ok {\n\t\thttp.Error(w, \"streaming unsupported\", http.StatusInternalServerError)\n\n\t\treturn\n\t}\n\n\tw.Header().Set(\"Content-Type\", \"text/event-stream\")\n\tw.Header().Set(\"Cache-Control\", \"no-cache\")\n\tw.Header().Set(\"Connection\", \"keep-alive\")\n\tw.WriteHeader(http.StatusOK)\n\n\tbuffered, ch, cancel := s.hub.Subscribe(id)\n\tdefer cancel()\n\n\twrite := func(ev livelog.Event) bool {\n\t\tdata, err := json.Marshal(ev)\n\t\tif err != nil {\n\t\t\treturn true\n\t\t}\n\n\t\tif _, err := fmt.Fprintf(w, \"data: %s\\n\\n\", data); err != nil {\n\t\t\treturn false\n\t\t}\n\n\t\tflusher.Flush()\n\n\t\treturn true\n\t}\n\n\tfor _, ev := range buffered {\n\t\tif !write(ev) {\n\t\t\treturn\n\t\t}\n\t}\n\n\theartbeat := time.NewTicker(15 * time.Second)\n\tdefer heartbeat.Stop()\n\n\tfor {\n\t\tselect {\n\t\tcase ev, ok := \u003c-ch:\n\t\t\tif !ok {\n\t\t\t\treturn\n\t\t\t}\n\n\t\t\tif !write(ev) {\n\t\t\t\treturn\n\t\t\t}\n\n\t\t\tif ev.Type == livelog.Status {\n\t\t\t\treturn\n\t\t\t}\n\n\t\tcase \u003c-heartbeat.C:\n\t\t\tif _, err := fmt.Fprint(w, \": heartbeat\\n\\n\"); err != nil {\n\t\t\t\treturn\n\t\t\t}\n\n\t\t\tflusher.Flush()\n\n\t\tcase \u003c-r.Context().Done():\n\t\t\treturn\n\t\t}\n\t}\n}\n\nfunc (s *Server) render(w http.ResponseWriter, name string, data any) {\n\tw.Header().Set(\"Content-Type\", \"text/html; charset=utf-8\")\n\n\t// Every page here reflects live, fast-changing state (job status, log\n\t// blocks). Without this, browsers are free to serve a stale cached\n\t// copy of the job detail page on a same-URL navigation β including\n\t// the auto location.reload() a job's SSE stream triggers when it\n\t// finishes β which looks exactly like blocks that were visible while\n\t// streaming vanishing once the page reloads.\n\tw.Header().Set(\"Cache-Control\", \"no-store\")\n\n\tif err := s.tmpl.ExecuteTemplate(w, name, data); err != nil {\n\t\thttp.Error(w, err.Error(), http.StatusInternalServerError)\n\t}\n}","start_line":280,"end_line":374,"total_lines":374,"truncated":false}
Thinking
There's a `Status` event type for livelog. Let me check how the dashboard's JS handles events β I'll take a look at templates.go for the SSE client-side JS.
Tool
π§ bash
{"command":"cd /project && grep -n \"reasoning_start\\|content_start\\|status\\|EventSource\\|addEventListener\\|switch\" internal/web/templates.go | head -30"}
Arguments
{"command":"cd /project && grep -n \"reasoning_start\\|content_start\\|status\\|EventSource\\|addEventListener\\|switch\" internal/web/templates.go | head -30"}
Result
762: <p>All agent runs and their current status.</p>
957: var es = new EventSource("/jobs/" + jobID + "/events");
963: switch (ev.type) {
964: case "reasoning_start":
976: case "content_start":
989: case "status":
Thinking
Let me look at the JS in the job details template (around lines 940-1050) and understand the live rendering flow. This is to determine whether we should also emit a live event for retries.
Tool
π§ read_file
{"path":"internal/web/templates.go","start_line":900,"end_line":1080}
Arguments
{"path":"internal/web/templates.go","start_line":900,"end_line":1080}
Result
{"path":"internal/web/templates.go","content":" var body = document.createElement(\"div\");\n body.className = \"block-body\";\n div.appendChild(body);\n log.appendChild(div);\n return body;\n }\n\n function newToolBlock(ev) {\n var details = document.createElement(\"details\");\n details.className = \"block block-tool\" + (ev.error ? \" block-tool-error\" : \"\");\n\n var summary = document.createElement(\"summary\");\n\n var badge = document.createElement(\"span\");\n badge.className = \"tool-badge\";\n badge.textContent = \"Tool\";\n\n var text = document.createElement(\"span\");\n text.className = \"tool-summary-text\";\n\n var name = document.createElement(\"span\");\n name.className = \"tool-name\";\n name.textContent = \"π§ \" + ev.name;\n\n var preview = document.createElement(\"span\");\n preview.className = \"tool-args-preview\";\n preview.textContent = ev.arguments;\n\n text.appendChild(name);\n text.appendChild(preview);\n summary.appendChild(badge);\n summary.appendChild(text);\n details.appendChild(summary);\n\n var body = document.createElement(\"div\");\n body.className = \"block-body\";\n\n var argsLabel = document.createElement(\"div\");\n argsLabel.className = \"tool-section-label\";\n argsLabel.textContent = \"Arguments\";\n var argsPre = document.createElement(\"pre\");\n argsPre.textContent = ev.arguments;\n\n var resultLabel = document.createElement(\"div\");\n resultLabel.className = \"tool-section-label\";\n resultLabel.textContent = \"Result\";\n var resultPre = document.createElement(\"pre\");\n resultPre.textContent = ev.result;\n\n body.appendChild(argsLabel);\n body.appendChild(argsPre);\n body.appendChild(resultLabel);\n body.appendChild(resultPre);\n details.appendChild(body);\n log.appendChild(details);\n }\n\n var es = new EventSource(\"/jobs/\" + jobID + \"/events\");\n\n es.onmessage = function(e) {\n var ev = JSON.parse(e.data);\n var shouldScroll = nearBottom();\n\n switch (ev.type) {\n case \"reasoning_start\":\n reasoningBody = newBlock(\"reasoning\", \"Thinking\");\n break;\n case \"reasoning_delta\":\n if (reasoningBody) reasoningBody.textContent += ev.text;\n break;\n case \"reasoning_end\":\n // Model output routinely starts/ends with newlines; drop them\n // so the finished block matches the static (trimmed) view.\n if (reasoningBody) reasoningBody.textContent = reasoningBody.textContent.trim();\n reasoningBody = null;\n break;\n case \"content_start\":\n contentBody = newBlock(\"content\", null);\n break;\n case \"content_delta\":\n if (contentBody) contentBody.textContent += ev.text;\n break;\n case \"content_end\":\n if (contentBody) contentBody.textContent = contentBody.textContent.trim();\n contentBody = null;\n break;\n case \"tool\":\n newToolBlock(ev);\n break;\n case \"status\":\n es.close();\n location.reload();\n return;\n }\n\n if (shouldScroll) window.scrollTo(0, document.body.scrollHeight);\n };\n })();\n \u003c/script\u003e\n {{end}}\n\u003c/div\u003e\n\u003c/body\u003e\n\u003c/html\u003e\n{{end}}\n`","start_line":900,"end_line":1004,"total_lines":1004,"truncated":false}
Thinking
Now I understand the live rendering. The live JS handles the following: reasoning_start/delta/end, content_start/delta/end, tool, status (β reload).
For retry notifications, there are two options:
1. **Persist only**: `logAppend("system", ...)` β shown as a system block in the static view after the page reloads. But during the retry wait (which could last several minutes), the live view won't show anything β the job looks frozen.
2. **Persist + live**: Also emit a livelog event. To render it live, the JS needs a new case. We could reuse the "tool" event type⦠no, that's hacky. We could add a new event type "retry" (or "system") and a JS case that renders a system block.
Hmm, scope check. The core request is retry logic. Live dashboard rendering of retry notifications is a nice-to-have. But without it, a job that stalls for 4 minutes on retries looks dead in the dashboard (no output), which is exactly the situation where visibility matters most.
Let me think about the cleanest, most minimal approach:
- In the OnRetry hook: reset the buffer (discard partial content from the failed attempt), emit a `livelog` event, and `logAppend` a system line.
- Add a new livelog Type `Retry` = "retry" (or reuse `Status`? no β status triggers a page reload, which is wrong).
- Add a JS case "retry" that renders a system block with text like "LLM endpoint failed (attempt 2): ... β retrying".
Wait, but there's a subtle point about the replay buffer and Checkpoint: the retry event is added to the replay buffer. If the retried turn succeeds, OnTurnEnd persists the new block and Checkpoints β which clears the buffer including the retry event. So a subscriber who connects after the checkpoint won't see the retry event live, but they will see the persisted "system" log line in the static view.
But wait β there's a subtle double-rendering issue: the retry event is in the replay buffer, and it's also persisted as a "system" line. A subscriber who connected before the retry sees it live. A subscriber who connects after the checkpoint sees the persisted line. A subscriber who connects between the retry event and the checkpointβ¦ sees the replayed retry event live, and after the next page reload (triggered by job completion) sees the persisted line β but the live view is replaced by the static view on reload (status event β location.reload), so there's no double rendering on the same page. This matches the existing design (the same pattern as tool events: persisted + buffered, checkpointed after persist).
Hmm wait, actually there's an existing subtlety: tool events are published to the hub and then persisted, and then Checkpoint. So the same event is in the buffer until the checkpoint. A subscriber who joins in between sees the tool event live, and thenβ¦ the page doesn't reload until the job ends. When the job ends, status event β reload β static view β the tool block appears from the persisted log. The live tool block disappears on reload. No double rendering. OK, so the same pattern works for retry events.
Now β what about the partial content from the failed attempt? In the live view, a "Thinking" block (reasoningBody) with partial text is already rendered. When the retry happens, the live JS still has the old reasoningBody. The retried stream emits reasoning_start again β newBlock creates a new block, and the old partial block remains in the DOM with its partial text. Hmm, that's a cosmetic issue: the user sees a truncated "Thinking" block followed by a full block.
To handle this cleanly in the live view, the "retry" event should tell the JS to remove the current partial block (reasoningBody/contentBody) and reset them. Let me add that: in the JS case "retry": remove the current reasoningBody/contentBody elements from the DOM (if they exist), reset to null, and render a system block with the retry message.
That's clean. And on the server side, the OnRetry hook resets the buffer/state, so the persisted log matches the live view (no partial block persisted).
Wait, one more check: the OnTurnEnd hook persists the buffer if reasoningOpen/contentOpen. If a turn fails mid-stream and we retry, OnTurnEnd doesn't fire (only on success). On the retry's success, OnTurnEnd fires with the fresh buffer (which was reset on OnRetry) β only the retried content is persisted.
But hmm β what if the partial content from the failed attempt was already checkpointed? No β the checkpoint only happens on OnTurnEnd/OnTool, and those don't fire mid-turn. So the partial content is only in the in-memory buffer + hub replay. Resetting on OnRetry discards it completely.
Now, the "system" log line: which stream name? buildBlocks treats any stream that isn't reasoning/content/tool as a "system" block. So `logAppend("system", line)` works. Let me check if there's precedent for a "system" streamβ¦ web_test.go uses "stdout". Both render as system blocks. I'll use "system" β it's clearer.
Hmm, wait. Actually, let me reconsider whether we should persist the retry notification at all, or just log it via slog. The job log in the dashboard is the primary UX for run visibility. abrander's failure message ("zoo agent run failed: ...") is posted as an issue comment by some other mechanism (probably a scheduler). For a successful run with retries, the dashboard's job view is where humans look. Persisting a "system" line is the right call.
Now, let me also double-check that the scheduler failure path posts the error as a comment β the "zoo agent run failed: agent loop: chat completion: status: 429" comment. Let me look at scheduler.go to see how the failure surfaces.
Let me check the scheduler to see how run failures surface (the "zoo agent run failed" comments).
Tool
π§ read_file
{"path":"internal/scheduler/scheduler.go","start_line":1,"end_line":200}
Arguments
{"path":"internal/scheduler/scheduler.go","start_line":1,"end_line":200}
Result
{"path":"internal/scheduler/scheduler.go","content":"// Package scheduler resolves incoming Forgejo events to configured\n// agents and runs them, bounded by max_live_agents.\npackage scheduler\n\nimport (\n\t\"context\"\n\t\"errors\"\n\t\"fmt\"\n\t\"log/slog\"\n\t\"sync\"\n\n\t\"github.com/google/uuid\"\n\n\t\"github.com/abrander/zoo/internal/config\"\n\t\"github.com/abrander/zoo/internal/forgejo\"\n\t\"github.com/abrander/zoo/internal/livelog\"\n\t\"github.com/abrander/zoo/internal/store\"\n)\n\n// forgejoActions is the narrow slice of Client the scheduler needs for\n// its own failure-reporting side effects (defined here, not in\n// internal/forgejo, so tests can inject a fake).\ntype forgejoActions interface {\n\tCreateIssueComment(owner, repo string, index int64, body string) error\n\tAddLabel(owner, repo string, index int64, name string) error\n}\n\n// FailureLabel is applied to the triggering issue/PR, alongside a\n// comment, whenever an agent run fails or times out.\nconst FailureLabel = \"zoo:failed\"\n\n// Runner runs a single agent invocation to completion. Implemented by\n// internal/agentrun.Run; a narrow interface here so the scheduler is\n// testable without Docker.\ntype Runner interface {\n\tRun(ctx context.Context, jobID string, agent config.AgentConfig, llm config.LLM, dockerImage string, ev forgejo.Event) error\n}\n\ntype Scheduler struct {\n\tcfg *config.Config\n\tstore *store.Store\n\tforgejo forgejoActions\n\trunner Runner\n\thub *livelog.Hub\n\tlogger *slog.Logger\n\n\tsem chan struct{}\n\twg sync.WaitGroup\n}\n\nfunc New(cfg *config.Config, st *store.Store, fg forgejoActions, runner Runner, hub *livelog.Hub, logger *slog.Logger) *Scheduler {\n\treturn \u0026Scheduler{\n\t\tcfg: cfg,\n\t\tstore: st,\n\t\tforgejo: fg,\n\t\trunner: runner,\n\t\thub: hub,\n\t\tlogger: logger,\n\t\tsem: make(chan struct{}, cfg.Environment.MaxLive),\n\t}\n}\n\n// resolveAgent returns the name of the agent that should handle ev, if\n// any. Two kinds resolve dynamically: issue:assigned to the agent whose\n// config label matches the Forgejo assignee's username, and pr:review to\n// the agent whose config label matches the pull request author's\n// username (the agent that opened the PR reacts to the review of it).\n// Every other event kind uses the static event-\u003eagent mapping from\n// config.\nfunc resolveAgent(cfg *config.Config, ev forgejo.Event) (string, bool) {\n\tswitch ev.Kind {\n\tcase config.EventIssueAssigned:\n\t\tif _, ok := cfg.AgentByName(ev.Assignee); ok {\n\t\t\treturn ev.Assignee, true\n\t\t}\n\n\t\treturn \"\", false\n\n\tcase config.EventPRReview:\n\t\tif _, ok := cfg.AgentByName(ev.PRAuthor); ok {\n\t\t\treturn ev.PRAuthor, true\n\t\t}\n\n\t\treturn \"\", false\n\t}\n\n\treturn cfg.EventAgent(ev.Kind)\n}\n\n// Run consumes events until ctx is canceled or the channel closes,\n// dispatching each to its resolved agent and blocking on the\n// max_live_agents semaphore before starting a run.\nfunc (s *Scheduler) Run(ctx context.Context, events \u003c-chan forgejo.Event) {\n\tfor {\n\t\tselect {\n\t\tcase \u003c-ctx.Done():\n\t\t\treturn\n\n\t\tcase ev, ok := \u003c-events:\n\t\t\tif !ok {\n\t\t\t\treturn\n\t\t\t}\n\n\t\t\ts.handle(ctx, ev)\n\t\t}\n\t}\n}\n\nfunc (s *Scheduler) handle(ctx context.Context, ev forgejo.Event) {\n\tagentName, ok := resolveAgent(s.cfg, ev)\n\tif !ok {\n\t\ts.logger.Debug(\"no agent resolved for event, dropping\", \"kind\", ev.Kind, \"owner\", ev.Owner, \"repo\", ev.Repo, \"index\", ev.Index)\n\t\treturn\n\t}\n\n\t// An agent's own actions (e.g. a comment posted via the `comment`\n\t// tool, authenticated with its own per-agent token) can themselves\n\t// show up as new events. Don't let an agent trigger itself off its\n\t// own activity β that's a self-reinforcing loop, not new work.\n\tif ev.Author != \"\" \u0026\u0026 ev.Author == agentName {\n\t\ts.logger.Debug(\"dropping event authored by the agent it would trigger\", \"kind\", ev.Kind, \"agent\", agentName, \"owner\", ev.Owner, \"repo\", ev.Repo, \"index\", ev.Index)\n\t\treturn\n\t}\n\n\tagent, ok := s.cfg.AgentByName(agentName)\n\tif !ok {\n\t\ts.logger.Error(\"resolved agent not declared in config\", \"agent\", agentName)\n\t\treturn\n\t}\n\n\tllm, ok := s.cfg.LLMByName(agent.LLM)\n\tif !ok {\n\t\ts.logger.Error(\"agent references undeclared llm\", \"agent\", agentName, \"llm\", agent.LLM)\n\t\treturn\n\t}\n\n\tjobID := uuid.NewString()\n\n\tif err := s.store.CreateJob(ctx, store.Job{\n\t\tID: jobID,\n\t\tEventKind: ev.Kind,\n\t\tAgent: agentName,\n\t\tOwner: ev.Owner,\n\t\tRepo: ev.Repo,\n\t\tIssueIndex: ev.Index,\n\t\tTitle: ev.Title,\n\t}); err != nil {\n\t\ts.logger.Error(\"failed to record job\", \"job\", jobID, \"error\", err)\n\t\treturn\n\t}\n\n\tselect {\n\tcase s.sem \u003c- struct{}{}:\n\n\tcase \u003c-ctx.Done():\n\t\treturn\n\t}\n\n\ts.wg.Add(1)\n\n\tgo func() {\n\t\tdefer s.wg.Done()\n\t\tdefer func() { \u003c-s.sem }()\n\n\t\ts.run(ctx, jobID, agent, llm, ev)\n\t}()\n}\n\nfunc (s *Scheduler) run(ctx context.Context, jobID string, agent config.AgentConfig, llm config.LLM, ev forgejo.Event) {\n\tlogger := s.logger.With(\"job\", jobID, \"agent\", agent.Name, \"event\", ev.Kind, \"owner\", ev.Owner, \"repo\", ev.Repo, \"index\", ev.Index)\n\n\t// Job status writes use a context detached from ctx, not ctx itself:\n\t// ctx is canceled on daemon shutdown to unwind the in-flight run, and\n\t// an already-canceled ctx would make these UPDATEs fail instantly,\n\t// leaving the job stuck at \"running\" forever even though the process\n\t// has exited.\n\tif err := s.store.MarkJobStarted(context.Background(), jobID); err != nil {\n\t\tlogger.Error(\"failed to mark job started\", \"error\", err)\n\t}\n\n\tlogger.Info(\"agent run starting\")\n\n\terr := s.runner.Run(ctx, jobID, agent, llm, s.cfg.Environment.DockerImage, ev)\n\n\tstatus := store.JobSucceeded\n\terrMsg := \"\"\n\n\tif err != nil {\n\t\terrMsg = err.Error()\n\t\t// The daemon-wide ctx passed in here isn't what times a run out\n\t\t// (agentrun.Runner.Run applies its own per-run deadline\n\t\t// internally); a timed-out run surfaces as a wrapped\n\t\t// context.DeadlineExceeded in the returned error instead.\n\t\tif errors.Is(err, context.DeadlineExceeded) {\n\t\t\tstatus = store.JobTimedOut\n\t\t} else {\n\t\t\tstatus = store.JobFailed\n\t\t}\n\n\t\tlogger.Error(\"agent run failed\", \"status\", status, \"error\", err)","start_line":1,"end_line":200,"total_lines":252,"truncated":true}
Thinking
OK, I now have the full picture. Time to design and implement.
## Implementation Plan
### 1. `internal/llm/retry.go` (new file)
```go
package llm
import (
"context"
"errors"
"fmt"
"math/rand/v2"
"net/http"
"strconv"
"time"
)
// APIError is a non-2xx response from the LLM endpoint.
type APIError struct {
Status int
Body string
// retryAfter is the raw Retry-After header, if the server sent one.
retryAfter string
}
func (e *APIError) Error() string {
return fmt.Sprintf("status: %d: %s", e.Status, e.Body)
}
// RetryAfter returns the server-requested retry delay, parsed from the
// Retry-After header (delta-seconds or HTTP date), if present.
func (e *APIError) RetryAfter() (time.Duration, bool) {
...
}
// IsRetryable reports whether err is a transient failure worth
// retrying.
//
// HTTP responses: 408, 429, and 5xx are transient; other statuses
// (400, 401, 403, ...) are client errors that retrying the same
// request will not fix. Anything that is not an APIError β connection
// failures, timeouts, a stream that dropped mid-response β is treated
// as transient. A canceled or deadline-exceeded context is never
// retried: the run is over, and waiting would only delay reporting it.
func IsRetryable(err error) bool {
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
return false
}
var apiErr *APIError
if errors.As(err, &apiErr) {
return apiErr.Status == http.StatusRequestTimeout ||
apiErr.Status == http.StatusTooManyRequests ||
apiErr.Status >= 500
}
return true
}
// Retry configures Do.
type Retry struct {
// MaxAttempts is the total number of attempts, including the first.
MaxAttempts int
// BaseDelay is the wait before the second attempt; each subsequent
// wait doubles, up to MaxDelay.
BaseDelay time.Duration
// MaxDelay caps any single wait.
MaxDelay time.Duration
// BeforeRetry, if non-nil, is called with the 1-based index of the
// attempt that just failed and its error, before the backoff wait.
BeforeRetry func(attempt int, err error)
// Sleep, if non-nil, is used instead of a context-bound time.Sleep
// (tests).
Sleep func(ctx context.Context, d time.Duration) error
}
// DefaultRetry is the policy the agent loop uses: enough attempts and
// headroom to ride out typical LLM-gateway cooldowns (rate limits,
// model restarts) without stalling a run for more than a few minutes.
func DefaultRetry() Retry {
return Retry{
MaxAttempts: 10,
BaseDelay: 2 * time.Second,
MaxDelay: 60 * time.Second,
}
}
// Do runs fn, retrying it while it fails with a transient error
// (IsRetryable) and attempts remain. Waits between attempts grow
// exponentially (BaseDelay * 2^n, capped at MaxDelay) with jitter, and
// honor a Retry-After hint carried by an *APIError when it exceeds the
// computed backoff. ctx bounds both the attempts and the waits.
func Do[T any](ctx context.Context, r Retry, fn func() (T, error)) (T, error) {
if r.MaxAttempts < 1 {
r.MaxAttempts = 1
}
if r.BaseDelay <= 0 {
r.BaseDelay = 2 * time.Second
}
if r.MaxDelay < r.BaseDelay {
r.MaxDelay = r.BaseDelay
}
sleep := r.Sleep
if sleep == nil {
sleep = func(ctx context.Context, d time.Duration) error {
t := time.NewTimer(d)
defer t.Stop()
select {
case <-ctx.Done():
return ctx.Err()
case <-t.C:
return nil
}
}
}
var zero T
var lastErr error
for attempt := 1; ; attempt++ {
if err := ctx.Err(); err != nil {
return zero, err
}
result, err := fn()
if err == nil {
return result, nil
}
lastErr = err
if attempt >= r.MaxAttempts || !IsRetryable(err) {
if attempt > 1 {
return zero, fmt.Errorf("after %d attempts: %w", attempt, err)
}
return zero, err
}
if r.BeforeRetry != nil {
r.BeforeRetry(attempt, err)
}
delay := backoffDelay(r.BaseDelay, r.MaxDelay, attempt)
if ra, ok := retryAfter(err); ok && ra > delay {
delay = ra
}
if err := sleep(ctx, delay); err != nil {
return zero, err
}
}
_ = lastErr // hmm, not needed
}
```
Wait, `lastErr` is unused β remove it.
`backoffDelay(base, max time.Duration, attempt int) time.Duration`:
```go
// backoffDelay is the wait before retry number attempt (1-based):
// base * 2^(attempt-1) capped at max, with equal jitter β half the
// delay is fixed and half is random β so concurrent runs don't all
// retry the same endpoint in lockstep.
func backoffDelay(base, max time.Duration, attempt int) time.Duration {
d := base
for i := 1; i < attempt && d < max; i++ {
d *= 2
}
if d > max || d < 0 {
d = max
}
half := d / 2
return half + time.Duration(rand.Int63n(int64(half)+1))
}
```
Wait, check: attempt=1 β loop doesn't run (i=1 < 1 is false) β d = base. Good. attempt=2 β d = base*2. Good. Overflow: d *= 2 could overflow if attempt is huge (e.g. 2^63). With MaxAttempts=10 and base=2s it can't happen, but defensively: `d < 0` after overflow β set to max. Also `d < max` in the loop condition stops doubling once we hit max. If base > max... we normalize MaxDelay >= BaseDelay above, so no issue.
`retryAfter`:
```go
func retryAfter(err error) (time.Duration, bool) {
var apiErr *APIError
if errors.As(err, &apiErr) {
return apiErr.RetryAfter()
}
return 0, false
}
```
`APIError.RetryAfter`:
```go
func (e *APIError) RetryAfter() (time.Duration, bool) {
if e.retryAfter == "" {
return 0, false
}
if secs, err := strconv.Atoi(e.retryAfter); err == nil {
if secs < 0 {
return 0, false
}
return time.Duration(secs) * time.Second, true
}
if t, err := http.ParseTime(e.retryAfter); err == nil {
if d := time.Until(t); d > 0 {
return d, true
}
}
return 0, false
}
```
### 2. `internal/llm/streamer.go` β return *APIError
```go
func NewStreamer(resp *http.Response) (*Streamer, error) {
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
resp.Body.Close()
return nil, &APIError{
Status: resp.StatusCode,
Body: string(body),
retryAfter: resp.Header.Get("Retry-After"),
}
}
...
}
```
Error text is unchanged: "status: 429: {...}".
### 3. `internal/agentrun/loop.go` β retry the turn
Add `OnRetry` to Hooks, and restructure runLoop:
```go
// Hooks lets a caller observe a runLoop invocation live, ...
type Hooks struct {
OnReasoningDelta func(delta string)
OnContentDelta func(delta string)
OnTurnEnd func()
OnTool func(name, arguments, result string, toolErr bool)
// OnRetry fires before the loop waits to retry a failed turn:
// attempt is the 1-based index of the attempt that just failed, and
// err its error. Deltas from the failed attempt have already been
// delivered (OnReasoningDelta/OnContentDelta fired for them), but
// OnTurnEnd did not; implementers should discard the failed
// attempt's partial output so the retried turn is recorded fresh.
OnRetry func(attempt int, err error)
}
```
runLoop:
```go
// runLoop is a headless port of ../a's App.generate(): send messages +
// tool defs, get a completion, run any tool_calls and append their
// results, repeat until a plain finish or ctx is done.
//
// Each model turn β the request and the consumption of its stream β is
// retried with llm.DefaultRetry while it fails with a transient error
// (rate limit, server error, dropped connection). The messages slice
// is only extended after a turn completes, so a retried turn is always
// a clean re-send of the same request.
func runLoop(ctx context.Context, client *llm.Client, toolsCtx tools.Context, messages []llm.Message, hooks Hooks) ([]llm.Message, error) {
retry := llm.DefaultRetry()
retry.BeforeRetry = func(attempt int, err error) {
if hooks.OnRetry != nil {
hooks.OnRetry(attempt, err)
}
}
for {
if err := ctx.Err(); err != nil {
return messages, err
}
var (
completion *llm.ChatCompletion
prevContent string
prevReason string
)
err := llm.Do(ctx, retry, func() error {
// A retried turn re-streams from scratch: reset the delta
// cursors so hooks only see the new attempt's suffixes.
prevContent, prevReason = "", ""
completion = nil
streamer, err := client.StreamChatCompletion(ctx, &llm.ChatCompletionRequest{
Messages: messages,
Stream: true,
Tools: tools.All(),
})
if err != nil {
return fmt.Errorf("chat completion: %w", err)
}
for {
c, err := streamer.Get()
if err == io.EOF {
break
}
if err != nil {
return fmt.Errorf("stream completion: %w", err)
}
completion = c
if len(c.Choices) > 0 {
msg := c.Choices[0].Message
if hooks.OnReasoningDelta != nil && len(msg.ReasoningContent) > len(prevReason) {
hooks.OnReasoningDelta(msg.ReasoningContent[len(prevReason):])
}
prevReason = msg.ReasoningContent
if hooks.OnContentDelta != nil && len(msg.Content) > len(prevContent) {
hooks.OnContentDelta(msg.Content[len(prevContent):])
}
prevContent = msg.Content
}
}
if completion == nil || len(completion.Choices) == 0 {
return fmt.Errorf("model returned an empty completion")
}
return nil
})
if err != nil {
return messages, err
}
if hooks.OnTurnEnd != nil {
hooks.OnTurnEnd()
}
choice := completion.Choices[0]
messages = append(messages, choice.Message)
if choice.FinishReason != "tool_calls" || len(choice.Message.ToolCalls) == 0 {
return messages, nil
}
for _, tc := range choice.Message.ToolCalls {
... unchanged ...
}
}
}
```
Hmm wait, one issue: the "model returned an empty completion" error is now inside Do β IsRetryable β true (not an APIError, not a ctx error) β it gets retried. Is that OK? Yes, I decided that earlier. A 200 response with no choices is anomalous, and a retry is reasonable.
But wait, there's a subtlety with the error wrapping in Do: "after %d attempts: %w". So the final error from runLoop would be, e.g.:
"agent loop: after 10 attempts: chat completion: status: 502: engine not healthy"
Good, informative.
Another subtlety: the `err` returned from the closure is wrapped by Do only if attempt > 1. If the first attempt fails with a non-retryable error (e.g. 400), Do returns it as-is β "agent loop: chat completion: status: 400: ..." β same as today.
### 4. `internal/agentrun/run.go` β implement OnRetry in streamHooks
```go
OnRetry: func(attempt int, err error) {
// Discard the failed attempt's partial output so the retried turn
// is recorded as a fresh block, and note the retry in the job log.
reasoningBuf.Reset()
contentBuf.Reset()
reasoningOpen = false
contentOpen = false
line := fmt.Sprintf("llm endpoint failed (attempt %d): %s β retrying", attempt, err)
logAppend("system", line)
r.hub.Publish(jobID, livelog.Event{Type: livelog.Retry, Text: line})
},
```
Hmm wait, but `err` here is the raw error from the closure ("chat completion: status: 502: engine not healthy"). For the user-facing line, maybe I should strip the "chat completion: " prefix? The message "llm endpoint failed (attempt 1): chat completion: status: 502: engine not healthy β retrying" is a bit noisy. Options: pass a cleaner message. The BeforeRetry callback receives the wrapped error. I could use errors.As to extract the APIError and format it: "status 502: engine not healthy". Hmm.
Simpler: in the OnRetry hook, just use err.Error() as-is. It's a log line, and the prefix is informative (distinguishes request failure from mid-stream failure). Keep it.
Wait, actually β do I want to include the full error body? A litellm 429 body is a huge JSON blob (the one in the issue comment is ~1KB). Persisting that into the job log for every failed attempt is noisy. Hmm. For APIError, the body can be long. Let me truncate the line to a reasonable length, say 200 chars? Or just use the status code + first line?
Let me think: for APIError, a compact form is "status 429" + first line of body. For transport errors, the full error is short anyway ("Post https://...: EOF").
I could add a helper method to APIError: `Short() string` that returns "status 429: <first line of body, truncated>". Hmm, or handle it in the hook:
```go
func retryLine(err error) string {
var apiErr *llm.APIError
if errors.As(err, &apiErr) {
body := apiErr.Body
if i := strings.IndexByte(body, '\n'); i > 0 {
body = body[:i]
}
if len(body) > 160 {
body = body[:160] + "β¦"
}
return fmt.Sprintf("status %d: %s", apiErr.Status, body)
}
return err.Error()
}
```
Hmm, this is getting complicated. Let me simplify: put the compact formatting on APIError itself as a method in the llm package:
Actually β cleaner idea: give APIError a `Summary()` method... hmm, naming. Let me think about what's most natural in Go.
Alternative: keep it simple β just use err.Error() and truncate the whole line to, say, 300 chars in the hook. The 429 body in the issue is one long line (JSON). Truncated at 300 chars it would show "status: 429: {\"error\":{\"message\":\"litellm.BadRequestError: OpenAIException - Prefill server error (400 Bad Request): {\"object\":\"error\",\"message\":\"Tool 0 function has invalid 'parameters' schema: '' is not valid under any of the given schemas\n\nFailed validating..." β hmm, that reveals the real problem (invalid tool schema!) but is truncated mid-sentence.
Hmm wait, actually that 429 is interesting β the underlying error is "Tool 0 function has invalid 'parameters' schema: '' is not valid under any of the given schemas" β a 400 wrapped in a 429 by litellm (cooldown after the 400). Retrying that won't fix it (the tool schema is broken β probably a tool with no parameters... `Parameters` with `Type: ""`? Actually, looking at types.go: `Parameters{Type string json:"type"}` β if a tool has no properties, Type stays "" and gets serialized as `"type":""`... wait, no omitempty on Type. And `Function.Parameters Parameters json:"parameters,omitempty"` β a struct value is never omitted. So a tool with no properties sends `"parameters":{"type":""}` β invalid schema β 400! That's the bug in the issue's 429 comment!
Wait, is that in scope? The issue is about retry logic. But the 429 in the comment is caused by an invalid tool schema (a bug in zoo's tool definitions), not by rate limiting. Retrying won't fix it β it'll retry 10 times and fail. Hmm!
Let me verify: which tools have no properties? Let me look at tools/registry.go and the tool definitions. For example, `close_issue` probably takes no arguments. Let me check.
If a tool with no parameters sends `"parameters": {"type": ""}`, then any model call with that tool in the list gets a 400 from strict validators (litellm validates the schema). That's a real bug that causes the exact failure in the comment. But β is it in scope for "retry logic"? The issue is specifically about retry logic when the endpoint is down. The 429/502 errors are what triggered the issue. The 429's root cause is the invalid schema, but...
Hmm, wait. Let me re-read the 429: "Prefill server error (400 Bad Request): ... Tool 0 function has invalid 'parameters' schema: '' is not valid under any of the given schemas ... Received Model Group=Qwen/Qwen3.8-27B ... Error doing the fallback: No deployments available for selected model, Try again in 152.0 seconds."
So: the request was rejected with 400 (invalid tool schema), litellm put the model in cooldown, and returned 429 "try again in 152 seconds". Even if we retry after 152s, the same 400 β 429 cycle repeats. So retry logic alone won't fix this specific failure β but it will fix the 502 "engine not healthy" case, and transient 429s/5xx in general.
Should I also fix the tool schema bug? It's adjacent, but it's the root cause of the observed 429. As a senior architect, I'd point it out. Let me verify whether tools with no parameters actually produce `"parameters":{"type":""}`.
Let me look at the tool definitions.
Now let me look at the tools registry β the 429 in the issue comments mentions an invalid `'parameters'` schema with empty `type`, which may be a real bug in how parameterless tools are serialized.
Tool
π§ read_file
{"path":"internal/tools/registry.go","start_line":1,"end_line":120}
Arguments
{"path":"internal/tools/registry.go","start_line":1,"end_line":120}
Result
{"path":"internal/tools/registry.go","content":"package tools\n\nimport (\n\t\"encoding/json\"\n\t\"fmt\"\n\n\t\"github.com/abrander/zoo/internal/llm\"\n)\n\ntype toolFunc[P any] func(Context, P) (string, error)\n\ntype tool struct {\n\ttool llm.Tool\n\tcall func(ctx Context, args string) (string, error)\n}\n\nvar registry = map[string]tool{}\n\n// All lists all known tools.\nfunc All() []llm.Tool {\n\tts := make([]llm.Tool, 0, len(registry))\n\n\tfor _, t := range registry {\n\t\tts = append(ts, t.tool)\n\t}\n\n\treturn ts\n}\n\n// Register registers a tool function with the given definition. The\n// registry is a stateless name-\u003efunc map shared across concurrent runs;\n// all per-run state lives in the Context passed to Call.\nfunc Register[P any](def llm.Tool, fun toolFunc[P]) {\n\tif def.Type != \"function\" {\n\t\tpanic(\"tool must be a function\")\n\t}\n\n\tif def.Function.Name == \"\" {\n\t\tpanic(\"tool function must have a name\")\n\t}\n\n\tif _, found := registry[def.Function.Name]; found {\n\t\tpanic(fmt.Sprintf(\"tool with name %q already registered\", def.Function.Name))\n\t}\n\n\tregistry[def.Function.Name] = tool{\n\t\ttool: def,\n\t\tcall: func(ctx Context, args string) (string, error) {\n\t\t\tvar params P\n\n\t\t\tif err := json.Unmarshal([]byte(args), \u0026params); err != nil {\n\t\t\t\treturn \"\", err\n\t\t\t}\n\n\t\t\treturn fun(ctx, params)\n\t\t},\n\t}\n}\n\n// Call invokes a registered tool by name with the given (raw JSON)\n// arguments, against the given run's Context.\nfunc Call(ctx Context, fcall llm.FunctionCall) (string, error) {\n\tt, found := registry[fcall.Name]\n\tif !found {\n\t\treturn \"\", fmt.Errorf(\"unknown function: %s\", fcall.Name)\n\t}\n\n\treturn t.call(ctx, fcall.Arguments)\n}","start_line":1,"end_line":69,"total_lines":69,"truncated":false}
Tool
π§ grep_search
{"query":"NewTool\\(","glob":"internal/tools/*.go","context_after":3}
Arguments
{"query":"NewTool\\(","glob":"internal/tools/*.go","context_after":3}
Result
No matches found.
Thinking
grep with glob specified didn't work again β there seems to be some issue with the glob parameter. Let me try using bash's grep instead.
Tool
π§ bash
{"command":"cd /project && grep -rn \"NewTool(\" internal/tools/ | head -30"}
Arguments
{"command":"cd /project && grep -rn \"NewTool(\" internal/tools/ | head -30"}
Result
internal/tools/forgejo.go:15: tool := llm.NewTool(
internal/tools/forgejo.go:40: tool := llm.NewTool(
internal/tools/forgejo.go:65: tool := llm.NewTool(
internal/tools/forgejo.go:88: tool := llm.NewTool(
internal/tools/forgejo.go:115: tool := llm.NewTool(
internal/tools/forgejo.go:133: tool := llm.NewTool(
internal/tools/forgejo.go:155: tool := llm.NewTool(
internal/tools/forgejo.go:175: tool := llm.NewTool(
internal/tools/forgejo.go:191: tool := llm.NewTool(
internal/tools/grep_search.go:22: tool := llm.NewTool(
internal/tools/remove_file.go:15: tool := llm.NewTool(
internal/tools/write_file.go:15: tool := llm.NewTool(
internal/tools/move_file.go:16: tool := llm.NewTool(
internal/tools/read_file.go:29: tool := llm.NewTool(
internal/tools/edit-file.go:19: tool := llm.NewTool(
internal/tools/list_files.go:16: tool := llm.NewTool(
internal/tools/bash.go:19: tool := llm.NewTool(
Thinking
Let me check one of the tools that likely has no parameters (e.g., close_issue in forgejo.go) and see whether Parameters.Type gets left as "".
Tool
π§ read_file
{"path":"internal/tools/forgejo.go","start_line":1,"end_line":120}
Arguments
{"path":"internal/tools/forgejo.go","start_line":1,"end_line":120}
Result
{"path":"internal/tools/forgejo.go","content":"package tools\n\nimport (\n\t\"fmt\"\n\t\"strings\"\n\n\t\"github.com/abrander/zoo/internal/llm\"\n)\n\ntype commentParams struct {\n\tBody string `json:\"body\"`\n}\n\nfunc init() {\n\ttool := llm.NewTool(\n\t\t\"comment\",\n\t\t\"Post a comment on the issue or pull request that triggered this run.\")\n\n\ttool.AddStringProperty(\"body\", \"The comment text, in Markdown\", true)\n\n\tRegister(tool, comment)\n}\n\nfunc comment(ctx Context, params commentParams) (string, error) {\n\tif err := ctx.Forgejo().Comment(params.Body); err != nil {\n\t\treturn \"\", err\n\t}\n\n\treturn \"Comment posted\", nil\n}\n\ntype openPullRequestParams struct {\n\tHead string `json:\"head\"`\n\tBase string `json:\"base\"`\n\tTitle string `json:\"title\"`\n\tBody string `json:\"body\"`\n}\n\nfunc init() {\n\ttool := llm.NewTool(\n\t\t\"open_pull_request\",\n\t\t\"Open a pull request from a branch you've pushed into the repository's default branch (or another base branch).\")\n\n\ttool.AddStringProperty(\"head\", \"The branch containing your changes (usually the run's own branch)\", true)\n\ttool.AddStringProperty(\"base\", \"The branch to merge into, e.g. the repository's default branch\", true)\n\ttool.AddStringProperty(\"title\", \"The pull request title\", true)\n\ttool.AddStringProperty(\"body\", \"The pull request description, in Markdown\", true)\n\n\tRegister(tool, openPullRequest)\n}\n\nfunc openPullRequest(ctx Context, params openPullRequestParams) (string, error) {\n\tif err := ctx.Forgejo().OpenPullRequest(params.Head, params.Base, params.Title, params.Body); err != nil {\n\t\treturn \"\", err\n\t}\n\n\treturn \"Pull request opened\", nil\n}\n\ntype requestReviewParams struct {\n\tReviewers []string `json:\"reviewers\"`\n}\n\nfunc init() {\n\ttool := llm.NewTool(\n\t\t\"request_review\",\n\t\t\"Request review of the current pull request from one or more Forgejo users.\")\n\n\ttool.AddStringArrayProperty(\"reviewers\", \"Usernames to request review from\", true)\n\n\tRegister(tool, requestReview)\n}\n\nfunc requestReview(ctx Context, params requestReviewParams) (string, error) {\n\tif err := ctx.Forgejo().RequestReview(params.Reviewers); err != nil {\n\t\treturn \"\", err\n\t}\n\n\treturn \"Review requested\", nil\n}\n\ntype submitReviewParams struct {\n\tState string `json:\"state\"`\n\tBody string `json:\"body\"`\n}\n\nfunc init() {\n\ttool := llm.NewTool(\n\t\t\"submit_review\",\n\t\t\"Submit a review on the pull request that triggered this run: approve it, request changes, or leave review feedback. A submitted review fires a pr:review event that triggers a run, so use it to react to a PR (e.g. after addressing earlier feedback) or to review someone else's PR.\")\n\n\ttool.AddEnumProperty(\"state\", \"The review verdict\", []string{\"approved\", \"changes_requested\", \"comment\"}, true)\n\ttool.AddStringProperty(\"body\", \"The review text, in Markdown. Required unless state is 'approved'.\", true)\n\n\tRegister(tool, submitReview)\n}\n\nfunc submitReview(ctx Context, params submitReviewParams) (string, error) {\n\tif params.State != \"approved\" \u0026\u0026 strings.TrimSpace(params.Body) == \"\" {\n\t\treturn \"\", fmt.Errorf(\"body is required unless state is 'approved'\")\n\t}\n\n\tif err := ctx.Forgejo().SubmitReview(params.State, params.Body); err != nil {\n\t\treturn \"\", err\n\t}\n\n\treturn \"Review submitted\", nil\n}\n\ntype labelParams struct {\n\tName string `json:\"name\"`\n}\n\nfunc init() {\n\ttool := llm.NewTool(\n\t\t\"add_label\",\n\t\t\"Add a label to the issue or pull request that triggered this run, creating it on the repo first if needed.\")\n\n\ttool.AddStringProperty(\"name\", \"The label name\", true)\n","start_line":1,"end_line":120,"total_lines":204,"truncated":true}
Thinking
Let me check the remaining tools in forgejo.go (close_issue, reopen_issue, assign_issue) β these probably don't have parameters.
Tool
π§ read_file
{"path":"internal/tools/forgejo.go","start_line":120,"end_line":204}
Arguments
{"path":"internal/tools/forgejo.go","start_line":120,"end_line":204}
Result
{"path":"internal/tools/forgejo.go","content":"\n\tRegister(tool, addLabel)\n}\n\nfunc addLabel(ctx Context, params labelParams) (string, error) {\n\tif err := ctx.Forgejo().AddLabel(params.Name); err != nil {\n\t\treturn \"\", err\n\t}\n\n\treturn \"Label added\", nil\n}\n\nfunc init() {\n\ttool := llm.NewTool(\n\t\t\"remove_label\",\n\t\t\"Remove a label from the issue or pull request that triggered this run.\")\n\n\ttool.AddStringProperty(\"name\", \"The label name\", true)\n\n\tRegister(tool, removeLabel)\n}\n\nfunc removeLabel(ctx Context, params labelParams) (string, error) {\n\tif err := ctx.Forgejo().RemoveLabel(params.Name); err != nil {\n\t\treturn \"\", err\n\t}\n\n\treturn \"Label removed\", nil\n}\n\ntype assignIssueParams struct {\n\tAgent string `json:\"agent\"`\n}\n\nfunc init() {\n\ttool := llm.NewTool(\n\t\t\"assign_issue\",\n\t\t\"Hand the issue or pull request that triggered this run off to another agent, by name. This assigns it to that agent's Forgejo account, which triggers their run. Use this when the task is better suited to another agent's described role than your own.\")\n\n\ttool.AddStringProperty(\"agent\", \"The name of the agent to hand this off to, e.g. one listed in the \\\"Other agents\\\" section of your system prompt\", true)\n\n\tRegister(tool, assignIssue)\n}\n\nfunc assignIssue(ctx Context, params assignIssueParams) (string, error) {\n\tif err := ctx.Forgejo().AssignIssue(params.Agent); err != nil {\n\t\treturn \"\", err\n\t}\n\n\treturn fmt.Sprintf(\"Assigned to %s\", params.Agent), nil\n}\n\ntype noParams struct{}\n\nfunc init() {\n\ttool := llm.NewTool(\n\t\t\"close_issue\",\n\t\t\"Close the issue or pull request that triggered this run.\")\n\n\tRegister(tool, closeIssue)\n}\n\nfunc closeIssue(ctx Context, _ noParams) (string, error) {\n\tif err := ctx.Forgejo().CloseIssue(); err != nil {\n\t\treturn \"\", err\n\t}\n\n\treturn \"Issue closed\", nil\n}\n\nfunc init() {\n\ttool := llm.NewTool(\n\t\t\"reopen_issue\",\n\t\t\"Reopen the issue or pull request that triggered this run.\")\n\n\tRegister(tool, reopenIssue)\n}\n\nfunc reopenIssue(ctx Context, _ noParams) (string, error) {\n\tif err := ctx.Forgejo().ReopenIssue(); err != nil {\n\t\treturn \"\", err\n\t}\n\n\treturn \"Issue reopened\", nil\n}","start_line":120,"end_line":204,"total_lines":204,"truncated":false}
Thinking
Confirmed: `close_issue` and `reopen_issue` are registered without any properties. That means `Function.Parameters` stays as `Parameters{Type: ""}`. And in `types.go`:
```go
type Function struct {
Name string `json:"name"`
Description string `json:"description,omitempty"`
Parameters Parameters `json:"parameters,omitempty"`
Strict bool `json:"strict"`
}
```
`Parameters` is a struct, so `omitempty` doesn't omit it. It gets serialized as `"parameters":{"type":""}` β exactly the invalid schema in the 429 error: "Tool 0 function has invalid 'parameters' schema: '' is not valid under any of the given schemas".
So the 429 in the issue is caused by this bug: every chat completion request includes all tools (tools.All()), including close_issue/reopen_issue with `"parameters":{"type":""}`. litellm (with pre-call-checks) rejects the request with a 400, puts the model group into cooldown, and surfaces it as a 429.
This is important: the retry logic will retry this 429 ten times, and it'll still fail after that. The issue is "retry logic for when the LLM endpoint fails" β the 502 case is a genuine endpoint failure that retry will fix. The 429 case in the comment is a client-side bug that retry can't fix.
Should I fix the tool schema bug as part of this change? It's a separate bug, but it's the root cause of one of the two reported failures. As a senior architect, I think the right move is:
1. Implement retry logic (the actual ask).
2. Fix the parameterless tool schema bug (small, surgical: set `Parameters.Type = "object"` in NewTool, and/or omit parameters when there are no properties) β because without this, the retry logic will just turn a fast failure into a slow failure for every run (every run sends all tools, so every run hits the 429!).
Wait, actually β hold on. Every run sends all tools. So every run against litellm with pre-call-checks enabled gets a 400/429? But the agents have been working (this very run is happening, and previous PRs have been merged by greg). Hmm. So maybe litellm's pre-call-checks only validate whenβ¦? The 429 message says "pre-call-checks=True". If pre-call-checks were always enabled and always failing, no run would ever work. But this run is working (I'm running!). So either:
- pre-call-checks was recently enabled, or
- validation only happens on some code path, or
- the model group changed (Qwen3.8-27B is a new model with a fallback group).
Whatever the reason, the error is real, and it's in the issue's comments as evidence of the failure. The invalid schema is objectively wrong per the JSON Schema metaschema (`type` must be a string from the simpleTypes list, or an array thereof β "" is invalid). Fixing it is correct regardless of whether it's the direct cause.
Minimal fix in `NewTool`:
```go
func NewTool(name, description string) Tool {
return Tool{
Type: "function",
Function: Function{
Name: name,
Description: description,
Parameters: Parameters{Type: "object"},
Strict: true,
},
}
}
```
This makes parameterless tools serialize as `"parameters":{"type":"object"}` β valid schema (an object with no properties). And tools with properties already have Type setβ¦ wait, do they? `AddProperty` sets `t.Function.Parameters.Properties[name] = ...`, but never sets `Type`! Let me recheck:
```go
func (t *Tool) AddProperty(name string, property Property, required bool) *Tool {
if t.Function.Parameters.Properties == nil {
t.Function.Parameters.Properties = make(map[string]Property)
} else if _, exists := t.Function.Parameters.Properties[name]; exists {
panic("property already exists: " + name)
}
t.Function.Parameters.Properties[name] = property
...
}
```
Nope β `Type` is never set anywhere! So every tool, even ones with properties, serializes as `"parameters":{"type":"","properties":{...},"required":[...]}`. The `type` field is always ""!
Wait, then every request should fail validation⦠but the agents have been working. Hmm. Let me double-check: is `Type` set somewhere? Let me grep for `Parameters{` and `.Type =`.
Actually, wait β in the `Parameters` struct, `Type string json:"type"` β no omitempty. So it always serializes, even as "". So every tool definition in every request has `"type":""`. And yet, the agents have been running successfully for weeks (many merged PRs). So the LLM endpoint in normal operation doesn't validate `type` strictly β but litellm with pre-call-checks=True does (it validates against the JSON Schema metaschema).
Hmm, OK, so here's the deal: this deployment's endpoint (litellm in front of Qwen) has pre-call-checks that reject `"type":""`. This is a real compatibility bug in zoo's tool serialization. The fix (set Type: "object" in NewTool) is small, surgical, and directly addresses the failure in the issue's comment.
But wait β is this in scope? The issue says "I need some kind of retry logic in the agent loop for when the LLM endpoint is failing". The comments show two failures: a 429 (caused by the schema bug + cooldown) and a 502 (engine not healthy). Retry logic addresses the 502 and transient 429s. The schema bug is a separate root cause.
Decision: implement retry logic as the primary change, and also fix the tool schema bug as a small companion fix. Reasons:
1. It's the root cause of the 429 in the issue's comments.
2. Without it, the retry logic will just convert an immediate failure into a 5-minute failure for every run on that endpoint.
3. It's a 1-line fix + a test.
I'll call this out clearly in the PR/issue comments so abrander can split it out if he wants.
Hmm, wait. Let me think again. "Make surgical changes: only change what's necessary for the task in the triggering event". The task is retry logic. The schema fix isβ¦ necessary for the retry logic to be effective against the observed 429. I'll include it, but keep it minimal and well-documented. Actually β let me verify the claim that Type is never set with a quick test: marshal a tool and look at the JSON. I'll do that in a scratch test.
Also, `Strict: true` β with strict mode, OpenAI requires the schema to be fully strict (additionalProperties: false, etc.). llama.cpp and others may handle it differently. Changing strict is out of scope.
One more consideration on the schema fix: setting `Type: "object"` in NewTool β for tools with properties, the serialized JSON becomes `{"type":"object","properties":{...},"required":[...]}` β the standard shape.
Now β what about `Function.Parameters` for tools with no properties: `{"type":"object"}` β valid.
OK. Now, let me also double-check that there aren't other places that construct tools/parameters⦠`tools.All()` returns registered tools. All go through NewTool. Good.
Let me also reconsider the "after %d attempts" wrapping. Hmm, one more thought: Do wraps the error when attempt > 1. But what if the final attempt fails with a non-retryable error after several retries? E.g., attempts 1β9 fail with 429 (retryable), attempt 10 fails with 401 (non-retryable). Then Do returns "after 10 attempts: chat completion: status: 401: ...". Accurate enough.
What if ctx is canceled during the wait? sleep returns ctx.Err() β Do returns it unwrapped. Good β the scheduler checks errors.Is(err, context.DeadlineExceeded) for the JobTimedOut status. Important: if we wrap a DeadlineExceeded with "after N attempts: %w", errors.Is still works through %w. And if ctx is canceled during the wait, we return the raw ctx.Err(). Either way errors.Is works.
But wait β there's a subtle issue with the run timeout: run.go's `ctx` has a 120-minute timeout (DefaultTimeout). If the LLM endpoint is down and we retry for ~5 minutes total, that's well within 120 minutes. No problem.
Now, one more design question: **where does the retry go β inside runLoop or in a Client method?** I decided on runLoop + llm.Do. Let me also double-check that no other callers of StreamChatCompletion need retryβ¦ only runLoop calls it. Good.
Next, the livelog `Retry` event type + JS rendering. Let me design:
hub.go:
```go
const (
...
Tool Type = "tool"
Retry Type = "retry"
Status Type = "status"
)
```
JS in templates.go:
```js
case "retry":
if (reasoningBody) { reasoningBody.parentElement.remove(); reasoningBody = null; }
if (contentBody) { contentBody.parentElement.remove(); contentBody = null; }
newSystemBlock(ev.text);
break;
```
I need a newSystemBlock helper similar to newBlock. Let me look at the whole JS section to see newBlock and the CSS classes for blocks. Let me read the template around lines 850β905.
Let me see the full job-detail JS (block creation helpers) and the block CSS to plan the live retry rendering.
Tool
π§ read_file
{"path":"internal/web/templates.go","start_line":820,"end_line":905}
Arguments
{"path":"internal/web/templates.go","start_line":820,"end_line":905}
Result
{"path":"internal/web/templates.go","content":" {{.Job.Status}}\n \u003c/span\u003e\n \u003c/span\u003e\n \u003c/div\u003e\n \u003cdiv class=\"job-meta-item\"\u003e\n \u003cspan class=\"label\"\u003eEvent\u003c/span\u003e\n \u003cspan class=\"value\"\u003e{{.Job.EventKind}} on \u003ccode\u003e{{.Job.Owner}}/{{.Job.Repo}}#{{.Job.IssueIndex}}\u003c/code\u003e\u003c/span\u003e\n \u003c/div\u003e\n \u003cdiv class=\"job-meta-item\"\u003e\n \u003cspan class=\"label\"\u003eAgent\u003c/span\u003e\n \u003cspan class=\"value\"\u003e\n \u003cspan class=\"agent\"\u003e\n {{if .AvatarURL}}\u003cimg class=\"agent-avatar\" src=\"{{.AvatarURL}}\" alt=\"{{.Job.Agent}}\" title=\"{{.Job.Agent}}\"\u003e{{end}}\n \u003cstrong\u003e{{.Job.Agent}}\u003c/strong\u003e\n \u003c/span\u003e\n \u003c/span\u003e\n \u003c/div\u003e\n {{if .Job.Error}}\n \u003cdiv class=\"job-meta-item\"\u003e\n \u003cspan class=\"label\"\u003eError\u003c/span\u003e\n \u003cspan class=\"value error-text\"\u003e{{.Job.Error}}\u003c/span\u003e\n \u003c/div\u003e\n {{end}}\n \u003c/div\u003e\n\n \u003ch2\u003eLog\u003c/h2\u003e\n \u003cdiv class=\"log-container\" id=\"log\"\u003e\n {{range .Blocks}}\n {{if eq .Kind \"reasoning\"}}\n \u003cdiv class=\"block block-reasoning\"\u003e\n \u003cdiv class=\"block-label\"\u003eThinking\u003c/div\u003e\n \u003cdiv class=\"block-body\"\u003e{{.Text}}\u003c/div\u003e\n \u003c/div\u003e\n {{else if eq .Kind \"content\"}}\n \u003cdiv class=\"block block-content\"\u003e\n \u003cdiv class=\"block-body\"\u003e{{.Text}}\u003c/div\u003e\n \u003c/div\u003e\n {{else if eq .Kind \"tool\"}}\n \u003cdetails class=\"block block-tool{{if .Tool.Error}} block-tool-error{{end}}\"\u003e\n \u003csummary\u003e\n \u003cspan class=\"tool-badge\"\u003eTool\u003c/span\u003e\n \u003cspan class=\"tool-summary-text\"\u003e\n \u003cspan class=\"tool-name\"\u003eπ§ {{.Tool.Name}}\u003c/span\u003e\n \u003cspan class=\"tool-args-preview\"\u003e{{.Tool.Arguments}}\u003c/span\u003e\n \u003c/span\u003e\n \u003c/summary\u003e\n \u003cdiv class=\"block-body\"\u003e\n \u003cdiv class=\"tool-section-label\"\u003eArguments\u003c/div\u003e\n \u003cpre\u003e{{.Tool.Arguments}}\u003c/pre\u003e\n \u003cdiv class=\"tool-section-label\"\u003eResult\u003c/div\u003e\n \u003cpre\u003e{{.Tool.Result}}\u003c/pre\u003e\n \u003c/div\u003e\n \u003c/details\u003e\n {{else}}\n \u003cdiv class=\"block block-system\"\u003e{{.Text}}\u003c/div\u003e\n {{end}}\n {{end}}\n \u003c/div\u003e\n\n {{if .Live}}\n \u003cscript\u003e\n (function() {\n var jobID = {{.Job.ID}};\n var log = document.getElementById(\"log\");\n var reasoningBody = null;\n var contentBody = null;\n\n function nearBottom() {\n return (window.innerHeight + window.scrollY) \u003e= (document.body.offsetHeight - 80);\n }\n\n function newBlock(kind, label) {\n var div = document.createElement(\"div\");\n div.className = \"block block-\" + kind;\n if (label) {\n var l = document.createElement(\"div\");\n l.className = \"block-label\";\n l.textContent = label;\n div.appendChild(l);\n }\n var body = document.createElement(\"div\");\n body.className = \"block-body\";\n div.appendChild(body);\n log.appendChild(div);\n return body;\n }","start_line":820,"end_line":905,"total_lines":1004,"truncated":true}
Tool
π§ bash
{"command":"cd /project && grep -n \"block-system\\|block block\\|\\.block\" internal/web/templates.go | head -30"}
Arguments
{"command":"cd /project && grep -n \"block-system\\|block block\\|\\.block\" internal/web/templates.go | head -30"}
Result
421: children (.block-tool) gives those children an automatic min-height
424: .log-container .block + .block {
440: .block-label {
449: .block-body {
458: .block-reasoning,
459: .block-content {
464: .block-reasoning {
469: .block-reasoning .block-body {
474: .block-content {
479: .block-system {
486: .block-tool {
494: .block-tool summary {
505: .block-tool summary::-webkit-details-marker { display: none; }
507: .block-tool summary::before {
516: .block-tool[open] summary::before { transform: rotate(90deg); }
537: .block-tool .tool-name {
543: .block-tool .tool-args-preview {
554: .block-tool-error { border-left-color: #ef4444; }
555: .block-tool-error .tool-badge { background: rgba(239, 68, 68, 0.15); color: #ef4444; }
557: .block-tool .block-body {
562: .block-body's white-space: pre-wrap β that would render the
568: .block-tool .tool-section-label {
577: .block-tool pre {
849: <div class="block block-reasoning">
854: <div class="block block-content">
858: <details class="block block-tool{{if .Tool.Error}} block-tool-error{{end}}">
874: <div class="block block-system">{{.Text}}</div>
893: div.className = "block block-" + kind;
909: details.className = "block block-tool" + (ev.error ? " block-tool-error" : "");
Thinking
The `block-system` CSS class already exists. So the live "retry" event can be rendered as a system block:
```js
case "retry":
// Discard the partial block from the failed attempt.
if (reasoningBody) { reasoningBody.parentElement.remove(); reasoningBody = null; }
if (contentBody) { contentBody.parentElement.remove(); contentBody = null; }
var sys = document.createElement("div");
sys.className = "block block-system";
sys.textContent = ev.text;
log.appendChild(sys);
break;
```
Alright, let me write out the complete implementation plan:
### Files to change:
1. **`internal/llm/retry.go`** (new): `APIError`, `IsRetryable`, `Retry` policy, `Do`, backoff, Retry-After parsing.
2. **`internal/llm/streamer.go`**: `NewStreamer` returns `*APIError` (preserving the text format).
3. **`internal/llm/client.go`**: no changes? Actually β hmm, wait. Should I also retry inside the client? No β the loop-level retry covers it.
4. **`internal/llm/types.go`**: `NewTool` sets `Parameters: Parameters{Type: "object"}` (schema bug fix).
5. **`internal/agentrun/loop.go`**: `Hooks.OnRetry`; wrap the turn in `llm.Do`.
6. **`internal/agentrun/run.go`**: implement `OnRetry` in `streamHooks` (reset state, logAppend "system", hub.Publish Retry).
7. **`internal/livelog/hub.go`**: add `Retry Type = "retry"`.
8. **`internal/web/templates.go`**: JS case for "retry".
9. **Tests**:
- `internal/llm/retry_test.go`: IsRetryable cases, Do retry behavior (succeeds after transient failure, stops on non-retryable, respects MaxAttempts, respects Retry-After, respects ctx cancellation during sleep), backoff bounds.
- `internal/llm/streamer_test.go`: extend TestStreamerNonOKStatus to check the APIError type + Retry-After parsing.
- `internal/llm/types_test.go`? A test that NewTool produces a valid schema JSON (Type "object").
- `internal/agentrun/loop_test.go` (new): test runLoop's retry with a fake... hmm, runLoop takes *llm.Client (a concrete type). To test retry, I need an HTTP server that fails N times then succeeds. llm.Client takes an endpoint from config.LLM β I can point it at an httptest.Server! `llm.NewClient(config.LLM{OpenAI: server.URL, ...})`.
runLoop test: a server that returns 500 twice, then returns a valid SSE stream with finish_reason "stop". Verify runLoop succeeds and the server saw 3 requests. Also test that non-retryable (400) fails immediately (1 request). And that OnRetry fires with the correct attempt number. And that the messages slice is correct.
But runLoop calls `tools.All()` β that's fine (the registry is initialized in the tools package's init; runLoop passes tools.All() to the request; the test server ignores the body). And toolsCtx β runLoop only uses toolsCtx when executing tool calls; for a "stop" finish, no tool calls. Pass a zero-value tools.Context? Let me check tools.Context β it's a struct; zero value should be fine as long as no tools are called. Let me verify.
Also, runLoop is unexported β the test goes in package agentrun. There's already a seed_test.go in agentrun. Let me check how it tests (what package name).
For the retry test I need to control the backoff: runLoop uses llm.DefaultRetry() β with a base of 2 seconds, a 3-attempt test takes 2+4=6 seconds. Too slow. I need to make the retry policy injectable into runLoop. Options:
- A package variable `var loopRetry = llm.DefaultRetry()` that the test can override. Hmm, a mutable global β but it's a common pattern in Go for testability.
- Or runLoop takes the policy as a parameter β but that changes the signature (the only caller is run.go). I could add a parameter... it's an internal function, so it's fine. But it pollutes the signature with a test-only concern.
- Or a field on Hooks? No.
A package-level variable with a clear name is practical:
```go
// retryPolicy is the retry policy runLoop applies to each model turn.
// A package variable (not a parameter) so tests can shorten the
// backoff; production always uses the default.
var retryPolicy = llm.DefaultRetry()
```
Hmm, wait, but then run.go's streamHooks OnRetry wiring... no, that's separate. OK.
Actually β hmm, let me reconsider. A cleaner approach: `runLoop` stays as-is and uses `llm.DefaultRetry()`, and I add an unexported variant for testing? No β the package variable is the simplest and most honest. Many Go codebases do this. Let's go with that.
Actually, wait, an even simpler alternative for the agentrun test: override `retryPolicy.Sleep` to a no-op in the test:
```go
retryPolicy.Sleep = func(ctx context.Context, d time.Duration) error { return nil }
```
Then the actual delays are skipped.
10. **`zoo.hcl.example`**: no changes (no new config).
Wait β one more thing to verify: `tools.Context` β let me make sure the zero value is OK for the test. Let me read tools/context.go.
Also, check the package name of seed_test.go in agentrun (agentrun or agentrun_test).
Let me also reconsider the default values: MaxAttempts 10, BaseDelay 2s, MaxDelay 60s.
Total worst-case wait time: attempts 2..10 β delays 2,4,8,16,32,60,60,60,60 = 302s β 5 minutes. Plus request time. Within the 120-minute run window, no problem.
Expected total with equal jitter (0.75Γ): about 3.8 minutes.
For litellm's 152-second cooldown (no Retry-After header): cumulative worst case: 2,6,14,30,62,122,182 β attempt 7 starts at t=122s (before cooldown ends), attempt 8 at t=182s (after). So it recovers on attempt 7 or 8. Good.
If Retry-After: 152 is present β first wait = max(backoff 2s, 152s) = 152s β attempt 2 at t=152s. Right at the boundary β hmm, "try again in 152 seconds" means it's available at t=152. A request exactly at 152 might still be in cooldown (cooldown end is computed from the failure time, and our wait starts after the response). To be safe, add a small margin? If Retry-After is present, delay = ra + 1s? Over-engineering; litellm's cooldown check is `time.time() < cooldown_until`. A 1-second margin is cheap insurance:
Hmm, let's keep it simple: use ra as-is. If it fails, the next backoff (4s) covers it. Not a problem.
Now β the `Do` signature with generics. Go 1.26, no problem.
One more consideration: **should `Do` retry when fn fails with a non-retryable error?** No β return immediately. **Should it check IsRetryable before the first attempt?** No, the first attempt always happens.
**Edge case: fn returns (zero, nil)** β success, return. No problem.
Now, the `BeforeRetry` callback: it's called before the sleep. In run.go's OnRetry implementation, it does logAppend (store write, context.Background) + hub publish. It's synchronous and fast. OK.
But wait β there's a subtle issue with the ordering of OnRetry and the hub's replay buffer. Sequence for a failed attempt:
1. Deltas are published (reasoning_start, reasoning_deltaΓn) β into the buffer.
2. Turn fails β OnRetry: reset the server-side buffer; publish the Retry event β into the buffer.
3. Retry the turn: new deltas β into the buffer.
4. Turn succeeds β OnTurnEnd: persist content, Checkpoint (clear buffer).
A live viewer sees: partial thinking block, then the retry system block (the JS removes the partial block), then the new thinking block.
A viewer who connects between 2 and 4: the buffered events = [Retry, new deltas...] β the partial deltas are gone (they were in the buffer before the Retry, but... wait, no! Publish appends to the buffer; OnRetry does not clear the hub buffer. The partial deltas are still in the buffer!
Hmm. Let me re-check: the hub buffer accumulates events until Checkpoint. OnTurnEnd/OnTool checkpoint. A mid-turn failure doesn't checkpoint. So after OnRetry, buffer = [partial deltas..., Retry event, ...]. A late subscriber replays the partial deltas, then the Retry (the JS removes the partial block), then the new deltas. OK β the JS removal handles it! The partial block gets rendered and then removed by the retry event.
But what if the retry also fails, and the run ultimately fails? Then no checkpoint happens; the job ends; the scheduler publishes the Status event; the page reloads β static view β the persisted log has: the "system" retry line (persisted via logAppend) β the partial content is never persisted (the buffer was reset on the server side, and OnTurnEnd never fired). The static view shows the retry line, no partial block. Consistent.
And on the live view before the reload: partial block, retry line, partial block 2, retry line 2... the JS removes each partial on its retry event. Consistent.
One more: the server-side `reasoningBuf`/`contentBuf` are reset in OnRetry β but the hub buffer still has the old delta events (unavoidable without a new hub API; the JS handles it). Hmm, wait β actually, is there a problem if the same job later... no, the buffer is per-job, and the job is done. No problem.
Actually, wait, let me reconsider: should I add a hub method to clear the buffer without checkpointing? A `Reset` that's like Checkpoint... Checkpoint is exactly "clear the buffer". The semantics of Checkpoint: "call this right after the block derived from buffered events has been durably written to storage". Calling Checkpoint on retry would be a semantic misuse (nothing has been persisted) β but the effect we want is to discard the stale partial events so late subscribers don't replay them (the JS removal makes the replay harmless, but it's cleaner not to replay).
Hmm. If I Checkpoint on OnRetry: late subscribers see only [Retry, new deltas]. Early subscribers (already connected) receive the Retry live β the JS removes the partial. Both consistent. And the persisted log has the retry line. That's cleaner!
But wait β is there a risk: Checkpoint clears the buffer; then the retry's deltas get appended; if the retry succeeds, OnTurnEnd persists + checkpoints. No problem. If the retry fails β OnRetry again β Checkpoint again. No problem.
But Checkpoint's docstring says "so that subscribers connecting later β and therefore already seeing that block via the persisted history β are never replayed it a second time". Using it for "discard stale partials" is a different purpose. I could add a dedicated method with a clear name... or just use Checkpoint with a comment. Hmm.
Let me think about whether the replay of partial deltas is actually harmful without the Checkpoint: a late subscriber sees the partial thinking block get rendered, then the retry event removes it, then the full block. A slight flicker, but correct. With the Checkpoint: no flicker.
I'll go with calling Checkpoint in OnRetry with a clear comment. It's the existing primitive that does exactly "clear the replay buffer". Actually β hmm, wait. One edge case: what if OnRetry fires and then the process crashes before the retry? The buffer is cleared, the partial is lost from the live view, but the persisted log has the retry line. The job shows as failed with an error. No problem.
OK, wait, I realize I should also double-check the `Status` event: who publishes it? Let me grep for livelog.Status usage in agentrun.
Let me check a few remaining details: `tools.Context`, who publishes `livelog.Status`, and the existing agentrun test conventions.
Tool
π§ bash
{"command":"cd /project && grep -rn \"livelog.Status\\|Type: livelog\" internal/agentrun/ | head && echo ---- && head -40 internal/tools/context.go && echo ---- && head -30 internal/agentrun/seed_test.go && echo ---- && grep -n \"package\" internal/agentrun/*_test.go"}
Arguments
{"command":"cd /project && grep -rn \"livelog.Status\\|Type: livelog\" internal/agentrun/ | head && echo ---- && head -40 internal/tools/context.go && echo ---- && head -30 internal/agentrun/seed_test.go && echo ---- && grep -n \"package\" internal/agentrun/*_test.go"}
Result
internal/agentrun/run.go:298: r.hub.Publish(jobID, livelog.Event{Type: livelog.ReasoningStart})
internal/agentrun/run.go:302: r.hub.Publish(jobID, livelog.Event{Type: livelog.ReasoningDelta, Text: delta})
internal/agentrun/run.go:308: r.hub.Publish(jobID, livelog.Event{Type: livelog.ContentStart})
internal/agentrun/run.go:312: r.hub.Publish(jobID, livelog.Event{Type: livelog.ContentDelta, Text: delta})
internal/agentrun/run.go:316: r.hub.Publish(jobID, livelog.Event{Type: livelog.ReasoningEnd})
internal/agentrun/run.go:323: r.hub.Publish(jobID, livelog.Event{Type: livelog.ContentEnd})
----
// Package tools is the registry of functions an agent's tool-calling
// loop can invoke, ported from ../a/tools but retargeted at a Docker
// container (via Exec) and the Forgejo issue/PR that triggered the run
// (via Forgejo()) instead of a local sandboxed filesystem.
package tools
import "context"
// Context is bound to a single agent run (job), unlike ../a's single
// global Context, since multiple agents run concurrently in zoo.
type Context interface {
// Exec runs command inside the run's container via `sh -c` and
// returns combined stdout+stderr.
Exec(ctx context.Context, command string) (string, error)
// Forgejo returns the actions bound to the issue/PR that triggered
// this run, so tools don't need to be told which repo/issue to act
// on.
Forgejo() ForgejoActions
}
// ForgejoActions are the actions available to an agent against the
// triggering issue/PR. Implemented by internal/agentrun, backed by
// internal/forgejo.Client plus the run's owner/repo/issue.
type ForgejoActions interface {
Comment(body string) error
OpenPullRequest(head, base, title, body string) error
RequestReview(reviewers []string) error
AddLabel(name string) error
RemoveLabel(name string) error
CloseIssue() error
ReopenIssue() error
// AssignIssue hands the triggering issue/PR off to another agent by
// their name (which doubles as their Forgejo username); the
// resulting issue:assigned event triggers that agent's run.
AssignIssue(agentName string) error
// SubmitReview submits a review on the triggering pull request with
// the given verdict (forgejo.ReviewStateApproved,
----
package agentrun
import (
"strings"
"testing"
"time"
"github.com/abrander/zoo/internal/forgejo"
)
func TestSeedMessageIncludesInstructions(t *testing.T) {
ev := forgejo.Event{Kind: "issue:comment", Owner: "acme", Repo: "widgets", Index: 5, Title: "t", Body: "b", Raw: []byte(`{}`)}
msg := seedMessage(ev, "issue-5-leon", "main", "Please review the comment and respond appropriately.", nil, nil)
if !strings.Contains(msg, "Instructions for this event, from zoo.hcl:\nPlease review the comment and respond appropriately.") {
t.Fatalf("expected instructions section, got: %s", msg)
}
}
func TestSeedMessageOmitsEmptyInstructions(t *testing.T) {
ev := forgejo.Event{Kind: "pr:new", Owner: "acme", Repo: "widgets", Index: 9, Raw: []byte(`{}`)}
msg := seedMessage(ev, "issue-9-greg", "main", "", nil, nil)
if strings.Contains(msg, "Instructions for this event") {
t.Fatalf("expected no instructions section, got: %s", msg)
}
}
----
internal/agentrun/docker_integration_test.go:5:package agentrun
internal/agentrun/roster_test.go:1:package agentrun
internal/agentrun/sandboxgit_test.go:1:package agentrun
internal/agentrun/seed_test.go:1:package agentrun
Thinking
`livelog.Status` is not used in agentrun β probably used in the scheduler. Let me check right awayβ¦ actually not important. Let me check where Status is published:
No, let's not worry about it. The JS treats "status" as "job ended β reload".
`tools.Context` is an interface β in the runLoop test, we can pass nil. As long as no tool calls are executed (finish_reason "stop").
Tests are in `package agentrun` (internal test package). Good β we can access the `retryPolicy` variable and `runLoop`.
Now, before writing code, let me verify the tool schema bug with a quick experiment. I'll write a small Go test that prints the JSON of tools.All().
Let me verify the tool-schema serialization bug with a quick experiment before writing the fix.
Tool
π§ bash
{"command":"cd /project && cat > /tmp/schema_check_test.go <<'EOF'\npackage tools\n\nimport (\n\t\"encoding/json\"\n\t\"strings\"\n\t\"testing\"\n)\n\nfunc TestSchemaCheck(t *testing.T) {\n\tfor _, tool := range All() {\n\t\tb, _ := json.Marshal(tool)\n\t\tif strings.Contains(string(b), `\"type\":\"\"`) {\n\t\t\tt.Logf(\"INVALID: %s -> %s\", tool.Function.Name, string(b))\n\t\t}\n\t}\n\tb, _ := json.Marshal(NewToolByName(\"close_issue\"))\n\tt.Logf(\"close_issue: %s\", b)\n}\nEOF\n# NewToolByName doesn't exist; just marshal via the registry\nsed -i 's/b, _ := json.Marshal(NewToolByName(\"close_issue\"))\\n\\tt.Logf(\"close_issue: %s\", b)//' /tmp/schema_check_test.go\ncp /tmp/schema_check_test.go internal/tools/zz_schema_check_test.go\n# fix: remove the nonexistent helper call\npython3 - <<'PY'\nimport re\np = 'internal/tools/zz_schema_check_test.go'\ns = open(p).read()\ns = s.replace('''\tb, _ := json.Marshal(NewToolByName(\"close_issue\"))\n\tt.Logf(\"close_issue: %s\", b)\n''', '')\nopen(p, 'w').write(s)\nPY\ngo test ./internal/tools/ -run TestSchemaCheck -v 2>&1 | head -30\nrm internal/tools/zz_schema_check_test.go"}
Arguments
{"command":"cd /project && cat > /tmp/schema_check_test.go <<'EOF'\npackage tools\n\nimport (\n\t\"encoding/json\"\n\t\"strings\"\n\t\"testing\"\n)\n\nfunc TestSchemaCheck(t *testing.T) {\n\tfor _, tool := range All() {\n\t\tb, _ := json.Marshal(tool)\n\t\tif strings.Contains(string(b), `\"type\":\"\"`) {\n\t\t\tt.Logf(\"INVALID: %s -> %s\", tool.Function.Name, string(b))\n\t\t}\n\t}\n\tb, _ := json.Marshal(NewToolByName(\"close_issue\"))\n\tt.Logf(\"close_issue: %s\", b)\n}\nEOF\n# NewToolByName doesn't exist; just marshal via the registry\nsed -i 's/b, _ := json.Marshal(NewToolByName(\"close_issue\"))\\n\\tt.Logf(\"close_issue: %s\", b)//' /tmp/schema_check_test.go\ncp /tmp/schema_check_test.go internal/tools/zz_schema_check_test.go\n# fix: remove the nonexistent helper call\npython3 - <<'PY'\nimport re\np = 'internal/tools/zz_schema_check_test.go'\ns = open(p).read()\ns = s.replace('''\tb, _ := json.Marshal(NewToolByName(\"close_issue\"))\n\tt.Logf(\"close_issue: %s\", b)\n''', '')\nopen(p, 'w').write(s)\nPY\ngo test ./internal/tools/ -run TestSchemaCheck -v 2>&1 | head -30\nrm internal/tools/zz_schema_check_test.go"}
Result
go: downloading github.com/hashicorp/hcl/v2 v2.24.0
go: downloading github.com/zclconf/go-cty v1.16.3
go: downloading github.com/agext/levenshtein v1.2.1
go: downloading github.com/apparentlymart/go-textseg/v15 v15.0.0
go: downloading github.com/mitchellh/go-wordwrap v1.0.1
go: downloading github.com/google/go-cmp v0.7.0
go: downloading golang.org/x/text v0.40.0
=== RUN TestSchemaCheck
zz_schema_check_test.go:13: INVALID: bash -> {"type":"function","function":{"name":"bash","description":"Execute a shell command in the project container. Use this to run shell commands and scripts.","parameters":{"type":"","properties":{"command":{"type":"string","description":"The shell command to execute"},"timeout":{"type":"integer","description":"Optional timeout in seconds. Must not exceed 300 seconds."}},"required":["command"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: comment -> {"type":"function","function":{"name":"comment","description":"Post a comment on the issue or pull request that triggered this run.","parameters":{"type":"","properties":{"body":{"type":"string","description":"The comment text, in Markdown"}},"required":["body"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: open_pull_request -> {"type":"function","function":{"name":"open_pull_request","description":"Open a pull request from a branch you've pushed into the repository's default branch (or another base branch).","parameters":{"type":"","properties":{"base":{"type":"string","description":"The branch to merge into, e.g. the repository's default branch"},"body":{"type":"string","description":"The pull request description, in Markdown"},"head":{"type":"string","description":"The branch containing your changes (usually the run's own branch)"},"title":{"type":"string","description":"The pull request title"}},"required":["head","base","title","body"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: list_files -> {"type":"function","function":{"name":"list_files","description":"List files in a directory in the project container, optionally recursively. Hidden (dot) entries are skipped.","parameters":{"type":"","properties":{"path":{"type":"string","description":"The path to list, relative to the project root, e.g. . or src/"},"recursive":{"type":"boolean","description":"Whether to list files recursively (default: false)"}},"required":["path"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: edit_file -> {"type":"function","function":{"name":"edit_file","description":"Edit a file by replacing text. The tool will read the file, replace the specified text, and write the changes back to the file.","parameters":{"type":"","properties":{"find":{"type":"string","description":"The text to find in the file"},"path":{"type":"string","description":"The path to the file to edit"},"replace":{"type":"string","description":"The text to replace the found text with"},"replace_all":{"type":"boolean","description":"Whether to replace all occurrences of the found text (true) or just the first occurrence (false)"}},"required":["path","find","replace"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: request_review -> {"type":"function","function":{"name":"request_review","description":"Request review of the current pull request from one or more Forgejo users.","parameters":{"type":"","properties":{"reviewers":{"type":"array","description":"Usernames to request review from","items":{"type":"string"}}},"required":["reviewers"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: submit_review -> {"type":"function","function":{"name":"submit_review","description":"Submit a review on the pull request that triggered this run: approve it, request changes, or leave review feedback. A submitted review fires a pr:review event that triggers a run, so use it to react to a PR (e.g. after addressing earlier feedback) or to review someone else's PR.","parameters":{"type":"","properties":{"body":{"type":"string","description":"The review text, in Markdown. Required unless state is 'approved'."},"state":{"type":"string","description":"The review verdict","enum":["approved","changes_requested","comment"]}},"required":["state","body"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: remove_label -> {"type":"function","function":{"name":"remove_label","description":"Remove a label from the issue or pull request that triggered this run.","parameters":{"type":"","properties":{"name":{"type":"string","description":"The label name"}},"required":["name"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: move_file -> {"type":"function","function":{"name":"move_file","description":"Rename or move a file or directory in the project container.","parameters":{"type":"","properties":{"from_path":{"type":"string","description":"The current path to the file or directory"},"to_path":{"type":"string","description":"The new path for the file or directory"}},"required":["from_path","to_path"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: read_file -> {"type":"function","function":{"name":"read_file","description":"Read lines from a file in the project container. Returns the content of the file along with metadata such as total lines.","parameters":{"type":"","properties":{"end_line":{"type":"integer","description":"Optional line number to stop reading at (inclusive, 1-indexed)"},"path":{"type":"string","description":"The path to the file to read"},"start_line":{"type":"integer","description":"Optional line number to start reading from (1-indexed)"}},"required":["path","start_line","end_line"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: remove_file -> {"type":"function","function":{"name":"remove_file","description":"Remove an existing file in the project container.","parameters":{"type":"","properties":{"path":{"type":"string","description":"The path to the file to remove"}},"required":["path"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: write_file -> {"type":"function","function":{"name":"write_file","description":"Write a new file with the given content in the project container. Use this to create new files. If the file already exists, this will fail.","parameters":{"type":"","properties":{"content":{"type":"string","description":"The content to write to the file"},"path":{"type":"string","description":"The path to the file to create"}},"required":["path","content"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: add_label -> {"type":"function","function":{"name":"add_label","description":"Add a label to the issue or pull request that triggered this run, creating it on the repo first if needed.","parameters":{"type":"","properties":{"name":{"type":"string","description":"The label name"}},"required":["name"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: close_issue -> {"type":"function","function":{"name":"close_issue","description":"Close the issue or pull request that triggered this run.","parameters":{"type":""},"strict":true}}
zz_schema_check_test.go:13: INVALID: grep_search -> {"type":"function","function":{"name":"grep_search","description":"Search for an extended regular expression (ERE, e.g. 'foo|bar', 'func\\(') in project files (in the project container). Supports context lines, a glob filter, and case sensitivity control.","parameters":{"type":"","properties":{"case_sensitive":{"type":"boolean","description":"Whether the search should be case sensitive (default: false)"},"context_after":{"type":"integer","description":"Number of lines of context to show after each match (default: 0)"},"context_before":{"type":"integer","description":"Number of lines of context to show before each match (default: 0)"},"glob":{"type":"string","description":"Optional glob pattern to filter files, e.g. '*.go'"},"max_results":{"type":"integer","description":"Maximum number of matching lines to return (default: 100, 0 means no limit)"},"query":{"type":"string","description":"The regex pattern to search for"}},"required":["query"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: assign_issue -> {"type":"function","function":{"name":"assign_issue","description":"Hand the issue or pull request that triggered this run off to another agent, by name. This assigns it to that agent's Forgejo account, which triggers their run. Use this when the task is better suited to another agent's described role than your own.","parameters":{"type":"","properties":{"agent":{"type":"string","description":"The name of the agent to hand this off to, e.g. one listed in the \"Other agents\" section of your system prompt"}},"required":["agent"]},"strict":true}}
zz_schema_check_test.go:13: INVALID: reopen_issue -> {"type":"function","function":{"name":"reopen_issue","description":"Reopen the issue or pull request that triggered this run.","parameters":{"type":""},"strict":true}}
--- PASS: TestSchemaCheck (0.00s)
PASS
ok github.com/abrander/zoo/internal/tools 0.003s
Thinking
Confirmed: **all** tools are serialized as `"parameters":{"type":""}` β the `Type` field is never set anywhere. This is exactly what litellm's pre-call check is rejecting: "Tool 0 function has invalid 'parameters' schema: '' is not valid under any of the given schemas".
This is the root cause of the 429 in the issue. Fix: set `Parameters: Parameters{Type: "object"}` in `NewTool`.
Now I have the complete plan. Let's write the code.
### Implementation
**1. `internal/llm/retry.go`** (new):
```go
package llm
import (
"context"
"errors"
"fmt"
"math/rand/v2"
"net/http"
"strconv"
"time"
)
// APIError is a non-2xx response from the LLM endpoint.
type APIError struct {
// Status is the HTTP status code.
Status int
// Body is the raw response body, for diagnostics.
Body string
// retryAfter is the raw Retry-After header value, if the server
// sent one.
retryAfter string
}
func (e *APIError) Error() string {
return fmt.Sprintf("status: %d: %s", e.Status, e.Body)
}
// RetryAfter returns the delay the server asked for via the
// Retry-After header (delta-seconds or HTTP date), if present and
// parseable.
func (e *APIError) RetryAfter() (time.Duration, bool) {
if e.retryAfter == "" {
return 0, false
}
if secs, err := strconv.Atoi(e.retryAfter); err == nil {
if secs < 0 {
return 0, false
}
return time.Duration(secs) * time.Second, true
}
if t, err := http.ParseTime(e.retryAfter); err == nil {
if d := time.Until(t); d > 0 {
return d, true
}
}
return 0, false
}
// IsRetryable reports whether err is a transient failure worth
// retrying.
//
// HTTP responses: 408, 429, and 5xx are transient. Other statuses
// (400, 401, 403, ...) are client errors that retrying the same
// request will not fix. Anything that is not an APIError β connection
// failures, timeouts, a stream that dropped mid-response β is treated
// as transient. A canceled or deadline-exceeded context is never
// retried: the run is over, and waiting would only delay reporting it.
func IsRetryable(err error) bool {
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
return false
}
var apiErr *APIError
if errors.As(err, &apiErr) {
return apiErr.Status == http.StatusRequestTimeout ||
apiErr.Status == http.StatusTooManyRequests ||
apiErr.Status >= 500
}
return true
}
// Retry configures Do.
type Retry struct {
// MaxAttempts is the total number of attempts, including the first.
MaxAttempts int
// BaseDelay is the wait before the second attempt; each subsequent
// wait doubles, up to MaxDelay.
BaseDelay time.Duration
// MaxDelay caps any single wait.
MaxDelay time.Duration
// BeforeRetry, if non-nil, is called with the 1-based index of the
// attempt that just failed and its error, before the backoff wait.
BeforeRetry func(attempt int, err error)
// Sleep, if non-nil, is used instead of a context-bound sleep
// (tests).
Sleep func(ctx context.Context, d time.Duration) error
}
// DefaultRetry is the policy the agent loop applies to each model
// turn: enough attempts and headroom to ride out typical LLM-gateway
// hiccups (rate limits, model restarts, brief outages) without
// stalling a run for more than a few minutes.
func DefaultRetry() Retry {
return Retry{
MaxAttempts: 10,
BaseDelay: 2 * time.Second,
MaxDelay: 60 * time.Second,
}
}
// Do runs fn, retrying it while it fails with a transient error
// (IsRetryable) and attempts remain. Waits between attempts grow
// exponentially (BaseDelay * 2^(n-1), capped at MaxDelay) with jitter,
// and honor a Retry-After hint carried by an *APIError when it exceeds
// the computed backoff. ctx bounds both the attempts and the waits.
//
// If fn keeps failing until the attempts are exhausted, the returned
// error is the last failure wrapped with the number of attempts made.
func Do[T any](ctx context.Context, r Retry, fn func() (T, error)) (T, error) {
if r.MaxAttempts < 1 {
r.MaxAttempts = 1
}
if r.BaseDelay <= 0 {
r.BaseDelay = 2 * time.Second
}
if r.MaxDelay < r.BaseDelay {
r.MaxDelay = r.BaseDelay
}
sleep := r.Sleep
if sleep == nil {
sleep = func(ctx context.Context, d time.Duration) error {
t := time.NewTimer(d)
defer t.Stop()
select {
case <-ctx.Done():
return ctx.Err()
case <-t.C:
return nil
}
}
}
var zero T
for attempt := 1; ; attempt++ {
if err := ctx.Err(); err != nil {
return zero, err
}
result, err := fn()
if err == nil {
return result, nil
}
if attempt >= r.MaxAttempts || !IsRetryable(err) {
if attempt > 1 {
return zero, fmt.Errorf("after %d attempts: %w", attempt, err)
}
return zero, err
}
if r.BeforeRetry != nil {
r.BeforeRetry(attempt, err)
}
delay := backoffDelay(r.BaseDelay, r.MaxDelay, attempt)
if ra, ok := retryAfter(err); ok && ra > delay {
delay = ra
}
if err := sleep(ctx, delay); err != nil {
return zero, err
}
}
}
// backoffDelay is the wait after the attempt-1..n failure: base *
// 2^(attempt-1) capped at max, with equal jitter β half the delay is
// fixed and half is random β so concurrent runs don't all retry the
// same endpoint in lockstep.
func backoffDelay(base, max time.Duration, attempt int) time.Duration {
d := base
for i := 1; i < attempt && d < max; i++ {
d *= 2
}
if d < 0 || d > max {
d = max
}
half := d / 2
return half + time.Duration(rand.Int63n(int64(half)+1))
}
func retryAfter(err error) (time.Duration, bool) {
var apiErr *APIError
if errors.As(err, &apiErr) {
return apiErr.RetryAfter()
}
return 0, false
}
```
Wait, there's a problem with the `d < 0` overflow check: if base is a huge value and d *= 2 overflows... with the normalization (MaxDelay >= BaseDelay) and the loop condition d < max, d will never exceed max before the check. Actually, if base <= max, then the loop only doubles while d < max, so after the loop d <= max*2 (one extra doubling past max). No overflow can happen (max is at most about 2^63/2... well, a time.Duration is int64 ns; max = 292 years; d < max means d < 2^63/2^... hmm, if max is close to the int64 max, then max*2 could overflow. Realistically max is 60 seconds. The `d < 0` guard is just belt-and-suspenders. OK.
Hmm, one more thing: `rand.Int63n(int64(half)+1)` β if half is huge (close to int64 max), int64(half)+1 overflows. Not a realistic case. No problem.
**2. Changes to `internal/llm/streamer.go`**:
```go
func NewStreamer(resp *http.Response) (*Streamer, error) {
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
resp.Body.Close()
return nil, &APIError{
Status: resp.StatusCode,
Body: string(body),
retryAfter: resp.Header.Get("Retry-After"),
}
}
...
}
```
**3. `internal/llm/types.go`** β NewTool:
```go
func NewTool(name, description string) Tool {
return Tool{
Type: "function",
Function: Function{
Name: name,
Description: description,
// The parameters schema must carry a valid type even when
// the tool takes no arguments: strict validators (e.g.
// litellm's pre-call checks) reject an empty one.
Parameters: Parameters{Type: "object"},
Strict: true,
},
}
}
```
**4. `internal/agentrun/loop.go`**:
Add OnRetry to Hooks + restructure runLoop as planned.
**5. `internal/agentrun/run.go`** β OnRetry in streamHooks:
```go
OnRetry: func(attempt int, err error) {
// Discard the failed attempt's partial output so the retried turn
// is recorded as a fresh block, and drop its (now stale) events
// from the live replay buffer.
reasoningBuf.Reset()
contentBuf.Reset()
reasoningOpen = false
contentOpen = false
r.hub.Checkpoint(jobID)
line := fmt.Sprintf("llm endpoint failed (attempt %d): %s β retrying", attempt, err)
logAppend("system", line)
r.hub.Publish(jobID, livelog.Event{Type: livelog.Retry, Text: line})
},
```
Hmm, wait β about the error text in the line. err is something like "chat completion: status: 429: {big JSON}". For the log line, the body could be huge. Let's truncate the body. Options: use errors.As to get the APIError and format it compactly. Let's write a helper in run.go:
```go
// retryLine renders a failed attempt's error for the job log: API
// errors as "status N: first line of body" (bodies can be long JSON),
// everything else verbatim.
func retryLine(err error) string {
var apiErr *llm.APIError
if !errors.As(err, &apiErr) {
return err.Error()
}
body := apiErr.Body
if i := strings.IndexByte(body, '\n'); i >= 0 {
body = body[:i]
}
if len(body) > 200 {
body = body[:200] + "β¦"
}
return fmt.Sprintf("status %d: %s", apiErr.Status, body)
}
```
Hmm, is this overengineering? The 429 body in the issue is about 1.2KB of JSON. Persisting that per failed attempt (up to 10 times) bloats the job log. Truncation is justified. I'll include it β it's small.
Wait, actually, hmm: "status 429: {\"error\":{\"message\":\"litellm.BadRequestError: ...\"" β even the first 200 chars of JSON is ugly but informative. OK.
**6. `internal/livelog/hub.go`**: add `Retry Type = "retry"`.
**7. `internal/web/templates.go`** JS:
```js
case "retry":
// The failed attempt's partial block is stale: drop it so the
// retried turn renders as a fresh block.
if (reasoningBody) { reasoningBody.parentElement.remove(); reasoningBody = null; }
if (contentBody) { contentBody.parentElement.remove(); contentBody = null; }
var sys = document.createElement("div");
sys.className = "block block-system";
sys.textContent = ev.text;
log.appendChild(sys);
break;
```
**8. Tests**:
a) `internal/llm/retry_test.go`:
- TestIsRetryable: table of cases: APIError 429 β true, 500/502/503 β true, 408 β true, 400/401/403/404 β false, wrapped APIError (fmt.Errorf %w) β true, context.Canceled β false, context.DeadlineExceeded β false, url.Error (connection refused) β true, plain error β true.
- TestDoRetriesTransientErrors: fn fails twice with APIError 500, then succeeds; Sleep is a no-op that records delays; assert 3 calls, result is correct, BeforeRetry called twice with attempts 1,2.
- TestDoStopsOnPermanentError: fn fails with 400 β 1 call, error is unwrapped (no "after" prefix).
- TestDoExhaustsAttempts: always fails with 500, MaxAttempts 3 β 3 calls, error contains "after 3 attempts".
- TestDoHonorsRetryAfter: fn fails with APIError{Status:429, retryAfter:"30"}; Sleep records the delay; assert delay == 30s.
- TestDoContextCanceledDuringSleep: ctx with cancel; Sleep returns ctx.Err()... actually, test that if ctx is canceled during sleep, Do returns ctx.Err. Use a real Sleep with a short delay and cancel after 1ms? Simpler: inject a Sleep that cancels the ctx and returns ctx.Err()? That would be testing my own fake. Better: use the real default sleep (nil) with BaseDelay 50ms, cancel the ctx after 10ms via a goroutine... timing-dependent. Alternative: inject a Sleep that does `select { case <-ctx.Done(): return ctx.Err(); case <-time.After(0): }`... hmm.
Cleaner: a test where Sleep returns a context error:
```go
r.Sleep = func(ctx context.Context, d time.Duration) error {
return context.Canceled
}
```
This verifies Do propagates the sleep error. Plus a separate test for the real timer path:
```go
ctx, cancel := context.WithCancel(context.Background())
r := Retry{MaxAttempts: 5, BaseDelay: time.Hour} // long
go func() { time.Sleep(20*time.Millisecond); cancel() }()
start := time.Now()
_, err := Do(ctx, r, func() (int, error) { return 0, &APIError{Status: 500} })
// err should be context.Canceled, and elapsed < 1s
```
With BaseDelay 1h, if the timer weren't ctx-aware, the test would hang for an hour. The default sleep uses a timer + select, so cancel β ctx.Err(). Elapsed about 20ms. This is a real behavior test.
- TestBackoffDelayBounds: for attempts 1..12 with base 2s, max 60s: delay is in [d/2, d], where d = min(60s, 2s*2^(a-1)); monotonic non-decreasing caps.
- TestAPIErrorRetryAfter: seconds form, date form, invalid, negative, empty.
b) `internal/llm/streamer_test.go`: extend TestStreamerNonOKStatus:
```go
func TestStreamerNonOKStatus(t *testing.T) {
resp := newTestResponse("boom")
resp.StatusCode = http.StatusTooManyRequests
resp.Header = http.Header{"Retry-After": ["30"]}
_, err := NewStreamer(resp)
var apiErr *APIError
if !errors.As(err, &apiErr) { t.Fatal(...) }
if apiErr.Status != 429 || apiErr.Body != "boom" { ... }
if !IsRetryable(err) { ... }
ra, ok := apiErr.RetryAfter(); if !ok || ra != 30*time.Second { ... }
}
```
Keep the 500 case too.
c) `internal/llm/types_test.go` (new): TestNewToolSchema:
```go
func TestNewToolParameterSchemaIsValid(t *testing.T) {
tool := NewTool("close_issue", "Close the issue.")
b, err := json.Marshal(tool)
...
var parsed struct {
Function struct {
Parameters struct {
Type string `json:"type"`
} `json:"parameters"`
} `json:"function"`
}
json.Unmarshal(b, &parsed)
if parsed.Function.Parameters.Type != "object" { ... }
// and with a property
tool2 := NewTool("x", "y")
tool2.AddStringProperty("a", "b", true)
... type is still "object"
}
```
d) `internal/agentrun/loop_test.go` (new): test runLoop's retry behavior against an httptest server.
```go
package agentrun
import (
"context"
"fmt"
"io"
"net/http"
"net/http/httptest"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/abrander/zoo/internal/config"
"github.com/abrander/zoo/internal/llm"
)
// sseStop is a minimal valid streaming completion that ends the turn.
const sseStop = `data: {"choices":[{"index":0,"delta":{"role":"assistant","content":"done"},"finish_reason":"stop"}]}` + "\n\n" + "data: [DONE]\n\n"
func newTestLLMClient(t *testing.T, handler http.HandlerFunc) *llm.Client {
t.Helper()
srv := httptest.NewServer(handler)
t.Cleanup(srv.Close)
return llm.NewClient(config.LLM{Name: "test", OpenAI: srv.URL, Token: "test", Model: "test"})
}
func TestRunLoopRetriesTransientErrors(t *testing.T) {
var calls int32
client := newTestLLMClient(t, func(w http.ResponseWriter, r *http.Request) {
if atomic.AddInt32(&calls, 1) <= 2 {
w.WriteHeader(http.StatusBadGateway)
io.WriteString(w, "engine not healthy")
return
}
w.Header().Set("Content-Type", "text/event-stream")
io.WriteString(w, sseStop)
})
// Shorten the backoff so the test doesn't actually wait.
oldPolicy := retryPolicy
retryPolicy = llm.DefaultRetry()
retryPolicy.Sleep = func(ctx context.Context, d time.Duration) error { return nil }
t.Cleanup(func() { retryPolicy = oldPolicy })
var retries []int
hooks := Hooks{OnRetry: func(attempt int, err error) { retries = append(retries, attempt) }}
messages, err := runLoop(context.Background(), client, nil, []llm.Message{{Role: "user", Content: "hi"}}, hooks)
if err != nil { t.Fatalf("expected retries to succeed, got: %v", err) }
if got := atomic.LoadInt32(&calls); got != 3 { t.Fatalf("expected 3 attempts, got %d", got) }
if len(retries) != 2 || retries[0] != 1 || retries[1] != 2 { t.Fatalf("unexpected retry callbacks: %v", retries) }
if len(messages) != 2 || messages[1].Content != "done" { t.Fatalf("unexpected messages: %+v", messages) }
}
```
Wait β runLoop's signature: `runLoop(ctx, client, toolsCtx tools.Context, messages, hooks)`. toolsCtx is an interface; passing nil is fine (no tool calls happen).
Also, the OnRetry hook in the test: check the error text contains "status: 502".
- TestRunLoopStopsOnPermanentError: server always returns 400 β 1 call, error contains "status: 400", no "after N attempts".
- TestRunLoopExhaustsRetries: server always returns 500, MaxAttempts 3 β 3 calls, error contains "after 3 attempts".
- TestRunLoopRetriesMidStreamFailure: server: first request β 200 with one chunk, then connection close (no [DONE]); second request β full sseStop. How to simulate a mid-stream drop: write one chunk, then hijack and close? Simpler: use http.ResponseWriter... to cut the connection: `hijacker, ok := w.(http.Hijacker); c, _, _ := hjacker.Hijack(); c.Close()`. After writing the chunk and flushing. The client's scanner will get a read error or a clean EOF. A clean EOF β scanner.Scan() returns false, Err is nil β Get() proceeds to json.Unmarshal("") β decode error β retryable. A hijack close might give "EOF" or "connection reset" β either way, retryable.
Hmm, wait, let me actually trace Get() on a clean EOF:
```go
for s.scanner.Scan() { ... } // exits on EOF, buf = last line (chunk line?)
```
Wait, the loop: `buf = bytes.TrimSpace(s.scanner.Bytes())` β on the last successful Scan, buf = the chunk line. If the chunk line has the "data: " prefix β break. So if the server sent a valid chunk and then closed, Get() returns that chunk successfully (break on the data: prefix). The next Get() call: scanner.Scan() immediately returns false (EOF), buf is still the previous chunk's line (stale!). Then TrimPrefix β chunk JSON, not [DONE] β unmarshal succeeds β Merge the same chunk again β returns the same accumulated completion, no error!! Infinite loop??
Wait, no β let me re-read Get():
```go
func (s *Streamer) Get() (*ChatCompletion, error) {
if s.eof {
return nil, io.EOF
}
prefix := []byte("data: ")
var buf []byte
for s.scanner.Scan() {
if err := s.scanner.Err(); err != nil {
return nil, err
}
buf = bytes.TrimSpace(s.scanner.Bytes())
if len(buf) == 0 {
continue
}
if bytes.HasPrefix(buf, prefix) {
break
}
}
buf = bytes.TrimPrefix(buf, prefix)
if bytes.Equal(buf, []byte("[DONE]")) {
s.resp.Body.Close()
s.eof = true
return nil, io.EOF
}
var completion ChatCompletion
if err := json.Unmarshal(buf, &completion); err != nil {
return nil, fmt.Errorf("decode completion chunk: %w: %s", err, buf)
}
s.completion.Merge(&completion)
return &s.completion, nil
}
```
Scenario: the server sends chunk1\n\n and then closes.
- Get() #1: Scan β chunk1 line (buf=chunk1), break. Unmarshal OK. Return completion1.
- Get() #2: Scan β false (EOF). buf is still chunk1 (stale from the previous call? no wait β buf is a local variable, `var buf []byte`, fresh each call!).
Ah right, `buf` is local. So Get() #2: the loop doesn't execute, buf = nil. TrimPrefix(nil) = nil. Not [DONE]. json.Unmarshal(nil) β "unexpected end of JSON input" β error "decode completion chunk: unexpected end of JSON input: ". Retryable (not an APIError, not a ctx error).
So a clean mid-stream EOF β decode error β retried. A mid-stream connection reset β scanner.Err() β retryable. Both covered.
So my mid-stream test: request 1 β write a chunk, flush, then hijack+close (or just return from the handler β when the handler returns, net/http closes the body... actually, if the handler returns without writing [DONE], the connection is closed after the handler returns (with HTTP/1.1's default keep-alive, it stays open, but the body is done β the client sees a clean EOF). Either way, the client sees an EOF.
So:
```go
var calls int32
handler: if calls==1 { w.WriteHeader(200); w.Write(chunk1); w.(http.Flusher).Flush(); return } else { w.Write(sseStop) }
```
When the handler returns, the response body is complete (no [DONE]) β the client's next Get() β decode error β retry.
But wait β there's a subtlety with the hooks: the first (failed) attempt streams a content delta "partial". OnRetry fires (the test hook records it). The second attempt streams "done". The final messages: [user, assistant "done"]? Let me trace: attempt 1's completion is partial ("partial" content, no finish reason). The stream "completes" with an error β the closure returns an error β the completion variable is not used (it's reset to nil at the start of the next attempt). Attempt 2: full sseStop β completion "done", finish stop. messages = [user, {assistant done}].
Also, verify the delta behavior: OnContentDelta fires with "partial" (attempt 1), then "done" (attempt 2 β the cursor is reset). Good, the test can assert the sequence of deltas: ["partial", "done"].
e) Also, a test that the empty completion is retried? Server returns 200 with `data: [DONE]` only β no choices β "model returned an empty completion" β retryable β retry. Worth a quick test? It pins down the behavior. I'll include it.
Now β the `retryPolicy` package variable in loop.go:
```go
// retryPolicy is the retry policy runLoop applies to each model turn.
// A package variable rather than a parameter so tests can shorten the
// backoff; production always uses the default.
var retryPolicy = llm.DefaultRetry()
```
And runLoop uses `retry := retryPolicy; retry.BeforeRetry = ...`.
Hmm, wait, one consideration: runLoop is called once per run; the policy variable is read at the start of runLoop. Tests set it before calling runLoop. No problem.
Also double-check the `llm.Do` generic instantiation in loop.go: `llm.Do(ctx, retry, func() error {...})` β T is inferred as... the closure returns `error` only β T is inferred from the return: `func() error` doesn't match `func() (T, error)`. Need an explicit type parameter or a two-value return. Options: `llm.Do[error](ctx, retry, func() (error, error) {...})` β ugly. Better: make the closure return (struct{}, error)? Also ugly.
Alternative: have Do take `fn func() error` and return only an error... but then the turn's result (completion) has to escape via a closure variable β which is what we're already doing (the `completion` variable is captured). So a simple `Do(ctx, retry, func() error {...}) error` suffices! The turn's result is assigned to the captured `completion` variable.
Should Do be generic at all? The only consumer needs `func() error`. A generic adds API surface with no current benefit. KISS: non-generic `Do(ctx context.Context, r Retry, fn func() error) error`.
Hmm, but a generic Do is more reusable for the future... YAGNI. Non-generic. Decision:
```go
func Do(ctx context.Context, r Retry, fn func() error) error
```
And in loop.go:
```go
err := llm.Do(ctx, retry, func() error {
...
return nil
})
if err != nil {
return messages, err
}
```
Clean.
Now, let me also reconsider the "after %d attempts" wrapping in Do with the non-generic signature. Same as before.
One more point about `IsRetryable` and errors from `http.Client.Do`: e.g. `Post "http://...": dial tcp: connection refused` β a *url.Error wrapping *net.OpError β not an APIError, not a ctx error β retryable. Good. A TLS handshake failure β retryable (transient). A DNS failure β retryable (may be transient). OK.
What about `*net.DNSError` with IsNotFound (a typo'd endpoint)? It'd be retried 10 times over ~5 minutes before failing. A bit wasteful, but the endpoint comes from config; a typo fails fast enough... 5 minutes is not great, but acceptable. We could special-case DNSError.IsNotFound β not retryable. Hmm, that's a nice touch:
```go
var dnsErr *net.DNSError
if errors.As(err, &dnsErr) && dnsErr.IsNotFound {
return false
}
```
I'll include it β a nonexistent host will never come back; retrying just delays the obvious failure. Wait, but a transient DNS outage also returns a DNSError... with IsNotFound? No β a transient DNS failure is a timeout/refused (IsNotFound is false). IsNotFound is specifically "no such host". OK, include it.
Hmm, actually, let's keep IsRetryable focused. Adding the DNSError case is 4 lines and clearly correct. Include.
Now β the run.go changes. Current streamHooks:
```go
func (r *Runner) streamHooks(jobID string, logAppend func(stream, line string)) Hooks {
var reasoningBuf, contentBuf strings.Builder
reasoningOpen, contentOpen := false, false
return Hooks{
OnReasoningDelta: ...,
OnContentDelta: ...,
OnTurnEnd: ...,
OnTool: ...,
}
}
```
Add OnRetry. run.go imports: need "errors" (for retryLine) β check current imports: context, embed, encoding/json, fmt, log/slog, os, path/filepath, strings, sync, time, config, forgejo, livelog, llm, store. Need to add "errors". llm is already imported.
Now, the line format: "llm endpoint failed (attempt 1): status 502: engine not healthy β retrying". Hmm, em-dash in the codebase? The codebase uses "β" in comments (run.go: "so its git activity is attributed..."). No problem.
Wait, actually β let me reconsider the wording. "attempt 1 failed" β attempt is the index of the failed attempt. "attempt 1: ... β retrying" is clear.
Also, should OnRetry log via slog too? The runner's logger has a job context. streamHooks doesn't take a logger... it takes logAppend. The persisted line is enough (it shows in the dashboard). The scheduler's error log will show the final failure. OK.
Now β the web JS. Also need to double-check: the `retry` event is published to the hub and buffered; if the job ends while a retry is pending (ctx canceled during sleep), the Status event is published β the page reloads β the static view shows the persisted retry line. Good.
Where does the scheduler publish Status? Let me check the rest of scheduler.go (lines 200β252) to confirm the hub publish on job end.
I have the full picture now. Let me check the tail of the scheduler (job-end hub publish) before writing the implementation.
Tool
π§ read_file
{"path":"internal/scheduler/scheduler.go","start_line":200,"end_line":252}
Arguments
{"path":"internal/scheduler/scheduler.go","start_line":200,"end_line":252}
Result
{"path":"internal/scheduler/scheduler.go","content":"\t\tlogger.Error(\"agent run failed\", \"status\", status, \"error\", err)\n\t\ts.reportFailure(ev, status, err)\n\t} else {\n\t\tlogger.Info(\"agent run succeeded\")\n\t}\n\n\tif err := s.store.MarkJobFinished(context.Background(), jobID, status, errMsg); err != nil {\n\t\tlogger.Error(\"failed to mark job finished\", \"error\", err)\n\t}\n\n\t// Closing the hub here, only after the terminal status is durably\n\t// written, is what lets a dashboard tab's auto-reload (triggered by\n\t// this Close) see the finished status immediately. Closing it inside\n\t// Runner.Run instead β before this status write β raced a fast\n\t// reload against MarkJobFinished: the reloaded page could still see\n\t// \"running\", open a fresh SSE subscription against a hub entry that\n\t// no longer exists, and be stuck showing a stale status forever.\n\ts.hub.Close(jobID)\n}\n\n// reportFailure implements the locked-in failure policy: comment on the\n// triggering issue/PR with an error summary, and apply FailureLabel.\nfunc (s *Scheduler) reportFailure(ev forgejo.Event, status store.JobStatus, runErr error) {\n\tbody := fmt.Sprintf(\"zoo agent run %s: %s\", status, runErr.Error())\n\n\tif err := s.forgejo.CreateIssueComment(ev.Owner, ev.Repo, ev.Index, body); err != nil {\n\t\ts.logger.Error(\"failed to post failure comment\", \"owner\", ev.Owner, \"repo\", ev.Repo, \"index\", ev.Index, \"error\", err)\n\t}\n\n\tif err := s.forgejo.AddLabel(ev.Owner, ev.Repo, ev.Index, FailureLabel); err != nil {\n\t\ts.logger.Error(\"failed to add failure label\", \"owner\", ev.Owner, \"repo\", ev.Repo, \"index\", ev.Index, \"error\", err)\n\t}\n}\n\n// Shutdown stops accepting new work is the caller's responsibility (stop\n// feeding Run's channel); Shutdown just waits for in-flight runs, up to\n// ctx's deadline.\nfunc (s *Scheduler) Shutdown(ctx context.Context) error {\n\tdone := make(chan struct{})\n\n\tgo func() {\n\t\ts.wg.Wait()\n\t\tclose(done)\n\t}()\n\n\tselect {\n\tcase \u003c-done:\n\t\treturn nil\n\n\tcase \u003c-ctx.Done():\n\t\treturn ctx.Err()\n\t}\n}","start_line":200,"end_line":252,"total_lines":252,"truncated":false}