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 [] 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() ); CREATE TABLE IF NOT EXISTS instruments ( code text PRIMARY KEY CHECK (code ~ '^[0-9]{6}$'), name text NOT NULL, fund_type text NULL, source text NOT NULL, source_revision text NOT NULL, source_collected_at timestamptz NOT NULL, source_payload_hash text NOT NULL, first_seen_at timestamptz NOT NULL DEFAULT now(), last_seen_at timestamptz NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS fund_nav_observations ( instrument_code text NOT NULL REFERENCES instruments(code), nav_date date NOT NULL, published_at timestamptz NULL, nav numeric(28, 8) NOT NULL, accumulated_nav numeric(28, 8) NULL, daily_return numeric(20, 8) NULL, source text NOT NULL, source_revision text NOT NULL, source_collected_at timestamptz NOT NULL, source_payload_hash text NOT NULL, first_seen_at timestamptz NOT NULL DEFAULT now(), last_seen_at timestamptz NOT NULL DEFAULT now(), PRIMARY KEY (instrument_code, nav_date) ); CREATE INDEX IF NOT EXISTS fund_nav_observations_date_idx ON fund_nav_observations (instrument_code, nav_date); """ 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 dateTimeOffsetFromReader (reader: DbDataReader) index = reader.GetFieldValue(index) let optionalDateTimeOffsetFromReader (reader: DbDataReader) index = if reader.IsDBNull(index) then None else Some(dateTimeOffsetFromReader reader index) let optionalDecimalFromReader (reader: DbDataReader) index = if reader.IsDBNull(index) then None else Some(reader.GetDecimal(index)) let instrumentRecordFromReader (reader: DbDataReader) : MarketDataInstrumentRecord = { Code = reader.GetString(0) Name = reader.GetString(1) FundType = if reader.IsDBNull(2) then None else Some(reader.GetString(2)) Source = reader.GetString(3) SourceRevision = reader.GetString(4) SourceCollectedAt = dateTimeOffsetFromReader reader 5 SourcePayloadHash = reader.GetString(6) FirstSeenAt = dateTimeOffsetFromReader reader 7 LastSeenAt = dateTimeOffsetFromReader reader 8 } let navRecordFromReader (reader: DbDataReader) : MarketDataNavRecord = { Code = reader.GetString(0) NavDate = reader.GetFieldValue(1) PublishedAt = optionalDateTimeOffsetFromReader reader 2 Nav = reader.GetDecimal(3) AccumulatedNav = optionalDecimalFromReader reader 4 DailyReturn = optionalDecimalFromReader reader 5 Source = reader.GetString(6) SourceRevision = reader.GetString(7) SourceCollectedAt = dateTimeOffsetFromReader reader 8 SourcePayloadHash = reader.GetString(9) FirstSeenAt = dateTimeOffsetFromReader reader 10 LastSeenAt = dateTimeOffsetFromReader reader 11 } 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 optionalParameterValue value = match value with | Some actual -> box actual | None -> box DBNull.Value 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 upsertInstrument connection transaction (payload: MarketDataSearchPayload) payloadHash (instrument: MarketDataInstrument) = use command = commandWithTransaction connection transaction """ INSERT INTO instruments (code, name, fund_type, source, source_revision, source_collected_at, source_payload_hash) VALUES (@code, @name, @fund_type, @source, @source_revision, @source_collected_at, @source_payload_hash) ON CONFLICT (code) DO UPDATE SET name = EXCLUDED.name, fund_type = EXCLUDED.fund_type, source = EXCLUDED.source, source_revision = EXCLUDED.source_revision, source_collected_at = EXCLUDED.source_collected_at, source_payload_hash = EXCLUDED.source_payload_hash, last_seen_at = now() """ addParameter command "code" NpgsqlDbType.Text (box instrument.Code) |> ignore addParameter command "name" NpgsqlDbType.Text (box instrument.Name) |> ignore addParameter command "fund_type" NpgsqlDbType.Text (optionalParameterValue instrument.FundType) |> ignore addParameter command "source" NpgsqlDbType.Text (box payload.Source) |> ignore addParameter command "source_revision" NpgsqlDbType.Text (box payload.SourceRevision) |> ignore addParameter command "source_collected_at" NpgsqlDbType.TimestampTz (box payload.CollectedAt) |> ignore addParameter command "source_payload_hash" NpgsqlDbType.Text (box payloadHash) |> ignore command.ExecuteNonQuery() |> ignore let upsertNavObservation connection transaction (payload: MarketDataNavPayload) payloadHash (observation: MarketDataObservation) = use command = commandWithTransaction connection transaction """ INSERT INTO fund_nav_observations (instrument_code, nav_date, published_at, nav, accumulated_nav, daily_return, source, source_revision, source_collected_at, source_payload_hash) VALUES (@instrument_code, @nav_date, @published_at, @nav, @accumulated_nav, @daily_return, @source, @source_revision, @source_collected_at, @source_payload_hash) ON CONFLICT (instrument_code, nav_date) DO UPDATE SET published_at = EXCLUDED.published_at, nav = EXCLUDED.nav, accumulated_nav = EXCLUDED.accumulated_nav, daily_return = EXCLUDED.daily_return, source = EXCLUDED.source, source_revision = EXCLUDED.source_revision, source_collected_at = EXCLUDED.source_collected_at, source_payload_hash = EXCLUDED.source_payload_hash, last_seen_at = now() """ addParameter command "instrument_code" NpgsqlDbType.Text (box payload.Code) |> ignore addParameter command "nav_date" NpgsqlDbType.Date (box observation.NavDate) |> ignore addParameter command "published_at" NpgsqlDbType.TimestampTz (optionalParameterValue observation.PublishedAt) |> ignore addParameter command "nav" NpgsqlDbType.Numeric (box observation.Nav) |> ignore addParameter command "accumulated_nav" NpgsqlDbType.Numeric (optionalParameterValue observation.AccumulatedNav) |> ignore addParameter command "daily_return" NpgsqlDbType.Numeric (optionalParameterValue observation.DailyReturn) |> ignore addParameter command "source" NpgsqlDbType.Text (box payload.Source) |> ignore addParameter command "source_revision" NpgsqlDbType.Text (box payload.SourceRevision) |> ignore addParameter command "source_collected_at" NpgsqlDbType.TimestampTz (box payload.CollectedAt) |> ignore addParameter command "source_payload_hash" NpgsqlDbType.Text (box payloadHash) |> 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 _.UpsertInstruments(payload: MarketDataSearchPayload, payloadHash: string) = use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try for instrument in payload.Instruments do upsertInstrument connection (Some transaction) payload payloadHash instrument transaction.Commit() with error -> try transaction.Rollback() with _ -> () raise error member _.GetInstrument(code: string) = use connection = new NpgsqlConnection(connectionString) connection.Open() use command = commandWithTransaction connection None """ SELECT code, name, fund_type, source, source_revision, source_collected_at, source_payload_hash, first_seen_at, last_seen_at FROM instruments WHERE code = @code """ addParameter command "code" NpgsqlDbType.Text (box code) |> ignore use reader = command.ExecuteReader() if reader.Read() then Some(instrumentRecordFromReader reader) else None member _.UpsertNavObservations(payload: MarketDataNavPayload, payloadHash: string) = use connection = new NpgsqlConnection(connectionString) connection.Open() use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted) try for observation in payload.Observations do upsertNavObservation connection (Some transaction) payload payloadHash observation transaction.Commit() with error -> try transaction.Rollback() with _ -> () raise error member _.GetNav(code: string, fromDate: DateOnly option, toDate: DateOnly option) = use connection = new NpgsqlConnection(connectionString) connection.Open() let conditions = ResizeArray() conditions.Add("instrument_code = @instrument_code") if fromDate.IsSome then conditions.Add("nav_date >= @from_date") if toDate.IsSome then conditions.Add("nav_date <= @to_date") use command = commandWithTransaction connection None (sprintf """ SELECT instrument_code, nav_date, published_at, nav, accumulated_nav, daily_return, source, source_revision, source_collected_at, source_payload_hash, first_seen_at, last_seen_at FROM fund_nav_observations WHERE %s ORDER BY nav_date ASC """ (String.concat " AND " conditions)) addParameter command "instrument_code" NpgsqlDbType.Text (box code) |> ignore match fromDate with | Some value -> addParameter command "from_date" NpgsqlDbType.Date (box value) |> ignore | None -> () match toDate with | Some value -> addParameter command "to_date" NpgsqlDbType.Date (box value) |> ignore | None -> () use reader = command.ExecuteReader() let records = ResizeArray() while reader.Read() do records.Add(navRecordFromReader reader) records |> Seq.toList 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