pool

package
v0.13.0 Latest Latest
Warning

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

Go to latest
Published: Oct 10, 2026 License: MIT Imports: 12 Imported by: 0

Documentation

Overview

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. The pool imports the root module and weft/thread, never a persistence module (ADR 0011 §1's module rule, one level down), and changes nothing in the core: it is a tool middleware, session entries, and pool-owned goroutines.

The bound

A Pool admits at most max child runs doing work at once — model calls and tool calls. The bound is the Pool value's: two pools are two bounds, and a process that wants one bound shares one Pool. A run that is only waiting holds no slot: a child whose sync delegation is running a child of its own gives its slot up for the wait and queues for one again before it works again, and a child parked at an approval holds none. That hand-off is what makes fan-out times depth safe on any max, 1 included; the depth of a delegation chain has its own limit (MaxDepth).

Admission is first in, first out, in the order work queued: Submit queues before it returns, a wrapped tool queues when it is called, and a parent taking its slot back queues behind whatever was already waiting.

Costs

A pool child is a session, not a tool-call child of the core's (core.Subagent), so the core's rule "a child's usage is yours" does not reach it: its usage is on neither the parent's RunResult.Usage nor StepRecord.SubagentUsage. It is on the receipt, and thread.Session.Usage sums the settled receipts into its Delegated bucket. That is the pool's one exception to the core's roll-up (AGENTS.md rule 13), and a budget that must cover delegated work reads Delegated.

Lineage

A child is linked to its parent by reference — its header names the parent session and the delegating call; the parent's receipts name the child — and by nothing else. Deleting the parent deletes no child: Children and Descendants enumerate them so an application can cascade, children first. A child still running when its parent is deleted runs to its end and cannot record its settlement (the failure is logged). A Branch in the parent moves no child: receipts are ledger, read from every line of the tree, and a parked child's requests stay pending wherever the leaf is; a delegating call abandoned by the branch is simply no longer there to resolve, and the child's answer stays on its receipt. A Fork copies the ledger but not the delegations: the fork's copies of unsettled receipts are settled canceled and mirrored requests are left out (the core's rule), and the children stay the origin's — Children does not list them for the fork.

Index

Examples

Constants

View Source
const (
	// CodeSubagentDepth marks a delegation refused because it would
	// exceed the pool's MaxDepth — refused before any child is
	// created. The core has no depth limit (ADR 0014); the code is
	// the pool's.
	CodeSubagentDepth = "SUBAGENT_DEPTH"
	// CodeSubagentCanceled marks a sync delegation whose child the
	// pool canceled — Cancel on its receipt, or the pool's Close —
	// as against the parent's own cancellation, which ends the
	// parent's run instead.
	CodeSubagentCanceled = "SUBAGENT_CANCELED"
)

The model-visible texts the pool produces (AGENTS.md rule 5: a contract; ADR 0022's amendment lists them, TestModelVisibleTexts pins them byte for byte). Everything else a wrapped tool returns is the child's own answer.

View Source
const DefaultMaxDepth = 8

DefaultMaxDepth is the delegation depth a Pool allows when MaxDepth is not given: a child of a child, eight levels down.

Variables

View Source
var ErrClosed = errors.New("thread/pool: pool is closed")

ErrClosed is returned for work handed to a pool whose Close has run — Submit, Decide, Recover, and a wrapped tool's call, in or outside a session: an accepted delegation outlives the pool that would run it, so a closing pool takes no new work.

View Source
var ErrCycle = errors.New("thread/pool: agent is already running in this call chain")

ErrCycle is returned by Submit for a delegation to an agent that is already running above the delegating run; a wrapped tool reports the same refusal to the model as SUBAGENT_CYCLE.

View Source
var ErrDepth = errors.New("thread/pool: delegation depth exceeds the pool's limit")

ErrDepth is returned by Submit for a delegation that would exceed the pool's MaxDepth; a wrapped tool reports the same refusal to the model as SUBAGENT_DEPTH.

View Source
var ErrDuplicateWrap = errors.New("thread/pool: wrap name already in use")

ErrDuplicateWrap is returned by Wrap for a name the pool already wrapped a different agent under: the name is how a restarted process finds the agent a parked child resumes with, so one name names one agent.

View Source
var ErrNoAgent = errors.New("thread/pool: no agent for child session")

ErrNoAgent is returned by Decide, Recover and Cancel when a child session must be reopened and the pool holds no agent for it: no Wrap of this pool carries the name the child was made under, and nothing was registered for the session. Register the agent, or Wrap it again under its name, and call again.

View Source
var ErrNotRunning = errors.New("thread/pool: receipt is not running")

ErrNotRunning is what every *StateError matches under errors.Is: the receipt exists, and its child is not in the state the call needs. The StateError says which state it is in.

View Source
var ErrUnknownReceipt = errors.New("thread/pool: unknown receipt")

ErrUnknownReceipt is returned by Cancel, Forward and Wait for a receipt id the parent session's ledger does not hold — a wrong id, or another session's.

Functions

func Children

func Children(ctx context.Context, parent *thread.Session) ([]string, error)

Children returns the sessions parent delegated to, in acceptance order: the child of every receipt in its ledger whose stored header names parent as its lineage. The check is what keeps a fork honest — a Fork copies its origin's receipts, and the children stay the origin's — and a child that has been deleted is left out. Deleting a session deletes none of its children; this is the list to cascade over (Descendants gives the whole subtree).

func Descendants

func Descendants(ctx context.Context, parent *thread.Session) ([]string, error)

Descendants returns every session delegated from parent, directly or through a child, deepest first: each session appears after all of its own children, so deleting in the order returned — and the parent last — never leaves a child whose parent is already gone half-way through. It reads each descendant from parent's storage; a session that has been deleted is left out, with whatever was below it.

Example

Deleting a session deletes none of its children. Descendants lists the whole subtree, deepest first — the order to delete in, the parent last.

package main

import (
	"context"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

func main() {
	ctx := context.Background()
	st := thread.Memory()
	p := pool.New(1)
	leaf := weft.New(wefttest.Script(wefttest.Say("leaf answer")))
	mid := weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "leaf", ID: "call-leaf", Args: `{"prompt":"go"}`}),
		wefttest.Say("mid answer"),
	), p.MustWrap("leaf", "", leaf))
	top := weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "mid", ID: "call-mid", Args: `{"prompt":"go"}`}),
		wefttest.Say("done"),
	), p.MustWrap("mid", "", mid))
	s, _ := thread.Create(ctx, st, top)
	turn, _ := s.Send(ctx, weft.User("go"))
	if _, err := turn.Wait(); err != nil {
		panic(err)
	}
	children, _ := pool.Children(ctx, s)
	all, err := pool.Descendants(ctx, s)
	if err != nil {
		panic(err)
	}
	fmt.Println("children:", len(children), "descendants:", len(all))
	fmt.Println("deepest first:", all[len(all)-1] == children[0])

	if err := s.Close(ctx); err != nil {
		panic(err)
	}
	for _, id := range append(all, s.ID()) {
		if err := thread.Delete(ctx, st, id); err != nil {
			panic(err)
		}
	}
	page, _ := thread.List(ctx, st, thread.Query{})
	fmt.Println("sessions left:", page.Total)
}
Output:
children: 1 descendants: 2
deepest first: true
sessions left: 0

