Documentation
¶
Overview ¶
Package taskengine orchestrates an agent: it drives LLM turns, tool calls, and routing in a loop, defined as a JSON chain you version in git. The unit of execution is the conversation; the TaskEvent stream is the contract clients consume (see docs/development/engine-events.md). Data is shaped into a turn or produced by a tool call, never mutated invisibly where no event would see it.
Index ¶
- Constants
- Variables
- func ApprovalVerdictFromContext(ctx context.Context, approvalID string) (approved bool, ok bool)
- func ConvertToType(value interface{}, dataType DataType) (interface{}, error)
- func EdgeCountsFromContext(ctx context.Context) map[string]int
- func ExportedResolveToolsNames(ctx context.Context, allowlist []string, provider ToolsProvider) ([]string, error)
- func ExtractJSONArray(s string) string
- func ExtractJSONObject(s string) string
- func GetPrimaryModel(llmCall *LLMExecutionConfig) string
- func HasCheckpointSaver(ctx context.Context) bool
- func IsAssistantProseHandler(handler string) bool
- func IsToolBearingHandler(handler string) bool
- func LintChain(chain *TaskChainDefinition, entryTypes ...DataType) error
- func MarshalCheckpoint(cp *Checkpoint) ([]byte, error)
- func MergeTemplateVars(ctx context.Context, overlay map[string]string) context.Context
- func RequestedContextLengthFromContext(ctx context.Context) int
- func RuntimeToolsAllowlistFromContext(ctx context.Context) ([]string, bool)
- func StateSubject(reqID string) string
- func StripCodeFences(s string) string
- func SupportedOperators() []string
- func TaskEventRequestSubject(requestID string) string
- func TemplateVarsFromContext(ctx context.Context) (map[string]string, error)
- func ToolCallSuspendable(ctx context.Context) bool
- func ToolsArgsFromContext(ctx context.Context, toolsName string) map[string]string
- func ToolsToolsUnavailable(toolsName string, cause error) error
- func WithApprovalVerdicts(ctx context.Context, verdicts map[string]bool) context.Context
- func WithAttentionAnswers(ctx context.Context, answers map[string]AttentionAnswer) context.Context
- func WithCheckpointSaver(ctx context.Context, saver CheckpointSaver) context.Context
- func WithEdgeCounts(ctx context.Context, counts map[string]int) context.Context
- func WithRequestedContextLength(ctx context.Context, contextLength int) context.Context
- func WithResumeCheckpoint(ctx context.Context, cp *Checkpoint) context.Context
- func WithRetryOutcomeSink(ctx context.Context, sink *RetryOutcomeSink) context.Context
- func WithRuntimeToolsAllowlist(ctx context.Context, allowlist []string) context.Context
- func WithSuspendableToolCall(ctx context.Context) context.Context
- func WithTaskEventScope(ctx context.Context, scope TaskEventScope) context.Context
- func WithTaskEventSink(ctx context.Context, sink TaskEventSink) context.Context
- func WithTemplateVars(ctx context.Context, vars map[string]string) context.Context
- func WithToolsArgs(ctx context.Context, toolsName string, args map[string]string) context.Context
- type ApprovalPendingError
- type AttentionAnswer
- type BusInspector
- type BusTaskEventSink
- type CapturedPayloadSummary
- type CapturedStateUnit
- type ChainContext
- type ChainSuspendedError
- type ChainTerms
- type ChatHistory
- type Checkpoint
- type CheckpointSaver
- type DataType
- type EnvExecutor
- type ErrorResponse
- type EventScope
- type FunctionCall
- type FunctionCallObject
- type FunctionTool
- type HandlerOutputMode
- type HandlerSignature
- type ImagePart
- type Inspector
- type KVInspector
- type KVJournalTaskEventSink
- type LLMExecutionConfig
- type MacroEnv
- type Message
- type MockTaskExecutor
- type NoopTaskEventSink
- type OperatorTerm
- type PendingToolCall
- type RetryOutcomeSink
- type SimpleEnv
- type SimpleExec
- type SimpleStackTrace
- type StackTrace
- type TaskChainDefinition
- type TaskDefinition
- type TaskEvent
- type TaskEventKind
- type TaskEventScope
- type TaskEventSink
- type TaskExecutor
- type TaskHandler
- type TaskTransition
- type TokenUsage
- type Tool
- type ToolCall
- type ToolWithResolution
- type ToolsCall
- type ToolsProvider
- type ToolsRegistry
- type ToolsRepo
- type ToolsWithSchema
- type TransitionBranch
Constants ¶
const ( ContextKeyOutputByteLimit contextKey = "output_byte_limit" ContextKeyToolCallID contextKey = "tool_call_id" )
const ( // TransitionExecuted: a chat_completion turn finished with no tool calls. TransitionExecuted = "executed" // TransitionToolCall: a chat_completion turn requested one or more tool // calls. Snake_case to match the "tool_call" task-event kind. TransitionToolCall = "tool_call" // TransitionNoop: the noop handler ran, or execute_tool_calls saw empty history. TransitionNoop = "noop" // TransitionNoCallsFound: the model's last message carried no tool calls to run. TransitionNoCallsFound = "no_calls_found" // TransitionToolsExecuted: a tools task ran its tool successfully. TransitionToolsExecuted = "tools_executed" // TransitionFailed: a tools task failed. TransitionFailed = "failed" )
Transition-eval tokens are the control values a handler emits as its transition "eval"; a TransitionBranch matches them via its When field (with the default Operator, exact string equality). These are part of the DSL contract — branch on these constants, not the model's free text:
- chat_completion → TransitionToolCall (model requested tools) | TransitionExecuted (finished, no tool calls)
- execute_tool_calls → TransitionNoop (empty history) | TransitionNoCallsFound (model produced no tool calls) | TransitionToolsExecuted | TransitionFailed
- tools → TransitionToolsExecuted | TransitionFailed (or, when OutputTemplate is set, its rendered text)
- noop → TransitionNoop
To branch on the model's actual text, use the `route` handler, whose eval IS the model's chosen label.
const CheckpointSchemaVersion = 1
CheckpointSchemaVersion is the wire version MarshalCheckpoint writes. Bump it together with a migration entry in checkpointMigrations, never alone.
const TaskEventSubjectAll = "taskengine.events"
const (
TermEnd = "end"
)
Variables ¶
var ErrChainLint = errors.New("chain failed load-time validation")
ErrChainLint marks every defect the load-time chain linter reports, so services can distinguish "this chain is invalid" (disable it, teach the author) from I/O failures (retry, propagate).
var ErrCheckpointVersion = errors.New("taskengine: unsupported checkpoint schema version")
ErrCheckpointVersion reports a checkpoint whose schema version this binary cannot load (newer than it knows, or older with no migration registered).
var ErrContextLengthExceeded = errors.New("exceeds context length")
ErrContextLengthExceeded is returned when the input or chat history exceeds the allowed context length.
var ErrToolsNotFound = errors.New("tools not found")
ErrToolsNotFound is returned when a named tools is not registered in any repo.
ErrToolsToolsUnavailable is returned when a tools is registered but its tool list cannot be loaded (e.g. MCP server unreachable or list-tools failed). ExecEnv treats this like a missing tools for tool preload: skip tools, continue the chain.
var ErrUnsupportedTaskType = errors.New("executor does not support the task type")
ErrUnsupportedTaskType indicates unrecognized task type
Functions ¶
func ApprovalVerdictFromContext ¶
ApprovalVerdictFromContext reports the pre-loaded verdict for approvalID, ok=false when none was injected.
func ConvertToType ¶
ConvertToType converts a value to the specified DataType
func EdgeCountsFromContext ¶
EdgeCountsFromContext returns the edge counts attached via WithEdgeCounts, or nil if not set. A nil map is safe to read (lookup returns zero).
func ExportedResolveToolsNames ¶
func ExportedResolveToolsNames(ctx context.Context, allowlist []string, provider ToolsProvider) ([]string, error)
ExportedResolveToolsNames is a test-only export of resolveToolsNames.
func ExtractJSONArray ¶
ExtractJSONArray scans s for the outermost [...] block and returns it. It first strips code fences, then skips any preamble text the LLM may have placed before the JSON array to be robust to inconsistent model output.
func ExtractJSONObject ¶
ExtractJSONObject scans s for the outermost {...} block and returns it. It strips code fences first, same spirit as ExtractJSONArray.
func GetPrimaryModel ¶
func GetPrimaryModel(llmCall *LLMExecutionConfig) string
GetPrimaryModel returns llmCall's primary model name, falling back to "default" for token counting when none is configured.
func HasCheckpointSaver ¶
HasCheckpointSaver reports whether a durable checkpoint sink is installed on this run — the precondition for any park-and-release ask.
func IsAssistantProseHandler ¶
IsAssistantProseHandler reports whether a TaskEventStepChunk carrying this handler's streamed output is user-visible assistant narration — text and reasoning a chat surface should render into the reply — as opposed to control flow that merely happens to stream (e.g. a route task's streamed output is its routing label, not prose). A chunk this rejects is still journaled, not dropped; every event translator consumes this judgement.
func IsToolBearingHandler ¶
IsToolBearingHandler reports whether this handler already reports its own work through the dedicated tool-call events (TaskEventToolCallPending / TaskEventToolCall) or a chat surface's own tool rendering. Event translators use it to suppress the generic step-lifecycle card (TaskEventStepStarted / StepCompleted / StepFailed) for those handlers, so one action is not double-rendered as both a step card and a tool card.
func LintChain ¶
func LintChain(chain *TaskChainDefinition, entryTypes ...DataType) error
LintChain vets a chain at load time. entryTypes are the DataTypes the caller may feed the chain's entry task; when omitted the chain input is treated as DataTypeAny (unknown, runtime-checked). All defects are collected and joined so one vet run teaches everything at once; the result wraps ErrChainLint.
func MarshalCheckpoint ¶
func MarshalCheckpoint(cp *Checkpoint) ([]byte, error)
MarshalCheckpoint encodes cp as the current schema version. A var that cannot be marshalled fails the whole encode, rather than resuming a different run than the one that suspended.
func MergeTemplateVars ¶
MergeTemplateVars overlays keys onto any template vars already in ctx, then attaches the combined map. Use this when a nested step must add request_id / previous_output without dropping caller-supplied vars like model and provider.
func RequestedContextLengthFromContext ¶
RequestedContextLengthFromContext returns the positive context window attached by WithRequestedContextLength, or 0 when the caller did not request one.
func RuntimeToolsAllowlistFromContext ¶
RuntimeToolsAllowlistFromContext returns (allowlist, true) when an allowlist was attached via WithRuntimeToolsAllowlist. The returned slice follows the same grammar as TaskDefinition.Tools. Returns (nil, false) when no runtime allowlist is attached — callers should treat this as "no restriction".
func StateSubject ¶
func StripCodeFences ¶
func SupportedOperators ¶
func SupportedOperators() []string
func TaskEventRequestSubject ¶
func TemplateVarsFromContext ¶
TemplateVarsFromContext returns the template variables map from the context. Returns nil if not set; a nil map is safe to read (key lookup returns false). MacroEnv will return an error for any {{var:key}} whose key is absent.
func ToolCallSuspendable ¶
ToolCallSuspendable reports whether the current tool call may suspend the run — see WithSuspendableToolCall.
func ToolsArgsFromContext ¶
ToolsArgsFromContext returns the args previously stored for toolsName, or nil if none were set. The returned map must not be mutated by the caller.
func ToolsToolsUnavailable ¶
ToolsToolsUnavailable wraps cause as ErrToolsToolsUnavailable for toolsName (for errors.Is).
func WithApprovalVerdicts ¶
WithApprovalVerdicts pre-loads human verdicts keyed by approval ID for a resumed run: the HITL wrapper executes (true) or denies (false) a call whose verdict is already recorded, without gating it again. The map is copied.
func WithAttentionAnswers ¶
WithAttentionAnswers pre-loads operator answers keyed by ask ID for a resumed run, exactly as WithApprovalVerdicts pre-loads verdicts. The map is copied.
func WithCheckpointSaver ¶
func WithCheckpointSaver(ctx context.Context, saver CheckpointSaver) context.Context
WithCheckpointSaver installs the durable checkpoint sink for runs under ctx. Without one, a run that hits an approval park fails with a teaching error instead of suspending, never a silent loss.
func WithEdgeCounts ¶
WithEdgeCounts attaches the in-flight edge traversal counts for the current chain run to ctx. SimpleEnv updates this between each task step so handlers and step-time macros (e.g. {{edge_count:from->to}}) can read counts that reflect the loop iteration they are in, not the chain-start zero.
Layout of the map: keys are "<fromTaskID>-><toTaskID>", values are the number of times that edge has been traversed in this run.
func WithRequestedContextLength ¶
WithRequestedContextLength attaches a per-request context window. Task execution uses it as the resolver's minimum without replacing the chain's token_limit guardrail.
func WithResumeCheckpoint ¶
func WithResumeCheckpoint(ctx context.Context, cp *Checkpoint) context.Context
WithResumeCheckpoint marks the run under ctx as a resume of cp: ExecEnv restores vars/edgeCounts, re-enters at cp.TaskID, and feeds the checkpointed history verbatim into that task's first attempt.
func WithRetryOutcomeSink ¶
func WithRetryOutcomeSink(ctx context.Context, sink *RetryOutcomeSink) context.Context
WithRetryOutcomeSink attaches sink to ctx so chat_completion tasks can append outcomes via [appendRetryOutcome].
func WithRuntimeToolsAllowlist ¶
WithRuntimeToolsAllowlist attaches a caller-supplied tools allowlist to ctx that is intersected with each task's own allowlist inside resolveToolsNames. A caller can only further restrict — never expand — what a chain JSON permits. Grammar matches TaskDefinition.Tools: nil/[]/["*"]/exact names/["*","!name"].
Use this when a host must enforce per-call policy (such as disabling local_shell for a step) regardless of what the chain JSON declares. Absent key means "no runtime restriction" — behavior matches pre-feature code.
func WithSuspendableToolCall ¶
WithSuspendableToolCall marks the current tool call as one the engine can suspend on. Set only at the model-batch execution site, whose task output is the ChatHistory suspendRun requires; a tools-handler call never carries it, so askers gate park-and-release behavior on this marker.
func WithTaskEventScope ¶
func WithTaskEventScope(ctx context.Context, scope TaskEventScope) context.Context
func WithTaskEventSink ¶
func WithTaskEventSink(ctx context.Context, sink TaskEventSink) context.Context
func WithTemplateVars ¶
WithTemplateVars attaches a map of template variables to the context. MacroEnv expands {{var:name}} from this map. The engine never reads os.Getenv; callers (e.g. Contenox CLI, API) build the map and attach it here.
func WithToolsArgs ¶
WithToolsArgs stores a copy of args for the named tools in ctx.
The map is copied on entry so the stored value is immutable — callers must not rely on mutating the original map being visible to tools implementations. This ensures no data races when the same context is read concurrently (e.g. during tool-list construction in ExecEnv).
Types ¶
type ApprovalPendingError ¶
type ApprovalPendingError struct {
// ApprovalID is the durable approval row's ID and the checkpoint key.
ApprovalID string
ToolName string
}
ApprovalPendingError is the HITL wrapper's third outcome beside allow/deny: the fast-path park elapsed with no human verdict. The durable approval row exists before this is returned; the executor persists a checkpoint and releases the run after.
func (*ApprovalPendingError) Error ¶
func (e *ApprovalPendingError) Error() string
type AttentionAnswer ¶
AttentionAnswer is the text twin of an approval verdict: the operator's words, or Answered=false meaning the asking tool runs its blocker fallback.
func AttentionAnswerFromContext ¶
func AttentionAnswerFromContext(ctx context.Context, askID string) (ans AttentionAnswer, ok bool)
AttentionAnswerFromContext reports the pre-loaded answer for askID, ok=false when none was injected.
type BusInspector ¶
type BusInspector struct {
// contains filtered or unexported fields
}
func NewBusInspector ¶
func NewBusInspector(inner Inspector, bus libbus.Messenger, tracker libtracker.ActivityTracker) *BusInspector
func (*BusInspector) Start ¶
func (i *BusInspector) Start(ctx context.Context) StackTrace
type BusTaskEventSink ¶
type BusTaskEventSink struct {
// contains filtered or unexported fields
}
func NewBusTaskEventSink ¶
func NewBusTaskEventSink(bus libbus.Messenger) *BusTaskEventSink
func (*BusTaskEventSink) PublishTaskEvent ¶
func (s *BusTaskEventSink) PublishTaskEvent(ctx context.Context, event TaskEvent) error
func (*BusTaskEventSink) Wants ¶
func (s *BusTaskEventSink) Wants(TaskEventKind) bool
Wants accepts every kind while a bus is attached (the pre-Wants Enabled() behavior: bus sinks consumed all kinds).
type CapturedPayloadSummary ¶
type CapturedPayloadSummary struct {
Truncated bool `json:"truncated"`
Reason string `json:"reason"`
OriginalType string `json:"originalType,omitempty"`
OriginalJSONBytes int `json:"originalJsonBytes,omitempty"`
SHA256 string `json:"sha256,omitempty"`
Preview string `json:"preview,omitempty"`
PreviewBytes int `json:"previewBytes,omitempty"`
}
CapturedPayloadSummary replaces oversized or non-JSON-marshallable captured payloads in persisted/streamed state. In-memory execution history keeps the original values; this type is for observability storage only.
type CapturedStateUnit ¶
type CapturedStateUnit struct {
// Scope is the unit's hierarchical address (chain + task); the same
// address contract TaskEvent carries, so checkpoints and replay name the
// position of captured state without re-deriving it from loose IDs.
// Additive JSON field; TaskID below stays populated identically.
Scope EventScope `json:"scope,omitzero"`
TaskID string `json:"taskID" example:"validate_input"`
TaskHandler string `json:"taskHandler" example:"chat_completion"`
InputType DataType `json:"inputType" example:"string" openapi_include_type:"string"`
OutputType DataType `json:"outputType" example:"string" openapi_include_type:"string"`
Transition string `json:"transition" example:"valid_input"`
Duration time.Duration `json:"duration" example:"452000000"`
Error ErrorResponse `json:"error" openapi_include_type:"taskengine.ErrorResponse"`
Input any `json:"input,omitempty"`
Output any `json:"output,omitempty"`
InputVar string `json:"inputVar" example:"input"`
RetryIndex int `json:"retryIndex"`
Cancelled bool `json:"cancelled,omitempty"`
TimedOut bool `json:"timedOut,omitempty"`
ProviderType string `json:"providerType,omitempty"`
ModelName string `json:"modelName,omitempty"`
ToolNames []string `json:"toolNames,omitempty"`
TokenUsage *TokenUsage `json:"tokenUsage,omitempty"`
// FinishReason is the provider's verbatim finish reason for the model call
// this step captured, when its output was a chat history ("length"-class
// values mean truncation). Consumers normalize; the engine records.
FinishReason string `json:"finishReason,omitempty"`
}
type ChainContext ¶
type ChainContext struct {
Tools map[string]ToolWithResolution
ClientTools []Tool
Debug bool
}
type ChainSuspendedError ¶
type ChainSuspendedError struct {
ApprovalID string
// Scope is the hierarchical address of the interrupt point.
Scope EventScope
}
ChainSuspendedError is ExecEnv's terminal for a suspended run: the checkpoint is persisted and answering the approval resumes the chain. It is a typed outcome, not a failure.
func (*ChainSuspendedError) Error ¶
func (e *ChainSuspendedError) Error() string
type ChainTerms ¶
type ChainTerms string
type ChatHistory ¶
type ChatHistory struct {
// Messages is the list of messages in the conversation.
Messages []Message `json:"messages"`
// Model is the name of the model to use for the conversation.
Model string `json:"model" example:"llama2:7b"`
// InputTokens will be filled by the engine and will hold the number of tokens used for the input.
InputTokens int `json:"inputTokens" example:"15"`
// OutputTokens will be filled by the engine and will hold the number of tokens used for the output.
OutputTokens int `json:"outputTokens" example:"10"`
// FinishReason is the provider's verbatim finish reason for the LAST model
// call that produced this history ("length"-class values mean the response
// was truncated). Filled by the engine; empty when no model call ran or
// the provider reported none.
FinishReason string `json:"finishReason,omitempty" example:"stop"`
}
ChatHistory represents a conversation history with an LLM.
type Checkpoint ¶
type Checkpoint struct {
ApprovalID string
PendingCalls []PendingToolCall
// Chain is the full definition as executed: chains are data, so a
// reference could dangle across a restart, but the definition cannot.
Chain *TaskChainDefinition
TaskID string
RetryIndex int
Scope EventScope
// Vars/VarTypes round-trip through the closed DataType enum (see
// decodeCheckpointVar) — no reflection.
Vars map[string]any
VarTypes map[string]DataType
EdgeCounts map[string]int
// History ends with the unanswered tool calls (plus any results
// produced before the gate) — the shape the tool-pairing repair path
// re-enters on.
History ChatHistory
TemplateVars map[string]string
ToolsAllowlist []string
HasToolsAllowlist bool
ContextLength int
SessionID string
MissionID string
RequestID string
ChainRef string
CreatedAt time.Time
}
Checkpoint is everything a suspended run needs to resume in any process: position, state (vars, edge counts, chat history), pending tool calls, and the request-scoped identity the service layer re-injects.
A checkpoint is keyed by the one call awaiting a verdict; PendingCalls also lists the batch's not-yet-started calls. Resume injects that one verdict and re-enters execute_tool_calls, which runs the remaining calls through the normal path — each may gate again and suspend afresh under its own key.
func UnmarshalCheckpoint ¶
func UnmarshalCheckpoint(raw []byte) (*Checkpoint, error)
UnmarshalCheckpoint decodes raw, migrating older schema versions forward. A version this binary cannot reach errors with ErrCheckpointVersion rather than risking a mis-decoded, silently corrupted resume.
type CheckpointSaver ¶
type CheckpointSaver interface {
SaveCheckpoint(ctx context.Context, cp *Checkpoint) error
}
CheckpointSaver persists a suspension checkpoint. The service layer installs an implementation on the run context; Save must be atomic per call — either the checkpoint is readable under its approval ID afterward, or the run fails instead of suspending into a lost run.
type DataType ¶
type DataType int
DataType represents the type of data passed between tasks.
func DataTypeFromString ¶
DataTypeFromString converts a string to DataType.
func InferDataType ¶
InferDataType picks the narrowest concrete DataType for a runtime value.
func NormalizeDataType ¶
NormalizeDataType upgrades DataTypeAny to a concrete type and coerces the value with ConvertToType.
func (DataType) MarshalJSON ¶
func (DataType) MarshalYAML ¶
func (*DataType) UnmarshalJSON ¶
type EnvExecutor ¶
type EnvExecutor interface {
ExecEnv(ctx context.Context, chain *TaskChainDefinition, input any, dataType DataType) (any, DataType, []CapturedStateUnit, error)
}
EnvExecutor executes complete task chains with input and environment management.
func NewEnv ¶
func NewEnv( ctx context.Context, tracker libtracker.ActivityTracker, exec TaskExecutor, inspector Inspector, toolsProvider ToolsRepo, ) (EnvExecutor, error)
NewEnv creates a new SimpleEnv with the given tracker and task executor.
func NewMacroEnv ¶
func NewMacroEnv(inner EnvExecutor, toolsProvider ToolsRepo) (EnvExecutor, error)
NewMacroEnv wraps an existing EnvExecutor with macro expansion.
type ErrorResponse ¶
type EventScope ¶
type EventScope struct {
Chain string `json:"chain,omitempty"`
Task string `json:"task,omitempty"`
ToolCall string `json:"tool_call,omitempty"`
}
EventScope is the hierarchical address of an event or captured state unit: which chain / which task / which tool call produced it. It is the address contract checkpoints and nested consumers name positions with — structured, never re-parsed out of loose strings. Fields are filled top-down: Chain is set on every event of a run, Task on every event emitted inside a task attempt, ToolCall only on events attributable to one tool invocation. The legacy flat TaskEvent fields (ChainID, TaskID) remain populated for wire compatibility; Scope is the additive, authoritative form.
type FunctionCall ¶
type FunctionCall struct {
Name string `json:"name" example:"get_current_weather"`
Arguments string `json:"arguments" example:"{\n \"location\": \"San Francisco, CA\",\n \"unit\": \"celsius\"\n}"`
}
FunctionCall specifies the function name and arguments for a tool call.
type FunctionCallObject ¶
type FunctionTool ¶
type FunctionTool struct {
Name string `json:"name"`
Description string `json:"description,omitempty"`
Parameters interface{} `json:"parameters,omitempty"` // JSON Schema object
}
FunctionTool defines the schema for a function-type tool.
type HandlerOutputMode ¶
type HandlerOutputMode int
HandlerOutputMode names how a handler's success-output type derives from its input type.
const ( // HandlerOutputFixed: the handler always produces HandlerSignature.Output // on success, regardless of input type. HandlerOutputFixed HandlerOutputMode = iota // HandlerOutputPassthrough: the handler returns its input unchanged — // output type equals input type (noop, and route, whose product is the // transition label, not the data). HandlerOutputPassthrough // HandlerOutputDynamic: the output type is decided at runtime by the tool // that executes (the tools handler). The linter must treat it as // DataTypeAny — unless the task's OutputTemplate forces a rendered string. HandlerOutputDynamic // HandlerOutputNone: the handler never succeeds (raise_error). Only its // on_failure edge can ever be taken; success branches are dead. HandlerOutputNone )
type HandlerSignature ¶
type HandlerSignature struct {
// Inputs is the closed set of DataTypes the handler accepts, in the order
// teaching errors should name them. Empty means the handler accepts every
// DataType.
Inputs []DataType
// Mode says how the success-output type derives from the input type.
Mode HandlerOutputMode
// Output is the produced type when Mode == HandlerOutputFixed.
Output DataType
// SuccessEvals is the closed transition-eval vocabulary the handler can
// emit on success — the only values a TransitionBranch can ever match for
// this handler. Nil means the vocabulary is open (route emits the model's
// chosen label; a tools task's OutputTemplate replaces the eval with its
// rendered text) and the linter cannot prove a branch dead.
SuccessEvals []string
}
HandlerSignature is the frozen I/O contract of one task handler.
func HandlerSignatureFor ¶
func HandlerSignatureFor(h TaskHandler) (HandlerSignature, bool)
HandlerSignatureFor returns the frozen contract for h. The bool is false for a handler the table does not know — which validateChain already rejects.
func (HandlerSignature) AcceptsInput ¶
func (s HandlerSignature) AcceptsInput(dt DataType) bool
AcceptsInput reports whether the handler's closed input set admits dt. DataTypeAny is never "admitted" here — an Any-typed value is unknown at load time and stays the runtime backstop's problem; callers must special-case it.
type ImagePart ¶
type ImagePart struct {
// Data is the raw image bytes. JSON encoding carries it as standard base64.
Data []byte `json:"data"`
// MimeType is the image media type.
MimeType string `json:"mime_type" example:"image/png"`
}
ImagePart is a binary image attachment on a Message.
type Inspector ¶
type Inspector interface {
Start(ctx context.Context) StackTrace
}
func NewSimpleInspector ¶
func NewSimpleInspector() Inspector
type KVInspector ¶
type KVInspector struct {
// contains filtered or unexported fields
}
func NewKVInspector ¶
func NewKVInspector(inner Inspector, kv libkv.KVManager, tracker libtracker.ActivityTracker) *KVInspector
func (*KVInspector) GetExecutionStateByRequestID ¶
func (i *KVInspector) GetExecutionStateByRequestID(ctx context.Context, reqID string) ([]CapturedStateUnit, error)
func (*KVInspector) GetStatefulRequests ¶
func (i *KVInspector) GetStatefulRequests(ctx context.Context) ([]string, error)
func (*KVInspector) Start ¶
func (i *KVInspector) Start(ctx context.Context) StackTrace
type KVJournalTaskEventSink ¶
type KVJournalTaskEventSink struct {
// contains filtered or unexported fields
}
KVJournalTaskEventSink journals task events durably per request ID into the KV store (alongside the CapturedStateUnit stream the KVInspector persists), after forwarding them to the wrapped sink. This is what makes a run's work log — tool calls, diffs, approvals — re-renderable after the 5-minute bus TTL and across restarts.
step_chunk events are deliberately not journaled: they are streaming detail whose final text is already persisted with the chat history. step_stream_end is journaled — the durable stream bracket (chunk count, finish reason, usage) — so a replayed run can tell streaming happened and how it ended even though the chunks themselves are gone. Every other kind is journaled.
func NewKVJournalTaskEventSink ¶
func NewKVJournalTaskEventSink(inner TaskEventSink, kv libkv.KVManager, tracker libtracker.ActivityTracker) *KVJournalTaskEventSink
func (*KVJournalTaskEventSink) PublishTaskEvent ¶
func (s *KVJournalTaskEventSink) PublishTaskEvent(ctx context.Context, event TaskEvent) error
func (*KVJournalTaskEventSink) Wants ¶
func (s *KVJournalTaskEventSink) Wants(kind TaskEventKind) bool
Wants defers to the wrapped sink (matching the pre-Wants Enabled() delegation); with no inner sink the journal itself consumes every kind.
type LLMExecutionConfig ¶
type LLMExecutionConfig struct {
// Model is the primary model: it is placed first in the candidate list and is
// the model used for token counting (see GetPrimaryModel). When both Model and
// Models are set, Model plus Models form the candidate set (Model first);
// the resolver then picks a reachable one — so set exactly Model for a single
// pinned model, or use Models for an explicit candidate pool.
Model string `yaml:"model" json:"model" example:"llama2:7b"`
// Models is an additional candidate pool, considered alongside Model.
Models []string `yaml:"models,omitempty" json:"models,omitempty" example:"[\"gpt-4\", \"gpt-3.5-turbo\"]"`
// Provider is the primary provider, placed first in the candidate list;
// Providers supplies additional candidates.
Provider string `yaml:"provider,omitempty" json:"provider,omitempty" example:"ollama"`
Providers []string `yaml:"providers,omitempty" json:"providers,omitempty" example:"[\"ollama\", \"openai\"]"`
// Temperature is the sampling temperature; pointer so "unset" (nil) is
// distinguishable from an explicit 0.0. When set it is honored everywhere.
// When unset: chat_completion uses the provider default; the prompt/route
// handlers use 0.0 (route depends on deterministic single-label output).
Temperature *float32 `yaml:"temperature,omitempty" json:"temperature,omitempty" example:"0.7"`
// Tools is the allowlist of registry tool names this task may invoke
// (client-passed tools are governed separately by PassClientsTools):
// absent/null/[] exposes none; ["*"] exposes all; ["a","b"] exposes only
// those names (unknown ones ignored); ["*","!name"] excludes name(s),
// meaningful only combined with "*".
Tools []string `yaml:"tools,omitempty" json:"tools,omitempty" example:"[\"local_shell\", \"nws\"]"`
// HideTools suppresses specific tools by (namespaced) name from BOTH the
// registry tools selected via Tools and the client-passed tools.
HideTools []string `yaml:"hide_tools,omitempty" json:"hide_tools,omitempty" example:"[\"tool1\", \"tools_name1.tool1\"]"`
// ToolsPolicies carries per-tools policy overrides for this task (tools
// name -> policy key -> value), injected into the context before
// GetToolsForToolsByName so tools can render dynamic descriptions and
// enforce the policy at Exec time.
ToolsPolicies map[string]map[string]string `yaml:"tools_policies,omitempty" json:"tools_policies,omitempty"`
PassClientsTools bool `yaml:"pass_clients_tools" json:"pass_clients_tools"`
// Think controls reasoning mode: auto, off, minimal, low, medium, high,
// xhigh, or a boolean-style alias. Empty uses the provider default.
Think string `yaml:"think,omitempty" json:"think,omitempty" example:"high"`
// MaxTokens caps the model's output tokens. Unset sends no explicit cap
// and never falls back to the chain's TokenLimit, which bounds the
// input+output window, not the output alone (conflating the two trips
// per-model output limits, e.g. Vertex Gemini 2.5 Pro's 65536 cap).
MaxTokens *int `yaml:"max_tokens,omitempty" json:"max_tokens,omitempty" example:"8192"`
// MaxTokensTemplate stores a string max_tokens macro from chain JSON until
// MacroEnv expands it into MaxTokens. It is not emitted as a separate field.
MaxTokensTemplate string `yaml:"-" json:"-"`
// Shift allows the context window to slide on overflow instead of erroring.
Shift bool `yaml:"shift,omitempty" json:"shift,omitempty"`
// RetryPolicy wraps the underlying chat/prompt call with classified retry
// (rate-limit / server-error / timeout) and an optional model fallback.
// Nil or zero-value disables retry — current default. See [llmretry.Do].
RetryPolicy *llmretry.RetryPolicy `yaml:"retry_policy,omitempty" json:"retry_policy,omitempty"`
}
LLMExecutionConfig represents configuration for executing tasks using Large Language Models (LLMs).
func (LLMExecutionConfig) MarshalJSON ¶
func (c LLMExecutionConfig) MarshalJSON() ([]byte, error)
func (*LLMExecutionConfig) UnmarshalJSON ¶
func (c *LLMExecutionConfig) UnmarshalJSON(data []byte) error
func (*LLMExecutionConfig) UnmarshalYAML ¶
func (c *LLMExecutionConfig) UnmarshalYAML(value *yaml.Node) error
type MacroEnv ¶
type MacroEnv struct {
// contains filtered or unexported fields
}
MacroEnv is a transparent decorator around EnvExecutor that expands special macros in task templates before execution. Supported macros:
- {{toolservice:list}} -> JSON map of tools name -> tool names
- {{toolservice:tools}} -> JSON array of tools names
- {{toolservice:tools <tools_name>}} -> JSON array of tool names for that tools
- {{var:<name>}} -> value from context template vars (set by caller via WithTemplateVars; engine never reads env); errors if key is missing
- {{var:<name>|<fallback>}} -> value from context template vars, or fallback when the var is missing or empty
- {{var:<name>|var:<fallback-name>}} -> value from context template vars, or another template var when the first is missing or empty
- {{date}} or {{date:<layout>}} -> current local date (default 2006-01-02)
- {{now}} or {{now:<layout>}} -> current time (default RFC3339; layout e.g. 2006-01-02)
- {{chain:id}} -> chain ID of the chain being executed
The engine does not expand any env:VAR-style macro; var:* is populated only by the caller.
type Message ¶
type Message struct {
// ID is not used by the engine; it is useful for tracking messages and
// computing differences of histories before storage.
ID string `json:"id" example:"msg_123456"`
Role string `json:"role" example:"user"`
Content string `json:"content,omitempty" example:"What is the capital of France?"`
// Images travel with Content to the model; requests with images
// resolve only to vision-capable models and fail with a typed error
// when none is available.
Images []ImagePart `json:"images,omitempty" openapi_include_type:"taskengine.ImagePart"`
// Thinking is the model's internal reasoning trace, populated only when
// thinking is enabled; never sent back to the model as history.
Thinking string `json:"thinking,omitempty"`
ToolCallID string `json:"tool_call_id,omitempty"`
CallTools []ToolCall `json:"callTools,omitempty"`
Timestamp time.Time `json:"timestamp" example:"2023-11-15T14:30:45Z"`
// RequestID and ChainRef are turn provenance (the run's X-Request-ID and
// the chain path that ran it), not used by the engine.
RequestID string `json:"requestId,omitempty"`
ChainRef string `json:"chainRef,omitempty"`
}
Message represents a single message in a chat conversation.
func SynthesizeHistory ¶
func SynthesizeHistory(prior []Message, units []CapturedStateUnit, chainErr error) []Message
SynthesizeHistory rebuilds a conversation transcript from a chain run's captured step stream (units, from Inspector.GetExecutionHistory), so hard-failed turns (errors, timeouts, cancellations, denied HITL gates) make it into the persisted ChatHistory — unlike the chain's returned ChatHistory, which only contains messages from steps that completed successfully. prior is the session history sent into the chain; chainErr is the chain runner's error, if any.
Messages are deduped by identity (Message.ID, or a content hash for pre-ID messages), never by index, since handlers legitimately mutate the message list between a unit's input and output. Engine-injected system messages are excluded, since task system instructions are re-applied from the task definition on every run. The result satisfies the tool-call pairing invariant (see repairToolCallPairing) and is a candidate []Message ready for chatservice.PersistDiff, which dedupes by ID.
type MockTaskExecutor ¶
type MockTaskExecutor struct {
// Single value responses
MockOutput any
MockTransitionValue string
MockError error
// Sequence responses
MockOutputSequence []any
MockTaskTypeSequence []DataType
MockTransitionValueSequence []string
ErrorSequence []error
// Tracking
CalledWithTask *TaskDefinition
CalledWithInput any
CalledWithPrompt string
// contains filtered or unexported fields
}
MockTaskExecutor is a mock implementation of taskengine.TaskExecutor.
func (*MockTaskExecutor) CallCount ¶
func (m *MockTaskExecutor) CallCount() int
CallCount returns how many times TaskExec was called
func (*MockTaskExecutor) Reset ¶
func (m *MockTaskExecutor) Reset()
Reset clears all mock state between tests
func (*MockTaskExecutor) TaskExec ¶
func (m *MockTaskExecutor) TaskExec(ctx context.Context, startingTime time.Time, tokenLimit int, chainContext *ChainContext, currentTask *TaskDefinition, input any, dataType DataType) (any, DataType, string, error)
TaskExec is the mock implementation of the TaskExec method.
type NoopTaskEventSink ¶
type NoopTaskEventSink struct{}
func (NoopTaskEventSink) PublishTaskEvent ¶
func (NoopTaskEventSink) PublishTaskEvent(context.Context, TaskEvent) error
func (NoopTaskEventSink) Wants ¶
func (NoopTaskEventSink) Wants(TaskEventKind) bool
type OperatorTerm ¶
type OperatorTerm string
OperatorTerm represents logical operators used for task transition evaluation
const ( OpEquals OperatorTerm = "equals" OpContains OperatorTerm = "contains" OpStartsWith OperatorTerm = "starts_with" OpEndsWith OperatorTerm = "ends_with" OpDefault OperatorTerm = "default" // OpEdgeTraversedAtLeast fires when the edge specified by TransitionBranch.Edge // (formatted "fromTaskID->toTaskID") has been traversed at least the integer // in TransitionBranch.When times during the current chain run. Reads engine // state, not task output. Use it to bound workflow loops: // // { "operator": "edge_traversed_at_least", // "edge": "chat->run_tools", "when": "20", "goto": "summarise_failure" } // // Place this branch ahead of the normal loop branch so it intercepts before // the next loop iteration fires. OpEdgeTraversedAtLeast OperatorTerm = "edge_traversed_at_least" )
func ToOperatorTerm ¶
func ToOperatorTerm(s string) (OperatorTerm, error)
func (OperatorTerm) String ¶
func (t OperatorTerm) String() string
type PendingToolCall ¶
type PendingToolCall struct {
CallID string `json:"call_id"`
Name string `json:"name"`
Arguments string `json:"arguments"`
}
PendingToolCall records one model-requested tool call the suspended run has not answered yet: the awaiting-verdict call first, then any calls of the batch that never started.
type RetryOutcomeSink ¶
type RetryOutcomeSink struct {
// contains filtered or unexported fields
}
RetryOutcomeSink collects per-call retry outcomes from chat_completion tasks running inside one chain invocation. It is safe for concurrent appenders.
A host attaches a sink via WithRetryOutcomeSink before running a chain so it can observe whether any chat call retried, used fallback, or hit a non-retryable class (e.g. capacity).
func (*RetryOutcomeSink) Append ¶
func (s *RetryOutcomeSink) Append(o llmretry.Outcome)
Append records one outcome. Safe for concurrent use.
func (*RetryOutcomeSink) LastErrorClass ¶
func (s *RetryOutcomeSink) LastErrorClass() llmretry.ErrorClass
LastErrorClass returns the class of the most recent recorded outcome, or llmretry.ClassNone if no outcomes were recorded.
func (*RetryOutcomeSink) Outcomes ¶
func (s *RetryOutcomeSink) Outcomes() []llmretry.Outcome
Outcomes returns a snapshot of recorded outcomes in append order.
type SimpleEnv ¶
type SimpleEnv struct {
// contains filtered or unexported fields
}
SimpleEnv is the default implementation of EnvExecutor.
type SimpleExec ¶
type SimpleExec struct {
// contains filtered or unexported fields
}
SimpleExec is a basic implementation of TaskExecutor. It executes chat completion, tools, route, raise_error, and noop tasks.
func (*SimpleExec) Prompt ¶
func (exe *SimpleExec) Prompt(ctx context.Context, systemInstruction string, llmCall LLMExecutionConfig, prompt string, ctxLength int) (string, error)
Prompt resolves a model client using the resolver policy and sends the prompt to be executed. Returns the trimmed response string or an error.
func (*SimpleExec) TaskExec ¶
func (exe *SimpleExec) TaskExec(taskCtx context.Context, startingTime time.Time, ctxLength int, chainContext *ChainContext, currentTask *TaskDefinition, input any, dataType DataType) (any, DataType, string, error)
TaskExec dispatches execution to currentTask.Handler (chat_completion, execute_tool_calls, tools, route, raise_error, noop) and returns its output plus a transition-eval string (see TaskTransition).
type SimpleStackTrace ¶
type SimpleStackTrace struct {
// contains filtered or unexported fields
}
func (*SimpleStackTrace) GetExecutionHistory ¶
func (s *SimpleStackTrace) GetExecutionHistory() []CapturedStateUnit
func (*SimpleStackTrace) RecordStep ¶
func (s *SimpleStackTrace) RecordStep(step CapturedStateUnit)
type StackTrace ¶
type StackTrace interface {
RecordStep(step CapturedStateUnit)
GetExecutionHistory() []CapturedStateUnit
}
type TaskChainDefinition ¶
type TaskChainDefinition struct {
ID string `yaml:"id" json:"id"`
// Debug enables capturing user input and output.
Debug bool `yaml:"debug" json:"debug"`
Description string `yaml:"description" json:"description"`
Tasks []TaskDefinition `yaml:"tasks" json:"tasks" openapi_include_type:"taskengine.TaskDefinition"`
// TokenLimit is the token limit for the context window used during execution.
TokenLimit int64 `yaml:"token_limit" json:"token_limit"`
}
TaskChainDefinition describes a sequence of tasks to execute in order, with branching logic, retry policies, and model preferences.
type TaskDefinition ¶
type TaskDefinition struct {
ID string `yaml:"id" json:"id" example:"validate_input"`
Description string `yaml:"description" json:"description" example:"Validates user input meets quality requirements"`
// Handler determines how the LLM output (or tools) will be interpreted.
Handler TaskHandler `yaml:"handler" json:"handler" example:"chat_completion" openapi_include_type:"string"`
SystemInstruction string `` /* 158-byte string literal not displayed */
ExecuteConfig *LLMExecutionConfig `yaml:"execute_config,omitempty" json:"execute_config,omitempty" openapi_include_type:"taskengine.LLMExecutionConfig"`
// Tools defines an external action to run: required for Tools tasks,
// must be nil/omitted for all other types.
Tools *ToolsCall `yaml:"tools,omitempty" json:"tools,omitempty" openapi_include_type:"taskengine.ToolsCall"`
// Print optionally formats the output for display/logging, supporting
// template variables from previous task outputs.
Print string `yaml:"print,omitempty" json:"print,omitempty" example:"Validation result: {{.validate_input}}"`
// PromptTemplate, when set, overrides the resolved input as the prompt
// sent to the LLM. Supports template variables from previous task outputs.
PromptTemplate string `yaml:"prompt_template,omitempty" json:"prompt_template,omitempty" example:"Is this input valid? {{.input}}"`
// OutputTemplate, when set, renders a tools task's JSON output through
// this go template; the rendered string becomes the task's output.
OutputTemplate string `yaml:"output_template,omitempty" json:"output_template,omitempty" example:"Tools result: {{.status}}"`
// InputVar names the variable to use as this task's input; each task
// stores its own output in a variable named after its task id.
InputVar string `yaml:"input_var,omitempty" json:"input_var,omitempty" example:"input"`
// InputMaxBytes caps oversized string/chat-history inputs before this task
// runs, for recovery/summarization tasks that should explain a failure
// without re-feeding the huge input that caused it.
InputMaxBytes int `yaml:"input_max_bytes,omitempty" json:"input_max_bytes,omitempty" example:"8192"`
// Transition defines what to do after this task completes.
Transition TaskTransition `yaml:"transition" json:"transition" openapi_include_type:"taskengine.TaskTransition"`
// Timeout is a Go duration string ("10s", "2m") bounding task execution.
Timeout string `yaml:"timeout,omitempty" json:"timeout,omitempty" example:"30s"`
// RetryOnFailure sets how many times to retry this task on failure (all
// task types); 0 means no retries.
RetryOnFailure int `yaml:"retry_on_failure,omitempty" json:"retry_on_failure,omitempty" example:"2"`
}
type TaskEvent ¶
type TaskEvent struct {
Kind TaskEventKind `json:"kind"`
Timestamp time.Time `json:"timestamp"`
RequestID string `json:"request_id,omitempty"`
// Scope is the event's hierarchical address (chain/task/tool-call). It is
// additive on the wire: consumers that predate it keep reading the flat
// ChainID/TaskID fields below, which stay populated identically.
Scope EventScope `json:"scope,omitzero"`
ChainID string `json:"chain_id,omitempty"`
TaskID string `json:"task_id,omitempty"`
TaskHandler string `json:"task_handler,omitempty"`
Retry int `json:"retry"`
ModelName string `json:"model_name,omitempty"`
ProviderType string `json:"provider_type,omitempty"`
BackendID string `json:"backend_id,omitempty"`
OutputType string `json:"output_type,omitempty"`
Transition string `json:"transition,omitempty"`
Content string `json:"content,omitempty"`
Thinking string `json:"thinking,omitempty"`
Error string `json:"error,omitempty"`
ApprovalID string `json:"approval_id,omitempty"`
HookName string `json:"hook_name,omitempty"`
ToolName string `json:"tool_name,omitempty"`
ApprovalArgs map[string]any `json:"approval_args,omitempty"`
ApprovalDiff string `json:"approval_diff,omitempty"`
HITLAction string `json:"hitl_action,omitempty"`
HITLReason string `json:"hitl_reason,omitempty"`
HITLPolicyName string `json:"hitl_policy_name,omitempty"`
HITLPolicyPath string `json:"hitl_policy_path,omitempty"`
HITLArgsSummary string `json:"hitl_args_summary,omitempty"`
HITLMatchedRule *int `json:"hitl_matched_rule,omitempty"`
HITLTimeoutS int `json:"hitl_timeout_s,omitempty"`
HITLApprovalRequested *bool `json:"hitl_approval_requested,omitempty"`
ToolDiffPath string `json:"tool_diff_path,omitempty"`
ToolDiffOldText string `json:"tool_diff_old_text,omitempty"`
ToolDiffNewText string `json:"tool_diff_new_text,omitempty"`
// For token_usage
TokenUsed int `json:"token_used,omitempty"`
TokenSize int `json:"token_size,omitempty"`
// For step_stream_end: the terminal bracket of one model stream.
// ChunkCount is the number of parcels that carried visible content or
// thinking (i.e. the number of step_chunk events a subscribed sink saw —
// counted independently of Wants, so the bracket is truthful even when no
// sink consumed chunks). FinishReason is the provider's verbatim finish
// reason ("" on the non-streaming fallback, which reports none). Usage is
// provider-reported token usage when available.
ChunkCount int `json:"chunk_count,omitempty"`
FinishReason string `json:"finish_reason,omitempty"`
Usage *TokenUsage `json:"usage,omitempty"`
}
func GetJournaledEvents ¶
GetJournaledEvents returns the durably journaled events of a run, in arrival order. A request with no journal yields an empty slice.
func NewTaskEvent ¶
func NewTaskEvent(ctx context.Context, kind TaskEventKind) TaskEvent
NewTaskEvent builds an event of the given kind addressed from ctx: the request ID, the chain/task scope installed by the executor (ExecEnv wraps the whole run in a chain-level scope; each task attempt overrides it with the full task scope), and — when the emission site runs inside one tool invocation (ContextKeyToolCallID) — the tool-call address. Emission sites that know a more precise tool-call ID than the context carries set Scope.ToolCall explicitly after construction.
type TaskEventKind ¶
type TaskEventKind string
const ( TaskEventChainStarted TaskEventKind = "chain_started" TaskEventStepStarted TaskEventKind = "step_started" TaskEventStepChunk TaskEventKind = "step_chunk" TaskEventStepStreamEnd TaskEventKind = "step_stream_end" TaskEventStepCompleted TaskEventKind = "step_completed" TaskEventStepFailed TaskEventKind = "step_failed" TaskEventChainCompleted TaskEventKind = "chain_completed" TaskEventChainFailed TaskEventKind = "chain_failed" // TaskEventChainSuspended terminates a run segment that parked on a human // approval past the fast window: the checkpoint is persisted, the // goroutine is released, and answering the approval resumes the chain as // a fresh run segment under the same request ID. Carries the hierarchical // address of the interrupt point ({chain, task, tool_call}) and // approval_id (== the checkpoint key). TaskEventChainSuspended TaskEventKind = "chain_suspended" TaskEventApprovalRequested TaskEventKind = "approval_requested" TaskEventHITLDecision TaskEventKind = "hitl_decision" TaskEventToolCallPending TaskEventKind = "tool_call_pending" TaskEventToolCall TaskEventKind = "tool_call" TaskEventPrint TaskEventKind = "print" TaskEventTokenUsage TaskEventKind = "token_usage" )
func AllTaskEventKinds ¶
func AllTaskEventKinds() []TaskEventKind
AllTaskEventKinds enumerates every kind the engine can emit. Contract tests iterate it so that adding a kind without updating the documented matrix and its consumers fails CI rather than drifting silently.
type TaskEventScope ¶
type TaskEventSink ¶
type TaskEventSink interface {
PublishTaskEvent(ctx context.Context, event TaskEvent) error
Wants(kind TaskEventKind) bool
}
TaskEventSink receives engine observation events. Wants is per-kind so a sink can decline event kinds it does not consume — and ONLY gates whether events are built and published. It must never select an execution path: execution semantics are identical whether every Wants returns true or false (streaming is observation, not a mode).
type TaskExecutor ¶
type TaskExecutor interface {
// TaskExec executes currentTask against input/dataType and returns its
// output, output type, and transition-eval string (see TaskTransition).
// ctxLength bounds LLM token usage; startingTime anchors timing for the
// whole chain run.
TaskExec(ctx context.Context, startingTime time.Time, ctxLength int, chainContext *ChainContext, currentTask *TaskDefinition, input any, dataType DataType) (any, DataType, string, error)
}
TaskExecutor executes individual tasks within a workflow. Implementations must handle every TaskHandler and return appropriate outputs.
func NewExec ¶
func NewExec( ctx context.Context, repo llmrepo.ModelRepo, toolsProvider ToolsRepo, tracker libtracker.ActivityTracker, ) (TaskExecutor, error)
NewExec creates a new SimpleExec instance
type TaskHandler ¶
type TaskHandler string
TaskHandler defines how task outputs are processed and interpreted.
const ( HandleRaiseError TaskHandler = "raise_error" HandleRoute TaskHandler = "route" HandleChatCompletion TaskHandler = "chat_completion" HandleExecuteToolCalls TaskHandler = "execute_tool_calls" HandleNoop TaskHandler = "noop" HandleTools TaskHandler = "tools" )
func (TaskHandler) String ¶
func (t TaskHandler) String() string
type TaskTransition ¶
type TaskTransition struct {
// OnFailure is the task ID to jump to in case of failure.
OnFailure string `yaml:"on_failure" json:"on_failure" example:"error_handler"`
// Branches defines conditional branches for successful task completion.
Branches []TransitionBranch `yaml:"branches" json:"branches" openapi_include_type:"taskengine.TransitionBranch"`
}
TaskTransition defines what happens after a task completes, including which task to go to next and how to handle errors.
type TokenUsage ¶
type Tool ¶
type Tool struct {
Type string `json:"type"`
Function FunctionTool `json:"function"`
}
Tool represents a tool that can be called by the model.
type ToolCall ¶
type ToolCall struct {
ID string `json:"id" example:"call_abc123"`
Type string `json:"type" example:"function"`
Function FunctionCall `json:"function" openapi_include_type:"taskengine.FunctionCall"`
// ProviderMeta carries opaque provider-specific data (e.g. Gemini thought_signature)
// that must be round-tripped back on the next turn.
ProviderMeta map[string]string `json:"provider_meta,omitempty" example:"{\"thought_signature\":\"123456\"}"`
}
ToolCall represents a tool call requested by the model.
type ToolWithResolution ¶
type ToolsCall ¶
type ToolsCall struct {
// Name is the registered tools-PROVIDER (the service/server, e.g. "slack"),
// not the tool. Required.
Name string `yaml:"name" json:"name" example:"slack"`
// ToolName is the specific TOOL to invoke on that provider
// (e.g. "send_slack_notification").
ToolName string `yaml:"tool_name" json:"tool_name" example:"send_slack_notification"`
// Args are key-value pairs passed to the tool call.
Args map[string]string `` /* 126-byte string literal not displayed */
}
ToolsCall configures a `tools` task — a direct, deterministic call to one tool of one registered tools-provider (e.g. an MCP server), distinct from the model-driven tool calls of chat_completion/execute_tool_calls.
type ToolsProvider ¶
type ToolsProvider interface {
ToolsRegistry
ToolsWithSchema
}
type ToolsRegistry ¶
type ToolsRepo ¶
type ToolsRepo interface {
Exec(ctx context.Context, startingTime time.Time, input any, debug bool, args *ToolsCall) (any, DataType, error)
ToolsRegistry
ToolsWithSchema
}
ToolsRepo defines interface for external system integrations and side effects.
type ToolsWithSchema ¶
type TransitionBranch ¶
type TransitionBranch struct {
// Operator defines how to compare the task's transition eval to When. It is
// REQUIRED and must be one of SupportedOperators() — an empty or unknown
// operator is rejected at chain validation (at runtime it would never match,
// a silent dead branch). The comparison for equals/contains/starts_with/
// ends_with is byte-exact and CASE-SENSITIVE with no trimming — a trailing
// newline (common in multiline template literals) will not match. Only the
// `route` handler normalizes its answer.
Operator OperatorTerm `yaml:"operator,omitempty" json:"operator,omitempty" example:"equals" openapi_include_type:"string"`
// When is the value this branch matches against the task's transition eval.
// What the eval is depends on the handler:
// - chat_completion / execute_tool_calls / tools / noop → a control token,
// one of the Transition* constants (e.g. "tool_call", "executed",
// "tools_executed", "no_calls_found", "noop", "failed"). You CANNOT branch
// on the model's free text here — use the `route` handler for that.
// - route → the model's chosen label (one of the declared branch targets).
// - edge_traversed_at_least → an integer threshold (see that operator).
When string `yaml:"when" json:"when" example:"tool_call"`
// Goto specifies the target task ID if this branch is taken.
// Leave empty or use taskengine.TermEnd to end the chain.
Goto string `yaml:"goto" json:"goto" example:"positive_response"`
// Edge identifies a graph edge "fromTaskID->toTaskID" whose traversal
// count is consulted by edge-state operators (e.g. edge_traversed_at_least).
// Required when Operator is one of those; ignored otherwise.
Edge string `yaml:"edge,omitempty" json:"edge,omitempty" example:"chat->run_tools"`
}
TransitionBranch defines a single possible path in the workflow, selected when the task's output matches the specified condition.
Source Files
¶
- chainlint.go
- checkpoint.go
- context.go
- converter.go
- errors.go
- events.go
- events_journal.go
- handler_signatures.go
- input_cap.go
- inspector.go
- inspector_bus.go
- inspector_kv.go
- llmutil.go
- macroenv.go
- mockexec.go
- retry_outcome.go
- state_sanitize.go
- step_macros.go
- synthesizer.go
- taskengine.go
- taskenv.go
- taskenv_suspend.go
- taskexec.go
- tasktype.go
- toolpairing.go
- tools.go
- toolsctx.go