Repository navigation
pineforge-feed: serve fan-out, segmented journal, REST catch-up and backtest bar export - #19
Merged
Merged
Conversation
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>
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.
Summary
Adds a fan-out server, a segmented journal, a restructured REST catch-up and a backtest bar export to the native
pineforge-feedadapter innative/feed/. Public market data only; the Python package and wheel are unchanged.serve: one producer, many runners.GET /v1/status,GET /v1/snapshot(a complete prefix of at most 4 MiB), andws://…/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.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.--qty-stepgrid, or the engine's compensated sum without a step.Verification
uv run pytest -qandruff check .pass.serveproducer feeding two runners on bars, with a runner restart from its committed cursor and a producer SIGKILL plus resume;servein tick mode;All pass. Bars and prints equal REST, and actions equal the batch.
🤖 Generated with Claude Code