From bea92eba1357ef81ebb9e9ad962e1d17f0470811 Mon Sep 17 00:00:00 2001 From: mloubout Date: Tue, 29 Sep 2026 16:20:16 -0400 Subject: [PATCH 1/3] compiler: Make every reader of a prefetched buffer wait for it memcpy_prefetch attached the WaitLock to the first Cluster reading the buffer in the order it sees them. That Cluster need not run first: a HaloTouch is only a placeholder, later scheduled beside the stencil that needs the halo, and later passes reorder and fuse Clusters. Under the CUDA backend a reader of the buffer was then launched before the wait and read the previous snapshot, e.g. the mass term of a reflection-FWI adjoint while its halo-reading neighbours waited correctly. Every reader now carries the wait; one on a released lock is free, and Clusters with the same syncs fuse again, which also removes a kernel launch in that case. The prefetch writes the buffer, so data dependences keep it after all readers. --- devito/passes/clusters/asynchrony.py | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/devito/passes/clusters/asynchrony.py b/devito/passes/clusters/asynchrony.py index a64ccf6122..1aafd6a39a 100644 --- a/devito/passes/clusters/asynchrony.py +++ b/devito/passes/clusters/asynchrony.py @@ -345,13 +345,12 @@ def _actions_from_update_memcpy(c, d, bounds, clusters, actions, sregistry): pc = c.rebuild(exprs=expr, ispace=ispace, guards=guards, syncs=syncs) - # Wait before the first, prefetch after the last access to `target`, then + # Wait before every, prefetch after the last access to `target`, then # drop the memcpy `c`. `c` may be toposorted amid the readers, so scan them # all; count only reads over the streamed `pd`, not the buffer-init loop. reads = [c1 for c1 in clusters if c1 is not c and target in c1.scope.reads and pd in c1.ispace.itdims] assert reads - first = reads[0] last = reads[-1] # Advance `last` past its loop nest so the prefetch follows it rather than @@ -362,7 +361,13 @@ def _actions_from_update_memcpy(c, d, bounds, clusters, actions, sregistry): break last = c1 - actions[first].syncs[d].append(WaitLock(handle, target)) + # Every reader waits, not only the first in this order: later passes + # reorder and fuse Clusters, and a HaloTouch is a placeholder scheduled + # beside the stencil that needs the halo, so the Cluster first here need + # not run first. A wait on an already released lock is free. The + # prefetch writes `target`, so data dependences keep it after every reader + for c1 in reads: + actions[c1].syncs[d].append(WaitLock(handle, target)) actions[last].insert.append(pc) actions[c].drop = True From 2d8bfe395cf989a2389ee9680dffbdaf9846029a Mon Sep 17 00:00:00 2001 From: mloubout Date: Tue, 29 Sep 2026 22:04:07 -0400 Subject: [PATCH 2/3] compiler: Keep a prefetch reader's scalar definitions in its wait Waiting on the reader alone split an interpolation's point loop: the Clusters defining the positions posx/posy, right before the reader in the same IterationSpace, had no sync, landed in a separate nest and left the reader with undeclared scalars. The Clusters that only define scalars right before a reader, with no syncs of their own, now share its wait. --- devito/passes/clusters/asynchrony.py | 22 +++++++++++++++++++++- 1 file changed, 21 insertions(+), 1 deletion(-) diff --git a/devito/passes/clusters/asynchrony.py b/devito/passes/clusters/asynchrony.py index 1aafd6a39a..0d7af4c4a2 100644 --- a/devito/passes/clusters/asynchrony.py +++ b/devito/passes/clusters/asynchrony.py @@ -365,8 +365,18 @@ def _actions_from_update_memcpy(c, d, bounds, clusters, actions, sregistry): # reorder and fuse Clusters, and a HaloTouch is a placeholder scheduled # beside the stencil that needs the halo, so the Cluster first here need # not run first. A wait on an already released lock is free. The - # prefetch writes `target`, so data dependences keep it after every reader + # prefetch writes `target`, so data dependences keep it after every reader. + # A reader shares its wait with the Clusters right before it defining the + # scalars it uses, e.g. an interpolation's positions: with different + # syncs they would land in separate nests, out of the reader's scope + waiting = [] for c1 in reads: + i = clusters.index(c1) + while i > 0 and defines_scalars_for(clusters[i-1], c1): + i -= 1 + waiting.extend(c2 for c2 in clusters[i:clusters.index(c1)+1] + if c2 not in waiting) + for c1 in waiting: actions[c1].syncs[d].append(WaitLock(handle, target)) actions[last].insert.append(pc) actions[c].drop = True @@ -382,6 +392,16 @@ def __init__(self, drop=False, syncs=None, insert=None): self.insert = insert or [] +def defines_scalars_for(c0, c1): + """ + True if `c0` only defines scalars, in `c1`'s very IterationSpace, and + carries no syncs of its own. + """ + return (c0.ispace.itdims == c1.ispace.itdims and + not c0.syncs and + not any(f.is_AbstractFunction for f in c0.scope.writes)) + + def wraps_memcpy(cluster): return len(cluster.exprs) == 1 and is_memcpy(cluster.exprs[0]) From 751250e0d40e973473cd9f4d3de97cce388817bb Mon Sep 17 00:00:00 2001 From: mloubout Date: Tue, 29 Sep 2026 22:26:24 -0400 Subject: [PATCH 3/3] compiler: Schedule the readers of a waited-on Function after the wait Fusion's toposort only ordered ClusterGroups by data hazards and fences. Two readers of a prefetched buffer carry no hazard between them, so a reader could be scheduled ahead of the Cluster holding the WaitLock on it, and read the previous snapshot: under the CUDA backend the mass term of a reflection-FWI adjoint ran before the wait while its halo readers ran after. The DAG now orders every later reader of a waited-on Function after the wait, and memcpy_prefetch is back to a single WaitLock on the first reader. --- devito/passes/clusters/asynchrony.py | 31 +++------------------------- devito/passes/clusters/fusion.py | 7 +++++++ 2 files changed, 10 insertions(+), 28 deletions(-) diff --git a/devito/passes/clusters/asynchrony.py b/devito/passes/clusters/asynchrony.py index 0d7af4c4a2..a64ccf6122 100644 --- a/devito/passes/clusters/asynchrony.py +++ b/devito/passes/clusters/asynchrony.py @@ -345,12 +345,13 @@ def _actions_from_update_memcpy(c, d, bounds, clusters, actions, sregistry): pc = c.rebuild(exprs=expr, ispace=ispace, guards=guards, syncs=syncs) - # Wait before every, prefetch after the last access to `target`, then + # Wait before the first, prefetch after the last access to `target`, then # drop the memcpy `c`. `c` may be toposorted amid the readers, so scan them # all; count only reads over the streamed `pd`, not the buffer-init loop. reads = [c1 for c1 in clusters if c1 is not c and target in c1.scope.reads and pd in c1.ispace.itdims] assert reads + first = reads[0] last = reads[-1] # Advance `last` past its loop nest so the prefetch follows it rather than @@ -361,23 +362,7 @@ def _actions_from_update_memcpy(c, d, bounds, clusters, actions, sregistry): break last = c1 - # Every reader waits, not only the first in this order: later passes - # reorder and fuse Clusters, and a HaloTouch is a placeholder scheduled - # beside the stencil that needs the halo, so the Cluster first here need - # not run first. A wait on an already released lock is free. The - # prefetch writes `target`, so data dependences keep it after every reader. - # A reader shares its wait with the Clusters right before it defining the - # scalars it uses, e.g. an interpolation's positions: with different - # syncs they would land in separate nests, out of the reader's scope - waiting = [] - for c1 in reads: - i = clusters.index(c1) - while i > 0 and defines_scalars_for(clusters[i-1], c1): - i -= 1 - waiting.extend(c2 for c2 in clusters[i:clusters.index(c1)+1] - if c2 not in waiting) - for c1 in waiting: - actions[c1].syncs[d].append(WaitLock(handle, target)) + actions[first].syncs[d].append(WaitLock(handle, target)) actions[last].insert.append(pc) actions[c].drop = True @@ -392,16 +377,6 @@ def __init__(self, drop=False, syncs=None, insert=None): self.insert = insert or [] -def defines_scalars_for(c0, c1): - """ - True if `c0` only defines scalars, in `c1`'s very IterationSpace, and - carries no syncs of its own. - """ - return (c0.ispace.itdims == c1.ispace.itdims and - not c0.syncs and - not any(f.is_AbstractFunction for f in c0.scope.writes)) - - def wraps_memcpy(cluster): return len(cluster.exprs) == 1 and is_memcpy(cluster.exprs[0]) diff --git a/devito/passes/clusters/fusion.py b/devito/passes/clusters/fusion.py index bbd4694614..2a699551ad 100644 --- a/devito/passes/clusters/fusion.py +++ b/devito/passes/clusters/fusion.py @@ -305,10 +305,17 @@ def _build_dag(self, cgroups, prefix): # Track whether there is any fence between `cg0` and the current `cg1`. fenced = cg0.scope.has_barrier + # Functions `cg0` waits on, e.g. a prefetched buffer: reading them + # is no data hazard, but no reader may be scheduled before the wait + waited = {s.target for s in flatten(cg0.syncs.values()) + if isinstance(s, WaitLock)} + for n1, cg1 in enumerate(cgroups[n+1:], start=n+1): fenced = fenced or cg1.scope.has_barrier hazard = _fusion_hazards(cg0.scope, cg1.scope, prefix) + if waited & set(cg1.scope.reads): + dag.add_edge(cg0, cg1) if not (hazard or fenced): continue