Documentation
¶
Overview ¶
Package thread gives weft sessions: a conversation as an append-only tree of entries, durable through a Storage backend, with turns, branching, compaction, approvals and delegation all built on that one tree. ADR 0011 is the design record; docs/thread-operations.md is the operator's page.
Sessions ¶
Create starts a session, Open loads one, List pages their headers and Delete removes one. A Session holds the tree in memory and is its session's only writer: every write is one Storage.Append, and nothing is ever rewritten or deleted in place. The leaf is the entry the next one attaches to; Session.Context is what the model sees — the messages on the path from a root to the leaf, repaired.
Open reads and nothing else. It validates the tree — unique valid ids, every parent an earlier entry — and refuses a file that fails with ErrCorrupt; it writes no entry, takes no lock and starts no run. Input a stopped writer had accepted and not settled — steers and queued sends — is restored to the queue and waits for the caller: Session.Queue lists it, the next Send runs behind it, Session.Continue runs it now, Session.ClearQueue drops it.
Session.Close ends a Session: new work is refused, the running turn and the queue drain, every later write fails with ErrClosed, and the session's writer lease is given up. Close every Session that wrote — until it does, no other Session can write that session.
The entry tree and the format ¶
The entry kinds are a sealed set, like core.Event, one wire discriminator each, restored by UnmarshalEntry: message, turn, compaction, branch_summary, leaf, label, info, custom, custom_message; the approval kinds approval_request, approval_decision, approval_audit, grant, grant_revoked; receipt (steering and queued sends) and pool_receipt (delegation). Every entry carries an id, a parent and a time; a message entry embeds the core's message wire (ADR 0001) verbatim.
A stored session is a header ({"weft":1,"type":"session",…} — the envelope integer every weft wire document carries, ADR 0005) and then its entries in append order. The reader's rules: an entry kind added after format 1, or an entry carrying a field an older reader would misread, is written with "v":N, the minimum reader version; a reader that meets a kind or a version it does not know fails with ErrNewerFormat instead of skipping it; every other unknown key is ignored. Golden files in testdata pin each format version this build reads. The format and the API are not frozen: before 1.0 a release may change either, and says so in the CHANGELOG (ADR 0011's format reference documents the current one).
The tree branches without losing anything: Session.Branch appends a leaf entry that moves the leaf to any earlier entry (SummarizeLeft records what the abandoned branch held), and Session.Fork copies a path into a new, self-contained session whose header names its origin. Session.Label, Session.SetInfo, Session.Custom and Session.CustomMessage append the bookkeeping kinds.
Storage ¶
Storage is the backend interface — Create, Append, Load, List, Delete — and package threadtest is the conformance table every backend runs. Three ship: Memory, in process; package jsonl, one JSON-lines file per session, readable with jq and backed up with cp; and package sqlite, in its own module so its driver stays out of this one. A backend's Append is atomic against the writer's death, and what Load returns never aliases what is stored. Optional capabilities are small interfaces found by type assertion: Flusher (the FsyncOnFlush durability cadence), Releaser (the backend's hold on a session, ended), Leaser (the per-Session writer lease) and Watcher, the live tail — entries yielded as they are appended, for a reader that follows a session another process writes.
A crash mid-append leaves at most a torn final line. A load drops it and the next writer removes it before appending; a damaged line elsewhere fails the load with ErrCorrupt unless the backend was opened with Salvage. None of it is silent: Session.LoadReport says what was dropped, skipped, and orphaned by the skip.
Turns and busy policies ¶
Session.Send appends the prompt — durable before the run starts — runs the session's agent over the leaf's context under the run id <session>-t<n>, and returns a Turn at once: its receipt and its handle. Turn.Wait blocks for the result, Turn.WaitContext bounds the wait, Turn.Done is the channel form, Turn.Events streams the run, and Turn.Outcome names the end — answered, parked, delivered, deferred, dropped, failed or canceled. The run's messages are appended as they join the run, each step exactly once, and a turn entry closes the ledger whether the run succeeded, failed or was canceled.
One run per session at a time. A Send that meets a running turn follows the busy policy, set with BusyPolicy or per Send with As: Queue (the default) runs it next; Reject fails it with ErrBusy; Steer delivers it into the running turn at its next drain point (ADR 0019); Interrupt cancels the running turn and runs the message next; Rollback also branches back to before the interrupted turn. Whatever is accepted is durable at acceptance: a queued send and a steer each write a receipt entry before Send returns, so neither is lost to a crash, and Session.Queue lists both.
A turn is decided before the session's between-turn work — the automatic compaction — so Wait never waits on a summarizer. Session.WaitIdle waits for that work too.
Compaction, approvals, delegation ¶
Compaction (ADR 0020) keeps a long conversation inside the model's window without deleting anything: a compaction entry carries a summary and the id of the first entry kept raw, and the context reads the summary in place of what came before. ContextWindow arms the automatic trigger; Session.Compact, Session.PreviewCompaction and Session.ApplyCompaction are the manual path; a turn that fails with core.ErrContextOverflow compacts and runs once more.
Approvals (ADR 0021) make the core's approval boundary durable: a parked call is an approval_request entry, Session.Pending lists them across restarts, Session.Decide records a decision and resumes the turn. Grants, a live Approver, quorum, expiry and signed decisions (Keyring, DecideSigned, RequireSigned) form the chain a call passes before it parks and the doors a decision comes through.
Package pool (ADR 0022) runs bounded concurrent child sessions for a parent: each child is a session with a Lineage, its cost lands in the parent's Usage.Delegated, its journey is pool_receipt entries on the parent, and its parked calls mirror onto the parent's Pending.
Identity ¶
Every run a session starts carries weft.session.id and weft.turn as run metadata, the fork origin or pool lineage when the header names one, and weft.public_id when the session was created with PublicID — an opaque handle safe to show a browser (ADR 0024 S5). The public id lives in the header, is matched by List's Query.Meta, and cannot be changed after Create.
Concurrency ¶
The whole contract, in one place; the Session type's godoc carries the detail.
- A Session is safe for concurrent use. One mutex guards its tree and is held across each storage write, so a returned write is durable and the tree in memory equals the stored one.
- One writer per session. Backends refuse a second writer from another process or another Storage value with ErrLocked. Between Session values on one Storage value the rule is a lease: a Session takes it with its first write (Create and Fork are one) and holds it until its Close; another Session's writes fail with ErrLocked and change nothing.
- Reading takes nothing. Open, List, Load, Watch and every read of a Session work while another writer holds the session.
- A Session that has fallen behind does not write: when the stored session holds entries it never loaded, its write fails with ErrStale. Open the session again.
- One run at a time; concurrent Sends are accepted in the order they take the lock and follow their busy policy. Branch, Compact, ApplyCompaction and Uncompact are between-turns operations and fail with ErrBusy while a turn runs; Fork, the reads and the bookkeeping writes do not wait for a turn.
- Hooks and callbacks run without the session's lock and may call the session. The IDs and Clock functions are the exception: they run under it and must not.
- Delete does not look for open Sessions. Close a session before deleting it.
Errors ¶
Failures are sentinel errors, wrapped with context and matched with errors.Is. Three classes tell a caller what to do next:
- Retry: ErrBusy (a turn is running) and ErrLocked (another writer holds the session). The same call succeeds once the turn has ended or the other writer has closed.
- Reopen: ErrStale (the session moved on behind this Session) and ErrClosed (this Session's Close has run). The Session value is done writing; Open the session again.
- Terminal for the stored data as this build reads it: ErrCorrupt (carried by *CorruptError, naming the line and the entry) and ErrNewerFormat (written by a newer weft). Retrying changes nothing; Salvage skips corrupt lines, never newer ones.
The rest name a refused call — ErrNotFound, ErrExists, ErrCreateOnly, ErrReservedKey, ErrNotPending, ErrDelegated, ErrInvalidDecision, the signed-decision and compaction sentinels — or how a turn ended: the error Turn.Wait returns wraps ErrNotRun, ErrDropped, ErrNotPersisted or ErrTurnPanicked, or is the run's own *core.RunError.
Index ¶
- Constants
- Variables
- func Delete(ctx context.Context, st Storage, id string) error
- func NewEntryID() string
- func NewSessionID() string
- func ValidID(id string) bool
- type ApprovalAuditEntry
- type ApprovalDecisionEntry
- type ApprovalRequestEntry
- type Approver
- type Arg
- type BranchOption
- type BranchSummaryEntry
- type CompactOption
- type Compaction
- type CompactionEntry
- type Compactor
- type CorruptError
- type CustomEntry
- type CustomMessageEntry
- type Decision
- type Entry
- type Estimator
- type Flusher
- type Grant
- type GrantEntry
- type GrantRevokedEntry
- type GrantStore
- type Header
- type InfoEntry
- type Key
- type Keyring
- type LabelEntry
- type LeafEntry
- type Leaser
- type Lineage
- type LoadReport
- type MessageEntry
- type MirroredRequest
- type NativeCompactor
- type OpenOption
- type OpenReport
- type Outcome
- type Page
- type ParentRef
- type Policy
- type PoolReceiptEntry
- type Preparation
- type Query
- type QueuedSteer
- type Reason
- type ReceiptEntry
- type Releaser
- type Request
- type SendOption
- type Session
- func (s *Session) AppendApprovalRequests(ctx context.Context, reqs ...ApprovalRequestEntry) ([]ApprovalRequestEntry, error)
- func (s *Session) AppendPoolReceipt(ctx context.Context, e PoolReceiptEntry) (PoolReceiptEntry, error)
- func (s *Session) ApplyCompaction(ctx context.Context, c *Compaction) error
- func (s *Session) Audit() []Entry
- func (s *Session) Branch(ctx context.Context, entryID string, opts ...BranchOption) error
- func (s *Session) CancelDelegated(ctx context.Context, reason string) error
- func (s *Session) ClearQueue(ctx context.Context) (int, error)
- func (s *Session) Close(ctx context.Context) error
- func (s *Session) Compact(ctx context.Context, opts ...CompactOption) error
- func (s *Session) Context() []core.Message
- func (s *Session) Continue(ctx context.Context) (*Turn, error)
- func (s *Session) Custom(ctx context.Context, kind string, data json.RawMessage) error
- func (s *Session) CustomMessage(ctx context.Context, kind string, msg core.Message) error
- func (s *Session) Decide(ctx context.Context, ds ...Decision) (*Turn, error)
- func (s *Session) DecideSigned(ctx context.Context, sd SignedDecision) (*Turn, error)
- func (s *Session) DenyMirrored(ctx context.Context, child, reason string) error
- func (s *Session) Entries() []Entry
- func (s *Session) Fork(ctx context.Context, entryID string, opts ...SessionOption) (*Session, error)
- func (s *Session) Grant(ctx context.Context, g Grant) error
- func (s *Session) ID() string
- func (s *Session) Label(ctx context.Context, entryID, name string) error
- func (s *Session) Leaf() string
- func (s *Session) Lineage() Lineage
- func (s *Session) LoadReport() *OpenReport
- func (s *Session) Meta() map[string]string
- func (s *Session) MirroredRequests(ctx context.Context) ([]MirroredRequest, error)
- func (s *Session) Path(entryID string) ([]Entry, error)
- func (s *Session) Pending() []Request
- func (s *Session) Pin(ctx context.Context, entryID string) error
- func (s *Session) PreviewCompaction(ctx context.Context, opts ...CompactOption) (*Compaction, error)
- func (s *Session) Queue() []QueuedSteer
- func (s *Session) ReplayDecisions(ctx context.Context, ds ...ApprovalDecisionEntry) (*Turn, error)
- func (s *Session) Request(callID string) (Request, error)
- func (s *Session) ResolveDelegation(ctx context.Context, callID, child, content string, isError bool) (*Turn, error)
- func (s *Session) Resume(ctx context.Context) (*Turn, error)
- func (s *Session) Revoke(ctx context.Context, grantID string) error
- func (s *Session) Send(ctx context.Context, msg core.Message, opts ...SendOption) (*Turn, error)
- func (s *Session) SetInfo(ctx context.Context, title string, meta map[string]string) error
- func (s *Session) Storage() Storage
- func (s *Session) Title() string
- func (s *Session) Uncompact(ctx context.Context) error
- func (s *Session) Usage() Usage
- func (s *Session) WaitIdle(ctx context.Context) error
- func (s *Session) WatchMirrors(owner any, fn func(log *slog.Logger))
- type SessionOption
- func AfterCompact(fn func(ctx context.Context, e CompactionEntry)) SessionOption
- func AutoResume(on bool) SessionOption
- func BeforeCompact(fn func(ctx context.Context, p *Preparation) (Verdict, error)) SessionOption
- func BusyPolicy(p Policy) SessionOption
- func CheckSummary(fn func(Summary) error) SessionOption
- func ClearOldToolResults(keepLast int) SessionOption
- func Clock(now func() time.Time) SessionOption
- func CompactFailed(fn func(ctx context.Context, r Reason, err error)) SessionOption
- func ContextWindow(n int64) SessionOption
- func IDs(id func() string) SessionOption
- func InheritApprovals(parent *Session) SessionOption
- func KeepRecent(n int64) SessionOption
- func MaxPerSession(n int) SessionOption
- func MinTurnsBetween(n int) SessionOption
- func ModelReserves(reserves map[core.ModelInfo]int64) SessionOption
- func ModelWindows(windows map[core.ModelInfo]int64) SessionOption
- func NoAutoCompact() SessionOption
- func OnRequest(fn func(Request)) SessionOption
- func PreferNative() SessionOption
- func PublicID(id string) SessionOption
- func Quorum(n int) SessionOption
- func ReRunOnOverflow(on bool) SessionOption
- func RequestExpiry(d time.Duration) SessionOption
- func RequireSigned() SessionOption
- func Reserve(n int64) SessionOption
- func SummaryFocus(focus string) SessionOption
- func SummaryMaxTokens(n int64) SessionOption
- func SummaryModel(m core.Model) SessionOption
- func SummaryPrompt(tmpl string) SessionOption
- func TriggerFunc(fn func(TriggerInput) bool) SessionOption
- func WithApprover(a Approver, timeout time.Duration) SessionOption
- func WithCompactor(c Compactor) SessionOption
- func WithEstimator(est Estimator) SessionOption
- func WithGrantStore(gs GrantStore) SessionOption
- func WithKeyring(r *Keyring) SessionOption
- func WithLineage(parentSession, call string) SessionOption
- func WithMeta(meta map[string]string) SessionOption
- func WithSummarizer(s Summarizer) SessionOption
- func WithTrimmer(t Trimmer) SessionOption
- type SharedGrant
- type SignedDecision
- type Storage
- type Summarizer
- type Summary
- type SummaryInput
- type TriggerInput
- type TrimRecord
- type TrimStub
- type Trimmer
- type Turn
- func (t *Turn) Done() <-chan struct{}
- func (t *Turn) Events() iter.Seq2[core.Event, error]
- func (t *Turn) ID() string
- func (t *Turn) Next() *Turn
- func (t *Turn) Outcome() TurnOutcome
- func (t *Turn) RunID() string
- func (t *Turn) Wait() (*core.RunResult, error)
- func (t *Turn) WaitContext(ctx context.Context) (*core.RunResult, error)
- type TurnEntry
- type TurnOutcome
- type Usage
- type Verdict
- type Watcher
Examples ¶
- As
- BeforeCompact
- BusyPolicy
- CheckSummary
- ClearOldToolResults
- Clock
- Create
- Delete
- Header
- List
- Memory
- MessageEntry
- NewSessionID
- Open
- PublicID
- Query
- Query (KeysetCursor)
- Quorum
- ReRunOnOverflow
- Releaser
- RequestExpiry
- RunOptions
- Salvage
- Session.ApplyCompaction
- Session.Audit
- Session.Branch
- Session.Close
- Session.Close (HandOver)
- Session.Compact
- Session.Continue
- Session.Decide
- Session.DecideSigned
- Session.Fork
- Session.Grant
- Session.Label
- Session.Path
- Session.Pin
- Session.PreviewCompaction
- Session.Queue
- Session.Resume
- Session.Send
- Session.Usage
- SummarizeLeft
- SummaryModel
- TriggerFunc
- Turn.Done
- Turn.Events
- Turn.Next
- Turn.WaitContext
- Watcher
- WithApprover
- WithMeta
- WithSummarizer
Constants ¶
const ( StepGrant = "grant" // a matching live grant decided; GrantID names it StepApprover = "approver" // the live chain step, consulted and bounded StepPark = "park" // the request persisted, the turn ended pending StepExpiry = "expiry" // an expired request denied StepResume = "resume" // the boundary's resume run: one entry "started" before the run, one "completed" or "failed" with its end StepSigned = "signed" // a signed decision refused; Detail names why, no decision recorded )
The audit entry's steps (ADR 0021 §2): every step the session takes over a parked call leaves an approval_audit entry naming its step.
const ( // ReceiptQueued marks a steer's acceptance: the message is held // for the running turn's steering drain. ReceiptQueued = "queued" // ReceiptAccepted marks a queued send's acceptance: the message is // held for a turn of its own, whose prompt entry will take the id // in Turn. ReceiptAccepted = "accepted" // ReceiptDelivered marks a message the running run drained: it is // an ordinary transcript message of that run, named by RunID. ReceiptDelivered = "delivered" // ReceiptDeferred marks a steer that became a follow-up turn // (Turn names its prompt entry): a StopWhen end, an open approval // boundary, or the run ended before the drain. ReceiptDeferred = "deferred" // ReceiptDropped marks a message removed before it reached a // model: ClearQueue, or an interrupting send that was refused. ReceiptDropped = "dropped" )
Receipt statuses — the wire values, pinned by the format-3 goldens. "accepted" joined them without a format bump: a reader from before it restores only "queued" receipts and reads an accepted one as a receipt it has nothing to do for — the queued send is not run by that reader, exactly as before the status existed, and nothing is misread.
const ( // PoolAccepted marks the delegation recorded and queued for a slot. PoolAccepted = "accepted" // PoolRunning marks the slot acquired and the child session's turn // started — the wait between acceptance and running is the pool's // queue, visible. PoolRunning = "running" // PoolParked marks a child whose run ended at an approval boundary // (ADR 0022 §7): it holds no slot and waits for decisions on the // requests mirrored onto this session. The next "running" entry is // its resume. A status, not a new kind: a reader from before it // sees one more unsettled state, which is what it is. PoolParked = "parked" // PoolDone marks a child that ran to its intended end; Stop is its // answer, Usage its total cost. PoolDone = "done" // PoolFailed marks a child whose run failed; Stop is the cause. PoolFailed = "failed" // PoolCanceled marks a child canceled — by an explicit Cancel, by // the pool's Close, or, for a sync child, with the delegating call // it ran under — never by the submitting turn's own end, which an // async child survives by design (ADR 0022 D4). PoolCanceled = "canceled" // PoolCapped marks a child that died on a budget — ErrMaxSteps or // ErrUsageLimit (DeerFlow's token-capped/turn-capped collapsed; the // stop reason distinguishes them). PoolCapped = "capped" )
Pool receipt statuses — the wire values, pinned by the format-4 goldens. The machine is accepted → running ⇄ parked → exactly one of done, failed, canceled, capped; a delegation canceled or failed before its child ever ran settles straight from accepted.
const FormatVersion = 1
FormatVersion is the wire version of every document the thread module writes (ADR 0011 §6) — the integer every weft wire document carries ("weft": 1, ADR 0005). It moves only for a layout change a reader cannot handle additively: a new entry kind, a new optional key, a new entry version never move it. A header carrying a higher number fails with ErrNewerFormat; nothing is ever rewritten on open.
The format is not frozen. This build reads every file an earlier thread release wrote — the committed goldens of every format pin that, and a release that could not read one would fail its own tests — but no compatibility promise beyond what the tests hold is made before the module's API and format are declared stable; there is no migration tool, and none is needed while every change so far has been additive.
Variables ¶
var ( // ErrNothingToCompact is the refusal of a compaction that would // summarize nothing: the tail fits inside KeepRecent, nothing new // sits past the previous boundary, or the plan names the boundary // the previous compaction already left. Nothing is written. ErrNothingToCompact = errors.New("thread: nothing to compact") // ErrCompactCanceled is returned when a BeforeCompact hook answered // Cancel. Nothing is written. ErrCompactCanceled = errors.New("thread: compaction canceled") // ErrNoEntry is returned when an operation names an entry the // session does not hold — Pin's target, a compaction's FirstKept — // or asks for one the leaf's path does not have (Uncompact with no // compaction to undo). ErrNoEntry = errors.New("thread: session holds no such entry") // ErrSummaryTruncated is the failure of a summary the model cut // off at its output cap (a max_tokens finish). A cut summary is // never stored: the attempt is retried once, the chain falls back, // and with no fallback left the compaction fails wrapping this. // Raise SummaryMaxTokens (or Reserve, which sizes the default cap). ErrSummaryTruncated = errors.New("thread: summary truncated at the output cap") // ErrInvalidCompaction is returned for a Compaction that cannot be // written as given: a FirstKept off the leaf's path or before the // previous boundary, a summary-less plan that is not a trim, a trim // without its record, a trim record naming a result the path does // not hold — and for a Trimmer whose output the trim record cannot // represent. ErrInvalidCompaction = errors.New("thread: invalid compaction") // ErrCompactConfig is returned by Create and Open when the // compaction knobs cannot work together: with a known window, // Reserve must be below it and KeepRecent below window − Reserve. ErrCompactConfig = errors.New("thread: invalid compaction configuration") // ErrNotPinnable is returned by Pin for an entry that contributes // no message to the context (a turn, a label, a custom entry, …): // there is nothing a compaction could keep showing. ErrNotPinnable = errors.New("thread: entry cannot be pinned") // ErrAwaitingApproval is returned by Compact and ApplyCompaction // while approval requests are pending: the parked tail must stay // raw for its decisions to resolve. Decide, or branch away, first. ErrAwaitingApproval = errors.New("thread: approval requests pending") )
The compaction sentinels. Match with errors.Is; every error this file returns for one of these conditions wraps its sentinel.
var ( // ErrBusy is returned by Send under the Reject busy policy when the // session is already running a turn, and by the between-turns // operations while a turn is in flight — Branch, Compact, // ApplyCompaction, Uncompact and CustomMessage: the call was not // accepted and nothing was written (ADR 0011 §4). Retryable — the // same call succeeds once the turn has ended. ErrBusy = errors.New("thread: session is busy with another turn") // ErrNotFound is returned by Load, Append, and Delete for a session // id the storage does not hold. ErrNotFound = errors.New("thread: session not found") // ErrExists is returned by Create for a session id the storage // already holds — in this process, or in any other sharing the // backend. A session is never silently replaced (ADR 0011 §5). ErrExists = errors.New("thread: session already exists") // ErrLocked is returned when a session is already held by another // writer (the one-writer rule, ADR 0011 §5): by a backend, for a // writer in another process or another Storage in this one; and by // a Session's writes, for another Session value on the same // Storage (the Leaser capability). Readers never lock: Load, List // and Open always work. Retryable — the same call succeeds once // the other writer has let go, which for a Session is its Close. ErrLocked = errors.New("thread: session is locked by another writer") // ErrStale is returned by a Session's write when the stored // session is not the one the Session loaded: another writer // appended to it after this Session was opened — a write now // would attach to a leaf that is no longer the session's, a fork // nobody asked for — or the session was deleted and created again // under the same id. The check is the stored header's Created and // the number of entry lines, not a comparison of contents. Nothing // is written and the Session's tree is unchanged. Terminal for the // Session value: Open the session again to write from what it now // holds. ErrStale = errors.New("thread: session changed since it was opened") // ErrCorrupt wraps the failures a backend reports for stored data // it cannot decode: a malformed line that is not a torn tail (a // torn final line is a crash, dropped and reported through Load's // LoadReport; data from a newer weft is ErrNewerFormat, never // skipped). Salvage, the open option every backend accepts, // downgrades this to a skip reported in the LoadReport. ErrCorrupt = errors.New("thread: stored session data is corrupt") // ErrNewerFormat wraps the decode failures UnmarshalEntry and // Header decoding report for data written by a newer weft than // this build: an entry kind this build does not know, an entry // carrying a "v" above the version this build reads of its kind, // or a session header whose envelope is ahead of FormatVersion // (ADR 0011 §5–§6). Loud over silent: a session must decode to // exactly what was written, and an older weft says so instead of // guessing. List, which reads headers only, still works — an older // reader sees a session whose entries are newer and Load names why // it cannot open it. The one thing List cannot show is a session // whose header itself is newer: the envelope gates the whole file, // so such a session is invisible to an older List — visible again // as soon as its directory is read by a weft that knows the format. ErrNewerFormat = errors.New("thread: session format newer than this build") // ErrNotPending is returned by Decide and DecideSigned for a // decision addressing a call that is not pending — decided already, // resumed already, or never parked — or, signed, carrying a // signature for another occurrence of the call; by Request for a // call that is not pending; by ResolveDelegation on a call that // delegates to no child; and by Resume with no open boundary. The // core's rule (ADR 0007: a decision names a pending call) made // strict and raised before any entry lands or any run starts: a // session never records a decision it cannot apply. ErrNotPending = errors.New("thread: call is not pending") // ErrUnknownKey is returned by Keyring.Sign when the ring does not // hold the challenge's key. DecideSigned never returns it: a // signature under a key id the session's ring does not hold is // ErrBadSignature, so a caller probing key ids learns nothing. ErrUnknownKey = errors.New("thread: signing key not in the keyring") // ErrBadSignature is returned by DecideSigned when the signed // decision does not verify against the session's keyring and the // pending request: a MAC mismatch, an unknown key id, an outcome // that is none of the four, a session other than this one, an // empty nonce or one this session never issued for the request, // or an expiry or a tool that is not the request's. ErrBadSignature = errors.New("thread: decision signature does not verify") // ErrExpired means the request is past its expiry: returned by // Decide, DecideSigned and Request. No decision can approve a // lapsed request — the session denies it on its own (ADR 0021 §5). ErrExpired = errors.New("thread: approval request expired") // ErrReplay is returned by DecideSigned when the nonce already // answered a recorded decision: a signature decides once, across // restarts. ErrReplay = errors.New("thread: decision signature replayed") // ErrArgsChanged is returned by DecideSigned when the signed // arguments hash is not the pending request's: the signer approved // a different call. ErrArgsChanged = errors.New("thread: request arguments changed under the signature") // ErrSignatureRequired is returned by Decide on a session opened // with RequireSigned: the unsigned door is closed, and only // DecideSigned records caller-held decisions (ADR 0021 §3). ErrSignatureRequired = errors.New("thread: this session requires signed decisions") // ErrClosed is returned by Send, Continue and every write on a // Session whose Close has run: a closed session accepts no new // work and writes nothing. Reads keep answering from the tree the // session held when it closed. Terminal for the Session value — // Open the session again to continue it. ErrClosed = errors.New("thread: session is closed") // ErrNotPersisted is wrapped by the error Turn.Wait returns when // the turn ran but its end could not be written: the storage // refused the batch that carries the turn entry (and whatever // messages the steps had not already written), so the session's // tree holds no record of how the turn ended and its usage is in // no ledger. The run's result still rides beside the error — the // model answered; the storage did not keep it. Messages the steps // wrote as they joined are in the tree. ErrNotPersisted = errors.New("thread: turn end not persisted") // ErrNotRun is wrapped by the error Turn.Wait returns for a turn // whose run never started: the context its Send carried ended // first (the error also wraps context.Canceled or // context.DeadlineExceeded), its prompt entry could not be written // (the storage's error is wrapped too), or — a resume — the // boundary it was armed for was gone (ErrNotPending) or its audit // entry could not be written. No model was called. ErrNotRun = errors.New("thread: turn did not run") // ErrDropped is wrapped by the error Turn.Wait returns for a // message ClearQueue removed before it reached a model — a queued // steer or a queued send — and for a send an interrupting Send had // to take back. The receipt entry records the drop. ErrDropped = errors.New("thread: message dropped from the queue") // ErrTurnPanicked is wrapped by the error Turn.Wait returns when // the session's own turn machinery panicked — which includes the // caller's IDs and Clock functions, called under the session's // lock. The panic is contained: the turn ends with this error, the // lock is released, and the session keeps working. Tool and model // panics never surface here; the core reports them as run errors. ErrTurnPanicked = errors.New("thread: turn panicked") // ErrCreateOnly is returned by Open when it is handed an option // that only means something while a session's header is being // written — WithMeta, PublicID, WithLineage. The header is // immutable once stored, so Open refuses the option instead of // ignoring it; Create and Fork honour it. ErrCreateOnly = errors.New("thread: option applies only when a session is created") // ErrReservedKey is returned by SetInfo for a metadata key under // the reserved "weft." prefix: those keys are the session's // identity (weft.public_id), set once in the header at Create — // WithMeta, PublicID — and never edited afterwards. ErrReservedKey = errors.New("thread: metadata key is reserved") )
var ErrDelegated = errors.New("thread: call is delegated to a child session")
ErrDelegated is returned by Decide, DecideSigned and Request for a call that delegates to a thread/pool child session (ADR 0022 §7): a call some mirrored child request names as its Wrapper. Such a call is parked because its child is, and it completes with the child's answer — never by a decision of its own: approving it would run the delegation a second time, in a second child session. The session's own decision chain keeps the same rule without an error: a grant or an Approver is never consulted for such a call. Decide the child's requests (Pending lists them, Child naming the session); the pool resolves the call when the child ends.
var ErrInvalidDecision = errors.New("thread: invalid decision")
ErrInvalidDecision is returned by Decide for a batch that cannot be recorded as written: no decisions at all, a decision without an outcome, or two decisions naming the same call — one call takes one verdict per Decide, and a batch that both approves and denies it says nothing the session could apply. Nothing is recorded.
Functions ¶
func Delete ¶
Delete removes a session and its entries from st; an id the storage does not hold fails with ErrNotFound. History is removed with the session, never rewritten (ADR 0011 §5).
Delete does not look for open Session values, and a Session's lease does not stop it: through the Storage value the session's writer uses, Delete always removes (Storage.Delete's rule; through another Storage value it fails with ErrLocked while the writer holds the session). A Session already loaded keeps its in-memory tree, its next write fails with ErrNotFound, and a turn it is running loses the entries it has yet to write — Close the session first.
Example ¶
Delete removes a session and its entries. Close the Session first: Delete does not look for open ones.
package main
import (
"context"
"errors"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// exampleIDs mints the given ids in order — the IDs option's usual
// shape in an example: the session id first, then one per entry.
func exampleIDs(ids ...string) thread.SessionOption {
next := 0
return thread.IDs(func() string {
id := ids[next]
next++
return id
})
}
func main() {
ctx := context.Background()
st := thread.Memory()
agent := weft.New(wefttest.Script())
s, err := thread.Create(ctx, st, agent, exampleIDs("s_scratch"))
if err != nil {
fmt.Println(err)
return
}
if err := s.Close(ctx); err != nil {
fmt.Println(err)
return
}
fmt.Println(thread.Delete(ctx, st, "s_scratch"))
_, err = thread.Open(ctx, st, "s_scratch", agent)
fmt.Println(errors.Is(err, thread.ErrNotFound))
fmt.Println(errors.Is(thread.Delete(ctx, st, "s_scratch"), thread.ErrNotFound))
}
Output: <nil> true true
func NewEntryID ¶
func NewEntryID() string
NewEntryID returns a new entry id: "e_" plus 26 time-sortable random characters (see NewSessionID).
func NewSessionID ¶
func NewSessionID() string
NewSessionID returns a new session id: "s_" plus 26 characters — a millisecond-precision timestamp then random bits, Crockford base32, the ULID shape — so ids sort with the sessions they name. Sorting is to the millisecond only: ids born in the same millisecond order randomly among themselves, and a session's order of record is its entry order, never id order.
Example ¶
NewSessionID and NewEntryID mint time-sortable ids: "s_" and "e_" plus 26 characters, so a directory of sessions lists in creation order and a session's entries debug readably.
package main
import (
"fmt"
"github.com/weftgo/weft/thread"
)
func main() {
fmt.Println(len(thread.NewSessionID()), len(thread.NewEntryID()))
}
Output: 28 28
func ValidID ¶
ValidID reports whether id is safe as a session or entry id: 1 to 128 characters of letters, digits, '_' or '-', so an id is always exactly one path component — never a separator, a traversal, or a dotfile (ADR 0011 §5: backends vet the ids they are given). weft's own ids satisfy it by construction; the ids a caller supplies must stay inside it too.
Types ¶
type ApprovalAuditEntry ¶
type ApprovalAuditEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
CallID string `json:"call_id,omitempty"`
Step string `json:"step"`
Outcome string `json:"outcome,omitempty"`
Detail string `json:"detail,omitempty"`
RunID string `json:"run_id,omitempty"`
// GrantID names the grant a StepGrant entry matched — a session
// grant's entry id, or with GrantShared the GrantStore's own id
// for it. A session grant's MaxUses is counted from these fields:
// one entry is one use.
GrantID string `json:"grant_id,omitempty"`
// KeyID names the keyring key a refused signed decision claimed
// (StepSigned), when the ring holds it.
KeyID string `json:"key_id,omitempty"`
// Decisions lists the decision entries a resume applied
// (StepResume, outcome "started"). A decision is spent by the
// resume that applied it: it resolves its call on that resume's
// own line of the tree and nowhere else, so a Branch back to the
// decided boundary asks for a new decision instead of running the
// call again on the old one.
Decisions []string `json:"decisions,omitempty"`
}
ApprovalAuditEntry is the chain's own trail (ADR 0021 §2): every step the decision chain takes over a call — a grant matched, the Approver consulted (decided, declined, timed out), the park, an expiry denial, a signed decision refused, a resume started — leaves one of these, including automatic approvals, so s.Audit() can tell the whole story from the file alone. Step and Outcome name the step and how it ended; Detail is prose for a human reader and never parsed — what the session reads back lives in the typed fields. It never enters the model's context.
func (ApprovalAuditEntry) MarshalJSON ¶
func (e ApprovalAuditEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator and "v":2.
type ApprovalDecisionEntry ¶
type ApprovalDecisionEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
CallID string `json:"call_id"`
Outcome Outcome `json:"outcome"`
Reason string `json:"reason,omitempty"`
Content string `json:"content,omitempty"`
Who string `json:"who,omitempty"`
Via string `json:"via,omitempty"`
RunID string `json:"run_id,omitempty"`
// Nonce and KeyID are the signed-decision replay guard's record
// (ADR 0021 §3): a decision that arrived signed carries the nonce
// it answered and the key that vouched for it, so a replayed
// signature is detectable from the file alone — across restarts.
// KeyID is also the signed decision's identity under Quorum: one
// key is one approver, whatever Who says. Empty on the in-process
// paths, which mint no challenge.
Nonce string `json:"nonce,omitempty"`
KeyID string `json:"key_id,omitempty"`
// RequestID is the id of the request entry the decision answers —
// the occurrence of the call, which a call id alone does not name.
// Set on every decision recorded over a parked request — Decide,
// DecideSigned (whose challenge is bound to it), an expiry or
// interrupt denial; empty on a chain step's decision (a grant, the
// Approver), written in the same append as the request or without
// one, and RunID names the occurrence either way.
RequestID string `json:"request_id,omitempty"`
// Always records that the decision was an "approve and always
// allow" (ApproveAlways): the grant is minted when the call's
// effective verdict becomes approve — at once without a quorum,
// with the completing approval under one — and never when the
// verdict is anything else.
Always bool `json:"always,omitempty"`
}
ApprovalDecisionEntry is one decision over a parked call (ADR 0021 §1): approve, deny with a reason, or resolve with content computed outside the process (resolve_error marks it an error). It records Who decided, When (the entry's Created), Via which channel — "user" for Decide, "signed" for DecideSigned, "approver" for the live chain step, "grant" for a grant's approval, "expiry" for an expired request's automatic denial, "interrupt" for the denial an interrupting Send records, "child" for a pool delegation's resolution, "parent" for a decision a pool child's parent session took and the pool replayed into the child — and the run the decided request belonged to. It never enters the model's context; the model sees the decision only through the result the resumed run produces.
func (ApprovalDecisionEntry) MarshalJSON ¶
func (e ApprovalDecisionEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator and "v":2.
type ApprovalRequestEntry ¶
type ApprovalRequestEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
CallID string `json:"call_id"`
Tool string `json:"tool"`
Args json.RawMessage `json:"args,omitempty"`
ArgsSHA256 string `json:"args_sha256"`
RunID string `json:"run_id"`
Reason string `json:"reason,omitempty"`
Expiry time.Time `json:"expiry,omitzero"`
// Child names the delegated child session this request mirrors,
// and Wrapper the parent-side delegating call it parks under (ADR
// 0022 §7): a pool child's parked call is requested in its own
// session and mirrored here so the parent's Pending surfaces it
// with its lineage. Both empty on an ordinary request.
Child string `json:"child,omitempty"`
Wrapper string `json:"wrapper,omitempty"`
}
ApprovalRequestEntry is a parked call made durable (ADR 0021 §1): a call the core's approval boundary left unexecuted, recorded with the arguments and their SHA-256 so a decision can name exactly what it decided, the run that parked it, why it parked, and an optional expiry — a request past its expiry takes no decision and is denied with the stated reason the next time the session looks at the boundary (ADR 0021 §5). It is written in the same Append as the turn that parked it, so no window exists where the turn is durable and the request is not. Its ID names this one occurrence of the call: call ids repeat across turns, the entry id never does, and a signed decision is bound to it. It never enters the model's context; the pending call itself stays unresolved in the transcript until a decision resolves it.
func (ApprovalRequestEntry) MarshalJSON ¶
func (e ApprovalRequestEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator and "v":2.
type Approver ¶
Approver is the decision chain's live step (ADR 0021 §2): the "ask now" path for a terminal or an already-connected UI. It returns the Decision and true when it decided, or false to decline (a pipe, no TTY). It is never given unbounded time: the session consults it under WithApprover's timeout, and a timeout reads as a decline. It must not call back into the same Session — the chain consults it outside the session lock, but the turn is still running. The Request it is handed has no ID, Nonce or Created: the chain asks before the request is persisted.
type Arg ¶
type Arg struct {
// Pointer is the RFC 6901 JSON pointer to the argument: "/order_id",
// "/file/path", "/command".
Pointer string `json:"pointer"`
// Equals is the raw JSON the value must equal (ArgEquals).
Equals json.RawMessage `json:"equals,omitempty"`
// Prefix is the string prefix the value must have (ArgPrefix).
Prefix string `json:"prefix,omitempty"`
// Glob is the glob the string value must match (ArgGlob) — *
// spans separators, a command is not a path.
Glob string `json:"glob,omitempty"`
}
Arg is one argument predicate (ADR 0021 §4): a JSON pointer into the call's arguments and the test the value there must pass — equality with ArgEquals, a string prefix with ArgPrefix, a glob with ArgGlob. Build one with its constructor; a hand-made Arg with several tests set uses Equals first, then Prefix, then Glob.
func ArgEquals ¶
func ArgEquals(pointer string, value json.RawMessage) Arg
ArgEquals returns the predicate requiring args[pointer] to equal value as JSON: strings and booleans exactly, arrays and objects deeply (key order never matters), and numbers by this rule — two integers (no fraction, no exponent) compare exactly, digit for digit, however large; as soon as either side has a fraction or an exponent both compare as float64, so 1, 1.0 and 1e0 are equal and two values a float64 cannot tell apart are equal too. Pointer "" is the whole arguments document; a call that carried no arguments reads as the empty object {}.
func ArgGlob ¶
ArgGlob returns the predicate requiring args[pointer] to be a string matching glob — the "go test …" command shape. * matches any run of characters, the empty run included, and ? exactly one character (one Unicode code point, never one byte of it); every other character is itself. There is no escape and there are no classes: a literal * or ? cannot be asked for — use ArgEquals or ArgPrefix when the value holds one. A command is not a path: * spans separators, so "go test*" matches "go test ./..." (and anything after it on the line — anchor the glob's tail when that matters).
func ArgPrefix ¶
ArgPrefix returns the predicate requiring args[pointer] to be a string with prefix — the "path stays under the workspace" shape. A prefix is a plain string prefix, nothing more: end it with the separator or it matches siblings too ("/ws" matches "/ws-evil") — the prefix's one sharp edge, documented rather than hidden. An empty prefix matches nothing (an Arg with no test never matches, the same rule a hand-made Arg follows): the tool-wide grant is no Args at all, not an empty prefix.
type BranchOption ¶
type BranchOption interface {
// contains filtered or unexported methods
}
BranchOption configures one Branch call — the per-call layer over the session's own SessionOptions. There is one: SummarizeLeft.
func SummarizeLeft ¶
func SummarizeLeft() BranchOption
SummarizeLeft returns the BranchOption that summarizes the branch being left — back to the common ancestor of the branch and the new branch point — into a branch_summary entry the new branch's context carries in the abandoned branch's place (ADR 0020 §6). The summarizer runs under the session's compaction chain (SummaryModel or the session's own model, the skeleton prompt) with the lock released, then the navigation and the summary land as one atomic batch.
Example ¶
SummarizeLeft leaves a branch behind as a summary: the new line's context carries it in the abandoned branch's place.
ctx := context.Background()
agent := weft.New(&recordingModel{reply: "Tried the carrier API; it rate-limits."})
st := thread.Memory()
s, _ := thread.Create(ctx, st, agent)
now := time.Now().UTC()
if err := st.Append(ctx, s.ID(),
thread.MessageEntry{ID: "e_ask", Created: now, Message: weft.User("Where is order 1234?")},
thread.MessageEntry{ID: "e_try", ParentID: "e_ask", Created: now, Message: weft.Assistant("Trying the carrier API…")},
thread.MessageEntry{ID: "e_fail", ParentID: "e_try", Created: now, Message: weft.Assistant("The carrier API is rate-limited.")},
); err != nil {
fmt.Println(err)
return
}
// The entries went in behind the Session's back, and it is the
// session's writer since Create: close it, and open one that
// sees them.
_ = s.Close(ctx)
s, _ = thread.Open(ctx, st, s.ID(), agent)
// Back to the question, summarizing the detour on the way out.
if err := s.Branch(ctx, "e_ask", thread.SummarizeLeft()); err != nil {
fmt.Println(err)
return
}
for _, m := range s.Context() {
fmt.Printf("%s: %q\n", m.Role, m.Text())
}
Output: user: "Where is order 1234?" user: "<weft-summary>\nTried the carrier API; it rate-limits.\n</weft-summary>"
type BranchSummaryEntry ¶
type BranchSummaryEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Summary string `json:"summary"`
FromEntry string `json:"from_entry"`
}
BranchSummaryEntry summarizes the branch a Session.Branch leaves behind (ADR 0020 §6): Summary is the text and FromEntry is the entry the abandoned branch grew from, so the context shows the summary in place of that branch's messages. Branch summaries share no cache prefix with the main line; that cost is documented, not hidden.
func (BranchSummaryEntry) MarshalJSON ¶
func (e BranchSummaryEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator.
type CompactOption ¶
type CompactOption interface {
// contains filtered or unexported methods
}
CompactOption configures one compaction call.
func SummaryInstructions ¶
func SummaryInstructions(text string) CompactOption
SummaryInstructions returns the CompactOption appending per-call instructions to this compaction's summary prompt: s.Compact(ctx, thread.SummaryInstructions("focus on the API design")). They guide the summarizer for this one call; the agent's own instructions (core.Instructions) are untouched.
type Compaction ¶
type Compaction struct {
// Summary is the summarizer's text, shown in the context behind
// the fixed marker from the entry's first kept entry onward. Empty
// only on a trim.
Summary string
// FirstKept is the id of the first entry the context still shows
// raw after this compaction.
FirstKept string
// TokensBefore is the estimated size of the context the model was
// shown at compaction time — the compacted view (the previous
// summary, pinned entries, the kept tail with trims applied), not
// the raw path.
TokensBefore int64
// Reason is why it ran: manual, threshold, overflow, from_hook, or
// trim. ApplyCompaction records an empty Reason as manual.
Reason Reason
// SummarizerUsage and SummarizerModel record what the summary
// cost and which model made it — the cost ledger's inputs.
SummarizerUsage core.Usage
SummarizerModel core.ModelInfo
// FilesRead lists the file URLs the summarized range carried,
// sorted and deduplicated.
FilesRead []string
// Pinned lists the pinned entry ids below FirstKept: the context
// re-includes each of them, raw, after the summary.
Pinned []string
// RangeHash is the SHA-256 of the serialized range the summarizer
// was fed — the audit that a summary summarizes exactly this.
RangeHash string
// Trim is the trim record: required when Reason is trim, forbidden
// otherwise. The session's Trimmer pre-pass builds it; see
// TrimRecord.
Trim *TrimRecord
// contains filtered or unexported fields
}
Compaction is one computed compaction: the plan PreviewCompaction returns, ApplyCompaction writes as a compaction entry, and Compact does both. It carries everything the entry records (ADR 0020 §1) — nothing else: applying writes exactly these fields, nothing is recomputed.
type CompactionEntry ¶
type CompactionEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Summary string `json:"summary,omitempty"`
FirstKept string `json:"first_kept"`
// TokensBefore is the estimated size of the context the model was
// shown when the compaction ran — the compacted view, not the raw
// path.
TokensBefore int64 `json:"tokens_before"`
Reason Reason `json:"reason,omitempty"`
// SummarizerUsage and SummarizerModel name what the summary cost
// and which model made it (the cost ledger, ADR 0020 §4) — absent
// on a trim, which summarizes nothing.
SummarizerUsage core.Usage `json:"summarizer_usage,omitzero"`
SummarizerModel core.ModelInfo `json:"summarizer_model,omitzero"`
FilesRead []string `json:"files_read,omitempty"`
// FilesModified is a format-1 wire field kept readable: no build
// writes it (the sandbox write log it was reserved for was
// abandoned), and a file that carries it round-trips unchanged.
FilesModified []string `json:"files_modified,omitempty"`
Pinned []string `json:"pinned,omitempty"`
RangeHash string `json:"range_hash,omitempty"`
// Trim is the trim record: present exactly on a trim, absent on a
// summary compaction. An entry carrying one is written with "v":5
// — a reader that does not know the record would replay the trim
// wrongly, so it must fail loudly instead (ADR 0011 §6).
Trim *TrimRecord `json:"trim,omitempty"`
}
CompactionEntry records one compaction (ADR 0020 §1): the summary text, the id of the first entry kept raw after it, the estimated size of the model's context before it, the reason it ran, the summarizer's cost and identity, the details (files read, pinned entries kept through), and a hash of the summarized range. Nothing is deleted — the summarized entries stay in the file, and the context at a leaf is the latest summary compaction's summary on the path, the entries it pinned, then the entries from its first kept id onward.
A trim that needed no summary is the same entry with an empty Summary, Reason "trim" and a Trim record naming exactly what was stubbed. A trim never moves the boundary: the latest summary compaction below it keeps governing the context, the trim's stubs layer over the kept range, and its FirstKept only repeats the boundary in force when it landed.
func (CompactionEntry) MarshalJSON ¶
func (e CompactionEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator. A summary compaction is a format-1 line and carries no "v"; an entry holding a trim record carries "v":5, its minimum reader version.
type Compactor ¶
type Compactor interface {
Compact(ctx context.Context, p Preparation) (*Compaction, error)
}
Compactor is the whole compaction algorithm, swapped in whole: it receives the Preparation (after BeforeCompact's edits) and returns the Compaction to write. The default serializes, summarizes and hashes; a replacement may ignore the suggested cut and keep any FirstKept on the leaf's path past the previous boundary. A returned Compaction with no Reason takes the Preparation's. It runs without the session lock and may call the session; a panic in it fails the compaction.
type CorruptError ¶
type CorruptError struct {
Session string
Line int // 1-based; 0 when the failure is not one line's
Entry string // the entry id the failure names; empty when it has none
Err error
}
CorruptError is the typed shape ErrCorrupt takes when the failure belongs to one place in a stored session (ADR 0011 §5: "ErrCorrupt naming the line"): it carries the session, the 1-based line number — the header is line 1 — and, when the failure is an entry the tree cannot hold (Open's validation: an empty, invalid or duplicate id, a parent the session does not hold before it), that entry's id, so a caller, a log, or a UI can point at the place. Match the class with errors.Is(err, ErrCorrupt), take the place with errors.As, and reach the cause through Err or errors.Is — Unwrap exposes both. Backends construct it directly through their own decode-failure paths.
func (*CorruptError) Error ¶
func (e *CorruptError) Error() string
Error names the session and the place — the line, the entry, or neither — and then the cause.
func (*CorruptError) Unwrap ¶
func (e *CorruptError) Unwrap() []error
Unwrap returns the class and the cause: errors.Is(err, ErrCorrupt) is true for every CorruptError, and errors.Is and errors.As also reach whatever Err wraps — a JSON syntax error, a sentinel of the backend's own.
type CustomEntry ¶
type CustomEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Kind string `json:"kind"`
Data json.RawMessage `json:"data,omitempty"`
}
CustomEntry carries application state: a caller-chosen Kind and opaque JSON Data. It survives every compaction (ADR 0020 §4) and never enters the model's context. A nil Data writes no "data" key — omitempty keeps nil and absent the same value both ways, so a round trip through the wire never turns "no data" into "null".
func (CustomEntry) MarshalJSON ¶
func (e CustomEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator.
type CustomMessageEntry ¶
type CustomMessageEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Kind string `json:"kind"`
Message core.Message `json:"message"`
}
CustomMessageEntry carries an application message: a caller-chosen Kind and a core.Message that is always in the model's context — how an application injects a note the model must see without attributing it to the user.
func (CustomMessageEntry) MarshalJSON ¶
func (e CustomMessageEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator.
type Decision ¶
type Decision struct {
// CallID is the pending call the decision addresses.
CallID string
// Kind is the outcome: approve, deny, resolve or resolve_error.
Kind Outcome
// Reason is the deny reason — the text the model sees after
// "DENIED: ".
Reason string
// Content is the resolve payload, verbatim; resolve_error marks it
// an error.
Content string
// Who decided, for the audit trail. Through Decide it is a
// declaration — whatever the caller wrote, verified by nobody —
// and it is the identity Quorum counts for unsigned decisions.
// Through DecideSigned it is a label only: the signing key is the
// identity.
Who string
// Via names the channel the decision arrived by. The session fills
// it on every path it records — Decide always records "user",
// DecideSigned "signed", whatever this field holds; only an
// Approver's answer keeps a Via of its own ("approver" when
// empty).
Via string
// Always approves and grants the same thing for the future: the
// tool plus the call's exact arguments become a session grant
// ("approve and always allow this", ADR 0021 §4), minted when the
// call's effective verdict is approve — never beside a denial.
// Build it with ApproveAlways.
Always bool
}
Decision is one call's outcome, the value Decide records (ADR 0021 §1). Build one with Approve, Deny, Resolve or ResolveError; Who is the caller's to fill.
func Approve ¶
Approve returns a Decision that resumes the call: the run executes it through the ordinary tool chain with Call.Approved set (ADR 0007).
func ApproveAlways ¶
ApproveAlways returns a Decision that approves the call and grants the same thing for the future (ADR 0021 §4's "approve and always allow this command"): the tool plus the call's exact arguments become a session grant — so the next such call never parks. The grant is minted only when the call's effective verdict is approve: in the same append as the decision when one approval resolves the call, and with the approval that completes the quorum under Quorum — a call that ends denied leaves no grant behind. A richer grant (a command glob, a path prefix) is s.Grant's to make.
func Deny ¶
Deny returns a Decision that resolves the call without running it: the model sees "DENIED: <reason>" and the loop continues.
func Resolve ¶
Resolve returns a Decision that resumes the call with a result computed outside the process — the human-as-tool-executor shape. The handler never runs; the content becomes the call's result verbatim.
func ResolveError ¶
ResolveError is Resolve with the result marked as an error.
type Entry ¶
type Entry interface {
// contains filtered or unexported methods
}
Entry is the sealed set of session entry kinds (ADR 0011 §2): a session is an append-only tree of these, and the leaf is the entry the next one attaches to. External types cannot join, so switches over entries stay exhaustively lintable; weft adds kinds additively under the "v" rule below. On the wire every entry carries a "type" discriminator and UnmarshalEntry restores it — the same rule and the same compatibility contract as the core's events and message parts (ADR 0001, ADR 0004).
Every entry, of every kind, carries an ID, a ParentID (empty for the session's root entry) and a Created time. Whether a kind's content enters the model's context is part of the kind's contract: message and custom_message do; the bookkeeping kinds never do; compaction and branch_summary contribute their summary text.
func UnmarshalEntry ¶
UnmarshalEntry decodes one entry line, dispatching on its "type" discriminator. The loud rules of the format (ADR 0011 §5–§6): an entry kind this build does not know fails with ErrNewerFormat — it was written by a newer weft and is never silently skipped; a known kind carrying a "v" above the version this build reads of it fails the same way; every other unknown key on the line is additive and ignored.
type Estimator ¶
Estimator estimates the token weight of messages — the trigger's delta, the cut's walk and the entry's TokensBefore. The default is a quarter of the wire bytes; a provider-aware implementation can do better. Estimate returns the total for msgs, in tokens, the unit every other number here uses; it is called with one message when the walk weighs entries one by one and with a batch otherwise. It runs without the session lock and may call the session.
type Flusher ¶
Flusher is the optional Storage capability that completes buffered durability work — the small-interface rule (ADR 0011 §5: capabilities are discovered by type assertion, so Storage never grows). With FsyncOnFlush, an append writes without fsyncing and Flush makes the session durable; with the default policy Flush is a no-op that still returns any error the storage holds for the session. Flushing a session the storage does not hold fails with ErrNotFound; a session another writer holds but this one never buffered flushes nothing.
type Grant ¶
type Grant struct {
// Tool is the tool name the grant matches, exactly.
Tool string `json:"tool"`
// Args are the argument predicates; every one must match. Empty
// matches any arguments, which is the tool-wide grant Crush
// taught the field to want more than (ADR 0021 §4's rejected).
Args []Arg `json:"args,omitempty"`
// Deny makes the grant a standing refusal: the chain denies the
// call at once, with Reason the model sees.
Deny bool `json:"deny,omitempty"`
// Reason is the refusal's text on a deny-grant — the model-visible
// bytes, "denied by grant" when empty — or the record's note on an
// approval grant.
Reason string `json:"reason,omitempty"`
// Expiry is when the grant lapses; zero means never. Expired
// strictly after, the requests' rule.
Expiry time.Time `json:"expiry,omitzero"`
// MaxUses bounds how many times a session-scoped grant may match
// (its audit entries count the uses); 0 means unlimited. A shared
// grant's uses are the store's own business — the session cannot
// see other sessions' matches.
MaxUses int `json:"max_uses,omitempty"`
}
A Grant approves — or, with Deny, refuses — future calls without asking (ADR 0021 §4): one tool, and argument predicates every one of which the call's arguments must satisfy, so "allow run_command for go test …" is expressible, not just "allow run_command". Scope, expiry, uses and revocation are the grant's lifetime: a session-scoped grant is an entry (s.Grant), a shared grant lives behind the GrantStore interface; an expiry passes; MaxUses bounds the times a session-scoped grant may match (counted from its audit trail); s.Revoke ends one with an entry.
type GrantEntry ¶
type GrantEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Grant
}
GrantEntry is a session-scoped grant made durable (ADR 0021 §4): a standing approval — or, with Deny, a standing refusal — for future calls of one tool whose arguments match every predicate. It is the decision chain's first step: a matching live grant decides at once, audited, including the automatic approval. Liveness is derived, never stored: a GrantRevokedEntry naming the grant ends it, an Expiry passes, or its MaxUses is reached — uses counted from the audit entries the chain writes when it matches (their GrantID). It never enters the model's context; the model sees a grant only through the result of the call it allowed or refused.
func (GrantEntry) MarshalJSON ¶
func (e GrantEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator and "v":2.
type GrantRevokedEntry ¶
type GrantRevokedEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
GrantID string `json:"grant_id"`
}
GrantRevokedEntry ends a grant (ADR 0021 §4): revocation is an append, never a rewrite — the grant entry stays, the walk reads the revocation after it, and the audit trail keeps both.
func (GrantRevokedEntry) MarshalJSON ¶
func (e GrantRevokedEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator and "v":2.
type GrantStore ¶
type GrantStore interface {
// Grants returns the store's live grants, in the order they are
// tried (first match wins). The chain calls it once for every call
// about to park that no session grant matched — so once per gated
// call per turn, never on a timer — on the session's runner
// goroutine, with the turn's persistence context, holding no
// session lock. The turn waits for the answer: a slow store is a
// slow turn, so answer from memory or bound the lookup with ctx.
// The returned slice is read, never kept or changed. An error
// reads as no match, with a warning through the agent's logger — a
// broken store never parks nothing silently, and never decides
// either.
Grants(ctx context.Context) ([]SharedGrant, error)
}
A GrantStore is the application-wide grant scope (ADR 0021 §4): grants that outlive any one session, consulted by every session handed the store. The store owns their lifetime — expiry, revocation, use counts — because only it can see every session's matches; the session's entry-scoped grants keep their own rules, and a SharedGrant's Expiry and MaxUses are the store's to enforce: the session matches whatever Grants returns.
An implementation must be safe for concurrent use: every session handed the store calls it from its own runner goroutine, and sessions run concurrently.
type Header ¶
type Header struct {
// Weft is the envelope integer. The zero value writes and means
// FormatVersion; Create rejects anything else, so a caller never
// has to know the number to construct a valid header.
Weft int `json:"weft"`
ID string `json:"id"`
Created time.Time `json:"created"`
Parent *ParentRef `json:"parent,omitempty"`
Lineage *Lineage `json:"lineage,omitempty"`
Meta map[string]string `json:"meta,omitempty"`
}
Header is a session's first line: its identity, its creation time, its fork origin and caller metadata. A Fork names the session and entry it grew from in Parent (ADR 0011 §3), so the new session is self-contained but traceable; beyond that the header is immutable — title and metadata edits are info entries appended to the file, never rewrites of it.
Example ¶
A session's header is its first line: the weft envelope integer, the session's id and creation time, and — for a fork — where it came from (ADR 0011 §3, §6).
package main
import (
"encoding/json"
"fmt"
"time"
"github.com/weftgo/weft/thread"
)
func main() {
h := thread.Header{
ID: "s_01J8X9M2K7QW4R5N8T6V2B3C4Q",
Created: time.Date(2026, 9, 28, 12, 0, 0, 123456789, time.UTC),
Parent: &thread.ParentRef{Session: "s_01J8X9M2K7QW4R5N8T6V2B3C4D", Entry: "e_01J8X9M2K7QW4R5N8T6V2B3C4G"},
}
b, err := json.Marshal(h)
if err != nil {
fmt.Println(err)
return
}
fmt.Println(string(b))
}
Output: {"type":"session","weft":1,"id":"s_01J8X9M2K7QW4R5N8T6V2B3C4Q","created":"2026-09-28T12:00:00.123456789Z","parent":{"session":"s_01J8X9M2K7QW4R5N8T6V2B3C4D","entry":"e_01J8X9M2K7QW4R5N8T6V2B3C4G"}}
func (Header) MarshalJSON ¶
MarshalJSON encodes the header with its "type":"session" discriminator. A zero Weft writes FormatVersion — the zero Header is a valid, current-format header.
func (*Header) UnmarshalJSON ¶
UnmarshalJSON decodes a session header and enforces the envelope rule (ADR 0011 §6): a header from a newer format fails with ErrNewerFormat instead of being read with this build's tags, and a line without the session shape was never a thread session's first line. Unknown keys are additive and ignored.
type InfoEntry ¶
type InfoEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Title string `json:"title,omitempty"`
Meta map[string]string `json:"meta,omitempty"`
}
InfoEntry edits the session's title and metadata as an append, never a rewrite: the current title is the last info entry's Title, and the current metadata is every info entry's Meta merged in order.
func (InfoEntry) MarshalJSON ¶
MarshalJSON encodes the entry with its "type" discriminator.
type Key ¶
Key is one key in a Keyring: an id, the secret bytes (any non-empty length; HMAC-SHA256 derives its key from them), and whether it is the ring's active key — the one challenges are minted under. At most one Key in a ring may be active. A Key is also what an approver holds on the signing side: Sign signs a decision as this key.
func (Key) Sign ¶
func (k Key) Sign(req Request, d Decision) (SignedDecision, error)
Sign signs d over the challenge req carries as this key — the approver's half of the exchange when each approver holds a key of their own: the signature names k.ID, whichever key the challenge was minted under, so under Quorum it counts as this approver and no other. The session verifying it must hold k in its ring. req must be a challenge, the value Session.Request returned: Sign fails on one without a nonce, and on a key with no id or no secret.
type Keyring ¶
type Keyring struct {
// contains filtered or unexported fields
}
A Keyring holds HMAC keys by id (ADR 0021 §3). Every key verifies; at most one is active, and the active key is the one new challenges are minted under — rotation without invalidating requests in flight: the new active key signs new challenges while the old keys vouch for decisions still arriving. A ring with no active key is verify-only: it checks decisions over challenges minted earlier and mints none, so Session.Request fails on it and RequireSigned refuses it. Keys never touch the session file; only key ids do, so a leaked file cannot forge a decision.
One key is one approver: under Quorum a signed approval counts as its key, whatever Who it carries — give each approver their own key and the session a ring holding all of them.
Construct one with NewKeyring and hand it to the sessions that verify signed decisions (WithKeyring). The signing side signs with what it holds: an approver with a key of their own signs with Key.Sign; a process holding the whole ring signs with Keyring.Sign under the key the challenge names; SignDecision is the same under a bare secret. A Keyring is immutable and safe for concurrent use.
func NewKeyring ¶
NewKeyring builds a ring from keys: at least one key, ids unique and non-empty, secrets non-empty, at most one key active. A ring with no active key is verify-only (see Keyring). The ring is immutable from here — rotation is a new ring handed to new sessions, holding the old keys so the requests they minted keep verifying.
func (*Keyring) Sign ¶
func (r *Keyring) Sign(req Request, d Decision) (SignedDecision, error)
Sign signs d over the challenge req carries — the signing side's half of the exchange for a caller that holds the ring: the key is the one req.KeyID names (the key the challenge was minted under), looked up here, so no caller maps key ids to secrets by hand. req must be a challenge, the value Session.Request returned (a Pending view carries no nonce): Sign fails on one without a nonce, and with ErrUnknownKey when the ring does not hold the challenge's key — the ring's holder may know which keys it has; DecideSigned tells a remote caller nothing of the kind.
type LabelEntry ¶
type LabelEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Entry string `json:"entry"`
Name string `json:"name"`
}
LabelEntry names an entry — bookmarks, checkpoints, a UI's anchors. The label never enters the model's context.
func (LabelEntry) MarshalJSON ¶
func (e LabelEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator.
type LeafEntry ¶
type LeafEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Entry string `json:"entry"`
}
LeafEntry moves the session's leaf to an existing entry — branch navigation without a rewrite (ADR 0011 §3). The next entry attaches to Entry, and the model's context is rebuilt from there.
func (LeafEntry) MarshalJSON ¶
MarshalJSON encodes the entry with its "type" discriminator.
type Leaser ¶
type Leaser interface {
// Acquire takes the session's writer lease for holder and reports
// how many complete entry lines the storage holds for the session
// at that moment: every line a Load accounts for — the entries it
// returns and the lines it skips under Salvage — and never a torn
// tail. A writer that knows how many it has loaded and written
// compares: a different number means the session changed behind
// its view (ErrStale is the Session's answer).
//
// created is the stored header's Created — which session this is,
// not only how long: a session deleted and created again under
// the same id carries a new one, so a writer that loaded the old
// session is told apart even when the two happen to hold the same
// number of entries. It is the zero time when the storage cannot
// read the header's (a header line this build does not decode);
// the writer then has only the count to go by.
//
// Acquire takes the backend's cross-instance lock exactly as a
// first Append does — ErrLocked when another Storage or process
// holds the session, a torn tail repaired — and then the lease:
// ErrLocked naming the session when a different holder has it on
// this Storage value. For the holder that already has it Acquire
// is idempotent and cheap: no I/O. A session the storage does not
// hold fails with ErrNotFound, a nil holder with a plain error,
// and nothing is taken either way.
//
// The lease ends when its holder yields, when the session is
// released or deleted through this Storage value, or with the
// process.
Acquire(ctx context.Context, session string, holder any) (entries int, created time.Time, err error)
// Yield ends holder's lease and the storage's hold with it, as
// Release does. It is Release for a writer that must not let go
// of what is not its own: when a different holder has the lease,
// Yield releases nothing and returns nil. When no holder has it —
// the session is held by this Storage value's direct use, or not
// held — Yield is Release. A session the storage does not hold
// fails with ErrNotFound.
Yield(ctx context.Context, session string, holder any) error
}
Leaser is the optional Storage capability that makes the one-writer rule hold between Session values sharing one Storage value — the same small-interface rule. The backend's own lock tells Storage values and processes apart; it cannot tell two Sessions on one Storage value apart, because Append names a session and not who is writing. A lease can: the writer names itself with holder, an opaque comparable token — a pointer the writer owns — and the storage remembers which holder has the session.
A Session takes the lease before every write, so the first write is what makes it the session's writer, and it stays so until its Close yields. Reading never takes it: Load, List and Watch — and so Open and every read of a Session — are not refused by a lease and do not stand in a writer's way.
The lease is bookkeeping over the backend's lock, not a second lock: the Storage methods do not consult it. A Storage value used directly, without Sessions, enforces one writer per instance; Sessions enforce one writer per Session.
type Lineage ¶
type Lineage struct {
Session string `json:"parent_session"`
Call string `json:"parent_call_id,omitempty"`
}
Lineage names a pool child's origin (ADR 0022 §3): the parent session and, for a wrapped delegation, the call that delegated. It is a reference, not a copy — unlike Parent, a fork's self-contained path, the child's file holds only its own entries, and the link is how a nested approval finds the session that must resume the child.
type LoadReport ¶
type LoadReport struct {
// Torn is the 1-based line number of the incomplete final line the
// load dropped — bytes after the last newline, left by a writer cut
// mid-write — or 0 when the session ended on a complete line.
Torn int
// Skipped lists the 1-based line numbers of malformed lines the
// load skipped under Salvage (the open option), in file
// order. Without Salvage, a malformed line fails the load with
// ErrCorrupt instead.
Skipped []int
}
LoadReport names what a load had to drop or skip to return a session — a repair is never silent (ADR 0011 §5). A clean load returns a nil report.
type MessageEntry ¶
type MessageEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Message core.Message `json:"message"`
// RunID is set on a turn's prompt entry: the run id minted for the
// turn the message starts (<session>-t<n>), written before the run
// does anything — so the id is on the record even when the turn's
// end never lands, and a reopened session never mints it again.
// Empty on every other message entry.
RunID string `json:"run_id,omitempty"`
}
MessageEntry is one core.Message in the transcript, embedded with the core's message wire (ADR 0001) verbatim. It is how conversation content is stored, and it is always in the model's context.
Example ¶
A session entry marshals with its "type" discriminator — one JSON line per entry is the on-disk format (ADR 0011 §2) — and UnmarshalEntry restores the sealed type.
package main
import (
"encoding/json"
"fmt"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
)
func main() {
e := thread.MessageEntry{
ID: "e_01J8X9M2K7QW4R5N8T6V2B3C4E",
Created: time.Date(2026, 9, 28, 12, 0, 0, 123456789, time.UTC),
Message: weft.User("Where is order 1234?"),
}
b, err := json.Marshal(e)
if err != nil {
fmt.Println(err)
return
}
fmt.Println(string(b))
restored, err := thread.UnmarshalEntry(b)
if err != nil {
fmt.Println(err)
return
}
fmt.Println(restored.(thread.MessageEntry).Message.Role)
}
Output: {"type":"message","id":"e_01J8X9M2K7QW4R5N8T6V2B3C4E","created":"2026-09-28T12:00:00.123456789Z","message":{"role":"user","content":[{"type":"text","text":"Where is order 1234?"}]}} user
func (MessageEntry) MarshalJSON ¶
func (e MessageEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator.
type MirroredRequest ¶
type MirroredRequest struct {
// Request is the mirror: CallID in the namespaced form decisions
// address, Child the session the call is parked in, RunID the
// child run that parked it — the occurrence's name.
Request Request
// Wrapper is the parent-side delegating call; empty for an async
// delegation.
Wrapper string
// Decisions are the decisions in force for this occurrence, in
// the order they were recorded.
Decisions []ApprovalDecisionEntry
// Decided reports whether those decisions reach an effective
// verdict under this session's Quorum — the moment the pool
// replays them into the child.
Decided bool
}
MirroredRequest is one mirrored child request as the parent session holds it (ADR 0022 §7): the request, the delegating call it parks under, and what this session has decided about it.
type NativeCompactor ¶
type NativeCompactor interface {
CompactNative(ctx context.Context, req core.ModelRequest, instructions string) (core.Message, core.Usage, error)
}
NativeCompactor is implemented by adapters whose provider compacts server-side. Root types only, so an adapter implements it without importing thread (ADR 0020 §7).
type OpenOption ¶
type OpenOption interface {
// contains filtered or unexported methods
}
OpenOption configures a Storage backend at open, one value per concern, applied over the defaults. The options live here — in the package that owns the Storage contract — so every backend accepts the same vocabulary and a caller never learns a backend to say Salvage (ADR 0011 §5 names it thread.Salvage). Memory, jsonl.Open and sqlite.Open all take them; a backend an option means nothing to accepts it and says so. The set is sealed: the options are the ones this package returns. Backend authors resolve them with thread/backend.Resolve.
func FsyncEveryAppend ¶
func FsyncEveryAppend() OpenOption
FsyncEveryAppend returns the open option that restores the default durability: every Append fsyncs before returning, so an accepted entry is durable before anything replies on it — the rule behind "the prompt is durable before the run starts" (ADR 0011 §4).
func FsyncOnFlush ¶
func FsyncOnFlush() OpenOption
FsyncOnFlush returns the open option that defers the fsync to the Flusher capability — the turn-end cadence, cheaper than an fsync per append. A crash between appends and the flush can lose the tail of a turn, never a synced one; the file stays readable (a torn final line is dropped and reported). The caller who flushes at turn ends — a Session does — owns the durability window; a plain Flush after the prompt keeps ADR 0011 §4's promise.
func NoLock ¶
func NoLock() OpenOption
NoLock returns the open option that opens a file backend without its cross-process writer lock. It exists for platforms with no advisory file lock to take — there jsonl.Open fails unless NoLock says the caller knows — and for filesystems whose locks cannot be trusted. With it, one-writer-per-session (ADR 0011 §5) is the caller's promise instead of the backend's check: goroutines of one Storage are still serialized, but a second Storage or a second process writing the same session is not refused with ErrLocked and can interleave its lines with the first's. Backends whose lock needs no platform support (sqlite's lock row, Memory) accept the option and keep locking.
func OpenLogger ¶
func OpenLogger(l *slog.Logger) OpenOption
OpenLogger returns the open option that names where a backend reports what it repairs on its own: a torn tail a crashed writer left, truncated before the next append; a dead holder's lock, taken over. Each is one Warn line carrying the session id. The default is slog.Default() as it stands at open; a nil logger keeps the default. What a load had to drop is still the LoadReport's to say — the logger covers the write path, where no report is returned.
func Salvage ¶
func Salvage() OpenOption
Salvage returns the open option that skips malformed lines instead of failing the load: each skip is reported in Load's LoadReport (Skipped, in file order). A torn final line is always dropped and reported — that is a crash, not damage — and data from a newer weft still fails with ErrNewerFormat: salvage repairs what a crash wrote, never what it cannot read (ADR 0011 §5).
Example ¶
Salvage loads a session whose file has a damaged line instead of refusing it, and LoadReport says exactly what that cost: the line skipped, and the entries left without the parent it held. They are kept — the context starts at the orphan.
package main
import (
"context"
"errors"
"fmt"
"os"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/thread/jsonl"
"github.com/weftgo/weft/wefttest"
)
// exampleIDs mints the given ids in order — the IDs option's usual
// shape in an example: the session id first, then one per entry.
func exampleIDs(ids ...string) thread.SessionOption {
next := 0
return thread.IDs(func() string {
id := ids[next]
next++
return id
})
}
func main() {
ctx := context.Background()
dir, err := os.MkdirTemp("", "weft-salvage-example-")
if err != nil {
fmt.Println(err)
return
}
defer func() { _ = os.RemoveAll(dir) }()
agent := weft.New(wefttest.Script())
at := time.Date(2026, 10, 1, 9, 0, 0, 0, time.UTC)
// A session of three messages, the middle line destroyed on disk.
st, err := jsonl.Open(dir)
if err != nil {
fmt.Println(err)
return
}
if _, err := thread.Create(ctx, st, agent, exampleIDs("s_damaged")); err != nil {
fmt.Println(err)
return
}
damage := st.(interface {
Inject(ctx context.Context, session string, data []byte) error
})
_ = st.Append(ctx, "s_damaged", thread.MessageEntry{ID: "e_1", Created: at, Message: weft.User("one")})
_ = damage.Inject(ctx, "s_damaged", []byte("{\"type\":\"message\",\"id\":\"e_2\",#!\n"))
_ = st.Append(ctx, "s_damaged", thread.MessageEntry{ID: "e_3", ParentID: "e_2", Created: at, Message: weft.User("three")})
_, err = thread.Open(ctx, st, "s_damaged", agent)
fmt.Println("without Salvage:", errors.Is(err, thread.ErrCorrupt))
salvaging, err := jsonl.Open(dir, thread.Salvage())
if err != nil {
fmt.Println(err)
return
}
s, err := thread.Open(ctx, salvaging, "s_damaged", agent)
if err != nil {
fmt.Println(err)
return
}
report := s.LoadReport()
fmt.Println("skipped lines:", report.Skipped)
fmt.Println("orphaned entries:", report.Orphaned)
fmt.Println("entries kept:", len(s.Entries()))
for _, m := range s.Context() {
fmt.Println("context:", m.Text())
}
}
Output: without Salvage: true skipped lines: [3] orphaned entries: [e_3] entries kept: 2 context: three
type OpenReport ¶
type OpenReport struct {
LoadReport
// Orphaned lists, in append order, the ids of the entries whose
// link names an entry the file does not hold — possible only when
// a line was skipped under Salvage, which is the only time Open
// accepts it (otherwise the same shape fails with ErrCorrupt).
//
// An orphaned entry is kept, never dropped: it stays in Entries,
// its turns stay in Usage, and it is the root of what survives of
// its line — Path and Context walk back to it and end there. So a
// session whose leaf descends from an orphaned entry has a context
// that starts at that entry: everything recorded before the
// skipped line is no longer on the path. An orphaned leaf entry —
// a navigation whose target was skipped — is ignored: the leaf
// stays where it was before the navigation.
Orphaned []string
}
OpenReport names what Open had to drop, skip or accept to return a session — a repair is never silent (ADR 0011 §5). It is the storage's LoadReport (the torn final line, the lines skipped under Salvage) plus what the tree validation found among the entries that survived. Session.LoadReport returns it; a clean load has none.
type Outcome ¶
type Outcome string
Outcome is a decision's kind (ADR 0021 §1): approve (run the call through the ordinary chain), deny with a reason, or resolve with content computed outside the process, resolve_error marking it an error. The wire values are the decision entry's "outcome" field: stored bytes, so a value this build writes is one it keeps reading (the format goldens pin them); a decision entry carrying any other value never resolves a call.
const ( // OutcomeApprove runs the call: the resume executes it through the // ordinary tool chain with Call.Approved set (ADR 0007). OutcomeApprove Outcome = "approve" // OutcomeDeny resolves the call without running it: the model sees // "DENIED: " followed by the decision's Reason. OutcomeDeny Outcome = "deny" // OutcomeResolve resolves the call with the decision's Content as // its result, verbatim; the handler never runs. OutcomeResolve Outcome = "resolve" // OutcomeResolveError is OutcomeResolve with the result marked as // an error. OutcomeResolveError Outcome = "resolve_error" )
type Page ¶
type Page struct {
// Sessions is the page, newest first by Created (ties by id,
// descending), headers only — Load returns the entries.
Sessions []Header
// Total is the number of sessions matching the query's filters,
// ignoring the cursor (Before, BeforeID) and Limit — the "how many
// pages are there" number.
Total int
}
Page is one List result.
func List ¶
List returns a page of session headers from st — headers only, never entries; see Query for the cursor and the limit. Listing needs no Session: it is the storage's List, nil storage aside.
Example ¶
List pages session headers, newest first — headers only, never entries.
package main
import (
"context"
"fmt"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// exampleIDs mints the given ids in order — the IDs option's usual
// shape in an example: the session id first, then one per entry.
func exampleIDs(ids ...string) thread.SessionOption {
next := 0
return thread.IDs(func() string {
id := ids[next]
next++
return id
})
}
// exampleClock ticks one second per reading from a fixed instant, so
// every stored timestamp in an example is known.
func exampleClock() thread.SessionOption {
now := time.Date(2026, 10, 1, 9, 0, 0, 0, time.UTC)
return thread.Clock(func() time.Time {
now = now.Add(time.Second)
return now
})
}
func main() {
ctx := context.Background()
st := thread.Memory()
agent := weft.New(wefttest.Script())
clock := exampleClock()
for _, id := range []string{"s_monday", "s_tuesday", "s_wednesday"} {
if _, err := thread.Create(ctx, st, agent, exampleIDs(id), clock); err != nil {
fmt.Println(err)
return
}
}
page, err := thread.List(ctx, st, thread.Query{Limit: 2})
if err != nil {
fmt.Println(err)
return
}
fmt.Println("total:", page.Total)
for _, h := range page.Sessions {
fmt.Println(h.ID, h.Created.Format(time.TimeOnly))
}
// The next page starts before the last header seen.
next, err := thread.List(ctx, st, thread.Query{Limit: 2, Before: page.Sessions[1].Created})
if err != nil {
fmt.Println(err)
return
}
for _, h := range next.Sessions {
fmt.Println(h.ID, h.Created.Format(time.TimeOnly))
}
}
Output: total: 3 s_wednesday 09:00:03 s_tuesday 09:00:02 s_monday 09:00:01
type ParentRef ¶
ParentRef names a fork's origin: the session the new session was forked from, and the entry its copied path ends at (ADR 0011 §3).
type Policy ¶
type Policy int
Policy is what a Send does when the session is already running a turn — the busy policy (ADR 0011 §4). Queue is the default.
The constant Queue names a policy; the method Session.Queue lists what is waiting. They share a word, not a meaning.
const ( // Queue holds the follow-up and runs it when the current turn // ends, in acceptance order: the send is accepted — an accepted // receipt entry (ReceiptAccepted) carrying the message is appended // and flushed, so the send survives a crash — its Turn is returned // at once, and its prompt entry is written when its turn starts. // The default. Queue Policy = iota // Reject fails the send with ErrBusy: one run per session at a // time, and a busy session says so instead of holding work. Reject // Steer delivers the send into the running turn instead of // waiting for it (ADR 0019): the message is accepted at once — // its queued receipt entry durable, flushed — and handed to the // run's next steering drain point: after the step's tool batch, // every call paired with its result, or at what would have been // the final step, which the delivery redirects into one more // step. A steer that meets an intended end (StopWhen) or an open // approval boundary is never drained: it becomes a deferred // follow-up that runs as the next turn, the receipt recording // the fate. The Send's Turn is the receipt: it ends when the steer // reaches its final state — Turn.Outcome says which: delivered, // deferred (Turn.Next is the follow-up turn) or dropped — with a // nil result. A Steer on a session that is not busy has nothing to // steer into: it runs as a plain turn. Steer // Interrupt cancels the running turn and runs the message next: // the in-flight run's context is canceled, the calls its partial // transcript left without a result carry the golden interruption // text, an approval boundary that holds the session is denied with // the interrupted reason, and the message runs as the next turn. // The interrupted turn's entries stay on the tree — evidence, // never deleted. A boundary holding nested approvals (thread/pool, // ADR 0022 §7) is denied whole — the delegating call and the // child's mirrored requests: the pool that made the child replays // the denial into it and the child runs to its end, its answer on // its receipt; in a process whose pool has not taken the session // over yet (after a restart, before Recover) the child stays // parked until it does. Interrupt // Rollback is an Interrupt that also branches the leaf back to // before the interrupted turn's receipt entry: the follow-up runs // as though the interrupted turn never happened, while its // entries keep their own line of the tree (nothing lost). Rollback )
type PoolReceiptEntry ¶
type PoolReceiptEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Receipt string `json:"receipt,omitempty"`
Status string `json:"status"`
Child string `json:"child,omitempty"`
Call string `json:"call,omitempty"`
Prompt string `json:"prompt,omitempty"`
Stop string `json:"stop,omitempty"`
Usage core.Usage `json:"usage,omitzero"`
}
PoolReceiptEntry is the pool receipt (ADR 0022 §4): the journey of one delegation a thread/pool started for this session. One entry records acceptance — Status "accepted", the child session on Child, the task on Prompt — and every later one links back to it by Receipt: "running" when a slot is acquired and the child's run starts, "parked" when that run ends at an approval boundary (the two alternate, once per park and resume), and exactly one settlement — "done" (Stop carries the child's answer), "failed" (Stop the cause), "canceled" (Cancel, or the pool's Close), or "capped" (the child died on a budget — MaxSteps or a usage limit). A settlement's Usage is the child session's whole cost — every run it made for the delegation, the ones before a park included. Call names the delegating tool call for wrapped delegations. Pool receipt entries never enter the model's context: the answer reaches the model as the delegating call's result (a sync delegation) or however the application delivers it (an async one); the entry is the ledger, not the channel.
func (PoolReceiptEntry) MarshalJSON ¶
func (e PoolReceiptEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator and "v":4.
type Preparation ¶
type Preparation struct {
// Reason is why the compaction runs.
Reason Reason
// Context is the context the model is shown now — the compacted
// view at the leaf: the previous summary, pinned entries, the kept
// tail with recorded trims applied.
Context []core.Message
// Messages is the range this compaction summarizes: the entries
// from the previous boundary up to FirstKept, as stored. A split
// turn's prefix is part of it (one pass, ADR 0020's 2026-09-29
// amendment).
Messages []core.Message
// PrevSummary is the iterative chain's last link: the latest
// summary compaction's text on the path, trims skipped.
PrevSummary string
// FirstKept is the id of the first entry kept raw.
FirstKept string
// TokensBefore is the estimated size of Context.
TokensBefore int64
// Instructions is the per-call text (SummaryInstructions).
Instructions string
// Pinned lists the pinned entry ids below the cut that the context
// keeps showing after the summary.
Pinned []string
}
Preparation is a computed compaction on its way to becoming one: what the algorithm, the BeforeCompact hook and a custom Compactor see.
The hook may edit four fields, and the compaction proceeds with the edited values: Messages (the range to summarize — redact here), Instructions, Pinned, and FirstKept. An edited FirstKept must name a user or assistant message on the leaf's path, past the previous compaction's boundary, that does not split a tool call from its result; when the hook moves it and leaves Messages alone, the range is rebuilt for the new boundary. Reason, Context, PrevSummary and TokensBefore are the session's facts: edits to them are ignored.
type Query ¶
type Query struct {
// Before and BeforeID are the paging cursor, a keyset over List's
// own order (Created descending, ties by ID descending): pass the
// last session of the previous page — its Created as Before, its
// ID as BeforeID — and the next page starts right after it, however
// many sessions share that creation time. A zero Before means start
// at the newest (BeforeID is then ignored). Before alone, with an
// empty BeforeID, returns only sessions created strictly before it:
// correct when creation times are distinct, but it skips the rest
// of a group of sessions sharing the cursor's time — set BeforeID
// to walk through ties. Offsets are deliberately absent: they drift
// under concurrent inserts (ADR 0010 §0.1's rule, inherited here).
Before time.Time
BeforeID string
// Limit caps the page: 0 means 50, a negative value reads as 0
// (50), and values above 500 clamp to 500.
Limit int
// Meta filters by session metadata: every key must be present and
// match its value exactly. A session matches when its header's
// create-time Meta holds every pair — what WithMeta (and PublicID,
// its sugar) set at Create, which Load returns as is. No backend merges the info
// entries' Meta into this view (that merged view is Session.Meta's,
// the runs' runMetadata), so a key a later SetInfo added never
// matches here. Backends answer this from the header alone — the
// cheap path; the title filter below is the one that can cost more.
Meta map[string]string
// TitleSearch filters by the session's current title — the last
// info entry carrying a non-empty Title, what Session.Title
// returns — matching case-insensitively as a substring. A title is
// entry state, not header state, so this filter is the one shape of
// List that may read beyond headers (a backend without a title
// column scans the session's info entries): opt-in by the query,
// priced accordingly, and never paid by a query without it.
TitleSearch string
}
Query selects sessions for List. The zero value lists every session, newest first, 50 at a time.
Example ¶
Query filters the session list: metadata pairs match the header exactly, and the title search matches the session's current title — the last info entry's — case-insensitively as a substring. Total counts the matches; the cursor and limit page them.
package main
import (
"context"
"fmt"
"time"
"github.com/weftgo/weft/thread"
)
func main() {
ctx := context.Background()
st := thread.Memory()
created := time.Date(2026, 9, 29, 12, 0, 0, 0, time.UTC)
titled := func(id string, meta map[string]string, title string) {
h := thread.Header{ID: id, Created: created, Meta: meta}
if err := st.Create(ctx, h); err != nil {
fmt.Println(err)
return
}
if title != "" {
if err := st.Append(ctx, id, thread.InfoEntry{
ID: "e_" + id, Created: created, Title: title,
}); err != nil {
fmt.Println(err)
}
}
}
titled("s_prod", map[string]string{"env": "prod"}, "Checkout bug")
titled("s_dev", map[string]string{"env": "dev"}, "login flow")
titled("s_prod_2", map[string]string{"env": "prod"}, "checkout again")
p, err := st.List(ctx, thread.Query{Meta: map[string]string{"env": "prod"}})
if err != nil {
fmt.Println(err)
return
}
fmt.Println("env=prod:", p.Total)
p, err = st.List(ctx, thread.Query{TitleSearch: "CHECKOUT"})
if err != nil {
fmt.Println(err)
return
}
for _, h := range p.Sessions {
fmt.Println("title match:", h.ID)
}
}
Output: env=prod: 2 title match: s_prod_2 title match: s_prod
Example (KeysetCursor) ¶
Page through every session with the keyset cursor: hand List the last session of the previous page — its Created as Before, its ID as BeforeID — and the next page starts right after it, even when sessions share a creation time. Before alone would skip the rest of such a group.
package main
import (
"context"
"fmt"
"time"
"github.com/weftgo/weft/thread"
)
func main() {
ctx := context.Background()
st := thread.Memory()
imported := time.Date(2026, 9, 29, 9, 0, 0, 0, time.UTC)
for _, id := range []string{"s_a", "s_b", "s_c", "s_d", "s_e"} {
// A bulk import: five sessions, one creation time.
if err := st.Create(ctx, thread.Header{ID: id, Created: imported}); err != nil {
fmt.Println(err)
return
}
}
q := thread.Query{Limit: 2}
for {
page, err := thread.List(ctx, st, q)
if err != nil {
fmt.Println(err)
return
}
if len(page.Sessions) == 0 {
break
}
for _, h := range page.Sessions {
fmt.Print(h.ID, " ")
}
fmt.Println("of", page.Total)
last := page.Sessions[len(page.Sessions)-1]
q.Before, q.BeforeID = last.Created, last.ID
}
}
Output: s_e s_d of 5 s_c s_b of 5 s_a of 5
type QueuedSteer ¶
type QueuedSteer struct {
// Receipt is the id the message's Turn reports as its own (Turn.ID):
// for a steer, its queued receipt entry; for a queued send, the id
// its prompt entry takes when its turn starts — the accepted
// receipt entry names it in ReceiptEntry.Turn.
Receipt string
// Msg is the message held.
Msg core.Message
// Policy says how it waits: Steer — held for the running turn's
// next drain point — or Queue — an accepted send (or a deferred
// steer's follow-up) waiting for a turn of its own.
Policy Policy
}
QueuedSteer is one message the session has accepted and not yet given to a model, as Queue reports it.
type Reason ¶
type Reason string
Reason is why a compaction ran (ADR 0020 §1) — the value a CompactionEntry records on the wire.
const ( // ReasonManual marks a compaction the caller asked for: Compact, // or ApplyCompaction over a previewed plan. ReasonManual Reason = "manual" // ReasonThreshold marks a compaction the configured trigger // started: the measured context crossed the window minus the // reserve (ADR 0020 §2). ReasonThreshold Reason = "threshold" // ReasonTrim marks a trim: a trimmer pre-pass brought the context // under the line, so no summary was made — the entry's Summary is // empty. ReasonTrim Reason = "trim" // ReasonFromHook marks a compaction whose plan a BeforeCompact // hook replaced. ReasonFromHook Reason = "from_hook" // ReasonOverflow marks the compaction that follows a turn failing // with core.ErrContextOverflow, before the turn's one re-run // (ADR 0020 §5). ReasonOverflow Reason = "overflow" )
The reasons a compaction entry records.
type ReceiptEntry ¶
type ReceiptEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
Receipt string `json:"receipt,omitempty"`
Status string `json:"status"`
Msg *core.Message `json:"msg,omitempty"`
RunID string `json:"run_id,omitempty"`
Turn string `json:"turn,omitempty"`
// Unanswered, on a delivered receipt, marks a steer the run took
// and never answered (ADR 0019 §5): the run failed — a budget, a
// model error, a cancellation — before another model step
// completed. The message is in that run's recorded transcript, so
// the next turn's model sees it; no reply to it exists.
Unanswered bool `json:"unanswered,omitempty"`
}
ReceiptEntry is the receipt of a message accepted while the session was busy (ADR 0019, ADR 0011 §4): the durable record that the session took it, and of what became of it. Receipt entries never enter the model's context: the message reaches the model through the run that delivered it or the turn that ran it, exactly once.
A steer — a Send under the Steer policy — is two entries. One records acceptance: Status "queued", the message on Msg. A second, linked by Receipt, records the fate: "delivered" (the running run's steering drain took it; RunID names the run, the message landed in that run's transcript, and Unanswered says the run ended before the model answered it), "deferred" (it runs as the next turn instead — a StopWhen end, an open approval boundary, or the run ended before the drain; Turn names the follow-up turn's prompt entry), or "dropped" (ClearQueue).
A queued send — a Send that waits for a turn of its own: the Queue policy on a busy session, an interrupting send, a deferred steer's follow-up — is one entry: Status "accepted", the message on Msg, Turn the id its prompt entry will take, RunID the run id minted for it. It is settled by that prompt entry landing (the turn started), or by a "dropped" entry linked by Receipt (ClearQueue, a refused interrupt). An accepted receipt with neither is a send the writer never got to: Open restores it to the queue.
func (ReceiptEntry) MarshalJSON ¶
func (e ReceiptEntry) MarshalJSON() ([]byte, error)
MarshalJSON encodes the entry with its "type" discriminator and "v":3.
type Releaser ¶
Releaser is the optional Storage capability that ends this writer's hold on a session — the same small-interface rule. A backend takes the one-writer lock on a session's first write and, without Release, keeps it for the life of the process (and, on jsonl, an open file with it). Release flushes what the writer buffered and lets go: from then on another Storage — in this process or another — may write the session. The hold is a lease the writer renews by writing: a later Append through the releasing Storage re-acquires the lock, and fails with ErrLocked if another writer took the session in between. Releasing a session this Storage does not hold is a no-op; one the storage does not hold at all fails with ErrNotFound. Release must not race the session's own Append — it is the last call of a writer that is done. Release speaks for the whole Storage value: it ends a Leaser lease whoever holds it, which is why a Session — one writer among possibly several on the value — closes through Yield instead when the backend offers it.
Example ¶
Release a session when its writer is done with it: a backend that locks (jsonl, sqlite) lets go of the session so another Storage or process may write it; one that does not (Memory) only checks the session exists. The capability is discovered by type assertion.
package main
import (
"context"
"fmt"
"time"
"github.com/weftgo/weft/thread"
)
func main() {
ctx := context.Background()
st := thread.Memory()
if err := st.Create(ctx, thread.Header{ID: "s_done", Created: time.Date(2026, 9, 29, 9, 0, 0, 0, time.UTC)}); err != nil {
fmt.Println(err)
return
}
if r, ok := st.(thread.Releaser); ok {
fmt.Println("released:", r.Release(ctx, "s_done"))
fmt.Println("unknown session:", r.Release(ctx, "s_never"))
}
}
Output: released: <nil> unknown session: thread: session not found: s_never
type Request ¶
type Request struct {
// ID is the request entry's id: the name of this one occurrence of
// the call. Call ids repeat across turns; this never does, and a
// signed decision is bound to it. Empty in the two places no entry
// exists yet or at all — the Request an Approver is handed (the
// chain asks before it persists), and a dangling call with no
// request entry.
ID string
// Session is the session the request lives in.
Session string
// CallID is the pending call's id — the key decisions address and
// the core's Approve/Deny/Resolve key on.
CallID string
// Tool is the call's tool name.
Tool string
// Args is the call's arguments, verbatim as the model wrote them.
Args json.RawMessage
// ArgsSHA256 is the hex SHA-256 of Args — a decision names what it
// decided by this hash, and a signed decision carries it in the
// challenge.
ArgsSHA256 string
// RunID is the run that parked the call.
RunID string
// Reason is why the call parked: "tool requires approval" for a
// RequireApproval tool, "middleware required approval" when the
// tool chain parked it.
Reason string
// Expiry is when the request lapses — zero means never. A request
// is expired strictly after this time, never at it: from then on
// it takes no decision and is denied with the stated reason (ADR
// 0021 §5).
Expiry time.Time
// Created is when the request was persisted.
Created time.Time
// Nonce is the challenge a signed decision answers (ADR 0021 §3):
// minted by Session.Request, single-use, bound to this request
// entry — and empty on the Pending view, which mints no challenge.
Nonce string
// KeyID names the keyring key the challenge was minted under — the
// ring's active key, and the key Keyring.Sign and SignDecision sign
// with. An approver holding a key of their own signs with Key.Sign
// instead, which names that key. Set by Session.Request beside
// Nonce; empty on the Pending view.
KeyID string
// Child names the delegated session this call parks in, when the
// call is a pool child's mirrored onto this parent (ADR 0022 §7):
// the lineage a decision routes by — decide it through the pool,
// which resumes the child and then completes the delegation. Empty
// on an ordinary request.
Child string
}
Request is a parked call awaiting a decision (ADR 0021 §1): what the model asked for, hashed and named so a decision can state exactly what it decided. Session.Pending returns the undecided ones; they survive restarts because they are entries.
type SendOption ¶
type SendOption interface {
// contains filtered or unexported methods
}
SendOption configures one Send: RunOptions carries extra run options into the turn's run, As overrides the busy policy.
func As ¶
func As(p Policy) SendOption
As returns the SendOption overriding the session's busy policy for this one Send (ADR 0019): As(Steer) steers a message into the running turn on a Queue session; As(Queue) holds a message for the next turn on a Steer session. The policy is captured when Send is called, with the turn's other settings, and recorded on the entry of the turn the send runs as (TurnEntry.Policy). On a session that is not busy every policy does the same thing — the send runs as a plain turn: there is nothing to steer into, interrupt or queue behind.
Example ¶
As steers one Send into a running turn on a Queue session: the message is accepted at once — its queued receipt durable — and the run's next drain point delivers it after the tool batch, the receipt recording the fate (ADR 0019).
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
release := make(chan struct{})
wait := weft.Tool("wait", "blocks until released", func(ctx context.Context, _ struct{}) (string, error) {
select {
case <-release:
return "ok", nil
case <-ctx.Done():
return "", ctx.Err()
}
})
model := wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "wait"}),
wefttest.Say("done, in metric units"),
)
var s *thread.Session
agent := weft.New(model, wait, weft.Tap(func(_ context.Context, ev weft.Event) {
if _, ok := ev.(weft.ToolStart); ok {
if _, err := s.Send(ctx, weft.User("use metric units"), thread.As(thread.Steer)); err != nil {
fmt.Println(err)
}
}
}))
s, _ = thread.Create(ctx, thread.Memory(), agent)
t, _ := s.Send(ctx, weft.User("convert this"))
close(release)
res, err := t.Wait()
if err != nil {
fmt.Println(err)
return
}
fmt.Println(res.Text())
for _, e := range s.Entries() {
r, ok := e.(thread.ReceiptEntry)
if !ok {
continue
}
if r.Receipt == "" {
fmt.Println("receipt queued:", r.Msg.Text()) // acceptance, durable
continue
}
fmt.Println("receipt", r.Status) // the fate, linked by Receipt
}
}
Output: done, in metric units receipt queued: use metric units receipt delivered
func RunOptions ¶
func RunOptions(opts ...core.RunOption) SendOption
RunOptions returns the SendOption carrying extra core.RunOptions into this turn's run — budgets, taps, thinking, metadata. Three things are the session's own, and an option that would set one is rejected by Send with an error wrapping core.ErrInvalidRunOption:
- the transcript and the run id — core.Messages, core.Prompt, core.RunID: Send builds the transcript from the session's tree and mints <session>-t<n>; either option would quietly detach the run from the tree (ADR 0011 §4);
- the steering source — core.Steering: the session owns the steer queue (Send under the Steer policy);
- approval decisions — core.Approve, core.Deny, core.Resolve, core.ResolveError: a Send never runs while an approval boundary is open (it queues behind it), so a decision passed here would reach no parked call. Decisions are recorded with Session.Decide and applied by the boundary's resume run (ADR 0021 §1).
The options bind every run the session starts on the turn's behalf: the resume run of a boundary the turn parked (after Decide or Resume), its overflow re-run, and the follow-up turn a steer aimed at it becomes when the steer cannot join the run (it meets the approval boundary or a StopWhen end, or arrives after the last drain point) — the follow-up runs under these options, the steer's own after them. A steer that finds no turn in flight and no boundary open runs as a plain turn under its own options. The same holds for the values on the turn's context (a ParkAllExcept list or metadata a delegating run handed down): those runs see every value their own context lacks. Options and context are process state — a reopened session's restored sends and boundaries run without them.
Example ¶
RunOptions carries per-run configuration into one turn's run — metadata here. What the session owns cannot be overridden: the transcript, the run id, the steering source and approval decisions are refused before anything is written.
package main
import (
"context"
"errors"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
agent := weft.New(wefttest.Script(wefttest.Say("Hello, Acme.")),
weft.Tap(func(ctx context.Context, ev weft.Event) {
if _, ok := ev.(weft.RunStart); ok {
fmt.Println("run for tenant:", weft.MetadataFromContext(ctx)["tenant"])
}
}))
s, _ := thread.Create(ctx, thread.Memory(), agent)
_, err := s.Send(ctx, weft.User("Hi."), thread.RunOptions(weft.Prompt("something else")))
fmt.Println("refused:", errors.Is(err, weft.ErrInvalidRunOption))
turn, err := s.Send(ctx, weft.User("Hi."),
thread.RunOptions(weft.Metadata(map[string]string{"tenant": "acme"})))
if err != nil {
fmt.Println(err)
return
}
res, _ := turn.Wait()
fmt.Println(res.Text())
}
Output: refused: true run for tenant: acme Hello, Acme.
type Session ¶
type Session struct {
// contains filtered or unexported fields
}
A Session is a loaded conversation: the append-only entry tree (ADR 0011 §2) held in memory, every write going through the storage's Append, the leaf tracked as the entry the next one attaches to. Create, Open and Session.Fork return one.
Concurrency ¶
A Session is safe for concurrent use by any number of goroutines: one mutex guards the tree, and every method takes it.
One writer. The Session is its session's only writer: entries appended to the storage behind its back are invisible to it until the next Open. Backends refuse a second writer from another process or another Storage value with ErrLocked. Between Session values on one Storage value the rule is a lease (the Leaser capability, which Memory, jsonl and sqlite implement): a Session takes it with its first write — Create is one — and holds it until Close. While it is held, every write of another Session on that Storage value fails with ErrLocked and changes nothing, in its tree or in the storage. Reading takes no lease: any number of Sessions may be open on a session to read it, and Open never fails for a writer elsewhere.
A Session whose view has fallen behind does not write either: when the storage holds entries the Session never loaded — another Session wrote and closed since this one was opened — its write fails with ErrStale instead of attaching to a leaf the session has moved past; so does a write to a session that was deleted and created again under the same id since. The check compares two things the lease reports — the stored header's Created and the number of entry lines — not contents. Open the session again. On a backend without the Leaser capability neither check exists, and a second Session on the same Storage value is the caller's to avoid.
One run at a time. A Send while a turn runs — or while an approval boundary is open — follows the busy policy captured at that Send: Queue holds it and runs it next, in acceptance order; Reject fails it with ErrBusy; Steer delivers it into the running turn at its next drain point, or defers it to a follow-up turn; Interrupt cancels the running turn and runs the message next; Rollback also branches back to before the interrupted turn. Concurrent Sends are serialised by the lock: the order they acquire it in is the acceptance order. Branch fails with ErrBusy while a turn is in flight; Fork, the reads, and the bookkeeping writes (Label, SetInfo, Custom, CustomMessage, Grant, Decide) do not wait for a turn, and a write made during one lands on the line the turn is extending.
Callbacks. The session calls the caller's hooks — an Approver, OnRequest, the compaction hooks and Summarizer, an Estimator, the agent's own taps and tools — without holding its lock, so they may call back into the Session. The two exceptions are the IDs and Clock functions, which run under the lock and must not.
The lock is held across the storage's Append and Flush: a write — and under FsyncEveryAppend that is an fsync — stalls every other method, reads included, until it returns. That is what makes a returned write durable and the tree in memory equal to the file.
Close. Close stops new work at once — Send and Continue fail with ErrClosed — then waits for the running turn and the queue to drain, seals the session so that every later write fails with ErrClosed, and releases the storage's hold on the session. Reads keep answering from the tree the session held. See Close for what a canceled wait leaves behind.
Delete. Delete through the Storage value a Session writes with removes the stored session whether or not a Session is open on it, lease or no lease (through another Storage value it fails with ErrLocked while the writer holds the session). An open Session keeps its in-memory tree; its next write fails with ErrNotFound, and a turn running at that moment loses its remaining entries (the failure is logged through the agent's logger). Close a session before deleting it.
func Create ¶
func Create(ctx context.Context, st Storage, agent *core.Agent, opts ...SessionOption) (*Session, error)
Create starts a new session in st: a fresh header under a new time-sortable id (or the IDs option's), stamped with the session's clock, carrying the WithMeta, PublicID and WithLineage options; the tree is empty. agent is the session's own — Send runs it and compaction summarizes with its model — and must not be nil, like the core's New. The header is all that is written: nothing else lands in the storage until the first append. An id the storage already holds fails with ErrExists. Create is the new Session's first write: it holds the session's writer lease from here (see Session, One writer).
Example ¶
A session carries the conversation: application state survives outside the model's context, application messages ride inside it, and the title and metadata are editable appends. The IDs option pins deterministic ids so the output is stable.
package main
import (
"context"
"encoding/json"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
agent := weft.New(wefttest.Script()) // no run happens here: the example only appends
st := thread.Memory()
next := 0
ids := []string{"s_demo1", "e_demo2", "e_demo3", "e_demo4", "e_demo5"}
deterministic := thread.IDs(func() string { id := ids[next]; next++; return id })
s, err := thread.Create(ctx, st, agent, deterministic)
if err != nil {
fmt.Println(err)
return
}
fmt.Println("session", s.ID())
if err := s.SetInfo(ctx, "Order support", map[string]string{"team": "ops"}); err != nil {
fmt.Println(err)
return
}
if err := s.Custom(ctx, "cart", json.RawMessage(`{"items":2}`)); err != nil {
fmt.Println(err)
return
}
if err := s.CustomMessage(ctx, "note", weft.User("The customer's quote covers two items.")); err != nil {
fmt.Println(err)
return
}
again, err := thread.Open(ctx, st, s.ID(), agent, thread.IDs(func() string { return "e_demo6" }))
if err != nil {
fmt.Println(err)
return
}
fmt.Println(again.Title(), again.Meta()["team"])
for _, m := range again.Context() {
fmt.Println(m.Role, ":", m.Text())
}
}
Output: session s_demo1 Order support ops user : The customer's quote covers two items.
func Open ¶
func Open(ctx context.Context, st Storage, id string, agent *core.Agent, opts ...SessionOption) (*Session, error)
Open loads an existing session from st, positioned at its leaf: the entry the next append attaches to, recovered by replaying the entries in append order (a leaf entry redirects; any other entry leaves the leaf at itself).
Open validates the tree as it indexes it, and a file that fails is refused with a *CorruptError (errors.Is ErrCorrupt) naming the line and the entry: every entry has a non-empty id that passes ValidID and appears once; every non-empty parent names an entry earlier in the file — append-only order makes that rule exclude cycles and parents the file does not hold — and every leaf entry navigates to the root or to an earlier entry. Several roots are legal: Branch to the root starts a new one. Nothing is repaired by guessing; a context is never quietly cut short where a link is broken.
A load that had to drop a torn tail or skip a salvaged line still opens — the entries that survived are the session — and says so twice: one warning through the agent's logger, and the report Session.LoadReport returns. Under Salvage an entry orphaned by a skipped line is kept and listed (OpenReport.Orphaned) rather than refused.
Open reads and nothing else: it writes no entry and starts no run. Input the file shows accepted but never settled — the writer stopped between the two — is restored to the queue (Queue lists it) and waits there: a steer is delivered by the next turn the session runs, a queued send runs ahead of the next Send, Continue runs either now, ClearQueue drops them.
The header options (WithMeta, PublicID, WithLineage) fail Open with ErrCreateOnly: the stored header is what the session has.
Open reads and takes nothing: it succeeds while another Session — or another process — is the session's writer, and stands in no writer's way. The Session it returns becomes the writer with its first write, which fails with ErrLocked while another holds the session and with ErrStale once the session has moved past what this Open loaded (see Session, One writer).
Example ¶
Open loads a stored session at its leaf. It reads and nothing else: no entry is written and no run starts, whatever the file holds.
package main
import (
"context"
"errors"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// exampleIDs mints the given ids in order — the IDs option's usual
// shape in an example: the session id first, then one per entry.
func exampleIDs(ids ...string) thread.SessionOption {
next := 0
return thread.IDs(func() string {
id := ids[next]
next++
return id
})
}
func main() {
ctx := context.Background()
st := thread.Memory()
model := wefttest.Script(wefttest.Say("Order 1234 shipped Tuesday."))
agent := weft.New(model)
s, err := thread.Create(ctx, st, agent, exampleIDs("s_orders", "e_prompt", "e_reply", "e_turn"))
if err != nil {
fmt.Println(err)
return
}
turn, err := s.Send(ctx, weft.User("Where is order 1234?"))
if err != nil {
fmt.Println(err)
return
}
if _, err := turn.Wait(); err != nil {
fmt.Println(err)
return
}
if err := s.Close(ctx); err != nil { // one Session per session id
fmt.Println(err)
return
}
again, err := thread.Open(ctx, st, "s_orders", agent)
if err != nil {
fmt.Println(err)
return
}
fmt.Println("leaf:", again.Leaf())
for _, m := range again.Context() {
fmt.Printf("%s: %s\n", m.Role, m.Text())
}
fmt.Println("model calls:", len(model.Requests()))
_, err = thread.Open(ctx, st, "s_missing", agent)
fmt.Println(errors.Is(err, thread.ErrNotFound))
}
Output: leaf: e_turn user: Where is order 1234? assistant: Order 1234 shipped Tuesday. model calls: 1 true
func SessionFromContext ¶
SessionFromContext returns the session whose run ctx carries it, or nil outside a session's run — a wrapped tool invoked through a bare Generate has no parent session, and the pool falls back to the ordinary subagent path (ADR 0022 §2).
func (*Session) AppendApprovalRequests ¶
func (s *Session) AppendApprovalRequests(ctx context.Context, reqs ...ApprovalRequestEntry) ([]ApprovalRequestEntry, error)
AppendApprovalRequests appends mirrored approval requests in one atomic batch (ADR 0022 §7): the pool writes a child session's parked calls onto the parent's tree — Child naming the session they park in, Wrapper the delegating call they park under — so the parent's Pending surfaces them and a decision records like any other. Every request must name its Child: an entry without one would be an ordinary request with no parked call behind it. The entries' tree fields are minted here, each id vetted like every other the session mints (valid, and new to the tree and to the batch); the stored entries return. Mirrors are ledger until decided: they never join the model's context, and their resolution is the pool's to route. A mirror stops being pending when it is decided — by an approver, by its expiry, or by DenyMirrored when its delegation ends — or when a later mirror reuses its call id.
func (*Session) AppendPoolReceipt ¶
func (s *Session) AppendPoolReceipt(ctx context.Context, e PoolReceiptEntry) (PoolReceiptEntry, error)
AppendPoolReceipt appends one pool receipt entry (ADR 0022 §4) and returns it as stored, its minted ID the receipt handle a later entry links back to with Receipt. The pool calls this for every state its delegations pass through — acceptance, each start, each park, the settlement — under the rule the entry kind's contract states: pool receipts are ledger, never model context, and a child's answer reaches the model only through its delegating call's result or the application. Session.Usage sums the Usage of settlement entries into its Delegated bucket, one per entry: writing exactly one settlement per receipt is the caller's contract.
func (*Session) ApplyCompaction ¶
func (s *Session) ApplyCompaction(ctx context.Context, c *Compaction) error
ApplyCompaction writes a computed Compaction as the session's next compaction entry, appended at the leaf, and then runs the AfterCompact hook with the entry. The entry lands whatever else happened between its computation and this call — the tree only grew, so the kept range stays correct — provided the plan still fits the leaf's path:
- c.FirstKept must name an entry on the leaf's path (ErrNoEntry when the session does not hold it, ErrInvalidCompaction when it sits on another branch), at or after the previous summary compaction's boundary — an iterative compaction never summarizes what a summary already replaced (ADR 0020 §1).
- A summary compaction naming the very boundary the previous one left adds nothing and is refused with ErrNothingToCompact.
- A compaction without a summary must be a trim (Reason trim), and a trim must carry its TrimRecord, every stub naming a tool result on the path; a summary compaction must carry none.
Like Compact, it is a between-turns operation: ErrBusy while a turn is in flight, ErrAwaitingApproval while approval requests are pending. Nothing is written on any error; a validation or storage failure also reaches the CompactFailed hook. An empty c.Reason is recorded as manual.
Example ¶
ApplyCompaction validates a plan against the leaf's path: the same boundary twice is refused — errors.Is tells the refusals apart.
package main
import (
"context"
"errors"
"fmt"
"strings"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// longSession creates a session whose history is three long user
// messages (ids e_0, e_1, e_2) — more than the default KeepRecent
// keeps raw — opened on agent with opts.
func longSession(agent *weft.Agent, opts ...thread.SessionOption) *thread.Session {
ctx := context.Background()
st := thread.Memory()
s, err := thread.Create(ctx, st, agent, opts...)
if err != nil {
panic(err)
}
parent := ""
var batch []thread.Entry
for i, topic := range []string{"order ", "invoice ", "refund "} {
id := fmt.Sprintf("e_%d", i)
batch = append(batch, thread.MessageEntry{ID: id, ParentID: parent, Created: time.Now().UTC(),
Message: weft.User(strings.Repeat(topic, 30_000/len(topic)))})
parent = id
}
if err := st.Append(ctx, s.ID(), batch...); err != nil {
panic(err)
}
if err := s.Close(ctx); err != nil {
panic(err)
}
s, err = thread.Open(ctx, st, s.ID(), agent, opts...)
if err != nil {
panic(err)
}
return s
}
func main() {
ctx := context.Background()
s := longSession(weft.New(wefttest.Script()))
// A hand-made plan: your own summary, keeping from e_2.
plan := &thread.Compaction{Summary: "Orders and invoices were reviewed.", FirstKept: "e_2"}
if err := s.ApplyCompaction(ctx, plan); err != nil {
fmt.Println(err)
return
}
for _, m := range s.Context() {
fmt.Println(m.Role, len(m.Text()) < 100)
}
err := s.ApplyCompaction(ctx, plan)
fmt.Println("again:", errors.Is(err, thread.ErrNothingToCompact))
err = s.ApplyCompaction(ctx, &thread.Compaction{Summary: "x", FirstKept: "e_missing"})
fmt.Println("unknown entry:", errors.Is(err, thread.ErrNoEntry))
}
Output: user true user false again: true unknown entry: true
func (*Session) Audit ¶
Audit returns the session's approval trail (ADR 0021 §5): every request, chain step, decision, grant and revocation entry, in append order, each carrying its entry id and time — the whole file, not only the leaf's path, because the trail is what happened, not what the leaf remembers. An expiry denial reads as its audit entry plus the decision with Via "expiry"; a grant match as the audit entry naming the grant in GrantID; a refused signed decision as a StepSigned audit entry with no decision beside it.
A resume reads as two StepResume audit entries: "started", written before the run starts and listing the decisions it applies, and "completed" or "failed" (Detail carrying the error), written in the same atomic append as the resume's turn entry. Each carries the run id it was written under; they agree unless the resume re-ran after a context overflow, where the second names the re-run. A started entry with no second entry after it is a resume that never finished: a crash, or a run still in flight. Turn entries are not part of the trail — the turn's ledger is Entries'.
The trail is an index of the session's log, not evidence that stands on its own: entries are plain appended lines, unsigned and unchained, so whoever can write the session's storage can add, change or remove them, and Who on an unsigned decision is whatever the caller declared. It answers "what did this session record"; tamper-evidence, where it is needed, belongs to the storage (an append-only store, a signed export).
Example ¶
Audit is the session's approval trail: every request, chain step, decision, grant and revocation, in order — and for a resume, the entry that says it started and the turn entry that says how it ended. It indexes the session's log; it is not tamper-evident on its own.
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// gatedDeploy is the agent of the approval examples: its one tool,
// deploy, needs a decision before it runs, and the scripted model
// plays turns.
func gatedDeploy(turns ...wefttest.Turn) *weft.Agent {
return weft.New(wefttest.Script(turns...),
weft.Tool("deploy", "Deploy the service.",
func(ctx context.Context, in struct {
Env string `json:"env"`
}) (string, error) {
return "deployed to " + in.Env, nil
},
weft.RequireApproval()))
}
// deployTo is a model turn asking to deploy to env.
func deployTo(env string) wefttest.Turn {
return wefttest.ToolCalls(wefttest.Call{Name: "deploy", Args: `{"env":"` + env + `"}`})
}
// settle waits for a turn and, when its boundary resumed on its own,
// for the resume too; it returns the last result.
func settle(t *thread.Turn) *weft.RunResult {
res, err := t.Wait()
if err != nil {
fmt.Println("turn failed:", err)
return nil
}
if next := t.Next(); next != nil {
return settle(next)
}
return res
}
func main() {
ctx := context.Background()
agent := gatedDeploy(deployTo("prod"), wefttest.Say("Understood."))
s, _ := thread.Create(ctx, thread.Memory(), agent)
t1, _ := s.Send(ctx, weft.User("Deploy to prod."))
settle(t1)
no := thread.Deny(s.Pending()[0].CallID, "change freeze")
no.Who = "avi"
resume, _ := s.Decide(ctx, no)
settle(resume)
for _, e := range s.Audit() {
switch e := e.(type) {
case thread.ApprovalRequestEntry:
fmt.Println("request:", e.Tool, "-", e.Reason)
case thread.ApprovalAuditEntry:
fmt.Println("step:", e.Step, e.Outcome)
case thread.ApprovalDecisionEntry:
fmt.Println("decision:", e.Outcome, "by", e.Who, "via", e.Via, "-", e.Reason)
}
}
}
Output: request: deploy - tool requires approval step: park parked decision: deny by avi via user - change freeze step: resume started step: resume completed
func (*Session) Branch ¶
Branch navigates the session to entryID: it appends a leaf entry (never a rewrite, ADR 0011 §3), and the next write attaches there — the model's context is rebuilt from that path alone, and nothing appended on the abandoned branch reaches it again (pi's invariant). entryID may name any entry the session holds, on any branch, including the current leaf (a recorded no-op navigation); "" is the root, restarting the conversation from nothing while the file keeps everything — the next entry is then a new root, and the tree holds several. An id the session does not hold is an error.
Branching is a between-turns operation: while a turn runs the session holds the line the run's transcript must land on, and a navigation underneath it would strand the run's messages on a branch whose context the model never saw — so a Branch while a turn is in flight fails with ErrBusy, like a Send under the Reject policy. A parked approval boundary is not a running turn: branching away from it is the documented way out of an unwanted boundary (the armed resume reads ErrNotPending). Fork is the operation that works mid-turn: it copies and never moves this session's leaf.
Example ¶
Branch navigates the tree and Fork copies a path into a new session: the branch's context drops the abandoned entries, the fork carries the whole copied path and names its origin in its header.
package main
import (
"context"
"fmt"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
agent := weft.New(wefttest.Script())
st := thread.Memory()
next := 0
ids := []string{"s_nav", "e_nav"}
s, err := thread.Create(ctx, st, agent, thread.IDs(func() string { id := ids[next]; next++; return id }))
if err != nil {
fmt.Println(err)
return
}
// Two messages of history, appended the way Send appends them.
if err := st.Append(ctx, s.ID(),
thread.MessageEntry{ID: "e_1", Created: time.Now().UTC(), Message: weft.User("draft the intro")},
thread.MessageEntry{ID: "e_2", ParentID: "e_1", Created: time.Now().UTC(), Message: weft.Assistant("done")},
); err != nil {
fmt.Println(err)
return
}
// The entries went in behind the Session's back, and it is the
// session's writer since Create: close it, and open one that
// sees them.
_ = s.Close(ctx)
s, err = thread.Open(ctx, st, s.ID(), agent, thread.IDs(func() string { id := ids[next]; next++; return id }))
if err != nil {
fmt.Println(err)
return
}
// Navigate back to the first entry: the abandoned reply drops out
// of the context, and the file keeps it.
if err := s.Branch(ctx, "e_1"); err != nil {
fmt.Println(err)
return
}
fmt.Println("branch context:", len(s.Context()))
// Fork the full path into a new session: self-contained, its
// header naming where it grew from.
f, err := s.Fork(ctx, "e_2", thread.IDs(func() string { return "s_fork" }))
if err != nil {
fmt.Println(err)
return
}
page, _ := thread.List(ctx, st, thread.Query{})
for _, h := range page.Sessions {
if h.ID == "s_fork" && h.Parent != nil {
fmt.Println("fork of", h.Parent.Session, "at", h.Parent.Entry, "carries", len(f.Context()), "messages")
}
}
}
Output: branch context: 1 fork of s_nav at e_2 carries 2 messages
func (*Session) CancelDelegated ¶
CancelDelegated denies every pending request of a pool child with reason and resumes nothing (ADR 0022 §6): the pool's Cancel of a child parked at an approval. The denials are recorded with Who "thread/pool" and Via "parent", so the parked calls can never run on a later approval; the child's run is not started — a canceled delegation does no more work — and whoever opens the session later finds a decided boundary whose resume shows the model the denial. A session with no pool lineage fails; one with nothing pending records nothing.
func (*Session) ClearQueue ¶
ClearQueue drops everything Queue lists — every steer not yet delivered and every send not yet started — marking each receipt dropped in one atomic append: the messages never reach the model, and a reopen does not restore them. Each dropped message's Turn ends with an error wrapping ErrDropped (Outcome TurnDropped). Steers already delivered or deferred are untouched — the runs that took them own them — though the follow-up a deferred steer became is a queued send like any other, and is dropped. A resume waiting for the runner is not queue: it is the boundary's resolution, and stays. ClearQueue returns how many messages it dropped.
func (*Session) Close ¶
Close quiesces the session and lets go of its storage. In order:
- New work stops at once: from the moment Close is called, Send and Continue fail with ErrClosed.
- Close waits for the session to drain — the running turn, the sends queued behind it, and any approval resume those turns arm. Sends queued behind an approval boundary nobody is going to decide cannot drain: once no turn is running their Turns end with ErrClosed. Their accepted receipts stay in the file, and the next Open restores them to the queue (Queue lists them, ClearQueue drops them).
- The session is sealed: every later write — Label, SetInfo, Branch, Decide, Compact, any of them — fails with ErrClosed. Reads (Entries, Path, Context, Pending, Usage, …) keep answering from the tree the session held, and Fork still works: it writes another session.
- The storage's hold on the session is released, when the backend offers that, so another writer can take it: the Session's own lease through the Leaser capability (Yield — a Session that never held the lease lets go of nobody else's), or the storage's hold through Releaser. Its error is Close's result.
If ctx ends while Close is waiting, Close gives up: it cancels the running turn (which records as canceled, its unanswered calls carrying the interruption text), ends every queued Turn with ErrClosed, and returns ctx.Err(). The session then stays closed to new work, and nothing can start another run — but the canceled turn may still be writing its end, so the session is not yet sealed and the storage not yet released. Call Close again to finish: it waits for that turn to land, seals, and releases. Steers still queued at that moment keep their queued receipts in the file, and the next Open restores them.
Close is idempotent and safe to call from several goroutines: once a Close has completed, every call returns what the release returned. It does not wait for thread/pool children the session delegated to — close the pool first. It must not be called from inside the session's own run (a tool, a hook): it would wait for the turn that is calling it.
Example ¶
Close quiesces a session: new Sends are refused at once, the running turn and the queue behind it finish, and then every write fails with ErrClosed while reads keep answering.
package main
import (
"context"
"errors"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
st := thread.Memory()
agent := weft.New(wefttest.Script(wefttest.Say("One."), wefttest.Say("Two.")))
s, err := thread.Create(ctx, st, agent)
if err != nil {
fmt.Println(err)
return
}
first, _ := s.Send(ctx, weft.User("first"))
second, _ := s.Send(ctx, weft.User("second")) // queued behind the first
// Close waits for both turns; a deadline on ctx bounds the wait.
fmt.Println("close:", s.Close(ctx))
for _, turn := range []*thread.Turn{first, second} {
res, err := turn.Wait()
fmt.Println(res.Text(), err)
}
_, err = s.Send(ctx, weft.User("third"))
fmt.Println(errors.Is(err, thread.ErrClosed))
fmt.Println(errors.Is(s.Label(ctx, first.ID(), "late"), thread.ErrClosed))
fmt.Println("still readable:", len(s.Context()), "messages")
fmt.Println("close again:", s.Close(ctx))
}
Output: close: <nil> One. <nil> Two. <nil> true true still readable: 4 messages close again: <nil>
Example (HandOver) ¶
One Session writes a session. A second Session on the same storage opens and reads, but its writes are refused until the writer closes — and by then the session has moved on, so the reader opens it again to write from what it now holds.
package main
import (
"context"
"errors"
"fmt"
"github.com/weftgo/weft/core"
"github.com/weftgo/weft/core/wefttest"
"github.com/weftgo/weft/thread"
)
func main() {
ctx := context.Background()
st := thread.Memory()
agent := core.New(wefttest.Script())
writer, _ := thread.Create(ctx, st, agent)
_ = writer.SetInfo(ctx, "Order 1234", nil)
reader, _ := thread.Open(ctx, st, writer.ID(), agent)
fmt.Println("the reader sees:", reader.Title())
err := reader.SetInfo(ctx, "Order 1234 (refunded)", nil)
fmt.Println("locked while the writer is open:", errors.Is(err, thread.ErrLocked))
_ = writer.SetInfo(ctx, "Order 1234 — shipped", nil)
_ = writer.Close(ctx)
err = reader.SetInfo(ctx, "Order 1234 (refunded)", nil)
fmt.Println("stale after the writer wrote and closed:", errors.Is(err, thread.ErrStale))
_ = reader.Close(ctx)
again, _ := thread.Open(ctx, st, writer.ID(), agent)
fmt.Println("reopened:", again.Title())
fmt.Println("writes:", again.SetInfo(ctx, "Order 1234 (refunded)", nil))
}
Output: the reader sees: Order 1234 locked while the writer is open: true stale after the writer wrote and closed: true reopened: Order 1234 — shipped writes: <nil>
func (*Session) Compact ¶
func (s *Session) Compact(ctx context.Context, opts ...CompactOption) error
Compact computes and applies the session's next compaction in one call: PreviewCompaction, then ApplyCompaction. Nothing is deleted — the summarized entries stay in the file, and the context at the leaf becomes the summary, the pinned entries, then the entries from the compaction's first kept entry onward.
Compaction is a between-turns operation: while a turn is in flight Compact fails with ErrBusy before any model call, and while approval requests are pending with ErrAwaitingApproval. A failure past those gates — other than ErrNothingToCompact and ErrCompactCanceled — also reaches the CompactFailed hook.
Example ¶
Compaction summarizes the old part of a session and keeps the recent part raw — nothing is deleted, and Uncompact branches right back.
ctx := context.Background()
// The same model summarizes; a script makes it deterministic.
rec := &recordingModel{reply: "Goal: ship the order service."}
agent := weft.New(rec)
st := thread.Memory()
s, _ := thread.Create(ctx, st, agent)
// A long history, appended the way Send does.
msgs := []string{strings.Repeat("order ", 12_000), strings.Repeat("invoice ", 12_000), strings.Repeat("refund ", 12_000)}
parent := ""
var batch []thread.Entry
for i, text := range msgs {
id := fmt.Sprintf("e_%d", i)
batch = append(batch, thread.MessageEntry{ID: id, ParentID: parent, Created: time.Now().UTC(), Message: weft.User(text)})
parent = id
}
if err := st.Append(ctx, s.ID(), batch...); err != nil {
fmt.Println(err)
return
}
// The entries went in behind the Session's back, and it is the
// session's writer since Create: close it, and open one that
// sees them.
_ = s.Close(ctx)
s, err := thread.Open(ctx, st, s.ID(), agent)
if err != nil {
fmt.Println(err)
return
}
if err := s.Compact(ctx); err != nil {
fmt.Println(err)
return
}
fmt.Println("after Compact:", len(s.Context()), "messages")
fmt.Println("summary rides first:", strings.HasPrefix(s.Context()[0].Text(), "<weft-summary>"))
fmt.Println("kept raw:", len(s.Entries()) > 3)
if err := s.Uncompact(ctx); err != nil {
fmt.Println(err)
return
}
fmt.Println("after Uncompact:", len(s.Context()), "messages")
Output: after Compact: 2 messages summary rides first: true kept raw: true after Uncompact: 3 messages
func (*Session) Context ¶
Context returns the messages the model sees at the session's leaf, in conversation order, with core.Repair applied last — every call the transcript shows has a result (ADR 0011 §2, ADR 0001). Entries of the bookkeeping kinds never reach it; a custom entry's whole point is to survive outside it. A call left pending by its turn is shown repaired here — the caller's view; the run Send starts repairs pending calls itself, so the decision options can resolve them.
Without a compaction on the path that is every message and custom_message entry, and every branch_summary as its marked summary. With one, it is the compacted view (ADR 0020 §1 and its 2026-10-01 amendment), in this order: the latest summary compaction's summary behind the fixed marker; the entries that compaction pinned, raw, in path order; then the entries from its first kept entry onward — with every trim record above the compaction replayed (the stubs the record names, nothing else), and signed reasoning stripped from entries recorded before the latest compaction or trim.
func (*Session) Continue ¶
Continue starts the work the session holds but is not running, and returns the first turn that will run — nil when nothing is waiting. It is the explicit counterpart of a rule Open keeps: opening a session never runs anything.
What can be waiting on an idle session:
- steers restored by Open — accepted before a crash, never settled (Queue lists them). Continue defers each to a follow-up turn, recorded on its receipt, and runs them in acceptance order under ctx;
- sends restored by Open — accepted while the session was busy, their turns never started (Queue lists them too). The Send that accepted each is gone, and its context with it: Continue runs them under ctx, so canceling ctx cancels a restored turn, as it does a restored steer's follow-up;
- sends queued behind an approval boundary that has since been cleared or decided, each still under its own Send's context;
- an approval boundary whose every call is decided but whose resume never ran (the writer died between the two): with AutoResume on, Continue arms the resume, and the queue follows it.
Without Continue the same work starts with the session's next Send, Decide or Resume: a restored steer is delivered into the turn that runs, at its first drain point, as any queued steer is. Continue is for the caller who wants the accepted input to run now, with no new message. While an undecided approval boundary holds the session the returned turn waits behind it, as a Send's would. On a session that is already running, Continue starts nothing and returns the next queued turn, if any. A closed session fails with ErrClosed.
Example ¶
Continue runs what a session holds but is not running — here a steer a crashed writer had accepted and never settled. Open restores it to the queue and runs nothing; Continue is the caller saying go.
package main
import (
"context"
"fmt"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// exampleIDs mints the given ids in order — the IDs option's usual
// shape in an example: the session id first, then one per entry.
func exampleIDs(ids ...string) thread.SessionOption {
next := 0
return thread.IDs(func() string {
id := ids[next]
next++
return id
})
}
func main() {
ctx := context.Background()
st := thread.Memory()
model := wefttest.Script(wefttest.Say("Switched to metric."))
agent := weft.New(model)
s, err := thread.Create(ctx, st, agent, exampleIDs("s_crashed"))
if err != nil {
fmt.Println(err)
return
}
// What the crash left: a queued receipt, and no entry settling it.
steer := weft.User("Use metric units.")
if err := st.Append(ctx, s.ID(), thread.ReceiptEntry{
ID: "e_steer", Created: time.Date(2026, 10, 1, 9, 0, 0, 0, time.UTC),
Status: thread.ReceiptQueued, Msg: &steer,
}); err != nil {
fmt.Println(err)
return
}
// The crashed process took its hold on the session with it.
_ = s.Close(ctx)
open, err := thread.Open(ctx, st, "s_crashed", agent, exampleIDs("e_followup", "e_deferred", "e_accepted", "e_reply", "e_turn"))
if err != nil {
fmt.Println(err)
return
}
for _, q := range open.Queue() {
fmt.Println("restored:", q.Receipt, q.Msg.Text())
}
fmt.Println("model calls after Open:", len(model.Requests()))
turn, err := open.Continue(ctx)
if err != nil {
fmt.Println(err)
return
}
res, err := turn.Wait()
if err != nil {
fmt.Println(err)
return
}
fmt.Println(res.Text())
fmt.Println("queue:", len(open.Queue()))
}
Output: restored: e_steer Use metric units. model calls after Open: 0 Switched to metric. queue: 0
func (*Session) Custom ¶
Custom appends application state: a caller-chosen Kind and opaque JSON Data, surviving every compaction and never entering the model's context (ADR 0011 §2, ADR 0020 §4). data is stored verbatim — nil stores no data key at all; bytes that are not JSON are rejected, the same failure the storage's Append would raise, caught before anything is written. An empty kind is rejected: state nobody can name is not state.
func (*Session) CustomMessage ¶
CustomMessage appends an application message: a caller-chosen Kind and a core.Message that is always in the model's context — how an application puts a note the model must see without attributing it to the user (ADR 0011 §2). An empty kind is rejected. The message is copied: the caller's value is not retained.
It is a between-turns write, like Branch and Compact: while a turn is in flight CustomMessage fails with ErrBusy and writes nothing. A message that joined the context mid-run would land between the run's own messages — between an assistant's calls and their results — where the running model never saw it and every later run would read a transcript no run produced. Append it before the Send, or after the turn (Turn.Wait); to reach a running turn, Send with the Steer policy. A parked approval boundary is not a running turn: a note written there is accepted and reads after the boundary's results. Custom, which never enters the context, is not restricted.
func (*Session) Decide ¶
Decide records decisions over the session's pending calls, durably, and resumes the boundary when they complete it (ADR 0021 §1). It is the unsigned, in-process door: every decision is recorded with Via "user" — whatever its Via field holds — and Who exactly as the caller declared it; a session under RequireSigned rejects the call with ErrSignatureRequired.
The batch is validated whole before anything is recorded, and recorded in one atomic append: no decisions, a decision without an outcome, or two decisions for one call fail with ErrInvalidDecision; a decision for a call that is not pending fails with ErrNotPending; one for a delegating call with ErrDelegated; a decision for a request past its expiry fails with ErrExpired. On any of them none of the batch is recorded and no run starts on its account.
A call that delegates to a thread/pool child (ADR 0022 §7) takes no decision: a batch naming one fails with ErrDelegated. The rule is the same on every session — a pool child included, whose parent's decisions reach it through the pool's own replay, not this door.
Expiry is resolved first, on every call: each pending request strictly past its expiry is denied on the spot with the stated reason (an expiry audit step and a decision with Via "expiry"), before the batch is looked at — which is why a decision for one is ErrExpired, never an approval of a lapsed request. Those denials stand even when Decide then returns an error, and when they complete the boundary under AutoResume its resume starts: reach it through the parked turn's Next, or Resume.
When, after recording, every pending call of the boundary has a decision and AutoResume is on (the default), the resume run starts and Decide returns its Turn; otherwise the return is (nil, nil) and the caller drives Resume. The undecided calls a Resume-driven run carries are denied with the core's "no decision" text.
Example ¶
Approvals make a parked call durable: the turn ends pending, the request survives a restart as an entry, and Decide records the decision and resumes the conversation on its own (ADR 0021).
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
agent := weft.New(
wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "refund", ID: "call_1", Args: `{"order_id":"1234"}`}),
wefttest.Say("Refund issued."),
),
weft.Tool("refund", "Refund an order.",
func(ctx context.Context, in struct {
OrderID string `json:"order_id"`
}) (string, error) {
return "refunded " + in.OrderID, nil
},
weft.RequireApproval()),
)
next := 0
ids := []string{"s_appr", "e_1", "e_2", "e_3", "e_4", "e_5", "e_6", "e_7", "e_8", "e_9", "e_10", "e_11", "e_12"}
s, _ := thread.Create(ctx, thread.Memory(), agent,
thread.IDs(func() string { id := ids[next]; next++; return id }))
t1, err := s.Send(ctx, weft.User("Refund order 1234."))
if err != nil {
fmt.Println(err)
return
}
res, err := t1.Wait()
if err != nil {
fmt.Println(err) // a pending turn is a success
return
}
fmt.Println("parked:", res.Pending[0].Name)
for _, r := range s.Pending() {
fmt.Println("pending request:", r.CallID, "on", r.Tool)
}
rt, err := s.Decide(ctx, thread.Approve("call_1"))
if err != nil {
fmt.Println(err)
return
}
if _, err := rt.Wait(); err != nil {
fmt.Println(err)
return
}
fmt.Println("pending after decide:", len(s.Pending()))
for _, m := range s.Context() {
if r, ok := lastResult(m); ok {
fmt.Println("result:", r)
}
}
}
// lastResult reports the message's last tool result text, if any.
func lastResult(m weft.Message) (string, bool) {
for i := len(m.Content) - 1; i >= 0; i-- {
if r, ok := m.Content[i].(weft.ToolResultPart); ok {
return r.Content, true
}
}
return "", false
}
Output: parked: refund pending request: call_1 on refund pending after decide: 0 result: refunded 1234
func (*Session) DecideSigned ¶
DecideSigned verifies sd fail-closed, then records it as the request's decision — with the nonce and key id, so a replay of the same signature is detectable from the file — and resumes when the boundary completes (AutoResume, like Decide). The decision is recorded with Via "signed"; its identity under Quorum is the key.
No decision is recorded until every check passes. The checks, in order, and the error each fails with (ADR 0021 §3):
- the MAC, compared in constant time over the canonical encoding of every claimed field, under the key sd.KeyID names; an outcome that is none of the four; a session other than this one; an empty nonce — ErrBadSignature. A key the ring does not hold fails the same way, with the same words: DecideSigned never says which key ids exist.
- a nonce a recorded decision already answered — ErrReplay, whatever became of the call since.
- the request past its own expiry — ErrExpired. The session enforces the request's expiry, not the one the signature states: the lapsed request is denied on the spot with the stated reason, like on every other path.
- no pending call with that id, or a pending one that is another occurrence (its request entry or run differ from the ones signed — the call id was parked again by a later run) — ErrNotPending.
- a nonce no key of the ring minted for this request, an expiry or a tool that differ from the request's — ErrBadSignature.
- an arguments hash that differs from the request's — ErrArgsChanged.
A signature that passes every check and names a call delegating to a pool child fails with ErrDelegated, unrecorded and unaudited: the call completes with its child's answer (ADR 0022 §7).
A refusal is audited: a StepSigned entry names the reason, the call and, when the ring holds it, the key — and carries no decision. Two bounds keep the log from growing at a stranger's will: a signature that fails its MAC and names no pending call leaves nothing, and at most 16 refusals are recorded per request.
Example ¶
The signed flow, for decisions that cross a process boundary: the session mints a challenge over the pending request, the signing side signs a decision over it, and DecideSigned verifies before anything is recorded. RequireSigned closes the unsigned door, durably — the session's header keeps the rule.
package main
import (
"context"
"errors"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// gatedDeploy is the agent of the approval examples: its one tool,
// deploy, needs a decision before it runs, and the scripted model
// plays turns.
func gatedDeploy(turns ...wefttest.Turn) *weft.Agent {
return weft.New(wefttest.Script(turns...),
weft.Tool("deploy", "Deploy the service.",
func(ctx context.Context, in struct {
Env string `json:"env"`
}) (string, error) {
return "deployed to " + in.Env, nil
},
weft.RequireApproval()))
}
// deployTo is a model turn asking to deploy to env.
func deployTo(env string) wefttest.Turn {
return wefttest.ToolCalls(wefttest.Call{Name: "deploy", Args: `{"env":"` + env + `"}`})
}
// settle waits for a turn and, when its boundary resumed on its own,
// for the resume too; it returns the last result.
func settle(t *thread.Turn) *weft.RunResult {
res, err := t.Wait()
if err != nil {
fmt.Println("turn failed:", err)
return nil
}
if next := t.Next(); next != nil {
return settle(next)
}
return res
}
func main() {
ctx := context.Background()
ring, err := thread.NewKeyring(thread.Key{ID: "ops-2026", Secret: []byte("a secret the session file never sees"), Active: true})
if err != nil {
fmt.Println(err)
return
}
agent := gatedDeploy(deployTo("prod"), wefttest.Say("Deployed."))
s, err := thread.Create(ctx, thread.Memory(), agent, thread.WithKeyring(ring), thread.RequireSigned())
if err != nil {
fmt.Println(err)
return
}
t1, _ := s.Send(ctx, weft.User("Deploy to prod."))
settle(t1)
call := s.Pending()[0].CallID
_, err = s.Decide(ctx, thread.Approve(call))
fmt.Println("unsigned:", errors.Is(err, thread.ErrSignatureRequired))
challenge, err := s.Request(call) // travels to the signing side
if err != nil {
fmt.Println(err)
return
}
signed, err := ring.Sign(challenge, thread.Approve(call)) // travels back
if err != nil {
fmt.Println(err)
return
}
resume, err := s.DecideSigned(ctx, signed)
if err != nil {
fmt.Println(err)
return
}
fmt.Println("signed:", settle(resume).Text())
_, err = s.DecideSigned(ctx, signed)
fmt.Println("replayed:", errors.Is(err, thread.ErrReplay))
}
Output: unsigned: true signed: Deployed. replayed: true
func (*Session) DenyMirrored ¶
DenyMirrored denies every still-pending mirrored request of child with reason (ADR 0022 §7): what the pool records when a delegation ends with requests nobody decided — a canceled child, a child that failed or is gone. Nothing can resume the child through this session any more, so its requests must stop being offered, and the record says why: a deny with Who "thread/pool" and Via "child". The denials are never replayed anywhere. When they complete this session's own boundary its resume is armed, as after Decide. A child with nothing pending records nothing.
func (*Session) Entries ¶
Entries returns the session's whole tree in append order — every entry ever written, abandoned branches included. A tree may hold several roots (entries with no parent): Branch to the root starts a new one. The slice is fresh and every entry in it is a deep copy — its maps, slices, raw JSON, and the parts of the messages it carries: mutating what comes back never reaches the session or the storage.
func (*Session) Fork ¶
func (s *Session) Fork(ctx context.Context, entryID string, opts ...SessionOption) (*Session, error)
Fork copies the session's path root → entryID into a new session in the same storage, and returns it positioned at the end of that path: the new session is self-contained — its file holds the whole copied path, the entries keeping the ids they were born with — and traceable, its header naming the session and entry it grew from (ADR 0011 §3). entryID follows Branch's rule — any held entry, or "" for the root (a fork of an empty session is an empty session with a parent) — with one exception: a leaf entry is a navigation, not a position, and Fork rejects it; fork the entry it navigates to.
Options ¶
opts configure the new session, not this one, exactly as they would at Create: IDs mints the fork's session id and its later entries, Clock stamps its header, WithMeta and PublicID fill its header metadata, WithLineage its pool lineage; the busy policy, compaction and approval options are the fork's own. Nothing is inherited from this session's options or header — not its public id, not its header metadata: a fork is another session — with one exception, the signing rule: a fork of a session that requires signed decisions (RequireSigned) requires them too, recorded in its own header, and takes this session's keyring when opts give it none. A copy must not be a way around the rule its origin's approvals were held to.
What a fork inherits ¶
Everything recorded on the copied path, because the path is copied whole:
- the conversation, its compactions and branch summaries: the fork's Context at entryID is this session's;
- the title and the SetInfo metadata — the info entries on the path (the fork starts with this session's title, and renames itself with SetInfo);
- labels, custom entries, the turn ledger (Usage counts the copied turns, and the fork's run ids continue after them), and the compaction trigger's last measurement;
- grants and their revocations, with the uses the path recorded;
- an open approval boundary: if the path ends on parked calls the fork parks on them too, with their requests and whatever decisions the path holds. Each session then resolves its own copy — a call approved in both runs in both.
And what it does not — the fork copies nothing that would start a run, and nothing that reaches into this session's work:
- the running turn: it lives in this Session value, not in the tree;
- queued sends and queued steers: a send accepted and still waiting for its turn, or a steer still waiting, on the path is recorded as dropped in the fork (one receipt entry each, appended after the copy), so the fork's file says what became of it and reopening the fork never restores it — the message belongs to this session, which runs it;
- mirrored child approval requests (thread/pool): they are handles on this session's children, so they are left out of the copy — the one case where a copied entry's parent link is rewritten, to the nearest entry the fork does hold. The fork cannot decide another session's children;
- pool delegations still unsettled on the path: each is settled in the fork as canceled, with a Stop naming the fork. A call parked on one — a sync delegation's wrapper — stays as an ordinary parked call the fork decides itself;
- in a session opened under Salvage, an orphaned entry at the head of the path is copied as the fork's root: the fork's file is whole, and loads without Salvage.
When the fork had to append such settling entries its Leaf is the last of them rather than entryID; its Context is the same either way.
During a turn ¶
Unlike Branch, Fork is allowed while a turn runs: it takes a snapshot under the session's lock and writes somewhere else, so it never disturbs the run. The snapshot is of what has been appended: forking at the running turn's latest entry copies a transcript the run is still extending, and calls the run has issued but not yet answered read in the fork as an open boundary (Pending lists them; Resume denies the undecided ones). Fork also works on a closed session — it reads this session and writes another.
The fork is durable when Fork returns. If its entries cannot be written, the half-made session is deleted again and the error returned.
Example ¶
Fork copies a path into a new, self-contained session. The fork takes its own header options and inherits what the copied path records — here the conversation and the title.
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// exampleIDs mints the given ids in order — the IDs option's usual
// shape in an example: the session id first, then one per entry.
func exampleIDs(ids ...string) thread.SessionOption {
next := 0
return thread.IDs(func() string {
id := ids[next]
next++
return id
})
}
func main() {
ctx := context.Background()
st := thread.Memory()
agent := weft.New(wefttest.Script(
wefttest.Say("Draft one."),
wefttest.Say("Draft two, in the fork."),
))
s, err := thread.Create(ctx, st, agent,
exampleIDs("s_draft", "e_prompt", "e_reply", "e_turn", "e_title"), thread.PublicID("pub-original"))
if err != nil {
fmt.Println(err)
return
}
turn, _ := s.Send(ctx, weft.User("Draft the intro."))
if _, err := turn.Wait(); err != nil {
fmt.Println(err)
return
}
if err := s.SetInfo(ctx, "Intro drafts", nil); err != nil {
fmt.Println(err)
return
}
fork, err := s.Fork(ctx, s.Leaf(),
exampleIDs("s_whatif", "e_fprompt", "e_freply", "e_fturn"), thread.PublicID("pub-whatif"))
if err != nil {
fmt.Println(err)
return
}
fmt.Println("fork:", fork.ID(), "title:", fork.Title(), "public id:", fork.Meta()["weft.public_id"])
turn, _ = fork.Send(ctx, weft.User("Try another angle."))
if _, err := turn.Wait(); err != nil {
fmt.Println(err)
return
}
fmt.Println("run:", turn.RunID())
fmt.Println("fork context:", len(fork.Context()), "messages; original:", len(s.Context()))
page, _ := thread.List(ctx, st, thread.Query{Meta: map[string]string{"weft.public_id": "pub-whatif"}})
h := page.Sessions[0]
fmt.Println("forked from:", h.Parent.Session, "at", h.Parent.Entry)
}
Output: fork: s_whatif title: Intro drafts public id: pub-whatif run: s_whatif-t2 fork context: 4 messages; original: 2 forked from: s_draft at e_title
func (*Session) Grant ¶
Grant appends a session-scoped grant (ADR 0021 §4): the decision chain consults it — newest grant first — for every call about to park, from the next boundary on. The tool name must not be empty: a grant that matches every tool is a policy this shape does not have.
Example ¶
A grant approves future calls without asking: here every deploy whose env matches a glob. The grant is an entry, so is its revocation — after it the same call parks again.
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// gatedDeploy is the agent of the approval examples: its one tool,
// deploy, needs a decision before it runs, and the scripted model
// plays turns.
func gatedDeploy(turns ...wefttest.Turn) *weft.Agent {
return weft.New(wefttest.Script(turns...),
weft.Tool("deploy", "Deploy the service.",
func(ctx context.Context, in struct {
Env string `json:"env"`
}) (string, error) {
return "deployed to " + in.Env, nil
},
weft.RequireApproval()))
}
// deployTo is a model turn asking to deploy to env.
func deployTo(env string) wefttest.Turn {
return wefttest.ToolCalls(wefttest.Call{Name: "deploy", Args: `{"env":"` + env + `"}`})
}
// settle waits for a turn and, when its boundary resumed on its own,
// for the resume too; it returns the last result.
func settle(t *thread.Turn) *weft.RunResult {
res, err := t.Wait()
if err != nil {
fmt.Println("turn failed:", err)
return nil
}
if next := t.Next(); next != nil {
return settle(next)
}
return res
}
func main() {
ctx := context.Background()
agent := gatedDeploy(
deployTo("staging-eu"), wefttest.Say("Staging is live."),
deployTo("staging-us"),
)
s, _ := thread.Create(ctx, thread.Memory(), agent)
err := s.Grant(ctx, thread.Grant{
Tool: "deploy",
Args: []thread.Arg{thread.ArgGlob("/env", "staging-*")},
})
if err != nil {
fmt.Println(err)
return
}
t1, _ := s.Send(ctx, weft.User("Deploy to staging-eu."))
fmt.Println("granted:", settle(t1).Text())
fmt.Println("pending:", len(s.Pending()))
// The audit entry of the match names the grant; revoke it.
for _, e := range s.Audit() {
if a, ok := e.(thread.ApprovalAuditEntry); ok && a.Step == thread.StepGrant {
if err := s.Revoke(ctx, a.GrantID); err != nil {
fmt.Println(err)
return
}
}
}
t2, _ := s.Send(ctx, weft.User("Deploy to staging-us."))
settle(t2)
fmt.Println("pending after the revocation:", len(s.Pending()))
}
Output: granted: Staging is live. pending: 0 pending after the revocation: 1
func (*Session) Label ¶
Label names an entry — bookmarks, checkpoints, the anchors a UI lists. The label is an appended entry, never a rewrite, and never enters the model's context. entryID must name an entry the session holds; the name must not be empty.
Example ¶
Label names an entry: a bookmark a UI lists and a later Branch or Fork can return to. Labels are entries; they never reach the model.
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// exampleIDs mints the given ids in order — the IDs option's usual
// shape in an example: the session id first, then one per entry.
func exampleIDs(ids ...string) thread.SessionOption {
next := 0
return thread.IDs(func() string {
id := ids[next]
next++
return id
})
}
func main() {
ctx := context.Background()
agent := weft.New(wefttest.Script(wefttest.Say("Shipped Tuesday.")))
s, err := thread.Create(ctx, thread.Memory(), agent,
exampleIDs("s_labels", "e_prompt", "e_reply", "e_turn", "e_label"))
if err != nil {
fmt.Println(err)
return
}
turn, _ := s.Send(ctx, weft.User("Where is order 1234?"))
if _, err := turn.Wait(); err != nil {
fmt.Println(err)
return
}
if err := s.Label(ctx, "e_reply", "the shipping answer"); err != nil {
fmt.Println(err)
return
}
for _, e := range s.Entries() {
if l, ok := e.(thread.LabelEntry); ok {
fmt.Printf("%s: %q names %s\n", l.ID, l.Name, l.Entry)
}
}
fmt.Println("context:", len(s.Context()), "messages")
fmt.Println(s.Label(ctx, "e_nowhere", "x"))
}
Output: e_label: "the shipping answer" names e_reply context: 2 messages thread: session s_labels holds no entry "e_nowhere"
func (*Session) Leaf ¶
Leaf returns the entry the next appended entry attaches to: the id of the last entry, the target of a trailing leaf entry, or "" while the session holds no entries (the next entry is the root).
func (*Session) Lineage ¶
Lineage returns the session's pool lineage (ADR 0022 §3): the parent session and delegating call it was started from, as recorded in its header by WithLineage. The zero value means the session is nobody's child.
func (*Session) LoadReport ¶
func (s *Session) LoadReport() *OpenReport
LoadReport returns what Open had to drop, skip or accept to load the session — the torn final line, the lines skipped under Salvage, the entries those skips orphaned — or nil when the load was clean, which is every session Create or Fork returned. The value is a copy the caller owns. Open also logs it once, as a warning through the agent's logger; this accessor is how code, not an operator, learns that the session it holds is a repaired one.
func (*Session) Meta ¶
Meta returns the session's current metadata as a fresh map the caller owns: the header's Meta as the base layer, overlaid by every info entry's Meta in append order. nil before anything set any.
Keys under the reserved "weft." prefix are the exception to the overlay: the header's value stands, whatever an info entry says, and a key the header lacks keeps the first value an info entry gave it. SetInfo refuses those keys, so only a file written before that rule, or by other hands, can hold such an entry — and even there the session's identity never changes mid-life.
func (*Session) MirroredRequests ¶
func (s *Session) MirroredRequests(ctx context.Context) ([]MirroredRequest, error)
MirroredRequests returns the mirrored child requests in force on this session, decided or not, in append order (ADR 0022 §7) — the pool's read of what to replay. Expiry is resolved first, as on every path that looks at the boundary: a mirror strictly past its expiry is denied on the spot (Via "expiry") and returns Decided, which is how a lapsed nested request reaches its child. A mirror superseded by a later one for the same call id is not in force and not returned.
func (*Session) Path ¶
Path returns the entries from a root to entryID, inclusive, in conversation order — the chain of parent links. entryID "" is the root and returns no entries; an id the session does not hold is an error. Path names tree structure: a leaf entry's own parent is where it was appended, not the entry it navigated to. The entries are deep copies, like Entries'.
The path starts at whichever root entryID descends from — a session may hold several — or, in a session opened under Salvage, at an entry LoadReport lists as orphaned: what preceded the skipped line is not on the path. A parent link that is broken in any other way cannot survive Open; were one found here the walk fails with ErrCorrupt rather than returning a path cut short.
Example ¶
Entries is the whole tree in append order, abandoned branches included; Path is one line of it, from a root to an entry. Both return copies.
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// exampleIDs mints the given ids in order — the IDs option's usual
// shape in an example: the session id first, then one per entry.
func exampleIDs(ids ...string) thread.SessionOption {
next := 0
return thread.IDs(func() string {
id := ids[next]
next++
return id
})
}
func main() {
ctx := context.Background()
agent := weft.New(wefttest.Script(wefttest.Say("First draft."), wefttest.Say("Second draft.")))
s, err := thread.Create(ctx, thread.Memory(), agent, exampleIDs("s_tree",
"e_p1", "e_r1", "e_t1", // turn 1
"e_nav", // the branch back to the prompt
"e_p2", "e_r2", "e_t2", // turn 2, on the new branch
))
if err != nil {
fmt.Println(err)
return
}
turn, _ := s.Send(ctx, weft.User("Draft it."))
if _, err := turn.Wait(); err != nil {
fmt.Println(err)
return
}
// Wait returns when the turn is decided; the runner may still be
// between items. A Send in that window is accepted first — one
// more entry, one more id — so an example that prints its ids
// waits for the session to be idle.
if err := s.WaitIdle(ctx); err != nil {
fmt.Println(err)
return
}
if err := s.Branch(ctx, "e_p1"); err != nil {
fmt.Println(err)
return
}
turn, _ = s.Send(ctx, weft.User("Shorter."))
if _, err := turn.Wait(); err != nil {
fmt.Println(err)
return
}
fmt.Println("entries:", len(s.Entries()))
path, err := s.Path(s.Leaf())
if err != nil {
fmt.Println(err)
return
}
for _, e := range path {
switch e := e.(type) {
case thread.MessageEntry:
fmt.Printf("%s message %s: %s\n", e.ID, e.Message.Role, e.Message.Text())
case thread.TurnEntry:
fmt.Printf("%s turn %s\n", e.ID, e.RunID)
}
}
}
Output: entries: 7 e_p1 message user: Draft it. e_p2 message user: Shorter. e_r2 message assistant: Second draft. e_t2 turn s_tree-t2
func (*Session) Pending ¶
Pending returns the session's undecided requests, in call order: the calls the leaf's path leaves dangling that have an approval request and no decision yet. It works after a restart — requests and decisions are entries — and it never shows a call whose boundary a resume already resolved.
func (*Session) Pin ¶
Pin marks an entry to stay in the model's context through every compaction — a requirement, a key decision — by appending a pin record (a custom entry, which survives compaction the way everything custom does). A pin does not constrain the cut: compactions summarize past a pinned entry like any other, record its id in the entry's Pinned list, and the context re-includes the pinned message, raw, right after the summary — so a pin near the root never holds the whole context raw.
Only an entry that contributes a message to the context can be pinned — a message, a custom_message or a branch_summary; any other kind fails with ErrNotPinnable, and an id the session does not hold with ErrNoEntry. A pin is read when a compaction is computed: one made on an entry already below the boundary takes effect at the next compaction, not immediately.
Example ¶
Pin keeps an entry in the context through every compaction — the requirement, the key decision — recorded as a reserved custom entry that survives compaction the way all custom state does.
package main
import (
"context"
"fmt"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
agent := weft.New(wefttest.Script())
st := thread.Memory()
s, _ := thread.Create(ctx, st, agent)
if err := st.Append(ctx, s.ID(), thread.MessageEntry{
ID: "e_req", Created: time.Now().UTC(),
Message: weft.User("THE REQUIREMENT: ship by Friday"),
}); err != nil {
fmt.Println(err)
return
}
// The entries went in behind the Session's back, and it is the
// session's writer since Create: close it, and open one that
// sees them.
_ = s.Close(ctx)
open, _ := thread.Open(ctx, st, s.ID(), agent)
if err := open.Pin(ctx, "e_req"); err != nil {
fmt.Println(err)
return
}
fmt.Println("pinned")
}
Output: pinned
func (*Session) PreviewCompaction ¶
func (s *Session) PreviewCompaction(ctx context.Context, opts ...CompactOption) (*Compaction, error)
PreviewCompaction computes the session's next compaction without writing anything: the cut, the summarized range, the summary text, and the entry fields a Compact would record. It runs the whole algorithm — the BeforeCompact hook, then the Compactor or the summary chain (SummaryModel, then the session's own model, ADR 0020 §2) with the skeleton prompt, the previous summary (iterative compaction) as the range's first message, no cache hints, and the output capped at SummaryMaxTokens or 0.8 × Reserve.
A session whose tail fits inside KeepRecent has nothing to compact and gets ErrNothingToCompact, not a no-op plan; a hook's Cancel is ErrCompactCanceled; a summary cut off at the cap with no fallback left is ErrSummaryTruncated. A dry run is not a failed compaction: CompactFailed does not run for a PreviewCompaction error. It may be called while a turn runs — it reads a snapshot — but the plan's ApplyCompaction then waits for the turn (ErrBusy).
Example ¶
PreviewCompaction computes the next compaction without writing it — the cut and the summary, to inspect or approve — and ApplyCompaction writes the plan as it stands.
ctx := context.Background()
s := longSession(weft.New(&recordingModel{reply: "Goal: ship the order service."}))
plan, err := s.PreviewCompaction(ctx, thread.SummaryInstructions("focus on the refund"))
if err != nil {
fmt.Println(err)
return
}
fmt.Println("summary:", plan.Summary)
fmt.Println("first kept:", plan.FirstKept)
fmt.Println("written by the preview:", len(s.Entries()) != 3)
if err := s.ApplyCompaction(ctx, plan); err != nil {
fmt.Println(err)
return
}
got := s.Context()
fmt.Println("context leads with the summary:", strings.HasPrefix(got[0].Text(), "<weft-summary>"))
fmt.Println("then the kept entries:", len(got)-1)
Output: summary: Goal: ship the order service. first kept: e_1 written by the preview: false context leads with the summary: true then the kept entries: 2
func (*Session) Queue ¶
func (s *Session) Queue() []QueuedSteer
Queue returns what the session has accepted and not yet given to a model, in the order it will get there: first the steers waiting for the running turn's next drain point (QueuedSteer.Policy is Steer), then the sends waiting for a turn of their own (Policy is Queue) — sends accepted while the session was busy, an interrupting send's follow-up, the follow-up a deferred steer became. Delivered steers and turns that have started are not listed — their stories live on the entries (Entries).
On a session that was just opened, the queue holds what Open restored: steers and sends accepted by a writer that stopped before settling them. They wait here — Open runs nothing — until the session's next Send (the restored sends run ahead of it, in order, and a restored steer is delivered into the first turn that runs), Continue runs them now, or ClearQueue drops them.
Queue lists messages; the Policy constant of the same name is the busy policy that queues them.
Example ¶
Queue lists what the session has accepted and not yet given to the model — steers waiting for the running turn, sends waiting for a turn of their own. Both are durable from the moment Send returns; ClearQueue drops them, and their Turns say so.
package main
import (
"context"
"errors"
"fmt"
"sync"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// slowTool returns a tool that blocks until release is closed or its
// context ends, and a channel closed when its first call has started —
// the examples' stand-in for a long piece of work.
func slowTool(name string) (tool *weft.ToolDef, started, release chan struct{}) {
started, release = make(chan struct{}), make(chan struct{})
var once sync.Once
tool = weft.Tool(name, "Takes a while.", func(ctx context.Context, _ struct{}) (string, error) {
once.Do(func() { close(started) })
select {
case <-release:
return name + " finished", nil
case <-ctx.Done():
return "", ctx.Err()
}
})
return tool, started, release
}
func main() {
ctx := context.Background()
work, started, release := slowTool("work")
model := wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "work", ID: "call_1"}),
wefttest.Say("Done."),
)
s, _ := thread.Create(ctx, thread.Memory(), weft.New(model, work))
running, _ := s.Send(ctx, weft.User("Do the long job."))
<-started
later, _ := s.Send(ctx, weft.User("Then email me the result.")) // Queue, the default
nudge, _ := s.Send(ctx, weft.User("Skip the appendix."), thread.As(thread.Steer)) // into the running turn
for _, q := range s.Queue() {
fmt.Printf("%s: %s\n", q.Policy, q.Msg.Text())
}
n, err := s.ClearQueue(ctx)
fmt.Println("dropped:", n, err)
close(release)
res, _ := running.Wait()
fmt.Println("running turn:", res.Text())
_, err = later.Wait()
fmt.Println("queued send:", later.Outcome(), "-", errors.Is(err, thread.ErrDropped))
fmt.Println("steer:", nudge.Outcome())
fmt.Println("model calls:", len(model.Requests()))
}
Output: steer: Skip the appendix. queue: Then email me the result. dropped: 2 <nil> running turn: Done. queued send: dropped - true steer: dropped model calls: 2
func (*Session) ReplayDecisions ¶
ReplayDecisions records, in a pool child, the decisions its parent session took over the child's parked calls, and resumes the child when they complete its boundary (ADR 0022 §7). It is the child's only door for them: not Decide — the batch rules of the user's door do not apply, several decisions for one call (a parent's Quorum) record in the order given — and not DecideSigned, though it records on a RequireSigned child, because each decision was verified in the parent, where it was taken.
Each entry is a decision the parent recorded, its CallID rewritten to the child's own call id. The child's record keeps the decider's identity — Who, and KeyID, so a signed approver counts by key under the child's Quorum as it did under the parent's — and names the channel honestly: Via "parent", never "signed", for the signature itself lives in the parent's entry, not here.
Replay is idempotent by construction. The child's own expiry is resolved first; then every decision for a call that is no longer pending — resolved by an earlier replay, by the child's expiry, by a resume already run — is dropped rather than refused, so a replay repeated after a crash, or one carrying the decisions of an earlier park, records nothing twice. With nothing left to record the call still arms a resume that a completed boundary is owed.
The returned Turn is the resume run, started under ctx; nil when the boundary is not complete. A session with no pool lineage fails: replay is not a way around Decide.
func (*Session) Request ¶
Request returns the pending request for callID with a fresh challenge minted under the keyring's active key: a nonce, the key's id, and everything a signer needs to vouch for a decision over exactly this request (ADR 0021 §3). The challenge is bound to the request entry — this occurrence of the call, not its call id — and is single-use: DecideSigned records the nonce, and a replay of the same signature fails with ErrReplay, across restarts, because recorded nonces are entries. Request writes nothing; it may be called any number of times, each call minting another valid challenge for the same request.
It fails with ErrNotPending when callID is not pending, with ErrDelegated when it is a call delegating to a pool child (which takes no decision), with ErrExpired when the request is past its expiry (no decision could be accepted for it), and when the session has no keyring with an active key.
func (*Session) ResolveDelegation ¶
func (s *Session) ResolveDelegation(ctx context.Context, callID, child, content string, isError bool) (*Turn, error)
ResolveDelegation records the outcome of a pool delegation as the resolution of the parent-side call that delegated it (ADR 0022 §7): content becomes the call's result, marked as an error when isError, recorded with Who "thread/pool" and Via "child", and the boundary resumes when that completes it, like Decide. It is the one way a delegating call is resolved — Decide refuses such a call with ErrDelegated. callID must be the delegating call of a mirrored request of child: the current occurrence of the call, named as Wrapper by a request entry whose Child is child. Anything else — a call that delegates to nobody, or to another child, as a call id reused by a later turn does — fails with ErrNotPending. A delegation's answer is not an approval, so the call records under RequireSigned too (the approvals that let the child run were signed where they were decided), and it is not subject to a request expiry: the child ran, and its answer is the result.
func (*Session) Resume ¶
Resume starts the boundary's resume run now (ADR 0021 §1): the decided calls resolve under their decisions — approved calls run with Call.Approved set, denied ones show their reason, resolved ones their content — and every undecided call is denied with the core's "no decision" text. Requests strictly past their expiry are denied first, on the spot, with the stated reason — the same sweep Decide and every other arming path runs. The parked boundary must still be open: Resume with nothing to resume fails with ErrNotPending. The returned Turn is the resume run's receipt and handle, like any Send's; a Resume while the boundary's resume is already armed returns that turn.
A decision is spent by the resume that applies it: it resolves its call on that resume's line of the tree only. After a Branch back to the decided boundary the calls are pending again and need new decisions — the approved call never runs twice on one approval.
Example ¶
With AutoResume off the caller drives the boundary: decisions are recorded as they arrive, and Resume runs it — the decided calls under their decisions, every undecided one denied as "no decision".
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// gatedDeploy is the agent of the approval examples: its one tool,
// deploy, needs a decision before it runs, and the scripted model
// plays turns.
func gatedDeploy(turns ...wefttest.Turn) *weft.Agent {
return weft.New(wefttest.Script(turns...),
weft.Tool("deploy", "Deploy the service.",
func(ctx context.Context, in struct {
Env string `json:"env"`
}) (string, error) {
return "deployed to " + in.Env, nil
},
weft.RequireApproval()))
}
// settle waits for a turn and, when its boundary resumed on its own,
// for the resume too; it returns the last result.
func settle(t *thread.Turn) *weft.RunResult {
res, err := t.Wait()
if err != nil {
fmt.Println("turn failed:", err)
return nil
}
if next := t.Next(); next != nil {
return settle(next)
}
return res
}
func main() {
ctx := context.Background()
agent := gatedDeploy(
wefttest.ToolCalls(
wefttest.Call{Name: "deploy", ID: "call_eu", Args: `{"env":"eu"}`},
wefttest.Call{Name: "deploy", ID: "call_us", Args: `{"env":"us"}`},
),
wefttest.Say("One region is live."),
)
s, _ := thread.Create(ctx, thread.Memory(), agent, thread.AutoResume(false))
t1, _ := s.Send(ctx, weft.User("Deploy everywhere."))
settle(t1)
turn, err := s.Decide(ctx, thread.Approve("call_eu"))
fmt.Println("decided one of two; resumed:", turn != nil, err)
resume, err := s.Resume(ctx)
if err != nil {
fmt.Println(err)
return
}
res := settle(resume)
for _, m := range s.Context() {
for _, p := range m.Content {
if r, ok := p.(weft.ToolResultPart); ok {
fmt.Println(r.CallID+":", r.Content)
}
}
}
fmt.Println(res.Text())
}
Output: decided one of two; resumed: false <nil> call_eu: deployed to eu call_us: DENIED: no decision One region is live.
func (*Session) Revoke ¶
Revoke ends a session-scoped grant with an entry (ADR 0021 §4): revocation is an append — the grant entry stays, the walk reads the revocation after it. grantID must name a grant entry the session holds; a shared grant's revocation is its store's own operation.
func (*Session) Send ¶
Send appends the prompt and runs the session's agent on the leaf's context, returning the turn's receipt at once. The prompt entry is appended and flushed before the run starts (ADR 0011 §4: accepted input is durable input), the run carries the session's context plus msg and the run id <session>-t<n>, and at the end — success, failure or cancellation — the run's new messages and a turn entry are appended in one atomic batch under a context.WithoutCancel window, so a caller who walks away still leaves a complete session behind.
Settings are captured when Send is called: the agent, the extra run options, and the busy policy. On a busy session the policy decides — Queue (the default) holds the follow-up, in order, durably: an accepted receipt entry carries the message until its turn starts, so a queued send survives a crash and a reopened session finds it in its queue (Session.Queue; Continue runs it). Reject fails with ErrBusy; Steer, Interrupt and Rollback are described on Policy.
A Send while approval requests are pending is queued, not run (ADR 0021 §5): the parked boundary must resolve first — Decide (which resumes on its own when AutoResume is on) or Resume — and the queued follow-up then runs with the boundary's transcript completed. Under Reject a pending boundary reads as busy.
Between a turn's end and the session's next piece of work — the between-turn compaction, when one is due — nothing is running to steer into, interrupt or reject for: a Send in that window is accepted under every policy and runs next.
Example ¶
Send runs a turn: the prompt is durable before the run starts, the reply and the turn's ledger land after it, and a reopen sees the whole conversation.
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
// A scripted model makes the example deterministic and offline.
agent := weft.New(wefttest.Script(wefttest.Say("It shipped Tuesday."), wefttest.Say("Order 1234, two items.")))
st := thread.Memory()
next := 0
ids := []string{"s_demo", "e_1", "e_2", "e_3", "e_4", "e_5", "e_6"}
s, _ := thread.Create(ctx, st, agent, thread.IDs(func() string { id := ids[next]; next++; return id }))
t1, err := s.Send(ctx, weft.User("Where is order 1234?"))
if err != nil {
fmt.Println(err)
return
}
res, err := t1.Wait()
if err != nil {
fmt.Println(err)
return
}
fmt.Println("run", t1.RunID(), "replied:", res.Text())
// The next Send could go at once — it would queue behind the
// runner's between-turn work, with an accepted receipt entry of its
// own. The example prints its ids, so it waits for the session to
// be idle and the Send takes the idle path.
if err := s.WaitIdle(ctx); err != nil {
fmt.Println(err)
return
}
t2, err := s.Send(ctx, weft.User("And what was in it?"))
if err != nil {
fmt.Println(err)
return
}
if _, err := t2.Wait(); err != nil {
fmt.Println(err)
return
}
again, err := thread.Open(ctx, st, s.ID(), agent)
if err != nil {
fmt.Println(err)
return
}
for _, m := range again.Context() {
if m.Text() != "" {
fmt.Println(m.Role, ":", m.Text())
}
}
}
Output: run s_demo-t1 replied: It shipped Tuesday. user : Where is order 1234? assistant : It shipped Tuesday. user : And what was in it? assistant : Order 1234, two items.
func (*Session) SetInfo ¶
SetInfo edits the session's title and metadata as an appended info entry, never a rewrite of the header (ADR 0011 §2): the current title is the last info entry's non-empty Title, the current metadata the header's overlaid by every info entry's Meta in append order — Title and Meta resolve them.
An empty title keeps the current one, so a title can be replaced but not cleared. A nil or empty meta changes nothing; a non-empty one overwrites the keys it names and leaves the rest (a key is set to "" to blank it; there is no removal). A call that names neither a title nor a key has nothing to record: it returns nil and writes no entry.
Keys under the reserved "weft." prefix are the session's identity (weft.public_id) and belong to the header: SetInfo rejects them with ErrReservedKey and writes nothing. Set them at Create, with WithMeta or PublicID.
func (*Session) Storage ¶
Storage returns the storage the session was created or opened on — the handle thread/pool holds to create child sessions in the same place (ADR 0022 §3). It is the application's own storage returned, read-only by convention: every write to this session goes through the Session, never behind it. A closed session still returns it.
func (*Session) Title ¶
Title returns the session's current title — the last info entry's (SetInfo's rule), "" before any SetInfo named one. A fork carries the info entries on its copied path, and so their title.
func (*Session) Uncompact ¶
Uncompact undoes the latest compaction entry on the leaf's path — a summary compaction or a trim, whichever landed last — the only way the format allows: a branch back to the entry before it (ADR 0020 §1), a leaf entry, no rewrite, nothing deleted. The context is exactly what it was before that entry landed; the entry and everything after it stay in the file, off the path. After a compaction followed by a trim, the first Uncompact undoes the trim and a second one the compaction.
It is a navigation, so Branch's rule applies: ErrBusy while a turn is in flight. With no compaction on the path it fails wrapping ErrNoEntry.
func (*Session) Usage ¶
Usage returns the session's cost ledger over every entry in the file, not only the leaf's path: turns from the per-turn ledger, summarizer costs in their own bucket (ADR 0011 §2, ADR 0020 §4).
Example ¶
The tree in memory and the session in storage agree, and the leaf is where the next entry attaches: Entries walks the whole tree, Path one root-to-entry chain, and Usage reads the cost ledger.
package main
import (
"context"
"fmt"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
agent := weft.New(wefttest.Script())
st := thread.Memory()
s, err := thread.Create(ctx, st, agent, thread.IDs(func() string { return "s_ledger" }))
if err != nil {
fmt.Println(err)
return
}
// A finished turn, recorded the way Send records one.
if err := st.Append(ctx, s.ID(),
thread.MessageEntry{ID: "e_q", Created: time.Now().UTC(), Message: weft.User("Summarize the plan.")},
thread.TurnEntry{
ID: "e_t", ParentID: "e_q", Created: time.Now().UTC(),
RunID: "s_ledger-t1", StopReason: weft.StopEndTurn,
Usage: weft.Usage{InputTokens: 120, OutputTokens: 30},
},
); err != nil {
fmt.Println(err)
return
}
open, err := thread.Open(ctx, st, "s_ledger", agent)
if err != nil {
fmt.Println(err)
return
}
fmt.Println("entries:", len(open.Entries()), "leaf:", open.Leaf())
fmt.Println("path:", len(mustPath(open, "e_t")))
u := open.Usage()
fmt.Println("turn input tokens:", u.Turns.InputTokens, "summary tokens:", u.Summaries.InputTokens)
}
func mustPath(s *thread.Session, id string) []thread.Entry {
path, err := s.Path(id)
if err != nil {
panic(err)
}
return path
}
Output: entries: 2 leaf: e_t path: 2 turn input tokens: 120 summary tokens: 0
func (*Session) WaitIdle ¶
WaitIdle blocks until the session has nothing in progress — no turn in flight, no between-turn work, nothing runnable left in its queue — and returns nil; it returns ctx.Err() when ctx ends first.
A turn's Wait returns when the turn is decided. What the session does between turns — the automatic compaction a turn can trigger (ADR 0020 §2) — runs after that, before the next queued turn. A caller that needs the session's state after that housekeeping (a test, a tool that prints the tree, a pause that is not a Close) waits here; a caller that only sends the next message need not — the Send queues behind the housekeeping on its own.
Idle is not empty: sends held by an open approval boundary wait for a decision, not for the runner, and steers restored by Open wait for a turn. WaitIdle returns while they are still queued.
func (*Session) WatchMirrors ¶
WatchMirrors installs fn as owner's notification that this session recorded a decision for a mirrored child request (ADR 0022 §7) — by any path: Decide and DecideSigned, the denial an interrupting Send gives a parked boundary, the expiry sweep. The pool installs one per parent it delegates from, so a decision reaches the child it is for without a caller having to pump: the session records, the pool replays.
fn is called after the decision is durable, with the session's lock held: it must not call the session and must return at once — the pool hands the work to a goroutine of its own. It is handed the session's logger, the place its failures are reported. One owner holds one notification: installing again under the same owner replaces it, and a nil fn removes it. The notification is the live Session's, never the file's: a session opened again has none until its pool attaches (Recover, Decide, a new delegation).
type SessionOption ¶
type SessionOption interface {
// contains filtered or unexported methods
}
SessionOption configures a session at Create, Open or Fork — one value per concern, folded over the defaults, the constructor convention of the core (ADR 0011 §6). Every option family of the package joins this one interface: ids and the clock (IDs, Clock), the busy policy (BusyPolicy), compaction (ContextWindow and its layers), approvals (WithApprover, Quorum, WithKeyring, …).
Three options only mean something while a session's header is being written — WithMeta, PublicID and WithLineage. Create and Fork honour them; Open refuses them with ErrCreateOnly, because the header it loads is immutable and an option it silently dropped would read as applied.
func AfterCompact ¶
func AfterCompact(fn func(ctx context.Context, e CompactionEntry)) SessionOption
AfterCompact returns the SessionOption setting the hook that runs after a compaction entry lands, with the entry — the durable record, not the plan. It runs for every entry, however it was written: Compact, ApplyCompaction, the automatic trigger, the overflow re-run, and a trim (e.Reason is ReasonTrim and e.Trim holds the record). A panic in fn is contained and logged.
fn runs without the session lock; it may call the session — the Context it reads already shows the compaction.
func AutoResume ¶
func AutoResume(on bool) SessionOption
AutoResume sets whether the session resumes on its own once every pending call of the parked boundary has a decision (ADR 0021 §1). The default is on; AutoResume(false) parks the decided boundary until the caller calls Resume. Either way a plain Send never runs while a boundary is open — only the resume run resolves it.
func BeforeCompact ¶
func BeforeCompact(fn func(ctx context.Context, p *Preparation) (Verdict, error)) SessionOption
BeforeCompact returns the SessionOption setting the hook that sees every computed summary compaction before the summarizer runs, told the reason, and decides: Proceed (with its edits to the Preparation — see Preparation for the editable fields; a redaction hook rewrites p.Messages), Cancel, or Replace with a hook-made summary. A hook error or panic fails the compaction with that error. A trim does not pass through it: a trim summarizes nothing.
fn runs without the session lock; it may call the session.
Example ¶
BeforeCompact sees every compaction before the summarizer runs and may edit the plan: a redaction hook rewrites p.Messages, and the summarizer is fed the redacted range.
ctx := context.Background()
summarizer := &summaryRecorder{reply: "Goal: resolve the customer's ticket."}
redact := thread.BeforeCompact(func(ctx context.Context, p *thread.Preparation) (thread.Verdict, error) {
clean := make([]weft.Message, len(p.Messages))
for i, m := range p.Messages {
clean[i] = weft.Message{Role: m.Role, Content: []weft.Part{
weft.TextPart{Text: strings.ReplaceAll(m.Text(), "ada@example.com", "[email]")},
}}
}
p.Messages = clean
return thread.Proceed(), nil
})
agent := weft.New(summarizer)
st := thread.Memory()
s, _ := thread.Create(ctx, st, agent, redact)
now := time.Now().UTC()
if err := st.Append(ctx, s.ID(),
thread.MessageEntry{ID: "e_0", Created: now, Message: weft.User("I am ada@example.com. " + strings.Repeat("My order is late. ", 5_000))},
thread.MessageEntry{ID: "e_1", ParentID: "e_0", Created: now, Message: weft.User(strings.Repeat("Any news? ", 4_000))},
); err != nil {
fmt.Println(err)
return
}
// The entries went in behind the Session's back, and it is the
// session's writer since Create: close it, and open one that
// sees them.
_ = s.Close(ctx)
s, _ = thread.Open(ctx, st, s.ID(), agent, redact)
if err := s.Compact(ctx); err != nil {
fmt.Println(err)
return
}
fed := summarizer.saw()[0].Messages[0].Text()
fmt.Println("summarizer saw the address:", strings.Contains(fed, "ada@example.com"))
fmt.Println("summarizer saw:", fed[:len("I am [email].")])
Output: summarizer saw the address: false summarizer saw: I am [email].
func BusyPolicy ¶
func BusyPolicy(p Policy) SessionOption
BusyPolicy returns the SessionOption setting what Send does when the session is busy: Queue (the default), Reject, Steer, Interrupt, or Rollback. A single Send overrides it with As.
Example ¶
A session under the Reject busy policy says ErrBusy instead of holding a follow-up while a turn runs; once the turn has ended the same Send is accepted.
package main
import (
"context"
"errors"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
started, release := make(chan struct{}), make(chan struct{})
agent := weft.New(
wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "work", ID: "call_1"}),
wefttest.Say("Done."),
wefttest.Say("Here it is."),
),
weft.Tool("work", "Takes a while.", func(ctx context.Context, _ struct{}) (string, error) {
close(started)
<-release
return "ok", nil
}),
)
s, _ := thread.Create(ctx, thread.Memory(), agent, thread.BusyPolicy(thread.Reject))
turn, _ := s.Send(ctx, weft.User("Do the long job."))
<-started
_, err := s.Send(ctx, weft.User("One more thing."))
fmt.Println("while the turn runs:", errors.Is(err, thread.ErrBusy))
close(release)
if _, err := turn.Wait(); err != nil {
fmt.Println(err)
return
}
next, err := s.Send(ctx, weft.User("One more thing."))
if err != nil {
fmt.Println(err)
return
}
res, _ := next.Wait()
fmt.Println("after it:", res.Text())
}
Output: while the turn runs: true after it: Here it is.
func CheckSummary ¶
func CheckSummary(fn func(Summary) error) SessionOption
CheckSummary returns the SessionOption validating each summary — the headings-present, length-floor checks a careful caller writes. A failure retries the same summarizer once, then falls back down the chain (SummaryModel → session model → no compaction: the error is returned and reported through CompactFailed). A summary cut off at the output cap never reaches fn: it is rejected as ErrSummaryTruncated on the same retry-and-fallback road.
fn runs without the session lock; it may call the session. A panic in fn counts as a rejection.
Example ¶
CheckSummary validates each summary: a rejected one is retried once on the same model, then the chain falls back to the session's own.
ctx := context.Background()
cheap := &summaryRecorder{reply: "it went fine"}
session := &summaryRecorder{reply: "Goal: ship the order service."}
s := longSession(weft.New(session),
thread.SummaryModel(cheap),
thread.CheckSummary(func(sum thread.Summary) error {
if !strings.Contains(sum.Text, "Goal:") {
return errors.New("summary has no Goal heading")
}
return nil
}),
)
if err := s.Compact(ctx); err != nil {
fmt.Println(err)
return
}
fmt.Println("cheap model attempts:", len(cheap.saw()))
fmt.Println("session model attempts:", len(session.saw()))
fmt.Println(s.Context()[0].Text())
Output: cheap model attempts: 2 session model attempts: 1 <weft-summary> Goal: ship the order service. </weft-summary>
func ClearOldToolResults ¶
func ClearOldToolResults(keepLast int) SessionOption
ClearOldToolResults returns the SessionOption turning on the built-in trimmer: when the trigger fires, every tool result in the context except the newest keepLast is replaced with the stub naming the call ("[cleared tool result NAME CALLID]"). If the trimmed context fits window − Reserve, no summary is made and a compaction entry with Reason trim lands instead, its Trim record naming each stubbed result — the record, not this option, is what later context builds replay, so reopening with another keepLast (or none) never changes what a recorded trim shows. A negative keepLast is ignored.
Example ¶
ClearOldToolResults is the cheap pre-pass: when the trigger fires, old tool results are stubbed, and if that is enough no summary is made — a trim record lands instead, and the file keeps the results.
package main
import (
"context"
"fmt"
"strings"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
agent := weft.New(wefttest.Script(
wefttest.Say("Both files read.").WithUsage(weft.Usage{InputTokens: 90_000, OutputTokens: 5}),
))
opts := []thread.SessionOption{
thread.ContextWindow(100_000), // arms the trigger at 100,000 − Reserve
thread.ClearOldToolResults(1), // keep the newest result raw
}
st := thread.Memory()
s, _ := thread.Create(ctx, st, agent, opts...)
now := time.Now().UTC()
if err := st.Append(ctx, s.ID(),
thread.MessageEntry{ID: "e_ask", Created: now, Message: weft.User("read both files")},
thread.MessageEntry{ID: "e_calls", ParentID: "e_ask", Created: now, Message: weft.Message{Role: weft.RoleAssistant,
Content: []weft.Part{
weft.ToolCallPart{ID: "call_1", Name: "read", Args: []byte(`{"path":"a.go"}`)},
weft.ToolCallPart{ID: "call_2", Name: "read", Args: []byte(`{"path":"b.go"}`)},
}}},
thread.MessageEntry{ID: "e_results", ParentID: "e_calls", Created: now, Message: weft.Message{Role: weft.RoleTool,
Content: []weft.Part{
weft.ToolResultPart{CallID: "call_1", Name: "read", Content: strings.Repeat("a", 40_000)},
weft.ToolResultPart{CallID: "call_2", Name: "read", Content: strings.Repeat("b", 40_000)},
}}},
); err != nil {
fmt.Println(err)
return
}
// The entries went in behind the Session's back, and it is the
// session's writer since Create: close it, and open one that
// sees them.
_ = s.Close(ctx)
s, _ = thread.Open(ctx, st, s.ID(), agent, opts...)
// The turn reports 90,000 input tokens: over the line, so the
// trimmer runs after it.
turn, err := s.Send(ctx, weft.User("thanks"))
if err != nil {
fmt.Println(err)
return
}
if _, err := turn.Wait(); err != nil {
fmt.Println(err)
return
}
// The trigger runs between turns, after the turn is decided: wait
// for the session to go quiet before reading what it did.
if err := s.WaitIdle(ctx); err != nil {
fmt.Println(err)
return
}
for _, e := range s.Entries() {
if c, ok := e.(thread.CompactionEntry); ok {
fmt.Println("compaction entry:", c.Reason, "— stubbed", len(c.Trim.Stubs), "result")
}
}
for _, m := range s.Context() {
for _, p := range m.Content {
if r, ok := p.(weft.ToolResultPart); ok {
fmt.Println(r.CallID, "shows", len(r.Content), "bytes:", r.Content[:min(40, len(r.Content))])
}
}
}
}
Output: compaction entry: trim — stubbed 1 result call_1 shows 33 bytes: [cleared tool result read call_1] call_2 shows 40000 bytes: bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb
func Clock ¶
func Clock(now func() time.Time) SessionOption
Clock returns the SessionOption that makes now the session's time source, for everything the session reads a time for: the header's Created at Create and Fork, the Created of every entry the Session appends, and the approval machinery's arithmetic — when a request or a grant expires, and whether a signed decision arrived in time (all converted to UTC). The default is time.Now. Tests and examples pin time with it, the way IDs pins ids; a nil now is ignored. A thread/pool child inherits its parent's (InheritApprovals), so a nested request lapses on the same clock.
Like the IDs function, now may be called with the session's lock held, and from the runner's goroutine: it must be safe for concurrent use, return quickly, and never call back into the Session. Entry ids are not read from it: the default ids are time-sortable by the wall clock (IDs replaces them).
Example ¶
Clock pins the session's time source, the way IDs pins its ids: the header and every entry the session appends read it.
package main
import (
"context"
"fmt"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// exampleIDs mints the given ids in order — the IDs option's usual
// shape in an example: the session id first, then one per entry.
func exampleIDs(ids ...string) thread.SessionOption {
next := 0
return thread.IDs(func() string {
id := ids[next]
next++
return id
})
}
// exampleClock ticks one second per reading from a fixed instant, so
// every stored timestamp in an example is known.
func exampleClock() thread.SessionOption {
now := time.Date(2026, 10, 1, 9, 0, 0, 0, time.UTC)
return thread.Clock(func() time.Time {
now = now.Add(time.Second)
return now
})
}
func main() {
ctx := context.Background()
st := thread.Memory()
s, err := thread.Create(ctx, st, weft.New(wefttest.Script()),
exampleIDs("s_clock", "e_note"), exampleClock())
if err != nil {
fmt.Println(err)
return
}
if err := s.Custom(ctx, "note", []byte(`{"pinned":true}`)); err != nil {
fmt.Println(err)
return
}
page, _ := thread.List(ctx, st, thread.Query{})
fmt.Println("header:", page.Sessions[0].Created.Format(time.RFC3339))
fmt.Println("entry: ", s.Entries()[0].(thread.CustomEntry).Created.Format(time.RFC3339))
}
Output: header: 2026-10-01T09:00:01Z entry: 2026-10-01T09:00:02Z
func CompactFailed ¶
func CompactFailed(fn func(ctx context.Context, r Reason, err error)) SessionOption
CompactFailed returns the SessionOption setting the hook that runs when a compaction that was to be written is not — the session is unchanged, and the caller is told why, with the reason it was attempted for. That covers Compact and the automatic paths failing at any stage after every fallback (the hook, the Compactor, the summary chain, an unrepresentable trim), and ApplyCompaction failing to validate or store its entry. It does not run for PreviewCompaction — a dry run writes nothing and returns its error to its caller — nor for the refusals that attempt nothing: ErrNothingToCompact, ErrCompactCanceled, ErrBusy and ErrAwaitingApproval. A panic in fn is contained and logged.
fn runs without the session lock; it may call the session.
func ContextWindow ¶
func ContextWindow(n int64) SessionOption
ContextWindow returns the SessionOption setting the context window the trigger budgets against, in tokens. Zero or negative is ignored (the window stays unknown: no automatic compaction, one warning). ModelWindows overrides it per model. A known window must leave room for the other two knobs — Reserve < window and KeepRecent < window − Reserve — or Create and Open fail with ErrCompactConfig: a window too small for its reserve would otherwise never compact, silently.
func IDs ¶
func IDs(id func() string) SessionOption
IDs returns the SessionOption that takes every id a session mints — the session id Create or Fork generates and the entry id of each entry the Session appends — from id. Tests and examples pin deterministic ids with it. A nil id is ignored, leaving the default time-sortable ids; a value that fails ValidID, or repeats an id the session already holds, fails the write that would carry it.
The session calls id with its lock held — ids are minted inside the write they name, so their order is the entries' order. The function must return quickly and must not call back into the Session: any Session method called from it deadlocks.
func InheritApprovals ¶
func InheritApprovals(parent *Session) SessionOption
InheritApprovals returns the SessionOption that gives a pool child the parts of its parent's approval policy that nesting needs (ADR 0022 §7): the parent's RequestExpiry, so a nested request lapses like one of the parent's own; its Quorum, so the replayed approvals count in the child as they did in the parent; its keyring and its RequireSigned rule, so nobody holding the child session can decide its parked calls through the unsigned door the parent closed; and its Clock, the time those expiries are read against. The pool passes it when it creates a child and again when it reopens one — only RequireSigned is stored in the header; the rest is configuration, as on any session. Options listed after it override what it set; a nil parent is ignored.
func KeepRecent ¶
func KeepRecent(n int64) SessionOption
KeepRecent returns the SessionOption setting roughly how many estimated tokens of the tail stay raw on either side of the cut. Zero or negative keeps the default 20,000.
func MaxPerSession ¶
func MaxPerSession(n int) SessionOption
MaxPerSession returns the SessionOption capping how many automatic compactions (every reason but manual, trims included) may sit on the leaf's path — zero (the default) means no cap; at the cap the trigger stops firing and manual Compact keeps working. The count runs along the leaf's path: an abandoned branch's compactions are not this line's.
func MinTurnsBetween ¶
func MinTurnsBetween(n int) SessionOption
MinTurnsBetween returns the SessionOption rate-limiting automatic compaction: no automatic compaction or trim within n turns of the last compaction entry on the leaf's path (manual Compact always works). Zero, the default, means no limit. The count runs along the leaf's path, so a compaction on another branch does not hold this one back. The stop against a context that re-crosses the line every turn.
func ModelReserves ¶
func ModelReserves(reserves map[core.ModelInfo]int64) SessionOption
ModelReserves returns the SessionOption setting per-model reserves, keyed by ModelInfo, overriding Reserve for the session agent's model — the small-window model that needs a different headroom than the default. Resolved once at Create or Open, like ModelWindows.
func ModelWindows ¶
func ModelWindows(windows map[core.ModelInfo]int64) SessionOption
ModelWindows returns the SessionOption setting per-model context windows, keyed by ModelInfo (provider + name): the session agent's model is looked up once, at Create or Open, and its entry overrides ContextWindow for the life of that Session value — one options list shared by sessions that run different agents. The lookup is not repeated per turn: a run that swaps its model (core.UseModel) keeps the window resolved for the agent's own model, and a session reopened with another agent resolves again for that agent.
func NoAutoCompact ¶
func NoAutoCompact() SessionOption
NoAutoCompact returns the SessionOption turning automatic compaction off entirely — no trigger, no trim, no no-window warning: the session that compacts only by hand. Manual Compact, PreviewCompaction, ApplyCompaction and Uncompact are unaffected, and so is the overflow compaction a failed turn runs (ReRunOnOverflow governs that one).
func OnRequest ¶
func OnRequest(fn func(Request)) SessionOption
OnRequest sets the notification fired when a request parks (ADR 0021 §5): called after the request entry is durable, once per parked call, in call order, so an application can push it anywhere. fn runs synchronously on the session's runner goroutine, before the parking turn's Wait returns and before the next turn can start: a slow fn holds the session, so hand the request to a queue or a goroutine and return. It holds no session lock — fn may call Pending or Request — but it must not wait for the turn that is calling it. The session does no I/O of its own for this; a panic in fn is contained and logged, never failing the turn. A nil fn is ignored.
func PreferNative ¶
func PreferNative() SessionOption
PreferNative returns the SessionOption asking thread to use the provider's own compaction when the summarizer model (or the one it wraps) implements NativeCompactor — found by following Unwrap() through middleware. Only the returned message's text is kept: it is stored as the summary, shown behind the same marker as any other, and the entry records the model that made it. An opaque provider-native compaction item is not stored or replayed — that part of ADR 0020 §7 is not implemented. On a model without the interface, an error (logged) or an empty text, the text summary runs: the seam, with the fallback.
func PublicID ¶
func PublicID(id string) SessionOption
PublicID sets the session's public id: an opaque, browser-safe handle stamped on every run of the session as weft.public_id, beside weft.session.id and weft.turn (ADR 0024 S5). It is stored as the header metadata key "weft.public_id" — exactly what WithMeta(map[string]string{"weft.public_id": id}) stores — so List's Query.Meta filter finds the session by it: backends match header metadata only.
The public id is fixed for the session's life. It is a header option — Create and Fork write it, Open fails with ErrCreateOnly when given it — and no later edit can rotate it: SetInfo rejects every key under the reserved "weft." prefix with ErrReservedKey, and Session.Meta never lets an info entry override a "weft." key (a session file written before that rule keeps the first value it recorded). A fork does not inherit it: a fork is another session, and takes its own PublicID or none.
Example ¶
PublicID sets the session's public id: an opaque, browser-safe handle for the session (WEFT-OTEL-DATA-ARCHITECTURE §5), stamped on every run as weft.public_id — beside the session id and turn number the session adds — and read back, like every metadata pair, with weft.MetadataFromContext inside the run.
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
var public, session, turn string
agent := weft.New(wefttest.Script(wefttest.Say("a reply")),
weft.Tap(func(ctx context.Context, ev weft.Event) {
if _, ok := ev.(weft.RunStart); ok {
md := weft.MetadataFromContext(ctx)
public, session, turn = md["weft.public_id"], md["weft.session.id"], md["weft.turn"]
}
}))
s, _ := thread.Create(ctx, thread.Memory(), agent, thread.PublicID("support-42"))
t, _ := s.Send(ctx, weft.User("hello"))
if _, err := t.Wait(); err != nil {
fmt.Println(err)
return
}
fmt.Println(public, session == s.ID(), turn)
}
Output: support-42 true 1
func Quorum ¶
func Quorum(n int) SessionOption
Quorum sets how many approvals from distinct approver identities a call needs before it resolves (ADR 0021 §5). Values below 2 read as the default — one decision resolves. A deny resolves alone, and conflicting decisions resolve to deny with the pinned reason.
What an identity is depends on the door the approval came through. A signed approval (DecideSigned) is its key: one keyring key is one approver, whatever Who the signature carries — give every approver their own key. An unsigned approval (Decide, an Approver's answer) is its Who, which is only what the caller declared: one caller can write two names. A grant's approval is one identity of its own. So a quorum that must hold against the deciding process itself needs RequireSigned; without it Quorum is a workflow rule among callers that are trusted to say who they are.
Example ¶
Under Quorum a call needs approvals from distinct approvers. Signed, an approver is a key: each approver signs with their own, and one key signing twice — whatever names it gives — is still one approver.
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// gatedDeploy is the agent of the approval examples: its one tool,
// deploy, needs a decision before it runs, and the scripted model
// plays turns.
func gatedDeploy(turns ...wefttest.Turn) *weft.Agent {
return weft.New(wefttest.Script(turns...),
weft.Tool("deploy", "Deploy the service.",
func(ctx context.Context, in struct {
Env string `json:"env"`
}) (string, error) {
return "deployed to " + in.Env, nil
},
weft.RequireApproval()))
}
// deployTo is a model turn asking to deploy to env.
func deployTo(env string) wefttest.Turn {
return wefttest.ToolCalls(wefttest.Call{Name: "deploy", Args: `{"env":"` + env + `"}`})
}
// settle waits for a turn and, when its boundary resumed on its own,
// for the resume too; it returns the last result.
func settle(t *thread.Turn) *weft.RunResult {
res, err := t.Wait()
if err != nil {
fmt.Println("turn failed:", err)
return nil
}
if next := t.Next(); next != nil {
return settle(next)
}
return res
}
func main() {
ctx := context.Background()
alice := thread.Key{ID: "alice", Secret: []byte("alice's secret"), Active: true}
bob := thread.Key{ID: "bob", Secret: []byte("bob's secret")}
ring, _ := thread.NewKeyring(alice, bob)
agent := gatedDeploy(deployTo("prod"), wefttest.Say("Deployed."))
s, _ := thread.Create(ctx, thread.Memory(), agent,
thread.WithKeyring(ring), thread.RequireSigned(), thread.Quorum(2))
t1, _ := s.Send(ctx, weft.User("Deploy to prod."))
settle(t1)
call := s.Pending()[0].CallID
approve := func(k thread.Key) *thread.Turn {
challenge, err := s.Request(call) // one challenge per signature
if err != nil {
fmt.Println(err)
return nil
}
signed, err := k.Sign(challenge, thread.Approve(call))
if err != nil {
fmt.Println(err)
return nil
}
turn, err := s.DecideSigned(ctx, signed)
if err != nil {
fmt.Println(err)
}
return turn
}
approve(alice)
approve(alice) // the same key again adds nothing
fmt.Println("after alice, twice:", len(s.Pending()), "pending")
resume := approve(bob)
fmt.Println("after bob:", settle(resume).Text())
}
Output: after alice, twice: 1 pending after bob: Deployed.
func ReRunOnOverflow ¶
func ReRunOnOverflow(on bool) SessionOption
ReRunOnOverflow sets whether a turn that fails with core.ErrContextOverflow compacts (reason overflow) and re-runs once over the shrunken path (ADR 0020 §5) — on by default. A second overflow fails the turn with both errors joined; with the re-run off, the first overflow fails it directly.
Example ¶
ReRunOnOverflow arms the overflow recovery (ADR 0020 §5): a turn failing with weft.ErrContextOverflow compacts — reason overflow — and runs once more over the shrunken path, so the caller sees one turn with the re-run's answer.
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
model := wefttest.Script(
wefttest.Say("a first answer"), // a prior turn, so the compaction has a cut
wefttest.Fail(weft.ErrContextOverflow), // the overflowing request
wefttest.Say("the summary of what came before"), // the compaction's summarizer call
wefttest.Say("recovered after compaction"), // the re-run
)
s, _ := thread.Create(ctx, thread.Memory(), weft.New(model), thread.KeepRecent(1))
t0, _ := s.Send(ctx, weft.User("a first question"))
if _, err := t0.Wait(); err != nil {
fmt.Println(err)
return
}
t1, _ := s.Send(ctx, weft.User("a prompt that overflows"))
res, err := t1.Wait()
fmt.Println(err)
fmt.Println(res.Text())
}
Output: <nil> recovered after compaction
func RequestExpiry ¶
func RequestExpiry(d time.Duration) SessionOption
RequestExpiry gives every request the session parks a lifetime: a request strictly past Created plus d takes no decision any more — Decide and DecideSigned refuse one with ErrExpired — and is denied with the stated reason the next time the session looks at the boundary: a Decide, a Resume, a Send, or the runner's own pickup (ADR 0021 §5). Nothing fires on a timer: an idle session denies an expired request when it is next touched. Values <= 0 (the default) mean requests never expire.
Example ¶
RequestExpiry gives every parked request a lifetime, and OnRequest says when one parks. A request past its expiry takes no decision: Decide refuses it with ErrExpired and the request is denied with a reason the model sees.
ctx := context.Background()
agent := gatedDeploy(deployTo("prod"), wefttest.Say("Noted: the request lapsed."))
s, _ := thread.Create(ctx, thread.Memory(), agent,
thread.RequestExpiry(10*time.Millisecond),
thread.OnRequest(func(r thread.Request) {
// Runs on the session's runner: hand the request to a queue
// or a goroutine and return.
fmt.Println("parked:", r.Tool, "expires:", !r.Expiry.IsZero())
}))
t1, _ := s.Send(ctx, weft.User("Deploy to prod."))
settle(t1)
req := s.Pending()[0]
time.Sleep(time.Until(req.Expiry) + 5*time.Millisecond) // nobody decided in time
_, err := s.Decide(ctx, thread.Approve(req.CallID))
fmt.Println("too late:", errors.Is(err, thread.ErrExpired))
// The expiry's denial completed the boundary; the session resumed.
settle(t1)
for _, m := range s.Context() {
if r, ok := lastResult(m); ok {
fmt.Println("the model saw:", strings.SplitN(r, " before ", 2)[0])
}
}
Output: parked: deploy expires: true too late: true the model saw: DENIED: expired: no decision
func RequireSigned ¶
func RequireSigned() SessionOption
RequireSigned returns the SessionOption that closes the in-process door: Decide — the unsigned path — is rejected with ErrSignatureRequired, and only DecideSigned records a caller's decision. The session's own machinery still records, each path with its own audit trail: grants, the Approver, expiry, the denial an interrupting Send makes, a pool delegation's resolution.
The rule is the session's, not the process's: Create writes it into the session's header (metadata key "weft.require_signed"), so every later Open enforces it whether or not it passes the option. Open with RequireSigned on a session created without it tightens that one Session value; nothing loosens a session created with it.
It needs a way to decide: Create and Open fail when RequireSigned is in force and the session has no keyring (WithKeyring), or a ring with no active key — the session could mint no challenge, and no call it parked could ever be decided.
func Reserve ¶
func Reserve(n int64) SessionOption
Reserve returns the SessionOption setting the trigger's headroom: compaction fires before the context is within Reserve tokens of the window. Zero or negative keeps the default 16,384. It also sizes the summarizer's output cap (0.8 × Reserve unless SummaryMaxTokens says otherwise).
func SummaryFocus ¶
func SummaryFocus(focus string) SessionOption
SummaryFocus returns the SessionOption appending instructions to the summary prompt (a replacement or the skeleton): "keep every file path" — the caller's standing concern, appended after the skeleton's own rules.
func SummaryMaxTokens ¶
func SummaryMaxTokens(n int64) SessionOption
SummaryMaxTokens returns the SessionOption overriding the summarizer output cap, in tokens (the default 0.8 × Reserve). Zero or negative keeps the default. A summary that reaches the cap is not stored: it fails as ErrSummaryTruncated through the retry-and-fallback chain.
func SummaryModel ¶
func SummaryModel(m core.Model) SessionOption
SummaryModel returns the SessionOption setting a different model for summaries — a cheap one. When it fails — an error, a summary CheckSummary rejects twice, or one cut off at the output cap twice — the session's own model takes over (the chain: SummaryModel → session model → no compaction, the error returned and reported through CompactFailed).
Example ¶
SummaryModel summarizes with a different, cheaper model; its cost lands in the ledger's own bucket, never mixed with the turns'.
ctx := context.Background()
cheap := &summaryRecorder{reply: "Goal: ship the order service."}
s := longSession(weft.New(wefttest.Script()), thread.SummaryModel(cheap))
if err := s.Compact(ctx); err != nil {
fmt.Println(err)
return
}
fmt.Println("summaries made by the cheap model:", len(cheap.saw()))
u := s.Usage()
fmt.Println("summary tokens:", u.Summaries.InputTokens, "in,", u.Summaries.OutputTokens, "out")
fmt.Println("turn tokens:", u.Turns.InputTokens, "in")
Output: summaries made by the cheap model: 1 summary tokens: 7 in, 3 out turn tokens: 0 in
func SummaryPrompt ¶
func SummaryPrompt(tmpl string) SessionOption
SummaryPrompt returns the SessionOption replacing the skeleton system prompt whole. The replacement's bytes are the caller's; the default skeleton stays pinned by its golden.
func TriggerFunc ¶
func TriggerFunc(fn func(TriggerInput) bool) SessionOption
TriggerFunc returns the SessionOption replacing the trigger's condition — fire when fn says so. It is consulted only when a window is known and a provider-reported input describes the current path: the no-estimated-signal rule stands under a custom trigger too, and so does the stand-down after a compaction (the trigger waits for the next provider report). fn runs without the session lock; it may call the session. A panic in fn is contained and reads as "do not fire".
Example ¶
TriggerFunc replaces the trigger's condition: here, compact as soon as the reported input passes half the window.
ctx := context.Background()
agent := weft.New(wefttest.Script(
wefttest.Say("Noted.").WithUsage(weft.Usage{InputTokens: 60_000, OutputTokens: 5}),
))
s := longSession(agent,
thread.ContextWindow(100_000),
thread.SummaryModel(&recordingModel{reply: "Goal: ship the order service."}),
thread.TriggerFunc(func(in thread.TriggerInput) bool {
return in.LastInput+in.Estimated > in.Window/2
}),
)
turn, err := s.Send(ctx, weft.User("one more thing"))
if err != nil {
fmt.Println(err)
return
}
if _, err := turn.Wait(); err != nil {
fmt.Println(err)
return
}
if err := s.WaitIdle(ctx); err != nil { // the trigger runs between turns
fmt.Println(err)
return
}
for _, e := range s.Entries() {
if c, ok := e.(thread.CompactionEntry); ok {
fmt.Println("compacted:", c.Reason)
}
}
Output: compacted: threshold
func WithApprover ¶
func WithApprover(a Approver, timeout time.Duration) SessionOption
WithApprover sets the decision chain's live step (ADR 0021 §2) and the time it is given. The Approver is consulted for a call about to park, after grants and before parking, under a context with the timeout as its deadline; a timeout or a panic reads as a decline, audited, and the call parks. The timeout is part of the option because an Approver without one would either never be consulted or block the turn forever: it must be positive — Create and Open fail on a non-positive timeout rather than leave the Approver silently unconsulted. A session with no interactive channel sets no Approver at all. A nil approver is ignored.
Example ¶
An Approver is the "ask now" step for a UI that is already connected: consulted before a call parks, under the timeout it is given. It declines what it will not decide, and that call parks like any other.
package main
import (
"context"
"fmt"
"strings"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// gatedDeploy is the agent of the approval examples: its one tool,
// deploy, needs a decision before it runs, and the scripted model
// plays turns.
func gatedDeploy(turns ...wefttest.Turn) *weft.Agent {
return weft.New(wefttest.Script(turns...),
weft.Tool("deploy", "Deploy the service.",
func(ctx context.Context, in struct {
Env string `json:"env"`
}) (string, error) {
return "deployed to " + in.Env, nil
},
weft.RequireApproval()))
}
// deployTo is a model turn asking to deploy to env.
func deployTo(env string) wefttest.Turn {
return wefttest.ToolCalls(wefttest.Call{Name: "deploy", Args: `{"env":"` + env + `"}`})
}
// settle waits for a turn and, when its boundary resumed on its own,
// for the resume too; it returns the last result.
func settle(t *thread.Turn) *weft.RunResult {
res, err := t.Wait()
if err != nil {
fmt.Println("turn failed:", err)
return nil
}
if next := t.Next(); next != nil {
return settle(next)
}
return res
}
func main() {
ctx := context.Background()
approver := func(ctx context.Context, r thread.Request) (thread.Decision, bool) {
if strings.Contains(string(r.Args), `"staging"`) {
d := thread.Approve(r.CallID)
d.Who = "terminal"
return d, true
}
return thread.Decision{}, false // not mine to decide: park it
}
agent := gatedDeploy(
deployTo("staging"), wefttest.Say("Staging is live."),
deployTo("prod"),
)
s, err := thread.Create(ctx, thread.Memory(), agent, thread.WithApprover(approver, 30*time.Second))
if err != nil {
fmt.Println(err)
return
}
t1, _ := s.Send(ctx, weft.User("Deploy to staging."))
fmt.Println("staging:", settle(t1).Text())
t2, _ := s.Send(ctx, weft.User("Now prod."))
settle(t2)
for _, r := range s.Pending() {
fmt.Println("parked:", r.Tool, string(r.Args))
}
}
Output: staging: Staging is live. parked: deploy {"env":"prod"}
func WithCompactor ¶
func WithCompactor(c Compactor) SessionOption
WithCompactor returns the SessionOption replacing the whole algorithm past the cut: the summary layers (SummaryModel, WithSummarizer, CheckSummary, PreferNative) do not run under it. A nil c is ignored.
func WithEstimator ¶
func WithEstimator(est Estimator) SessionOption
WithEstimator returns the SessionOption replacing the token estimator. A nil est is ignored.
func WithGrantStore ¶
func WithGrantStore(gs GrantStore) SessionOption
WithGrantStore returns the SessionOption adding the shared grant scope (ADR 0021 §4): grants many sessions consult, application-wide, behind the interface the application owns. The chain consults it after the session's own grants. A nil store is ignored.
func WithKeyring ¶
func WithKeyring(r *Keyring) SessionOption
WithKeyring returns the SessionOption making the session verify signed decisions under r — and mint signed challenges from r's active key (s.Request). A nil ring is ignored.
func WithLineage ¶
func WithLineage(parentSession, call string) SessionOption
WithLineage sets a new session's pool lineage (ADR 0022 §3): the parent session and the delegating call it grew from. A header option — Create and Fork write it into the header, which is immutable once stored; Open fails with ErrCreateOnly when given it (the lineage is read from the file, see Session.Lineage).
func WithMeta ¶
func WithMeta(meta map[string]string) SessionOption
WithMeta adds keys to a new session's header metadata — the caller metadata List's Query.Meta filter matches (ADR 0011 §1). A header option, merged over earlier ones: Create and Fork write it; Open fails with ErrCreateOnly when given it. SetInfo is the later-life edit, for every key outside the reserved "weft." prefix — those are set here, once, or never.
Example ¶
WithMeta writes header metadata at Create — what List's Query.Meta matches. SetInfo overlays later edits in Session.Meta, except under the reserved "weft." prefix, which only the header may set.
package main
import (
"context"
"errors"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// exampleIDs mints the given ids in order — the IDs option's usual
// shape in an example: the session id first, then one per entry.
func exampleIDs(ids ...string) thread.SessionOption {
next := 0
return thread.IDs(func() string {
id := ids[next]
next++
return id
})
}
func main() {
ctx := context.Background()
st := thread.Memory()
agent := weft.New(wefttest.Script())
s, err := thread.Create(ctx, st, agent, exampleIDs("s_acme", "e_info"),
thread.WithMeta(map[string]string{"tenant": "acme", "tier": "trial"}),
thread.PublicID("pub-7f3a"))
if err != nil {
fmt.Println(err)
return
}
if err := s.SetInfo(ctx, "", map[string]string{"tier": "paid"}); err != nil {
fmt.Println(err)
return
}
meta := s.Meta()
fmt.Println(meta["tenant"], meta["tier"], meta["weft.public_id"])
err = s.SetInfo(ctx, "", map[string]string{"weft.public_id": "pub-other"})
fmt.Println(errors.Is(err, thread.ErrReservedKey))
// List matches the header: the create-time values.
trial, _ := thread.List(ctx, st, thread.Query{Meta: map[string]string{"tier": "trial"}})
paid, _ := thread.List(ctx, st, thread.Query{Meta: map[string]string{"tier": "paid"}})
fmt.Println(trial.Total, paid.Total)
// A header option handed to Open is refused, not dropped.
_, err = thread.Open(ctx, st, "s_acme", agent, thread.WithMeta(map[string]string{"tenant": "other"}))
fmt.Println(errors.Is(err, thread.ErrCreateOnly))
}
Output: acme paid pub-7f3a true 1 0 true
func WithSummarizer ¶
func WithSummarizer(s Summarizer) SessionOption
WithSummarizer returns the SessionOption replacing text production — just the text; the cut and the entry stay the session's. SummaryModel and PreferNative do not apply under it; CheckSummary does (one retry, then the compaction fails). A nil s is ignored.
Example ¶
WithSummarizer swaps text production — just the text: the cut, the serialization and the entry stay the session's.
package main
import (
"context"
"fmt"
"strings"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// longSession creates a session whose history is three long user
// messages (ids e_0, e_1, e_2) — more than the default KeepRecent
// keeps raw — opened on agent with opts.
func longSession(agent *weft.Agent, opts ...thread.SessionOption) *thread.Session {
ctx := context.Background()
st := thread.Memory()
s, err := thread.Create(ctx, st, agent, opts...)
if err != nil {
panic(err)
}
parent := ""
var batch []thread.Entry
for i, topic := range []string{"order ", "invoice ", "refund "} {
id := fmt.Sprintf("e_%d", i)
batch = append(batch, thread.MessageEntry{ID: id, ParentID: parent, Created: time.Now().UTC(),
Message: weft.User(strings.Repeat(topic, 30_000/len(topic)))})
parent = id
}
if err := st.Append(ctx, s.ID(), batch...); err != nil {
panic(err)
}
if err := s.Close(ctx); err != nil {
panic(err)
}
s, err = thread.Open(ctx, st, s.ID(), agent, opts...)
if err != nil {
panic(err)
}
return s
}
// bulletSummarizer is a Summarizer that runs no model: it keeps the
// first words of every message in the range.
type bulletSummarizer struct{}
func (bulletSummarizer) Summarize(ctx context.Context, in thread.SummaryInput) (thread.Summary, error) {
var b strings.Builder
if in.PrevSummary != "" {
b.WriteString(in.PrevSummary + "\n")
}
for _, m := range in.Messages {
fmt.Fprintf(&b, "- %s: %s…\n", m.Role, strings.TrimSpace(m.Text()[:min(12, len(m.Text()))]))
}
return thread.Summary{Text: strings.TrimSpace(b.String())}, nil
}
func main() {
ctx := context.Background()
s := longSession(weft.New(wefttest.Script()), thread.WithSummarizer(bulletSummarizer{}))
if err := s.Compact(ctx); err != nil {
fmt.Println(err)
return
}
fmt.Println(s.Context()[0].Text())
}
Output: <weft-summary> - user: order order… </weft-summary>
func WithTrimmer ¶
func WithTrimmer(t Trimmer) SessionOption
WithTrimmer returns the SessionOption setting a custom trimmer. The default trimmer is off; ClearOldToolResults turns a built-in one on. A nil t is ignored.
type SharedGrant ¶
type SharedGrant struct {
}
A SharedGrant is one grant of a GrantStore, with the store's id for the audit trail.
type SignedDecision ¶
type SignedDecision struct {
Session string
// RequestID and RunID name the occurrence the challenge was minted
// for: the request entry's id and the run that parked the call. A
// signature answers that occurrence only — the same call id parked
// again by a later run is a different request.
RequestID string
RunID string
CallID string
Tool string
ArgsSHA256 string
// Expiry is the request's expiry as the challenge stated it. It
// must equal the request's own: the session enforces the request's
// expiry, never a time the signer chose.
Expiry time.Time
Kind Outcome
Reason string
Content string
// Who is a label for the audit trail. The decision's identity —
// what Quorum counts — is KeyID.
Who string
// Always approves and grants the same thing for the future, the
// signed ApproveAlways (Decision.Always carries it through
// SignDecision).
Always bool
Nonce string
KeyID string
MAC []byte
}
SignedDecision is a decision that crossed a process boundary (ADR 0021 §3): what the client decided over the challenge it was handed, vouched by the MAC the challenge's key computed. Build one with Keyring.Sign or SignDecision on the signing side; hand it to Session.DecideSigned on the session side. The fields are the challenge's claims — the verifier checks them against the pending request, so a stale or retargeted signature fails loudly.
func SignDecision ¶
func SignDecision(key []byte, r Request, d Decision) SignedDecision
SignDecision signs d over r's challenge under key — the client-side half of the exchange for a signer that holds a bare secret (Keyring.Sign is the same thing with the key looked up, Key.Sign the same as a key of the signer's own). r is the Request the session minted (its Nonce names the challenge) and key must be the secret of the key r.KeyID names; d is what the human decided; the MAC covers the challenge fields and the decision together, so neither can be altered in flight without failing verification. An empty d.Who is signed as the key's id.
type Storage ¶
type Storage interface {
// Create writes a session's header as its first line. The header's
// ID must satisfy ValidID — an id is a path component in some
// backend, so every backend vets it — a zero Weft means the current
// FormatVersion, any other envelope is rejected, and creating a
// session that already exists fails with ErrExists: a session is
// never silently replaced.
Create(ctx context.Context, h Header) error
// Append adds entries to a session in arrival order, after the
// entries it already holds. A batch is validated whole — every
// entry must encode — before anything is written, and is atomic
// against the writer's death: after a crash, a later Load returns
// all of the batch or none of it, never a prefix that decodes as
// fewer entries. What it is not is isolated from a reader in
// another process on a file backend: such a reader can catch the
// write in flight and see a torn final line, which Load drops and
// reports (LoadReport.Torn) and Watch waits out. A writer that
// finds a torn tail left by a crashed predecessor removes it before
// its first append, so new entries never join a dead writer's
// half-line. Appending to a session the storage does not hold fails
// with ErrNotFound; one another writer holds, with ErrLocked.
Append(ctx context.Context, session string, entries ...Entry) error
// Load returns a session's header and every entry in append order,
// and a LoadReport naming anything the load had to drop or skip —
// nil when the load was clean. The returned values never alias the
// storage: they are the caller's to keep and to mutate. An unknown
// session fails with ErrNotFound; data the storage cannot decode
// fails loudly — ErrNewerFormat for the unknown kind and the newer
// version, ErrCorrupt for the malformed line (ADR 0011 §5).
Load(ctx context.Context, session string) (Header, []Entry, *LoadReport, error)
// List returns a page of session headers — headers only, never
// entries: a list body that read whole sessions would be the
// storage bloat every surveyed store had to walk back (ADR 0010's
// listTracesLight lesson, the same shape here). See Query for the
// cursor and the limit.
List(ctx context.Context, q Query) (Page, error)
// Delete removes a session and its entries; an unknown session
// fails with ErrNotFound, and one another Storage or process holds
// as its writer, with ErrLocked. A hold or a lease of this Storage
// value's own does not refuse it (see One writer). History is
// otherwise forever — nothing else in the interface removes data.
Delete(ctx context.Context, session string) error
}
Storage is a session backend: Memory in process, jsonl on disk, sqlite in one database file (its own module) — all running the threadtest conformance table. Implementations must be safe for concurrent use: one writer per session is the policy (ADR 0011 §5), and readers may read at any time. context.Context comes first in every signature; a canceled context fails the call before anything is written.
One writer ¶
Who "one writer" is depends on who is asking. A Storage value used directly, without Sessions, enforces one writer per instance: its methods name a session, not a caller, so every goroutine using the value is the same writer, and the writer a backend refuses with ErrLocked is another Storage value or another process. Sessions enforce one writer per Session: a Session takes the session's lease through the optional Leaser capability before it writes, and a second Session on the same Storage value is refused. A backend that does not implement Leaser gives Sessions no such check — two Sessions on one value of it both write, each blind to the other — and keeps whatever rule its own lock enforces.
Delete is not refused by a lease: through the Storage value that holds a session — the one its writer uses — Delete removes the session and ends the lease with it, and that writer's next write fails with ErrNotFound. Through any other Storage value, Delete of a session a writer holds fails with ErrLocked.
func Memory ¶
func Memory(opts ...OpenOption) Storage
Memory returns a ready-to-use in-process session storage: sessions in a map behind a mutex, gone when the process is. For tests, examples, and as the reference behaviour the durable backends are compared against — every backend runs the same threadtest table, corruption rows included (Memory implements the table's threadtest.RawInjector hooks by holding the raw bytes).
opts are the open vocabulary every backend takes. Salvage and OpenLogger mean here what they mean on disk: a malformed line is skipped and reported instead of failing the load, and the one repair Memory makes on its own — a torn tail removed before an append — is reported to the given logger (slog.Default() as it stands at the call, without the option). The fsync options and NoLock are accepted and are no-ops: nothing here is buffered, and there is no cross-process lock to turn off.
There is no second Storage or process to refuse — one map, one process — so the Storage methods never answer ErrLocked; the one writer Memory tells from another is a lease holder (the Leaser capability, which is how two Sessions on one Memory are kept to one writer), and its Release (the Releaser capability) ends that lease.
Example ¶
Memory is the in-process Storage: create a session, append entries, load them back — the shape every backend shares (threadtest pins the rest).
package main
import (
"context"
"fmt"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
)
func main() {
ctx := context.Background()
st := thread.Memory()
h := thread.Header{
ID: "s_demo",
Created: time.Date(2026, 9, 28, 12, 0, 0, 0, time.UTC),
Meta: map[string]string{"room": "table-1"},
}
if err := st.Create(ctx, h); err != nil {
fmt.Println(err)
return
}
if err := st.Append(ctx, h.ID, thread.MessageEntry{
ID: "e_demo1",
Created: h.Created.Add(time.Second),
Message: weft.User("Where is order 1234?"),
}); err != nil {
fmt.Println(err)
return
}
header, entries, report, err := st.Load(ctx, h.ID)
if err != nil {
fmt.Println(err)
return
}
fmt.Println(header.ID, header.Meta["room"], len(entries), report)
for _, e := range entries {
switch e := e.(type) {
case thread.MessageEntry:
fmt.Println(e.Message.Role, "asks:", e.Message.Text())
}
}
}
Output: s_demo table-1 1 <nil> user asks: Where is order 1234?
type Summarizer ¶
type Summarizer interface {
Summarize(ctx context.Context, in SummaryInput) (Summary, error)
}
Summarizer produces the summary text and nothing else: the cut, the serialization, the entry and the retry chain stay the session's. The default runs a model (SummaryModel, then the session's own). A Summarizer that knows its output was cut short returns an error wrapping ErrSummaryTruncated. It runs without the session lock and may call the session; a panic in it fails the compaction.
type Summary ¶
type Summary struct {
// Text is the summary. Empty text is a failed summary.
Text string
// Usage and Model are what the summary cost and which model made
// it — the cost ledger's inputs; a custom Summarizer that runs no
// model leaves them zero.
Usage core.Usage
Model core.ModelInfo
}
Summary is a summarizer's output: the text, and what it cost.
type SummaryInput ¶
type SummaryInput struct {
Messages []core.Message
PrevSummary string
Instructions string
MaxTokens int64
Reason Reason
SystemPrompt string
}
SummaryInput is what a Summarizer is asked to summarize: Messages is the range, already serialized for summarizing (tool results capped, signed reasoning dropped, files as names); PrevSummary is the previous summary when one exists — the iterative chain's last link, which the new summary must carry forward; Instructions is the per-call text (SummaryInstructions, possibly edited by BeforeCompact); MaxTokens the output cap; Reason why the compaction runs. SystemPrompt is the prompt the default summarizer would send — the skeleton or SummaryPrompt, then SummaryFocus and Instructions — handed over so a custom Summarizer can reuse it; it may ignore it.
type TriggerInput ¶
TriggerInput is what a TriggerFunc decides on: the provider-reported input of the last model step, the estimated tokens added since, and the window and reserve in force. The reported number is the signal; only the delta is estimated (ADR 0020 §2).
type TrimRecord ¶
type TrimRecord struct {
Stubs []TrimStub `json:"stubs"`
}
TrimRecord is what a trim did, persisted so the context replays it from the file alone (ADR 0020, amendment 2026-10-01): every tool result the trimmer replaced, with the replacement text. The walk applies exactly these stubs — it never consults the session's configured Trimmer — so a recorded trim reads the same under any options, in any process.
type TrimStub ¶
type TrimStub struct {
Entry string `json:"entry"`
CallID string `json:"call_id"`
Content string `json:"content"`
IsError bool `json:"is_error,omitempty"`
}
TrimStub is one stubbed tool result: the entry that holds it, the call it answers, and the content the model sees in its place. The stored result is untouched; IsError is what the stubbed part reports (false for the built-in trimmer's stub).
type Trimmer ¶
Trimmer is the cheap pre-pass: before summarizing, replace old tool results with a stub, and the context may fit again (Anthropic clear_tool_uses, Vercel pruneMessages). It runs on the automatic path only — a manual Compact summarizes.
Trim receives the context the model is shown (a copy it may edit) and returns the trimmed context. The session diffs the two and persists the difference as the compaction entry's trim record, which is what every later context build replays — the Trimmer itself is never consulted on read. The representable change is exactly one: replace a tool result's Content (and IsError); the messages, their order, every other part, and each result's CallID and Name must come back as given. Any other change fails the trim with ErrInvalidCompaction — reported through the logger and CompactFailed — and the summary compaction runs instead. A trim that changes nothing, or that does not bring the estimated context under window − Reserve, writes nothing and the summary compaction runs.
Trim runs without the session lock and may call the session; an error or a panic is logged and the summary compaction runs.
type Turn ¶
type Turn struct {
// contains filtered or unexported fields
}
A Turn is the receipt and the handle of one accepted Send — or of a resume run over a parked approval boundary, which has no prompt of its own: its receipt names its turn entry. It is safe for concurrent use.
func (*Turn) Done ¶
func (t *Turn) Done() <-chan struct{}
Done returns a channel that is closed when the turn ends — the select-shaped form of Wait:
select {
case <-turn.Done():
res, err := turn.Wait() // returns at once
case <-ctx.Done():
// stop waiting; the turn keeps running
}
The turn is decided when the channel closes: its entries have landed and Outcome is final. Work the session does between turns — the automatic compaction — happens after, and never delays it.
Example ¶
Done is the channel form of Wait, for a select: wait for the turn or for something else, whichever comes first. Giving up the wait does not stop the turn.
package main
import (
"context"
"fmt"
"sync"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// slowTool returns a tool that blocks until release is closed or its
// context ends, and a channel closed when its first call has started —
// the examples' stand-in for a long piece of work.
func slowTool(name string) (tool *weft.ToolDef, started, release chan struct{}) {
started, release = make(chan struct{}), make(chan struct{})
var once sync.Once
tool = weft.Tool(name, "Takes a while.", func(ctx context.Context, _ struct{}) (string, error) {
once.Do(func() { close(started) })
select {
case <-release:
return name + " finished", nil
case <-ctx.Done():
return "", ctx.Err()
}
})
return tool, started, release
}
func main() {
ctx := context.Background()
work, started, release := slowTool("work")
agent := weft.New(wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "work", ID: "call_1"}),
wefttest.Say("Finished."),
), work)
s, _ := thread.Create(ctx, thread.Memory(), agent)
turn, _ := s.Send(ctx, weft.User("Do the long job."))
<-started
select {
case <-turn.Done():
fmt.Println("done already")
case <-time.After(10 * time.Millisecond):
fmt.Println("still running:", turn.Outcome())
}
close(release)
<-turn.Done()
res, err := turn.Wait() // returns at once
fmt.Println(res.Text(), err)
}
Output: still running: running Finished. <nil>
func (*Turn) Events ¶
Events yields the turn's run events in emission order, forwarded by the session as it observes the run. Unlike the core's Run, a Turn's Events may be ranged more than once — the session, not the consumer, drives the run, so abandoning the stream mid-way (breaking out of the range) never cancels it; the turn's end wakes a suspended range and ends it. A failed run delivers its error exactly once as the final element, the core's rule. A turn with no run of its own — a steer, a turn that never started — yields nothing.
Example ¶
Events yields the turn's run events as the session observes them. The session drives the run, not the reader: breaking out of the range cancels nothing, and Wait still returns the result.
package main
import (
"context"
"fmt"
"strings"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
agent := weft.New(
wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "lookup", ID: "call_1"}),
wefttest.Say("Order 1234 shipped."),
),
weft.Tool("lookup", "Look up an order.", func(context.Context, struct{}) (string, error) {
return "shipped", nil
}),
)
s, _ := thread.Create(ctx, thread.Memory(), agent)
turn, _ := s.Send(ctx, weft.User("Where is order 1234?"))
var text strings.Builder
for ev, err := range turn.Events() {
if err != nil {
fmt.Println("run failed:", err) // at most once, as the last element
break
}
switch ev := ev.(type) {
case weft.ToolStart:
fmt.Println("tool:", ev.Name)
case weft.ToolFinish:
fmt.Println("result:", ev.Content)
case weft.TextDelta:
text.WriteString(ev.Text)
case weft.RunFinish:
fmt.Println("steps:", ev.Steps)
}
}
fmt.Println("text:", text.String())
res, _ := turn.Wait()
fmt.Println("same answer from Wait:", res.Text() == text.String())
}
Output: tool: lookup result: shipped steps: 2 text: Order 1234 shipped. same answer from Wait: true
func (*Turn) ID ¶
ID returns the turn's receipt id: the id of the entry that makes the Send durable and names the turn in the tree.
- A send that started at once: its prompt entry, appended and flushed before Send returned.
- A send accepted while the session was busy (Queue, Interrupt, Rollback, and a steer's follow-up): the id its prompt entry takes when its turn starts. Until then no entry has this id; the accepted receipt entry — durable from the moment Send returned — records it in ReceiptEntry.Turn.
- A steer (Send under the Steer policy on a busy session): its queued receipt entry.
- A resume: its turn entry, written when the resume ends (a resume writes no prompt of its own).
func (*Turn) Next ¶
Next returns the turn that continues this one, or nil when there is none: the resume the session started for the approval boundary this turn parked (TurnParked), or the follow-up turn a deferred steer became (TurnDeferred). It is nil for every other outcome, and for a parked turn whose boundary has not been resumed yet. Following Next is how a caller watches an approval flow through: Send's turn parks, and the resume links here whichever path armed it — the decision chain, a Decide or DecideSigned that completed the boundary, Resume — so the Turn those calls return is this same Turn, and its Wait is the conversation's continuation. One boundary resumes at most once and a steer defers at most once, so the link is set at most once. The link is the live Session's: after a restart the parked Turn value is gone, and the resume is reached through what Decide, Resume or Continue return.
Example ¶
A steer that cannot be delivered — here the session is waiting on an approval, and a steer never decides a parked call — becomes a follow-up turn. Turn.Next is that turn: it runs once the boundary resolves.
package main
import (
"context"
"fmt"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
func main() {
ctx := context.Background()
agent := weft.New(
wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "deploy", ID: "call_1"}),
wefttest.Say("Not deployed."),
wefttest.Say("Rolled back to v41."),
),
weft.Tool("deploy", "Deploy the service.", func(context.Context, struct{}) (string, error) {
return "deployed", nil
}, weft.RequireApproval()),
)
s, _ := thread.Create(ctx, thread.Memory(), agent)
parked, _ := s.Send(ctx, weft.User("Deploy v42."))
if _, err := parked.Wait(); err != nil {
fmt.Println(err)
return
}
fmt.Println("first turn:", parked.Outcome())
steer, _ := s.Send(ctx, weft.User("Actually, roll back to v41."), thread.As(thread.Steer))
_, _ = steer.Wait()
fmt.Println("steer:", steer.Outcome())
followUp := steer.Next()
resume, _ := s.Decide(ctx, thread.Deny("call_1", "superseded"))
if _, err := resume.Wait(); err != nil {
fmt.Println(err)
return
}
res, err := followUp.Wait()
if err != nil {
fmt.Println(err)
return
}
fmt.Println("follow-up:", res.Text())
}
Output: first turn: parked steer: deferred follow-up: Rolled back to v41.
func (*Turn) Outcome ¶
func (t *Turn) Outcome() TurnOutcome
Outcome reports how the turn ended, or TurnRunning while it has not. It is what tells apart the ends Wait reports alike: a delivered steer from a deferred one (both nil, nil), a canceled turn from a failed one.
func (*Turn) RunID ¶
RunID returns the run's id, <session>-t<n> — the key the run store holds the run's records under, unique across a reopen. A turn that re-ran after an overflow reports the re-run's id (the run that produced its outcome). A steer has no run of its own: it reports the id minted for it at acceptance, which no run uses.
func (*Turn) Wait ¶
Wait blocks until the turn ends and returns its result and error. The session drains the run itself, so Wait alone observes every event; it is safe to call after or while ranging over Events, from any number of goroutines, any number of times.
The result is the run's on success — including a run that ended at an approval boundary (RunResult.Pending) — and nil for a turn with no run of its own (a steer: see Outcome).
The error says what kind of end it was; match with errors.Is and errors.As:
- *core.RunError — the run started and failed, exactly as the core's Wait reports it: Result carries the partial transcript, Unwrap the cause (a canceled run's is context.Canceled). After a failed overflow re-run the error joins both attempts'.
- ErrNotRun — the turn's run never started: its context ended first (the error also wraps context.Canceled or context.DeadlineExceeded), its prompt entry could not be written (the storage's error), or — a resume — the boundary was gone (ErrNotPending) or its audit entry could not be written.
- ErrClosed — the session closed before the turn ran.
- ErrDropped — ClearQueue removed the message.
- ErrTurnPanicked — the session's own turn machinery, or the caller's IDs or Clock function, panicked.
- ErrNotPersisted — the run ended but its end could not be written. It is joined to the run's own error when there is one; when the run succeeded the result is returned beside it.
func (*Turn) WaitContext ¶
WaitContext is Wait bounded by ctx: it returns the turn's result and error when the turn ends, or nil and ctx.Err() as soon as ctx does. A turn that has already ended is reported whatever ctx says — the answer is there, and no waiting is needed for it. Giving up the wait changes nothing for the turn — it keeps running, and a later Wait or WaitContext still returns its end. To stop the turn itself, cancel the context its Send was given.
Example ¶
WaitContext bounds the wait, not the turn: when the context ends first the caller gets its error back and the turn keeps running.
package main
import (
"context"
"errors"
"fmt"
"sync"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/wefttest"
)
// slowTool returns a tool that blocks until release is closed or its
// context ends, and a channel closed when its first call has started —
// the examples' stand-in for a long piece of work.
func slowTool(name string) (tool *weft.ToolDef, started, release chan struct{}) {
started, release = make(chan struct{}), make(chan struct{})
var once sync.Once
tool = weft.Tool(name, "Takes a while.", func(ctx context.Context, _ struct{}) (string, error) {
once.Do(func() { close(started) })
select {
case <-release:
return name + " finished", nil
case <-ctx.Done():
return "", ctx.Err()
}
})
return tool, started, release
}
func main() {
ctx := context.Background()
work, started, release := slowTool("work")
agent := weft.New(wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "work", ID: "call_1"}),
wefttest.Say("Finished."),
), work)
s, _ := thread.Create(ctx, thread.Memory(), agent)
turn, _ := s.Send(ctx, weft.User("Do the long job."))
<-started
short, cancel := context.WithTimeout(ctx, 10*time.Millisecond)
defer cancel()
_, err := turn.WaitContext(short)
fmt.Println("gave up waiting:", errors.Is(err, context.DeadlineExceeded))
close(release)
res, err := turn.WaitContext(ctx)
fmt.Println(res.Text(), err)
}
Output: gave up waiting: true Finished. <nil>
type TurnEntry ¶
type TurnEntry struct {
ID string `json:"id"`
ParentID string `json:"parent,omitempty"`
Created time.Time `json:"created"`
RunID string `json:"run_id"`
StopReason core.StopReason `json:"stop_reason,omitempty"`
Usage core.Usage `json:"usage"`
Steps int `json:"steps,omitempty"`
Err string `json:"err,omitempty"`
Pending []core.ToolCallPart `json:"pending,omitempty"`
// Canceled marks a turn whose context ended — canceled (the caller
// walked away, an Interrupt or Rollback send, Close) or past its
// deadline. Err says which.
Canceled bool `json:"canceled,omitempty"`
// Policy is the busy policy the turn's Send was called under —
// "queue", "reject", "steer", "interrupt" or "rollback"
// (Policy.String) — captured at Send, the session's or the Send's
// own (As). A steer's follow-up turn records "steer". Empty on a
// resume turn, which no Send started, and on entries written
// before the field existed.
Policy string `json:"policy,omitempty"`
// ReRun, on the entry of an overflow attempt that was re-run, is
// the run id of the re-run — written before the re-run starts.
ReRun string `json:"rerun,omitempty"`
// LastInput is the compaction trigger's baseline for the turn:
// the provider-reported input of the run's final model step, plus
// the estimated tokens of the tail that report cannot cover (the
// final step's own messages) — one number, recovered identically
// by a live session and a reopen (ADR 0020 §2: reported tokens
// are the signal; only what the report cannot cover is estimated).
// Absent on turns that ran no step.
LastInput int64 `json:"last_input,omitempty"`
// LateSteps counts the per-step appends that failed while the turn
// ran (ADR 0011 §7): each batch was held and written late — by the
// next step's append or by this entry's own batch — so nothing was
// lost, but for that long a crash would have lost messages the run
// had already emitted. Zero, and absent on the wire, for a turn
// whose every step landed as it joined.
LateSteps int `json:"late_steps,omitempty"`
}
TurnEntry is the per-turn ledger: the run's id (<session>-t<n>), its stop reason, usage and step count, the error text on failure, the calls left pending on the approval boundary, and whether the turn was canceled. It records the turn's outcome, never its content — the messages are MessageEntries — and it does not enter the model's context.
One run, one turn entry. A Send is one entry — except a turn that overflowed and re-ran (ADR 0020 §5), which leaves two: the failed attempt's, under the attempt's run id, with the overflow in Err, the usage of the steps it completed, and ReRun naming the run that followed; then the re-run's own. Session.Usage sums both: the attempt's tokens were spent.
func (TurnEntry) MarshalJSON ¶
MarshalJSON encodes the entry with its "type" discriminator.
type TurnOutcome ¶
type TurnOutcome int
TurnOutcome is how a Turn ended — what Turn.Outcome reports. A Turn stands for one accepted Send, whatever became of it: a run of its own, a message delivered into another turn's run, a follow-up, or nothing at all.
const ( // TurnRunning is the outcome of a turn that has not ended: queued, // running, or — a steer — waiting for its fate. TurnRunning TurnOutcome = iota // TurnAnswered: the turn's run completed; Wait returns its result. TurnAnswered // TurnParked: the turn's run ended at an approval boundary — a // success whose result lists the calls awaiting a decision // (RunResult.Pending). Turn.Next is the resume, once one is armed. TurnParked // TurnDelivered: a steer the running turn's run drained. The // message is in that run's transcript, and the answer is that // turn's, not this one's: Wait returns a nil result and a nil // error. A delivery the run never got to answer — it failed // before another model step — is marked on the receipt entry // (ReceiptEntry.Unanswered). TurnDelivered // TurnDeferred: a steer that became a follow-up turn — it met an // intended end, an approval boundary, or a run that ended before // its drain. Wait returns nil, nil; Turn.Next is the follow-up. TurnDeferred // TurnDropped: ClearQueue removed the message before it reached a // model. Wait returns an error wrapping ErrDropped. TurnDropped // TurnFailed: the turn ended with an error that is not a // cancellation — its run failed, it could not start, or its end // could not be persisted. Wait returns the error. TurnFailed // TurnCanceled: the turn's context ended (canceled, or past its // deadline), an Interrupt or Rollback send felled it, or the // session closed before it could run. Wait returns the error. TurnCanceled )
func (TurnOutcome) String ¶
func (o TurnOutcome) String() string
String returns the outcome's name: "running", "answered", "parked", "delivered", "deferred", "dropped", "failed" or "canceled".
type Usage ¶
type Usage struct {
// Turns sums every turn entry's usage — every run the session ever
// made, abandoned branches included: the tokens were spent.
Turns core.Usage
// Summaries sums the summarizer usage of every compaction entry —
// the cost of keeping the context small, in its own bucket.
Summaries core.Usage
// Delegated sums the usage every settled pool receipt carries —
// the cost of work handed to thread/pool children (ADR 0022 D3),
// in its own bucket: a delegated child's tokens are not this
// session's turns, and the ledger never mixes kinds.
Delegated core.Usage
}
Usage is a session's cost ledger (ADR 0020 §4): what its turns cost and what its summaries cost, never mixed.
type Verdict ¶
type Verdict struct {
// contains filtered or unexported fields
}
Verdict is what a BeforeCompact hook returns: Proceed with the computed compaction, Cancel it, or Replace it with a hook-made one (recorded from_hook). The zero Verdict proceeds.
func Cancel ¶
func Cancel() Verdict
Cancel returns the Verdict that stops the compaction: nothing is written, and the call fails with ErrCompactCanceled (the automatic path stays quiet about it).
func Proceed ¶
func Proceed() Verdict
Proceed returns the Verdict that runs the compaction as prepared — with whatever the hook edited on the Preparation.
func Replace ¶
func Replace(c *Compaction) Verdict
Replace returns the Verdict that writes c instead — the hook's own summary, recorded with Reason from_hook. c.FirstKept must name an entry the session holds; ApplyCompaction's rules decide the rest.
type Watcher ¶
type Watcher interface {
Watch(ctx context.Context, session string, after string) (iter.Seq2[Entry, error], error)
}
Watcher is the optional Storage capability that tails a session as it is appended to — the same small-interface rule; jsonl and sqlite implement it. Watch yields the session's entries in arrival order, each exactly once, starting after the entry named by after (empty — from the beginning), then keeps yielding as appends land, and ends when ctx is done. A session the storage does not hold fails with ErrNotFound before the first yield, as does an after the session does not hold (a plain error). The stream ends with one terminal error when the session is deleted, or deleted and created again, under the watcher (ErrNotFound: the session being tailed is gone), or when it meets data it cannot decode (ErrNewerFormat, or ErrCorrupt naming the line — skipped instead under Salvage). A torn final line is a write in flight: it is yielded once complete, never half. Readers never lock: a watcher is a reader that waits, and the consumer is free to call the same Storage from inside the loop.
Example ¶
A watcher tails a session as it is appended to — the live tail a second process reads while the first writes. The stream yields the backlog in arrival order, then each new entry exactly once, and ends when its context is canceled.
package main
import (
"context"
"fmt"
"os"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/thread"
"github.com/weftgo/weft/thread/jsonl"
)
func main() {
ctx := context.Background()
dir, err := os.MkdirTemp("", "weft-watch")
if err != nil {
fmt.Println(err)
return
}
defer func() { _ = os.RemoveAll(dir) }()
st, err := jsonl.Open(dir)
if err != nil {
fmt.Println(err)
return
}
if err := st.Create(ctx, thread.Header{ID: "s_demo", Created: time.Now().UTC()}); err != nil {
fmt.Println(err)
return
}
if err := st.Append(ctx, "s_demo",
thread.MessageEntry{ID: "e_1", Created: time.Now().UTC(), Message: weft.User("one")}); err != nil {
fmt.Println(err)
return
}
watch := st.(thread.Watcher)
wctx, cancel := context.WithCancel(ctx)
defer cancel()
seq, err := watch.Watch(wctx, "s_demo", "")
if err != nil {
fmt.Println(err)
return
}
go func() {
for e, err := range seq {
if err != nil {
return
}
fmt.Println("watched:", e.(thread.MessageEntry).Message.Text())
}
}()
if err := st.Append(ctx, "s_demo",
thread.MessageEntry{ID: "e_2", Created: time.Now().UTC(), Message: weft.User("two")}); err != nil {
fmt.Println(err)
return
}
time.Sleep(600 * time.Millisecond) // let the tail catch both
cancel()
time.Sleep(50 * time.Millisecond)
}
Output: watched: one watched: two
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package backend is for authors of thread.Storage backends: it resolves the open options an application passes — thread.Salvage, thread.FsyncOnFlush, thread.NoLock, thread.OpenLogger — into the configuration a backend acts on.
|
Package backend is for authors of thread.Storage backends: it resolves the open options an application passes — thread.Salvage, thread.FsyncOnFlush, thread.NoLock, thread.OpenLogger — into the configuration a backend acts on. |
|
examples
|
|
|
approvals
command
Command approvals walks one weft/thread session through the approval flow (ADR 0021): a gated call parks, the process "restarts" — the session is reopened from the JSONL file — a decision arrives signed over the challenge the session minted, and the conversation resumes under it.
|
Command approvals walks one weft/thread session through the approval flow (ADR 0021): a gated call parks, the process "restarts" — the session is reopened from the JSONL file — a decision arrives signed over the challenge the session minted, and the conversation resumes under it. |
|
session
command
Command session walks one weft/thread session through its whole life: two turns, a label, a branch off the first answer, a fork of the branch, a manual compaction with preview, a close, and a reopen from disk — everything on a JSONL backend, everything offline through a scripted model, every id deterministic.
|
Command session walks one weft/thread session through its whole life: two turns, a label, a branch off the first answer, a fork of the branch, a manual compaction with preview, a close, and a reopen from disk — everything on a JSONL backend, everything offline through a scripted model, every id deterministic. |
|
internal
|
|
|
carry
Package carry builds the context a run continues on when it runs on behalf of another one — a resume over a parked boundary, a deferred steer's follow-up, an async pool child: cancellation and deadline from one context, and, for the keys that context does not carry, the values of the context the work began on.
|
Package carry builds the context a run continues on when it runs on behalf of another one — a resume over a parked boundary, a deferred steer's follow-up, an async pool child: cancellation and deadline from one context, and, for the keys that context does not carry, the values of the context the work began on. |
|
opencfg
Package opencfg holds the resolved form of thread's open options — the one piece of the option vocabulary that thread (which declares the options), its in-tree backends (Memory) and thread/backend (which publishes the resolution to backend authors) all need, kept here so that none of them has to import another to share it.
|
Package opencfg holds the resolved form of thread's open options — the one piece of the option vocabulary that thread (which declares the options), its in-tree backends (Memory) and thread/backend (which publishes the resolution to backend authors) all need, kept here so that none of them has to import another to share it. |
|
rules
Package rules holds the small rules every thread.Storage backend must answer identically: the List limit, the metadata and title filters, the paging cursor, and the line discipline of a session's bytes.
|
Package rules holds the small rules every thread.Storage backend must answer identically: the List limit, the metadata and title filters, the paging cursor, and the line discipline of a session's bytes. |
|
Package jsonl is the default durable thread.Storage: one directory, one <id>.jsonl file per session — the header line first, then one entry per line in append order (ADR 0011).
|
Package jsonl is the default durable thread.Storage: one directory, one <id>.jsonl file per session — the header line first, then one entry per line in append order (ADR 0011). |
|
Package pool provides bounded concurrent child runs for thread sessions (ADR 0022): a FIFO semaphore per Pool value, subagent tools whose children run as sessions of their own, linked to the parent session and call, and receipts recording every delegation's journey.
|
Package pool provides bounded concurrent child runs for thread sessions (ADR 0022): a FIFO semaphore per Pool value, subagent tools whose children run as sessions of their own, linked to the parent session and call, and receipts recording every delegation's journey. |
|
Package sqlite is the thread's second durable backend: every session in one SQLite file on the CGO-free modernc.org/sqlite driver — the store's choice (store/sqlite, Crush's before it), reused so one dependency serves both modules — with WAL and embedded migrations.
|
Package sqlite is the thread's second durable backend: every session in one SQLite file on the CGO-free modernc.org/sqlite driver — the store's choice (store/sqlite, Crush's before it), reused so one dependency serves both modules — with WAL and embedded migrations. |
|
Package threadtest is the shared conformance table for thread.Storage backends — the executable form of ADR 0011 §5's promises.
|
Package threadtest is the shared conformance table for thread.Storage backends — the executable form of ADR 0011 §5's promises. |