From be1afd248e4cc9b03a1e21090602a945af6988b2 Mon Sep 17 00:00:00 2001 From: "Somhairle H. Marisol" Date: Mon, 21 Sep 2026 01:40:10 +0800 Subject: feat(app): 接入 AKShare 行情采集、净值查询与前端净值曲线 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit [变更性质] - 本提交完成 3b 行情闭环:真实 AKShare 搜索与净值采集、PostgreSQL 持久化、鉴权 API 与 Fable 前端净值曲线展示。 [新增功能] - API 新增 /api/instruments/search、/api/instruments/{code}/nav 与 nav/refresh 接口,经 Bearer 鉴权调用 AKShare Python 采集器并幂等落库。 - 前端提供基金搜索、精确代码选择、历史净值刷新/重读与 SVG 折线图,token 变更即清空私有结果并失效在途请求。 [实现方案] - MarketData 以 requiredProperty 校验采集器 payload,净值观测按 (code, nav_date) 幂等 upsert 并保留来源与哈希。 - 前端以 requestId 序列守卫 SearchCompleted/SearchFailed/NavCompleted/NavFailed,TokenChanged 同时递增 searchSeq/navSeq 拒绝过期响应;边界解码兼容 F# option 的 {"case":"Some"} 线格式。 [影响范围] - 新增 FundLab.Web.Tests(边界解码、序列失效、图表纯函数共 8 项),扩展 API 测试至 20 项;vite dev 代理 /api 至本地 API。 - 真实 AKShare 端到端依赖本机 akshare 环境,由负责人另行验收;本提交不包含 docs/overnight-progress.md 的现有修改。 --- src/FundLab.Api/Persistence.fs | 272 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 272 insertions(+) (limited to 'src/FundLab.Api/Persistence.fs') 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(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 @@ -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() + 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" -- cgit v1.2.3