krow-workforce-agent has declared five subagents since it was written and
answered every question by itself. Everything for §6 existed except the
delegation: the parser read `subagents:`, runtime.Agent carried them, the
loader populated them, agent_runs had a parent_run_id column with a
self-reference and a no-self-parent constraint, and budget.go's comments
already described sharing a budget with subagents. Nothing called any of it.
A subagent is offered to the parent's model as a tool, because §6 says that is
what delegation is from the parent's side. Three rules are enforced rather than
assumed, each with a test that fails if it stops holding:
I1 The subagent runs as the ORIGINAL caller. It cannot read anything the
person could not read directly.
§6 It SHARES the parent's budget. The test sets MaxSteps to 1, spends it in
the parent, and asserts the child terminates BudgetExceeded — an
assertion that only passes when the budget is shared, and that a fresh
budget would quietly turn green.
§3 Depth is capped at 2. At the cap no subagent is loaded or offered, so a
cycle reaching run time is bounded rather than unbounded.
I4 survives too: a write a SUBAGENT wants approved still stops the whole run
and asks a person, rather than being performed because it happened one level
down.
Two bugs found by running it rather than by reading it:
- delegate() read the error before the result. finish returns a non-nil
error for every termination that is not Completed, INCLUDING
ConfirmationPending — which is not a failure but a run that stopped to ask
a question. Reading the error first discarded the result and with it the
confirmation, so a subagent's write silently never happened and nobody was
asked.
- Delegated trajectories were never persisted at all. parent_run_id is a
foreign key and a subagent finishes BEFORE the run that delegated to it,
so every child insert named a parent row that did not exist yet. The
database refused it; finish deliberately does not fail a run over a sink
error; and the entry recording that the trajectory could not be saved was
itself in the trajectory that was not saved. Children are now buffered and
written by finish after the parent's own row, each arriving with its
descendants already ordered behind it, so one pass writes a whole tree
parent-first. The regression test asserts on save ORDER, because a
MemorySink has no foreign key and will pass either way.
Verified end to end against a live model: an agent with no tools of its own and
one subagent produced
delegation-probe run=run_16622d7de6 parent=(root)
talent-pool-agent run=run_64160684b1 parent=run_16622d7de6
with the subagent's answer reaching the parent's model. Full suite green, only
TestLive* skipped.
Not addressed: §3's publish-time cycle detection, which needs the whole agent
set in hand. The depth cap is what holds without it, and is the half that
matters at run time.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PJvibeSc1JYXjatankqM1g
302 lines
10 KiB
Go
302 lines
10 KiB
Go
package runtime
|
|
|
|
import (
|
|
"context"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/krow/krow-backend/go-api/internal/knowledge"
|
|
)
|
|
|
|
// EntryKind is what one line of a trajectory records.
|
|
type EntryKind string
|
|
|
|
const (
|
|
EntryMessage EntryKind = "message"
|
|
EntryToolCall EntryKind = "tool_call"
|
|
EntryToolResult EntryKind = "tool_result"
|
|
EntryBudget EntryKind = "budget"
|
|
EntryError EntryKind = "error"
|
|
// EntryConfirmation is a write that was described and not performed. It is
|
|
// recorded because "what was this person asked to approve, and when" is the
|
|
// question an audit of an agent-initiated write actually asks.
|
|
EntryConfirmation EntryKind = "confirmation"
|
|
// EntryRetrieval is what the knowledge layer returned for this run. The
|
|
// chunk ids are recorded, never the chunk text: a trajectory is already the
|
|
// most sensitive row in the database, and duplicating the corpus into it
|
|
// would mean a retention policy on runs quietly became a retention policy on
|
|
// every document too.
|
|
EntryRetrieval EntryKind = "retrieval"
|
|
)
|
|
|
|
// Entry is one recorded moment in a run.
|
|
//
|
|
// Deliberately flat and deliberately typed as data rather than prose: an eval
|
|
// asserts on `tools_called`, a debugger reads the budget line before the
|
|
// dispatch that overran, and neither can do that against a log string.
|
|
type Entry struct {
|
|
Seq int `json:"seq"`
|
|
At time.Time `json:"at"`
|
|
Kind EntryKind `json:"kind"`
|
|
Role string `json:"role,omitempty"`
|
|
Name string `json:"name,omitempty"`
|
|
Text string `json:"text,omitempty"`
|
|
Data any `json:"data,omitempty"`
|
|
Budget *Snapshot `json:"budget,omitempty"`
|
|
ErrorCode string `json:"errorCode,omitempty"`
|
|
|
|
// Effect is the tool's declared effect, on tool_call and tool_result
|
|
// entries. Recorded because "did this run change anything" is not answerable
|
|
// from a tool name — an eval reading the trajectory would otherwise have to
|
|
// keep its own list of which tools write, and that list would go stale on
|
|
// the first tool anyone added.
|
|
//
|
|
// It records what the RUNTIME BELIEVED, which is what the gate acted on. A
|
|
// tool that declares itself a read and writes anyway is invisible here, and
|
|
// is meant to be: the trajectory cannot be the check on a tool lying about
|
|
// itself. That check is the database.
|
|
Effect string `json:"effect,omitempty"`
|
|
|
|
// Failed says a tool result carried an error. A refusal is recorded like
|
|
// any other result, and an eval that could not tell the two apart would
|
|
// read every denial as a successful call.
|
|
Failed bool `json:"failed,omitempty"`
|
|
}
|
|
|
|
// Trajectory is the full record of one run.
|
|
//
|
|
// §6 is explicit that this is not optional telemetry: it is what makes
|
|
// debugging and evals possible at all. A run whose trajectory was dropped
|
|
// because the sink was busy is a run nobody can explain afterwards, which is
|
|
// why recording never blocks on persistence — see Recorder.
|
|
type Trajectory struct {
|
|
RunID string `json:"runId"`
|
|
ParentRunID string `json:"parentRunId,omitempty"`
|
|
OrgID string `json:"orgId"`
|
|
UserID string `json:"userId"`
|
|
AgentID string `json:"agentId"`
|
|
AgentVersion int `json:"agentVersion"`
|
|
Tier string `json:"tier"`
|
|
Model string `json:"model,omitempty"`
|
|
StartedAt time.Time `json:"startedAt"`
|
|
EndedAt time.Time `json:"endedAt"`
|
|
Termination Termination `json:"termination"`
|
|
Entries []Entry `json:"entries"`
|
|
Usage RunUsage `json:"usage"`
|
|
}
|
|
|
|
// RunUsage is what a whole run cost, across every call it made.
|
|
type RunUsage struct {
|
|
InputTokens int64 `json:"inputTokens"`
|
|
OutputTokens int64 `json:"outputTokens"`
|
|
CachedTokens int64 `json:"cachedTokens"`
|
|
TotalTokens int64 `json:"totalTokens"`
|
|
ModelCalls int `json:"modelCalls"`
|
|
}
|
|
|
|
// Sink persists a finished trajectory.
|
|
//
|
|
// An interface with one method so the eval harness can hold runs in memory and
|
|
// the service can write them to Postgres without either knowing about the
|
|
// other. A sink that fails must not fail the run: the answer was already
|
|
// produced, and losing the record is worse than losing nothing but is not
|
|
// worth discarding a correct answer over.
|
|
type Sink interface {
|
|
Save(ctx context.Context, t *Trajectory) error
|
|
}
|
|
|
|
// DiscardSink drops trajectories. The default, so a service wired without a
|
|
// store still runs — and so tests that do not care about persistence say so by
|
|
// choosing it rather than by leaving a nil that panics.
|
|
type DiscardSink struct{}
|
|
|
|
func (DiscardSink) Save(context.Context, *Trajectory) error { return nil }
|
|
|
|
// MemorySink keeps trajectories in memory. For tests and the eval harness.
|
|
type MemorySink struct {
|
|
mu sync.Mutex
|
|
Runs []*Trajectory
|
|
}
|
|
|
|
func (m *MemorySink) Save(_ context.Context, t *Trajectory) error {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
m.Runs = append(m.Runs, t)
|
|
return nil
|
|
}
|
|
|
|
// Last returns the most recent trajectory, or nil.
|
|
func (m *MemorySink) Last() *Trajectory {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if len(m.Runs) == 0 {
|
|
return nil
|
|
}
|
|
return m.Runs[len(m.Runs)-1]
|
|
}
|
|
|
|
// Recorder accumulates a trajectory during a run.
|
|
//
|
|
// Entries are held in memory and written once at the end rather than streamed
|
|
// per line. A run is short and bounded by construction — I3 guarantees it —
|
|
// so the whole record fits, and one write means a trajectory is either wholly
|
|
// there or wholly absent, never a half-run that reads as a run that stopped.
|
|
//
|
|
// Safe for concurrent use: a subagent records into its own recorder, but tool
|
|
// calls within one run may be dispatched in parallel.
|
|
type Recorder struct {
|
|
mu sync.Mutex
|
|
t *Trajectory
|
|
|
|
// children are trajectories of runs delegated beneath this one, already
|
|
// in parent-before-child order.
|
|
//
|
|
// They are held rather than saved as they finish because agent_runs.
|
|
// parent_run_id is a FOREIGN KEY, and a subagent always finishes before
|
|
// the run that delegated to it. Saving in completion order means every
|
|
// child insert names a parent row that does not exist yet, which the
|
|
// database refuses — and finish deliberately does not fail a run over a
|
|
// sink error, so the whole thing vanished silently. The trajectory
|
|
// recording that it could not be saved was itself the one not saved.
|
|
children []*Trajectory
|
|
}
|
|
|
|
// AddChildren queues trajectories delegated beneath this run, to be written
|
|
// after this run's own row exists.
|
|
func (r *Recorder) AddChildren(ts []*Trajectory) {
|
|
if len(ts) == 0 {
|
|
return
|
|
}
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
r.children = append(r.children, ts...)
|
|
}
|
|
|
|
// Children returns the queued delegated trajectories.
|
|
func (r *Recorder) Children() []*Trajectory {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
return append([]*Trajectory(nil), r.children...)
|
|
}
|
|
|
|
// NewRecorder begins recording a run.
|
|
func NewRecorder(t *Trajectory) *Recorder {
|
|
if t.StartedAt.IsZero() {
|
|
t.StartedAt = time.Now()
|
|
}
|
|
return &Recorder{t: t}
|
|
}
|
|
|
|
func (r *Recorder) append(e Entry) {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
e.Seq = len(r.t.Entries) + 1
|
|
if e.At.IsZero() {
|
|
e.At = time.Now()
|
|
}
|
|
r.t.Entries = append(r.t.Entries, e)
|
|
}
|
|
|
|
// RunID is the run being recorded, so a tool handler can be told which
|
|
// trajectory its call belongs to.
|
|
func (r *Recorder) RunID() string {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
return r.t.RunID
|
|
}
|
|
|
|
// Message records one turn of the conversation.
|
|
func (r *Recorder) Message(role, text string) {
|
|
r.append(Entry{Kind: EntryMessage, Role: role, Text: text})
|
|
}
|
|
|
|
// Budget records what was left before a dispatch.
|
|
//
|
|
// Called *before* the call it precedes, so the last budget line in a trajectory
|
|
// is the state that permitted the dispatch that ended the run — which is the
|
|
// line anyone debugging an overrun actually wants.
|
|
func (r *Recorder) Budget(s Snapshot) {
|
|
r.append(Entry{Kind: EntryBudget, Budget: &s})
|
|
}
|
|
|
|
// ToolCall records a dispatch to a tool.
|
|
func (r *Recorder) ToolCall(name, effect string, input any) {
|
|
r.append(Entry{Kind: EntryToolCall, Name: name, Effect: effect, Data: input})
|
|
}
|
|
|
|
// ToolResult records what a tool returned, and whether it failed.
|
|
func (r *Recorder) ToolResult(name, effect string, failed bool, output any) {
|
|
r.append(Entry{Kind: EntryToolResult, Name: name, Effect: effect, Failed: failed, Data: output})
|
|
}
|
|
|
|
// Retrieval records what the knowledge layer returned.
|
|
//
|
|
// Ids and ranks, not text. "Which chunks grounded this answer" is the question
|
|
// an eval and a debugger both ask, and it is answerable from ids alone — while
|
|
// copying the text in would make every trajectory a partial copy of the corpus,
|
|
// with all of the corpus's access rules and none of its retention.
|
|
func (r *Recorder) Retrieval(res *knowledge.Results) {
|
|
if res == nil {
|
|
return
|
|
}
|
|
cited := make([]map[string]any, 0, len(res.Chunks))
|
|
for _, c := range res.Chunks {
|
|
cited = append(cited, map[string]any{
|
|
"chunkId": c.ChunkID, "documentId": c.DocumentID,
|
|
"source": c.Source, "score": c.Score,
|
|
"denseRank": c.DenseRank, "sparseRank": c.SparseRank,
|
|
})
|
|
}
|
|
data := map[string]any{"chunks": cited, "tokens": res.TotalTokens}
|
|
if res.DenseSkipped != "" {
|
|
data["degraded"] = res.DenseSkipped
|
|
}
|
|
r.append(Entry{Kind: EntryRetrieval, Name: "knowledge", Data: data})
|
|
}
|
|
|
|
// Confirmation records a write that was described and is awaiting approval.
|
|
func (r *Recorder) Confirmation(name string, c any) {
|
|
r.append(Entry{Kind: EntryConfirmation, Name: name, Data: c})
|
|
}
|
|
|
|
// Error records a failure with the code that classified it.
|
|
func (r *Recorder) Error(code, message string) {
|
|
r.append(Entry{Kind: EntryError, ErrorCode: code, Text: message})
|
|
}
|
|
|
|
// ChargeUsage adds one model call's cost to the run total.
|
|
func (r *Recorder) ChargeUsage(input, output, cached int64) {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
r.t.Usage.InputTokens += input
|
|
r.t.Usage.OutputTokens += output
|
|
r.t.Usage.CachedTokens += cached
|
|
r.t.Usage.TotalTokens += input + output + cached
|
|
r.t.Usage.ModelCalls++
|
|
}
|
|
|
|
// SetModel records which model actually answered, as opposed to the tier that
|
|
// was requested.
|
|
func (r *Recorder) SetModel(model string) {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
r.t.Model = model
|
|
}
|
|
|
|
// Finish closes the trajectory with its termination reason and returns it.
|
|
//
|
|
// A reason that is not one of the six is recorded as ToolFailure rather than
|
|
// stored as-is: an unrecognised termination is a bug in the loop, and writing
|
|
// it verbatim would let that bug propagate into every eval and dashboard that
|
|
// groups by this column.
|
|
func (r *Recorder) Finish(t Termination) *Trajectory {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
if !t.Valid() {
|
|
t = TerminationToolFailure
|
|
}
|
|
r.t.Termination = t
|
|
r.t.EndedAt = time.Now()
|
|
return r.t
|
|
}
|