summaryrefslogtreecommitdiff
path: root/src/FundLab.Api/MarketDataService.fs
diff options
context:
space:
mode:
Diffstat (limited to 'src/FundLab.Api/MarketDataService.fs')
-rw-r--r--src/FundLab.Api/MarketDataService.fs21
1 files changed, 13 insertions, 8 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"