lash/advanced

Storage boundaries, session lifecycle, process work, durable workflows, subagents, MCP servers, the advanced runtime surface, and a complete example.

App Storage

Product storage stays separate from runtime storage.

App state

Chat tables, accounts, frontend state, auth, transport. The example app uses SQLite plus NDJSON to the browser.

Runtime state

Pass a SessionStoreFactory to LashCoreBuilder::store_factory, or a concrete store to SessionBuilder::store, for durable runtime state across process restarts.

Explicit Runtime Facets

Mode presets install protocol and plugin defaults only. Every core build names the host-owned runtime facets explicitly: an effect host, a Lashlang artifact store, and an attachment store. Add a session-store factory when sessions must survive process restarts or when the inline process worker is enabled.

use std::sync::Arc;

let factory = lash::rlm::RlmProtocolPluginFactory::new(
    lash::rlm::RlmProtocolPluginConfig::builder()
        .instruction_limit(lash::rlm::InstructionBound::instructions(1_000_000))
        .wall_clock(lash::rlm::WallClockBound::secs(30))
        .memory_limit(lash::rlm::MemoryBound::mebibytes(64))
        .build(),
    Arc::new(lash::persistence::InMemoryLashlangArtifactStore::new()),
);
let core = lash::LashCore::rlm_builder(lash::TurnBudget::Unbounded, factory)
    .provider(provider)
    .model(model)
    .effect_host(Arc::new(lash::durability::InlineEffectHost::default()))
    .attachment_store(Arc::new(lash::persistence::InMemoryAttachmentStore::new()))
    .process_env_store(Arc::new(
        lash::persistence::InMemoryProcessExecutionEnvStore::new(),
    ))
    // Start bounded; tune both limits for your backend's latency envelope.
    .commit_budget(lash::CommitBudget::bounded(1024 * 1024, 512))
    .queued_work_batching(lash::QueuedWorkBatchingConfig::new(1024))
    .build(crate::example_process_owner())?;

The supplied owner is the host identity: stable owner_id, fresh per-boot incarnation_id. Lash mints a fresh executor_id for every runtime open. Lease reentry requires all three values to match, so resuming the same parked open reenters but a second concurrent open under the same core observes Busy and cannot revoke the first executor's fence.

use std::sync::Arc;

use lash::persistence::FileAttachmentStore;
use lash_sqlite_store::{SqliteSessionStoreFactory, Store};

let store_factory = Arc::new(SqliteSessionStoreFactory::new(data_dir.join("sessions")));
let artifact_store = Arc::new(Store::open(&data_dir.join("artifacts.db")).await?);

let factory = lash::rlm::RlmProtocolPluginFactory::new(
    lash::rlm::RlmProtocolPluginConfig::builder()
        .instruction_limit(lash::rlm::InstructionBound::instructions(1_000_000))
        .wall_clock(lash::rlm::WallClockBound::secs(30))
        .memory_limit(lash::rlm::MemoryBound::mebibytes(64))
        .build(),
    artifact_store.clone(),
);
let core = lash::LashCore::rlm_builder(lash::TurnBudget::Unbounded, factory)
    .provider(provider)
    .model(model)
    .store_factory(store_factory)
    .effect_host(Arc::new(lash::durability::InlineEffectHost::default()))
    .attachment_store(Arc::new(FileAttachmentStore::new(
        data_dir.join("attachments"),
    )))
    .process_env_store(artifact_store)
    // Start bounded; tune both limits for your backend's latency envelope.
    .commit_budget(lash::CommitBudget::bounded(1024 * 1024, 512))
    .queued_work_batching(lash::QueuedWorkBatchingConfig::new(1024))
    .build(crate::example_process_owner())?;

Session Lifecycle

Two states: active (resident, ready) and parked (flushed; only id + store reference held). Opening an existing id rehydrates from the store: one read of the graph head plus referenced snapshots.

Sessions park automatically when the handle drops. Long-running servers can hold thousands of session ids cheaply. Reopening a session loads its one durable leaf-to-root history.

