From a3dddc28f33254317d07323037ba5d1db35df3d6 Mon Sep 17 00:00:00 2001 From: "Somhairle H. Marisol" Date: Mon, 21 Sep 2026 13:28:02 +0800 Subject: Add rebalance execution through shared order pipeline (3d-7) --- src/FundLab.Api/Persistence.fs | 447 ++++++++++++++++++++++++++++++++++++++++- 1 file changed, 446 insertions(+), 1 deletion(-) (limited to 'src/FundLab.Api/Persistence.fs') diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs index f714267..d64fe98 100644 --- a/src/FundLab.Api/Persistence.fs +++ b/src/FundLab.Api/Persistence.fs @@ -311,6 +311,48 @@ type SipAdvanceResult = Plans: SipPlanAdvanceResult list } +type RebalanceTarget = RebalancePolicy.TargetAllocation + +type RebalancePlanCommand = + { + Targets: RebalanceTarget list + } + +type RebalancePlanRecord = + { + Id: Guid + FundId: Guid + Targets: RebalanceTarget list + Status: string + IsSynthetic: bool + CreatedAt: DateTimeOffset + } + +type RebalanceWriteResult = + | RebalancePlanCreated of RebalancePlanRecord + | RebalancePlanReplayed of RebalancePlanRecord + | RebalanceIdempotencyConflict + | RebalanceInvalid of string + | RebalanceFundNotFound + | RebalanceInstrumentNotFound + +type RebalanceOrderOutcome = + { + InstrumentCode: string + Action: string + Amount: string + Status: string + OrderId: Guid option + PendingReason: string option + } + +type RebalanceExecutionResult = + { + PlanId: Guid + RunDate: DateOnly + Outcomes: RebalanceOrderOutcome list + } + type CapitalDepositCommand = { Amount: decimal @@ -561,6 +603,21 @@ type FundRepository(connectionString: string) = executed_at timestamptz NOT NULL DEFAULT now(), PRIMARY KEY (plan_id, trade_date) ); + + CREATE TABLE IF NOT EXISTS rebalance_plans ( + id uuid PRIMARY KEY, + fund_id uuid NOT NULL REFERENCES funds(id), + status text NOT NULL, + is_synthetic boolean NOT NULL, + created_at timestamptz NOT NULL DEFAULT now() + ); + + CREATE TABLE IF NOT EXISTS rebalance_targets ( + plan_id uuid NOT NULL REFERENCES rebalance_plans(id), + instrument_code text NOT NULL REFERENCES instruments(code), + target_percent numeric(9, 2) NOT NULL CHECK (target_percent > 0 AND target_percent <= 100), + PRIMARY KEY (plan_id, instrument_code) + ); """ let statusText status = @@ -569,6 +626,10 @@ type FundRepository(connectionString: string) = | 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 @@ -1489,6 +1550,151 @@ type FundRepository(connectionString: string) = Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) + let rebalancePlanRecordFromReader (reader: DbDataReader) : RebalancePlanRecord = + { + Id = reader.GetGuid(0) + FundId = reader.GetGuid(1) + Targets = [] + Status = reader.GetString(2) + IsSynthetic = reader.GetBoolean(3) + CreatedAt = reader.GetFieldValue(4) + } + + let rebalanceTargetsFromReader (reader: DbDataReader) : RebalanceTarget list = + let targets = ResizeArray() + + while reader.Read() do + targets.Add( + { + InstrumentCode = reader.GetString(0) + TargetPercent = reader.GetDecimal(1) + } + ) + + targets |> Seq.toList + + let rebalancePlanWithTargets connection transaction (plan: RebalancePlanRecord) = + use targetsCommand = + commandWithTransaction + connection + transaction + "SELECT instrument_code, target_percent FROM rebalance_targets WHERE plan_id = @plan_id ORDER BY instrument_code" + + addParameter targetsCommand "plan_id" NpgsqlDbType.Uuid (box plan.Id) |> ignore + + use reader = targetsCommand.ExecuteReader() + let targets = rebalanceTargetsFromReader reader + { plan with Targets = targets } + + let findRebalancePlan connection transaction planId = + use command = + commandWithTransaction + connection + transaction + "SELECT id, fund_id, status, is_synthetic, created_at FROM rebalance_plans WHERE id = @plan_id" + + addParameter command "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore + + use reader = command.ExecuteReader() + + if reader.Read() then + let plan = rebalancePlanRecordFromReader reader + reader.Close() + Some(rebalancePlanWithTargets connection transaction plan) + else + None + + let findRebalanceIdempotency connection transaction key = + use command = + commandWithTransaction + connection + transaction + "SELECT request_hash, fund_id, plan_id FROM rebalance_idempotencies WHERE idempotency_key = @idempotency_key" + + addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore + + use reader = command.ExecuteReader() + if reader.Read() then + Some(reader.GetString(0), reader.GetGuid(1), reader.GetGuid(2)) + else + None + + let insertRebalancePlan connection transaction (plan: RebalancePlanRecord) = + use planCommand = + commandWithTransaction + connection + transaction + """ + INSERT INTO rebalance_plans (id, fund_id, status, is_synthetic) + VALUES (@id, @fund_id, @status, @is_synthetic) + RETURNING created_at + """ + + addParameter planCommand "id" NpgsqlDbType.Uuid (box plan.Id) |> ignore + addParameter planCommand "fund_id" NpgsqlDbType.Uuid (box plan.FundId) |> ignore + addParameter planCommand "status" NpgsqlDbType.Text (box plan.Status) |> ignore + addParameter planCommand "is_synthetic" NpgsqlDbType.Boolean (box plan.IsSynthetic) |> ignore + + use reader = planCommand.ExecuteReader() + reader.Read() |> ignore + let createdAt = reader.GetFieldValue(0) + reader.Close() + + for target in plan.Targets do + use targetCommand = + commandWithTransaction + connection + transaction + """ + INSERT INTO rebalance_targets (plan_id, instrument_code, target_percent) + VALUES (@plan_id, @instrument_code, @target_percent) + """ + + addParameter targetCommand "plan_id" NpgsqlDbType.Uuid (box plan.Id) |> ignore + addParameter targetCommand "instrument_code" NpgsqlDbType.Text (box target.InstrumentCode) |> ignore + addParameter targetCommand "target_percent" NpgsqlDbType.Numeric (box target.TargetPercent) |> ignore + targetCommand.ExecuteNonQuery() |> ignore + + createdAt + + let insertRebalanceIdempotency connection transaction key requestHash planId fundId = + use command = + commandWithTransaction + connection + transaction + """ + INSERT INTO rebalance_idempotencies (idempotency_key, request_hash, plan_id, fund_id) + VALUES (@idempotency_key, @request_hash, @plan_id, @fund_id) + """ + + addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore + addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore + addParameter command "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore + addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + command.ExecuteNonQuery() |> ignore + + let rebalanceRequestHash (fundId: Guid) (command: RebalancePlanCommand) = + let invariant = CultureInfo.InvariantCulture + let encoded (value: string) = sprintf "%d:%s" value.Length value + + let payload = + String.concat + "|" + ([ "rebalance-plan"; encoded (fundId.ToString("D")) ] + @ (command.Targets + |> List.sortBy (fun target -> target.InstrumentCode) + |> List.collect (fun target -> + [ + encoded target.InstrumentCode + encoded (target.TargetPercent.ToString("G29", invariant)) + ]))) + + Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) + + let validateRebalanceCommand (command: RebalancePlanCommand) = + match RebalancePolicy.validateTargets command.Targets with + | Error message -> Error message + | Ok() -> Ok() member _.EnsureSchema() = use connection = new NpgsqlConnection(connectionString) connection.Open() @@ -2190,6 +2396,245 @@ type FundRepository(connectionString: string) = records |> Seq.toList + member _.CreateRebalancePlan(idempotencyKey: string, fundId: Guid, command: RebalancePlanCommand) : RebalanceWriteResult = + if String.IsNullOrWhiteSpace idempotencyKey then + RebalanceWriteResult.RebalanceInvalid "idempotency key cannot be empty" + else + match validateRebalanceCommand command with + | Error message -> RebalanceWriteResult.RebalanceInvalid message + | Ok() -> + let fingerprint = rebalanceRequestHash fundId command + use connection = new NpgsqlConnection(connectionString) + connection.Open() + use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) + + try + use lockCommand = + commandWithTransaction + connection + (Some transaction) + "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" + + addParameter lockCommand "lock_key" NpgsqlDbType.Text (box idempotencyKey) |> ignore + lockCommand.ExecuteNonQuery() |> ignore + + match findRebalanceIdempotency connection (Some transaction) idempotencyKey with + | Some(existingHash, existingFundId, planId) + when existingHash = fingerprint && existingFundId = fundId -> + match findRebalancePlan connection (Some transaction) planId with + | Some plan -> + transaction.Commit() + RebalanceWriteResult.RebalancePlanReplayed plan + | None -> + transaction.Rollback() + RebalanceWriteResult.RebalanceInvalid "idempotency record references a missing plan" + | Some _ -> + transaction.Rollback() + RebalanceWriteResult.RebalanceIdempotencyConflict + | None -> + match lockFundForOrder connection (Some transaction) fundId with + | None -> + transaction.Rollback() + RebalanceWriteResult.RebalanceFundNotFound + | Some isSynthetic -> + let missingTarget = + command.Targets + |> List.tryFind (fun target -> not (instrumentExists connection (Some transaction) target.InstrumentCode)) + + match missingTarget with + | Some target -> + transaction.Rollback() + RebalanceWriteResult.RebalanceInstrumentNotFound + | None -> + let plan: RebalancePlanRecord = + { + Id = Guid.NewGuid() + FundId = fundId + Targets = command.Targets + Status = "active" + IsSynthetic = isSynthetic + CreatedAt = DateTimeOffset.UtcNow + } + + let createdAt = insertRebalancePlan connection (Some transaction) plan + insertRebalanceIdempotency connection (Some transaction) idempotencyKey fingerprint plan.Id fundId + transaction.Commit() + RebalanceWriteResult.RebalancePlanCreated { plan with CreatedAt = createdAt } + with error -> + try + transaction.Rollback() + with _ -> + () + + raise error + + member _.GetRebalancePlans(fundId: Guid) = + use connection = new NpgsqlConnection(connectionString) + connection.Open() + + use command = + commandWithTransaction + connection + None + "SELECT id, fund_id, status, is_synthetic, created_at FROM rebalance_plans WHERE fund_id = @fund_id ORDER BY created_at DESC, id" + + addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + + use reader = command.ExecuteReader() + let plans = ResizeArray() + + while reader.Read() do + plans.Add(rebalancePlanRecordFromReader reader) + + reader.Close() + plans |> Seq.toList |> List.map (rebalancePlanWithTargets connection None) + + member this.ExecuteRebalancePlan(planId: Guid) : Result = + use connection = new NpgsqlConnection(connectionString) + connection.Open() + + let plan = + use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) + + try + use lockCommand = + commandWithTransaction + connection + (Some transaction) + "SELECT pg_advisory_xact_lock(hashtext(@lock_key))" + + addParameter lockCommand "lock_key" NpgsqlDbType.Text (box (sprintf "rebalance-execute:%O" planId)) |> ignore + lockCommand.ExecuteNonQuery() |> ignore + + let plan = findRebalancePlan connection (Some transaction) planId + transaction.Commit() + plan + with error -> + try + transaction.Rollback() + with _ -> + () + + raise error + + match plan with + | None -> Error "rebalance plan was not found" + | Some plan -> + let fundId = plan.FundId + let runDate = ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow + + // snapshot: fund cash plus each holding valued at its latest valuation NAV + let fund = this.GetFund fundId + + match fund with + | None -> Error "fund was not found" + | Some fund -> + let positions = this.GetFundPositions fundId + + let buildSnapshot (position: FundPositionRecord) : RebalancePolicy.RebalancePositionSnapshot = + let marketValue = + match position.ValuationNav with + | Some nav -> Decimal.Round(position.Units * nav, 2) + | None -> 0m + + { + RebalancePolicy.RebalancePositionSnapshot.InstrumentCode = position.InstrumentCode + RebalancePolicy.RebalancePositionSnapshot.MarketValue = marketValue + RebalancePolicy.RebalancePositionSnapshot.Units = position.Units + RebalancePolicy.RebalancePositionSnapshot.AvailableUnits = position.Units - position.ReservedUnits + RebalancePolicy.RebalancePositionSnapshot.ValuationNav = position.ValuationNav + } + + let snapshots = positions |> List.map buildSnapshot + + let diffs = + RebalancePolicy.computeOrders plan.Targets snapshots fund.AvailableCash + + match diffs with + | Error message -> Error message + | Ok diffs -> + let outcomes = ResizeArray() + + // sells first so freed cash can fund the buys + let ordered = + diffs + |> List.sortBy (fun diff -> + match diff.Action with + | RebalancePolicy.Sell -> 0 + | RebalancePolicy.Buy -> 1 + | RebalancePolicy.Hold -> 2) + + for diff: RebalancePolicy.RebalanceDiff in ordered do + match diff.Action with + | RebalancePolicy.Hold -> () + | RebalancePolicy.Buy -> + let orderKey = RebalancePolicy.orderKey plan.Id runDate diff.InstrumentCode + + match + this.CreateSubscriptionOrder( + orderKey, + fundId, + { FundCode = diff.InstrumentCode; Amount = diff.Amount; FeeAmount = 0m } + ) + with + | SubscriptionOrderWriteResult.OrderCreated order + | SubscriptionOrderWriteResult.OrderReplayed order -> + let confirmKey = RebalancePolicy.confirmKey plan.Id runDate diff.InstrumentCode + + match this.ConfirmSubscriptionOrder(confirmKey, fundId, order.Id) with + | SubscriptionConfirmResult.OrderConfirmed _ + | SubscriptionConfirmResult.ConfirmReplayed _ -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "succeeded"; OrderId = Some order.Id; PendingReason = None }) + | SubscriptionConfirmResult.ConfirmPendingNav record -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "pending_nav"; OrderId = Some order.Id; PendingReason = record.PendingReason }) + | SubscriptionConfirmResult.ConfirmIdempotencyConflict -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "idempotency_conflict"; OrderId = Some order.Id; PendingReason = None }) + | other -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "failed"; OrderId = Some order.Id; PendingReason = Some (sprintf "%A" other) }) + | SubscriptionOrderWriteResult.OrderInsufficientFunds -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "insufficient_cash"; OrderId = None; PendingReason = Some "available cash is not enough for the rebalance buy" }) + | other -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "buy"; Amount = cashText diff.Amount; Status = "failed"; OrderId = None; PendingReason = Some (sprintf "%A" other) }) + | RebalancePolicy.Sell -> + match diff.Units with + | None -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "skipped_no_valuation"; OrderId = None; PendingReason = Some "holding has no valuation NAV to price the sell" }) + | Some units when units <= 0m -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "skipped_no_available_units"; OrderId = None; PendingReason = Some "no available units to redeem" }) + | Some units -> + let redeemKey = RebalancePolicy.redemptionKey plan.Id runDate diff.InstrumentCode + + match + this.CreateRedemptionOrder( + redeemKey, + fundId, + { InstrumentCode = diff.InstrumentCode; Units = units; FeeAmount = 0m } + ) + with + | RedemptionWriteResult.RedemptionCreated order + | RedemptionWriteResult.RedemptionReplayed order -> + let confirmKey = RebalancePolicy.redemptionConfirmKey plan.Id runDate diff.InstrumentCode + + match this.ConfirmRedemptionOrder(confirmKey, fundId, order.Id) with + | RedemptionConfirmResult.RedemptionConfirmed _ + | RedemptionConfirmResult.RedemptionConfirmReplayed _ -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "succeeded"; OrderId = Some order.Id; PendingReason = None }) + | RedemptionConfirmResult.RedemptionPendingNav record -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "pending_nav"; OrderId = Some order.Id; PendingReason = record.PendingReason }) + | other -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "failed"; OrderId = Some order.Id; PendingReason = Some (sprintf "%A" other) }) + | RedemptionWriteResult.RedemptionInsufficientUnits -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "insufficient_units"; OrderId = None; PendingReason = Some "available units are not enough for the rebalance sell" }) + | other -> + outcomes.Add({ InstrumentCode = diff.InstrumentCode; Action = "sell"; Amount = cashText diff.Amount; Status = "failed"; OrderId = None; PendingReason = Some (sprintf "%A" other) }) + + Ok + { + PlanId = plan.Id + RunDate = runDate + Outcomes = outcomes |> Seq.toList + } + member _.CreateCapitalDeposit(idempotencyKey: string, fundId: Guid, command: CapitalDepositCommand) : CapitalDepositWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then CapitalDepositWriteResult.CapitalDepositInvalid "idempotency key cannot be empty" @@ -3058,7 +3503,7 @@ type FundRepository(connectionString: string) = raise error - member _.GetFundPositions(fundId: Guid) = + member _.GetFundPositions(fundId: Guid) : FundPositionRecord list = use connection = new NpgsqlConnection(connectionString) connection.Open() -- cgit v1.2.3