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 +++++++++++++++++++++----- tests/FundLab.Api.Tests/OrderTests.fs | 165 ++++++++++++++++++++++++++++++++-- 2 files changed, 262 insertions(+), 28 deletions(-) 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 diff --git a/tests/FundLab.Api.Tests/OrderTests.fs b/tests/FundLab.Api.Tests/OrderTests.fs index d975462..0c0d71a 100644 --- a/tests/FundLab.Api.Tests/OrderTests.fs +++ b/tests/FundLab.Api.Tests/OrderTests.fs @@ -693,10 +693,11 @@ type SubscriptionConfirmationTests(fixture: PostgresFixture) = DateTimeOffset(utc.Ticks - (utc.Ticks % 10L), TimeSpan.Zero) let insertQuoteOnDate (code: string) (nav: decimal) (navDate: DateOnly) (collectedAt: DateTimeOffset) (publishedAt: DateTimeOffset option) = + let revision = sprintf "akshare-test/%O" (Guid.NewGuid()) let payload: MarketDataNavPayload = { Source = "akshare" - SourceRevision = sprintf "akshare-test/%O" (Guid.NewGuid()) + SourceRevision = revision CollectedAt = collectedAt Code = code Observations = @@ -711,7 +712,7 @@ type SubscriptionConfirmationTests(fixture: PostgresFixture) = ] } - repository().UpsertNavObservations(payload, "confirm-quote-hash") + repository().UpsertNavObservations(payload, sprintf "confirm-quote-hash/%s" revision) let insertQuote (code: string) (nav: decimal) (collectedAt: DateTimeOffset) (publishedAt: DateTimeOffset option) = insertQuoteOnDate code nav (ConfirmationPolicy.tradeDateFor DateTimeOffset.Now) collectedAt publishedAt @@ -752,6 +753,19 @@ type SubscriptionConfirmationTests(fixture: PostgresFixture) = command.ExecuteScalar() :?> string + let scalarTimestamp (sql: string) (parameters: (string * obj * NpgsqlTypes.NpgsqlDbType) list) = + use connection = new NpgsqlConnection(fixture.ConnectionString) + connection.Open() + use command = connection.CreateCommand() + command.CommandText <- sql + + for name, value, dbType in parameters do + let parameter = command.Parameters.Add(name, dbType) + parameter.Value <- value + + let value = command.ExecuteScalar() :?> DateTime + DateTimeOffset(DateTime.SpecifyKind(value, DateTimeKind.Utc)) + let orderFundCode orderId = scalarText "SELECT fund_code FROM subscription_orders WHERE id = @order_id" [ "order_id", box orderId, NpgsqlTypes.NpgsqlDbType.Uuid ] @@ -827,14 +841,26 @@ type SubscriptionConfirmationTests(fixture: PostgresFixture) = insertQuote code 2.5m collected None let key = fixture.Key("confirm-replay-key") - match confirmKey fundId orderId key with - | OrderConfirmed _ -> () - | other -> failwithf "unexpected first confirm result: %A" other + let originalHash, originalFirstSeen = + match confirmKey fundId orderId key with + | OrderConfirmed confirmed -> + match confirmed.ConfirmedQuote with + | Some quote -> quote.PayloadHash, quote.FirstSeenAt + | None -> failwith "confirmed quote evidence missing" + | other -> failwithf "unexpected first confirm result: %A" other let cashAfterFirst = availableCash fundId insertQuote code 9.9m (truncateMicroseconds (DateTimeOffset.Now.AddSeconds(-8.0))) None + let supersededHash = + scalarText + "SELECT source_payload_hash FROM fund_nav_observations WHERE instrument_code = @code AND nav_date = @nav_date" + [ "code", box code, NpgsqlTypes.NpgsqlDbType.Text + "nav_date", box (ConfirmationPolicy.tradeDateFor DateTimeOffset.Now), NpgsqlTypes.NpgsqlDbType.Date ] + + Assert.NotEqual(originalHash, supersededHash) + match confirmKey fundId orderId key with | ConfirmReplayed replayed -> Assert.Equal("confirmed", replayed.Status) @@ -844,6 +870,8 @@ type SubscriptionConfirmationTests(fixture: PostgresFixture) = Assert.Equal(2.5m, quote.Nav) Assert.Equal(collected, quote.CollectedAt) Assert.Equal(None, quote.PublishedAt) + Assert.Equal(originalHash, quote.PayloadHash) + Assert.Equal(originalFirstSeen, quote.FirstSeenAt) | None -> failwith "replayed quote evidence missing" Assert.Equal(400.00000000m, replayed.ConfirmedUnits |> Option.get) @@ -954,6 +982,133 @@ type SubscriptionConfirmationTests(fixture: PostgresFixture) = | None -> failwith "late confirmed quote evidence missing" | other -> failwithf "unexpected late confirm result: %A" other + [] + member _.``a revision first observed after the confirmation attempt stays pending without moving cash or holdings``() = + let fundId = createFund 10000.00m + let orderId = postOrder fundId "confirm-evidence-1" "1000.00" "0.00" + let code = orderFundCode orderId + let collected = truncateMicroseconds (DateTimeOffset.Now.AddSeconds(-10.0)) + insertQuote code 2.5m collected None + + execute + "UPDATE fund_nav_observation_evidence SET first_seen_at = @future WHERE instrument_code = @code AND nav_date = @nav_date" + [ "future", box (truncateMicroseconds (DateTimeOffset.Now.AddSeconds(5.0))), NpgsqlTypes.NpgsqlDbType.TimestampTz + "code", box code, NpgsqlTypes.NpgsqlDbType.Text + "nav_date", box (ConfirmationPolicy.tradeDateFor DateTimeOffset.Now), NpgsqlTypes.NpgsqlDbType.Date ] + + match confirmKey fundId orderId (fixture.Key("confirm-evidence-key")) with + | ConfirmPendingNav pending -> + Assert.Equal("pending_nav", pending.Status) + Assert.True(pending.PendingReason.IsSome) + Assert.True(pending.PendingReason.Value.Contains("first observed after")) + Assert.Equal(None, pending.ConfirmIdempotencyKey) + Assert.Equal(None, pending.ConfirmedAt) + Assert.Equal(None, pending.ConfirmedQuote) + | other -> failwithf "unexpected evidence-deferral result: %A" other + + Assert.Equal(9000.00m, availableCash fundId) + Assert.Equal(1000.00m, reservedCash fundId) + Assert.Equal(0, (repository().GetFundPositions(fundId)).Length) + Assert.Equal(0L, PersistenceTestHelpers.queryCount fixture.ConnectionString "SELECT count(*) FROM subscription_order_events WHERE order_id = @order_id" [ "order_id", box orderId, NpgsqlTypes.NpgsqlDbType.Uuid ]) + Assert.Equal(0L, PersistenceTestHelpers.queryCount fixture.ConnectionString "SELECT count(*) FROM subscription_confirm_idempotencies WHERE order_id = @order_id" [ "order_id", box orderId, NpgsqlTypes.NpgsqlDbType.Uuid ]) + + [] + member _.``reobserving the same revision keeps the earliest first seen and a changed revision is not backdated``() = + let code = seedInstrument () + let navDate = ConfirmationPolicy.tradeDateFor DateTimeOffset.Now + let observationRow = + "SELECT first_seen_at FROM fund_nav_observations WHERE instrument_code = @code AND nav_date = @nav_date" + let rowParameters = + [ "code", box code, NpgsqlTypes.NpgsqlDbType.Text + "nav_date", box navDate, NpgsqlTypes.NpgsqlDbType.Date ] + + let buildPayload (revision: string) (collectedAt: DateTimeOffset) (nav: decimal) : MarketDataNavPayload = + { + Source = "akshare" + SourceRevision = revision + CollectedAt = collectedAt + Code = code + Observations = + [ + { + NavDate = navDate + PublishedAt = None + Nav = nav + AccumulatedNav = Some nav + DailyReturn = Some 0.0m + } + ] + } + + repository().UpsertNavObservations(buildPayload "akshare-test/first-seen-r1" (truncateMicroseconds (DateTimeOffset.Now.AddHours(-6.0))) 2.5m, "confirm-first-seen-hash-r1") + + let earliest = truncateMicroseconds (DateTimeOffset.Now.AddHours(-2.0)) + + execute + "UPDATE fund_nav_observations SET first_seen_at = @first_seen WHERE instrument_code = @code AND nav_date = @nav_date" + [ "first_seen", box earliest, NpgsqlTypes.NpgsqlDbType.TimestampTz + "code", box code, NpgsqlTypes.NpgsqlDbType.Text + "nav_date", box navDate, NpgsqlTypes.NpgsqlDbType.Date ] + + repository().UpsertNavObservations(buildPayload "akshare-test/first-seen-r1" (truncateMicroseconds (DateTimeOffset.Now.AddHours(-3.0))) 2.5m, "confirm-first-seen-hash-r1") + Assert.Equal(earliest, scalarTimestamp observationRow rowParameters) + + let marker = truncateMicroseconds DateTimeOffset.UtcNow + let r2Collected = truncateMicroseconds (DateTimeOffset.Now.AddHours(-5.0)) + repository().UpsertNavObservations(buildPayload "akshare-test/first-seen-r2" r2Collected 3.3m, "confirm-first-seen-hash-r2") + + Assert.Equal("confirm-first-seen-hash-r2", scalarText "SELECT source_payload_hash FROM fund_nav_observations WHERE instrument_code = @code AND nav_date = @nav_date" rowParameters) + Assert.Equal(3.3m, scalarDecimal "SELECT nav FROM fund_nav_observations WHERE instrument_code = @code AND nav_date = @nav_date" rowParameters) + Assert.Equal(r2Collected, scalarTimestamp "SELECT source_collected_at FROM fund_nav_observations WHERE instrument_code = @code AND nav_date = @nav_date" rowParameters) + + let firstSeenAfterChange = scalarTimestamp observationRow rowParameters + Assert.True( + firstSeenAfterChange >= marker.AddSeconds(-1.0), + sprintf "expected non-backdated first_seen_at after revision change, got %O (marker %O, earliest %O)" firstSeenAfterChange marker earliest + ) + + [] + member _.``re-observing a superseded revision preserves its earliest first seen evidence``() = + let code = seedInstrument () + let navDate = ConfirmationPolicy.tradeDateFor DateTimeOffset.Now + let rowParameters = + [ "code", box code, NpgsqlTypes.NpgsqlDbType.Text + "nav_date", box navDate, NpgsqlTypes.NpgsqlDbType.Date ] + + let buildPayload (revision: string) (collectedAt: DateTimeOffset) (nav: decimal) : MarketDataNavPayload = + { + Source = "akshare" + SourceRevision = revision + CollectedAt = collectedAt + Code = code + Observations = + [ + { + NavDate = navDate + PublishedAt = None + Nav = nav + AccumulatedNav = Some nav + DailyReturn = Some 0.0m + } + ] + } + + repository().UpsertNavObservations(buildPayload "akshare-test/aba-r1" (truncateMicroseconds (DateTimeOffset.Now.AddHours(-6.0))) 2.5m, "confirm-aba-hash-r1") + + let evidenceAFirst = + scalarTimestamp + "SELECT first_seen_at FROM fund_nav_observation_evidence WHERE instrument_code = @code AND nav_date = @nav_date AND source_payload_hash = @hash" + (rowParameters @ [ "hash", box "confirm-aba-hash-r1", NpgsqlTypes.NpgsqlDbType.Text ]) + + repository().UpsertNavObservations(buildPayload "akshare-test/aba-r2" (truncateMicroseconds (DateTimeOffset.Now.AddHours(-3.0))) 3.3m, "confirm-aba-hash-r2") + Assert.Equal(3.3m, scalarDecimal "SELECT nav FROM fund_nav_observations WHERE instrument_code = @code AND nav_date = @nav_date" rowParameters) + + repository().UpsertNavObservations(buildPayload "akshare-test/aba-r1" (truncateMicroseconds (DateTimeOffset.Now.AddHours(-6.0))) 2.5m, "confirm-aba-hash-r1") + + Assert.Equal(evidenceAFirst, scalarTimestamp "SELECT first_seen_at FROM fund_nav_observation_evidence WHERE instrument_code = @code AND nav_date = @nav_date AND source_payload_hash = @hash" (rowParameters @ [ "hash", box "confirm-aba-hash-r1", NpgsqlTypes.NpgsqlDbType.Text ])) + Assert.Equal("confirm-aba-hash-r1", scalarText "SELECT source_payload_hash FROM fund_nav_observations WHERE instrument_code = @code AND nav_date = @nav_date" rowParameters) + Assert.Equal(2.5m, scalarDecimal "SELECT nav FROM fund_nav_observations WHERE instrument_code = @code AND nav_date = @nav_date" rowParameters) + [] member _.``retrying after the cutoff keeps the original order trade date``() = let fundId = createFund 10000.00m -- cgit v1.2.3