Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -450,6 +450,23 @@ frames, then restore. Unplanned failover loses whatever never shipped.
> succeeds — see [`SPEC.md`](docs/SPEC.md) §8b and
> [litelink#75](https://github.com/nhobin219/litelink/issues/75).

### Retiring a stream

A stream that just stops being written leaves its last rows on the box, unpublished.
`Stream.retire` finishes it: every row published, the log refusing writes for good, and
the retirement recorded in its metadata. It is still served, read-only — subscribers
replay and catch up, publishers are refused with 4410. `Stream.restore(..., revive=True)`
undoes it on any box, continuing at exactly the retired end, so a planned move is
`retire` on the old box and `revive` on the new one, losing nothing and skipping no
offsets:

```python
streamcast.Stream.retire("trades", root="data") # old box, stopped
stream = streamcast.Stream.restore(
"trades", root="data", published="s3://market-data/prod", revive=True,
) # new box
```

### Serving the whole history

`replay_published=True` with `max_replay=None` makes the server a complete gateway to the
Expand Down
14 changes: 7 additions & 7 deletions SECURITY.md
Original file line number Diff line number Diff line change
Expand Up @@ -75,17 +75,17 @@ policy is what decides who may read them, not this library. There is no
switch to leave the location out of the greeting; if it is itself sensitive,
the port should not be reachable by anyone who should not learn it.

**A server that allows publishing accepts writes from anyone who can reach
it.** `serve(..., publish=True)` is opt-in for that reason, and off by default
so an upgrade cannot make a server writable on its own. With it on, the
authentication above stops being optional: a published row is durable, every
subscriber sees it, and no consumer cursor undoes it. `process_request` is the
**A server accepts writes from anyone who can reach it.** Every served stream
takes publishers — only a retired one refuses them — so the authentication
above is not optional on a reachable port: a published row is durable, every
subscriber sees it, and no consumer cursor undoes it. (Before 0.14.0 this was
opt-in, `serve(..., publish=True)`; that switch is gone, and an upgraded server
accepts publishers.) `process_request` is the
hook — it sees the request before the WebSocket opens, including the
`?publish` query that distinguishes a publisher from a subscriber, so a policy
can allow reads and refuse writes on the same port.

What IS in scope: anything that lets a subscriber see messages from a stream it
did not subscribe to, receive a stream that silently differs from another
subscriber's, make the server exhaust memory through a path `max_backlog` is
supposed to bound, or publish into a stream on a server where `publish=True`
was not set.
supposed to bound, or publish into a retired stream.
10 changes: 10 additions & 0 deletions docs/API.md
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,14 @@ streamcast.Stream.new(name="", *, root, schema, # creates or opens the
sort_by=None, config=None, published=None, s3_options=None,
replay_published=False, group_commit=True,
max_replay=100_000)

streamcast.Stream.migrate(name="", *, root, schema, ...) # a new log, a new schema
streamcast.Stream.restore(name="", *, root, published, # on a box without the log
s3_options=None, schema=None, sort_by=None, config=None,
replica_reserve=None, published_reserve=None,
revive=False, ...)
streamcast.Stream.retire(name="", *, root, s3_options=None) # for good; returns the
# record: .at, .end_offset
```

`name` is where it is served: `"trades"` at `/trades`, `""` at `/`. It is the name's only
Expand Down Expand Up @@ -1006,6 +1014,8 @@ StreamcastError
├── StreamNotFound nothing served at that path; `.serves` lists what is
├── NotReplayable that offset cannot be served; `.why` says which of five
├── CatchUpUnavailable `catch_up=True` could not read the gap; says what to change
├── Rejected a published row the schema refuses; nothing committed
├── StreamRetired the stream is retired and takes no rows; `.at`, `.end_offset`
└── TooSlow dropped for falling behind; `.offset` is where to resume
```

Expand Down
34 changes: 25 additions & 9 deletions docs/SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -996,7 +996,7 @@ where it asked — a hole at the join, which is the one wrong answer a resume mu
never give.

**An empty scan below the frontier is not always a hole.** A restore fences
2**20 offsets that were never issued, so a replay inside the fence reads
2**20 offsets that were never issued (2**40 with no WAL replica), so a replay inside the fence reads
nothing and there is nothing to refuse. Rows evicted from staging read nothing
too. The published tier tells them apart: if it holds rows from the requested
offset on, the rows exist and staging has dropped them, so the replay is
Expand Down Expand Up @@ -1349,8 +1349,8 @@ places an offset arrives from somewhere TCP does not cover:
The first is the only one reachable today. What both would otherwise be is
silent: processing a stream whose offsets went backwards means skipping data
once a cursor is involved. `<=` rather than `!= previous + 1`, because
litelink's offset space has legitimate gaps — a `restore` fences 2**20 of them
— so a jump forward is ordinary and only a step backwards is wrong.
litelink's offset space has legitimate gaps — a `restore` fences 2**20 of them,
or 2**40 with no WAL replica — so a jump forward is ordinary and only a step backwards is wrong.

It does **not** span a reconnect: a new `Subscription` starts with no previous
offset, so nothing is compared across the gap. The log is what makes resuming
Expand All @@ -1374,6 +1374,7 @@ against the source. I3 and I4 are checked end to end. I5 is litelink's.
| a replay outruns `max_backlog` | the subscriber is dropped right after catching up. Size the two together (§4) |
| two publishers on one log | litelink refuses: one writer per log. A second server on the same directory fails to open |
| the server is restored from a replica | offsets are fenced by litelink and jump; a consumer resuming into the fence gets `ahead` rather than silence |
| a publisher writes to a retired stream | refused with 4410 (`StreamRetired`); its readers are unaffected |
| a consumer was down past `max_replay` | refused with `too_old`; `catch_up=True` reads the gap from the published tables and then connects (§5) |
| a catching-up consumer has no credentials | `CatchUpUnavailable` at `connect`, naming the endpoint, the credential source, and four ways out |

Expand All @@ -1382,24 +1383,31 @@ against the source. I3 and I4 are checked end to end. I5 is litelink's.
## 8b. Recovering a server

`Stream.restore(name, root=…, published=…)` stands a stream up on a box that
never held its log: litelink rebuilds it from the published table and the
replicated WAL, and the result serves and appends like any other.
never held its log: litelink rebuilds it from its replicated WAL when there is
one, and otherwise from its published table alone (litelink 0.10), and the
result serves and appends like any other. Its schema and `sort_by` come from
the stream's metadata, exactly — an Iceberg schema keeps no Arrow field
metadata, where a binary column's encoding lives — and litelink checks them
against the replica or the table.

**Offsets are fenced, not reissued**, and that is what makes the move safe for
consumers. litelink burns 2**20 offsets, so the restored stream resumes above
anything the dead machine may have served. No offset a consumer holds is ever
consumers. litelink skips 2**20 past what a replica recorded, or 2**40 past what
the published table says the log issued — rows written after the last publish
are gone with the machine — so the restored stream resumes above anything the
dead machine may have served. No offset a consumer holds is ever
handed out again carrying different data — the one thing a resume cannot
survive. `recv` permits a forward jump for exactly this reason (I4).

**The fence is a million offsets wide, and it does not strand anyone,
**The fence is a million offsets wide, or a trillion without a replica, and it
does not strand anyone,
because `max_replay` counts ROWS rather than offset distance.** A consumer
150 rows behind a failed-over producer is 150 rows behind; measuring it as
1,048,746 was a property of the proxy, not of the work.

`max_replay` exists to bound what a replay costs, and that cost is rows.
Offset distance is a proxy for it and an exact one only while the offset space
is dense — which litelink's is not, by design: a `restore` fences 2**20
offsets that were never issued, and I4 already says a forward jump is
offsets that were never issued (2**40 with no WAL replica), and I4 already says a forward jump is
ordinary. So the distance check runs first, free, and only a subscribe it
would REFUSE pays to find out what the replay actually costs:

Expand Down Expand Up @@ -1447,6 +1455,14 @@ and both published to the same table. Nothing refused, nothing warned.
requirement, not something either library enforces, and it is
[tracked upstream](https://github.com/nhobin219/litelink/issues/75).

**A planned move needs no fence: retire, then revive.** `Stream.retire` fixes
the log's end — litelink stops appends and publishes every row before it — and
records it in the metadata. `Stream.restore(..., revive=True)` then continues
the stream on its next log at exactly that end, on any box: nothing after it
was ever issued, so nothing is skipped, and nothing was left unpublished, so
nothing is lost. Until revived, a retired stream is served read-only, from its
published table.

---

## 9. Open
Expand Down
3 changes: 2 additions & 1 deletion src/streamcast/_catchup.py
Original file line number Diff line number Diff line change
Expand Up @@ -288,7 +288,8 @@ async def stream(self) -> AsyncGenerator[tuple[int, int | None, dict], None]:
f"keeps moving. Publish more often, raise the server's "
f"`max_replay`, or pass a larger `catch_up_retries`.\n"
f" * The server was RESTORED onto another machine. litelink "
f"fences offsets on a restore — 2**20 of them — so the range "
f"fences offsets on a restore — 2**20 of them, 2**40 with no WAL "
f"replica — so the range "
f"below its window was never issued and no table will ever hold "
f"it. Only raising `max_replay` (or `None`) helps. Restore a "
f"failed-over producer with `max_replay=None` if existing "
Expand Down
3 changes: 2 additions & 1 deletion src/streamcast/_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -522,7 +522,8 @@ async def _read(self) -> Row:
# what makes that safe, not this.
#
# `<=` rather than `!= previous + 1`: litelink's offset space has
# legitimate GAPS (a `restore` fences 2**20 of them), so a jump
# legitimate GAPS (a `restore` fences 2**20 of them, or 2**40 with no
# WAL replica), so a jump
# forward is ordinary and only a step backwards is wrong.
if (
offset is not None
Expand Down
2 changes: 1 addition & 1 deletion src/streamcast/_limits.py
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,7 @@
live view stops, each saying the published tables are too far behind, rather
than the reader running out of memory waiting on a publisher that has
stalled. Counted in rows read, not offset distance, which a restore fence
stretches by 2**20.
stretches by 2**20, or 2**40 with no WAL replica.
"""


Expand Down
2 changes: 1 addition & 1 deletion src/streamcast/_log.py
Original file line number Diff line number Diff line change
Expand Up @@ -336,7 +336,7 @@ def rows_from(log: LogHandle, offset: int, *, published: bool = False) -> int:
**`max_replay` bounds the WORK a replay costs, and that work is rows.**
Offset distance is a proxy for it, and an exact one only while the offset
space is dense — which litelink's is not. A `restore` fences 2**20
offsets that were never issued, so a consumer 150 rows behind measures as
offsets that were never issued (2**40 without a WAL replica), so a consumer 150 rows behind measures as
a million and a bounded server refuses a replay it could serve instantly.

Asked only when the cheap proxy has already said "too far", so the common
Expand Down
11 changes: 7 additions & 4 deletions src/streamcast/_stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -1281,6 +1281,8 @@ def schema(self) -> dict[str, object] | None:
def retirement(self) -> _metadata.Retirement | None:
"""When the stream was retired and where it ended, or None if it is not.

Not `retired`, which lists the logs a migration sealed.

A retired stream is served read-only: it replays, snapshots and
catches up, and refuses every send with `StreamRetired`. See
`Stream.retire`.
Expand Down Expand Up @@ -1329,9 +1331,10 @@ def ensure_metadata(self) -> None:
def retired(self) -> tuple[tuple[Path, str], ...]:
"""The `(root, name)` of each retired log still on this disk.

Empty for a stream that has never migrated. `serve` maintains these
beside the current log, so their local retention keeps running; it
never writes to them.
The logs a MIGRATION sealed, not whether this stream is retired —
that is `retirement`. Empty for a stream that has never migrated.
`serve` maintains these beside the current log, so their local
retention keeps running; it never writes to them.
"""
return self._retired

Expand Down Expand Up @@ -1960,7 +1963,7 @@ async def _resolve(self, requested: int | None) -> tuple[LogHandle, int] | None:
# `max_replay` bounds the work a replay does, and that work is
# rows — offset distance is a proxy, exact only while the offset
# space is dense. It is not: a `restore` fences 2**20 offsets
# that were never issued, so a consumer 150 rows behind a
# that were never issued (2**40 with no WAL replica), so a consumer 150 rows behind a
# failed-over producer measures as a million and is refused a
# replay the server could serve instantly.
#
Expand Down
Loading