Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 37 additions & 6 deletions parallel.ml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 =
Expand All @@ -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
Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down
2 changes: 2 additions & 0 deletions parallel.mli
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down