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/App.fs | 92 ++++++++++ src/FundLab.Api/Persistence.fs | 369 ++++++++++++++++++++++++++++++++++++++++- 2 files changed, 455 insertions(+), 6 deletions(-) (limited to 'src/FundLab.Api') diff --git a/src/FundLab.Api/App.fs b/src/FundLab.Api/App.fs index 9f88468..9badcfa 100644 --- a/src/FundLab.Api/App.fs +++ b/src/FundLab.Api/App.fs @@ -142,6 +142,8 @@ type SipPlanResponse = status: string anchorDate: string nextTradeDate: string + lastExecutionStatus: string option + lastExecutionDate: string option isSynthetic: bool createdAt: string } @@ -306,6 +308,8 @@ module App = status = plan.Status anchorDate = dateText plan.AnchorDate nextTradeDate = dateText plan.NextTradeDate + lastExecutionStatus = plan.LastExecutionStatus + lastExecutionDate = plan.LastExecutionDate |> Option.map dateText isSynthetic = plan.IsSynthetic createdAt = timestampText plan.CreatedAt } @@ -810,6 +814,93 @@ module App = json (plans |> List.map sipPlanResponse) next ctx with _ -> errorResponse 500 "PERSISTENCE_ERROR" "sip plan persistence failed" next ctx + + let private parseSipAdvanceCommand (body: string) = + try + if String.IsNullOrWhiteSpace body then + Ok(None, 24) + else + 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 endDate = + match tryStringProperty root "endDate" with + | None -> None + | Some text -> + match DateOnly.TryParseExact(text, "yyyy-MM-dd", CultureInfo.InvariantCulture, DateTimeStyles.None) with + | true, date -> Some date + | _ -> failwith "endDate must be yyyy-MM-dd" + + let limit = + match tryStringProperty root "limit" with + | None -> 24 + | Some text -> + match Int32.TryParse(text, CultureInfo.InvariantCulture) with + | true, value when value > 0 -> value + | _ -> failwith "limit must be a positive integer" + + Ok(endDate, limit) + with + | :? JsonException -> Error "request body must be valid JSON" + + let private sipPlanAdvanceResponse (result: SipPlanAdvanceResult) = + {| + planId = result.PlanId + instrumentCode = result.InstrumentCode + amount = cashText result.Amount + frequency = SipPolicy.frequencyText result.Frequency + nextTradeDate = dateText result.NextTradeDate + executions = + result.Executions + |> List.map (fun outcome -> + {| + tradeDate = dateText outcome.TradeDate + status = outcome.Status + orderId = outcome.OrderId |> Option.map (fun id -> id.ToString("D")) + pendingReason = outcome.PendingReason + |}) + |} + + let private advanceSipPlans (repository: FundRepository) (fundIdText: string) : HttpHandler = + fun next ctx -> + task { + match Guid.TryParse fundIdText with + | false, _ -> + return! invokeHandler (errorResponse 400 "INVALID_SIP_REQUEST" "fund id must be a UUID") next ctx + | true, fundId -> + use reader = new StreamReader(ctx.Request.Body) + let! body = reader.ReadToEndAsync() + + match parseSipAdvanceCommand body with + | Error message -> + return! invokeHandler (errorResponse 400 "INVALID_SIP_REQUEST" message) next ctx + | Ok(endDateText, limit) -> + let endDate = + endDateText + |> Option.defaultWith (fun () -> ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow) + + try + match repository.GetFund fundId with + | None -> + return! invokeHandler (errorResponse 404 "FUND_NOT_FOUND" "fund was not found") next ctx + | Some _ -> + let advanced = repository.AdvanceSipPlans(fundId, endDate, limit) + + return! + json + {| + fundId = fundId + endDate = dateText endDate + plans = advanced.Plans |> List.map sipPlanAdvanceResponse + |} + next + ctx + with _ -> + return! invokeHandler (errorResponse 500 "PERSISTENCE_ERROR" "sip advance failed") next ctx + } let private marketDataError (failure: MarketDataFailure) : HttpHandler = let status, error, message = match failure with @@ -916,6 +1007,7 @@ module App = GET >=> routef "/funds/%s/capital/deposits" (getCapitalDeposits repository) POST >=> routef "/funds/%s/sip/plans" (createSipPlan repository) GET >=> routef "/funds/%s/sip/plans" (getSipPlans repository) + POST >=> routef "/funds/%s/sip/advance" (advanceSipPlans repository) GET >=> routef "/funds/%s" (getFund repository) ] @ (marketData |> Option.map marketDataRoutes |> Option.defaultValue []) 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