From ac44797721121cb7e00b380ee5ed5d464b74da6f Mon Sep 17 00:00:00 2001 From: "Somhairle H. Marisol" Date: Mon, 21 Sep 2026 08:13:13 +0800 Subject: Implement subscription order confirmation with locked trade dates and bound quote evidence --- src/FundLab.Api/Persistence.fs | 656 ++++++++++++++++++++++++++++++++++++++++- 1 file changed, 651 insertions(+), 5 deletions(-) (limited to 'src/FundLab.Api') 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 = + 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 + } [] 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(14) + Source = reader.GetString(15) + Revision = reader.GetString(16) + CollectedAt = reader.GetFieldValue(17) + PublishedAt = + if reader.IsDBNull(18) then + None + else + Some(reader.GetFieldValue(18)) + PayloadHash = + if reader.IsDBNull(19) then + "" + else + reader.GetString(19) + FirstSeenAt = + if reader.IsDBNull(20) then + DateTimeOffset.UnixEpoch + else + reader.GetFieldValue(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(8) + TradeDate = reader.GetFieldValue(9) + ConfirmIdempotencyKey = readStringOption reader 10 + PendingReason = readStringOption reader 11 + ConfirmedAt = + if reader.IsDBNull(12) then + None + else + Some(reader.GetFieldValue(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(1) + Source = quoteReader.GetString(2) + Revision = quoteReader.GetString(3) + CollectedAt = quoteReader.GetFieldValue(4) + PublishedAt = + if quoteReader.IsDBNull(5) then + None + else + Some(quoteReader.GetFieldValue(5)) + PayloadHash = quoteReader.GetString(6) + FirstSeenAt = quoteReader.GetFieldValue(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() + + 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(5)) + + let valuationCollectedAt = + if reader.IsDBNull(6) then None else Some(reader.GetFieldValue(6)) + + records.Add( + { + FundId = fundId + InstrumentCode = reader.GetString(0) + Units = reader.GetDecimal(1) + CostCash = reader.GetDecimal(2) + LastConfirmedAt = reader.GetFieldValue(3) + ValuationNav = valuationNav + ValuationNavDate = valuationNavDate + ValuationCollectedAt = valuationCollectedAt + } + ) + + records |> Seq.toList -- cgit v1.2.3