Plugin background work (host-managed process handles, post-turn observation summaries, MCP warm-ups) continues past turn return. To drain it before returning (one-shot processes that need observations persisted before exit):

session.refresh_background_graph().await?;

Exposed by the CLI as --await-background-work on lash --print. Embedders rarely need it outside short-lived processes.

Process Work

Lashlang process blocks need a process plane. The registry records process intent, handles, wakes, events, and cancellation; the process work driver consumes those records, executes background work, and owns process waits.

inline runner

Call .process_registry(...). The facade builds the default inline DurableProcessWorkerConfig at build() and starts process work from captured execution environments; no origin session has to be reopened for tool or Lashlang process execution. .process_execution_concurrency(...) sets the per-worker inline execution bound (default 64, minimum 1). Runs release their slot while parked on another process or external completion and reacquire before resuming. The bound is not registry-global: two workers over one registry can execute twice the configured number.

external runner

Call .process_work_driver(...) when a deployment worker, queue, or workflow runtime owns process execution. The driver's registry becomes the core process registry and Lash does not spawn the inline runner.

observability

core.processes() is the single host-level surface — start, get, list, events, signal, await_output, cancel, cancel_all, transfer, prune, and request_abandon — with two distinct list filters (ADR 0019): list_observed_by is the observer lens (addressability — what a session may see) and list_originated_by is the provenance lens (origin — what a session created; a process it started then transferred away still matches here, one merely observed by it does not). core.process_registry() is the point-read/write registry for storage integrations, and session.processes() is thin observer-scoped sugar over the global surface for in-session await/cancel/list.

waiting & retention

ProcessWorkDriver::await_terminal is the one way to wait on a started work item (ADR 0016) — an engine-native durable promise or hub-plus-backoff point reads, never a store poll loop; bound every wait with tokio::time::timeout. An optional ProcessEventSink, installed at construction (new_with_sink / LashCoreBuilder::process_event_sink), pushes appended events best-effort for a live feed — freshness, never truth (ADR 0017), so reconcile from events_after and take terminal state from await_terminal. core.processes().prune(cutoff_epoch_ms, watermark) is the host retention lever for terminal rows; pass ProjectionWatermark::UpTo(cursor) or explicitly choose NoProjector. Use core.processes().compact_tombstones for delivery-aware tombstone reclamation.

let factory = lash::rlm::RlmProtocolPluginFactory::new(
    lash::rlm::RlmProtocolPluginConfig::builder()
        .instruction_limit(lash::rlm::InstructionBound::instructions(1_000_000))
        .wall_clock(lash::rlm::WallClockBound::secs(30))
        .memory_limit(lash::rlm::MemoryBound::mebibytes(64))
        .build(),
    artifact_store,
);
let core = lash::LashCore::rlm_builder(lash::TurnBudget::Unbounded, factory)
    .provider(provider)
    .model(model)
    .store_factory(store_factory)
    .process_registry(process_registry)
    .effect_host(Arc::new(lash::durability::InlineEffectHost::default()))
    .attachment_store(attachment_store)
    .process_env_store(process_env_store)
    // Start bounded; tune both limits for your backend's latency envelope.
    .commit_budget(lash::CommitBudget::bounded(1024 * 1024, 512))
    .queued_work_batching(lash::QueuedWorkBatchingConfig::new(1024))
    .build(crate::example_process_owner())?;

Lashlang Modules

Hosts that accept or ship Lashlang modules should compile through compile_module. The facade returns the persisted artifact refs, stable introspection, exported process definition identities, and structured diagnostics; trigger registration screens can run compatibility checks before saving a subscription.

use std::collections::BTreeMap;

let mut resources = lashlang::LashlangHostCatalog::new();
resources.add_trigger_source_constructor(
    ["app", "button"],
    lashlang::TypeExpr::Object(Vec::new()),
    lashlang::NamedDataType::object(
        "app.ButtonPressed",
        vec![lashlang::TypeField {
            name: "color".into(),
            ty: lashlang::TypeExpr::Str,
            optional: false,
        }],
    )?,
)?;

let environment = lashlang::LashlangHostEnvironment {
    resources,
    abilities: lashlang::LashlangAbilities::default().with_processes(),
    ..lashlang::LashlangHostEnvironment::default()
};

