1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
open Effect
open Effect.Deep
let opt_sequential = ref false
type channels = in_channel * out_channel * in_channel
type open_result = Unix.process_status * string * string
type _ Effect.t += Open_process_full : (string * string Array.t * string option) -> open_result t
module ParUnix = struct
let open_process_full cmd env to_stdin = perform (Open_process_full (cmd, env, to_stdin))
end
let recommended_parallelism () = Domain.recommended_domain_count ()
let read_all ch =
let buf = Buffer.create 4096 in
let chunk = Bytes.create 4096 in
( try
while true do
let n = input ch chunk 0 4096 in
if n = 0 then raise End_of_file;
Buffer.add_subbytes buf chunk 0 n
done
with End_of_file -> ()
);
Buffer.contents buf
let parmap (type a b) ~parallelism f (xs : a list) : b list =
let xs = Array.of_list xs in
let len = Array.length xs in
let tmp : b option Array.t = Array.make len None in
let waiting : (int * int * channels * (open_result, unit) continuation) Queue.t = Queue.create () in
let next = ref 0 in
let run i action =
match action () with
| effect Open_process_full (cmd, env, to_stdin), cont ->
let ((_, stdin, _) as channels) = Unix.open_process_full cmd env in
(match to_stdin with None -> () | Some str -> output_string stdin str);
let pid = Unix.process_full_pid channels in
close_out stdin;
Queue.add (i, pid, channels, cont) waiting
| y -> tmp.(i) <- Some y
in
let rec go () =
if not (!next = len && Queue.is_empty waiting) then (
let ready =
let pending = Queue.length waiting in
let rec scan k =
if k = 0 then None
else (
let i, pid, channels, cont = Queue.pop waiting in
let res, status = Unix.waitpid [Unix.WNOHANG] pid in
if res = 0 then (
Queue.add (i, pid, channels, cont) waiting;
scan (k - 1)
)
else (
let out_chan, _, err_chan = channels in
let out_str = read_all out_chan in
let err_str = read_all err_chan in
close_in out_chan;
close_in err_chan;
Some ((status, out_str, err_str), cont)
)
)
in
scan pending
in
( match ready with
| Some (result, cont) -> continue cont result
| None ->
if Queue.length waiting >= parallelism then ()
else (
let i = !next in
if i = len then ()
else (
incr next;
run i (fun () -> f xs.(i))
)
)
);
go ()
)
in
go ();
Option.get (Util.option_all (Array.to_list tmp))
let map ~parallelism f xs = if !opt_sequential then List.map f xs else parmap ~parallelism f xs
let run_process cmd env to_stdin =
let out_chan, in_chan, err_chan = Unix.open_process_full cmd env in
(match to_stdin with None -> () | Some str -> output_string in_chan str);
close_out in_chan;
let stdout_str = read_all out_chan in
let stderr_str = read_all err_chan in
let status = Unix.close_process_full (out_chan, in_chan, err_chan) in
(status, stdout_str, stderr_str)
let toplevel_handler f =
let rec run action =
match action () with
| effect Open_process_full (cmd, env, to_stdin), cont -> run (fun () -> continue cont (run_process cmd env to_stdin))
| x -> x
in
run f