summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-21 12:38:56 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-21 12:38:56 +0800
commitd2a547b6a60ea769092894fff1f2b5105954f226 (patch)
tree5dfd493d926891509258ef133c0a54b0fac2363f
parentf4ab0b08d7648914f6fd049fa1887c0422f69318 (diff)
downloadfund-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.fs92
-rw-r--r--src/FundLab.Api/Persistence.fs369
-rw-r--r--src/FundLab.Domain/Sip.fs19
-rw-r--r--src/FundLab.Web/App.fs21
-rw-r--r--tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj1
-rw-r--r--tests/FundLab.Api.Tests/SipAdvanceTests.fs339
-rw-r--r--tests/FundLab.Domain.Tests/DomainTests.fs22
-rw-r--r--tests/FundLab.Web.Tests/BoundaryTests.fs2
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>]