summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-21 23:21:54 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-21 23:21:54 +0800
commitd4b539a2791a0cd4d098fecfece0039073a9f7e2 (patch)
treedfd314b37fa24379ef13dcd66a1de76e2a923bf2
parent12c4d3625458c830a2746a1f277fb684e9c498cb (diff)
downloadfund-lab-d4b539a2791a0cd4d098fecfece0039073a9f7e2.tar.gz
Add scheduled investment plans (3d-10)
-rw-r--r--docs/qa-snapshots/qa-2026-09-21.pngbin0 -> 99746 bytes
-rw-r--r--qa/driver/browser-test.js111
-rw-r--r--src/FundLab.Api/App.fs208
-rw-r--r--src/FundLab.Api/Persistence.fs602
-rw-r--r--src/FundLab.Domain/FundLab.Domain.fsproj1
-rw-r--r--src/FundLab.Domain/InvestmentPlan.fs97
-rw-r--r--src/FundLab.Web/App.fs358
-rw-r--r--src/FundLab.Web/src/api.js12
-rw-r--r--src/FundLab.Web/src/styles.css6
-rw-r--r--tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj1
-rw-r--r--tests/FundLab.Api.Tests/InvestmentPlanTests.fs322
-rw-r--r--tests/FundLab.Domain.Tests/FundLab.Domain.Tests.fsproj1
-rw-r--r--tests/FundLab.Domain.Tests/InvestmentPlanTests.fs81
13 files changed, 1798 insertions, 2 deletions
diff --git a/docs/qa-snapshots/qa-2026-09-21.png b/docs/qa-snapshots/qa-2026-09-21.png
new file mode 100644
index 0000000..2130a37
--- /dev/null
+++ b/docs/qa-snapshots/qa-2026-09-21.png
Binary files differ
diff --git a/qa/driver/browser-test.js b/qa/driver/browser-test.js
index 6403dc5..6b2d2ec 100644
--- a/qa/driver/browser-test.js
+++ b/qa/driver/browser-test.js
@@ -26,6 +26,8 @@ let lastOrdersResponse = null;
let lastConfirmResponse = null;
let lastRedemptionConfirmResponse = null;
let lastReturnsResponse = null;
+let lastPlanRunResponse = null;
+let lastPlansResponse = null;
function check(name, ok, detail) {
results.push({ name, ok, detail: detail || "" });
@@ -135,6 +137,16 @@ async function summaryLabelExists(page, label, name) {
lastReturnsResponse = await r.json();
} catch {}
}
+ if (r.url().includes("/investment-plans/run") && r.request().method() === "POST" && r.status() < 400) {
+ try {
+ lastPlanRunResponse = await r.json();
+ } catch {}
+ }
+ if (/\/investment-plans$/i.test(new URL(r.url()).pathname) && r.request().method() === "GET" && r.status() < 400) {
+ try {
+ lastPlansResponse = await r.json();
+ } catch {}
+ }
if (r.status() >= 400 && !/\/api\/instruments\//.test(r.url()) && !/\/orders/.test(r.url()) && !/\/redemptions/.test(r.url()) && !/\/capital/.test(r.url())) {
consoleErrors.push("resource " + r.status() + ": " + r.url());
}
@@ -182,6 +194,7 @@ async function summaryLabelExists(page, label, name) {
await capitalScenario(page);
await dividendScenario(page);
await returnsScenario(page);
+ await plansScenario(page);
} finally {
check("G1 无浏览器控制台/页面错误", consoleErrors.length === 0, consoleErrors.slice(0, 3).join(" | "));
await browser.close();
@@ -463,6 +476,104 @@ async function returnsScenario(page) {
await page.screenshot({ path: SHOTS + "/16-returns-curve.png" });
}
+async function postInvestmentPlanViaApi(page, fundId, token) {
+ return page.evaluate(
+ async ({ fundId, token }) => {
+ const randomKey = () => {
+ if (window.crypto && window.crypto.randomUUID) return window.crypto.randomUUID().replace(/-/g, "");
+ return Math.random().toString(16).slice(2).padEnd(32, "0").slice(0, 32);
+ };
+ const response = await fetch(`/api/funds/${fundId}/investment-plans`, {
+ method: "POST",
+ headers: {
+ Authorization: `Bearer ${token}`,
+ "Idempotency-Key": randomKey(),
+ "Content-Type": "application/json",
+ },
+ body: JSON.stringify({ instrumentCode: "000001", amount: "100.00", frequency: "daily" }),
+ });
+ return response.status;
+ },
+ { fundId, token }
+ );
+}
+
+async function plansScenario(page) {
+ // P 系列: 定投计划执行最小闭环(接 R 场景:现金 5442.99,持仓 50.00000000)
+ await page.waitForSelector(".plans-panel", { timeout: 10000 });
+ const fundId = await currentFundId(page);
+ const token = process.env.QA_TOKEN || "qa-token";
+
+ const createStatus = await postInvestmentPlanViaApi(page, fundId, token);
+ await page.click(".plans-refresh-action");
+ const planRowAppeared = await page
+ .waitForSelector(".plans-list .order-row", { timeout: 15000 })
+ .then(() => true)
+ .catch(() => false);
+ const planListText = ((await page.textContent(".plans-list")) || "");
+
+ lastPlanRunResponse = null;
+ await page.click(".plans-run-action");
+ for (let i = 0; i < 200 && lastPlanRunResponse === null; i++) {
+ await page.waitForTimeout(100);
+ }
+ const run = lastPlanRunResponse || {};
+ const runPlans = Array.isArray(run.plans) ? run.plans : [];
+ const firstPlan = runPlans.length > 0 ? runPlans[0] : {};
+ const runs = Array.isArray(firstPlan.runs) ? firstPlan.runs : [];
+ const firstRun = runs.length > 0 ? runs[0] : {};
+
+ for (let i = 0; i < 200 && (lastPositionsResponse === null || lastPositionsResponse.availableCash !== "5343.00"); i++) {
+ await page.waitForTimeout(100);
+ }
+
+ for (let i = 0; i < 200 && lastPlansResponse === null; i++) {
+ await page.waitForTimeout(100);
+ }
+
+ const plans = Array.isArray(lastPlansResponse) ? lastPlansResponse : [];
+ const plan = plans.length > 0 ? plans[0] : {};
+
+ check(
+ "P1 到期定投计划按当日净值扣款买入",
+ createStatus === 201 &&
+ planRowAppeared &&
+ planListText.includes("000001") &&
+ planListText.includes("金额 100.00") &&
+ planListText.includes("频率 每日") &&
+ plan.amount === "100.00" &&
+ plan.frequency === "daily" &&
+ run.processingDate === "2026-09-21" &&
+ firstRun.runDate === "2026-09-21" &&
+ firstRun.status === "succeeded" &&
+ lastPositionsResponse &&
+ lastPositionsResponse.availableCash === "5343.00" &&
+ lastPositionsResponse.positions?.[0]?.units === "134.01949252",
+ `status=${createStatus} run=${JSON.stringify(run).slice(0, 240)} cash=${lastPositionsResponse?.availableCash}`
+ );
+
+ lastPlanRunResponse = null;
+ await page.click(".plans-run-action");
+ for (let i = 0; i < 200 && lastPlanRunResponse === null; i++) {
+ await page.waitForTimeout(100);
+ }
+ const replay = lastPlanRunResponse || {};
+ const replayPlans = Array.isArray(replay.plans) ? replay.plans : [];
+ const replayRuns = replayPlans.length > 0 && Array.isArray(replayPlans[0].runs) ? replayPlans[0].runs : [];
+
+ check(
+ "P2 重复执行幂等且不产生第二笔扣款",
+ replayRuns.length === 1 &&
+ replayRuns[0].status === "succeeded" &&
+ lastPositionsResponse &&
+ lastPositionsResponse.availableCash === "5343.00" &&
+ lastPositionsResponse.positions?.[0]?.units === "134.01949252",
+ `runs=${JSON.stringify(replayRuns).slice(0, 200)} cash=${lastPositionsResponse?.availableCash}`
+ );
+
+ await page.screenshot({ path: SHOTS + "/17-investment-plans.png" });
+}
+
const failed = results.filter((r) => !r.ok);
console.log(`\n==== ${results.length - failed.length}/${results.length} passed ====`);
process.exit(failed.length > 0 ? 1 : 0);
diff --git a/src/FundLab.Api/App.fs b/src/FundLab.Api/App.fs
index fc174e5..b559ac2 100644
--- a/src/FundLab.Api/App.fs
+++ b/src/FundLab.Api/App.fs
@@ -148,6 +148,47 @@ type SipPlanResponse =
createdAt: string
}
+type InvestmentPlanResponse =
+ {
+ id: Guid
+ fundId: Guid
+ instrumentCode: string
+ amount: string
+ frequency: string
+ status: string
+ anchorDate: string
+ nextRunDate: string
+ lastRunStatus: string option
+ lastRunDate: string option
+ isSynthetic: bool
+ createdAt: string
+ }
+
+type InvestmentPlanRunOutcomeResponse =
+ {
+ runDate: string
+ status: string
+ orderId: string option
+ pendingReason: string option
+ }
+
+type InvestmentPlanRunPlanResponse =
+ {
+ planId: Guid
+ instrumentCode: string
+ amount: string
+ frequency: string
+ runs: InvestmentPlanRunOutcomeResponse list
+ nextRunDate: string
+ }
+
+type InvestmentPlanRunResponse =
+ {
+ fundId: Guid
+ processingDate: string
+ plans: InvestmentPlanRunPlanResponse list
+ }
+
type DividendResponse =
{
id: Guid
@@ -370,6 +411,47 @@ module App =
createdAt = timestampText plan.CreatedAt
}
+ let private investmentPlanResponse (plan: InvestmentPlanRecord) : InvestmentPlanResponse =
+ {
+ id = plan.Id
+ fundId = plan.FundId
+ instrumentCode = plan.InstrumentCode
+ amount = cashText plan.Amount
+ frequency = InvestmentPlanPolicy.frequencyText plan.Frequency
+ status = plan.Status
+ anchorDate = dateText plan.AnchorDate
+ nextRunDate = dateText plan.NextRunDate
+ lastRunStatus = plan.LastRunStatus
+ lastRunDate = plan.LastRunDate |> Option.map dateText
+ isSynthetic = plan.IsSynthetic
+ createdAt = timestampText plan.CreatedAt
+ }
+
+ let private investmentPlanRunOutcomeResponse (outcome: InvestmentPlanRunOutcome) : InvestmentPlanRunOutcomeResponse =
+ {
+ runDate = dateText outcome.RunDate
+ status = outcome.Status
+ orderId = outcome.OrderId |> Option.map (fun id -> id.ToString("D"))
+ pendingReason = outcome.PendingReason
+ }
+
+ let private investmentPlanRunPlanResponse (result: InvestmentPlanRunPlanResult) : InvestmentPlanRunPlanResponse =
+ {
+ planId = result.PlanId
+ instrumentCode = result.InstrumentCode
+ amount = cashText result.Amount
+ frequency = InvestmentPlanPolicy.frequencyText result.Frequency
+ runs = result.Runs |> List.map investmentPlanRunOutcomeResponse
+ nextRunDate = dateText result.NextRunDate
+ }
+
+ let private investmentPlanRunResponse (result: InvestmentPlanRunResult) : InvestmentPlanRunResponse =
+ {
+ fundId = result.FundId
+ processingDate = dateText result.ProcessingDate
+ plans = result.Plans |> List.map investmentPlanRunPlanResponse
+ }
+
let private rebalancePlanResponse (plan: RebalancePlanRecord) : RebalancePlanResponse =
{
id = plan.Id
@@ -532,7 +614,7 @@ module App =
with
| :? JsonException -> Error "request body must be valid JSON"
- let private parseSipPlanCommand (body: string) =
+ let private parseSipPlanCommand (body: string) : Result<SipPlanCommand, string> =
try
use document = JsonDocument.Parse(body)
let root = document.RootElement
@@ -557,6 +639,31 @@ module App =
with
| :? JsonException -> Error "request body must be valid JSON"
+ let private parseInvestmentPlanCommand (body: string) : Result<InvestmentPlanCommand, string> =
+ try
+ use document = JsonDocument.Parse(body)
+ let root = document.RootElement
+
+ if root.ValueKind <> JsonValueKind.Object then
+ Error "request body must be a JSON object"
+ else
+ match tryStringProperty root "instrumentCode", tryStringProperty root "amount", tryStringProperty root "frequency" with
+ | Some code, Some amountText, Some frequencyText ->
+ match tryDecimal "amount" amountText, InvestmentPlanPolicy.parseFrequency frequencyText with
+ | Ok amount, Some frequency ->
+ Ok
+ {
+ InstrumentCode = code
+ Amount = amount
+ Frequency = frequency
+ }
+ | Error message, _ -> Error message
+ | _, None -> Error "frequency must be one of daily, weekly or monthly"
+ | _ ->
+ Error "instrumentCode, amount and frequency are required"
+ with
+ | :? JsonException -> Error "request body must be valid JSON"
+
let private parseRebalancePlanCommand (body: string) =
try
use document = JsonDocument.Parse(body)
@@ -1078,6 +1185,102 @@ module App =
with _ ->
errorResponse 500 "PERSISTENCE_ERROR" "sip plan persistence failed" next ctx
+ let private createInvestmentPlan (repository: FundRepository) (fundIdText: string) : HttpHandler =
+ fun next ctx ->
+ task {
+ match Guid.TryParse fundIdText with
+ | false, _ ->
+ return! invokeHandler (errorResponse 400 "INVALID_INVESTMENT_PLAN_REQUEST" "fund id must be a UUID") next ctx
+ | true, fundId ->
+ use reader = new StreamReader(ctx.Request.Body)
+ let! body = reader.ReadToEndAsync()
+ let idempotencyKey = ctx.Request.Headers["Idempotency-Key"].ToString()
+
+ match parseInvestmentPlanCommand body with
+ | Error message ->
+ return! invokeHandler (errorResponse 400 "INVALID_INVESTMENT_PLAN_REQUEST" message) next ctx
+ | Ok command ->
+ try
+ match repository.CreateInvestmentPlan(idempotencyKey, fundId, command) with
+ | InvestmentPlanWriteResult.InvestmentPlanCreated plan ->
+ return! invokeHandler (setStatusCode 201 >=> json (investmentPlanResponse plan)) next ctx
+ | InvestmentPlanWriteResult.InvestmentPlanReplayed plan ->
+ return! invokeHandler (json (investmentPlanResponse plan)) next ctx
+ | InvestmentPlanWriteResult.InvestmentPlanIdempotencyConflict ->
+ return! invokeHandler (errorResponse 409 "IDEMPOTENCY_CONFLICT" "idempotency key was used with a different request") next ctx
+ | InvestmentPlanWriteResult.InvestmentPlanInvalid message ->
+ return! invokeHandler (errorResponse 400 "INVALID_INVESTMENT_PLAN_REQUEST" message) next ctx
+ | InvestmentPlanWriteResult.InvestmentPlanFundNotFound ->
+ return! invokeHandler (errorResponse 404 "FUND_NOT_FOUND" "fund was not found") next ctx
+ | InvestmentPlanWriteResult.InvestmentPlanInstrumentNotFound ->
+ return! invokeHandler (errorResponse 404 "INSTRUMENT_NOT_FOUND" "instrument code was not found in the instrument catalog") next ctx
+ with _ ->
+ return! invokeHandler (errorResponse 500 "PERSISTENCE_ERROR" "investment plan persistence failed") next ctx
+ }
+
+ let private getInvestmentPlans (repository: FundRepository) (fundIdText: string) : HttpHandler =
+ fun next ctx ->
+ match Guid.TryParse fundIdText with
+ | false, _ -> errorResponse 400 "INVALID_FUND_ID" "fund id must be a UUID" next ctx
+ | true, fundId ->
+ try
+ match repository.GetFund fundId with
+ | None -> errorResponse 404 "FUND_NOT_FOUND" "fund was not found" next ctx
+ | Some _ ->
+ let plans = repository.GetInvestmentPlans fundId
+ json (plans |> List.map investmentPlanResponse) next ctx
+ with _ ->
+ errorResponse 500 "PERSISTENCE_ERROR" "investment plan persistence failed" next ctx
+
+ let private parseInvestmentPlanRunCommand (body: string) =
+ try
+ if String.IsNullOrWhiteSpace body then
+ Ok None
+ 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
+ match tryStringProperty root "processingDate" with
+ | None -> Ok None
+ | Some text ->
+ match DateOnly.TryParseExact(text, "yyyy-MM-dd", CultureInfo.InvariantCulture, DateTimeStyles.None) with
+ | true, date -> Ok(Some date)
+ | _ -> Error "processingDate must be yyyy-MM-dd"
+ with
+ | :? JsonException -> Error "request body must be valid JSON"
+
+ let private runInvestmentPlans (repository: FundRepository) (fundIdText: string) : HttpHandler =
+ fun next ctx ->
+ task {
+ match Guid.TryParse fundIdText with
+ | false, _ ->
+ return! invokeHandler (errorResponse 400 "INVALID_INVESTMENT_PLAN_REQUEST" "fund id must be a UUID") next ctx
+ | true, fundId ->
+ use reader = new StreamReader(ctx.Request.Body)
+ let! body = reader.ReadToEndAsync()
+
+ match parseInvestmentPlanRunCommand body with
+ | Error message ->
+ return! invokeHandler (errorResponse 400 "INVALID_INVESTMENT_PLAN_REQUEST" message) next ctx
+ | Ok processingDateText ->
+ let processingDate =
+ processingDateText
+ |> 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 result = repository.RunInvestmentPlans(fundId, processingDate)
+ return! invokeHandler (json (investmentPlanRunResponse result)) next ctx
+ with _ ->
+ return! invokeHandler (errorResponse 500 "PERSISTENCE_ERROR" "investment plan run failed") next ctx
+ }
+
let private parseSipAdvanceCommand (body: string) =
try
if String.IsNullOrWhiteSpace body then
@@ -1347,6 +1550,9 @@ module App =
POST >=> routef "/funds/%s/dividends" (createDividend repository)
GET >=> routef "/funds/%s/dividends" (getDividends repository)
GET >=> routef "/funds/%s/returns" (getFundReturns repository)
+ POST >=> routef "/funds/%s/investment-plans/run" (runInvestmentPlans repository)
+ POST >=> routef "/funds/%s/investment-plans" (createInvestmentPlan repository)
+ GET >=> routef "/funds/%s/investment-plans" (getInvestmentPlans 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 a8abd94..e1cbb6d 100644
--- a/src/FundLab.Api/Persistence.fs
+++ b/src/FundLab.Api/Persistence.fs
@@ -325,6 +325,62 @@ type SipAdvanceResult =
Plans: SipPlanAdvanceResult list
}
+type InvestmentPlanCommand =
+ {
+ InstrumentCode: string
+ Amount: decimal
+ Frequency: InvestmentFrequency
+ }
+
+type InvestmentPlanRecord =
+ {
+ Id: Guid
+ FundId: Guid
+ InstrumentCode: string
+ Amount: decimal
+ Frequency: InvestmentFrequency
+ Status: string
+ IsSynthetic: bool
+ AnchorDate: DateOnly
+ NextRunDate: DateOnly
+ CreatedAt: DateTimeOffset
+ LastRunStatus: string option
+ LastRunDate: DateOnly option
+ }
+
+type InvestmentPlanWriteResult =
+ | InvestmentPlanCreated of InvestmentPlanRecord
+ | InvestmentPlanReplayed of InvestmentPlanRecord
+ | InvestmentPlanIdempotencyConflict
+ | InvestmentPlanInvalid of string
+ | InvestmentPlanFundNotFound
+ | InvestmentPlanInstrumentNotFound
+
+type InvestmentPlanRunOutcome =
+ {
+ RunDate: DateOnly
+ Status: string
+ OrderId: Guid option
+ PendingReason: string option
+ }
+
+type InvestmentPlanRunPlanResult =
+ {
+ PlanId: Guid
+ InstrumentCode: string
+ Amount: decimal
+ Frequency: InvestmentFrequency
+ Runs: InvestmentPlanRunOutcome list
+ NextRunDate: DateOnly
+ }
+
+type InvestmentPlanRunResult =
+ {
+ FundId: Guid
+ ProcessingDate: DateOnly
+ Plans: InvestmentPlanRunPlanResult list
+ }
+
type RebalanceTarget = RebalancePolicy.TargetAllocation
type RebalancePlanCommand =
@@ -733,6 +789,39 @@ type FundRepository(connectionString: string) =
fund_id uuid NOT NULL REFERENCES funds(id),
created_at timestamptz NOT NULL DEFAULT now()
);
+
+ CREATE TABLE IF NOT EXISTS investment_plans (
+ id uuid PRIMARY KEY,
+ fund_id uuid NOT NULL REFERENCES funds(id),
+ instrument_code text NOT NULL REFERENCES instruments(code),
+ amount numeric(20, 2) NOT NULL CHECK (amount > 0),
+ frequency text NOT NULL,
+ status text NOT NULL,
+ is_synthetic boolean NOT NULL,
+ anchor_date date NOT NULL,
+ next_run_date date NOT NULL,
+ created_at timestamptz NOT NULL DEFAULT now()
+ );
+
+ CREATE TABLE IF NOT EXISTS investment_plan_idempotencies (
+ idempotency_key text PRIMARY KEY,
+ request_hash text NOT NULL,
+ plan_id uuid NOT NULL REFERENCES investment_plans(id),
+ fund_id uuid NOT NULL REFERENCES funds(id),
+ created_at timestamptz NOT NULL DEFAULT now()
+ );
+
+ CREATE TABLE IF NOT EXISTS investment_plan_runs (
+ plan_id uuid NOT NULL REFERENCES investment_plans(id),
+ run_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, run_date)
+ );
"""
let statusText status =
@@ -1667,6 +1756,121 @@ type FundRepository(connectionString: string) =
Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload)))
+ let investmentPlanRecordFromReader (reader: DbDataReader) : InvestmentPlanRecord =
+ {
+ Id = reader.GetGuid(0)
+ FundId = reader.GetGuid(1)
+ InstrumentCode = reader.GetString(2)
+ Amount = reader.GetDecimal(3)
+ Frequency =
+ match InvestmentPlanPolicy.parseFrequency (reader.GetString(4)) with
+ | Some frequency -> frequency
+ | None -> failwith "investment plan frequency is invalid"
+ Status = reader.GetString(5)
+ IsSynthetic = reader.GetBoolean(6)
+ AnchorDate = reader.GetFieldValue<DateOnly>(7)
+ NextRunDate = reader.GetFieldValue<DateOnly>(8)
+ CreatedAt = reader.GetFieldValue<DateTimeOffset>(9)
+ LastRunStatus = readStringOption reader 10
+ LastRunDate = if reader.IsDBNull(11) then None else Some(reader.GetFieldValue<DateOnly>(11))
+ }
+
+ let findInvestmentPlan connection transaction planId =
+ use command =
+ commandWithTransaction
+ connection
+ transaction
+ """
+ SELECT p.id, p.fund_id, p.instrument_code, p.amount, p.frequency, p.status, p.is_synthetic,
+ p.anchor_date, p.next_run_date, p.created_at,
+ r.status, r.run_date
+ FROM investment_plans p
+ LEFT JOIN LATERAL (
+ SELECT status, run_date FROM investment_plan_runs
+ WHERE plan_id = p.id ORDER BY executed_at DESC, run_date DESC LIMIT 1
+ ) r ON true
+ WHERE p.id = @plan_id
+ """
+
+ addParameter command "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore
+
+ use reader = command.ExecuteReader()
+ if reader.Read() then Some(investmentPlanRecordFromReader reader) else None
+
+ let findInvestmentPlanIdempotency connection transaction key =
+ use command =
+ commandWithTransaction
+ connection
+ transaction
+ "SELECT request_hash, fund_id, plan_id FROM investment_plan_idempotencies WHERE idempotency_key = @idempotency_key"
+
+ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore
+
+ use reader = command.ExecuteReader()
+ if reader.Read() then
+ Some(reader.GetString(0), reader.GetGuid(1), reader.GetGuid(2))
+ else
+ None
+
+ let insertInvestmentPlan connection transaction (plan: InvestmentPlanRecord) =
+ use command =
+ commandWithTransaction
+ connection
+ transaction
+ """
+ INSERT INTO investment_plans (id, fund_id, instrument_code, amount, frequency, status, is_synthetic, anchor_date, next_run_date)
+ VALUES (@id, @fund_id, @instrument_code, @amount, @frequency, @status, @is_synthetic, @anchor_date, @next_run_date)
+ RETURNING created_at
+ """
+
+ addParameter command "id" NpgsqlDbType.Uuid (box plan.Id) |> ignore
+ addParameter command "fund_id" NpgsqlDbType.Uuid (box plan.FundId) |> ignore
+ addParameter command "instrument_code" NpgsqlDbType.Text (box plan.InstrumentCode) |> ignore
+ addParameter command "amount" NpgsqlDbType.Numeric (box plan.Amount) |> ignore
+ addParameter command "frequency" NpgsqlDbType.Text (box (InvestmentPlanPolicy.frequencyText plan.Frequency)) |> ignore
+ addParameter command "status" NpgsqlDbType.Text (box plan.Status) |> ignore
+ addParameter command "is_synthetic" NpgsqlDbType.Boolean (box plan.IsSynthetic) |> ignore
+ addParameter command "anchor_date" NpgsqlDbType.Date (box plan.AnchorDate) |> ignore
+ addParameter command "next_run_date" NpgsqlDbType.Date (box plan.NextRunDate) |> ignore
+
+ use reader = command.ExecuteReader()
+ reader.Read() |> ignore
+ reader.GetFieldValue<DateTimeOffset>(0)
+
+ let insertInvestmentPlanIdempotency connection transaction key requestHash planId fundId =
+ use command =
+ commandWithTransaction
+ connection
+ transaction
+ """
+ INSERT INTO investment_plan_idempotencies (idempotency_key, request_hash, plan_id, fund_id)
+ VALUES (@idempotency_key, @request_hash, @plan_id, @fund_id)
+ """
+
+ addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore
+ addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore
+ addParameter command "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore
+ addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore
+ command.ExecuteNonQuery() |> ignore
+
+ let investmentPlanRequestHash (fundId: Guid) (command: InvestmentPlanCommand) =
+ let invariant = CultureInfo.InvariantCulture
+ let encoded (value: string) = sprintf "%d:%s" value.Length value
+ let code = if isNull command.InstrumentCode then "" else command.InstrumentCode
+
+ let payload =
+ String.concat
+ "|"
+ [
+ "investment-plan"
+ encoded (fundId.ToString("D"))
+ encoded code
+ (encoded (command.Amount.ToString("G29", invariant)))
+ (encoded (InvestmentPlanPolicy.frequencyText command.Frequency))
+ ]
+
+ Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload)))
+
let rebalancePlanRecordFromReader (reader: DbDataReader) : RebalancePlanRecord =
{
Id = reader.GetGuid(0)
@@ -2662,6 +2866,404 @@ type FundRepository(connectionString: string) =
records |> Seq.toList
+ member _.CreateInvestmentPlan(idempotencyKey: string, fundId: Guid, command: InvestmentPlanCommand, ?anchorOverride: DateOnly) : InvestmentPlanWriteResult =
+ if String.IsNullOrWhiteSpace idempotencyKey then
+ InvestmentPlanWriteResult.InvestmentPlanInvalid "idempotency key cannot be empty"
+ else
+ match InvestmentPlanPolicy.validateAmount command.Amount with
+ | Error message -> InvestmentPlanWriteResult.InvestmentPlanInvalid message
+ | Ok() ->
+ let fingerprint = investmentPlanRequestHash fundId command
+ use connection = new NpgsqlConnection(connectionString)
+ connection.Open()
+ use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted)
+
+ try
+ use lockCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ "SELECT pg_advisory_xact_lock(hashtext(@lock_key))"
+
+ addParameter lockCommand "lock_key" NpgsqlDbType.Text (box idempotencyKey) |> ignore
+ lockCommand.ExecuteNonQuery() |> ignore
+
+ match findInvestmentPlanIdempotency connection (Some transaction) idempotencyKey with
+ | Some(existingHash, existingFundId, planId)
+ when existingHash = fingerprint && existingFundId = fundId ->
+ match findInvestmentPlan connection (Some transaction) planId with
+ | Some plan ->
+ transaction.Commit()
+ InvestmentPlanWriteResult.InvestmentPlanReplayed plan
+ | None ->
+ transaction.Rollback()
+ InvestmentPlanWriteResult.InvestmentPlanInvalid "idempotency record references a missing plan"
+ | Some _ ->
+ transaction.Rollback()
+ InvestmentPlanWriteResult.InvestmentPlanIdempotencyConflict
+ | None ->
+ match lockFundForOrder connection (Some transaction) fundId with
+ | None ->
+ transaction.Rollback()
+ InvestmentPlanWriteResult.InvestmentPlanFundNotFound
+ | Some isSynthetic ->
+ if instrumentExists connection (Some transaction) command.InstrumentCode then
+ let anchorDate =
+ anchorOverride
+ |> Option.defaultWith (fun () -> ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow)
+
+ let plan: InvestmentPlanRecord =
+ {
+ Id = Guid.NewGuid()
+ FundId = fundId
+ InstrumentCode = command.InstrumentCode
+ Amount = command.Amount
+ Frequency = command.Frequency
+ Status = "active"
+ IsSynthetic = isSynthetic
+ AnchorDate = anchorDate
+ NextRunDate = InvestmentPlanPolicy.nextRunDate command.Frequency anchorDate anchorDate
+ CreatedAt = DateTimeOffset.UtcNow
+ LastRunStatus = None
+ LastRunDate = None
+ }
+
+ let createdAt = insertInvestmentPlan connection (Some transaction) plan
+ insertInvestmentPlanIdempotency connection (Some transaction) idempotencyKey fingerprint plan.Id fundId
+ transaction.Commit()
+ InvestmentPlanWriteResult.InvestmentPlanCreated { plan with CreatedAt = createdAt }
+ else
+ transaction.Rollback()
+ InvestmentPlanWriteResult.InvestmentPlanInstrumentNotFound
+ with error ->
+ try
+ transaction.Rollback()
+ with _ ->
+ ()
+
+ raise error
+
+ member _.GetInvestmentPlans(fundId: Guid) =
+ use connection = new NpgsqlConnection(connectionString)
+ connection.Open()
+
+ use command =
+ commandWithTransaction
+ connection
+ None
+ """
+ SELECT p.id, p.fund_id, p.instrument_code, p.amount, p.frequency, p.status, p.is_synthetic,
+ p.anchor_date, p.next_run_date, p.created_at,
+ r.status, r.run_date
+ FROM investment_plans p
+ LEFT JOIN LATERAL (
+ SELECT status, run_date FROM investment_plan_runs
+ WHERE plan_id = p.id ORDER BY executed_at DESC, run_date DESC LIMIT 1
+ ) r 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
+
+ use reader = command.ExecuteReader()
+ let records = ResizeArray<InvestmentPlanRecord>()
+
+ while reader.Read() do
+ records.Add(investmentPlanRecordFromReader reader)
+
+ records |> Seq.toList
+
+ member this.RunInvestmentPlans(fundId: Guid, processingDate: DateOnly) : InvestmentPlanRunResult =
+ 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 "investment-plan-run:%O" fundId)) |> ignore
+ lockCommand.ExecuteNonQuery() |> ignore
+
+ // 0. retry phase: runs that were pending_nav get re-confirmed once their NAV landed
+ use retryCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ SELECT r.plan_id, r.run_date, r.order_id
+ FROM investment_plan_runs r
+ JOIN investment_plans p ON p.id = r.plan_id
+ WHERE p.fund_id = @fund_id AND r.status = 'pending_nav'
+ ORDER BY r.run_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 "investment-plan-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 investment_plan_runs SET status = 'succeeded', pending_reason = NULL, executed_at = now() WHERE plan_id = @plan_id AND run_date = @run_date"
+
+ addParameter doneCommand "plan_id" NpgsqlDbType.Uuid (box retryPlanId) |> ignore
+ addParameter doneCommand "run_date" NpgsqlDbType.Date (box retryDate) |> ignore
+ doneCommand.ExecuteNonQuery() |> ignore
+ | _ ->
+ // still pending: keep the run as pending_nav for the next drive
+ ()
+
+ use plansCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ "SELECT id, instrument_code, amount, frequency, anchor_date, next_run_date FROM investment_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 * InvestmentFrequency * DateOnly * DateOnly>()
+
+ while plansReader.Read() do
+ plans.Add(
+ plansReader.GetGuid(0),
+ plansReader.GetString(1),
+ plansReader.GetDecimal(2),
+ (match InvestmentPlanPolicy.parseFrequency (plansReader.GetString(3)) with
+ | Some frequency -> frequency
+ | None -> failwith "investment plan frequency is invalid"),
+ plansReader.GetFieldValue<DateOnly>(4),
+ plansReader.GetFieldValue<DateOnly>(5)
+ )
+
+ plansReader.Close()
+
+ let readRunsUpTo (planId: Guid) (upTo: DateOnly) =
+ use replayCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ SELECT run_date, status, order_id, pending_reason
+ FROM investment_plan_runs
+ WHERE plan_id = @plan_id AND run_date <= @up_to
+ ORDER BY run_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<InvestmentPlanRunOutcome>()
+
+ while reader.Read() do
+ rows.Add(
+ {
+ RunDate = 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 runPlanRow (planId: Guid, code: string, amount: decimal, frequency: InvestmentFrequency, anchor: DateOnly, nextDate: DateOnly) =
+ let dueDates = InvestmentPlanPolicy.dueDates frequency anchor nextDate processingDate
+ let outcomes = ResizeArray<InvestmentPlanRunOutcome>()
+
+ let claimRun (runDate: DateOnly) =
+ use insertCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ INSERT INTO investment_plan_runs (plan_id, run_date, amount, fee_amount, status)
+ VALUES (@plan_id, @run_date, @amount, @fee_amount, 'processing')
+ ON CONFLICT (plan_id, run_date) DO NOTHING
+ RETURNING run_date
+ """
+
+ addParameter insertCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore
+ addParameter insertCommand "run_date" NpgsqlDbType.Date (box runDate) |> 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 setRun (runDate: DateOnly) (status: string) (orderId: Guid option) (reason: string option) =
+ use updateCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ UPDATE investment_plan_runs
+ SET status = @status,
+ order_id = @order_id,
+ pending_reason = @reason,
+ executed_at = now()
+ WHERE plan_id = @plan_id AND run_date = @run_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 "run_date" NpgsqlDbType.Date (box runDate) |> ignore
+ updateCommand.ExecuteNonQuery() |> ignore
+
+ let findExistingRun (runDate: DateOnly) =
+ use selectCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ "SELECT status, order_id, pending_reason FROM investment_plan_runs WHERE plan_id = @plan_id AND run_date = @run_date"
+
+ addParameter selectCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore
+ addParameter selectCommand "run_date" NpgsqlDbType.Date (box runDate) |> 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 (runDate: DateOnly) =
+ let orderKey = sprintf "investment-plan:%O:%s" planId (runDate.ToString("yyyy-MM-dd"))
+
+ match this.CreateSubscriptionOrder(orderKey, fundId, { FundCode = code; Amount = amount; FeeAmount = 0m }, runDate) with
+ | SubscriptionOrderWriteResult.OrderCreated order -> Some order.Id
+ | SubscriptionOrderWriteResult.OrderReplayed order -> Some order.Id
+ | SubscriptionOrderWriteResult.OrderInsufficientFunds -> None
+ | other -> failwithf "unexpected investment plan order result: %A" other
+
+ for runDate in dueDates do
+ // 1. claim the slot atomically: same plan + same run date executes once
+ if claimRun runDate then
+ // 2. place the order through the shared pipeline with a deterministic key
+ match orderIdFor runDate with
+ | None ->
+ setRun runDate "insufficient_cash" None (Some "available cash is not enough for the scheduled amount")
+
+ outcomes.Add(
+ {
+ RunDate = runDate
+ 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 "investment-plan-confirm:%O:%s" planId (runDate.ToString("yyyy-MM-dd"))
+
+ match this.ConfirmSubscriptionOrder(confirmKey, fundId, orderId) with
+ | SubscriptionConfirmResult.OrderConfirmed _
+ | SubscriptionConfirmResult.ConfirmReplayed _ ->
+ setRun runDate "succeeded" (Some orderId) None
+ outcomes.Add({ RunDate = runDate; Status = "succeeded"; OrderId = Some orderId; PendingReason = None })
+ | SubscriptionConfirmResult.ConfirmPendingNav record ->
+ setRun runDate "pending_nav" (Some orderId) record.PendingReason
+ outcomes.Add({ RunDate = runDate; Status = "pending_nav"; OrderId = Some orderId; PendingReason = record.PendingReason })
+ | other ->
+ setRun runDate "failed" (Some orderId) (Some(sprintf "%A" other))
+ outcomes.Add({ RunDate = runDate; Status = "failed"; OrderId = Some orderId; PendingReason = Some(sprintf "%A" other) })
+ else
+ match findExistingRun runDate with
+ | Some("succeeded", orderId, reason) ->
+ outcomes.Add({ RunDate = runDate; Status = "succeeded"; OrderId = orderId; PendingReason = reason })
+ | Some(status, orderId, reason) ->
+ outcomes.Add({ RunDate = runDate; Status = status; OrderId = orderId; PendingReason = reason })
+ | None ->
+ outcomes.Add({ RunDate = runDate; Status = "unknown"; OrderId = None; PendingReason = None })
+
+ // 4. roll the plan pointer forward past the processed window
+ let rolled =
+ match dueDates with
+ | [] -> nextDate
+ | lastDueDates -> InvestmentPlanPolicy.nextRunDate frequency anchor ((List.last lastDueDates).AddDays 1)
+
+ if not (List.isEmpty dueDates) then
+ use rollCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ "UPDATE investment_plans SET next_run_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 processing date: replay what was already executed
+ readRunsUpTo planId processingDate
+ else
+ outcomes |> Seq.toList
+
+ {
+ PlanId = planId
+ InstrumentCode = code
+ Amount = amount
+ Frequency = frequency
+ Runs = replayedOutcomes
+ NextRunDate = rolled
+ }
+
+ let planResults = plans |> Seq.map runPlanRow |> Seq.toList
+ transaction.Commit()
+
+ { FundId = fundId; ProcessingDate = processingDate; Plans = planResults }
+ with error ->
+ try
+ transaction.Rollback()
+ with _ ->
+ ()
+
+ raise error
+
member _.CreateRebalancePlan(idempotencyKey: string, fundId: Guid, command: RebalancePlanCommand) : RebalanceWriteResult =
if String.IsNullOrWhiteSpace idempotencyKey then
RebalanceWriteResult.RebalanceInvalid "idempotency key cannot be empty"
diff --git a/src/FundLab.Domain/FundLab.Domain.fsproj b/src/FundLab.Domain/FundLab.Domain.fsproj
index ef72c53..1ecbd5d 100644
--- a/src/FundLab.Domain/FundLab.Domain.fsproj
+++ b/src/FundLab.Domain/FundLab.Domain.fsproj
@@ -11,6 +11,7 @@
<Compile Include="Redemption.fs" />
<Compile Include="Capital.fs" />
<Compile Include="Sip.fs" />
+ <Compile Include="InvestmentPlan.fs" />
<Compile Include="Rebalance.fs" />
<Compile Include="Dividend.fs" />
<Compile Include="Performance.fs" />
diff --git a/src/FundLab.Domain/InvestmentPlan.fs b/src/FundLab.Domain/InvestmentPlan.fs
new file mode 100644
index 0000000..ecc3425
--- /dev/null
+++ b/src/FundLab.Domain/InvestmentPlan.fs
@@ -0,0 +1,97 @@
+namespace FundLab.Domain
+
+open System
+
+type InvestmentFrequency =
+ | Daily
+ | Weekly
+ | Monthly
+
+module InvestmentPlanPolicy =
+ let cashMaximum = 999999999999999999.99m
+
+ let validateAmount (amount: decimal) : Result<unit, string> =
+ if amount <= 0m then
+ Error "investment plan amount must be positive"
+ elif Decimal.Round(amount, 2) <> amount then
+ Error "investment plan amount exceeds cash precision"
+ elif amount > cashMaximum then
+ Error "investment plan amount exceeds database precision"
+ else
+ Ok()
+
+ let parseFrequency (text: string) : InvestmentFrequency option =
+ match text with
+ | null -> None
+ | "daily" -> Some InvestmentFrequency.Daily
+ | "weekly" -> Some InvestmentFrequency.Weekly
+ | "monthly" -> Some InvestmentFrequency.Monthly
+ | _ -> None
+
+ let frequencyText (frequency: InvestmentFrequency) : string =
+ match frequency with
+ | InvestmentFrequency.Daily -> "daily"
+ | InvestmentFrequency.Weekly -> "weekly"
+ | InvestmentFrequency.Monthly -> "monthly"
+
+ let private isWeekend (date: DateOnly) =
+ date.DayOfWeek = DayOfWeek.Saturday || date.DayOfWeek = DayOfWeek.Sunday
+
+ let private rollToWeekday (date: DateOnly) : DateOnly =
+ let mutable candidate = date
+
+ while isWeekend candidate do
+ candidate <- candidate.AddDays 1
+
+ candidate
+
+ let private addMonthsPreservingDay (anchor: DateOnly) (months: int) : DateOnly =
+ let targetMonthIndex = anchor.Year * 12 + (anchor.Month - 1) + months
+ let year = targetMonthIndex / 12
+ let month = targetMonthIndex % 12 + 1
+ let daysInTargetMonth = DateTime.DaysInMonth(year, month)
+ DateOnly(year, month, min anchor.Day daysInTargetMonth)
+
+ /// First scheduled run date on or after `fromDate`, rolled forward off weekends.
+ /// The anchor date is the plan's first scheduled run date. Calendar approximation:
+ /// only weekends are skipped; public holidays are not modeled yet (no fabrication).
+ let nextRunDate (frequency: InvestmentFrequency) (anchorDate: DateOnly) (fromDate: DateOnly) : DateOnly =
+ let baseDate = if fromDate < anchorDate then anchorDate else fromDate
+
+ match frequency with
+ | InvestmentFrequency.Daily -> rollToWeekday baseDate
+ | InvestmentFrequency.Weekly ->
+ let steps = max 0 ((baseDate.DayNumber - anchorDate.DayNumber + 6) / 7)
+ rollToWeekday (anchorDate.AddDays(steps * 7))
+ | InvestmentFrequency.Monthly ->
+ let monthsBetween = (baseDate.Year - anchorDate.Year) * 12 + (baseDate.Month - anchorDate.Month)
+
+ let candidateMonths =
+ if baseDate.Day <= anchorDate.Day || monthsBetween < 0 then
+ monthsBetween
+ else
+ monthsBetween + 1
+
+ rollToWeekday (addMonthsPreservingDay anchorDate (max 0 candidateMonths))
+
+ /// Due run dates from `dueFrom` through `endDate` inclusive, skipping weekends.
+ /// Empty when nothing is due in the window. This is the "calendar truncation":
+ /// a run never fires past `endDate`, and a repeat drive over the same window
+ /// yields the same dates.
+ let dueDates
+ (frequency: InvestmentFrequency)
+ (anchorDate: DateOnly)
+ (dueFrom: DateOnly)
+ (endDate: DateOnly)
+ : DateOnly list =
+ let rec loop (current: DateOnly) (acc: DateOnly list) =
+ if current > endDate then
+ List.rev acc
+ else
+ loop (nextRunDate frequency anchorDate (current.AddDays 1)) (current :: acc)
+
+ if dueFrom > endDate then
+ []
+ else
+ let start = nextRunDate frequency anchorDate dueFrom
+ if start > endDate then [] else loop start []
diff --git a/src/FundLab.Web/App.fs b/src/FundLab.Web/App.fs
index f1e6df6..1386b01 100644
--- a/src/FundLab.Web/App.fs
+++ b/src/FundLab.Web/App.fs
@@ -235,6 +235,47 @@ type RawSipPlan =
lastExecutionDate: obj
}
+type RawInvestmentPlan =
+ {
+ id: string
+ fundId: string
+ instrumentCode: string
+ amount: string
+ frequency: string
+ status: string
+ anchorDate: string
+ nextRunDate: string
+ lastRunStatus: obj
+ lastRunDate: obj
+ isSynthetic: bool
+ createdAt: string
+ }
+
+type RawInvestmentPlanRunOutcome =
+ {
+ runDate: string
+ status: string
+ orderId: obj
+ pendingReason: obj
+ }
+
+type RawInvestmentPlanRunPlan =
+ {
+ planId: string
+ instrumentCode: string
+ amount: string
+ frequency: string
+ runs: RawInvestmentPlanRunOutcome array
+ nextRunDate: string
+ }
+
+type RawInvestmentPlanRun =
+ {
+ fundId: string
+ processingDate: string
+ plans: RawInvestmentPlanRunPlan array
+ }
+
type RawRedemption =
{
id: string
@@ -393,6 +434,44 @@ type SipPlan =
lastExecutionDate: string option
}
+type InvestmentPlan =
+ {
+ id: string
+ instrumentCode: string
+ amount: string
+ frequency: string
+ status: string
+ anchorDate: string
+ nextRunDate: string
+ lastRunStatus: string option
+ lastRunDate: string option
+ }
+
+type InvestmentPlanRunOutcome =
+ {
+ runDate: string
+ status: string
+ orderId: string option
+ pendingReason: string option
+ }
+
+type InvestmentPlanRunPlan =
+ {
+ planId: string
+ instrumentCode: string
+ amount: string
+ frequency: string
+ runs: InvestmentPlanRunOutcome list
+ nextRunDate: string
+ }
+
+type InvestmentPlanRun =
+ {
+ fundId: string
+ processingDate: string
+ plans: InvestmentPlanRunPlan list
+ }
+
type RebalanceTarget =
{
instrumentCode: string
@@ -576,6 +655,12 @@ module Api =
[<Import("getSipPlans", "./src/api.js")>]
let getSipPlans (token: string) (fundId: string) : JS.Promise<RawSipPlan array> = jsNative
+ [<Import("getInvestmentPlans", "./src/api.js")>]
+ let getInvestmentPlans (token: string) (fundId: string) : JS.Promise<RawInvestmentPlan array> = jsNative
+
+ [<Import("runInvestmentPlans", "./src/api.js")>]
+ let runInvestmentPlans (token: string) (fundId: string) : JS.Promise<RawInvestmentPlanRun> = jsNative
+
let decodeOptionalText (raw: obj) : string option =
if isNull raw then
None
@@ -685,6 +770,44 @@ module Api =
lastExecutionDate = decodeOptionalText raw.lastExecutionDate
}
+ let decodeInvestmentPlan (raw: RawInvestmentPlan) : InvestmentPlan =
+ {
+ id = raw.id
+ instrumentCode = raw.instrumentCode
+ amount = raw.amount
+ frequency = raw.frequency
+ status = raw.status
+ anchorDate = raw.anchorDate
+ nextRunDate = raw.nextRunDate
+ lastRunStatus = decodeOptionalText raw.lastRunStatus
+ lastRunDate = decodeOptionalText raw.lastRunDate
+ }
+
+ let decodeInvestmentPlanRunOutcome (raw: RawInvestmentPlanRunOutcome) : InvestmentPlanRunOutcome =
+ {
+ runDate = raw.runDate
+ status = raw.status
+ orderId = decodeOptionalText raw.orderId
+ pendingReason = decodeOptionalText raw.pendingReason
+ }
+
+ let decodeInvestmentPlanRunPlan (raw: RawInvestmentPlanRunPlan) : InvestmentPlanRunPlan =
+ {
+ planId = raw.planId
+ instrumentCode = raw.instrumentCode
+ amount = raw.amount
+ frequency = raw.frequency
+ runs = raw.runs |> Array.map decodeInvestmentPlanRunOutcome |> Array.toList
+ nextRunDate = raw.nextRunDate
+ }
+
+ let decodeInvestmentPlanRun (raw: RawInvestmentPlanRun) : InvestmentPlanRun =
+ {
+ fundId = raw.fundId
+ processingDate = raw.processingDate
+ plans = raw.plans |> Array.map decodeInvestmentPlanRunPlan |> Array.toList
+ }
+
let decodeRedemption (raw: RawRedemption) : RedemptionDetail =
{
id = raw.id
@@ -812,6 +935,12 @@ type Model =
returnsReadSeq: int
returnsInFlight: bool
returns: FundReturns option
+ planReadSeq: int
+ planInFlight: bool
+ planRunSeq: int
+ planRunInFlight: bool
+ investmentPlans: InvestmentPlan list
+ lastPlanRun: InvestmentPlanRun option
error: string option
}
@@ -893,6 +1022,12 @@ type Msg =
| ReturnsReadRequested
| ReturnsReadCompleted of requestId: int * returns: RawReturns
| ReturnsReadFailed of requestId: int * message: string
+ | InvestmentPlansReadRequested
+ | InvestmentPlansReadCompleted of requestId: int * plans: RawInvestmentPlan array
+ | InvestmentPlansReadFailed of requestId: int * message: string
+ | InvestmentPlansRunRequested
+ | InvestmentPlansRunCompleted of requestId: int * fundId: string * run: RawInvestmentPlanRun
+ | InvestmentPlansRunFailed of requestId: int * fundId: string * message: string
let defaultInitialUnitNav = "1.00000000"
@@ -994,6 +1129,12 @@ let init () =
returnsReadSeq = 0
returnsInFlight = false
returns = None
+ planReadSeq = 0
+ planInFlight = false
+ planRunSeq = 0
+ planRunInFlight = false
+ investmentPlans = []
+ lastPlanRun = None
error = None
}
@@ -1143,6 +1284,20 @@ let private readReturnsCommand token fundId requestId =
(fun returns -> ReturnsReadCompleted(requestId, returns))
(fun error -> ReturnsReadFailed(requestId, errorText error))
+let private readInvestmentPlansCommand token fundId requestId =
+ Cmd.OfPromise.either
+ (fun () -> Api.getInvestmentPlans token fundId)
+ ()
+ (fun plans -> InvestmentPlansReadCompleted(requestId, plans))
+ (fun error -> InvestmentPlansReadFailed(requestId, errorText error))
+
+let private runInvestmentPlansCommand token fundId requestId =
+ Cmd.OfPromise.either
+ (fun () -> Api.runInvestmentPlans token fundId)
+ ()
+ (fun run -> InvestmentPlansRunCompleted(requestId, fundId, run))
+ (fun error -> InvestmentPlansRunFailed(requestId, fundId, errorText error))
+
let update message model =
match message with
| TokenChanged token ->
@@ -1217,6 +1372,12 @@ let update message model =
returnsReadSeq = model.returnsReadSeq + 1
returnsInFlight = false
returns = None
+ planReadSeq = model.planReadSeq + 1
+ planInFlight = false
+ planRunSeq = model.planRunSeq + 1
+ planRunInFlight = false
+ investmentPlans = []
+ lastPlanRun = None
error = None
},
Cmd.none
@@ -1417,9 +1578,15 @@ let update message model =
returnsReadSeq = model.returnsReadSeq + 1
returnsInFlight = false
returns = None
+ planReadSeq = model.planReadSeq + 1
+ planInFlight = false
+ planRunSeq = model.planRunSeq + 1
+ planRunInFlight = false
+ investmentPlans = []
+ lastPlanRun = None
error = None
},
- Cmd.batch [ Cmd.ofMsg OrdersReadRequested; Cmd.ofMsg ReturnsReadRequested ]
+ Cmd.batch [ Cmd.ofMsg OrdersReadRequested; Cmd.ofMsg ReturnsReadRequested; Cmd.ofMsg InvestmentPlansReadRequested ]
else
model, Cmd.none
| FundCreateFailed (requestId, message) ->
@@ -2120,7 +2287,79 @@ let update message model =
{ model with returnsInFlight = false; error = Some message }, Cmd.none
else
model, Cmd.none
+ | InvestmentPlansReadRequested ->
+ match model.createdFund with
+ | Some fund when not (String.IsNullOrWhiteSpace model.token) ->
+ let requestId = model.planReadSeq + 1
+
+ {
+ model with
+ planReadSeq = requestId
+ planInFlight = true
+ error = None
+ },
+ readInvestmentPlansCommand model.token fund.id requestId
+ | Some _ -> { model with error = Some "请输入 API token" }, Cmd.none
+ | None -> model, Cmd.none
+ | InvestmentPlansReadCompleted (requestId, plans) ->
+ if requestId = model.planReadSeq then
+ {
+ model with
+ investmentPlans = plans |> Array.map Api.decodeInvestmentPlan |> Array.toList
+ planInFlight = false
+ error = None
+ },
+ Cmd.none
+ else
+ model, Cmd.none
+ | InvestmentPlansReadFailed (requestId, message) ->
+ if requestId = model.planReadSeq then
+ { model with planInFlight = false; error = Some message }, Cmd.none
+ else
+ model, Cmd.none
+ | InvestmentPlansRunRequested ->
+ match model.createdFund with
+ | Some fund when not (String.IsNullOrWhiteSpace model.token) ->
+ if model.planRunInFlight then
+ model, Cmd.none
+ else
+ let requestId = model.planRunSeq + 1
+ {
+ model with
+ planRunSeq = requestId
+ planRunInFlight = true
+ error = None
+ },
+ runInvestmentPlansCommand model.token fund.id requestId
+ | Some _ -> { model with error = Some "请输入 API token" }, Cmd.none
+ | None -> model, Cmd.none
+ | InvestmentPlansRunCompleted (requestId, fundId, run) ->
+ if requestId = model.planRunSeq
+ && (match model.createdFund with Some fund -> fund.id = fundId | None -> false) then
+ let decoded = Api.decodeInvestmentPlanRun run
+
+ {
+ model with
+ planRunInFlight = false
+ lastPlanRun = Some decoded
+ error = None
+ },
+ Cmd.batch
+ [
+ Cmd.ofMsg InvestmentPlansReadRequested
+ Cmd.ofMsg FundReadRequested
+ Cmd.ofMsg PositionsReadRequested
+ Cmd.ofMsg ReturnsReadRequested
+ ]
+ else
+ model, Cmd.none
+ | InvestmentPlansRunFailed (requestId, fundId, message) ->
+ if requestId = model.planRunSeq
+ && (match model.createdFund with Some fund -> fund.id = fundId | None -> false) then
+ { model with planRunInFlight = false; error = Some message }, Cmd.none
+ else
+ model, Cmd.none
let private navText (text: string) =
match Decimal.TryParse(text, NumberStyles.Float, CultureInfo.InvariantCulture) with
@@ -3389,6 +3628,122 @@ let private returnsPanel model dispatch =
]
]
+let private investmentPlanFrequencyText (frequency: string) =
+ if frequency = "daily" then "每日"
+ elif frequency = "weekly" then "每周"
+ elif frequency = "monthly" then "每月"
+ else frequency
+
+let private investmentPlanRunStatusText (status: string) =
+ if status = "succeeded" then "已执行"
+ elif status = "pending_nav" then "待净值"
+ elif status = "insufficient_cash" then "现金不足"
+ elif status = "failed" then "失败"
+ else status
+
+let private investmentPlanRow (plan: InvestmentPlan) =
+ Html.div [
+ prop.className "order-row"
+ prop.children [
+ Html.span [ prop.className "order-code"; prop.text plan.instrumentCode ]
+ Html.span [ prop.className "order-cell"; prop.text (sprintf "金额 %s" plan.amount) ]
+ Html.span [ prop.className "order-cell"; prop.text (sprintf "频率 %s" (investmentPlanFrequencyText 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.nextRunDate) ]
+
+ let lastRun =
+ match plan.lastRunStatus with
+ | Some status ->
+ sprintf "最近执行 %s(%s)" (investmentPlanRunStatusText status) (plan.lastRunDate |> Option.defaultValue "—")
+ | None -> "最近执行 暂无"
+
+ Html.span [ prop.className "order-cell"; prop.text lastRun ]
+ Html.span [ prop.className "order-status"; prop.text (if plan.status = "active" then "进行中" else plan.status) ]
+ ]
+ ]
+
+let private investmentPlanRunRow (plan: InvestmentPlanRunPlan) =
+ Html.div [
+ prop.className "order-row"
+ prop.children [
+ Html.span [ prop.className "order-code"; prop.text plan.instrumentCode ]
+ Html.span [ prop.className "order-cell"; prop.text (sprintf "金额 %s" plan.amount) ]
+ Html.span [ prop.className "order-cell"; prop.text (sprintf "频率 %s" (investmentPlanFrequencyText plan.frequency)) ]
+
+ let runs =
+ if List.isEmpty plan.runs then
+ "本次无到期扣款"
+ else
+ plan.runs
+ |> List.map (fun run -> sprintf "%s %s" run.runDate (investmentPlanRunStatusText run.status))
+ |> String.concat "、"
+
+ Html.span [ prop.className "order-cell"; prop.text runs ]
+ Html.span [ prop.className "order-cell"; prop.text (sprintf "下次执行 %s" plan.nextRunDate) ]
+ ]
+ ]
+
+let private investmentPlansPanel model dispatch =
+ Html.section [
+ prop.className "panel plans-panel"
+ prop.children [
+ Html.div [
+ prop.className "section-heading"
+ prop.children [
+ Html.div [
+ Html.p [ prop.className "eyebrow"; prop.text "11 / PLANS" ]
+ Html.h2 "定投计划执行"
+ ]
+ Html.span [ prop.className "section-note"; prop.text "scheduled contribution" ]
+ ]
+ ]
+ Html.p [
+ prop.className "hint"
+ prop.text "执行到期计划:按当日净值扣款买入,重复执行不产生第二笔;存款与买入都不计入收益。"
+ ]
+ Html.div [
+ prop.className "sip-plans plans-list"
+ prop.children [
+ if List.isEmpty model.investmentPlans then
+ Html.p [ prop.className "hint"; prop.text "暂无定投计划" ]
+ else
+ yield! (model.investmentPlans |> List.map investmentPlanRow)
+ ]
+ ]
+ match model.lastPlanRun with
+ | Some run ->
+ Html.div [
+ prop.className "plans-run-result"
+ prop.children [
+ Html.p [ prop.className "hint"; prop.text (sprintf "最近执行处理日 %s" run.processingDate) ]
+
+ if List.isEmpty run.plans then
+ Html.p [ prop.className "hint"; prop.text "本次没有可执行的计划" ]
+ else
+ yield! (run.plans |> List.map investmentPlanRunRow)
+ ]
+ ]
+ | None -> Html.none
+ Html.div [
+ prop.className "panel-actions"
+ prop.children [
+ Html.button [
+ prop.className "primary-action plans-run-action"
+ prop.disabled model.planRunInFlight
+ prop.onClick (fun _ -> dispatch InvestmentPlansRunRequested)
+ prop.text ((if model.planRunInFlight then "执行中..." else "执行到期计划"): string)
+ ]
+ Html.button [
+ prop.className "secondary-action plans-refresh-action"
+ prop.disabled model.planInFlight
+ prop.onClick (fun _ -> dispatch InvestmentPlansReadRequested)
+ prop.text ((if model.planInFlight then "读取中..." else "刷新计划"): string)
+ ]
+ ]
+ ]
+ ]
+ ]
+
let view model dispatch =
Html.main [
prop.className "app-shell"
@@ -3451,6 +3806,7 @@ let view model dispatch =
rebalancePanel model dispatch
dividendPanel model dispatch
returnsPanel model dispatch
+ investmentPlansPanel model dispatch
Html.footer [ prop.className "footer-note"; prop.text "SOURCE · AKShare / STORAGE · PostgreSQL / LEDGER · CREATE & READ & SUBSCRIBE" ]
]
]
diff --git a/src/FundLab.Web/src/api.js b/src/FundLab.Web/src/api.js
index 7606e02..4f3f47c 100644
--- a/src/FundLab.Web/src/api.js
+++ b/src/FundLab.Web/src/api.js
@@ -143,6 +143,18 @@ export function getSipPlans(token, fundId) {
return requestJson(`/api/funds/${encodeURIComponent(fundId)}/sip/plans`, token);
}
+export function getInvestmentPlans(token, fundId) {
+ return requestJson(`/api/funds/${encodeURIComponent(fundId)}/investment-plans`, token);
+}
+
+export function runInvestmentPlans(token, fundId) {
+ return requestJson(`/api/funds/${encodeURIComponent(fundId)}/investment-plans/run`, token, {
+ method: "POST",
+ headers: { "Content-Type": "application/json" },
+ body: "{}"
+ });
+}
+
export function createRebalancePlan(token, fundId, payload) {
const body = `{"targets":${JSON.stringify(payload.targets)}}`;
return requestJson(`/api/funds/${encodeURIComponent(fundId)}/rebalance/plans`, token, {
diff --git a/src/FundLab.Web/src/styles.css b/src/FundLab.Web/src/styles.css
index 016178d..de00388 100644
--- a/src/FundLab.Web/src/styles.css
+++ b/src/FundLab.Web/src/styles.css
@@ -554,6 +554,12 @@ h2 {
height: 300px;
}
+.plans-panel .plans-run-result {
+ margin-top: 10px;
+ padding-top: 8px;
+ border-top: 1px dashed #e2e8f0;
+}
+
.footer-note {
padding: 16px 0 34px;
font-family: "SFMono-Regular", Consolas, monospace;
diff --git a/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj b/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj
index 0024deb..ca5bde8 100644
--- a/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj
+++ b/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj
@@ -27,6 +27,7 @@
<Compile Include="SipAdvanceTests.fs" />
<Compile Include="DividendTests.fs" />
<Compile Include="ReturnsTests.fs" />
+ <Compile Include="InvestmentPlanTests.fs" />
<Compile Include="Program.fs" />
</ItemGroup>
</Project>
diff --git a/tests/FundLab.Api.Tests/InvestmentPlanTests.fs b/tests/FundLab.Api.Tests/InvestmentPlanTests.fs
new file mode 100644
index 0000000..77a90ab
--- /dev/null
+++ b/tests/FundLab.Api.Tests/InvestmentPlanTests.fs
@@ -0,0 +1,322 @@
+namespace FundLab.Api.Tests
+
+open System
+open System.Text.Json
+open Xunit
+open FundLab.Api
+open FundLab.Domain
+
+[<Collection("postgres")>]
+type InvestmentPlanTests(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, "investment-plan-test-hash")
+ code
+
+ let createFund (initialCash: decimal) =
+ let command =
+ {
+ Name = "定投测试 FOF"
+ InitialCash = initialCash
+ InitialUnitNav = 1.00000000m
+ IsSynthetic = true
+ }
+
+ let key = fixture.Key(sprintf "investment-plan-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 "investment-plan-hash/%s" revision)
+
+ let createPlanViaApi fundId (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/investment-plans" fundId)
+ [
+ "Authorization", "Bearer test-token"
+ "Idempotency-Key", fixture.Key(sprintf "investment-plan-%s" (Guid.NewGuid().ToString("N")))
+ ]
+ body
+
+ let createPlanWithKey fundId (code: string) (amount: string) (frequency: string) idempotencyKey =
+ let body =
+ sprintf "{\"instrumentCode\":\"%s\",\"amount\":\"%s\",\"frequency\":\"%s\"}" code amount frequency
+
+ PersistenceTestHelpers.invoke
+ (app ())
+ "POST"
+ (sprintf "/api/funds/%O/investment-plans" fundId)
+ [ "Authorization", "Bearer test-token"; "Idempotency-Key", idempotencyKey ]
+ body
+
+ let listPlansViaApi fundId =
+ PersistenceTestHelpers.invoke
+ (app ())
+ "GET"
+ (sprintf "/api/funds/%O/investment-plans" fundId)
+ [ "Authorization", "Bearer test-token" ]
+ ""
+
+ let runPlansViaApi fundId =
+ PersistenceTestHelpers.invoke
+ (app ())
+ "POST"
+ (sprintf "/api/funds/%O/investment-plans/run" fundId)
+ [ "Authorization", "Bearer test-token" ]
+ ""
+
+ let runPlansOnViaApi fundId (processingDate: DateOnly) =
+ PersistenceTestHelpers.invoke
+ (app ())
+ "POST"
+ (sprintf "/api/funds/%O/investment-plans/run" fundId)
+ [ "Authorization", "Bearer test-token" ]
+ (sprintf "{\"processingDate\":\"%s\"}" (processingDate.ToString("yyyy-MM-dd")))
+
+ let createPlanWithAnchor (fundId: Guid) (code: string) (amount: decimal) (frequency: InvestmentFrequency) (anchor: DateOnly) =
+ let key = fixture.Key(sprintf "investment-plan-anchor-%s" (Guid.NewGuid().ToString("N")))
+
+ match repository().CreateInvestmentPlan(key, fundId, { InstrumentCode = code; Amount = amount; Frequency = frequency }, anchor) with
+ | InvestmentPlanWriteResult.InvestmentPlanCreated plan -> plan.Id
+ | other -> failwithf "unexpected investment plan result: %A" other
+
+ let today = ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow
+
+ let availableCash fundId =
+ match repository().GetFund fundId with
+ | Some fund -> fund.AvailableCash
+ | None -> failwith "fund was not found"
+
+ let unitsFor fundId (code: string) =
+ repository().GetFundPositions fundId
+ |> List.tryFind (fun position -> position.InstrumentCode = code)
+ |> Option.map (fun position -> position.Units)
+ |> Option.defaultValue 0m
+
+ let planRuns (body: string) =
+ use document = JsonDocument.Parse(body)
+
+ document.RootElement.GetProperty("plans").EnumerateArray()
+ |> Seq.collect (fun plan -> plan.GetProperty("runs").EnumerateArray())
+ |> Seq.map (fun run -> run.Clone())
+ |> Seq.toList
+
+ [<Fact>]
+ member _.``create is idempotent, lists the plan and rejects an unknown instrument``() =
+ let fundId = createFund 10000.00m
+ let code = seedInstrument ()
+ let key = fixture.Key(sprintf "investment-plan-%s" (Guid.NewGuid().ToString("N")))
+
+ let firstStatus, firstBody = createPlanWithKey fundId code "100.00" "daily" key
+ Assert.Equal(201, firstStatus)
+
+ use firstDocument = JsonDocument.Parse(firstBody)
+ let planId = firstDocument.RootElement.GetProperty("id").GetString()
+ Assert.Equal(code, firstDocument.RootElement.GetProperty("instrumentCode").GetString())
+ Assert.Equal("100.00", firstDocument.RootElement.GetProperty("amount").GetString())
+ Assert.Equal("daily", firstDocument.RootElement.GetProperty("frequency").GetString())
+ Assert.Equal(today.ToString("yyyy-MM-dd"), firstDocument.RootElement.GetProperty("nextRunDate").GetString())
+
+ let replayStatus, replayBody = createPlanWithKey fundId code "100.00" "daily" key
+ Assert.Equal(200, replayStatus)
+
+ use replayDocument = JsonDocument.Parse(replayBody)
+ Assert.Equal(planId, replayDocument.RootElement.GetProperty("id").GetString())
+
+ let listStatus, listBody = listPlansViaApi fundId
+ Assert.Equal(200, listStatus)
+ Assert.Contains(planId, listBody)
+
+ let missingStatus, missingBody = createPlanViaApi fundId "999999" "100.00" "daily"
+ Assert.Equal(404, missingStatus)
+ Assert.Contains("INSTRUMENT_NOT_FOUND", missingBody)
+
+ [<Fact>]
+ member _.``create rejects a bad amount and a bad frequency``() =
+ let fundId = createFund 10000.00m
+ let code = seedInstrument ()
+
+ let badAmountStatus, badAmountBody = createPlanViaApi fundId code "0.00" "daily"
+ Assert.Equal(400, badAmountStatus)
+ Assert.Contains("INVALID_INVESTMENT_PLAN_REQUEST", badAmountBody)
+
+ let badFrequencyStatus, badFrequencyBody = createPlanViaApi fundId code "100.00" "hourly"
+ Assert.Equal(400, badFrequencyStatus)
+ Assert.Contains("frequency must be one of daily, weekly or monthly", badFrequencyBody)
+
+ [<Fact>]
+ member _.``run executes a due plan through the shared order pipeline``() =
+ let fundId = createFund 10000.00m
+ let code = seedInstrument ()
+ insertQuoteOnDate code 2.0m today
+ let _ = createPlanViaApi fundId code "100.00" "daily"
+
+ let status, body = runPlansViaApi fundId
+ Assert.Equal(200, status)
+
+ let runs = planRuns body
+ Assert.Single(runs) |> ignore
+ Assert.Equal("succeeded", runs.Head.GetProperty("status").GetString())
+ Assert.Equal(today.ToString("yyyy-MM-dd"), runs.Head.GetProperty("runDate").GetString())
+
+ Assert.Equal(50.00000000m, unitsFor fundId code)
+ Assert.Equal(9900.00m, availableCash fundId)
+
+ let returnsStatus, returnsBody =
+ PersistenceTestHelpers.invoke
+ (app ())
+ "GET"
+ (sprintf "/api/funds/%O/returns" fundId)
+ [ "Authorization", "Bearer test-token" ]
+ ""
+
+ Assert.Equal(200, returnsStatus)
+ Assert.Contains("\"cumulativeReturn\":\"0.00\"", returnsBody)
+
+ [<Fact>]
+ member _.``run is idempotent across repeated drives``() =
+ let fundId = createFund 10000.00m
+ let code = seedInstrument ()
+ insertQuoteOnDate code 2.0m today
+ let _ = createPlanViaApi fundId code "100.00" "daily"
+
+ let firstStatus, _ = runPlansViaApi fundId
+ Assert.Equal(200, firstStatus)
+
+ let secondStatus, secondBody = runPlansViaApi fundId
+ Assert.Equal(200, secondStatus)
+ Assert.Single(planRuns secondBody) |> ignore
+
+ Assert.Equal(50.00000000m, unitsFor fundId code)
+ Assert.Equal(9900.00m, availableCash fundId)
+
+ [<Fact>]
+ member _.``run defers a plan when nav is missing and retries once it lands``() =
+ let fundId = createFund 10000.00m
+ let code = seedInstrument ()
+ let _ = createPlanViaApi fundId code "100.00" "daily"
+
+ let pendingStatus, pendingBody = runPlansViaApi fundId
+ Assert.Equal(200, pendingStatus)
+
+ let pendingRuns = planRuns pendingBody
+ Assert.Single(pendingRuns) |> ignore
+ Assert.Equal("pending_nav", pendingRuns.Head.GetProperty("status").GetString())
+ // the pending order freezes the scheduled amount without spending it
+ Assert.Equal(9900.00m, availableCash fundId)
+ Assert.Equal(0m, unitsFor fundId code)
+
+ insertQuoteOnDate code 2.0m today
+
+ let retryStatus, retryBody = runPlansViaApi fundId
+ Assert.Equal(200, retryStatus)
+
+ let retryRuns = planRuns retryBody
+ Assert.Single(retryRuns) |> ignore
+ Assert.Equal("succeeded", retryRuns.Head.GetProperty("status").GetString())
+ Assert.Equal(50.00000000m, unitsFor fundId code)
+ Assert.Equal(9900.00m, availableCash fundId)
+
+ [<Fact>]
+ member _.``run records insufficient cash without debiting twice``() =
+ let fundId = createFund 50.00m
+ let code = seedInstrument ()
+ insertQuoteOnDate code 2.0m today
+ let _ = createPlanViaApi fundId code "100.00" "daily"
+
+ let firstStatus, firstBody = runPlansViaApi fundId
+ Assert.Equal(200, firstStatus)
+
+ let firstRuns = planRuns firstBody
+ Assert.Single(firstRuns) |> ignore
+ Assert.Equal("insufficient_cash", firstRuns.Head.GetProperty("status").GetString())
+ Assert.Equal(50.00m, availableCash fundId)
+ Assert.Equal(0m, unitsFor fundId code)
+
+ let secondStatus, secondBody = runPlansViaApi fundId
+ Assert.Equal(200, secondStatus)
+ Assert.Single(planRuns secondBody) |> ignore
+ Assert.Equal(50.00m, availableCash fundId)
+
+ [<Fact>]
+ member _.``run rolls the plan pointer across multiple due dates``() =
+ let fundId = createFund 10000.00m
+ let code = seedInstrument ()
+ insertQuoteOnDate code 2.0m (DateOnly(2026, 9, 15))
+ insertQuoteOnDate code 2.0m (DateOnly(2026, 9, 16))
+ insertQuoteOnDate code 2.0m (DateOnly(2026, 9, 17))
+ let _ = createPlanWithAnchor fundId code 100.00m InvestmentFrequency.Daily (DateOnly(2026, 9, 15))
+
+ let status, body = runPlansOnViaApi fundId (DateOnly(2026, 9, 17))
+ Assert.Equal(200, status)
+
+ let runs = planRuns body
+ Assert.Equal(3, runs.Length)
+ Assert.All(runs, fun run -> Assert.Equal("succeeded", run.GetProperty("status").GetString()))
+ Assert.Equal(150.00000000m, unitsFor fundId code)
+ Assert.Equal(9700.00m, availableCash fundId)
+
+ use document = JsonDocument.Parse(body)
+
+ let nextRunDate =
+ document.RootElement.GetProperty("plans").EnumerateArray()
+ |> Seq.head
+ |> fun plan -> plan.GetProperty("nextRunDate").GetString()
+
+ Assert.Equal("2026-09-18", nextRunDate)
+
+ [<Fact>]
+ member _.``run for an unknown fund is not found``() =
+ let status, body = runPlansViaApi (Guid.NewGuid())
+ Assert.Equal(404, status)
+ Assert.Contains("FUND_NOT_FOUND", body)
diff --git a/tests/FundLab.Domain.Tests/FundLab.Domain.Tests.fsproj b/tests/FundLab.Domain.Tests/FundLab.Domain.Tests.fsproj
index 2ba20b5..a9c6fc0 100644
--- a/tests/FundLab.Domain.Tests/FundLab.Domain.Tests.fsproj
+++ b/tests/FundLab.Domain.Tests/FundLab.Domain.Tests.fsproj
@@ -20,6 +20,7 @@
</ItemGroup>
<ItemGroup>
<Compile Include="DomainTests.fs" />
+ <Compile Include="InvestmentPlanTests.fs" />
<Compile Include="Program.fs" />
</ItemGroup>
</Project>
diff --git a/tests/FundLab.Domain.Tests/InvestmentPlanTests.fs b/tests/FundLab.Domain.Tests/InvestmentPlanTests.fs
new file mode 100644
index 0000000..9a5e2b9
--- /dev/null
+++ b/tests/FundLab.Domain.Tests/InvestmentPlanTests.fs
@@ -0,0 +1,81 @@
+namespace FundLab.Domain.Tests
+
+module InvestmentPlanTests =
+
+ open System
+ open Xunit
+ open FundLab.Domain
+
+ let private anchor = DateOnly(2026, 9, 21)
+
+ [<Fact>]
+ let ``frequency parsing and text roundtrip`` () =
+ Assert.Equal(Some InvestmentFrequency.Daily, InvestmentPlanPolicy.parseFrequency "daily")
+ Assert.Equal(Some InvestmentFrequency.Weekly, InvestmentPlanPolicy.parseFrequency "weekly")
+ Assert.Equal(Some InvestmentFrequency.Monthly, InvestmentPlanPolicy.parseFrequency "monthly")
+ Assert.Equal(None, InvestmentPlanPolicy.parseFrequency "biweekly")
+ Assert.Equal(None, InvestmentPlanPolicy.parseFrequency null)
+ Assert.Equal("daily", InvestmentPlanPolicy.frequencyText InvestmentFrequency.Daily)
+ Assert.Equal("weekly", InvestmentPlanPolicy.frequencyText InvestmentFrequency.Weekly)
+ Assert.Equal("monthly", InvestmentPlanPolicy.frequencyText InvestmentFrequency.Monthly)
+
+ [<Fact>]
+ let ``amount validation requires positive two-decimal cash`` () =
+ Assert.Equal(Ok(), InvestmentPlanPolicy.validateAmount 200.00m)
+ Assert.Equal(Error "investment plan amount must be positive", InvestmentPlanPolicy.validateAmount 0m)
+ Assert.Equal(Error "investment plan amount must be positive", InvestmentPlanPolicy.validateAmount -1.00m)
+ Assert.Equal(Error "investment plan amount exceeds cash precision", InvestmentPlanPolicy.validateAmount 1.005m)
+
+ [<Fact>]
+ let ``daily schedule runs every weekday and skips weekends`` () =
+ Assert.Equal(DateOnly(2026, 9, 21), InvestmentPlanPolicy.nextRunDate InvestmentFrequency.Daily anchor anchor)
+ Assert.Equal(DateOnly(2026, 9, 22), InvestmentPlanPolicy.nextRunDate InvestmentFrequency.Daily anchor (DateOnly(2026, 9, 22)))
+ // Friday -> next scheduled run is Monday
+ Assert.Equal(DateOnly(2026, 9, 28), InvestmentPlanPolicy.nextRunDate InvestmentFrequency.Daily anchor (DateOnly(2026, 9, 26)))
+
+ let due =
+ InvestmentPlanPolicy.dueDates InvestmentFrequency.Daily (DateOnly(2026, 9, 25)) (DateOnly(2026, 9, 25)) (DateOnly(2026, 9, 28))
+
+ Assert.Equal<DateOnly list>([ DateOnly(2026, 9, 25); DateOnly(2026, 9, 28) ], due)
+
+ [<Fact>]
+ let ``weekly schedule keeps the seven day cadence`` () =
+ Assert.Equal(DateOnly(2026, 9, 21), InvestmentPlanPolicy.nextRunDate InvestmentFrequency.Weekly anchor anchor)
+ Assert.Equal(DateOnly(2026, 9, 28), InvestmentPlanPolicy.nextRunDate InvestmentFrequency.Weekly anchor (DateOnly(2026, 9, 22)))
+
+ let due =
+ InvestmentPlanPolicy.dueDates InvestmentFrequency.Weekly anchor anchor (DateOnly(2026, 10, 5))
+
+ Assert.Equal<DateOnly list>(
+ [ DateOnly(2026, 9, 21); DateOnly(2026, 9, 28); DateOnly(2026, 10, 5) ],
+ due
+ )
+
+ [<Fact>]
+ let ``monthly schedule clamps to month end and rolls off weekends`` () =
+ // 2026-02-15 is a Sunday, so the February run rolls to Monday 2026-02-16
+ Assert.Equal(DateOnly(2026, 2, 16), InvestmentPlanPolicy.nextRunDate InvestmentFrequency.Monthly (DateOnly(2026, 1, 15)) (DateOnly(2026, 2, 1)))
+
+ // 2026-02 has 28 days, so a 31st anchor clamps to the 28th then rolls to Monday 2026-03-02
+ Assert.Equal(DateOnly(2026, 3, 2), InvestmentPlanPolicy.nextRunDate InvestmentFrequency.Monthly (DateOnly(2026, 1, 31)) (DateOnly(2026, 2, 1)))
+
+ let due =
+ InvestmentPlanPolicy.dueDates InvestmentFrequency.Monthly anchor anchor (DateOnly(2026, 11, 30))
+
+ // 2026-11-21 is a Saturday, so the November run rolls to Monday 2026-11-23
+ Assert.Equal<DateOnly list>(
+ [ DateOnly(2026, 9, 21); DateOnly(2026, 10, 21); DateOnly(2026, 11, 23) ],
+ due
+ )
+
+ [<Fact>]
+ let ``due dates are truncated at the processing window`` () =
+ Assert.Equal<DateOnly list>([], InvestmentPlanPolicy.dueDates InvestmentFrequency.Daily anchor anchor (DateOnly(2026, 9, 20)))
+
+ let truncated =
+ InvestmentPlanPolicy.dueDates InvestmentFrequency.Daily anchor anchor (DateOnly(2026, 9, 22))
+
+ Assert.Equal<DateOnly list>([ DateOnly(2026, 9, 21); DateOnly(2026, 9, 22) ], truncated)
+
+ // a repeat drive over the same window yields the same dates (idempotent schedule)
+ Assert.Equal<DateOnly list>(truncated, InvestmentPlanPolicy.dueDates InvestmentFrequency.Daily anchor anchor (DateOnly(2026, 9, 22)))