summaryrefslogtreecommitdiff
path: root/src/FundLab.Api
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-21 08:13:13 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-21 08:13:13 +0800
commitac44797721121cb7e00b380ee5ed5d464b74da6f (patch)
tree9bad6f1b4ecbaccb11adce190f4c9cfb88c52b1d /src/FundLab.Api
parent64a4ca5c0a523ffc51c861eac1291671b90a2f3a (diff)
downloadfund-lab-ac44797721121cb7e00b380ee5ed5d464b74da6f.tar.gz
Implement subscription order confirmation with locked trade dates and bound quote evidence
Diffstat (limited to 'src/FundLab.Api')
-rw-r--r--src/FundLab.Api/Persistence.fs656
1 files changed, 651 insertions, 5 deletions
diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs
index 6c60bdd..de8061b 100644
--- a/src/FundLab.Api/Persistence.fs
+++ b/src/FundLab.Api/Persistence.fs
@@ -9,6 +9,105 @@ open System.Text
open FundLab.Domain
open Npgsql
open NpgsqlTypes
+open System.Text.Json
+
+module ConfirmationPolicy =
+ let cutoffTimeOfDay = TimeSpan(15, 0, 0)
+ let unitsMaximum = 99999999999999999999.99999999m
+ let amountMaximum = 999999999999999999.99m
+
+ let private shanghaiZone =
+ lazy (
+ try
+ TimeZoneInfo.FindSystemTimeZoneById("Asia/Shanghai")
+ with _ ->
+ TimeZoneInfo.CreateCustomTimeZone(
+ "Asia/Shanghai (UTC+8)",
+ TimeSpan(8, 0, 0),
+ "Asia/Shanghai (UTC+8)",
+ "Asia/Shanghai (UTC+8)"
+ )
+ )
+
+ let shanghaiDate (moment: DateTimeOffset) : DateOnly =
+ DateOnly.FromDateTime(TimeZoneInfo.ConvertTime(moment, shanghaiZone.Value).Date)
+
+ let private isWeekend (date: DateOnly) =
+ date.DayOfWeek = DayOfWeek.Saturday || date.DayOfWeek = DayOfWeek.Sunday
+
+ let rec private rollToWeekday (date: DateOnly) =
+ if isWeekend date then rollToWeekday (date.AddDays 1) else date
+
+ let isModeledTradingDay (date: DateOnly) : bool = not (isWeekend date)
+
+ let tradeDateFor (submittedAt: DateTimeOffset) : DateOnly =
+ let local = TimeZoneInfo.ConvertTime(submittedAt, shanghaiZone.Value)
+ let date = DateOnly.FromDateTime(local.Date)
+ let candidate = if local.TimeOfDay >= cutoffTimeOfDay then date.AddDays 1 else date
+ rollToWeekday candidate
+
+ type NavQuote =
+ { NavDate: DateOnly
+ Nav: decimal
+ CollectedAt: DateTimeOffset
+ PublishedAt: DateTimeOffset option }
+
+ type ConfirmationComputation =
+ {
+ Units: decimal
+ InvestedCash: decimal
+ ResidualCash: decimal
+ }
+
+ let navDeferralReason
+ (quote: NavQuote)
+ (tradeDate: DateOnly)
+ (today: DateOnly)
+ (now: DateTimeOffset)
+ : string option =
+ let tradeDateText = tradeDate.ToString("yyyy-MM-dd")
+
+ if quote.Nav <= 0m then
+ Some(sprintf "nav for trade date %s is not positive" tradeDateText)
+ elif quote.NavDate > today then
+ Some(sprintf "nav for trade date %s is dated in the future" tradeDateText)
+ elif quote.CollectedAt > now then
+ Some(sprintf "nav for trade date %s is not yet collected" tradeDateText)
+ elif quote.PublishedAt |> Option.exists (fun publishedAt -> publishedAt > now) then
+ Some(sprintf "nav for trade date %s is not yet published" tradeDateText)
+ elif quote.NavDate <> tradeDate then
+ Some(sprintf "nav for trade date %s is not available yet" tradeDateText)
+ else
+ None
+
+ let compute (amount: decimal) (nav: decimal) : Result<ConfirmationComputation, string> =
+ if amount <= 0m then
+ Error "amount must be positive"
+ elif amount > amountMaximum then
+ Error "amount exceeds supported precision"
+ elif nav <= 0m then
+ Error "unit nav must be positive"
+ elif nav < 0.00000001m then
+ Error "unit nav is below database precision"
+ else
+ let rawUnits = amount / nav
+
+ if rawUnits > unitsMaximum then
+ Error "unit amount exceeds database precision"
+ else
+ let units = Decimal.Truncate(rawUnits * 100000000m) / 100000000m
+
+ if units <= 0m then
+ Error "amount converts to zero units at this unit nav"
+ else
+ let invested = Decimal.Truncate(units * nav * 100m) / 100m
+
+ Ok
+ {
+ Units = units
+ InvestedCash = invested
+ ResidualCash = amount - invested
+ }
[<CLIMutable>]
type FundCreateCommand =
@@ -46,6 +145,18 @@ type SubscriptionOrderCommand =
FeeAmount: decimal
}
+type ConfirmedQuoteEvidence =
+ {
+ Nav: decimal
+ NavDate: DateOnly
+ Source: string
+ Revision: string
+ CollectedAt: DateTimeOffset
+ PublishedAt: DateTimeOffset option
+ PayloadHash: string
+ FirstSeenAt: DateTimeOffset
+ }
+
type SubscriptionOrderRecord =
{
Id: Guid
@@ -57,6 +168,14 @@ type SubscriptionOrderRecord =
Status: string
SubmittedAt: DateTimeOffset
IsSynthetic: bool
+ TradeDate: DateOnly
+ ConfirmIdempotencyKey: string option
+ PendingReason: string option
+ ConfirmedAt: DateTimeOffset option
+ ConfirmedQuote: ConfirmedQuoteEvidence option
+ ConfirmedUnits: decimal option
+ ConfirmedInvestedCash: decimal option
+ ConfirmedResidualCash: decimal option
}
type SubscriptionOrderWriteResult =
@@ -68,6 +187,28 @@ type SubscriptionOrderWriteResult =
| OrderInstrumentNotFound
| OrderInsufficientFunds
+type SubscriptionConfirmResult =
+ | OrderConfirmed of SubscriptionOrderRecord
+ | ConfirmReplayed of SubscriptionOrderRecord
+ | ConfirmPendingNav of SubscriptionOrderRecord
+ | ConfirmIdempotencyConflict
+ | ConfirmAlreadyConfirmed
+ | ConfirmOrderNotFound
+ | ConfirmInvalidStatus
+ | ConfirmInvalid of string
+
+type FundPositionRecord =
+ {
+ FundId: Guid
+ InstrumentCode: string
+ Units: decimal
+ CostCash: decimal
+ LastConfirmedAt: DateTimeOffset
+ ValuationNav: decimal option
+ ValuationNavDate: DateOnly option
+ ValuationCollectedAt: DateTimeOffset option
+ }
+
type FundRepository(connectionString: string) =
let cashMaximum = 999999999999999999.99m
let unitNavMaximum = 99999999999999999999.99999999m
@@ -148,6 +289,53 @@ type FundRepository(connectionString: string) =
fund_id uuid NOT NULL REFERENCES funds(id),
created_at timestamptz NOT NULL DEFAULT now()
);
+
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirm_idempotency_key text NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS pending_reason text NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_at timestamptz NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_nav numeric(28, 8) NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_nav_date date NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_nav_source text NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_nav_revision text NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_nav_collected_at timestamptz NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_nav_published_at timestamptz NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_units numeric(28, 8) NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_invested_cash numeric(20, 2) NULL;
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_residual_cash numeric(20, 2) NULL;
+
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_nav_payload_hash text NULL;
+
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS confirmed_nav_first_seen_at timestamptz NULL;
+
+ ALTER TABLE subscription_orders ADD COLUMN IF NOT EXISTS trade_date date NULL;
+
+ UPDATE subscription_orders SET trade_date = (submitted_at AT TIME ZONE 'Asia/Shanghai')::date WHERE trade_date IS NULL;
+
+ CREATE TABLE IF NOT EXISTS fund_positions (
+ fund_id uuid NOT NULL REFERENCES funds(id),
+ instrument_code text NOT NULL REFERENCES instruments(code),
+ units numeric(28, 8) NOT NULL CHECK (units > 0),
+ cost_cash numeric(20, 2) NOT NULL CHECK (cost_cash >= 0),
+ first_confirmed_at timestamptz NOT NULL,
+ last_confirmed_at timestamptz NOT NULL,
+ PRIMARY KEY (fund_id, instrument_code)
+ );
+
+ CREATE TABLE IF NOT EXISTS subscription_order_events (
+ id bigserial PRIMARY KEY,
+ order_id uuid NOT NULL REFERENCES subscription_orders(id),
+ event_type text NOT NULL,
+ detail jsonb NOT NULL,
+ created_at timestamptz NOT NULL DEFAULT now()
+ );
+
+ CREATE TABLE IF NOT EXISTS subscription_confirm_idempotencies (
+ idempotency_key text PRIMARY KEY,
+ request_hash text NOT NULL,
+ order_id uuid NOT NULL REFERENCES subscription_orders(id),
+ fund_id uuid NOT NULL REFERENCES funds(id),
+ created_at timestamptz NOT NULL DEFAULT now()
+ );
"""
let statusText status =
@@ -311,7 +499,41 @@ type FundRepository(connectionString: string) =
let orderStatusText = "submitted"
+ let readStringOption (reader: DbDataReader) ordinal =
+ if reader.IsDBNull ordinal then None else Some(reader.GetString ordinal)
+
+ let readDecimalOption (reader: DbDataReader) ordinal =
+ if reader.IsDBNull ordinal then None else Some(reader.GetDecimal ordinal)
+
let orderRecordFromReader (reader: DbDataReader) : SubscriptionOrderRecord =
+ let confirmedQuote =
+ if reader.IsDBNull(13) then
+ None
+ else
+ Some
+ {
+ Nav = reader.GetDecimal(13)
+ NavDate = reader.GetFieldValue<DateOnly>(14)
+ Source = reader.GetString(15)
+ Revision = reader.GetString(16)
+ CollectedAt = reader.GetFieldValue<DateTimeOffset>(17)
+ PublishedAt =
+ if reader.IsDBNull(18) then
+ None
+ else
+ Some(reader.GetFieldValue<DateTimeOffset>(18))
+ PayloadHash =
+ if reader.IsDBNull(19) then
+ ""
+ else
+ reader.GetString(19)
+ FirstSeenAt =
+ if reader.IsDBNull(20) then
+ DateTimeOffset.UnixEpoch
+ else
+ reader.GetFieldValue<DateTimeOffset>(20)
+ }
+
{
Id = reader.GetGuid(0)
FundId = reader.GetGuid(1)
@@ -322,6 +544,18 @@ type FundRepository(connectionString: string) =
Status = reader.GetString(6)
IsSynthetic = reader.GetBoolean(7)
SubmittedAt = reader.GetFieldValue<DateTimeOffset>(8)
+ TradeDate = reader.GetFieldValue<DateOnly>(9)
+ ConfirmIdempotencyKey = readStringOption reader 10
+ PendingReason = readStringOption reader 11
+ ConfirmedAt =
+ if reader.IsDBNull(12) then
+ None
+ else
+ Some(reader.GetFieldValue<DateTimeOffset>(12))
+ ConfirmedUnits = readDecimalOption reader 21
+ ConfirmedQuote = confirmedQuote
+ ConfirmedInvestedCash = readDecimalOption reader 22
+ ConfirmedResidualCash = readDecimalOption reader 23
}
let findOrder connection transaction orderId =
@@ -331,7 +565,13 @@ type FundRepository(connectionString: string) =
transaction
"""
SELECT id, fund_id, fund_code, amount, fee_amount,
- reserved_total, status, is_synthetic, submitted_at
+ reserved_total, status, is_synthetic, submitted_at,
+ trade_date,
+ confirm_idempotency_key, pending_reason, confirmed_at,
+ confirmed_nav, confirmed_nav_date, confirmed_nav_source,
+ confirmed_nav_revision, confirmed_nav_collected_at, confirmed_nav_published_at,
+ confirmed_nav_payload_hash, confirmed_nav_first_seen_at,
+ confirmed_units, confirmed_invested_cash, confirmed_residual_cash
FROM subscription_orders
WHERE id = @order_id
"""
@@ -341,6 +581,21 @@ type FundRepository(connectionString: string) =
use reader = command.ExecuteReader()
if reader.Read() then Some(orderRecordFromReader reader) else None
+ let findConfirmIdempotency connection transaction key =
+ use command =
+ commandWithTransaction
+ connection
+ transaction
+ "SELECT request_hash, fund_id, order_id FROM subscription_confirm_idempotencies WHERE idempotency_key = @idempotency_key"
+
+ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore
+
+ use reader = command.ExecuteReader()
+ if reader.Read() then
+ Some(reader.GetString(0), reader.GetGuid(1), reader.GetGuid(2))
+ else
+ None
+
let findOrderIdempotency connection transaction key =
use command =
commandWithTransaction
@@ -388,10 +643,10 @@ type FundRepository(connectionString: string) =
"""
INSERT INTO subscription_orders
(id, fund_id, fund_code, amount, fee_amount,
- reserved_total, status, is_synthetic)
+ reserved_total, status, is_synthetic, trade_date)
VALUES
(@id, @fund_id, @fund_code, @amount, @fee_amount,
- @reserved_total, @status, @is_synthetic)
+ @reserved_total, @status, @is_synthetic, @trade_date)
RETURNING submitted_at
"""
@@ -403,6 +658,7 @@ type FundRepository(connectionString: string) =
addParameter command "reserved_total" NpgsqlDbType.Numeric (box order.ReservedTotal) |> ignore
addParameter command "status" NpgsqlDbType.Text (box order.Status) |> ignore
addParameter command "is_synthetic" NpgsqlDbType.Boolean (box order.IsSynthetic) |> ignore
+ addParameter command "trade_date" NpgsqlDbType.Date (box order.TradeDate) |> ignore
use reader = command.ExecuteReader()
reader.Read() |> ignore
@@ -825,6 +1081,7 @@ type FundRepository(connectionString: string) =
transaction.Rollback()
OrderInsufficientFunds
else
+ let submittedAt = DateTimeOffset.UtcNow
let order: SubscriptionOrderRecord =
{
Id = Guid.NewGuid()
@@ -834,8 +1091,16 @@ type FundRepository(connectionString: string) =
FeeAmount = command.FeeAmount
ReservedTotal = reservedTotal
Status = orderStatusText
- SubmittedAt = DateTimeOffset.UnixEpoch
+ SubmittedAt = submittedAt
IsSynthetic = isSynthetic
+ TradeDate = ConfirmationPolicy.tradeDateFor submittedAt
+ ConfirmIdempotencyKey = None
+ PendingReason = None
+ ConfirmedAt = None
+ ConfirmedQuote = None
+ ConfirmedUnits = None
+ ConfirmedInvestedCash = None
+ ConfirmedResidualCash = None
}
let submittedAt = insertSubscriptionOrder connection (Some transaction) order
@@ -863,7 +1128,13 @@ type FundRepository(connectionString: string) =
None
"""
SELECT id, fund_id, fund_code, amount, fee_amount,
- reserved_total, status, is_synthetic, submitted_at
+ reserved_total, status, is_synthetic, submitted_at,
+ trade_date,
+ confirm_idempotency_key, pending_reason, confirmed_at,
+ confirmed_nav, confirmed_nav_date, confirmed_nav_source,
+ confirmed_nav_revision, confirmed_nav_collected_at, confirmed_nav_published_at,
+ confirmed_nav_payload_hash, confirmed_nav_first_seen_at,
+ confirmed_units, confirmed_invested_cash, confirmed_residual_cash
FROM subscription_orders
WHERE fund_id = @fund_id
ORDER BY submitted_at DESC, id
@@ -878,3 +1149,378 @@ type FundRepository(connectionString: string) =
records.Add(orderRecordFromReader reader)
records |> Seq.toList
+
+ member _.ConfirmSubscriptionOrder(idempotencyKey: string, fundId: Guid, orderId: Guid) =
+ if String.IsNullOrWhiteSpace idempotencyKey then
+ ConfirmInvalid "idempotency key cannot be empty"
+ else
+ use connection = new NpgsqlConnection(connectionString)
+ connection.Open()
+ use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted)
+
+ try
+ let confirmedAt = DateTimeOffset.UtcNow
+ let today = ConfirmationPolicy.shanghaiDate confirmedAt
+
+ let requestHash =
+ Convert
+ .ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(sprintf "confirm:%O:%O" fundId orderId)))
+ .ToLowerInvariant()
+
+ use lockCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ "SELECT pg_advisory_xact_lock(hashtext(@lock_key))"
+
+ addParameter lockCommand "lock_key" NpgsqlDbType.Text (box (sprintf "confirm-order:%O" orderId)) |> ignore
+ lockCommand.ExecuteNonQuery() |> ignore
+
+ match findOrder connection (Some transaction) orderId with
+ | None ->
+ transaction.Rollback()
+ ConfirmOrderNotFound
+ | Some order when order.FundId <> fundId ->
+ transaction.Rollback()
+ ConfirmOrderNotFound
+ | Some order ->
+ if order.Status = "confirmed" then
+ match order.ConfirmIdempotencyKey with
+ | Some storedKey when storedKey = idempotencyKey ->
+ transaction.Commit()
+ ConfirmReplayed order
+ | _ ->
+ match findConfirmIdempotency connection (Some transaction) idempotencyKey with
+ | Some(_, _, storedOrderId) when storedOrderId = orderId ->
+ transaction.Commit()
+ ConfirmReplayed order
+ | Some _ ->
+ transaction.Rollback()
+ ConfirmIdempotencyConflict
+ | None ->
+ transaction.Rollback()
+ ConfirmAlreadyConfirmed
+ elif order.Status <> "submitted" && order.Status <> "pending_nav" then
+ transaction.Rollback()
+ ConfirmInvalidStatus
+ else
+ match findConfirmIdempotency connection (Some transaction) idempotencyKey with
+ | Some _ ->
+ transaction.Rollback()
+ ConfirmIdempotencyConflict
+ | None ->
+ let tradeDate = order.TradeDate
+ let tradeDateText = tradeDate.ToString("yyyy-MM-dd")
+
+ use quoteCommand =
+ commandWithTransaction
+ 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
+ LIMIT 1
+ """
+
+ addParameter quoteCommand "code" NpgsqlDbType.Text (box order.FundCode) |> ignore
+ addParameter quoteCommand "trade_date" NpgsqlDbType.Date (box tradeDate) |> ignore
+
+ use quoteReader = quoteCommand.ExecuteReader()
+ let quoteFound = quoteReader.Read()
+
+ let selectedQuote =
+ if quoteFound then
+ Some
+ {
+ Nav = quoteReader.GetDecimal(0)
+ NavDate = quoteReader.GetFieldValue<DateOnly>(1)
+ Source = quoteReader.GetString(2)
+ Revision = quoteReader.GetString(3)
+ CollectedAt = quoteReader.GetFieldValue<DateTimeOffset>(4)
+ PublishedAt =
+ if quoteReader.IsDBNull(5) then
+ None
+ else
+ Some(quoteReader.GetFieldValue<DateTimeOffset>(5))
+ PayloadHash = quoteReader.GetString(6)
+ FirstSeenAt = quoteReader.GetFieldValue<DateTimeOffset>(7)
+ }
+ else
+ None
+
+ quoteReader.Close()
+
+ 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
+
+ let markPending reason =
+ use pendingCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ UPDATE subscription_orders
+ SET status = 'pending_nav',
+ pending_reason = @reason,
+ confirm_idempotency_key = NULL,
+ confirmed_at = NULL,
+ confirmed_nav = NULL,
+ confirmed_nav_date = NULL,
+ confirmed_nav_source = NULL,
+ confirmed_nav_revision = NULL,
+ confirmed_nav_collected_at = NULL,
+ confirmed_nav_published_at = NULL,
+ confirmed_nav_payload_hash = NULL,
+ confirmed_nav_first_seen_at = NULL,
+ confirmed_units = NULL,
+ confirmed_invested_cash = NULL,
+ confirmed_residual_cash = NULL
+ WHERE id = @order_id
+ """
+
+ addParameter pendingCommand "reason" NpgsqlDbType.Text (box reason) |> ignore
+ addParameter pendingCommand "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore
+ pendingCommand.ExecuteNonQuery() |> ignore
+
+ let readPendingAfterCommit () =
+ match findOrder connection None orderId with
+ | Some pendingOrder -> ConfirmPendingNav pendingOrder
+ | None -> failwith "pending order disappeared after confirmation deferral"
+
+ match deferralReason with
+ | Some reason ->
+ markPending reason
+ transaction.Commit()
+ readPendingAfterCommit ()
+ | None ->
+ let quote = selectedQuote |> Option.get
+ let source = quote.Source
+ let revision = quote.Revision
+
+ match ConfirmationPolicy.compute order.Amount quote.Nav with
+ | Error message ->
+ markPending message
+ transaction.Commit()
+ readPendingAfterCommit ()
+ | Ok computation ->
+ use cashCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ UPDATE funds
+ SET available_cash = available_cash + @residual,
+ reserved_cash = reserved_cash - @reserved_total
+ WHERE id = @fund_id AND reserved_cash >= @reserved_total
+ """
+
+ addParameter cashCommand "residual" NpgsqlDbType.Numeric (box computation.ResidualCash) |> ignore
+ addParameter cashCommand "reserved_total" NpgsqlDbType.Numeric (box order.ReservedTotal) |> ignore
+ addParameter cashCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore
+
+ if cashCommand.ExecuteNonQuery() = 0 then
+ failwith "insufficient reserved cash for order confirmation"
+
+ let publishedParameter =
+ match quote.PublishedAt with
+ | Some published -> box published
+ | None -> box DBNull.Value
+
+ use confirmCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ UPDATE subscription_orders
+ SET status = 'confirmed',
+ pending_reason = NULL,
+ confirm_idempotency_key = @idempotency_key,
+ confirmed_at = @confirmed_at,
+ confirmed_nav = @nav,
+ confirmed_nav_date = @nav_date,
+ confirmed_nav_source = @source,
+ confirmed_nav_revision = @revision,
+ confirmed_nav_collected_at = @collected_at,
+ confirmed_nav_published_at = @published_at,
+ confirmed_nav_payload_hash = @payload_hash,
+ confirmed_nav_first_seen_at = @first_seen_at,
+ confirmed_units = @units,
+ confirmed_invested_cash = @invested,
+ confirmed_residual_cash = @residual
+ WHERE id = @order_id
+ """
+
+ addParameter confirmCommand "idempotency_key" NpgsqlDbType.Text (box idempotencyKey) |> ignore
+ addParameter confirmCommand "confirmed_at" NpgsqlDbType.TimestampTz (box confirmedAt) |> ignore
+ addParameter confirmCommand "nav" NpgsqlDbType.Numeric (box quote.Nav) |> ignore
+ addParameter confirmCommand "nav_date" NpgsqlDbType.Date (box quote.NavDate) |> ignore
+ addParameter confirmCommand "source" NpgsqlDbType.Text (box source) |> ignore
+ addParameter confirmCommand "revision" NpgsqlDbType.Text (box revision) |> ignore
+ addParameter confirmCommand "collected_at" NpgsqlDbType.TimestampTz (box quote.CollectedAt) |> ignore
+ addParameter confirmCommand "published_at" NpgsqlDbType.TimestampTz publishedParameter |> ignore
+ addParameter confirmCommand "payload_hash" NpgsqlDbType.Text (box quote.PayloadHash) |> ignore
+ addParameter confirmCommand "first_seen_at" NpgsqlDbType.TimestampTz (box quote.FirstSeenAt) |> ignore
+ addParameter confirmCommand "units" NpgsqlDbType.Numeric (box computation.Units) |> ignore
+ addParameter confirmCommand "invested" NpgsqlDbType.Numeric (box computation.InvestedCash) |> ignore
+ addParameter confirmCommand "residual" NpgsqlDbType.Numeric (box computation.ResidualCash) |> ignore
+ addParameter confirmCommand "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore
+ confirmCommand.ExecuteNonQuery() |> ignore
+
+ use positionCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ INSERT INTO fund_positions
+ (fund_id, instrument_code, units, cost_cash, first_confirmed_at, last_confirmed_at)
+ VALUES
+ (@fund_id, @code, @units, @invested, @confirmed_at, @confirmed_at)
+ ON CONFLICT (fund_id, instrument_code) DO UPDATE
+ SET units = fund_positions.units + EXCLUDED.units,
+ cost_cash = fund_positions.cost_cash + EXCLUDED.cost_cash,
+ last_confirmed_at = EXCLUDED.last_confirmed_at
+ """
+
+ addParameter positionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore
+ addParameter positionCommand "code" NpgsqlDbType.Text (box order.FundCode) |> ignore
+ addParameter positionCommand "units" NpgsqlDbType.Numeric (box computation.Units) |> ignore
+ addParameter positionCommand "invested" NpgsqlDbType.Numeric (box computation.InvestedCash) |> ignore
+ addParameter positionCommand "confirmed_at" NpgsqlDbType.TimestampTz (box confirmedAt) |> ignore
+ positionCommand.ExecuteNonQuery() |> ignore
+
+ let publishedJson =
+ match quote.PublishedAt with
+ | Some published -> sprintf "\"%O\"" published
+ | None -> "null"
+
+ let detail =
+ sprintf
+ "{\"fund_id\":\"%O\",\"order_id\":\"%O\",\"units\":\"%M\",\"invested_cash\":\"%M\",\"residual_cash\":\"%M\",\"quote\":{\"nav_date\":\"%s\",\"nav\":\"%M\",\"source\":%s,\"source_revision\":%s,\"source_collected_at\":\"%O\",\"published_at\":%s,\"source_payload_hash\":%s,\"first_seen_at\":\"%O\"}}"
+ fundId
+ orderId
+ computation.Units
+ computation.InvestedCash
+ computation.ResidualCash
+ (quote.NavDate.ToString("yyyy-MM-dd"))
+ quote.Nav
+ (JsonSerializer.Serialize source)
+ (JsonSerializer.Serialize revision)
+ quote.CollectedAt
+ publishedJson
+ (JsonSerializer.Serialize quote.PayloadHash)
+ quote.FirstSeenAt
+
+ use eventCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ INSERT INTO subscription_order_events (order_id, event_type, detail)
+ VALUES (@order_id, 'confirmed', @detail)
+ """
+
+ addParameter eventCommand "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore
+ addParameter eventCommand "detail" NpgsqlDbType.Jsonb (box detail) |> ignore
+ eventCommand.ExecuteNonQuery() |> ignore
+
+ use idempotencyCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ INSERT INTO subscription_confirm_idempotencies
+ (idempotency_key, request_hash, order_id, fund_id)
+ VALUES
+ (@idempotency_key, @request_hash, @order_id, @fund_id)
+ """
+
+ addParameter idempotencyCommand "idempotency_key" NpgsqlDbType.Text (box idempotencyKey) |> ignore
+ addParameter idempotencyCommand "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore
+ addParameter idempotencyCommand "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore
+ addParameter idempotencyCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore
+ idempotencyCommand.ExecuteNonQuery() |> ignore
+
+ transaction.Commit()
+
+ match findOrder connection None orderId with
+ | Some confirmedOrder -> OrderConfirmed confirmedOrder
+ | None -> failwith "confirmed order disappeared after commit"
+ with error ->
+ try
+ transaction.Rollback()
+ with _ ->
+ ()
+
+ raise error
+
+ member _.GetFundPositions(fundId: Guid) =
+ use connection = new NpgsqlConnection(connectionString)
+ connection.Open()
+
+ use command =
+ commandWithTransaction
+ connection
+ None
+ """
+ SELECT p.instrument_code, p.units, p.cost_cash, p.last_confirmed_at,
+ q.nav, q.nav_date, q.source_collected_at
+ FROM fund_positions p
+ LEFT JOIN LATERAL (
+ SELECT nav, nav_date, source_collected_at
+ FROM fund_nav_observations
+ WHERE instrument_code = p.instrument_code
+ AND nav_date <= (now() AT TIME ZONE 'Asia/Shanghai')::date
+ AND nav > 0
+ AND source_collected_at <= now()
+ AND (published_at IS NULL OR published_at <= now())
+ ORDER BY nav_date DESC, source_collected_at DESC
+ LIMIT 1
+ ) q ON true
+ WHERE p.fund_id = @fund_id
+ ORDER BY p.instrument_code
+ """
+
+ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore
+
+ use reader = command.ExecuteReader()
+ let records = ResizeArray<FundPositionRecord>()
+
+ while reader.Read() do
+ let valuationNav =
+ if reader.IsDBNull(4) then None else Some(reader.GetDecimal(4))
+
+ let valuationNavDate =
+ if reader.IsDBNull(5) then None else Some(reader.GetFieldValue<DateOnly>(5))
+
+ let valuationCollectedAt =
+ if reader.IsDBNull(6) then None else Some(reader.GetFieldValue<DateTimeOffset>(6))
+
+ records.Add(
+ {
+ FundId = fundId
+ InstrumentCode = reader.GetString(0)
+ Units = reader.GetDecimal(1)
+ CostCash = reader.GetDecimal(2)
+ LastConfirmedAt = reader.GetFieldValue<DateTimeOffset>(3)
+ ValuationNav = valuationNav
+ ValuationNavDate = valuationNavDate
+ ValuationCollectedAt = valuationCollectedAt
+ }
+ )
+
+ records |> Seq.toList