Skip to content

Commit b54188d

Browse files
committed
tests: pin the async handshake's storage, and keep the reproducers
`debug-scripts/async_handshake_deadlock.py` drives the bug through devito: a checkpointed wavefield streamed to disk, written by a forward operator and read back by an adjoint one, six workers at a time. Against the unfixed compiler a worker stops making progress within minutes, and CvxCompress trips its own assertion on a buffer that was never filled -- the same lost update wearing its other face. With the fix the same run reports no deadlock in 900 s. `debug-scripts/handshake_mfe.c` is that protocol extracted verbatim, 130 lines and no devito, where ThreadSanitizer reports the race on every run and the -DATOMIC build is clean. The codegen assertion in devito's own suite needs a device, so devitopro carries the equivalent one for the streaming path, which runs anywhere.
1 parent 38c949c commit b54188d

4 files changed

Lines changed: 374 additions & 0 deletions

File tree

Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
"""The async task handshake loses an update and both threads wait forever.
2+
3+
A checkpointed wavefield streamed to disk is written by a forward operator and
4+
read back by an adjoint one. Each pair hands off through the generated
5+
`lock0`/`flag` protocol roughly once per time step. Declared `volatile int`
6+
those two objects are neither atomic nor ordered, so a delivery can be lost --
7+
the lock handed back after the value that replaced it -- and then the compute
8+
thread waits for data it has already spent the request for while the task waits
9+
to be asked. Nothing moves, both threads spin at 100% CPU.
10+
11+
python debug-scripts/async_handshake_deadlock.py # 6 workers, 15 min
12+
python debug-scripts/async_handshake_deadlock.py 4 300 # 4 workers, 5 min
13+
14+
Needs contention: one worker on its own rarely hits the window. A worker that
15+
goes quiet for over 90 s while still burning CPU is the deadlock; the script
16+
says so and prints what the generated operator declared.
17+
18+
The same protocol extracted to plain C is in handshake_mfe.c, where
19+
ThreadSanitizer reports the race on every run.
20+
"""
21+
22+
import os
23+
import sys
24+
import time
25+
from multiprocessing import Process, Value
26+
27+
from devitopro import TimeFunction
28+
from devitopro.types.enriched import Disk
29+
30+
from devito import ConditionalDimension, Eq, Function, Grid, Inc, Operator
31+
32+
SHAPE = (240, 200)
33+
SO = 8
34+
NT = 1200 # ~1200 handshakes per operator call, as in a real gradient
35+
FACTOR = 5
36+
37+
38+
def build():
39+
grid = Grid(shape=SHAPE, extent=(12000., 10000.))
40+
t_sub = ConditionalDimension(name='t_sub', parent=grid.time_dim,
41+
factor=FACTOR)
42+
b = Function(name='b', grid=grid, space_order=SO)
43+
b.data[:] = 0.5
44+
u = TimeFunction(name='u', grid=grid, space_order=SO, time_order=2)
45+
v = TimeFunction(name='v', grid=grid, space_order=SO, time_order=2)
46+
g = Function(name='g', grid=grid)
47+
usave = TimeFunction(name='usave', grid=grid, space_order=SO, time_order=2,
48+
time_dim=t_sub, save=NT // FACTOR + 2, layers=Disk,
49+
compression='cvxcompress')
50+
51+
fwd = Operator([Eq(u.forward, (b * u).laplace + 2 * u - u.backward),
52+
Eq(usave, u.forward)])
53+
adj = Operator([Eq(v.backward, (b * v).laplace + 2 * v - v.forward),
54+
Inc(g, usave * v)])
55+
return fwd, adj, u, v, usave
56+
57+
58+
def declaration(op):
59+
"""What the generated operator declared the two shared objects as."""
60+
code = str(op.ccode)
61+
return [line.strip() for line in code.splitlines()
62+
if ('lock0[1]' in line or 'int flag' in line)]
63+
64+
65+
def worker(beat, cycles, index):
66+
fwd, adj, u, v, usave = build()
67+
if index == 0: # one worker reports the code
68+
print(f' generated: {"; ".join(declaration(adj))}', flush=True)
69+
for _ in range(cycles):
70+
u.data[:] = 0.
71+
fwd.apply(time_M=NT, dt=1.)
72+
v.data[:] = 0.
73+
adj.apply(time_M=NT - FACTOR - 2, dt=1.)
74+
for f in usave.values() if hasattr(usave, 'values') else [usave]:
75+
f._reset()
76+
beat.value = int(time.time()) # heartbeat
77+
78+
79+
def main():
80+
nworkers = int(sys.argv[1]) if len(sys.argv) > 1 else 6
81+
budget = int(sys.argv[2]) if len(sys.argv) > 2 else 900
82+
os.environ.setdefault('OMP_NUM_THREADS', '4')
83+
84+
beats = [Value('l', 0) for _ in range(nworkers)]
85+
procs = [Process(target=worker, args=(b, 10 ** 6, i))
86+
for i, b in enumerate(beats)]
87+
for p in procs:
88+
p.start()
89+
print(f'{nworkers} workers, {budget}s budget', flush=True)
90+
91+
start = time.time()
92+
stalled = []
93+
while time.time() - start < budget and len(stalled) == 0:
94+
time.sleep(15)
95+
now = time.time()
96+
for i, b in enumerate(beats):
97+
if b.value and now - b.value > 90 and procs[i].is_alive():
98+
stalled.append(i)
99+
for p in procs:
100+
p.terminate()
101+
102+
if stalled:
103+
print(f'DEADLOCK: worker(s) {stalled} stopped making progress while '
104+
f'still running', flush=True)
105+
return 1
106+
print(f'no deadlock in {int(time.time() - start)}s', flush=True)
107+
return 0
108+
109+
110+
if __name__ == '__main__':
111+
sys.exit(main())