Types

type Option

type Option interface {
	// contains filtered or unexported methods
}

An Option configures a Pool at New.

func IDs

func IDs(id func() string) Option

IDs sets the function minting ids for the child sessions the pool creates, in place of the thread module's own time-sortable ids — the determinism knob tests and examples use. It is handed to each child as its thread.IDs option, so it mints the child's session id and then every entry id inside that child; ids in the parent session stay the parent's own. It is called from many goroutines and under session locks: it must be safe for concurrent use, return quickly, and never repeat an id.

func MaxDepth

func MaxDepth(n int) Option

MaxDepth sets how deep a chain of delegations through the pool may go: 1 allows children, 2 children of children. A delegation that would go deeper is refused before any child is created — a wrapped tool's call with SUBAGENT_DEPTH, Submit with ErrDepth. The default is DefaultMaxDepth. The limit is its own bound, unrelated to max: waiting parents hold no slot, so depth costs sessions and tokens, not capacity. New panics when n is below 1.

Example

MaxDepth bounds how deep delegation goes, separately from how many children run at once. Past it a wrapped tool refuses with SUBAGENT_DEPTH — data the delegating model reads.

package main

import (
	"context"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

func main() {
	ctx := context.Background()
	st := thread.Memory()
	p := pool.New(1, pool.MaxDepth(1)) // children, but no children of children
	leaf := weft.New(wefttest.Script(wefttest.Say("never asked")))
	mid := weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "deeper", ID: "call-deeper", Args: `{"prompt":"go on"}`}),
		wefttest.Say("could not delegate further"),
	), p.MustWrap("deeper", "", leaf))
	top := weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "delegate", ID: "call-delegate", Args: `{"prompt":"go"}`}),
		wefttest.Say("done"),
	), p.MustWrap("delegate", "", mid))
	s, _ := thread.Create(ctx, st, top)
	turn, _ := s.Send(ctx, weft.User("go"))
	if _, err := turn.Wait(); err != nil {
		panic(err)
	}
	// What the middle agent's model read when it tried to go deeper.
	midSession, err := thread.Open(ctx, st, pool.Receipts(s)[0].Child, mid)
	if err != nil {
		panic(err)
	}
	for _, m := range midSession.Context() {
		for _, part := range m.Content {
			if result, ok := part.(weft.ToolResultPart); ok {
				fmt.Println(result.Content)
			}
		}
	}
}
Output:
SUBAGENT_DEPTH: delegation depth 2 exceeds the pool's limit 1

type Pool

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

