Source file proverThread.ml
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
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
[%%import "vsrocq_config.mlh"]
open Types
let preempt = ref false
let set_options ~preempt:x = preempt := x
let (Log log) = Log.mk_log "proverThread"
type retry = bool
module Queue = struct
type 'a t = 'a list ref
let create () = ref []
let is_empty q = [] = !q
let pop q =
match !q with
| [] -> assert false
| x :: xs -> q := xs; x
let enqueue x q = q := !q @ [x]
let add_head x q = q := x :: !q
let iter f q = List.iter f !q
end
type rocq_job =
| Job :
document_id
* string
* (Memprof_limits.Token.t -> 'a interruptible_result)
* 'a interruptible_result Sel.Promise.handler
* Memprof_limits.Token.t
* retry ref
-> rocq_job
type jobs = { mutable running : rocq_job option; queue : rocq_job Queue.t }
let jobs : jobs = { running = None; queue = Queue.create () }
let jobs_mutex = Mutex.create ()
let jobs_condition = Condition.create ()
let _runner =
Thread.create
(fun () ->
while true do
let Job (doc_id , name, task, resolver, token, retry) =
Mutex.lock jobs_mutex;
while Queue.is_empty jobs.queue do
Condition.wait jobs_condition jobs_mutex
done;
let rc = Queue.pop jobs.queue in
jobs.running <- Some rc;
Condition.signal jobs_condition;
Mutex.unlock jobs_mutex;
rc
in
log (fun () -> Printf.sprintf "runner: job begins: %s" name);
match task token with
| Interrupted when !retry ->
Mutex.lock jobs_mutex;
log (fun () -> Printf.sprintf "runner: postponing running job: %s" name);
jobs.running <- None;
Queue.enqueue (Job (doc_id , name, task, resolver, Memprof_limits.Token.create (), ref false)) jobs.queue;
Condition.signal jobs_condition;
Mutex.unlock jobs_mutex
| x ->
Mutex.lock jobs_mutex;
log (fun () -> Printf.sprintf "runner: job ends: %s" name);
Sel.Promise.fulfill resolver x;
jobs.running <- None;
Condition.signal jobs_condition;
Mutex.unlock jobs_mutex
done)
()
let interrupt_job_if ~doc_id (Job(id,_,_,_,token,_)) = if id = doc_id then Memprof_limits.Token.set token
let postpone_job (Job(_,name,_,_,token,retry)) =
if not !preempt then ()
else begin
log (fun () -> Printf.sprintf "main: postponing running job: %s" name);
retry := true;
Memprof_limits.Token.set token
end
let interrupt ~doc_id =
Mutex.lock jobs_mutex;
Option.iter (interrupt_job_if ~doc_id) jobs.running;
Queue.iter (interrupt_job_if ~doc_id) jobs.queue;
Mutex.unlock jobs_mutex
let limit f token =
match Terminated (Memprof_limits.limit_with_token ~token f) with
| Aborted _ | Interrupted -> assert false
| Terminated (Error _) -> Interrupted
| Terminated (Ok x) -> Terminated x
| exception e ->
let e, info = Exninfo.capture e in
Aborted (CErrors.iprint (e, info))
let busy_wait timeout p token =
let timeout = Unix.gettimeofday () +. timeout in
while not (Sel.Promise.is_resolved p) && Unix.gettimeofday () < timeout do
Unix.sleepf 0.01;
done;
Memprof_limits.Token.set token;
try
match Sel.Promise.get p with
| Sel.Promise.Fulfilled x -> x
| Sel.Promise.Rejected e -> raise e
with Failure _ -> Aborted (Pp.str "Rocq times out")
let try_run ~doc_id ~name ~timeout f =
log (fun () -> "main: run");
let token = Memprof_limits.Token.create () in
let promise, r = Sel.Promise.make () in
Mutex.lock jobs_mutex;
Option.iter postpone_job jobs.running;
Queue.add_head (Job (doc_id, name, limit f, r, token, ref false)) jobs.queue;
Condition.signal jobs_condition;
Mutex.unlock jobs_mutex;
busy_wait timeout promise token
let eventually_run ~doc_id ~name f =
log (fun () -> "main: eventually_run");
let token = Memprof_limits.Token.create () in
let promise, r = Sel.Promise.make () in
Mutex.lock jobs_mutex;
Queue.enqueue (Job (doc_id, name, limit f, r, token, ref false)) jobs.queue;
Condition.signal jobs_condition;
Mutex.unlock jobs_mutex;
promise
let run ~doc_id ~name f =
log (fun () -> "main: run");
let token = Memprof_limits.Token.create () in
let promise, r = Sel.Promise.make () in
Mutex.lock jobs_mutex;
Option.iter postpone_job jobs.running;
Queue.add_head (Job (doc_id, name, limit f, r, token, ref false)) jobs.queue;
Condition.signal jobs_condition;
Mutex.unlock jobs_mutex;
match busy_wait 99999. promise token with
| Interrupted -> assert false
| Aborted pp -> Result.Error pp
| Terminated x -> Result.Ok x