LaserData Cloud
Laser SDKAdvanced

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.

SettingRust Agent::builder()TypeScript Agent.builder()Python spawn_agent(..)
Identityid(..)id(..)First argument, a name or None
Input topiclisten_on(..)listenOn(..)Second argument
Handlerhandler(..)handler(..)Third argument
Reply topic for respondrespond_on(..)respondOn(..)respond_on=
Consumer groupconsumer_group(..)consumerGroup(..)consumer_group=
Capabilitiescapabilities(vec![..])capabilities([..])capabilities=[..]
Pickup statusack_on_pickup(true)ackOnPickup()ack_on_pickup=True
Served operationsoperations(vec![..])operations([..])operations=[..]
Session configurationsessions(SessionConfig)sessions(config)sessions=laser.sessions(..)
Default route for directed sendsinbox_route(..)inboxRoute(..)fixed_inbox=
Poll intervalpoll_interval(Duration)pollInterval(ms)poll_interval_ms=
Shutdown graceshutdown_grace(Duration)shutdownGrace(ms)shutdown_grace_ms=
Retryretry(RetryPolicy)retry({ maxAttempts, baseDelayMs })retry_max_attempts=, retry_base_delay_ms=
Dedup windowdedup_window(n)dedupWindow(n)dedup_window=
Custom deduplicatordeduplicator(..)deduplicator(..)dedup=
Load dedup keys on startwarm_dedup(true)warmDedup()warm_dedup=True
Partition lanesconcurrency(ConcurrencyPolicy::SerialPerPartition { max_partitions })concurrency({ kind: "serial-per-partition", maxPartitions })max_partitions=
Queue boundsmax_queued_records(n), max_queued_bytes(n)maxQueuedRecords(n), maxQueuedBytes(n)max_queued_records=, max_queued_bytes=
Understood feature bitsunderstood_features(bits)understoodFeatures(bits)understood_features=
Middlewaremiddleware(vec![..])middleware(..), once per hook objectmiddleware=[..]
Dead-letter callbackon_dead_letter(..)onDeadLetter(..)dead_letter=
Governorgovernor((governor, mode))governor([governor, mode])governor=, governor_mode=
Governor retentiongovernor_retention(GovernorRetention { capacity, idle_ttl })governorRetention({ capacity, idleTtlMs })governor_retention=(capacity, idle_ttl_ms)
Periodic consolidationconsolidate_every(..), consolidator(..)consolidateEvery(ms), consolidator(..)consolidate_every_ms=, consolidator=
Signing and verificationsigning_key(..), verifier(..)signingKey(..), verifier(..)signing_key=, verifier=

Defaults in all three SDKs:

SettingDefault
Consumer groupThe agent ID
Poll interval10 ms
Shutdown grace30 seconds
Retry5 attempts in total, 200 ms before the first retry, doubling each time
Dedup window10,000 keys in memory
ConcurrencySerial, one message at a time across all partitions
Queue bounds4,096 records and 64 MiB across partition lanes
Pickup statusOff
Served operationsAll

More on individual settings:

  • respond_on is required for respond and fan_out. Without it, respond fails 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_on topic) 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_pickup makes the agent report a Working status 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, and load. The health= keyword sets one value (healthy, degraded, or unavailable) on every advertised skill.
  • Pass None as the identity and an explicit consumer_group to 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_features is 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:

CallWhat 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_dedup to 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.dlq and 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 DeadlineExceeded before 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: InvalidError and TypeError dead-letter at once, PermissionError maps to unauthorized, the built-in TimeoutError is retryable, and a plain ValueError is 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.control across all its partitions, outside its consumer group, so every instance sees each pause and cancel request. A handler reads them through ctx.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 budget and never reaches the handler.
  • Dead letters keep the record's conversation. An agent whose session configuration turns on fail_on_dead_letter fails 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.

ResultMeaning
ConsumedCommitted past the position. Carries the committed offset and the partition head
NotYetConsumedStill behind, by behind_by records
SkippedThe 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.

OptionRust and TypeScriptPythonDefault
Senderfrom(agent)source=Required
Task bodypayload(..)Second argumentEmpty
Inbox routeinbox_route(..)fixed_inbox=The advertised inbox
Reply topicreply_on(topic)reply_on=agent.sessions
Consumption expiryexpire_if_not_consumed(d)expire_if_not_consumed_ms=None
Reply deadlinedeadline(d)deadline_ms=30 seconds
Conversationconversation(id)conversation=A new one per send
Fence tokenfence(token)fence=None
Child sessionparent(parent, root)parent=, root=None
Route policyIn the routerpolicy=Any
Required principalRouter::to_principal, routeToPrincipalprincipal=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 with ack_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.sessions before 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. respond addresses 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.

PolicyRustTypeScriptPython
Any candidateRoutePolicy::Any{ kind: "any" }"any" or None
Lowest cost classRoutePolicy::Cheapest{ kind: "cheapest" }"cheapest"
Lowest latency classRoutePolicy::Fastest{ kind: "fastest" }"fastest"
Lowest loadRoutePolicy::LeastLoaded{ kind: "leastLoaded" }"least_loaded"
Prefer one agentRoutePolicy::Sticky(agent){ kind: "sticky", agent }"sticky:<agent>"
Your own rankingRoutePolicy::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 policyRustTypeScriptPython
Wait for every branchGatherPolicy::RequireAll{ kind: "requireAll" }policy="require_all", the default
Stop after n successesGatherPolicy::Quorum(n){ kind: "quorum", needed: n }policy="quorum", quorum=n
Take what landed by the deadlineGatherPolicy::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

OptionRustTypeScriptPython
Depend on a stepafter(label)after(label)after=[..]
Verify the outputverify_with(..)verifyWith(..)verify=
Compensate on a later failurecompensate_with(..)compensateWith(..)compensate=
Hold a fenced leaseexclusive(), exclusive_in(namespace)exclusive(), exclusiveIn(namespace)exclusive=True, fence_namespace=
On timeouton_timeout(OnTimeout::Fail) or OnTimeout::ReassignonTimeout("fail") or "reassign"on_timeout="fail" or "reassign"
TargetRouter::to, Router::to_capable, Router::all_capablerouteTo, 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.control for a cancel request. On a cancel it runs the compensations and fails with a cancelled error. Operators cancel a run with laser.sessions().control(stream, run_id).as_operator(op).cancel(). Python's CancelledError also subclasses asyncio.CancelledError. TypeScript also accepts run({ signal }).
  • run_id(id) (TypeScript: runId) resumes an earlier run from its journal on agent.workflow_journal. The engine skips the steps recorded as complete and dispatches the rest. Rust and TypeScript take a ConversationId, 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 build or verify method, and a compensation object also needs build. 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 Python BudgetExceededError), and ends the run session failed with reason budget.

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

CallWhat 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.

On this page