Pool is a bound on concurrent child runs (ADR 0022 §1, D2), per Pool value: at most max of the children it started — wrapped or submitted — are doing work at any moment. A child holds a slot while its run works and none while it waits: not while a sync delegation of its own runs a child (the slot is handed back for the wait and re-acquired, through the same queue, before the run's next model call), and not while it is parked at an approval. Admission is FIFO in queueing order. A session's fan-out is already bounded by its agent's Parallelism, so the pool adds no per-session quota; a delegation chain's depth is bounded separately (MaxDepth).

A Pool is safe for concurrent use. The zero value is not usable; New constructs. One Pool should own a storage's delegations at a time: Recover treats an unsettled receipt no child of this pool runs for as left behind by a process that is gone.

func New

func New(max int, opts ...Option) *Pool

New returns a pool admitting at most max child runs at work at once. The bound is the pool's reason to exist, so max below 1 panics — an unbounded pool is no pool (DeerFlow's max_running, the field's evidence for the shape) — as does a MaxDepth below 1.

func (*Pool) Cancel

func (p *Pool) Cancel(ctx context.Context, parent *thread.Session, receiptID string) error

Cancel ends the delegation behind receiptID, a receipt of parent's (D4). What that means depends on where the child is:

  • queued or running: its run is canceled. The child session records the canceled turn, the receipt settles canceled on the run's goroutine — Cancel does not wait for it; Wait does — and a sync delegation's parent model reads SUBAGENT_CANCELED as the call's result.
  • parked at an approval: the child's pending requests are denied in the child session (so its parked calls can never run on a later approval), its run is not resumed, the receipt settles canceled before Cancel returns, the mirrored requests leave the parent's Pending, and a sync delegation's parked call resolves with SUBAGENT_CANCELED.
  • unsettled in the ledger with no child of this pool running for it — a previous process's: it is recovered first (see Recover) and then canceled if recovery left it parked.

A settled receipt fails with a *StateError (ErrNotRunning under errors.Is) carrying its state; an id the ledger does not hold with ErrUnknownReceipt.

Example

Cancel ends one delegation: the running child's run is canceled, its receipt settles canceled, and a second Cancel says the receipt is no longer running — and what it is instead.

package main

import (
	"context"
	"errors"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

// waiting builds a child agent that blocks inside a tool until
// release closes — a child caught mid-run, for the examples that need
// one. started closes when the tool is entered.
func waiting(started, release chan struct{}, then string) *weft.Agent {
	wait := weft.Tool("wait", "", func(ctx context.Context, _ struct{}) (string, error) {
		close(started)
		select {
		case <-release:
			return "released", nil
		case <-ctx.Done():
			return "", ctx.Err()
		}
	})
	return weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "wait", ID: "call-wait"}),
		wefttest.Say(then),
	), wait)
}

func main() {
	ctx := context.Background()
	s, _ := thread.Create(ctx, thread.Memory(), weft.New(wefttest.Script()))
	p := pool.New(1)
	started, release := make(chan struct{}), make(chan struct{})
	r, err := p.Submit(ctx, s, waiting(started, release, "never said"), "work")
	if err != nil {
		panic(err)
	}
	<-started // the child is mid-run
	if err := p.Cancel(ctx, s, r.ID); err != nil {
		panic(err)
	}
	final, _ := p.Wait(ctx, s, r.ID)
	fmt.Println("receipt:", final.State, "settled:", final.Settled())

	err = p.Cancel(ctx, s, r.ID)
	var se *pool.StateError
	fmt.Println("again:", errors.Is(err, pool.ErrNotRunning), errors.As(err, &se), se.State)
	fmt.Println("wrong id:", errors.Is(p.Cancel(ctx, s, "e_unknown"), pool.ErrUnknownReceipt))
}
Output:
receipt: canceled settled: true
again: true true canceled
wrong id: true

func (*Pool) Close

func (p *Pool) Close(ctx context.Context) error

Close stops taking work and drains: every queued or running child is canceled — async ones on their pool-owned contexts, sync ones through their registered cancels, their parents' models reading SUBAGENT_CANCELED — wrapped calls running outside any session are canceled too, and Close waits for all of them to settle their receipts and end, every child's session closed before the process gives up the storage (DeerFlow's gateway drain). A child parked at an approval is not canceled: its session is closed, its receipt stays parked, and a later pool resumes it (Recover, Decide).

ctx bounds the wait; its error returns with the drain still progressing. A second Close returns nil — the drain is already done or underway.

Example

Close drains the pool: the running child and the one still queued behind it are canceled and settled before Close returns, and the pool takes no more work.

package main

import (
	"context"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

// waiting builds a child agent that blocks inside a tool until
// release closes — a child caught mid-run, for the examples that need
// one. started closes when the tool is entered.
func waiting(started, release chan struct{}, then string) *weft.Agent {
	wait := weft.Tool("wait", "", func(ctx context.Context, _ struct{}) (string, error) {
		close(started)
		select {
		case <-release:
			return "released", nil
		case <-ctx.Done():
			return "", ctx.Err()
		}
	})
	return weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "wait", ID: "call-wait"}),
		wefttest.Say(then),
	), wait)
}

func main() {
	ctx := context.Background()
	s, _ := thread.Create(ctx, thread.Memory(), weft.New(wefttest.Script()))
	p := pool.New(1) // one slot: the second child queues
	started, release := make(chan struct{}), make(chan struct{})
	if _, err := p.Submit(ctx, s, waiting(started, release, "never said"), "first"); err != nil {
		panic(err)
	}
	<-started
	if _, err := p.Submit(ctx, s, weft.New(wefttest.Script(wefttest.Say("never ran"))), "second"); err != nil {
		panic(err)
	}
	if err := p.Close(ctx); err != nil {
		panic(err)
	}
	for _, r := range pool.Receipts(s) {
		fmt.Println("receipt:", r.State)
	}
	_, err := p.Submit(ctx, s, weft.New(wefttest.Script()), "third")
	fmt.Println("after Close:", err)
}
Output:
receipt: canceled
receipt: canceled
after Close: thread/pool: pool is closed

func (*Pool) Decide

func (p *Pool) Decide(ctx context.Context, parent *thread.Session, ds ...thread.Decision) error

Decide records decisions over a parent session's nested approvals and arms the children they complete (ADR 0022 §7). It does not wait for any child to run.

