Skip to content

feat(stream): retire a stream, serve it read-only, revive it on restore - #97

Merged
nhobin219 merged 3 commits into
mainfrom
stream-retire
Oct 4, 2026
Merged

nhobin219 merged 3 commits into
mainfrom
stream-retire

Conversation

@nhobin219

@nhobin219 nhobin219 commented Oct 4, 2026 •

Copy link
Copy Markdown
Owner

Closes #96. Three commits, reviewable separately.

1. Stream.retire, read-only serving, and restore(revive=True) (5869e8b)

streamcast.Stream.retire("trades", root="data")   # server stopped
stream = streamcast.Stream.restore(
    "trades", root="data", published="s3://bucket/prefix", revive=True,
)                                                  # undone, here or on another box

Retiring:

  • Makes the published table complete. litelink's retire() seals and publishes every row, including the trailing run a plain publish holds back. A stream that just stops being written leaves that tail on the box.
  • Refuses writers in litelink, for good.
  • Records the retirement in metadata.json, locally and beside the tables. It holds at (int64 UTC µs), end_offset, sort_by and the streamcast_ts span, so a revive needs nothing from the old box. The log's statistics join the manifest.
  • Safe to rerun: an already retired stream returns its record, and a retire that died after litelink's step finishes.
  • Older builds refuse it: a retired stream's metadata.json is version 3. A build from before retirement would otherwise take the retired log for a migration that died, and quietly start -v2. Reviving writes version 2 again.

A retired stream is read-only, and still served. Whether a stream can be written is the stream's to decide, not serve's:

  • new and migrate open its log read-only and set stream.retirement;
  • subscribers replay and catch up, always from the published table, because retiring evicted staging and the table is the only copy; a live view just sees nothing new;
  • a local send raises StreamRetired;
  • a publisher is refused with a new close code, 4410 (Close.RETIRED); the client raises StreamRetired with at and end_offset;
  • serve and asgi start no maintainer or litestream sidecar for it.

Reviving:

  • restore(..., revive=True) continues the stream on its next log (-vN), with the same columns and sort_by, starting at exactly the retired end, so the offsets stay one dense sequence.
  • Works on another box, reading only the published metadata, or on the same box with local tables.
  • Without revive, restore refuses a retired stream with StreamRetired.
  • A revive that died after creating its log carries on, if that log is empty at the seam.

litelink 0.10 (now >=0.10.1,<0.11, see commit 3):

  • restore now rebuilds a log that never had WAL replication from its published table (litelink#144).
  • Stream.restore passes through replica_reserve and published_reserve. They default to litelink's own (2^20 and 2^40), passed on only when given.

2. publish= removed from serve and asgi (ccccddc, breaking)

  • Every served stream takes publishers, and a retired one refuses them itself (4410). The flag guarded writes while reads stayed open to anyone who could reach the port, so it protected little; authentication is the fix for both.
  • Callers: 39 publish=True calls are removed from tests and examples, and the docs are updated.
  • Older servers: the client keeps its publish_disabled message for a server from before this change, now worded to name upgrading.
  • Tests: the two "publishing is off unless allowed" tests become "every served stream takes publishers", plus a test that the older server's refusal still reads correctly.

3. sort_by in the metadata, and restore with the exact shape (632bb9b)

litelink 0.10.1 (#148) checks a restore's schema and sort_by against the replica or the published table. For a table no 0.10 publish stamped, it requires them, because an Iceberg schema alone loses Arrow field metadata, which is where a binary column's encoding lives.

  • metadata.json records each log's sort_by beside its schema.
    • new, migrate and revive write it, and serve's startup check fills it in for an existing stream from its live log.
    • A log not on the machine stays null (unknown). No version bump: older builds ignore the field.
  • _metadata.shape(entry) rebuilds a log's Arrow schema exactly: declared columns with their encodings, plus only the system columns that log has.
  • Stream.restore(schema=, sort_by=, config=):
    • by default the shape comes from the metadata;
    • schema= (JSON Schema) and sort_by= override, or cover streams whose metadata predates recording them;
    • config= sets the restored log's policy, or the new log's with revive=True.
  • The restore docstring now covers both fences, the shape and revive.
  • Tests: tests/test_restore_shape.py, 7 tests.
    • The S3 restore uses a base16 binary column with a sort_by. litelink checks the shape exactly, and the test records the arguments litelink's restore received.
    • Falsified: each of these breaks fails a test: stripped field encodings, sort_by not passed, serve not filling it in, config dropped on restore, config dropped on revive.
  • On litelink 0.10.1: the restore, retire, shape, metadata and migrate tests pass (86). The full suite passed, 705 tests, against litelink main at the 0.10.1 commit.

Tests

  • tests/test_retire.py, 13 tests, about 6 s:
    • retiring publishes the tail a plain stop leaves behind;
    • the metadata record and its version 3;
    • rerunning, and finishing a half-done retire;
    • read-only opening via new and migrate;
    • served read-only: replay works, a publisher gets 4410 with at and end_offset;
    • no maintainer;
    • restore refuses without revive;
    • revive on the same box: dense offsets, catch-up across the seam, metadata back to version 2;
    • revive crash-resume;
    • on S3: retire on one box, delete it, revive on another, with the published metadata flipped;
    • on S3: a stream without WAL replication restores from its published table, with published_reserve honoured.
  • Falsified: each of these breaks fails a test:
    • retire without publishing;
    • sends not refused;
    • publishers not refused;
    • a retired log maintained;
    • revive at the wrong offset.
  • Locally, on litelink 0.10: test_retire, test_restore, test_migrate and test_metadata (79), plus every test file touched by commit 2 (198).

🤖 Generated with Claude Code

https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi

nhobin219 and others added 3 commits October 4, 2026 17:54
Stream.retire publishes every row (the trailing run a plain publish holds
back included), retires the log in litelink, and records the retirement
in metadata.json, locally and beside the tables, at version 3 so an older
build refuses it rather than starting a new log.

A retired stream opens read-only and is still served: it replays from
the published table, its only copy; a local send raises StreamRetired, a
publisher is refused with the new close code 4410, and it gets no
maintainer or sidecar. restore(revive=True) continues it on the next log
at exactly the retired end, on any box. Requires litelink 0.10, whose
restore also rebuilds a log with no WAL replica; its replica_reserve and
published_reserve pass through.

Closes #96.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi
Every served stream takes publishers. Whether a stream can be written
is the stream's to say: a retired one refuses publishers with 4410. The
flag guarded writes while reads stayed open to anyone who could reach
the port, so it protected little; authentication is the fix for both.

The client keeps its message for publish_disabled, which a server from
before this still sends.

BREAKING CHANGE: serve(publish=) and asgi(publish=) are removed; passing
it raises TypeError.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi
metadata.json records each log's sort_by beside its schema; serve fills
it in for an existing stream from its live log. Stream.restore passes
litelink 0.10.1 the log's shape from the metadata: its exact schema,
binary encodings and system columns included, which an Iceberg schema
does not keep, and its sort_by. schema= and sort_by= override, and
config= sets the restored log's policy, or the new log's on revive.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi
@nhobin219
nhobin219 merged commit 5fefae6 into main Oct 4, 2026
6 checks passed
@nhobin219 nhobin219 mentioned this pull request Oct 4, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Stream.retire: finish a stream for good, and a restore flag that undoes it

1 participant