From 0aad89fc76aa77524625964ebad896a6bc0d3660 Mon Sep 17 00:00:00 2001 From: "Somhairle H. Marisol" Date: Mon, 21 Sep 2026 08:36:20 +0800 Subject: Bind subscription confirmation to append-only nav observation evidence --- src/FundLab.Api/Persistence.fs | 125 +++++++++++++++++++++++++++++++++-------- 1 file changed, 102 insertions(+), 23 deletions(-) (limited to 'src/FundLab.Api') diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs index de8061b..6080beb 100644 --- a/src/FundLab.Api/Persistence.fs +++ b/src/FundLab.Api/Persistence.fs @@ -265,6 +265,19 @@ type FundRepository(connectionString: string) = CREATE INDEX IF NOT EXISTS fund_nav_observations_date_idx ON fund_nav_observations (instrument_code, nav_date); + CREATE TABLE IF NOT EXISTS fund_nav_observation_evidence ( + instrument_code text NOT NULL REFERENCES instruments(code), + nav_date date NOT NULL, + source_payload_hash text NOT NULL, + source_revision text NOT NULL, + source text NOT NULL, + source_collected_at timestamptz NOT NULL, + published_at timestamptz NULL, + nav numeric(28, 8) NOT NULL, + first_seen_at timestamptz NOT NULL DEFAULT now(), + PRIMARY KEY (instrument_code, nav_date, source_payload_hash) + ); + ALTER TABLE funds ADD COLUMN IF NOT EXISTS reserved_cash numeric(20, 2) NOT NULL DEFAULT 0; CREATE TABLE IF NOT EXISTS subscription_orders ( @@ -741,7 +754,13 @@ type FundRepository(connectionString: string) = source_revision = EXCLUDED.source_revision, source_collected_at = EXCLUDED.source_collected_at, source_payload_hash = EXCLUDED.source_payload_hash, - last_seen_at = now() + last_seen_at = now(), + first_seen_at = + CASE + WHEN fund_nav_observations.source_payload_hash = EXCLUDED.source_payload_hash + THEN fund_nav_observations.first_seen_at + ELSE now() + END """ addParameter command "instrument_code" NpgsqlDbType.Text (box payload.Code) |> ignore @@ -776,6 +795,37 @@ type FundRepository(connectionString: string) = addParameter command "source_payload_hash" NpgsqlDbType.Text (box payloadHash) |> ignore command.ExecuteNonQuery() |> ignore + use evidenceCommand = + commandWithTransaction + connection + transaction + """ + INSERT INTO fund_nav_observation_evidence + (instrument_code, nav_date, source_payload_hash, source_revision, source, + source_collected_at, published_at, nav) + VALUES + (@instrument_code, @nav_date, @source_payload_hash, @source_revision, @source, + @source_collected_at, @published_at, @nav) + ON CONFLICT (instrument_code, nav_date, source_payload_hash) DO NOTHING + """ + + addParameter evidenceCommand "instrument_code" NpgsqlDbType.Text (box payload.Code) |> ignore + addParameter evidenceCommand "nav_date" NpgsqlDbType.Date (box observation.NavDate) |> ignore + addParameter evidenceCommand "source_payload_hash" NpgsqlDbType.Text (box payloadHash) |> ignore + addParameter evidenceCommand "source_revision" NpgsqlDbType.Text (box payload.SourceRevision) |> ignore + addParameter evidenceCommand "source" NpgsqlDbType.Text (box payload.Source) |> ignore + addParameter evidenceCommand "source_collected_at" NpgsqlDbType.TimestampTz (box payload.CollectedAt) |> ignore + + addParameter + evidenceCommand + "published_at" + NpgsqlDbType.TimestampTz + (optionalParameterValue observation.PublishedAt) + |> ignore + + addParameter evidenceCommand "nav" NpgsqlDbType.Numeric (box observation.Nav) |> ignore + evidenceCommand.ExecuteNonQuery() |> ignore + let requestHash (command: FundCreateCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value @@ -1217,11 +1267,15 @@ type FundRepository(connectionString: string) = connection (Some transaction) """ - SELECT nav, nav_date, source, source_revision, source_collected_at, published_at, - source_payload_hash, first_seen_at - FROM fund_nav_observations - WHERE instrument_code = @code AND nav_date = @trade_date AND nav > 0 - ORDER BY source_collected_at DESC, published_at DESC NULLS LAST, source_revision DESC + SELECT o.nav, o.nav_date, o.source, o.source_revision, o.source_collected_at, o.published_at, + o.source_payload_hash, o.first_seen_at, e.first_seen_at + FROM fund_nav_observations o + LEFT JOIN fund_nav_observation_evidence e + ON e.instrument_code = o.instrument_code + AND e.nav_date = o.nav_date + AND e.source_payload_hash = o.source_payload_hash + WHERE o.instrument_code = @code AND o.nav_date = @trade_date AND o.nav > 0 + ORDER BY o.source_collected_at DESC, o.published_at DESC NULLS LAST, o.source_revision DESC LIMIT 1 """ @@ -1231,8 +1285,14 @@ type FundRepository(connectionString: string) = use quoteReader = quoteCommand.ExecuteReader() let quoteFound = quoteReader.Read() - let selectedQuote = + let selectedQuote, evidenceFirstSeen = if quoteFound then + let evidenceFirstSeen = + if quoteReader.IsDBNull(8) then + None + else + Some(quoteReader.GetFieldValue(8)) + Some { Nav = quoteReader.GetDecimal(0) @@ -1247,26 +1307,45 @@ type FundRepository(connectionString: string) = Some(quoteReader.GetFieldValue(5)) PayloadHash = quoteReader.GetString(6) FirstSeenAt = quoteReader.GetFieldValue(7) - } + }, + evidenceFirstSeen else - None + None, None quoteReader.Close() + let boundQuote = + match selectedQuote, evidenceFirstSeen with + | Some quote, Some firstSeen -> Some { quote with FirstSeenAt = firstSeen } + | _ -> None + let deferralReason = - match selectedQuote with - | None -> Some(sprintf "nav for trade date %s is not available yet" tradeDateText) - | Some quote -> - ConfirmationPolicy.navDeferralReason - { - Nav = quote.Nav - NavDate = quote.NavDate - CollectedAt = quote.CollectedAt - PublishedAt = quote.PublishedAt - } - tradeDate - today - confirmedAt + if selectedQuote.IsSome && evidenceFirstSeen.IsNone then + Some( + sprintf + "nav revision for trade date %s has no observation evidence recorded" + tradeDateText + ) + else + match boundQuote with + | None -> Some(sprintf "nav for trade date %s is not available yet" tradeDateText) + | Some quote when quote.FirstSeenAt > confirmedAt -> + Some( + sprintf + "nav revision for trade date %s was first observed after the confirmation attempt" + tradeDateText + ) + | Some quote -> + ConfirmationPolicy.navDeferralReason + { + Nav = quote.Nav + NavDate = quote.NavDate + CollectedAt = quote.CollectedAt + PublishedAt = quote.PublishedAt + } + tradeDate + today + confirmedAt let markPending reason = use pendingCommand = @@ -1308,7 +1387,7 @@ type FundRepository(connectionString: string) = transaction.Commit() readPendingAfterCommit () | None -> - let quote = selectedQuote |> Option.get + let quote = boundQuote |> Option.get let source = quote.Source let revision = quote.Revision -- cgit v1.2.3