mirror of
https://github.com/multica-ai/multica.git
synced 2026-07-28 14:09:22 +02:00
* fix(agent): attribute Grok usage from the turn's own model id A resumed Grok session with no configured model recorded its entire spend under the model id "unknown", which matches no pricing row — so the task reported $0 cost instead of its real spend. grok.go only learned the model from the session handshake, and ACP's `session/load` carries no model id (only `session/new` does). When neither the agent nor MULTICA_GROK_MODEL pins a model, `daemon.go` legitimately passes an empty model, leaving nothing to attribute the usage to. Every Grok turn stamps `result._meta.modelId` with what it actually billed against. Parse it in the shared ACP result parser and use it as the fallback in grok.go. Other ACP backends are untouched — they keep whatever the handshake gave them. Co-authored-by: J <agent@multica.ai> Co-authored-by: multica-agent <github@multica.ai> * fix(metrics): price the Grok catalog in server-side cost metrics server/internal/metrics/pricing.go carried no Grok rows at all, so RecordLLMUsage took the unpriced branch for every Grok turn: llm_cost_usd reported zero Grok spend while the tokens accumulated in llm_unpriced_tokens. Internal cost monitoring simply could not see Grok. Add the six SKUs xAI publishes rates for, mirroring the frontend table in packages/views/runtimes/utils.ts. Aliases are anchored exact matches like the gpt-5.6 rows, so `grok-composer-*` (in the catalog, absent from the price sheet) stays unmapped instead of inheriting a guessed rate. Short-context tier on purpose: xAI bills a request at 2x once its prompt reaches 200K tokens, but a usage record aggregates every model call in a turn and cannot say which tier an individual request hit. A regression test re-derives the cost of a real grok 0.2.106 turn from the table and checks it against the costUsdTicks xAI returned for that turn. Co-authored-by: J <agent@multica.ai> Co-authored-by: multica-agent <github@multica.ai> * docs(changelog): scope the Grok cost claim to what was actually fixed The v0.4.9 entry promised "accurate cost" in all four languages, but the fix corrected catalog pricing and cached-input double-counting — it did not implement xAI's 2x long-context tier, so a turn whose requests reach 200K prompt tokens still under-reports by up to 50%. Say what was fixed instead. Also correct two stale claims in the pricing comment: the daemon tags usage rows with the runtime provider `grok`, not `xai` (the bare `grok-*` keys are what make them resolve), and record why thresholding the long-context tier on an aggregated row would be worse than not pricing it at all. Co-authored-by: J <agent@multica.ai> Co-authored-by: multica-agent <github@multica.ai> * feat(usage): carry the provider's own cost through to the usage record Cost has always been derived client-side as tokens x a static rate, which cannot express request-level pricing rules. xAI bills a Grok request at 2x once its prompt reaches 200K tokens, and a task_usage row aggregates every model call in a turn — so the stored token counts genuinely cannot say which tier any individual request hit. Thresholding on the aggregate would be worse than the status quo: it turns a bounded 50% under-estimate into an unbounded over-estimate for turns made of many short requests. Grok already reports what it charged, per turn, in `_meta.usage.costUsdTicks`. Parse it, carry it through agent -> daemon -> API, and store it on task_usage as a nullable BIGINT of 1e-10 USD ticks (integer, so sub-cent turns stay exact end to end). NULL means the provider reported no cost — every pre-existing row and every provider that doesn't return one. No backfill: there is no authoritative figure to recover for those, and inventing one is the guess this removes. A single hourly bucket can mix rows that carry a cost with rows that don't, so task_usage_hourly gains both halves: `cost_usd_ticks` sums the authoritative side, and `uncosted_*_tokens` carry exactly the tokens that still need a rate-table estimate. Consumers report authoritative + estimate(uncosted), which degrades to today's behaviour when nothing in the bucket is authoritative. The existing token columns keep covering every row, so token displays are untouched. The new columns are additive with defaults, so the unique key, the dirty-queue shape, and migration 102's triggers are unaffected. Co-authored-by: J <agent@multica.ai> Co-authored-by: multica-agent <github@multica.ai> * feat(usage): prefer the provider's own cost over the rate table With the authoritative figure now stored, both cost consumers use it: the usage dashboard (estimateCost / estimateCostBreakdown) and the server-side llm_cost_usd metric. Each reports `authoritative + estimate(uncosted tokens)`, so a row or bucket that mixes priced and unpriced sources stays whole. The static rate tables remain, but for Grok they are now a fallback — they still price usage recorded by a daemon too old to report cost, and every provider that reports none. Custom pricing overrides likewise apply only to the estimated half: they are a user's guess at a rate, and the authoritative half is not a guess. A model with no rate-table row but a provider-reported cost now also drops out of the "unmapped models" banner, since asking the user to supply a rate for it would invite overriding a real bill. llm_cost_usd is labelled by token_type and the provider reports one number per turn, so the charge is distributed across the buckets in the rate table's own proportions. Only the total is authoritative; the split stays an estimate, which is why this scales the existing buckets rather than inventing a label. estimateCostBreakdown does the same, keeping the stacked chart summing to the headline figure instead of silently under-drawing every Grok row. Co-authored-by: J <agent@multica.ai> Co-authored-by: multica-agent <github@multica.ai> * docs(changelog): say Grok cost now follows xAI's actual charge The earlier wording scoped the claim down to catalog pricing and cached input because the long-context tier was still unhandled. It is handled now — the cost comes from what xAI charged for the turn — so the entry can say so. Co-authored-by: J <agent@multica.ai> Co-authored-by: multica-agent <github@multica.ai> * fix(usage): keep the provider's cost when the model has no rate row Both cost consumers bailed out before reading the authoritative figure when the rate table had no row for the model. A `grok-composer-*` turn — in the Grok Build catalog, absent from xAI's price sheet — was therefore reported as $0 spend even though xAI told us exactly what it charged. Worse on the client: estimateCost returned the real cost while estimateCostBreakdown returned zeros, so the headline and the stacked chart disagreed on precisely the rows whose cost is exact — and the unmapped-models banner was (correctly) hidden, so nothing explained the discrepancy. Handle the charge before the rate lookup in both places. Without rates there is nothing to split a total by, so it lands whole in the `input` bucket, the same fallback distributeAuthoritativeCost already uses when it has no shape to scale. Tokens with no rate keep going to llm_unpriced_tokens: "unpriced" describes the rate table, not the money. Co-authored-by: J <agent@multica.ai> Co-authored-by: multica-agent <github@multica.ai> * perf(usage): drop the historical rewrite from the cost-split migration Migration 213 rewrote every existing task_usage_hourly row to seed the uncosted counters. That is a full-table UPDATE inside a schema migration — lock time, WAL and bloat all scaling with table size — for rows this issue explicitly does not care about. Deleting the UPDATE alone would have zeroed historical cost: with `NOT NULL DEFAULT 0`, an untouched row asserts "nothing here needs estimating", so every pre-split bucket would report $0 until the rollup happened to touch it. Make the uncosted columns nullable with no default instead. NULL means "never recomputed since the split existed", readers COALESCE it to the row's own token total ("estimate all of it"), and the pre-split behaviour is preserved exactly — with nothing to seed, so no rewrite. A bare ADD COLUMN is metadata-only, so this is now fast DDL. Rows heal into the split naturally as the rollup recomputes their buckets. Verified on a fresh database: a legacy-shaped row reads back as its full tokens to estimate, and a group mixing legacy and post-split buckets sums to the authoritative cost plus both rows' estimable tokens. Co-authored-by: J <agent@multica.ai> Co-authored-by: multica-agent <github@multica.ai> --------- Co-authored-by: Bohan-J <bohan@devv.ai> Co-authored-by: J <agent@multica.ai> Co-authored-by: multica-agent <github@multica.ai>
2359 lines
80 KiB
Go
2359 lines
80 KiB
Go
package agent
|
|
|
|
import (
|
|
"bufio"
|
|
"bytes"
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"os"
|
|
"os/exec"
|
|
"path/filepath"
|
|
"regexp"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
)
|
|
|
|
// hermesBlockedArgs are flags hardcoded by the daemon that must not be
|
|
// overridden by user-configured custom_args. `acp` is the protocol
|
|
// subcommand that drives the ACP JSON-RPC transport; overriding it
|
|
// would break the daemon↔Hermes communication contract.
|
|
//
|
|
// `-p`/`--profile` are NOT stripped unconditionally: a skill-less Hermes task
|
|
// has no overlay, so its profile selection must pass through to Hermes
|
|
// unchanged. The daemon strips the selected occurrence via StripHermesProfileArgs
|
|
// only when it actually built the per-task overlay (see the daemon's launch-arg
|
|
// handling), so the flag can't re-point HERMES_HOME past the overlay while
|
|
// leaving no-overlay tasks' behavior untouched.
|
|
var hermesBlockedArgs = map[string]blockedArgMode{
|
|
"acp": blockedStandalone,
|
|
}
|
|
|
|
// hermesArgProfileRe mirrors the space-form guard in Hermes'
|
|
// hermes_cli.main._apply_profile_override step 1b: a `-p <value>` whose value
|
|
// doesn't match the profile-id shape is not a profile selection at all (e.g.
|
|
// pytest's `-p no:xdist`), so it is ignored rather than consumed. The inline
|
|
// `--profile=<value>` form is NOT guarded here — Hermes forwards it verbatim to
|
|
// resolve_profile_env, which validates and hard-fails on an invalid value; that
|
|
// validation lives in the daemon-side resolver (execenv.ResolveHermesProfile).
|
|
var hermesArgProfileRe = regexp.MustCompile(`^[a-z0-9][a-z0-9_-]{0,63}$`)
|
|
|
|
// hermesValueFlags and hermesOptionalValueFlags mirror the value-taking flags
|
|
// Hermes skips while scanning argv for -p/--profile, so a value like the `coder`
|
|
// in `-m coder -p research` is never misread as the profile. Kept in sync with
|
|
// _apply_profile_override.value_flags / optional_value_flags.
|
|
var hermesValueFlags = map[string]struct{}{
|
|
"-z": {}, "--oneshot": {}, "-m": {}, "--model": {}, "--provider": {},
|
|
"-t": {}, "--toolsets": {}, "-r": {}, "--resume": {}, "-s": {},
|
|
"--skills": {}, "--usage-file": {},
|
|
}
|
|
var hermesOptionalValueFlags = map[string]struct{}{"-c": {}, "--continue": {}}
|
|
|
|
// HermesProfileSelection is the profile selection parsed out of custom_args by
|
|
// ParseHermesProfileArgs. It carries the exact argv occurrence to consume so the
|
|
// daemon-side resolver and the launch-arg stripping act on one authoritative
|
|
// parse instead of each re-approximating Hermes' argv handling.
|
|
type HermesProfileSelection struct {
|
|
Name string // the selected value; "" for the empty inline `--profile=` value
|
|
Found bool // a -p/--profile selection with a value was matched
|
|
Inline bool // matched the `--profile=<value>` form (value validated downstream)
|
|
ArgFrom int // index of the first token to strip, or -1 when nothing matched
|
|
ArgLen int // tokens to strip: 2 for `-p <value>`, 1 for `--profile=<value>`
|
|
}
|
|
|
|
// ParseHermesProfileArgs finds the first Hermes profile selection in custom_args,
|
|
// mirroring hermes_cli.main._apply_profile_override step 1/1b: it scans for the
|
|
// first `-p`/`--profile <value>` or `--profile=<value>`, skipping value-taking
|
|
// flags and stopping at a `--` sentinel or an `mcp add --args` command-argv
|
|
// passthrough region. A space-form value that fails the profile-id shape is
|
|
// ignored (matches Hermes discarding it). Args are unquoted with the same helper
|
|
// as filterCustomArgs so quoting is handled consistently.
|
|
func ParseHermesProfileArgs(args []string) HermesProfileSelection {
|
|
none := HermesProfileSelection{ArgFrom: -1}
|
|
i := 0
|
|
for i < len(args) {
|
|
arg := unshellQuoteArg(args[i])
|
|
if arg == "--" {
|
|
break
|
|
}
|
|
if arg == "--args" && hermesInsideMcpAdd(args, i) {
|
|
break
|
|
}
|
|
if arg == "-p" || arg == "--profile" {
|
|
if i+1 < len(args) {
|
|
val := unshellQuoteArg(args[i+1])
|
|
if !hermesArgProfileRe.MatchString(val) {
|
|
return none // step 1b: not a valid profile value, ignore
|
|
}
|
|
return HermesProfileSelection{Name: val, Found: true, ArgFrom: i, ArgLen: 2}
|
|
}
|
|
return none // trailing flag with no value
|
|
}
|
|
if v, ok := strings.CutPrefix(arg, "--profile="); ok {
|
|
return HermesProfileSelection{Name: v, Found: true, Inline: true, ArgFrom: i, ArgLen: 1}
|
|
}
|
|
if _, ok := hermesValueFlags[arg]; ok && i+1 < len(args) {
|
|
i += 2
|
|
continue
|
|
}
|
|
if _, ok := hermesOptionalValueFlags[arg]; ok && i+1 < len(args) &&
|
|
!strings.HasPrefix(unshellQuoteArg(args[i+1]), "-") {
|
|
i += 2
|
|
continue
|
|
}
|
|
i++
|
|
}
|
|
return none
|
|
}
|
|
|
|
// hermesInsideMcpAdd reports whether argv position index sits inside an
|
|
// `mcp add ... --args <child argv>` passthrough region, where flags belong to
|
|
// the child MCP command and must not be read as Hermes' own profile selector.
|
|
func hermesInsideMcpAdd(args []string, index int) bool {
|
|
mcp := -1
|
|
for j := 0; j < index; j++ {
|
|
if unshellQuoteArg(args[j]) == "mcp" {
|
|
mcp = j
|
|
break
|
|
}
|
|
}
|
|
if mcp < 0 {
|
|
return false
|
|
}
|
|
for j := mcp + 1; j < index; j++ {
|
|
if unshellQuoteArg(args[j]) == "add" {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// StripHermesProfileArgs removes exactly the argv occurrence ParseHermesProfileArgs
|
|
// selected. The daemon calls this only when it built the per-task overlay, so
|
|
// Hermes uses the overlay's HERMES_HOME instead of re-resolving the profile —
|
|
// while a skill-less task keeps its flags untouched.
|
|
func StripHermesProfileArgs(args []string, sel HermesProfileSelection) []string {
|
|
if !sel.Found || sel.ArgFrom < 0 || sel.ArgLen <= 0 {
|
|
return args
|
|
}
|
|
end := sel.ArgFrom + sel.ArgLen
|
|
if end > len(args) {
|
|
end = len(args)
|
|
}
|
|
out := make([]string, 0, len(args)-(end-sel.ArgFrom))
|
|
out = append(out, args[:sel.ArgFrom]...)
|
|
out = append(out, args[end:]...)
|
|
return out
|
|
}
|
|
|
|
// hermesBackend implements Backend by spawning `hermes acp` and communicating
|
|
// via the ACP (Agent Communication Protocol) JSON-RPC 2.0 over stdin/stdout.
|
|
// This is the same pattern as Codex but with the ACP protocol instead of
|
|
// the Codex-specific JSON-RPC methods.
|
|
type hermesBackend struct {
|
|
cfg Config
|
|
}
|
|
|
|
var (
|
|
hermesReaderDrainGrace = 2 * time.Second
|
|
hermesNotificationQuietTime = 250 * time.Millisecond
|
|
)
|
|
|
|
func (b *hermesBackend) Execute(ctx context.Context, prompt string, opts ExecOptions) (*Session, error) {
|
|
execPath := b.cfg.ExecutablePath
|
|
if execPath == "" {
|
|
execPath = "hermes"
|
|
}
|
|
if _, err := exec.LookPath(execPath); err != nil {
|
|
return nil, fmt.Errorf("hermes executable not found at %q: %w", execPath, err)
|
|
}
|
|
|
|
// Translate the agent's mcp_config (Claude-style object of objects)
|
|
// into the array shape ACP `session/new` expects. Fail closed on
|
|
// malformed JSON so the launch surfaces the real error instead of
|
|
// silently dropping all MCP servers.
|
|
mcpServers, err := buildACPMcpServers(opts.McpConfig, b.cfg.Logger)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("hermes: invalid mcp_config: %w", err)
|
|
}
|
|
|
|
timeout := opts.Timeout
|
|
runCtx, cancel := runContext(ctx, timeout)
|
|
|
|
hermesArgs := append([]string{"acp"}, filterCustomArgs(opts.CustomArgs, hermesBlockedArgs, b.cfg.Logger)...)
|
|
cmd := exec.CommandContext(runCtx, execPath, hermesArgs...)
|
|
hideAgentWindow(cmd)
|
|
b.cfg.Logger.Info("agent command", "exec", execPath, "args", hermesArgs)
|
|
agentsMDPresent := false
|
|
if opts.Cwd != "" {
|
|
cmd.Dir = opts.Cwd
|
|
if _, err := os.Stat(filepath.Join(opts.Cwd, "AGENTS.md")); err == nil {
|
|
agentsMDPresent = true
|
|
}
|
|
}
|
|
b.cfg.Logger.Info("hermes acp starting", "cwd", opts.Cwd, "agents_md_present", agentsMDPresent)
|
|
if opts.SystemPrompt != "" {
|
|
b.cfg.Logger.Debug("hermes ignoring ExecOptions.SystemPrompt; using cwd-scoped context files", "cwd", opts.Cwd)
|
|
}
|
|
|
|
env := buildEnv(b.cfg.Env)
|
|
// Enable yolo mode so Hermes auto-approves all tool executions.
|
|
env = append(env, "HERMES_YOLO_MODE=1")
|
|
cmd.Env = env
|
|
|
|
stdout, err := cmd.StdoutPipe()
|
|
if err != nil {
|
|
cancel()
|
|
return nil, fmt.Errorf("hermes stdout pipe: %w", err)
|
|
}
|
|
stdin, err := cmd.StdinPipe()
|
|
if err != nil {
|
|
cancel()
|
|
return nil, fmt.Errorf("hermes stdin pipe: %w", err)
|
|
}
|
|
// Forward stderr to the daemon log *and* sniff provider-level
|
|
// errors out of it so we can surface them in the task result.
|
|
// Hermes' session/prompt still reports stopReason=end_turn when
|
|
// the underlying HTTP call to the LLM returns 4xx/5xx, so
|
|
// without this we'd report a misleading "empty output" and hide
|
|
// the real cause (wrong model for the current provider, bad
|
|
// credentials, rate limit, …) in the daemon log.
|
|
//
|
|
// We use StderrPipe + an explicit copier goroutine instead of
|
|
// `cmd.Stderr = io.MultiWriter(...)` so we have a join point
|
|
// (`stderrDone`) before the failure-promotion decision. With the
|
|
// MultiWriter form, exec's internal copy goroutine is only
|
|
// joined by `cmd.Wait()`, which runs in the deferred cleanup —
|
|
// after `promoteACPResultOnProviderError` already consulted the
|
|
// sniffer. That race lost the 429 / usage-limit message under
|
|
// CI load and surfaced as a flaky test
|
|
// (TestHermesBackendPromotesProviderErrorWithNonEmptyOutput).
|
|
providerErr := newACPProviderErrorSniffer("hermes")
|
|
stderr, err := cmd.StderrPipe()
|
|
if err != nil {
|
|
cancel()
|
|
return nil, fmt.Errorf("hermes stderr pipe: %w", err)
|
|
}
|
|
|
|
if err := cmd.Start(); err != nil {
|
|
cancel()
|
|
return nil, fmt.Errorf("start hermes: %w", err)
|
|
}
|
|
|
|
stderrSink := io.MultiWriter(newLogWriter(b.cfg.Logger, "[hermes:stderr] "), providerErr)
|
|
stderrDone := make(chan struct{})
|
|
go func() {
|
|
defer close(stderrDone)
|
|
_, _ = io.Copy(stderrSink, stderr)
|
|
}()
|
|
|
|
b.cfg.Logger.Info("hermes acp started", "pid", cmd.Process.Pid, "cwd", opts.Cwd)
|
|
|
|
msgCh := make(chan Message, 256)
|
|
resCh := make(chan Result, 1)
|
|
|
|
var outputMu sync.Mutex
|
|
var output strings.Builder
|
|
// streamingCurrentTurn gates all session updates so that history
|
|
// replay (Hermes sends full prior-turn transcripts on session/resume,
|
|
// and may flush queued chunks before our session/prompt response
|
|
// streams) is dropped instead of duplicating the previous answer
|
|
// into output. We flip it to true only after session/prompt is sent.
|
|
var streamingCurrentTurn atomic.Bool
|
|
|
|
promptDone := make(chan hermesPromptResult, 1)
|
|
activity := make(chan struct{}, 1)
|
|
|
|
c := &hermesClient{
|
|
cfg: b.cfg,
|
|
stdin: stdin,
|
|
pending: make(map[int]*pendingRPC),
|
|
pendingTools: make(map[string]*pendingToolCall),
|
|
acceptNotification: func(string) bool {
|
|
return streamingCurrentTurn.Load()
|
|
},
|
|
onActivity: func() {
|
|
select {
|
|
case activity <- struct{}{}:
|
|
default:
|
|
}
|
|
},
|
|
onMessage: func(msg Message) {
|
|
if !streamingCurrentTurn.Load() {
|
|
return
|
|
}
|
|
if msg.Type == MessageText {
|
|
outputMu.Lock()
|
|
output.WriteString(msg.Content)
|
|
outputMu.Unlock()
|
|
}
|
|
trySend(msgCh, msg)
|
|
},
|
|
onPromptDone: func(result hermesPromptResult) {
|
|
if !streamingCurrentTurn.Load() {
|
|
return
|
|
}
|
|
select {
|
|
case promptDone <- result:
|
|
default:
|
|
}
|
|
},
|
|
}
|
|
|
|
// Start reading stdout in background.
|
|
readerDone := make(chan struct{})
|
|
go func() {
|
|
defer close(readerDone)
|
|
scanner := bufio.NewScanner(stdout)
|
|
scanner.Buffer(make([]byte, 0, 1024*1024), 10*1024*1024)
|
|
for scanner.Scan() {
|
|
line := strings.TrimSpace(scanner.Text())
|
|
if line == "" {
|
|
continue
|
|
}
|
|
c.handleLine(line)
|
|
}
|
|
c.closeAllPending(fmt.Errorf("hermes process exited"))
|
|
}()
|
|
|
|
// Drive the ACP session lifecycle in a goroutine.
|
|
go func() {
|
|
defer close(msgCh)
|
|
defer close(resCh)
|
|
defer func() {
|
|
stdin.Close()
|
|
// Cancellation must be reachable before Wait. A pathological child
|
|
// can close stdout/stderr (so the pipe drain succeeds) but keep the
|
|
// process alive; waiting first would then block until the overall
|
|
// task timeout and make a later deferred cancel ineffective.
|
|
cancel()
|
|
_ = cmd.Wait()
|
|
}()
|
|
|
|
startTime := time.Now()
|
|
finalStatus := "completed"
|
|
var finalError string
|
|
var sessionID string
|
|
// Set when the ACP runtime refuses the session we asked to
|
|
// resume. Only that is curable by starting a fresh session, so
|
|
// handshake/network failures below must leave it false.
|
|
var resumeRejected bool
|
|
effectiveModel := strings.TrimSpace(opts.Model)
|
|
// The model id the runtime reports as current right after
|
|
// session/new or session/resume. Used to skip a redundant
|
|
// session/set_model when we would otherwise re-select the model the
|
|
// session is already on (see the set_model gate below).
|
|
var sessionCurrentModel string
|
|
|
|
// 1. Initialize handshake.
|
|
initResult, err := c.request(runCtx, "initialize", map[string]any{
|
|
"protocolVersion": 1,
|
|
"clientInfo": map[string]any{
|
|
"name": "multica-agent-sdk",
|
|
"version": "0.2.0",
|
|
},
|
|
"clientCapabilities": map[string]any{},
|
|
})
|
|
if err != nil {
|
|
finalStatus = "failed"
|
|
finalError = fmt.Sprintf("hermes initialize failed: %v", err)
|
|
resCh <- Result{Status: finalStatus, Error: finalError, DurationMs: time.Since(startTime).Milliseconds()}
|
|
return
|
|
}
|
|
|
|
// Drop MCP entries whose remote transport the runtime didn't
|
|
// advertise. ACP requires the client to honour
|
|
// agentCapabilities.mcpCapabilities; sending an http/sse entry to
|
|
// a runtime that says it only supports stdio reliably rejects the
|
|
// whole session/new request.
|
|
mcpServers = filterACPMcpServersByCapability(mcpServers, extractACPMcpCapabilities(initResult), "hermes", b.cfg.Logger)
|
|
|
|
// 2. Create or resume a session.
|
|
cwd := opts.Cwd
|
|
if cwd == "" {
|
|
cwd = "."
|
|
}
|
|
|
|
if opts.ResumeSessionID != "" {
|
|
// Per ACP Session Setup, session/resume accepts mcpServers and
|
|
// the runtime re-connects them as part of the resume. Without
|
|
// this, a resumed Hermes task lost access to MCP tools that a
|
|
// fresh task on the same agent would have.
|
|
result, err := c.request(runCtx, "session/resume", map[string]any{
|
|
"cwd": cwd,
|
|
"sessionId": opts.ResumeSessionID,
|
|
"mcpServers": mcpServers,
|
|
})
|
|
if err != nil {
|
|
finalStatus = "failed"
|
|
finalError = fmt.Sprintf("hermes session/resume failed: %v", err)
|
|
resCh <- Result{Status: finalStatus, Error: finalError, DurationMs: time.Since(startTime).Milliseconds()}
|
|
return
|
|
}
|
|
var changed bool
|
|
sessionID, changed = resolveResumedSessionID(opts.ResumeSessionID, result)
|
|
if changed {
|
|
b.cfg.Logger.Warn("agent returned a different session id on resume — original was likely lost; continuing with the new id",
|
|
"backend", "hermes",
|
|
"requested", opts.ResumeSessionID,
|
|
"actual", sessionID,
|
|
)
|
|
}
|
|
sessionCurrentModel = extractACPCurrentModelID(result)
|
|
if effectiveModel == "" {
|
|
effectiveModel = sessionCurrentModel
|
|
}
|
|
} else {
|
|
result, err := c.request(runCtx, "session/new", buildHermesSessionParams(cwd, opts.Model, mcpServers))
|
|
if err != nil {
|
|
finalStatus = "failed"
|
|
finalError = fmt.Sprintf("hermes session/new failed: %v", err)
|
|
resCh <- Result{Status: finalStatus, Error: finalError, DurationMs: time.Since(startTime).Milliseconds()}
|
|
return
|
|
}
|
|
sessionID = extractACPSessionID(result)
|
|
if sessionID == "" {
|
|
finalStatus = "failed"
|
|
finalError = "hermes session/new returned no session ID"
|
|
resCh <- Result{Status: finalStatus, Error: finalError, DurationMs: time.Since(startTime).Milliseconds()}
|
|
return
|
|
}
|
|
sessionCurrentModel = extractACPCurrentModelID(result)
|
|
if effectiveModel == "" {
|
|
effectiveModel = sessionCurrentModel
|
|
}
|
|
}
|
|
|
|
c.sessionID = sessionID
|
|
b.cfg.Logger.Info("hermes session created", "session_id", sessionID)
|
|
|
|
// 3. If the caller picked a model (via agent.model from the
|
|
// UI dropdown), ask hermes to switch the session to it
|
|
// before we send any prompt. Hermes' _build_model_state
|
|
// exposes modelId as `provider:model` — we pass that
|
|
// through verbatim. This MUST fail the task on error:
|
|
// if we silently fell back to hermes' default model the
|
|
// user would think their pick was honoured while the
|
|
// task actually ran on something else.
|
|
//
|
|
// Skip the call when the session already reports this exact model as
|
|
// current. Hermes' set_model re-runs provider auto-detection on the
|
|
// model id, and for a `provider:model` id whose parsed provider equals
|
|
// the session's current provider it can mis-route to a different
|
|
// provider (e.g. custom:deepseek-v4-pro → OpenRouter) and fail with an
|
|
// auth error. Re-selecting the model the session is already on is pure
|
|
// downside. An empty sessionCurrentModel (older runtime or unparsable
|
|
// state) falls through and still sends set_model, preserving prior
|
|
// behaviour. See MUL-5029 / NousResearch/hermes-agent#59089.
|
|
if opts.Model != "" && effectiveModel == sessionCurrentModel {
|
|
b.cfg.Logger.Info("hermes session already on requested model; skipping redundant set_model",
|
|
"model", opts.Model,
|
|
"session_id", sessionID,
|
|
)
|
|
} else if opts.Model != "" {
|
|
if _, err := c.request(runCtx, "session/set_model", map[string]any{
|
|
"sessionId": sessionID,
|
|
"modelId": opts.Model,
|
|
}); err != nil {
|
|
b.cfg.Logger.Warn("hermes set_session_model failed", "error", err, "requested_model", opts.Model)
|
|
finalStatus = "failed"
|
|
finalError = fmt.Sprintf("hermes could not switch to model %q: %v", opts.Model, err)
|
|
if opts.ResumeSessionID != "" && isACPSessionNotFound(err) {
|
|
// On a resumed session with a model override, the dead
|
|
// session surfaces here instead of at session/prompt.
|
|
// Same fix as the prompt path below: clear the id so
|
|
// the daemon's resume-failure fallback retries fresh.
|
|
b.cfg.Logger.Warn("resumed session not found at set_model time; clearing session id so the daemon retries fresh",
|
|
"backend", "hermes",
|
|
"session_id", sessionID,
|
|
)
|
|
sessionID = ""
|
|
resumeRejected = true
|
|
}
|
|
resCh <- Result{
|
|
Status: finalStatus,
|
|
Error: finalError,
|
|
DurationMs: time.Since(startTime).Milliseconds(),
|
|
SessionID: sessionID,
|
|
ResumeRejected: resumeRejected,
|
|
}
|
|
return
|
|
}
|
|
b.cfg.Logger.Info("hermes session model set", "model", opts.Model)
|
|
}
|
|
|
|
// 4. Send the prompt and wait for PromptResponse.
|
|
//
|
|
// Do NOT prepend opts.SystemPrompt here. Hermes ACP loads project/context
|
|
// files from cwd (AGENTS.md, .agent_context, etc.) itself; duplicating the
|
|
// full runtime brief in the user prompt makes the request much larger and
|
|
// has triggered upstream safety filters on otherwise ordinary tasks.
|
|
// Flip the gate
|
|
// just before the request so any history replay flushed during
|
|
// initialize / session setup stays dropped, but every notification
|
|
// belonging to this turn is processed.
|
|
streamingCurrentTurn.Store(true)
|
|
_, err = c.request(runCtx, "session/prompt", map[string]any{
|
|
"sessionId": sessionID,
|
|
"prompt": []map[string]any{
|
|
{"type": "text", "text": prompt},
|
|
},
|
|
})
|
|
if err != nil {
|
|
// If the request itself failed (not just context cancelled),
|
|
// check if the context was cancelled/timed out.
|
|
if runCtx.Err() == context.DeadlineExceeded {
|
|
finalStatus = "timeout"
|
|
finalError = fmt.Sprintf("hermes timed out after %s", timeout)
|
|
} else if runCtx.Err() == context.Canceled {
|
|
finalStatus = "aborted"
|
|
finalError = "execution cancelled"
|
|
} else {
|
|
finalStatus = "failed"
|
|
finalError = fmt.Sprintf("hermes session/prompt failed: %v", err)
|
|
if opts.ResumeSessionID != "" && isACPSessionNotFound(err) {
|
|
// The agent no longer knows the session we resumed.
|
|
// Hermes echoes the requested id back from
|
|
// session/resume even when the session is gone, so
|
|
// resolveResumedSessionID can't catch this — it only
|
|
// surfaces here, at prompt time. Return an empty
|
|
// SessionID so the daemon's resume-failure fallback
|
|
// retries with a fresh session and stores the
|
|
// replacement id; keeping the stale id makes every
|
|
// future dispatch on this (agent, issue) fail the
|
|
// same way.
|
|
b.cfg.Logger.Warn("resumed session not found at prompt time; clearing session id so the daemon retries fresh",
|
|
"backend", "hermes",
|
|
"session_id", sessionID,
|
|
)
|
|
sessionID = ""
|
|
resumeRejected = true
|
|
}
|
|
}
|
|
} else {
|
|
// The prompt completed. Check if we got a promptDone result
|
|
// from the response parsing.
|
|
select {
|
|
case pr := <-promptDone:
|
|
if pr.stopReason == "cancelled" {
|
|
finalStatus = "aborted"
|
|
finalError = "hermes cancelled the prompt"
|
|
}
|
|
// Merge usage from the PromptResponse.
|
|
c.usageMu.Lock()
|
|
c.usage.InputTokens += pr.usage.InputTokens
|
|
c.usage.OutputTokens += pr.usage.OutputTokens
|
|
c.usage.CacheReadTokens += pr.usage.CacheReadTokens
|
|
c.usageMu.Unlock()
|
|
default:
|
|
}
|
|
waitForHermesNotificationQuiescence(runCtx, activity, readerDone)
|
|
}
|
|
|
|
duration := time.Since(startTime)
|
|
b.cfg.Logger.Info("hermes finished", "pid", cmd.Process.Pid, "status", finalStatus, "duration", duration.Round(time.Millisecond).String())
|
|
|
|
// Close stdin first so Hermes can observe EOF and exit cleanly. Keep the
|
|
// process alive while stdout/stderr drain; cancelling at the prompt
|
|
// response boundary can truncate final notifications that arrive just
|
|
// after the response.
|
|
stdin.Close()
|
|
|
|
// Wait for the stdout reader and stderr copier so all output is
|
|
// accumulated and the provider-error sniffer
|
|
// has every byte the child wrote before we consult it for failure
|
|
// promotion. Skipping this leaves a small race where stopReason=
|
|
// end_turn arrives over stdout while the stderr 429 / usage-limit
|
|
// lines are still in transit, causing the promoted error message
|
|
// to fall through to the synthetic agent-text fallback. If Hermes does
|
|
// not honor stdin EOF within the bound, cancel it and join both readers
|
|
// before accessing their buffers.
|
|
if !waitForHermesPipeDrain(readerDone, stderrDone, hermesReaderDrainGrace) {
|
|
b.cfg.Logger.Warn("hermes did not close output pipes after stdin EOF; forcing shutdown",
|
|
"pid", cmd.Process.Pid,
|
|
"grace", hermesReaderDrainGrace.String(),
|
|
)
|
|
cancel()
|
|
<-readerDone
|
|
<-stderrDone
|
|
}
|
|
streamingCurrentTurn.Store(false)
|
|
|
|
outputMu.Lock()
|
|
finalOutput := output.String()
|
|
outputMu.Unlock()
|
|
|
|
// Hermes reports stopReason=end_turn even when the upstream
|
|
// LLM call ultimately fails (HTTP 429 rate-limit, expired
|
|
// token, ...). promoteACPResultOnProviderError flips the
|
|
// status to "failed" when either the stderr sniffer saw a
|
|
// *terminal* failure marker (not just a transient per-attempt
|
|
// warning), the agent text stream contains the synthetic
|
|
// "API call failed after N retries..." turn the adapter
|
|
// injects on give-up, or there's no output to fall back on.
|
|
finalStatus, finalError = promoteACPResultOnProviderError(finalStatus, finalError, finalOutput, providerErr)
|
|
|
|
// Build usage map.
|
|
c.usageMu.Lock()
|
|
u := c.usage
|
|
c.usageMu.Unlock()
|
|
|
|
var usageMap map[string]TokenUsage
|
|
if u.InputTokens > 0 || u.OutputTokens > 0 || u.CacheReadTokens > 0 || u.CacheWriteTokens > 0 {
|
|
model := effectiveModel
|
|
if model == "" {
|
|
model = "unknown"
|
|
}
|
|
usageMap = map[string]TokenUsage{model: u}
|
|
}
|
|
|
|
resCh <- Result{
|
|
Status: finalStatus,
|
|
Output: finalOutput,
|
|
Error: finalError,
|
|
DurationMs: duration.Milliseconds(),
|
|
SessionID: sessionID,
|
|
ResumeRejected: resumeRejected,
|
|
Usage: usageMap,
|
|
}
|
|
}()
|
|
|
|
return &Session{Messages: msgCh, Result: resCh}, nil
|
|
}
|
|
|
|
// waitForHermesNotificationQuiescence gives the stdout reader a bounded chance
|
|
// to consume session updates emitted just after session/prompt returns. Hermes
|
|
// may deliver the final agent_message_chunk after the response; closing stdin
|
|
// or cancelling immediately at that boundary loses the user-visible answer.
|
|
func waitForHermesNotificationQuiescence(ctx context.Context, activity <-chan struct{}, readerDone <-chan struct{}) {
|
|
quiet := time.NewTimer(hermesNotificationQuietTime)
|
|
defer quiet.Stop()
|
|
hard := time.NewTimer(hermesReaderDrainGrace)
|
|
defer hard.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-activity:
|
|
if !quiet.Stop() {
|
|
select {
|
|
case <-quiet.C:
|
|
default:
|
|
}
|
|
}
|
|
quiet.Reset(hermesNotificationQuietTime)
|
|
case <-quiet.C:
|
|
return
|
|
case <-readerDone:
|
|
return
|
|
case <-hard.C:
|
|
return
|
|
case <-ctx.Done():
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
func waitForHermesPipeDrain(readerDone, stderrDone <-chan struct{}, timeout time.Duration) bool {
|
|
timer := time.NewTimer(timeout)
|
|
defer timer.Stop()
|
|
|
|
for readerDone != nil || stderrDone != nil {
|
|
select {
|
|
case <-readerDone:
|
|
readerDone = nil
|
|
case <-stderrDone:
|
|
stderrDone = nil
|
|
case <-timer.C:
|
|
return false
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
|
|
// ── hermesClient: ACP JSON-RPC 2.0 transport ──
|
|
|
|
type hermesPromptResult struct {
|
|
stopReason string
|
|
usage TokenUsage
|
|
// modelID is the model the agent actually billed this turn against, as
|
|
// reported on `result._meta.modelId`. Empty for agents that don't report
|
|
// it. Backends use it to attribute usage when the session handshake
|
|
// didn't surface a model id (see grok.go).
|
|
modelID string
|
|
}
|
|
|
|
type hermesClient struct {
|
|
cfg Config
|
|
stdin interface{ Write([]byte) (int, error) }
|
|
writeMu sync.Mutex // serialises stdin.Write calls across goroutines
|
|
mu sync.Mutex
|
|
nextID int
|
|
pending map[int]*pendingRPC
|
|
sessionID string
|
|
onMessage func(Message)
|
|
onPromptDone func(hermesPromptResult)
|
|
// onActivity observes accepted ACP session updates. Hermes and Grok use it
|
|
// to retain a short post-response drain window; other ACP backends leave it
|
|
// nil and keep their existing lifecycle behavior.
|
|
onActivity func()
|
|
// acceptNotification can drop ACP session updates before dispatching to
|
|
// handlers that mutate client state such as usage or pending tool calls.
|
|
acceptNotification func(updateType string) bool
|
|
|
|
// pendingTools buffers the args for tool calls whose input streams in
|
|
// across multiple ACP tool_call_update messages (kimi does this —
|
|
// tokens from the LLM arrive one at a time, and each update carries
|
|
// the cumulative args JSON so far). We defer emitting MessageToolUse
|
|
// until we either see status=completed/failed or have a full arg set,
|
|
// so the UI never sees a half-written command like `{"comma`.
|
|
toolMu sync.Mutex
|
|
pendingTools map[string]*pendingToolCall
|
|
|
|
usageMu sync.Mutex
|
|
usage TokenUsage
|
|
}
|
|
|
|
// pendingToolCall buffers state for a tool call while its arguments
|
|
// are streaming in. One entry per ACP toolCallId.
|
|
type pendingToolCall struct {
|
|
toolName string // already mapped via hermesToolNameFromTitle
|
|
input map[string]any // from rawInput when the agent sends it up front (hermes)
|
|
argsText string // accumulated `content[].text` args (kimi, cumulative)
|
|
emitted bool // whether we've already sent MessageToolUse
|
|
}
|
|
|
|
// writeLine serialises concurrent JSON-RPC writes so request() (main
|
|
// goroutine) and handleAgentRequest() (reader goroutine) don't
|
|
// interleave frames. The pipe itself is atomic for small writes, but
|
|
// we also want deterministic ordering under contention.
|
|
func (c *hermesClient) writeLine(data []byte) error {
|
|
c.writeMu.Lock()
|
|
defer c.writeMu.Unlock()
|
|
_, err := c.stdin.Write(data)
|
|
return err
|
|
}
|
|
|
|
func (c *hermesClient) request(ctx context.Context, method string, params any) (json.RawMessage, error) {
|
|
c.mu.Lock()
|
|
id := c.nextID
|
|
c.nextID++
|
|
pr := &pendingRPC{ch: make(chan rpcResult, 1), method: method}
|
|
c.pending[id] = pr
|
|
c.mu.Unlock()
|
|
|
|
msg := map[string]any{
|
|
"jsonrpc": "2.0",
|
|
"id": id,
|
|
"method": method,
|
|
"params": params,
|
|
}
|
|
data, err := json.Marshal(msg)
|
|
if err != nil {
|
|
c.mu.Lock()
|
|
delete(c.pending, id)
|
|
c.mu.Unlock()
|
|
return nil, err
|
|
}
|
|
data = append(data, '\n')
|
|
if err := c.writeLine(data); err != nil {
|
|
c.mu.Lock()
|
|
delete(c.pending, id)
|
|
c.mu.Unlock()
|
|
return nil, fmt.Errorf("write %s: %w", method, err)
|
|
}
|
|
|
|
select {
|
|
case res := <-pr.ch:
|
|
return res.result, res.err
|
|
case <-ctx.Done():
|
|
c.mu.Lock()
|
|
delete(c.pending, id)
|
|
c.mu.Unlock()
|
|
return nil, ctx.Err()
|
|
}
|
|
}
|
|
|
|
func (c *hermesClient) closeAllPending(err error) {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
for id, pr := range c.pending {
|
|
pr.ch <- rpcResult{err: err}
|
|
delete(c.pending, id)
|
|
}
|
|
}
|
|
|
|
func (c *hermesClient) handleLine(line string) {
|
|
var raw map[string]json.RawMessage
|
|
if err := json.Unmarshal([]byte(line), &raw); err != nil {
|
|
return
|
|
}
|
|
|
|
// Agent → client request: has id + method (no result / error yet).
|
|
// Kimi and Hermes both use session/request_permission; if we don't
|
|
// answer, the agent blocks for its internal timeout and the task
|
|
// hangs. HERMES_YOLO_MODE=1 only suppresses Hermes' dangerous-shell-
|
|
// command prompts (tools/approval.py); its ACP edit-approval guard
|
|
// (acp_adapter/edit_approval.py) still asks before every file write,
|
|
// so we must handle these requests for Hermes too.
|
|
if _, hasID := raw["id"]; hasID {
|
|
if _, hasResult := raw["result"]; hasResult {
|
|
c.handleResponse(raw)
|
|
return
|
|
}
|
|
if _, hasError := raw["error"]; hasError {
|
|
c.handleResponse(raw)
|
|
return
|
|
}
|
|
if _, hasMethod := raw["method"]; hasMethod {
|
|
c.handleAgentRequest(raw)
|
|
return
|
|
}
|
|
}
|
|
|
|
// Notification (no id, has method) — session updates from Hermes.
|
|
if _, hasMethod := raw["method"]; hasMethod {
|
|
c.handleNotification(raw)
|
|
}
|
|
}
|
|
|
|
// handleAgentRequest replies to JSON-RPC requests the agent sends
|
|
// us (agent → client direction). The only one we care about today is
|
|
// `session/request_permission`: the daemon is headless and cannot
|
|
// actually prompt a user, so we answer it ourselves — granting when a
|
|
// safe option is offered, otherwise declining just this action or
|
|
// failing closed (see below).
|
|
//
|
|
// The reply MUST select one of the optionIds the agent actually
|
|
// offered — the ACP permission contract is "pick from these options",
|
|
// and an id the agent never offered is treated as a denial. Hermes'
|
|
// edit-approval path offers only ["allow_once","deny"] and rejects
|
|
// anything but exactly "allow_once" (acp_adapter/edit_approval.py), so
|
|
// the previous hardcoded "approve_for_session" silently blocked every
|
|
// file write on the Hermes ACP runtime (GitHub multica#5300).
|
|
// selectACPPermissionOption picks an option the agent offered — a safe
|
|
// grant when one exists, otherwise an offered single-use reject to deny
|
|
// just this action — and we fail closed with a protocol error when the
|
|
// request offers nothing safely selectable, never a permanent grant or a
|
|
// whole-turn "cancelled".
|
|
func (c *hermesClient) handleAgentRequest(raw map[string]json.RawMessage) {
|
|
var method string
|
|
_ = json.Unmarshal(raw["method"], &method)
|
|
|
|
rawID, ok := raw["id"]
|
|
if !ok {
|
|
return
|
|
}
|
|
|
|
var resp map[string]any
|
|
switch method {
|
|
case "session/request_permission":
|
|
optionID, grant, ok := selectACPPermissionOption(raw["params"])
|
|
if ok {
|
|
// Select an offered option — either a safe grant (approve) or,
|
|
// when no safe grant exists, an offered reject_once (deny THIS
|
|
// action). Both are ACP "selected" outcomes; we deliberately do
|
|
// NOT reply "cancelled" here, which means the whole prompt turn
|
|
// was cancelled — other ACP backends sharing this client (kimi,
|
|
// kiro, ...) would abort the entire task, not just this action.
|
|
resp = map[string]any{
|
|
"jsonrpc": "2.0",
|
|
"id": json.RawMessage(rawID),
|
|
"result": map[string]any{
|
|
"outcome": map[string]any{
|
|
"outcome": "selected",
|
|
"optionId": optionID,
|
|
},
|
|
},
|
|
}
|
|
if grant {
|
|
c.cfg.Logger.Debug("auto-approved agent permission request", "method", method, "optionId", optionID)
|
|
} else {
|
|
c.cfg.Logger.Warn("no safe grant offered; selecting offered reject option", "method", method, "optionId", optionID)
|
|
}
|
|
} else {
|
|
// The request offered nothing we can safely select: no safe grant
|
|
// and no single-use reject_once (empty, malformed, permanent-only,
|
|
// or reject_always-only). Return a protocol error rather than
|
|
// fabricate an un-offered id or a whole-turn "cancelled".
|
|
resp = map[string]any{
|
|
"jsonrpc": "2.0",
|
|
"id": json.RawMessage(rawID),
|
|
"error": map[string]any{
|
|
"code": -32603,
|
|
"message": "no auto-selectable permission option offered",
|
|
},
|
|
}
|
|
c.cfg.Logger.Warn("no safely selectable permission option offered; returning error", "method", method)
|
|
}
|
|
default:
|
|
// Unknown agent→client method — reply with standard "method
|
|
// not found" so the agent doesn't block waiting for us. Better
|
|
// than silence: the agent can decide how to proceed.
|
|
resp = map[string]any{
|
|
"jsonrpc": "2.0",
|
|
"id": json.RawMessage(rawID),
|
|
"error": map[string]any{
|
|
"code": -32601,
|
|
"message": "method not found: " + method,
|
|
},
|
|
}
|
|
c.cfg.Logger.Debug("unhandled agent→client request", "method", method)
|
|
}
|
|
|
|
data, err := json.Marshal(resp)
|
|
if err != nil {
|
|
c.cfg.Logger.Warn("marshal agent-request response", "method", method, "error", err)
|
|
return
|
|
}
|
|
data = append(data, '\n')
|
|
if err := c.writeLine(data); err != nil {
|
|
c.cfg.Logger.Warn("write agent-request response", "method", method, "error", err)
|
|
}
|
|
}
|
|
|
|
// acpPermissionOption is one entry in a session/request_permission
|
|
// request's `options` array. `kind` is the ACP-level classification
|
|
// (allow_once / allow_always / reject_once / reject_always); `optionId`
|
|
// is the agent-defined opaque string we echo back to select that option.
|
|
type acpPermissionOption struct {
|
|
OptionID string `json:"optionId"`
|
|
Kind string `json:"kind"`
|
|
}
|
|
|
|
// ACP v1 PermissionOptionKind values. Only the two allow kinds grant;
|
|
// any other or unknown kind is treated as non-granting so a future or
|
|
// abnormal kind can never be auto-approved. reject_once denies a single
|
|
// action (the only reject we auto-select — reject_always would persist a
|
|
// denial the way allow_always persists a grant).
|
|
// https://agentclientprotocol.com/protocol/v1/schema#permissionoptionkind
|
|
const (
|
|
acpKindAllowOnce = "allow_once"
|
|
acpKindAllowAlways = "allow_always"
|
|
acpKindRejectOnce = "reject_once"
|
|
)
|
|
|
|
// acpSessionScopedOptionIDs are optionIds known to grant for the current
|
|
// session only, without persisting a decision. Both Hermes' "allow_session"
|
|
// and its permanent "allow_always" option carry ACP kind "allow_always"
|
|
// (ACP has no session-scoped kind), so kind alone cannot tell them apart —
|
|
// we recognise the session-scoped ones by id. "approve_for_session" is the
|
|
// equivalent id other ACP backends use.
|
|
var acpSessionScopedOptionIDs = []string{"allow_session", "approve_for_session"}
|
|
|
|
// selectACPPermissionOption decides how to auto-answer a
|
|
// session/request_permission. It returns the offered optionId to select,
|
|
// whether that selection grants (true) or denies (false) the action, and
|
|
// ok=false when the request offers nothing safely selectable (the caller
|
|
// then returns a protocol error rather than fabricate an outcome).
|
|
//
|
|
// It only ever returns an id the agent actually offered. Per review of
|
|
// GitHub multica#5300 it refuses to auto-select a permanent "allow_always"
|
|
// grant — on Hermes that persists to the runtime owner's on-disk allowlist
|
|
// and would outlive the task (ACP v1 allow_always "remembers the choice").
|
|
// Grant nature is decided purely by the explicit ACP kind, never by the
|
|
// opaque optionId, so unknown kinds fail closed. Order of preference:
|
|
//
|
|
// 1. a known session-scoped grant id;
|
|
// 2. any single-use (kind=allow_once) grant;
|
|
// 3. an offered single-use reject_once — deny just this action rather than
|
|
// reply "cancelled", which other ACP backends read as cancelling the
|
|
// whole prompt turn.
|
|
func selectACPPermissionOption(params json.RawMessage) (optionID string, grant bool, ok bool) {
|
|
var p struct {
|
|
Options []acpPermissionOption `json:"options"`
|
|
}
|
|
if len(params) > 0 {
|
|
if err := json.Unmarshal(params, &p); err != nil {
|
|
return "", false, false
|
|
}
|
|
}
|
|
|
|
// 1. A known session-scoped grant id, if actually offered with a grant kind.
|
|
for _, want := range acpSessionScopedOptionIDs {
|
|
for _, opt := range p.Options {
|
|
if opt.OptionID == want && isACPGrantKind(opt.Kind) {
|
|
return opt.OptionID, true, true
|
|
}
|
|
}
|
|
}
|
|
// 2. Any single-use grant. kind=allow_once is inherently scoped to this
|
|
// one action, so it is safe regardless of the (opaque) optionId — this
|
|
// also covers agents that use non-standard option ids.
|
|
for _, opt := range p.Options {
|
|
if opt.OptionID != "" && strings.EqualFold(strings.TrimSpace(opt.Kind), acpKindAllowOnce) {
|
|
return opt.OptionID, true, true
|
|
}
|
|
}
|
|
// 3. No safe grant: deny THIS action by selecting an offered reject_once.
|
|
for _, opt := range p.Options {
|
|
if opt.OptionID != "" && strings.EqualFold(strings.TrimSpace(opt.Kind), acpKindRejectOnce) {
|
|
return opt.OptionID, false, true
|
|
}
|
|
}
|
|
// 4. Nothing safely selectable (empty, malformed, permanent-only, or
|
|
// reject_always-only). Signal the caller to return a protocol error.
|
|
return "", false, false
|
|
}
|
|
|
|
// isACPGrantKind reports whether an ACP PermissionOptionKind grants the
|
|
// action. Only the two current allow kinds qualify; every other or unknown
|
|
// value is non-granting, so grant detection fails closed.
|
|
func isACPGrantKind(kind string) bool {
|
|
switch strings.ToLower(strings.TrimSpace(kind)) {
|
|
case acpKindAllowOnce, acpKindAllowAlways:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
// acpRPCError is a JSON-RPC error frame returned by the agent process.
|
|
// It renders exactly like the flat string handleResponse used to build
|
|
// with fmt.Errorf, so logs and surfaced task errors are unchanged, but
|
|
// keeps the code and message structured so callers can branch on the
|
|
// error class (see isACPSessionNotFound) instead of parsing text.
|
|
type acpRPCError struct {
|
|
Method string
|
|
Code int
|
|
Message string
|
|
Data string
|
|
}
|
|
|
|
func (e *acpRPCError) Error() string {
|
|
if e.Data != "" {
|
|
return fmt.Sprintf("%s: %s (code=%d, data=%s)", e.Method, e.Message, e.Code, e.Data)
|
|
}
|
|
return fmt.Sprintf("%s: %s (code=%d)", e.Method, e.Message, e.Code)
|
|
}
|
|
|
|
// isACPSessionNotFound reports whether err is the agent rejecting a
|
|
// session id it no longer knows. Runtimes signal this with codes and
|
|
// wording that vary — Hermes says "Session not found" under -32603
|
|
// (Internal error), Kiro puts "No session found with id ..." in
|
|
// `data` under -32603, and kimi-cli raises invalid_params (-32602)
|
|
// with {"session_id": "Session not found"} in `data` for every
|
|
// unknown-session path (src/kimi_cli/acp/server.py) — so neither the
|
|
// code nor the text alone is discriminating and both are matched.
|
|
func isACPSessionNotFound(err error) bool {
|
|
var rpcErr *acpRPCError
|
|
if !errors.As(err, &rpcErr) {
|
|
return false
|
|
}
|
|
if rpcErr.Code != -32603 && rpcErr.Code != -32602 {
|
|
return false
|
|
}
|
|
text := strings.ToLower(rpcErr.Message + " " + rpcErr.Data)
|
|
return strings.Contains(text, "session not found") ||
|
|
strings.Contains(text, "no session found")
|
|
}
|
|
|
|
func (c *hermesClient) handleResponse(raw map[string]json.RawMessage) {
|
|
var id int
|
|
if err := json.Unmarshal(raw["id"], &id); err != nil {
|
|
// Try float (JSON numbers are floats by default).
|
|
var fid float64
|
|
if err := json.Unmarshal(raw["id"], &fid); err != nil {
|
|
return
|
|
}
|
|
id = int(fid)
|
|
}
|
|
|
|
c.mu.Lock()
|
|
pr, ok := c.pending[id]
|
|
if ok {
|
|
delete(c.pending, id)
|
|
}
|
|
c.mu.Unlock()
|
|
|
|
if !ok {
|
|
return
|
|
}
|
|
|
|
if errData, hasErr := raw["error"]; hasErr {
|
|
var rpcErr struct {
|
|
Code int `json:"code"`
|
|
Message string `json:"message"`
|
|
Data json.RawMessage `json:"data"`
|
|
}
|
|
_ = json.Unmarshal(errData, &rpcErr)
|
|
// JSON-RPC `data` carries the provider-specific reason (e.g. Kiro
|
|
// returns "No session found with id" for code=-32603). Surface it
|
|
// in the wrapped error so daemon logs / UI can show *why* the
|
|
// agent failed instead of a bare "Internal error". `data` may be
|
|
// any JSON value: render strings unquoted, everything else as raw
|
|
// JSON.
|
|
detail := ""
|
|
if len(rpcErr.Data) > 0 && string(rpcErr.Data) != "null" {
|
|
var s string
|
|
if err := json.Unmarshal(rpcErr.Data, &s); err == nil {
|
|
detail = s
|
|
} else {
|
|
detail = string(rpcErr.Data)
|
|
}
|
|
}
|
|
pr.ch <- rpcResult{err: &acpRPCError{Method: pr.method, Code: rpcErr.Code, Message: rpcErr.Message, Data: detail}}
|
|
} else {
|
|
// If this is a prompt response, extract usage and stop reason.
|
|
if pr.method == "session/prompt" {
|
|
c.extractPromptResult(raw["result"])
|
|
}
|
|
pr.ch <- rpcResult{result: raw["result"]}
|
|
}
|
|
}
|
|
|
|
func (c *hermesClient) extractPromptResult(data json.RawMessage) {
|
|
var resp struct {
|
|
StopReason string `json:"stopReason"`
|
|
Usage json.RawMessage `json:"usage"`
|
|
Meta json.RawMessage `json:"_meta"`
|
|
}
|
|
if err := json.Unmarshal(data, &resp); err != nil {
|
|
return
|
|
}
|
|
|
|
pr := hermesPromptResult{
|
|
stopReason: resp.StopReason,
|
|
modelID: parseACPModelIDFromMeta(resp.Meta),
|
|
}
|
|
if len(resp.Usage) > 0 && string(resp.Usage) != "null" {
|
|
pr.usage = parseACPTokenUsage(resp.Usage)
|
|
}
|
|
// Prefer the standard top-level ACP `usage` field when present. Some
|
|
// agents (notably xAI Grok Build) put per-turn metering only under
|
|
// result._meta — either as `_meta.usage` or as flat token counters on
|
|
// `_meta` itself. Without this fallback, tasks complete with an empty
|
|
// usage map and Multica's Usage/cost dashboards stay at zero.
|
|
if !acpTokenUsagePresent(pr.usage) {
|
|
if metaUsage := parseACPTokenUsageFromMeta(resp.Meta); acpTokenUsagePresent(metaUsage) {
|
|
pr.usage = metaUsage
|
|
}
|
|
}
|
|
|
|
if c.onPromptDone != nil {
|
|
c.onPromptDone(pr)
|
|
}
|
|
}
|
|
|
|
// acpTokenUsagePresent reports whether any token counter is non-zero.
|
|
func acpTokenUsagePresent(u TokenUsage) bool {
|
|
return u.InputTokens > 0 || u.OutputTokens > 0 || u.CacheReadTokens > 0 || u.CacheWriteTokens > 0
|
|
}
|
|
|
|
// parseACPTokenUsageFromMeta extracts token usage from an ACP result `_meta`
|
|
// object. Grok Build returns shapes like:
|
|
//
|
|
// {"inputTokens":…,"outputTokens":…,"cachedReadTokens":…,"usage":{…}}
|
|
//
|
|
// Prefer the nested `usage` object when it carries counters; otherwise parse
|
|
// the flat `_meta` fields with the same alias rules as top-level usage.
|
|
func parseACPTokenUsageFromMeta(meta json.RawMessage) TokenUsage {
|
|
if len(meta) == 0 || string(meta) == "null" {
|
|
return TokenUsage{}
|
|
}
|
|
var envelope struct {
|
|
Usage json.RawMessage `json:"usage"`
|
|
}
|
|
if err := json.Unmarshal(meta, &envelope); err == nil {
|
|
if len(envelope.Usage) > 0 && string(envelope.Usage) != "null" {
|
|
if u := parseACPTokenUsage(envelope.Usage); acpTokenUsagePresent(u) {
|
|
return u
|
|
}
|
|
}
|
|
}
|
|
return parseACPTokenUsage(meta)
|
|
}
|
|
|
|
// parseACPModelIDFromMeta pulls the model id off an ACP result `_meta`
|
|
// object. Grok Build stamps every turn with `_meta.modelId`, which is the
|
|
// only authoritative statement of what the turn was billed against — the
|
|
// session handshake reports a model id on `session/new` but NOT on
|
|
// `session/load`, so a resumed session has no other source.
|
|
func parseACPModelIDFromMeta(meta json.RawMessage) string {
|
|
if len(meta) == 0 || string(meta) == "null" {
|
|
return ""
|
|
}
|
|
var r struct {
|
|
ModelID string `json:"modelId"`
|
|
ModelIDSnake string `json:"model_id"`
|
|
}
|
|
if err := json.Unmarshal(meta, &r); err != nil {
|
|
return ""
|
|
}
|
|
if id := strings.TrimSpace(r.ModelID); id != "" {
|
|
return id
|
|
}
|
|
return strings.TrimSpace(r.ModelIDSnake)
|
|
}
|
|
|
|
func (c *hermesClient) handleNotification(raw map[string]json.RawMessage) {
|
|
var method string
|
|
_ = json.Unmarshal(raw["method"], &method)
|
|
|
|
if method != "session/update" && method != "session/notification" {
|
|
return
|
|
}
|
|
|
|
var params struct {
|
|
SessionID string `json:"sessionId"`
|
|
Update json.RawMessage `json:"update"`
|
|
}
|
|
if p, ok := raw["params"]; ok {
|
|
_ = json.Unmarshal(p, ¶ms)
|
|
}
|
|
if len(params.Update) == 0 {
|
|
return
|
|
}
|
|
|
|
updateType, updateData := normalizeACPUpdate(params.Update)
|
|
if c.acceptNotification != nil && !c.acceptNotification(updateType) {
|
|
return
|
|
}
|
|
if c.onActivity != nil {
|
|
c.onActivity()
|
|
}
|
|
|
|
switch updateType {
|
|
case "agent_message_chunk":
|
|
c.handleAgentMessage(updateData)
|
|
case "agent_thought_chunk":
|
|
c.handleAgentThought(updateData)
|
|
case "tool_call":
|
|
c.handleToolCallStart(updateData)
|
|
case "tool_call_update":
|
|
c.handleToolCallUpdate(updateData)
|
|
case "usage_update":
|
|
c.handleUsageUpdate(updateData)
|
|
case "turn_end":
|
|
c.extractPromptResult(updateData)
|
|
}
|
|
}
|
|
|
|
func normalizeACPUpdate(data json.RawMessage) (string, json.RawMessage) {
|
|
var updateType struct {
|
|
SessionUpdate string `json:"sessionUpdate"`
|
|
Type string `json:"type"`
|
|
}
|
|
_ = json.Unmarshal(data, &updateType)
|
|
if updateType.SessionUpdate != "" {
|
|
return normalizeACPUpdateType(updateType.SessionUpdate), data
|
|
}
|
|
if updateType.Type != "" {
|
|
return normalizeACPUpdateType(updateType.Type), data
|
|
}
|
|
|
|
// Some ACP implementations serialize enum variants as an externally
|
|
// tagged object: {"agentMessageChunk": {"content": ...}}.
|
|
var wrapper map[string]json.RawMessage
|
|
if err := json.Unmarshal(data, &wrapper); err == nil && len(wrapper) == 1 {
|
|
for k, v := range wrapper {
|
|
return normalizeACPUpdateType(k), v
|
|
}
|
|
}
|
|
|
|
return "", data
|
|
}
|
|
|
|
func normalizeACPUpdateType(t string) string {
|
|
key := strings.ToLower(strings.ReplaceAll(strings.ReplaceAll(strings.TrimSpace(t), "_", ""), "-", ""))
|
|
switch key {
|
|
case "agentmessagechunk":
|
|
return "agent_message_chunk"
|
|
case "agentthoughtchunk":
|
|
return "agent_thought_chunk"
|
|
case "toolcall":
|
|
return "tool_call"
|
|
case "toolcallupdate":
|
|
return "tool_call_update"
|
|
case "usageupdate":
|
|
return "usage_update"
|
|
case "turnend", "endturn":
|
|
return "turn_end"
|
|
default:
|
|
return ""
|
|
}
|
|
}
|
|
|
|
func (c *hermesClient) handleAgentMessage(data json.RawMessage) {
|
|
var msg struct {
|
|
Content struct {
|
|
Type string `json:"type"`
|
|
Text string `json:"text"`
|
|
} `json:"content"`
|
|
}
|
|
if err := json.Unmarshal(data, &msg); err != nil || msg.Content.Text == "" {
|
|
return
|
|
}
|
|
if c.onMessage != nil {
|
|
c.onMessage(Message{Type: MessageText, Content: msg.Content.Text})
|
|
}
|
|
}
|
|
|
|
func (c *hermesClient) handleAgentThought(data json.RawMessage) {
|
|
var msg struct {
|
|
Content struct {
|
|
Type string `json:"type"`
|
|
Text string `json:"text"`
|
|
} `json:"content"`
|
|
}
|
|
if err := json.Unmarshal(data, &msg); err != nil || msg.Content.Text == "" {
|
|
return
|
|
}
|
|
if c.onMessage != nil {
|
|
c.onMessage(Message{Type: MessageThinking, Content: msg.Content.Text})
|
|
}
|
|
}
|
|
|
|
func (c *hermesClient) handleToolCallStart(data json.RawMessage) {
|
|
var msg struct {
|
|
ToolCallID string `json:"toolCallId"`
|
|
Name string `json:"name"`
|
|
Title string `json:"title"`
|
|
Kind string `json:"kind"`
|
|
RawInput map[string]any `json:"rawInput"`
|
|
Input map[string]any `json:"input"`
|
|
Parameters map[string]any `json:"parameters"`
|
|
Content []json.RawMessage `json:"content"`
|
|
}
|
|
if err := json.Unmarshal(data, &msg); err != nil {
|
|
return
|
|
}
|
|
|
|
toolName := hermesToolNameFromTitle(msg.Title, msg.Kind)
|
|
if toolName == "" {
|
|
toolName = msg.Name
|
|
}
|
|
rawInput := msg.RawInput
|
|
if rawInput == nil {
|
|
rawInput = msg.Input
|
|
}
|
|
if rawInput == nil {
|
|
rawInput = msg.Parameters
|
|
}
|
|
|
|
// Hermes pre-populates rawInput on the initial tool_call — emit
|
|
// MessageToolUse immediately so the UI can show the tool invocation
|
|
// live. Record the emission so handleToolCallUpdate doesn't re-emit
|
|
// on completion.
|
|
if rawInput != nil {
|
|
c.trackTool(msg.ToolCallID, &pendingToolCall{
|
|
toolName: toolName,
|
|
input: rawInput,
|
|
emitted: true,
|
|
})
|
|
if c.onMessage != nil {
|
|
c.onMessage(Message{
|
|
Type: MessageToolUse,
|
|
Tool: toolName,
|
|
CallID: msg.ToolCallID,
|
|
Input: rawInput,
|
|
})
|
|
}
|
|
return
|
|
}
|
|
|
|
// Kimi streams args token-by-token across tool_call_update messages;
|
|
// the initial tool_call often carries an empty content block. Buffer
|
|
// the tool and defer MessageToolUse emission to avoid the UI seeing
|
|
// a command with `{""` as its input.
|
|
c.trackTool(msg.ToolCallID, &pendingToolCall{
|
|
toolName: toolName,
|
|
argsText: extractACPToolCallText(msg.Content),
|
|
emitted: false,
|
|
})
|
|
}
|
|
|
|
func (c *hermesClient) handleToolCallUpdate(data json.RawMessage) {
|
|
var msg struct {
|
|
ToolCallID string `json:"toolCallId"`
|
|
Status string `json:"status"`
|
|
Name string `json:"name"`
|
|
Title string `json:"title"`
|
|
Kind string `json:"kind"`
|
|
RawInput map[string]any `json:"rawInput"`
|
|
Input map[string]any `json:"input"`
|
|
Parameters map[string]any `json:"parameters"`
|
|
RawOutput json.RawMessage `json:"rawOutput"`
|
|
Output json.RawMessage `json:"output"`
|
|
Content []json.RawMessage `json:"content"`
|
|
}
|
|
if err := json.Unmarshal(data, &msg); err != nil {
|
|
return
|
|
}
|
|
|
|
rawInput := msg.RawInput
|
|
if rawInput == nil {
|
|
rawInput = msg.Input
|
|
}
|
|
if rawInput == nil {
|
|
rawInput = msg.Parameters
|
|
}
|
|
title := msg.Title
|
|
if title == "" {
|
|
title = msg.Name
|
|
}
|
|
|
|
// Mid-stream: only buffer updates. Kimi emits many of these per
|
|
// tool call, each carrying the cumulative args JSON so far.
|
|
if msg.Status != "completed" && msg.Status != "failed" {
|
|
if pending := c.getPendingTool(msg.ToolCallID); pending != nil && !pending.emitted {
|
|
if text := extractACPToolCallText(msg.Content); text != "" {
|
|
// kimi streams the full cumulative args on every frame;
|
|
// overwrite rather than concatenate.
|
|
pending.argsText = text
|
|
}
|
|
}
|
|
return
|
|
}
|
|
|
|
// Completion: emit any deferred MessageToolUse first, then the result.
|
|
pending := c.takePendingTool(msg.ToolCallID)
|
|
c.emitDeferredToolUse(pending, msg.ToolCallID, title, msg.Kind, rawInput)
|
|
|
|
output := acpRawText(msg.RawOutput)
|
|
if output == "" {
|
|
output = acpRawText(msg.Output)
|
|
}
|
|
if output == "" {
|
|
output = extractACPToolCallText(msg.Content)
|
|
}
|
|
if c.onMessage != nil {
|
|
c.onMessage(Message{
|
|
Type: MessageToolResult,
|
|
CallID: msg.ToolCallID,
|
|
Output: output,
|
|
Status: msg.Status,
|
|
})
|
|
}
|
|
}
|
|
|
|
// trackTool stores pending-tool state for a given callID. Lazy-inits
|
|
// the map so zero-value hermesClient values (common in tests) don't
|
|
// panic on the first tool call.
|
|
func (c *hermesClient) trackTool(callID string, p *pendingToolCall) {
|
|
c.toolMu.Lock()
|
|
defer c.toolMu.Unlock()
|
|
if c.pendingTools == nil {
|
|
c.pendingTools = make(map[string]*pendingToolCall)
|
|
}
|
|
c.pendingTools[callID] = p
|
|
}
|
|
|
|
// getPendingTool returns the pending entry (may be nil) without
|
|
// removing it. Safe to call on a zero-value hermesClient.
|
|
func (c *hermesClient) getPendingTool(callID string) *pendingToolCall {
|
|
c.toolMu.Lock()
|
|
defer c.toolMu.Unlock()
|
|
if c.pendingTools == nil {
|
|
return nil
|
|
}
|
|
return c.pendingTools[callID]
|
|
}
|
|
|
|
// takePendingTool removes and returns the pending entry, or nil if
|
|
// none was tracked (e.g. the tool completed before we saw its start,
|
|
// or we missed the start frame).
|
|
func (c *hermesClient) takePendingTool(callID string) *pendingToolCall {
|
|
c.toolMu.Lock()
|
|
defer c.toolMu.Unlock()
|
|
if c.pendingTools == nil {
|
|
return nil
|
|
}
|
|
p := c.pendingTools[callID]
|
|
delete(c.pendingTools, callID)
|
|
return p
|
|
}
|
|
|
|
// emitDeferredToolUse emits a buffered MessageToolUse right before the
|
|
// matching MessageToolResult. Handles three cases:
|
|
// - hermes tool: already emitted on tool_call → skip
|
|
// - kimi tool with streamed args → parse accumulated JSON as Input
|
|
// - unknown tool (completed arrived without a start frame) →
|
|
// synthesize minimal info from the update's own fields
|
|
func (c *hermesClient) emitDeferredToolUse(
|
|
p *pendingToolCall,
|
|
callID, updateTitle, updateKind string,
|
|
updateRawInput map[string]any,
|
|
) {
|
|
if p != nil && p.emitted {
|
|
return
|
|
}
|
|
|
|
var toolName string
|
|
var input map[string]any
|
|
|
|
switch {
|
|
case p != nil && p.input != nil:
|
|
// Pre-buffered rawInput path — shouldn't happen because we set
|
|
// emitted=true in that case, but handle defensively.
|
|
toolName = p.toolName
|
|
input = p.input
|
|
case p != nil:
|
|
toolName = p.toolName
|
|
input = parseToolArgsJSON(p.argsText)
|
|
default:
|
|
// No record of the start frame — fall back to the update's own
|
|
// title/kind/rawInput so the UI at least sees the tool name.
|
|
toolName = hermesToolNameFromTitle(updateTitle, updateKind)
|
|
input = updateRawInput
|
|
}
|
|
|
|
if c.onMessage == nil {
|
|
return
|
|
}
|
|
c.onMessage(Message{
|
|
Type: MessageToolUse,
|
|
Tool: toolName,
|
|
CallID: callID,
|
|
Input: input,
|
|
})
|
|
}
|
|
|
|
// parseToolArgsJSON turns kimi's accumulated args string into the
|
|
// structured map the UI expects under Message.Input. Kimi sends args
|
|
// as a JSON-encoded object (`{"command":"echo hi"}`), so a full JSON
|
|
// parse recovers the original tool-arg shape. On malformed input
|
|
// (streaming glitch, non-JSON tool) we preserve the raw text under a
|
|
// `text` key so the UI still has something to render.
|
|
func parseToolArgsJSON(argsText string) map[string]any {
|
|
argsText = strings.TrimSpace(argsText)
|
|
if argsText == "" {
|
|
return nil
|
|
}
|
|
var m map[string]any
|
|
if err := json.Unmarshal([]byte(argsText), &m); err == nil {
|
|
return m
|
|
}
|
|
return map[string]any{"text": argsText}
|
|
}
|
|
|
|
// extractACPToolCallText concatenates the rendered text of every ACP
|
|
// block in a tool_call / tool_call_update's `content` array.
|
|
//
|
|
// Handles the two block types kimi emits:
|
|
// - {type:"content", content:{type:"text", text:"..."}} — plain text
|
|
// (shell output, tool args). Text is concatenated verbatim.
|
|
// - {type:"diff", path, oldText, newText} — FileEdit output. Rendered
|
|
// as a minimal unified-diff header so the UI distinguishes writes
|
|
// from reads without needing a diff viewer.
|
|
//
|
|
// acpRawText renders an ACP output field (rawOutput / output) that may arrive
|
|
// as either a JSON string or a structured value. Some model adapters — notably
|
|
// Kiro's GPT-5.6 Sol path — send the completed tool_call_update's rawOutput as
|
|
// an object like {"items":[{"Json":{...}}]} rather than a string. Declaring
|
|
// that field as a Go string made json.Unmarshal fail, which made
|
|
// handleToolCallUpdate return early and silently DROP the entire update —
|
|
// including its status:"completed" — so the completion signal was lost and the
|
|
// task was wrongly marked failed (issue #5509 / MUL-4860). Accept both shapes.
|
|
func acpRawText(raw json.RawMessage) string {
|
|
if len(raw) == 0 {
|
|
return ""
|
|
}
|
|
var s string
|
|
if err := json.Unmarshal(raw, &s); err == nil {
|
|
return s
|
|
}
|
|
// Non-string (object / array / number): keep the raw JSON as text so the
|
|
// output is preserved rather than discarded.
|
|
return string(raw)
|
|
}
|
|
|
|
// Terminal blocks ({type:"terminal", terminalId}) reference a remote
|
|
// terminal the client would normally subscribe to via terminal/output;
|
|
// we don't advertise terminal capability so we never receive those in
|
|
// practice, but if one slips through we skip it (nothing useful to
|
|
// surface from a bare ID).
|
|
func extractACPToolCallText(blocks []json.RawMessage) string {
|
|
var b strings.Builder
|
|
appendPiece := func(piece string) {
|
|
if piece == "" {
|
|
return
|
|
}
|
|
if b.Len() > 0 {
|
|
b.WriteByte('\n')
|
|
}
|
|
b.WriteString(piece)
|
|
}
|
|
for _, raw := range blocks {
|
|
var kind struct {
|
|
Type string `json:"type"`
|
|
}
|
|
if err := json.Unmarshal(raw, &kind); err != nil {
|
|
continue
|
|
}
|
|
switch kind.Type {
|
|
case "content":
|
|
var outer struct {
|
|
Content json.RawMessage `json:"content"`
|
|
}
|
|
if err := json.Unmarshal(raw, &outer); err != nil || len(outer.Content) == 0 {
|
|
continue
|
|
}
|
|
var inner struct {
|
|
Type string `json:"type"`
|
|
Text string `json:"text"`
|
|
}
|
|
if err := json.Unmarshal(outer.Content, &inner); err != nil {
|
|
continue
|
|
}
|
|
if inner.Type != "text" {
|
|
continue
|
|
}
|
|
appendPiece(inner.Text)
|
|
case "diff":
|
|
var diff struct {
|
|
Path string `json:"path"`
|
|
OldText string `json:"oldText"`
|
|
NewText string `json:"newText"`
|
|
}
|
|
if err := json.Unmarshal(raw, &diff); err != nil || diff.Path == "" {
|
|
continue
|
|
}
|
|
// Keep it tiny — a full unified diff can be huge and we're
|
|
// really just recording "this tool wrote to this file".
|
|
// The UI can re-read the file if it needs the actual content.
|
|
var piece strings.Builder
|
|
piece.WriteString("--- ")
|
|
piece.WriteString(diff.Path)
|
|
piece.WriteString("\n+++ ")
|
|
piece.WriteString(diff.Path)
|
|
if diff.OldText == "" {
|
|
piece.WriteString("\n(new file, ")
|
|
piece.WriteString(strconv.Itoa(len(diff.NewText)))
|
|
piece.WriteString(" bytes)")
|
|
} else {
|
|
piece.WriteString("\n(edited: ")
|
|
piece.WriteString(strconv.Itoa(len(diff.OldText)))
|
|
piece.WriteString(" → ")
|
|
piece.WriteString(strconv.Itoa(len(diff.NewText)))
|
|
piece.WriteString(" bytes)")
|
|
}
|
|
appendPiece(piece.String())
|
|
default:
|
|
// terminal blocks, image blocks, unknown future types —
|
|
// ignore. We have no way to inline-render them.
|
|
}
|
|
}
|
|
return b.String()
|
|
}
|
|
|
|
func (c *hermesClient) handleUsageUpdate(data json.RawMessage) {
|
|
var msg struct {
|
|
Usage json.RawMessage `json:"usage"`
|
|
}
|
|
if err := json.Unmarshal(data, &msg); err != nil {
|
|
return
|
|
}
|
|
usage := parseACPTokenUsage(msg.Usage)
|
|
|
|
c.usageMu.Lock()
|
|
// Usage updates from ACP are cumulative snapshots, so take the latest.
|
|
if usage.InputTokens > c.usage.InputTokens {
|
|
c.usage.InputTokens = usage.InputTokens
|
|
}
|
|
if usage.OutputTokens > c.usage.OutputTokens {
|
|
c.usage.OutputTokens = usage.OutputTokens
|
|
}
|
|
if usage.CacheReadTokens > c.usage.CacheReadTokens {
|
|
c.usage.CacheReadTokens = usage.CacheReadTokens
|
|
}
|
|
if usage.CacheWriteTokens > c.usage.CacheWriteTokens {
|
|
c.usage.CacheWriteTokens = usage.CacheWriteTokens
|
|
}
|
|
if usage.CostUSDTicks > c.usage.CostUSDTicks {
|
|
c.usage.CostUSDTicks = usage.CostUSDTicks
|
|
}
|
|
c.usageMu.Unlock()
|
|
}
|
|
|
|
func parseACPTokenUsage(data json.RawMessage) TokenUsage {
|
|
if len(data) == 0 || string(data) == "null" {
|
|
return TokenUsage{}
|
|
}
|
|
var fields map[string]json.RawMessage
|
|
if err := json.Unmarshal(data, &fields); err != nil {
|
|
return TokenUsage{}
|
|
}
|
|
usage := TokenUsage{
|
|
InputTokens: acpUsageInt64(fields, "inputTokens", "input_tokens"),
|
|
OutputTokens: acpUsageInt64(fields, "outputTokens", "output_tokens"),
|
|
CacheReadTokens: acpUsageInt64(fields,
|
|
"cachedReadTokens",
|
|
"cacheReadTokens",
|
|
"cached_input_tokens",
|
|
"cache_read_tokens",
|
|
"cache_read_input_tokens",
|
|
),
|
|
CacheWriteTokens: acpUsageInt64(fields,
|
|
"cachedWriteTokens",
|
|
"cacheWriteTokens",
|
|
"cache_write_tokens",
|
|
"cache_creation_input_tokens",
|
|
),
|
|
// The provider's own price for this turn, already inclusive of
|
|
// request-level pricing rules we cannot reconstruct from token
|
|
// counts (see TokenUsage.CostUSDTicks).
|
|
CostUSDTicks: acpUsageInt64(fields, "costUsdTicks", "cost_usd_ticks"),
|
|
}
|
|
return excludeACPCachedInput(usage, acpUsageInt64(fields, "totalTokens", "total_tokens"))
|
|
}
|
|
|
|
// excludeACPCachedInput re-buckets a usage record whose `inputTokens` already
|
|
// contains `cachedReadTokens`, so the persisted buckets stay mutually
|
|
// exclusive and dashboard cost math does not charge the cached prefix twice
|
|
// (same normalization codex.go applies via codexUncachedInputTokens).
|
|
//
|
|
// ACP does not specify whether cached reads are counted inside inputTokens.
|
|
// Grok Build counts them inside: a real `grok 0.2.106` turn reports
|
|
// inputTokens=12929, cachedReadTokens=10880, outputTokens=29,
|
|
// totalTokens=12958 — i.e. total == input + output, so the cached prefix is
|
|
// counted once, within input. The same payload's costUsdTicks=75360000
|
|
// ($0.007536) matches exactly (12929-10880) uncached input + 10880 cached
|
|
// read + 29 output at xAI's published grok-4.5 rates, confirming how xAI
|
|
// bills it. Kept raw, that turn is priced as if 12929 tokens were uncached —
|
|
// ~4x the real spend on a cache-heavy turn.
|
|
//
|
|
// `totalTokens` is the only self-describing signal available, so the
|
|
// re-bucketing only happens when it is present and equals input + output.
|
|
// Agents that report exclusive buckets (total == input + cached + output) or
|
|
// omit totalTokens keep their counters untouched.
|
|
func excludeACPCachedInput(usage TokenUsage, totalTokens int64) TokenUsage {
|
|
if totalTokens <= 0 || usage.CacheReadTokens <= 0 || usage.CacheReadTokens > usage.InputTokens {
|
|
return usage
|
|
}
|
|
if totalTokens != usage.InputTokens+usage.OutputTokens {
|
|
return usage
|
|
}
|
|
usage.InputTokens -= usage.CacheReadTokens
|
|
return usage
|
|
}
|
|
|
|
func acpUsageInt64(fields map[string]json.RawMessage, names ...string) int64 {
|
|
for _, name := range names {
|
|
raw, ok := fields[name]
|
|
if !ok {
|
|
continue
|
|
}
|
|
var n json.Number
|
|
dec := json.NewDecoder(bytes.NewReader(raw))
|
|
dec.UseNumber()
|
|
if err := dec.Decode(&n); err == nil {
|
|
if v, err := n.Int64(); err == nil {
|
|
return v
|
|
}
|
|
if f, err := n.Float64(); err == nil {
|
|
return int64(f)
|
|
}
|
|
}
|
|
var s string
|
|
if err := json.Unmarshal(raw, &s); err == nil {
|
|
if v, err := strconv.ParseInt(strings.TrimSpace(s), 10, 64); err == nil {
|
|
return v
|
|
}
|
|
}
|
|
}
|
|
return 0
|
|
}
|
|
|
|
// ── Helpers ──
|
|
|
|
// extractACPSessionID pulls `sessionId` out of a session/new or
|
|
// session/resume response. Shared by all ACP backends (hermes, kimi, kiro,
|
|
// and anything else that follows the standard ACP schema).
|
|
func extractACPSessionID(result json.RawMessage) string {
|
|
var r struct {
|
|
SessionID string `json:"sessionId"`
|
|
}
|
|
if err := json.Unmarshal(result, &r); err != nil {
|
|
return ""
|
|
}
|
|
return r.SessionID
|
|
}
|
|
|
|
// extractACPAuthMethods returns the `authMethods` ids advertised in an ACP
|
|
// `initialize` response, in the order the agent listed them. Agents that
|
|
// require authentication (e.g. xAI's Grok Build) enumerate the accepted
|
|
// methods here; per the ACP flow the client MUST send `authenticate` with one
|
|
// of these ids before `session/new` / `session/load`. Agents that need no
|
|
// explicit auth omit the field, so an empty slice means "skip authenticate".
|
|
// A malformed response degrades to an empty slice (fail open on parsing so we
|
|
// don't wedge agents that never needed the step).
|
|
func extractACPAuthMethods(result json.RawMessage) []string {
|
|
var r struct {
|
|
AuthMethods []struct {
|
|
ID string `json:"id"`
|
|
} `json:"authMethods"`
|
|
}
|
|
if err := json.Unmarshal(result, &r); err != nil {
|
|
return nil
|
|
}
|
|
ids := make([]string, 0, len(r.AuthMethods))
|
|
for _, m := range r.AuthMethods {
|
|
if id := strings.TrimSpace(m.ID); id != "" {
|
|
ids = append(ids, id)
|
|
}
|
|
}
|
|
return ids
|
|
}
|
|
|
|
// extractACPCurrentModelID pulls the model selected by the ACP runtime out of
|
|
// a session/new or session/resume response. Hermes returns this when it uses
|
|
// its own default model, so token usage can still be attributed to a real model
|
|
// even when Multica did not pass an explicit agent.model override.
|
|
func extractACPCurrentModelID(result json.RawMessage) string {
|
|
var r struct {
|
|
Models struct {
|
|
CurrentModelID string `json:"currentModelId"`
|
|
CurrentModelIDSnake string `json:"current_model_id"`
|
|
} `json:"models"`
|
|
CurrentModelID string `json:"currentModelId"`
|
|
CurrentModelIDSnake string `json:"current_model_id"`
|
|
}
|
|
if err := json.Unmarshal(result, &r); err != nil {
|
|
return ""
|
|
}
|
|
for _, candidate := range []string{
|
|
r.Models.CurrentModelID,
|
|
r.Models.CurrentModelIDSnake,
|
|
r.CurrentModelID,
|
|
r.CurrentModelIDSnake,
|
|
} {
|
|
if model := strings.TrimSpace(candidate); model != "" {
|
|
return model
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
// resolveResumedSessionID picks which session id we should treat as live
|
|
// after a `session/resume` round-trip. Hermes (and other ACP servers)
|
|
// return the canonical sessionId in the response — when the local
|
|
// state.db has been wiped, the server silently creates a brand-new
|
|
// session and returns its new id rather than failing. If we keep using
|
|
// our requested id in that case, every subsequent session/prompt is
|
|
// addressed to a session the server doesn't know about and fails with
|
|
// JSON-RPC -32603. Returns (chosenID, changed). When the response is
|
|
// malformed or omits sessionId we fall back to the requested id so the
|
|
// happy path keeps working against older / non-conforming servers.
|
|
func resolveResumedSessionID(requested string, response json.RawMessage) (string, bool) {
|
|
got := extractACPSessionID(response)
|
|
if got == "" {
|
|
return requested, false
|
|
}
|
|
return got, got != requested
|
|
}
|
|
|
|
// buildHermesSessionParams constructs the params map for the ACP `session/new`
|
|
// request. The `model` field is only included when non-empty so Hermes falls
|
|
// back to its default only when no explicit model was configured.
|
|
//
|
|
// mcpServers should be the ACP-shaped array produced by buildACPMcpServers
|
|
// from the agent's mcp_config; a nil slice is normalised to an empty array
|
|
// so the wire request always carries the field (ACP requires it).
|
|
func buildHermesSessionParams(cwd, model string, mcpServers []any) map[string]any {
|
|
if mcpServers == nil {
|
|
mcpServers = []any{}
|
|
}
|
|
params := map[string]any{
|
|
"cwd": cwd,
|
|
"mcpServers": mcpServers,
|
|
}
|
|
if model != "" {
|
|
params["model"] = model
|
|
}
|
|
return params
|
|
}
|
|
|
|
// buildACPMcpServers translates an agent's Claude-style mcp_config
|
|
// (`{"mcpServers": {"<name>": {...}}}`) into the array shape that ACP's
|
|
// `session/new` and `session/load` requests expect.
|
|
//
|
|
// Each Claude-style entry maps to one of:
|
|
//
|
|
// - Stdio: `{name, command, args, env: [{name,value}, ...]}` —
|
|
// when the entry has a `command` field. No `type` field is emitted;
|
|
// ACP treats untagged entries as stdio.
|
|
// - HTTP / SSE: `{type, name, url, headers: [{name,value}, ...]}` —
|
|
// when the entry has a `url` field. `type` defaults to "http"; Claude's
|
|
// "sse" and "streamable-http" / "http_streamable" aliases are accepted.
|
|
//
|
|
// Empty / null input returns an empty slice — the launch proceeds with no
|
|
// MCP servers (the existing default for ACP backends). Malformed top-level
|
|
// JSON returns an error so the launch fails closed, mirroring codex's
|
|
// `renderCodexMcpServersBlock` contract. Individual entries that have
|
|
// neither `command` nor `url` are skipped with a warning rather than
|
|
// failing the whole launch, so a single bad entry can't kill the agent.
|
|
//
|
|
// Output entries are sorted by name and each entry's env / headers are
|
|
// sorted by key, so the wire request is deterministic across reruns —
|
|
// useful for tests, log diffs, and reproducibility.
|
|
func buildACPMcpServers(raw json.RawMessage, logger *slog.Logger) ([]any, error) {
|
|
trimmed := bytes.TrimSpace(raw)
|
|
if len(trimmed) == 0 || bytes.Equal(trimmed, []byte("null")) {
|
|
return []any{}, nil
|
|
}
|
|
var parsed struct {
|
|
McpServers map[string]json.RawMessage `json:"mcpServers"`
|
|
}
|
|
if err := json.Unmarshal(trimmed, &parsed); err != nil {
|
|
return nil, fmt.Errorf("parse mcp_config json: %w", err)
|
|
}
|
|
if len(parsed.McpServers) == 0 {
|
|
return []any{}, nil
|
|
}
|
|
|
|
names := make([]string, 0, len(parsed.McpServers))
|
|
for name := range parsed.McpServers {
|
|
names = append(names, name)
|
|
}
|
|
sort.Strings(names)
|
|
|
|
out := make([]any, 0, len(names))
|
|
for _, name := range names {
|
|
entry, err := convertACPMcpServer(name, parsed.McpServers[name])
|
|
if err != nil {
|
|
if logger != nil {
|
|
logger.Warn("skipping invalid mcp_config entry", "name", name, "error", err)
|
|
}
|
|
continue
|
|
}
|
|
out = append(out, entry)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// convertACPMcpServer converts a single Claude-style entry into the ACP
|
|
// McpServer wire shape. Returns an error for entries that can't be
|
|
// classified (no command and no url).
|
|
func convertACPMcpServer(name string, raw json.RawMessage) (map[string]any, error) {
|
|
var entry struct {
|
|
Type string `json:"type"`
|
|
Command string `json:"command"`
|
|
Args []string `json:"args"`
|
|
Env map[string]string `json:"env"`
|
|
URL string `json:"url"`
|
|
Headers map[string]string `json:"headers"`
|
|
}
|
|
if err := json.Unmarshal(raw, &entry); err != nil {
|
|
return nil, fmt.Errorf("parse entry: %w", err)
|
|
}
|
|
|
|
command := strings.TrimSpace(entry.Command)
|
|
url := strings.TrimSpace(entry.URL)
|
|
|
|
if command != "" {
|
|
args := entry.Args
|
|
if args == nil {
|
|
args = []string{}
|
|
}
|
|
envArr := make([]map[string]any, 0, len(entry.Env))
|
|
for _, k := range sortedStringMapKeys(entry.Env) {
|
|
envArr = append(envArr, map[string]any{
|
|
"name": k,
|
|
"value": entry.Env[k],
|
|
})
|
|
}
|
|
return map[string]any{
|
|
"name": name,
|
|
"command": command,
|
|
"args": args,
|
|
"env": envArr,
|
|
}, nil
|
|
}
|
|
|
|
if url != "" {
|
|
t := strings.ToLower(strings.TrimSpace(entry.Type))
|
|
switch t {
|
|
case "sse":
|
|
t = "sse"
|
|
case "", "http", "streamable-http", "http_streamable":
|
|
t = "http"
|
|
default:
|
|
// Unknown remote transport — degrade to "http" rather than fail.
|
|
// ACP servers that don't recognise the type will reject the
|
|
// session/new request and surface a real error to the user.
|
|
t = "http"
|
|
}
|
|
headerArr := make([]map[string]any, 0, len(entry.Headers))
|
|
for _, k := range sortedStringMapKeys(entry.Headers) {
|
|
headerArr = append(headerArr, map[string]any{
|
|
"name": k,
|
|
"value": entry.Headers[k],
|
|
})
|
|
}
|
|
return map[string]any{
|
|
"type": t,
|
|
"name": name,
|
|
"url": url,
|
|
"headers": headerArr,
|
|
}, nil
|
|
}
|
|
|
|
return nil, fmt.Errorf("entry has neither command nor url")
|
|
}
|
|
|
|
func sortedStringMapKeys(m map[string]string) []string {
|
|
keys := make([]string, 0, len(m))
|
|
for k := range m {
|
|
keys = append(keys, k)
|
|
}
|
|
sort.Strings(keys)
|
|
return keys
|
|
}
|
|
|
|
// acpMcpTransportCapabilities reports which remote MCP transports the ACP
|
|
// runtime advertised in its `initialize` response. Stdio is always
|
|
// supported (it's the baseline transport and the spec does not gate it),
|
|
// so it's not represented here.
|
|
type acpMcpTransportCapabilities struct {
|
|
HTTP bool
|
|
SSE bool
|
|
}
|
|
|
|
// extractACPMcpCapabilities reads `agentCapabilities.mcpCapabilities.http`
|
|
// and `.sse` out of an ACP `initialize` response. Missing or false fields
|
|
// stay false, matching the spec default: the runtime must opt-in to
|
|
// remote MCP transports. Unparseable responses degrade to "neither
|
|
// supported" so we fail closed on remote entries.
|
|
//
|
|
// See https://agentclientprotocol.com/protocol/initialization — clients
|
|
// MUST NOT send `mcpServers` entries with a type the agent did not
|
|
// advertise support for.
|
|
func extractACPMcpCapabilities(result json.RawMessage) acpMcpTransportCapabilities {
|
|
var r struct {
|
|
AgentCapabilities struct {
|
|
McpCapabilities struct {
|
|
HTTP bool `json:"http"`
|
|
SSE bool `json:"sse"`
|
|
} `json:"mcpCapabilities"`
|
|
} `json:"agentCapabilities"`
|
|
}
|
|
if err := json.Unmarshal(result, &r); err != nil {
|
|
return acpMcpTransportCapabilities{}
|
|
}
|
|
return acpMcpTransportCapabilities{
|
|
HTTP: r.AgentCapabilities.McpCapabilities.HTTP,
|
|
SSE: r.AgentCapabilities.McpCapabilities.SSE,
|
|
}
|
|
}
|
|
|
|
// filterACPMcpServersByCapability drops remote MCP entries whose transport
|
|
// the runtime didn't advertise in its initialize response. Stdio entries
|
|
// (no `type` field) always pass through.
|
|
//
|
|
// Sending an http/sse entry to a runtime that doesn't support it is a
|
|
// protocol violation per the ACP spec, and Hermes / Kimi observed in
|
|
// practice reject the whole session/new request with a JSON-RPC error.
|
|
// Dropping the offending entries with a warning lets the rest of the
|
|
// session start and surfaces the problem in the daemon log instead of
|
|
// tanking every task on that agent.
|
|
func filterACPMcpServersByCapability(
|
|
servers []any,
|
|
caps acpMcpTransportCapabilities,
|
|
backend string,
|
|
logger *slog.Logger,
|
|
) []any {
|
|
if len(servers) == 0 {
|
|
return servers
|
|
}
|
|
filtered := make([]any, 0, len(servers))
|
|
for _, raw := range servers {
|
|
entry, ok := raw.(map[string]any)
|
|
if !ok {
|
|
filtered = append(filtered, raw)
|
|
continue
|
|
}
|
|
transport, _ := entry["type"].(string)
|
|
switch transport {
|
|
case "http":
|
|
if !caps.HTTP {
|
|
if logger != nil {
|
|
logger.Warn("dropping http MCP server: runtime did not advertise mcpCapabilities.http",
|
|
"backend", backend, "name", entry["name"])
|
|
}
|
|
continue
|
|
}
|
|
case "sse":
|
|
if !caps.SSE {
|
|
if logger != nil {
|
|
logger.Warn("dropping sse MCP server: runtime did not advertise mcpCapabilities.sse",
|
|
"backend", backend, "name", entry["name"])
|
|
}
|
|
continue
|
|
}
|
|
}
|
|
filtered = append(filtered, entry)
|
|
}
|
|
return filtered
|
|
}
|
|
|
|
// hermesToolNameFromTitle extracts a tool name from the ACP tool call title.
|
|
// Hermes ACP titles look like "terminal: ls -la", "read: /path/to/file", etc.
|
|
// Some titles have no colon (e.g. "execute code").
|
|
func hermesToolNameFromTitle(title string, kind string) string {
|
|
// Check exact-match titles first (no colon).
|
|
switch title {
|
|
case "execute code":
|
|
return "execute_code"
|
|
}
|
|
|
|
// Try to extract the tool name from before the first colon.
|
|
if idx := strings.Index(title, ":"); idx > 0 {
|
|
name := strings.TrimSpace(title[:idx])
|
|
// Map common ACP title prefixes back to tool names.
|
|
// Some titles include mode info like "patch (replace)", so check prefix.
|
|
switch {
|
|
case name == "terminal":
|
|
return "terminal"
|
|
case name == "read":
|
|
return "read_file"
|
|
case name == "write":
|
|
return "write_file"
|
|
case strings.HasPrefix(name, "patch"):
|
|
return "patch"
|
|
case name == "search":
|
|
return "search_files"
|
|
case name == "web search":
|
|
return "web_search"
|
|
case name == "extract":
|
|
return "web_extract"
|
|
case name == "delegate":
|
|
return "delegate_task"
|
|
case name == "analyze image":
|
|
return "vision_analyze"
|
|
}
|
|
return name
|
|
}
|
|
|
|
// Fall back to kind.
|
|
switch kind {
|
|
case "read":
|
|
return "read_file"
|
|
case "edit":
|
|
return "write_file"
|
|
case "execute":
|
|
return "terminal"
|
|
case "search":
|
|
return "search_files"
|
|
case "fetch":
|
|
return "web_search"
|
|
case "think":
|
|
return "thinking"
|
|
default:
|
|
// Preserve a non-empty title when we can't classify it: kimi
|
|
// emits bare titles like "Shell" or "Read file" without any
|
|
// `kind`, so returning an empty string here drops the tool
|
|
// name entirely before kimiToolNameFromTitle can map it.
|
|
// Hermes titles always carry a colon, so hermes never reaches
|
|
// this branch with a non-empty title.
|
|
if title != "" {
|
|
return title
|
|
}
|
|
return kind
|
|
}
|
|
}
|
|
|
|
// ── Provider-error sniffing ──
|
|
//
|
|
// ACP agents (hermes, kimi, …) all have the same failure mode:
|
|
// session/prompt reports stopReason=end_turn even when the underlying
|
|
// HTTP call to the configured LLM endpoint returned an error — the
|
|
// actionable detail only appears on stderr (e.g.
|
|
// `⚠️ API call failed (attempt 1/3): BadRequestError [HTTP 400]` and
|
|
// `Error: HTTP 400: Error code: 400 - {'detail': "The '...' model
|
|
// is not supported when using Codex with a ChatGPT account."}`).
|
|
// The sniffer scans for those patterns so the daemon can surface a
|
|
// real failure instead of a generic "empty output".
|
|
//
|
|
// Parameterised by provider name so both hermes and kimi can share
|
|
// the transport: the regexes match format-level signals (HTTP status,
|
|
// error-kind tags, "API call failed" banner) that both runtimes emit.
|
|
//
|
|
// The sniffer distinguishes *transient* per-attempt warnings (e.g.
|
|
// "API call failed (attempt 1/3): RateLimitError [HTTP 429]" — followed
|
|
// by a successful retry) from *terminal* exhausted failures (e.g.
|
|
// "API call failed after 3 retries: ..." or "❌ ... Non-retryable"):
|
|
// `message()` returns whichever was last seen, while `terminalMessage()`
|
|
// returns non-empty only when a terminal-failure marker was matched.
|
|
// Promotion to status="failed" must use `terminalMessage()`, otherwise
|
|
// a successful retry following an early per-attempt warning would be
|
|
// wrongly marked as failed.
|
|
type acpProviderErrorSniffer struct {
|
|
provider string
|
|
mu sync.Mutex
|
|
remains []byte // buffer for a partial trailing line across writes
|
|
lines []string // captured error lines, bounded
|
|
seen map[string]bool
|
|
terminal bool // sticky: at least one line matched acpTerminalErrorRe
|
|
}
|
|
|
|
// acpErrorHeaderRe matches the first line of an API-error block.
|
|
// ACP agents typically prefix these with ⚠️ / ❌ and include an HTTP
|
|
// status code or a non-retryable-error tag.
|
|
var acpErrorHeaderRe = regexp.MustCompile(`(?:⚠️|❌|\[ERROR\]).*(?:BadRequestError|AuthenticationError|RateLimitError|HTTP [0-9]{3}|Non-retryable|API call failed)`)
|
|
|
|
// acpErrorDetailRe pulls the most useful single-line messages out of
|
|
// the subsequent lines of the error block (the one whose "Error:" or
|
|
// "Details:" tag actually spells out what happened).
|
|
var acpErrorDetailRe = regexp.MustCompile(`(?:Error:|detail:|Details:)\s*(.+)`)
|
|
|
|
// acpTerminalErrorRe matches markers that only appear when the
|
|
// adapter has *given up* on the upstream call — either after
|
|
// exhausting retries ("after N retries"), or because the error is
|
|
// classified as non-retryable up front (Non-retryable, BadRequest /
|
|
// Authentication errors, ❌ / [ERROR] log levels). Per-attempt
|
|
// warnings ("(attempt 1/3)") deliberately do NOT match this pattern.
|
|
var acpTerminalErrorRe = regexp.MustCompile(`(?:❌|\[ERROR\]|after \d+ retr|Non-retryable|BadRequestError|AuthenticationError)`)
|
|
|
|
// acpAgentOutputTerminalRe matches the synthetic agent-text turn that
|
|
// hermes-style ACP adapters inject when they exhaust retries against
|
|
// the upstream LLM ("API call failed after 3 retries: HTTP 429..."),
|
|
// surfaced via session/update agent_message_chunk and ending up in the
|
|
// final output buffer. Per-attempt warnings (which only go to stderr
|
|
// and use "(attempt N/M)" phrasing) won't match.
|
|
var acpAgentOutputTerminalRe = regexp.MustCompile(`API call failed after \d+ retr(?:y|ies)`)
|
|
|
|
const acpMaxErrorLines = 8
|
|
|
|
// newACPProviderErrorSniffer returns a sniffer that tags its messages
|
|
// with the given provider name (e.g. "hermes", "kimi") so failure
|
|
// strings make it obvious which runtime produced the error.
|
|
func newACPProviderErrorSniffer(provider string) *acpProviderErrorSniffer {
|
|
return &acpProviderErrorSniffer{provider: provider, seen: map[string]bool{}}
|
|
}
|
|
|
|
// Write implements io.Writer so the sniffer can sit behind an
|
|
// io.MultiWriter next to the normal stderr log forwarder.
|
|
func (s *acpProviderErrorSniffer) Write(p []byte) (int, error) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
data := append(s.remains, p...)
|
|
// Keep the final partial line (no trailing newline) for the
|
|
// next write so multi-line error blocks aren't split.
|
|
nl := strings.LastIndexByte(string(data), '\n')
|
|
var complete string
|
|
if nl < 0 {
|
|
s.remains = append(s.remains[:0], data...)
|
|
return len(p), nil
|
|
}
|
|
complete = string(data[:nl])
|
|
s.remains = append(s.remains[:0], data[nl+1:]...)
|
|
|
|
for _, line := range strings.Split(complete, "\n") {
|
|
line = strings.TrimSpace(line)
|
|
if line == "" {
|
|
continue
|
|
}
|
|
if !(acpErrorHeaderRe.MatchString(line) || acpErrorDetailRe.MatchString(line)) {
|
|
continue
|
|
}
|
|
if acpTerminalErrorRe.MatchString(line) {
|
|
s.terminal = true
|
|
}
|
|
if s.seen[line] {
|
|
continue
|
|
}
|
|
s.seen[line] = true
|
|
s.lines = append(s.lines, line)
|
|
if len(s.lines) > acpMaxErrorLines {
|
|
s.lines = s.lines[len(s.lines)-acpMaxErrorLines:]
|
|
}
|
|
}
|
|
return len(p), nil
|
|
}
|
|
|
|
// message returns a single-line summary suitable for the task
|
|
// error field. Prefers the most specific "Error:" / "detail:"
|
|
// fragment; falls back to the first captured header line; empty
|
|
// when nothing useful was seen.
|
|
//
|
|
// NOTE: a non-empty message() can describe a *transient* per-attempt
|
|
// warning that was followed by a successful retry. Code that flips
|
|
// task status to "failed" must instead use terminalMessage() — see
|
|
// the type doc above.
|
|
func (s *acpProviderErrorSniffer) message() string {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
return s.messageLocked()
|
|
}
|
|
|
|
// terminalMessage returns the same single-line summary as message()
|
|
// but only when the sniffer has seen at least one line matching
|
|
// acpTerminalErrorRe — i.e. the adapter has given up retrying. This
|
|
// is the signal callers should use to decide whether to promote a
|
|
// run from "completed" to "failed". Returns empty if all captured
|
|
// lines look like transient retry warnings.
|
|
func (s *acpProviderErrorSniffer) terminalMessage() string {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
if !s.terminal {
|
|
return ""
|
|
}
|
|
return s.messageLocked()
|
|
}
|
|
|
|
// messageLocked is the lock-held implementation shared by message()
|
|
// and terminalMessage(). Caller must hold s.mu.
|
|
func (s *acpProviderErrorSniffer) messageLocked() string {
|
|
prefix := s.provider + " provider error: "
|
|
for _, line := range s.lines {
|
|
if m := acpErrorDetailRe.FindStringSubmatch(line); m != nil {
|
|
detail := strings.TrimSpace(m[1])
|
|
if detail != "" {
|
|
return prefix + detail
|
|
}
|
|
}
|
|
}
|
|
for _, line := range s.lines {
|
|
if acpErrorHeaderRe.MatchString(line) {
|
|
return prefix + line
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
// promoteACPResultOnProviderError flips finalStatus to "failed" if
|
|
// either (a) the stderr sniffer captured a terminal-failure marker,
|
|
// (b) the adapter injected a synthetic "API call failed after N
|
|
// retries..." turn into the agent text stream, or (c) output was
|
|
// empty AND the sniffer captured anything at all (no real result to
|
|
// fall back on, even from a transient-only sequence). Returns the
|
|
// updated (status, error) pair; callers should overwrite their
|
|
// locals with the result.
|
|
//
|
|
// This is the shared post-processing step for hermes/kimi/kiro.
|
|
// Without it, runs that exhaust retries against the upstream LLM
|
|
// (HTTP 429, expired token, …) silently report as "completed"
|
|
// because session/prompt still ends with stopReason=end_turn — see
|
|
// GitHub multica#1952.
|
|
func promoteACPResultOnProviderError(finalStatus, finalError, finalOutput string, sniffer *acpProviderErrorSniffer) (string, string) {
|
|
if finalStatus != "completed" {
|
|
return finalStatus, finalError
|
|
}
|
|
if msg := sniffer.terminalMessage(); msg != "" {
|
|
return "failed", msg
|
|
}
|
|
if acpAgentOutputTerminalRe.MatchString(finalOutput) {
|
|
msg := sniffer.message()
|
|
if msg == "" {
|
|
msg = sniffer.provider + " provider error: " + acpAgentOutputTerminalRe.FindString(finalOutput)
|
|
}
|
|
return "failed", msg
|
|
}
|
|
if finalOutput == "" {
|
|
if msg := sniffer.message(); msg != "" {
|
|
return "failed", msg
|
|
}
|
|
}
|
|
return finalStatus, finalError
|
|
}
|