steploop

package
v0.6.1-rc.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 4, 2026 License: Apache-2.0 Imports: 16 Imported by: 0

Documentation

Overview

Package steploop drives one turn of the agent. Run loads the persisted transcript, issues one model request per step, settles the tool calls the model made, and repeats until the model stops calling tools, the step cap is reached, or the processor asks to stop. When the transcript records a pending compaction, the step is handed to the TaskController instead.

Seams:

  • Store is the message persistence service. Messages must return fresh, chronological values; the loop rebuilds its view from them every step.
  • LLMClient/PartStream are the narrow model-client seam; scripted tests return an in-memory stream.
  • ToolExecutor executes one already-resolved tool call. Tool discovery, permission checks and the concrete tools stay outside this package and arrive in RunOptions.Tools.
  • ModelResolver names the model for the latest user message.
  • TaskController owns compaction: overflow detection, boundary creation, processing and pruning.
  • The clock and the ascending ID factories are package seams with SetNowForTesting and SetIDFactoryForTesting restore closures.

Behaviour worth knowing:

  • The natural exit is: a non-empty finish other than "tool-calls", no non-provider-executed tool part on the matching persisted assistant, and the last user message older than the last assistant.
  • Tool cleanup snapshots all outstanding calls, waits for each with an independent short timeout, then force-writes every survivor as "Tool execution aborted".

Index

Constants

View Source
const MaxStepsPrompt = `` /* 750-byte string literal not displayed */

MaxStepsPrompt is appended to the prompt once the step cap is reached.

Variables

This section is empty.

Functions

func CompactedAfter

func CompactedAfter(msgs []msgmodel.WithParts, assistant msgmodel.Assistant) bool

CompactedAfter reports whether a completed compaction summary is newer (by ascending message ID) than the given assistant. The token count recorded on that assistant is then stale: the boundary has already replaced the context it measured. Re-checking it would create a boundary loop, because the verbatim tail keeps that very message -- and its count -- after every compaction, and the projection places the tail after the summary, where a newest-first scan finds it first.

func NewAscendingID

func NewAscendingID(prefix string) string

NewAscendingID lets production adapters mint session-layer records from the same ordered sequence as loop-owned messages and parts.

func SetIDFactoryForTesting

func SetIDFactoryForTesting(f func(prefix string) string) func()

SetIDFactoryForTesting swaps MessageID/PartID ascending generation. Prefix is "msg" or "prt". It returns a restore closure.

func SetNowForTesting

func SetNowForTesting(f func() uint64) func()

SetNowForTesting swaps the millisecond clock and returns a restore closure.

func ShouldExit

func ShouldExit(lastUser *msgmodel.User, lastAssistant *msgmodel.Assistant, msgs []msgmodel.WithParts) bool

ShouldExit reports the natural exit: the last assistant finished for a reason other than tool calls, none of its persisted tool parts is pending on this side (provider-executed tool parts do not count), and the last user message is older than it.

func ToolMessagesFromContext

func ToolMessagesFromContext(ctx context.Context) []msgmodel.WithParts

ToolMessagesFromContext returns the history attached by the live executor.

func WithToolMessages

func WithToolMessages(ctx context.Context, messages []msgmodel.WithParts) context.Context

WithToolMessages gives a tool the persisted session history that was current when execution began. Read uses it to deduplicate nested instruction files.

func WrapLateUserText

func WrapLateUserText(msgs []msgmodel.WithParts, lastFinished msgmodel.Assistant)

WrapLateUserText wraps the text of any user message that arrived after the last finished assistant in a system reminder, so the model treats it as an interjection. The input must be a fresh store load because this mutates text-part values in place.

Types

type BackScanResult

type BackScanResult struct {
	LastUser      *msgmodel.User
	LastAssistant *msgmodel.Assistant
	LastFinished  *msgmodel.Assistant
	Tasks         []msgmodel.Part
}

BackScanResult is what a newest-first scan of the transcript finds.

func BackScan

func BackScan(msgs []msgmodel.WithParts) BackScanResult

BackScan scans chronological filtered messages from newest to oldest.

type LLMClient

type LLMClient interface {
	Stream(ctx context.Context, params orclient.RequestParams) (PartStream, error)
}

LLMClient starts one request. Multi-turn behavior belongs to Loop.Run.

type Loop

type Loop struct {
	Store    Store
	Client   LLMClient
	Models   ModelResolver
	Executor ToolExecutor
	Tasks    TaskController

	// ProcessorWaitTimeout bounds the tool drain after an aborted stream.
	// Zero means 250ms.
	ProcessorWaitTimeout time.Duration
}

Loop wires the outer state machine's service seams.

func (*Loop) Run

func (l *Loop) Run(ctx context.Context, opts RunOptions) (msgmodel.Assistant, error)

Run drives the step loop until its natural exit or a processor stop.

type Model

type Model struct {
	Message msgmodel.Model
	Calc    calc.Model
	Request orclient.RequestParams
}

Model is the resolved model in the three projections the loop needs: the message-conversion view, the budget view and the request parameters.

type ModelResolver

type ModelResolver interface {
	Resolve(ctx context.Context, user msgmodel.User) (Model, error)
}

ModelResolver resolves the model named by the latest user message.

type ModelResolverFunc

type ModelResolverFunc func(context.Context, msgmodel.User) (Model, error)

ModelResolverFunc adapts a function to ModelResolver.

func (ModelResolverFunc) Resolve

func (f ModelResolverFunc) Resolve(ctx context.Context, user msgmodel.User) (Model, error)

type PartStream

