diff options
| author | Somhairle H. Marisol <[email protected]> | 2026-09-21 12:38:56 +0800 |
|---|---|---|
| committer | Somhairle H. Marisol <[email protected]> | 2026-09-21 12:38:56 +0800 |
| commit | d2a547b6a60ea769092894fff1f2b5105954f226 (patch) | |
| tree | 5dfd493d926891509258ef133c0a54b0fac2363f | |
| parent | f4ab0b08d7648914f6fd049fa1887c0422f69318 (diff) | |
| download | fund-lab-d2a547b6a60ea769092894fff1f2b5105954f226.tar.gz | |
Drive SIP plans through the shared order pipeline (3d-6)
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.
| -rw-r--r-- | src/FundLab.Api/App.fs | 92 | ||||
| -rw-r--r-- | src/FundLab.Api/Persistence.fs | 369 | ||||
| -rw-r--r-- | src/FundLab.Domain/Sip.fs | 19 | ||||
| -rw-r--r-- | src/FundLab.Web/App.fs | 21 | ||||
| -rw-r--r-- | tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj | 1 | ||||
| -rw-r--r-- | tests/FundLab.Api.Tests/SipAdvanceTests.fs | 339 | ||||
| -rw-r--r-- | tests/FundLab.Domain.Tests/DomainTests.fs | 22 | ||||
| -rw-r--r-- | tests/FundLab.Web.Tests/BoundaryTests.fs | 2 |
8 files changed, 859 insertions, 6 deletions
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<DateOnly>(7) NextTradeDate = reader.GetFieldValue<DateOnly>(8) CreatedAt = reader.GetFieldValue<DateTimeOffset>(9) + LastExecutionStatus = readStringOption reader 10 + LastExecutionDate = + if reader.IsDBNull(11) then None else Some(reader.GetFieldValue<DateOnly>(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<Guid * DateOnly * Guid>() + + while retryReader.Read() do + pendingRetries.Add( + retryReader.GetGuid(0), + retryReader.GetFieldValue<DateOnly>(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<Guid * string * decimal * SipFrequency * DateOnly * DateOnly>() + + 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<DateOnly>(4), + plansReader.GetFieldValue<DateOnly>(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<SipExecutionOutcome>() + + while reader.Read() do + rows.Add( + { + TradeDate = reader.GetFieldValue<DateOnly>(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<SipExecutionOutcome>() + + 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 diff --git a/src/FundLab.Domain/Sip.fs b/src/FundLab.Domain/Sip.fs index 2db7fbd..28e0896 100644 --- a/src/FundLab.Domain/Sip.fs +++ b/src/FundLab.Domain/Sip.fs @@ -82,6 +82,25 @@ module SipPolicy = let months = max 0 candidateMonths rollToWeekday (addMonthsPreservingDay anchorDate months) + /// Due trade dates a plan must run through, from its current nextTradeDate up to and + /// including endDate. Calendar approximation: only weekends are skipped (Chinese fund + /// trading calendar approximated by weekdays; public holidays are not modeled here — + /// noted per the no-fabrication rule, holidays may be added to the calendar later). + let advancePlan + (frequency: SipFrequency) + (anchorDate: DateOnly) + (dueFrom: DateOnly) + (endDate: DateOnly) + : DateOnly list = + let rec loop (current: DateOnly) (acc: DateOnly list) = + if current > endDate then + List.rev acc + else + let following = nextTradeDate frequency anchorDate (current.AddDays 1) + loop following (current :: acc) + + if dueFrom > endDate then [] else loop dueFrom [] + /// Execution is out of scope for this slice. These pure helpers define the /// due/settlement contract that the (future) executor will implement: when a /// plan is due, the net amount debited equals the plan amount exactly, and no diff --git a/src/FundLab.Web/App.fs b/src/FundLab.Web/App.fs index 6fe1189..5d091dc 100644 --- a/src/FundLab.Web/App.fs +++ b/src/FundLab.Web/App.fs @@ -169,6 +169,8 @@ type RawSipPlan = status: string anchorDate: string nextTradeDate: string + lastExecutionStatus: obj + lastExecutionDate: obj } type RawRedemption = @@ -276,6 +278,8 @@ type SipPlan = status: string anchorDate: string nextTradeDate: string + lastExecutionStatus: string option + lastExecutionDate: string option } type OrderDetail = @@ -458,6 +462,8 @@ module Api = status = raw.status anchorDate = raw.anchorDate nextTradeDate = raw.nextTradeDate + lastExecutionStatus = decodeOptionalText raw.lastExecutionStatus + lastExecutionDate = decodeOptionalText raw.lastExecutionDate } let decodeRedemption (raw: RawRedemption) : RedemptionDetail = @@ -2253,6 +2259,21 @@ let private sipPlanRow (plan: SipPlan) = Html.span [ prop.className "order-cell"; prop.text (sprintf "频率 %s" (sipFrequencyText plan.frequency)) ] Html.span [ prop.className "order-cell"; prop.text (sprintf "起投日 %s" plan.anchorDate) ] Html.span [ prop.className "order-cell"; prop.text (sprintf "下次扣款 %s" plan.nextTradeDate) ] + + let lastExecution = + match plan.lastExecutionStatus with + | Some status -> + let label = + if status = "succeeded" then "已执行" + elif status = "pending_nav" then "待净值" + elif status = "insufficient_cash" then "现金不足" + elif status = "failed" then "失败" + else status + + sprintf "最近执行 %s(%s)" label (plan.lastExecutionDate |> Option.defaultValue "—") + | None -> "最近执行 暂无" + + Html.span [ prop.className "order-cell"; prop.text lastExecution ] Html.span [ prop.className "order-status"; prop.text (if plan.status = "active" then "进行中" else plan.status) ] ] ] diff --git a/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj b/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj index 21778a3..bd45410 100644 --- a/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj +++ b/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj @@ -24,6 +24,7 @@ <Compile Include="ProcessCollectorTests.fs" /> <Compile Include="PersistenceTests.fs" /> <Compile Include="OrderTests.fs" /> + <Compile Include="SipAdvanceTests.fs" /> <Compile Include="Program.fs" /> </ItemGroup> </Project> diff --git a/tests/FundLab.Api.Tests/SipAdvanceTests.fs b/tests/FundLab.Api.Tests/SipAdvanceTests.fs new file mode 100644 index 0000000..7ac305d --- /dev/null +++ b/tests/FundLab.Api.Tests/SipAdvanceTests.fs @@ -0,0 +1,339 @@ +namespace FundLab.Api.Tests + +open System +open Npgsql +open Xunit +open FundLab.Api +open FundLab.Domain + +[<Collection("postgres")>] +type SipAdvanceTests(fixture: PostgresFixture) = + let sharedRepository = + lazy + let value = FundRepository(fixture.ConnectionString) + value.EnsureSchema() + value + + let repository () = sharedRepository.Value + + let seedInstrument () = + let code = Random.Shared.Next(0, 1000000).ToString("D6") + let payload = + { + Source = "akshare" + SourceRevision = "akshare-test/eastmoney" + CollectedAt = DateTimeOffset(2026, 9, 21, 8, 0, 0, TimeSpan.Zero) + Instruments = [ { Code = code; Name = "定投执行测试基金"; FundType = None } ] + } + + repository().UpsertInstruments(payload, "sip-advance-test-hash") + code + + let createFund (initialCash: decimal) = + let command = + { + Name = "定投执行测试 FOF" + InitialCash = initialCash + InitialUnitNav = 1.00000000m + IsSynthetic = true + } + + let key = fixture.Key(sprintf "sip-adv-fund-%s" (Guid.NewGuid().ToString("N"))) + + match repository().CreateFund(key, command) with + | FundWriteResult.Created fund -> fund.Id + | other -> failwithf "unexpected fund creation result: %A" other + + let app () = App.createApplication (repository ()) + + let truncateMicroseconds (moment: DateTimeOffset) = + let utc = moment.ToUniversalTime() + DateTimeOffset(utc.Ticks - (utc.Ticks % 10L), TimeSpan.Zero) + + let insertQuoteOnDate (code: string) (nav: decimal) (navDate: DateOnly) = + let revision = sprintf "akshare-test/%O" (Guid.NewGuid()) + let payload: MarketDataNavPayload = + { + Source = "akshare" + SourceRevision = revision + CollectedAt = truncateMicroseconds (DateTimeOffset.Now.AddSeconds(-10.0)) + Code = code + Observations = + [ + { + NavDate = navDate + PublishedAt = None + Nav = nav + AccumulatedNav = Some nav + DailyReturn = Some 0.0m + } + ] + } + + repository().UpsertNavObservations(payload, sprintf "sip-advance-hash/%s" revision) + + let postPlan (fundId: Guid) (code: string) (amount: string) (frequency: string) = + let body = + sprintf "{\"instrumentCode\":\"%s\",\"amount\":\"%s\",\"frequency\":\"%s\"}" code amount frequency + + PersistenceTestHelpers.invoke + (app ()) + "POST" + (sprintf "/api/funds/%O/sip/plans" fundId) + [ + "Authorization", "Bearer test-token" + "Idempotency-Key", fixture.Key(sprintf "sip-plan-%s" (Guid.NewGuid().ToString("N"))) + ] + body + |> fun (status, response) -> + if status <> 201 then failwithf "unexpected plan status %d: %s" status response + PersistenceTestHelpers.responseId response + + // plans anchored in the past so due periods fall on dates whose NAV the tests can seed + let createPlanWithAnchor (fundId: Guid) (code: string) (amount: decimal) (frequency: SipFrequency) = + let key = fixture.Key(sprintf "sip-plan-%s" (Guid.NewGuid().ToString("N"))) + + match repository().CreateSipPlan(key, fundId, { InstrumentCode = code; Amount = amount; Frequency = frequency }, DateOnly(2026, 9, 7)) with + | SipPlanWriteResult.SipPlanCreated plan -> plan.Id + | other -> failwithf "unexpected sip plan result: %A" other + + let postManualOrder (fundId: Guid) (code: string) (amount: string) = + let body = + sprintf "{\"fundCode\":\"%s\",\"amount\":\"%s\",\"feeAmount\":\"0.00\"}" code amount + + PersistenceTestHelpers.invoke + (app ()) + "POST" + (sprintf "/api/funds/%O/orders" fundId) + [ + "Authorization", "Bearer test-token" + "Idempotency-Key", fixture.Key(sprintf "manual-%s" (Guid.NewGuid().ToString("N"))) + ] + body + |> fun (status, response) -> + if status <> 201 then failwithf "unexpected manual order status %d: %s" status response + PersistenceTestHelpers.responseId response + + let confirmManualOrder (fundId: Guid) (orderId: Guid) = + PersistenceTestHelpers.invoke + (app ()) + "POST" + (sprintf "/api/funds/%O/orders/%O/confirm" fundId orderId) + [ + "Authorization", "Bearer test-token" + "Idempotency-Key", fixture.Key(sprintf "manual-confirm-%s" (Guid.NewGuid().ToString("N"))) + ] + "" + |> fun (status, _) -> + if status <> 200 then failwithf "unexpected manual confirm status %d" status + + let advance (fundId: Guid) (body: string) = + PersistenceTestHelpers.invoke + (app ()) + "POST" + (sprintf "/api/funds/%O/sip/advance" fundId) + [ "Authorization", "Bearer test-token" ] + body + + let scalarDecimal (sql: string) (parameters: (string * obj * NpgsqlTypes.NpgsqlDbType) list) = + use connection = new NpgsqlConnection(fixture.ConnectionString) + connection.Open() + use command = connection.CreateCommand() + command.CommandText <- sql + + for name, value, dbType in parameters do + let parameter = command.Parameters.Add(name, dbType) + parameter.Value <- value + + command.ExecuteScalar() :?> decimal + + let availableCash fundId = + scalarDecimal "SELECT available_cash FROM funds WHERE id = @fund_id" [ "fund_id", box fundId, NpgsqlTypes.NpgsqlDbType.Uuid ] + + let executionCount fundId = + scalarDecimal + "SELECT count(*)::numeric FROM sip_executions e JOIN sip_plans p ON p.id = e.plan_id WHERE p.fund_id = @fund_id" + [ "fund_id", box fundId, NpgsqlTypes.NpgsqlDbType.Uuid ] + |> int64 + + let confirmedOrderCount fundId = + scalarDecimal + "SELECT count(*)::numeric FROM subscription_orders WHERE fund_id = @fund_id AND status = 'confirmed'" + [ "fund_id", box fundId, NpgsqlTypes.NpgsqlDbType.Uuid ] + |> int64 + + let sipOrderSnapshot fundId = + use connection = new NpgsqlConnection(fixture.ConnectionString) + connection.Open() + use command = connection.CreateCommand() + command.CommandText <- + "SELECT o.confirmed_units, o.confirmed_invested_cash, o.confirmed_residual_cash FROM subscription_orders o JOIN sip_executions e ON e.order_id = o.id WHERE o.fund_id = @fund_id ORDER BY o.submitted_at, o.id" + let parameter = command.Parameters.Add("fund_id", NpgsqlTypes.NpgsqlDbType.Uuid) + parameter.Value <- box fundId + use reader = command.ExecuteReader() + let rows = ResizeArray<decimal * decimal * decimal>() + + while reader.Read() do + rows.Add(reader.GetDecimal(0), reader.GetDecimal(1), reader.GetDecimal(2)) + + rows |> Seq.toList + + [<Fact>] + member _.``sip advance executes orders through the shared confirmation pipeline``() = + let fundId = createFund 10000.00m + let code = seedInstrument () + let planId = createPlanWithAnchor fundId code 100.00m Weekly + + insertQuoteOnDate code 2.5m (DateOnly(2026, 9, 14)) + insertQuoteOnDate code 2.5m (DateOnly(2026, 9, 21)) + + let status, response = advance fundId "{\"endDate\":\"2026-09-21\"}" + + Assert.Equal(200, status) + Assert.Contains("\"status\":\"succeeded\"", response) + Assert.Contains("\"tradeDate\":\"2026-09-14\"", response) + Assert.Contains("\"tradeDate\":\"2026-09-21\"", response) + Assert.Contains("\"nextTradeDate\":\"2026-09-28\"", response) + + Assert.Equal(9800.00m, availableCash fundId) + Assert.Equal(2L, executionCount fundId) + Assert.Equal(2L, confirmedOrderCount fundId) + + Assert.All(sipOrderSnapshot fundId, fun (units, invested, residual) -> + Assert.Equal(40.00000000m, units) + Assert.Equal(100.00m, invested) + Assert.Equal(0.00m, residual)) + + [<Fact>] + member _.``sip orders match a manual order with the same parameters exactly``() = + let fundId = createFund 10000.00m + let code = seedInstrument () + let planId = createPlanWithAnchor fundId code 100.00m Weekly + insertQuoteOnDate code 2.5m (DateOnly(2026, 9, 14)) + insertQuoteOnDate code 2.5m (DateOnly(2026, 9, 21)) + + let _, _ = advance fundId "{\"endDate\":\"2026-09-14\"}" + + let manualOrderId = postManualOrder fundId code "100.00" + confirmManualOrder fundId manualOrderId + + use connection = new NpgsqlConnection(fixture.ConnectionString) + connection.Open() + use command = connection.CreateCommand() + command.CommandText <- + "SELECT confirmed_units, confirmed_invested_cash, confirmed_residual_cash FROM subscription_orders WHERE id = @order_id" + let parameter = command.Parameters.Add("order_id", NpgsqlTypes.NpgsqlDbType.Uuid) + parameter.Value <- box manualOrderId + use reader = command.ExecuteReader() + reader.Read() |> ignore + let manualUnits = reader.GetDecimal(0) + let manualInvested = reader.GetDecimal(1) + let manualResidual = reader.GetDecimal(2) + + let sipSnapshots = sipOrderSnapshot fundId + + Assert.Equal(1, sipSnapshots.Length) + + let sipUnits, sipInvested, sipResidual = sipSnapshots[0] + + Assert.Equal(manualUnits, sipUnits) + Assert.Equal(manualInvested, sipInvested) + Assert.Equal(manualResidual, sipResidual) + + [<Fact>] + member _.``sip advance reruns are idempotent for the same trade dates``() = + let fundId = createFund 10000.00m + let code = seedInstrument () + let planId = createPlanWithAnchor fundId code 100.00m Weekly + insertQuoteOnDate code 2.5m (DateOnly(2026, 9, 14)) + insertQuoteOnDate code 2.5m (DateOnly(2026, 9, 21)) + + let firstStatus, firstResponse = advance fundId "{\"endDate\":\"2026-09-21\"}" + let firstCash = availableCash fundId + let firstCount = executionCount fundId + let firstOrders = confirmedOrderCount fundId + + let secondStatus, secondResponse = advance fundId "{\"endDate\":\"2026-09-21\"}" + + Assert.Equal(200, firstStatus) + Assert.Equal(200, secondStatus) + Assert.Equal(firstResponse, secondResponse) + Assert.Equal(firstCash, availableCash fundId) + Assert.Equal(firstCount, executionCount fundId) + Assert.Equal(firstOrders, confirmedOrderCount fundId) + + [<Fact>] + member _.``insufficient cash marks the period without placing any order``() = + let fundId = createFund 1.00m + let code = seedInstrument () + let planId = createPlanWithAnchor fundId code 100.00m Weekly + insertQuoteOnDate code 2.5m (DateOnly(2026, 9, 14)) + + let status, response = advance fundId "{\"endDate\":\"2026-09-14\"}" + + Assert.Equal(200, status) + Assert.Contains("\"status\":\"insufficient_cash\"", response) + Assert.Contains("\"orderId\":null", response) + Assert.Equal(1L, executionCount fundId) + Assert.Equal(0L, confirmedOrderCount fundId) + Assert.Equal(1.00m, availableCash fundId) + + // rerun keeps the failed marker and never debits cash + let rerunStatus, rerunResponse = advance fundId "{\"endDate\":\"2026-09-14\"}" + Assert.Equal(200, rerunStatus) + Assert.Equal(response, rerunResponse) + Assert.Equal(1.00m, availableCash fundId) + + [<Fact>] + member _.``pending nav executions are retried once the nav lands``() = + let fundId = createFund 10000.00m + let code = seedInstrument () + let planId = createPlanWithAnchor fundId code 100.00m Weekly + + let status, response = advance fundId "{\"endDate\":\"2026-09-14\"}" + + Assert.Equal(200, status) + Assert.Contains("\"status\":\"pending_nav\"", response) + + // frozen cash is held by the submitted order until confirmation + let frozen = availableCash fundId + Assert.Equal(9900.00m, frozen) + + insertQuoteOnDate code 2.5m (DateOnly(2026, 9, 14)) + + let retryStatus, retryResponse = advance fundId "{\"endDate\":\"2026-09-14\"}" + + Assert.Equal(200, retryStatus) + Assert.Contains("\"status\":\"succeeded\"", retryResponse) + // freeze happened on first drive; confirmation converts frozen cash to cost, no extra flow + Assert.Equal(9900.00m, availableCash fundId) + Assert.Equal(1L, executionCount fundId) + + [<Fact>] + member _.``redemption frozen shares do not bypass the cash check``() = + let fundId = createFund 10000.00m + let code = seedInstrument () + let planId = createPlanWithAnchor fundId code 100.00m Weekly + + // drain available cash through a subscription freeze, mirroring held positions + let body = sprintf "{\"fundCode\":\"%s\",\"amount\":\"10000.00\",\"feeAmount\":\"0.00\"}" code + + PersistenceTestHelpers.invoke + (app ()) + "POST" + (sprintf "/api/funds/%O/orders" fundId) + [ + "Authorization", "Bearer test-token" + "Idempotency-Key", fixture.Key(sprintf "hold-%s" (Guid.NewGuid().ToString("N"))) + ] + body + |> fun (status, _) -> Assert.Equal(201, status) + + insertQuoteOnDate code 2.5m (DateOnly(2026, 9, 14)) + + let status, response = advance fundId "{\"endDate\":\"2026-09-14\"}" + + Assert.Equal(200, status) + Assert.Contains("\"status\":\"insufficient_cash\"", response) + Assert.Equal(0L, confirmedOrderCount fundId) + Assert.Equal(0.00m, availableCash fundId) diff --git a/tests/FundLab.Domain.Tests/DomainTests.fs b/tests/FundLab.Domain.Tests/DomainTests.fs index 13be13f..42dd2a5 100644 --- a/tests/FundLab.Domain.Tests/DomainTests.fs +++ b/tests/FundLab.Domain.Tests/DomainTests.fs @@ -621,3 +621,25 @@ module SipPolicyTests = Assert.Equal(Ok 200.00m, SipPolicy.executionNetAmount 200.00m) Assert.Equal(Error "sip amount must be positive", SipPolicy.executionNetAmount 0m) + + [<Fact>] + let ``advance plan enumerates due dates up to the end date`` () = + let anchor = d 2026 9 21 + + Assert.Equal<DateOnly list>( + [ d 2026 9 28; d 2026 10 5; d 2026 10 12 ], + SipPolicy.advancePlan Weekly anchor (d 2026 9 28) (d 2026 10 12) + ) + + Assert.Equal<DateOnly list>([ d 2026 9 28 ], SipPolicy.advancePlan Weekly anchor (d 2026 9 28) (d 2026 9 30)) + Assert.Equal<DateOnly list>([], SipPolicy.advancePlan Weekly anchor (d 2026 10 5) (d 2026 9 30)) + + [<Fact>] + let ``advance plan never emits weekend dates`` () = + let anchor = d 2026 1 31 + + let dates = SipPolicy.advancePlan Monthly anchor (d 2026 2 2) (d 2026 5 31) + + Assert.All(dates, fun (date: DateOnly) -> Assert.NotEqual(DayOfWeek.Saturday, date.DayOfWeek); Assert.NotEqual(DayOfWeek.Sunday, date.DayOfWeek)) + Assert.Contains(d 2026 3 2, dates) + Assert.Equal(4, dates.Length) diff --git a/tests/FundLab.Web.Tests/BoundaryTests.fs b/tests/FundLab.Web.Tests/BoundaryTests.fs index dbaf4e5..5ab327d 100644 --- a/tests/FundLab.Web.Tests/BoundaryTests.fs +++ b/tests/FundLab.Web.Tests/BoundaryTests.fs @@ -732,6 +732,8 @@ module SipBoundaryTests = status = "active" anchorDate = "2026-09-21" nextTradeDate = "2026-09-28" + lastExecutionStatus = null + lastExecutionDate = null } [<Fact>] |
