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.fs272
1 files changed, 272 insertions, 0 deletions
diff --git a/src/FundLab.Api/Persistence.fs b/src/FundLab.Api/Persistence.fs
index db0fb6f..e48d9b1 100644
--- a/src/FundLab.Api/Persistence.fs
+++ b/src/FundLab.Api/Persistence.fs
@@ -61,6 +61,37 @@ type FundRepository(connectionString: string) =
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 =
@@ -93,6 +124,44 @@ type FundRepository(connectionString: string) =
Status = reader.GetString(7)
}
+ let dateTimeOffsetFromReader (reader: DbDataReader) index =
+ reader.GetFieldValue<DateTimeOffset>(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<DateOnly>(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
@@ -108,6 +177,11 @@ type FundRepository(connectionString: string) =
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
@@ -176,6 +250,102 @@ type FundRepository(connectionString: string) =
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
@@ -222,6 +392,108 @@ type FundRepository(connectionString: string) =
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<string>()
+ 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<MarketDataNavRecord>()
+
+ 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"