The decisions are recorded in the parent first, through its own Decide — validation, quorum, audit and expiry all live there, and a RequireSigned parent refuses them with thread.ErrSignatureRequired (use DecideSigned). Then the pump runs: every child parked on this parent whose currently pending calls all hold an effective verdict in the parent — decided, or lapsed past their expiry and denied by the parent's sweep — is queued to resume. With no decisions Decide is the pump alone.

Decisions the parent session records by a path of its own reach the children the same way, without a call: the pool watches every parent it delegates from (and every one handed to Decide, DecideSigned or Recover), and a decision recorded for a mirrored request — Session.Decide or DecideSigned used directly, the denial an interrupting Send (thread.Interrupt, thread.Rollback) gives a parked boundary, a request lapsing past its expiry when the session is next touched — runs the pump on a pool goroutine. An interrupted child is therefore resumed with the denial ("DENIED: interrupted by a newer message"), runs to its end and settles; its delegating call, denied by the same interrupt, is no longer there to resolve, and the child's answer stays on its receipt. The watch is the live Session's: after a restart nothing is watched until Recover or Decide is called for the parent, which is also what pumps whatever was decided in between.

A resumed child is pool work like its first run: on a pool-owned goroutine and context, behind the queue, holding a slot while it runs, its receipt recording running again. There the parent's decisions for the child's pending calls are replayed into the child session (thread.Session.ReplayDecisions: the decider's identity kept, Via "parent") — only for calls pending in the child now, so the decisions of an earlier park are never replayed — and the child resumes exactly as a top-level session does (ADR 0021 §1). When it ends, the delegation completes: the receipt settles, and a sync delegation's parked call resolves with the child's answer, which resumes the parent. A child that parks again mirrors its new requests and waits for the next Decide.

The delegating call a child is parked under is never decided: it is not in Pending, and a decision naming it fails with thread.ErrDelegated — it completes with the child's answer. The parent session's own decision chain is bound by the same rule: no grant and no live Approver is consulted for a delegating call when it parks, whatever they would match — the call parks, and the child's mirrored requests are what there is to decide. Decisions for the parent's own, ordinary calls may ride in the same batch; a resume they arm is the parked turn's Next, as with Session.Decide.

Decide returns once the decisions are recorded and the ready children queued; ctx bounds that, not the children. Follow a child with Wait, Receipts, or the parent's own continuation (Turn.Next once Wait reports the receipt settled). The error is every failure joined (errors.Join): the parent's refusal of the decisions — nothing is recorded and nothing pumped — or, per child that could not be armed, why: ErrNoAgent for one the pool cannot reopen. The other children are armed regardless.

Example

A nested approval, end to end (ADR 0022 §7): the wrapped child parks at its gated tool, the parent's delegating call parks with it, the child's request surfaces on the parent's Pending with its lineage, and one decision through the pool resumes the child and completes the parent's call with the child's answer — the parent model reads it as the tool result and finishes its turn.

package main

import (
	"context"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

func main() {
	ctx := context.Background()
	gated := weft.Tool("refund", "", func(_ context.Context, in struct{ OrderID string }) (string, error) {
		return "refunded " + in.OrderID, nil
	}, weft.RequireApproval())
	child := weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "refund", Args: `{"order_id":"42"}`}),
		wefttest.Say("refund issued"),
	), gated)
	p := pool.New(2)
	parent := weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "research",
			Args: wefttest.Args(struct{ Prompt string }{"refund order 42"})}),
		wefttest.Say("handled"),
	), p.MustWrap("research", "delegates the refund flow", child))
	s, _ := thread.Create(ctx, thread.Memory(), parent)
	t1, _ := s.Send(ctx, weft.User("refund order 42"))
	if _, err := t1.Wait(); err != nil {
		panic(err)
	}
	pend := s.Pending()
	fmt.Println("pending:", pend[0].Tool)
	if err := p.Decide(ctx, s, thread.Approve(pend[0].CallID)); err != nil {
		panic(err)
	}
	// Decide queued the child's resume and returned; Wait follows the
	// delegation to rest — here its settlement, which also resolved
	// the parent's parked call.
	if _, err := p.Wait(ctx, s, pool.Receipts(s)[0].ID); err != nil {
		panic(err)
	}
	res, err := t1.Next().Wait() // the parent's parked run, resumed with the answer
	if err != nil {
		panic(err)
	}
	fmt.Println("parent:", res.Text())
	for _, r := range pool.Receipts(s) {
		fmt.Println("receipt:", r.State, "-", r.Stop)
	}
}
Output:
pending: refund
parent: handled
receipt: done - refund issued

func (*Pool) DecideSigned

func (p *Pool) DecideSigned(ctx context.Context, parent *thread.Session, sd thread.SignedDecision) error

DecideSigned is Decide for one signed decision: sd is verified and recorded by the parent session (thread.Session.DecideSigned, every check its own), then the pump runs as in Decide.

func (*Pool) Forward

func (p *Pool) Forward(ctx context.Context, parent *thread.Session, receiptID string, msg core.Message) (*thread.Turn, error)

