Skip to main content

Open data-platform SDK over Apache Iggy: typed streaming, declared projections and a query DSL, a key-value store, copy-on-write forks, and an optional agent runtime with the Agent Data Exchange Protocol (AGDX).

Project description

laser-sdk (Python)

The LaserData SDK for Python: an open data-platform SDK over Apache Iggy. Native bindings to the Rust SDK via PyO3, so the wire contract, codecs, and runtime are the same ones the Rust client uses.

Rust and Python are one v1 contract. Every public primitive, builder option, validation rule, error classification, capability, and transport limitation ships in both SDKs with matched examples and shared BDD coverage where the behavior is language-neutral.

spawn_agent(agent_id, ..., consumer_group=None) keeps logical identity separate from Iggy replica topology. The group defaults to the agent id spelling, set it explicitly when deployment grouping differs.

One Apache Iggy connection gives you typed streaming, declared projections and a query DSL, a key-value store, copy-on-write forks of the read model, and an optional agent runtime with the Agent Data Exchange Protocol (AGDX): publish, request/reply, and a consumer that drives your async def handler with at-least-once delivery, dedup, retry, and a dead-letter queue.

Apache Iggy is the underlying streaming core. Projections, the query layer, the key-value store, and forks are served by LaserData Cloud over that same connection. Against raw Apache Iggy those calls raise UnsupportedError.

Install

pip install laser-sdk

Wheels ship for Linux (x86_64, aarch64) and macOS (Intel, Apple Silicon), Python 3.10 through 3.13.

Connect

import asyncio
from laser_sdk import Laser

async def main():
    laser = await Laser.connect("iggy://iggy:iggy@127.0.0.1:8090", stream="agents")
    caps = await laser.capabilities()
    print(caps)

asyncio.run(main())