let compiled = lashlang::compile_module(lashlang::ModuleCompileRequest {
    source: r#"
process on_button(event: app.ButtonPressed) {
  finish event.color
}
source = app.button({})
finish source
"#,
    environment: &environment,
    artifact_store: Some(artifact_store.as_ref()),
})
.await?;

let process = compiled
    .introspection
    .exported_processes
    .iter()
    .find(|process| process.definition.process_name == "on_button")
    .expect("compiled module exports on_button");

let inputs = lashlang::TriggerInputTemplate::new(BTreeMap::from([(
    "event".to_string(),
    lashlang::TriggerInputBinding::Event,
)]));
let compatibility =
    lashlang::check_trigger_compatibility(lashlang::TriggerCompatibilityRequest {
        artifact: &compiled.artifact,
        definition: &process.definition,
        source_type: "app.button",
        inputs: &inputs,
    })?;

println!(
    "compiled {} and trigger emits {}",
    compiled.module_ref,
    compatibility.event_type.name()
);

Durable Workflows

Standard LashCore persists at the turn boundary. For durable in-flight effects (expensive LLM calls, long tool calls, process admin, mid-turn worker migration), wire an EffectHost on the builder and run each durable handler turn with .turn_id(...).effects(&controller).run() or .turn_id(...).effects(&controller).stream_to(&sink). The facade creates the turn ExecutionScope internally, while the handler still supplies the workflow controller explicitly. The EffectHost is the stable boundary an external workflow runtime (Temporal, Restate) would wrap; the default InlineEffectHost runs in process and is not in-flight durable.

Shape: the runtime driver polls TurnMachine. For each reply-producing turn effect it builds a RuntimeEffectEnvelope with a typed RuntimeInvocation, invokes the active ScopedEffectController, and feeds the returned outcome back as the matching response. The InlineEffectHost runs locally and is not in-flight durable. A durable workflow EffectHost such as Restate records effects in workflow history, uses durable timers for sleeps, and reruns the handler with the same turn id before Lash retries the final idempotent commit.

Do not wrap a durable handler turn in another persistent submitted/running work-item state machine. The workflow key and turn_id are the durable in-flight identity; Lash's session execution lease, TurnWorkDriver keyed promises, and idempotent commit are the runtime correctness boundary; host tables should store product rows, reconnect outbox rows, exact turn addresses, or terminal output.

Full contract: Architecture → Durability. Embedders without workflow orchestration use the regular session-store path.

Subagents

Same SessionSpec shape as root sessions. The factory keeps the capability registry; child policy resolves from the live parent snapshot.

use std::sync::Arc;

use lash::{SessionSpec, plugins::PluginFactory};
use lash_subagents::{SubagentsPluginFactory, default_registry};

let registry = Arc::new(default_registry(&tier_models));

let child_spec = SessionSpec::inherit().turn_budget(lash::TurnBudget::bounded(8));
let _configured_turn_budget = child_spec.turn_budget;
let subagents = SubagentsPluginFactory::new(registry).with_session_spec(child_spec);

let factory = lash::rlm::RlmProtocolPluginFactory::new(
    lash::rlm::RlmProtocolPluginConfig::builder()
        .instruction_limit(lash::rlm::InstructionBound::instructions(1_000_000))
        .wall_clock(lash::rlm::WallClockBound::secs(30))
        .memory_limit(lash::rlm::MemoryBound::mebibytes(64))
        .build(),
    Arc::new(lash::persistence::InMemoryLashlangArtifactStore::new()),
);
let core = lash::LashCore::rlm_builder(lash::TurnBudget::Unbounded, factory)
    .provider(provider)
    .model(
        lash::ModelSpec::builder(model.clone())
            .context_window_tokens(200_000)
            .build()
            .expect("valid model metadata"),
    )
    .effect_host(Arc::new(lash::durability::InlineEffectHost::default()))
    .attachment_store(Arc::new(lash::persistence::InMemoryAttachmentStore::new()))
    .process_env_store(Arc::new(
        lash::persistence::InMemoryProcessExecutionEnvStore::new(),
    ))
    // Start bounded; tune both limits for your backend's latency envelope.
    .commit_budget(lash::CommitBudget::bounded(1024 * 1024, 512))
    .queued_work_batching(lash::QueuedWorkBatchingConfig::new(1024))
    .plugin(Arc::new(subagents) as Arc<dyn PluginFactory>)
    .build(crate::example_process_owner())?;

