diff options
Diffstat (limited to 'src/FundLab.Api')
| -rw-r--r-- | src/FundLab.Api/App.fs | 104 | ||||
| -rw-r--r-- | src/FundLab.Api/Persistence.fs | 245 |
2 files changed, 349 insertions, 0 deletions
diff --git a/src/FundLab.Api/App.fs b/src/FundLab.Api/App.fs index 55a29f4..9f88468 100644 --- a/src/FundLab.Api/App.fs +++ b/src/FundLab.Api/App.fs @@ -4,6 +4,7 @@ open System open System.Globalization open System.IO open System.Text.Json +open FundLab.Domain open Giraffe open Microsoft.AspNetCore.Http open Microsoft.Extensions.DependencyInjection @@ -131,6 +132,20 @@ type FundPositionsResponse = positions: FundPositionResponse list } +type SipPlanResponse = + { + id: Guid + fundId: Guid + instrumentCode: string + amount: string + frequency: string + status: string + anchorDate: string + nextTradeDate: string + isSynthetic: bool + createdAt: string + } + type CapitalDepositResponse = { id: Guid @@ -281,6 +296,20 @@ module App = createdAt = timestampText deposit.CreatedAt } + let private sipPlanResponse (plan: SipPlanRecord) : SipPlanResponse = + { + id = plan.Id + fundId = plan.FundId + instrumentCode = plan.InstrumentCode + amount = cashText plan.Amount + frequency = SipPolicy.frequencyText plan.Frequency + status = plan.Status + anchorDate = dateText plan.AnchorDate + nextTradeDate = dateText plan.NextTradeDate + isSynthetic = plan.IsSynthetic + createdAt = timestampText plan.CreatedAt + } + let private errorResponse status error message : HttpHandler = setStatusCode status >=> json ({ @@ -411,8 +440,34 @@ module App = with | :? JsonException -> Error "request body must be valid JSON" + let private parseSipPlanCommand (body: string) = + try + use document = JsonDocument.Parse(body) + let root = document.RootElement + + if root.ValueKind <> JsonValueKind.Object then + Error "request body must be a JSON object" + else + match tryStringProperty root "instrumentCode", tryStringProperty root "amount", tryStringProperty root "frequency" with + | Some code, Some amountText, Some frequencyText -> + match tryDecimal "amount" amountText, SipPolicy.parseFrequency frequencyText with + | Ok amount, Some frequency -> + Ok + { + InstrumentCode = code + Amount = amount + Frequency = frequency + } + | Error message, _ -> Error message + | _, None -> Error "frequency must be one of weekly, biweekly or monthly" + | _ -> + Error "instrumentCode, amount and frequency are required" + with + | :? JsonException -> Error "request body must be valid JSON" + let private invokeHandler handler next ctx = handler next ctx + let private unauthorized : HttpHandler = setStatusCode 401 >=> setHttpHeader "WWW-Authenticate" "Bearer" @@ -708,6 +763,53 @@ module App = json (deposits |> List.map capitalDepositResponse) next ctx with _ -> errorResponse 500 "PERSISTENCE_ERROR" "capital deposit persistence failed" next ctx + + let private createSipPlan (repository: FundRepository) (fundIdText: string) : HttpHandler = + fun next ctx -> + task { + match Guid.TryParse fundIdText with + | false, _ -> + return! invokeHandler (errorResponse 400 "INVALID_SIP_REQUEST" "fund id must be a UUID") next ctx + | true, fundId -> + use reader = new StreamReader(ctx.Request.Body) + let! body = reader.ReadToEndAsync() + let idempotencyKey = ctx.Request.Headers["Idempotency-Key"].ToString() + + match parseSipPlanCommand body with + | Error message -> + return! invokeHandler (errorResponse 400 "INVALID_SIP_REQUEST" message) next ctx + | Ok command -> + try + match repository.CreateSipPlan(idempotencyKey, fundId, command) with + | SipPlanWriteResult.SipPlanCreated plan -> + return! invokeHandler (setStatusCode 201 >=> json (sipPlanResponse plan)) next ctx + | SipPlanWriteResult.SipPlanReplayed plan -> + return! invokeHandler (json (sipPlanResponse plan)) next ctx + | SipPlanWriteResult.SipPlanIdempotencyConflict -> + return! invokeHandler (errorResponse 409 "IDEMPOTENCY_CONFLICT" "idempotency key was used with a different request") next ctx + | SipPlanWriteResult.SipPlanInvalid message -> + return! invokeHandler (errorResponse 400 "INVALID_SIP_REQUEST" message) next ctx + | SipPlanWriteResult.SipPlanFundNotFound -> + return! invokeHandler (errorResponse 404 "FUND_NOT_FOUND" "fund was not found") next ctx + | SipPlanWriteResult.SipPlanInstrumentNotFound -> + return! invokeHandler (errorResponse 404 "INSTRUMENT_NOT_FOUND" "instrument code was not found in the instrument catalog") next ctx + with _ -> + return! invokeHandler (errorResponse 500 "PERSISTENCE_ERROR" "sip plan persistence failed") next ctx + } + + let private getSipPlans (repository: FundRepository) (fundIdText: string) : HttpHandler = + fun next ctx -> + match Guid.TryParse fundIdText with + | false, _ -> errorResponse 400 "INVALID_FUND_ID" "fund id must be a UUID" next ctx + | true, fundId -> + try + match repository.GetFund fundId with + | None -> errorResponse 404 "FUND_NOT_FOUND" "fund was not found" next ctx + | Some _ -> + let plans = repository.GetSipPlans fundId + json (plans |> List.map sipPlanResponse) next ctx + with _ -> + errorResponse 500 "PERSISTENCE_ERROR" "sip plan persistence failed" next ctx let private marketDataError (failure: MarketDataFailure) : HttpHandler = let status, error, message = match failure with @@ -812,6 +914,8 @@ module App = GET >=> routef "/funds/%s/positions" (getPositions repository) POST >=> routef "/funds/%s/capital/deposit" (createCapitalDeposit repository) GET >=> routef "/funds/%s/capital/deposits" (getCapitalDeposits repository) + POST >=> routef "/funds/%s/sip/plans" (createSipPlan repository) + GET >=> routef "/funds/%s/sip/plans" (getSipPlans repository) GET >=> routef "/funds/%s" (getFund repository) ] @ (marketData |> Option.map marketDataRoutes |> Option.defaultValue []) 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" |