type PartStream interface {
	Next() (orclient.StreamPart, error)
	Close() error
}

PartStream is one one-request OpenRouter response.

type Processor

type Processor struct {
	// contains filtered or unexported fields
}

Processor persists one assistant turn and owns its in-flight tool registry.

func NewProcessor

func NewProcessor(opts ProcessorOptions) *Processor

func (*Processor) CompleteToolCall

func (p *Processor) CompleteToolCall(ctx context.Context, toolCallID string, output ToolResult) error

CompleteToolCall persists a tool result and settles the call. A non-running or missing call is a no-op and is not settled.

func (*Processor) FailToolCall

func (p *Processor) FailToolCall(ctx context.Context, toolCallID string, failure error) (bool, error)

FailToolCall persists a tool failure and settles the call. A non-running or missing call is a no-op.

func (*Processor) Message

func (p *Processor) Message() msgmodel.Assistant

Message returns a copy of the current assistant message.

func (*Processor) Process

func (p *Processor) Process(ctx context.Context, stream PartStream) (result Result, err error)

Process drains one provider stream, synthesizing start-step/finish-step and tool-result/error events around orclient's lower-level parts.

func (*Processor) UpdateToolCall

func (p *Processor) UpdateToolCall(ctx context.Context, toolCallID string, update func(msgmodel.ToolPart) msgmodel.ToolPart) (*msgmodel.ToolPart, error)

UpdateToolCall rewrites and persists one in-flight tool part. The callback runs while the processor lock is held.

type ProcessorOptions

type ProcessorOptions struct {
	Store       Store
	Assistant   msgmodel.Assistant
	Model       Model
	Tools       []ToolDefinition
	Executor    ToolExecutor
	WaitTimeout time.Duration
}

ProcessorOptions configure the processor for one assistant turn.

type Result

type Result string

Result is what a processed turn asks the loop to do next.

const (
	ResultCompact  Result = "compact"
	ResultStop     Result = "stop"
	ResultContinue Result = "continue"
)

type RunOptions

type RunOptions struct {
	SessionID string
	ParentID  string
	Workspace string
	Worktree  string

	// MaxSteps is agent.steps. Nil means Infinity.
	MaxSteps *float64
	Tools    []ToolDefinition

	// InjectReminders may add in-memory-only reminder parts to the prompt
	// before each request. It must return fresh values.
	InjectReminders func(context.Context, []msgmodel.WithParts, msgmodel.User) ([]msgmodel.WithParts, error)

	// AfterAssistant runs after each processor turn, including stop and
	// error turns, with the assistant message ID; the instruction tracker
	// uses it to release the claims made for that message.
	AfterAssistant func(context.Context, string)

	// AfterTurn observes a fully persisted assistant turn and may stop the loop
	// before another provider call.
	AfterTurn func(context.Context, msgmodel.Assistant, msgmodel.Parts) error
}

RunOptions are the session and agent values one Run needs.

type SliceStream

type SliceStream struct {
	Parts   []orclient.StreamPart
	Failure error
	// contains filtered or unexported fields
}

SliceStream is a deterministic in-memory PartStream useful to embedders and tests. Failure is returned after all Parts; Close is idempotent.

func (*SliceStream) Close

func (s *SliceStream) Close() error

func (*SliceStream) Next

func (s *SliceStream) Next() (orclient.StreamPart, error)

type Store

type Store interface {
	Messages(ctx context.Context, sessionID string) ([]msgmodel.WithParts, error)
	UpdateMessage(ctx context.Context, info msgmodel.Info) error
	UpdatePart(ctx context.Context, part msgmodel.Part) error
}

Store is the MessageV2 persistence slice used by the loop and processor. Messages returns chronological order and fresh part values.

type TaskController

type TaskController interface {
	ProcessCompaction(ctx context.Context, input TaskInput, task msgmodel.CompactionPart) (Result, error)
	IsOverflow(ctx context.Context, assistant msgmodel.Assistant, model Model) (bool, error)
	CreateCompaction(ctx context.Context, sessionID string, user msgmodel.User, overflow bool) error
	Prune(ctx context.Context, sessionID string) error
}

TaskController is the narrow seam to the compaction service. The step loop owns when each operation runs; the controller owns the operations.

type TaskInput

type TaskInput struct {
	SessionID string
	Messages  []msgmodel.WithParts
	User      msgmodel.User
	Model     Model
}

TaskInput is the context handed to the compaction branch of the loop.

type ToolCall

type ToolCall struct {
	ID        string          `json:"id"`
	Name      string          `json:"name"`
	Input     json.RawMessage `json:"input"`
	SessionID string          `json:"sessionID"`
	MessageID string          `json:"messageID"`
	Agent     string          `json:"agent,omitempty"`
	ModelID   string          `json:"modelID,omitempty"`
}

ToolCall is the already-repaired, already-validated call handed to a real tool implementation.

type ToolDefinition

type ToolDefinition struct {
	Provider      orclient.Tool
	Validate      func(input json.RawMessage) error
	WaitForResult bool
}

ToolDefinition is a provider declaration plus the optional validation seam consumed by orclient.ParseToolCall.

type ToolExecutor

type ToolExecutor interface {
	Execute(ctx context.Context, call ToolCall) (ToolResult, error)
}

ToolExecutor is intentionally the minimal real-tool seam. Resolution and declarations are RunOptions.Tools; this interface only performs a call.

type ToolResult

type ToolResult struct {
	Title       string
	Metadata    msgmodel.RawObject
	Output      string
	Attachments *[]msgmodel.FilePart
}

ToolResult is what a tool returns when it completes.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL