summaryrefslogtreecommitdiff
path: root/src/FundLab.Api/Persistence.fs
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-21 13:28:02 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-21 13:28:02 +0800
commita3dddc28f33254317d07323037ba5d1db35df3d6 (patch)
treec2a13404bda8a8b62d07559c2ab95b5491034f70 /src/FundLab.Api/Persistence.fs
parentd2a547b6a60ea769092894fff1f2b5105954f226 (diff)
downloadfund-lab-a3dddc28f33254317d07323037ba5d1db35df3d6.tar.gz
Add rebalance execution through shared order pipeline (3d-7)
Diffstat (limited to 'src/FundLab.Api/Persistence.fs')
-rw-r--r--src/FundLab.Api/Persistence.fs447
1 files changed, 446 insertions, 1 deletions
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()