diff --git a/parallel.ml b/parallel.ml index c0043b8..f737a0f 100644 --- a/parallel.ml +++ b/parallel.ml @@ -4,6 +4,20 @@ open Prelude let log = Log.from "parallel" +(* counters are process-wide, as in Httpev : several run_forks calls are not told apart. + Add Var attributes here if that is ever needed. *) +let stats = new Var.typ "parallel.forks" "event" + +let nr_spawned = stats#count "spawned" +let nr_exited = stats#count "exited" +let nr_vanished = stats#count "vanished" +let nr_killed = stats#count "killed" + +(* a child inherits a snapshot of the counters, but it did not do those forks *) +let () = Nix.register_on_fork begin fun () -> + nr_spawned := 0; nr_exited := 0; nr_vanished := 0; nr_killed := 0 +end + type revive_mode = | Never | On_failure @@ -123,6 +137,7 @@ let worker (execute : task -> result) = Unix.set_close_on_exec main_read; let cout = Unix.out_channel_of_descr main_write in let cin = Unix.in_channel_of_descr main_read in + incr nr_spawned; { ch = Some (cin, cout); pid; } let create execute n = @@ -141,9 +156,15 @@ let stop ?wait t = let l = t.running |> List.map (fun w -> close_ch w; w.pid) in Nix.sleep 0.1; (* let idle workers detect EOF and exit peacefully (frequent io-in-signal-handler deadlock problem) *) t.running <- []; + let total = List.length l in match do_stop ?wait l with - | `Done -> log #info "Stopped %d workers properly%s" (List.length l) (gone ()) - | `Killed killed -> log #info "Timeouted, killing %d (of %d) workers with SIGKILL%s" killed (List.length l) (gone ()) + | `Done -> + nr_exited += total; + log #info "Stopped %d workers properly%s" total (gone ()) + | `Killed killed -> + nr_killed += killed; + nr_exited += (total - killed); + log #info "Timeouted, killing %d (of %d) workers with SIGKILL%s" killed total (gone ()) let perform t ?(autoexit=false) tasks finish = match t.running with @@ -176,6 +197,7 @@ let perform t ?(autoexit=false) tasks finish = | exception exn -> log #warn ~exn "no result from PID %d" w.pid; t.gone <- t.gone + 1; + incr nr_vanished; decr workers; (* close pipes and forget dead child, do not reap zombie so that premature exit is visible in process list *) close_ch w; @@ -191,7 +213,7 @@ let perform t ?(autoexit=false) tasks finish = end; Some answer with - | exn -> log #warn ~exn "perform (from PID %d)" w.pid; decr workers; None + | exn -> log #warn ~exn "perform (from PID %d)" w.pid; incr nr_vanished; decr workers; None end in List.iter finish answers; @@ -261,10 +283,11 @@ let run_forks_simple ?(revive=Never) ?wait_stop f args = log #error ~exn ~backtrace:true "worker failed"; exit 1 end - | `Forked pid -> Hashtbl.add workers pid x; pid + | `Forked pid -> incr nr_spawned; Hashtbl.add workers pid x; pid in args |> List.iter (fun x -> let (_:int) = launch f x in ()); let pids () = Hashtbl.keys workers |> List.of_enum in + let account = List.iter (function (_, Some (Unix.WEXITED 0)) -> incr nr_exited | _ -> incr nr_vanished) in let maybe_revive ~always dead = dead |> List.iter begin fun (pid, result) -> match Hashtbl.find workers pid with @@ -291,20 +314,28 @@ let run_forks_simple ?(revive=Never) ?wait_stop f args = | true -> log #info "Stopping %d workers" total; begin match do_stop ?wait:wait_stop (Hashtbl.keys workers |> List.of_enum) with - | `Done -> log #info "Stopped %d workers" total - | `Killed n -> log #info "Killed %d (of %d) workers with SIGKILL" n total + | `Done -> + nr_exited += total; + log #info "Stopped %d workers" total + | `Killed n -> + nr_killed += n; + nr_exited += (total - n); + log #info "Killed %d (of %d) workers with SIGKILL" n total end | false -> let (dead,_live) = reap (pids ()) in match dead, revive with | [], _ -> loop (max 1. (pause /. 2.)) | dead, Always -> + account dead; maybe_revive ~always:true dead; loop (min 10. (pause *. 1.5)) | dead, On_failure -> + account dead; maybe_revive ~always:false dead; loop (min 10. (pause *. 1.5)) | dead, Never -> + account dead; log #info "%d child workers exited (PIDs: %s)" (List.length dead) (Stre.list (string_of_int $ fst) dead); List.iter (Hashtbl.remove workers $ fst) dead; loop pause diff --git a/parallel.mli b/parallel.mli index eb3d677..a6672a7 100644 --- a/parallel.mli +++ b/parallel.mli @@ -14,6 +14,8 @@ val launch_forks : ('a -> unit) -> 'a list -> unit (** Launch forks for each element of the list and wait for all workers to finish. Pass exit signals to the workers, see {!Forks.stop} for the description of [wait_stop] parameter. + Workers are counted process-wide in the [parallel.forks] {!Var} family, per [event] : + [spawned], [exited], [vanished] (crashed or exited with non-zero status) and [killed] (with SIGKILL). @param revive to keep workers running (restarting with same param if exited) [default: Never] *) val run_forks : ?wait_stop:int -> ?revive:revive_mode -> ?wait:int -> ?workers:int -> ('a -> unit) -> 'a list -> unit