namespace FundLab.Api open System open System.Diagnostics open System.IO open System.Runtime.InteropServices open System.Security.Cryptography open System.Text open System.Threading open System.Threading.Tasks type MarketDataFailure = | InvalidMarketDataRequest of string | MarketDataCollectorUnavailable of string | InvalidMarketDataPayload of string | MarketDataPersistenceFailure of string type IMarketDataCollector = abstract Search: query: string * CancellationToken -> Result abstract FetchNav: code: string * CancellationToken -> Result abstract FetchBondQuote: code: string * CancellationToken -> Result abstract FetchStockQuote: code: string * CancellationToken -> Result type IMarketDataService = abstract Search: query: string * CancellationToken -> Result abstract RefreshNav: code: string * CancellationToken -> Result abstract GetNav: code: string * fromDate: DateOnly option * toDate: DateOnly option -> Result module private CollectorNative = let signalKill = 9 [] extern int kill(int pid, int signal) let resolveSetsid () = if OperatingSystem.IsLinux() then [ "/bin/setsid"; "/usr/bin/setsid" ] |> List.tryFind File.Exists else None type ProcessMarketDataCollector(pythonExecutable: string, scriptPath: string, pythonPath: string option, timeoutSeconds: int) = let logCollectorStderr (error: string) = if not (String.IsNullOrWhiteSpace error) then let trimmed = error.Trim() let bounded = if trimmed.Length > 2000 then trimmed.Substring(0, 2000) else trimmed Console.Error.WriteLine("fund-lab collector stderr: {0}", bounded) let execute (token: CancellationToken) arguments = if token.IsCancellationRequested then Error "collector run was cancelled" elif not (File.Exists scriptPath) then Error "collector script is not deployed" else let startInfo = ProcessStartInfo() let setsid = CollectorNative.resolveSetsid() match setsid with | Some setsidPath -> startInfo.FileName <- setsidPath startInfo.ArgumentList.Add(pythonExecutable) | None -> 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) pythonPath |> Option.iter (fun value -> startInfo.Environment["PYTHONPATH"] <- value) use child = new Process() child.StartInfo <- startInfo let ownsProcessGroup = setsid.IsSome let killOwnedGroup () = if ownsProcessGroup then try CollectorNative.kill(-child.Id, CollectorNative.signalKill) |> ignore with _ -> () try if not child.HasExited then child.Kill(entireProcessTree = true) with _ -> () let clock = Stopwatch.StartNew() let deadline = TimeSpan.FromSeconds(float timeoutSeconds) try if not (child.Start()) then Error "could not start the AKShare collector" else let outputTask = child.StandardOutput.ReadToEndAsync(token) let errorTask = child.StandardError.ReadToEndAsync(token) use timeoutSource = new CancellationTokenSource(deadline) use linkedSource = CancellationTokenSource.CreateLinkedTokenSource(token, timeoutSource.Token) let mutable expired = false try child.WaitForExitAsync(linkedSource.Token).GetAwaiter().GetResult() with _ -> if not token.IsCancellationRequested then expired <- true killOwnedGroup () try child.WaitForExit(5000) |> ignore with _ -> () let mutable eofRemaining = deadline - clock.Elapsed let mutable eofExpired = eofRemaining <= TimeSpan.Zero let mutable readsDone = false try if not eofExpired then readsDone <- Task.WhenAll(outputTask, errorTask) .Wait(int eofRemaining.TotalMilliseconds, linkedSource.Token) with _ -> () if not readsDone && not token.IsCancellationRequested then eofExpired <- true if not readsDone then killOwnedGroup () try Task.WhenAll(outputTask, errorTask).Wait(2000) |> ignore with _ -> () let snapshot (task: Task) = try if task.Status = TaskStatus.RanToCompletion then task.Result else "" with _ -> "" let output = snapshot outputTask let error = snapshot errorTask logCollectorStderr error if expired || eofExpired then Error(sprintf "collector did not finish within %d seconds" timeoutSeconds) elif token.IsCancellationRequested then Error "collector run was cancelled" elif child.ExitCode = 0 && not (String.IsNullOrWhiteSpace output) then Ok output elif child.ExitCode <> 0 then Error(sprintf "collector exited with code %d" child.ExitCode) else Error "collector produced no output" with | :? OperationCanceledException -> killOwnedGroup () Error "collector run was cancelled" | error -> killOwnedGroup () Console.Error.WriteLine("fund-lab collector failed to run: {0}", error.Message) Error "could not run the AKShare collector" static member FromEnvironment() = let maxTimeoutSeconds = 600 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) let timeout = let mutable parsed = 0 let raw = Environment.GetEnvironmentVariable("FUND_LAB_COLLECTOR_TIMEOUT_SECONDS") if String.IsNullOrWhiteSpace raw || not (Int32.TryParse(raw, &parsed)) || parsed < 1 || parsed > maxTimeoutSeconds then 30 else parsed ProcessMarketDataCollector(python, script, pythonPath, timeout) interface IMarketDataCollector with member _.Search(query: string, token: CancellationToken) = execute token [ "--operation"; "search"; "--query"; query ] member _.FetchNav(code: string, token: CancellationToken) = execute token [ "--operation"; "nav"; "--code"; code ] member _.FetchBondQuote(code: string, token: CancellationToken) = execute token [ "--operation"; "bond-quote"; "--code"; code ] member _.FetchStockQuote(code: string, token: CancellationToken) = execute token [ "--operation"; "stock-quote"; "--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, token: CancellationToken) = 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, token) 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, token: CancellationToken) = 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, token) 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, token) 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, token) = this.Search(query, token) member this.RefreshNav(code, token) = this.RefreshNav(code, token) member this.GetNav(code, fromDate, toDate) = this.GetNav(code, fromDate, toDate)