Skip to content

pineforge-feed: serve fan-out, segmented journal, REST catch-up and backtest bar export - #19

Merged
luisleo526 merged 15 commits into
mainfrom
lv/feed-serve
Oct 4, 2026
Merged

luisleo526 merged 15 commits into
mainfrom
lv/feed-serve

Conversation

@luisleo526

Copy link
Copy Markdown
Contributor

Summary

Adds a fan-out server, a segmented journal, a restructured REST catch-up and a backtest bar export to the native pineforge-feed adapter in native/feed/. Public market data only; the Python package and wheel are unchanged.

  • serve: one producer, many runners.
    • Embedded CivetWeb 1.16, pinned by SHA-256 and attributed in NOTICE. It is built only behind a build option, with TLS, CGI, Lua, WebDAV, SSI and file serving disabled, and it listens on loopback by default (IPv4 and IPv6). A non-loopback bind needs an explicit flag; put it behind a reverse proxy for TLS.
    • GET /v1/status, GET /v1/snapshot (a complete prefix of at most 4 MiB), and ws://…/v1/stream?epoch=E&from=I, which streams from message index I with immutable message boundaries. Cursor negotiation is HTTP, never an event inside the runner input.
    • Bounds:
      • replay and per-client queues are bounded;
      • a slow client is closed (1008) without blocking the producer or other clients;
      • request and handshake timeouts are bounded, and workers are reserved, so status and new runner connects stay available;
      • client data frames close the connection.
    • Quiet symbols: the server sends an unsolicited PONG every 5 s, which keeps the runner's WebSocket idle deadline from firing.
  • Segmented journal. The single capped log becomes sealed segments with hashed checkpoints, retained by bytes and by venue time. A cursor older than retention fails explicitly (exit 22, HTTP 410) with recovery guidance. Persist-before-publish and group commit hold across seal and expiry; crash points are exercised by kill-and-resume tests.
  • REST catch-up. A long outage now reconnects and heals incrementally instead of stopping with exit 22. Binance REST is paced by the venue's used-weight header (up to 90%, per process), and quiet streams get a liveness probe instead of a reconnect. Completeness rules are unchanged.
  • export: 1-minute bars built from prints. The bars equal the runner's own tick-built bars, value for value and bit for bit, so a tick-mode deployment can be backtested on exactly the bars its forward run builds.
    • The source is REST within retention, or the venue's public daily aggregate-trade archive (checksum-verified).
    • Volume follows the engine's rule: the exact sum on a --qty-step grid, or the engine's compensated sum without a step.
    • A quiet first minute carries the warmup close, as the runner does.

Verification

  • Tests: unit and mock suites pass in Release, ASan+UBSan and TSan, including a TSan stress with many clients connecting and disconnecting during segment rotation, with no sanitizer suppressions. Kill-and-resume fuzzing gives an exact tape in every variant. uv run pytest -q and ruff check . pass.
  • Live, public data, about 17 minutes per round:
    • one serve producer feeding two runners on bars, with a runner restart from its committed cursor and a producer SIGKILL plus resume;
    • serve in tick mode;
    • the export gate against the runner's bars.
      All pass. Bars and prints equal REST, and actions equal the batch.
  • Independent review: MERGE-WITH-FIXES, then a re-review: MERGE.

🤖 Generated with Claude Code

luisleo526 and others added 15 commits October 4, 2026 08:31
Journal: the normalized log is now a sequence of segments under journal/,
each with a hashed checkpoint of the verified state before its first
message. A full segment is sealed after a commit (successor checkpoint,
then the empty successor, then the cursor that names it), and sealed
segments expire by retained bytes (--replay-bytes) and venue time
(--replay-age-seconds), never the open segment or one holding the
protected window (open minute and closed-minute proofs). Resume verifies
the retained journal from its oldest checkpoint, removes the residue of
an interrupted seal or expiry, and refuses an expired --output-from with
recovery guidance. The 256 MiB terminal budget (--max-log-bytes) is gone.
Group commit, persist-before-publish and message boundaries are kept.
PINEFORGE_FEED_CRASH_AT is a test hook that kills the process at a named
durable step.

