The Python SDK - foreign/python/, PyO3 bindings over the Rust laser-sdk crate. Use when changing the Python surface (PyLaser, the publish/query/kv/fork builders, the agent runtime and the async-callback handler consumer, errors, stubs), the maturin packaging, the .pyi stubs, the pytest suite, or the Python BDD runner under bdd/python/.
Installs into .claude/skills of the current project.
Are you the author of Python Bindings?
Add the live security badge to your README. It updates with every re-scan.
[](https://www.skillsdirectory.com/skills/laserdata-python-bindings)
---
name: python-bindings
description: The Python SDK - foreign/python/, PyO3 bindings over the Rust laser-sdk crate. Use when changing the Python surface (PyLaser, the publish/query/kv/fork builders, the agent runtime and the async-callback handler consumer, errors, stubs), the maturin packaging, the .pyi stubs, the pytest suite, or the Python BDD runner under bdd/python/.
---
# Python bindings (foreign/python)
Python uses PyO3 bindings over the Rust `laser-sdk` crate. Keep data encoding and runtime behavior in Rust. The native TypeScript implementation is maintained separately under `foreign/typescript/`.
## Layout
- `foreign/python/Cargo.toml` - the binding crate (`laser_sdk_py` lib, `cdylib`), depends on `laser-sdk` with the managed and agent surfaces, and uses Iggy's native VSR transport. The crate is excluded from the root workspace and has its own lock.
- `pyproject.toml` - maturin build, PyPI name `laser-sdk`, import name `laser_sdk` (set via `tool.maturin.module-name`, since the cdylib lib name is `laser_sdk_py` to avoid clashing with the Rust crate's `laser_sdk` lib).
- `src/client.rs` and `src/stream.rs` bind `Laser`, capabilities, the stream/topic accessor grammar, typed topics, publish/batch/replay/ensure, and `producer`/`consumer`/`consumer_group`.
- `src/transport.rs` binds `Producer`, `Consumer`, and `ConsumerMessage`. Preserve routing, batching, retry, polling, group, header, and commit behavior. It also provides offset inspection and continuous asynchronous reads.
- Direct and fluent sends return immutable `SendMessagesResponse` and `SendMessagesConfirmationResponse` classes. Keep their stream, topic, partition, and base-offset fields in parity with Rust, regenerate the stub after changes, and never turn an empty confirmation list into a client error.
- `src/publish.rs`, `reader.rs`, `typed.rs`, `schema.rs`, `query.rs`, `watch.rs`, `kv.rs`, `fork.rs`, `graph.rs`, and `filters.rs` bind the data-platform write and read surfaces. `filters.rs` binds `ConsumerFilter`, `FilterExpr`, the filtered reader, and the catalog client, and must stay in parity with `sdk/src/filters/`. Complex managed shapes cross through serde instead of hand-maintained mirror classes.
- `src/batching.rs` binds `Topic.batching()` and `BatchingProducer` (`send`, `flush`, `close`). `Topic.send(payload, *, headers, partition_key)` and `Topic.batch(messages, *, partition_key)` mirror the Rust raw sends.
- `src/registry.rs` binds `Laser.agent_registry()` (`AgentRegistry` with `refresh`, `refresh_presence`, `agents`, `lookup`, `resolve`, `is_quarantined`, `inbox_for`, `inbox_for_principal`, `principal_for`), `publish_card`, `quarantine_signed`, `unquarantine_signed`, `advertise_presence`, `clear_presence`, `client_metadata` (a `ClientMetadataPage` with `clients` and `next_cursor`), and `Laser.agent(id)` returning `AgentScope` (`send`, `ask`, `contract`, `publish_card`, `advertise`, `id`). `agents`, `lookup`, and `resolve` return `RegisteredCard` objects (`agent`, `card`, `observed_at_micros`, `is_fresh`, `serves`, `available_for`), and `resolve_targets` mirrors Rust. The card itself and presence cross as dicts. Route scorers receive `RouteCandidate` objects.
- `src/agent.rs`, `agdx.rs`, `agent_runtime.rs`, `workflow.rs`, `runs.rs`, `rbac.rs`, `context.rs`, `memory.rs`, `state_store.rs`, and `sign.rs` expose agent behavior. Keep these bindings consistent with Rust. `laser.agdx(..., signing_key=)` signs producer envelopes, and `spawn_agent(signing_key=)` signs replies through `PyAgentCtx`. `verifier=` selects receive-side keys. The `agent_ctx` test helper also accepts `signing_key=`.
`laser.sessions(*, stream=None, topics=None, memory_namespace=None, context_turns=None, context_tokens=None)` returns `Sessions`. Its `create(id)`, `start()`, and `open(conversation_id)` methods return a `Session`. `append(kind, data)` records a turn. `context(*, last_n, token_budget)` returns `list[SessionTurn]`, with `kind`, `text()`, and `message`, a `ContextMessage` (`id`, `provenance`, `payload`, `envelope`, `topic`).
`memory(memory=None)` returns `ScopedMemory`, and `memory_in(namespace)` opens another namespace. `ScopedMemory` provides `recall(folded=, agent=, user=, application=)`, `search(query, folded=)`, `remember(kind=, durable=, dedup=, ...)`, `block(token_budget=None)`, `consolidate(max_items)` returning `ConsolidationReport`, `forget(id)`, `improve(target, weight, *, note)`, and a `handle` property. `checkpoint()` returns a `Checkpoint` with `to_json` and `from_json` methods. `turns_at` and `turns_since` read around its saved offsets. `state_at(checkpoint, init, fold)` and `replay(checkpoint, init, fold)` rebuild state. `topics` maps each turn kind to a distinct topic name. Two kinds that share a topic raise `InvalidError`.
`laser.topic(name).replay()` returns `Cursor`. `Laser.assemble_context` reads a conversation with topic, role, last-count, and token-budget controls. `laser.context(conv)` returns `ContextScope`. Its `fetch` and `block` default to `agent.commands` and `agent.responses` and map `n` and `token_budget` to Rust `Chain(LastN, TokenBudget)`. Context reads, folds, policies, and estimators receive `ContextMessage`, not `AgentMessage`. `fetch_with(topics, policy)` takes the `LastN`, `TokenBudget(max_tokens, estimator=None)`, `RoleFilter(agents)`, and `Chain` classes. `state(topics, init, fold, ...)`, `state_with(store, topics, init, fold)`, and `checkpoint` fold the conversation log like Rust, as do the free `context_checkpoint(laser, topics)` and `ConversationState.load`/`load_with`. `memory()` accepts a namespace or a handle, `memory_with(namespace, backend, *, embedder)` picks a backend, and the module constant `CONTEXT_READ_WINDOW` is the 10,000-record read window. RBAC includes `bind_roles(principal, roles, expect_revision=)`, `get_bindings(principal)`, `Grant(feature, action, *, effect, resource=ResourcePattern.prefix(..))`, and `authz_history("all" | {"role": n} | {"binding": {"user_id": id}}, after_revision=, limit=)`.
`A2aBridge(laser, source, request_topic, reply_topic, *, capabilities=, signing_key=)` provides `submit(params_json)`, task, cancel, card, and `signed_card(key)` operations. `McpBridge(laser, source, tool_topic, reply_topic, server_name, *, memory_tools=False, timeout_secs=None)` provides initialization, `call_tool(name, params_json)`, tools, resources, and prompts through configured `McpPrompt` values. AG-UI methods include `publish_state_snapshot`, `publish_state_delta`, `reconstruct_state`, and `agui_events`. `AgentCtx.respond_input` answers an interrupt. Python applications host HTTP endpoints through their own framework rather than the Rust axum `router()`.
`Laser.memory`, `Laser.memory_on_topic(name)`, `Laser.memory_topic(name, partitions=, ttl_secs=)`, and `Laser.memory_with(namespace, backend="auto"|"log"|"vector", embedder=None)` return a `MemoryHandle`. `MemoryHandle.remember` takes `agent`, `conversation`, `stream`, `user`, `application`, `kind="fact"`, `durable=False`, and `dedup=False`, and `context`, `consolidate`, `forget(id)`, `append(id, ...)`, and `improve(target, weight)` mirror the Rust handle. `recall(..., block=True)` returns the rendered block. Topic expiry defaults to 30 days. The subclasses `LogMemory(laser)` (`in_namespace`, `on_topic`, `on_stream_topic`, `*_named` state ops), `VectorMemory(embedder)`, and `RerankedMemory(inner, reranker)` build each backend. `VectorMemory(embedder)` is standalone, while `VectorMemory.governed(laser, embedder)` inherits the Laser governor. `to_context_block`, `memory_id_content`, `memory_kind_class`, and `memory_kind_code` are module functions. `MemoryItem.text()` is a method and `MemoryItem.signals` returns `RecallSignal` objects. Embedders can return vectors directly or return awaitables.
Default `recall` reads the managed KV view. `recall(folded=True)` and `fetch_folded` read the log locally. `MemoryItem.source` identifies the original record. The durable handle also provides `set`, `fetch`, `update`, and `remove` by name. Vector memory returns `UnsupportedError` for those named-state operations. `InMemoryStore` and `FileStore` implement `StateStore`.
`Laser.graph(name)` returns `GraphHandle` with `upsert`, `link`, `unlink`, `relink`, `neighbors(node, dir, edge_type, depth)`, and a fluent traversal: `start_ids`, `start_match(Filter)`, or `start_nearest(embedding, k)`, then `out`, `incoming`, `both`, `return_edges`/`return_triplets`/`return_paths`, `limit`, `as_of`, `conversation`, and `fetch()`. Nodes, edges, results, and `SourceRef` cross as serde dicts, with attributes as `[name, value]` pairs. `node_id_content`, `edge_id_content`, `graph_node_entity`, and `graph_edge_relate` construct content IDs and records. `register_graph` defines a graph projection, as demonstrated in `examples/python/memory.py`.
Graph traversals, queries, and KV scans scope with `.conversation(id)`. Build borrowed Rust handles for each Python call and retain owned state between calls. The local vector backend shares its items through `Arc`.
- `src/bin/stub_gen.rs` - generates `laser_sdk.pyi` and appends the exception hierarchy (which `create_exception!` does not expose to the stub gatherer).
- `tests/` - pytest. `test_smoke.py`, `test_parity.py`, `test_parity_helpers.py`, `test_coordination_client.py`, and `test_readme_snippets.py` need no server. `test_integration.py`, `test_filters.py`, `test_publish_failures.py`, `test_resume_parity.py`, and the callback suites (`test_trait_hooks.py`, `test_snapshot_callbacks.py`, `test_workflow_callbacks.py`) use the shared native Iggy fixture. `test_connect_timeout.py`, `test_rbac.py`, and `test_reconnect.py` cover connection behavior.
- `bdd/python/` - the Python BDD runner over the shared Gherkin in `bdd/scenarios/`. Streaming, capabilities, provenance, and agent run against Iggy. Query and key-value CAS use the transport-free reference engine.
## Data stack parity
`src/query.rs` binds Query, `QueryBuilder`, and the `Filter` predicate class (the Rust struct is `PyQueryFilter`, apart from the consumer-filter `FilterExpr`), with `filter`, `having`, `agg_as`, `max_rows`, `rows`, `rows_typed`, `deadline(seconds)`, and `execution_id`, explicit operational and lakehouse targets, snapshot selection, typed raw SQL parameters, cursor paging, execution status and cancellation, and rich result context. Query results expose `fields`, positional tagged `Row.values`, and `page` (a `Page` with `total`, `has_more`, `next_cursor`, `at_least`, `total_pages`). Never restore `row.headers`. Use `QueryResult.value`, `value_text`, `value_u64`, and `value_i64` for field-name access.
`src/destinations.rs` binds revision-guarded destination mutations, bounded checkpoint reads, and explicit query routes. `client.rs` exposes `Laser.destinations()` and `Laser.query_lakehouse()`.
`src/publish.rs` exposes `arrow_ipc` and `add_arrow_ipc`. Metadata crosses through serde and exact payload length is checked before I/O. Keep the self-contained Arrow stream policy in stubs and README examples.
After a data-contract change, regenerate `laser_sdk.pyi` and update stub tests, Python tests, shared `@data_stack` scenarios, and README examples. Use serde-derived dictionaries for large managed types rather than separate Python data definitions.
## How the binding works
- Every async Rust method becomes a Python awaitable via `pyo3_async_runtimes::tokio::future_into_py`. `Laser` is `Clone` (Arc inside), so each method clones it into the async block.
- Lifetime-bound builders are not exposed directly. The Rust builders borrow `&Laser`, so the Python builder classes hold owned, accumulated state and reconstruct + run the Rust builder inside the async block at the terminal (`send`/`fetch`). Fluent setters take `PyRefMut<'_, Self>` and mutate in place.
- Complex managed inputs/outputs ride serde. A Python dict deserializes straight into a wire type (`projection`, `binding`, `schema source`, dead-letter capsule) via `pythonize::depythonize`, and structured replies serialize back to dicts. This binds the whole registry/control surface without a class per type.
- `spawn_agent` calls Python `async def handle(ctx, message)` through the Rust consumer. Capture the event-loop locals and run within `scope(locals, ...)`. `into_future` then schedules the coroutine on the correct loop.
- Map `LaserError` to the hierarchy in `errors.rs` and retain classification attributes such as `code`, `retryable`, and `unsupported`. `TimeoutError` also inherits `builtins.TimeoutError`. `CancelledError` also inherits `asyncio.CancelledError`. Cache these generated exception types in `OnceLock` so both SDK and standard exception handlers recognize them.
- `Cursor`, `WatchReader`, and `TypedRecords` provide `__aiter__` and `__anext__`. They yield buffered records and raise `StopAsyncIteration` when caught up. A later iteration resumes from saved offsets. Keep `poll()` and typed `next()` available.
`Consumer` instead waits for new records until `shutdown()`. It supports automatic and explicit commits, and `next_within(wait_secs)` raises `TimeoutError` like Rust. `commit_interval_ms` defaults to 0, so the default policy is `Polling` like Rust. `laser.watch()` raises `UnsupportedError` when the change feed is unavailable, and `WatchReader.from_offsets` resumes a reader. `header_kinds` reports exact Iggy types. Use `(kind, value)`, such as `("uint16", 7)`, for an explicit numeric width.
- The live consumer avoids an intermediate payload copy. `PyConsumerMessage` retains Apache Iggy's reference-counted `Bytes` until its `payload` getter creates the one Python-owned `bytes` object required by the language boundary.
- `Laser` provides `__aenter__` and `__aexit__` for `async with`. The connection closes when the last shared handle is dropped. `with_default_stream`, `with_ops_stream`, `with_control_topic`, `with_dlq_topic`, and `with_changes_topic` share that connection. `Laser.capabilities()` and `refresh_capabilities()` use Rust discovery. `Capabilities.query`, `kv`, `filters`, and `destinations` are `QueryCaps`, `KvCaps`, `FilterCaps`, and `DestinationCaps` objects, never flat booleans. Preserve unavailable-backend handling and explicit handle overrides.
- `spawn_agent` exposes `max_partitions`, `shutdown_grace_ms`, `retry_max_attempts`, `retry_base_delay_ms`, `dead_letter`, and `middleware`. The dead-letter callback receives message, the full wire capsule dict, and a typed publication error or None through `PyDeadLetterSink`. Middleware `after_handle` receives the message, an {ok, error} result dict, and attempt. Both observers accept direct values or awaitables. `Laser.contract` returns a `Contract` (`Completed`, `Failed`, `NotConsumed`, `TimedOut`) and `Laser.scatter_report(...)` returns a `ScatterReport`. `Provenance` keeps `correlation_id` separately from `idempotency_key`. `AgentMessage` exposes `id` (a `MessageId`), `provenance`, `payload`, `envelope`, and `body()`, with no per-field shortcuts.
- `AgentCtx.fan_out(skill, payload, *, policy, quorum, deadline_ms, fixed_inbox, principal, route_policy)` returns a `Gather` (`ok`, `failures`, `replies()`). `approval_gate` waits for a decision through a temporary `laser_sdk::testing::agent_ctx`. The same approach supports the owned Python wrapper over borrowed Rust context. Module functions `agent_message` and `agent_ctx` let tests call handlers without a live consumer.
- Governance delegates to Rust. Preserve mandatory-voter rules and rejection of invalid configurations. `SwappableGovernor.swap` returns the previous policy, and `current` returns the active one. `Intent`, `Vote`, and `Decision` use typed topic encoding. Invalid construction, voting, or folding raises `IntentError`, an `InvalidError` subclass with a `kind`. `ActionDecision.verdict` is a `Verdict`, `with_policy` takes a `PolicyRef`, and `GovernedAction.counters` holds the `ActionCounters`.
Call `Decision.authorizes(intent)` before the effect. Names remain claims unless signatures or topic permissions establish authorship. `SwarmActivity` preserves unknown verdict names and ignores repeated evidence. `CrashContext` exposes `journal`, `dead_letter`, and `last_decision` and escapes control characters in summaries. Keep parity tests and stubs current.
- Contracts default to a 30-second `deadline_ms` like Rust and TypeScript, and take `agent=` (named-agent routing), `expire_if_not_consumed_ms`, `reply_on`, `conversation`, `fence`, and `registered`. `Laser.connect_env`, `local`, `connect_with_stream`, and `with_governor_retention` mirror Rust. `NoStreamError` and `NoRespondTopicError` subclass `ConfigError`. Key-value `send()` raises `InvalidError` when a precondition was set, `KvSetRequest.bytes(payload)` sets a raw body, `Kv.exists` returns `KvMetadata`, `copy_to`/`move_to` return an awaitable `KvCopyRequest`, and `cas_fenced` returns an awaitable `KvCasFencedRequest` builder. Query lifecycle `execute_query`, `query_page`, `cancel_query`, `query_status`, `execute_checkpoint`, `reassemble_channel`, and `consumed` are on `Laser`.
- `Laser.execute_batch` accepts Rust `BatchItem` dictionaries and returns exact reply bytes. `Agdx.status`, `Agdx.fail`, and `AgdxStream.fail` retain their typed data. Convert floating-point durations through `Duration::try_from_secs_f64`. Negative, non-finite, and out-of-range values must raise `InvalidError` rather than panic.
- `Kv.lease` and `Kv.renew_lease` return `Lease` with token, granted lifetime, and position. Python inherits the dedicated Rust coordination connection, validation, and timeout reset. An uncertain acquisition waits through its requested lifetime before raising.
Pass `MutationPosition { topic_generation, partition, offset }` to `Kv.get_entry_at_least` for takeover reads. An unmet position returns stale rather than absent data. `LaserError.ambiguous_mutation` requires operation-specific recovery and must not trigger generic handler retries.
## Versioning and naming
- The Python package is `laser-sdk` on PyPI, imported as `laser_sdk`. The internal Rust crate is `laser-sdk-python` (`publish = false`) with cdylib lib `laser_sdk_py`, named to avoid clashing with the `laser_sdk` dependency crate. Maturin renames the built module to `laser_sdk` via `module-name`.
- Python follows the shared workspace version, currently `0.6.0`. Its dependency must select the matching Rust `laser-sdk` crate.
## Working on it
- Build into a venv: `maturin develop` (the venv lives at `foreign/python/.venv` in local dev).
- Regenerate stubs after any surface change: `cargo run --bin stub_gen`, then check `laser_sdk.pyi` is current.
- `cargo check` for fast type-checking. The crate is outside the workspace, so the workspace clippy/test gates do not cover it: run them here too.
- Lint and format with ruff (configuration in the repo-root `ruff.toml`): `ruff check` and `ruff format --check` over `foreign/python`, `bdd/python`, and `examples/python`. The generated `.pyi` is excluded.
- Tests: `pytest -q` against the versioned Iggy server, plus the BDD suite in `bdd/python`. `LASER_TEST_IGGY_SERVER` selects a local Iggy binary for development.
- After a Rust API or wire change, update Python bindings and regenerate stubs. Update the corresponding tests and documentation in the same authorized change.
- After any public surface change, regenerate `docs/parity.md` with `python3 scripts/check-parity.py --write`, then run `just parity-check`.
## Publish recovery
`Laser.connect` takes `connect_timeout_ms`, `publish_timeout_ms`, `publish_max_retries`, and `publish_retry_backoff_ms`, which override the `LASER_CONNECT_TIMEOUT_MS` and `LASER_PUBLISH_*` variables. `Topic.producer(retries=None, retry_interval_ms=None)` inherits the connection's retry count and delay, and `retries=0` disables resends. A publish that gives up raises `PublishFailedError` with `stream`, `topic`, `committed`, `unconfirmed`, and the cause as `__cause__`. Each Rust error variant maps to its own exception class, and a handler rejection raises `RejectedError`. Defaults and outage handling are in [publish recovery](../../../docs/publish-recovery.md) and [connect timeout and cleanup](../../../docs/connect-timeout.md).
## Consumer-group filters
Group handles, filter setup, scan budgets, progress, and cross-language test parity follow the [consumer-filters](../consumer-filters/SKILL.md) skill.
## 0.6.0 parity additions
`Laser.memory_custom(backend)` wraps a Python object with `remember`, `recall`, `improve`, `forget` (sync or async, dict scopes and queries). `MemoryHandle.reranker(callable)` returns a reranked handle, and the callable runs only when a recall carries `semantic` text. `Topic.cbor(cls)` and `await Topic.schema(schema_id, cls)` are typed topics. Avro and JSON Schema publication encodes values. Protobuf publication requires encoded bytes through the plain raw builder with a schema ID. `PublishRequest.claim_check(store, threshold_bytes)` uses the blob store hooks. `Topic.producer(background=True)` or `background=BackgroundConfig(...)` switches to Iggy's buffered mode and `Producer.shutdown()` waits for sends in flight, flushes, and closes. Hooks take a plain callable or an object with the hook method, sync or async. Async hooks run on the calling event loop with its context variables. Sync hooks run on an SDK worker thread and must not block. `ForkHandle.create(continuous=True)` refuses `severed` at the same time. `Workflow.run_id(id)` resumes a run. `Capabilities.is_open_only()`, `serves_consistency(level)`, and `with_capabilities(query_execution=, versions=OpVersions(..), backends=[..])`. `AgdxStream` gains `channel`, `with_deadline_micros`, `with_target`, `content_type`, `buffered`, `flush`. `Intent.validate()`. `SigningKey.sign` and `sign_with_context` return a signature dict, and `KeyRegistry` gains `enroll_record`, `verify`, `verify_at`, `verify_observed_at`. `Session.context_with(policy)` and `Session.graph(name)`.
## Coordination and callbacks
`src/coordination.rs` binds `FencedLeaseClient`, `PreparedMutation`, and `DedicatedKvTransport`, and `src/memory_handler.rs` binds `MemoryHandler`. Exceptions raised through a callback keep their SDK class. Cancelling an SDK task retires the asynchronous callback it is running. Context policies stay synchronous. The upgrade steps from 0.5, including the new middleware and dead-letter callback arguments, are in [client behavior](../../../docs/client-behavior.md).
Query `rows()` and `rows_typed()` return lazy asynchronous iterators under `max_rows`, consumed with `async for`. The managed key registry takes `enroll(principal, verifying)` for an agent key and `enroll_record(record)` for lifecycle metadata. Both bridges expose asynchronous `handle_rpc(request)` returning response dicts.