diff options
| author | Somhairle H. Marisol <[email protected]> | 2026-09-21 01:40:10 +0800 |
|---|---|---|
| committer | Somhairle H. Marisol <[email protected]> | 2026-09-21 01:40:10 +0800 |
| commit | be1afd248e4cc9b03a1e21090602a945af6988b2 (patch) | |
| tree | d31bd6e6c1736ab4ad784b52ebc5010ade074439 /src/FundLab.Api | |
| parent | 2ada70d6467aec11b45328112d598454ca50f2d6 (diff) | |
| download | fund-lab-be1afd248e4cc9b03a1e21090602a945af6988b2.tar.gz | |
feat(app): 接入 AKShare 行情采集、净值查询与前端净值曲线
[变更性质]
- 本提交完成 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 的现有修改。
Diffstat (limited to 'src/FundLab.Api')
| -rw-r--r-- | src/FundLab.Api/App.fs | 152 | ||||
| -rw-r--r-- | src/FundLab.Api/FundLab.Api.fsproj | 5 | ||||
| -rw-r--r-- | src/FundLab.Api/MarketData.fs | 294 | ||||
| -rw-r--r-- | src/FundLab.Api/MarketDataService.fs | 181 | ||||
| -rw-r--r-- | src/FundLab.Api/Persistence.fs | 272 | ||||
| -rw-r--r-- | src/FundLab.Api/Program.fs | 4 | ||||
| -rw-r--r-- | src/FundLab.Api/akshare_collector.py | 171 |
7 files changed, 1072 insertions, 7 deletions
diff --git a/src/FundLab.Api/App.fs b/src/FundLab.Api/App.fs index 0eb367d..a4e0a00 100644 --- a/src/FundLab.Api/App.fs +++ b/src/FundLab.Api/App.fs @@ -31,11 +31,51 @@ type ApiErrorResponse = message: string } +type MarketDataInstrumentApiResponse = + { + code: string + name: string + fundType: string option + } + +type MarketDataSearchApiResponse = + { + source: string + sourceRevision: string + collectedAt: string + instruments: MarketDataInstrumentApiResponse list + } + +type MarketDataObservationApiResponse = + { + code: string + navDate: string + publishedAt: string option + nav: string + accumulatedNav: string option + dailyReturn: string option + source: string + sourceRevision: string + sourceCollectedAt: string + sourcePayloadHash: string + firstSeenAt: string + lastSeenAt: string + } + +type MarketDataNavApiResponse = + { + code: string + observations: MarketDataObservationApiResponse list + } + module App = let private invariant = CultureInfo.InvariantCulture let private cashText (value: decimal) = value.ToString("0.00", invariant) let private unitNavText (value: decimal) = value.ToString("0.00000000", invariant) + let private decimalText (value: decimal) = value.ToString("0.00000000", invariant) + let private dateText (value: DateOnly) = value.ToString("yyyy-MM-dd", invariant) + let private timestampText (value: DateTimeOffset) = value.ToString("O", invariant) let private fundResponse (fund: FundRecord) : FundApiResponse = { @@ -176,16 +216,116 @@ module App = with _ -> errorResponse 500 "PERSISTENCE_ERROR" "fund persistence failed" next ctx - let createApplication (repository: FundRepository) : HttpHandler = + let private marketDataError (failure: MarketDataFailure) : HttpHandler = + let status, error, message = + match failure with + | InvalidMarketDataRequest message -> 400, "INVALID_MARKET_DATA_REQUEST", message + | MarketDataCollectorUnavailable message -> 503, "MARKET_DATA_UNAVAILABLE", message + | InvalidMarketDataPayload message -> 502, "INVALID_MARKET_DATA_PAYLOAD", message + | MarketDataPersistenceFailure _ -> 500, "PERSISTENCE_ERROR", "market data persistence failed" + + errorResponse status error message + + let private marketDataInstrumentResponse (instrument: MarketDataInstrument) = + { + code = instrument.Code + name = instrument.Name + fundType = instrument.FundType + } + + let private marketDataSearchResponse (payload: MarketDataSearchPayload) = + { + source = payload.Source + sourceRevision = payload.SourceRevision + collectedAt = timestampText payload.CollectedAt + instruments = payload.Instruments |> List.map marketDataInstrumentResponse + } + + let private marketDataObservationResponse (observation: MarketDataNavRecord) = + { + code = observation.Code + navDate = dateText observation.NavDate + publishedAt = observation.PublishedAt |> Option.map timestampText + nav = decimalText observation.Nav + accumulatedNav = observation.AccumulatedNav |> Option.map decimalText + dailyReturn = observation.DailyReturn |> Option.map decimalText + source = observation.Source + sourceRevision = observation.SourceRevision + sourceCollectedAt = timestampText observation.SourceCollectedAt + sourcePayloadHash = observation.SourcePayloadHash + firstSeenAt = timestampText observation.FirstSeenAt + lastSeenAt = timestampText observation.LastSeenAt + } + + let private marketDataNavResponse code observations = + { + code = code + observations = observations |> List.map marketDataObservationResponse + } + + let private searchInstruments (marketData: IMarketDataService) : HttpHandler = + fun next ctx -> + match marketData.Search(ctx.Request.Query["q"].ToString()) with + | Ok payload -> json (marketDataSearchResponse payload) next ctx + | Error failure -> marketDataError failure next ctx + + let private refreshNav (marketData: IMarketDataService) (code: string) : HttpHandler = + fun next ctx -> + match marketData.RefreshNav code with + | Ok observations -> json (marketDataNavResponse code observations) next ctx + | Error failure -> marketDataError failure next ctx + + let private queryDate (name: string) (ctx: HttpContext) = + let value = ctx.Request.Query[name].ToString() + + if String.IsNullOrWhiteSpace value then + Ok None + else + let mutable date = DateOnly.MinValue + + if DateOnly.TryParseExact(value, "yyyy-MM-dd", invariant, DateTimeStyles.None, &date) then + Ok(Some date) + else + Error(sprintf "%s must be an ISO date" name) + + let private getNav (marketData: IMarketDataService) (code: string) : HttpHandler = + fun next ctx -> + match queryDate "from" ctx, queryDate "to" ctx with + | Ok fromDate, Ok toDate -> + match marketData.GetNav(code, fromDate, toDate) with + | Ok observations -> json (marketDataNavResponse code observations) next ctx + | Error failure -> marketDataError failure next ctx + | Error message, _ + | _, Error message -> + marketDataError (InvalidMarketDataRequest message) next ctx + + let private marketDataRoutes (marketData: IMarketDataService) = + [ + GET >=> route "/instruments/search" >=> searchInstruments marketData + POST >=> routef "/instruments/%s/nav/refresh" (refreshNav marketData) + GET >=> routef "/instruments/%s/nav" (getNav marketData) + ] + + let private createApplicationInternal (repository: FundRepository) (marketData: IMarketDataService option) : HttpHandler = + let apiRoutes = + [ + GET >=> route "/portfolio/summary" >=> emptyPortfolio + POST >=> route "/funds" >=> createFund repository + GET >=> routef "/funds/%s" (getFund repository) + ] + @ (marketData |> Option.map marketDataRoutes |> Option.defaultValue []) + choose [ GET >=> route "/health" >=> health subRoute "/api" ( requireBearer - >=> choose [ - GET >=> route "/portfolio/summary" >=> emptyPortfolio - POST >=> route "/funds" >=> createFund repository - GET >=> routef "/funds/%s" (getFund repository) - ] + >=> choose apiRoutes ) setStatusCode 404 >=> text "Not Found" ] + + let createApplicationWithMarketData (repository: FundRepository) (marketData: IMarketDataService) : HttpHandler = + createApplicationInternal repository (Some marketData) + + let createApplication (repository: FundRepository) : HttpHandler = + createApplicationInternal repository None diff --git a/src/FundLab.Api/FundLab.Api.fsproj b/src/FundLab.Api/FundLab.Api.fsproj index 8a3f590..04724c0 100644 --- a/src/FundLab.Api/FundLab.Api.fsproj +++ b/src/FundLab.Api/FundLab.Api.fsproj @@ -15,8 +15,13 @@ <ItemGroup> <Compile Include="Authentication.fs" /> <Compile Include="Health.fs" /> + <Compile Include="MarketData.fs" /> <Compile Include="Persistence.fs" /> + <Compile Include="MarketDataService.fs" /> <Compile Include="App.fs" /> <Compile Include="Program.fs" /> </ItemGroup> + <ItemGroup> + <None Include="akshare_collector.py" CopyToOutputDirectory="PreserveNewest" /> + </ItemGroup> </Project> diff --git a/src/FundLab.Api/MarketData.fs b/src/FundLab.Api/MarketData.fs new file mode 100644 index 0000000..e396a7d --- /dev/null +++ b/src/FundLab.Api/MarketData.fs @@ -0,0 +1,294 @@ +namespace FundLab.Api + +open System +open System.Globalization +open System.Text.Json + +type MarketDataInstrument = + { + Code: string + Name: string + FundType: string option + } + +type MarketDataObservation = + { + NavDate: DateOnly + PublishedAt: DateTimeOffset option + Nav: decimal + AccumulatedNav: decimal option + DailyReturn: decimal option + } + +type MarketDataSearchPayload = + { + Source: string + SourceRevision: string + CollectedAt: DateTimeOffset + Instruments: MarketDataInstrument list + } + +type MarketDataNavPayload = + { + Source: string + SourceRevision: string + CollectedAt: DateTimeOffset + Code: string + Observations: MarketDataObservation list + } + +type MarketDataInstrumentRecord = + { + Code: string + Name: string + FundType: string option + Source: string + SourceRevision: string + SourceCollectedAt: DateTimeOffset + SourcePayloadHash: string + FirstSeenAt: DateTimeOffset + LastSeenAt: DateTimeOffset + } + +type MarketDataNavRecord = + { + Code: string + NavDate: DateOnly + PublishedAt: DateTimeOffset option + Nav: decimal + AccumulatedNav: decimal option + DailyReturn: decimal option + Source: string + SourceRevision: string + SourceCollectedAt: DateTimeOffset + SourcePayloadHash: string + FirstSeenAt: DateTimeOffset + LastSeenAt: DateTimeOffset + } + +module MarketData = + type private ResultBuilder() = + member _.Bind(value, binder) = Result.bind binder value + member _.Return(value) = Ok value + member _.ReturnFrom(value) = value + member _.Zero() : Result<unit, string> = Ok () + member _.Combine(first: Result<unit, string>, second: unit -> Result<'a, string>) = + match first with + | Ok () -> second () + | Error message -> Error message + member _.Delay(generator: unit -> Result<'a, string>) = generator + member _.Run(generator: unit -> Result<'a, string>) = generator () + + let private result = ResultBuilder() + let private invariant = CultureInfo.InvariantCulture + + let private requiredProperty (root: JsonElement) (name: string) = + let mutable property = Unchecked.defaultof<JsonElement> + + if root.ValueKind <> JsonValueKind.Object then + Error "payload root must be a JSON object" + elif root.TryGetProperty(name, &property) then + Ok property + else + Error(sprintf "payload property '%s' is required" name) + + let private requiredString (root: JsonElement) name = + requiredProperty root name + |> Result.bind (fun property -> + if property.ValueKind = JsonValueKind.String then + let value = property.GetString() + + if String.IsNullOrWhiteSpace value then + Error(sprintf "payload property '%s' must not be empty" name) + else + Ok value + else + Error(sprintf "payload property '%s' must be a string" name)) + + let private optionalString (root: JsonElement) name = + requiredProperty root name + |> Result.bind (fun property -> + if property.ValueKind = JsonValueKind.Null then + Ok None + elif property.ValueKind = JsonValueKind.String then + let value = property.GetString() + + if String.IsNullOrWhiteSpace value then + Ok None + else + Ok(Some value) + else + Error(sprintf "payload property '%s' must be null or a string" name)) + + let private validateEnvelope root operation = + result { + let! schemaVersion = requiredString root "schema_version" + let! actualOperation = requiredString root "operation" + let! source = requiredString root "source" + let! sourceRevision = requiredString root "source_revision" + let! collectedAtText = requiredString root "collected_at" + + if schemaVersion <> "fund-lab.akshare.v1" then + return! Error(sprintf "unsupported schema_version '%s'" schemaVersion) + + if actualOperation <> operation then + return! Error(sprintf "payload operation must be '%s'" operation) + + match DateTimeOffset.TryParse(collectedAtText, invariant, DateTimeStyles.RoundtripKind) with + | true, collectedAt -> + return source, sourceRevision, collectedAt + | false, _ -> + return! Error "collected_at must be an ISO-8601 timestamp" + } + + let private isFundCode (value: string) = + value.Length = 6 && value |> Seq.forall Char.IsDigit + + let private parseFundCode root = + requiredString root "code" + |> Result.bind (fun code -> + if isFundCode code then + Ok code + else + Error "fund code must contain exactly six digits") + + let private parseDate text = + let mutable date = DateOnly.MinValue + + if DateOnly.TryParseExact(text, "yyyy-MM-dd", invariant, DateTimeStyles.None, &date) then + Ok date + else + Error "nav_date must be an ISO date" + + let private parseDecimal label (property: JsonElement) = + if property.ValueKind <> JsonValueKind.String then + Error(sprintf "%s must be a decimal string" label) + else + let text = property.GetString() + + if not (String.IsNullOrWhiteSpace text) then + match Decimal.TryParse(text, NumberStyles.AllowLeadingSign ||| NumberStyles.AllowDecimalPoint, invariant) with + | true, value -> Ok value + | false, _ -> Error(sprintf "%s must be a decimal string" label) + else + Error(sprintf "%s must be a decimal string" label) + + let private optionalDecimal label (property: JsonElement) = + if property.ValueKind = JsonValueKind.Null then + Ok None + else + parseDecimal label property |> Result.map Some + + let private optionalPublishedAt (property: JsonElement) = + if property.ValueKind = JsonValueKind.Null then + Ok None + elif property.ValueKind <> JsonValueKind.String then + Error "published_at must be null or an ISO-8601 timestamp" + else + let text = property.GetString() + + if String.IsNullOrWhiteSpace text then + Error "published_at must be null or an ISO-8601 timestamp" + else + match DateTimeOffset.TryParse(text, invariant, DateTimeStyles.RoundtripKind) with + | true, value -> Ok(Some value) + | false, _ -> Error "published_at must be null or an ISO-8601 timestamp" + + let private parseArray parser (property: JsonElement) = + if property.ValueKind <> JsonValueKind.Array then + Error "payload collection must be a JSON array" + else + property.EnumerateArray() + |> Seq.fold + (fun state item -> + match state, parser item with + | Ok values, Ok value -> Ok(value :: values) + | Error message, _ -> Error message + | _, Error message -> Error message) + (Ok []) + |> Result.map List.rev + + let private parseSearchInstrument root = + result { + let! code = parseFundCode root + let! name = requiredString root "name" + let! fundType = optionalString root "fund_type" + + return + { + Code = code + Name = name + FundType = fundType + } + } + + let private parseObservation root = + result { + let! navDateText = requiredString root "nav_date" + let! navDate = parseDate navDateText + let! publishedAtProperty = requiredProperty root "published_at" + let! publishedAt = optionalPublishedAt publishedAtProperty + let! navProperty = requiredProperty root "nav" + let! nav = parseDecimal "nav" navProperty + let! accumulatedNavProperty = requiredProperty root "accumulated_nav" + let! accumulatedNav = optionalDecimal "accumulated_nav" accumulatedNavProperty + let! dailyReturnProperty = requiredProperty root "daily_return" + let! dailyReturn = optionalDecimal "daily_return" dailyReturnProperty + + return + { + NavDate = navDate + PublishedAt = publishedAt + Nav = nav + AccumulatedNav = accumulatedNav + DailyReturn = dailyReturn + } + } + + let parseSearchPayload (json: string) : Result<MarketDataSearchPayload, string> = + try + use document = JsonDocument.Parse(json) + let root = document.RootElement + + result { + let! source, sourceRevision, collectedAt = validateEnvelope root "search" + let! instrumentsProperty = requiredProperty root "instruments" + let! instruments = parseArray parseSearchInstrument instrumentsProperty + + return + { + Source = source + SourceRevision = sourceRevision + CollectedAt = collectedAt + Instruments = instruments + } + } + with + | :? JsonException -> Error "payload must be valid JSON" + + let parseNavPayload (json: string) : Result<MarketDataNavPayload, string> = + try + use document = JsonDocument.Parse(json) + let root = document.RootElement + + result { + let! source, sourceRevision, collectedAt = validateEnvelope root "nav" + let! instrumentProperty = requiredProperty root "instrument" + let! code = parseFundCode instrumentProperty + let! observationsProperty = requiredProperty root "observations" + let! observations = parseArray parseObservation observationsProperty + + if List.isEmpty observations then + return! Error "observations must not be empty" + + return + { + Source = source + SourceRevision = sourceRevision + CollectedAt = collectedAt + Code = code + Observations = observations + } + } + with + | :? JsonException -> Error "payload must be valid JSON" diff --git a/src/FundLab.Api/MarketDataService.fs b/src/FundLab.Api/MarketDataService.fs new file mode 100644 index 0000000..8047071 --- /dev/null +++ b/src/FundLab.Api/MarketDataService.fs @@ -0,0 +1,181 @@ +namespace FundLab.Api + +open System +open System.Diagnostics +open System.IO +open System.Security.Cryptography +open System.Text + +type MarketDataFailure = + | InvalidMarketDataRequest of string + | MarketDataCollectorUnavailable of string + | InvalidMarketDataPayload of string + | MarketDataPersistenceFailure of string + +type IMarketDataCollector = + abstract Search: query: string -> Result<string, string> + abstract FetchNav: code: string -> Result<string, string> + +type IMarketDataService = + abstract Search: query: string -> Result<MarketDataSearchPayload, MarketDataFailure> + abstract RefreshNav: code: string -> Result<MarketDataNavRecord list, MarketDataFailure> + abstract GetNav: code: string * fromDate: DateOnly option * toDate: DateOnly option -> Result<MarketDataNavRecord list, MarketDataFailure> + +type ProcessMarketDataCollector(pythonExecutable: string, scriptPath: string, pythonPath: string option) = + let execute arguments = + if not (File.Exists scriptPath) then + Error(sprintf "collector script was not found at %s" scriptPath) + else + let startInfo = ProcessStartInfo() + startInfo.FileName <- pythonExecutable + startInfo.UseShellExecute <- false + startInfo.RedirectStandardOutput <- true + startInfo.RedirectStandardError <- true + startInfo.WorkingDirectory <- Path.GetDirectoryName(scriptPath) + + startInfo.ArgumentList.Add(scriptPath) + + for argument in arguments do + startInfo.ArgumentList.Add(argument) + + match pythonPath with + | Some value -> startInfo.Environment["PYTHONPATH"] <- value + | None -> () + + use child = new Process() + child.StartInfo <- startInfo + + try + if not (child.Start()) then + Error "could not start the AKShare collector" + else + let output = child.StandardOutput.ReadToEnd() + let error = child.StandardError.ReadToEnd() + child.WaitForExit() + + if child.ExitCode = 0 && not (String.IsNullOrWhiteSpace output) then + Ok output + elif String.IsNullOrWhiteSpace error then + Error(sprintf "collector exited with code %d" child.ExitCode) + else + Error(error.Trim()) + with error -> + Error(error.Message) + + static member FromEnvironment() = + let value name fallback = + Environment.GetEnvironmentVariable(name) + |> Option.ofObj + |> Option.filter (String.IsNullOrWhiteSpace >> not) + |> Option.defaultValue fallback + + let python = value "FUND_LAB_AKSHARE_PYTHON" "python3" + let script = Path.Combine(AppContext.BaseDirectory, "akshare_collector.py") + + let pythonPath = + Environment.GetEnvironmentVariable("FUND_LAB_AKSHARE_PYTHONPATH") + |> Option.ofObj + |> Option.filter (String.IsNullOrWhiteSpace >> not) + + ProcessMarketDataCollector(python, script, pythonPath) + + interface IMarketDataCollector with + member _.Search(query: string) = + execute [ "--operation"; "search"; "--query"; query ] + + member _.FetchNav(code: string) = + execute [ "--operation"; "nav"; "--code"; code ] + +type MarketDataService(repository: FundRepository, collector: IMarketDataCollector) = + let codePattern = Text.RegularExpressions.Regex("^[0-9]{6}$", Text.RegularExpressions.RegexOptions.Compiled) + + let payloadHash (raw: string) = + Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(raw))) + + let collectorFailure error = MarketDataCollectorUnavailable error + let payloadFailure error = InvalidMarketDataPayload error + + let parseSearch raw = + match MarketData.parseSearchPayload raw with + | Ok payload -> Ok(payload, payloadHash raw) + | Error message -> Error(payloadFailure message) + + let parseNav raw = + match MarketData.parseNavPayload raw with + | Ok payload -> Ok(payload, payloadHash raw) + | Error message -> Error(payloadFailure message) + + let persistSearch (payload: MarketDataSearchPayload) hash = + try + repository.UpsertInstruments(payload, hash) + Ok payload + with error -> + Error(MarketDataPersistenceFailure error.Message) + + let persistNav (payload: MarketDataNavPayload) hash = + try + repository.UpsertNavObservations(payload, hash) + Ok(repository.GetNav(payload.Code, None, None)) + with error -> + Error(MarketDataPersistenceFailure error.Message) + + let validCode code = + not (String.IsNullOrWhiteSpace code) && codePattern.IsMatch(code) + + member _.Search(query: string) = + let normalizedQuery = if isNull query then "" else query.Trim() + + if String.IsNullOrWhiteSpace normalizedQuery || normalizedQuery.Length > 80 then + Error(InvalidMarketDataRequest "search query must contain 1 to 80 characters") + else + match collector.Search normalizedQuery with + | Error message -> Error(collectorFailure message) + | Ok raw -> + match parseSearch raw with + | Error failure -> Error failure + | Ok(payload, hash) -> persistSearch payload hash + + member _.RefreshNav(code: string) = + let normalizedCode = if isNull code then "" else code.Trim() + + if not (validCode normalizedCode) then + Error(InvalidMarketDataRequest "fund code must contain exactly six digits") + else + match collector.Search normalizedCode with + | Error message -> Error(collectorFailure message) + | Ok searchRaw -> + match parseSearch searchRaw with + | Error failure -> Error failure + | Ok(searchPayload, searchHash) -> + match searchPayload.Instruments |> List.tryFind (fun instrument -> instrument.Code = normalizedCode) with + | None -> Error(InvalidMarketDataRequest "fund code was not found in the source catalog") + | Some _ -> + match persistSearch searchPayload searchHash with + | Error failure -> Error failure + | Ok _ -> + match collector.FetchNav normalizedCode with + | Error message -> Error(collectorFailure message) + | Ok navRaw -> + match parseNav navRaw with + | Error failure -> Error failure + | Ok(navPayload, navHash) when navPayload.Code <> normalizedCode -> + Error(payloadFailure "NAV payload code does not match the requested fund code") + | Ok(navPayload, navHash) -> persistNav navPayload navHash + + member _.GetNav(code: string, fromDate: DateOnly option, toDate: DateOnly option) = + let normalizedCode = if isNull code then "" else code.Trim() + + if not (validCode normalizedCode) then + Error(InvalidMarketDataRequest "fund code must contain exactly six digits") + elif fromDate.IsSome && toDate.IsSome && fromDate.Value > toDate.Value then + Error(InvalidMarketDataRequest "from date must not be after to date") + else + try + Ok(repository.GetNav(normalizedCode, fromDate, toDate)) + with error -> + Error(MarketDataPersistenceFailure error.Message) + + interface IMarketDataService with + member this.Search(query) = this.Search(query) + member this.RefreshNav(code) = this.RefreshNav(code) + member this.GetNav(code, fromDate, toDate) = this.GetNav(code, fromDate, toDate) 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" diff --git a/src/FundLab.Api/Program.fs b/src/FundLab.Api/Program.fs index 4c757d6..d4f2c76 100644 --- a/src/FundLab.Api/Program.fs +++ b/src/FundLab.Api/Program.fs @@ -16,8 +16,10 @@ let main argv = builder.Services.AddGiraffe() |> ignore let repository = FundRepository(requiredEnvironment "FUND_LAB_DATABASE_URL") repository.EnsureSchema() + let collector = ProcessMarketDataCollector.FromEnvironment() :> IMarketDataCollector + let marketData = MarketDataService(repository, collector) :> IMarketDataService let app = builder.Build() - app.UseGiraffe(App.createApplication(repository)) + app.UseGiraffe(App.createApplicationWithMarketData repository marketData) app.Run() 0 diff --git a/src/FundLab.Api/akshare_collector.py b/src/FundLab.Api/akshare_collector.py new file mode 100644 index 0000000..f988a61 --- /dev/null +++ b/src/FundLab.Api/akshare_collector.py @@ -0,0 +1,171 @@ +#!/usr/bin/env python3 +"""Small raw-data adapter for the AKShare endpoints used by Fund Lab.""" + +import argparse +import datetime as dt +import json +import re +import sys +from decimal import Decimal, InvalidOperation + +import akshare as ak +import pandas as pd + + +SCHEMA_VERSION = "fund-lab.akshare.v1" + + +def collected_at(): + return dt.datetime.now(dt.timezone.utc).isoformat().replace("+00:00", "Z") + + +def source_revision(): + return f"akshare-{getattr(ak, '__version__', 'unknown')}/eastmoney" + + +def is_missing(value): + try: + return bool(pd.isna(value)) + except (TypeError, ValueError): + return False + + +def text(value): + if is_missing(value): + return None + value = str(value).strip() + return value or None + + +def fund_code(value): + value = text(value) + if value is None: + return None + if value.isdigit() and len(value) < 6: + value = value.zfill(6) + return value if re.fullmatch(r"\d{6}", value) else None + + +def decimal_text(value): + value = text(value) + if value is None: + return None + try: + return format(Decimal(value), "f") + except InvalidOperation: + return None + + +def date_text(value): + if is_missing(value): + return None + if isinstance(value, (dt.datetime, dt.date)): + return value.strftime("%Y-%m-%d") + parsed = pd.to_datetime(value, errors="coerce") + if pd.isna(parsed): + return None + return parsed.strftime("%Y-%m-%d") + + +def search(query): + query = query.strip() + if not query: + raise ValueError("search query cannot be empty") + + frame = ak.fund_name_em() + query_lower = query.casefold() + rows = [] + + for _, row in frame.iterrows(): + code = fund_code(row.get("基金代码")) + name = text(row.get("基金简称")) + pinyin = text(row.get("拼音缩写")) + full_pinyin = text(row.get("拼音全称")) + + if code is None or name is None: + continue + + searchable = [code, name, pinyin or "", full_pinyin or ""] + if not any(query_lower in value.casefold() for value in searchable): + continue + + rows.append( + { + "code": code, + "name": name, + "fund_type": text(row.get("基金类型")), + "rank": 0 if code == query else (1 if name == query else 2), + } + ) + + rows.sort(key=lambda item: (item.pop("rank"), item["code"])) + return { + "schema_version": SCHEMA_VERSION, + "operation": "search", + "source": "akshare", + "source_revision": source_revision(), + "collected_at": collected_at(), + "instruments": rows[:100], + } + + +def nav(code): + if not re.fullmatch(r"\d{6}", code): + raise ValueError("fund code must contain exactly six digits") + + frame = ak.fund_open_fund_info_em( + symbol=code, + indicator="单位净值走势", + period="成立来", + ) + observations = [] + + for _, row in frame.iterrows(): + nav_date = date_text(row.get("净值日期")) + nav_value = decimal_text(row.get("单位净值")) + if nav_date is None or nav_value is None: + continue + + observations.append( + { + "nav_date": nav_date, + "published_at": None, + "nav": nav_value, + "accumulated_nav": None, + "daily_return": decimal_text(row.get("日增长率")), + } + ) + + if not observations: + raise ValueError(f"AKShare returned no usable NAV observations for {code}") + + return { + "schema_version": SCHEMA_VERSION, + "operation": "nav", + "source": "akshare", + "source_revision": source_revision(), + "collected_at": collected_at(), + "instrument": {"code": code}, + "observations": observations, + } + + +def main(): + parser = argparse.ArgumentParser() + parser.add_argument("--operation", choices=("search", "nav"), required=True) + parser.add_argument("--query") + parser.add_argument("--code") + args = parser.parse_args() + + try: + payload = search(args.query) if args.operation == "search" else nav(args.code) + json.dump(payload, sys.stdout, ensure_ascii=False, separators=(",", ":")) + sys.stdout.write("\n") + return 0 + except Exception as error: + print(f"AKShare collector failed: {error}", file=sys.stderr) + return 2 + + +if __name__ == "__main__": + raise SystemExit(main()) |
