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