serve: the same producer publishes to WebSocket clients over CivetWeb
1.16 (MIT, fetched by release archive and pinned by SHA-256, HTTP and
WebSocket core only, no TLS, CGI, Lua, WebDAV, SSI or file serving,
behind PINEFORGE_FEED_SERVE). GET /v1/status, GET /v1/snapshot (complete
prefix from index 0, at most 4 MiB) and ws://HOST/v1/stream?epoch=E&from=I
with cursor negotiation in the HTTP upgrade (400/409/410/416/503). Each
client reads the journal to the head, then drains a bounded queue
(--client-queue-bytes); an overflowing client is closed with 1008 "slow
consumer" and never blocks the producer or other clients. An unsolicited
PONG after 5 s without a message keeps the runner's 15 s idle deadline.
Loopback by default; another interface needs --allow-remote-listen.

REST catch-up: a long outage no longer ends in exit 22: when the reader's
queue is full it drops the connection, waits for the session to drain and
reconnects, and REST heals the rest after the overlap check. Binance REST
is paced by the venue's X-MBX-USED-WEIGHT-1M count up to 90% of the
one-minute ceiling instead of a fixed half-rate floor. A quiet Binance
stream is proven alive with a LIST_SUBSCRIPTIONS probe instead of a
reconnect every 75 s (--silence-seconds).

Follow-ups: the OKX multiplier test waits on answered instrument reads
instead of a pause and also runs in tick mode; the USD-M quiet start
searches hour windows up to the newest venue time (48 h at most) and has
its own test; Venue::predecessor_if_ready is pure virtual; Bybit's 403
stop names a regional block and a ten-minute backoff.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
pineforge-feed export writes the bars a tick-mode runner builds from the
same prints, in the warmup CSV format: per minute the first, highest,
lowest and last price and the exact decimal volume sum, and a quiet
minute that repeats the previous close with zero volume. The window is
proven by a contiguous ID chain from the last print before --start to
the first print at or after --end; anything less stops with 20 and
writes nothing.

Sources: venue REST within its history (USD-M aggTrades by fromId within
48 h, Binance spot historicalTrades, OKX history-trades), or a local copy
of the Binance USD-M daily aggTrades archive as published: the zip is
inflated in-process (no new dependency) and CRC-checked, or the extracted
CSV is read, streamed either way; --checksum verifies the published
CHECKSUM file's hash and name. A manifest records the source, window,
fence prints, counts and the CSV's SHA-256.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
README: the serve endpoints, cursor negotiation, resume with from=S+N and
--from-input N, client bounds, the keepalive PONG and the TLS-by-reverse-
proxy rule; the segmented journal with retention by bytes and venue time,
crash safety at every seal and expiry step and the expired-cursor stop;
export of prints-built bars from REST or the daily archive; header-based
Binance pacing, the quiet-stream probe and reader backpressure; and the
follow-up wording for OKX re-checks, Bybit 403 and quiet starts.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The sanitizer jobs set only the C++ compiler; serve also compiles
CivetWeb's C source, which must be instrumented by the same toolchain.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
ThreadSanitizer found CivetWeb 1.16's stop flag written by mg_stop while
its master thread reads it. STOP_FLAG_NEEDS_LOCK makes those accesses
atomic; the one remaining report is STOP_FLAG_ASSIGN's plain read of the
expected value inside its compare-and-swap loop, which tests/tsan.supp
suppresses by that function name only. The TSan CTest rows use it.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The USD-M start lookup searches hour windows forward only up to the
newest venue time it has seen, which export (no WebSocket) never set:
a window whose first aggregate came more than an hour after --start
stopped 20. Venue::horizon lets export pass the clock, bounded by the
48-hour history. README throughput adds the segmented-journal and serve
bench figures.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
serve: a client holds a journal segment only while it catches up (one at
the head opens none), so retained segments are not pinned by live
clients; serve raises its open-file limit and refuses to start below
2 x --max-clients + 64; a client whose segment expires while it catches
up is closed with 4410 at once instead of after a silent 30 s; the
stopping check no longer sends a close frame under the hub mutex and is
read per line while catching up; an EMFILE is an I/O error, never an
expired cursor; no exception unwinds into CivetWeb callbacks.

