summaryrefslogtreecommitdiff
path: root/src/FundLab.Api/Persistence.fs
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-22 05:56:14 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-22 05:56:14 +0800
commit79600c15aef9bd411abe88dcf36286a102b5abea (patch)
treee5e3bfbff6f60a48a3c0ff8b78b48da98c8f238a /src/FundLab.Api/Persistence.fs
parent6abd4c206fbfe8369b3fde4c905619d08f13c5d1 (diff)
downloadfund-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.fs130
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"