summaryrefslogtreecommitdiff
path: root/src/FundLab.Api/Persistence.fs
diff options
context:
space:
mode:
Diffstat (limited to 'src/FundLab.Api/Persistence.fs')
-rw-r--r--src/FundLab.Api/Persistence.fs245
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"