summaryrefslogtreecommitdiff
path: root/src/FundLab.Api
diff options
context:
space:
mode:
Diffstat (limited to 'src/FundLab.Api')
-rw-r--r--src/FundLab.Api/App.fs152
-rw-r--r--src/FundLab.Api/FundLab.Api.fsproj5
-rw-r--r--src/FundLab.Api/MarketData.fs294
-rw-r--r--src/FundLab.Api/MarketDataService.fs181
-rw-r--r--src/FundLab.Api/Persistence.fs272
-rw-r--r--src/FundLab.Api/Program.fs4
-rw-r--r--src/FundLab.Api/akshare_collector.py171
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())