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
2 changes: 1 addition & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ All notable changes are recorded here. Versions follow
The 0.1.0 entry describes what the library is rather than what changed, since
there was nothing to have changed from. Everything above it is ordinary.

## Unreleased
## 0.14.1 — 2026-10-04

### Fixed

Expand Down
29 changes: 13 additions & 16 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,24 +41,20 @@ only writer.
The usual shape is a message log in one system and an analytical store in another, with a
pipeline extracting between them — two copies of every row and a job that keeps them in
step. There is no extraction step here and no second copy. A litelink log **is** an Iceberg
table, so the Parquet your messages were appended to is the Parquet an analytical engine
table, so the log your messages were appended to is the table an analytical engine
reads:

```python
import duckdb
import litelink
import streamcast

# Published through streamcast, live.
async with streamcast.publish(uri) as producer:
await producer.send({"event_ts": 1790038800123456, "price": 85565.0, "side": 1})

# The same bytes as a table, on the box that holds the log.
with litelink.open("data", "trades", read_only=True) as log:
log.sql("SELECT count(*), max(price) FROM log WHERE side = 1").read_all()

# Or from anywhere, from the stream's published tables: every log it has
# been through, as of a point, with no server involved.
# The same rows as a table, from anywhere, from the stream's published
# tables: every log it has been through, as of a point, with no server
# involved.
async with await streamcast.Stream.snapshot(
"s3://bucket/prefix/trades.metadata.json"
) as snapshot:
Expand All @@ -73,8 +69,8 @@ duckdb.sql("""
```

Rows land in a SQLite buffer first and seal into Parquet behind it, so the newest messages
are in the buffer and the rest are columnar — `log.sql` reads across both, and a reader
anywhere else reads what the log has published. That is one store with tiers, not a
are in the buffer and the rest are columnar — a snapshot given `broker=` reads across
both, and a reader anywhere else reads what the log has published. That is one store with tiers, not a
transactional copy and an analytical copy that have to be reconciled.

The tiering, the published-table layout, consistency guarantees, and costs are
Expand Down Expand Up @@ -254,12 +250,13 @@ should not be. `Maintain(dedicated=("trades",))` gives a named log a set of its
`Stream.new` creates or opens the log; `Stream(log=handle)` takes one you opened yourself
and does no I/O. `streamcast.to_arrow(SCHEMA)` is the `pa.schema` if you want it.

What it captures is a table, queryable without streamcast:
What it captures is a table, queried from its published tables:

```python
log.sql("SELECT count(*), max(price), sum(amount) FROM log").read_all()
log.scan(columns=["litelink_offset", "price"], where="side = 1") # prunes on statistics
log.sql("SELECT max(streamcast_ts - event_ts) FROM log").read_all() # feed latency, us
uri = stream.metadata_uri
await streamcast.Stream.sql(uri, "SELECT count(*), max(price), sum(amount) FROM log")
await streamcast.Stream.scan(uri, columns=["litelink_offset", "price"], where="side = 1")
await streamcast.Stream.sql(uri, "SELECT max(streamcast_ts - event_ts) FROM log") # feed latency, us
```

The table has two columns you did not declare. `litelink_offset` is the offset every frame
Expand Down Expand Up @@ -809,11 +806,11 @@ Offsets are per server and are not translated between hops.

## What it is not

- **Not a message broker.** No topics beyond a name, no consumer groups, and no consumer
- **Not a message queue.** No topics beyond a name, no consumer groups, and no consumer
acknowledgements — where a subscriber has got to is its own cursor, not state the server
keeps. A publisher does get an ack, the offset once the row is durable; nothing tracks
what a subscriber has consumed. A subscriber needing at-least-once with server-side acks
wants a queue.
wants one.
- **Not tuned for high fan-out across a WAN.** `compression` costs CPU per subscriber
while the encode is shared, so it defaults off — see the API section. Turn it on for
few subscribers over a WAN.
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[project]
name = "streamcast"
version = "0.14.0"
version = "0.14.1"
description = "A JSON WebSocket pubsub framework for structured data. Every message lands in an Iceberg table: no ETL, no second copy."
readme = "README.md"
authors = [
Expand Down
2 changes: 1 addition & 1 deletion uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading