diff options
Diffstat (limited to 'src/FundLab.Api/Persistence.fs')
| -rw-r--r-- | src/FundLab.Api/Persistence.fs | 344 |
1 files changed, 344 insertions, 0 deletions
diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs index 9ccb0c7..b22fbd9 100644 --- a/src/FundLab.Api/Persistence.fs +++ b/src/FundLab.Api/Persistence.fs @@ -358,6 +358,41 @@ type StockSellWriteResult = | StockSellInsufficientHoldings of string | StockSellFundNotFound +type StockCashflowCommand = + { + InstrumentCode: string + StockName: string option + /// "buy" (cash out), "sell" (cash in) or "dividend" (cash in). + EventType: string + EventDate: DateOnly + Quantity: decimal + /// Absolute cash magnitude; the sign is implied by the event type. + Amount: decimal + Note: string option + } + +type StockCashflowRecord = + { + Id: Guid + FundId: Guid + InstrumentCode: string + StockName: string option + EventType: string + EventDate: DateOnly + Quantity: decimal + Amount: decimal + Note: string option + IsSynthetic: bool + CreatedAt: DateTimeOffset + } + +type StockCashflowWriteResult = + | StockCashflowCreated of StockCashflowRecord + | StockCashflowReplayed of StockCashflowRecord + | StockCashflowIdempotencyConflict + | StockCashflowInvalid of string + | StockCashflowFundNotFound + type BondTradeCommand = { InstrumentCode: string @@ -1100,6 +1135,28 @@ type FundRepository(connectionString: string) = executed_at timestamptz NOT NULL ); + CREATE TABLE IF NOT EXISTS stock_cashflow_events ( + id uuid PRIMARY KEY, + fund_id uuid NOT NULL REFERENCES funds(id), + instrument_code text NOT NULL, + stock_name text NULL, + event_type text NOT NULL CHECK (event_type IN ('buy', 'sell', 'dividend')), + event_date date NOT NULL, + quantity numeric(28, 8) NOT NULL CHECK (quantity >= 0), + amount numeric(20, 2) NOT NULL CHECK (amount >= 0), + note text NULL, + is_synthetic boolean NOT NULL, + created_at timestamptz NOT NULL + ); + + CREATE TABLE IF NOT EXISTS stock_cashflow_idempotencies ( + idempotency_key text PRIMARY KEY, + request_hash text NOT NULL, + event_id uuid NOT NULL REFERENCES stock_cashflow_events(id), + fund_id uuid NOT NULL REFERENCES funds(id), + created_at timestamptz NOT NULL DEFAULT now() + ); + CREATE TABLE IF NOT EXISTS stock_sell_idempotencies ( idempotency_key text PRIMARY KEY, request_hash text NOT NULL, @@ -2312,6 +2369,165 @@ type FundRepository(connectionString: string) = Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) + let stockCashflowRecordFromReader (reader: DbDataReader) : StockCashflowRecord = + { + Id = reader.GetGuid(0) + FundId = reader.GetGuid(1) + InstrumentCode = reader.GetString(2) + StockName = if reader.IsDBNull(3) then None else Some(reader.GetString(3)) + EventType = reader.GetString(4) + EventDate = reader.GetFieldValue<DateOnly>(5) + Quantity = reader.GetDecimal(6) + Amount = reader.GetDecimal(7) + Note = readStringOption reader 8 + IsSynthetic = reader.GetBoolean(9) + CreatedAt = reader.GetFieldValue<DateTimeOffset>(10) + } + + let stockCashflowColumns = + "id, fund_id, instrument_code, stock_name, event_type, event_date, quantity, amount, note, is_synthetic, created_at" + + let insertStockCashflow connection transaction (record: StockCashflowRecord) = + use command = + commandWithTransaction + connection + transaction + """ + INSERT INTO stock_cashflow_events + (id, fund_id, instrument_code, stock_name, event_type, event_date, quantity, amount, note, is_synthetic, created_at) + VALUES + (@id, @fund_id, @instrument_code, @stock_name, @event_type, @event_date, @quantity, @amount, @note, @is_synthetic, @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 + + let nameParameter = + match record.StockName with + | Some name -> box name + | None -> box DBNull.Value + + addParameter command "stock_name" NpgsqlDbType.Text nameParameter |> ignore + addParameter command "event_type" NpgsqlDbType.Text (box record.EventType) |> ignore + addParameter command "event_date" NpgsqlDbType.Date (box record.EventDate) |> ignore + addParameter command "quantity" NpgsqlDbType.Numeric (box record.Quantity) |> ignore + addParameter command "amount" NpgsqlDbType.Numeric (box record.Amount) |> ignore + + let noteParameter = + match record.Note with + | Some note -> box note + | None -> box DBNull.Value + + addParameter command "note" NpgsqlDbType.Text noteParameter |> ignore + addParameter command "is_synthetic" NpgsqlDbType.Boolean (box record.IsSynthetic) |> ignore + addParameter command "created_at" NpgsqlDbType.TimestampTz (box record.CreatedAt) |> ignore + command.ExecuteNonQuery() |> ignore + + let insertStockCashflowIdempotency connection transaction key requestHash eventId fundId = + use command = + commandWithTransaction + connection + transaction + """ + INSERT INTO stock_cashflow_idempotencies (idempotency_key, request_hash, event_id, fund_id) + VALUES (@idempotency_key, @request_hash, @event_id, @fund_id) + """ + + addParameter command "idempotency_key" NpgsqlDbType.Text (box key) |> ignore + addParameter command "request_hash" NpgsqlDbType.Text (box requestHash) |> ignore + addParameter command "event_id" NpgsqlDbType.Uuid (box eventId) |> ignore + addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + command.ExecuteNonQuery() |> ignore + + let findStockCashflowIdempotency connection transaction key = + use command = + commandWithTransaction + connection + transaction + """ + SELECT request_hash, fund_id, event_id + FROM stock_cashflow_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 findStockCashflow connection transaction eventId = + use command = + commandWithTransaction + connection + transaction + $""" + SELECT {stockCashflowColumns} + FROM stock_cashflow_events + WHERE id = @id + """ + + addParameter command "id" NpgsqlDbType.Uuid (box eventId) |> ignore + + use reader = command.ExecuteReader() + + if reader.Read() then + Some(stockCashflowRecordFromReader reader) + else + None + + let stockCashflowRequestHash (fundId: Guid) (command: StockCashflowCommand) = + let invariant = CultureInfo.InvariantCulture + let encoded (value: string) = sprintf "%d:%s" value.Length value + let name = command.StockName |> Option.defaultValue "" + let note = command.Note |> Option.defaultValue "" + + let payload = + String.concat + "|" + [ + "stock-cashflow" + encoded (fundId.ToString("D")) + encoded command.InstrumentCode + encoded name + encoded command.EventType + encoded (command.EventDate.ToString("yyyy-MM-dd", invariant)) + encoded (command.Quantity.ToString("G29", invariant)) + encoded (command.Amount.ToString("G29", invariant)) + encoded note + ] + + Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(payload))) + + let stockCashflowEvent + (fundId: Guid) + (instrumentCode: string) + (stockName: string option) + (eventType: string) + (eventDate: DateOnly) + (quantity: decimal) + (amount: decimal) + (note: string option) + (isSynthetic: bool) + : StockCashflowRecord = + { + Id = Guid.NewGuid() + FundId = fundId + InstrumentCode = instrumentCode + StockName = stockName + EventType = eventType + EventDate = eventDate + Quantity = quantity + Amount = amount + Note = note + IsSynthetic = isSynthetic + CreatedAt = DateTimeOffset.UtcNow + } + let bondTradeRecordFromReader (reader: DbDataReader) : BondTradeRecord = { Id = reader.GetGuid(0) @@ -5392,6 +5608,20 @@ type FundRepository(connectionString: string) = addParameter positionCommand "last_traded_at" NpgsqlDbType.TimestampTz (box executedAt) |> ignore positionCommand.ExecuteNonQuery() |> ignore + let buyEvent = + stockCashflowEvent + fundId + normalized.InstrumentCode + normalized.StockName + "buy" + (DateOnly.FromDateTime executedAt.UtcDateTime) + normalized.Quantity + costCash + None + isSynthetic + + insertStockCashflow connection (Some transaction) buyEvent + let persisted = { trade with ExecutedAt = executedAt } transaction.Commit() StockTradeWriteResult.StockTradeCreated persisted @@ -5487,6 +5717,105 @@ type FundRepository(connectionString: string) = records |> Seq.toList + member _.RecordStockCashflow(idempotencyKey: string, fundId: Guid, command: StockCashflowCommand) : StockCashflowWriteResult = + let code = if isNull command.InstrumentCode then "" else command.InstrumentCode.Trim() + let eventType = if isNull command.EventType then "" else command.EventType.Trim().ToLowerInvariant() + + if String.IsNullOrWhiteSpace idempotencyKey then + StockCashflowWriteResult.StockCashflowInvalid "idempotency key cannot be empty" + elif code.Length <> 6 || not (code |> Seq.forall Char.IsDigit) then + StockCashflowWriteResult.StockCashflowInvalid "stock code must contain exactly six digits" + elif eventType <> "buy" && eventType <> "sell" && eventType <> "dividend" then + StockCashflowWriteResult.StockCashflowInvalid "event type must be buy, sell or dividend" + elif command.Quantity < 0m then + StockCashflowWriteResult.StockCashflowInvalid "quantity cannot be negative" + elif command.Amount < 0m then + StockCashflowWriteResult.StockCashflowInvalid "amount cannot be negative" + else + let normalized = { command with InstrumentCode = code; EventType = eventType } + let fingerprint = stockCashflowRequestHash fundId normalized + 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 findStockCashflowIdempotency connection (Some transaction) idempotencyKey with + | Some(existingHash, existingFundId, eventId) + when existingHash = fingerprint && existingFundId = fundId -> + match findStockCashflow connection (Some transaction) eventId with + | Some record -> + transaction.Commit() + StockCashflowWriteResult.StockCashflowReplayed record + | None -> + transaction.Rollback() + StockCashflowWriteResult.StockCashflowInvalid "idempotency record references a missing event" + | Some _ -> + transaction.Rollback() + StockCashflowWriteResult.StockCashflowIdempotencyConflict + | None -> + match lockFundForOrder connection (Some transaction) fundId with + | None -> + transaction.Rollback() + StockCashflowWriteResult.StockCashflowFundNotFound + | Some isSynthetic -> + let record = + stockCashflowEvent + fundId + normalized.InstrumentCode + normalized.StockName + normalized.EventType + normalized.EventDate + normalized.Quantity + normalized.Amount + normalized.Note + isSynthetic + + insertStockCashflow connection (Some transaction) record + insertStockCashflowIdempotency connection (Some transaction) idempotencyKey fingerprint record.Id fundId + transaction.Commit() + StockCashflowWriteResult.StockCashflowCreated record + with error -> + try + transaction.Rollback() + with _ -> + () + + raise error + + member _.GetStockCashflows(fundId: Guid) : StockCashflowRecord list = + use connection = new NpgsqlConnection(connectionString) + connection.Open() + + use command = + commandWithTransaction + connection + None + $""" + SELECT {stockCashflowColumns} + FROM stock_cashflow_events + WHERE fund_id = @fund_id + ORDER BY event_date, created_at, id + """ + + addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore + + use reader = command.ExecuteReader() + let records = ResizeArray<StockCashflowRecord>() + + while reader.Read() do + records.Add(stockCashflowRecordFromReader reader) + + records |> Seq.toList + member _.CreateStockSell(idempotencyKey: string, fundId: Guid, command: StockSellCommand, ?executedAtOverride: DateTimeOffset) : StockSellWriteResult = if String.IsNullOrWhiteSpace idempotencyKey then StockSellWriteResult.StockSellInvalid "idempotency key cannot be empty" @@ -5646,6 +5975,21 @@ type FundRepository(connectionString: string) = insertStockSell connection (Some transaction) sell insertStockSellIdempotency connection (Some transaction) idempotencyKey fingerprint sell.Id fundId + + let sellEvent = + stockCashflowEvent + fundId + normalized.InstrumentCode + normalized.StockName + "sell" + (DateOnly.FromDateTime executedAt.UtcDateTime) + normalized.Quantity + proceeds + None + isSynthetic + + insertStockCashflow connection (Some transaction) sellEvent + transaction.Commit() StockSellWriteResult.StockSellCreated sell with error -> |
