From bcb160f36f4c3ca858b50cfdd1df918ce86aa322 Mon Sep 17 00:00:00 2001 From: Emil Lerch Date: Fri, 28 Aug 2026 14:11:59 -0700 Subject: [PATCH] add assemble/record commands --- README.md | 153 +++++++++++++++++--- build.zig | 6 + build.zig.zon | 12 ++ data/observed.srf | 3 + src/assemble.zig | 344 ++++++++++++++++++++++++++++++++++++++++++++ src/cache_files.zig | 282 ++++++++++++++++++++++++++++++++++++ src/data.zig | 206 ++++++++++++++++++++++++++ src/main.zig | 273 +++++++++++++++++++++++++++++++++-- src/series.zig | 41 ++++-- src/srf_num.zig | 81 +++++++++++ tools/gen_model.py | 32 ++++- 11 files changed, 1382 insertions(+), 51 deletions(-) create mode 100644 data/observed.srf create mode 100644 src/assemble.zig create mode 100644 src/cache_files.zig create mode 100644 src/srf_num.zig diff --git a/README.md b/README.md index 51b2e41..d363537 100644 --- a/README.md +++ b/README.md @@ -4,8 +4,10 @@ Daily unit-value history for the two Oregon College Savings Plan (Embark) portfolios held in the finance repo: **ORCBI** and **ORC42**. ``` -zfin-vestwell reconstruct # rebuild data/ORCBI.srf and data/ORC42.srf -zfin-vestwell verify # check them against every evidence source +zfin-vestwell reconstruct # rebuild data/ORCBI.srf and data/ORC42.srf +zfin-vestwell verify # check them against every evidence source +zfin-vestwell record # append today's published value to the log +zfin-vestwell assemble --out DIR # write zfin cache files for it to serve ``` ## The problem @@ -104,6 +106,68 @@ Those values are rounded to 2dp at source, which is already +/-0.005 of noise. the plan's published cost table (0.234%/yr for ORCBI, 0.239% for ORC42); a fee is not identifiable over a short gap, so per-gap figures are noisy by nature. +## Serving it to zfin + +`ORC42` and `ORCBI` have no provider coverage, so zfin can only see them if +something puts files in its candle cache. That is what `assemble` does. + +``` +zfin-vestwell record # feed -> data/observed.srf (idempotent) +zfin-vestwell assemble --out "$CACHE" # series + log + feed -> cache files +``` + +`assemble` writes `//candles_daily.srf` and `candles_meta.srf`, and +needs no zfin cache of its own -- only the committed series and the observation +log. Run it on whatever host owns the cache zfin reads. On ZFIN_SERVER that makes +the symbols available to every client, because a read never triggers a refetch: +the server serves a present file as-is. + +Verified against a real zfin: + +``` +$ zfin diagnose ORC42 +local newest 2026-08-27, TTL still in the future, external +adj basis 1970-01-01 - no dividends or splits cached, nothing to restate +tracked NO - no normal run fetches this symbol +server the server will never refresh this symbol on its own +provider tiingo FAILED: NotFound +``` + +Every line there is load-bearing. The cache is fresh and tagged `external`; the +adjustment-restatement path cannot fire because there are no corporate actions; +the symbol is untracked so no refresh sweep touches it; and no provider carries it, +which is precisely why the cache has to be populated this way. + +### Deployment order matters + +`assemble` writes `provider::external`. A zfin that does not know that enum +variant treats the meta file as a cache **miss**, takes the cold-start path, gets +a unanimous 404, and writes a negative-cache marker **over +`candles_daily.srf`** -- the history is then gone. + +So deploy a zfin carrying the variant to both the client and ZFIN_SERVER *before* +populating a cache they read. That is an operational precondition, not something +this program checks: the enum is verified at compile time against the zfin this +links against, which says nothing about the zfin on the consuming side. + +There is deliberately no default for `--out`, so nothing is written anywhere by +accident. + + +### Steady state + +`record` is idempotent per navDate, so a daily job is safe -- the feed republishes +the same value all weekend. Suggested cron/Cronicle shape: + +``` +zfin-vestwell record && zfin-vestwell assemble --out "$ZFIN_CACHE_DIR" +``` + +`assemble` also fetches the feed itself, so the cache is current even on the run +where `record` had nothing new to add. Observed values always beat reconstructed +ones for the same date, and `assemble` reports any date where they disagreed by +more than 0.1% -- a sustained count there means the model has drifted. + ## Limitations - **The series is derived, not observed**, except at the nine anchors. Every row @@ -153,6 +217,7 @@ is not identifiable over a short gap, so per-gap figures are noisy by nature. |------------------------------------|---------------------------------------------------------------| | `data/anchors.srf` | The nine observed values, with provenance. **Irreplaceable.** | | `data/model.srf` | Fund weights per era, prefilled to 2047. Generated. | +| `data/observed.srf` | Append-only log of values read from the feed. Grows daily. | | `data/recorded.srf` | Hand-typed values, used only to verify. | | `data/ORCBI.srf`, `data/ORC42.srf` | Generated. Daily unit values, per-row provenance. | @@ -208,7 +273,7 @@ library. ``` zig build # build the exe -- run this, see below -zig build test # 55 tests +zig build test # 87 tests zig build coverage -Dcoverage-threshold=80 # coverage floor zlint --deny-warnings --fix zig fmt @@ -217,21 +282,77 @@ zig fmt Reading the underlying prices requires zfin's cache to be populated for VSMPX, VTPSX, VBMPX, VIPIX and VTIFX. `zfin quote ` fills it. +### Depending on zfin + +`src/cache_files.zig` imports zfin and uses `zfin.Candle`, `zfin.cache.Store` and +`Store.cacheCandles` rather than re-declaring those types and serializing them +here. Two reasons, both about the destructive path: + +- **`provider` is `Store.CandleProvider.external`, a compile-checked enum value**, + not the string `"external"`. A rename upstream is now a build error instead of a + cache file that destroys itself on next read. +- **zfin's own writer produces the bytes.** The directive block, field separators, + type tags, trailing newline and atomic rename are all its code, so the format + cannot drift from what its reader expects. + +What zfin does not do is validate the series or report a write failure -- +`cacheCandles` returns void and logs. Both gaps are covered: bars are checked +before the call, and afterwards the result is read back through zfin's own +`readCandleMeta`, which proves the bytes parse as well as exist. + +The `zfin` library module is lean -- `srf`, `zeit`, `build_info` -- so none of the +CLI/TUI dependency tree comes with it. + +This checks the producing side only. An *older* zfin reading the cache is still a +hazard, and the installed client and the deployed ZFIN_SERVER are upgraded +separately from this repo -- see "Deployment order matters" above. + +### Writing SRF + +Every file this project writes goes through `srf.fmt` from the `srf` library, not +through format strings. That is a correctness matter, not style: SRF +length-prefixes a string value containing a comma (`key:24:some, value`) because a +bare comma reads as a field separator. Hand-formatted records silently lost fields +whenever a value contained one -- which is how a comment written into an `evidence` +field first broke `anchors.srf`. + +The library also owns the `#!srfv1` line, the `#!expires=`/`#!created=` +directives, the `:num:` type tags and the per-record trailing newline. For +`candles_meta.srf` those are exactly the things whose malformation makes zfin +treat the file as a cache miss and overwrite `candles_daily.srf`, so they are +much better handled by the library than by this project. + +`srf.FormatOptions` has no precision control and Zig's `{d}` is +shortest-round-trip, which would spend 17 significant digits on a reconstructed +value. `src/srf_num.zig` solves that with SRF's own extension point: a one-field +wrapper whose `srfFormat` method renders fixed decimals. It is the only place in +the project that writes SRF field syntax by hand, and it does so because that is +the interface the library asks for. + +`tools/gen_model.py` is the exception -- Python has no binding for the library -- +so instead of emulating the escaping it refuses to emit any value that would need +it. + ### Coverage, and what is deliberately not covered -Currently 87.86%. The split is not uniform, on purpose: +Currently 88.84%. The split is not uniform, on purpose: -| Module | Coverage | -|------------------------------------------------|----------| -| recon, verify, data, series, civil | 100% | -| candles (`parse` tested, the file read is not) | 95% | -| feed (both parsers tested, `fetchBody` is not) | 82% | -| main (CLI wiring and file I/O) | 26% | +| Module | Coverage | +|--------------------------------------------------------------------|----------| +| recon, verify, data, series, assemble, cache_files, srf_num, civil | 100% | +| candles (`parse` tested, the file read is not) | 95% | +| feed (both parsers tested, `fetchBody` is not) | 82% | +| main (CLI wiring; `assemble` covered end to end) | 38% | -Every line of *logic* is tested. What is not tested is the I/O boundary: reading -files, the one HTTP call, and CLI plumbing. That end of the program is verified by -running it against real data and checking `verify` exits zero, which is stronger -evidence than a fixture test would give. +Every line of *logic* is tested. What is not tested is mostly the I/O boundary: +file reads, the one HTTP call, and CLI plumbing. That end of the program is also +verified by running it against real data and checking `verify` exits zero. + +`cmdAssemble` is the exception and does have an integration test, against a +temp-directory fixture. It is the one function that writes into a cache nothing +else can rebuild, and the test reads the result back with `zfin.cache.Store` +rather than with our own parser -- proving the bytes parse, not merely that they +exist. `src/main.zig` calls `std.testing.refAllDecls` so that `main` and the command functions compile in the test binary. Without it they are dead-code eliminated, @@ -241,5 +362,5 @@ coverage but makes the untested surface visible in the number rather than hiding it. **Always run `zig build`, not just `zig build test`.** The pre-commit hook runs both. -A temp-directory integration harness for `cmdReconstruct` and `cmdVerify ---offline` would close most of the remaining gap and is a reasonable follow-up. +A similar harness for `cmdReconstruct`, `cmdVerify` and `cmdRecord` would close +most of the remaining gap and is a reasonable follow-up. diff --git a/build.zig b/build.zig index cac95a4..6d6bc49 100644 --- a/build.zig +++ b/build.zig @@ -11,8 +11,14 @@ pub fn build(b: *std.Build) void { }); const srf_mod = srf_dep.module("srf"); + const zfin_dep = b.dependency("zfin", .{ + .target = target, + .optimize = optimize, + }); + const imports: []const std.Build.Module.Import = &.{ .{ .name = "srf", .module = srf_mod }, + .{ .name = "zfin", .module = zfin_dep.module("zfin") }, }; const exe_mod = b.createModule(.{ diff --git a/build.zig.zon b/build.zig.zon index c9b8c6c..a994d1d 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -8,6 +8,18 @@ .url = "git+https://git.lerch.org/lobo/srf#ea2c35825d652691e6a22526d76e2a06f61d70a5", .hash = "srf-0.0.0-qZj578QeAgCDjih2ii5soqz02fj3g8OhKwUkFh4ReK56", }, + // zfin, for `Candle`, `cache.Store` and `CandleProvider`. Depending on it + // rather than re-declaring those types means an upstream field or enum + // change breaks this build instead of silently producing a cache file + // zfin cannot parse -- and a parse failure is the path that overwrites + // candles_daily.srf with a negative-cache marker. + // + // The `zfin` library module is lean: srf, zeit and build_info only. None + // of the CLI/TUI dependency tree comes with it. + .zfin = .{ + .url = "git+https://git.lerch.org/lobo/zfin#df204a3eb21902b69fba154361e0213536570f46", + .hash = "zfin-0.0.0-J-B21hxqXQAXJig3PDBS-7cP9GkeDIXjkNanMyAmY3hc", + }, }, .paths = .{ "build.zig", diff --git a/data/observed.srf b/data/observed.srf new file mode 100644 index 0000000..d6fdfa9 --- /dev/null +++ b/data/observed.srf @@ -0,0 +1,3 @@ +#!srfv1 +symbol::ORCBI,date::2026-08-27,unit_value:num:19.092721,recorded_at::2026-08-28 +symbol::ORC42,date::2026-08-27,unit_value:num:16.808106,recorded_at::2026-08-28 diff --git a/src/assemble.zig b/src/assemble.zig new file mode 100644 index 0000000..09dd8d1 --- /dev/null +++ b/src/assemble.zig @@ -0,0 +1,344 @@ +//! Merge the reconstructed series with observed feed values into the bar series +//! zfin's cache wants. +//! +//! Two sources, and a clear precedence: an observed value always beats a +//! reconstructed one for the same date. Reconstruction exists only to cover +//! dates nobody recorded, so as the observation log grows the derived portion +//! shrinks and eventually stops mattering for anything recent. + +const std = @import("std"); +const cache_files = @import("cache_files.zig"); +const civil = @import("civil.zig"); +const data = @import("data.zig"); +const recon = @import("recon.zig"); + +pub const Counts = struct { + total: usize = 0, + /// Dates taken from the observation log or the live feed. + observed: usize = 0, + /// Dates that exist only in the reconstruction. + reconstructed: usize = 0, + /// Dates present in both, where the observed value differed enough to be + /// worth reporting. A large number here means the reconstruction has drifted + /// or the model is wrong. + corrected: usize = 0, + /// Largest relative disagreement seen on an overlapping date. + worst_correction: f64 = 0, +}; + +pub const Result = struct { + bars: []const cache_files.Bar, + counts: Counts, +}; + +/// Relative difference above which an overlap counts as a correction worth +/// reporting. The reconstruction is cross-validated to under 0.08%, and observed +/// anchors are reproduced exactly, so anything past 0.1% is a signal rather than +/// rounding. +const correction_threshold = 0.001; + +/// Where a candidate value came from. The numeric order IS the precedence: an +/// observation is a published fact and always beats a derived value for the same +/// date. +const Rank = enum(u8) { observed = 0, derived = 1 }; + +/// One candidate value for one date, before precedence is resolved. +const Entry = struct { + date: []const u8, + unit_value: f64, + rank: Rank, + /// Input position, used only to break a tie between two entries of equal + /// date and rank. The log is append-only, so later means fresher. + seq: usize, +}; + +/// Sort so that, within a date, the winner comes first. +/// +/// Date ascending because `candles_daily.srf` must be ascending anyway; then rank, +/// so an observation precedes a derived value; then sequence DESCENDING, so the +/// most recently appended of two equal-ranked entries wins. +fn entryLessThan(_: void, a: Entry, b: Entry) bool { + if (!std.mem.eql(u8, a.date, b.date)) return civil.lessThan(a.date, b.date); + if (a.rank != b.rank) return @intFromEnum(a.rank) < @intFromEnum(b.rank); + return a.seq > b.seq; +} + +/// Build the merged bar series for one symbol. +/// +/// Collect every candidate into one array, sort it so the winner of each date +/// leads its group, then scan the groups once. Precedence is expressed as a sort +/// key rather than as the order in which hash-map writes happen to overwrite each +/// other, and a correction becomes a local comparison between two adjacent +/// entries instead of state carried across a loop. +/// +/// `extra` carries a value fetched live this run, not yet in the log. It is +/// appended last, so it outranks a logged observation for the same date -- they +/// should agree, and if they do not, the fresher read wins. +pub fn merge( + arena: std.mem.Allocator, + symbol: []const u8, + series_points: []const recon.Point, + observations: []const data.Observation, + extra: ?data.Observation, +) !Result { + var entries: std.ArrayList(Entry) = .empty; + for (series_points) |p| { + try entries.append(arena, .{ + .date = p.date, + .unit_value = p.unit_value, + .rank = .derived, + .seq = entries.items.len, + }); + } + for (observations) |o| { + if (!std.mem.eql(u8, o.symbol, symbol)) continue; + try entries.append(arena, .{ + .date = o.date, + .unit_value = o.unit_value, + .rank = .observed, + .seq = entries.items.len, + }); + } + if (extra) |o| { + if (std.mem.eql(u8, o.symbol, symbol)) { + try entries.append(arena, .{ + .date = o.date, + .unit_value = o.unit_value, + .rank = .observed, + .seq = entries.items.len, + }); + } + } + + const sorted = try entries.toOwnedSlice(arena); + std.mem.sort(Entry, sorted, {}, entryLessThan); + + var bars: std.ArrayList(cache_files.Bar) = .empty; + var counts: Counts = .{}; + + var i: usize = 0; + while (i < sorted.len) { + const start = i; + while (i < sorted.len and std.mem.eql(u8, sorted[i].date, sorted[start].date)) i += 1; + const group = sorted[start..i]; + const winner = group[0]; + + // An observation that displaced a derived value for the same date is + // worth reporting when they disagree materially: a rising count there + // means the model has drifted. + if (winner.rank == .observed) { + for (group[1..]) |e| { + if (e.rank != .derived) continue; + const rel = @abs(winner.unit_value - e.unit_value) / winner.unit_value; + if (rel > correction_threshold) { + counts.corrected += 1; + counts.worst_correction = @max(counts.worst_correction, rel); + } + break; + } + counts.observed += 1; + } else { + counts.reconstructed += 1; + } + + try bars.append(arena, .{ .date = winner.date, .unit_value = winner.unit_value }); + } + + counts.total = bars.items.len; + return .{ .bars = try bars.toOwnedSlice(arena), .counts = counts }; +} + +const testing = std.testing; + +fn pt(date: []const u8, v: f64) recon.Point { + return .{ .date = date, .unit_value = v, .source = .reconstructed }; +} + +fn ob(symbol: []const u8, date: []const u8, v: f64) data.Observation { + return .{ .symbol = symbol, .date = date, .unit_value = v, .recorded_at = "2026-08-28" }; +} + +test "merge prefers an observed value over a reconstructed one" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + const r = try merge( + arena, + "ORC42", + &.{ pt("2026-08-26", 16.70), pt("2026-08-27", 16.75) }, + &.{ob("ORC42", "2026-08-27", 16.808106)}, + null, + ); + + try testing.expectEqual(@as(usize, 2), r.bars.len); + try testing.expectApproxEqAbs(@as(f64, 16.70), r.bars[0].unit_value, 1e-9); + try testing.expectApproxEqAbs(@as(f64, 16.808106), r.bars[1].unit_value, 1e-9); + try testing.expectEqual(@as(usize, 1), r.counts.observed); + try testing.expectEqual(@as(usize, 1), r.counts.reconstructed); + // 16.75 -> 16.808106 is about 0.35%, well past the reporting threshold. + try testing.expectEqual(@as(usize, 1), r.counts.corrected); + try testing.expect(r.counts.worst_correction > 0.003); +} + +test "merge extends the series past the reconstruction" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + // This is the steady state: the committed series ends where the backfill + // stopped, and the log carries everything since. + const r = try merge( + arena, + "ORC42", + &.{pt("2026-08-27", 16.808106)}, + &.{ ob("ORC42", "2026-08-28", 16.9), ob("ORC42", "2026-08-31", 17.0) }, + null, + ); + try testing.expectEqual(@as(usize, 3), r.bars.len); + try testing.expectEqualStrings("2026-08-27", r.bars[0].date); + try testing.expectEqualStrings("2026-08-31", r.bars[2].date); + try testing.expectEqual(@as(usize, 2), r.counts.observed); + try testing.expectEqual(@as(usize, 0), r.counts.corrected); +} + +test "merge treats a live fetch exactly like a logged observation" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + const r = try merge( + arena, + "ORC42", + &.{pt("2026-08-27", 16.80)}, + &.{}, + ob("ORC42", "2026-08-28", 16.9), + ); + try testing.expectEqual(@as(usize, 2), r.bars.len); + try testing.expectApproxEqAbs(@as(f64, 16.9), r.bars[1].unit_value, 1e-9); + try testing.expectEqual(@as(usize, 1), r.counts.observed); +} + +test "merge ignores other symbols in a shared log" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + const r = try merge( + arena, + "ORC42", + &.{pt("2026-08-27", 16.80)}, + &.{ ob("ORCBI", "2026-08-28", 19.1), ob("ORC42", "2026-08-28", 16.9) }, + ob("ORCBI", "2026-08-31", 19.2), + ); + try testing.expectEqual(@as(usize, 2), r.bars.len); + try testing.expectApproxEqAbs(@as(f64, 16.9), r.bars[1].unit_value, 1e-9); +} + +test "merge output is always ascending regardless of log order" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + // An append-only log can be out of order if a backfill was pasted in. + const r = try merge( + arena, + "X", + &.{pt("2026-01-05", 1.0)}, + &.{ ob("X", "2026-03-01", 3.0), ob("X", "2026-01-02", 0.5), ob("X", "2026-02-01", 2.0) }, + null, + ); + try testing.expectEqual(@as(usize, 4), r.bars.len); + for (r.bars[1..], 0..) |b, i| { + try testing.expect(civil.lessThan(r.bars[i].date, b.date)); + } + try testing.expectEqualStrings("2026-01-02", r.bars[0].date); + try testing.expectEqualStrings("2026-03-01", r.bars[3].date); +} + +test "merge keeps the last value when a log repeats a date" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + // `record` is idempotent, but a hand edit could still duplicate a date. + // Last wins, and the series must not gain a duplicate bar -- zfin's readers + // binary-search it. + const r = try merge( + arena, + "X", + &.{}, + &.{ ob("X", "2026-01-02", 1.0), ob("X", "2026-01-02", 1.5) }, + null, + ); + try testing.expectEqual(@as(usize, 1), r.bars.len); + try testing.expectApproxEqAbs(@as(f64, 1.5), r.bars[0].unit_value, 1e-9); +} + +test "merge does not flag a rounding-scale overlap as a correction" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + // The reconstruction is cross-validated to under 0.08%, so an overlap at + // that scale is expected and must stay quiet or the report becomes noise. + const r = try merge( + arena, + "X", + &.{pt("2026-01-02", 16.8000)}, + &.{ob("X", "2026-01-02", 16.8080)}, + null, + ); + try testing.expectEqual(@as(usize, 0), r.counts.corrected); + try testing.expectApproxEqAbs(@as(f64, 16.8080), r.bars[0].unit_value, 1e-9); +} + +test "merge on an empty reconstruction yields only observations" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + const r = try merge(arena, "X", &.{}, &.{ob("X", "2026-01-02", 1.0)}, null); + try testing.expectEqual(@as(usize, 1), r.counts.total); + try testing.expectEqual(@as(usize, 1), r.counts.observed); + try testing.expectEqual(@as(usize, 0), r.counts.reconstructed); +} + +test "a live fetch outranks a logged observation for the same date" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + // They should agree. If they do not, the value read this run is the fresher + // one, and the sort makes that precedence explicit rather than incidental. + const r = try merge( + arena, + "X", + &.{pt("2026-01-02", 1.0)}, + &.{ob("X", "2026-01-02", 2.0)}, + ob("X", "2026-01-02", 3.0), + ); + try testing.expectEqual(@as(usize, 1), r.bars.len); + try testing.expectApproxEqAbs(@as(f64, 3.0), r.bars[0].unit_value, 1e-12); + try testing.expectEqual(@as(usize, 1), r.counts.observed); +} + +test "entryLessThan orders by date, then precedence, then recency" { + const a: Entry = .{ .date = "2026-01-01", .unit_value = 1, .rank = .derived, .seq = 0 }; + const b: Entry = .{ .date = "2026-01-02", .unit_value = 1, .rank = .observed, .seq = 1 }; + // Date dominates, even though b outranks a. + try testing.expect(entryLessThan({}, a, b)); + try testing.expect(!entryLessThan({}, b, a)); + + // Same date: an observation leads. + const derived: Entry = .{ .date = "2026-01-01", .unit_value = 1, .rank = .derived, .seq = 9 }; + const observed: Entry = .{ .date = "2026-01-01", .unit_value = 1, .rank = .observed, .seq = 0 }; + try testing.expect(entryLessThan({}, observed, derived)); + try testing.expect(!entryLessThan({}, derived, observed)); + + // Same date and rank: the later sequence leads, so last-appended wins. + const early: Entry = .{ .date = "2026-01-01", .unit_value = 1, .rank = .observed, .seq = 1 }; + const late: Entry = .{ .date = "2026-01-01", .unit_value = 1, .rank = .observed, .seq = 2 }; + try testing.expect(entryLessThan({}, late, early)); + try testing.expect(!entryLessThan({}, early, late)); +} diff --git a/src/cache_files.zig b/src/cache_files.zig new file mode 100644 index 0000000..ebc366f --- /dev/null +++ b/src/cache_files.zig @@ -0,0 +1,282 @@ +//! Populate zfin's candle cache for the plan portfolios. +//! +//! ## Why this delegates to zfin instead of writing SRF itself +//! +//! zfin treats a `candles_meta.srf` it cannot parse as a cache MISS, not as an +//! error. A miss sends `getCandles` down the cold-start path, which asks every +//! provider for the symbol, gets a unanimous 404 (nothing carries a 529 plan's +//! internal unit values), and writes a negative-cache marker -- 23 bytes -- OVER +//! `candles_daily.srf`. The history is then gone until someone notices. +//! +//! Re-declaring `Candle` and `CandleMeta` here to hand-serialize them made that +//! outcome one upstream field addition away, silently. So this module depends on +//! zfin and uses `Store.cacheCandles`: the writer, the directive block, the field +//! separators, the type tags, the trailing newline and the atomic rename are all +//! zfin's own code, and the format cannot drift from what zfin reads. +//! +//! Two consequences worth naming: +//! +//! * `provider` is `Store.CandleProvider.external`, a compile-checked enum +//! value rather than the string `"external"`. A rename or removal upstream is +//! now a build error instead of a cache file that destroys itself. +//! * `cacheCandles` sets `adj_basis` to the newest bar's date. That is safe +//! here: with no dividends or splits cached, `newestCorporateAction` returns +//! null and `adjustmentBasisStale` short-circuits false regardless. +//! +//! What zfin does NOT do is validate the series, and it reports failure only to +//! its log. Both gaps are covered here: bars are checked before the call, and the +//! result is read back through zfin's own reader afterwards. +//! +//! ## Bar shape +//! +//! A unit value is a single number per day, so `open`, `high`, `low`, `close` and +//! `adj_close` all take it and `volume` is 0. That is exactly what zfin already +//! stores for mutual funds, so nothing downstream sees anything unusual. +//! `adj_close == close` is correct rather than lazy: these portfolios never +//! distribute, so there is nothing to adjust for. + +const std = @import("std"); +const zfin = @import("zfin"); +const civil = @import("civil.zig"); + +const Store = zfin.cache.Store; + +/// Provider tag recorded in `candles_meta.srf`. +/// +/// `external` means "produced and managed outside zfin". +/// +/// This is an enum value, not a string, so the variant is checked at compile +/// time against the zfin this links against. What that cannot check is the zfin +/// on the *consuming* side: an older client or server reading a cache tagged +/// `external` fails to parse the meta, treats the symbol as a cache miss, and +/// overwrites `candles_daily.srf` with a negative-cache marker. Deploying the +/// client and server before populating a cache is an operational precondition, +/// not something this program can verify. +pub const provider: Store.CandleProvider = .external; + +/// Hour (UTC) at which the cache is considered possibly-stale. +/// +/// The plan republishes its feed around 12:00 UTC with the prior business day's +/// value, so nothing new can appear before then. Expiring just after that means +/// a client asks the server roughly once a day and otherwise serves locally. +/// +/// Deliberately not `zfin.market.nextCandleExpiry`: that models exchange hours, +/// 16:55 ET for an equity or 03:25 ET for a mutual fund, and this is neither. The +/// binding schedule is one recordkeeper's publishing job. +/// +/// A weekend or holiday costs one redundant sync, which is a read from the +/// server's disk and never a provider call, so modelling the market calendar here +/// would buy nothing. +const expiry_hour_utc = 13; + +pub const Bar = struct { + date: []const u8, + unit_value: f64, +}; + +/// Convert to zfin's own `Candle`, rejecting anything that would corrupt a cache. +/// +/// zfin will happily serialize a descending or duplicated series, and its readers +/// binary-search the result, so an unsorted series is silently wrong rather than +/// loudly broken. Checking here is the only place it gets checked. +fn toCandles(arena: std.mem.Allocator, bars: []const Bar) ![]zfin.Candle { + if (bars.len == 0) return error.NoBars; + const out = try arena.alloc(zfin.Candle, bars.len); + for (bars, 0..) |b, i| { + if (!civil.isValidIso(b.date)) return error.InvalidBarDate; + if (!std.math.isFinite(b.unit_value) or !(b.unit_value > 0)) return error.InvalidBarValue; + if (i > 0 and !civil.lessThan(bars[i - 1].date, b.date)) return error.BarsNotAscending; + const date = zfin.Date.parse(b.date) catch return error.InvalidBarDate; + out[i] = .{ + .date = date, + .open = b.unit_value, + .high = b.unit_value, + .low = b.unit_value, + .close = b.unit_value, + .adj_close = b.unit_value, + .volume = 0, + }; + } + return out; +} + +/// The next `expiry_hour_utc` strictly after `now_s`, as a unix timestamp. +fn nextExpiry(now_s: i64) i64 { + const day = @divFloor(now_s, std.time.s_per_day); + const today_at = day * std.time.s_per_day + expiry_hour_utc * std.time.s_per_hour; + return if (today_at > now_s) today_at else today_at + std.time.s_per_day; +} + +/// Write `candles_daily.srf` and `candles_meta.srf` for one symbol under +/// `out_dir`, then verify the result by reading it back. +/// +/// `cacheCandles` returns void and logs its failures, which is reasonable for a +/// cache zfin can always refill from a provider but not for one nothing else can +/// rebuild. The read-back closes that: it goes through zfin's own +/// `readCandleMeta`, so it proves the bytes parse as well as exist -- and a meta +/// file that does not parse is precisely what triggers the destructive path. +pub fn writeSymbol( + arena: std.mem.Allocator, + io: std.Io, + out_dir: []const u8, + symbol: []const u8, + bars: []const Bar, + now_s: i64, +) !void { + const candles = try toCandles(arena, bars); + + var store: Store = .init(io, arena, out_dir); + store.cacheCandles(symbol, candles, .{ .provider = provider }, nextExpiry(now_s)); + + const newest = candles[candles.len - 1]; + const read = store.readCandleMeta(symbol) orelse return error.CacheWriteUnreadable; + if (!read.meta.last_date.eql(newest.date)) return error.CacheWriteMismatch; + if (read.meta.last_close != newest.close) return error.CacheWriteMismatch; + if (read.meta.provider != provider) return error.CacheWriteMismatch; +} + +const testing = std.testing; + +/// Path to a name inside a `std.testing.tmpDir`, relative to cwd. +/// +/// `Io.Dir` has no `realpath` in 0.16, and the functions under test take paths +/// because that is what the CLI hands them. `std.testing.tmpDir` documents its +/// own location as `.zig-cache/tmp/`, so reconstruct it. +fn tmpPath(arena: std.mem.Allocator, tmp: *const std.testing.TmpDir, name: []const u8) ![]const u8 { + return std.fs.path.join(arena, &.{ ".zig-cache", "tmp", &tmp.sub_path, name }); +} + +test "toCandles fills every OHLC field with the unit value" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + const candles = try toCandles(arena, &.{ + .{ .date = "2026-08-26", .unit_value = 16.743025 }, + .{ .date = "2026-08-27", .unit_value = 16.808106 }, + }); + + try testing.expectEqual(@as(usize, 2), candles.len); + const c = candles[1]; + try testing.expectEqualStrings("2026-08-27", try std.fmt.allocPrint(arena, "{f}", .{c.date})); + try testing.expectApproxEqAbs(@as(f64, 16.808106), c.open, 1e-12); + try testing.expectApproxEqAbs(@as(f64, 16.808106), c.high, 1e-12); + try testing.expectApproxEqAbs(@as(f64, 16.808106), c.low, 1e-12); + try testing.expectApproxEqAbs(@as(f64, 16.808106), c.close, 1e-12); + // adj_close must equal close: these portfolios never distribute, and a zero + // here is zfin's "unusable" sentinel that analytics discards. + try testing.expectApproxEqAbs(@as(f64, 16.808106), c.adj_close, 1e-12); + try testing.expectEqual(@as(u64, 0), c.volume); +} + +test "toCandles refuses input that would corrupt a cache silently" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + try testing.expectError(error.NoBars, toCandles(arena, &.{})); + try testing.expectError(error.InvalidBarDate, toCandles(arena, &.{ + .{ .date = "2026-8-27", .unit_value = 1.0 }, + })); + try testing.expectError(error.InvalidBarValue, toCandles(arena, &.{ + .{ .date = "2026-08-27", .unit_value = 0 }, + })); + try testing.expectError(error.InvalidBarValue, toCandles(arena, &.{ + .{ .date = "2026-08-27", .unit_value = std.math.nan(f64) }, + })); + // Descending, and duplicated: zfin's readers binary-search, so both are fatal. + try testing.expectError(error.BarsNotAscending, toCandles(arena, &.{ + .{ .date = "2026-08-27", .unit_value = 1.0 }, + .{ .date = "2026-08-26", .unit_value = 1.0 }, + })); + try testing.expectError(error.BarsNotAscending, toCandles(arena, &.{ + .{ .date = "2026-08-27", .unit_value = 1.0 }, + .{ .date = "2026-08-27", .unit_value = 1.0 }, + })); +} + +test "nextExpiry lands on the next 13:00 UTC" { + const day: i64 = std.time.s_per_day; + const h: i64 = std.time.s_per_hour; + + // Epoch day, 00:00 -> same day 13:00. + try testing.expectEqual(13 * h, nextExpiry(0)); + // 12:59 -> still today. + try testing.expectEqual(13 * h, nextExpiry(12 * h + 59 * 60)); + // Exactly 13:00 -> must move on, or the file would already be expired. + try testing.expectEqual(13 * h + day, nextExpiry(13 * h)); + // 13:01 -> tomorrow. + try testing.expectEqual(13 * h + day, nextExpiry(13 * h + 60)); + // Always strictly in the future. + for ([_]i64{ 0, 1, 12345, 1787932800, 1787932800 + 7 * day }) |t| { + try testing.expect(nextExpiry(t) > t); + } +} + +test "writeSymbol round-trips through zfin's own reader" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + var threaded: std.Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + const base = try tmpPath(arena, &tmp, "."); + const now: i64 = 1787932800; + + try writeSymbol(arena, io, base, "ORC42", &.{ + .{ .date = "2026-08-26", .unit_value = 16.743025 }, + .{ .date = "2026-08-27", .unit_value = 16.808106 }, + }, now); + + // The strongest available check: read it with zfin, not with our own parser. + var store: Store = .init(io, arena, base); + const read = store.readCandleMeta("ORC42").?; + try testing.expectEqual(Store.CandleProvider.external, read.meta.provider); + try testing.expectApproxEqAbs(@as(f64, 16.808106), read.meta.last_close, 1e-9); + try testing.expectEqualStrings("2026-08-27", try std.fmt.allocPrint(arena, "{f}", .{read.meta.last_date})); + try testing.expectEqual(@as(?f64, 16.808106), store.readLastClose("ORC42")); + + // And the bars themselves come back through zfin. + const back = store.read(arena, zfin.Candle, "ORC42", null, .any).?; + try testing.expectEqual(@as(usize, 2), back.data.len); + try testing.expectApproxEqAbs(@as(f64, 16.743025), back.data[0].close, 1e-9); + try testing.expectApproxEqAbs(@as(f64, 16.808106), back.data[1].adj_close, 1e-9); + + // Rerunning overwrites rather than appending or failing on the existing dir. + try writeSymbol(arena, io, base, "ORC42", &.{ + .{ .date = "2026-08-28", .unit_value = 16.9 }, + }, now); + const again = store.read(arena, zfin.Candle, "ORC42", null, .any).?; + try testing.expectEqual(@as(usize, 1), again.data.len); +} + +test "writeSymbol reports a cache it could not write" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + var threaded: std.Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + + // A regular file standing where the cache directory should be: every path + // under it fails with NotDir, promptly and portably. + const blocker = try tmpPath(arena, &tmp, "not-a-dir"); + try std.Io.Dir.cwd().writeFile(io, .{ .sub_path = blocker, .data = "x" }); + + // `cacheCandles` only logs a write failure, so without the read-back this + // would silently "succeed" against a location it never wrote to. + try testing.expectError(error.CacheWriteUnreadable, writeSymbol( + arena, + io, + blocker, + "ORC42", + &.{.{ .date = "2026-08-27", .unit_value = 16.8 }}, + 1787932800, + )); +} diff --git a/src/data.zig b/src/data.zig index 2f328e6..6e84c97 100644 --- a/src/data.zig +++ b/src/data.zig @@ -8,6 +8,7 @@ const std = @import("std"); const srf = @import("srf"); +const srf_num = @import("srf_num.zig"); const srf_opts = @import("srf_opts.zig"); const civil = @import("civil.zig"); @@ -282,6 +283,104 @@ pub fn parseRecorded(arena: std.mem.Allocator, data: []const u8) ![]const Record return out.toOwnedSlice(arena); } +// ── observed.srf ───────────────────────────────────────────────────── + +const ObservationRow = struct { + symbol: []const u8, + date: []const u8, + unit_value: f64, + recorded_at: []const u8 = "", +}; + +/// One unit value read straight from the plan's feed. +/// +/// These are observations, not derivations, and they accumulate forever. Once +/// this log is dense the reconstruction stops mattering for recent dates: `assemble` +/// overlays observations on top of the reconstructed series, so an observed value +/// always wins. +pub const Observation = struct { + symbol: []const u8, + /// The feed's own navDate, which is the prior business day. + date: []const u8, + unit_value: f64, + /// When the fetch happened. Provenance only; nothing keys off it. + recorded_at: []const u8, +}; + +pub fn parseObserved(arena: std.mem.Allocator, data: []const u8) ![]const Observation { + var reader = std.Io.Reader.fixed(data); + var it = srf.iterator(&reader, arena, .{ .parse_allocator = .none }) catch + return error.InvalidObservedFile; + defer it.deinit(); + + var out: std.ArrayList(Observation) = .empty; + while (try it.next()) |fields| { + const row = try fields.to(ObservationRow, srf_opts.user_edited); + if (!civil.isValidIso(row.date)) return error.InvalidObservedDate; + if (!(row.unit_value > 0)) return error.InvalidObservedValue; + try out.append(arena, .{ + .symbol = try arena.dupe(u8, row.symbol), + .date = try arena.dupe(u8, row.date), + .unit_value = row.unit_value, + .recorded_at = try arena.dupe(u8, row.recorded_at), + }); + } + return out.toOwnedSlice(arena); +} + +/// Whether this log already holds a value for `symbol` on `date`. +/// +/// Makes `record` idempotent: the feed republishes the same navDate all weekend +/// and across holidays, so a daily job re-reads the same value repeatedly and +/// must not append it repeatedly. +pub fn hasObservation(obs: []const Observation, symbol: []const u8, date: []const u8) bool { + for (obs) |o| { + if (std.mem.eql(u8, o.symbol, symbol) and std.mem.eql(u8, o.date, date)) return true; + } + return false; +} + +/// The newest observation for `symbol`, or null when there is none. +pub fn newestObservation(obs: []const Observation, symbol: []const u8) ?Observation { + var best: ?Observation = null; + for (obs) |o| { + if (!std.mem.eql(u8, o.symbol, symbol)) continue; + if (best == null or civil.lessThan(best.?.date, o.date)) best = o; + } + return best; +} + +/// Write-side observation row. `unit_value` is fixed-precision so the log keeps +/// the feed's six published decimals rather than shortest-round-trip digits. +const ObservationOut = struct { + symbol: []const u8, + date: []const u8, + unit_value: srf_num.Fixed(6), + recorded_at: []const u8, +}; + +/// Render one observation as an SRF record line, newline included. +/// +/// Goes through `srf.fmt` rather than a format string. That is not tidiness: SRF +/// length-prefixes a string value containing a comma (`key:LEN:value`), because a +/// bare comma would otherwise read as a field separator. Hand-formatting this +/// line silently produced a corrupt record for any value with a comma in it -- +/// the same mistake that broke `anchors.srf` when a comment was written into an +/// `evidence` field. +pub fn formatObservation(arena: std.mem.Allocator, o: Observation) ![]const u8 { + const rows = [_]ObservationOut{.{ + .symbol = o.symbol, + .date = o.date, + .unit_value = .init(o.unit_value), + .recorded_at = o.recorded_at, + }}; + // `emit_directives = false`: this line is appended to a file that already has + // its `#!srfv1` header. + return std.fmt.allocPrint(arena, "{f}", .{ + srf.fmt(ObservationOut, &rows, .{ .emit_directives = false }), + }); +} + const testing = std.testing; test "parseModel groups rows into eras and scales percentages to fractions" { @@ -504,3 +603,110 @@ test "tickersThrough excludes funds only a future era needs" { // Before any era, nothing is needed. try testing.expectEqual(@as(usize, 0), (try m.tickersThrough(arena, "2020-01-01")).len); } + +test "parseObserved reads the feed log and rejects unusable rows" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + const obs = try parseObserved(arena, + \\#!srfv1 + \\symbol::ORCBI,date::2026-08-27,unit_value:num:19.092721,recorded_at::2026-08-28 + \\symbol::ORC42,date::2026-08-27,unit_value:num:16.808106,recorded_at::2026-08-28 + \\ + ); + try testing.expectEqual(@as(usize, 2), obs.len); + try testing.expectEqualStrings("ORCBI", obs[0].symbol); + try testing.expectEqualStrings("2026-08-27", obs[0].date); + try testing.expectEqualStrings("2026-08-28", obs[0].recorded_at); + try testing.expectApproxEqAbs(@as(f64, 16.808106), obs[1].unit_value, 1e-9); + + // recorded_at is provenance only, so an older log without it still parses. + const bare = try parseObserved(arena, "#!srfv1\nsymbol::X,date::2026-01-01,unit_value:num:1.5\n"); + try testing.expectEqual(@as(usize, 1), bare.len); + try testing.expectEqualStrings("", bare[0].recorded_at); + + // An empty log is normal on a fresh checkout. + try testing.expectEqual(@as(usize, 0), (try parseObserved(arena, "#!srfv1\n")).len); + + try testing.expectError(error.InvalidObservedDate, parseObserved(arena, "#!srfv1\nsymbol::X,date::2026-1-1,unit_value:num:1\n")); + // A zero must never enter as a price; the feed publishes 0 for closed funds. + try testing.expectError(error.InvalidObservedValue, parseObserved(arena, "#!srfv1\nsymbol::X,date::2026-01-01,unit_value:num:0\n")); +} + +test "hasObservation makes recording idempotent per symbol and date" { + const obs: []const Observation = &.{ + .{ .symbol = "ORCBI", .date = "2026-08-27", .unit_value = 19.09, .recorded_at = "" }, + .{ .symbol = "ORC42", .date = "2026-08-27", .unit_value = 16.81, .recorded_at = "" }, + }; + try testing.expect(hasObservation(obs, "ORCBI", "2026-08-27")); + try testing.expect(hasObservation(obs, "ORC42", "2026-08-27")); + // Same date, different symbol, and same symbol, different date: both misses. + try testing.expect(!hasObservation(obs, "OTHER", "2026-08-27")); + try testing.expect(!hasObservation(obs, "ORCBI", "2026-08-28")); + try testing.expect(!hasObservation(&.{}, "ORCBI", "2026-08-27")); +} + +test "newestObservation picks the latest date, not the last row" { + // An append-only log can be out of order if a backfill was pasted in. + const obs: []const Observation = &.{ + .{ .symbol = "X", .date = "2026-08-27", .unit_value = 3.0, .recorded_at = "" }, + .{ .symbol = "Y", .date = "2026-12-31", .unit_value = 9.0, .recorded_at = "" }, + .{ .symbol = "X", .date = "2026-08-31", .unit_value = 4.0, .recorded_at = "" }, + .{ .symbol = "X", .date = "2026-08-20", .unit_value = 1.0, .recorded_at = "" }, + }; + const n = newestObservation(obs, "X").?; + try testing.expectEqualStrings("2026-08-31", n.date); + try testing.expectApproxEqAbs(@as(f64, 4.0), n.unit_value, 1e-12); + try testing.expectEqual(@as(?Observation, null), newestObservation(obs, "MISSING")); + try testing.expectEqual(@as(?Observation, null), newestObservation(&.{}, "X")); +} + +test "formatObservation round-trips through parseObserved" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + const line = try formatObservation(arena, .{ + .symbol = "ORC42", + .date = "2026-08-27", + .unit_value = 16.808106, + .recorded_at = "2026-08-28", + }); + try testing.expect(std.mem.endsWith(u8, line, "\n")); + + const body = try std.fmt.allocPrint(arena, "#!srfv1\n{s}", .{line}); + const back = try parseObserved(arena, body); + try testing.expectEqual(@as(usize, 1), back.len); + try testing.expectEqualStrings("ORC42", back[0].symbol); + try testing.expectEqualStrings("2026-08-27", back[0].date); + // Six decimals must survive: these are the exact published values. + try testing.expectApproxEqAbs(@as(f64, 16.808106), back[0].unit_value, 1e-9); +} + +test "formatObservation survives a comma in a string value" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + // A bare comma in a value reads as a field separator, so SRF length-prefixes + // the string instead. Hand-formatting this line did not, which produced a + // record that either failed to parse or silently lost fields -- exactly how + // an `evidence` comment once broke anchors.srf. Round-tripping is the proof. + const o: Observation = .{ + .symbol = "ORC42", + .date = "2026-08-27", + .unit_value = 16.808106, + .recorded_at = "backfilled 2026-08-28, from the archive", + }; + + const line = try formatObservation(arena, o); + const body = try std.fmt.allocPrint(arena, "#!srfv1\n{s}", .{line}); + const back = try parseObserved(arena, body); + + try testing.expectEqual(@as(usize, 1), back.len); + try testing.expectEqualStrings(o.symbol, back[0].symbol); + try testing.expectEqualStrings(o.date, back[0].date); + try testing.expectApproxEqAbs(o.unit_value, back[0].unit_value, 1e-9); + try testing.expectEqualStrings(o.recorded_at, back[0].recorded_at); +} diff --git a/src/main.zig b/src/main.zig index 335688d..49bce6e 100644 --- a/src/main.zig +++ b/src/main.zig @@ -10,12 +10,16 @@ const std = @import("std"); +const assemble = @import("assemble.zig"); +const cache_files = @import("cache_files.zig"); const candles = @import("candles.zig"); const civil = @import("civil.zig"); const data = @import("data.zig"); const feed = @import("feed.zig"); const recon = @import("recon.zig"); const series = @import("series.zig"); +const srf_num = @import("srf_num.zig"); +const zfin = @import("zfin"); const verify = @import("verify.zig"); const symbols = [_][]const u8{ "ORCBI", "ORC42" }; @@ -26,29 +30,42 @@ const usage = \\Usage: \\ zfin-vestwell reconstruct [--data-dir DIR] \\ zfin-vestwell verify [--data-dir DIR] [--offline] + \\ zfin-vestwell record [--data-dir DIR] + \\ zfin-vestwell assemble --out DIR [--data-dir DIR] [--offline] \\ \\Commands: \\ reconstruct Rebuild data/.srf from data/model.srf and \\ data/anchors.srf, using the underlying funds' adjusted - \\ closes from zfin's candle cache. + \\ closes from zfin's candle cache. Run rarely, and only where + \\ that cache is populated. \\ verify Re-read the series from disk and check it: anchor exactness, \\ leave-one-out cross-validation, agreement with hand-typed \\ values, structural sanity, and a live feed cross-check. \\ Exits non-zero on any failure. + \\ record Fetch the feed and append anything new to + \\ data/observed.srf. Idempotent per navDate, so running it + \\ daily is safe. + \\ assemble Merge data/.srf with data/observed.srf and the live + \\ feed, then write candles_daily.srf + candles_meta.srf into + \\ // for zfin to serve. Needs no zfin cache. \\ \\Options: - \\ --data-dir DIR Where the SRF inputs and outputs live (default: ./data) - \\ --offline Skip the live feed cross-check in `verify` - \\ -h, --help Show this message + \\ --data-dir DIR Where the SRF inputs and outputs live (default: ./data) + \\ --out DIR Cache directory `assemble` writes into. Required; there + \\ is deliberately no default. + \\ --offline Skip the live feed + \\ -h, --help Show this message \\ ; -const Command = enum { reconstruct, verify, help }; +const Command = enum { reconstruct, verify, record, assemble, help }; const Options = struct { cmd: Command, data_dir: []const u8 = "data", offline: bool = false, + /// Target cache directory for `assemble`. No default on purpose. + out: ?[]const u8 = null, }; fn parseArgs(args: []const []const u8) !Options { @@ -58,6 +75,10 @@ fn parseArgs(args: []const []const u8) !Options { .reconstruct else if (std.mem.eql(u8, args[1], "verify")) .verify + else if (std.mem.eql(u8, args[1], "record")) + .record + else if (std.mem.eql(u8, args[1], "assemble")) + .assemble else if (std.mem.eql(u8, args[1], "-h") or std.mem.eql(u8, args[1], "--help") or std.mem.eql(u8, args[1], "help")) @@ -75,10 +96,17 @@ fn parseArgs(args: []const []const u8) !Options { i += 1; if (i >= args.len) return error.MissingArgument; o.data_dir = args[i]; + } else if (std.mem.eql(u8, a, "--out")) { + i += 1; + if (i >= args.len) return error.MissingArgument; + o.out = args[i]; } else { return error.UnknownFlag; } } + // `assemble` writes into a cache directory. Refusing to guess one is the + // point: a default would eventually write somewhere nobody intended. + if (o.cmd == .assemble and o.out == null) return error.MissingOutDir; return o; } @@ -88,6 +116,7 @@ const Loaded = struct { models: []const data.Model, anchor_sets: []const data.AnchorSet, recorded: []const data.Recorded, + observed: []const data.Observation, tickers: []const candles.Ticker, fn inputsFor(self: Loaded, symbol: []const u8) !recon.Inputs { @@ -122,6 +151,15 @@ fn load( const anchor_sets = try data.parseAnchors(arena, try readDataFile(arena, io, dir, "anchors.srf")); const recorded = try data.parseRecorded(arena, try readDataFile(arena, io, dir, "recorded.srf")); + // The observation log may not exist yet on a fresh checkout; that is not an + // error, it just means nothing has been recorded. + const observed_path = try std.fs.path.join(arena, &.{ dir, "observed.srf" }); + const observed_raw = std.Io.Dir.cwd().readFileAlloc(io, observed_path, arena, .limited(32 * 1024 * 1024)) catch |err| switch (err) { + error.FileNotFound => "#!srfv1\n", + else => return err, + }; + const observed = try data.parseObserved(arena, observed_raw); + // Union of every underlying fund across every symbol, loaded once. // // A fund that is not in zfin's cache is skipped rather than being fatal. @@ -165,6 +203,7 @@ fn load( .models = models, .anchor_sets = anchor_sets, .recorded = recorded, + .observed = observed, .tickers = try loaded.toOwnedSlice(arena), }; } @@ -183,7 +222,7 @@ fn cmdReconstruct( const points = try recon.build(arena, in); var aw: std.Io.Writer.Allocating = .init(arena); - try series.write(&aw.writer, sym, points); + try series.write(arena, &aw.writer, sym, points); const path = try std.fs.path.join(arena, &.{ o.data_dir, try std.fmt.allocPrint(arena, "{s}.srf", .{sym}) }); try std.Io.Dir.cwd().writeFile(io, .{ .sub_path = path, .data = aw.written() }); @@ -265,17 +304,34 @@ fn cmdVerify( if (feed_body) |body| { const fund_id = feed.fundIdFor(sym) orelse return error.NoFundIdForSymbol; const q = try feed.quoteFrom(arena, body, fund_id); - const last = points[points.len - 1]; - const matches_date = std.mem.eql(u8, last.date, q.date); - const err = @abs(last.unit_value - q.unit_value) / q.unit_value; - if (!matches_date or err > 1e-9) { - try w.print(" live feed MISMATCH: feed {d:.6} @ {s}, series ends {d:.6} @ {s}\n", .{ - q.unit_value, q.date, last.unit_value, last.date, + + // Compare the feed against the newest value ANY source has, not just + // the committed series. Once `record` is running daily the series + // deliberately stops at the last anchor and `observed.srf` carries + // everything since, so checking the series alone would report a + // false staleness every day. + var newest_date = points[points.len - 1].date; + var newest_value = points[points.len - 1].unit_value; + var newest_from: []const u8 = "series"; + if (data.newestObservation(l.observed, sym)) |obs| { + if (civil.lessThan(newest_date, obs.date)) { + newest_date = obs.date; + newest_value = obs.unit_value; + newest_from = "observed.srf"; + } + } + + const err = @abs(newest_value - q.unit_value) / q.unit_value; + if (!std.mem.eql(u8, newest_date, q.date) or err > 1e-9) { + try w.print(" live feed STALE: feed {d:.6} @ {s}, newest local {d:.6} @ {s} ({s})\n", .{ + q.unit_value, q.date, newest_value, newest_date, newest_from, }); - try w.print(" the series is stale; add the current value to anchors.srf and re-run reconstruct\n", .{}); + try w.print(" run `zfin-vestwell record` to capture it\n", .{}); failed = true; } else { - try w.print(" live feed matches ({d:.6} @ {s})\n", .{ q.unit_value, q.date }); + try w.print(" live feed matches ({d:.6} @ {s}, from {s})\n", .{ + q.unit_value, q.date, newest_from, + }); } } else { try w.print(" live feed skipped\n", .{}); @@ -299,6 +355,117 @@ fn cmdVerify( return 0; } +/// Fetch the feed and append anything new to the observation log. +/// +/// Idempotent by navDate: the feed republishes the same value all weekend and +/// through holidays, so a daily job re-reads it repeatedly and must not append it +/// repeatedly. Appends rather than rewrites, so the log is only ever added to. +fn cmdRecord(arena: std.mem.Allocator, w: *std.Io.Writer, io: std.Io, o: Options) !u8 { + const path = try std.fs.path.join(arena, &.{ o.data_dir, "observed.srf" }); + const existing = std.Io.Dir.cwd().readFileAlloc(io, path, arena, .limited(32 * 1024 * 1024)) catch |err| switch (err) { + error.FileNotFound => "#!srfv1\n", + else => return err, + }; + const obs = try data.parseObserved(arena, existing); + + const body = try feed.fetchBody(arena, io, feed.url); + var today_buf: [10]u8 = undefined; + const now = try civil.today(io, &today_buf); + + var appended: std.ArrayList(u8) = .empty; + var added: usize = 0; + for (symbols) |sym| { + const fund_id = feed.fundIdFor(sym) orelse return error.NoFundIdForSymbol; + const q = try feed.quoteFrom(arena, body, fund_id); + if (data.hasObservation(obs, sym, q.date)) { + try w.print("{s:<6} {s} already recorded ({d:.6})\n", .{ sym, q.date, q.unit_value }); + continue; + } + const line = try data.formatObservation(arena, .{ + .symbol = sym, + .date = q.date, + .unit_value = q.unit_value, + .recorded_at = now, + }); + try appended.appendSlice(arena, line); + added += 1; + try w.print("{s:<6} {s} recorded {d:.6}\n", .{ sym, q.date, q.unit_value }); + } + + if (added == 0) { + try w.print("\nnothing new to record\n", .{}); + return 0; + } + + // Read-modify-write rather than an append-mode handle: the file is tiny, and + // rewriting the whole thing keeps a partial write from leaving a torn record. + var full: std.ArrayList(u8) = .empty; + try full.appendSlice(arena, existing); + if (existing.len > 0 and existing[existing.len - 1] != '\n') try full.append(arena, '\n'); + try full.appendSlice(arena, appended.items); + try std.Io.Dir.cwd().writeFile(io, .{ .sub_path = path, .data = full.items }); + + try w.print("\nappended {d} observation(s) to {s}\n", .{ added, path }); + return 0; +} + +/// Merge reconstruction + observations into zfin cache files. +fn cmdAssemble( + arena: std.mem.Allocator, + w: *std.Io.Writer, + io: std.Io, + o: Options, +) !u8 { + const out_dir = o.out.?; + + const observed_path = try std.fs.path.join(arena, &.{ o.data_dir, "observed.srf" }); + const observed_raw = std.Io.Dir.cwd().readFileAlloc(io, observed_path, arena, .limited(32 * 1024 * 1024)) catch |err| switch (err) { + error.FileNotFound => "#!srfv1\n", + else => return err, + }; + const obs = try data.parseObserved(arena, observed_raw); + + const body: ?[]const u8 = if (o.offline) null else try feed.fetchBody(arena, io, feed.url); + const now_s = std.Io.Timestamp.now(io, .real).toSeconds(); + var today_buf: [10]u8 = undefined; + const now = try civil.today(io, &today_buf); + + for (symbols) |sym| { + const raw = try readDataFile(arena, io, o.data_dir, try std.fmt.allocPrint(arena, "{s}.srf", .{sym})); + const points = try series.read(arena, raw); + + var extra: ?data.Observation = null; + if (body) |b| { + const fund_id = feed.fundIdFor(sym) orelse return error.NoFundIdForSymbol; + const q = try feed.quoteFrom(arena, b, fund_id); + extra = .{ .symbol = sym, .date = q.date, .unit_value = q.unit_value, .recorded_at = now }; + } + + const merged = try assemble.merge(arena, sym, points, obs, extra); + try cache_files.writeSymbol(arena, io, out_dir, sym, merged.bars, now_s); + + const newest = merged.bars[merged.bars.len - 1]; + try w.print("{s:<6} {d:>5} bars {s} .. {s} last {d:.6} ({d} observed, {d} reconstructed)\n", .{ + sym, + merged.counts.total, + merged.bars[0].date, + newest.date, + newest.unit_value, + merged.counts.observed, + merged.counts.reconstructed, + }); + if (merged.counts.corrected > 0) { + try w.print(" {d} date(s) where an observation corrected the reconstruction, worst {d:.3}%\n", .{ + merged.counts.corrected, + merged.counts.worst_correction * 100, + }); + } + } + + try w.print("\nwrote candles_daily.srf + candles_meta.srf under {s}\n", .{out_dir}); + return 0; +} + pub fn main(init: std.process.Init) !u8 { var arena_state: std.heap.ArenaAllocator = .init(init.gpa); defer arena_state.deinit(); @@ -339,6 +506,8 @@ fn run(arena: std.mem.Allocator, w: *std.Io.Writer, io: std.Io, init: std.proces }, .reconstruct => cmdReconstruct(arena, w, io, init.environ_map, o), .verify => cmdVerify(arena, w, io, init.environ_map, o), + .record => cmdRecord(arena, w, io, o), + .assemble => cmdAssemble(arena, w, io, o), }; } @@ -392,5 +561,81 @@ test { _ = feed; _ = recon; _ = series; + _ = srf_num; _ = verify; } + +/// Path to a name inside a `std.testing.tmpDir`, relative to cwd. +/// `Io.Dir` has no `realpath` in 0.16, and `std.testing.tmpDir` documents its +/// own location, so reconstruct it. +fn tmpPath(arena: std.mem.Allocator, tmp: *const std.testing.TmpDir, name: []const u8) ![]const u8 { + return std.fs.path.join(arena, &.{ ".zig-cache", "tmp", &tmp.sub_path, name }); +} + +/// Lay down a minimal but complete data directory for the integration tests. +fn fixtureDataDir(arena: std.mem.Allocator, io: std.Io, dir: []const u8) !void { + try std.Io.Dir.cwd().createDirPath(io, dir); + const files = [_]struct { name: []const u8, body: []const u8 }{ + .{ .name = "ORCBI.srf", .body = "#!srfv1\n" ++ + "date::2026-08-26,unit_value:num:19.000000,source::reconstructed\n" ++ + "date::2026-08-27,unit_value:num:19.092721,source::anchor\n" }, + .{ .name = "ORC42.srf", .body = "#!srfv1\n" ++ + "date::2026-08-26,unit_value:num:16.743025,source::reconstructed\n" ++ + "date::2026-08-27,unit_value:num:16.808106,source::anchor\n" }, + .{ .name = "observed.srf", .body = "#!srfv1\n" ++ + "symbol::ORCBI,date::2026-08-28,unit_value:num:19.150000,recorded_at::2026-08-29\n" ++ + "symbol::ORC42,date::2026-08-28,unit_value:num:16.900000,recorded_at::2026-08-29\n" }, + }; + for (files) |f| { + try std.Io.Dir.cwd().writeFile(io, .{ + .sub_path = try std.fs.path.join(arena, &.{ dir, f.name }), + .data = f.body, + }); + } +} + +test "cmdAssemble writes a cache zfin itself can read back" { + var arena_state: std.heap.ArenaAllocator = .init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + var threaded: std.Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + const data_dir = try tmpPath(arena, &tmp, "data"); + const out_dir = try tmpPath(arena, &tmp, "cache"); + try fixtureDataDir(arena, io, data_dir); + + var aw: std.Io.Writer.Allocating = .init(arena); + const code = try cmdAssemble(arena, &aw.writer, io, .{ + .cmd = .assemble, + .data_dir = data_dir, + .out = out_dir, + .offline = true, + }); + try testing.expectEqual(@as(u8, 0), code); + try testing.expect(std.mem.indexOf(u8, aw.written(), "ORC42") != null); + + // Read the result with zfin, not with our own parser: that proves the bytes + // parse as well as exist, and a meta file that does not parse is exactly what + // sends zfin down the path that overwrites candles_daily.srf. + var store: zfin.cache.Store = .init(io, arena, out_dir); + for ([_][]const u8{ "ORCBI", "ORC42" }) |sym| { + const meta = store.readCandleMeta(sym) orelse { + std.debug.print("{s}: candles_meta.srf did not parse\n", .{sym}); + return error.TestUnexpectedResult; + }; + try testing.expectEqual(zfin.cache.Store.CandleProvider.external, meta.meta.provider); + // The observation extends the series past where the backfill stopped, so + // the meta must point at it rather than at the last reconstructed bar. + try testing.expectEqualStrings( + "2026-08-28", + try std.fmt.allocPrint(arena, "{f}", .{meta.meta.last_date}), + ); + + const bars = store.read(arena, zfin.Candle, sym, null, .any).?; + try testing.expectEqual(@as(usize, 3), bars.data.len); + } +} diff --git a/src/series.zig b/src/series.zig index a885ab4..90ce53f 100644 --- a/src/series.zig +++ b/src/series.zig @@ -11,6 +11,7 @@ const std = @import("std"); const srf = @import("srf"); +const srf_num = @import("srf_num.zig"); const srf_opts = @import("srf_opts.zig"); const recon = @import("recon.zig"); const civil = @import("civil.zig"); @@ -29,19 +30,29 @@ const Row = struct { /// so anchors round-trip exactly, and keeps reconstructed values well inside /// their own error bars -- the sixth decimal is worth about a hundredth of a /// cent on the largest position. -/// -/// Rows are formatted explicitly rather than through `srf.fmt` because SRF has -/// no decimal-precision option, and its default shortest-round-trip float -/// formatting would emit up to 17 significant digits. That is exact but makes a -/// 2000-row committed file noisy, and makes its diffs jump on the last digit -/// whenever an upstream adjusted price shifts imperceptibly. const decimals = 6; +/// Write-side row. `unit_value` is a fixed-precision wrapper rather than a plain +/// `f64` so SRF renders six decimals instead of shortest-round-trip; see +/// `srf_num`. Field order here is the column order in the file. +const OutRow = struct { + date: []const u8, + unit_value: srf_num.Fixed(decimals), + source: []const u8, +}; + pub fn write( + arena: std.mem.Allocator, w: *std.Io.Writer, symbol: []const u8, points: []const recon.Point, ) !void { + // The magic line is written here rather than by `srf.fmt` only because the + // comment block has to follow it: `#!srfv1` must be the first line of the + // file, and SRF emits its directives immediately before the first record. + // The records themselves go through the library, so field separators, + // comma-safe string escaping and the trailing newline are not this + // function's problem. try w.print( \\#!srfv1 \\# Daily unit-value series for {s}, Oregon College Savings Plan (Embark). @@ -61,15 +72,15 @@ pub fn write( \\ , .{symbol}); - // Built at comptime so `decimals` stays the single source of truth for the - // on-disk precision, rather than being duplicated in a literal. - const row_fmt = std.fmt.comptimePrint( - "date::{{s}},unit_value:num:{{d:.{d}}},source::{{s}}\n", - .{decimals}, - ); - for (points) |p| { - try w.print(row_fmt, .{ p.date, p.unit_value, @tagName(p.source) }); + const rows = try arena.alloc(OutRow, points.len); + for (points, 0..) |p, i| { + rows[i] = .{ + .date = p.date, + .unit_value = .init(p.unit_value), + .source = @tagName(p.source), + }; } + try w.print("{f}", .{srf.fmt(OutRow, rows, .{ .emit_directives = false })}); } /// Read a series file back. Used by `verify` so it checks the artifact on disk @@ -110,7 +121,7 @@ test "write then read round-trips values, dates and provenance" { var aw: std.Io.Writer.Allocating = .init(testing.allocator); defer aw.deinit(); - try write(&aw.writer, "ORCBI", points); + try write(arena, &aw.writer, "ORCBI", points); const bytes = aw.written(); try testing.expect(std.mem.startsWith(u8, bytes, "#!srfv1\n")); diff --git a/src/srf_num.zig b/src/srf_num.zig new file mode 100644 index 0000000..9afee6a --- /dev/null +++ b/src/srf_num.zig @@ -0,0 +1,81 @@ +//! Fixed-precision float rendering for SRF output. +//! +//! SRF renders a plain `f64` field as `key:num:{d}`, and Zig's `{d}` is +//! shortest-round-trip: exact, but it will spend up to 17 significant digits on a +//! value like 10.012345678901234. That is unreadable in a 2000-row committed file +//! and makes its diffs jump on the last digit whenever an upstream adjusted price +//! moves imperceptibly. +//! +//! `srf.FormatOptions` has no precision knob, but SRF does support custom +//! rendering: a struct field whose type has an `srfFormat` method controls its own +//! `key:type:value` output entirely. So a one-field wrapper gets fixed decimals +//! while leaving the library in charge of everything that actually matters -- +//! record framing, field separators, comma-safe string escaping, the directive +//! block and the trailing newline. +//! +//! Note the `:num:` in the emitted text. It is not decoration: zfin's candle +//! reader dispatches on the value type, and a numeric field written as `::` +//! silently reads back as 0 rather than failing. + +const std = @import("std"); + +/// An `f64` that renders with exactly `decimals` places after the point. +pub fn Fixed(comptime decimals: comptime_int) type { + return struct { + v: f64, + + const Self = @This(); + const value_fmt = std.fmt.comptimePrint("{{d:.{d}}}", .{decimals}); + + pub fn init(v: f64) Self { + return .{ .v = v }; + } + + /// SRF custom-serialization hook. `key` is the field name, and this is + /// responsible for the whole `key:num:value` triple. + pub fn srfFormat( + self: Self, + comptime key: []const u8, + w: *std.Io.Writer, + ) std.Io.Writer.Error!void { + try w.print(key ++ ":num:" ++ value_fmt, .{self.v}); + } + }; +} + +const testing = std.testing; + +test "Fixed renders the requested number of decimals with a num type tag" { + var aw: std.Io.Writer.Allocating = .init(testing.allocator); + defer aw.deinit(); + + try Fixed(6).init(16.808106).srfFormat("unit_value", &aw.writer); + try testing.expectEqualStrings("unit_value:num:16.808106", aw.written()); +} + +test "Fixed pads and rounds rather than emitting shortest-round-trip" { + var aw: std.Io.Writer.Allocating = .init(testing.allocator); + defer aw.deinit(); + + // A whole number still gets its decimals, so columns line up. + try Fixed(6).init(10.0).srfFormat("v", &aw.writer); + try testing.expectEqualStrings("v:num:10.000000", aw.written()); + + // And a value with more precision than requested is rounded, not truncated + // to 17 significant digits the way `{d}` would render it. + aw.clearRetainingCapacity(); + try Fixed(6).init(10.0123456789012345).srfFormat("v", &aw.writer); + try testing.expectEqualStrings("v:num:10.012346", aw.written()); +} + +test "Fixed honours other precisions" { + var aw: std.Io.Writer.Allocating = .init(testing.allocator); + defer aw.deinit(); + + try Fixed(2).init(19.092721).srfFormat("x", &aw.writer); + try testing.expectEqualStrings("x:num:19.09", aw.written()); + + aw.clearRetainingCapacity(); + try Fixed(0).init(19.6).srfFormat("x", &aw.writer); + try testing.expectEqualStrings("x:num:20", aw.written()); +} diff --git a/tools/gen_model.py b/tools/gen_model.py index cdcf72e..abba996 100644 --- a/tools/gen_model.py +++ b/tools/gen_model.py @@ -45,12 +45,32 @@ def num(v): def rows_for(symbol, era_start, basis, pairs): - return [ - f"symbol::{symbol},era_start::{era_start},basis::{basis}," - f"ticker::{t},weight:num:{num(w)}" - for t, w in pairs - if w > 0 - ] + """Render weight rows, refusing anything SRF would mis-parse. + + This script hand-formats SRF because Python has no binding for the library + the Zig side uses. That means the escaping the library would do for free is + absent here, so instead of emulating it, reject the input that would need it. + A bare comma in a value reads as a field separator; a newline ends the record. + Every value written here is an identifier, an ISO date or a number, so a hit + means the booklet layout changed under us rather than a quoting bug. + """ + out = [] + for ticker, weight in pairs: + if weight <= 0: + continue + for field, value in ( + ("symbol", symbol), + ("era_start", era_start), + ("basis", basis), + ("ticker", ticker), + ): + if any(c in value for c in ",\n\r"): + sys.exit(f"refusing to emit unescaped SRF: {field}={value!r}") + out.append( + f"symbol::{symbol},era_start::{era_start},basis::{basis}," + f"ticker::{ticker},weight:num:{num(weight)}" + ) + return out def main():