diff options
| -rw-r--r-- | src/FundLab.Api/MarketDataService.fs | 21 | ||||
| -rw-r--r-- | tests/FundLab.Api.Tests/ProcessCollectorTests.fs | 67 |
2 files changed, 73 insertions, 15 deletions
diff --git a/src/FundLab.Api/MarketDataService.fs b/src/FundLab.Api/MarketDataService.fs index 0c4714d..e9993f3 100644 --- a/src/FundLab.Api/MarketDataService.fs +++ b/src/FundLab.Api/MarketDataService.fs @@ -46,7 +46,9 @@ type ProcessMarketDataCollector(pythonExecutable: string, scriptPath: string, py Console.Error.WriteLine("fund-lab collector stderr: {0}", bounded) let execute (token: CancellationToken) arguments = - if not (File.Exists scriptPath) then + 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() @@ -98,12 +100,12 @@ type ProcessMarketDataCollector(pythonExecutable: string, scriptPath: string, py use timeoutSource = new CancellationTokenSource(deadline) use linkedSource = CancellationTokenSource.CreateLinkedTokenSource(token, timeoutSource.Token) - let mutable timedOut = false + let mutable expired = false try child.WaitForExitAsync(linkedSource.Token).GetAwaiter().GetResult() with _ -> - if not token.IsCancellationRequested then timedOut <- true + if not token.IsCancellationRequested then expired <- true killOwnedGroup () @@ -112,15 +114,18 @@ type ProcessMarketDataCollector(pythonExecutable: string, scriptPath: string, py with _ -> () let mutable eofRemaining = deadline - clock.Elapsed - - if eofRemaining <= TimeSpan.Zero then eofRemaining <- TimeSpan.FromSeconds(2.0) - + let mutable eofExpired = eofRemaining <= TimeSpan.Zero let mutable readsDone = false try - readsDone <- Task.WhenAll(outputTask, errorTask).Wait(int eofRemaining.TotalMilliseconds) + 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 () @@ -138,7 +143,7 @@ type ProcessMarketDataCollector(pythonExecutable: string, scriptPath: string, py logCollectorStderr error - if timedOut then + if expired || eofExpired then Error(sprintf "collector did not finish within %d seconds" timeoutSeconds) elif token.IsCancellationRequested then Error "collector run was cancelled" diff --git a/tests/FundLab.Api.Tests/ProcessCollectorTests.fs b/tests/FundLab.Api.Tests/ProcessCollectorTests.fs index 32009bf..ee3f016 100644 --- a/tests/FundLab.Api.Tests/ProcessCollectorTests.fs +++ b/tests/FundLab.Api.Tests/ProcessCollectorTests.fs @@ -6,6 +6,7 @@ module ProcessCollectorTests = open System.Diagnostics open System.IO open System.Threading + open System.Threading.Tasks open Xunit open FundLab.Api @@ -127,25 +128,77 @@ module ProcessCollectorTests = Directory.Delete(dir, true) [<Fact>] - let ``collector recovers payload and reaps pipe-holding orphan after parent exit`` () = + let ``collector rejects orphan-held pipes at the total exit and EOF deadline`` () = let dir = fixtureDir () try let script = orphanScript dir - let collector = ProcessMarketDataCollector("/bin/sh", script, None, 5) :> IMarketDataCollector + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 2) :> 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)) + | Ok raw -> Assert.Fail(sprintf "recovered payload after expired deadline: %s" raw) | Error message -> - Assert.Fail(sprintf "expected recovered payload, got error %s" message) + Assert.Equal("collector did not finish within 2 seconds", message) + Assert.DoesNotContain(dir, message) + Assert.DoesNotContain(secretSentinel, message) + + Assert.True(stopwatch.Elapsed >= TimeSpan.FromSeconds(1.8)) + Assert.True(stopwatch.Elapsed < TimeSpan.FromSeconds(15.0)) + Assert.True(waitForPidDeath (recordedPid dir), "pipe-holding orphan survived the EOF deadline kill") + finally + Directory.Delete(dir, true) + + [<Fact>] + let ``cancellation during EOF wait interrupts promptly and kills descendants`` () = + let dir = fixtureDir () + + try + let script = orphanScript dir + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 30) :> IMarketDataCollector + use cts = new CancellationTokenSource() + + let searchTask = Task.Run(fun () -> collector.Search("000001", cts.Token)) + Thread.Sleep(700) + + let stopwatch = Stopwatch.StartNew() + cts.Cancel() + let completed = searchTask.Wait(10000) + stopwatch.Stop() + + Assert.True(completed, "search did not return promptly after cancellation") + + match searchTask.Result with + | Error message -> Assert.Equal("collector run was cancelled", message) + | Ok raw -> Assert.Fail(sprintf "expected cancellation, got %s" raw) + + Assert.True(stopwatch.Elapsed < TimeSpan.FromSeconds(5.0), "cancellation during EOF wait was not prompt") + Assert.True(waitForPidDeath (recordedPid dir), "pipe-holding orphan survived EOF-wait cancellation") + finally + Directory.Delete(dir, true) + + [<Fact>] + let ``pre-cancelled request does not launch the collector`` () = + let dir = fixtureDir () + + try + let script = + writeScript dir "recording.sh" ("echo $$ > pidfile\nprintf '%s' '" + searchEnvelope + "'\n") + + let collector = ProcessMarketDataCollector("/bin/sh", script, None, 30) :> IMarketDataCollector + use cts = new CancellationTokenSource() + cts.Cancel() + + match collector.Search("000001", cts.Token) with + | Error message -> + Assert.Equal("collector run was cancelled", message) + Assert.DoesNotContain(dir, message) + | Ok raw -> Assert.Fail(sprintf "expected pre-cancel, got %s" raw) - Assert.True(waitForPidDeath (recordedPid dir), "pipe-holding orphan survived parent exit") + Assert.False(File.Exists(Path.Combine(dir, "pidfile")), "collector process was launched despite pre-cancelled request") finally Directory.Delete(dir, true) |
