Metadata-Version: 2.5
Name: statewire
Version: 0.8.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.8,>=0.7.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 envelope streams:

- `GET /stream` (SSE) and `/ws` (WebSocket twin) — every message is an envelope
  `{"ops"?, "res"?, "ack"?, "syn"?, "fin"?}`. The first envelope on attach is a
  full snapshot (`ops: [{"op": "replace", "path": [], "value": <state>}]`) with
  the `syn` handshake `{"lastSeq", "lease"?}`: the client's last admitted
  command seq (`-1` unknown client) and the attach's writer-lease token.
  Sessions are single-writer per client id: attaching takes the lease, and any
  other live stream of that client id ends with `fin {"reason": "superseded"}`.
- `POST /commands` — one command at a time, a method call
  `{"method": <name>, "params": [...]}` routed to the `@command` handler
  registered under that name. Commands carry `Statewire-Client-Id`, a
  monotonic `Statewire-Command-Seq`, and the current `Statewire-Lease`. The HTTP
  response is receipt only (`200 {}` | `400` malformed | `409` seq gap | `412`
  unknown client | `423` stale lease); the protocol statuses carry a
  discriminator body (`{"error": "seq-gap" | "unknown-client" |
  "stale-lease", "message"?}`) so clients can tell statewire's verdict from a
  middleware-minted bare status; verdicts ride the stream. Over WS the same
  commands arrive as `{"method": <name>, "params": [...], "seq": <int>}` frames —
  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. `ack`
is the cumulative watermark for the issuing client, emitted after the command's
effect ops. `res` carries responses (`{"seq", "type": "accepted" | "rejected" |
"pending" | "crashed" | "result-unavailable", "message"?, "payload"?}`): the envelope
that first covers a seq states the command's fate — a terminal response, a
`pending` response (result follows in a later envelope), or nothing = void. `fin`
(`{"reason": "evicted" | "error" | "gone" | "superseded", "message"?}`) is always
the last envelope.

Handler outcomes: the return value becomes the terminal `accepted` payload;
`raise StatewireReject` becomes `rejected`; returning a `CommandExecution` keeps the
command open past the handler's return — its `ack()` flushes the covering ack
with a `pending` response, its `resolve`/`reject` produce the late terminal
response (post-ack failures are terminal `rejected`, not stream faults).
Duplicate seqs never re-run: they are answered
from a ~30s result cache / in-flight registry — in the POST body over HTTP
(`200 {"res": ...}`), as a `res` on the stream over WS.

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 => ack
```

`@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` response on the stream,
not an HTTP error.

`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 `/statewire`, `/stream`, `/commands`, 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.

`GET /statewire` is the meta endpoint: `{"protocol": 1, "commands": [<names>]}`.

## Layering

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

## Develop

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