From ea3a5026ad76028562ee4aa6b70c79c94fb17d4b Mon Sep 17 00:00:00 2001 From: "Somhairle H. Marisol" Date: Mon, 21 Sep 2026 03:03:59 +0800 Subject: fix(api): 加固行情采集子进程超时、取消与进程组清理并脱敏错误 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - MarketDataService: setsid 独立会话启动采集器,硬超时覆盖子进程退出与管道 EOF 两个阶段;超时、取消与异常路径按进程组杀灭(含持有管道的后台子进程),不遗留孤儿进程 - App.ts/Search/RefreshNav 传入 ctx.RequestAborted,客户端断开即取消采集 - 对外错误仅返回统一短消息(503 MARKET_DATA_UNAVAILABLE 等),不泄露路径、stderr 或密钥;stderr 仅记录服务端日志 - 新增 9 项子进程回归(成功/stderr 洪泛/超时杀灭/取消杀灭/孤儿管道恢复/畸形载荷/非零退出/启动失败/缺脚本)与 HTTP 脱敏回归;全文套件 66 通过(Domain 19 / Web 17 / API 30) - 实机验证:挂起采集器 503 2.4s 无进程残留;真实 akshare 1.18.96 搜索 000001 返回 200,净值刷新持久化 6010 条观测 - README 与 qa/README 补充环境变量、回归命令、G1 过滤断言口径与真实建档标签辨析 --- src/FundLab.Api/MarketDataService.fs | 174 ++++++++++++++++++++++++++++------- 1 file changed, 143 insertions(+), 31 deletions(-) (limited to 'src/FundLab.Api/MarketDataService.fs') diff --git a/src/FundLab.Api/MarketDataService.fs b/src/FundLab.Api/MarketDataService.fs index 8047071..0c4714d 100644 --- a/src/FundLab.Api/MarketDataService.fs +++ b/src/FundLab.Api/MarketDataService.fs @@ -3,8 +3,11 @@ 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 @@ -13,21 +16,48 @@ type MarketDataFailure = | MarketDataPersistenceFailure of string type IMarketDataCollector = - abstract Search: query: string -> Result - abstract FetchNav: code: string -> Result + abstract Search: query: string * CancellationToken -> Result + abstract FetchNav: code: string * CancellationToken -> Result type IMarketDataService = - abstract Search: query: string -> Result - abstract RefreshNav: code: string -> Result + abstract Search: query: string * CancellationToken -> Result + abstract RefreshNav: code: string * CancellationToken -> Result abstract GetNav: code: string * fromDate: DateOnly option * toDate: DateOnly option -> Result -type ProcessMarketDataCollector(pythonExecutable: string, scriptPath: string, pythonPath: string option) = - let execute arguments = +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 not (File.Exists scriptPath) then - Error(sprintf "collector script was not found at %s" scriptPath) + Error "collector script is not deployed" else let startInfo = ProcessStartInfo() - startInfo.FileName <- pythonExecutable + 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 @@ -38,31 +68,98 @@ type ProcessMarketDataCollector(pythonExecutable: string, scriptPath: string, py for argument in arguments do startInfo.ArgumentList.Add(argument) - match pythonPath with - | Some value -> startInfo.Environment["PYTHONPATH"] <- value - | None -> () + 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 output = child.StandardOutput.ReadToEnd() - let error = child.StandardError.ReadToEnd() - child.WaitForExit() + 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 timedOut = false + + try + child.WaitForExitAsync(linkedSource.Token).GetAwaiter().GetResult() + with _ -> + if not token.IsCancellationRequested then timedOut <- true - if child.ExitCode = 0 && not (String.IsNullOrWhiteSpace output) then + killOwnedGroup () + + try + child.WaitForExit(5000) |> ignore + with _ -> () + + let mutable eofRemaining = deadline - clock.Elapsed + + if eofRemaining <= TimeSpan.Zero then eofRemaining <- TimeSpan.FromSeconds(2.0) + + let mutable readsDone = false + + try + readsDone <- Task.WhenAll(outputTask, errorTask).Wait(int eofRemaining.TotalMilliseconds) + with _ -> () + + 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 timedOut 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 String.IsNullOrWhiteSpace error then + elif child.ExitCode <> 0 then Error(sprintf "collector exited with code %d" child.ExitCode) else - Error(error.Trim()) - with error -> - Error(error.Message) + 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 @@ -77,14 +174,29 @@ type ProcessMarketDataCollector(pythonExecutable: string, scriptPath: string, py |> Option.ofObj |> Option.filter (String.IsNullOrWhiteSpace >> not) - ProcessMarketDataCollector(python, script, pythonPath) + 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) = - execute [ "--operation"; "search"; "--query"; query ] + member _.Search(query: string, token: CancellationToken) = + execute token [ "--operation"; "search"; "--query"; query ] - member _.FetchNav(code: string) = - execute [ "--operation"; "nav"; "--code"; code ] + member _.FetchNav(code: string, token: CancellationToken) = + execute token [ "--operation"; "nav"; "--code"; code ] type MarketDataService(repository: FundRepository, collector: IMarketDataCollector) = let codePattern = Text.RegularExpressions.Regex("^[0-9]{6}$", Text.RegularExpressions.RegexOptions.Compiled) @@ -122,26 +234,26 @@ type MarketDataService(repository: FundRepository, collector: IMarketDataCollect let validCode code = not (String.IsNullOrWhiteSpace code) && codePattern.IsMatch(code) - member _.Search(query: string) = + 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 with + 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) = + 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 with + match collector.Search(normalizedCode, token) with | Error message -> Error(collectorFailure message) | Ok searchRaw -> match parseSearch searchRaw with @@ -153,7 +265,7 @@ type MarketDataService(repository: FundRepository, collector: IMarketDataCollect match persistSearch searchPayload searchHash with | Error failure -> Error failure | Ok _ -> - match collector.FetchNav normalizedCode with + match collector.FetchNav(normalizedCode, token) with | Error message -> Error(collectorFailure message) | Ok navRaw -> match parseNav navRaw with @@ -176,6 +288,6 @@ type MarketDataService(repository: FundRepository, collector: IMarketDataCollect Error(MarketDataPersistenceFailure error.Message) interface IMarketDataService with - member this.Search(query) = this.Search(query) - member this.RefreshNav(code) = this.RefreshNav(code) + 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) -- cgit v1.2.3