summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--README.md13
-rw-r--r--qa/README.md4
-rw-r--r--src/FundLab.Api/App.fs4
-rw-r--r--src/FundLab.Api/MarketDataService.fs174
-rw-r--r--tests/FundLab.Api.Tests/FundLab.Api.Tests.fsproj1
-rw-r--r--tests/FundLab.Api.Tests/PersistenceTests.fs51
-rw-r--r--tests/FundLab.Api.Tests/ProcessCollectorTests.fs232
7 files changed, 441 insertions, 38 deletions
diff --git a/README.md b/README.md
index 1fcd88b..73350e4 100644
--- a/README.md
+++ b/README.md
@@ -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)