Journal: in tick modes the segment holding the newest print never
expires, so a quiet fence-venue stream can still re-read its recent
prints on resume; the predecessor overlap is checked only before the
first print. Recovery removes only real crash leftovers (the empty
successor of an interrupted seal, the oldest checkpoint of an
interrupted expiry), stops 21 on anything else, truncates the open
segment only after the journal verified, and maps filesystem errors to
22. Expiry erases its entry and syncs the directory between the two
unlinks; a failed seal or expiry poisons the journal. Lookups stop at
the first match, tick sessions compare a duplicate candle with the proofs
only, and a retained duplicate is checked before the expiry rule applies.

Binance probe: a reply must list every stream of the connection; an
error reply or a partial list reconnects instead of being decoded as
data (23). REST weight windows are keyed to the request's send time. A
--max-queue-bytes below one connection marker stops 22 instead of
waiting forever.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
CivetWeb runs one worker per connection and every streaming client keeps
its worker, so incomplete requests could take the few spare workers and
block status and new runner connects for 30 s. Run max-clients + 16
workers and drop a request that is not complete within 10 s, which also
bounds a blocked write to a dead client.

A client never sends messages: any frame it sends after serve has closed
it now ends the connection instead of being read on. The snapshot fails
at once when a segment expires mid-read instead of polling, and
--listen accepts the IPv6 loopback ([::1]:PORT; CivetWeb built with
USE_IPV6). A client cut off for overflowing its queue is logged.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The runner's carry-forward bar after warmup repeats the last warmup
close, not the last print before the start. export now takes --warmup
(the runner's warmup CSV, ending at the minute before --start) and a
quiet first minute carries its close; without it a quiet first minute
stops 20 instead of guessing. The manifest records warmup_close.

An archive window whose predecessor or fence lies in a neighbouring
daily file stops 20 and names that file.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The export gate compared with bars the harness built from the tape. It
now replays the runner's committed ledger inputs through the observed
build of the strategy, requires the runner's recorded state hash after
every input, and reads the bar the strategy saw at each time event
through the engine observer. Open, high, low and close must be equal;
volume may differ by one ulp, since without a qty_step the runner sums
tick volume in compensated doubles.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Without REENTRANT_TIME, CivetWeb formats the Date header with gmtime(),
whose static result two workers overwrite when they answer at once, as
the request timeouts of several incomplete requests do (found by TSan
in the half-open scenario). The mock's serve reader now keeps non-JSON
stderr lines, so a sanitizer report shows in the failure.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The missing predecessor is in the day before the archive's first row,
and the missing fence in the day after its last row. Deriving them from
the window named the wrong file when the archive's first print came
after the window start, or when the window reached past the day.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The export wrote the exact decimal sum of the quantities, while the
runner builds bar volume with the engine's tick-volume rule, so a bar
could differ from the runner's by one ulp. Port that rule operation for
operation: with a decimal quantity step 1/10^k the exact sum of grid
units, rounded once; otherwise, or after an off-grid quantity or an
overflow, the compensated binary64 sum in print order. Quantities are
read with std::stod as the runner reads the feed, and the volume is
written as the shortest round-trip decimal.

export takes --qty-step to match a runner deployed with
--syminfo qty_step=V; without it the rule is the runner's default. The
manifest records qty_step and volume_rule. The unit suite runs the
engine's own tick-volume vectors against the port.

The live export gate now requires every field to equal the runner's
double, with no tolerance, and the harness passes --qty-step to the
runner, the export, the replay and the batch probe.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@luisleo526
luisleo526 merged commit 05bf258 into main Oct 4, 2026
17 checks passed
@luisleo526
luisleo526 deleted the lv/feed-serve branch October 4, 2026 04:42
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.

1 participant