Capabilities return SessionSpec overlays. StaticCapability pins exact child authority; TierCapability backs the built-in explore and peer tiers. Do not construct SessionPolicy directly; it's the resolved runtime artifact.

Built-in tiers

explore

Read-only investigator. RLM by default. Cannot spawn further subagents. For scan / summarise / verify tasks.

peer

Concurrent self. Inherits parent execution mode and tool authority, including recursive spawning. For a sibling branch of the same work.

Both tiers consult tier_models for an override, then fall back to the parent's model. Pass overrides through default_registry(&tier_models).

Capability resolution

Capability::build_session_request(ctx) sees the live parent SessionPolicy and returns a SessionCreateRequest. The runtime composes the capability's SessionSpec against parent policy to produce effective child config. The spec is what's persisted, so a resumed subagent re-derives the same authority.

Interactive-only tools such as ask and showcase would block a subagent indefinitely; nothing is listening on the other side. Hide them from every child surface with SubagentsPluginFactory::with_hidden_tools(...) or a custom SessionToolAccess; the CLI hides its five interactive tools this way. Recursion is bounded separately: at the maximum spawn depth (5) the runtime hides spawn_agent itself.

Usage attribution

Child sessions emit TurnEvent::ChildUsage on the parent's stream, tagged with source ("subagent", "compaction", "observer") and child session_id. TurnReport.children_usage rolls up per (source, model); total_usage() sums parent + children. Finer splits: fold ChildUsage from the stream directly.

MCP Servers

lash-plugin-mcp wraps the rmcp SDK with two transports: stdio and streamable_http. SSE-capable servers are reached through streamable_http, which negotiates SSE responses itself. HTTP headers are static across reconnects: Lash does not enable rmcp OAuth or token refresh, so the host must supply and rotate credentials. The plugin keeps a shared connection pool; stdio servers spawn once per process, not per session.

use std::collections::BTreeMap;
use std::time::Duration;

use lash_plugin_mcp::{MCP_PROTOCOL_VERSION, McpPluginFactory, McpServerConfig};

println!("Lash MCP handshake protocol: {MCP_PROTOCOL_VERSION}");

let mut servers = BTreeMap::new();
servers.insert(
    "docs".to_string(),
    McpServerConfig::stdio("uvx", vec!["mcp-server-docs".into()])
        .with_env([("DOCS_INDEX", "/srv/docs")]),
);
servers.insert(
    "web".to_string(),
    // Headers are static across reconnects: Lash does not perform OAuth
    // or token refresh, so the host must rotate credentials itself.
    McpServerConfig::streamable_http("https://mcp.example.com/rpc")
        .with_headers([("Authorization", "Bearer host-managed-token")])
        .with_timeouts(
            Duration::from_secs(10),
            Duration::from_secs(60),
            Duration::from_secs(600),
        ),
);

// Capabilities are advertised only for handlers installed here. Lash
// routes server requests; the host still owns model, UI, and root policy.
let mcp = McpPluginFactory::builder(servers)
    .sampling_handler(sampling)
    .elicitation_handler(elicitation)
    .roots_provider(roots)
    .build()
    .await?;

let core = configured_mcp_core(provider, model, mcp)?;
Ok(())
}

