diff options
| author | Somhairle H. Marisol <[email protected]> | 2026-09-21 00:06:45 +0800 |
|---|---|---|
| committer | Somhairle H. Marisol <[email protected]> | 2026-09-21 00:06:45 +0800 |
| commit | 2ada70d6467aec11b45328112d598454ca50f2d6 (patch) | |
| tree | 5deab7d12095dfd0599c109edc2200a38258040e /src/FundLab.Api/Persistence.fs | |
| parent | c60905e9e7f992a7f8c79c3812e92a44b2d606f5 (diff) | |
| download | fund-lab-2ada70d6467aec11b45328112d598454ca50f2d6.tar.gz | |
feat(api): 接入基金 PostgreSQL 持久化与幂等接口
[变更性质]
- 本提交新增基金创建与读取的 PostgreSQL 持久化能力,并完成 3a API 收口验证。
[新增功能]
- 提供带事务和幂等键的基金创建、重放、冲突检测及基金读取接口。
- 覆盖认证、JSON 输入、数值边界、进程重启持久化和 SQL 故障回滚测试。
[实现方案]
- 使用 Npgsql 建立 `funds` 与 `fund_idempotencies` 表,并以 advisory lock 串行化同一幂等键。
- 通过真实 Kestrel 子进程和临时 PostgreSQL 验证 HTTP 行为;同步更新环境、构建和交接文档。
[影响范围]
- API 新增 PostgreSQL 配置要求 `FUND_LAB_DATABASE_URL`,前端构建入口和本地验证命令同步明确。
- 本提交不包含 `docs/overnight-progress.md` 的现有 Hermes 修改。
Diffstat (limited to 'src/FundLab.Api/Persistence.fs')
| -rw-r--r-- | src/FundLab.Api/Persistence.fs | 282 |
1 files changed, 282 insertions, 0 deletions
diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs new file mode 100644 index 0000000..db0fb6f --- /dev/null +++ b/src/FundLab.Api/Persistence.fs @@ -0,0 +1,282 @@ +namespace FundLab.Api + +open System +open System.Data +open System.Data.Common +open System.Globalization +open System.Security.Cryptography +open System.Text +open FundLab.Domain +open Npgsql +open NpgsqlTypes + +[<CLIMutable>] +type FundCreateCommand = + { + Name: string + InitialCash: decimal + InitialUnitNav: decimal + IsSynthetic: bool + } + +type FundRecord = + { + Id: Guid + Name: string + Currency: string + InitialCash: decimal + InitialUnitNav: decimal + IsSynthetic: bool + AvailableCash: decimal + Status: string + } + +type FundWriteResult = + | Created of FundRecord + | Replayed of FundRecord + | IdempotencyConflict + | Invalid of string + +type FundRepository(connectionString: string) = + let cashMaximum = 999999999999999999.99m + let unitNavMaximum = 99999999999999999999.99999999m + + let schema = + """ + CREATE TABLE IF NOT EXISTS funds ( + id uuid PRIMARY KEY, + name text NOT NULL, + currency text NOT NULL, + initial_cash numeric(20, 2) NOT NULL, + initial_unit_nav numeric(28, 8) NOT NULL, + is_synthetic boolean NOT NULL, + available_cash numeric(20, 2) NOT NULL, + status text NOT NULL, + created_at timestamptz NOT NULL DEFAULT now() + ); + + CREATE TABLE IF NOT EXISTS fund_idempotencies ( + idempotency_key text PRIMARY KEY, + request_hash text NOT NULL, + fund_id uuid NOT NULL REFERENCES funds(id), + created_at timestamptz NOT NULL DEFAULT now() + ); + """ + + let statusText status = + match status with + | FundStatus.Empty -> "empty" + | FundStatus.Active -> "active" + | FundStatus.ZeroUnits -> "zero_units" + + let recordFromLedger (fund: LedgerFund) = + { + Id = fund.Id + Name = fund.Name + Currency = fund.Currency + InitialCash = fund.InitialCash + InitialUnitNav = fund.InitialUnitNav + IsSynthetic = fund.IsSynthetic + AvailableCash = fund.AvailableCash + Status = statusText fund.Status + } + + let recordFromReader (reader: DbDataReader) = + { + Id = reader.GetGuid(0) + Name = reader.GetString(1) + Currency = reader.GetString(2) + InitialCash = reader.GetDecimal(3) + InitialUnitNav = reader.GetDecimal(4) + IsSynthetic = reader.GetBoolean(5) + AvailableCash = reader.GetDecimal(6) + Status = reader.GetString(7) + } + + let commandWithTransaction (connection: NpgsqlConnection) (transaction: NpgsqlTransaction option) sql = + let command = connection.CreateCommand() + command.CommandText <- sql + + match transaction with + | Some value -> command.Transaction <- value + | None -> () + + command + + let addParameter (command: NpgsqlCommand) name dbType (value: obj) = + let parameter = command.Parameters.Add(name, dbType) + parameter.Value <- value + parameter + + let findFund connection transaction fundId = + use command = + commandWithTransaction + connection + transaction + """ + SELECT id, name, currency, initial_cash, initial_unit_nav, + is_synthetic, available_cash, status + FROM funds + WHERE id = @fund_id + """ + + addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + + use reader = command.ExecuteReader() + if reader.Read() then Some(recordFromReader reader) else None + + let findIdempotency connection transaction key = + use command = + commandWithTransaction + connection + transaction + "SELECT request_hash, fund_id FROM fund_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)) else None + + let insertFund connection transaction (fund: LedgerFund) = + use command = + commandWithTransaction + connection + transaction + """ + INSERT INTO funds + (id, name, currency, initial_cash, initial_unit_nav, + is_synthetic, available_cash, status) + VALUES + (@id, @name, @currency, @initial_cash, @initial_unit_nav, + @is_synthetic, @available_cash, @status) + """ + + addParameter command "id" NpgsqlDbType.Uuid (box fund.Id) |> ignore + addParameter command "name" NpgsqlDbType.Text (box fund.Name) |> ignore + addParameter command "currency" NpgsqlDbType.Text (box fund.Currency) |> ignore + addParameter command "initial_cash" NpgsqlDbType.Numeric (box fund.InitialCash) |> ignore + addParameter command "initial_unit_nav" NpgsqlDbType.Numeric (box fund.InitialUnitNav) |> ignore + addParameter command "is_synthetic" NpgsqlDbType.Boolean (box fund.IsSynthetic) |> ignore + addParameter command "available_cash" NpgsqlDbType.Numeric (box fund.AvailableCash) |> ignore + addParameter command "status" NpgsqlDbType.Text (box (statusText fund.Status)) |> ignore + command.ExecuteNonQuery() |> ignore + + let insertIdempotency connection transaction key requestHash fundId = + use command = + commandWithTransaction + connection + transaction + """ + INSERT INTO fund_idempotencies (idempotency_key, request_hash, fund_id) + VALUES (@idempotency_key, @request_hash, @fund_id) + """ + + addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore + addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore + addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + command.ExecuteNonQuery() |> ignore + + let requestHash (command: FundCreateCommand) = + let invariant = CultureInfo.InvariantCulture + let encoded (value: string) = sprintf "%d:%s" value.Length value + let name = if isNull command.Name then "" else command.Name + let payload = + String.concat + "|" + [ + "fund-create" + encoded name + (encoded (command.InitialCash.ToString("G29", invariant))) + (encoded (command.InitialUnitNav.ToString("G29", invariant))) + (encoded (command.IsSynthetic.ToString())) + ] + + Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) + + let validateStorageRange (command: FundCreateCommand) = + if command.InitialCash > cashMaximum then + Error "initial cash exceeds database precision" + elif command.InitialUnitNav > unitNavMaximum then + Error "initial unit NAV exceeds database precision" + else + Ok() + + let ledgerErrorMessage error = + match error with + | InvalidIdentifier label -> sprintf "%s is invalid" label + | InvalidAmount label -> sprintf "%s is invalid" label + | InvalidPrecision label -> sprintf "%s has invalid precision" label + | InvalidState message -> message + | FundAlreadyExists fundId -> sprintf "fund %O already exists" fundId + | other -> sprintf "%A" other + + member _.EnsureSchema() = + use connection = new NpgsqlConnection(connectionString) + connection.Open() + use command = connection.CreateCommand() + command.CommandText <- schema + command.ExecuteNonQuery() |> ignore + + member _.GetFund(fundId: Guid) = + use connection = new NpgsqlConnection(connectionString) + connection.Open() + findFund connection None fundId + + member _.CreateFund(idempotencyKey: string, command: FundCreateCommand) = + if String.IsNullOrWhiteSpace idempotencyKey then + FundWriteResult.Invalid "idempotency key cannot be empty" + else + match validateStorageRange command with + | Error message -> FundWriteResult.Invalid message + | Ok() -> + let fingerprint = requestHash 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 findIdempotency connection (Some transaction) idempotencyKey with + | Some(existingHash, fundId) when existingHash = fingerprint -> + match findFund connection (Some transaction) fundId with + | Some fund -> + transaction.Commit() + FundWriteResult.Replayed fund + | None -> + transaction.Rollback() + FundWriteResult.Invalid "idempotency record references a missing fund" + | Some _ -> + transaction.Rollback() + FundWriteResult.IdempotencyConflict + | None -> + let fundId = Guid.NewGuid() + + match Ledger.initializeFund fundId command.Name command.InitialCash command.InitialUnitNav command.IsSynthetic Ledger.empty with + | Error error -> + transaction.Rollback() + FundWriteResult.Invalid(ledgerErrorMessage error) + | Ok state -> + match Ledger.getFund fundId state with + | Error error -> + transaction.Rollback() + FundWriteResult.Invalid(ledgerErrorMessage error) + | Ok fund -> + insertFund connection (Some transaction) fund + insertIdempotency connection (Some transaction) idempotencyKey fingerprint fundId + transaction.Commit() + FundWriteResult.Created(recordFromLedger fund) + with error -> + try + transaction.Rollback() + with _ -> + () + + raise error |