Forward steers a running child with msg, explicitly (ADR 0022 §8, ADR 0019 §7): the message is sent to the child's session under the Steer policy — delivered at the run's next steering drain point, or deferred to a follow-up turn when the run ends on an intended stop or an open approval boundary (the follow-up runs before the delegation's session closes, the receipt keeping the answer of the run it steered); the child's session records the steer's receipt either way. Nothing is ever forwarded implicitly. The returned Turn is the steer's receipt turn.

receiptID is a receipt of parent's. A child that is not running refuses with a *StateError (ErrNotRunning under errors.Is) naming the state it is in — accepted (queued for a slot: a steer now would run as the session's first turn, the task's prompt behind it), parked (its fate is the decision), or settled; an id the ledger does not hold fails with ErrUnknownReceipt.

Example

Forward steers a running child, explicitly: the message joins the child's run at its next drain point. A child that is not running refuses, naming its state.

package main

import (
	"context"
	"errors"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

// waiting builds a child agent that blocks inside a tool until
// release closes — a child caught mid-run, for the examples that need
// one. started closes when the tool is entered.
func waiting(started, release chan struct{}, then string) *weft.Agent {
	wait := weft.Tool("wait", "", func(ctx context.Context, _ struct{}) (string, error) {
		close(started)
		select {
		case <-release:
			return "released", nil
		case <-ctx.Done():
			return "", ctx.Err()
		}
	})
	return weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "wait", ID: "call-wait"}),
		wefttest.Say(then),
	), wait)
}

func main() {
	ctx := context.Background()
	st := thread.Memory()
	s, _ := thread.Create(ctx, st, weft.New(wefttest.Script()))
	p := pool.New(1)
	started, release := make(chan struct{}), make(chan struct{})
	child := waiting(started, release, "converted to euros")
	r, err := p.Submit(ctx, s, child, "convert the totals")
	if err != nil {
		panic(err)
	}
	<-started // mid-run: inside its tool
	if _, err := p.Forward(ctx, s, r.ID, weft.User("use euros, not dollars")); err != nil {
		panic(err)
	}
	close(release)
	final, _ := p.Wait(ctx, s, r.ID)
	fmt.Println("receipt:", final.State, "-", final.Stop)

	// The steer is in the child's transcript, after its prompt.
	childSession, err := thread.Open(ctx, st, r.Child, child)
	if err != nil {
		panic(err)
	}
	for _, m := range childSession.Context() {
		if m.Role == weft.RoleUser {
			fmt.Println("child read:", m.Text())
		}
	}
	_, err = p.Forward(ctx, s, r.ID, weft.User("too late"))
	fmt.Println("after the end:", err != nil, errors.Is(err, pool.ErrNotRunning))
}
Output:
receipt: done - converted to euros
child read: convert the totals
child read: use euros, not dollars
after the end: true true

func (*Pool) MustWrap

func (p *Pool) MustWrap(name, description string, agent *core.Agent, opts ...WrapOption) *core.ToolDef

MustWrap is Wrap for package-level and constructor wiring, where the arguments are the program's own: it panics on the errors Wrap returns — a nil agent, an empty name, a name already wrapping a different agent.

func (*Pool) Recover

func (p *Pool) Recover(ctx context.Context, parent *thread.Session) error

Recover brings parent's delegations back under this pool after a restart (ADR 0022 §7): the ledger is entries, the pool's memory is not, and a process can die between any two of them. For every receipt of parent's that is unsettled and that no child of this pool is working, Recover reads the child session and makes the ledger say what is true:

  • the child is parked at an approval: the bridge is rebuilt (the child reopened under its registered or wrap-named agent), any pending request the parent holds no mirror for is mirrored — the crash window between a child's park and its mirror — and the receipt records parked. The child resumes when its requests are decided, as if nothing had happened.
  • the child finished its run: the receipt settles from the child's own ledger — done with its final text, failed, capped or canceled with its recorded cause — with the child's usage.
  • the child never finished a run — queued, or mid-run when the process died: the receipt settles failed, saying so. The run is not started again: what it had done, it had done.
  • the child session is gone: the receipt settles failed.

A sync delegation's call still parked in the parent is resolved with the settlement, including for receipts that settled but whose resolution a crash lost. Last, the pump runs (Decide with no decisions): children whose requests were already decided resume.

Call it once per parent session a new process takes over, after the agents are wrapped or registered. A child whose agent the pool does not hold is reported with ErrNoAgent and left for a later call; the rest are recovered. The error is every failure joined. Recover on a session with nothing to recover does nothing; it is safe to call on a live session too — receipts whose children this pool is working are left alone.

Example

Recover after a restart, for wrapped delegations: wrapping the same agent under the same name is all the new process needs — the name is in the child's header — and Recover reattaches what was parked. Without the agent, Recover says which child it could not reopen.

package main

