summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-21 08:36:20 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-21 08:36:20 +0800
commit0aad89fc76aa77524625964ebad896a6bc0d3660 (patch)
treebbdca8e2b70efe67570b7ba55b1437dbcca0abc2
parentac44797721121cb7e00b380ee5ed5d464b74da6f (diff)
downloadfund-lab-0aad89fc76aa77524625964ebad896a6bc0d3660.tar.gz
Bind subscription confirmation to append-only nav observation evidence
-rw-r--r--src/FundLab.Api/Persistence.fs125
-rw-r--r--tests/FundLab.Api.Tests/OrderTests.fs165
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<DateTimeOffset>(8))
+
Some
{
Nav = quoteReader.GetDecimal(0)
@@ -1247,26 +1307,45 @@ type FundRepository(connectionString: string) =
Some(quoteReader.GetFieldValue<DateTimeOffset>(5))
PayloadHash = quoteReader.GetString(6)
FirstSeenAt = quoteReader.GetFieldValue<DateTimeOffset>(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<string>(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)
@@ -955,6 +983,133 @@ type SubscriptionConfirmationTests(fixture: PostgresFixture) =
| other -> failwithf "unexpected late confirm result: %A" other
[<Fact>]
+ 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 ])
+
+ [<Fact>]
+ 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
+ )
+
+ [<Fact>]
+ 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)
+
+ [<Fact>]
member _.``retrying after the cutoff keeps the original order trade date``() =
let fundId = createFund 10000.00m
let orderId = postOrder fundId "confirm-overnight-1" "1000.00" "0.00"