Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -167,7 +167,7 @@ Tag PR bodies with `Relates to NewGraphEnvironment/sred-2025-2026#16` (the issue
sliver-settles-less result in 6 of 8 bands. Neither geometric leg generalises.
- [`inst/notes/temporal-qa-groups.md`](inst/notes/temporal-qa-groups.md) — `dft_rast_break_class()` across the four published seven-year groups (#62): sustained-break share 20-31% of 2017-2023 change, flicker 40-49%, 2017 the odd endpoint everywhere; the shape proxies and why the per-stream segment layer was rejected. Q6 (#73) adds distance from permanent water: the gradient is real, the null does not separate it from any other class boundary, the near-band `stable` share is biased down by construction because the reference excises the permanent-water pixels, and `break_year` runs the opposite way to a migrating bank. Read before quoting a `transition_2017_2023` hectare as change.
- [`inst/notes/temporal-qa-disturbance.md`](inst/notes/temporal-qa-disturbance.md) — does dated fire/harvest corroborate the temporal QA (#67): the published `transition_vector.gpkg` is drift's own change layer after a 1 ha class-agnostic sieve and a sub-basin clip, reproduced exactly in all four groups (0 of 53/44/41/46 classes differ), which closes the two-totals problem; `break_year` lands at lag 0 or +1 for 86% of discriminating fire patches that broke and 72% of harvest (70.7% / 58.4% of all tagged), mode +1; flicker is LOWER in disturbed patches, so it does not read as succession. Read before treating drift's patch count and the published one as the same population.
- [`inst/notes/gdalcubes-pc-gotchas.md`](inst/notes/gdalcubes-pc-gotchas.md) — non-obvious gdalcubes 0.7.3 + Planetary Computer Sentinel-2 gotchas (filter_geom segfault, reduce_time worker closures, terra↔gdalcubes NetCDF round-trip, the S2 +1000 DN offset boundary at 2022-01-25, PC pagination). Read before touching the continuous pipeline.
- [`inst/notes/gdalcubes-pc-gotchas.md`](inst/notes/gdalcubes-pc-gotchas.md) — non-obvious gdalcubes 0.7.3 + Planetary Computer Sentinel-2 gotchas (filter_geom segfault, reduce_time worker closures, terra↔gdalcubes NetCDF round-trip, the S2 +1000 DN offset boundary at 2022-01-25, PC pagination, and what `chunk_status` does and does not record about a failed read, #87/#99). Read before touching the continuous pipeline.


<!-- BEGIN SOUL CONVENTIONS — DO NOT EDIT BELOW THIS LINE -->
Expand Down
1 change: 1 addition & 0 deletions DESCRIPTION
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ Suggests:
leaflet,
leaflet.extras,
mapgl,
ncdf4,
rmarkdown,
testthat (>= 3.2.0),
withr,
Expand Down
165 changes: 154 additions & 11 deletions R/dft_stac_cube.R
Original file line number Diff line number Diff line change
Expand Up @@ -401,7 +401,9 @@ cube_view_choice_check <- function(x, arg, allowed, fallback, class,
#' parallel = 1 measured finer chunks 45% SLOWER.
#'
#' GDAL cloud-read tuning for /vsicurl COG streaming (biggest win:
#' DISABLE_READDIR_ON_OPEN avoids a remote directory listing on every open).
#' DISABLE_READDIR_ON_OPEN avoids a remote directory listing on every open), and
#' HTTP retries (#87). Environment variables, so gdalcubes' worker processes
#' inherit them.
#' @noRd
stac_cube_session <- function(parallel) {
old_parallel <- gdalcubes::gdalcubes_options()$parallel
Expand All @@ -413,7 +415,11 @@ stac_cube_session <- function(parallel) {
GDAL_HTTP_MULTIPLEX = "YES",
GDAL_HTTP_VERSION = "2",
VSI_CACHE = "TRUE",
CPL_VSIL_CURL_ALLOWED_EXTENSIONS = ".tif"
CPL_VSIL_CURL_ALLOWED_EXTENSIONS = ".tif",
# A transient 429 or 5xx now aborts a whole read (#87), and one that strikes
# after an image has opened is a hole nothing can detect, so retry them.
GDAL_HTTP_MAX_RETRY = "3",
GDAL_HTTP_RETRY_DELAY = "5"
)
old_cfg <- Sys.getenv(names(gdal_cfg), unset = NA)
do.call(Sys.setenv, as.list(gdal_cfg))
Expand Down Expand Up @@ -549,10 +555,9 @@ stac_cube_items <- function(cfg, aoi_wgs84, datetime, cloud_cover_max, months,
#' expired token and replaces the `sig` parameter of an already-signed href, so
#' re-signing before each extent keeps every read inside a live token (verified
#' live: features with a corrupted `sig` read 102,364 cells after re-signing, 0
#' without). A single extent whose read outlives one token is still exposed, and
#' gdalcubes reports a failed chunk only on the worker's stderr, which R cannot
#' reliably capture (#87). A `fetched` without `sign_fn` (a test stub) passes
#' through unchanged.
#' without). A single extent whose read outlives one token is still exposed; an
#' image that then fails to open aborts the read through cube_write_ncdf() (#87).
#' A `fetched` without `sign_fn` (a test stub) passes through unchanged.
#' @noRd
stac_features_resign <- function(fetched, feats) {
if (is.null(fetched$sign_fn) || is.null(fetched$items) || !length(feats)) {
Expand Down Expand Up @@ -639,7 +644,8 @@ stac_cube_assemble <- function(fetched, cfg, aoi_target, target_crs, t0, t1,
)
out <- pixel_fn(cube, offset_use)
tmp <- tempfile(fileext = ".nc")
gdalcubes::write_ncdf(out, tmp, overwrite = TRUE)
# aborts on a failed chunk read, before any caller caches the result (#87)
cube_write_ncdf(out, tmp)
terra::rast(tmp)
}

Expand All @@ -661,10 +667,12 @@ stac_cube_assemble <- function(fetched, cfg, aoi_target, target_crs, t0, t1,
# re-sign here, per extent, so no read outlives its token
feats <- stac_features_resign(fetched, features)
if (any(is_pre) && !all(is_pre)) {
terra::cover(
build_stack(feats[is_pre], offset_before, v),
build_stack(feats[!is_pre], offset, v)
)
# Both sides built BEFORE terra::cover(): built as its arguments, they are
# evaluated inside S4 method selection, which re-raises an abort as a plain
# "error in evaluating the argument" and drops its class (#87).
pre <- build_stack(feats[is_pre], offset_before, v)
post <- build_stack(feats[!is_pre], offset, v)
terra::cover(pre, post)
} else {
build_stack(feats, if (all(is_pre)) offset_before else offset, v)
}
Expand Down Expand Up @@ -720,6 +728,141 @@ cube_check_nonempty <- function(stk, collection, datetime, cached) {
}


#' Write a gdalcubes cube to NetCDF, aborting if any chunk failed to read
#'
#' The one place drift calls [gdalcubes::write_ncdf()], shared by
#' stac_cube_assemble() (so [dft_stac_cube()] and [dft_stac_composite()]) and
#' fetch_extent_to() ([dft_stac_fetch()]). Both call it before anything derived
#' from the file is cached, so a failed read aborts with nothing cached.
#'
#' gdalcubes signals a failed chunk in two ways, and this watches both:
#' * a chunk whose merge into the output throws raises an R warning, `Chunk N
#' could not be added to output`, and the write carries on (gdalcubes 0.7.5,
#' `src/multiprocess.cpp`). That chunk holds no status at all, so
#' cube_check_chunks() alone would pass it.
#' * a chunk whose image would not open is recorded in the file; see
#' cube_check_chunks().
#'
#' The warning is only NOTED in the handler, and the abort comes after
#' `write_ncdf()` returns. gdalcubes raises it with `Rcpp::warning()`, so an
#' abort from inside the handler would longjmp through its C++ frames and skip
#' their destructors: the worker shutdown, the work directory, the output
#' file's close (Rcpp's own header warns of exactly this). Other gdalcubes
#' warnings pass through untouched.
#'
#' Not every failed merge reaches either signal: a worker chunk file the main
#' process cannot open is skipped with status OK and no warning, a C++ stderr
#' line only (round-2 review of #87; drift#99).
#' @noRd
cube_write_ncdf <- function(cube, out) {
unmerged <- character(0)
withCallingHandlers(
gdalcubes::write_ncdf(cube, out, overwrite = TRUE),
warning = function(w) {
if (grepl("could not be added to output", conditionMessage(w),
fixed = TRUE)) {
unmerged <<- c(unmerged, conditionMessage(w))
invokeRestart("muffleWarning")
}
}
)
if (length(unmerged)) {
cli::cli_abort(class = "drift_incomplete_cube", c(
"gdalcubes could not merge {length(unmerged)} chunk{?s} into its output.",
"x" = "{unmerged[[1]]}",
"i" = "Nothing was cached."
))
}
cube_check_chunks(out)
}


#' Abort on a gdalcubes NetCDF whose chunk reads failed
#'
#' gdalcubes does not raise when it cannot read a chunk. It prints `[WARNING] n
#' out of m chunks have repoprted errors / incompleteness` from C++ stdio and
#' writes the chunk as `NA`. That line never reaches R from a worker process, and
#' only sometimes from the main one, so capturing it was measured unreliable
#' (#79, #87).
#'
#' What gdalcubes does write, at every worker count, is a `chunk_status` variable
#' in the output NetCDF, one integer per chunk: `0` OK, `1` ERROR, `2`
#' INCOMPLETE, `128` UNKNOWN (gdalcubes 0.7.5, `src/gdalcubes/src/cube.cpp`). A
#' worker carries the status to the main process inside its chunk file.
#'
#' **What it records is narrower than "the read failed".** Measured offline
#' (#87): a band image that fails to OPEN marks its chunks INCOMPLETE, which is
#' the expired or refused signed URL. A band image that opens and then fails
#' mid-read (a corrupted tile here; a throttled range request on the network)
#' leaves the status OK, and so does a mask (SCL) image that fails to open, which
#' silently drops that scene. A mask image that opens and then fails mid-read is
#' worse: the scene is used unmasked. gdalcubes sets the status only when a band
#' image will not open; every other I/O step drops its return code, so none of
#' these is visible to this check or to anything else R can see (drift#99;
#' inst/notes/gdalcubes-pc-gotchas.md).
#'
#' **The netCDF integer fill value means no status was written for that id**,
#' not that the chunk was empty: a worker skips writing a chunk that is OK and
#' all-`NA` (no image, or fully masked), and a chunk whose merge threw keeps fill
#' too. That second case raises an R warning, which cube_write_ncdf() turns into
#' an abort, so fill is read as OK here. A worker chunk file the main process
#' cannot open is recorded as OK, with neither signal (drift#99). Anything that is neither OK nor fill fails, so a
#' status code a later gdalcubes adds aborts rather than passes.
#'
#' For an open failure, this is the only check that sees it. A median over two
#' scenes where one failed fills every cell from the survivor (measured: 1600 of
#' 1600 cells set), so no test on the pixels, cube_check_nonempty() included, can
#' tell that cube from a complete one.
#'
#' A file with no `chunk_status` variable also aborts: a guard that cannot find
#' its evidence must not report the read clean.
#'
#' Read with ncdf4, which gdalcubes imports: GDAL does not list the
#' one-dimensional variable as a subdataset, so terra cannot open it by name.
#' @noRd
cube_check_chunks <- function(nc_file) {
failed <- chunk_status_failed(nc_file)
if (is.na(failed[["failed"]])) {
cli::cli_abort(class = "drift_incomplete_cube", c(
"The gdalcubes output has no {.field chunk_status}, so its chunk reads \\
cannot be checked.",
"i" = "Nothing was cached. drift needs a gdalcubes that records per-chunk \\
status in {.fn gdalcubes::write_ncdf} output (0.7.5 does)."
))
}
if (failed[["failed"]] == 0L) return(invisible(nc_file))
cli::cli_abort(class = "drift_incomplete_cube", c(
"gdalcubes could not read {failed[['failed']]} of {failed[['chunks']]} \\
chunk{?s}: an image would not open.",
"x" = "Those cells are {.val NA}, or computed from only the scenes that \\
did open.",
"i" = "Usually an expired or refused signed URL. Nothing was cached.",
"i" = "Re-running signs the request again. If the same read fails every \\
time, an asset may be missing from the catalogue."
))
}


#' Count the failed chunks recorded in a gdalcubes NetCDF
#'
#' The reader behind cube_check_chunks() and the cache gate. Returns integers
#' `failed` and `chunks`, with `failed` `NA` when the file has no
#' `chunk_status` variable. Status `0` and the netCDF integer fill
#' (`NC_FILL_INT`, no chunk merged for that id) count as OK; see
#' cube_check_chunks().
#' @noRd
chunk_status_failed <- function(nc_file) {
nc <- ncdf4::nc_open(nc_file)
on.exit(ncdf4::nc_close(nc), add = TRUE)
if (!"chunk_status" %in% names(nc$var)) {
return(c(failed = NA_integer_, chunks = NA_integer_))
}
status <- as.vector(ncdf4::ncvar_get(nc, "chunk_status", raw_datavals = TRUE))
failed <- !is.na(status) & status != 0L & status != -2147483647L
c(failed = sum(failed), chunks = length(status))
}


#' Resolve the gdalcubes worker count for a cube read
#'
#' `NULL` means auto: `min(4, cores - 1)`, capped at 4 rather than the full core
Expand Down
44 changes: 38 additions & 6 deletions R/dft_stac_fetch.R
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,18 @@ dft_stac_fetch <- function(aoi,
aggregation <- aggregation_check(aggregation)
resampling <- resampling_check(resampling)

# Retry a transient 429 or 5xx rather than abort the read on it: a failed
# open now aborts (#87). Tiled or not, restored on exit.
gdal_retry <- c(GDAL_HTTP_MAX_RETRY = "3", GDAL_HTTP_RETRY_DELAY = "5")
old_retry <- Sys.getenv(names(gdal_retry), unset = NA)
do.call(Sys.setenv, as.list(gdal_retry))
on.exit({
set_again <- old_retry[!is.na(old_retry)]
if (length(set_again)) do.call(Sys.setenv, as.list(set_again))
unset <- names(old_retry)[is.na(old_retry)]
if (length(unset)) Sys.unsetenv(unset)
}, add = TRUE)

# Normalize tile_size ONCE so the path gate (is.null) and the cache key derive
# from the same snapped scalar. When tiling, tune GDAL for the many extra
# per-item COG opens (restored on exit so the caller's session is untouched).
Expand Down Expand Up @@ -226,14 +238,17 @@ dft_stac_fetch <- function(aoi,
})
} else {
message(" ", yr, ": fetching ", length(tiles), " tile(s)...")
# Named and registered for removal BEFORE any read: on.exit, not a
# trailing unlink(), because a failed tile (#87) or a mosaic_tiles() error
# would otherwise strand every tile already written.
tile_files <- vapply(seq_along(tiles), function(i) {
fetch_extent_to(col, tiles[[i]], t0, t1, target_crs, res, dt,
aggregation, resampling,
tempfile(sprintf("drift_tile%d_", i), fileext = ".nc"))
tempfile(sprintf("drift_tile%d_", i), fileext = ".nc")
}, character(1))
# on.exit, not a trailing unlink(): a mosaic_tiles() error would otherwise
# strand every tile file for the life of the session.
on.exit(unlink(tile_files), add = TRUE)
for (i in seq_along(tiles)) {
fetch_extent_to(col, tiles[[i]], t0, t1, target_crs, res, dt,
aggregation, resampling, tile_files[[i]])
}
cache_write_atomic(cache_file, function(out) mosaic_tiles(tile_files, out))
}

Expand Down Expand Up @@ -629,7 +644,9 @@ fetch_extent_to <- function(col, ext, t0, t1, target_crs, res, dt,
aggregation = aggregation, resampling = resampling
)
cube <- gdalcubes::raster_cube(col, v)
gdalcubes::write_ncdf(cube, out_nc, overwrite = TRUE)
# aborts on a failed chunk read: inside cache_write_atomic() untiled, before
# the mosaic tiled, so nothing is cached (#87)
cube_write_ncdf(cube, out_nc)
out_nc
}

Expand Down Expand Up @@ -802,6 +819,21 @@ cache_invalid_reason <- function(path, probe = cache_probe_last_row) {
}
)
if (!ok) return("its pixel data could not be read cleanly")

# An untiled dft_stac_fetch() cache IS the gdalcubes NetCDF, so it still holds
# the chunk_status it was written with, and one cached before #87 can record an
# image that never opened. A miss re-fetches it. A .nc with no chunk_status
# (written by something other than gdalcubes) has nothing to say and passes:
# the read arms above already vouched for it.
if (identical(tolower(tools::file_ext(path)), "nc") &&
requireNamespace("ncdf4", quietly = TRUE)) {
failed <- tryCatch(chunk_status_failed(path)[["failed"]],
error = function(e) NA_integer_)
if (!is.na(failed) && failed > 0L) {
return(paste0("records ", failed, " chunk read(s) that failed when it ",
"was fetched"))
}
}
NA_character_
}

Expand Down
18 changes: 18 additions & 0 deletions data-raw/logs/probe_chunk_status_live/20261002T140832Z.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
14:08:32Z drift 0.21.0 (source tree), gdalcubes 0.7.5, cores 10, auto parallel 4
14:08:37Z PASS | composite untiled, bad tokens | expect drift_incomplete_cube | got drift_incomplete_cube | cache files 0 | 5 s | gdalcubes parallel after 1
14:08:37Z gdalcubes could not read 9 of 9 chunks: an image would not open. ✖ Those cells are "NA", or computed from only the scenes that did open. ℹ Usually an expired or refused signed URL. Nothing was cached. ℹ Re-running signs the request again. If the same read fails every time, an asset may be missing fr
14:08:41Z PASS | composite tiled 1000 m, bad tokens | expect drift_incomplete_cube | got drift_incomplete_cube | cache files 0 | 4 s | gdalcubes parallel after 1
14:08:41Z gdalcubes could not read 4 of 4 chunks: an image would not open. ✖ Those cells are "NA", or computed from only the scenes that did open. ℹ Usually an expired or refused signed URL. Nothing was cached. ℹ Re-running signs the request again. If the same read fails every time, an asset may be missing fr
14:08:52Z PASS | cube untiled, bad tokens | expect drift_incomplete_cube | got drift_incomplete_cube | cache files 0 | 10 s | gdalcubes parallel after 1
14:08:52Z gdalcubes could not read 9 of 9 chunks: an image would not open. ✖ Those cells are "NA", or computed from only the scenes that did open. ℹ Usually an expired or refused signed URL. Nothing was cached. ℹ Re-running signs the request again. If the same read fails every time, an asset may be missing fr
14:08:54Z PASS | fetch untiled, bad tokens, session parallel | expect drift_incomplete_cube | got drift_incomplete_cube | cache files 0 | 2 s | gdalcubes parallel after 1
14:08:54Z gdalcubes could not read 4 of 4 chunks: an image would not open. ✖ Those cells are "NA", or computed from only the scenes that did open. ℹ Usually an expired or refused signed URL. Nothing was cached. ℹ Re-running signs the request again. If the same read fails every time, an asset may be missing fr
14:08:56Z PASS | fetch tiled 1000 m, bad tokens, session parallel | expect drift_incomplete_cube | got drift_incomplete_cube | cache files 0 | 2 s | gdalcubes parallel after 1
14:08:56Z gdalcubes could not read 1 of 1 chunk: an image would not open. ✖ Those cells are "NA", or computed from only the scenes that did open. ℹ Usually an expired or refused signed URL. Nothing was cached. ℹ Re-running signs the request again. If the same read fails every time, an asset may be missing fro
14:13:22Z PASS | control: cube months 6:9, 2023, parallel 4 | expect returned | got returned | cache files 1 | 266 s | gdalcubes parallel after 1
14:13:22Z SpatRaster
14:15:45Z PASS | control: composite 2023 July, tiled 1000 m, parallel 4 | expect returned | got returned | cache files 1 | 142 s | gdalcubes parallel after 1
14:15:45Z list
14:15:52Z PASS | control: fetch 2023 untiled | expect returned | got returned | cache files 1 | 8 s | gdalcubes parallel after 1
14:15:52Z list
14:15:52Z 8 of 8 as expected
Loading
Loading