import (
	"context"
	"errors"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

func main() {
	ctx := context.Background()
	st := thread.Memory()
	gated := weft.Tool("refund", "", func(_ context.Context, in struct{ OrderID string }) (string, error) {
		return "refunded " + in.OrderID, nil
	}, weft.RequireApproval())

	// The first process parks a sync delegation and stops.
	p := pool.New(1)
	child := weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "refund", ID: "call-refund", Args: `{"order_id":"42"}`}),
	), gated)
	parent := weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "refunds", ID: "call-delegate", Args: `{"prompt":"refund order 42"}`}),
	), p.MustWrap("refunds", "delegates the refund flow", child))
	s, _ := thread.Create(ctx, st, parent)
	t1, _ := s.Send(ctx, weft.User("refund order 42"))
	if _, err := t1.Wait(); err != nil {
		panic(err)
	}
	receipt := pool.Receipts(s)[0]
	if err := p.Close(ctx); err != nil {
		panic(err)
	}
	if err := s.Close(ctx); err != nil {
		panic(err)
	}

	// The second process: the same agents wrapped under the same names.
	p2 := pool.New(1)
	child2 := weft.New(wefttest.Script(wefttest.Say("refund issued")), gated)
	parent2 := weft.New(wefttest.Script(wefttest.Say("handled")),
		p2.MustWrap("refunds", "delegates the refund flow", child2))
	reopened, err := thread.Open(ctx, st, s.ID(), parent2)
	if err != nil {
		panic(err)
	}
	// A pool that never wrapped "refunds" cannot reopen the child, and
	// says so.
	fmt.Println("no agent yet:", errors.Is(pool.New(1).Recover(ctx, reopened), pool.ErrNoAgent))

	if err := p2.Recover(ctx, reopened); err != nil {
		panic(err)
	}
	pend := reopened.Pending()
	if err := p2.Decide(ctx, reopened, thread.Approve(pend[0].CallID)); err != nil {
		panic(err)
	}
	final, _ := p2.Wait(ctx, reopened, receipt.ID)
	fmt.Println("receipt:", final.State, "-", final.Stop)
	// The delegating call resolved with the child's answer, which
	// armed the parent's own resume; Close waits for it.
	if err := reopened.Close(ctx); err != nil {
		panic(err)
	}
	last := reopened.Context()
	fmt.Println("parent:", last[len(last)-1].Text())
}
Output:
no agent yet: true
receipt: done - refund issued
parent: handled

func (*Pool) Register

func (p *Pool) Register(sessionID string, agent *core.Agent) error

Register attaches an agent to a child session the pool holds none for — the restart hook (ADR 0022 §7): a parked child's boundary outlives the process, and resuming it needs the agent that runs it, which no file carries. Children a Wrap made need no call: the wrap's name, recorded in the child's header metadata, is looked up among the pool's wrapped agents. Register is for the rest — Submit children — and for wraps renamed between runs; it takes precedence over the wrap name. The registration is dropped when the session's delegation settles. An empty session id or a nil agent is an error. Register alone recovers nothing: Recover does, and Decide resumes.

Example

After a restart: a Submit child parked at an approval when the process stopped. The new process registers the child's agent — nothing on disk carries it — recovers the parent's ledger, and the decision resumes the child where it parked.

package main

import (
	"context"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

func main() {
	ctx := context.Background()
	st := thread.Memory()
	gated := weft.Tool("refund", "", func(_ context.Context, in struct{ OrderID string }) (string, error) {
		return "refunded " + in.OrderID, nil
	}, weft.RequireApproval())

	// The first process: the child parks, and the process stops.
	s, _ := thread.Create(ctx, st, weft.New(wefttest.Script()))
	p := pool.New(1)
	r, err := p.Submit(ctx, s, weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "refund", ID: "call-refund", Args: `{"order_id":"42"}`}),
	), gated), "refund order 42")
	if err != nil {
		panic(err)
	}
	parked, _ := p.Wait(ctx, s, r.ID)
	fmt.Println("before the restart:", parked.State)
	if err := p.Close(ctx); err != nil { // closes the parked child's session
		panic(err)
	}
	if err := s.Close(ctx); err != nil {
		panic(err)
	}

	// The second process: same storage, new pool, new session values.
	reopened, err := thread.Open(ctx, st, s.ID(), weft.New(wefttest.Script()))
	if err != nil {
		panic(err)
	}
	p2 := pool.New(1)
	resumed := weft.New(wefttest.Script(wefttest.Say("refund issued")), gated)
	if err := p2.Register(r.Child, resumed); err != nil {
		panic(err)
	}
	if err := p2.Recover(ctx, reopened); err != nil {
		panic(err)
	}
	pend := reopened.Pending()
	fmt.Println("pending after the restart:", pend[0].Tool)
	if err := p2.Decide(ctx, reopened, thread.Approve(pend[0].CallID)); err != nil {
		panic(err)
	}
	final, _ := p2.Wait(ctx, reopened, r.ID)
	fmt.Println("after the decision:", final.State, "-", final.Stop)
}
Output:
before the restart: parked
pending after the restart: refund
after the decision: done - refund issued

func (*Pool) Submit

func (p *Pool) Submit(ctx context.Context, parent *thread.Session, agent *core.Agent, prompt string) (*Receipt, error)

Submit hands a task to a child session of parent and returns at once (ADR 0022 §2): the child session is created in the parent's storage with a lineage naming the parent and the parent's approval policy (thread.InheritApprovals), the acceptance receipt is durable in the parent session and the delegation has its place in the pool's queue before Submit returns, and the child runs on a pool-owned context — independent of the submitting caller's cancellation by design (D4), canceled explicitly by Cancel or Close.

The receipt's Child names the session to read, reopen or tail. A live tail needs a storage with the thread.Watcher capability — thread/jsonl and thread/sqlite have it, thread.Memory does not; on any storage Wait blocks until the delegation settles or parks, and Receipts reads the ledger.

parent and agent must not be nil. Called from inside a pool child's run — a tool handler passing its context — Submit counts the delegation's depth and ancestry from there, and fails with ErrDepth or ErrCycle as a wrapped tool would refuse. A closed pool fails with ErrClosed.

Example

An async delegation: Submit returns at once with the acceptance, the child runs on the pool's context, and Wait follows it to its settlement.

package main

