Repository navigation
Expand file tree
/
Copy pathSnakefile
More file actions
431 lines (385 loc) · 19.5 KB
/
Copy pathSnakefile
File metadata and controls
431 lines (385 loc) · 19.5 KB
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
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
# The Snakemake build — ONE DAG, ONE entry (this file). See docs/plans/2026-07-14-snakemake-build.md.
# Per-source knobs (crs/nodata/negate/datum_offset_m/offset_surface/clamp_positive/unpack)
# live in sources/<id>/metadata.json. Run from the repo root:
# uv run snakemake sources [--config source=<id>] [-n]
# uv run snakemake catalogs # + masks, covering, coverage
# uv run snakemake check --config source=<id>
# uv run snakemake publish [--config source=<id>] # per-source R2 push
# uv run snakemake bundles [publish_mosaic] # stage 2/3: mosaic → forks → bundles
# ./docker.sh snakemake sources # same, in the toolchain container
# Pipeline-code edits don't invalidate outputs; force explicitly (-R prep_source / -F).
#
# `cover` is a CHECKPOINT: the covering (the stem inventory the whole build scopes from) is a
# runtime product, so stage 2/3 targets can't enumerate their per-stem jobs at parse. Their
# input functions derive STEMS through the checkpoint (see pipelines/build.smk), so ONE
# invocation walks sources → covering → mosaic → forks → bundles with no second `-s` entry.
include: "pipelines/common.smk"
ALL_SOURCES = pipeline_config.sources()
# raw: true ⇒ register remote objects 1:1 (header reads + object mirror); everything else preps locally
RAW = [s for s in ALL_SOURCES if pipeline_config.load_metadata(s).get("raw")]
PROCESSED = [s for s in ALL_SOURCES if s not in RAW]
# listed ⇒ the enumerate checkpoint lists a bucket/urllist (vs echoing a committed file_list);
# these re-enumerate on the weekly cron. Orthogonal to processed/raw: nz_coastal is listed
# AND processed, the raw collections are listed AND raw, the rest are static.
LISTED = [s for s in ALL_SOURCES if "filter" in pipeline_config.load_metadata(s)]
# ── streaming ("source access resolves per source, local wins") ──────────────────────
# --config stream=1 (laptop preview): a source with no local prep evidence
# (store/source/<id>/raw/) fetches its published catalog.json (the whole registration) from
# R2 instead of prepping, and its COGs stream via SOURCE_VSI_BASE.
# A locally-prepped source keeps its prep rules — local wins. Off (the box, the
# sources workflow): STREAMED is empty and nothing changes.
STREAM_BASE = config.get("stream_base", "https://data.openwaters.io/bathymetry/source")
_wd = Path(config.get("workdir", str(SCRIPTS)))
STREAMED = ([s for s in ALL_SOURCES if not (_wd / "store/source" / s / "raw").is_dir()]
if config.get("stream") else [])
LOCAL_PROCESSED = [s for s in PROCESSED if s not in STREAMED]
LOCAL_RAW = [s for s in RAW if s not in STREAMED]
LOCAL_LISTED = [s for s in LISTED if s not in STREAMED]
# offset_surface: <name> ⇒ that source's prep subtracts store/datum/<name>.tif per pixel —
# the chart datum's height in the source's own vertical frame, for the tidal separations no
# scalar datum_offset_m can express. config.DATUM_BUILDERS names the module that composes each;
# a surface neither registered nor bespoke is a hard error, because the alternative is one
# builder writing another's filename and silently correcting a coast onto the wrong
# continent's datum.
DATUM_BUILDERS = pipeline_config.DATUM_BUILDERS
DATUM_SURFACES = sorted({name for name in
(pipeline_config.load_metadata(s).get("offset_surface")
for s in ALL_SOURCES) if name})
# Composed by their own rule instead of the registry (datum_surface_dgm_w: the builder lives in
# sources/dgm_w/ with its own gauge-profile input, a shape the generic rule can't express).
BESPOKE_SURFACES = {"dgm_w_lowwater"}
REGISTRY_SURFACES = [n for n in DATUM_SURFACES if n not in BESPOKE_SURFACES]
_unbuildable = [n for n in REGISTRY_SURFACES if n not in DATUM_BUILDERS]
if _unbuildable:
raise WorkflowError(f"offset_surface with no builder in config.DATUM_BUILDERS: "
f"{_unbuildable}")
ONLY = config.get("source")
if ONLY and ONLY not in PROCESSED + RAW:
raise WorkflowError(f"unknown source {ONLY!r} — known: {PROCESSED + RAW}")
TARGETS = [ONLY] if ONLY else PROCESSED + RAW
include: "pipelines/publish.smk" # needs the source lists above
rule sources:
input:
expand("store/source/{source}/catalog.json", source=TARGETS),
expand("store/polygon/{source}.gpkg", source=[s for s in TARGETS if s in PROCESSED]),
# Enumerate every source's fetchable items through ONE code path (source_enumerate). A
# CHECKPOINT: items.txt is a runtime product (a listed source's items come from a live bucket
# listing), so the per-item fetch jobs can't be enumerated at parse — the DAG re-evaluates
# once items.txt lands (raw_assets reads it via checkpoints.enumerate.get). Static sources
# echo their committed file_list; listed sources list + filter + dedupe. write-if-changed, so
# an unchanged enumeration doesn't cascade.
checkpoint enumerate:
input:
str(SOURCES_DIR / "{source}/file_list.txt"),
metadata=str(SOURCES_DIR / "{source}/metadata.json"),
output:
"store/source/{source}/items.txt"
wildcard_constraints:
source=pat(LOCAL_PROCESSED + LOCAL_RAW)
log:
f"{TMP}/logs/enumerate/{{source}}.log"
shell:
"{PY}/source_enumerate.py {wildcards.source} 2> {log}"
# One download per enumerated item; the raw name is the URL hash (source_fetch), so inserting
# an item can't re-key the others and a warm legacy raw/<index> self-migrates without refetch.
# items.txt is `ancient` — it must exist (built by enumerate), but a re-enumeration must not
# invalidate raws whose URL is unchanged.
#
# temp(): a raw is deleted as soon as prep_source has staged it. The durable artifacts are the
# staged/normalized COGs and catalog.json — the raw archive is a means, and holding every
# source's upstream bytes on the volume costs more than refetching the rare source that needs
# re-staging (Litto3D alone is 191 GB of 7z for ~1 GB of grids). Snakemake does not re-fetch a
# missing temp input whose downstream output is up to date, so a normal run stays a no-op; a
# forced re-prep re-fetches on demand, then discards again.
rule fetch_item:
input:
ancient("store/source/{source}/items.txt")
output:
temp("store/source/{source}/raw/{hash}")
wildcard_constraints:
source=pat(LOCAL_PROCESSED), hash=r"[0-9a-f]{16}"
retries: 2
benchmark:
f"{TMP}/bench/fetch/{{source}}-{{hash}}.tsv"
# stderr per job for forensics (a failed job's diagnostics stay isolated in its own log);
# stdout keeps flowing to the run log so monitors and Actions heartbeats parse progress.
log:
f"{TMP}/logs/fetch/{{source}}-{{hash}}.log"
shell:
"{PY}/source_fetch.py {wildcards.source} {wildcards.hash} 2> {log}"
# Stage → datum → normalize → register, one job per source. source_catalog scans the
# normalized COGs into the item's seascape:files, so catalog.json is the SINGLE registration
# artifact (the separate bounds/catalog steps collapsed into this one).
rule prep_source:
input:
raw_assets,
metadata=str(SOURCES_DIR / "{source}/metadata.json"),
recipe=recipe_files, # everything --hash-recipe hashes, so any recipe edit restamps
surface=offset_surface, # the datum reference, for a source that declares one
output:
# staged tif names aren't knowable at parse; catalog.json is the declared artifact
"store/source/{source}/catalog.json"
wildcard_constraints:
source=pat(LOCAL_PROCESSED)
params:
version=1, # increment to force a rebuild
# The staged COGs' shoal pyramids branch on the cap (utils._block_reduce), and the
# aggregation warp reads those pyramids back — a cap change must re-prep.
drying_cap=pipeline_config.DRYING_CAP,
priority: source_priority
# One worker per thread over the source's staged files (source_prep.DEFAULT_WORKERS holds
# the workers x GDAL-threads arithmetic). 4 is the box's own half-the-vCPUs figure.
threads: 4
resources:
# Staging still bounds this: a zip member is read whole into memory, and asc-mosaic
# holds a raster. The fan-out fits inside it — 4 workers measured 2.5 GB peak RSS on
# 1/9" CUDEM (both transforms stripe, so a worker scales with raster WIDTH, not size).
mem_gb=8
benchmark:
f"{TMP}/bench/prep/{{source}}.tsv"
log:
f"{TMP}/logs/prep/{{source}}.log"
shell:
"( {PY}/source_prep.py {wildcards.source} {threads} && "
"{PY}/source_catalog.py {wildcards.source} --hash-recipe ) 2> {log}"
# The weekly forced source refresh — run as its OWN invocation before catalogs/publish:
# a forced producer schedules all dependents at plan time, but across an invocation
# boundary the engine's checksums cure unchanged registrations, so the main invocation
# cascades only on real upstream drift. Scoped to the LISTED sources: they are the ones whose
# enumeration can drift week to week (a bucket/urllist listing); static sources change only
# when their committed file_list.txt does. The workflow runs `refresh -R enumerate`, so a
# re-listing that produces an unchanged items.txt cures downstream at the boundary.
# The objects push rides here too: mirror_objects is on the mirror.txt branch, NOT the
# catalog cascade the boundary suppresses, so pulling it into this invocation overlaps each
# source's ~190 GB copy with the next source's re-listing instead of stranding it behind the
# barrier. Objects stay additive and land before catalog.json, so cross-boundary atomicity holds.
rule refresh:
input:
expand("store/source/{source}/catalog.json", source=LOCAL_LISTED),
expand("store/meta/publish/{source}.objects",
source=[s for s in LOCAL_LISTED if s in RAW]),
# Raw collections register objects/<key> rows off the public bucket, reading the
# enumerate checkpoint's items.txt. source_mirror header-reads (with carry-forward) into the
# .rows.csv handoff, then source_catalog folds it into catalog.json — the SINGLE registration
# artifact. Re-listing on cadence is the enumerate checkpoint's job (the weekly workflow runs
# `refresh -R enumerate`); this step just materializes the registration from items.txt.
rule register:
input:
"store/source/{source}/items.txt",
recipe=recipe_files, # any recipe edit restamps the item's hash
output:
# `update` keeps the previous catalog.json in place so source_mirror can carry forward
# its seascape:files without a full re-read. mirror.txt/mirror-bucket.txt keep their
# names — they're the object-copy store contract, out of scope to rename.
catalog=update("store/source/{source}/catalog.json"),
mirror="store/source/{source}/mirror.txt",
bucket="store/source/{source}/mirror-bucket.txt",
wildcard_constraints:
source=pat(LOCAL_RAW)
# Priority BANDS, separated by orders of magnitude so no byte-weighted prep (raw MB,
# realistically <1,000,000) can cross into a higher band, and so the values dominate
# the scheduler's packing objective (which otherwise favors count-maximizing selection of
# light jobs): masks 10M > registrations 5M > preps (raw MB).
priority: 5_000_000 # long serial job with thousands of network header reads; start early
retries: 2
resources:
mem_gb=2 # header reads + list bookkeeping, no raster in memory
benchmark:
f"{TMP}/bench/register/{{source}}.tsv"
log:
f"{TMP}/logs/register/{{source}}.log"
shell:
"( {PY}/source_mirror.py {wildcards.source} && "
"{PY}/source_catalog.py {wildcards.source} --hash-recipe ) 2> {log}"
rule polygon:
input:
"store/source/{source}/catalog.json",
output:
"store/polygon/{source}.gpkg"
wildcard_constraints:
source=pat(LOCAL_PROCESSED)
threads: 4
log:
f"{TMP}/logs/polygon/{{source}}.log"
shell:
"{PY}/source_polygonize.py {wildcards.source} {threads} 2> {log}"
# The streamed half of the catalogs invocation (--config stream=1): fetch a not-locally-
# prepped source's published catalog.json — the whole registration. Absent-only (engine:
# outputs exist ⇒ done); -R fetch_catalog refreshes. tmp+mv so a 404 never leaves a truncated
# file. Fallback: a published item from before phase 4 has no seascape:files, so ALSO fetch
# the sibling bounds.csv that config.source_files falls back to (drop when every source has
# re-registered — mirrors the config.source_files fallback, self-checked there).
rule fetch_catalog:
output:
catalog="store/source/{source}/catalog.json",
wildcard_constraints:
source=pat(STREAMED)
retries: 2
log:
f"{TMP}/logs/fetch_catalog/{{source}}.log"
shell:
"( set -e; c={output.catalog}; b=store/source/{wildcards.source}/bounds.csv; "
"curl -fsS {STREAM_BASE}/{wildcards.source}/catalog.json -o $c.tmp && mv $c.tmp $c; "
"if ! grep -q '\"seascape:files\"' $c; then "
" curl -fsS {STREAM_BASE}/{wildcards.source}/bounds.csv -o $b.tmp && mv $b.tmp $b; "
"fi ) 2> {log} || {{ rm -f {output.catalog}.tmp store/source/{wildcards.source}/bounds.csv.tmp; exit 1; }}"
# The chart-datum reference a prep subtracts (source_datum --offset-surface) — a support
# artifact like the landmask, NOT a sources/ entry (everything under sources/ enters the merge).
# Composing one downloads its pinned upstream (NOAA's 3.2 GB VDatum bundle; Shom's BATHYELLI
# plus IGN's geoid grids) and caches it beside the output, so the store's copy IS the cache.
# Keyed on the module that holds both the pin and the composition: a formula fix must not ship
# under the old grid.
rule datum_surface:
input:
lambda wc: str(SCRIPTS / DATUM_BUILDERS[wc.name]),
output:
"store/datum/{name}.tif"
params:
builder=lambda wc: DATUM_BUILDERS[wc.name]
wildcard_constraints:
name=pat(REGISTRY_SURFACES)
priority: 5_000_000 # a source prep waits on it; same band as the registrations
retries: 2 # the bundle fetch is one 3.2 GB stream from vdatum.noaa.gov
resources:
mem_gb=8 # a region's grids in flight + the coarse fill pass
benchmark:
f"{TMP}/bench/datum_surface/{{name}}.tsv"
log:
f"{TMP}/logs/datum_surface/{{name}}.log"
shell:
"{PY}/{params.builder} --out {output} 2> {log}"
# dgm_w's reference is bespoke (the BSH SKN grid + checked-in gauge profiles composed by the
# source's own script), so it gets its own rule under the same store/datum/ contract.
rule datum_surface_dgm_w:
input:
script=str(SOURCES_DIR / "dgm_w" / "build_reference.py"),
gauges=str(SOURCES_DIR / "dgm_w" / "tideelbe_skn.csv"),
output:
"store/datum/dgm_w_lowwater.tif"
priority: 5_000_000 # a source prep waits on it; same band as the registrations
retries: 2 # the BSH SKN grid fetch is one stream from gdi.bsh.de
resources:
mem_gb=8 # the widened SKN canvas in memory
benchmark:
f"{TMP}/bench/datum_surface/dgm_w_lowwater.tsv"
log:
f"{TMP}/logs/datum_surface/dgm_w_lowwater.log"
shell:
"{PY}/../sources/dgm_w/build_reference.py --out {output} 2> {log}"
# Masks rebuild only when forced (-R landmask): pinned snapshot/release ⇒ no data drift.
rule landmask:
output:
"store/landmask/land.fgb"
priority: 10_000_000 # top band (see register): long single-threaded, ready at t=0 — overlap, don't tail
retries: 2
resources:
mem_gb=4
benchmark:
f"{TMP}/bench/landmask.tsv"
log:
f"{TMP}/logs/landmask.log"
shell:
"{PY}/landmask.py prep 2> {log}"
# Effective land (land ∖ water) rasterized once onto the z8 render grid — the overview
# terrain stems' mask (a continental vector clip per coarse stem would re-read the whole
# multi-GB FGB; this is windowed instead).
rule landraster:
input:
"store/landmask/land.fgb",
"store/landmask/water.fgb",
output:
"store/landmask/land-z8.tif"
priority: 10_000_000 # see landmask
resources:
mem_gb=8 # planet gdal_rasterize + mode overviews
benchmark:
f"{TMP}/bench/landraster.tsv"
log:
f"{TMP}/logs/landraster.log"
shell:
"{PY}/landmask.py prep-raster 2> {log}"
rule watermask:
output:
"store/landmask/water.fgb"
params:
version=2, # increment to force a rebuild
priority: 10_000_000 # see landmask
# No retries: the planet read is ~9 h, and every failure seen here has been deterministic
retries: 0
threads: 8 # the planet read is tiled + parallel (landmask._water_tile); IO-bound S3 reads
resources:
mem_gb=8 # the planet Overture-water reproject; refine from the benchmark
benchmark:
f"{TMP}/bench/watermask.tsv"
log:
f"{TMP}/logs/watermask.log"
shell:
"{PY}/landmask.py prep-water {threads} 2> {log}"
# The covering — the seam between stage 1 and stage 2/3, as a CHECKPOINT: it writes the stem
# inventory the build scopes from, so the engine re-evaluates the DAG (per-stem mosaic/fork/
# terrain jobs) AFTER it runs. BBOX rides as a param so a changed window reruns it (write-if-
# changed prunes in-window stale tiles, keeps out-of-window). Body/inputs/outputs are the plain
# rule's — only the `checkpoint` keyword differs.
checkpoint cover:
input:
expand("store/source/{source}/catalog.json", source=PROCESSED + RAW),
output:
"store/aggregation/covering.txt"
params:
version=2, # increment to force a re-derivation
bbox=os.environ.get("BBOX", ""),
max_child_z=pipeline_config.MAX_CHILD_Z, # ceiling changes re-derive the covering
log:
f"{TMP}/logs/cover.log"
shell:
"{PY}/aggregation_covering.py --stable 2> {log}"
# Source-coverage provenance tileset.
rule coverage:
input:
expand("store/polygon/{source}.gpkg", source=PROCESSED),
expand("store/source/{source}/catalog.json", source=PROCESSED + RAW),
"store/aggregation/covering.txt", # sequencing only; coverage reads footprints + catalogs
output:
"store/bundle/coverage.pmtiles"
log:
f"{TMP}/logs/coverage.log"
shell:
"{PY}/contour_run.py coverage 2> {log}"
# Serve-time land-mask tileset. Rebuilds only when the masks change (-R landmask/watermask
# upstream), not per DEM build — like coverage, a catalogs product stage_build ships from disk.
rule land:
input:
land="store/landmask/land.fgb",
water="store/landmask/water.fgb",
output:
"store/bundle/land.pmtiles"
log:
f"{TMP}/logs/land.log"
shell:
"{PY}/landmask.py tiles 2> {log}"
# The complete stage-1 product set (what the stage-2/3 invocation parses from).
rule catalogs:
input:
expand("store/source/{source}/catalog.json", source=PROCESSED + RAW),
"store/landmask/land.fgb",
"store/landmask/water.fgb",
"store/aggregation/covering.txt",
"store/bundle/coverage.pmtiles",
"store/bundle/land.pmtiles",
# Validate one source's contract (requires --config source=<id>).
def check_input(wc):
if not ONLY:
raise WorkflowError("check requires --config source=<id>")
return f"store/source/{ONLY}/catalog.json"
rule check:
input:
check_input,
params:
source=lambda wc: ONLY,
log:
f"{TMP}/logs/check.log"
shell:
"{PY}/source_check.py {params.source} 2> {log}"
# Stage 2/3 (mosaic → cartographic forks → bundles → publish), gated on the `cover` checkpoint
# above. Kept in its own file for grouping — but it is INCLUDED here, so there is one entry.
include: "pipelines/build.smk"