diff options
Diffstat (limited to 'src/FundLab.Api/Persistence.fs')
| -rw-r--r-- | src/FundLab.Api/Persistence.fs | 245 |
1 files changed, 245 insertions, 0 deletions
diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs index 444f578..8ca85e4 100644 --- a/src/FundLab.Api/Persistence.fs +++ b/src/FundLab.Api/Persistence.fs @@ -256,6 +256,35 @@ type RedemptionConfirmResult = | RedemptionInvalidStatus | RedemptionInvalid of string +type SipPlanCommand = + { + InstrumentCode: string + Amount: decimal + Frequency: SipFrequency + } + +type SipPlanRecord = + { + Id: Guid + FundId: Guid + InstrumentCode: string + Amount: decimal + Frequency: SipFrequency + Status: string + IsSynthetic: bool + AnchorDate: DateOnly + NextTradeDate: DateOnly + CreatedAt: DateTimeOffset + } + +type SipPlanWriteResult = + | SipPlanCreated of SipPlanRecord + | SipPlanReplayed of SipPlanRecord + | SipPlanIdempotencyConflict + | SipPlanInvalid of string + | SipPlanFundNotFound + | SipPlanInstrumentNotFound + type CapitalDepositCommand = { Amount: decimal @@ -473,6 +502,27 @@ type FundRepository(connectionString: string) = fund_id uuid NOT NULL REFERENCES funds(id), created_at timestamptz NOT NULL DEFAULT now() ); + + CREATE TABLE IF NOT EXISTS sip_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_trade_date date NOT NULL, + created_at timestamptz NOT NULL DEFAULT now() + ); + + CREATE TABLE IF NOT EXISTS sip_plan_idempotencies ( + idempotency_key text PRIMARY KEY, + request_hash text NOT NULL, + plan_id uuid NOT NULL REFERENCES sip_plans(id), + fund_id uuid NOT NULL REFERENCES funds(id), + created_at timestamptz NOT NULL DEFAULT now() + ); """ let statusText status = @@ -1285,6 +1335,109 @@ type FundRepository(connectionString: string) = Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) + let sipPlanRecordFromReader (reader: DbDataReader) : SipPlanRecord = + { + Id = reader.GetGuid(0) + FundId = reader.GetGuid(1) + InstrumentCode = reader.GetString(2) + Amount = reader.GetDecimal(3) + Frequency = + match SipPolicy.parseFrequency (reader.GetString(4)) with + | Some frequency -> frequency + | None -> failwith "sip plan frequency is invalid" + Status = reader.GetString(5) + IsSynthetic = reader.GetBoolean(6) + AnchorDate = reader.GetFieldValue<DateOnly>(7) + NextTradeDate = reader.GetFieldValue<DateOnly>(8) + CreatedAt = reader.GetFieldValue<DateTimeOffset>(9) + } + + let findSipPlan connection transaction planId = + use command = + commandWithTransaction + connection + transaction + "SELECT id, fund_id, instrument_code, amount, frequency, status, is_synthetic, anchor_date, next_trade_date, created_at FROM sip_plans WHERE id = @plan_id" + + addParameter command "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore + + use reader = command.ExecuteReader() + if reader.Read() then Some(sipPlanRecordFromReader reader) else None + + let findSipPlanIdempotency connection transaction key = + use command = + commandWithTransaction + connection + transaction + "SELECT request_hash, fund_id, plan_id FROM sip_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 insertSipPlan connection transaction (plan: SipPlanRecord) = + use command = + commandWithTransaction + connection + transaction + """ + INSERT INTO sip_plans (id, fund_id, instrument_code, amount, frequency, status, is_synthetic, anchor_date, next_trade_date) + VALUES (@id, @fund_id, @instrument_code, @amount, @frequency, @status, @is_synthetic, @anchor_date, @next_trade_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 (SipPolicy.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_trade_date" NpgsqlDbType.Date (box plan.NextTradeDate) |> ignore + + use reader = command.ExecuteReader() + reader.Read() |> ignore + reader.GetFieldValue<DateTimeOffset>(0) + + let insertSipPlanIdempotency connection transaction key requestHash planId fundId = + use command = + commandWithTransaction + connection + transaction + """ + INSERT INTO sip_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 sipPlanRequestHash (fundId: Guid) (command: SipPlanCommand) = + 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 + "|" + [ + "sip-plan" + encoded (fundId.ToString("D")) + encoded code + (encoded (command.Amount.ToString("G29", invariant))) + (encoded (SipPolicy.frequencyText command.Frequency)) + ] + + Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) + member _.EnsureSchema() = use connection = new NpgsqlConnection(connectionString) connection.Open() @@ -1588,6 +1741,98 @@ type FundRepository(connectionString: string) = records |> Seq.toList + member _.CreateSipPlan(idempotencyKey: string, fundId: Guid, command: SipPlanCommand) : SipPlanWriteResult = + if String.IsNullOrWhiteSpace idempotencyKey then + SipPlanWriteResult.SipPlanInvalid "idempotency key cannot be empty" + else + match SipPolicy.validateAmount command.Amount with + | Error message -> SipPlanWriteResult.SipPlanInvalid message + | Ok() -> + let fingerprint = sipPlanRequestHash 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 findSipPlanIdempotency connection (Some transaction) idempotencyKey with + | Some(existingHash, existingFundId, planId) + when existingHash = fingerprint && existingFundId = fundId -> + match findSipPlan connection (Some transaction) planId with + | Some plan -> + transaction.Commit() + SipPlanWriteResult.SipPlanReplayed plan + | None -> + transaction.Rollback() + SipPlanWriteResult.SipPlanInvalid "idempotency record references a missing plan" + | Some _ -> + transaction.Rollback() + SipPlanWriteResult.SipPlanIdempotencyConflict + | None -> + match lockFundForOrder connection (Some transaction) fundId with + | None -> + transaction.Rollback() + SipPlanWriteResult.SipPlanFundNotFound + | Some isSynthetic -> + if instrumentExists connection (Some transaction) command.InstrumentCode then + let anchorDate = ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow + let plan: SipPlanRecord = + { + Id = Guid.NewGuid() + FundId = fundId + InstrumentCode = command.InstrumentCode + Amount = command.Amount + Frequency = command.Frequency + Status = "active" + IsSynthetic = isSynthetic + AnchorDate = anchorDate + NextTradeDate = SipPolicy.nextTradeDate command.Frequency anchorDate (anchorDate.AddDays 1) + CreatedAt = DateTimeOffset.UtcNow + } + + let createdAt = insertSipPlan connection (Some transaction) plan + insertSipPlanIdempotency connection (Some transaction) idempotencyKey fingerprint plan.Id fundId + transaction.Commit() + SipPlanWriteResult.SipPlanCreated { plan with CreatedAt = createdAt } + else + transaction.Rollback() + SipPlanWriteResult.SipPlanInstrumentNotFound + with error -> + try + transaction.Rollback() + with _ -> + () + + raise error + + member _.GetSipPlans(fundId: Guid) = + use connection = new NpgsqlConnection(connectionString) + connection.Open() + + use command = + commandWithTransaction + connection + None + "SELECT id, fund_id, instrument_code, amount, frequency, status, is_synthetic, anchor_date, next_trade_date, created_at FROM sip_plans WHERE fund_id = @fund_id ORDER BY created_at DESC, id" + + addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + + use reader = command.ExecuteReader() + let records = ResizeArray<SipPlanRecord>() + + while reader.Read() do + records.Add(sipPlanRecordFromReader 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" |