import (
	"context"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

func main() {
	ctx := context.Background()
	s, _ := thread.Create(ctx, thread.Memory(), weft.New(wefttest.Script()))
	p := pool.New(1)
	child := weft.New(wefttest.Script(wefttest.Say("background answer")))
	r, err := p.Submit(ctx, s, child, "work in the background")
	if err != nil {
		panic(err)
	}
	fmt.Println("accepted:", r.State)
	final, err := p.Wait(ctx, s, r.ID)
	if err != nil {
		panic(err)
	}
	fmt.Println("settled:", final.State, "-", final.Stop)
	if err := p.Close(ctx); err != nil {
		panic(err)
	}
}
Output:
accepted: accepted
settled: done - background answer
Example (Watch)

Watching an async child live: the receipt names the child session, and a storage with the Watcher capability — jsonl here; Memory has none — tails its entries as the child writes them.

package main

import (
	"context"
	"fmt"
	"os"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/jsonl"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

func main() {
	ctx := context.Background()
	dir, err := os.MkdirTemp("", "weft-pool-example")
	if err != nil {
		panic(err)
	}
	defer func() { _ = os.RemoveAll(dir) }()
	st, err := jsonl.Open(dir)
	if err != nil {
		panic(err)
	}
	s, _ := thread.Create(ctx, st, weft.New(wefttest.Script()))
	p := pool.New(1)
	r, err := p.Submit(ctx, s, weft.New(wefttest.Script(wefttest.Say("background answer"))), "work in the background")
	if err != nil {
		panic(err)
	}
	watcher, ok := thread.Storage(st).(thread.Watcher)
	if !ok {
		panic("this storage cannot tail a session")
	}
	entries, err := watcher.Watch(ctx, r.Child, "")
	if err != nil {
		panic(err)
	}
	for e, err := range entries {
		if err != nil {
			panic(err)
		}
		switch e := e.(type) {
		case thread.MessageEntry:
			fmt.Printf("%s: %s\n", e.Message.Role, e.Message.Text())
		case thread.TurnEntry:
			fmt.Println("turn ended:", e.StopReason)
		}
		if _, ended := e.(thread.TurnEntry); ended {
			break // the child's one turn is over; a tail would wait for more
		}
	}
	if err := p.Close(ctx); err != nil {
		panic(err)
	}
}
Output:
user: work in the background
assistant: background answer
turn ended: stop

func (*Pool) Wait

func (p *Pool) Wait(ctx context.Context, parent *thread.Session, receiptID string) (Receipt, error)

Wait blocks until the delegation behind receiptID is at rest in this pool — settled, or parked at an approval — and returns its receipt as the ledger then reads. When it returns for a sync delegation that settled after a park, the parent's parked call is already resolved and its resume armed. A receipt no child of this pool is working returns at once, as it stands. ctx bounds the wait; an id the ledger does not hold fails with ErrUnknownReceipt.

func (*Pool) Wrap

func (p *Pool) Wrap(name, description string, agent *core.Agent, opts ...WrapOption) (*core.ToolDef, error)

Wrap returns a subagent tool whose children the pool runs (ADR 0022 §1–§2): a delegation tool over agent, built on the core's Subagent, run as a child session of the session whose run is calling it — linked by lineage, receipted in the parent, bounded by the pool.

Inside a session's run (the parent found through the run context) the child is a session of its own in the parent's storage: sync by default — the call waits and the result is the child's answer; with Async the result is the receipt and the child runs on. While a sync call waits, the run that made it holds no slot (see Pool). A wrapped tool called outside any session — a bare Generate — has no parent to receipt into and falls back to the ordinary subagent path under a slot, ADR 0014's semantics unchanged, with the same hand-off when such calls nest.

What the calling model reads, besides the child's answer:

  • SUBAGENT_CYCLE when agent is already running above the call — the delegation chain's real ancestry, checked before any child is created;
  • SUBAGENT_DEPTH when the delegation would exceed MaxDepth;
  • SUBAGENT_FAILED when the child's run fails, the core's text;
  • SUBAGENT_CANCELED when the pool canceled the child (Cancel, Close). The delegating call's own cancellation or timeout is not that: the ordinary machinery reports it.

name is also the resume key: children the wrap makes record it, and a restarted process that wraps the same agents under the same names can resume them. One name therefore names one agent — wrapping a different agent under a name the pool already holds fails with ErrDuplicateWrap; wrapping the same agent again returns another tool for it. A nil agent or an empty name is an error. MustWrap is Wrap for wiring that cannot fail.

Example

A sync delegation through a session: the wrapped tool's child runs as a child session under a pool slot, the call waits, the result is the child's answer, and the parent's ledger carries the child's cost in its Delegated bucket.

package main

import (
	"context"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

func main() {
	ctx := context.Background()
	st := thread.Memory()
	child := weft.New(wefttest.Script(wefttest.Say("the answer is 42")))
	p := pool.New(2)
	parent := weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "research",
			Args: wefttest.Args(struct{ Prompt string }{"find the answer"})}),
		wefttest.Say("done"),
	), p.MustWrap("research", "delegates research to a child session", child))
	s, _ := thread.Create(ctx, st, parent)
	turn, err := s.Send(ctx, weft.User("what is the answer?"))
	if err != nil {
		panic(err)
	}
	res, err := turn.Wait()
	if err != nil {
		panic(err)
	}
	fmt.Println("reply:", res.Text())
	for _, r := range pool.Receipts(s) {
		fmt.Println("receipt:", r.State, "-", r.Stop)
	}
	fmt.Println("delegated output tokens:", s.Usage().Delegated.OutputTokens)
}
Output:
reply: done
receipt: done - the answer is 42
delegated output tokens: 5

