Agents
Every agent setting, delivery rule, routing policy, and workflow option
This page is the full reference for agents, contracts, discovery, and workflows. For a short introduction, read Agents.
Agent settings
Rust and TypeScript build an agent with Agent::builder() and Agent.builder(), then call spawn(laser). Python calls laser.spawn_agent(..). Spawning returns a handle, and ready() on the handle waits until the agent reads its topic.
The handler is a Rust type that implements AgentHandler, a TypeScript object with a handle(message, ctx) method, or a Python callable or object with handle(ctx, message). Note the argument order in Python.
| Setting | Rust Agent::builder() | TypeScript Agent.builder() | Python spawn_agent(..) |
|---|---|---|---|
| Identity | id(..) | id(..) | First argument, a name or None |
| Input topic | listen_on(..) | listenOn(..) | Second argument |
| Handler | handler(..) | handler(..) | Third argument |
Reply topic for respond | respond_on(..) | respondOn(..) | respond_on= |
| Consumer group | consumer_group(..) | consumerGroup(..) | consumer_group= |
| Capabilities | capabilities(vec![..]) | capabilities([..]) | capabilities=[..] |
| Pickup status | ack_on_pickup(true) | ackOnPickup() | ack_on_pickup=True |
| Served operations | operations(vec![..]) | operations([..]) | operations=[..] |
| Session configuration | sessions(SessionConfig) | sessions(config) | sessions=laser.sessions(..) |
| Default route for directed sends | inbox_route(..) | inboxRoute(..) | fixed_inbox= |
| Poll interval | poll_interval(Duration) | pollInterval(ms) | poll_interval_ms= |
| Shutdown grace | shutdown_grace(Duration) | shutdownGrace(ms) | shutdown_grace_ms= |
| Retry | retry(RetryPolicy) | retry({ maxAttempts, baseDelayMs }) | retry_max_attempts=, retry_base_delay_ms= |
| Dedup window | dedup_window(n) | dedupWindow(n) | dedup_window= |
| Custom deduplicator | deduplicator(..) | deduplicator(..) | dedup= |
| Load dedup keys on start | warm_dedup(true) | warmDedup() | warm_dedup=True |
| Partition lanes | concurrency(ConcurrencyPolicy::SerialPerPartition { max_partitions }) | concurrency({ kind: "serial-per-partition", maxPartitions }) | max_partitions= |
| Queue bounds | max_queued_records(n), max_queued_bytes(n) | maxQueuedRecords(n), maxQueuedBytes(n) | max_queued_records=, max_queued_bytes= |
| Understood feature bits | understood_features(bits) | understoodFeatures(bits) | understood_features= |
| Middleware | middleware(vec![..]) | middleware(..), once per hook object | middleware=[..] |
| Dead-letter callback | on_dead_letter(..) | onDeadLetter(..) | dead_letter= |
| Governor | governor((governor, mode)) | governor([governor, mode]) | governor=, governor_mode= |
| Governor retention | governor_retention(GovernorRetention { capacity, idle_ttl }) | governorRetention({ capacity, idleTtlMs }) | governor_retention=(capacity, idle_ttl_ms) |
| Periodic consolidation | consolidate_every(..), consolidator(..) | consolidateEvery(ms), consolidator(..) | consolidate_every_ms=, consolidator= |
| Signing and verification | signing_key(..), verifier(..) | signingKey(..), verifier(..) | signing_key=, verifier= |
Defaults in all three SDKs:
| Setting | Default |
|---|---|
| Consumer group | The agent ID |
| Poll interval | 10 ms |
| Shutdown grace | 30 seconds |
| Retry | 5 attempts in total, 200 ms before the first retry, doubling each time |
| Dedup window | 10,000 keys in memory |
| Concurrency | Serial, one message at a time across all partitions |
| Queue bounds | 4,096 records and 64 MiB across partition lanes |
| Pickup status | Off |
| Served operations | All |
More on individual settings:
respond_onis required forrespondandfan_out. Without it,respondfails with a no-respond-topic error.- An agent with capabilities publishes a capability card to the registry on spawn and advertises its inbox (its
listen_ontopic) as live presence. Cards published at spawn never expire. Presence belongs to a connection, so on a server that serves presence give each advertising agent its own connection. A second one on the same connection fails to start with a presence conflict error. ack_on_pickupmakes the agent report aWorkingstatus when it picks up a command, before the handler runs.- Python accepts capability names or descriptor dictionaries. A dictionary can carry
skill_id,input,output,cost_class,latency_class,max_concurrency,health, andload. Thehealth=keyword sets one value (healthy,degraded, orunavailable) on every advertised skill. - Pass
Noneas the identity and an explicitconsumer_groupto get an unscoped reliable consumer in Python. It has no agent identity, cannot advertise capabilities, and skips the session budget check. - Run more copies of the same agent to share its load. Never put two different agent IDs in one consumer group, because a record could land on the wrong member and be skipped there.
- With
SerialPerPartition, each partition gets its own ordered lane, so a slow or retrying message on one partition does not stall the others. Order within a partition stays strict. - A message that requires an AGDX feature bit the agent does not declare in
understood_featuresis dead-lettered before the handler runs. - Periodic consolidation needs both an interval and a consolidator. A zero interval is a configuration error when the agent starts. Each pass is scoped to the agent ID. Python accepts a callable or an object with
consolidate(scope). Memory covers consolidation. - An agent governor replaces the connection's governor for that agent. Governance covers governors.
Run bootstrap(partitions, retention) once per stream before agents join. It creates agent.sessions with your retention, agent.heartbeats with a one-hour expiry, and agent.streams, agent.memory, agent.dlq, agent.audit, and agent.workflow_journal. The registry topic appears with the first card. agent.control is not created, because only operators may send to it. bootstrap needs a default stream and is safe to run again. The seven topics hold seven times partitions partitions, so keep the count small on small tiers.
Inside a handler
The context passed to a handler acts as the agent:
| Call | What it does |
|---|---|
respond(payload) | Reply on respond_on. An AGDX command gets a typed response with its correlation, addressed to the requester. A plain request gets a plain reply matched by correlation |
reply_on(topic, payload) | Reply on another topic, chained to the handled message |
send(topic, payload, provenance) | Send with an explicit provenance |
request(request_topic, reply_topic, payload, provenance, timeout) | Send a request and wait for the correlated reply |
fan_out(selector, payload, policy, deadline) | Send one task to every agent with a capability and gather the replies. See Fan out from a handler |
approval_gate(reply_topic, prompt, timeout) | Wait for a human decision. See Interop |
respond_input(reply_topic, response) | Answer a human-input request |
session() | The handled record's session, writing as this agent |
spawn_subconversation() | A child conversation of the handled one |
TypeScript uses camelCase. The request and approval timeouts are required in every SDK and have no default. Python passes them as timeout_ms= in milliseconds.
ctx.session() is a lens on the handled record's session. It holds no lease. When a command carries the submitted marker, as submit and child contracts do, the runtime marks the session working before the handler runs and holds a lease until it returns. If that pickup write fails, the runtime logs a warning and still runs the handler. Advanced sessions covers submitted sessions.
To unit test a handler without a server, build a message with agent_message(payload, provenance) and a context with agent_ctx(laser, message, ..). Rust has them in laser_sdk::testing, TypeScript exports agentMessage and agentCtx (also from @laserdata/laser-sdk/testing), and Python has ls.agent_message and ls.agent_ctx. The laser only needs a live server for the context calls the handler makes.
Delivery and retries
- Delivery is at least once. The runtime commits a message only after the handler succeeds. If the handler crashes first, the agent receives the message again.
- Deduplication uses each message's idempotency key, scoped to the agent that produced it. Messages without a key are not deduplicated. The default backend is an in-memory window. Supply your own deduplicator for durable deduplication, and turn on
warm_dedupto refill the window from the partition tail on start. - External effects still need an idempotency key or a fenced write before commit.
- A retryable handler error is retried until the attempts run out. A non-retryable error, such as an invalid error you raise, dead-letters the message at once.
- A dead-lettered message goes to
agent.dlqand the consumer commits past it, so one bad message does not block the partition. If publishing the dead-letter record fails, the worker stops and leaves the offset uncommitted. - A command picked up after its deadline is dead-lettered with reason
DeadlineExceededbefore the handler runs. - In Rust, a panic in the handler dead-letters the record as rejected, and the agent keeps running.
- In Python, an SDK error keeps its class and retry classification. An exception you raise is classified by type:
InvalidErrorandTypeErrordead-letter at once,PermissionErrormaps to unauthorized, the built-inTimeoutErroris retryable, and a plainValueErroris a retryable handler error. - Python fills an unset retry setting with its default when you set only the other one.
- Publishes from agents use the connection's publish timeout and retries. See Connect.
Which records reach the handler
Agents share agent.sessions by default, so the runtime classifies every record before the handler sees it. Only a command addressed to this agent or to every agent, for an operation the handler serves, reaches the handler. Replies, status records, events, and records addressed to another agent are skipped and committed, and consumed reports them as skipped. The author never decides, so an agent can send work to itself.
On a server that resolves group policies and serves filtered reads, a client built from a connection string binds the agent's group to the addressee filter agdx.to In [<agent id>, "*"], so the server delivers only the agent's own and broadcast records. Otherwise the agent classifies every record itself, with the same result. Advanced sessions has the details.
Session rules the runtime applies:
- Every agent follows
agent.controlacross all its partitions, outside its consumer group, so every instance sees each pause and cancel request. A handler reads them throughctx.session().pending_control()and decides when to stop. While a session is paused, the runtime holds its new work and handles it after the resume. - On a deployment that indexes sessions, work for a session over its budget fails that session once with reason
budgetand never reaches the handler. - Dead letters keep the record's conversation. An agent whose session configuration turns on
fail_on_dead_letterfails the session of a dead-lettered record.
Middleware and dead letters
Middleware wraps every handler dispatch. before_handle runs once before the retry loop and can reject the message, which dead-letters it as rejected without running the handler. after_handle runs after every attempt with its result and a one-based attempt number. A dead-letter callback runs for every dead letter and reports whether the dead-letter record itself was published.
import { Agent, AgentId, AgentTopic } from "@laserdata/laser-sdk"
await using triage = Agent.builder()
.id(AgentId.new("triage"))
.listenOn(AgentTopic.Sessions)
.respondOn(AgentTopic.Sessions)
.retry({ maxAttempts: 3, baseDelayMs: 100 })
.middleware({
afterHandle: async (_message, result, attempt) => {
console.log(`attempt ${attempt} succeeded: ${result.kind === "ok"}`)
}
})
.onDeadLetter({
onDeadLetter: async (_message, capsule, publishError) => {
console.log(
`dead letter after ${capsule.attempts} attempts, published: ${publishError === undefined}`
)
}
})
.handler({ handle: (_message, ctx) => ctx.respond(new TextEncoder().encode("triaged")) })
.build()
.spawn(laser)use laser_sdk::prelude::full::*;
use laser_sdk::wire::agent::AgentDeadLetter;
use std::sync::Arc;
use std::time::Duration;
struct Triage;
impl AgentHandler for Triage {
async fn handle(&self, _message: &AgentMessage, ctx: &AgentCtx<'_>) -> Result<(), LaserError> {
ctx.respond("triaged").await
}
}
struct Metrics;
#[async_trait::async_trait]
impl AgentMiddleware for Metrics {
async fn after_handle(
&self,
_message: &AgentMessage,
result: &Result<(), LaserError>,
attempt: u32,
) {
println!("attempt {attempt} succeeded: {}", result.is_ok());
}
}
struct Alert;
#[async_trait::async_trait]
impl DeadLetterSink for Alert {
async fn on_dead_letter(
&self,
_message: Option<&AgentMessage>,
capsule: &AgentDeadLetter,
publish_result: &Result<(), LaserError>,
) {
println!(
"dead letter after {} attempts, published: {}",
capsule.attempts,
publish_result.is_ok()
);
}
}
let triage = Agent::builder()
.id("triage".parse::<AgentId>()?)
.listen_on(AgentTopic::Sessions)
.respond_on(AgentTopic::Sessions)
.retry(RetryPolicy::backoff(3, Duration::from_millis(100)))
.middleware(vec![Arc::new(Metrics)])
.on_dead_letter(Arc::new(Alert))
.handler(Triage)
.build()
.spawn(laser.clone());import laser_sdk as ls
async def handle(ctx, message):
await ctx.respond(b"triaged")
class Metrics:
async def after_handle(self, message, result, attempt):
print(f"attempt {attempt} succeeded: {result['ok']}")
async def alert(message, capsule, publish_error):
print(f"dead letter after {capsule['attempts']} attempts, published: {publish_error is None}")
triage = laser.spawn_agent(
"triage",
ls.AgentTopic.Sessions,
handle,
respond_on=ls.AgentTopic.Sessions,
retry_max_attempts=3,
retry_base_delay_ms=100,
middleware=[Metrics()],
dead_letter=alert,
)The Rust AgentMiddleware and DeadLetterSink traits use #[async_trait], so add the async-trait crate. The Python middleware result is a dictionary with ok and error, where error is the exception or None. The Python dead-letter callback receives the message (None when the record does not decode), the full capsule as a dictionary, and a typed SDK exception when publishing the dead letter failed. Python hooks can return directly or through an awaitable. An asynchronous hook runs on the event loop that spawned the agent, and a synchronous one runs on an SDK worker thread, so it must not block.
To reprocess a dead letter after a fix, redrive_dead_letter(capsule) (TypeScript: redriveDeadLetter) reads the original record and publishes it again unchanged.
Check consumption
consumed(target, at) reports whether a consumer has committed past a log position, such as a dead-letter capsule's source or an envelope's cause_at. The target is a consumer group or a named consumer.
| Result | Meaning |
|---|---|
Consumed | Committed past the position. Carries the committed offset and the partition head |
NotYetConsumed | Still behind, by behind_by records |
Skipped | The agent's group committed past the record without handling it. Carries the offsets and dispatch, the reason, such as foreign or reply |
A consumer with no stored offset reads as not consumed. The server gives the same answer for a read it cannot authorize, so hold topic read access before you act on the result.
const status = await laser.consumed({ kind: "group", name: "triage" }, capsule.source)
if (status.kind === "notYetConsumed") {
console.log(`behind by ${status.behindBy}`)
}use laser_sdk::agent::{ConsumerRef, ConsumptionStatus};
let status = laser
.consumed(ConsumerRef::Group("triage".parse()?), capsule.source)
.await?;
if let ConsumptionStatus::NotYetConsumed { behind_by } = status {
println!("behind by {behind_by}");
}status = await laser.consumed(
ls.ConsumerRef.Group("triage"),
ls.LogPosition.from_bytes(capsule["source"]),
)
if isinstance(status, ls.ConsumptionStatus.NotYetConsumed):
print(f"behind by {status.behind_by}")A position is a LogPosition of stream ID, topic ID, partition ID, and offset. Rust builds one with LogPosition::new and Python with ls.LogPosition(stream_id, topic_id, partition_id, offset). Both read the packed 20-byte form of a dead-letter source with from_bytes. TypeScript has newLogPosition(streamId, topicId, partitionId, offset) and logPositionFromBytes.
Contracts
laser.contract(router) sends one task to one agent and waits for its terminal reply. In Rust and TypeScript it returns a builder that ends with send(). Python calls laser.contract(skill, payload, source=.., ..) directly, with skill=None and agent= to name an agent. laser.agent(id).contract(..) sends as that agent.
| Option | Rust and TypeScript | Python | Default |
|---|---|---|---|
| Sender | from(agent) | source= | Required |
| Task body | payload(..) | Second argument | Empty |
| Inbox route | inbox_route(..) | fixed_inbox= | The advertised inbox |
| Reply topic | reply_on(topic) | reply_on= | agent.sessions |
| Consumption expiry | expire_if_not_consumed(d) | expire_if_not_consumed_ms= | None |
| Reply deadline | deadline(d) | deadline_ms= | 30 seconds |
| Conversation | conversation(id) | conversation= | A new one per send |
| Fence token | fence(token) | fence= | None |
| Child session | parent(parent, root) | parent=, root= | None |
| Route policy | In the router | policy= | Any |
| Required principal | Router::to_principal, routeToPrincipal | principal= | None |
A contract ends in one of four outcomes:
Completed: the agent replied within the deadline.Failed: the agent replied with a terminal error.NotConsumed: no pickup report and no reply arrived before the consumption expiry. This is reliable only when the agent runs withack_on_pickup. Without pickup reports, a slow handler also reads as not consumed.TimedOut: no terminal reply arrived within the deadline.
Python returns a Contract: Contract.Completed(reply) and Contract.Failed(reply) carry the reply message at index 0, and Contract.NotConsumed() and Contract.TimedOut() carry nothing. TypeScript returns { kind: "completed" | "failed" | "notConsumed" | "timedOut" }, with reply on the first two.
More contract rules:
- The consumption expiry rides the command as its deadline, so an agent that picks it up later dead-letters it instead of handling it.
- A fence token buys two guarantees. A consumer drops a command whose fence is below the highest it has seen for the conversation, which only works when every attempt uses the same pinned
conversation. For an at-most-once external effect across agents or replicas, the handler must commit through a fenced compare-and-swap with the token. - A child contract writes the child's submitted start on
agent.sessionsbefore the command and ends the child by the outcome: completed on a reply, failed otherwise. Advanced sessions covers child sessions. - A reply completes a contract only when it carries the request's correlation, belongs to the request's session, and is addressed to the requester when the request named one.
respondaddresses the requester for you. Under a verifier, the reply must be signed by the routed agent, or by the required principal. - A broadcast or all-capable route is refused, because a contract goes to exactly one agent.
Inbox routes and presence
By default a contract sends work to the inbox the agent advertises in its live presence. Presence needs a server that serves it and, in Rust, the query feature. Against a server without presence, or for an agent with no capabilities, use a fixed inbox: InboxRoute::Fixed(topic) in Rust, { kind: "fixed", topic } in TypeScript, or fixed_inbox= in Python. A route that resolves no inbox fails instead of falling back to a shared topic. A handler's directed sends, fan-out included, use the agent's inbox_route.
A route can also require the agent's live connection to authenticate as a given principal. Use Router::to_principal(agent, principal) or CapabilitySelector::new(skill, policy).principal(p) in Rust, routeToPrincipal(agent, principal) or capabilitySelector(skill, policy, principal) in TypeScript, and principal= in Python on contracts, scatter, fan-out, and workflow steps.
Route policies
A capability route picks one agent among those that advertise the skill. The policy decides which one.
| Policy | Rust | TypeScript | Python |
|---|---|---|---|
| Any candidate | RoutePolicy::Any | { kind: "any" } | "any" or None |
| Lowest cost class | RoutePolicy::Cheapest | { kind: "cheapest" } | "cheapest" |
| Lowest latency class | RoutePolicy::Fastest | { kind: "fastest" } | "fastest" |
| Lowest load | RoutePolicy::LeastLoaded | { kind: "leastLoaded" } | "least_loaded" |
| Prefer one agent | RoutePolicy::Sticky(agent) | { kind: "sticky", agent } | "sticky:<agent>" |
| Your own ranking | RoutePolicy::Custom(scorer) | { kind: "custom", scorer } | A callable or an object with select |
Python passes the policy as policy= on contracts, scatter, and workflow steps, and as route_policy= on handler fan-out. The built-in policies compare the advisory classes on the capability descriptor, where lower is better. An agent that does not advertise the ranked field sorts last. Sticky falls back to any candidate when the preferred agent is absent. Candidates are ordered by agent ID in byte order, so Any and ranking ties pick the same agent in every SDK.
A route that cannot address an agent fails with one of three errors. NoCapableAgent means no live agent advertises the skill. NoInbox means the chosen agent advertises no inbox. RoutePrincipalMismatch means the agent is not bound to the required principal. TypeScript and Python name them NoCapableAgentError, NoInboxError, and RoutePrincipalMismatchError. Python groups them under RoutingError. TypeScript has no shared base class, so catch the three classes or test with isNoCapableAgent(error).
A custom scorer receives the skill ID and the same candidate view as the built-in policies. Each candidate carries the agent, its registered card, and the descriptor it advertises for the skill, which can be missing. Return the index of the chosen candidate. Return nothing, or an index out of range, to refuse every candidate, which fails the route with no capable agent. Python accepts a synchronous callable or an object with a select(skill_id, candidates) method, and rejects an asynchronous scorer. An exception from a Python scorer is raised in place of the routing error.
import { type RouteScorer, routeToCapable } from "@laserdata/laser-sdk"
const preferPrimary: RouteScorer = {
select: (_skillId, candidates) => {
const index = candidates.findIndex((candidate) =>
candidate.agent.asStr().endsWith("-primary")
)
return index === -1 ? undefined : index
}
}
const router = routeToCapable("triage-ticket", { kind: "custom", scorer: preferPrimary })use std::sync::Arc;
struct PreferPrimary;
impl RouteScorer for PreferPrimary {
fn select(&self, _skill_id: &str, candidates: &[RouteCandidate<'_>]) -> Option<usize> {
candidates
.iter()
.position(|candidate| candidate.agent.as_str().ends_with("-primary"))
}
}
let router = Router::to_capable("triage-ticket", RoutePolicy::Custom(Arc::new(PreferPrimary)));def prefer_primary(skill_id, candidates):
for index, candidate in enumerate(candidates):
if candidate.agent.endswith("-primary"):
return index
return None
outcome = await laser.contract(
"triage-ticket",
b"disk full on node-7",
source="intake",
policy=prefer_primary,
fixed_inbox=ls.AgentTopic.Sessions,
)Discover agents
Agent objects come from the connection: laser.agent(id) returns an AgentScope, agent_registry() an AgentRegistry, and workflow(name) a Workflow. In Rust and TypeScript, contract(router) returns a ContractBuilder and a workflow step returns a StepHandle. Python returns results directly and changes a workflow in place.
laser.agent(id) acts as that agent with send, ask, contract, publish_card, and advertise (TypeScript: publishCard). advertise publishes a card and live presence the way a spawning agent does, and fails when the connection already advertises another agent.
The registry surface is agent_registry(), publish_card, advertise_presence, clear_presence, quarantine, unquarantine, quarantine_signed, unquarantine_signed, and client_metadata (TypeScript uses camelCase). A quarantined agent is excluded from routing until it is lifted. With a verifier on the registry, unsigned quarantine facts are dropped, so use the signed variants, which need the sign feature in Rust. Python passes cards and presence as dictionaries and returns registry entries as RegisteredCard objects. Python's client_metadata(..) returns a ClientMetadataPage (clients, next_cursor) and adds client_metadata_all, while Rust and TypeScript return a request builder. In Rust, presence and client metadata need the query feature.
resolve(skill, now) returns the cards that advertise a skill, are fresh, do not mark the skill unavailable, and are not quarantined. agents() and resolve list cards by agent ID in byte order, so a route that picks any agent, and a ranking tie, choose the same agent in every SDK. The card predicates apply the same checks to cards you read yourself. A card is fresh when it has no lifetime or its lifetime has not run out. serves checks that the card advertises the skill. available_for also checks that the skill is not marked unavailable.
import { SystemClock, cardAvailableFor, cardIsFresh } from "@laserdata/laser-sdk"
const registry = await laser.agentRegistry()
const now = new SystemClock().nowMicros()
await registry.refresh(now)
const ready = registry
.agents()
.filter((card) => cardIsFresh(card, now) && cardAvailableFor(card, "triage-ticket"))
.map((card) => card.agent.asStr())use laser_sdk::agent::{Clock, SystemClock};
let mut registry = laser.agent_registry()?;
let now = SystemClock.now_micros();
registry.refresh(now).await?;
let ready: Vec<String> = registry
.agents()
.filter(|card| card.is_fresh(now) && card.available_for("triage-ticket"))
.map(|card| card.agent.to_string())
.collect();registry = laser.agent_registry()
now = ls.SystemClock().now_micros()
await registry.refresh(now)
ready = [
card.agent
for card in registry.agents()
if card.is_fresh(now) and card.available_for("triage-ticket")
]Rust and Python call the predicates on RegisteredCard as is_fresh, serves, and available_for. TypeScript exports cardIsFresh, cardServes, and cardAvailableFor. The freshness check takes the time in epoch microseconds. Python's refresh() and resolve() default the time to now.
Scatter to every capable agent
scatter sends one task to every agent that advertises a capability, at once, and returns the reply bodies of those that completed. Unavailable and quarantined agents are left out, and no capable agent at all fails with no capable agent. Rust calls laser.scatter(source, &selector, payload, &route, deadline), TypeScript laser.scatter(source, selector, payload, inboxRoute, deadlineMs), and Python laser.scatter(skill, payload, source=.., deadline_ms=..). The deadline is required in all three.
scatter_report (TypeScript: scatterReport) returns a ScatterReport with each agent's own result in outcomes. completed() lists the agents that replied, with their replies. failures() lists only branches that errored before reaching a terminal outcome. An agent that answered Failed, TimedOut, or NotConsumed is in outcomes only.
Fan out from a handler
Inside a handler, fan_out(selector, payload, policy, deadline) (TypeScript: fanOut) sends a task to every agent with a capability, each on its own child conversation, and gathers the replies on the agent's respond_on topic. It returns a Gather with ok (agent and reply pairs) and failures (agent and error pairs). A target with no inbox is a failure entry, never rerouted.
| Gather policy | Rust | TypeScript | Python |
|---|---|---|---|
| Wait for every branch | GatherPolicy::RequireAll | { kind: "requireAll" } | policy="require_all", the default |
| Stop after n successes | GatherPolicy::Quorum(n) | { kind: "quorum", needed: n } | policy="quorum", quorum=n |
| Take what landed by the deadline | GatherPolicy::BestEffort | { kind: "bestEffort" } | policy="best_effort" |
Python's fan_out(skill, payload, ..) takes a required deadline_ms=, plus principal= and route_policy=, the parts that Rust and TypeScript carry in the capability selector. Branches follow the agent's inbox route in every SDK, and no call overrides it.
Workflows
laser.workflow(name) runs dependency-ordered steps. The name is the orchestrator identity the run publishes as, so it must be a valid agent ID, which is checked when the run starts. Each step targets one agent, one capable agent, or every capable agent, and builds its payload from the outputs of earlier steps.
import { AgentTopic, WorkflowBudget, routeToCapable } from "@laserdata/laser-sdk"
const text = new TextEncoder()
const outcome = await laser
.workflow("incident-response")
.inboxRoute({ kind: "fixed", topic: AgentTopic.Sessions })
.budget(WorkflowBudget.unlimited().invocations(8).wallClock(60_000))
.step("isolate", routeToCapable("isolate-host", { kind: "any" }), () =>
text.encode("isolate node-7")
)
.compensateWith(() => text.encode("release node-7"))
.step("patch", routeToCapable("patch-host", { kind: "any" }), ({ outputs }) => {
const isolated = new TextDecoder().decode(outputs.get("isolate") ?? new Uint8Array())
return text.encode(`patch after: ${isolated}`)
})
.after("isolate")
.verifyWith((output) => output.byteLength > 0)
.run()let outcome = laser
.workflow("incident-response")
.inbox_route(InboxRoute::Fixed(AgentTopic::Sessions))
.budget(WorkflowBudget::unlimited().invocations(8).wall_clock(Duration::from_secs(60)))
.step(
"isolate",
Router::to_capable("isolate-host", RoutePolicy::Any),
|_ctx: &StepContext<'_>| b"isolate node-7".to_vec(),
)
.compensate_with(|_ctx: &StepContext<'_>| b"release node-7".to_vec())
.step(
"patch",
Router::to_capable("patch-host", RoutePolicy::Any),
|ctx: &StepContext<'_>| {
let isolated = ctx.outputs.get("isolate").cloned().unwrap_or_default();
[b"patch after: ".to_vec(), isolated].concat()
},
)
.after("isolate")
.verify_with(|output: &[u8]| !output.is_empty())
.run()
.await?;workflow = laser.workflow("incident-response", fixed_inbox=ls.AgentTopic.Sessions)
workflow.budget(ls.WorkflowBudget.unlimited().invocations(8).wall_clock(60_000))
workflow.step(
"isolate",
to_capable="isolate-host",
build=lambda outputs: b"isolate node-7",
compensate=lambda outputs: b"release node-7",
)
workflow.step(
"patch",
to_capable="patch-host",
after=["isolate"],
build=lambda outputs: b"patch after: " + outputs["isolate"],
verify=lambda output: len(output) > 0,
)
outcome = await workflow.run()All three return a WorkflowOutcome with the outputs by step label and the run ID (outputs and run_id, TypeScript: runId).
Step options
| Option | Rust | TypeScript | Python |
|---|---|---|---|
| Depend on a step | after(label) | after(label) | after=[..] |
| Verify the output | verify_with(..) | verifyWith(..) | verify= |
| Compensate on a later failure | compensate_with(..) | compensateWith(..) | compensate= |
| Hold a fenced lease | exclusive(), exclusive_in(namespace) | exclusive(), exclusiveIn(namespace) | exclusive=True, fence_namespace= |
| On timeout | on_timeout(OnTimeout::Fail) or OnTimeout::Reassign | onTimeout("fail") or "reassign" | on_timeout="fail" or "reassign" |
| Target | Router::to, Router::to_capable, Router::all_capable | routeTo, routeToCapable, { kind: "allCapable", selector } | to=, to_capable=, all_capable= |
How a run behaves
- A run is a root session whose ID is the run ID. Each step and compensation is a child session. The run ends as completed, failed, or canceled.
- Steps run one at a time in dependency order. Independent steps do not run in parallel. An unknown dependency or a cycle fails the run.
- Each dispatch waits up to 30 seconds, cut to the wall-clock budget that remains. No SDK has a per-step deadline setting.
- An all-capable step scatters the task and joins the completed replies with newlines into one output.
- A step that fails its verifier, or ends as failed, timed out, or not consumed, fails the run with a handler configuration error after the compensations.
- Compensations run in reverse order when a later step fails. They are best effort: each goes to the step's router with a 30 second deadline and no fence, so a capability route may reach a different agent, and errors are ignored. An all-capable step's compensation never runs.
- Between steps the engine checks
agent.controlfor a cancel request. On a cancel it runs the compensations and fails with a cancelled error. Operators cancel a run withlaser.sessions().control(stream, run_id).as_operator(op).cancel(). Python'sCancelledErroralso subclassesasyncio.CancelledError. TypeScript also acceptsrun({ signal }). run_id(id)(TypeScript:runId) resumes an earlier run from its journal onagent.workflow_journal. The engine skips the steps recorded as complete and dispatches the rest. Rust and TypeScript take aConversationId, Python a string.- Builders, verifiers, and compensations can be asynchronous in all three SDKs. Builders and compensations return bytes and verifiers return a boolean. Python accepts a callable or an object with a
buildorverifymethod, and a compensation object also needsbuild. An SDK error raised in a callback keeps its classification, and cancelling the run cancels the active asynchronous callback.
Budgets
A WorkflowBudget caps tokens, wall-clock time, and invocations across the run. Build it with WorkflowBudget.unlimited() and chain wall_clock and invocations (TypeScript: wallClock). WorkflowBudget.tokens(n) starts a budget from a token cap. Rust takes the wall-clock time as a Duration, TypeScript and Python in milliseconds. Each chained call returns a new budget, and workflow.budget(budget) replaces the whole budget.
- Invocations are counted before each dispatch, and every limit is checked again after each step.
- The token cap counts only the usage that replies carry, so it is advisory. An all-capable step reports no tokens.
- The token cap is also written as the run session's budget and checked at every step boundary, on plain Apache Iggy too.
- A run with less than 100 ms of wall-clock budget left fails as budget exceeded instead of timing out.
- A breach runs the compensations and fails with a budget exceeded error (Rust
LaserError::BudgetExceeded, TypeScript and PythonBudgetExceededError), and ends the run session failed with reasonbudget.
Exclusive steps
An exclusive step holds a fenced lease, an ownership grant with an increasing token, while it runs. The lease key is the run ID, in the namespace agdx.workflow.fence unless you name another. The lease lasts 60 seconds, is renewed at half that, and is released on every exit. The handler's fenced write must use the same namespace and the run ID as its fence key, so the effect and the workflow check one fence sequence. Advanced key-value state explains fenced writes.
Reassign on timeout releases the lease, takes it again with a higher fence, and hands the task to a fresh holder, at most 2 times, so 3 attempts in total. It needs an exclusive step and is refused on an all-capable step.
Exclusive steps need the kv_fenced_leases capability, which Laser Stack and LaserData Cloud provide. Without it, the step fails with an unsupported error when it is dispatched, after earlier steps ran, so their compensations run.
Stop an agent
The handle owns the worker. Keep it until you stop or join the agent.
shutdown()stops new polls, drains the message in flight, and returns the consumer result.join()waits for the worker to exit on its own and returns its result. Periodic consolidation stays active until the worker exits.abort()stops at once.
import { Agent, AgentId, AgentTopic } from "@laserdata/laser-sdk"
const triage = Agent.builder()
.id(AgentId.new("triage"))
.listenOn(AgentTopic.Sessions)
.handler({ handle: async () => {} })
.shutdownGrace(10_000)
.build()
.spawn(laser)
await triage.ready()
await triage.shutdown()use laser_sdk::prelude::full::*;
use std::time::Duration;
struct Triage;
impl AgentHandler for Triage {
async fn handle(&self, _message: &AgentMessage, _ctx: &AgentCtx<'_>) -> Result<(), LaserError> {
Ok(())
}
}
let mut triage = Agent::builder()
.id("triage".parse::<AgentId>()?)
.listen_on(AgentTopic::Sessions)
.handler(Triage)
.shutdown_grace(Duration::from_secs(10))
.build()
.spawn(laser.clone());
triage.ready().await?;
triage.shutdown().await?;import laser_sdk as ls
async def handle(ctx, message):
pass
triage = laser.spawn_agent(
"triage",
ls.AgentTopic.Sessions,
handle,
shutdown_grace_ms=10_000,
)
await triage.ready()
await triage.shutdown()The shutdown grace sets how long pending SDK work may run after shutdown starts. When it expires, shutdown() fails with a timeout and the SDK cancels the work still in flight. The TypeScript ReliableConsumer takes the same limit as shutdownGraceMs.
Dropping the Rust handle signals a graceful shutdown and stops consolidation. Python does the same when the handle is collected. A dropped handle reports no result, so call shutdown() to see a drain failure. TypeScript has no destructor and stops through shutdown(), abort(), or await using. Python also supports async with laser.spawn_agent(..) as agent, which waits for readiness on entry and shuts the agent down on exit.
Handlers, middleware, deduplicators, dead-letter callbacks, and consolidators are your code. The SDK requests cancellation but cannot stop a callback that ignores it. In Python, let asyncio.CancelledError propagate out of every callback.
Layouts and many sessions in one stream
The stream is the isolation boundary. One stream can hold many independent root sessions, each with its own child sessions, and every agent topic name is the same in every stream. Give each application and environment its own stream.
On the default shared layout, agent.sessions carries all work, keyed by session, so different sessions can share a partition. Order holds within a partition, and causal links connect results across sources. Parent and root IDs do not merge state, add up budgets, or cancel descendants.
A layout is declared with the session configuration, laser.sessions_with(SessionConfig::new().layout(..)) in Rust, laser.sessions(new SessionConfig().layout(..)) in TypeScript, or laser.sessions(layout=..) in Python. The declaration holds for later sends on that stream through the same connection, so declare it in every process that writes to the stream. The per-agent partition layout moves commands to the addressee's declared partition and replies to the requester's. The single partition layout makes bootstrap create one partition per topic. The per-agent topic layout moves work addressed to a declared agent off the lane onto that agent's own topic. Agents keep listen_on set to agent.sessions and pass the same session configuration, and a declared agent then reads its own topic. Callers keep sending to agent.sessions, and the SDK routes at send time. Declare a topic per agent shows it end to end, and Advanced sessions compares the layouts.
Key operations
| Call | What it does |
|---|---|
bootstrap(partitions, retention) | Create the agent topics on a stream once |
Agent::builder()...spawn(laser) | Define and start an agent (Python: spawn_agent) |
ready() | Wait until the agent reads its topic |
contract(router) | Send one task to one agent and wait for the outcome |
laser.agent(id) | Act as one agent: send, ask, contract, publish a card, advertise presence |
agent_registry() | Read cards, presence, and quarantine facts |
quarantine(..), quarantine_signed(..) | Exclude an agent from routing |
scatter(..), scatter_report(..) | Send one task to every capable agent |
ctx.fan_out(..) | Gather replies from every capable agent inside a handler |
consumed(target, at) | Check whether a consumer committed past a position |
sessions().submit(agent, input) | Hand a new session to an agent |
workflow(name) | Run ordered steps with budgets, verifiers, and compensation |
shutdown(), join(), abort() | Stop an agent |
Where it runs
Agents, contracts, scatter, fan-out, and workflows work on plain Apache Iggy, Laser Stack, and LaserData Cloud. Live presence needs a server that serves it. Exclusive steps need kv_fenced_leases. Budget enforcement in agents needs the managed session index.
In Rust, the agent runtime needs the agent feature. Add query for live presence and client metadata, sign for signing keys, verifiers, and signed quarantine, and kv for exclusive steps. managed includes query and kv but not agent or sign.
Related
Agents
The short guide to agents
Advanced sessions
Session lifecycle, control, layouts, and reads
Governance
Roles, pre-action policy, budgets, and approvals
Interop
A2A, MCP, and AG-UI bridges in front of agents
Orchestra example
Discovery, workflows, quarantine, and deadline recovery in all three languages