diff options
| author | Somhairle H. Marisol <[email protected]> | 2026-09-21 21:47:59 +0800 |
|---|---|---|
| committer | Somhairle H. Marisol <[email protected]> | 2026-09-21 21:47:59 +0800 |
| commit | 722ec6c19c0c64ea5348c6bc3aefc0061bd96307 (patch) | |
| tree | 65fc127dcb808e05ab332d15332b5cb69ccf6620 /src/FundLab.Api/Persistence.fs | |
| parent | 98c41c37896a41c5f2f15ed187e55ab617735342 (diff) | |
| download | fund-lab-722ec6c19c0c64ea5348c6bc3aefc0061bd96307.tar.gz | |
Add dividend slice (3d-8)
Book cash dividends to available cash exactly once and reinvest dividends through the shared subscription pipeline, with deterministic scheme idempotency and cross-position isolation.
Diffstat (limited to 'src/FundLab.Api/Persistence.fs')
| -rw-r--r-- | src/FundLab.Api/Persistence.fs | 483 |
1 files changed, 479 insertions, 4 deletions
diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs index 6c0374a..e252d0a 100644 --- a/src/FundLab.Api/Persistence.fs +++ b/src/FundLab.Api/Persistence.fs @@ -40,11 +40,24 @@ module ConfirmationPolicy = let isModeledTradingDay (date: DateOnly) : bool = not (isWeekend date) + /// Test-only clock anchor: modules may pin the trading-date clock by exporting + /// FUND_LAB_TEST_TRADE_DATE so suites stay independent of the wall clock. Unset in + /// production, the real Shanghai rule applies unchanged. let tradeDateFor (submittedAt: DateTimeOffset) : DateOnly = - let local = TimeZoneInfo.ConvertTime(submittedAt, shanghaiZone.Value) - let date = DateOnly.FromDateTime(local.Date) - let candidate = if local.TimeOfDay >= cutoffTimeOfDay then date.AddDays 1 else date - rollToWeekday candidate + match Environment.GetEnvironmentVariable("FUND_LAB_TEST_TRADE_DATE") with + | value when not (String.IsNullOrWhiteSpace value) -> + match DateOnly.TryParseExact(value, "yyyy-MM-dd", CultureInfo.InvariantCulture, DateTimeStyles.None) with + | true, anchored -> anchored + | _ -> + let local = TimeZoneInfo.ConvertTime(submittedAt, shanghaiZone.Value) + let date = DateOnly.FromDateTime(local.Date) + let candidate = if local.TimeOfDay >= cutoffTimeOfDay then date.AddDays 1 else date + rollToWeekday candidate + | _ -> + let local = TimeZoneInfo.ConvertTime(submittedAt, shanghaiZone.Value) + let date = DateOnly.FromDateTime(local.Date) + let candidate = if local.TimeOfDay >= cutoffTimeOfDay then date.AddDays 1 else date + rollToWeekday candidate type NavQuote = { NavDate: DateOnly @@ -353,6 +366,50 @@ type RebalanceExecutionResult = Outcomes: RebalanceOrderOutcome list } +type DividendMode = DividendPolicy.DividendMode + +type DividendCommand = + { + InstrumentCode: string + NavDate: DateOnly + Dps: decimal + Mode: DividendMode + } + +type DividendRecord = + { + Id: Guid + FundId: Guid + InstrumentCode: string + NavDate: DateOnly + Dps: decimal + Mode: DividendMode + Status: string + IsSynthetic: bool + GrossCash: decimal option + CreditedUnits: decimal option + CreditedInvested: decimal option + OrderId: Guid option + PendingReason: string option + CreatedAt: DateTimeOffset + } + +type DividendWriteResult = + | DividendCredited of DividendRecord + | DividendReplayed of DividendRecord + | DividendPendingReinvest of DividendRecord + | DividendIdempotencyConflict + | DividendInvalid of string + | DividendFundNotFound + | DividendInstrumentNotFound + | DividendNoHoldings + +type DividendStageDecision = + | StageReplay of Guid + | StageFail of DividendWriteResult + | StageBooked of DividendRecord + + type CapitalDepositCommand = { Amount: decimal @@ -626,6 +683,31 @@ type FundRepository(connectionString: string) = target_percent numeric(9, 2) NOT NULL CHECK (target_percent > 0 AND target_percent <= 100), PRIMARY KEY (plan_id, instrument_code) ); + + CREATE TABLE IF NOT EXISTS dividend_records ( + id uuid PRIMARY KEY, + fund_id uuid NOT NULL REFERENCES funds(id), + instrument_code text NOT NULL REFERENCES instruments(code), + nav_date date NOT NULL, + dps numeric(28, 8) NOT NULL CHECK (dps > 0), + mode text NOT NULL, + status text NOT NULL, + is_synthetic boolean NOT NULL, + gross_cash numeric(20, 2) NULL, + credited_units numeric(28, 8) NULL, + credited_invested numeric(20, 2) NULL, + order_id uuid NULL, + pending_reason text NULL, + created_at timestamptz NOT NULL DEFAULT now() + ); + + CREATE TABLE IF NOT EXISTS dividend_idempotencies ( + idempotency_key text PRIMARY KEY, + request_hash text NOT NULL, + record_id uuid NOT NULL REFERENCES dividend_records(id), + fund_id uuid NOT NULL REFERENCES funds(id), + created_at timestamptz NOT NULL DEFAULT now() + ); """ let statusText status = @@ -1703,6 +1785,155 @@ type FundRepository(connectionString: string) = match RebalancePolicy.validateTargets command.Targets with | Error message -> Error message | Ok() -> Ok() + let dividendRecordFromReader (reader: DbDataReader) : DividendRecord = + { + Id = reader.GetGuid(0) + FundId = reader.GetGuid(1) + InstrumentCode = reader.GetString(2) + NavDate = reader.GetFieldValue<DateOnly>(3) + Dps = reader.GetDecimal(4) + Mode = + match DividendPolicy.parseMode (reader.GetString(5)) with + | Some mode -> mode + | None -> failwith "dividend mode is invalid" + Status = reader.GetString(6) + IsSynthetic = reader.GetBoolean(7) + GrossCash = readDecimalOption reader 8 + CreditedUnits = readDecimalOption reader 9 + CreditedInvested = readDecimalOption reader 10 + OrderId = if reader.IsDBNull(11) then None else Some(reader.GetGuid(11)) + PendingReason = readStringOption reader 12 + CreatedAt = reader.GetFieldValue<DateTimeOffset>(13) + } + + let findDividendRecord connection transaction recordId = + use command = + commandWithTransaction + connection + transaction + """ + SELECT id, fund_id, instrument_code, nav_date, dps, mode, status, is_synthetic, + gross_cash, credited_units, credited_invested, order_id, pending_reason, created_at + FROM dividend_records WHERE id = @record_id + """ + + addParameter command "record_id" NpgsqlDbType.Uuid (box recordId) |> ignore + + use reader = command.ExecuteReader() + if reader.Read() then Some(dividendRecordFromReader reader) else None + + let findDividendIdempotency connection transaction key = + use command = + commandWithTransaction + connection + transaction + "SELECT request_hash, fund_id, record_id FROM dividend_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 insertDividendRecord connection transaction (record: DividendRecord) = + use command = + commandWithTransaction + connection + transaction + """ + INSERT INTO dividend_records + (id, fund_id, instrument_code, nav_date, dps, mode, status, is_synthetic, gross_cash) + VALUES (@id, @fund_id, @instrument_code, @nav_date, @dps, @mode, @status, @is_synthetic, @gross_cash) + RETURNING created_at + """ + + addParameter command "id" NpgsqlDbType.Uuid (box record.Id) |> ignore + addParameter command "fund_id" NpgsqlDbType.Uuid (box record.FundId) |> ignore + addParameter command "instrument_code" NpgsqlDbType.Text (box record.InstrumentCode) |> ignore + addParameter command "nav_date" NpgsqlDbType.Date (box record.NavDate) |> ignore + addParameter command "dps" NpgsqlDbType.Numeric (box record.Dps) |> ignore + addParameter command "mode" NpgsqlDbType.Text (box (DividendPolicy.modeText record.Mode)) |> ignore + addParameter command "status" NpgsqlDbType.Text (box record.Status) |> ignore + addParameter command "is_synthetic" NpgsqlDbType.Boolean (box record.IsSynthetic) |> ignore + + let grossParameter = + match record.GrossCash with + | Some gross -> box gross + | None -> box DBNull.Value + + addParameter command "gross_cash" NpgsqlDbType.Numeric grossParameter |> ignore + + use reader = command.ExecuteReader() + reader.Read() |> ignore + reader.GetFieldValue<DateTimeOffset>(0) + + let insertDividendIdempotency connection transaction key requestHash recordId fundId = + use command = + commandWithTransaction + connection + transaction + """ + INSERT INTO dividend_idempotencies (idempotency_key, request_hash, record_id, fund_id) + VALUES (@idempotency_key, @request_hash, @record_id, @fund_id) + """ + + addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore + addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore + addParameter command "record_id" NpgsqlDbType.Uuid (box recordId) |> ignore + addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + command.ExecuteNonQuery() |> ignore + + let finishDividendRecord connection transaction (recordId: Guid) (status: string) (pendingReason: string option) = + use command = + commandWithTransaction + connection + transaction + "UPDATE dividend_records SET status = @status, pending_reason = @reason WHERE id = @record_id" + + addParameter command "status" NpgsqlDbType.Text (box status) |> ignore + + let reasonParameter = + match pendingReason with + | Some reason -> box reason + | None -> box DBNull.Value + + addParameter command "reason" NpgsqlDbType.Text reasonParameter |> ignore + addParameter command "record_id" NpgsqlDbType.Uuid (box recordId) |> ignore + command.ExecuteNonQuery() |> ignore + + let dividendRequestHash (fundId: Guid) (command: DividendCommand) = + let invariant = CultureInfo.InvariantCulture + let encoded (value: string) = sprintf "%d:%s" value.Length value + + let payload = + String.concat + "|" + [ + "dividend" + encoded (fundId.ToString("D")) + encoded command.InstrumentCode + encoded (command.NavDate.ToString("yyyy-MM-dd")) + (encoded (command.Dps.ToString("G29", invariant))) + (encoded (DividendPolicy.modeText command.Mode)) + ] + + Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) + + let dividendSchemeKey (fundId: Guid) (command: DividendCommand) = + sprintf "dividend-scheme:%O:%s:%s:%s" fundId command.InstrumentCode (command.NavDate.ToString("yyyy-MM-dd")) (DividendPolicy.modeText command.Mode) + + let validateDividendCommand (today: DateOnly) (command: DividendCommand) = + if String.IsNullOrWhiteSpace command.InstrumentCode then + Error "instrument code cannot be empty" + elif command.Dps <= 0m then + Error "dividend per unit must be positive" + else + match DividendPolicy.validateDps command.Dps, DividendPolicy.validateNavDate command.NavDate today with + | Ok(), Ok() -> Ok() + | Error message, _ + | _, Error message -> Error message member _.EnsureSchema() = use connection = new NpgsqlConnection(connectionString) connection.Open() @@ -2643,6 +2874,250 @@ type FundRepository(connectionString: string) = Outcomes = outcomes |> Seq.toList } + member this.RegisterDividend(idempotencyKey: string, fundId: Guid, command: DividendCommand) : DividendWriteResult = + if String.IsNullOrWhiteSpace idempotencyKey then + DividendWriteResult.DividendInvalid "idempotency key cannot be empty" + else + let today = ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow + + match validateDividendCommand today command with + | Error message -> DividendWriteResult.DividendInvalid message + | Ok() -> + let schemeKey = dividendSchemeKey fundId command + let fingerprint = dividendRequestHash fundId command + use connection = new NpgsqlConnection(connectionString) + connection.Open() + + // Phase 1: claim the scheme atomically, validate, and book the payout + let claimedStage = + 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 schemeKey) |> ignore + lockCommand.ExecuteNonQuery() |> ignore + + match findDividendIdempotency connection (Some transaction) schemeKey with + | Some(existingHash, existingFundId, recordId) when existingFundId = fundId -> + if existingHash = fingerprint then + transaction.Commit() + StageReplay recordId + else + transaction.Rollback() + StageFail DividendWriteResult.DividendIdempotencyConflict + | Some _ -> + transaction.Rollback() + StageFail DividendWriteResult.DividendIdempotencyConflict + | None -> + match lockFundForOrder connection (Some transaction) fundId with + | None -> + transaction.Rollback() + StageFail DividendWriteResult.DividendFundNotFound + | Some isSynthetic -> + if not (instrumentExists connection (Some transaction) command.InstrumentCode) then + transaction.Rollback() + StageFail DividendWriteResult.DividendInstrumentNotFound + else + use positionCommand = + commandWithTransaction + connection + (Some transaction) + "SELECT units FROM fund_positions WHERE fund_id = @fund_id AND instrument_code = @code FOR UPDATE" + + addParameter positionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + addParameter positionCommand "code" NpgsqlDbType.Text (box command.InstrumentCode) |> ignore + + use positionReader = positionCommand.ExecuteReader() + let positionFound = positionReader.Read() + let heldUnits = if positionFound then positionReader.GetDecimal(0) else 0m + positionReader.Close() + + if not positionFound || heldUnits <= 0m then + transaction.Rollback() + StageFail DividendWriteResult.DividendNoHoldings + else + match DividendPolicy.computeCashPayout heldUnits command.Dps with + | Error message -> + transaction.Rollback() + StageFail (DividendWriteResult.DividendInvalid message) + | Ok payout -> + use cashCommand = + commandWithTransaction + connection + (Some transaction) + "UPDATE funds SET available_cash = available_cash + @gross WHERE id = @fund_id" + + addParameter cashCommand "gross" NpgsqlDbType.Numeric (box payout.GrossCash) |> ignore + addParameter cashCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + + if cashCommand.ExecuteNonQuery() = 0 then + transaction.Rollback() + StageFail DividendWriteResult.DividendFundNotFound + else + let recordId = Guid.NewGuid() + + let record: DividendRecord = + { + Id = recordId + FundId = fundId + InstrumentCode = command.InstrumentCode + NavDate = command.NavDate + Dps = command.Dps + Mode = command.Mode + Status = (if command.Mode = DividendPolicy.Cash then "cash_credited" else "pending_nav") + IsSynthetic = isSynthetic + GrossCash = Some payout.GrossCash + CreditedUnits = None + CreditedInvested = None + OrderId = None + PendingReason = None + CreatedAt = DateTimeOffset.UtcNow + } + + let createdAt = insertDividendRecord connection (Some transaction) record + insertDividendIdempotency connection (Some transaction) schemeKey fingerprint recordId fundId + transaction.Commit() + StageBooked { record with CreatedAt = createdAt } + with error -> + try + transaction.Rollback() + with _ -> + () + + raise error + + // Confirmation of the reinvestment order settles the dividend record: the + // credited units/invested cash are persisted and the row is marked succeeded. + let finalizeDividend (record: DividendRecord) (confirmed: SubscriptionOrderRecord) = + let units = confirmed.ConfirmedUnits |> Option.defaultValue 0m + let invested = confirmed.ConfirmedInvestedCash |> Option.defaultValue 0m + + do + use finalizeConnection = new NpgsqlConnection(connectionString) + finalizeConnection.Open() + + use finalizeTransaction = finalizeConnection.BeginTransaction(IsolationLevel.ReadCommitted) + + use finalizeCommand = + commandWithTransaction + finalizeConnection + (Some finalizeTransaction) + """ + UPDATE dividend_records + SET status = 'succeeded', + credited_units = @units, + credited_invested = @invested, + pending_reason = NULL + WHERE id = @record_id + """ + + addParameter finalizeCommand "units" NpgsqlDbType.Numeric (box units) |> ignore + addParameter finalizeCommand "invested" NpgsqlDbType.Numeric (box invested) |> ignore + addParameter finalizeCommand "record_id" NpgsqlDbType.Uuid (box record.Id) |> ignore + finalizeCommand.ExecuteNonQuery() |> ignore + finalizeTransaction.Commit() + + { record with Status = "succeeded"; CreditedUnits = Some units; CreditedInvested = Some invested } + + // Phase 2: replays never re-book; pending reinvestments retry their confirmation + match claimedStage with + | StageReplay recordId -> + let record = + match findDividendRecord connection None recordId with + | Some record -> record + | None -> failwith "idempotency scheme references a missing dividend" + + match record.Status, record.OrderId with + | "pending_nav", Some orderId -> + let confirmKey = sprintf "dividend-confirm:%O:%s:%s:%s" fundId command.InstrumentCode (command.NavDate.ToString("yyyy-MM-dd")) (command.Dps.ToString("G29", CultureInfo.InvariantCulture)) + + match this.ConfirmSubscriptionOrder(confirmKey, fundId, orderId) with + | SubscriptionConfirmResult.OrderConfirmed confirmed + | SubscriptionConfirmResult.ConfirmReplayed confirmed -> + DividendCredited(finalizeDividend record confirmed) + | SubscriptionConfirmResult.ConfirmPendingNav _ -> + DividendReplayed record + | other -> + failwithf "unexpected dividend confirm replay result: %A" other + | _ -> + DividendReplayed record + | StageFail failure -> failure + | StageBooked record -> + if record.Mode = DividendPolicy.Cash then + DividendCredited record + else + let orderKey = sprintf "dividend:%O:%s:%s:%s" fundId command.InstrumentCode (command.NavDate.ToString("yyyy-MM-dd")) (command.Dps.ToString("G29", CultureInfo.InvariantCulture)) + let confirmKey = sprintf "dividend-confirm:%O:%s:%s:%s" fundId command.InstrumentCode (command.NavDate.ToString("yyyy-MM-dd")) (command.Dps.ToString("G29", CultureInfo.InvariantCulture)) + + match + this.CreateSubscriptionOrder( + orderKey, + fundId, + { FundCode = command.InstrumentCode; Amount = record.GrossCash |> Option.defaultValue 0m; FeeAmount = 0m }, + command.NavDate + ) + with + | SubscriptionOrderWriteResult.OrderCreated order + | SubscriptionOrderWriteResult.OrderReplayed order -> + do + use bindConnection = new NpgsqlConnection(connectionString) + bindConnection.Open() + + use bindTransaction = bindConnection.BeginTransaction(IsolationLevel.ReadCommitted) + + use bindCommand = + commandWithTransaction + bindConnection + (Some bindTransaction) + "UPDATE dividend_records SET order_id = @order_id WHERE id = @record_id" + + addParameter bindCommand "order_id" NpgsqlDbType.Uuid (box order.Id) |> ignore + addParameter bindCommand "record_id" NpgsqlDbType.Uuid (box record.Id) |> ignore + bindCommand.ExecuteNonQuery() |> ignore + bindTransaction.Commit() + + let recordWithOrder = { record with OrderId = Some order.Id } + + match this.ConfirmSubscriptionOrder(confirmKey, fundId, order.Id) with + | SubscriptionConfirmResult.OrderConfirmed confirmed + | SubscriptionConfirmResult.ConfirmReplayed confirmed -> + DividendCredited(finalizeDividend recordWithOrder confirmed) + | SubscriptionConfirmResult.ConfirmPendingNav pending -> + DividendPendingReinvest { recordWithOrder with PendingReason = pending.PendingReason } + | other -> + failwithf "unexpected dividend confirm result: %A" other + | other -> + failwithf "unexpected dividend order result: %A" other + + member _.GetDividendRecords(fundId: Guid) = + use connection = new NpgsqlConnection(connectionString) + connection.Open() + + use command = + commandWithTransaction + connection + None + """ + SELECT id, fund_id, instrument_code, nav_date, dps, mode, status, is_synthetic, + gross_cash, credited_units, credited_invested, order_id, pending_reason, created_at + FROM dividend_records WHERE fund_id = @fund_id ORDER BY created_at DESC, id + """ + + addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + + use reader = command.ExecuteReader() + let records = ResizeArray<DividendRecord>() + + while reader.Read() do + records.Add(dividendRecordFromReader reader) + + records |> Seq.toList + member _.CreateCapitalDeposit(idempotencyKey: string, fundId: Guid, command: CapitalDepositCommand) : CapitalDepositWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then CapitalDepositWriteResult.CapitalDepositInvalid "idempotency key cannot be empty" |
