feat(core): streams with consumers and the hopperui module (M8) - #10
Merged
Merged
Conversation
giraffesyo
force-pushed
the
feat/workflows
branch
from
September 26, 2026 03:25
7c21931 to
e4788ea
Compare
giraffesyo
force-pushed
the
feat/streams
branch
from
September 26, 2026 03:25
0eaaf97 to
acd2fe1
Compare
giraffesyo
force-pushed
the
feat/workflows
branch
from
September 26, 2026 14:49
e4788ea to
0e23dff
Compare
giraffesyo
force-pushed
the
feat/streams
branch
from
September 26, 2026 14:50
acd2fe1 to
fbce729
Compare
giraffesyo
force-pushed
the
feat/workflows
branch
from
September 26, 2026 14:51
0e23dff to
eede8e3
Compare
giraffesyo
force-pushed
the
feat/streams
branch
from
September 26, 2026 14:51
fbce729 to
1914fe2
Compare
giraffesyo
force-pushed
the
feat/streams
branch
from
September 26, 2026 14:54
1964e93 to
2e82dc5
Compare
giraffesyo
enabled auto-merge (squash)
September 26, 2026 14:55
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.
Streams, plus the
hopperuimodule: #11 was merged into this branch (its base), so this PR carries both and completes M8.hopperui (from #11)
hopperui.New(client, cfg) http.Handler(separate module): overview and queues, jobs and history with filters and paging, job detail, workflows as an SVG DAG, subscriptions, stream consumers with seek forms, clients. Actions (retry, cancel, pause, resume, seek) only withConfig.Authorize; read-only otherwise.hopperui/cmd/hopperuiis a standalone server.JobFilter.Afterfor paging. Makefile, CI, dependabot and docs cover the module.Streams
What changed
hopper_stream_events, append-only and partitioned by day (dropped afterConfig.StreamRetention, 7 days by default, by the same leader maintenance as history).client.Streams().Append/AppendTxwrite one event: topic, optional key, payload, headers, a database-generated message ID. Events are identified by(xid, seq).hopper.Consume(workers, hopper.Consumer{Name, Pattern, Queue, Start, MaxAttempts, Metadata}, fn). A consumer is a named position on the log plus a topic pattern; the leader turns the events after it into delivery jobs of kindstream:<name>(the event's key becomes the ordering key), so a consumer gets typed handlers, retries, dead-lettering and competing workers for free.MessagegainsStreamandPosition. New consumers start atStreamStartLatest(default) orStreamStartEarliest;Streams().Seekmoves one to a position, a time, the start or now;Streams().Readpages the log;Streams().Consumerslists them.pg_snapshotit has read to. A pump takes a fresh snapshot and delivers the events of transactions visible in the new snapshot but not the old (pg_visible_in_snapshot, with an xid range to keep the index scan bounded), in(xid, seq)order; a pump that stops part way keeps both snapshots and the cursor and continues the same delta. A transaction that commits late is delivered when it commits, whatever its xid, and a long-running transaction holds nothing back. The plan's original "readers see only events belowxmin" rule was tried first and dropped: any open writing transaction in the database stalled every consumer (the parallel test suite alone made appends invisible), and it would have made a long batch transaction in production delay all deliveries.hopper_streamthrough the coalescing notifier and poke the local leader; transactional appends notify from the transaction; the leader's stream loop also pumps everyPollInterval. Two leaders pumping one consumer serialize on its row (lock first, snapshot second).hopper streams consumers,hopper streams seek <consumer> -earliest|-latest|-time|-position, andhopper subscriptions list(promised in the plan for M6, never added).drivertestStreamson both drivers: ordering, reads by pattern and position, earliest/latest consumers, delivery metadata and ordering keys, partial pumps continuing a delta, seeks of every kind, a late commit delivered after younger transactions without being skipped or delivered twice, retention partitions. A client-level end-to-end test; CLI cases.hopperbench
The insert, claim and finalize paths are untouched by this PR (streams are new tables and a leader duty), so no benchmark run.