diff options
| author | Somhairle H. Marisol <[email protected]> | 2026-09-21 23:21:54 +0800 |
|---|---|---|
| committer | Somhairle H. Marisol <[email protected]> | 2026-09-21 23:21:54 +0800 |
| commit | d4b539a2791a0cd4d098fecfece0039073a9f7e2 (patch) | |
| tree | dfd314b37fa24379ef13dcd66a1de76e2a923bf2 /src/FundLab.Api | |
| parent | 12c4d3625458c830a2746a1f277fb684e9c498cb (diff) | |
| download | fund-lab-d4b539a2791a0cd4d098fecfece0039073a9f7e2.tar.gz | |
Add scheduled investment plans (3d-10)
Diffstat (limited to 'src/FundLab.Api')
| -rw-r--r-- | src/FundLab.Api/App.fs | 208 | ||||
| -rw-r--r-- | src/FundLab.Api/Persistence.fs | 602 |
2 files changed, 809 insertions, 1 deletions
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" |
