Files
krow_backend/go-api/internal/runtime/trajectory.go
2026-08-28 12:21:44 +05:30

272 lines
9.3 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
}
// 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
}