‎debug-scripts/handshake_mfe.c‎

Lines changed: 149 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,149 @@
1+
/* devito's async-task handshake, extracted verbatim from a generated operator.
2+
*
3+
* As generated, the lock and the flag are `volatile int`, written by both
4+
* threads. `volatile` re-issues the loads; it orders nothing and makes
5+
* nothing atomic, so the two threads race and an update can be lost:
6+
*
7+
* compute: lock0[0] = 0; (release_lock0) -- may be observed late
8+
* flag = 2; (activate0) -- may be observed first
9+
* task: sees the request, delivers lock0[0] = 2, flag = 1
10+
* compute: the late lock0[0] = 0 lands, wiping the delivery
11+
* the next release_lock0 waits for a 2 nobody will write again
12+
*
13+
* Both threads then spin forever: the compute waiting for data it has already
14+
* spent the request for, the task waiting to be asked.
15+
*
16+
* cc -O1 -g -fsanitize=thread -o mfe handshake_mfe.c -lpthread
17+
* cc -O1 -g -fsanitize=thread -o mfe_fixed handshake_mfe.c -lpthread -DATOMIC
18+
* ./mfe 300000 # ThreadSanitizer: data race, every run
19+
* ./mfe_fixed 300000 # clean
20+
*
21+
* Without the sanitizer the deadlock itself shows up too, but rarely -- the
22+
* window is nanoseconds wide. Built -O3 and run eight at a time it has taken
23+
* ~3e8 handshakes per event here:
24+
*
25+
* DEADLOCK wait_lock at step 42451768: lock=0 flag=1
26+
*
27+
* which is the state the real workload wedges in.
28+
*/
29+
#include <pthread.h>
30+
#include <stdio.h>
31+
#include <stdlib.h>
32+
#include <string.h>
33+
#include <unistd.h>
34+
35+
#ifdef ATOMIC
36+
#define SHARED _Atomic int /* the fix: atomic storage */
37+
#else
38+
#define SHARED volatile int /* as devito generates it */
39+
#endif
40+
41+
#define SPIN_CAP 2000000000L /* ~2 s of spinning before we look closer */
42+
43+
struct tsdata {
44+
SHARED time;
45+
SHARED flag;
46+
SHARED *lock0;
47+
};
48+
49+
/* lock0 is a stack array in the operator, flag/time live in the shared struct:
50+
* different cache lines, as generated */
51+
struct padded_lock { SHARED lock0[1]; char pad[128]; };
52+
53+
static SHARED *gflag = 0;
54+
55+
/* A spin past the cap is either a descheduled peer or a lost store. Sleeping
56+
* tells them apart: a peer gets scheduled, a wiped delivery never arrives. */
57+
static int confirm_stall(SHARED *lock0, int want)
58+
{
59+
for (int i = 0; i < 3; i++) {
60+
sleep(2);
61+
if (lock0[0] == want) return 0;
62+
}
63+
fprintf(stderr, " state after 6 s: lock=%d flag=%d\n", lock0[0], *gflag);
64+
return 1;
65+
}
66+
67+
static void *task(void *_sdata)
68+
{
69+
struct tsdata *sdata = (struct tsdata *)_sdata;
70+
SHARED *lock0 = sdata->lock0;
71+
72+
while (sdata->flag != 0) {
73+
if (sdata->flag == 2) {
74+
volatile int t = sdata->time;
75+
(void)t; /* stands in for the fetch */
76+
lock0[0] = 2;
77+
sdata->flag = 1;
78+
}
79+
}
80+
return NULL;
81+
}
82+
83+
/* while(lock0[0] == 0); */
84+
static int wait_lock(SHARED *lock0)
85+
{
86+
long spins = 0;
87+
while (lock0[0] == 0)
88+
if (++spins > SPIN_CAP) { if (confirm_stall(lock0, 2)) return -1; spins = 0; }
89+
return 0;
90+
}
91+
92+
/* while(lock0[0] != 2); lock0[0] = 0; */
93+
static int release_lock(SHARED *lock0)
94+
{
95+
long spins = 0;
96+
while (lock0[0] != 2)
97+
if (++spins > SPIN_CAP) { if (confirm_stall(lock0, 2)) return -1; spins = 0; }
98+
lock0[0] = 0;
99+
return 0;
100+
}
101+
102+
/* while(flag != 1); time = t; flag = 2; */
103+
static void activate(int t, struct tsdata *sdata)
104+
{
105+
while (sdata->flag != 1) ;
106+
sdata->time = t;
107+
sdata->flag = 2;
108+
}
109+
110+
int main(int argc, char **argv)
111+
{
112+
const long nsteps = (argc > 1) ? atol(argv[1]) : 20000000L;
113+
114+
struct padded_lock pl;
115+
pl.lock0[0] = 2;
116+
struct tsdata sdata;
117+
sdata.flag = 1;
118+
sdata.time = 0;
119+
sdata.lock0 = pl.lock0;
120+
gflag = &sdata.flag;
121+
122+
pthread_t th;
123+
if (pthread_create(&th, NULL, task, &sdata) != 0) return 2;
124+
125+
for (long step = nsteps; step >= 1; step--) {
126+
if ((step % 4) == 0 && wait_lock(pl.lock0) < 0) {
127+
printf("DEADLOCK wait_lock at step %ld: lock=%d flag=%d\n",
128+
step, pl.lock0[0], sdata.flag);
129+
return 1;
130+
}
131+
if (release_lock(pl.lock0) < 0) {
132+
printf("DEADLOCK release_lock at step %ld: lock=%d flag=%d\n",
133+
step, pl.lock0[0], sdata.flag);
134+
return 1;
135+
}
136+
activate((int)step, &sdata);
137+
}
138+
while (sdata.flag == 2) ;
139+
sdata.flag = 0;
140+
pthread_join(th, NULL);
141+
printf("%s: completed %ld steps\n",
142+
#ifdef ATOMIC
143+
"atomic (fixed)",
144+
#else
145+
"volatile (as generated)",
146+
#endif
147+
nsteps);
148+
return 0;
149+
}

