Metadata-Version: 2.5
Name: statewire
Version: 0.12.4
Summary: Replicate one JSON object over an SSE op stream + a command endpoint
Project-URL: Repository, https://github.com/assistant-ui/harness-sdk
License-Expression: MIT
Requires-Python: <4.0,>=3.12
Requires-Dist: fastapi>=0.115
Requires-Dist: httpx>=0.27
Requires-Dist: pinned<0.15,>=0.14.0
Requires-Dist: statepatch<0.3,>=0.2.0
Provides-Extra: langgraph
Requires-Dist: langchain-core>=0.3; extra == 'langgraph'
Provides-Extra: orjson
Requires-Dist: orjson>=3.9; extra == 'orjson'
Description-Content-Type: text/markdown

# statewire

The Statewire protocol for Python, on top of [pinned](https://github.com/Yonom/pinned).

`Statewire` is a `PinnedAPI` (one live instance per id cluster-wide) that replicates
one JSON object — any State — over packet streams:

- `POST /stream` (SSE) — the body is the first client frame
  (`{"headers", "ctx"?, "cmd"?}`), the response the packet stream. Every
  packet is an envelope `{"ops"?, "cmd"?, "syn"?, "idle"?, "fin"?}`. The
  first packet is a full snapshot (`ops: [{"op": "replace", "path": [],
  "value": <state>}]`) with the hello `syn` `{"seq", "ctxSeq", "idle"?}`:
  the client's admitted watermark per lane. Each attach's lease rides the
  `Statewire-Lease` response header. `GET /stream` is the degenerate
  empty-frame attach and `/ws` the WebSocket twin (frames ride the
  socket). A client id may hold any number of concurrent attaches: each
  gets its own hello and lease, and cmd answers and `dur` receipts fan
  out to every live attach of the id.
- `POST /frames` — submits follow-up frames `{"headers"?, "ctx"?, "cmd"?}`
  under `Statewire-Client-Id` and a live attach's `Statewire-Lease`. Each cmd
  statement is a method call `{"method": <name>, "params": [...], "seq":
  <int>}` routed to the `@command` handler registered under that name; the
  ctx lane replaces the client's context object. The HTTP response is
  receipt only: `200 {}` admitted whole, `409` seq gap and `429` backpressure
  answer the `{"syn": ...}` refusal body (the 429 also carries `Retry-After`),
  `423` dead lease (its attach is gone — re-attach and resend), `400`
  malformed (the stream also fins `error`). Verdicts ride the stream.
  Over WS the same frames ride the socket — no lease header, holding
  the connection is the lease.

`ops` are Immer-style deltas with array paths (object keys as strings, list
indices as ints): `replace` sets a value, `add` splices — its final int segment
indexes into the parent, so a list parent gains an element and a string parent
gains text at that offset — `remove` deletes. `cmd` carries the issuing
client's verdicts as an array of `{"seq", "type": "result" | "applied" |
"rejected" | "failed" | "crashed" | "expired", "code"?, "message"?,
"payload"?}` entries plus `{"seq", "dur": true}` durability markers; effect
ops precede the verdict, and `applied` states the broadcast state now
contains the command's effects while the result is still in flight. `fin`
(`{"reason": "evicted" | "error" | "gone" | "idle",
"message"?}`) is always the last packet.

Handler outcomes: the return value becomes the terminal `result` payload;
`raise StatewireReject` becomes `rejected`; any other raise becomes `failed`;
an involuntary death (cancellation at shutdown or eviction) becomes `crashed`.
`ctx.applied()` flushes the `applied` entry riding the ops packet that carries
the handler's effects so far; the terminal entry follows on a later flush.
Duplicate seqs never re-run: a replayed statement re-emits its stored verdict
on the stream (`expired` once the outcome outlived retention).

The protocol is generic: it says nothing about messages, queues, or agents — it
only replicates whatever `self.state` dict you assign and dispatches whatever
commands you declare. Domain-specific layers (see the `harness-sdk` package) sit
on top.

```python
from statewire import Statewire, command


class Thread(Statewire):
    async def lifespan(self):
        self.state = {"messages": []}  # yielding without setting self.state throws
        yield

    @command
    async def addMessage(self, message_id: str, content: str):
        self.state["messages"].append({"id": message_id, "content": content})
        self.create_task(self.run())  # long work outside the inbox; returning here => applied
```

`@command` registers the handler under the method's own name; `@command("name")`
registers it under an explicit wire name. Params are positional. An unknown
method or a params/signature mismatch is a `rejected` verdict on the stream,
not an HTTP error; the seq is consumed.

`self.state` is a change-tracking proxy: mutate it plainly and the ops replicate
to every attached stream. `+=` on a string becomes an end-offset text-insert
`add` op of the suffix; other mutations become narrow `replace` / `add` /
`remove` ops. Mutations within
one synchronous segment coalesce into a single envelope.

## Extra routes

Need an endpoint beyond the protocol trio (a health check, a file upload)?
Decorate a method with pinned's `route` escape hatch:

```python
from pinned import route
from statewire import Statewire


class Thread(Statewire):
    @route.get("/health")
    async def health(self, request):
        return {"ok": True}
```

The reserved protocol paths `/stream`, `/frames`, and `/ws` are Statewire's
own; a subclass that decorates a `@route` onto any of them raises at
class-definition time rather than silently shadowing the protocol.

## Layering

```
pinned        one live instance per id, with an HTTP surface
  └─ statewire   the Statewire protocol: /stream + /ws (packets) + /frames
       └─ harness-sdk   HarnessState types, queueing, deepagents
```

## Develop

```sh
uv sync
uv run pytest
```
