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.fs174
1 files changed, 143 insertions, 31 deletions
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)