From d2a547b6a60ea769092894fff1f2b5105954f226 Mon Sep 17 00:00:00 2001 From: "Somhairle H. Marisol" Date: Mon, 21 Sep 2026 12:38:56 +0800 Subject: Drive SIP plans through the shared order pipeline (3d-6) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Advance executes due periods as real subscription orders: each due trade date claims a sip_executions row (PRIMARY KEY plan_id+trade_date makes same-plan same-day runs atomic), places the order through CreateSubscriptionOrder with a deterministic idempotency key (sip:{plan}:{date}) and the period's trade date, then confirms through the shared ConfirmSubscriptionOrder pipeline with key sip-confirm:{plan}:{date} — identical freeze/fee/units math as manual orders, nothing bypasses confirmation. Outcomes: succeeded, pending_nav (cash stays frozen; later drives retry the confirmation once the NAV lands), insufficient_cash (no order placed, period fails visibly, reruns keep the marker), failed. Plans roll their next_trade_date past the processed window and repeats of the same endDate replay recorded executions instead of re-debiting. Domain adds SipPolicy.advancePlan (due-date enumeration; weekends-only trading-calendar approximation, noted in code). POST /funds/{id}/sip/advance {endDate?, limit?} drives execution on demand (no daemon). Plan responses/rows now carry last execution status/date; front-end shows 已执行/待净值/现金不足 per plan. Optional tradeDateOverride/anchorOverride keep manual flows untouched. --- src/FundLab.Api/Persistence.fs | 369 ++++++++++++++++++++++++++++++++++++++++- 1 file changed, 363 insertions(+), 6 deletions(-) (limited to 'src/FundLab.Api/Persistence.fs') diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs index 8ca85e4..f714267 100644 --- a/src/FundLab.Api/Persistence.fs +++ b/src/FundLab.Api/Persistence.fs @@ -275,6 +275,8 @@ type SipPlanRecord = AnchorDate: DateOnly NextTradeDate: DateOnly CreatedAt: DateTimeOffset + LastExecutionStatus: string option + LastExecutionDate: DateOnly option } type SipPlanWriteResult = @@ -285,6 +287,30 @@ type SipPlanWriteResult = | SipPlanFundNotFound | SipPlanInstrumentNotFound +type SipExecutionOutcome = + { + TradeDate: DateOnly + Status: string + OrderId: Guid option + PendingReason: string option + } + +type SipPlanAdvanceResult = + { + PlanId: Guid + InstrumentCode: string + Amount: decimal + Frequency: SipFrequency + Executions: SipExecutionOutcome list + NextTradeDate: DateOnly + } + +type SipAdvanceResult = + { + FundId: Guid + Plans: SipPlanAdvanceResult list + } + type CapitalDepositCommand = { Amount: decimal @@ -523,6 +549,18 @@ type FundRepository(connectionString: string) = 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) + ); """ let statusText status = @@ -1350,6 +1388,9 @@ type FundRepository(connectionString: string) = 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 = @@ -1357,7 +1398,17 @@ type FundRepository(connectionString: string) = commandWithTransaction connection transaction - "SELECT id, fund_id, instrument_code, amount, frequency, status, is_synthetic, anchor_date, next_trade_date, created_at FROM sip_plans WHERE id = @plan_id" + """ + 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 @@ -1611,7 +1662,7 @@ type FundRepository(connectionString: string) = raise error - member _.CreateSubscriptionOrder(idempotencyKey: string, fundId: Guid, command: SubscriptionOrderCommand) = + member _.CreateSubscriptionOrder(idempotencyKey: string, fundId: Guid, command: SubscriptionOrderCommand, ?tradeDateOverride: DateOnly) : SubscriptionOrderWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then OrderInvalid "idempotency key cannot be empty" else @@ -1673,6 +1724,10 @@ type FundRepository(connectionString: string) = OrderInsufficientFunds else let submittedAt = DateTimeOffset.UtcNow + let tradeDate = + tradeDateOverride + |> Option.defaultWith (fun () -> ConfirmationPolicy.tradeDateFor submittedAt) + let order: SubscriptionOrderRecord = { Id = Guid.NewGuid() @@ -1684,7 +1739,7 @@ type FundRepository(connectionString: string) = Status = orderStatusText SubmittedAt = submittedAt IsSynthetic = isSynthetic - TradeDate = ConfirmationPolicy.tradeDateFor submittedAt + TradeDate = tradeDate ConfirmIdempotencyKey = None PendingReason = None ConfirmedAt = None @@ -1741,7 +1796,7 @@ type FundRepository(connectionString: string) = records |> Seq.toList - member _.CreateSipPlan(idempotencyKey: string, fundId: Guid, command: SipPlanCommand) : SipPlanWriteResult = + member _.CreateSipPlan(idempotencyKey: string, fundId: Guid, command: SipPlanCommand, ?anchorOverride: DateOnly) : SipPlanWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then SipPlanWriteResult.SipPlanInvalid "idempotency key cannot be empty" else @@ -1783,7 +1838,10 @@ type FundRepository(connectionString: string) = SipPlanWriteResult.SipPlanFundNotFound | Some isSynthetic -> if instrumentExists connection (Some transaction) command.InstrumentCode then - let anchorDate = ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow + let anchorDate = + anchorOverride + |> Option.defaultWith (fun () -> ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow) + let plan: SipPlanRecord = { Id = Guid.NewGuid() @@ -1796,6 +1854,8 @@ type FundRepository(connectionString: string) = 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 @@ -1813,6 +1873,292 @@ type FundRepository(connectionString: string) = 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() @@ -1821,7 +2167,18 @@ type FundRepository(connectionString: string) = commandWithTransaction connection None - "SELECT id, fund_id, instrument_code, amount, frequency, status, is_synthetic, anchor_date, next_trade_date, created_at FROM sip_plans WHERE fund_id = @fund_id ORDER BY created_at DESC, id" + """ + 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 -- cgit v1.2.3