namespace FundLab.Api open System open System.Data open System.Data.Common open System.Globalization open System.Security.Cryptography 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 = { Name: string InitialCash: decimal InitialUnitNav: decimal IsSynthetic: bool } type FundRecord = { Id: Guid Name: string Currency: string InitialCash: decimal InitialUnitNav: decimal IsSynthetic: bool AvailableCash: decimal ReservedCash: decimal Status: string } type FundWriteResult = | Created of FundRecord | Replayed of FundRecord | IdempotencyConflict | Invalid of string [] type SubscriptionOrderCommand = { FundCode: string Amount: decimal 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 FundId: Guid FundCode: string Amount: decimal FeeAmount: decimal ReservedTotal: decimal 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 = | OrderCreated of SubscriptionOrderRecord | OrderReplayed of SubscriptionOrderRecord | OrderIdempotencyConflict | OrderInvalid of string | OrderFundNotFound | 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 ReservedUnits: decimal CostCash: decimal LastConfirmedAt: DateTimeOffset ValuationNav: decimal option ValuationNavDate: DateOnly option ValuationCollectedAt: DateTimeOffset option } type RedemptionCommand = { InstrumentCode: string Units: decimal FeeAmount: decimal } type RedemptionOrderRecord = { Id: Guid FundId: Guid InstrumentCode: string Units: decimal FeeAmount: decimal Status: string IsSynthetic: bool SubmittedAt: DateTimeOffset TradeDate: DateOnly ConfirmIdempotencyKey: string option PendingReason: string option ConfirmedAt: DateTimeOffset option ConfirmedNav: decimal option ConfirmedNavDate: DateOnly option ConfirmedProceeds: decimal option ConfirmedCostReleased: decimal option } type RedemptionWriteResult = | RedemptionCreated of RedemptionOrderRecord | RedemptionReplayed of RedemptionOrderRecord | RedemptionIdempotencyConflict | RedemptionInvalid of string | RedemptionFundNotFound | RedemptionInstrumentNotFound | RedemptionInsufficientUnits type RedemptionConfirmResult = | RedemptionConfirmed of RedemptionOrderRecord | RedemptionConfirmReplayed of RedemptionOrderRecord | RedemptionPendingNav of RedemptionOrderRecord | RedemptionConfirmIdempotencyConflict | RedemptionAlreadyConfirmed | RedemptionOrderNotFound | RedemptionInvalidStatus | RedemptionInvalid of string type SipPlanCommand = { InstrumentCode: string Amount: decimal Frequency: SipFrequency } type SipPlanRecord = { Id: Guid FundId: Guid InstrumentCode: string Amount: decimal Frequency: SipFrequency Status: string IsSynthetic: bool AnchorDate: DateOnly NextTradeDate: DateOnly CreatedAt: DateTimeOffset LastExecutionStatus: string option LastExecutionDate: DateOnly option } type SipPlanWriteResult = | SipPlanCreated of SipPlanRecord | SipPlanReplayed of SipPlanRecord | SipPlanIdempotencyConflict | SipPlanInvalid of string | SipPlanFundNotFound | SipPlanInstrumentNotFound type SipExecutionOutcome = { TradeDate: DateOnly Status: string OrderId: Guid option PendingReason: string option } type SipPlanAdvanceResult = { PlanId: Guid InstrumentCode: string Amount: decimal Frequency: SipFrequency Executions: SipExecutionOutcome list NextTradeDate: DateOnly } type SipAdvanceResult = { FundId: Guid Plans: SipPlanAdvanceResult list } type RebalanceTarget = RebalancePolicy.TargetAllocation type RebalancePlanCommand = { Targets: RebalanceTarget list } type RebalancePlanRecord = { Id: Guid FundId: Guid Targets: RebalanceTarget list Status: string IsSynthetic: bool CreatedAt: DateTimeOffset } type RebalanceWriteResult = | RebalancePlanCreated of RebalancePlanRecord | RebalancePlanReplayed of RebalancePlanRecord | RebalanceIdempotencyConflict | RebalanceInvalid of string | RebalanceFundNotFound | RebalanceInstrumentNotFound type RebalanceOrderOutcome = { InstrumentCode: string Action: string Amount: string Status: string OrderId: Guid option PendingReason: string option } type RebalanceExecutionResult = { PlanId: Guid RunDate: DateOnly Outcomes: RebalanceOrderOutcome list } type CapitalDepositCommand = { Amount: decimal Note: string option } type CapitalDepositRecord = { Id: Guid FundId: Guid Amount: decimal Note: string option IsSynthetic: bool CreatedAt: DateTimeOffset } type CapitalDepositWriteResult = | CapitalDepositCreated of CapitalDepositRecord | CapitalDepositReplayed of CapitalDepositRecord | CapitalDepositIdempotencyConflict | CapitalDepositInvalid of string | CapitalDepositFundNotFound type FundRepository(connectionString: string) = let cashMaximum = 999999999999999999.99m let unitNavMaximum = 99999999999999999999.99999999m let schema = """ CREATE TABLE IF NOT EXISTS funds ( id uuid PRIMARY KEY, name text NOT NULL, currency text NOT NULL, initial_cash numeric(20, 2) NOT NULL, initial_unit_nav numeric(28, 8) NOT NULL, is_synthetic boolean NOT NULL, available_cash numeric(20, 2) NOT NULL, status text NOT NULL, created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS fund_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS instruments ( code text PRIMARY KEY CHECK (code ~ '^[0-9]{6}$'), name text NOT NULL, fund_type text NULL, source text NOT NULL, source_revision text NOT NULL, source_collected_at timestamptz NOT NULL, source_payload_hash text NOT NULL, first_seen_at timestamptz NOT NULL DEFAULT now(), last_seen_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS fund_nav_observations ( instrument_code text NOT NULL REFERENCES instruments(code), nav_date date NOT NULL, published_at timestamptz NULL, nav numeric(28, 8) NOT NULL, accumulated_nav numeric(28, 8) NULL, daily_return numeric(20, 8) NULL, source text NOT NULL, source_revision text NOT NULL, source_collected_at timestamptz NOT NULL, source_payload_hash text NOT NULL, first_seen_at timestamptz NOT NULL DEFAULT now(), last_seen_at timestamptz NOT NULL DEFAULT now(), PRIMARY KEY (instrument_code, nav_date) ); 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 ( id uuid PRIMARY KEY, fund_id uuid NOT NULL REFERENCES funds(id), fund_code text NOT NULL, amount numeric(20, 2) NOT NULL CHECK (amount > 0), fee_amount numeric(20, 2) NOT NULL CHECK (fee_amount >= 0), reserved_total numeric(20, 2) NOT NULL CHECK (reserved_total > 0), status text NOT NULL, is_synthetic boolean NOT NULL, submitted_at timestamptz NOT NULL DEFAULT now() ); ALTER TABLE subscription_orders DROP CONSTRAINT IF EXISTS subscription_orders_fee_amount_check; ALTER TABLE subscription_orders ADD CONSTRAINT subscription_orders_fee_amount_check CHECK (fee_amount >= 0); CREATE TABLE IF NOT EXISTS subscription_order_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() ); 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() ); CREATE TABLE IF NOT EXISTS redemption_orders ( id uuid PRIMARY KEY, 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), fee_amount numeric(20, 2) NOT NULL CHECK (fee_amount >= 0), status text NOT NULL, is_synthetic boolean NOT NULL, submitted_at timestamptz NOT NULL DEFAULT now(), trade_date date NOT NULL, confirm_idempotency_key text NULL, pending_reason text NULL, confirmed_at timestamptz NULL, confirmed_nav numeric(28, 8) NULL, confirmed_nav_date date NULL, confirmed_proceeds numeric(20, 2) NULL, confirmed_cost_released numeric(20, 2) NULL ); CREATE TABLE IF NOT EXISTS redemption_order_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, order_id uuid NOT NULL REFERENCES redemption_orders(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS redemption_confirm_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, order_id uuid NOT NULL REFERENCES redemption_orders(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); ALTER TABLE fund_positions ADD COLUMN IF NOT EXISTS reserved_units numeric(28, 8) NOT NULL DEFAULT 0; CREATE TABLE IF NOT EXISTS fund_capital_deposits ( id uuid PRIMARY KEY, fund_id uuid NOT NULL REFERENCES funds(id), amount numeric(20, 2) NOT NULL CHECK (amount > 0), note text NULL, is_synthetic boolean NOT NULL, created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS fund_capital_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, deposit_id uuid NOT NULL REFERENCES fund_capital_deposits(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS sip_plans ( id uuid PRIMARY KEY, fund_id uuid NOT NULL REFERENCES funds(id), instrument_code text NOT NULL REFERENCES instruments(code), amount numeric(20, 2) NOT NULL CHECK (amount > 0), frequency text NOT NULL, status text NOT NULL, is_synthetic boolean NOT NULL, anchor_date date NOT NULL, next_trade_date date NOT NULL, created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS sip_plan_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, plan_id uuid NOT NULL REFERENCES sip_plans(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS sip_executions ( plan_id uuid NOT NULL REFERENCES sip_plans(id), trade_date date NOT NULL, amount numeric(20, 2) NOT NULL, fee_amount numeric(20, 2) NOT NULL, status text NOT NULL, order_id uuid NULL, pending_reason text NULL, executed_at timestamptz NOT NULL DEFAULT now(), PRIMARY KEY (plan_id, trade_date) ); CREATE TABLE IF NOT EXISTS rebalance_plans ( id uuid PRIMARY KEY, fund_id uuid NOT NULL REFERENCES funds(id), status text NOT NULL, is_synthetic boolean NOT NULL, created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS rebalance_targets ( plan_id uuid NOT NULL REFERENCES rebalance_plans(id), instrument_code text NOT NULL REFERENCES instruments(code), target_percent numeric(9, 2) NOT NULL CHECK (target_percent > 0 AND target_percent <= 100), PRIMARY KEY (plan_id, instrument_code) ); """ let statusText status = match status with | FundStatus.Empty -> "empty" | FundStatus.Active -> "active" | FundStatus.ZeroUnits -> "zero_units" let cashText (value: decimal) = value.ToString("0.00", CultureInfo.InvariantCulture) let decimalText (value: decimal) = value.ToString("0.00000000", CultureInfo.InvariantCulture) let recordFromLedger (fund: LedgerFund) = { Id = fund.Id Name = fund.Name Currency = fund.Currency InitialCash = fund.InitialCash InitialUnitNav = fund.InitialUnitNav IsSynthetic = fund.IsSynthetic AvailableCash = fund.AvailableCash ReservedCash = fund.FrozenCash Status = statusText fund.Status } let recordFromReader (reader: DbDataReader) = { Id = reader.GetGuid(0) Name = reader.GetString(1) Currency = reader.GetString(2) InitialCash = reader.GetDecimal(3) InitialUnitNav = reader.GetDecimal(4) IsSynthetic = reader.GetBoolean(5) AvailableCash = reader.GetDecimal(6) ReservedCash = reader.GetDecimal(7) Status = reader.GetString(8) } let dateTimeOffsetFromReader (reader: DbDataReader) index = reader.GetFieldValue(index) let optionalDateTimeOffsetFromReader (reader: DbDataReader) index = if reader.IsDBNull(index) then None else Some(dateTimeOffsetFromReader reader index) let optionalDecimalFromReader (reader: DbDataReader) index = if reader.IsDBNull(index) then None else Some(reader.GetDecimal(index)) let instrumentRecordFromReader (reader: DbDataReader) : MarketDataInstrumentRecord = { Code = reader.GetString(0) Name = reader.GetString(1) FundType = if reader.IsDBNull(2) then None else Some(reader.GetString(2)) Source = reader.GetString(3) SourceRevision = reader.GetString(4) SourceCollectedAt = dateTimeOffsetFromReader reader 5 SourcePayloadHash = reader.GetString(6) FirstSeenAt = dateTimeOffsetFromReader reader 7 LastSeenAt = dateTimeOffsetFromReader reader 8 } let navRecordFromReader (reader: DbDataReader) : MarketDataNavRecord = { Code = reader.GetString(0) NavDate = reader.GetFieldValue(1) PublishedAt = optionalDateTimeOffsetFromReader reader 2 Nav = reader.GetDecimal(3) AccumulatedNav = optionalDecimalFromReader reader 4 DailyReturn = optionalDecimalFromReader reader 5 Source = reader.GetString(6) SourceRevision = reader.GetString(7) SourceCollectedAt = dateTimeOffsetFromReader reader 8 SourcePayloadHash = reader.GetString(9) FirstSeenAt = dateTimeOffsetFromReader reader 10 LastSeenAt = dateTimeOffsetFromReader reader 11 } let commandWithTransaction (connection: NpgsqlConnection) (transaction: NpgsqlTransaction option) sql = let command = connection.CreateCommand() command.CommandText <- sql match transaction with | Some value -> command.Transaction <- value | None -> () command let addParameter (command: NpgsqlCommand) name dbType (value: obj) = let parameter = command.Parameters.Add(name, dbType) parameter.Value <- value parameter let optionalParameterValue value = match value with | Some actual -> box actual | None -> box DBNull.Value let findFund connection transaction fundId = use command = commandWithTransaction connection transaction """ SELECT id, name, currency, initial_cash, initial_unit_nav, is_synthetic, available_cash, reserved_cash, status FROM funds WHERE id = @fund_id """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(recordFromReader reader) else None let findIdempotency connection transaction key = use command = commandWithTransaction connection transaction "SELECT request_hash, fund_id FROM fund_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)) else None let insertFund connection transaction (fund: LedgerFund) = use command = commandWithTransaction connection transaction """ INSERT INTO funds (id, name, currency, initial_cash, initial_unit_nav, is_synthetic, available_cash, reserved_cash, status) VALUES (@id, @name, @currency, @initial_cash, @initial_unit_nav, @is_synthetic, @available_cash, @reserved_cash, @status) """ addParameter command "id" NpgsqlDbType.Uuid (box fund.Id) |> ignore addParameter command "name" NpgsqlDbType.Text (box fund.Name) |> ignore addParameter command "currency" NpgsqlDbType.Text (box fund.Currency) |> ignore addParameter command "initial_cash" NpgsqlDbType.Numeric (box fund.InitialCash) |> ignore addParameter command "initial_unit_nav" NpgsqlDbType.Numeric (box fund.InitialUnitNav) |> ignore addParameter command "is_synthetic" NpgsqlDbType.Boolean (box fund.IsSynthetic) |> ignore addParameter command "available_cash" NpgsqlDbType.Numeric (box fund.AvailableCash) |> ignore addParameter command "reserved_cash" NpgsqlDbType.Numeric (box fund.FrozenCash) |> ignore addParameter command "status" NpgsqlDbType.Text (box (statusText fund.Status)) |> ignore command.ExecuteNonQuery() |> ignore let insertIdempotency connection transaction key requestHash fundId = use command = commandWithTransaction connection transaction """ INSERT INTO fund_idempotencies (idempotency_key, request_hash, fund_id) VALUES (@idempotency_key, @request_hash, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore 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) FundCode = reader.GetString(2) Amount = reader.GetDecimal(3) FeeAmount = reader.GetDecimal(4) ReservedTotal = reader.GetDecimal(5) 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 = use command = commandWithTransaction connection transaction """ SELECT id, fund_id, fund_code, amount, fee_amount, 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 """ addParameter command "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore 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 connection transaction "SELECT request_hash, fund_id, order_id FROM subscription_order_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 instrumentExists connection transaction code = use command = commandWithTransaction connection transaction "SELECT 1 FROM instruments WHERE code = @code" addParameter command "code" NpgsqlDbType.Text (box code) |> ignore use reader = command.ExecuteReader() reader.Read() let lockFundForOrder connection transaction fundId = use command = commandWithTransaction connection transaction "SELECT is_synthetic FROM funds WHERE id = @fund_id FOR UPDATE" addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(reader.GetBoolean(0)) else None let insertSubscriptionOrder connection transaction (order: SubscriptionOrderRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO subscription_orders (id, fund_id, fund_code, amount, fee_amount, reserved_total, status, is_synthetic, trade_date) VALUES (@id, @fund_id, @fund_code, @amount, @fee_amount, @reserved_total, @status, @is_synthetic, @trade_date) RETURNING submitted_at """ addParameter command "id" NpgsqlDbType.Uuid (box order.Id) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box order.FundId) |> ignore addParameter command "fund_code" NpgsqlDbType.Text (box order.FundCode) |> ignore addParameter command "amount" NpgsqlDbType.Numeric (box order.Amount) |> ignore addParameter command "fee_amount" NpgsqlDbType.Numeric (box order.FeeAmount) |> ignore 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 reader.GetFieldValue(0) let insertOrderIdempotency connection transaction key requestHash orderId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO subscription_order_idempotencies (idempotency_key, request_hash, order_id, fund_id) VALUES (@idempotency_key, @request_hash, @order_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let upsertInstrument connection transaction (payload: MarketDataSearchPayload) payloadHash (instrument: MarketDataInstrument) = use command = commandWithTransaction connection transaction """ INSERT INTO instruments (code, name, fund_type, source, source_revision, source_collected_at, source_payload_hash) VALUES (@code, @name, @fund_type, @source, @source_revision, @source_collected_at, @source_payload_hash) ON CONFLICT (code) DO UPDATE SET name = EXCLUDED.name, fund_type = EXCLUDED.fund_type, source = EXCLUDED.source, source_revision = EXCLUDED.source_revision, source_collected_at = EXCLUDED.source_collected_at, source_payload_hash = EXCLUDED.source_payload_hash, last_seen_at = now() """ addParameter command "code" NpgsqlDbType.Text (box instrument.Code) |> ignore addParameter command "name" NpgsqlDbType.Text (box instrument.Name) |> ignore addParameter command "fund_type" NpgsqlDbType.Text (optionalParameterValue instrument.FundType) |> ignore addParameter command "source" NpgsqlDbType.Text (box payload.Source) |> ignore addParameter command "source_revision" NpgsqlDbType.Text (box payload.SourceRevision) |> ignore addParameter command "source_collected_at" NpgsqlDbType.TimestampTz (box payload.CollectedAt) |> ignore addParameter command "source_payload_hash" NpgsqlDbType.Text (box payloadHash) |> ignore command.ExecuteNonQuery() |> ignore let upsertNavObservation connection transaction (payload: MarketDataNavPayload) payloadHash (observation: MarketDataObservation) = use command = commandWithTransaction connection transaction """ INSERT INTO fund_nav_observations (instrument_code, nav_date, published_at, nav, accumulated_nav, daily_return, source, source_revision, source_collected_at, source_payload_hash) VALUES (@instrument_code, @nav_date, @published_at, @nav, @accumulated_nav, @daily_return, @source, @source_revision, @source_collected_at, @source_payload_hash) ON CONFLICT (instrument_code, nav_date) DO UPDATE SET published_at = EXCLUDED.published_at, nav = EXCLUDED.nav, accumulated_nav = EXCLUDED.accumulated_nav, daily_return = EXCLUDED.daily_return, source = EXCLUDED.source, source_revision = EXCLUDED.source_revision, source_collected_at = EXCLUDED.source_collected_at, source_payload_hash = EXCLUDED.source_payload_hash, 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 addParameter command "nav_date" NpgsqlDbType.Date (box observation.NavDate) |> ignore addParameter command "published_at" NpgsqlDbType.TimestampTz (optionalParameterValue observation.PublishedAt) |> ignore addParameter command "nav" NpgsqlDbType.Numeric (box observation.Nav) |> ignore addParameter command "accumulated_nav" NpgsqlDbType.Numeric (optionalParameterValue observation.AccumulatedNav) |> ignore addParameter command "daily_return" NpgsqlDbType.Numeric (optionalParameterValue observation.DailyReturn) |> ignore addParameter command "source" NpgsqlDbType.Text (box payload.Source) |> ignore addParameter command "source_revision" NpgsqlDbType.Text (box payload.SourceRevision) |> ignore addParameter command "source_collected_at" NpgsqlDbType.TimestampTz (box payload.CollectedAt) |> ignore 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 let name = if isNull command.Name then "" else command.Name let payload = String.concat "|" [ "fund-create" encoded name (encoded (command.InitialCash.ToString("G29", invariant))) (encoded (command.InitialUnitNav.ToString("G29", invariant))) (encoded (command.IsSynthetic.ToString())) ] Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let validateStorageRange (command: FundCreateCommand) = if command.InitialCash > cashMaximum then Error "initial cash exceeds database precision" elif command.InitialUnitNav > unitNavMaximum then Error "initial unit NAV exceeds database precision" else Ok() let orderRequestHash (fundId: Guid) (command: SubscriptionOrderCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let fundCode = if isNull command.FundCode then "" else command.FundCode let payload = String.concat "|" [ "subscription-order" encoded (fundId.ToString("D")) encoded fundCode (encoded (command.Amount.ToString("G29", invariant))) (encoded (command.FeeAmount.ToString("G29", invariant))) ] Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let validateOrderCommand (command: SubscriptionOrderCommand) = if String.IsNullOrWhiteSpace command.FundCode then Error "fund code cannot be empty" elif command.Amount <= 0m then Error "amount must be positive" elif command.FeeAmount < 0m then Error "fee amount cannot be negative" elif Decimal.Round(command.Amount, 2) <> command.Amount then Error "amount exceeds cash precision" elif Decimal.Round(command.FeeAmount, 2) <> command.FeeAmount then Error "fee amount exceeds cash precision" elif command.Amount > cashMaximum then Error "amount exceeds database precision" elif command.FeeAmount > cashMaximum then Error "fee amount exceeds database precision" elif command.Amount + command.FeeAmount > cashMaximum then Error "reserved total exceeds database precision" else Ok() let ledgerErrorMessage error = match error with | InvalidIdentifier label -> sprintf "%s is invalid" label | InvalidAmount label -> sprintf "%s is invalid" label | InvalidPrecision label -> sprintf "%s has invalid precision" label | InvalidState message -> message | FundAlreadyExists fundId -> sprintf "fund %O already exists" fundId | other -> sprintf "%A" other let redemptionRecordFromReader (reader: DbDataReader) : RedemptionOrderRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) InstrumentCode = reader.GetString(2) Units = reader.GetDecimal(3) FeeAmount = reader.GetDecimal(4) Status = reader.GetString(5) IsSynthetic = reader.GetBoolean(6) SubmittedAt = reader.GetFieldValue(7) TradeDate = reader.GetFieldValue(8) ConfirmIdempotencyKey = readStringOption reader 9 PendingReason = readStringOption reader 10 ConfirmedAt = optionalDateTimeOffsetFromReader reader 11 ConfirmedNav = readDecimalOption reader 12 ConfirmedNavDate = if reader.IsDBNull(13) then None else Some(reader.GetFieldValue(13)) ConfirmedProceeds = readDecimalOption reader 14 ConfirmedCostReleased = readDecimalOption reader 15 } let redemptionOrderColumns = """ SELECT id, fund_id, instrument_code, units, fee_amount, status, is_synthetic, submitted_at, trade_date, confirm_idempotency_key, pending_reason, confirmed_at, confirmed_nav, confirmed_nav_date, confirmed_proceeds, confirmed_cost_released FROM redemption_orders """ let findRedemptionOrder connection transaction orderId = use command = commandWithTransaction connection transaction (redemptionOrderColumns + " WHERE id = @order_id") addParameter command "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(redemptionRecordFromReader reader) else None let findRedemptionOrderIdempotency connection transaction key = use command = commandWithTransaction connection transaction "SELECT request_hash, fund_id, order_id FROM redemption_order_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 findRedemptionConfirmIdempotency connection transaction key = use command = commandWithTransaction connection transaction "SELECT request_hash, fund_id, order_id FROM redemption_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 insertRedemptionOrder connection transaction (order: RedemptionOrderRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO redemption_orders (id, fund_id, instrument_code, units, fee_amount, status, is_synthetic, trade_date) VALUES (@id, @fund_id, @instrument_code, @units, @fee_amount, @status, @is_synthetic, @trade_date) RETURNING submitted_at """ addParameter command "id" NpgsqlDbType.Uuid (box order.Id) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box order.FundId) |> ignore addParameter command "instrument_code" NpgsqlDbType.Text (box order.InstrumentCode) |> ignore addParameter command "units" NpgsqlDbType.Numeric (box order.Units) |> ignore addParameter command "fee_amount" NpgsqlDbType.Numeric (box order.FeeAmount) |> 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 reader.GetFieldValue(0) let insertRedemptionOrderIdempotency connection transaction key requestHash orderId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO redemption_order_idempotencies (idempotency_key, request_hash, order_id, fund_id) VALUES (@idempotency_key, @request_hash, @order_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let insertRedemptionConfirmIdempotency connection transaction key requestHash orderId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO redemption_confirm_idempotencies (idempotency_key, request_hash, order_id, fund_id) VALUES (@idempotency_key, @request_hash, @order_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let redemptionRequestHash (fundId: Guid) (command: RedemptionCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let code = if isNull command.InstrumentCode then "" else command.InstrumentCode let payload = String.concat "|" [ "redemption-order" encoded (fundId.ToString("D")) encoded code (encoded (command.Units.ToString("G29", invariant))) (encoded (command.FeeAmount.ToString("G29", invariant))) ] Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let validateRedemptionCommand (command: RedemptionCommand) = if String.IsNullOrWhiteSpace command.InstrumentCode then Error "instrument code cannot be empty" elif command.Units <= 0m then Error "units must be positive" elif command.FeeAmount < 0m then Error "fee amount cannot be negative" elif Decimal.Round(command.Units, 8) <> command.Units then Error "units exceed supported precision" elif Decimal.Round(command.FeeAmount, 2) <> command.FeeAmount then Error "fee amount exceeds cash precision" elif command.Units > RedemptionPolicy.unitsMaximum then Error "units exceed database precision" elif command.FeeAmount > cashMaximum then Error "fee amount exceeds database precision" else Ok() let capitalDepositRecordFromReader (reader: DbDataReader) : CapitalDepositRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) Amount = reader.GetDecimal(2) Note = readStringOption reader 3 IsSynthetic = reader.GetBoolean(4) CreatedAt = reader.GetFieldValue(5) } let findCapitalDeposit connection transaction depositId = use command = commandWithTransaction connection transaction "SELECT id, fund_id, amount, note, is_synthetic, created_at FROM fund_capital_deposits WHERE id = @deposit_id" addParameter command "deposit_id" NpgsqlDbType.Uuid (box depositId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(capitalDepositRecordFromReader reader) else None let findCapitalDepositIdempotency connection transaction key = use command = commandWithTransaction connection transaction "SELECT request_hash, fund_id, deposit_id FROM fund_capital_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 insertCapitalDeposit connection transaction (deposit: CapitalDepositRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO fund_capital_deposits (id, fund_id, amount, note, is_synthetic) VALUES (@id, @fund_id, @amount, @note, @is_synthetic) RETURNING created_at """ addParameter command "id" NpgsqlDbType.Uuid (box deposit.Id) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box deposit.FundId) |> ignore addParameter command "amount" NpgsqlDbType.Numeric (box deposit.Amount) |> ignore let noteParameter = match deposit.Note with | Some note -> box note | None -> box DBNull.Value addParameter command "note" NpgsqlDbType.Text noteParameter |> ignore addParameter command "is_synthetic" NpgsqlDbType.Boolean (box deposit.IsSynthetic) |> ignore use reader = command.ExecuteReader() reader.Read() |> ignore reader.GetFieldValue(0) let insertCapitalDepositIdempotency connection transaction key requestHash depositId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO fund_capital_idempotencies (idempotency_key, request_hash, deposit_id, fund_id) VALUES (@idempotency_key, @request_hash, @deposit_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "deposit_id" NpgsqlDbType.Uuid (box depositId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let capitalDepositRequestHash (fundId: Guid) (command: CapitalDepositCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let note = if command.Note.IsNone then "" else command.Note.Value let payload = String.concat "|" [ "capital-deposit" encoded (fundId.ToString("D")) (encoded (command.Amount.ToString("G29", invariant))) (encoded note) ] Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let sipPlanRecordFromReader (reader: DbDataReader) : SipPlanRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) InstrumentCode = reader.GetString(2) Amount = reader.GetDecimal(3) Frequency = match SipPolicy.parseFrequency (reader.GetString(4)) with | Some frequency -> frequency | None -> failwith "sip plan frequency is invalid" Status = reader.GetString(5) IsSynthetic = reader.GetBoolean(6) AnchorDate = reader.GetFieldValue(7) NextTradeDate = reader.GetFieldValue(8) CreatedAt = reader.GetFieldValue(9) LastExecutionStatus = readStringOption reader 10 LastExecutionDate = if reader.IsDBNull(11) then None else Some(reader.GetFieldValue(11)) } let findSipPlan connection transaction planId = use command = commandWithTransaction connection transaction """ SELECT p.id, p.fund_id, p.instrument_code, p.amount, p.frequency, p.status, p.is_synthetic, p.anchor_date, p.next_trade_date, p.created_at, e.status, e.trade_date FROM sip_plans p LEFT JOIN LATERAL ( SELECT status, trade_date FROM sip_executions WHERE plan_id = p.id ORDER BY executed_at DESC, trade_date DESC LIMIT 1 ) e ON true WHERE p.id = @plan_id """ addParameter command "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(sipPlanRecordFromReader reader) else None let findSipPlanIdempotency connection transaction key = use command = commandWithTransaction connection transaction "SELECT request_hash, fund_id, plan_id FROM sip_plan_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 insertSipPlan connection transaction (plan: SipPlanRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO sip_plans (id, fund_id, instrument_code, amount, frequency, status, is_synthetic, anchor_date, next_trade_date) VALUES (@id, @fund_id, @instrument_code, @amount, @frequency, @status, @is_synthetic, @anchor_date, @next_trade_date) RETURNING created_at """ addParameter command "id" NpgsqlDbType.Uuid (box plan.Id) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box plan.FundId) |> ignore addParameter command "instrument_code" NpgsqlDbType.Text (box plan.InstrumentCode) |> ignore addParameter command "amount" NpgsqlDbType.Numeric (box plan.Amount) |> ignore addParameter command "frequency" NpgsqlDbType.Text (box (SipPolicy.frequencyText plan.Frequency)) |> ignore addParameter command "status" NpgsqlDbType.Text (box plan.Status) |> ignore addParameter command "is_synthetic" NpgsqlDbType.Boolean (box plan.IsSynthetic) |> ignore addParameter command "anchor_date" NpgsqlDbType.Date (box plan.AnchorDate) |> ignore addParameter command "next_trade_date" NpgsqlDbType.Date (box plan.NextTradeDate) |> ignore use reader = command.ExecuteReader() reader.Read() |> ignore reader.GetFieldValue(0) let insertSipPlanIdempotency connection transaction key requestHash planId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO sip_plan_idempotencies (idempotency_key, request_hash, plan_id, fund_id) VALUES (@idempotency_key, @request_hash, @plan_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let sipPlanRequestHash (fundId: Guid) (command: SipPlanCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let code = if isNull command.InstrumentCode then "" else command.InstrumentCode let payload = String.concat "|" [ "sip-plan" encoded (fundId.ToString("D")) encoded code (encoded (command.Amount.ToString("G29", invariant))) (encoded (SipPolicy.frequencyText command.Frequency)) ] Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let rebalancePlanRecordFromReader (reader: DbDataReader) : RebalancePlanRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) Targets = [] Status = reader.GetString(2) IsSynthetic = reader.GetBoolean(3) CreatedAt = reader.GetFieldValue(4) } let rebalanceTargetsFromReader (reader: DbDataReader) : RebalanceTarget list = let targets = ResizeArray() while reader.Read() do targets.Add( { InstrumentCode = reader.GetString(0) TargetPercent = reader.GetDecimal(1) } ) targets |> Seq.toList let rebalancePlanWithTargets connection transaction (plan: RebalancePlanRecord) = use targetsCommand = commandWithTransaction connection transaction "SELECT instrument_code, target_percent FROM rebalance_targets WHERE plan_id = @plan_id ORDER BY instrument_code" addParameter targetsCommand "plan_id" NpgsqlDbType.Uuid (box plan.Id) |> ignore use reader = targetsCommand.ExecuteReader() let targets = rebalanceTargetsFromReader reader { plan with Targets = targets } let findRebalancePlan connection transaction planId = use command = commandWithTransaction connection transaction "SELECT id, fund_id, status, is_synthetic, created_at FROM rebalance_plans WHERE id = @plan_id" addParameter command "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore use reader = command.ExecuteReader() if reader.Read() then let plan = rebalancePlanRecordFromReader reader reader.Close() Some(rebalancePlanWithTargets connection transaction plan) else None let findRebalanceIdempotency connection transaction key = use command = commandWithTransaction connection transaction "SELECT request_hash, fund_id, plan_id FROM rebalance_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 insertRebalancePlan connection transaction (plan: RebalancePlanRecord) = use planCommand = commandWithTransaction connection transaction """ INSERT INTO rebalance_plans (id, fund_id, status, is_synthetic) VALUES (@id, @fund_id, @status, @is_synthetic) RETURNING created_at """ addParameter planCommand "id" NpgsqlDbType.Uuid (box plan.Id) |> ignore addParameter planCommand "fund_id" NpgsqlDbType.Uuid (box plan.FundId) |> ignore addParameter planCommand "status" NpgsqlDbType.Text (box plan.Status) |> ignore addParameter planCommand "is_synthetic" NpgsqlDbType.Boolean (box plan.IsSynthetic) |> ignore use reader = planCommand.ExecuteReader() reader.Read() |> ignore let createdAt = reader.GetFieldValue(0) reader.Close() for target in plan.Targets do use targetCommand = commandWithTransaction connection transaction """ INSERT INTO rebalance_targets (plan_id, instrument_code, target_percent) VALUES (@plan_id, @instrument_code, @target_percent) """ addParameter targetCommand "plan_id" NpgsqlDbType.Uuid (box plan.Id) |> ignore addParameter targetCommand "instrument_code" NpgsqlDbType.Text (box target.InstrumentCode) |> ignore addParameter targetCommand "target_percent" NpgsqlDbType.Numeric (box target.TargetPercent) |> ignore targetCommand.ExecuteNonQuery() |> ignore createdAt let insertRebalanceIdempotency connection transaction key requestHash planId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO rebalance_idempotencies (idempotency_key, request_hash, plan_id, fund_id) VALUES (@idempotency_key, @request_hash, @plan_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let rebalanceRequestHash (fundId: Guid) (command: RebalancePlanCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let payload = String.concat "|" ([ "rebalance-plan"; encoded (fundId.ToString("D")) ] @ (command.Targets |> List.sortBy (fun target -> target.InstrumentCode) |> List.collect (fun target -> [ encoded target.InstrumentCode encoded (target.TargetPercent.ToString("G29", invariant)) ]))) Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let validateRebalanceCommand (command: RebalancePlanCommand) = match RebalancePolicy.validateTargets command.Targets with | Error message -> Error message | Ok() -> Ok() member _.EnsureSchema() = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = connection.CreateCommand() command.CommandText <- schema command.ExecuteNonQuery() |> ignore member _.GetFund(fundId: Guid) = use connection = new NpgsqlConnection(connectionString) connection.Open() findFund connection None fundId member _.UpsertInstruments(payload: MarketDataSearchPayload, payloadHash: string) = use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try for instrument in payload.Instruments do upsertInstrument connection (Some transaction) payload payloadHash instrument transaction.Commit() with error -> try transaction.Rollback() with _ -> () raise error member _.GetInstrument(code: string) = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT code, name, fund_type, source, source_revision, source_collected_at, source_payload_hash, first_seen_at, last_seen_at FROM instruments WHERE code = @code """ addParameter command "code" NpgsqlDbType.Text (box code) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(instrumentRecordFromReader reader) else None member _.UpsertNavObservations(payload: MarketDataNavPayload, payloadHash: string) = use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try for observation in payload.Observations do upsertNavObservation connection (Some transaction) payload payloadHash observation transaction.Commit() with error -> try transaction.Rollback() with _ -> () raise error member _.GetNav(code: string, fromDate: DateOnly option, toDate: DateOnly option) = use connection = new NpgsqlConnection(connectionString) connection.Open() let conditions = ResizeArray() conditions.Add("instrument_code = @instrument_code") if fromDate.IsSome then conditions.Add("nav_date >= @from_date") if toDate.IsSome then conditions.Add("nav_date <= @to_date") use command = commandWithTransaction connection None (sprintf """ SELECT instrument_code, nav_date, published_at, nav, accumulated_nav, daily_return, source, source_revision, source_collected_at, source_payload_hash, first_seen_at, last_seen_at FROM fund_nav_observations WHERE %s ORDER BY nav_date ASC """ (String.concat " AND " conditions)) addParameter command "instrument_code" NpgsqlDbType.Text (box code) |> ignore match fromDate with | Some value -> addParameter command "from_date" NpgsqlDbType.Date (box value) |> ignore | None -> () match toDate with | Some value -> addParameter command "to_date" NpgsqlDbType.Date (box value) |> ignore | None -> () use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(navRecordFromReader reader) records |> Seq.toList member _.CreateFund(idempotencyKey: string, command: FundCreateCommand) = if String.IsNullOrWhiteSpace idempotencyKey then FundWriteResult.Invalid "idempotency key cannot be empty" else match validateStorageRange command with | Error message -> FundWriteResult.Invalid message | Ok() -> let fingerprint = requestHash command use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try use lockCommand = commandWithTransaction connection (Some transaction) "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" addParameter lockCommand "lock_key" NpgsqlDbType.Text (box idempotencyKey) |> ignore lockCommand.ExecuteNonQuery() |> ignore match findIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, fundId) when existingHash = fingerprint -> match findFund connection (Some transaction) fundId with | Some fund -> transaction.Commit() FundWriteResult.Replayed fund | None -> transaction.Rollback() FundWriteResult.Invalid "idempotency record references a missing fund" | Some _ -> transaction.Rollback() FundWriteResult.IdempotencyConflict | None -> let fundId = Guid.NewGuid() match Ledger.initializeFund fundId command.Name command.InitialCash command.InitialUnitNav command.IsSynthetic Ledger.empty with | Error error -> transaction.Rollback() FundWriteResult.Invalid(ledgerErrorMessage error) | Ok state -> match Ledger.getFund fundId state with | Error error -> transaction.Rollback() FundWriteResult.Invalid(ledgerErrorMessage error) | Ok fund -> insertFund connection (Some transaction) fund insertIdempotency connection (Some transaction) idempotencyKey fingerprint fundId transaction.Commit() FundWriteResult.Created(recordFromLedger fund) with error -> try transaction.Rollback() with _ -> () raise error member _.CreateSubscriptionOrder(idempotencyKey: string, fundId: Guid, command: SubscriptionOrderCommand, ?tradeDateOverride: DateOnly) : SubscriptionOrderWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then OrderInvalid "idempotency key cannot be empty" else match validateOrderCommand command with | Error message -> OrderInvalid message | Ok() -> let fingerprint = orderRequestHash fundId command let reservedTotal = command.Amount + command.FeeAmount use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try use lockCommand = commandWithTransaction connection (Some transaction) "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" addParameter lockCommand "lock_key" NpgsqlDbType.Text (box idempotencyKey) |> ignore lockCommand.ExecuteNonQuery() |> ignore match findOrderIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, orderId) when existingHash = fingerprint && existingFundId = fundId -> match findOrder connection (Some transaction) orderId with | Some order -> transaction.Commit() OrderReplayed order | None -> transaction.Rollback() OrderInvalid "idempotency record references a missing order" | Some _ -> transaction.Rollback() OrderIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() OrderFundNotFound | Some isSynthetic -> if instrumentExists connection (Some transaction) command.FundCode then use cashCommand = commandWithTransaction connection (Some transaction) """ UPDATE funds SET available_cash = available_cash - @reserved_total, reserved_cash = reserved_cash + @reserved_total WHERE id = @fund_id AND available_cash >= @reserved_total """ addParameter cashCommand "reserved_total" NpgsqlDbType.Numeric (box reservedTotal) |> ignore addParameter cashCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore if cashCommand.ExecuteNonQuery() = 0 then transaction.Rollback() OrderInsufficientFunds else let submittedAt = DateTimeOffset.UtcNow let tradeDate = tradeDateOverride |> Option.defaultWith (fun () -> ConfirmationPolicy.tradeDateFor submittedAt) let order: SubscriptionOrderRecord = { Id = Guid.NewGuid() FundId = fundId FundCode = command.FundCode Amount = command.Amount FeeAmount = command.FeeAmount ReservedTotal = reservedTotal Status = orderStatusText SubmittedAt = submittedAt IsSynthetic = isSynthetic TradeDate = tradeDate ConfirmIdempotencyKey = None PendingReason = None ConfirmedAt = None ConfirmedQuote = None ConfirmedUnits = None ConfirmedInvestedCash = None ConfirmedResidualCash = None } let submittedAt = insertSubscriptionOrder connection (Some transaction) order insertOrderIdempotency connection (Some transaction) idempotencyKey fingerprint order.Id fundId transaction.Commit() OrderCreated { order with SubmittedAt = submittedAt } else transaction.Rollback() OrderInstrumentNotFound with error -> try transaction.Rollback() with _ -> () raise error member _.GetSubscriptionOrders(fundId: Guid) = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT id, fund_id, fund_code, amount, fee_amount, 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 """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(orderRecordFromReader reader) records |> Seq.toList member _.CreateSipPlan(idempotencyKey: string, fundId: Guid, command: SipPlanCommand, ?anchorOverride: DateOnly) : SipPlanWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then SipPlanWriteResult.SipPlanInvalid "idempotency key cannot be empty" else match SipPolicy.validateAmount command.Amount with | Error message -> SipPlanWriteResult.SipPlanInvalid message | Ok() -> let fingerprint = sipPlanRequestHash fundId command use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try use lockCommand = commandWithTransaction connection (Some transaction) "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" addParameter lockCommand "lock_key" NpgsqlDbType.Text (box idempotencyKey) |> ignore lockCommand.ExecuteNonQuery() |> ignore match findSipPlanIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, planId) when existingHash = fingerprint && existingFundId = fundId -> match findSipPlan connection (Some transaction) planId with | Some plan -> transaction.Commit() SipPlanWriteResult.SipPlanReplayed plan | None -> transaction.Rollback() SipPlanWriteResult.SipPlanInvalid "idempotency record references a missing plan" | Some _ -> transaction.Rollback() SipPlanWriteResult.SipPlanIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() SipPlanWriteResult.SipPlanFundNotFound | Some isSynthetic -> if instrumentExists connection (Some transaction) command.InstrumentCode then let anchorDate = anchorOverride |> Option.defaultWith (fun () -> ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow) let plan: SipPlanRecord = { Id = Guid.NewGuid() FundId = fundId InstrumentCode = command.InstrumentCode Amount = command.Amount Frequency = command.Frequency Status = "active" IsSynthetic = isSynthetic AnchorDate = anchorDate NextTradeDate = SipPolicy.nextTradeDate command.Frequency anchorDate (anchorDate.AddDays 1) CreatedAt = DateTimeOffset.UtcNow LastExecutionStatus = None LastExecutionDate = None } let createdAt = insertSipPlan connection (Some transaction) plan insertSipPlanIdempotency connection (Some transaction) idempotencyKey fingerprint plan.Id fundId transaction.Commit() SipPlanWriteResult.SipPlanCreated { plan with CreatedAt = createdAt } else transaction.Rollback() SipPlanWriteResult.SipPlanInstrumentNotFound with error -> try transaction.Rollback() with _ -> () raise error member this.AdvanceSipPlans(fundId: Guid, endDate: DateOnly, limit: int) : SipAdvanceResult = use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try use lockCommand = commandWithTransaction connection (Some transaction) "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" addParameter lockCommand "lock_key" NpgsqlDbType.Text (box (sprintf "sip-advance:%O" fundId)) |> ignore lockCommand.ExecuteNonQuery() |> ignore // 0. retry phase: confirmations that were pending_nav get re-run once their NAV landed use retryCommand = commandWithTransaction connection (Some transaction) """ SELECT e.plan_id, e.trade_date, e.order_id FROM sip_executions e JOIN sip_plans p ON p.id = e.plan_id WHERE p.fund_id = @fund_id AND e.status = 'pending_nav' ORDER BY e.trade_date """ addParameter retryCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use retryReader = retryCommand.ExecuteReader() let pendingRetries = ResizeArray() while retryReader.Read() do pendingRetries.Add( retryReader.GetGuid(0), retryReader.GetFieldValue(1), retryReader.GetGuid(2) ) retryReader.Close() for (retryPlanId, retryDate, retryOrderId) in pendingRetries do let confirmKey = sprintf "sip-confirm:%O:%s" retryPlanId (retryDate.ToString("yyyy-MM-dd")) match this.ConfirmSubscriptionOrder(confirmKey, fundId, retryOrderId) with | SubscriptionConfirmResult.OrderConfirmed _ | SubscriptionConfirmResult.ConfirmReplayed _ -> use doneCommand = commandWithTransaction connection (Some transaction) "UPDATE sip_executions SET status = 'succeeded', pending_reason = NULL, executed_at = now() WHERE plan_id = @plan_id AND trade_date = @trade_date" addParameter doneCommand "plan_id" NpgsqlDbType.Uuid (box retryPlanId) |> ignore addParameter doneCommand "trade_date" NpgsqlDbType.Date (box retryDate) |> ignore doneCommand.ExecuteNonQuery() |> ignore | _ -> // still pending: keep the execution as pending_nav for the next drive () use plansCommand = commandWithTransaction connection (Some transaction) "SELECT id, instrument_code, amount, frequency, anchor_date, next_trade_date FROM sip_plans WHERE fund_id = @fund_id AND status = 'active' ORDER BY created_at" addParameter plansCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use plansReader = plansCommand.ExecuteReader() let plans = ResizeArray() while plansReader.Read() do plans.Add( plansReader.GetGuid(0), plansReader.GetString(1), plansReader.GetDecimal(2), (match SipPolicy.parseFrequency (plansReader.GetString(3)) with | Some frequency -> frequency | None -> failwith "sip plan frequency is invalid"), plansReader.GetFieldValue(4), plansReader.GetFieldValue(5) ) plansReader.Close() let readExecutionsUpTo (planId: Guid) (upTo: DateOnly) = use replayCommand = commandWithTransaction connection (Some transaction) """ SELECT trade_date, status, order_id, pending_reason FROM sip_executions WHERE plan_id = @plan_id AND trade_date <= @up_to ORDER BY trade_date """ addParameter replayCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore addParameter replayCommand "up_to" NpgsqlDbType.Date (box upTo) |> ignore use reader = replayCommand.ExecuteReader() let rows = ResizeArray() while reader.Read() do rows.Add( { TradeDate = reader.GetFieldValue(0) Status = reader.GetString(1) OrderId = (if reader.IsDBNull(2) then None else Some(reader.GetGuid(2))) PendingReason = readStringOption reader 3 } ) rows |> Seq.toList let advancePlanRow (planId: Guid, code: string, amount: decimal, frequency: SipFrequency, anchor: DateOnly, nextDate: DateOnly) = let dueDates = SipPolicy.advancePlan frequency anchor nextDate endDate |> List.truncate (max 1 limit) let outcomes = ResizeArray() let upsertExecution (tradeDate: DateOnly) = use insertCommand = commandWithTransaction connection (Some transaction) """ INSERT INTO sip_executions (plan_id, trade_date, amount, fee_amount, status) VALUES (@plan_id, @trade_date, @amount, @fee_amount, 'processing') ON CONFLICT (plan_id, trade_date) DO NOTHING RETURNING trade_date """ addParameter insertCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore addParameter insertCommand "trade_date" NpgsqlDbType.Date (box tradeDate) |> ignore addParameter insertCommand "amount" NpgsqlDbType.Numeric (box amount) |> ignore addParameter insertCommand "fee_amount" NpgsqlDbType.Numeric (box 0m) |> ignore use reader = insertCommand.ExecuteReader() let inserted = reader.Read() reader.Close() inserted let setExecution (tradeDate: DateOnly) (status: string) (orderId: Guid option) (reason: string option) = use updateCommand = commandWithTransaction connection (Some transaction) """ UPDATE sip_executions SET status = @status, order_id = @order_id, pending_reason = @reason, executed_at = now() WHERE plan_id = @plan_id AND trade_date = @trade_date """ addParameter updateCommand "status" NpgsqlDbType.Text (box status) |> ignore let orderParameter = match orderId with | Some value -> box value | None -> box DBNull.Value addParameter updateCommand "order_id" NpgsqlDbType.Uuid orderParameter |> ignore let reasonParameter = match reason with | Some value -> box value | None -> box DBNull.Value addParameter updateCommand "reason" NpgsqlDbType.Text reasonParameter |> ignore addParameter updateCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore addParameter updateCommand "trade_date" NpgsqlDbType.Date (box tradeDate) |> ignore updateCommand.ExecuteNonQuery() |> ignore let findExistingExecution (tradeDate: DateOnly) = use selectCommand = commandWithTransaction connection (Some transaction) "SELECT status, order_id, pending_reason FROM sip_executions WHERE plan_id = @plan_id AND trade_date = @trade_date" addParameter selectCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore addParameter selectCommand "trade_date" NpgsqlDbType.Date (box tradeDate) |> ignore use reader = selectCommand.ExecuteReader() if reader.Read() then Some (reader.GetString(0), (if reader.IsDBNull(1) then None else Some(reader.GetGuid(1))), readStringOption reader 2) else None let orderIdFor (tradeDate: DateOnly) = let orderKey = sprintf "sip:%O:%s" planId (tradeDate.ToString("yyyy-MM-dd")) match this.CreateSubscriptionOrder(orderKey, fundId, { FundCode = code; Amount = amount; FeeAmount = 0m }, tradeDate) with | SubscriptionOrderWriteResult.OrderCreated order -> Some order.Id | SubscriptionOrderWriteResult.OrderReplayed order -> Some order.Id | SubscriptionOrderWriteResult.OrderInsufficientFunds -> None | other -> failwithf "unexpected sip order result: %A" other for tradeDate in dueDates do // 1. claim the slot atomically: same plan + same trade date runs once if upsertExecution tradeDate then // 2. place the order through the shared pipeline with a deterministic key match orderIdFor tradeDate with | None -> setExecution tradeDate "insufficient_cash" None (Some "available cash is not enough for the scheduled amount") outcomes.Add({ TradeDate = tradeDate; Status = "insufficient_cash"; OrderId = None; PendingReason = Some "available cash is not enough for the scheduled amount" }) | Some orderId -> // 3. confirm through the shared confirmation pipeline let confirmKey = sprintf "sip-confirm:%O:%s" planId (tradeDate.ToString("yyyy-MM-dd")) match this.ConfirmSubscriptionOrder(confirmKey, fundId, orderId) with | SubscriptionConfirmResult.OrderConfirmed _ | SubscriptionConfirmResult.ConfirmReplayed _ -> setExecution tradeDate "succeeded" (Some orderId) None outcomes.Add({ TradeDate = tradeDate; Status = "succeeded"; OrderId = Some orderId; PendingReason = None }) | SubscriptionConfirmResult.ConfirmPendingNav record -> setExecution tradeDate "pending_nav" (Some orderId) record.PendingReason outcomes.Add({ TradeDate = tradeDate; Status = "pending_nav"; OrderId = Some orderId; PendingReason = record.PendingReason }) | other -> setExecution tradeDate "failed" (Some orderId) (Some (sprintf "%A" other)) outcomes.Add({ TradeDate = tradeDate; Status = "failed"; OrderId = Some orderId; PendingReason = Some (sprintf "%A" other) }) else match findExistingExecution tradeDate with | Some("succeeded", orderId, reason) -> outcomes.Add({ TradeDate = tradeDate; Status = "succeeded"; OrderId = orderId; PendingReason = reason }) | Some(status, orderId, reason) -> outcomes.Add({ TradeDate = tradeDate; Status = status; OrderId = orderId; PendingReason = reason }) | None -> outcomes.Add({ TradeDate = tradeDate; Status = "unknown"; OrderId = None; PendingReason = None }) // 4. roll the plan pointer forward past the processed window let rolled = match dueDates with | [] -> nextDate | lastDueDates -> SipPolicy.nextTradeDate frequency anchor ((List.last lastDueDates).AddDays 1) if not (List.isEmpty dueDates) then use rollCommand = commandWithTransaction connection (Some transaction) "UPDATE sip_plans SET next_trade_date = @next_date WHERE id = @plan_id" addParameter rollCommand "next_date" NpgsqlDbType.Date (box rolled) |> ignore addParameter rollCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore rollCommand.ExecuteNonQuery() |> ignore let replayedOutcomes = if List.isEmpty dueDates then // repeat drive with the same endDate: replay what was already processed readExecutionsUpTo planId endDate else outcomes |> Seq.toList { PlanId = planId InstrumentCode = code Amount = amount Frequency = frequency Executions = replayedOutcomes NextTradeDate = rolled } let planResults = plans |> Seq.map advancePlanRow |> Seq.toList transaction.Commit() { FundId = fundId; Plans = planResults } with error -> try transaction.Rollback() with _ -> () raise error member _.GetSipPlans(fundId: Guid) = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT p.id, p.fund_id, p.instrument_code, p.amount, p.frequency, p.status, p.is_synthetic, p.anchor_date, p.next_trade_date, p.created_at, e.status, e.trade_date FROM sip_plans p LEFT JOIN LATERAL ( SELECT status, trade_date FROM sip_executions WHERE plan_id = p.id ORDER BY executed_at DESC, trade_date DESC LIMIT 1 ) e ON true WHERE p.fund_id = @fund_id ORDER BY p.created_at DESC, p.id """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(sipPlanRecordFromReader reader) records |> Seq.toList member _.CreateRebalancePlan(idempotencyKey: string, fundId: Guid, command: RebalancePlanCommand) : RebalanceWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then RebalanceWriteResult.RebalanceInvalid "idempotency key cannot be empty" else match validateRebalanceCommand command with | Error message -> RebalanceWriteResult.RebalanceInvalid message | Ok() -> let fingerprint = rebalanceRequestHash fundId command use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try use lockCommand = commandWithTransaction connection (Some transaction) "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" addParameter lockCommand "lock_key" NpgsqlDbType.Text (box idempotencyKey) |> ignore lockCommand.ExecuteNonQuery() |> ignore match findRebalanceIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, planId) when existingHash = fingerprint && existingFundId = fundId -> match findRebalancePlan connection (Some transaction) planId with | Some plan -> transaction.Commit() RebalanceWriteResult.RebalancePlanReplayed plan | None -> transaction.Rollback() RebalanceWriteResult.RebalanceInvalid "idempotency record references a missing plan" | Some _ -> transaction.Rollback() RebalanceWriteResult.RebalanceIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() RebalanceWriteResult.RebalanceFundNotFound | Some isSynthetic -> let missingTarget = command.Targets |> List.tryFind (fun target -> not (instrumentExists connection (Some transaction) target.InstrumentCode)) match missingTarget with | Some target -> transaction.Rollback() RebalanceWriteResult.RebalanceInstrumentNotFound | None -> let plan: RebalancePlanRecord = { Id = Guid.NewGuid() FundId = fundId Targets = command.Targets Status = "active" IsSynthetic = isSynthetic CreatedAt = DateTimeOffset.UtcNow } let createdAt = insertRebalancePlan connection (Some transaction) plan insertRebalanceIdempotency connection (Some transaction) idempotencyKey fingerprint plan.Id fundId transaction.Commit() RebalanceWriteResult.RebalancePlanCreated { plan with CreatedAt = createdAt } with error -> try transaction.Rollback() with _ -> () raise error member _.GetRebalancePlans(fundId: Guid) = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None "SELECT id, fund_id, status, is_synthetic, created_at FROM rebalance_plans WHERE fund_id = @fund_id ORDER BY created_at DESC, id" addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let plans = ResizeArray() while reader.Read() do plans.Add(rebalancePlanRecordFromReader reader) reader.Close() plans |> Seq.toList |> List.map (rebalancePlanWithTargets connection None) member this.ExecuteRebalancePlan(planId: Guid) : Result = use connection = new NpgsqlConnection(connectionString) connection.Open() let plan = use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try use lockCommand = commandWithTransaction connection (Some transaction) "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" addParameter lockCommand "lock_key" NpgsqlDbType.Text (box (sprintf "rebalance-execute:%O" planId)) |> ignore lockCommand.ExecuteNonQuery() |> ignore let plan = findRebalancePlan connection (Some transaction) planId transaction.Commit() plan with error -> try transaction.Rollback() with _ -> () raise error match plan with | None -> Error "rebalance plan was not found" | Some plan -> let fundId = plan.FundId let runDate = ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow // snapshot: fund cash plus each holding valued at its latest valuation NAV let fund = this.GetFund fundId match fund with | None -> Error "fund was not found" | Some fund -> let positions = this.GetFundPositions fundId let buildSnapshot (position: FundPositionRecord) : RebalancePolicy.RebalancePositionSnapshot = let marketValue = match position.ValuationNav with | Some nav -> Decimal.Round(position.Units * nav, 2) | None -> 0m { RebalancePolicy.RebalancePositionSnapshot.InstrumentCode = position.InstrumentCode RebalancePolicy.RebalancePositionSnapshot.MarketValue = marketValue RebalancePolicy.RebalancePositionSnapshot.Units = position.Units RebalancePolicy.RebalancePositionSnapshot.AvailableUnits = position.Units - position.ReservedUnits RebalancePolicy.RebalancePositionSnapshot.ValuationNav = position.ValuationNav } let snapshots = positions |> List.map buildSnapshot let diffs = RebalancePolicy.computeOrders plan.Targets snapshots fund.AvailableCash match diffs with | Error message -> Error message | Ok diffs -> let outcomes = ResizeArray() // sells first so freed cash can fund the buys let ordered = diffs |> List.sortBy (fun diff -> match diff.Action with | RebalancePolicy.Sell -> 0 | RebalancePolicy.Buy -> 1 | RebalancePolicy.Hold -> 2) for diff: RebalancePolicy.RebalanceDiff in ordered do match diff.Action with | RebalancePolicy.Hold -> () | RebalancePolicy.Buy -> let orderKey = RebalancePolicy.orderKey plan.Id runDate diff.InstrumentCode match this.CreateSubscriptionOrder( orderKey, fundId, { FundCode = diff.InstrumentCode; Amount = diff.Amount; FeeAmount = 0m } ) with | SubscriptionOrderWriteResult.OrderCreated order | SubscriptionOrderWriteResult.OrderReplayed order -> let confirmKey = RebalancePolicy.confirmKey plan.Id runDate diff.InstrumentCode match this.ConfirmSubscriptionOrder(confirmKey, fundId, order.Id) with | SubscriptionConfirmResult.OrderConfirmed _ | SubscriptionConfirmResult.ConfirmReplayed _ -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "succeeded"; OrderId = Some order.Id; PendingReason = None }) | SubscriptionConfirmResult.ConfirmPendingNav record -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "pending_nav"; OrderId = Some order.Id; PendingReason = record.PendingReason }) | SubscriptionConfirmResult.ConfirmIdempotencyConflict -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "idempotency_conflict"; OrderId = Some order.Id; PendingReason = None }) | other -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "failed"; OrderId = Some order.Id; PendingReason = Some (sprintf "%A" other) }) | SubscriptionOrderWriteResult.OrderInsufficientFunds -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "insufficient_cash"; OrderId = None; PendingReason = Some "available cash is not enough for the rebalance buy" }) | other -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "failed"; OrderId = None; PendingReason = Some (sprintf "%A" other) }) | RebalancePolicy.Sell -> match diff.Units with | None -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "skipped_no_valuation"; OrderId = None; PendingReason = Some "holding has no valuation NAV to price the sell" }) | Some units when units <= 0m -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "skipped_no_available_units"; OrderId = None; PendingReason = Some "no available units to redeem" }) | Some units -> let redeemKey = RebalancePolicy.redemptionKey plan.Id runDate diff.InstrumentCode match this.CreateRedemptionOrder( redeemKey, fundId, { InstrumentCode = diff.InstrumentCode; Units = units; FeeAmount = 0m } ) with | RedemptionWriteResult.RedemptionCreated order | RedemptionWriteResult.RedemptionReplayed order -> let confirmKey = RebalancePolicy.redemptionConfirmKey plan.Id runDate diff.InstrumentCode match this.ConfirmRedemptionOrder(confirmKey, fundId, order.Id) with | RedemptionConfirmResult.RedemptionConfirmed _ | RedemptionConfirmResult.RedemptionConfirmReplayed _ -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "succeeded"; OrderId = Some order.Id; PendingReason = None }) | RedemptionConfirmResult.RedemptionPendingNav record -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "pending_nav"; OrderId = Some order.Id; PendingReason = record.PendingReason }) | other -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "failed"; OrderId = Some order.Id; PendingReason = Some (sprintf "%A" other) }) | RedemptionWriteResult.RedemptionInsufficientUnits -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "insufficient_units"; OrderId = None; PendingReason = Some "available units are not enough for the rebalance sell" }) | other -> outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "failed"; OrderId = None; PendingReason = Some (sprintf "%A" other) }) Ok { PlanId = plan.Id RunDate = runDate Outcomes = outcomes |> Seq.toList } member _.CreateCapitalDeposit(idempotencyKey: string, fundId: Guid, command: CapitalDepositCommand) : CapitalDepositWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then CapitalDepositWriteResult.CapitalDepositInvalid "idempotency key cannot be empty" else match CapitalPolicy.validateDeposit command.Amount with | Error message -> CapitalDepositWriteResult.CapitalDepositInvalid message | Ok() -> let fingerprint = capitalDepositRequestHash fundId command use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try use lockCommand = commandWithTransaction connection (Some transaction) "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" addParameter lockCommand "lock_key" NpgsqlDbType.Text (box idempotencyKey) |> ignore lockCommand.ExecuteNonQuery() |> ignore match findCapitalDepositIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, depositId) when existingHash = fingerprint && existingFundId = fundId -> match findCapitalDeposit connection (Some transaction) depositId with | Some deposit -> transaction.Commit() CapitalDepositWriteResult.CapitalDepositReplayed deposit | None -> transaction.Rollback() CapitalDepositWriteResult.CapitalDepositInvalid "idempotency record references a missing deposit" | Some _ -> transaction.Rollback() CapitalDepositWriteResult.CapitalDepositIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() CapitalDepositWriteResult.CapitalDepositFundNotFound | Some isSynthetic -> use cashCommand = commandWithTransaction connection (Some transaction) "UPDATE funds SET available_cash = available_cash + @amount WHERE id = @fund_id" addParameter cashCommand "amount" NpgsqlDbType.Numeric (box command.Amount) |> ignore addParameter cashCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore if cashCommand.ExecuteNonQuery() = 0 then transaction.Rollback() CapitalDepositWriteResult.CapitalDepositFundNotFound else let deposit: CapitalDepositRecord = { Id = Guid.NewGuid() FundId = fundId Amount = command.Amount Note = command.Note IsSynthetic = isSynthetic CreatedAt = DateTimeOffset.UtcNow } let createdAt = insertCapitalDeposit connection (Some transaction) deposit insertCapitalDepositIdempotency connection (Some transaction) idempotencyKey fingerprint deposit.Id fundId transaction.Commit() CapitalDepositWriteResult.CapitalDepositCreated { deposit with CreatedAt = createdAt } with error -> try transaction.Rollback() with _ -> () raise error member _.GetCapitalDeposits(fundId: Guid) = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None "SELECT id, fund_id, amount, note, is_synthetic, created_at FROM fund_capital_deposits WHERE fund_id = @fund_id ORDER BY created_at DESC, id" addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(capitalDepositRecordFromReader reader) records |> Seq.toList member _.CreateRedemptionOrder(idempotencyKey: string, fundId: Guid, command: RedemptionCommand) : RedemptionWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then RedemptionWriteResult.RedemptionInvalid "idempotency key cannot be empty" else match validateRedemptionCommand command with | Error message -> RedemptionWriteResult.RedemptionInvalid message | Ok() -> let fingerprint = redemptionRequestHash fundId command use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try use lockCommand = commandWithTransaction connection (Some transaction) "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" addParameter lockCommand "lock_key" NpgsqlDbType.Text (box idempotencyKey) |> ignore lockCommand.ExecuteNonQuery() |> ignore match findRedemptionOrderIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, orderId) when existingHash = fingerprint && existingFundId = fundId -> match findRedemptionOrder connection (Some transaction) orderId with | Some order -> transaction.Commit() RedemptionReplayed order | None -> transaction.Rollback() RedemptionWriteResult.RedemptionInvalid "idempotency record references a missing order" | Some _ -> transaction.Rollback() RedemptionIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() RedemptionFundNotFound | Some isSynthetic -> if instrumentExists connection (Some transaction) command.InstrumentCode then use freezeCommand = commandWithTransaction connection (Some transaction) """ UPDATE fund_positions SET reserved_units = reserved_units + @units WHERE fund_id = @fund_id AND instrument_code = @code AND units - reserved_units >= @units """ addParameter freezeCommand "units" NpgsqlDbType.Numeric (box command.Units) |> ignore addParameter freezeCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter freezeCommand "code" NpgsqlDbType.Text (box command.InstrumentCode) |> ignore if freezeCommand.ExecuteNonQuery() = 0 then transaction.Rollback() RedemptionInsufficientUnits else let submittedAt = DateTimeOffset.UtcNow let order: RedemptionOrderRecord = { Id = Guid.NewGuid() FundId = fundId InstrumentCode = command.InstrumentCode Units = command.Units FeeAmount = command.FeeAmount Status = "submitted" IsSynthetic = isSynthetic SubmittedAt = submittedAt TradeDate = ConfirmationPolicy.tradeDateFor submittedAt ConfirmIdempotencyKey = None PendingReason = None ConfirmedAt = None ConfirmedNav = None ConfirmedNavDate = None ConfirmedProceeds = None ConfirmedCostReleased = None } let submittedAt = insertRedemptionOrder connection (Some transaction) order insertRedemptionOrderIdempotency connection (Some transaction) idempotencyKey fingerprint order.Id fundId transaction.Commit() RedemptionCreated { order with SubmittedAt = submittedAt } else transaction.Rollback() RedemptionInstrumentNotFound with error -> try transaction.Rollback() with _ -> () raise error member _.GetRedemptionOrders(fundId: Guid) = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None (redemptionOrderColumns + " WHERE fund_id = @fund_id ORDER BY submitted_at DESC, id") addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(redemptionRecordFromReader reader) records |> Seq.toList member _.ConfirmRedemptionOrder(idempotencyKey: string, fundId: Guid, orderId: Guid) : RedemptionConfirmResult = if String.IsNullOrWhiteSpace idempotencyKey then RedemptionConfirmResult.RedemptionInvalid "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 use lockCommand = commandWithTransaction connection (Some transaction) "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" addParameter lockCommand "lock_key" NpgsqlDbType.Text (box (sprintf "confirm-redemption:%O" orderId)) |> ignore lockCommand.ExecuteNonQuery() |> ignore match findRedemptionOrder connection (Some transaction) orderId with | None -> transaction.Rollback() RedemptionOrderNotFound | Some order when order.FundId <> fundId -> transaction.Rollback() RedemptionOrderNotFound | Some order -> if order.Status = "confirmed" then match order.ConfirmIdempotencyKey with | Some storedKey when storedKey = idempotencyKey -> transaction.Commit() RedemptionConfirmReplayed order | _ -> match findRedemptionConfirmIdempotency connection (Some transaction) idempotencyKey with | Some(_, _, storedOrderId) when storedOrderId = orderId -> transaction.Commit() RedemptionConfirmReplayed order | Some _ -> transaction.Rollback() RedemptionConfirmIdempotencyConflict | None -> transaction.Rollback() RedemptionAlreadyConfirmed elif order.Status <> "submitted" && order.Status <> "pending_nav" then transaction.Rollback() RedemptionInvalidStatus else match findRedemptionConfirmIdempotency connection (Some transaction) idempotencyKey with | Some _ -> transaction.Rollback() RedemptionConfirmIdempotencyConflict | None -> let tradeDate = order.TradeDate let tradeDateText = tradeDate.ToString("yyyy-MM-dd") use quoteCommand = commandWithTransaction connection (Some transaction) """ 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 """ addParameter quoteCommand "code" NpgsqlDbType.Text (box order.InstrumentCode) |> ignore addParameter quoteCommand "trade_date" NpgsqlDbType.Date (box tradeDate) |> ignore use quoteReader = quoteCommand.ExecuteReader() let quoteFound = quoteReader.Read() let selectedQuote = if quoteFound then let evidenceFirstSeen = if quoteReader.IsDBNull(8) then None else Some(quoteReader.GetFieldValue(8)) 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) }, evidenceFirstSeen) else None quoteReader.Close() let boundQuote = match selectedQuote 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(_, evidenceFirstSeen) when evidenceFirstSeen.IsNone -> Some(sprintf "nav revision for trade date %s has no observation evidence recorded" 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 match deferralReason with | Some reason -> use pendingCommand = commandWithTransaction connection (Some transaction) """ UPDATE redemption_orders SET status = 'pending_nav', pending_reason = @reason, confirm_idempotency_key = NULL, confirmed_at = NULL, confirmed_nav = NULL, confirmed_nav_date = NULL, confirmed_proceeds = NULL, confirmed_cost_released = 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 transaction.Commit() match findRedemptionOrder connection None orderId with | Some pendingOrder -> RedemptionPendingNav pendingOrder | None -> failwith "pending redemption disappeared after confirmation deferral" | None -> let nav = boundQuote |> Option.get |> fun quote -> quote.Nav use positionCommand = commandWithTransaction connection (Some transaction) """ SELECT units, reserved_units, cost_cash FROM fund_positions WHERE fund_id = @fund_id AND instrument_code = @code FOR UPDATE """ addParameter positionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter positionCommand "code" NpgsqlDbType.Text (box order.InstrumentCode) |> ignore use positionReader = positionCommand.ExecuteReader() let positionFound = positionReader.Read() let positionUnits = if positionFound then positionReader.GetDecimal(0) else 0m let positionReserved = if positionFound then positionReader.GetDecimal(1) else 0m let positionCost = if positionFound then positionReader.GetDecimal(2) else 0m positionReader.Close() if not positionFound || positionReserved < order.Units then failwith "frozen units are missing for redemption confirmation" match RedemptionPolicy.compute order.Units nav order.FeeAmount positionCost positionUnits with | Error message -> use pendingCommand = commandWithTransaction connection (Some transaction) """ UPDATE redemption_orders SET status = 'pending_nav', pending_reason = @reason WHERE id = @order_id """ addParameter pendingCommand "reason" NpgsqlDbType.Text (box message) |> ignore addParameter pendingCommand "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore pendingCommand.ExecuteNonQuery() |> ignore transaction.Commit() match findRedemptionOrder connection None orderId with | Some pendingOrder -> RedemptionPendingNav pendingOrder | None -> failwith "pending redemption disappeared after computation failure" | Ok computation -> use settlePositionCommand = commandWithTransaction connection (Some transaction) """ UPDATE fund_positions SET units = units - @units, reserved_units = reserved_units - @units, cost_cash = cost_cash - @cost_released WHERE fund_id = @fund_id AND instrument_code = @code AND units > @units AND reserved_units >= @units AND cost_cash >= @cost_released """ addParameter settlePositionCommand "units" NpgsqlDbType.Numeric (box order.Units) |> ignore addParameter settlePositionCommand "cost_released" NpgsqlDbType.Numeric (box computation.RedeemedCost) |> ignore addParameter settlePositionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter settlePositionCommand "code" NpgsqlDbType.Text (box order.InstrumentCode) |> ignore if settlePositionCommand.ExecuteNonQuery() = 0 then use deletePositionCommand = commandWithTransaction connection (Some transaction) """ DELETE FROM fund_positions WHERE fund_id = @fund_id AND instrument_code = @code AND units = @units AND reserved_units >= @units AND cost_cash >= @cost_released """ addParameter deletePositionCommand "units" NpgsqlDbType.Numeric (box order.Units) |> ignore addParameter deletePositionCommand "cost_released" NpgsqlDbType.Numeric (box computation.RedeemedCost) |> ignore addParameter deletePositionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter deletePositionCommand "code" NpgsqlDbType.Text (box order.InstrumentCode) |> ignore if deletePositionCommand.ExecuteNonQuery() = 0 then failwith "position units changed during redemption confirmation" use cashCommand = commandWithTransaction connection (Some transaction) """ UPDATE funds SET available_cash = available_cash + @proceeds WHERE id = @fund_id """ addParameter cashCommand "proceeds" NpgsqlDbType.Numeric (box computation.Proceeds) |> ignore addParameter cashCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore if cashCommand.ExecuteNonQuery() = 0 then failwith "fund disappeared during redemption confirmation" use confirmCommand = commandWithTransaction connection (Some transaction) """ UPDATE redemption_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_proceeds = @proceeds, confirmed_cost_released = @cost_released 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 nav) |> ignore addParameter confirmCommand "nav_date" NpgsqlDbType.Date (box tradeDate) |> ignore addParameter confirmCommand "proceeds" NpgsqlDbType.Numeric (box computation.Proceeds) |> ignore addParameter confirmCommand "cost_released" NpgsqlDbType.Numeric (box computation.RedeemedCost) |> ignore addParameter confirmCommand "order_id" NpgsqlDbType.Uuid (box orderId) |> ignore confirmCommand.ExecuteNonQuery() |> ignore insertRedemptionConfirmIdempotency connection (Some transaction) idempotencyKey "" orderId fundId transaction.Commit() match findRedemptionOrder connection None orderId with | Some confirmedOrder -> RedemptionConfirmed confirmedOrder | None -> failwith "confirmed redemption disappeared after commit" with error -> try transaction.Rollback() with _ -> () raise error 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 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 """ 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, evidenceFirstSeen = if quoteFound then let evidenceFirstSeen = if quoteReader.IsDBNull(8) then None else Some(quoteReader.GetFieldValue(8)) 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) }, evidenceFirstSeen else None, None quoteReader.Close() let boundQuote = match selectedQuote, evidenceFirstSeen with | Some quote, Some firstSeen -> Some { quote with FirstSeenAt = firstSeen } | _ -> None let deferralReason = 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 = 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 = boundQuote |> 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) : FundPositionRecord list = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT p.instrument_code, p.units, p.reserved_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(5) then None else Some(reader.GetDecimal(5)) let valuationNavDate = if reader.IsDBNull(6) then None else Some(reader.GetFieldValue(6)) let valuationCollectedAt = if reader.IsDBNull(7) then None else Some(reader.GetFieldValue(7)) records.Add( { FundId = fundId InstrumentCode = reader.GetString(0) Units = reader.GetDecimal(1) ReservedUnits = reader.GetDecimal(2) CostCash = reader.GetDecimal(3) LastConfirmedAt = reader.GetFieldValue(4) ValuationNav = valuationNav ValuationNavDate = valuationNavDate ValuationCollectedAt = valuationCollectedAt } ) records |> Seq.toList