fn configured_mcp_core(
provider: ProviderHandle,
model: String,
mcp: McpPluginFactory,
) -> lash::Result<LashCore> {
let factory = lash::rlm::RlmProtocolPluginFactory::new(
    lash::rlm::RlmProtocolPluginConfig::builder()
        .instruction_limit(lash::rlm::InstructionBound::instructions(1_000_000))
        .wall_clock(lash::rlm::WallClockBound::secs(30))
        .memory_limit(lash::rlm::MemoryBound::mebibytes(64))
        .build(),
    std::sync::Arc::new(lash::persistence::InMemoryLashlangArtifactStore::new()),
);
lash::LashCore::rlm_builder(lash::TurnBudget::Unbounded, factory)
    .provider(provider)
    .model(
        lash::ModelSpec::builder(model.clone())
            .context_window_tokens(200_000)
            .build()
            .expect("valid model metadata"),
    )
    .effect_host(std::sync::Arc::new(
        lash::durability::InlineEffectHost::default(),
    ))
    .attachment_store(std::sync::Arc::new(
        lash::persistence::InMemoryAttachmentStore::new(),
    ))
    .process_env_store(std::sync::Arc::new(
        lash::persistence::InMemoryProcessExecutionEnvStore::new(),
    ))
    // Start bounded; tune both limits for your backend's latency envelope.
    .commit_budget(lash::CommitBudget::bounded(1024 * 1024, 512))
    .queued_work_batching(lash::QueuedWorkBatchingConfig::new(1024))
    .plugin(std::sync::Arc::new(mcp))
    .build(crate::example_process_owner())
}

Tools appear as mcp__<server>__<tool> with original schemas preserved. attach_server / detach_server mutate the live pool without rebuilding the core.

Server-to-client requests: host owns policy

MCP 2025-11-25 sampling, elicitation, and roots are available through three dyn-compatible host seams: McpSamplingHandler, McpElicitationHandler, and McpRootsProvider. Each callback receives one request struct whose sealed context identifies the server and exposes a cancellation token; interactive UI and model work must stop when that token fires. Lash validates accepted form content against the server's requested schema before writing it to the wire. Decline and cancel carry no content. URL-capable elicitation handlers also receive the matching notifications/elicitation/complete id.

Lash never picks a model, fabricates user input, or guesses workspace roots. The initialize handshake advertises each capability only when its handler was installed on McpPluginFactory::builder; an elicitation handler declares its exact form and URL modes. Call notify_roots_changed() after a dynamic roots provider changes. Connected servers receive notifications/roots/list_changed; a server that reconnects reads the current value through roots/list.

Sampling and elicitation execute as ordinary host I/O inside the open MCP tool attempt. Handler seams must not emit journaled effects or submit nested Lash commands: the outer attempt is the single durable journal unit. A sampling handler owns its model and billing. Its spend belongs in the host/provider trace only; it is not added to the Lash session usage ledger or TurnReport usage.

Lifecycle and failure isolation

Servers connect once at McpPluginFactory::new(...). Built session plugins clone the shared Arc<McpConnectionPool>, so startup cost is one-per-process, not one-per-chat.

Each pool entry has one lifecycle actor that exclusively owns its running service, stdio child, generation counter, reconnect backoff, and liveness timers. Tool calls do not pass through that actor: they clone the currently published peer and generation from a read-optimized cell. A stale failure can therefore report its generation without disconnecting a newer service. This is per-connection resource plumbing, not a pool-wide scheduler or host drain policy (ADR 0014).

After attach_server, new tools are visible to sessions on their next tool-catalog refresh (start of next turn). detach_server shuts the server down and removes its tools; in-flight calls fail rather than hang.

// Hot-swap a server at runtime.
mcp.attach_server(
    "new-tool".to_string(),
    McpServerConfig::stdio("uvx", vec!["mcp-server-new".into()]),
)
.await?;

mcp.detach_server("old-tool").await?;

// Call this after the host's dynamic roots change.
mcp.notify_roots_changed().await?;

Explicit pool shutdown sends every entry actor Shutdown, then joins all entries concurrently. Shutdown preempts handshake, liveness-probe, and host-request waits, leaving each entry's five-second total deadline to cover a three-second graceful child close, kill, a one-second bounded reap, and scheduling margin. A child can be abandoned if it survives that preemptive kill and reap or if the entry deadline expires mid-reap; the abort branch reports the live active_pid, and its literal PID and reason are recorded in status and tracing. There is deliberately no background waitpid sweep. Ordinary drop remains kill-only and cannot promise reap, so hosts should use the factory/core shutdown lever during an orderly exit.

