diff options
| author | Somhairle H. Marisol <[email protected]> | 2026-09-22 05:56:14 +0800 |
|---|---|---|
| committer | Somhairle H. Marisol <[email protected]> | 2026-09-22 05:56:14 +0800 |
| commit | 79600c15aef9bd411abe88dcf36286a102b5abea (patch) | |
| tree | e5e3bfbff6f60a48a3c0ff8b78b48da98c8f238a /src/FundLab.Api/Persistence.fs | |
| parent | 6abd4c206fbfe8369b3fde4c905619d08f13c5d1 (diff) | |
| download | fund-lab-79600c15aef9bd411abe88dcf36286a102b5abea.tar.gz | |
Add daily market refresh endpoint, snapshot-first valuation and UI refresh (3d-25)
Diffstat (limited to 'src/FundLab.Api/Persistence.fs')
| -rw-r--r-- | src/FundLab.Api/Persistence.fs | 130 |
1 files changed, 130 insertions, 0 deletions
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<string, InstrumentSnapshotRecord> = + 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<string, InstrumentSnapshotRecord>() + + while reader.Read() do + let record : InstrumentSnapshotRecord = + { InstrumentCode = reader.GetString(0) + AssetClass = reader.GetString(1) + SnapshotDate = reader.GetFieldValue<DateOnly>(2) + Price = reader.GetDecimal(3) + Source = reader.GetString(4) + SourceRevision = reader.GetString(5) + SourceCollectedAt = reader.GetFieldValue<DateTimeOffset>(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" |