‎debug-scripts/handshake_mfe.md‎

Lines changed: 109 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,109 @@
1+
# Async task handshake: lost update (devito)
2+
3+
Fixed in devito `fix-async-memory-ordering` (38c949c20).
4+
5+
## Symptom
6+
7+
A run streaming compressed checkpoints to disk stops making progress. Both
8+
threads sit at 100% CPU indefinitely: the main thread inside the generated
9+
operator, the task thread in its own idle loop. Nothing is logged, nothing is
10+
written, and the process never recovers.
11+
12+
## The protocol, as generated
13+
14+
compute (adjoint time loop) async task
15+
--------------------------- ----------
16+
read_snapshot0: while (flag != 0):
17+
while (lock0[0] == 0); if (flag == 2):
18+
release_lock0: <fetch>
19+
while (lock0[0] != 2); lock0[0] = 2;
20+
lock0[0] = 0; flag = 1;
21+
activate0:
22+
while (flag != 1);
23+
time = …; flag = 2;
24+
25+
`lock0` and `flag` were `volatile int`, written by both threads. `volatile`
26+
re-issues loads; it orders nothing and makes nothing atomic.
27+
28+
## The lost update
29+
30+
compute: lock0[0] = 0; (release_lock0) -- observed late
31+
flag = 2; (activate0) -- observed first
32+
task: sees the request, fetches, lock0[0] = 2, flag = 1
33+
compute: the late lock0[0] = 0 lands, wiping the delivery
34+
the next release_lock0 waits for a 2 nobody will write again
35+
36+
Deadlock: the compute waits for data whose request it has already spent, the
37+
task waits to be asked. `lock == 0 && flag == 1` is reachable only this way --
38+
the task sets the lock before the flag, so a completed cycle must leave the
39+
lock at 2.
40+
41+
## Reproducers
42+
43+
`async_handshake_deadlock.py` drives it through devito: a checkpointed
44+
wavefield streamed to disk, written by a forward operator and read back by an
45+
adjoint one, six workers at a time. Against the unfixed compiler a worker
46+
stops making progress within minutes,
47+
48+
generated: volatile int flag;; volatile int lock0[1] = {2};
49+
Assertion failed: (nx == nx_check) … nx=216, ny=256, nz=1, nx_check=0
50+
DEADLOCK: worker(s) [2] stopped making progress while still running
51+
52+
-- and the assertion is the same lost update wearing its other face, the
53+
compute decompressing a buffer the task never filled. With the fix, the same
54+
run reports
55+
56+
generated: _Atomic int flag;; _Atomic int lock0[1] = {2};
57+
no deadlock in 900s
58+
59+
`handshake_mfe.c` is the same protocol extracted verbatim, 130 lines and no
60+
devito, for looking at the mechanism on its own:
61+
62+
cc -O1 -g -fsanitize=thread -o mfe handshake_mfe.c -lpthread
63+
cc -O1 -g -fsanitize=thread -o mfe_fixed handshake_mfe.c -lpthread -DATOMIC
64+
./mfe 300000 # ThreadSanitizer: data race, every run
65+
./mfe_fixed 300000 # clean
66+
67+
Built `-O3` and run eight at a time it also deadlocks outright, at roughly one
68+
event per 3e8 handshakes:
69+
70+
DEADLOCK wait_lock at step 42451768: lock=0 flag=1
71+
72+
## Evidence from the real workload
73+
74+
Spin counters in the three waits, via `DEVITO_JIT_BACKDOOR=1`:
75+
76+
STUCK release_lock0: lock=0 flag=1 time=1107 spins=10x1e8 … 120x1e8
77+
STUCK read_snapshot0: lock=0
78+
79+
`pthread_create` was checked at the same time and never failed.
80+
81+
Six concurrent processes of `subst_state.py` (BP94-size model, one solver
82+
reused over forward/adjoint/`_reset` cycles), 12 minutes per arm:
83+
84+
| arm | stalled processes | cycles |
85+
|---|---|---|
86+
| as generated, run 1 | 1 / 6 | ~640 |
87+
| as generated, run 2 | 2 / 6 | ~440 |
88+
| fenced (first attempt), run 1 | 0 / 6 | 741 |
89+
| fenced (first attempt), run 2 | 0 / 6 | 730 |
90+
| **atomic (the fix), from source** | **0 / 6** | 205 |
91+
92+
## Fix
93+
94+
`is_atomic` qualifier in the code printer, set on `Lock` and `VolatileInt`, so
95+
the two objects are declared `_Atomic int`. Every access becomes an atomic
96+
operation; the generated statements are unchanged, so nothing downstream moves.
97+
Both clang and gcc-14 accept `_Atomic` under `-std=c99`, which is what devito
98+
compiles with.
99+
100+
A first attempt used explicit acquire/release fences instead. It fixed the
101+
stalls just as well but left the accesses non-atomic, so ThreadSanitizer still
102+
reported the race; the type change is the smaller and the standard-clean one.
103+
104+
## Notes
105+
106+
* A single process rarely hits it; contention is what makes it reproducible.
107+
* Killing a run mid-JIT leaves `devito-codepy-uid*/…/lock` behind, and codepy
108+
waits on it forever (it warns after 10 attempts and keeps sleeping). A hang
109+
straight after a killed run is usually that, not this bug.

‎tests/test_gpu_common.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -411,6 +411,11 @@ def test_tasking_in_isolation(self, opt):
411411
sections = FindNodes(Section).visit(op)
412412
assert len(sections) == 3
413413
assert str(sections[0].body[0].body[0].body[0].body[0]) == 'while(lock0[0] == 0);'
414+
# Both threads write the lock and the flag, so they are atomic rather
415+
# than merely volatile -- otherwise a delivery can be lost and the two
416+
# deadlock (see handshake_mfe.c)
417+
assert '_Atomic int lock0' in str(op.ccode)
418+
assert '_Atomic int flag' in str(op.ccode)
414419
body = op._func_table['release_lock0'].root.body
415420
assert str(body.body[0].condition) == 'Ne(lock0[0], 2)'
416421
assert str(body.body[1]) == 'lock0[0] = 0;'

0 commit comments

Comments
 (0)