The connection string scheme is optional (iggy:// is assumed). Pin a default stream= so laser.topic(name) is the one-word shortcut against it, or address any topic on any stream explicitly with laser.stream(name).topic(name). The accessors are free and synchronous, IO happens at the verbs (publish, replay, ensure), mirroring the Rust grammar one-to-one.

Publish and consume

await laser.topic("orders").ensure(partitions=4)

await (
    laser.topic("orders").publish()
    .index("customer_id", "alice")
    .index("total", "129")
    .inline_payload()
    .json({"id": "o-1", "customer": "alice", "amount": 129})
    .send()
)

Batch and any payload

A single publish is the simplest call, not the common one. publish_batch accumulates records and sends them in one network round-trip, the largest throughput lever the SDK offers, and reads mirror it: a topic(..).replay() cursor drains every record that arrived since the last poll in one call. Batching on both sides is what makes the path efficient.

The payload is yours, in any format. add_json / add_msgpack (and extend_json for a whole list) are conveniences over add_payload, which takes raw bytes the SDK never inspects, so a compressed blob or your own framing rides unchanged. Schema-first Avro and Protobuf bodies are below.

batch = laser.topic("orders").publish_batch().inline_payload()
batch.extend_json([{"id": "o-1", "amount": 129}, {"id": "o-2", "amount": 80}])
batch.add_payload(b"\x00any-bytes-any-format")  # raw bytes, untouched by the SDK
await batch.send()                              # the whole batch, one round-trip

Typed topics

One handle binds a topic to a class: pass cls= (a dataclass or pydantic model) and the topic encodes on the way in and decodes with the log position attached on the way out. publish(order) encodes the instance as JSON in one call, records(reader_name) is the typed reader over the same caller-owned offsets as replay(): next() yields the next record decoded into the class (None when caught up), and a record that does not decode raises TypedDecodeError naming its exact log position, then the reader moves past it.

from dataclasses import dataclass

@dataclass
class Order:
    customer: str
    amount: int

orders = laser.topic("orders", cls=Order)
await orders.publish(Order(customer="alice", amount=129)).send()

records = orders.records("billing")
while (record := await records.next()) is not None:
    order: Order = record.value            # an Order instance, record.position names the log slot

Schema-first bodies (Avro / Protobuf)

Compile a registered writer schema once, then publish raw datums under it. The body is encoded client-side, so a value that stops matching the schema fails before publishing rather than as a managed-side warning you cannot see. The managed plane resolves the schema by id and extracts indexed columns from the binary body.

from laser_sdk import CompiledSchema

source = {"kind": "avro", "schema": fill_avro_schema}
schema_id = await laser.register_schema(source, name="fill")
compiled = CompiledSchema.compile(source, id=schema_id)

batch = laser.topic("trades_avro").publish_batch().inline_payload()
for fill in fills:
    batch = batch.add_avro(compiled, schema_id, fill)
await batch.send()

CompiledSchema also offers validate / validate_value / decode, and the single-record builder has .avro(compiled, schema_id, value). For Protobuf or your own framing, encode the body yourself and ship it with .raw_bytes(bytes, "protobuf") (or batch .add_raw_bytes(..)). Writer schemas live on LaserData Cloud, so registration is a managed feature.

Query (managed)

rows = await (
    laser.query("orders")
    .where_eq("customer_id", "alice")
    .filter_gte("total", 100)
    .order_desc("total")
    .limit(10)
    .with_payload()
    .fetch_all()
)
for row in rows:
    print(row.headers, row.json())

Query, the key-value store, and forks are managed features served by LaserData Cloud. Against raw Apache Iggy they raise UnsupportedError.

Key-value

kv = laser.kv("sessions")
await kv.set("user:42").json({"state": "online"}).ttl(300).send()
state = await kv.get_typed("user:42")
values = await kv.get_many(["user:42", "user:43"])  # one round trip (the mixed-operation batch)
await kv.copy_to("user:42", "user:42:2026", to_namespace="archive")  # one backend transaction
await kv.move_to("plan:draft", "plan:current")  # copy plus source delete
await kv.delete("user:42")

Agents

from laser_sdk import Laser

async def handle(ctx, message):
    text = message.payload.decode()
    await ctx.respond(f"echo: {text}".encode())

laser = await Laser.connect("iggy://iggy:iggy@127.0.0.1:8090", stream="agents")
await laser.bootstrap(partitions=4)

handle_agent = laser.spawn_agent(
    "echo", "agent.commands", handle, respond_on="agent.responses"
)
await handle_agent.ready()

from laser_sdk import Provenance
reply = await laser.request(
    "agent.commands", "agent.responses", b"hello",
    Provenance(agent="caller"), timeout_secs=10,
)
print(reply.payload.decode())

await handle_agent.shutdown()

Signed, principal-bound contracts

Rust and Python use the same Ed25519 verifier and routing rules. Enroll keys before connecting, give an agent its signing key, and constrain sensitive capability routes to the server-authenticated principal. One connection may advertise one agent. Attempting to advertise another raises a typed conflict instead of replacing the first presence.

from laser_sdk import KeyRegistry, Laser, SigningKey

risk_key = SigningKey(bytes(range(32)))
keys = KeyRegistry()
keys.enroll("42", risk_key)

laser = await Laser.connect(connection, stream="agents", verifier=keys)
risk = laser.spawn_agent(
    "risk",
    "risk.commands",
    handle,
    capabilities=["screen-order"],
    signing_key=risk_key,
    verifier=keys,
)
await risk.ready()

result = await laser.contract_report(
    "screen-order",
    b'{"order":"o-1"}',
    source="orders",
    principal=42,
)
assert result["state"] == "completed"
assert result["verified_principal"] == "42"

contract and scatter remain body-only conveniences. Use contract_report or scatter_report when policy or UI code must inspect verified_principal. With a verifier configured, unsigned, invalid, and wrong-principal replies are ignored rather than returned with an empty identity.

For a human-in-the-loop pause, the typed AGDX producer's request_input publishes a prompt and blocks on the human's correlated reply, which a handler resolves with AgentCtx.respond_input:

decision = await laser.agdx("agent.human_input", "orchestrator", conversation_id).request_input(
    "agent.responses", b"approve a $500 refund?", timeout_secs=15
)

Govern what an agent does before the effect runs: a policy object decides per action (allow, observe, block, step_up, modify, defer), enforce or shadow mode, and every non-allow decision lands as a digest-chained evidence event on the audit topic. PolicyBlockedError / StepUpRequiredError / PolicyDeferredError are the typed refusals:

from laser_sdk import ActionDecision, PolicyBlockedError

class NoWires:
    async def decide(self, action):
        if action.payload.startswith(b"wire-funds"):
            return ActionDecision.block("wires need approval").with_policy("finance", "3", ["no-wires"])
        return ActionDecision.allow()

governed = laser.with_governor(NoWires(), mode="enforce")
try:
    await governed.send_agent("agent.commands", b"wire-funds", provenance)
except PolicyBlockedError as refused:
    print(refused)  # policy blocked: no wire transfers

# Per-agent: everything the handler publishes is governed too.
handle_agent = laser.spawn_agent("clerk", "agent.commands", handle, governor=NoWires())

QuorumGovernor composes several named voters under a policy (all, any, or at_least(n)) into one governor, so a deterministic safety voter and an LLM voter combine into a single decision instead of picking one. Every mandatory voter must return allow, observe, or modify. A denial or error cannot be bypassed by a permissive any policy:

from laser_sdk import QuorumGovernor, QuorumPolicy

quorum = QuorumGovernor(QuorumPolicy.at_least(2))
quorum.voter("safety", NoWires(), mandatory=True)
quorum.voter("llm_reviewer", llm_voter, mandatory=False)

governed = laser.with_governor(quorum, mode="enforce")

SwappableGovernor hot-swaps the active policy at runtime, driven by anything (an operator call, a config reload, a folded policy-update topic), without dropping enrolled clones or reconnecting. A swap only changes the next decision, never one already recorded:

from laser_sdk import SwappableGovernor

swappable = SwappableGovernor(NoWires())
governed = laser.with_governor(swappable, mode="enforce")
...
previous = swappable.swap(a_stricter_policy)  # returns the replaced policy

Durable approvals are native typed records. They publish and replay directly, while the SDK keeps log ownership explicit:

import time
from laser_sdk import Decision, Intent, IntentPolicy, Vote, decide

intent = Intent(
    conversation=conversation_id,
    proposer="planner",
    body=b"reserve inventory",
    eligible_voters=["safety"],
    policy=IntentPolicy.all(),
    policy_version=7,
    deadline_micros=time.time_ns() // 1_000 + 30_000_000,
)
await laser.topic("intents", cls=Intent).publish(intent).send()
vote = Vote.cast(intent, "safety", "allow")
decision = decide(intent, [vote], time.time_ns() // 1_000)
if decision and decision.authorizes(intent):
    await laser.topic("decisions", cls=Decision).publish(decision).send()

Construction, casting, and folding fail with InvalidError on malformed state. Mandatory voters must affirm, and ballots outside the intent's time window never count.

SwarmActivity is a supervisor's read model over governance evidence: fold PolicyEvidence records already read off the audit topic and ask "what has this agent been doing" without hand-rolled bookkeeping:

from laser_sdk import PolicyEvidence, SwarmActivity, Topics

swarm = SwarmActivity()
for message in await laser.assemble_context(conversation_id, topics=[Topics.AUDIT]):
    envelope = message.envelope
    if envelope and envelope.get("operation") == "policy_decision":
        swarm.observe(PolicyEvidence.decode(bytes(message.agdx_body)))

activity = swarm.agent("planner")
if activity:
    print(activity.decisions, activity.count("block"))

CrashContext is a recovery tool's one-call bundle: combine an already-read journal tail, the crashed message's dead-letter capsule (if any), and the conversation's most recent decision (if any) into one deterministic digest, never invoking a model itself:

from laser_sdk import CrashContext

journal = await laser.assemble_context(conversation_id, topics=[Topics.COMMANDS])
context = CrashContext(journal=journal, dead_letter=None, last_decision=activity.last_decision)
print(context.summarize())

Runs (managed)

The managed run registry answers "what happened to that task" without folding topics yourself. Gated on the agent_workflow capability, UnsupportedError elsewhere.

runs = laser.runs()
run = await runs.submit("diagnoser", b'{"incident": "INC-7"}')
info = await runs.status(run.run_id)
page = await runs.list(state="running", limit=25)
await runs.cancel(run.run_id)  # records the intent, the engine observes it

wf = laser.workflow("incident-response")
wf.registered()  # the run's lifecycle lands in the registry

# A fenced external effect must use the same namespace in the workflow lease
# and in the handler's kv("payments").cas_fenced(...) commit.
wf.step(
    "charge",
    to="charger",
    build=lambda outputs: b'{"order":"o-1"}',
    fence_namespace="payments",
    on_timeout="reassign",
)

Change feed (managed)

Await a view's advance instead of polling it blind. A projection binding built with notify makes the plane publish one change record per committed batch, and laser.watch() reads that feed. Gated on the watch capability, UnsupportedError elsewhere.

feed = laser.watch(index="orders_v1")
for change in await feed.poll():
    print(change.index, change.from_offset, change.to_offset, change.rows)
    rows = await laser.query("orders_v1").fetch_all()  # the record is a wakeup, the rows come from query
saved = feed.offsets  # persist to resume after a restart

Consume and replay

# A resumable reader over a topic. Each poll drains what is new. Persist the
# offsets to resume after a restart.
cursor = laser.topic("orders").replay()
for message in await cursor.poll():
    print(message.json())
saved = cursor.offsets

# Replay a conversation's history off the log (agent runtime).
history = await laser.assemble_context(conversation_id, last_n=50)

Memory and state

Agent memory shares one remember / recall / forget surface over four backends, and the kv-backed handle adds the named-item altitude: set(key, value) / fetch(key) / update(key, patch) / remove(key) for working notes addressed by name (UnsupportedError on the other backends). The log-backed default works on raw Apache Iggy. The in-process vector backend ranks recall by semantic similarity, embedding through your own async def embed(text) -> list[float]. The query and key-value backends are managed.

async def embed(text: str) -> list[float]:
    ...  # your model, or a deterministic stand-in

memory = laser.vector_memory(embed)
await memory.remember("checkout latency traces to the database pool", conversation=cid)
hits = await memory.recall(conversation=cid, semantic="why is checkout slow", limit=3)
print([item.text for item in hits])

# A vector memory created from a governed Laser applies the same pre-write policy.

# A durable key/value seam for agent state, the same vocabulary as the managed store.
store = ls.InMemoryStore()        # or ls.FileStore("/var/lib/agent")
await store.set("cursor", saved_bytes)
value = await store.get("cursor")

vector_memory inherits the governor enrolled on the Laser that creates it. A blocked write never mutates the local index, and a modified decision replaces the proposed memory body before embedding. Rust and Python therefore apply the same effect-boundary policy to local semantic memory.

Edge interop (A2A / MCP / AG-UI)

Reach an agent as an A2A task source or an MCP tool server, and render a conversation as AG-UI events, all over the durable log:

# A2A: submit a task, poll for the result.
a2a = laser.a2a_bridge("a2a-gateway", "agent.commands", "agent.responses")
task = await a2a.submit({"message": {"role": "user", "parts": [{"kind": "text", "text": "hi"}]}})
status = await a2a.task(task["id"])

# MCP: advertise tools, route tools/call to the agent.
mcp = laser.mcp_bridge(
    "mcp-gateway", "agent.tool_calls", "agent.tool_results", "laser-mcp",
    tools=[{"name": "ask", "input_schema": {"type": "object"}}],
)
tools = mcp.list_tools()
result = await mcp.call_tool("ask", {"q": "what is AGDX?"})

# An agent answers a bridge request from its handler:
async def handle(ctx, message):
    await ctx.respond_input("agent.responses", b"the answer")

# AG-UI: snapshot + deltas reconstruct shared state off the log.
await laser.publish_state_snapshot("agent.llm_io", "ui", conversation_id, {"count": 1})
state = await laser.reconstruct_state(conversation_id, "agent.llm_io")
events = await laser.agui_events(conversation_id, "agent.llm_io")

Host the actual HTTP endpoint with your Python web framework over these adapter methods.

Errors

Every failure raises a subclass of LaserError: QueryError, KvError, ForkError, GraphError, AuthzError, SignatureError, UnsupportedError, InvalidError, CodecError, TypedDecodeError, ProtocolError, TimeoutError, ConfigError, TransportError, BudgetExceededError, PolicyBlockedError, StepUpRequiredError, PolicyDeferredError, CancelledError. Each instance carries code, retryable, unsupported, not_found, version_skew, version_conflict, stale, permission_denied, stream_or_topic_not_found, no_capable_agent, lease_lost, fence_violation, budget_exceeded, and quarantined attributes so you can branch without matching on the type. TimeoutError also subclasses the builtin TimeoutError and CancelledError also subclasses asyncio.CancelledError, so stdlib-style except TimeoutError / except asyncio.CancelledError catch them too.

Reading

The readers are async-iterable: async for message in laser.topic("events").replay(), async for record in reader on a WatchReader, and async for record in topic.records(reader_name) on a typed reader all drain what is currently appended and stop when caught up. A fresh async for later resumes from the same offsets. poll() is still there for one batch at a time.

Lifecycle

Laser supports async with: async with await Laser.connect(conn) as laser:. The connection is reference-counted and closes when the last handle drops, and with_stream / with_ops_stream return aliasing clones that share it.

License

Apache-2.0.

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

laser_sdk-0.0.1rc13.tar.gz (662.6 kB view details)

Uploaded Source

Built Distributions

If you're not sure about the file name format, learn more about wheel file names.

laser_sdk-0.0.1rc13-cp310-abi3-win_amd64.whl (11.7 MB view details)

Uploaded CPython 3.10+Windows x86-64

laser_sdk-0.0.1rc13-cp310-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl (13.0 MB view details)

Uploaded CPython 3.10+manylinux: glibc 2.17+ x86-64

laser_sdk-0.0.1rc13-cp310-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl (13.0 MB view details)

Uploaded CPython 3.10+manylinux: glibc 2.17+ ARM64

laser_sdk-0.0.1rc13-cp310-abi3-macosx_10_12_x86_64.macosx_11_0_arm64.macosx_10_12_universal2.whl (24.3 MB view details)

Uploaded CPython 3.10+macOS 10.12+ universal2 (ARM64, x86-64)macOS 10.12+ x86-64macOS 11.0+ ARM64

File details

Details for the file laser_sdk-0.0.1rc13.tar.gz.

File metadata

  • Download URL: laser_sdk-0.0.1rc13.tar.gz
  • Upload date:
  • Size: 662.6 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/6.1.0 CPython/3.13.12

File hashes

Hashes for laser_sdk-0.0.1rc13.tar.gz
Algorithm Hash digest
SHA256 28ad317ea5d4c4a044c9bcb6c25946ae15bf006c7eab3817363870e69f8a165e
MD5 706da148205def0f0bc308a79bfe7d49
BLAKE2b-256 1e0bbe5bee926ad5704e565b99444ce48b3fdc9f0d31773141855721959edaf0

See more details on using hashes here.

Provenance

The following attestation bundles were made for laser_sdk-0.0.1rc13.tar.gz:

Publisher: ci-python.yml on laserdata/laser-sdk

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file laser_sdk-0.0.1rc13-cp310-abi3-win_amd64.whl.

File metadata

File hashes

Hashes for laser_sdk-0.0.1rc13-cp310-abi3-win_amd64.whl
Algorithm Hash digest
SHA256 7a75594da45763a2c6b6b18b34b749103b4087791542b5ff092df357f9036cf5
MD5 8467b6a0e8fe5a95d0edcf4e25def45e
BLAKE2b-256 9e5a62afeb8dee8f23ea2b3332bf4ec9039663fa218858003d51b98de2f19398

See more details on using hashes here.

Provenance

The following attestation bundles were made for laser_sdk-0.0.1rc13-cp310-abi3-win_amd64.whl:

Publisher: ci-python.yml on laserdata/laser-sdk

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file laser_sdk-0.0.1rc13-cp310-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl.

File metadata

File hashes

Hashes for laser_sdk-0.0.1rc13-cp310-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl
Algorithm Hash digest
SHA256 06296ca03fe5cd1d0e6cab067cb2dab0e3a9837d9580ebb7ff5c23e32c58d82c
MD5 c4738031474ca2823b722df5b46e1615
BLAKE2b-256 f3e25b2ccafdcc24a0fc4495711497329a5fb052c7cc8fa82cd328f27a9271e4

See more details on using hashes here.

Provenance

The following attestation bundles were made for laser_sdk-0.0.1rc13-cp310-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl:

Publisher: ci-python.yml on laserdata/laser-sdk

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file laser_sdk-0.0.1rc13-cp310-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl.

File metadata

File hashes

Hashes for laser_sdk-0.0.1rc13-cp310-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl
Algorithm Hash digest
SHA256 7857d44b8ed2a012a4f0c1cbee8295648e82806872370cbe9f751007214cd4de
MD5 1c4424a9e3455bca347c02daf26e77ea
BLAKE2b-256 5e1dc52163b7a0c61098202e4aecf119bd8b50f5cef67c29f54974ba5cae1fbf

See more details on using hashes here.

Provenance

The following attestation bundles were made for laser_sdk-0.0.1rc13-cp310-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl:

Publisher: ci-python.yml on laserdata/laser-sdk

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file laser_sdk-0.0.1rc13-cp310-abi3-macosx_10_12_x86_64.macosx_11_0_arm64.macosx_10_12_universal2.whl.

File metadata

File hashes

Hashes for laser_sdk-0.0.1rc13-cp310-abi3-macosx_10_12_x86_64.macosx_11_0_arm64.macosx_10_12_universal2.whl
Algorithm Hash digest
SHA256 24a9b2aff43cec86418ba366f4f522a4355705717298c05ac48632da2de395c3
MD5 2eb8c86deabb7eb2b98055e7f6d886e5
BLAKE2b-256 e50d7b5d12484e786f0f1bd358eec91da8062bea7f9a819f7b8ac1b8a83bb828

See more details on using hashes here.

Provenance

The following attestation bundles were made for laser_sdk-0.0.1rc13-cp310-abi3-macosx_10_12_x86_64.macosx_11_0_arm64.macosx_10_12_universal2.whl:

Publisher: ci-python.yml on laserdata/laser-sdk

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page