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) /// Test-only clock anchor: modules may pin the trading-date clock by exporting /// FUND_LAB_TEST_TRADE_DATE so suites stay independent of the wall clock. Unset in /// production, the real Shanghai rule applies unchanged. let tradeDateFor (submittedAt: DateTimeOffset) : DateOnly = match Environment.GetEnvironmentVariable("FUND_LAB_TEST_TRADE_DATE") with | value when not (String.IsNullOrWhiteSpace value) -> match DateOnly.TryParseExact(value, "yyyy-MM-dd", CultureInfo.InvariantCulture, DateTimeStyles.None) with | true, anchored -> anchored | _ -> 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 | _ -> 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 /// Calendar date used to bucket events (cash flows, holdings) on the returns /// timeline. Shares the FUND_LAB_TEST_TRADE_DATE anchor with tradeDateFor so a /// pinned suite does not drift when the host crosses midnight; production (no /// pin) is exactly the Shanghai calendar date of the moment. let eventDateFor (moment: DateTimeOffset) : DateOnly = match Environment.GetEnvironmentVariable("FUND_LAB_TEST_TRADE_DATE") with | value when not (String.IsNullOrWhiteSpace value) -> match DateOnly.TryParseExact(value, "yyyy-MM-dd", CultureInfo.InvariantCulture, DateTimeStyles.None) with | true, anchored -> anchored | _ -> shanghaiDate moment | _ -> shanghaiDate moment 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 CreatedAt: DateTimeOffset } 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 StockTradeCommand = { InstrumentCode: string StockName: string option Quantity: decimal Price: decimal } type StockTradeRecord = { Id: Guid FundId: Guid InstrumentCode: string StockName: string option Quantity: decimal Price: decimal CostCash: decimal IsSynthetic: bool ExecutedAt: DateTimeOffset } type StockPositionRecord = { FundId: Guid InstrumentCode: string StockName: string option Quantity: decimal CostCash: decimal LastTradedAt: DateTimeOffset } type StockTradeWriteResult = | StockTradeCreated of StockTradeRecord | StockTradeReplayed of StockTradeRecord | StockTradeIdempotencyConflict | StockTradeInvalid of string | StockTradeFundNotFound type StockSellCommand = { InstrumentCode: string StockName: string option Quantity: decimal Price: decimal FeeAmount: decimal } type StockSellRecord = { Id: Guid FundId: Guid InstrumentCode: string StockName: string option Quantity: decimal Price: decimal FeeAmount: decimal Proceeds: decimal IsSynthetic: bool ExecutedAt: DateTimeOffset } type StockSellWriteResult = | StockSellCreated of StockSellRecord | StockSellReplayed of StockSellRecord | StockSellIdempotencyConflict | StockSellInvalid of string | StockSellInsufficientHoldings of string | StockSellFundNotFound type BondTradeCommand = { InstrumentCode: string BondName: string option Quantity: decimal /// All-in (dirty/全价) execution price per 100 of face value. Price: decimal CleanPrice: decimal AccruedInterest: decimal ParValue: decimal SettlementDate: DateOnly CouponRate: decimal option ValueDate: DateOnly option MaturityDate: DateOnly option /// Explicit trade/valuation date supplied by the caller (wins over the /// quote's own date so tests and backfills stay deterministic). TradeDate: DateOnly option } type BondTradeRecord = { Id: Guid FundId: Guid InstrumentCode: string BondName: string option Quantity: decimal Price: decimal CleanPrice: decimal AccruedInterest: decimal ParValue: decimal SettlementDate: DateOnly CouponRate: decimal option ValueDate: DateOnly option MaturityDate: DateOnly option TradeDate: DateOnly CostCash: decimal IsSynthetic: bool ExecutedAt: DateTimeOffset } type BondCashflowCommand = { InstrumentCode: string BondName: string option /// "coupon" (付息) or "maturity" (到期). EventType: string EventDate: DateOnly Quantity: decimal /// Cash credited to the fund's available cash. Amount: decimal Note: string option } type BondCashflowRecord = { Id: Guid FundId: Guid InstrumentCode: string BondName: string option EventType: string EventDate: DateOnly Quantity: decimal Amount: decimal Note: string option IsSynthetic: bool CreatedAt: DateTimeOffset } type BondCashflowWriteResult = | BondCashflowCreated of BondCashflowRecord | BondCashflowReplayed of BondCashflowRecord | BondCashflowIdempotencyConflict | BondCashflowInvalid of string | BondCashflowFundNotFound | BondCashflowPositionNotFound type BondPositionRecord = { FundId: Guid InstrumentCode: string BondName: string option Quantity: decimal CostCash: decimal LastTradedAt: DateTimeOffset } type BondSellCommand = { InstrumentCode: string BondName: string option Quantity: decimal /// All-in (dirty/全价) sell price per 100 of face value. Price: decimal CleanPrice: decimal AccruedInterest: decimal ParValue: decimal SettlementDate: DateOnly TradeDate: DateOnly option FeeAmount: decimal } type BondSellRecord = { Id: Guid FundId: Guid InstrumentCode: string BondName: string option Quantity: decimal Price: decimal CleanPrice: decimal AccruedInterest: decimal ParValue: decimal SettlementDate: DateOnly TradeDate: DateOnly FeeAmount: decimal /// Cash credited to the fund after fees (dirty amount minus fee). Proceeds: decimal /// Average-cost basis removed from the position. CostReleased: decimal RealizedPnl: decimal IsSynthetic: bool ExecutedAt: DateTimeOffset } type BondSellWriteResult = | BondSellCreated of BondSellRecord | BondSellReplayed of BondSellRecord | BondSellIdempotencyConflict | BondSellInvalid of string | BondSellInsufficientHoldings of string | BondSellFundNotFound type InstrumentSnapshotRecord = { InstrumentCode: string AssetClass: string SnapshotDate: DateOnly Price: decimal Source: string SourceRevision: string SourceCollectedAt: DateTimeOffset SourcePayloadHash: string } type BondTradeWriteResult = | BondTradeCreated of BondTradeRecord | BondTradeReplayed of BondTradeRecord | BondTradeIdempotencyConflict | BondTradeInvalid of string | BondTradeFundNotFound 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 SipPlanStatusResult = | SipPlanStatusChanged of SipPlanRecord | SipPlanStatusNotFound 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 InvestmentPlanCommand = { InstrumentCode: string Amount: decimal Frequency: InvestmentFrequency } type InvestmentPlanRecord = { Id: Guid FundId: Guid InstrumentCode: string Amount: decimal Frequency: InvestmentFrequency Status: string IsSynthetic: bool AnchorDate: DateOnly NextRunDate: DateOnly CreatedAt: DateTimeOffset LastRunStatus: string option LastRunDate: DateOnly option } type InvestmentPlanWriteResult = | InvestmentPlanCreated of InvestmentPlanRecord | InvestmentPlanReplayed of InvestmentPlanRecord | InvestmentPlanIdempotencyConflict | InvestmentPlanInvalid of string | InvestmentPlanFundNotFound | InvestmentPlanInstrumentNotFound type InvestmentPlanRunOutcome = { RunDate: DateOnly Status: string OrderId: Guid option PendingReason: string option } type InvestmentPlanRunPlanResult = { PlanId: Guid InstrumentCode: string Amount: decimal Frequency: InvestmentFrequency Runs: InvestmentPlanRunOutcome list NextRunDate: DateOnly } type InvestmentPlanRunResult = { FundId: Guid ProcessingDate: DateOnly Plans: InvestmentPlanRunPlanResult 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 RebalanceArchiveResult = | RebalancePlanArchived of RebalancePlanRecord | RebalanceArchiveNotFound 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 RebalanceExecutionRecord = { PlanId: Guid RunDate: DateOnly InstrumentCode: string Action: string Amount: decimal Status: string OrderId: Guid option PendingReason: string option ExecutedAt: DateTimeOffset } type RebalancePreview = { PlanId: Guid RunDate: DateOnly AvailableCash: decimal Equity: decimal Rows: RebalancePolicy.RebalanceWeightRow list } type DividendMode = DividendPolicy.DividendMode type DividendCommand = { InstrumentCode: string NavDate: DateOnly Dps: decimal Mode: DividendMode } type DividendRecord = { Id: Guid FundId: Guid InstrumentCode: string NavDate: DateOnly Dps: decimal Mode: DividendMode Status: string IsSynthetic: bool GrossCash: decimal option CreditedUnits: decimal option CreditedInvested: decimal option OrderId: Guid option PendingReason: string option CreatedAt: DateTimeOffset } type DividendWriteResult = | DividendCredited of DividendRecord | DividendReplayed of DividendRecord | DividendPendingReinvest of DividendRecord | DividendIdempotencyConflict | DividendInvalid of string | DividendFundNotFound | DividendInstrumentNotFound | DividendNoHoldings type DividendStageDecision = | StageReplay of Guid | StageFail of DividendWriteResult | StageBooked of DividendRecord 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 /// One reconstructed end-of-day fund valuation. `TotalAssets`/`UnitNav`/ /// `HoldingsValue`/`CumulativeReturn` are `None` when a held instrument has no /// NAV observation dated on or before `Date`: the value is unknown, never zero. type FundReturnsPoint = { Date: DateOnly Pending: bool TotalAssets: decimal option UnitNav: decimal option Cash: decimal ReservedCash: decimal HoldingsValue: decimal option CumulativeReturn: decimal option NetExternalFlow: decimal } type FundReturns = { FundId: Guid Pending: bool DataUpdatedAt: DateTimeOffset option Points: FundReturnsPoint list } 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_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, plan_id uuid NOT NULL REFERENCES rebalance_plans(id), fund_id uuid NOT NULL REFERENCES funds(id), 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) ); CREATE TABLE IF NOT EXISTS rebalance_executions ( plan_id uuid NOT NULL REFERENCES rebalance_plans(id), run_date date NOT NULL, instrument_code text NOT NULL, action text NOT NULL, 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, run_date, instrument_code) ); CREATE TABLE IF NOT EXISTS dividend_records ( id uuid PRIMARY KEY, fund_id uuid NOT NULL REFERENCES funds(id), instrument_code text NOT NULL REFERENCES instruments(code), nav_date date NOT NULL, dps numeric(28, 8) NOT NULL CHECK (dps > 0), mode text NOT NULL, status text NOT NULL, is_synthetic boolean NOT NULL, gross_cash numeric(20, 2) NULL, credited_units numeric(28, 8) NULL, credited_invested numeric(20, 2) NULL, order_id uuid NULL, pending_reason text NULL, created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS stock_trades ( id uuid PRIMARY KEY, fund_id uuid NOT NULL REFERENCES funds(id), instrument_code text NOT NULL, stock_name text NULL, quantity numeric(28, 8) NOT NULL CHECK (quantity > 0), price numeric(20, 4) NOT NULL CHECK (price > 0), cost_cash numeric(20, 2) NOT NULL CHECK (cost_cash >= 0), is_synthetic boolean NOT NULL, executed_at timestamptz NOT NULL ); CREATE TABLE IF NOT EXISTS stock_trade_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, trade_id uuid NOT NULL REFERENCES stock_trades(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS stock_positions ( fund_id uuid NOT NULL REFERENCES funds(id), instrument_code text NOT NULL, stock_name text NULL, quantity numeric(28, 8) NOT NULL CHECK (quantity > 0), cost_cash numeric(20, 2) NOT NULL CHECK (cost_cash >= 0), last_traded_at timestamptz NOT NULL, PRIMARY KEY (fund_id, instrument_code) ); CREATE TABLE IF NOT EXISTS stock_sells ( id uuid PRIMARY KEY, fund_id uuid NOT NULL REFERENCES funds(id), instrument_code text NOT NULL, stock_name text NULL, quantity numeric(28, 8) NOT NULL CHECK (quantity > 0), price numeric(20, 4) NOT NULL CHECK (price > 0), fee_amount numeric(20, 2) NOT NULL CHECK (fee_amount >= 0), proceeds numeric(20, 2) NOT NULL CHECK (proceeds >= 0), is_synthetic boolean NOT NULL, executed_at timestamptz NOT NULL ); CREATE TABLE IF NOT EXISTS stock_sell_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, sell_id uuid NOT NULL REFERENCES stock_sells(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS bond_trades ( id uuid PRIMARY KEY, fund_id uuid NOT NULL REFERENCES funds(id), instrument_code text NOT NULL, bond_name text NULL, quantity numeric(28, 8) NOT NULL CHECK (quantity > 0), price numeric(20, 4) NOT NULL CHECK (price > 0), cost_cash numeric(20, 2) NOT NULL CHECK (cost_cash >= 0), is_synthetic boolean NOT NULL, executed_at timestamptz NOT NULL ); ALTER TABLE bond_trades ADD COLUMN IF NOT EXISTS clean_price numeric(20, 4) NOT NULL DEFAULT 0; ALTER TABLE bond_trades ADD COLUMN IF NOT EXISTS accrued_interest numeric(20, 4) NOT NULL DEFAULT 0; ALTER TABLE bond_trades ADD COLUMN IF NOT EXISTS par_value numeric(20, 4) NOT NULL DEFAULT 100; ALTER TABLE bond_trades ADD COLUMN IF NOT EXISTS settlement_date date NOT NULL DEFAULT CURRENT_DATE; ALTER TABLE bond_trades ADD COLUMN IF NOT EXISTS coupon_rate numeric(12, 6) NULL; ALTER TABLE bond_trades ADD COLUMN IF NOT EXISTS value_date date NULL; ALTER TABLE bond_trades ADD COLUMN IF NOT EXISTS maturity_date date NULL; ALTER TABLE bond_trades ADD COLUMN IF NOT EXISTS trade_date date NOT NULL DEFAULT CURRENT_DATE; CREATE TABLE IF NOT EXISTS bond_trade_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, trade_id uuid NOT NULL REFERENCES bond_trades(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS bond_cashflow_events ( id uuid PRIMARY KEY, fund_id uuid NOT NULL REFERENCES funds(id), instrument_code text NOT NULL, bond_name text NULL, event_type text NOT NULL CHECK (event_type IN ('coupon', 'maturity', 'redemption')), event_date date NOT NULL, quantity numeric(28, 8) NOT NULL CHECK (quantity > 0), amount numeric(20, 2) NOT NULL CHECK (amount >= 0), note text NULL, is_synthetic boolean NOT NULL, created_at timestamptz NOT NULL ); ALTER TABLE bond_cashflow_events DROP CONSTRAINT IF EXISTS bond_cashflow_events_event_type_check; ALTER TABLE bond_cashflow_events ADD CONSTRAINT bond_cashflow_events_event_type_check CHECK (event_type IN ('coupon', 'maturity', 'redemption')); CREATE TABLE IF NOT EXISTS bond_cashflow_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, event_id uuid NOT NULL REFERENCES bond_cashflow_events(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS bond_sells ( id uuid PRIMARY KEY, fund_id uuid NOT NULL REFERENCES funds(id), instrument_code text NOT NULL, bond_name text NULL, quantity numeric(28, 8) NOT NULL CHECK (quantity > 0), price numeric(20, 4) NOT NULL CHECK (price > 0), clean_price numeric(20, 4) NOT NULL DEFAULT 0, accrued_interest numeric(20, 4) NOT NULL DEFAULT 0, par_value numeric(20, 4) NOT NULL DEFAULT 100, settlement_date date NOT NULL DEFAULT CURRENT_DATE, trade_date date NOT NULL DEFAULT CURRENT_DATE, fee_amount numeric(20, 2) NOT NULL DEFAULT 0 CHECK (fee_amount >= 0), proceeds numeric(20, 2) NOT NULL DEFAULT 0, cost_released numeric(20, 2) NOT NULL DEFAULT 0, realized_pnl numeric(20, 2) NOT NULL DEFAULT 0, is_synthetic boolean NOT NULL, executed_at timestamptz NOT NULL ); CREATE TABLE IF NOT EXISTS bond_sell_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, sell_id uuid NOT NULL REFERENCES bond_sells(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS bond_positions ( fund_id uuid NOT NULL REFERENCES funds(id), instrument_code text NOT NULL, bond_name text NULL, quantity numeric(28, 8) NOT NULL CHECK (quantity > 0), cost_cash numeric(20, 2) NOT NULL CHECK (cost_cash >= 0), last_traded_at timestamptz NOT NULL, PRIMARY KEY (fund_id, instrument_code) ); CREATE TABLE IF NOT EXISTS instrument_snapshots ( instrument_code text NOT NULL, asset_class text NOT NULL CHECK (asset_class IN ('stock', 'bond')), snapshot_date date NOT NULL, price numeric(28, 8) NOT NULL CHECK (price > 0), 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, asset_class, snapshot_date) ); CREATE INDEX IF NOT EXISTS instrument_snapshots_date_idx ON instrument_snapshots (instrument_code, asset_class, snapshot_date DESC); CREATE TABLE IF NOT EXISTS dividend_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, record_id uuid NOT NULL REFERENCES dividend_records(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS investment_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_run_date date NOT NULL, created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS investment_plan_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, plan_id uuid NOT NULL REFERENCES investment_plans(id), fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS investment_plan_runs ( plan_id uuid NOT NULL REFERENCES investment_plans(id), run_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, run_date) ); """ 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 CreatedAt = DateTimeOffset.UtcNow } 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) CreatedAt = reader.GetFieldValue(9) } 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, created_at 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 stockTradeRecordFromReader (reader: DbDataReader) : StockTradeRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) InstrumentCode = reader.GetString(2) StockName = if reader.IsDBNull(3) then None else Some(reader.GetString(3)) Quantity = reader.GetDecimal(4) Price = reader.GetDecimal(5) CostCash = reader.GetDecimal(6) IsSynthetic = reader.GetBoolean(7) ExecutedAt = reader.GetFieldValue(8) } let insertStockTrade connection transaction (trade: StockTradeRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO stock_trades (id, fund_id, instrument_code, stock_name, quantity, price, cost_cash, is_synthetic, executed_at) VALUES (@id, @fund_id, @instrument_code, @stock_name, @quantity, @price, @cost_cash, @is_synthetic, @executed_at) """ addParameter command "id" NpgsqlDbType.Uuid (box trade.Id) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box trade.FundId) |> ignore addParameter command "instrument_code" NpgsqlDbType.Text (box trade.InstrumentCode) |> ignore let nameParameter = match trade.StockName with | Some name -> box name | None -> box DBNull.Value addParameter command "stock_name" NpgsqlDbType.Text nameParameter |> ignore addParameter command "quantity" NpgsqlDbType.Numeric (box trade.Quantity) |> ignore addParameter command "price" NpgsqlDbType.Numeric (box trade.Price) |> ignore addParameter command "cost_cash" NpgsqlDbType.Numeric (box trade.CostCash) |> ignore addParameter command "is_synthetic" NpgsqlDbType.Boolean (box trade.IsSynthetic) |> ignore addParameter command "executed_at" NpgsqlDbType.TimestampTz (box trade.ExecutedAt) |> ignore command.ExecuteNonQuery() |> ignore let insertStockTradeIdempotency connection transaction key requestHash tradeId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO stock_trade_idempotencies (idempotency_key, request_hash, trade_id, fund_id) VALUES (@idempotency_key, @request_hash, @trade_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "trade_id" NpgsqlDbType.Uuid (box tradeId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let findStockTradeIdempotency connection transaction key = use command = commandWithTransaction connection transaction """ SELECT request_hash, fund_id, trade_id FROM stock_trade_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 findStockTrade connection transaction tradeId = use command = commandWithTransaction connection transaction """ SELECT id, fund_id, instrument_code, stock_name, quantity, price, cost_cash, is_synthetic, executed_at FROM stock_trades WHERE id = @id """ addParameter command "id" NpgsqlDbType.Uuid (box tradeId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(stockTradeRecordFromReader reader) else None let stockTradeRequestHash (fundId: Guid) (command: StockTradeCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let name = command.StockName |> Option.defaultValue "" let payload = String.concat "|" [ "stock-trade" encoded (fundId.ToString("D")) encoded command.InstrumentCode encoded name encoded (command.Quantity.ToString("G29", invariant)) encoded (command.Price.ToString("G29", invariant)) ] Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let stockSellRecordFromReader (reader: DbDataReader) : StockSellRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) InstrumentCode = reader.GetString(2) StockName = if reader.IsDBNull(3) then None else Some(reader.GetString(3)) Quantity = reader.GetDecimal(4) Price = reader.GetDecimal(5) FeeAmount = reader.GetDecimal(6) Proceeds = reader.GetDecimal(7) IsSynthetic = reader.GetBoolean(8) ExecutedAt = reader.GetFieldValue(9) } let insertStockSell connection transaction (sell: StockSellRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO stock_sells (id, fund_id, instrument_code, stock_name, quantity, price, fee_amount, proceeds, is_synthetic, executed_at) VALUES (@id, @fund_id, @instrument_code, @stock_name, @quantity, @price, @fee_amount, @proceeds, @is_synthetic, @executed_at) """ addParameter command "id" NpgsqlDbType.Uuid (box sell.Id) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box sell.FundId) |> ignore addParameter command "instrument_code" NpgsqlDbType.Text (box sell.InstrumentCode) |> ignore let nameParameter = match sell.StockName with | Some name -> box name | None -> box DBNull.Value addParameter command "stock_name" NpgsqlDbType.Text nameParameter |> ignore addParameter command "quantity" NpgsqlDbType.Numeric (box sell.Quantity) |> ignore addParameter command "price" NpgsqlDbType.Numeric (box sell.Price) |> ignore addParameter command "fee_amount" NpgsqlDbType.Numeric (box sell.FeeAmount) |> ignore addParameter command "proceeds" NpgsqlDbType.Numeric (box sell.Proceeds) |> ignore addParameter command "is_synthetic" NpgsqlDbType.Boolean (box sell.IsSynthetic) |> ignore addParameter command "executed_at" NpgsqlDbType.TimestampTz (box sell.ExecutedAt) |> ignore command.ExecuteNonQuery() |> ignore let insertStockSellIdempotency connection transaction key requestHash sellId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO stock_sell_idempotencies (idempotency_key, request_hash, sell_id, fund_id) VALUES (@idempotency_key, @request_hash, @sell_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "sell_id" NpgsqlDbType.Uuid (box sellId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let findStockSellIdempotency connection transaction key = use command = commandWithTransaction connection transaction """ SELECT request_hash, fund_id, sell_id FROM stock_sell_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 findStockSell connection transaction sellId = use command = commandWithTransaction connection transaction """ SELECT id, fund_id, instrument_code, stock_name, quantity, price, fee_amount, proceeds, is_synthetic, executed_at FROM stock_sells WHERE id = @id """ addParameter command "id" NpgsqlDbType.Uuid (box sellId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(stockSellRecordFromReader reader) else None let stockSellRequestHash (fundId: Guid) (command: StockSellCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let name = command.StockName |> Option.defaultValue "" let payload = String.concat "|" [ "stock-sell" encoded (fundId.ToString("D")) encoded command.InstrumentCode encoded name encoded (command.Quantity.ToString("G29", invariant)) encoded (command.Price.ToString("G29", invariant)) encoded (command.FeeAmount.ToString("G29", invariant)) ] Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let bondTradeRecordFromReader (reader: DbDataReader) : BondTradeRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) InstrumentCode = reader.GetString(2) BondName = if reader.IsDBNull(3) then None else Some(reader.GetString(3)) Quantity = reader.GetDecimal(4) Price = reader.GetDecimal(5) CleanPrice = reader.GetDecimal(6) AccruedInterest = reader.GetDecimal(7) ParValue = reader.GetDecimal(8) SettlementDate = reader.GetFieldValue(9) CouponRate = readDecimalOption reader 10 ValueDate = if reader.IsDBNull(11) then None else Some(reader.GetFieldValue(11)) MaturityDate = if reader.IsDBNull(12) then None else Some(reader.GetFieldValue(12)) TradeDate = reader.GetFieldValue(13) CostCash = reader.GetDecimal(14) IsSynthetic = reader.GetBoolean(15) ExecutedAt = reader.GetFieldValue(16) } let bondTradeColumns = "id, fund_id, instrument_code, bond_name, quantity, price, clean_price, accrued_interest, par_value, settlement_date, coupon_rate, value_date, maturity_date, trade_date, cost_cash, is_synthetic, executed_at" let insertBondTrade connection transaction (trade: BondTradeRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO bond_trades (id, fund_id, instrument_code, bond_name, quantity, price, clean_price, accrued_interest, par_value, settlement_date, coupon_rate, value_date, maturity_date, trade_date, cost_cash, is_synthetic, executed_at) VALUES (@id, @fund_id, @instrument_code, @bond_name, @quantity, @price, @clean_price, @accrued_interest, @par_value, @settlement_date, @coupon_rate, @value_date, @maturity_date, @trade_date, @cost_cash, @is_synthetic, @executed_at) """ addParameter command "id" NpgsqlDbType.Uuid (box trade.Id) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box trade.FundId) |> ignore addParameter command "instrument_code" NpgsqlDbType.Text (box trade.InstrumentCode) |> ignore let nameParameter = match trade.BondName with | Some name -> box name | None -> box DBNull.Value addParameter command "bond_name" NpgsqlDbType.Text nameParameter |> ignore addParameter command "quantity" NpgsqlDbType.Numeric (box trade.Quantity) |> ignore addParameter command "price" NpgsqlDbType.Numeric (box trade.Price) |> ignore addParameter command "clean_price" NpgsqlDbType.Numeric (box trade.CleanPrice) |> ignore addParameter command "accrued_interest" NpgsqlDbType.Numeric (box trade.AccruedInterest) |> ignore addParameter command "par_value" NpgsqlDbType.Numeric (box trade.ParValue) |> ignore addParameter command "settlement_date" NpgsqlDbType.Date (box trade.SettlementDate) |> ignore let couponParameter = match trade.CouponRate with | Some value -> box value | None -> box DBNull.Value addParameter command "coupon_rate" NpgsqlDbType.Numeric couponParameter |> ignore let valueDateParameter = match trade.ValueDate with | Some value -> box value | None -> box DBNull.Value addParameter command "value_date" NpgsqlDbType.Date valueDateParameter |> ignore let maturityParameter = match trade.MaturityDate with | Some value -> box value | None -> box DBNull.Value addParameter command "maturity_date" NpgsqlDbType.Date maturityParameter |> ignore addParameter command "trade_date" NpgsqlDbType.Date (box trade.TradeDate) |> ignore addParameter command "cost_cash" NpgsqlDbType.Numeric (box trade.CostCash) |> ignore addParameter command "is_synthetic" NpgsqlDbType.Boolean (box trade.IsSynthetic) |> ignore addParameter command "executed_at" NpgsqlDbType.TimestampTz (box trade.ExecutedAt) |> ignore command.ExecuteNonQuery() |> ignore let insertBondTradeIdempotency connection transaction key requestHash tradeId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO bond_trade_idempotencies (idempotency_key, request_hash, trade_id, fund_id) VALUES (@idempotency_key, @request_hash, @trade_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "trade_id" NpgsqlDbType.Uuid (box tradeId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let findBondTradeIdempotency connection transaction key = use command = commandWithTransaction connection transaction """ SELECT request_hash, fund_id, trade_id FROM bond_trade_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 findBondTrade connection transaction tradeId = use command = commandWithTransaction connection transaction $""" SELECT {bondTradeColumns} FROM bond_trades WHERE id = @id """ addParameter command "id" NpgsqlDbType.Uuid (box tradeId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(bondTradeRecordFromReader reader) else None let bondTradeRequestHash (fundId: Guid) (command: BondTradeCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let name = command.BondName |> Option.defaultValue "" let dateText (value: DateOnly) = value.ToString("yyyy-MM-dd", invariant) let optionDate = command.ValueDate |> Option.map dateText |> Option.defaultValue "" let optionMaturity = command.MaturityDate |> Option.map dateText |> Option.defaultValue "" let optionCoupon = command.CouponRate |> Option.map (fun v -> v.ToString("G29", invariant)) |> Option.defaultValue "" let optionTradeDate = command.TradeDate |> Option.map dateText |> Option.defaultValue "" let payload = String.concat "|" [ "bond-trade" encoded (fundId.ToString("D")) encoded command.InstrumentCode encoded name encoded (command.Quantity.ToString("G29", invariant)) encoded (command.Price.ToString("G29", invariant)) encoded (command.CleanPrice.ToString("G29", invariant)) encoded (command.AccruedInterest.ToString("G29", invariant)) encoded (command.ParValue.ToString("G29", invariant)) encoded (dateText command.SettlementDate) encoded optionTradeDate encoded optionCoupon encoded optionDate encoded optionMaturity ] Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let bondCashflowRecordFromReader (reader: DbDataReader) : BondCashflowRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) InstrumentCode = reader.GetString(2) BondName = if reader.IsDBNull(3) then None else Some(reader.GetString(3)) EventType = reader.GetString(4) EventDate = reader.GetFieldValue(5) Quantity = reader.GetDecimal(6) Amount = reader.GetDecimal(7) Note = readStringOption reader 8 IsSynthetic = reader.GetBoolean(9) CreatedAt = reader.GetFieldValue(10) } let bondCashflowColumns = "id, fund_id, instrument_code, bond_name, event_type, event_date, quantity, amount, note, is_synthetic, created_at" let insertBondCashflow connection transaction (record: BondCashflowRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO bond_cashflow_events (id, fund_id, instrument_code, bond_name, event_type, event_date, quantity, amount, note, is_synthetic, created_at) VALUES (@id, @fund_id, @instrument_code, @bond_name, @event_type, @event_date, @quantity, @amount, @note, @is_synthetic, @created_at) """ addParameter command "id" NpgsqlDbType.Uuid (box record.Id) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box record.FundId) |> ignore addParameter command "instrument_code" NpgsqlDbType.Text (box record.InstrumentCode) |> ignore let nameParameter = match record.BondName with | Some name -> box name | None -> box DBNull.Value addParameter command "bond_name" NpgsqlDbType.Text nameParameter |> ignore addParameter command "event_type" NpgsqlDbType.Text (box record.EventType) |> ignore addParameter command "event_date" NpgsqlDbType.Date (box record.EventDate) |> ignore addParameter command "quantity" NpgsqlDbType.Numeric (box record.Quantity) |> ignore addParameter command "amount" NpgsqlDbType.Numeric (box record.Amount) |> ignore let noteParameter = match record.Note with | Some note -> box note | None -> box DBNull.Value addParameter command "note" NpgsqlDbType.Text noteParameter |> ignore addParameter command "is_synthetic" NpgsqlDbType.Boolean (box record.IsSynthetic) |> ignore addParameter command "created_at" NpgsqlDbType.TimestampTz (box record.CreatedAt) |> ignore command.ExecuteNonQuery() |> ignore let insertBondCashflowIdempotency connection transaction key requestHash eventId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO bond_cashflow_idempotencies (idempotency_key, request_hash, event_id, fund_id) VALUES (@idempotency_key, @request_hash, @event_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "event_id" NpgsqlDbType.Uuid (box eventId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let findBondCashflowIdempotency connection transaction key = use command = commandWithTransaction connection transaction """ SELECT request_hash, fund_id, event_id FROM bond_cashflow_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 findBondCashflow connection transaction eventId = use command = commandWithTransaction connection transaction $""" SELECT {bondCashflowColumns} FROM bond_cashflow_events WHERE id = @id """ addParameter command "id" NpgsqlDbType.Uuid (box eventId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(bondCashflowRecordFromReader reader) else None let bondCashflowRequestHash (fundId: Guid) (command: BondCashflowCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let name = command.BondName |> Option.defaultValue "" let note = command.Note |> Option.defaultValue "" let payload = String.concat "|" [ "bond-cashflow" encoded (fundId.ToString("D")) encoded command.InstrumentCode encoded name encoded command.EventType encoded (command.EventDate.ToString("yyyy-MM-dd", invariant)) encoded (command.Quantity.ToString("G29", invariant)) encoded (command.Amount.ToString("G29", invariant)) encoded note ] Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let bondSellRecordFromReader (reader: DbDataReader) : BondSellRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) InstrumentCode = reader.GetString(2) BondName = if reader.IsDBNull(3) then None else Some(reader.GetString(3)) Quantity = reader.GetDecimal(4) Price = reader.GetDecimal(5) CleanPrice = reader.GetDecimal(6) AccruedInterest = reader.GetDecimal(7) ParValue = reader.GetDecimal(8) SettlementDate = reader.GetFieldValue(9) TradeDate = reader.GetFieldValue(10) FeeAmount = reader.GetDecimal(11) Proceeds = reader.GetDecimal(12) CostReleased = reader.GetDecimal(13) RealizedPnl = reader.GetDecimal(14) IsSynthetic = reader.GetBoolean(15) ExecutedAt = reader.GetFieldValue(16) } let bondSellColumns = "id, fund_id, instrument_code, bond_name, quantity, price, clean_price, accrued_interest, par_value, settlement_date, trade_date, fee_amount, proceeds, cost_released, realized_pnl, is_synthetic, executed_at" let insertBondSell connection transaction (record: BondSellRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO bond_sells (id, fund_id, instrument_code, bond_name, quantity, price, clean_price, accrued_interest, par_value, settlement_date, trade_date, fee_amount, proceeds, cost_released, realized_pnl, is_synthetic, executed_at) VALUES (@id, @fund_id, @instrument_code, @bond_name, @quantity, @price, @clean_price, @accrued_interest, @par_value, @settlement_date, @trade_date, @fee_amount, @proceeds, @cost_released, @realized_pnl, @is_synthetic, @executed_at) """ addParameter command "id" NpgsqlDbType.Uuid (box record.Id) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box record.FundId) |> ignore addParameter command "instrument_code" NpgsqlDbType.Text (box record.InstrumentCode) |> ignore let nameParameter = match record.BondName with | Some name -> box name | None -> box DBNull.Value addParameter command "bond_name" NpgsqlDbType.Text nameParameter |> ignore addParameter command "quantity" NpgsqlDbType.Numeric (box record.Quantity) |> ignore addParameter command "price" NpgsqlDbType.Numeric (box record.Price) |> ignore addParameter command "clean_price" NpgsqlDbType.Numeric (box record.CleanPrice) |> ignore addParameter command "accrued_interest" NpgsqlDbType.Numeric (box record.AccruedInterest) |> ignore addParameter command "par_value" NpgsqlDbType.Numeric (box record.ParValue) |> ignore addParameter command "settlement_date" NpgsqlDbType.Date (box record.SettlementDate) |> ignore addParameter command "trade_date" NpgsqlDbType.Date (box record.TradeDate) |> ignore addParameter command "fee_amount" NpgsqlDbType.Numeric (box record.FeeAmount) |> ignore addParameter command "proceeds" NpgsqlDbType.Numeric (box record.Proceeds) |> ignore addParameter command "cost_released" NpgsqlDbType.Numeric (box record.CostReleased) |> ignore addParameter command "realized_pnl" NpgsqlDbType.Numeric (box record.RealizedPnl) |> ignore addParameter command "is_synthetic" NpgsqlDbType.Boolean (box record.IsSynthetic) |> ignore addParameter command "executed_at" NpgsqlDbType.TimestampTz (box record.ExecutedAt) |> ignore command.ExecuteNonQuery() |> ignore let insertBondSellIdempotency connection transaction key requestHash sellId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO bond_sell_idempotencies (idempotency_key, request_hash, sell_id, fund_id) VALUES (@idempotency_key, @request_hash, @sell_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "sell_id" NpgsqlDbType.Uuid (box sellId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let findBondSellIdempotency connection transaction key = use command = commandWithTransaction connection transaction """ SELECT request_hash, fund_id, sell_id FROM bond_sell_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 findBondSell connection transaction sellId = use command = commandWithTransaction connection transaction $""" SELECT {bondSellColumns} FROM bond_sells WHERE id = @id """ addParameter command "id" NpgsqlDbType.Uuid (box sellId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(bondSellRecordFromReader reader) else None let bondSellRequestHash (fundId: Guid) (command: BondSellCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let name = command.BondName |> Option.defaultValue "" let dateText (value: DateOnly) = value.ToString("yyyy-MM-dd", invariant) let optionTradeDate = command.TradeDate |> Option.map dateText |> Option.defaultValue "" let payload = String.concat "|" [ "bond-sell" encoded (fundId.ToString("D")) encoded command.InstrumentCode encoded name encoded (command.Quantity.ToString("G29", invariant)) encoded (command.Price.ToString("G29", invariant)) encoded (command.CleanPrice.ToString("G29", invariant)) encoded (command.AccruedInterest.ToString("G29", invariant)) encoded (command.ParValue.ToString("G29", invariant)) encoded (dateText command.SettlementDate) encoded optionTradeDate encoded (command.FeeAmount.ToString("G29", invariant)) ] 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 investmentPlanRecordFromReader (reader: DbDataReader) : InvestmentPlanRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) InstrumentCode = reader.GetString(2) Amount = reader.GetDecimal(3) Frequency = match InvestmentPlanPolicy.parseFrequency (reader.GetString(4)) with | Some frequency -> frequency | None -> failwith "investment plan frequency is invalid" Status = reader.GetString(5) IsSynthetic = reader.GetBoolean(6) AnchorDate = reader.GetFieldValue(7) NextRunDate = reader.GetFieldValue(8) CreatedAt = reader.GetFieldValue(9) LastRunStatus = readStringOption reader 10 LastRunDate = if reader.IsDBNull(11) then None else Some(reader.GetFieldValue(11)) } let findInvestmentPlan 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_run_date, p.created_at, r.status, r.run_date FROM investment_plans p LEFT JOIN LATERAL ( SELECT status, run_date FROM investment_plan_runs WHERE plan_id = p.id ORDER BY executed_at DESC, run_date DESC LIMIT 1 ) r 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(investmentPlanRecordFromReader reader) else None let findInvestmentPlanIdempotency connection transaction key = use command = commandWithTransaction connection transaction "SELECT request_hash, fund_id, plan_id FROM investment_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 insertInvestmentPlan connection transaction (plan: InvestmentPlanRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO investment_plans (id, fund_id, instrument_code, amount, frequency, status, is_synthetic, anchor_date, next_run_date) VALUES (@id, @fund_id, @instrument_code, @amount, @frequency, @status, @is_synthetic, @anchor_date, @next_run_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 (InvestmentPlanPolicy.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_run_date" NpgsqlDbType.Date (box plan.NextRunDate) |> ignore use reader = command.ExecuteReader() reader.Read() |> ignore reader.GetFieldValue(0) let insertInvestmentPlanIdempotency connection transaction key requestHash planId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO investment_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 investmentPlanRequestHash (fundId: Guid) (command: InvestmentPlanCommand) = 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 "|" [ "investment-plan" encoded (fundId.ToString("D")) encoded code (encoded (command.Amount.ToString("G29", invariant))) (encoded (InvestmentPlanPolicy.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() let rebalanceExecutionRecordFromReader (reader: DbDataReader) : RebalanceExecutionRecord = { PlanId = reader.GetGuid(0) RunDate = reader.GetFieldValue(1) InstrumentCode = reader.GetString(2) Action = reader.GetString(3) Amount = reader.GetDecimal(4) Status = reader.GetString(5) OrderId = if reader.IsDBNull(6) then None else Some(reader.GetGuid(6)) PendingReason = readStringOption reader 7 ExecutedAt = reader.GetFieldValue(8) } let upsertRebalanceExecution connection transaction (planId: Guid) (runDate: DateOnly) (outcome: RebalanceOrderOutcome) = let amount = match Decimal.TryParse(outcome.Amount, NumberStyles.Float, CultureInfo.InvariantCulture) with | true, value -> value | false, _ -> 0m use command = commandWithTransaction connection transaction """ INSERT INTO rebalance_executions (plan_id, run_date, instrument_code, action, amount, status, order_id, pending_reason) VALUES (@plan_id, @run_date, @instrument_code, @action, @amount, @status, @order_id, @pending_reason) ON CONFLICT (plan_id, run_date, instrument_code) DO UPDATE SET action = EXCLUDED.action, amount = EXCLUDED.amount, status = EXCLUDED.status, order_id = EXCLUDED.order_id, pending_reason = EXCLUDED.pending_reason, executed_at = now() """ addParameter command "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore addParameter command "run_date" NpgsqlDbType.Date (box runDate) |> ignore addParameter command "instrument_code" NpgsqlDbType.Text (box outcome.InstrumentCode) |> ignore addParameter command "action" NpgsqlDbType.Text (box outcome.Action) |> ignore addParameter command "amount" NpgsqlDbType.Numeric (box amount) |> ignore addParameter command "status" NpgsqlDbType.Text (box outcome.Status) |> ignore let orderParameter = match outcome.OrderId with | Some value -> box value | None -> box DBNull.Value addParameter command "order_id" NpgsqlDbType.Uuid orderParameter |> ignore let reasonParameter = match outcome.PendingReason with | Some value -> box value | None -> box DBNull.Value addParameter command "pending_reason" NpgsqlDbType.Text reasonParameter |> ignore command.ExecuteNonQuery() |> ignore let dividendRecordFromReader (reader: DbDataReader) : DividendRecord = { Id = reader.GetGuid(0) FundId = reader.GetGuid(1) InstrumentCode = reader.GetString(2) NavDate = reader.GetFieldValue(3) Dps = reader.GetDecimal(4) Mode = match DividendPolicy.parseMode (reader.GetString(5)) with | Some mode -> mode | None -> failwith "dividend mode is invalid" Status = reader.GetString(6) IsSynthetic = reader.GetBoolean(7) GrossCash = readDecimalOption reader 8 CreditedUnits = readDecimalOption reader 9 CreditedInvested = readDecimalOption reader 10 OrderId = if reader.IsDBNull(11) then None else Some(reader.GetGuid(11)) PendingReason = readStringOption reader 12 CreatedAt = reader.GetFieldValue(13) } let findDividendRecord connection transaction recordId = use command = commandWithTransaction connection transaction """ SELECT id, fund_id, instrument_code, nav_date, dps, mode, status, is_synthetic, gross_cash, credited_units, credited_invested, order_id, pending_reason, created_at FROM dividend_records WHERE id = @record_id """ addParameter command "record_id" NpgsqlDbType.Uuid (box recordId) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(dividendRecordFromReader reader) else None let findDividendIdempotency connection transaction key = use command = commandWithTransaction connection transaction "SELECT request_hash, fund_id, record_id FROM dividend_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 insertDividendRecord connection transaction (record: DividendRecord) = use command = commandWithTransaction connection transaction """ INSERT INTO dividend_records (id, fund_id, instrument_code, nav_date, dps, mode, status, is_synthetic, gross_cash) VALUES (@id, @fund_id, @instrument_code, @nav_date, @dps, @mode, @status, @is_synthetic, @gross_cash) RETURNING created_at """ addParameter command "id" NpgsqlDbType.Uuid (box record.Id) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box record.FundId) |> ignore addParameter command "instrument_code" NpgsqlDbType.Text (box record.InstrumentCode) |> ignore addParameter command "nav_date" NpgsqlDbType.Date (box record.NavDate) |> ignore addParameter command "dps" NpgsqlDbType.Numeric (box record.Dps) |> ignore addParameter command "mode" NpgsqlDbType.Text (box (DividendPolicy.modeText record.Mode)) |> ignore addParameter command "status" NpgsqlDbType.Text (box record.Status) |> ignore addParameter command "is_synthetic" NpgsqlDbType.Boolean (box record.IsSynthetic) |> ignore let grossParameter = match record.GrossCash with | Some gross -> box gross | None -> box DBNull.Value addParameter command "gross_cash" NpgsqlDbType.Numeric grossParameter |> ignore use reader = command.ExecuteReader() reader.Read() |> ignore reader.GetFieldValue(0) let insertDividendIdempotency connection transaction key requestHash recordId fundId = use command = commandWithTransaction connection transaction """ INSERT INTO dividend_idempotencies (idempotency_key, request_hash, record_id, fund_id) VALUES (@idempotency_key, @request_hash, @record_id, @fund_id) """ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore addParameter command "record_id" NpgsqlDbType.Uuid (box recordId) |> ignore addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore command.ExecuteNonQuery() |> ignore let finishDividendRecord connection transaction (recordId: Guid) (status: string) (pendingReason: string option) = use command = commandWithTransaction connection transaction "UPDATE dividend_records SET status = @status, pending_reason = @reason WHERE id = @record_id" addParameter command "status" NpgsqlDbType.Text (box status) |> ignore let reasonParameter = match pendingReason with | Some reason -> box reason | None -> box DBNull.Value addParameter command "reason" NpgsqlDbType.Text reasonParameter |> ignore addParameter command "record_id" NpgsqlDbType.Uuid (box recordId) |> ignore command.ExecuteNonQuery() |> ignore let dividendRequestHash (fundId: Guid) (command: DividendCommand) = let invariant = CultureInfo.InvariantCulture let encoded (value: string) = sprintf "%d:%s" value.Length value let payload = String.concat "|" [ "dividend" encoded (fundId.ToString("D")) encoded command.InstrumentCode encoded (command.NavDate.ToString("yyyy-MM-dd")) (encoded (command.Dps.ToString("G29", invariant))) (encoded (DividendPolicy.modeText command.Mode)) ] Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) let dividendSchemeKey (fundId: Guid) (command: DividendCommand) = DividendPolicy.schemeKey fundId command.InstrumentCode command.NavDate command.Mode let validateDividendCommand (today: DateOnly) (command: DividendCommand) = if String.IsNullOrWhiteSpace command.InstrumentCode then Error "instrument code cannot be empty" elif command.Dps <= 0m then Error "dividend per unit must be positive" else match DividendPolicy.validateDps command.Dps, DividendPolicy.validateNavDate command.NavDate today with | Ok(), Ok() -> Ok() | Error message, _ | _, Error message -> Error message 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 _.UpsertInstrumentSnapshots(records: InstrumentSnapshotRecord list) = if records.IsEmpty then () else use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try for record in records do use command = commandWithTransaction connection (Some transaction) """ INSERT INTO instrument_snapshots (instrument_code, asset_class, snapshot_date, price, source, source_revision, source_collected_at, source_payload_hash) VALUES (@instrument_code, @asset_class, @snapshot_date, @price, @source, @source_revision, @source_collected_at, @source_payload_hash) ON CONFLICT (instrument_code, asset_class, snapshot_date) DO UPDATE SET price = EXCLUDED.price, 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 instrument_snapshots.source_payload_hash = EXCLUDED.source_payload_hash THEN instrument_snapshots.first_seen_at ELSE now() END """ addParameter command "instrument_code" NpgsqlDbType.Text (box record.InstrumentCode) |> ignore addParameter command "asset_class" NpgsqlDbType.Text (box record.AssetClass) |> ignore addParameter command "snapshot_date" NpgsqlDbType.Date (box record.SnapshotDate) |> ignore addParameter command "price" NpgsqlDbType.Numeric (box record.Price) |> ignore addParameter command "source" NpgsqlDbType.Text (box record.Source) |> ignore addParameter command "source_revision" NpgsqlDbType.Text (box record.SourceRevision) |> ignore addParameter command "source_collected_at" NpgsqlDbType.TimestampTz (box record.SourceCollectedAt) |> ignore addParameter command "source_payload_hash" NpgsqlDbType.Text (box record.SourcePayloadHash) |> ignore command.ExecuteNonQuery() |> ignore transaction.Commit() with error -> try transaction.Rollback() with _ -> () raise error /// Latest snapshot price on or before `asOfDate` for every instrument of the /// given asset class held by the fund. Instruments without any snapshot are /// omitted so a caller can tell "no snapshot yet" from a stored price. member _.GetLatestSnapshots(fundId: Guid, assetClass: string, asOfDate: DateOnly) : Map = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT s.instrument_code, s.asset_class, s.snapshot_date, s.price, s.source, s.source_revision, s.source_collected_at, s.source_payload_hash FROM instrument_snapshots s JOIN ( SELECT instrument_code, MAX(snapshot_date) AS snapshot_date FROM instrument_snapshots WHERE asset_class = @asset_class AND snapshot_date <= @as_of_date GROUP BY instrument_code ) latest ON latest.instrument_code = s.instrument_code AND latest.snapshot_date = s.snapshot_date WHERE s.asset_class = @asset_class """ addParameter command "asset_class" NpgsqlDbType.Text (box assetClass) |> ignore addParameter command "as_of_date" NpgsqlDbType.Date (box asOfDate) |> ignore use reader = command.ExecuteReader() let records = System.Collections.Generic.Dictionary() while reader.Read() do let record : InstrumentSnapshotRecord = { InstrumentCode = reader.GetString(0) AssetClass = reader.GetString(1) SnapshotDate = reader.GetFieldValue(2) Price = reader.GetDecimal(3) Source = reader.GetString(4) SourceRevision = reader.GetString(5) SourceCollectedAt = reader.GetFieldValue(6) SourcePayloadHash = reader.GetString(7) } records.[record.InstrumentCode] <- record records |> Seq.map (fun pair -> pair.Key, pair.Value) |> Map.ofSeq 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 _.SetSipPlanStatus(fundId: Guid, planId: Guid, status: string) : SipPlanStatusResult = 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-plan-status:%O" planId)) |> ignore lockCommand.ExecuteNonQuery() |> ignore match findSipPlan connection (Some transaction) planId with | Some plan when plan.FundId = fundId -> if plan.Status = status then transaction.Commit() SipPlanStatusResult.SipPlanStatusChanged plan else use updateCommand = commandWithTransaction connection (Some transaction) "UPDATE sip_plans SET status = @status WHERE id = @plan_id" addParameter updateCommand "status" NpgsqlDbType.Text (box status) |> ignore addParameter updateCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore updateCommand.ExecuteNonQuery() |> ignore let updated = findSipPlan connection (Some transaction) planId transaction.Commit() match updated with | Some updatedPlan -> SipPlanStatusResult.SipPlanStatusChanged updatedPlan | None -> SipPlanStatusResult.SipPlanStatusNotFound | _ -> transaction.Rollback() SipPlanStatusResult.SipPlanStatusNotFound with error -> try transaction.Rollback() with _ -> () raise error member _.CreateInvestmentPlan(idempotencyKey: string, fundId: Guid, command: InvestmentPlanCommand, ?anchorOverride: DateOnly) : InvestmentPlanWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then InvestmentPlanWriteResult.InvestmentPlanInvalid "idempotency key cannot be empty" else match InvestmentPlanPolicy.validateAmount command.Amount with | Error message -> InvestmentPlanWriteResult.InvestmentPlanInvalid message | Ok() -> let fingerprint = investmentPlanRequestHash 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 findInvestmentPlanIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, planId) when existingHash = fingerprint && existingFundId = fundId -> match findInvestmentPlan connection (Some transaction) planId with | Some plan -> transaction.Commit() InvestmentPlanWriteResult.InvestmentPlanReplayed plan | None -> transaction.Rollback() InvestmentPlanWriteResult.InvestmentPlanInvalid "idempotency record references a missing plan" | Some _ -> transaction.Rollback() InvestmentPlanWriteResult.InvestmentPlanIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() InvestmentPlanWriteResult.InvestmentPlanFundNotFound | Some isSynthetic -> if instrumentExists connection (Some transaction) command.InstrumentCode then let anchorDate = anchorOverride |> Option.defaultWith (fun () -> ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow) let plan: InvestmentPlanRecord = { Id = Guid.NewGuid() FundId = fundId InstrumentCode = command.InstrumentCode Amount = command.Amount Frequency = command.Frequency Status = "active" IsSynthetic = isSynthetic AnchorDate = anchorDate NextRunDate = InvestmentPlanPolicy.nextRunDate command.Frequency anchorDate anchorDate CreatedAt = DateTimeOffset.UtcNow LastRunStatus = None LastRunDate = None } let createdAt = insertInvestmentPlan connection (Some transaction) plan insertInvestmentPlanIdempotency connection (Some transaction) idempotencyKey fingerprint plan.Id fundId transaction.Commit() InvestmentPlanWriteResult.InvestmentPlanCreated { plan with CreatedAt = createdAt } else transaction.Rollback() InvestmentPlanWriteResult.InvestmentPlanInstrumentNotFound with error -> try transaction.Rollback() with _ -> () raise error member _.GetInvestmentPlans(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_run_date, p.created_at, r.status, r.run_date FROM investment_plans p LEFT JOIN LATERAL ( SELECT status, run_date FROM investment_plan_runs WHERE plan_id = p.id ORDER BY executed_at DESC, run_date DESC LIMIT 1 ) r 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(investmentPlanRecordFromReader reader) records |> Seq.toList member this.RunInvestmentPlans(fundId: Guid, processingDate: DateOnly) : InvestmentPlanRunResult = 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 "investment-plan-run:%O" fundId)) |> ignore lockCommand.ExecuteNonQuery() |> ignore // 0. retry phase: runs that were pending_nav get re-confirmed once their NAV landed use retryCommand = commandWithTransaction connection (Some transaction) """ SELECT r.plan_id, r.run_date, r.order_id FROM investment_plan_runs r JOIN investment_plans p ON p.id = r.plan_id WHERE p.fund_id = @fund_id AND r.status = 'pending_nav' ORDER BY r.run_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 "investment-plan-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 investment_plan_runs SET status = 'succeeded', pending_reason = NULL, executed_at = now() WHERE plan_id = @plan_id AND run_date = @run_date" addParameter doneCommand "plan_id" NpgsqlDbType.Uuid (box retryPlanId) |> ignore addParameter doneCommand "run_date" NpgsqlDbType.Date (box retryDate) |> ignore doneCommand.ExecuteNonQuery() |> ignore | _ -> // still pending: keep the run as pending_nav for the next drive () use plansCommand = commandWithTransaction connection (Some transaction) "SELECT id, instrument_code, amount, frequency, anchor_date, next_run_date FROM investment_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 InvestmentPlanPolicy.parseFrequency (plansReader.GetString(3)) with | Some frequency -> frequency | None -> failwith "investment plan frequency is invalid"), plansReader.GetFieldValue(4), plansReader.GetFieldValue(5) ) plansReader.Close() let readRunsUpTo (planId: Guid) (upTo: DateOnly) = use replayCommand = commandWithTransaction connection (Some transaction) """ SELECT run_date, status, order_id, pending_reason FROM investment_plan_runs WHERE plan_id = @plan_id AND run_date <= @up_to ORDER BY run_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( { RunDate = 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 runPlanRow (planId: Guid, code: string, amount: decimal, frequency: InvestmentFrequency, anchor: DateOnly, nextDate: DateOnly) = let dueDates = InvestmentPlanPolicy.dueDates frequency anchor nextDate processingDate let outcomes = ResizeArray() let claimRun (runDate: DateOnly) = use insertCommand = commandWithTransaction connection (Some transaction) """ INSERT INTO investment_plan_runs (plan_id, run_date, amount, fee_amount, status) VALUES (@plan_id, @run_date, @amount, @fee_amount, 'processing') ON CONFLICT (plan_id, run_date) DO NOTHING RETURNING run_date """ addParameter insertCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore addParameter insertCommand "run_date" NpgsqlDbType.Date (box runDate) |> 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 setRun (runDate: DateOnly) (status: string) (orderId: Guid option) (reason: string option) = use updateCommand = commandWithTransaction connection (Some transaction) """ UPDATE investment_plan_runs SET status = @status, order_id = @order_id, pending_reason = @reason, executed_at = now() WHERE plan_id = @plan_id AND run_date = @run_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 "run_date" NpgsqlDbType.Date (box runDate) |> ignore updateCommand.ExecuteNonQuery() |> ignore let findExistingRun (runDate: DateOnly) = use selectCommand = commandWithTransaction connection (Some transaction) "SELECT status, order_id, pending_reason FROM investment_plan_runs WHERE plan_id = @plan_id AND run_date = @run_date" addParameter selectCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore addParameter selectCommand "run_date" NpgsqlDbType.Date (box runDate) |> 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 (runDate: DateOnly) = let orderKey = sprintf "investment-plan:%O:%s" planId (runDate.ToString("yyyy-MM-dd")) match this.CreateSubscriptionOrder(orderKey, fundId, { FundCode = code; Amount = amount; FeeAmount = 0m }, runDate) with | SubscriptionOrderWriteResult.OrderCreated order -> Some order.Id | SubscriptionOrderWriteResult.OrderReplayed order -> Some order.Id | SubscriptionOrderWriteResult.OrderInsufficientFunds -> None | other -> failwithf "unexpected investment plan order result: %A" other for runDate in dueDates do // 1. claim the slot atomically: same plan + same run date executes once if claimRun runDate then // 2. place the order through the shared pipeline with a deterministic key match orderIdFor runDate with | None -> setRun runDate "insufficient_cash" None (Some "available cash is not enough for the scheduled amount") outcomes.Add( { RunDate = runDate 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 "investment-plan-confirm:%O:%s" planId (runDate.ToString("yyyy-MM-dd")) match this.ConfirmSubscriptionOrder(confirmKey, fundId, orderId) with | SubscriptionConfirmResult.OrderConfirmed _ | SubscriptionConfirmResult.ConfirmReplayed _ -> setRun runDate "succeeded" (Some orderId) None outcomes.Add({ RunDate = runDate; Status = "succeeded"; OrderId = Some orderId; PendingReason = None }) | SubscriptionConfirmResult.ConfirmPendingNav record -> setRun runDate "pending_nav" (Some orderId) record.PendingReason outcomes.Add({ RunDate = runDate; Status = "pending_nav"; OrderId = Some orderId; PendingReason = record.PendingReason }) | other -> setRun runDate "failed" (Some orderId) (Some(sprintf "%A" other)) outcomes.Add({ RunDate = runDate; Status = "failed"; OrderId = Some orderId; PendingReason = Some(sprintf "%A" other) }) else match findExistingRun runDate with | Some("succeeded", orderId, reason) -> outcomes.Add({ RunDate = runDate; Status = "succeeded"; OrderId = orderId; PendingReason = reason }) | Some(status, orderId, reason) -> outcomes.Add({ RunDate = runDate; Status = status; OrderId = orderId; PendingReason = reason }) | None -> outcomes.Add({ RunDate = runDate; Status = "unknown"; OrderId = None; PendingReason = None }) // 4. roll the plan pointer forward past the processed window let rolled = match dueDates with | [] -> nextDate | lastDueDates -> InvestmentPlanPolicy.nextRunDate frequency anchor ((List.last lastDueDates).AddDays 1) if not (List.isEmpty dueDates) then use rollCommand = commandWithTransaction connection (Some transaction) "UPDATE investment_plans SET next_run_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 processing date: replay what was already executed readRunsUpTo planId processingDate else outcomes |> Seq.toList { PlanId = planId InstrumentCode = code Amount = amount Frequency = frequency Runs = replayedOutcomes NextRunDate = rolled } let planResults = plans |> Seq.map runPlanRow |> Seq.toList transaction.Commit() { FundId = fundId; ProcessingDate = processingDate; Plans = planResults } with error -> try transaction.Rollback() with _ -> () raise error 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 _.ArchiveRebalancePlan(fundId: Guid, planId: Guid) : RebalanceArchiveResult = 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 "rebalance-archive:%O" planId)) |> ignore lockCommand.ExecuteNonQuery() |> ignore match findRebalancePlan connection (Some transaction) planId with | Some plan when plan.FundId = fundId -> if plan.Status = "archived" then transaction.Commit() RebalanceArchiveResult.RebalancePlanArchived plan else use updateCommand = commandWithTransaction connection (Some transaction) "UPDATE rebalance_plans SET status = 'archived' WHERE id = @plan_id" addParameter updateCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore updateCommand.ExecuteNonQuery() |> ignore let updated = findRebalancePlan connection (Some transaction) planId transaction.Commit() match updated with | Some updatedPlan -> RebalanceArchiveResult.RebalancePlanArchived updatedPlan | None -> RebalanceArchiveResult.RebalanceArchiveNotFound | _ -> transaction.Rollback() RebalanceArchiveResult.RebalanceArchiveNotFound with error -> try transaction.Rollback() with _ -> () raise error 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 when plan.Status = "archived" -> Error "rebalance plan is archived" | 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) }) let outcomeList = outcomes |> Seq.toList use persistTransaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try for outcome in outcomeList do upsertRebalanceExecution connection (Some persistTransaction) plan.Id runDate outcome persistTransaction.Commit() with error -> try persistTransaction.Rollback() with _ -> () raise error Ok { PlanId = plan.Id RunDate = runDate Outcomes = outcomeList } member _.GetRebalanceExecutions(fundId: Guid) = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT e.plan_id, e.run_date, e.instrument_code, e.action, e.amount, e.status, e.order_id, e.pending_reason, e.executed_at FROM rebalance_executions e JOIN rebalance_plans p ON p.id = e.plan_id WHERE p.fund_id = @fund_id ORDER BY e.executed_at DESC, e.run_date DESC, e.instrument_code """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(rebalanceExecutionRecordFromReader reader) records |> Seq.toList member this.PreviewRebalancePlan(planId: Guid) : Result = use connection = new NpgsqlConnection(connectionString) connection.Open() match findRebalancePlan connection None planId with | None -> Error "rebalance plan was not found" | Some plan when plan.Status = "archived" -> Error "rebalance plan is archived" | Some plan -> match this.GetFund plan.FundId with | None -> Error "fund was not found" | Some fund -> let positions = this.GetFundPositions plan.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 match RebalancePolicy.weightRows plan.Targets snapshots fund.AvailableCash with | Error message -> Error message | Ok rows -> let equity = fund.AvailableCash + (snapshots |> List.sumBy (fun snapshot -> snapshot.MarketValue)) Ok { PlanId = plan.Id RunDate = ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow AvailableCash = fund.AvailableCash Equity = equity Rows = rows } member this.RegisterDividend(idempotencyKey: string, fundId: Guid, command: DividendCommand) : DividendWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then DividendWriteResult.DividendInvalid "idempotency key cannot be empty" else let today = ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow match validateDividendCommand today command with | Error message -> DividendWriteResult.DividendInvalid message | Ok() -> let schemeKey = dividendSchemeKey fundId command let fingerprint = dividendRequestHash fundId command use connection = new NpgsqlConnection(connectionString) connection.Open() // Phase 1: claim the scheme atomically, validate, and book the payout let claimedStage = 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 schemeKey) |> ignore lockCommand.ExecuteNonQuery() |> ignore match findDividendIdempotency connection (Some transaction) schemeKey with | Some(existingHash, existingFundId, recordId) when existingFundId = fundId -> if existingHash = fingerprint then transaction.Commit() StageReplay recordId else transaction.Rollback() StageFail DividendWriteResult.DividendIdempotencyConflict | Some _ -> transaction.Rollback() StageFail DividendWriteResult.DividendIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() StageFail DividendWriteResult.DividendFundNotFound | Some isSynthetic -> if not (instrumentExists connection (Some transaction) command.InstrumentCode) then transaction.Rollback() StageFail DividendWriteResult.DividendInstrumentNotFound else use positionCommand = commandWithTransaction connection (Some transaction) "SELECT units 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 command.InstrumentCode) |> ignore use positionReader = positionCommand.ExecuteReader() let positionFound = positionReader.Read() let heldUnits = if positionFound then positionReader.GetDecimal(0) else 0m positionReader.Close() if not positionFound || heldUnits <= 0m then transaction.Rollback() StageFail DividendWriteResult.DividendNoHoldings else match DividendPolicy.computeCashPayout heldUnits command.Dps with | Error message -> transaction.Rollback() StageFail (DividendWriteResult.DividendInvalid message) | Ok payout -> use cashCommand = commandWithTransaction connection (Some transaction) "UPDATE funds SET available_cash = available_cash + @gross WHERE id = @fund_id" addParameter cashCommand "gross" NpgsqlDbType.Numeric (box payout.GrossCash) |> ignore addParameter cashCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore if cashCommand.ExecuteNonQuery() = 0 then transaction.Rollback() StageFail DividendWriteResult.DividendFundNotFound else let recordId = Guid.NewGuid() let record: DividendRecord = { Id = recordId FundId = fundId InstrumentCode = command.InstrumentCode NavDate = command.NavDate Dps = command.Dps Mode = command.Mode Status = (if command.Mode = DividendPolicy.Cash then "cash_credited" else "pending_nav") IsSynthetic = isSynthetic GrossCash = Some payout.GrossCash CreditedUnits = None CreditedInvested = None OrderId = None PendingReason = None CreatedAt = DateTimeOffset.UtcNow } let createdAt = insertDividendRecord connection (Some transaction) record insertDividendIdempotency connection (Some transaction) schemeKey fingerprint recordId fundId transaction.Commit() StageBooked { record with CreatedAt = createdAt } with error -> try transaction.Rollback() with _ -> () raise error // Confirmation of the reinvestment order settles the dividend record: the // credited units/invested cash are persisted and the row is marked succeeded. let finalizeDividend (record: DividendRecord) (confirmed: SubscriptionOrderRecord) = let units = confirmed.ConfirmedUnits |> Option.defaultValue 0m let invested = confirmed.ConfirmedInvestedCash |> Option.defaultValue 0m do use finalizeConnection = new NpgsqlConnection(connectionString) finalizeConnection.Open() use finalizeTransaction = finalizeConnection.BeginTransaction(IsolationLevel.ReadCommitted) use finalizeCommand = commandWithTransaction finalizeConnection (Some finalizeTransaction) """ UPDATE dividend_records SET status = 'succeeded', credited_units = @units, credited_invested = @invested, pending_reason = NULL WHERE id = @record_id """ addParameter finalizeCommand "units" NpgsqlDbType.Numeric (box units) |> ignore addParameter finalizeCommand "invested" NpgsqlDbType.Numeric (box invested) |> ignore addParameter finalizeCommand "record_id" NpgsqlDbType.Uuid (box record.Id) |> ignore finalizeCommand.ExecuteNonQuery() |> ignore finalizeTransaction.Commit() { record with Status = "succeeded"; CreditedUnits = Some units; CreditedInvested = Some invested } // Phase 2: replays never re-book; pending reinvestments retry their confirmation match claimedStage with | StageReplay recordId -> let record = match findDividendRecord connection None recordId with | Some record -> record | None -> failwith "idempotency scheme references a missing dividend" match record.Status, record.OrderId with | "pending_nav", Some orderId -> let confirmKey = sprintf "dividend-confirm:%O:%s:%s:%s" fundId command.InstrumentCode (command.NavDate.ToString("yyyy-MM-dd")) (command.Dps.ToString("G29", CultureInfo.InvariantCulture)) match this.ConfirmSubscriptionOrder(confirmKey, fundId, orderId) with | SubscriptionConfirmResult.OrderConfirmed confirmed | SubscriptionConfirmResult.ConfirmReplayed confirmed -> DividendCredited(finalizeDividend record confirmed) | SubscriptionConfirmResult.ConfirmPendingNav _ -> DividendReplayed record | other -> failwithf "unexpected dividend confirm replay result: %A" other | _ -> DividendReplayed record | StageFail failure -> failure | StageBooked record -> if record.Mode = DividendPolicy.Cash then DividendCredited record else let orderKey = sprintf "dividend:%O:%s:%s:%s" fundId command.InstrumentCode (command.NavDate.ToString("yyyy-MM-dd")) (command.Dps.ToString("G29", CultureInfo.InvariantCulture)) let confirmKey = sprintf "dividend-confirm:%O:%s:%s:%s" fundId command.InstrumentCode (command.NavDate.ToString("yyyy-MM-dd")) (command.Dps.ToString("G29", CultureInfo.InvariantCulture)) match this.CreateSubscriptionOrder( orderKey, fundId, { FundCode = command.InstrumentCode; Amount = record.GrossCash |> Option.defaultValue 0m; FeeAmount = 0m }, command.NavDate ) with | SubscriptionOrderWriteResult.OrderCreated order | SubscriptionOrderWriteResult.OrderReplayed order -> do use bindConnection = new NpgsqlConnection(connectionString) bindConnection.Open() use bindTransaction = bindConnection.BeginTransaction(IsolationLevel.ReadCommitted) use bindCommand = commandWithTransaction bindConnection (Some bindTransaction) "UPDATE dividend_records SET order_id = @order_id WHERE id = @record_id" addParameter bindCommand "order_id" NpgsqlDbType.Uuid (box order.Id) |> ignore addParameter bindCommand "record_id" NpgsqlDbType.Uuid (box record.Id) |> ignore bindCommand.ExecuteNonQuery() |> ignore bindTransaction.Commit() let recordWithOrder = { record with OrderId = Some order.Id } match this.ConfirmSubscriptionOrder(confirmKey, fundId, order.Id) with | SubscriptionConfirmResult.OrderConfirmed confirmed | SubscriptionConfirmResult.ConfirmReplayed confirmed -> DividendCredited(finalizeDividend recordWithOrder confirmed) | SubscriptionConfirmResult.ConfirmPendingNav pending -> DividendPendingReinvest { recordWithOrder with PendingReason = pending.PendingReason } | other -> failwithf "unexpected dividend confirm result: %A" other | other -> failwithf "unexpected dividend order result: %A" other member _.GetDividendRecords(fundId: Guid) = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT id, fund_id, instrument_code, nav_date, dps, mode, status, is_synthetic, gross_cash, credited_units, credited_invested, order_id, pending_reason, created_at FROM dividend_records 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(dividendRecordFromReader reader) records |> 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 _.CreateStockTrade(idempotencyKey: string, fundId: Guid, command: StockTradeCommand, ?executedAtOverride: DateTimeOffset) : StockTradeWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then StockTradeWriteResult.StockTradeInvalid "idempotency key cannot be empty" else let code = if isNull command.InstrumentCode then "" else command.InstrumentCode.Trim() if code.Length <> 6 || not (code |> Seq.forall Char.IsDigit) then StockTradeWriteResult.StockTradeInvalid "stock code must contain exactly six digits" elif command.Quantity <= 0m then StockTradeWriteResult.StockTradeInvalid "quantity must be positive" elif command.Price <= 0m then StockTradeWriteResult.StockTradeInvalid "price must be positive" else let normalized = { command with InstrumentCode = code } let fingerprint = stockTradeRequestHash fundId normalized 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 findStockTradeIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, tradeId) when existingHash = fingerprint && existingFundId = fundId -> match findStockTrade connection (Some transaction) tradeId with | Some trade -> transaction.Commit() StockTradeWriteResult.StockTradeReplayed trade | None -> transaction.Rollback() StockTradeWriteResult.StockTradeInvalid "idempotency record references a missing trade" | Some _ -> transaction.Rollback() StockTradeWriteResult.StockTradeIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() StockTradeWriteResult.StockTradeFundNotFound | Some isSynthetic -> let executedAt = defaultArg executedAtOverride DateTimeOffset.UtcNow let costCash = Decimal.Round(normalized.Quantity * normalized.Price, 2, MidpointRounding.AwayFromZero) let trade: StockTradeRecord = { Id = Guid.NewGuid() FundId = fundId InstrumentCode = normalized.InstrumentCode StockName = normalized.StockName Quantity = normalized.Quantity Price = normalized.Price CostCash = costCash IsSynthetic = isSynthetic ExecutedAt = executedAt } insertStockTrade connection (Some transaction) trade insertStockTradeIdempotency connection (Some transaction) idempotencyKey fingerprint trade.Id fundId use positionCommand = commandWithTransaction connection (Some transaction) """ INSERT INTO stock_positions (fund_id, instrument_code, stock_name, quantity, cost_cash, last_traded_at) VALUES (@fund_id, @code, @name, @quantity, @cost_cash, @last_traded_at) ON CONFLICT (fund_id, instrument_code) DO UPDATE SET quantity = stock_positions.quantity + EXCLUDED.quantity, cost_cash = stock_positions.cost_cash + EXCLUDED.cost_cash, stock_name = COALESCE(EXCLUDED.stock_name, stock_positions.stock_name), last_traded_at = EXCLUDED.last_traded_at """ addParameter positionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter positionCommand "code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore let nameParameter = match normalized.StockName with | Some name -> box name | None -> box DBNull.Value addParameter positionCommand "name" NpgsqlDbType.Text nameParameter |> ignore addParameter positionCommand "quantity" NpgsqlDbType.Numeric (box normalized.Quantity) |> ignore addParameter positionCommand "cost_cash" NpgsqlDbType.Numeric (box costCash) |> ignore addParameter positionCommand "last_traded_at" NpgsqlDbType.TimestampTz (box executedAt) |> ignore positionCommand.ExecuteNonQuery() |> ignore let persisted = { trade with ExecutedAt = executedAt } transaction.Commit() StockTradeWriteResult.StockTradeCreated persisted with error -> try transaction.Rollback() with _ -> () raise error member _.GetStockTrades(fundId: Guid) : StockTradeRecord list = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT id, fund_id, instrument_code, stock_name, quantity, price, cost_cash, is_synthetic, executed_at FROM stock_trades WHERE fund_id = @fund_id ORDER BY executed_at, id """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(stockTradeRecordFromReader reader) records |> Seq.toList member _.GetStockPositions(fundId: Guid) : StockPositionRecord list = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT instrument_code, stock_name, quantity, cost_cash, last_traded_at FROM stock_positions WHERE fund_id = @fund_id ORDER BY instrument_code """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add( { FundId = fundId InstrumentCode = reader.GetString(0) StockName = if reader.IsDBNull(1) then None else Some(reader.GetString(1)) Quantity = reader.GetDecimal(2) CostCash = reader.GetDecimal(3) LastTradedAt = reader.GetFieldValue(4) } ) records |> Seq.toList member _.GetStockSells(fundId: Guid) : StockSellRecord list = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT id, fund_id, instrument_code, stock_name, quantity, price, fee_amount, proceeds, is_synthetic, executed_at FROM stock_sells WHERE fund_id = @fund_id ORDER BY executed_at, id """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(stockSellRecordFromReader reader) records |> Seq.toList member _.CreateStockSell(idempotencyKey: string, fundId: Guid, command: StockSellCommand, ?executedAtOverride: DateTimeOffset) : StockSellWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then StockSellWriteResult.StockSellInvalid "idempotency key cannot be empty" else let code = if isNull command.InstrumentCode then "" else command.InstrumentCode.Trim() if code.Length <> 6 || not (code |> Seq.forall Char.IsDigit) then StockSellWriteResult.StockSellInvalid "stock code must contain exactly six digits" elif command.Quantity <= 0m then StockSellWriteResult.StockSellInvalid "quantity must be positive" elif command.Price <= 0m then StockSellWriteResult.StockSellInvalid "price must be positive" elif command.FeeAmount < 0m then StockSellWriteResult.StockSellInvalid "fee amount cannot be negative" elif Decimal.Round(command.FeeAmount, 2) <> command.FeeAmount then StockSellWriteResult.StockSellInvalid "fee amount exceeds cash precision" else let normalized = { command with InstrumentCode = code } let gross = Decimal.Round(normalized.Quantity * normalized.Price, 2, MidpointRounding.AwayFromZero) if normalized.FeeAmount > gross then StockSellWriteResult.StockSellInvalid "fee amount cannot exceed the sale proceeds" else let fingerprint = stockSellRequestHash fundId normalized 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 findStockSellIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, sellId) when existingHash = fingerprint && existingFundId = fundId -> match findStockSell connection (Some transaction) sellId with | Some sell -> transaction.Commit() StockSellWriteResult.StockSellReplayed sell | None -> transaction.Rollback() StockSellWriteResult.StockSellInvalid "idempotency record references a missing sale" | Some _ -> transaction.Rollback() StockSellWriteResult.StockSellIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() StockSellWriteResult.StockSellFundNotFound | Some isSynthetic -> let position = use positionCommand = commandWithTransaction connection (Some transaction) """ SELECT quantity, cost_cash FROM stock_positions WHERE fund_id = @fund_id AND instrument_code = @instrument_code FOR UPDATE """ addParameter positionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter positionCommand "instrument_code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore use reader = positionCommand.ExecuteReader() if reader.Read() then Some(reader.GetDecimal(0), reader.GetDecimal(1)) else None match position with | None -> transaction.Rollback() StockSellWriteResult.StockSellInsufficientHoldings(sprintf "no stock position in %s to sell" normalized.InstrumentCode) | Some(heldQuantity, _) when heldQuantity < normalized.Quantity -> transaction.Rollback() StockSellWriteResult.StockSellInsufficientHoldings( sprintf "available holdings %s are not enough for the requested sale quantity %s" (heldQuantity.ToString("G29", CultureInfo.InvariantCulture)) (normalized.Quantity.ToString("G29", CultureInfo.InvariantCulture)) ) | Some(heldQuantity, heldCost) -> let executedAt = defaultArg executedAtOverride DateTimeOffset.UtcNow let proceeds = gross - normalized.FeeAmount let remainingQuantity = heldQuantity - normalized.Quantity let releasedCost = if remainingQuantity <= 0m then heldCost else Decimal.Round(heldCost * (normalized.Quantity / heldQuantity), 2, MidpointRounding.AwayFromZero) if remainingQuantity <= 0m then use deleteCommand = commandWithTransaction connection (Some transaction) "DELETE FROM stock_positions WHERE fund_id = @fund_id AND instrument_code = @instrument_code" addParameter deleteCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter deleteCommand "instrument_code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore deleteCommand.ExecuteNonQuery() |> ignore else use updateCommand = commandWithTransaction connection (Some transaction) """ UPDATE stock_positions SET quantity = @quantity, cost_cash = @cost_cash, last_traded_at = @last_traded_at WHERE fund_id = @fund_id AND instrument_code = @instrument_code """ addParameter updateCommand "quantity" NpgsqlDbType.Numeric (box remainingQuantity) |> ignore addParameter updateCommand "cost_cash" NpgsqlDbType.Numeric (box (max 0m (heldCost - releasedCost))) |> ignore addParameter updateCommand "last_traded_at" NpgsqlDbType.TimestampTz (box executedAt) |> ignore addParameter updateCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter updateCommand "instrument_code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore updateCommand.ExecuteNonQuery() |> ignore use cashCommand = commandWithTransaction connection (Some transaction) "UPDATE funds SET available_cash = available_cash + @proceeds WHERE id = @fund_id" addParameter cashCommand "proceeds" NpgsqlDbType.Numeric (box proceeds) |> ignore addParameter cashCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore cashCommand.ExecuteNonQuery() |> ignore let sell: StockSellRecord = { Id = Guid.NewGuid() FundId = fundId InstrumentCode = normalized.InstrumentCode StockName = normalized.StockName Quantity = normalized.Quantity Price = normalized.Price FeeAmount = normalized.FeeAmount Proceeds = proceeds IsSynthetic = isSynthetic ExecutedAt = executedAt } insertStockSell connection (Some transaction) sell insertStockSellIdempotency connection (Some transaction) idempotencyKey fingerprint sell.Id fundId transaction.Commit() StockSellWriteResult.StockSellCreated sell with error -> try transaction.Rollback() with _ -> () raise error member _.CreateBondTrade(idempotencyKey: string, fundId: Guid, command: BondTradeCommand, ?executedAtOverride: DateTimeOffset) : BondTradeWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then BondTradeWriteResult.BondTradeInvalid "idempotency key cannot be empty" else let code = if isNull command.InstrumentCode then "" else command.InstrumentCode.Trim() if code.Length <> 6 || not (code |> Seq.forall Char.IsDigit) then BondTradeWriteResult.BondTradeInvalid "bond code must contain exactly six digits" elif command.Quantity <= 0m then BondTradeWriteResult.BondTradeInvalid "quantity must be positive" elif command.Price <= 0m then BondTradeWriteResult.BondTradeInvalid "price must be positive" elif command.CleanPrice <= 0m then BondTradeWriteResult.BondTradeInvalid "clean price must be positive" elif command.ParValue <= 0m then BondTradeWriteResult.BondTradeInvalid "par value must be positive" elif command.AccruedInterest < 0m then BondTradeWriteResult.BondTradeInvalid "accrued interest cannot be negative" else let normalized = { command with InstrumentCode = code } let fingerprint = bondTradeRequestHash fundId normalized 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 findBondTradeIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, tradeId) when existingHash = fingerprint && existingFundId = fundId -> match findBondTrade connection (Some transaction) tradeId with | Some trade -> transaction.Commit() BondTradeWriteResult.BondTradeReplayed trade | None -> transaction.Rollback() BondTradeWriteResult.BondTradeInvalid "idempotency record references a missing trade" | Some _ -> transaction.Rollback() BondTradeWriteResult.BondTradeIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() BondTradeWriteResult.BondTradeFundNotFound | Some isSynthetic -> let executedAt = defaultArg executedAtOverride DateTimeOffset.UtcNow let costCash = Decimal.Round( normalized.Quantity * normalized.Price * normalized.ParValue / 100m, 2, MidpointRounding.AwayFromZero ) let trade: BondTradeRecord = { Id = Guid.NewGuid() FundId = fundId InstrumentCode = normalized.InstrumentCode BondName = normalized.BondName Quantity = normalized.Quantity Price = normalized.Price CleanPrice = normalized.CleanPrice AccruedInterest = normalized.AccruedInterest ParValue = normalized.ParValue SettlementDate = normalized.SettlementDate CouponRate = normalized.CouponRate ValueDate = normalized.ValueDate MaturityDate = normalized.MaturityDate TradeDate = normalized.TradeDate |> Option.defaultValue normalized.SettlementDate CostCash = costCash IsSynthetic = isSynthetic ExecutedAt = executedAt } insertBondTrade connection (Some transaction) trade insertBondTradeIdempotency connection (Some transaction) idempotencyKey fingerprint trade.Id fundId use positionCommand = commandWithTransaction connection (Some transaction) """ INSERT INTO bond_positions (fund_id, instrument_code, bond_name, quantity, cost_cash, last_traded_at) VALUES (@fund_id, @code, @name, @quantity, @cost_cash, @last_traded_at) ON CONFLICT (fund_id, instrument_code) DO UPDATE SET quantity = bond_positions.quantity + EXCLUDED.quantity, cost_cash = bond_positions.cost_cash + EXCLUDED.cost_cash, bond_name = COALESCE(EXCLUDED.bond_name, bond_positions.bond_name), last_traded_at = EXCLUDED.last_traded_at """ addParameter positionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter positionCommand "code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore let nameParameter = match normalized.BondName with | Some name -> box name | None -> box DBNull.Value addParameter positionCommand "name" NpgsqlDbType.Text nameParameter |> ignore addParameter positionCommand "quantity" NpgsqlDbType.Numeric (box normalized.Quantity) |> ignore addParameter positionCommand "cost_cash" NpgsqlDbType.Numeric (box costCash) |> ignore addParameter positionCommand "last_traded_at" NpgsqlDbType.TimestampTz (box executedAt) |> ignore positionCommand.ExecuteNonQuery() |> ignore transaction.Commit() BondTradeWriteResult.BondTradeCreated trade with error -> try transaction.Rollback() with _ -> () raise error member _.GetBondTrades(fundId: Guid) : BondTradeRecord list = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None $""" SELECT {bondTradeColumns} FROM bond_trades WHERE fund_id = @fund_id ORDER BY executed_at, id """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(bondTradeRecordFromReader reader) records |> Seq.toList member _.RecordBondCashflow(idempotencyKey: string, fundId: Guid, command: BondCashflowCommand) : BondCashflowWriteResult = let code = if isNull command.InstrumentCode then "" else command.InstrumentCode.Trim() let eventType = if isNull command.EventType then "" else command.EventType.Trim().ToLowerInvariant() if String.IsNullOrWhiteSpace idempotencyKey then BondCashflowWriteResult.BondCashflowInvalid "idempotency key cannot be empty" elif code.Length <> 6 || not (code |> Seq.forall Char.IsDigit) then BondCashflowWriteResult.BondCashflowInvalid "bond code must contain exactly six digits" elif eventType <> "coupon" && eventType <> "maturity" && eventType <> "redemption" then BondCashflowWriteResult.BondCashflowInvalid "event type must be coupon, maturity or redemption" elif command.Quantity <= 0m then BondCashflowWriteResult.BondCashflowInvalid "quantity must be positive" elif command.Amount < 0m then BondCashflowWriteResult.BondCashflowInvalid "amount cannot be negative" else let normalized = { command with InstrumentCode = code; EventType = eventType } let fingerprint = bondCashflowRequestHash fundId normalized 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 findBondCashflowIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, eventId) when existingHash = fingerprint && existingFundId = fundId -> match findBondCashflow connection (Some transaction) eventId with | Some record -> transaction.Commit() BondCashflowWriteResult.BondCashflowReplayed record | None -> transaction.Rollback() BondCashflowWriteResult.BondCashflowInvalid "idempotency record references a missing event" | Some _ -> transaction.Rollback() BondCashflowWriteResult.BondCashflowIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() BondCashflowWriteResult.BondCashflowFundNotFound | Some isSynthetic -> let hasPosition = use positionQuery = commandWithTransaction connection (Some transaction) "SELECT 1 FROM bond_positions WHERE fund_id = @fund_id AND instrument_code = @code" addParameter positionQuery "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter positionQuery "code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore use reader = positionQuery.ExecuteReader() reader.Read() if not hasPosition then transaction.Rollback() BondCashflowWriteResult.BondCashflowPositionNotFound else let record: BondCashflowRecord = { Id = Guid.NewGuid() FundId = fundId InstrumentCode = normalized.InstrumentCode BondName = normalized.BondName EventType = normalized.EventType EventDate = normalized.EventDate Quantity = normalized.Quantity Amount = normalized.Amount Note = normalized.Note IsSynthetic = isSynthetic CreatedAt = DateTimeOffset.UtcNow } insertBondCashflow connection (Some transaction) record insertBondCashflowIdempotency connection (Some transaction) idempotencyKey fingerprint record.Id fundId use cashCommand = commandWithTransaction connection (Some transaction) "UPDATE funds SET available_cash = available_cash + @amount WHERE id = @fund_id" addParameter cashCommand "amount" NpgsqlDbType.Numeric (box normalized.Amount) |> ignore addParameter cashCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore cashCommand.ExecuteNonQuery() |> ignore if normalized.EventType = "maturity" || normalized.EventType = "redemption" then use removeCommand = commandWithTransaction connection (Some transaction) "DELETE FROM bond_positions WHERE fund_id = @fund_id AND instrument_code = @code" addParameter removeCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter removeCommand "code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore removeCommand.ExecuteNonQuery() |> ignore transaction.Commit() BondCashflowWriteResult.BondCashflowCreated record with error -> try transaction.Rollback() with _ -> () raise error member _.GetBondCashflows(fundId: Guid) : BondCashflowRecord list = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None $""" SELECT {bondCashflowColumns} FROM bond_cashflow_events WHERE fund_id = @fund_id ORDER BY event_date, created_at, id """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(bondCashflowRecordFromReader reader) records |> Seq.toList member _.CreateBondSell(idempotencyKey: string, fundId: Guid, command: BondSellCommand, ?executedAtOverride: DateTimeOffset) : BondSellWriteResult = let code = if isNull command.InstrumentCode then "" else command.InstrumentCode.Trim() if String.IsNullOrWhiteSpace idempotencyKey then BondSellWriteResult.BondSellInvalid "idempotency key cannot be empty" elif code.Length <> 6 || not (code |> Seq.forall Char.IsDigit) then BondSellWriteResult.BondSellInvalid "bond code must contain exactly six digits" elif command.Quantity <= 0m then BondSellWriteResult.BondSellInvalid "quantity must be positive" elif command.Price <= 0m then BondSellWriteResult.BondSellInvalid "price must be positive" elif command.CleanPrice <= 0m then BondSellWriteResult.BondSellInvalid "clean price must be positive" elif command.ParValue <= 0m then BondSellWriteResult.BondSellInvalid "par value must be positive" elif command.AccruedInterest < 0m then BondSellWriteResult.BondSellInvalid "accrued interest cannot be negative" elif command.FeeAmount < 0m then BondSellWriteResult.BondSellInvalid "fee cannot be negative" else let normalized = { command with InstrumentCode = code } let fingerprint = bondSellRequestHash fundId normalized 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 findBondSellIdempotency connection (Some transaction) idempotencyKey with | Some(existingHash, existingFundId, sellId) when existingHash = fingerprint && existingFundId = fundId -> match findBondSell connection (Some transaction) sellId with | Some record -> transaction.Commit() BondSellWriteResult.BondSellReplayed record | None -> transaction.Rollback() BondSellWriteResult.BondSellInvalid "idempotency record references a missing sell" | Some _ -> transaction.Rollback() BondSellWriteResult.BondSellIdempotencyConflict | None -> match lockFundForOrder connection (Some transaction) fundId with | None -> transaction.Rollback() BondSellWriteResult.BondSellFundNotFound | Some isSynthetic -> let position = use positionQuery = commandWithTransaction connection (Some transaction) "SELECT quantity, cost_cash FROM bond_positions WHERE fund_id = @fund_id AND instrument_code = @code FOR UPDATE" addParameter positionQuery "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter positionQuery "code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore use reader = positionQuery.ExecuteReader() if reader.Read() then Some(reader.GetDecimal(0), reader.GetDecimal(1)) else None match position with | None -> transaction.Rollback() BondSellWriteResult.BondSellInsufficientHoldings "fund does not hold this bond" | Some(positionQuantity, positionCost) -> if normalized.Quantity > positionQuantity then transaction.Rollback() BondSellWriteResult.BondSellInsufficientHoldings( sprintf "cannot sell %O 张: only %O held" normalized.Quantity positionQuantity ) else let executedAt = defaultArg executedAtOverride DateTimeOffset.UtcNow let gross = Decimal.Round( normalized.Quantity * normalized.Price * normalized.ParValue / 100m, 2, MidpointRounding.AwayFromZero ) let proceeds = Decimal.Round(gross - normalized.FeeAmount, 2, MidpointRounding.AwayFromZero) if proceeds < 0m then transaction.Rollback() BondSellWriteResult.BondSellInvalid "fee exceeds gross proceeds" else let costReleased = if normalized.Quantity = positionQuantity then positionCost else Decimal.Round( positionCost * normalized.Quantity / positionQuantity, 2, MidpointRounding.AwayFromZero ) let realizedPnl = Decimal.Round(proceeds - costReleased, 2, MidpointRounding.AwayFromZero) let record: BondSellRecord = { Id = Guid.NewGuid() FundId = fundId InstrumentCode = normalized.InstrumentCode BondName = normalized.BondName Quantity = normalized.Quantity Price = normalized.Price CleanPrice = normalized.CleanPrice AccruedInterest = normalized.AccruedInterest ParValue = normalized.ParValue SettlementDate = normalized.SettlementDate TradeDate = normalized.TradeDate |> Option.defaultValue normalized.SettlementDate FeeAmount = normalized.FeeAmount Proceeds = proceeds CostReleased = costReleased RealizedPnl = realizedPnl IsSynthetic = isSynthetic ExecutedAt = executedAt } insertBondSell connection (Some transaction) record insertBondSellIdempotency connection (Some transaction) idempotencyKey fingerprint record.Id fundId use cashCommand = commandWithTransaction connection (Some transaction) "UPDATE funds SET available_cash = available_cash + @proceeds WHERE id = @fund_id" addParameter cashCommand "proceeds" NpgsqlDbType.Numeric (box proceeds) |> ignore addParameter cashCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore cashCommand.ExecuteNonQuery() |> ignore if normalized.Quantity = positionQuantity then use removeCommand = commandWithTransaction connection (Some transaction) "DELETE FROM bond_positions WHERE fund_id = @fund_id AND instrument_code = @code" addParameter removeCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter removeCommand "code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore removeCommand.ExecuteNonQuery() |> ignore else use positionCommand = commandWithTransaction connection (Some transaction) """ UPDATE bond_positions SET quantity = quantity - @quantity, cost_cash = cost_cash - @cost_released, last_traded_at = @last_traded_at WHERE fund_id = @fund_id AND instrument_code = @code """ addParameter positionCommand "quantity" NpgsqlDbType.Numeric (box normalized.Quantity) |> ignore addParameter positionCommand "cost_released" NpgsqlDbType.Numeric (box costReleased) |> ignore addParameter positionCommand "last_traded_at" NpgsqlDbType.TimestampTz (box executedAt) |> ignore addParameter positionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore addParameter positionCommand "code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore positionCommand.ExecuteNonQuery() |> ignore transaction.Commit() BondSellWriteResult.BondSellCreated record with error -> try transaction.Rollback() with _ -> () raise error member _.GetBondSells(fundId: Guid) : BondSellRecord list = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None $""" SELECT {bondSellColumns} FROM bond_sells WHERE fund_id = @fund_id ORDER BY executed_at, id """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(bondSellRecordFromReader reader) records |> Seq.toList member _.GetBondPositions(fundId: Guid) : BondPositionRecord list = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT instrument_code, bond_name, quantity, cost_cash, last_traded_at FROM bond_positions WHERE fund_id = @fund_id ORDER BY instrument_code """ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add( { FundId = fundId InstrumentCode = reader.GetString(0) BondName = if reader.IsDBNull(1) then None else Some(reader.GetString(1)) Quantity = reader.GetDecimal(2) CostCash = reader.GetDecimal(3) LastTradedAt = reader.GetFieldValue(4) } ) records |> Seq.toList 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 /// Reconstructs the fund's asset and unit-NAV history from confirmed ledger /// events plus persisted NAV observations. Unit NAV is rebuilt by issuing fund /// units for external capital deposits at the previously known unit NAV, so a /// deposit never moves the unit NAV. Values that cannot be known at a date stay /// `None`; nothing is filled with zero. member this.GetFundReturns(fundId: Guid) : FundReturns option = match this.GetFund fundId with | None -> None | Some fund -> let today = ConfirmationPolicy.eventDateFor DateTimeOffset.UtcNow let deposits = this.GetCapitalDeposits fundId let orders = this.GetSubscriptionOrders fundId let redemptions = this.GetRedemptionOrders fundId let dividends = this.GetDividendRecords fundId let heldCodes = [ yield! orders |> List.map (fun order -> order.FundCode) yield! redemptions |> List.map (fun order -> order.InstrumentCode) ] |> List.distinct let navByCode = heldCodes |> List.map (fun code -> code, this.GetNav(code, None, Some today)) |> Map.ofList let firstHoldingDate (code: string) = orders |> List.filter (fun order -> order.FundCode = code && order.Status = "confirmed") |> List.map (fun order -> order.TradeDate) |> List.sort |> List.tryHead let navDates = heldCodes |> List.collect (fun code -> match firstHoldingDate code with | None -> [] | Some startDate -> navByCode.[code] |> List.filter (fun observation -> observation.NavDate >= startDate && observation.NavDate <= today) |> List.map (fun observation -> observation.NavDate)) |> Set.ofList let creationDate = ConfirmationPolicy.eventDateFor fund.CreatedAt let dividendCredited = dividends |> List.filter (fun record -> record.Status = "cash_credited" || record.Status = "succeeded") let hasActivity = not (List.isEmpty deposits) || not (List.isEmpty orders) || not (List.isEmpty redemptions) || not (List.isEmpty dividendCredited) || not (Set.isEmpty navDates) if not hasActivity then Some { FundId = fundId Pending = false DataUpdatedAt = None Points = [] } else let eventDates = [ yield creationDate yield! deposits |> List.map (fun deposit -> ConfirmationPolicy.eventDateFor deposit.CreatedAt) yield! orders |> List.map (fun order -> order.TradeDate) yield! redemptions |> List.map (fun order -> order.TradeDate) yield! dividendCredited |> List.map (fun record -> record.NavDate) ] |> Set.ofList let dates = Set.union eventDates navDates |> Set.toList |> List.sort let dataUpdatedAt = navByCode |> Map.toList |> List.collect snd |> List.map (fun observation -> observation.LastSeenAt) |> List.sortDescending |> List.tryHead let latestNavOnOrBefore (code: string) (date: DateOnly) = match Map.tryFind code navByCode with | None -> None | Some observations -> observations |> List.filter (fun observation -> observation.NavDate <= date && observation.Nav > 0m) |> List.sortByDescending (fun observation -> observation.NavDate, observation.SourceCollectedAt) |> List.tryHead let mutable availableCash = fund.InitialCash let mutable reservedCash = 0m let mutable fundUnits = if fund.InitialUnitNav > 0m then fund.InitialCash / fund.InitialUnitNav else 0m let mutable lastKnownNav = if fund.InitialUnitNav > 0m then Some fund.InitialUnitNav else None let positions = System.Collections.Generic.Dictionary() let mutable cumulativeDeposits = 0m let points = ResizeArray() for date in dates do for deposit in deposits do if ConfirmationPolicy.eventDateFor deposit.CreatedAt = date then availableCash <- availableCash + deposit.Amount cumulativeDeposits <- cumulativeDeposits + deposit.Amount match lastKnownNav with | Some nav when nav > 0m -> fundUnits <- fundUnits + deposit.Amount / nav | _ -> () for order in orders do if order.TradeDate = date then if order.Status = "confirmed" then let residual = order.ConfirmedResidualCash |> Option.defaultValue 0m availableCash <- availableCash + residual - order.ReservedTotal let units = order.ConfirmedUnits |> Option.defaultValue 0m let current = match positions.TryGetValue order.FundCode with | true, value -> value | _ -> 0m positions.[order.FundCode] <- current + units else availableCash <- availableCash - order.ReservedTotal reservedCash <- reservedCash + order.ReservedTotal for order in redemptions do if order.TradeDate = date && order.Status = "confirmed" then let current = match positions.TryGetValue order.InstrumentCode with | true, value -> value | _ -> 0m positions.[order.InstrumentCode] <- current - order.Units availableCash <- availableCash + (order.ConfirmedProceeds |> Option.defaultValue 0m) for record in dividendCredited do if record.NavDate = date then availableCash <- availableCash + (record.GrossCash |> Option.defaultValue 0m) let mutable pending = false let mutable holdingsValue = 0m for KeyValue(code, units) in positions do if units > 0m then match latestNavOnOrBefore code date with | Some observation -> holdingsValue <- holdingsValue + units * observation.Nav | None -> pending <- true let netExternalFlow = fund.InitialCash + cumulativeDeposits if pending then points.Add( { Date = date Pending = true TotalAssets = None UnitNav = None Cash = availableCash ReservedCash = reservedCash HoldingsValue = None CumulativeReturn = None NetExternalFlow = netExternalFlow } ) else let totalAssets = availableCash + reservedCash + holdingsValue let unitNav = if fundUnits > 0m then Some(totalAssets / fundUnits) else None match unitNav with | Some nav -> lastKnownNav <- Some nav | None -> () points.Add( { Date = date Pending = false TotalAssets = Some totalAssets UnitNav = unitNav Cash = availableCash ReservedCash = reservedCash HoldingsValue = Some holdingsValue CumulativeReturn = Some(totalAssets - netExternalFlow) NetExternalFlow = netExternalFlow } ) Some { FundId = fundId Pending = points |> Seq.exists (fun point -> point.Pending) DataUpdatedAt = dataUpdatedAt Points = points |> Seq.toList }