type Receipt

type Receipt struct {
	ID    string
	State State
	Child string
	// Call is the parent-side tool call that delegated; empty for a
	// Submit.
	Call string
	Stop string
}

Receipt is a delegation's handle and what is known of it: the acceptance entry's id, the state the ledger records last, the child session, the delegating call when a wrapped tool made it, and on settlement the child's stop text — its answer, when the state is Done; the cause otherwise.

func Receipts

func Receipts(parent *thread.Session) []Receipt

Receipts returns parent's pool receipts, in acceptance order, each at the state its entries record last — the ledger read back after a restart the same as live (entries are the truth; the pool holds no receipt state of its own). The first settlement a receipt records is its state for good. A session that never delegated returns nil.

Example

Receipts reads the ledger: every delegation of a session, in acceptance order, at the state its entries record last. Settled tells the finished from the ones still in flight or parked.

package main

import (
	"context"
	"fmt"

	"github.com/weftgo/weft"
	"github.com/weftgo/weft/thread"
	"github.com/weftgo/weft/thread/pool"
	"github.com/weftgo/weft/wefttest"
)

func main() {
	ctx := context.Background()
	s, _ := thread.Create(ctx, thread.Memory(), weft.New(wefttest.Script()))
	p := pool.New(2)
	gated := weft.Tool("refund", "", func(context.Context, struct{}) (string, error) {
		return "refunded", nil
	}, weft.RequireApproval())
	quick, _ := p.Submit(ctx, s, weft.New(wefttest.Script(wefttest.Say("42"))), "compute")
	asks, _ := p.Submit(ctx, s, weft.New(wefttest.Script(
		wefttest.ToolCalls(wefttest.Call{Name: "refund", ID: "call-refund"}),
	), gated), "refund the order")
	for _, r := range []*pool.Receipt{quick, asks} {
		if _, err := p.Wait(ctx, s, r.ID); err != nil { // at rest: settled, or parked
			panic(err)
		}
	}
	for _, r := range pool.Receipts(s) {
		fmt.Printf("%s settled=%v stop=%q\n", r.State, r.Settled(), r.Stop)
	}
	fmt.Println("delegated output tokens:", s.Usage().Delegated.OutputTokens)
}
Output:
done settled=true stop="42"
parked settled=false stop=""
delegated output tokens: 5

func (Receipt) Settled

func (r Receipt) Settled() bool

Settled reports whether the receipt reached a final state.

type State

type State string

State is where a delegation is (ADR 0022 §4): accepted → running ⇄ parked → exactly one of done, failed, canceled, capped. Its values are the wire strings of thread.PoolReceiptEntry.Status.

const (
	// Accepted: recorded, and queued for a slot.
	Accepted State = thread.PoolAccepted
	// Running: the child holds a slot and its run is in flight — or
	// is waiting, slotless, on a sync delegation of its own.
	Running State = thread.PoolRunning
	// Parked: the child's run ended at an approval boundary; it holds
	// no slot and resumes when its mirrored requests are decided.
	Parked State = thread.PoolParked
	// Done: the child ran to its end; Stop is its answer.
	Done State = thread.PoolDone
	// Failed: the child's run failed, or the delegation could not
	// proceed; Stop is the cause.
	Failed State = thread.PoolFailed
	// Canceled: Cancel, the pool's Close, or — for a sync child — the
	// delegating call's own cancellation or timeout.
	Canceled State = thread.PoolCanceled
	// Capped: the child died on a budget (MaxSteps, a usage limit).
	Capped State = thread.PoolCapped
)

The receipt states.

func (State) Settled

func (s State) Settled() bool

Settled reports whether the state is final: done, failed, canceled or capped. An unknown state — a newer writer's — is not settled.

func (State) String

func (s State) String() string

String returns the state's wire string.

type StateError

type StateError struct {
	// Receipt is the receipt id the call named.
	Receipt string
	// State is the state the receipt is in.
	State State
	// Orphan reports an unsettled receipt no child of this pool runs
	// for: the ledger of an earlier process. Recover settles it.
	Orphan bool
}

A StateError reports a receipt whose child is not in the state the call serves — Forward to a child that is queued, parked or settled; Cancel of a settled one. It matches ErrNotRunning under errors.Is; State tells an operator "already done" from "parked".

func (*StateError) Error

func (e *StateError) Error() string

Error names the receipt and its state.

func (*StateError) Is

func (e *StateError) Is(target error) bool

Is reports whether target is ErrNotRunning.

type WrapOption

type WrapOption interface {
	// contains filtered or unexported methods
}

A WrapOption configures one wrapped delegation tool.

func Async

func Async() WrapOption

Async makes the wrapped tool's result an acceptance receipt instead of the child's answer (ADR 0022 D1): the middleware submits the child and returns at once — the receipt line is the tool result, model-visible bytes pinned by a golden — and the answer is delivered by the application in a later turn. Without it the delegation is sync: the call waits for the child session's turn and the result is the child's answer, as an ordinary subagent's is.

func ToolOptions

func ToolOptions(opts ...core.ToolOption) WrapOption

ToolOptions forwards weft tool options to the underlying subagent tool — Timeout, MaxResultBytes, a snippet — which otherwise the wrap would hide. RequireApproval composes as on any tool: it gates the act of delegating (ADR 0014).

Jump to

Keyboard shortcuts

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