diff options
| -rw-r--r-- | README.md | 13 | ||||
| -rw-r--r-- | qa/README.md | 4 | ||||
| -rw-r--r-- | src/FundLab.Api/App.fs | 4 | ||||
| -rw-r--r-- | src/FundLab.Api/MarketDataService.fs | 174 | ||||
| -rw-r--r-- | tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj | 1 | ||||
| -rw-r--r-- | tests/FundLab.Api.Tests/PersistenceTests.fs | 51 | ||||
| -rw-r--r-- | tests/FundLab.Api.Tests/ProcessCollectorTests.fs | 232 |
7 files changed, 441 insertions, 38 deletions
@@ -19,6 +19,13 @@ API 持久化测试默认启动临时 Docker PostgreSQL 容器;运行 `dotnet dotnet test tests/FundLab.Domain.Tests/FundLab.Domain.Tests.fsproj ``` +只运行行情采集器子进程回归(成功/stderr 洪泛/超时/取消/孤儿管道/非零退出/启动失败/缺脚本/HTTP 脱敏): + +```sh +dotnet test tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj --filter "FullyQualifiedName~ProcessCollectorTests" +dotnet test tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj --filter "FullyQualifiedName~market data API keeps collector failures" +``` + 构建前端: ```sh @@ -30,8 +37,12 @@ PATH=/home/somhairle/projects/fund-lab/.tools/node-v22.23.2/bin:$PATH npm run bu API 使用环境变量 `FUND_LAB_AUTH_TOKEN` 配置 Bearer token,使用 `FUND_LAB_DATABASE_URL` 配置 PostgreSQL 连接字符串;仓库只提供 `.env.example`,不保存真实凭证。`/health` 可匿名访问,业务 API 需要 `Authorization: Bearer <token>`。 +行情采集器子进程(AKShare)行为:`FUND_LAB_AKSHARE_PYTHON` 指定解释器(默认 `python3`),`FUND_LAB_AKSHARE_PYTHONPATH` 注入 `PYTHONPATH`(例如指向预装 akshare 的 site 目录),`FUND_LAB_COLLECTOR_TIMEOUT_SECONDS` 设定硬超时(默认 30 秒,有效范围 1–600)。超时覆盖子进程退出与管道 EOF 两个阶段;超时、取消和异常路径都会杀掉整个进程组(通过 `setsid` 建立独立会话),不会遗留持有管道的后台子进程;对外的 HTTP 错误只返回统一短消息(如 503 `MARKET_DATA_UNAVAILABLE` + "collector did not finish within N seconds"),不泄露路径、stderr 或密钥,stderr 仅记录在服务端日志。 + 基金创建接口为 `POST /api/funds`,需要 `Idempotency-Key` 和 decimal 字符串字段;基金读取接口为 `GET /api/funds/{id}`。首次启动会自动创建所需表。 -运行时数据库、行情缓存、凭证和构建产物不应进入 Git。当前 `origin` 为 `/home/somhairle/git/fund-lab.git`;基线提交已独立发布并核对,最新 3c-1 Web 建档切片修改在完成最终验证前保持在本地工作树。 +运行时数据库、行情缓存、凭证和构建产物不应进入 Git。当前 `origin` 为 `/home/somhairle/git/fund-lab.git`;基线提交已独立发布并核对,最新 3c-2 行情采集器子进程加固切片(超时/取消/进程组清理/错误脱敏)在完成最终验证后单独提交。 + +真实 AKShare 冒烟(可选):将 `FUND_LAB_AKSHARE_PYTHON` 指向 `/usr/bin/python3` 并用 `FUND_LAB_AKSHARE_PYTHONPATH` 指向装有 akshare 的 site 目录后启动 API,`GET /api/instruments/search?q=000001` 应返回 200 与真实基金名单,`POST /api/instruments/000001/nav/refresh` 应返回 200 并持久化净值观测。 3a 收口测试还会启动独立临时 PostgreSQL 和真实 Kestrel 进程,覆盖 HTTP 创建、幂等重放/冲突、认证、畸形 JSON、金额边界、停止重启后的读取,以及首次 SQL 写入后故障的完整事务回滚。 diff --git a/qa/README.md b/qa/README.md index 96aea72..3c4adf0 100644 --- a/qa/README.md +++ b/qa/README.md @@ -42,7 +42,7 @@ bash qa/lifecycle-test.sh - D 重新读取:恰好 1 次 GET,金额持久化复核。 - E token 失效:换 token 立即清空档案;拦截并延迟的旧 GET 响应落地后被忽略(stale 防护)。 - F 幂等重试:首 POST 被掐断显示错误横幅,重试复用同一 `Idempotency-Key` 成功。 -- G 全程无浏览器 console/page 错误。 +- G 全程无浏览器 console/page 错误。判定口径(见 `qa/driver/browser-test.js` 的 `pageerror`/`response`/`console` 监听):任何非预期 ≥400 的资源响应、任何 pageerror、以及未被显式过滤的 console error 都记为失败;仅预期内的行情探测 `401 (Unauthorized)`、favicon 噪音和非 API 的 "Failed to load resource" 被过滤。 ## 真实行情模式 @@ -52,7 +52,7 @@ QA 默认走合成桩。需要真实 AKShare 时,不运行 `run.sh`,而是� PYTHONPATH=/tmp/opencode/fund-lab-akshare-site /usr/bin/python3 -c "import akshare" # 1.18.96 ``` -真实模式输出的 envelope 不含 `_synthetic`,来源为真实数据源;前端“数据来源”行相应显示“真实建档”。 +真实模式输出的 envelope 不含 `_synthetic`,来源为真实数据源。注意两处口径不同:基金档案行的"数据来源:真实建档"描述的是基金建档本身(对应 `isSynthetic=false`,无论行情模式),而行情面板固定显示"AKShare · exact match"并始终只承载行情数据的来源;不要把档案标签当作行情模式的指示器。 ## 不变量 diff --git a/src/FundLab.Api/App.fs b/src/FundLab.Api/App.fs index a4e0a00..f61e760 100644 --- a/src/FundLab.Api/App.fs +++ b/src/FundLab.Api/App.fs @@ -265,13 +265,13 @@ module App = let private searchInstruments (marketData: IMarketDataService) : HttpHandler = fun next ctx -> - match marketData.Search(ctx.Request.Query["q"].ToString()) with + match marketData.Search(ctx.Request.Query["q"].ToString(), ctx.RequestAborted) 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 + match marketData.RefreshNav(code, ctx.RequestAborted) with | Ok observations -> json (marketDataNavResponse code observations) next ctx | Error failure -> marketDataError failure next ctx 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<string, string> - abstract FetchNav: code: string -> Result<string, string> + abstract Search: query: string * CancellationToken -> Result<string, string> + abstract FetchNav: code: string * CancellationToken -> Result<string, string> type IMarketDataService = - abstract Search: query: string -> Result<MarketDataSearchPayload, MarketDataFailure> - abstract RefreshNav: code: string -> Result<MarketDataNavRecord list, MarketDataFailure> + abstract Search: query: string * CancellationToken -> Result<MarketDataSearchPayload, MarketDataFailure> + abstract RefreshNav: code: string * CancellationToken -> 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 = +module private CollectorNative = + let signalKill = 9 + + [<DllImport("libc", SetLastError = true)>] + 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<string>) = + 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) diff --git a/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj b/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj index c976006..63f0750 100644 --- a/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj +++ b/tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj @@ -21,6 +21,7 @@ </ItemGroup> <ItemGroup> <Compile Include="ApiTests.fs" /> + <Compile Include="ProcessCollectorTests.fs" /> <Compile Include="PersistenceTests.fs" /> <Compile Include="Program.fs" /> </ItemGroup> diff --git a/tests/FundLab.Api.Tests/PersistenceTests.fs b/tests/FundLab.Api.Tests/PersistenceTests.fs index 5157a9c..64e857c 100644 --- a/tests/FundLab.Api.Tests/PersistenceTests.fs +++ b/tests/FundLab.Api.Tests/PersistenceTests.fs @@ -695,8 +695,8 @@ type PersistenceTests(fixture: PostgresFixture) = let marketData = { new IMarketDataService with - member _.Search _ = Ok searchPayload - member _.RefreshNav _ = Ok [ observation ] + member _.Search(_, _) = Ok searchPayload + member _.RefreshNav(_, _) = Ok [ observation ] member _.GetNav(_, _, _) = Ok [ observation ] } let app = App.createApplicationWithMarketData (repository ()) marketData @@ -726,6 +726,53 @@ type PersistenceTests(fixture: PostgresFixture) = Assert.Contains("\"source\":\"akshare\"", navBody) [<Fact>] + member _.``market data API keeps collector failures and process tree out of HTTP responses``() = + let dir = Path.Combine(Path.GetTempPath(), "fund-lab-collector-" + Guid.NewGuid().ToString("N")) + + Directory.CreateDirectory(dir) |> ignore + + try + let scriptPath = Path.Combine(dir, "hang.sh") + + File.WriteAllText( + scriptPath, + "printf 'secret=%s\\n' 'fund-lab-HTTP-SENTINEL' >&2\nsleep 60 & child=$!\necho $child > pidfile\nwait $child\n" + ) + + let collector = + ProcessMarketDataCollector("/bin/sh", scriptPath, None, 1) :> IMarketDataCollector + + let marketData = MarketDataService(repository (), collector) + let app = App.createApplicationWithMarketData (repository ()) marketData + + let status, responseBody = + PersistenceTestHelpers.invoke + app + "GET" + "/api/instruments/search?q=000001" + [ "Authorization", "Bearer test-token" ] + "" + + Assert.Equal(503, status) + Assert.Contains("MARKET_DATA_UNAVAILABLE", responseBody) + Assert.Contains("did not finish within 1 seconds", responseBody) + Assert.DoesNotContain(dir, responseBody) + Assert.DoesNotContain("fund-lab-HTTP-SENTINEL", responseBody) + + let recordedPid = Int32.Parse(File.ReadAllText(Path.Combine(dir, "pidfile")).Trim()) + let mutable dead = not (Directory.Exists(sprintf "/proc/%d" recordedPid)) + let mutable waited = 0 + + while not dead && waited < 10000 do + Thread.Sleep(100) + waited <- waited + 100 + dead <- not (Directory.Exists(sprintf "/proc/%d" recordedPid)) + + Assert.True(dead, "backgrounded subprocess survived the HTTP timeout") + finally + Directory.Delete(dir, true) + + [<Fact>] member _.``fund API rejects non-object JSON bodies``() = let app = App.createApplication (repository ()) diff --git a/tests/FundLab.Api.Tests/ProcessCollectorTests.fs b/tests/FundLab.Api.Tests/ProcessCollectorTests.fs new file mode 100644 index 0000000..32009bf --- /dev/null +++ b/tests/FundLab.Api.Tests/ProcessCollectorTests.fs @@ -0,0 +1,232 @@ +namespace FundLab.Api.Tests + +module ProcessCollectorTests = + + open System + open System.Diagnostics + open System.IO + open System.Threading + open Xunit + open FundLab.Api + + let private secretSentinel = "fund-lab-SECRET-SENTINEL" + + let private fixtureDir () = + let dir = Path.Combine(Path.GetTempPath(), "fund-lab-collector-" + Guid.NewGuid().ToString("N")) + + Directory.CreateDirectory(dir) |> ignore + + dir + + let private writeScript (dir: string) (name: string) (body: string) = + let path = Path.Combine(dir, name) + File.WriteAllText(path, body) + path + + let private searchEnvelope = + "{\"schema_version\":\"fund-lab.akshare.v1\",\"operation\":\"search\",\"source\":\"fixture\",\"source_revision\":\"fixture/1\",\"collected_at\":\"2026-09-21T00:00:00Z\",\"_synthetic\":true,\"instruments\":[{\"code\":\"000001\",\"name\":\"fixture fund\",\"fund_type\":null}]}" + + let private hangScript dir = + writeScript dir "hang.sh" ("sleep 60 & child=$!\necho $child > pidfile\nwait $child\n") + + let private orphanScript dir = + writeScript + dir + "orphan.sh" + ("sleep 30 & child=$!\necho $child > pidfile\nprintf '%s' '" + searchEnvelope + "'\nexit 0\n") + + let private recordedPid (dir: string) = + Int32.Parse(File.ReadAllText(Path.Combine(dir, "pidfile")).Trim()) + + let private pidAlive (pid: int) = Directory.Exists(sprintf "/proc/%d" pid) + + let private waitForPidDeath (pid: int) = + let mutable dead = not (pidAlive pid) + let mutable waited = 0 + + while not dead && waited < 10000 do + Thread.Sleep(100) + waited <- waited + 100 + dead <- not (pidAlive pid) + + dead + + [<Fact>] + let ``collector captures stdout from a successful subprocess`` () = + let dir = fixtureDir () + + try + let script = writeScript dir "success.sh" ("printf '%s' '" + searchEnvelope + "'") + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 30) :> IMarketDataCollector + + match collector.Search("000001", CancellationToken.None) with + | Ok raw -> + Assert.Contains("000001", raw) + Assert.Contains("\"_synthetic\":true", raw) + | Error message -> Assert.Fail(sprintf "expected success, got %s" message) + finally + Directory.Delete(dir, true) + + [<Fact>] + let ``collector drains concurrent stderr flood without deadlocking`` () = + let dir = fixtureDir () + + try + let body = "yes 'flood padding padding padding padding' | head -n 5000 >&2\nprintf '%s' '" + searchEnvelope + "'" + + let script = writeScript dir "flood.sh" body + + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 30) :> IMarketDataCollector + + match collector.Search("000001", CancellationToken.None) with + | Ok raw -> Assert.Contains("000001", raw) + | Error message -> Assert.Fail(sprintf "expected success, got %s" message) + finally + Directory.Delete(dir, true) + + [<Fact>] + let ``collector enforces deadline, kills owned processes and sanitizes error`` () = + let dir = fixtureDir () + + try + let script = hangScript dir + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 1) :> IMarketDataCollector + let stopwatch = Stopwatch.StartNew() + + match collector.Search("000001", CancellationToken.None) with + | Error message -> + stopwatch.Stop() + Assert.Contains("did not finish within 1 seconds", message) + Assert.DoesNotContain(dir, message) + Assert.DoesNotContain(secretSentinel, message) + | Ok raw -> Assert.Fail(sprintf "expected timeout error, got %s" raw) + + Assert.True(stopwatch.Elapsed < TimeSpan.FromSeconds(20.0)) + Assert.True(waitForPidDeath (recordedPid dir), "backgrounded subprocess survived the deadline kill") + finally + Directory.Delete(dir, true) + + [<Fact>] + let ``collector reports cancellation and kills owned processes`` () = + let dir = fixtureDir () + + try + let script = hangScript dir + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 30) :> IMarketDataCollector + use cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(500.0)) + + match collector.Search("000001", cts.Token) with + | Error message -> + Assert.Equal("collector run was cancelled", message) + Assert.DoesNotContain(dir, message) + Assert.DoesNotContain(secretSentinel, message) + | Ok raw -> Assert.Fail(sprintf "expected cancellation, got %s" raw) + + Assert.True(waitForPidDeath (recordedPid dir), "backgrounded subprocess survived cancellation") + finally + Directory.Delete(dir, true) + + [<Fact>] + let ``collector recovers payload and reaps pipe-holding orphan after parent exit`` () = + let dir = fixtureDir () + + try + let script = orphanScript dir + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 5) :> IMarketDataCollector + let stopwatch = Stopwatch.StartNew() + + let result = collector.Search("000001", CancellationToken.None) + stopwatch.Stop() + + match result with + | Ok raw -> + Assert.Equal(searchEnvelope, raw) + Assert.True(stopwatch.Elapsed < TimeSpan.FromSeconds(20.0)) + | Error message -> + Assert.Fail(sprintf "expected recovered payload, got error %s" message) + + Assert.True(waitForPidDeath (recordedPid dir), "pipe-holding orphan survived parent exit") + finally + Directory.Delete(dir, true) + + [<Fact>] + let ``collector surfaces malformed payload for parser to reject`` () = + let dir = fixtureDir () + + try + let script = writeScript dir "malformed.sh" "printf '%s' '{ not json at all'" + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 30) :> IMarketDataCollector + + match collector.Search("000001", CancellationToken.None) with + | Ok raw -> + Assert.Equal("{ not json at all", raw) + + match MarketData.parseSearchPayload raw with + | Error _ -> () + | Ok _ -> Assert.Fail("expected malformed payload to fail parsing") + | Error message -> Assert.Fail(sprintf "expected raw output, got error %s" message) + finally + Directory.Delete(dir, true) + + [<Fact>] + let ``nonzero exit keeps stderr sentinel and paths out of the public error`` () = + let dir = fixtureDir () + + try + let body = + "printf 'secret=%s\\n' '" + + secretSentinel + + "' >&2\nprintf 'script lives under " + + dir + + "\\n' >&2\nexit 3\n" + + let script = writeScript dir "failing.sh" body + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 30) :> IMarketDataCollector + + match collector.Search("000001", CancellationToken.None) with + | Error message -> + Assert.Equal("collector exited with code 3", message) + Assert.DoesNotContain(secretSentinel, message) + Assert.DoesNotContain(dir, message) + | Ok raw -> Assert.Fail(sprintf "expected failure, got %s" raw) + finally + Directory.Delete(dir, true) + + [<Fact>] + let ``start failure keeps executable path out of the public error`` () = + let dir = fixtureDir () + + try + let script = writeScript dir "success.sh" ("printf '%s' '" + searchEnvelope + "'") + let missingExecutable = "/nonexistent-" + secretSentinel + "/python3" + + let collector = ProcessMarketDataCollector(missingExecutable, script, None, 30) :> IMarketDataCollector + + match collector.Search("000001", CancellationToken.None) with + | Error message -> + let sanitized = + message = "could not run the AKShare collector" + || message.StartsWith("collector exited with code") + + Assert.True(sanitized, sprintf "unexpected error message: %s" message) + Assert.DoesNotContain(secretSentinel, message) + Assert.DoesNotContain(dir, message) + | Ok raw -> Assert.Fail(sprintf "expected failure, got %s" raw) + finally + Directory.Delete(dir, true) + + [<Fact>] + let ``missing collector script yields sanitized error without paths`` () = + let dir = fixtureDir () + + try + let script = Path.Combine(dir, "missing.py") + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 30) :> IMarketDataCollector + + match collector.Search("000001", CancellationToken.None) with + | Error message -> + Assert.Equal("collector script is not deployed", message) + Assert.DoesNotContain(dir, message) + | Ok raw -> Assert.Fail(sprintf "expected error, got %s" raw) + finally + Directory.Delete(dir, true) |
