Metadata-Version: 2.5
Name: stepledger
Version: 0.1.2
Summary: A Temporal plugin for LangGraph agents: one fenced Postgres row per node Activity execution, committed when the workflow accepts the result, plus a deduplicating External Storage driver for large state.
Project-URL: Homepage, https://github.com/Poojan6216/stepledger
Project-URL: Issues, https://github.com/Poojan6216/stepledger/issues
Author: Poojan Patel
License-Expression: Apache-2.0
License-File: LICENSE
Keywords: agents,external-storage,langgraph,ledger,postgres,temporal
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: Database
Classifier: Topic :: System :: Distributed Computing
Classifier: Typing :: Typed
Requires-Python: >=3.11
Requires-Dist: fastcdc>=1.7
Requires-Dist: langgraph<1.3,>=1.2
Requires-Dist: psycopg[binary,pool]>=3.2
Requires-Dist: pydantic-settings>=2.3
Requires-Dist: pydantic<3,>=2.7
Requires-Dist: pyyaml>=6
Requires-Dist: temporalio[langgraph]<1.34,>=1.33
Requires-Dist: typer>=0.12
Description-Content-Type: text/markdown

<!-- Generated by `uv run python -m bench.report` from bench/results/*.json. Edit bench/report.py, not this file. -->

# Stepledger

**Every LangGraph node running on Temporal, recorded once per Activity execution in your own Postgres, at any state size.**

Stepledger is one Temporal plugin that sits next to Temporal's `LangGraphPlugin`. It writes one fenced Postgres row per node Activity execution, commits it only once the workflow has accepted that result; `reconcile` and the test suite check every row against Temporal's own history. Its deduplicating External Storage driver keeps accumulating agent state under Temporal's payload and history limits. It changes no graph code and does not modify the LangGraph plugin. It was built in response to [temporalio/sdk-python#1894](https://github.com/temporalio/sdk-python/issues/1894).

**Status:** alpha, a research prototype (see [limitations](https://github.com/Poojan6216/stepledger/blob/main/docs/limitations.md)). Every number below was measured on one macOS laptop against the Temporal dev server and a local Postgres 16; the environment is recorded in each results file.

## The problem, measured

The LangGraph plugin sends each node's whole input state as its Activity input, and Temporal records every Activity input in history. For an agent whose state accumulates:

- **Persisting once at the end fails.** The final state is the first payload over the 2 MiB limit: the run gets stuck (SDK default) or the server terminates it (the #1894 error).
- **Writing each node's delta from its Activity moves the failure, it does not remove it.** At 60 KiB of new output per node the run still stops at node 35, because that node's own input crosses 2 MiB; history is already 35.42 MiB. With smaller outputs the 50 MiB history limit comes first.
- **Per-node writes are at-least-once.** Under injected crashes, a plain insert produced 99 duplicate and 85 divergent rows in 100 runs; an upsert still produced 25 divergent rows (a stale attempt overwriting the accepted answer). Retried nodes also repeated their external calls: 35 and 34 duplicate side effects reached the fake ticket and Slack targets in those runs.
- **External Storage fixes history, but storage then grows with the square of the run:** one object per payload, and every node input is a slightly longer copy of the last (history itself still grows, linearly, with references and sub-threshold payloads).

## Install

```bash
pip install stepledger
```

Requires Python 3.11 or later, `temporalio[langgraph]` 1.33.x and `langgraph` 1.2.x (the tested ranges: both upstream features are experimental in the SDK and Stepledger reads five private symbols, see the section on private APIs), and Postgres (tested on 16). For the local environment the tests and benches use (Postgres plus a Temporal dev server with explicit payload and history limits): `scripts/dev.sh up`.

## Integration

```python
from datetime import timedelta
import os
from temporalio.client import Client
from temporalio.contrib.langgraph import LangGraphPlugin
from temporalio.worker import Worker
from stepledger import StepledgerPlugin

lg = LangGraphPlugin(
    graphs={"investigate": build_graph()},
    default_activity_options={"start_to_close_timeout": timedelta(minutes=2)},
)
sl = StepledgerPlugin(dsn=os.environ["STEPLEDGER_DSN"], langgraph=lg)
client = await Client.connect("localhost:7233", plugins=[sl])
worker = Worker(client, task_queue="agents", workflows=[InvestigateWorkflow], plugins=[lg])
```

Then `stepledger init-db` once. Workers built from the client inherit the plugin; `StepledgerPlugin.from_config(langgraph=lg)` reads the same options from `stepledger.yaml`. With External Storage on, every client that reads results or histories needs the same data converter (`stepledger.cli.build_data_converter`); see [how it works](https://github.com/Poojan6216/stepledger/blob/main/docs/how-it-works.md).

## What beats it

| id | attack | measured | expected |
|---|---|---|---|
| 7.4 | Encrypted payloads | dedupe 11.0x plaintext vs 1.0x encrypted | dedupe is defeated (about 1x); the ledger stays correct; payloads counted as opaque |
| 7.6 | Effect whose tool ignores the key; crash before its record | direct: 3 duplicate tickets / 0 needed a person, once_blind: 0 duplicate tickets / 3 needed a person, once: 0 duplicate tickets / 0 needed a person | duplicates only without once(); blind once() stops for a person |
| 7.8 | Workflow reset (new run ID) | 6 duplicate effects across 3 resets | effects repeat unless the tool's own key dedupes |

- With the journal on, 17 billed model calls were still wasted: the worker died between the provider call and the journal write.
- With `on_ledger_error="warn"`, a database outage leaves rows missing until `stepledger reconcile` repairs them from history.
- **Two kinds of node never get a row:** `execute_in="workflow"` nodes and task-cache hits run no Activity. In Demo 4, 20 of 100 seeded runs were declared GAP for those reasons; `materialize()` names the position and never claims EXACT there.
- History still grows, linearly; unbounded runs still need continue-as-new.

See [docs/limitations.md](https://github.com/Poojan6216/stepledger/blob/main/docs/limitations.md) and [RESULTS.md](https://github.com/Poojan6216/stepledger/blob/main/RESULTS.md).

## What it does, measured

- **One row per node Activity execution, fenced and committed against history.** Across 100 crash-injected runs (the worker killed at fault points F1 to F5, and zombie attempts writing late at F6), Stepledger had 0 duplicate, 0 divergent, 0 lost and 0 orphan rows, each checked against the result Temporal recorded; with `once()` on the two effect nodes, 0 duplicate side effects reached the targets.
- **Past the wall.** With the dedup driver the 40-node run that stops B1 completes; the largest payload left in history is 62.4 KiB (nothing above the 64 KiB threshold) and history is 2.5 MiB.
- **Linear storage.** At 80 nodes x 100 KiB per node: 342.97 MB as one object per payload, 12.15 MB as dedup chunks (28.2x less).
- **The view equals the truth.** Over 100 seeded runs, `materialize()` rebuilt 80 runs EXACT and equal to the workflow's result, declared 20 gaps (workflow-side nodes, task-cache hits), and gave 0 wrong answers.
- **The retry bill.** With the LLM journal, simulated spend wasted on attempts Temporal did not accept (a fake model at a fixed token count, priced at claude-haiku-4-5 list rates) fell from USD 1.452 to USD 0.204 (86% less) under the same seeded crash plan.
- **Small overhead.** For the ledger write alone (storage driver off, 1 KiB nodes, 10 runs on a local dev server): p95 6.109 ms and about 1.4 ms of wall clock per node; about 90 bytes of headers per node in history, the same from 10 to 80 nodes (the commit header grows only with pending ids after failed carriers).

All numbers come from `bench/results/*.json` via the commands in [RESULTS.md](https://github.com/Poojan6216/stepledger/blob/main/RESULTS.md).

## The five demos

Each runs with one command against the local dev environment (`scripts/dev.sh up && uv run stepledger init-db`):

```bash
uv run python bench/demo.py --demo cliff        # reproduce #1894, and get past it
uv run python bench/demo.py --demo chaos        # pull the plug
uv run python bench/demo.py --demo growth       # the quiet quadratic
uv run python bench/demo.py --demo materialize  # the view equals the truth
uv run python bench/demo.py --demo cost         # the retry bill
```

### Demo 1: the cliff

| config | nodes | SDK default | payload check disabled | largest payload in history (bytes) | rejected payload (bytes) | history (MiB) | events |
|---|---|---|---|---|---|---|---|
| B0 bulk persist at end | 36 | STUCK at persist_all (PAYLOADS_TOO_LARGE) | TERMINATED at persist_all (BAD_SCHEDULE_ACTIVITY_ATTRIBUTES) | 2,067,992 | 2,129,481 | 39.37 | 218 |
| B1 naive per-node write | 40 | STUCK at node 35 (PAYLOADS_TOO_LARGE) | TERMINATED at node 35 (BAD_SCHEDULE_ACTIVITY_ATTRIBUTES) | 2,067,894 | 2,131,471 | 35.42 | 206 |
| B2 stepledger (no ext storage) | 40 | STUCK at node 35 (PAYLOADS_TOO_LARGE) | TERMINATED at node 35 (BAD_SCHEDULE_ACTIVITY_ATTRIBUTES) | 2,067,893 | 2,131,471 | 35.43 | 212 |
| B3 stepledger + whole-blob storage | 40 | COMPLETED | COMPLETED | 63,873 |  | 2.50 | 253 |
| B4 stepledger + dedup storage | 40 | COMPLETED | COMPLETED | 63,872 |  | 2.50 | 250 |

### Demo 2: pull the plug

| config | runs | faults injected | duplicate rows | divergent rows | lost rows | orphan rows | duplicate side effects |
|---|---|---|---|---|---|---|---|
| B1 naive per-node write | 100 | 100 | 99 | 85 | 0 | 0 | 35 |
| B1u naive upsert | 100 | 100 | 0 | 25 | 0 | 0 | 34 |
| Stepledger | 100 | 100 | 0 | 0 | 0 | 0 | 0 |

### Demo 3: the quiet quadratic

<picture>
  <source media="(prefers-color-scheme: dark)" srcset="https://raw.githubusercontent.com/Poojan6216/stepledger/main/bench/plots/history-dark.png">
  <img alt="Workflow history size against node count with and without External Storage" src="https://raw.githubusercontent.com/Poojan6216/stepledger/main/bench/plots/history-light.png">
</picture>

<picture>
  <source media="(prefers-color-scheme: dark)" srcset="https://raw.githubusercontent.com/Poojan6216/stepledger/main/bench/plots/storage-dark.png">
  <img alt="External store bytes per run: one object per payload against dedup chunks" src="https://raw.githubusercontent.com/Poojan6216/stepledger/main/bench/plots/storage-light.png">
</picture>

### Demo 4: the view equals the truth

100 seeded runs (parallel supersteps, `interrupt()` plus resume, effects, a cached continue-as-new, a workflow-side node, a within-run cache hit): equal 80, declared gap 20, unequal 0. The cached continue-as-new runs are EXACT with `chain=True` (8 of 8).

### Demo 5: the retry bill

|  | calls billed | USD billed (simulated) | wasted calls | wasted USD | journal replays |
|---|---|---|---|---|---|
| no-journal | 401 | 4.812 | 121 | 1.452 | 0 |
| journal | 297 | 3.564 | 17 | 0.204 | 79 |

The wasted-call numbers come from the fake model's billing hook in the bench, not from the ledger. Stepledger's own `sl_retry_waste` view counts only attempts that reached the ledger (4,000 tokens in this run), because a worker that dies right after the model call never writes a row; with the journal on, that call is in `sl_llm_calls`.

### The run ledger

The artifact itself: `stepledger ledger <workflow id>` for a crash-test run (attempt 1 of `scan_iam` wrote a row and died before reporting; attempt 3 won; one duplicate effect prevented):

```
STEPLEDGER  wf=chaos-SL-a8688a-13  run=01a0dec0-2178-7b33-88ef-c5103e24630f  status=COMPLETED  sealed  nodes=30 committed=30 abandoned=0
 seq node             step attempt status      output_hash   tokens retry_waste  note
   0 list_assets         1       2 COMMITTED   a359a4e3...      200           0  
   1 scan_iam            2       3 COMMITTED   fb2f7fb1...      200         400  attempt 1's row overwritten by attempt 3; divergent output discarded; attempt 2's row overwritten by attempt 3; divergent output discarded
   2 scan_network        2       2 COMMITTED   10f0eb1c...      200           0  
   3 scan_storage        2       3 COMMITTED   24f12319...      200           0  
   4 enrich_cve_0        3       1 COMMITTED   1441c0c2...      200           0  
   5 enrich_cve_1        4       1 COMMITTED   8ec67f84...      200           0  
   6 enrich_cve_2        5       1 COMMITTED   f776e34a...      200           0  
   7 enrich_cve_3        6       1 COMMITTED   2b6f347f...      200           0  
   8 enrich_cve_4        7       1 COMMITTED   4027abb7...      200           0  
   9 enrich_cve_5        8       1 COMMITTED   a3848e2f...      200           0  
  10 enrich_cve_6        9       1 COMMITTED   354580de...      200           0  
  11 enrich_cve_7       10       1 COMMITTED   4e84b420...      200           0  
  12 enrich_cve_8       11       1 COMMITTED   ba57c6fa...      200           0  
  13 enrich_cve_9       12       1 COMMITTED   a9593937...      200           0  
  14 enrich_cve_10      13       1 COMMITTED   7c303ec0...      200           0  
  15 enrich_cve_11      14       1 COMMITTED   aa52a8a0...      200           0  
  16 enrich_cve_12      15       1 COMMITTED   90ac84d8...      200           0  
  17 enrich_cve_13      16       1 COMMITTED   66ea2495...      200           0  
  18 enrich_cve_14      17       1 COMMITTED   cd75e6a1...      200           0  
  19 enrich_cve_15      18       1 COMMITTED   bdb790fc...      200           0  
  20 enrich_cve_16      19       1 COMMITTED   e1130e22...      200           0  
  21 enrich_cve_17      20       1 COMMITTED   892761a8...      200           0  
  22 enrich_cve_18      21       1 COMMITTED   66c23bd0...      200           0  
  23 enrich_cve_19      22       1 COMMITTED   44050f7f...      200           0  
  24 enrich_cve_20      23       1 COMMITTED   3a979dd7...      200           0  
  25 enrich_cve_21      24       1 COMMITTED   408c84fb...      200           0  
  26 score_risk         25       1 COMMITTED   b88ce4d1...      200           0  
  27 open_ticket        26       2 COMMITTED   2d8de359...        0           0  
  28 notify_slack       27       1 COMMITTED   8c4fb5a0...        0           0  
  29 summarize          28       1 COMMITTED   ad361f3c...      200           0  
effects: open_ticket x1 (key ec8dce...), notify_slack x1 (key 3338af...)    duplicates prevented: 1
```

## What this is not

- Not a LangGraph checkpointer, and not a replacement for Temporal's durability: Temporal stays the source of truth and the ledger is a projection of it.
- Not exactly-once for external effects: `once()` is at-least-once delivery with dedupe, and an unknown outcome stops for a person.
- Not a hosted service, a UI, or a new agent framework.

## Private APIs and upgrade policy

Both upstream features Stepledger builds on, `LangGraphPlugin` and External Storage, are marked experimental in temporalio 1.33. Stepledger reads five private symbols, all in [`_compat.py`](https://github.com/Poojan6216/stepledger/blob/main/src/stepledger/_compat.py): the plugin's `ActivityInput`/`ActivityOutput` and task-cache context variable, the SDK's activity definition lookup, and LangGraph's `task_path_str` and `MISSING`. `tests/unit/test_compat.py` checks each on every run, so an SDK release that moves one fails the test suite before it fails at runtime; the dependency ranges are the tested ones and are widened release by release. Enabling the plugin on in-flight runs is safe (`workflow.patched`); removing it is not, see [limitations](https://github.com/Poojan6216/stepledger/blob/main/docs/limitations.md).

## Docs

- [How it works](https://github.com/Poojan6216/stepledger/blob/main/docs/how-it-works.md): one step's life, commits, seal, reconcile, the read side
- [Keys and fencing](https://github.com/Poojan6216/stepledger/blob/main/docs/keys-and-fencing.md)
- [Storage and GC](https://github.com/Poojan6216/stepledger/blob/main/docs/storage-and-gc.md)
- [Effects, the LLM journal and the retry bill](https://github.com/Poojan6216/stepledger/blob/main/docs/effects.md)
- [Limitations](https://github.com/Poojan6216/stepledger/blob/main/docs/limitations.md)
- [SDK facts](https://github.com/Poojan6216/stepledger/blob/main/docs/sdk-facts.md): every SDK behavior relied on, with file and line

## Prior art, and where each stops

- **Temporal's LangGraph plugin** (`temporalio.contrib.langgraph`) runs nodes as Activities and caches task results across continue-as-new. It keeps no per-node record outside Temporal. Stepledger composes with it and never modifies it.
- **LangGraph checkpointers** (`PostgresSaver`) persist each superstep in-process; behind the Activity boundary they are bypassed.
- **Temporal External Storage and its S3 driver** are the claim-check pattern built into the SDK. The driver stores one object per payload; Stepledger adds a deduplicating Postgres driver and does not replace the mechanism.
- **DataDog's `temporal-large-payload-codec`** does claim-check as a codec plus a service, whole-blob.
- **Temporal's idempotency guidance** (key on the run ID plus the Activity ID) is guidance, not a mechanism; it does not cover divergent retries, zombies or commit visibility.
- **Fencing tokens** (Martin Kleppmann, "How to do distributed locking") are the idea behind the attempt fence.
- **The transactional outbox** is the pattern behind committing on the next node's transaction.
- **FastCDC** (Xia et al., USENIX ATC 2016), as used by restic and borg, through the `fastcdc` Python package.
- **SpecuNode**, the author's earlier project, for journal-before-use and idempotency keys for agent effects.

## License

Apache-2.0.
