Repository navigation
feat(stream): retire a stream, serve it read-only, revive it on restore - #97
Merged
Merged
Conversation
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
Merged
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Closes #96. Three commits, reviewable separately.
1.
Stream.retire, read-only serving, andrestore(revive=True)(5869e8b)Retiring:
retire()seals and publishes every row, including the trailing run a plainpublishholds back. A stream that just stops being written leaves that tail on the box.metadata.json, locally and beside the tables. It holdsat(int64 UTC µs),end_offset,sort_byand thestreamcast_tsspan, so a revive needs nothing from the old box. The log's statistics join the manifest.metadata.jsonis 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:newandmigrateopen its log read-only and setstream.retirement;liveview just sees nothing new;sendraisesStreamRetired;Close.RETIRED); the client raisesStreamRetiredwithatandend_offset;serveandasgistart no maintainer or litestream sidecar for it.Reviving:
restore(..., revive=True)continues the stream on its next log (-vN), with the same columns andsort_by, starting at exactly the retired end, so the offsets stay one dense sequence.revive,restorerefuses a retired stream withStreamRetired.litelink 0.10 (now
>=0.10.1,<0.11, see commit 3):restorenow rebuilds a log that never had WAL replication from its published table (litelink#144).Stream.restorepasses throughreplica_reserveandpublished_reserve. They default to litelink's own (2^20 and 2^40), passed on only when given.2.
publish=removed fromserveandasgi(ccccddc, breaking)publish=Truecalls are removed from tests and examples, and the docs are updated.publish_disabledmessage for a server from before this change, now worded to name upgrading.3.
sort_byin the metadata, and restore with the exact shape (632bb9b)litelink 0.10.1 (#148) checks a restore's
schemaandsort_byagainst 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.jsonrecords each log'ssort_bybeside its schema.new,migrateandrevivewrite it, andserve's startup check fills it in for an existing stream from its live log.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=):schema=(JSON Schema) andsort_by=override, or cover streams whose metadata predates recording them;config=sets the restored log's policy, or the new log's withrevive=True.restoredocstring now covers both fences, the shape andrevive.tests/test_restore_shape.py, 7 tests.sort_by. litelink checks the shape exactly, and the test records the arguments litelink'srestorereceived.sort_bynot passed,servenot filling it in,configdropped on restore,configdropped on revive.Tests
tests/test_retire.py, 13 tests, about 6 s:newandmigrate;atandend_offset;restorerefuses withoutrevive;published_reservehonoured.test_retire,test_restore,test_migrateandtest_metadata(79), plus every test file touched by commit 2 (198).🤖 Generated with Claude Code
https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi