From 79600c15aef9bd411abe88dcf36286a102b5abea Mon Sep 17 00:00:00 2001 From: "Somhairle H. Marisol" Date: Tue, 22 Sep 2026 05:56:14 +0800 Subject: Add daily market refresh endpoint, snapshot-first valuation and UI refresh (3d-25) --- src/FundLab.Api/Persistence.fs | 130 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 130 insertions(+) (limited to 'src/FundLab.Api/Persistence.fs') diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs index 379c10d..09e3e0e 100644 --- a/src/FundLab.Api/Persistence.fs +++ b/src/FundLab.Api/Persistence.fs @@ -358,6 +358,18 @@ type BondPositionRecord = LastTradedAt: DateTimeOffset } +type InstrumentSnapshotRecord = + { + InstrumentCode: string + AssetClass: string + SnapshotDate: DateOnly + Price: decimal + Source: string + SourceRevision: string + SourceCollectedAt: DateTimeOffset + SourcePayloadHash: string + } + type BondTradeWriteResult = | BondTradeCreated of BondTradeRecord | BondTradeReplayed of BondTradeRecord @@ -973,6 +985,23 @@ type FundRepository(connectionString: string) = PRIMARY KEY (fund_id, instrument_code) ); + CREATE TABLE IF NOT EXISTS instrument_snapshots ( + instrument_code text NOT NULL, + asset_class text NOT NULL CHECK (asset_class IN ('stock', 'bond')), + snapshot_date date NOT NULL, + price numeric(28, 8) NOT NULL CHECK (price > 0), + source text NOT NULL, + source_revision text NOT NULL, + source_collected_at timestamptz NOT NULL, + source_payload_hash text NOT NULL, + first_seen_at timestamptz NOT NULL DEFAULT now(), + last_seen_at timestamptz NOT NULL DEFAULT now(), + PRIMARY KEY (instrument_code, asset_class, snapshot_date) + ); + + CREATE INDEX IF NOT EXISTS instrument_snapshots_date_idx + ON instrument_snapshots (instrument_code, asset_class, snapshot_date DESC); + CREATE TABLE IF NOT EXISTS dividend_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, @@ -2768,6 +2797,107 @@ type FundRepository(connectionString: string) = records |> Seq.toList + member _.UpsertInstrumentSnapshots(records: InstrumentSnapshotRecord list) = + if records.IsEmpty then + () + else + use connection = new NpgsqlConnection(connectionString) + connection.Open() + use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) + + try + for record in records do + use command = + commandWithTransaction + connection + (Some transaction) + """ + INSERT INTO instrument_snapshots + (instrument_code, asset_class, snapshot_date, price, source, + source_revision, source_collected_at, source_payload_hash) + VALUES + (@instrument_code, @asset_class, @snapshot_date, @price, @source, + @source_revision, @source_collected_at, @source_payload_hash) + ON CONFLICT (instrument_code, asset_class, snapshot_date) DO UPDATE SET + price = EXCLUDED.price, + source = EXCLUDED.source, + source_revision = EXCLUDED.source_revision, + source_collected_at = EXCLUDED.source_collected_at, + source_payload_hash = EXCLUDED.source_payload_hash, + last_seen_at = now(), + first_seen_at = CASE + WHEN instrument_snapshots.source_payload_hash = EXCLUDED.source_payload_hash + THEN instrument_snapshots.first_seen_at + ELSE now() + END + """ + + addParameter command "instrument_code" NpgsqlDbType.Text (box record.InstrumentCode) |> ignore + addParameter command "asset_class" NpgsqlDbType.Text (box record.AssetClass) |> ignore + addParameter command "snapshot_date" NpgsqlDbType.Date (box record.SnapshotDate) |> ignore + addParameter command "price" NpgsqlDbType.Numeric (box record.Price) |> ignore + addParameter command "source" NpgsqlDbType.Text (box record.Source) |> ignore + addParameter command "source_revision" NpgsqlDbType.Text (box record.SourceRevision) |> ignore + addParameter command "source_collected_at" NpgsqlDbType.TimestampTz (box record.SourceCollectedAt) |> ignore + addParameter command "source_payload_hash" NpgsqlDbType.Text (box record.SourcePayloadHash) |> ignore + command.ExecuteNonQuery() |> ignore + + transaction.Commit() + with error -> + try + transaction.Rollback() + with _ -> + () + + raise error + + /// Latest snapshot price on or before `asOfDate` for every instrument of the + /// given asset class held by the fund. Instruments without any snapshot are + /// omitted so a caller can tell "no snapshot yet" from a stored price. + member _.GetLatestSnapshots(fundId: Guid, assetClass: string, asOfDate: DateOnly) : Map = + use connection = new NpgsqlConnection(connectionString) + connection.Open() + use command = + commandWithTransaction + connection + None + """ + SELECT s.instrument_code, s.asset_class, s.snapshot_date, s.price, s.source, + s.source_revision, s.source_collected_at, s.source_payload_hash + FROM instrument_snapshots s + JOIN ( + SELECT instrument_code, MAX(snapshot_date) AS snapshot_date + FROM instrument_snapshots + WHERE asset_class = @asset_class AND snapshot_date <= @as_of_date + GROUP BY instrument_code + ) latest + ON latest.instrument_code = s.instrument_code + AND latest.snapshot_date = s.snapshot_date + WHERE s.asset_class = @asset_class + """ + + addParameter command "asset_class" NpgsqlDbType.Text (box assetClass) |> ignore + addParameter command "as_of_date" NpgsqlDbType.Date (box asOfDate) |> ignore + use reader = command.ExecuteReader() + let records = System.Collections.Generic.Dictionary() + + while reader.Read() do + let record : InstrumentSnapshotRecord = + { InstrumentCode = reader.GetString(0) + AssetClass = reader.GetString(1) + SnapshotDate = reader.GetFieldValue(2) + Price = reader.GetDecimal(3) + Source = reader.GetString(4) + SourceRevision = reader.GetString(5) + SourceCollectedAt = reader.GetFieldValue(6) + SourcePayloadHash = reader.GetString(7) } + + records.[record.InstrumentCode] <- record + + records + |> Seq.map (fun pair -> pair.Key, pair.Value) + |> Map.ofSeq + member _.CreateFund(idempotencyKey: string, command: FundCreateCommand) = if String.IsNullOrWhiteSpace idempotencyKey then FundWriteResult.Invalid "idempotency key cannot be empty" -- cgit v1.2.3