summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--src/FundLab.Api/MarketDataService.fs21
-rw-r--r--tests/FundLab.Api.Tests/ProcessCollectorTests.fs67
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)