summaryrefslogtreecommitdiff
path: root/src/FundLab.Api/Persistence.fs
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-22 10:24:58 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-22 10:24:58 +0800
commit2677ea0053ba9d6b34fec278ad52657da390b30c (patch)
treeb754cb7ba256a527147789ee6fb177bf29d58bf5 /src/FundLab.Api/Persistence.fs
parent1262eea047d4d427ab37a5bf2003c2ab82b0d11c (diff)
downloadfund-lab-2677ea0053ba9d6b34fec278ad52657da390b30c.tar.gz
Fix stock sip plan instrument FK (3d-31 B4a)
Diffstat (limited to 'src/FundLab.Api/Persistence.fs')
-rw-r--r--src/FundLab.Api/Persistence.fs418
1 files changed, 337 insertions, 81 deletions
diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs
index b22fbd9..d1b8250 100644
--- a/src/FundLab.Api/Persistence.fs
+++ b/src/FundLab.Api/Persistence.fs
@@ -289,6 +289,16 @@ type SipPlanCommand =
Frequency: SipFrequency
}
+/// Minimal live-quote view used to drive a stock SIP period. The resolver is
+/// supplied by the API layer (which owns the market-data probe); the repository
+/// only interprets the snapshot.
+type StockSipQuote =
+ {
+ Price: decimal option
+ Name: string option
+ Suspended: bool option
+ }
+
type StockTradeCommand =
{
InstrumentCode: string
@@ -325,6 +335,7 @@ type StockTradeWriteResult =
| StockTradeReplayed of StockTradeRecord
| StockTradeIdempotencyConflict
| StockTradeInvalid of string
+ | StockTradeInsufficientFunds of string
| StockTradeFundNotFound
type StockSellCommand =
@@ -549,6 +560,7 @@ type SipPlanRecord =
Id: Guid
FundId: Guid
InstrumentCode: string
+ AssetClass: string
Amount: decimal
Frequency: SipFrequency
Status: string
@@ -584,6 +596,7 @@ type SipPlanAdvanceResult =
{
PlanId: Guid
InstrumentCode: string
+ AssetClass: string
Amount: decimal
Frequency: SipFrequency
Executions: SipExecutionOutcome list
@@ -1019,6 +1032,10 @@ type FundRepository(connectionString: string) =
created_at timestamptz NOT NULL DEFAULT now()
);
+ ALTER TABLE sip_plans ADD COLUMN IF NOT EXISTS asset_class text NOT NULL DEFAULT 'fund';
+
+ ALTER TABLE sip_plans DROP CONSTRAINT IF EXISTS sip_plans_instrument_code_fkey;
+
CREATE TABLE IF NOT EXISTS sip_plan_idempotencies (
idempotency_key text PRIMARY KEY,
request_hash text NOT NULL,
@@ -2986,19 +3003,20 @@ type FundRepository(connectionString: string) =
Id = reader.GetGuid(0)
FundId = reader.GetGuid(1)
InstrumentCode = reader.GetString(2)
- Amount = reader.GetDecimal(3)
+ AssetClass = reader.GetString(3)
+ Amount = reader.GetDecimal(4)
Frequency =
- match SipPolicy.parseFrequency (reader.GetString(4)) with
+ match SipPolicy.parseFrequency (reader.GetString(5)) 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)
- LastExecutionStatus = readStringOption reader 10
+ Status = reader.GetString(6)
+ IsSynthetic = reader.GetBoolean(7)
+ AnchorDate = reader.GetFieldValue<DateOnly>(8)
+ NextTradeDate = reader.GetFieldValue<DateOnly>(9)
+ CreatedAt = reader.GetFieldValue<DateTimeOffset>(10)
+ LastExecutionStatus = readStringOption reader 11
LastExecutionDate =
- if reader.IsDBNull(11) then None else Some(reader.GetFieldValue<DateOnly>(11))
+ if reader.IsDBNull(12) then None else Some(reader.GetFieldValue<DateOnly>(12))
}
let findSipPlan connection transaction planId =
@@ -3007,7 +3025,7 @@ type FundRepository(connectionString: string) =
connection
transaction
"""
- SELECT p.id, p.fund_id, p.instrument_code, p.amount, p.frequency, p.status, p.is_synthetic,
+ SELECT p.id, p.fund_id, p.instrument_code, p.asset_class, p.amount, p.frequency, p.status, p.is_synthetic,
p.anchor_date, p.next_trade_date, p.created_at,
e.status, e.trade_date
FROM sip_plans p
@@ -3044,14 +3062,15 @@ type FundRepository(connectionString: string) =
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)
+ INSERT INTO sip_plans (id, fund_id, instrument_code, asset_class, amount, frequency, status, is_synthetic, anchor_date, next_trade_date)
+ VALUES (@id, @fund_id, @instrument_code, @asset_class, @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 "asset_class" NpgsqlDbType.Text (box plan.AssetClass) |> 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
@@ -3079,7 +3098,7 @@ type FundRepository(connectionString: string) =
addParameter command "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore
command.ExecuteNonQuery() |> ignore
- let sipPlanRequestHash (fundId: Guid) (command: SipPlanCommand) =
+ let sipPlanRequestHash (fundId: Guid) (assetClass: string) (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
@@ -3090,6 +3109,7 @@ type FundRepository(connectionString: string) =
[
"sip-plan"
encoded (fundId.ToString("D"))
+ encoded assetClass
encoded code
(encoded (command.Amount.ToString("G29", invariant)))
(encoded (SipPolicy.frequencyText command.Frequency))
@@ -3978,14 +3998,21 @@ type FundRepository(connectionString: string) =
records |> Seq.toList
- member _.CreateSipPlan(idempotencyKey: string, fundId: Guid, command: SipPlanCommand, ?anchorOverride: DateOnly) : SipPlanWriteResult =
- if String.IsNullOrWhiteSpace idempotencyKey then
+ member _.CreateSipPlan(idempotencyKey: string, fundId: Guid, command: SipPlanCommand, ?anchorOverride: DateOnly, ?assetClass: string) : SipPlanWriteResult =
+ let resolvedAssetClass =
+ match assetClass with
+ | Some value when not (String.IsNullOrWhiteSpace value) -> value.Trim().ToLowerInvariant()
+ | _ -> "fund"
+
+ if resolvedAssetClass <> "fund" && resolvedAssetClass <> "stock" then
+ SipPlanWriteResult.SipPlanInvalid "assetClass must be fund or stock"
+ elif 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
+ let fingerprint = sipPlanRequestHash fundId resolvedAssetClass command
use connection = new NpgsqlConnection(connectionString)
connection.Open()
use transaction = connection.BeginTransaction(IsolationLevel.ReadCommitted)
@@ -4019,7 +4046,7 @@ type FundRepository(connectionString: string) =
transaction.Rollback()
SipPlanWriteResult.SipPlanFundNotFound
| Some isSynthetic ->
- if instrumentExists connection (Some transaction) command.InstrumentCode then
+ if resolvedAssetClass = "stock" || instrumentExists connection (Some transaction) command.InstrumentCode then
let anchorDate =
anchorOverride
|> Option.defaultWith (fun () -> ConfirmationPolicy.tradeDateFor DateTimeOffset.UtcNow)
@@ -4029,6 +4056,7 @@ type FundRepository(connectionString: string) =
Id = Guid.NewGuid()
FundId = fundId
InstrumentCode = command.InstrumentCode
+ AssetClass = resolvedAssetClass
Amount = command.Amount
Frequency = command.Frequency
Status = "active"
@@ -4055,7 +4083,7 @@ type FundRepository(connectionString: string) =
raise error
- member this.AdvanceSipPlans(fundId: Guid, endDate: DateOnly, limit: int) : SipAdvanceResult =
+ member this.AdvanceSipPlans(fundId: Guid, endDate: DateOnly, limit: int, ?stockQuote: string -> Result<StockSipQuote, string>) : SipAdvanceResult =
use connection = new NpgsqlConnection(connectionString)
connection.Open()
@@ -4121,23 +4149,24 @@ type FundRepository(connectionString: string) =
commandWithTransaction
connection
(Some transaction)
- "SELECT id, instrument_code, amount, frequency, anchor_date, next_trade_date FROM sip_plans WHERE fund_id = @fund_id AND status = 'active' ORDER BY created_at"
+ "SELECT id, instrument_code, asset_class, amount, frequency, anchor_date, next_trade_date FROM sip_plans WHERE fund_id = @fund_id AND status = 'active' ORDER BY created_at"
addParameter plansCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore
use plansReader = plansCommand.ExecuteReader()
- let plans = ResizeArray<Guid * string * decimal * SipFrequency * DateOnly * DateOnly>()
+ let plans = ResizeArray<Guid * string * string * decimal * SipFrequency * DateOnly * DateOnly>()
while plansReader.Read() do
plans.Add(
plansReader.GetGuid(0),
plansReader.GetString(1),
- plansReader.GetDecimal(2),
- (match SipPolicy.parseFrequency (plansReader.GetString(3)) with
+ plansReader.GetString(2),
+ plansReader.GetDecimal(3),
+ (match SipPolicy.parseFrequency (plansReader.GetString(4)) with
| Some frequency -> frequency
| None -> failwith "sip plan frequency is invalid"),
- plansReader.GetFieldValue<DateOnly>(4),
- plansReader.GetFieldValue<DateOnly>(5)
+ plansReader.GetFieldValue<DateOnly>(5),
+ plansReader.GetFieldValue<DateOnly>(6)
)
plansReader.Close()
@@ -4172,7 +4201,204 @@ type FundRepository(connectionString: string) =
rows |> Seq.toList
- let advancePlanRow (planId: Guid, code: string, amount: decimal, frequency: SipFrequency, anchor: DateOnly, nextDate: DateOnly) =
+ let rollPlanPointer (planId: Guid) (nextDate: DateOnly) (frequency: SipFrequency) (anchor: DateOnly) (dueDates: DateOnly list) =
+ let rolled =
+ match dueDates with
+ | [] -> nextDate
+ | _ -> SipPolicy.nextTradeDate frequency anchor ((List.last dueDates).AddDays 1)
+
+ if not (List.isEmpty dueDates) then
+ use rollCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ "UPDATE sip_plans SET next_trade_date = @next_date WHERE id = @plan_id"
+
+ addParameter rollCommand "next_date" NpgsqlDbType.Date (box rolled) |> ignore
+ addParameter rollCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore
+ rollCommand.ExecuteNonQuery() |> ignore
+
+ rolled
+
+ // Stock SIP: schedule, lot sizing and idempotency mirror the fund path, but each
+ // due period buys the target stock through the shared stock-trade pipeline using a
+ // live quote (suspension and missing prices stop the period).
+ let advanceStockPlanRow (planId: Guid, code: string, assetClass: string, amount: decimal, frequency: SipFrequency, anchor: DateOnly, nextDate: DateOnly) =
+ let dueDates =
+ SipPolicy.advancePlan frequency anchor nextDate endDate
+ |> List.truncate (max 1 limit)
+
+ let outcomes = ResizeArray<SipExecutionOutcome>()
+
+ let upsertExecution (tradeDate: DateOnly) =
+ use insertCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ INSERT INTO sip_executions (plan_id, trade_date, amount, fee_amount, status)
+ VALUES (@plan_id, @trade_date, @amount, @fee_amount, 'processing')
+ ON CONFLICT (plan_id, trade_date) DO NOTHING
+ RETURNING trade_date
+ """
+
+ addParameter insertCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore
+ addParameter insertCommand "trade_date" NpgsqlDbType.Date (box tradeDate) |> ignore
+ addParameter insertCommand "amount" NpgsqlDbType.Numeric (box amount) |> ignore
+ addParameter insertCommand "fee_amount" NpgsqlDbType.Numeric (box 0m) |> ignore
+
+ use reader = insertCommand.ExecuteReader()
+ let inserted = reader.Read()
+ reader.Close()
+ inserted
+
+ let setExecution (tradeDate: DateOnly) (status: string) (orderId: Guid option) (reason: string option) =
+ use updateCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ UPDATE sip_executions
+ SET status = @status,
+ order_id = @order_id,
+ pending_reason = @reason,
+ executed_at = now()
+ WHERE plan_id = @plan_id AND trade_date = @trade_date
+ """
+
+ addParameter updateCommand "status" NpgsqlDbType.Text (box status) |> ignore
+
+ let orderParameter =
+ match orderId with
+ | Some value -> box value
+ | None -> box DBNull.Value
+
+ addParameter updateCommand "order_id" NpgsqlDbType.Uuid orderParameter |> ignore
+
+ let reasonParameter =
+ match reason with
+ | Some value -> box value
+ | None -> box DBNull.Value
+
+ addParameter updateCommand "reason" NpgsqlDbType.Text reasonParameter |> ignore
+ addParameter updateCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore
+ addParameter updateCommand "trade_date" NpgsqlDbType.Date (box tradeDate) |> ignore
+ updateCommand.ExecuteNonQuery() |> ignore
+
+ let findExistingExecution (tradeDate: DateOnly) =
+ use selectCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ "SELECT status, order_id, pending_reason FROM sip_executions WHERE plan_id = @plan_id AND trade_date = @trade_date"
+
+ addParameter selectCommand "plan_id" NpgsqlDbType.Uuid (box planId) |> ignore
+ addParameter selectCommand "trade_date" NpgsqlDbType.Date (box tradeDate) |> ignore
+
+ use reader = selectCommand.ExecuteReader()
+
+ if reader.Read() then
+ Some
+ (reader.GetString(0),
+ (if reader.IsDBNull(1) then None else Some(reader.GetGuid(1))),
+ readStringOption reader 2)
+ else
+ None
+
+ let settle (tradeDate: DateOnly) (status: string) (orderId: Guid option) (reason: string option) =
+ setExecution tradeDate status orderId reason
+
+ outcomes.Add(
+ {
+ TradeDate = tradeDate
+ Status = status
+ OrderId = orderId
+ PendingReason = reason
+ }
+ )
+
+ for tradeDate in dueDates do
+ if upsertExecution tradeDate then
+ let quoteResult =
+ match stockQuote with
+ | Some resolve -> resolve code
+ | None -> Error "stock quote probe is not configured"
+
+ match quoteResult with
+ | Error reason -> settle tradeDate "market_data_unavailable" None (Some reason)
+ | Ok quote ->
+ match quote.Suspended with
+ | Some true -> settle tradeDate "suspended" None (Some "stock is suspended and cannot be traded")
+ | _ ->
+ match quote.Price with
+ | Some price when price > 0m ->
+ let lot = StockTerms.aShareDefault.MinUnit
+ let lots = Decimal.Floor(amount / (price * lot))
+ let quantity = lots * lot
+
+ if quantity <= 0m then
+ settle tradeDate "insufficient_cash" None (Some "scheduled amount is below one board lot at the current price")
+ else
+ let tradeKey = sprintf "sip:%O:%s" planId (tradeDate.ToString("yyyy-MM-dd"))
+ let executedAt = DateTimeOffset(tradeDate.ToDateTime(TimeOnly.MinValue), TimeSpan.Zero)
+
+ let command: StockTradeCommand =
+ {
+ InstrumentCode = code
+ StockName = quote.Name
+ Quantity = quantity
+ Price = price
+ }
+
+ match this.CreateStockTrade(tradeKey, fundId, command, executedAt, true) with
+ | StockTradeWriteResult.StockTradeCreated trade
+ | StockTradeWriteResult.StockTradeReplayed trade ->
+ settle tradeDate "succeeded" (Some trade.Id) None
+ | StockTradeWriteResult.StockTradeInsufficientFunds reason ->
+ settle tradeDate "insufficient_cash" None (Some reason)
+ | StockTradeWriteResult.StockTradeInvalid reason -> settle tradeDate "failed" None (Some reason)
+ | StockTradeWriteResult.StockTradeIdempotencyConflict ->
+ settle tradeDate "failed" None (Some "idempotency key was used with a different request")
+ | StockTradeWriteResult.StockTradeFundNotFound ->
+ settle tradeDate "failed" None (Some "fund was not found")
+ | _ -> settle tradeDate "market_data_unavailable" None (Some "stock quote did not include a price")
+ else
+ match findExistingExecution tradeDate with
+ | Some(status, orderId, reason) ->
+ outcomes.Add(
+ {
+ TradeDate = tradeDate
+ Status = status
+ OrderId = orderId
+ PendingReason = reason
+ }
+ )
+ | None ->
+ outcomes.Add(
+ {
+ TradeDate = tradeDate
+ Status = "unknown"
+ OrderId = None
+ PendingReason = None
+ }
+ )
+
+ let rolled = rollPlanPointer planId nextDate frequency anchor dueDates
+
+ let replayedOutcomes =
+ if List.isEmpty dueDates then readExecutionsUpTo planId endDate else outcomes |> Seq.toList
+
+ {
+ PlanId = planId
+ InstrumentCode = code
+ AssetClass = assetClass
+ Amount = amount
+ Frequency = frequency
+ Executions = replayedOutcomes
+ NextTradeDate = rolled
+ }
+
+ let advanceFundPlanRow (planId: Guid, code: string, assetClass: string, amount: decimal, frequency: SipFrequency, anchor: DateOnly, nextDate: DateOnly) =
let dueDates =
SipPolicy.advancePlan frequency anchor nextDate endDate
|> List.truncate (max 1 limit)
@@ -4323,12 +4549,19 @@ type FundRepository(connectionString: string) =
{
PlanId = planId
InstrumentCode = code
+ AssetClass = assetClass
Amount = amount
Frequency = frequency
Executions = replayedOutcomes
NextTradeDate = rolled
}
+ let advancePlanRow (planId: Guid, code: string, assetClass: string, amount: decimal, frequency: SipFrequency, anchor: DateOnly, nextDate: DateOnly) =
+ if assetClass = "stock" then
+ advanceStockPlanRow (planId, code, assetClass, amount, frequency, anchor, nextDate)
+ else
+ advanceFundPlanRow (planId, code, assetClass, amount, frequency, anchor, nextDate)
+
let planResults = plans |> Seq.map advancePlanRow |> Seq.toList
transaction.Commit()
@@ -4350,7 +4583,7 @@ type FundRepository(connectionString: string) =
connection
None
"""
- SELECT p.id, p.fund_id, p.instrument_code, p.amount, p.frequency, p.status, p.is_synthetic,
+ SELECT p.id, p.fund_id, p.instrument_code, p.asset_class, p.amount, p.frequency, p.status, p.is_synthetic,
p.anchor_date, p.next_trade_date, p.created_at,
e.status, e.trade_date
FROM sip_plans p
@@ -5512,7 +5745,7 @@ type FundRepository(connectionString: string) =
raise error
- member _.CreateStockTrade(idempotencyKey: string, fundId: Guid, command: StockTradeCommand, ?executedAtOverride: DateTimeOffset) : StockTradeWriteResult =
+ member _.CreateStockTrade(idempotencyKey: string, fundId: Guid, command: StockTradeCommand, ?executedAtOverride: DateTimeOffset, ?debitAvailableCash: bool) : StockTradeWriteResult =
if String.IsNullOrWhiteSpace idempotencyKey then
StockTradeWriteResult.StockTradeInvalid "idempotency key cannot be empty"
else
@@ -5563,68 +5796,91 @@ type FundRepository(connectionString: string) =
let executedAt = defaultArg executedAtOverride DateTimeOffset.UtcNow
let costCash = Decimal.Round(normalized.Quantity * normalized.Price, 2, MidpointRounding.AwayFromZero)
- let trade: StockTradeRecord =
- {
- Id = Guid.NewGuid()
- FundId = fundId
- InstrumentCode = normalized.InstrumentCode
- StockName = normalized.StockName
- Quantity = normalized.Quantity
- Price = normalized.Price
- CostCash = costCash
- IsSynthetic = isSynthetic
- ExecutedAt = executedAt
- }
+ let debited =
+ if defaultArg debitAvailableCash false then
+ use cashCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ "UPDATE funds SET available_cash = available_cash - @cost WHERE id = @fund_id AND available_cash >= @cost"
- insertStockTrade connection (Some transaction) trade
- insertStockTradeIdempotency connection (Some transaction) idempotencyKey fingerprint trade.Id fundId
+ addParameter cashCommand "cost" NpgsqlDbType.Numeric (box costCash) |> ignore
+ addParameter cashCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore
+ cashCommand.ExecuteNonQuery() = 1
+ else
+ true
- use positionCommand =
- commandWithTransaction
- connection
- (Some transaction)
- """
- INSERT INTO stock_positions
- (fund_id, instrument_code, stock_name, quantity, cost_cash, last_traded_at)
- VALUES (@fund_id, @code, @name, @quantity, @cost_cash, @last_traded_at)
- ON CONFLICT (fund_id, instrument_code) DO UPDATE
- SET quantity = stock_positions.quantity + EXCLUDED.quantity,
- cost_cash = stock_positions.cost_cash + EXCLUDED.cost_cash,
- stock_name = COALESCE(EXCLUDED.stock_name, stock_positions.stock_name),
- last_traded_at = EXCLUDED.last_traded_at
- """
+ if not debited then
+ transaction.Rollback()
- addParameter positionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore
- addParameter positionCommand "code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore
+ StockTradeWriteResult.StockTradeInsufficientFunds(
+ sprintf
+ "available cash is not enough for a stock purchase of %s"
+ (costCash.ToString("0.00", CultureInfo.InvariantCulture))
+ )
+ else
+ let trade: StockTradeRecord =
+ {
+ Id = Guid.NewGuid()
+ FundId = fundId
+ InstrumentCode = normalized.InstrumentCode
+ StockName = normalized.StockName
+ Quantity = normalized.Quantity
+ Price = normalized.Price
+ CostCash = costCash
+ IsSynthetic = isSynthetic
+ ExecutedAt = executedAt
+ }
- let nameParameter =
- match normalized.StockName with
- | Some name -> box name
- | None -> box DBNull.Value
+ insertStockTrade connection (Some transaction) trade
+ insertStockTradeIdempotency connection (Some transaction) idempotencyKey fingerprint trade.Id fundId
- addParameter positionCommand "name" NpgsqlDbType.Text nameParameter |> ignore
- addParameter positionCommand "quantity" NpgsqlDbType.Numeric (box normalized.Quantity) |> ignore
- addParameter positionCommand "cost_cash" NpgsqlDbType.Numeric (box costCash) |> ignore
- addParameter positionCommand "last_traded_at" NpgsqlDbType.TimestampTz (box executedAt) |> ignore
- positionCommand.ExecuteNonQuery() |> ignore
+ use positionCommand =
+ commandWithTransaction
+ connection
+ (Some transaction)
+ """
+ INSERT INTO stock_positions
+ (fund_id, instrument_code, stock_name, quantity, cost_cash, last_traded_at)
+ VALUES (@fund_id, @code, @name, @quantity, @cost_cash, @last_traded_at)
+ ON CONFLICT (fund_id, instrument_code) DO UPDATE
+ SET quantity = stock_positions.quantity + EXCLUDED.quantity,
+ cost_cash = stock_positions.cost_cash + EXCLUDED.cost_cash,
+ stock_name = COALESCE(EXCLUDED.stock_name, stock_positions.stock_name),
+ last_traded_at = EXCLUDED.last_traded_at
+ """
- let buyEvent =
- stockCashflowEvent
- fundId
- normalized.InstrumentCode
- normalized.StockName
- "buy"
- (DateOnly.FromDateTime executedAt.UtcDateTime)
- normalized.Quantity
- costCash
- None
- isSynthetic
+ addParameter positionCommand "fund_id" NpgsqlDbType.Uuid (box fundId) |> ignore
+ addParameter positionCommand "code" NpgsqlDbType.Text (box normalized.InstrumentCode) |> ignore
+
+ let nameParameter =
+ match normalized.StockName with
+ | Some name -> box name
+ | None -> box DBNull.Value
+
+ addParameter positionCommand "name" NpgsqlDbType.Text nameParameter |> ignore
+ addParameter positionCommand "quantity" NpgsqlDbType.Numeric (box normalized.Quantity) |> ignore
+ addParameter positionCommand "cost_cash" NpgsqlDbType.Numeric (box costCash) |> ignore
+ 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
+ insertStockCashflow connection (Some transaction) buyEvent
- let persisted = { trade with ExecutedAt = executedAt }
- transaction.Commit()
- StockTradeWriteResult.StockTradeCreated persisted
+ let persisted = { trade with ExecutedAt = executedAt }
+ transaction.Commit()
+ StockTradeWriteResult.StockTradeCreated persisted
with error ->
try
transaction.Rollback()