Failures are isolated: an exited child, a 500ing endpoint, or a dropped SSE stream surfaces as a structured unavailable ToolOutcome on the affected tool; other servers keep working. Both initial construction and attach_server register an unreachable server and retry it in the background. Configuration errors still fail immediately, and server_statuses() exposes connectivity and the latest error.

Advanced Runtime

.advanced() is reserved for low-level overrides: a custom PluginHost or a replaceable RuntimeHostConfig. Durable embedding concerns such as effect hosts, process registries, process work drivers, residency policy, and termination policy live on the normal LashCoreBuilder. The registry trait itself lives under lash::process::ProcessRegistry and stores state only; process waits live on the work driver. RLM/Lashlang turns that generic process plane into authored process/lifecycle abilities when the Lashlang runtime pieces are installed.

let artifact_store =
    std::sync::Arc::new(lash_sqlite_store::Store::open(&data_dir.join("artifacts.db")).await?);
let factory = lash::rlm::RlmProtocolPluginFactory::new(
    lash::rlm::RlmProtocolPluginConfig::builder()
        .instruction_limit(lash::rlm::InstructionBound::instructions(1_000_000))
        .wall_clock(lash::rlm::WallClockBound::secs(30))
        .memory_limit(lash::rlm::MemoryBound::mebibytes(64))
        .build(),
    artifact_store.clone(),
);
let core = lash::LashCore::rlm_builder(lash::TurnBudget::Unbounded, factory)
    .provider(provider)
    .model(
        lash::ModelSpec::builder("anthropic/claude-sonnet-4.6")
            .context_window_tokens(200_000)
            .build()
            .expect("valid model metadata"),
    )
    .store_factory(store_factory)
    .effect_host(std::sync::Arc::new(
        lash::durability::InlineEffectHost::default(),
    ))
    .attachment_store(std::sync::Arc::new(
        lash::persistence::FileAttachmentStore::new(data_dir.join("attachments")),
    ))
    .process_env_store(artifact_store)
    // Start bounded; tune both limits for your backend's latency envelope.
    .commit_budget(lash::CommitBudget::bounded(1024 * 1024, 512))
    .queued_work_batching(lash::QueuedWorkBatchingConfig::new(1024))
    .turn_budget(lash::TurnBudget::bounded(50))
    // The outer bound on one turn's work; the inner bound is how long a
    // turn may fail to do any. This one defaults to bounded — name it to
    // change it, or to opt out with `NoProgressBudget::Unbounded`.
    .no_progress_budget(lash::NoProgressBudget::bounded(12))
    .build(crate::example_process_owner())?;

Streaming is semantic: TurnBuilder::stream emits TurnActivity and resolves to a TurnReport. Raw runtime telemetry belongs in tracing, not the app surface.

QueuedWorkBatchingConfig also names how much queued work one wake drains. Lash decides which rows may share a turn — merge key, delivery boundary, authority, kind, row and age bounds — and the host decides how many of that strictly FIFO prefix actually drain. The default is DrainMode::OneAtATime: one queued row per drain, no token arithmetic. with_drain_mode(DrainMode::All) drains every compatible pending row for large-window hosts that want throughput, and with_drain_policy(...) takes a custom QueuedDrainPolicy for anything else; the policy sees each row's conservative projection and the window budget. A single row that renders larger than the whole context window can never drain, so it is refused with a typed SelectedQueuedWorkDrainRefusalCause::QueuedItemExceedsContextWindow naming the row and the window it needs, and stays queued for host recovery. The policy sizes automatic drains only: a selection your host names explicitly (queued_turn().batch_ids(...)) already states the composition, so it drains exactly the rows you asked for.

Complete Example

The browser example: application chat database, RLM mode, typed session plugin activation, app-owned board tools, semantic stream events, terminal rendering from TurnOutput, and optional Restate-backed turns.

OPENROUTER_API_KEY=... cargo run -p agent-service
# then open http://127.0.0.1:3000

Source: examples/agent-service. Walkthrough: Agent Service.

read on ·