diff options
| author | Somhairle H. Marisol <[email protected]> | 2026-09-21 13:28:02 +0800 |
|---|---|---|
| committer | Somhairle H. Marisol <[email protected]> | 2026-09-21 13:28:02 +0800 |
| commit | a3dddc28f33254317d07323037ba5d1db35df3d6 (patch) | |
| tree | c2a13404bda8a8b62d07559c2ab95b5491034f70 /src/FundLab.Api | |
| parent | d2a547b6a60ea769092894fff1f2b5105954f226 (diff) | |
| download | fund-lab-a3dddc28f33254317d07323037ba5d1db35df3d6.tar.gz | |
Add rebalance execution through shared order pipeline (3d-7)
Diffstat (limited to 'src/FundLab.Api')
| -rw-r--r-- | src/FundLab.Api/App.fs | 95 | ||||
| -rw-r--r-- | src/FundLab.Api/Persistence.fs | 447 |
2 files changed, 539 insertions, 3 deletions
diff --git a/src/FundLab.Api/App.fs b/src/FundLab.Api/App.fs index 9badcfa..3de514e 100644 --- a/src/FundLab.Api/App.fs +++ b/src/FundLab.Api/App.fs @@ -148,6 +148,22 @@ type SipPlanResponse = createdAt: string } +type RebalancePlanResponse = + { + id: Guid + fundId: Guid + targets: {| instrumentCode: string; targetPercent: string |} list + status: string + createdAt: string + } + +type RebalanceExecutionResponse = + { + planId: Guid + runDate: string + outcomes: {| instrumentCode: string; action: string; amount: string; status: string; orderId: string option; pendingReason: string option |} list + } + type CapitalDepositResponse = { id: Guid @@ -314,6 +330,38 @@ module App = createdAt = timestampText plan.CreatedAt } + let private rebalancePlanResponse (plan: RebalancePlanRecord) : RebalancePlanResponse = + { + id = plan.Id + fundId = plan.FundId + targets = + plan.Targets + |> List.map (fun target -> + {| + instrumentCode = target.InstrumentCode + targetPercent = cashText target.TargetPercent + |}) + status = plan.Status + createdAt = timestampText plan.CreatedAt + } + + let private rebalanceExecutionResponse (result: RebalanceExecutionResult) : RebalanceExecutionResponse = + { + planId = result.PlanId + runDate = dateText result.RunDate + outcomes = + result.Outcomes + |> List.map (fun outcome -> + {| + instrumentCode = outcome.InstrumentCode + action = outcome.Action + amount = outcome.Amount + status = outcome.Status + orderId = outcome.OrderId |> Option.map (fun id -> id.ToString("D")) + pendingReason = outcome.PendingReason + |}) + } + let private errorResponse status error message : HttpHandler = setStatusCode status >=> json ({ @@ -469,9 +517,52 @@ module App = with | :? JsonException -> Error "request body must be valid JSON" + let private parseRebalancePlanCommand (body: string) = + try + use document = JsonDocument.Parse(body) + let root = document.RootElement + + if root.ValueKind <> JsonValueKind.Object then + Error "request body must be a JSON object" + else + let mutable parsed = None + let mutable failure = None + use targetsElement = document.RootElement.Clone() + + if not (root.TryGetProperty("targets", &targetsElement)) then + failure <- Some "targets is required" + elif targetsElement.ValueKind <> JsonValueKind.Array then + failure <- Some "targets must be an array" + else + let rows = ResizeArray<RebalanceTarget>() + let mutable rowError = None + + for item in targetsElement.EnumerateArray() do + if rowError.IsNone then + match tryStringProperty item "instrumentCode", tryStringProperty item "targetPercent" with + | Some code, Some percentText -> + match tryDecimal "targetPercent" percentText with + | Ok percent -> rows.Add({ InstrumentCode = code; TargetPercent = percent }) + | Error message -> rowError <- Some message + | _ -> rowError <- Some "each target needs instrumentCode and targetPercent" + + failure <- rowError + + if failure.IsNone then parsed <- Some { Targets = rows |> Seq.toList } + + match failure with + | Some message -> Error message + | None -> + match parsed with + | Some command -> Ok command + | None -> Error "targets is required" + with + | :? JsonException -> Error "request body must be valid JSON" + let private invokeHandler handler next ctx = handler next ctx + let private unauthorized : HttpHandler = setStatusCode 401 >=> setHttpHeader "WWW-Authenticate" "Bearer" @@ -621,9 +712,9 @@ module App = match repository.GetFund fundId with | None -> errorResponse 404 "FUND_NOT_FOUND" "fund was not found" next ctx | Some fund -> - let positions = + let positions: FundPositionResponse list = repository.GetFundPositions fundId - |> List.map (fun position -> + |> List.map (fun (position: FundPositionRecord) -> { instrumentCode = position.InstrumentCode units = decimalText position.Units 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<DateTimeOffset>(4) + } + + let rebalanceTargetsFromReader (reader: DbDataReader) : RebalanceTarget list = + let targets = ResizeArray<RebalanceTarget>() + + 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<DateTimeOffset>(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<RebalancePlanRecord>() + + while reader.Read() do + plans.Add(rebalancePlanRecordFromReader reader) + + reader.Close() + plans |> Seq.toList |> List.map (rebalancePlanWithTargets connection None) + + member this.ExecuteRebalancePlan(planId: Guid) : Result<RebalanceExecutionResult, string> = + 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<RebalanceOrderOutcome>() + + // 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() |
