diff --git a/DESCRIPTION b/DESCRIPTION index fb46a9e..52d227c 100644 --- a/DESCRIPTION +++ b/DESCRIPTION @@ -1,7 +1,7 @@ Package: filearray Type: Package Title: File-Backed Array for Out-of-Memory Computation -Version: 0.2.3.1 +Version: 0.2.3.3 Language: en-US Encoding: UTF-8 License: LGPL-3 diff --git a/NAMESPACE b/NAMESPACE index 668c98b..ad8c1fc 100644 --- a/NAMESPACE +++ b/NAMESPACE @@ -38,6 +38,7 @@ export(fmap_element_wise) export(fwhich) export(mapreduce) export(typeof) +export(with_filelock) exportClasses(FileArray) exportClasses(FileArrayProxy) exportMethods(apply) diff --git a/NEWS.md b/NEWS.md index 31ff486..a67b059 100644 --- a/NEWS.md +++ b/NEWS.md @@ -2,6 +2,16 @@ * Partition files are written to a temporary file and then renamed into place, so other processes loading the same array no longer fail with "File size too small" while a partition is being created or filled * `filearray_load` only treats files named `.farr` as partitions +* Writes lock the partitions they change, so several processes (for example `parallel`, `callr` or `mirai` workers, or several applications) can write to one array without losing each other's data: sub-assignment (`x[...] <- value`), `fill_partition`, `initialize_partition`, `set_partition`, `fmap`, `fmap_element_wise` and `filearray_bind` wait while another process holds one of their partitions. Reading takes no locks, except that functions that must first create missing partitions (such as `max` or `sum`) lock the partitions they create +* Added `lock_partition(parts, timeout)` and `unlock_partition(parts)` methods to hold partitions over several operations (for example to check whether a partition has been computed and write it only if not); `lock_partition` waits and retries until the partitions are free or `timeout` seconds have passed +* Locks are byte-range locks on an empty file `.filearray.lock` in the array directory, released by the operating system when a process exits. `options(filearray.lock = FALSE)` (environment variable `FILEARRAY_LOCK`) turns off automatic locking; `options(filearray.lock.timeout = )` (`FILEARRAY_LOCK_TIMEOUT`) limits how long writes wait, by default without limit +* A missing partition is created only if no other process has created it in the meantime, and every partition write uses its own temporary file, so concurrent writers no longer replace each other's partitions or share a temporary file +* `fmap` stops when the mapping function raises an error or is interrupted, instead of warning "cannot finish map" and returning a partly written result +* `delete()` releases the locks the session holds on the array before removing it +* Added `with_filelock(expr, lockfile, timeout, on_unsupported)` to evaluate code while holding an exclusive lock on any lock file, so that only one R session at a time runs it; other sessions wait. It is not tied to file arrays, so other packages can use it. `on_unsupported` chooses what happens when the lock cannot be taken: warn (default), stop, or ignore +* Partition locks never make a write fail when the lock itself cannot be taken: any error while locking (for example a stale network file handle or too many open files, not only file systems without lock support) now warns once and the write proceeds, and an invalid `filearray.lock.timeout` option or `FILEARRAY_LOCK_TIMEOUT` value warns once and is ignored instead of stopping every write +* Fixed `mapreduce` (and so `sum`, `max`, `min`, `range` and `fwhich`) reading only the first slice of every partition after the first when `partition_size` is greater than 1, which gave wrong results and wrong `fwhich` indices +* `fmap2` reports the error raised by the mapping function instead of "Unknown error." # filearray 0.2.3 diff --git a/R/RcppExports.R b/R/RcppExports.R index 7ed797c..3428963 100644 --- a/R/RcppExports.R +++ b/R/RcppExports.R @@ -77,6 +77,42 @@ FARR_subset2 <- function(filebase, listOrEnv, reshape = NULL, drop = FALSE, use_ .Call(`_filearray_FARR_subset2`, filebase, listOrEnv, reshape, drop, use_dimnames, thread_buffer, split_dim, strict) } +FARR_lock_acquire <- function(lock_path, parts, scope_id) { + .Call(`_filearray_FARR_lock_acquire`, lock_path, parts, scope_id) +} + +FARR_lock_release_scope <- function(scope_id) { + .Call(`_filearray_FARR_lock_release_scope`, scope_id) +} + +FARR_lock_commit_scope <- function(scope_id) { + .Call(`_filearray_FARR_lock_commit_scope`, scope_id) +} + +FARR_lock_release_explicit <- function(lock_path, parts, all = FALSE) { + .Call(`_filearray_FARR_lock_release_explicit`, lock_path, parts, all) +} + +FARR_lock_forget <- function(lock_path) { + .Call(`_filearray_FARR_lock_forget`, lock_path) +} + +FARR_lock_release_all <- function() { + .Call(`_filearray_FARR_lock_release_all`) +} + +FARR_link_noreplace <- function(from, to) { + .Call(`_filearray_FARR_link_noreplace`, from, to) +} + +FARR_lock_status <- function(lock_path) { + .Call(`_filearray_FARR_lock_status`, lock_path) +} + +FARR_lock_registry <- function() { + .Call(`_filearray_FARR_lock_registry`) +} + FARR_buffer_map <- function(input_filebases, output_filebase, map, buffer_nelems, result_nelems = 0L) { .Call(`_filearray_FARR_buffer_map`, input_filebases, output_filebase, map, buffer_nelems, result_nelems) } diff --git a/R/bind.R b/R/bind.R index c2a9deb..7788b6c 100644 --- a/R/bind.R +++ b/R/bind.R @@ -194,7 +194,12 @@ filearray_bind <- function( unlink(filebase, recursive = TRUE) } re <- filearray_create(filebase = filebase, dimension = dim, type = type, partition_size = part_size) - + + # The partition files are copied or linked in below + lock_scope <- lock_scope_begin() + on.exit(lock_scope_end(lock_scope), add = TRUE) + lock_scope_acquire(lock_scope, re, re$.partition_info[, 1]) + start <- 1 end <- 1 diff --git a/R/class-filearray.R b/R/class-filearray.R index 0ac723a..49e3f5a 100644 --- a/R/class-filearray.R +++ b/R/class-filearray.R @@ -18,15 +18,74 @@ #' \item{\code{expand(n)}}{Expand array along the last margin; returns true if expanded; if the \code{dimnames} have been assigned prior to expansion, the last dimension names will be filled with \code{NA}} #' \item{\code{initialize_partition()}}{Make sure a partition file exists; if not, create one and fill with \code{NA}s or 0 (\code{type='raw'})} #' \item{\code{load(filebase, mode = c("readwrite", "readonly"))}}{Load file array from existing directory} +#' \item{\code{lock_partition(parts, timeout = NULL)}}{Lock partitions \code{parts} (default: all) for this \R session, so that other processes cannot write to them; waits and retries until all of them are locked or \code{timeout} seconds have passed (\code{NULL}: option \code{filearray.lock.timeout}, which defaults to waiting without limit; \code{0}: try once). Returns \code{TRUE} once locked, or \code{FALSE} after the time-out, with nothing locked; see section 'Partition locks'} #' \item{\code{partition_path(part)}}{Get partition file path} #' \item{\code{partition_size()}}{Get partition size; see \code{\link{filearray}}} #' \item{\code{set_partition(part, value, ..., strict = TRUE)}}{Set partition value} #' \item{\code{sexp_type()}}{Get data \code{SEXP} type; see R internal manuals} #' \item{\code{show()}}{Print information} #' \item{\code{type()}}{Get data type} +#' \item{\code{unlock_partition(parts)}}{Release one \code{lock_partition} lock on each of \code{parts}, or every \code{lock_partition} lock this \R session holds on the array if \code{parts} is missing; returns \code{TRUE} if a lock was released} #' \item{\code{valid()}}{Check if the array is valid.} #' } -#' @seealso \code{\link{filearray}} +#' @section Partition locks: +#' Several processes (for example \pkg{parallel}, \pkg{callr} or \pkg{mirai} +#' workers) can write to one array. Writes lock the partitions they change: +#' sub-assignment (\code{x[...] <- value}), \code{fill_partition}, +#' \code{initialize_partition}, \code{set_partition}, \code{\link{fmap}}, +#' \code{\link{fmap_element_wise}} and \code{\link{filearray_bind}} wait +#' while another process holds one of their partitions, so concurrent writers +#' never lose each other's data. Reading takes no locks, so a reader may see a +#' write that is still in progress. Functions that read but must first create +#' missing partitions (\code{max}, \code{min}, \code{range} and \code{sum} +#' with \code{na.rm = FALSE}, inputs of \code{fmap}, \code{filearray_bind}, +#' evaluating array operations) lock the partitions they create, so they can +#' wait too. +#' +#' To keep partitions to yourself over several operations, for example to +#' check whether a partition has been computed and write it only if not, lock +#' them explicitly: +#' \preformatted{ +#' if (!x$lock_partition(k, timeout = 60)) stop("partition ", k, " is busy") +#' on.exit(x$unlock_partition(k), add = TRUE) +#' if (is.na(x[1, 1, k])) x[, , k] <- compute(k) +#' } +#' Locks belong to the \R session, not to the object: every object of the +#' same array shares them, writes of this session to partitions it holds do +#' not wait, and each \code{lock_partition} call needs one +#' \code{unlock_partition} call. Lock everything you need in one call: two +#' sessions that each hold one partition and wait for the other's wait until +#' their time-outs. Likewise, do not hold partitions that forked workers +#' (\code{mclapply}) of the same session write. Non-interactive workers +#' should set a time-out. +#' +#' The locks are byte-range locks on an empty file \code{.filearray.lock} in +#' the array directory (byte \code{k} locks partition \code{k}). They are +#' advisory: older versions of this package and other programs ignore them. +#' The operating system releases them when a process exits, so a crashed or +#' killed process never leaves a stale lock. The locks guard writes but never +#' make them fail: when a lock cannot be taken for a reason other than another +#' process holding it (a file system without lock support such as some network +#' mounts, no permission to create the lock file, a network error), writes warn +#' once and proceed unlocked, and \code{lock_partition} warns and returns +#' \code{TRUE} with attribute \code{locked = FALSE}. Cloud-synchronized +#' folders lock on one computer only. Copying or archiving the array directory +#' in the session that holds its locks releases them on \verb{POSIX} systems. +#' To lock anything else the same way, see \code{\link{with_filelock}}. +#' +#' Options (each falls back to the environment variable, which worker +#' processes inherit): +#' \describe{ +#' \item{\code{filearray.lock} (\code{FILEARRAY_LOCK})}{set to \code{FALSE} +#' to stop writes from locking automatically; \code{lock_partition} still +#' locks} +#' \item{\code{filearray.lock.timeout} (\code{FILEARRAY_LOCK_TIMEOUT})}{seconds +#' a write waits for a locked partition before it fails with an error of class +#' \code{filearray_lock_timeout}; default \code{Inf}, wait until the +#' partition is free. A value that is not a non-negative number warns once and +#' is ignored} +#' } +#' @seealso \code{\link{filearray}}, \code{\link{with_filelock}} #' @exportClass FileArray NULL @@ -368,7 +427,12 @@ FileArray <- setRefClass( return(FALSE) } } - + + # Hold the partition from reading it to writing it back + lock_scope <- lock_scope_begin() + on.exit(lock_scope_end(lock_scope), add = TRUE) + lock_scope_acquire(lock_scope, .self, part) + arglen <- ...length() dim[[length(dim)]] <- .self$partition_size() if ( arglen > 1 ) { @@ -465,112 +529,43 @@ FileArray <- setRefClass( if (length(value) > 1) { quiet_warning("`fill_partition` value length coerced to first value") } - - value <- value[[1]] - - type <- .self$type() - switch( - type, - "complex" = { - value <- cplxToReal2(as.complex(value)) - }, - "float" = { - value <- realToFloat2(as.double(value)) - }, - { - storage.mode(value) <- type - } - ) - size <- get_elem_size(type) - file <- .self$partition_path(part) - - # Write to a temporary file, then rename it to the partition: - # another process that loads this array while the partition is - # being written would reject a file shorter than the 1024-byte - # header. `load()` lists only `.farr`, never `..tmp` - tmp <- sprintf("%s.tmp%s", gsub("farr$", "", file, ignore.case = TRUE), Sys.getpid()) - fid <- file(tmp, "wb") - fid_closed <- FALSE - - on.exit({ - # `tmp` still exists only if writing or renaming failed - if (file.exists(tmp)) { - if (!fid_closed) { - close(fid) - } - unlink(tmp) - } - }, add = TRUE) - - if ( part <= nrow(.self$.partition_info)) { - partition_size <- .self$.partition_info[part, 2] - } else { - partition_size <- .self$partition_size() - } - dimension <- .self$dimension() - dimension[[length(dimension)]] <- partition_size - part_len <- prod(dimension) - buffer_len <- get_buffer_size() / size - if (buffer_len > part_len) { - buffer_len <- part_len - } - - write_header( - fid = fid, - partition = part, - dimension = dimension, - type = type, - size = size - ) - seek(con = fid, where = HEADER_SIZE, rw = "write") - buf <- writeBin( - con = raw(), - object = rep(value, buffer_len), - size = size, - endian = ENDIANNESS - ) - nloop <- floor(part_len / buffer_len) - replicate(nloop, { - writeBin(con = fid, object = buf) - NULL - }) - rest <- part_len - buffer_len * nloop - if ( rest > 0 ) { - writeBin(con = fid, object = buf[seq_len(rest * size)]) - } - seek(con = fid, where = HEADER_SIZE - 8L, rw = "write") - writeBin(con = fid, object = part_len, size = 8L, endian = ENDIANNESS) - close(fid) - fid_closed <- TRUE - - # rename() replaces `file` in one step, so other processes see - # either the old partition or the new one, never a partial one - if (!file.rename(tmp, file)) { - stop("Cannot write partition file: ", file) - } + lock_scope <- lock_scope_begin() + on.exit(lock_scope_end(lock_scope), add = TRUE) + locked <- lock_scope_acquire(lock_scope, .self, part) + write_partition_file(.self, part, value, replace = TRUE, locked = locked) invisible() }, initialize_partition = function(parts) { if (!.self$valid()) { stop("Invalid file array") } - if (isTRUE(.self$.mode == "readonly")) { - .self$.mode <- "readwrite" - on.exit({ - .self$.mode <- "readonly" - }) - } if (missing(parts)) { parts <- .self$.partition_info[, 1] } parts <- parts[!is.na(parts)] + # Read APIs (max, min, range, sum, filearray_bind, ...) call this: + # when every partition exists, they return here without locking + parts <- parts[!file.exists(.self$partition_path(parts))] + if (!length(parts)) { + return(invisible()) + } + parts <- sort(unique(as.integer(parts))) + if (any(parts <= 0)) { + stop("NA or non-positive partition are invalid") + } + + lock_scope <- lock_scope_begin() + on.exit(lock_scope_end(lock_scope), add = TRUE) + locked <- lock_scope_acquire(lock_scope, .self, parts) for (part in parts) { - file <- .self$partition_path(part) - if (!file.exists(file)) { - .self$fill_partition(part, NA) + # Another process may have created the partition meanwhile; + # write_partition_file() never replaces one that exists + if (!file.exists(.self$partition_path(part))) { + write_partition_file(.self, part, NA, replace = FALSE, locked = locked) } } + invisible() }, can_write = function() { if (isTRUE(.self$.mode == "readonly")) { @@ -578,6 +573,35 @@ FileArray <- setRefClass( } return(TRUE) }, + lock_partition = function(parts, timeout = NULL) { + if (!.self$valid()) { + stop("Invalid file array") + } + nparts <- nrow(.self$.partition_info) + if (missing(parts)) { + parts <- seq_len(nparts) + } + parts <- validate_lock_parts(parts, nparts) + if (is.null(timeout)) { + timeout <- lock_timeout_default() + } else { + timeout <- validate_lock_timeout(timeout) + } + locked <- lock_partition_explicit(.self$.filebase, parts, timeout) + if (isTRUE(locked)) { + return(invisible(locked)) + } + locked + }, + unlock_partition = function(parts) { + if (missing(parts)) { + released <- lock_release_explicit(.self$.filebase) + } else { + released <- lock_release_explicit( + .self$.filebase, validate_lock_parts(parts)) + } + invisible(released > 0) + }, delete = function(force = FALSE) { if (!.self$valid()) { .self$.valid <- FALSE @@ -591,6 +615,9 @@ FileArray <- setRefClass( } } filebase <- .self$.filebase + # Drop this session's partition locks and close the lock file: + # on Windows an open file keeps the directory from being removed + lock_forget(filebase) if (dir.exists(filebase)) { unlink(filebase, recursive = TRUE, force = force) } diff --git a/R/filearray-package.R b/R/filearray-package.R index 501bb41..886f774 100644 --- a/R/filearray-package.R +++ b/R/filearray-package.R @@ -98,6 +98,12 @@ symlink_enabled <- local({ ns$NA_float_ <- get_float_na() } +.onUnload <- function(libpath) { + # A reloaded package starts a new lock registry, which could not see the + # locks of this one: release them all + try(FARR_lock_release_all(), silent = TRUE) +} + .onAttach <- function(libname, pkgname) { if (Sys.getenv("_R_CHECK_LIMIT_CORES_") == "TRUE") { packageStartupMessage( diff --git a/R/lock.R b/R/lock.R index e69de29..e96461f 100644 --- a/R/lock.R +++ b/R/lock.R @@ -0,0 +1,423 @@ +# File locks: partition locks and with_filelock() +# +# Protocol v1. This is a contract between filearray versions; src/lock.h +# documents the same: +# +# * A lock file is an empty file, created by the first process that locks. +# Nothing writes to it, it is never renamed or removed, and no other code +# opens it. Each array has one: "/.filearray.lock". +# * Byte 0 is the lock that with_filelock() takes. On an array's lock file it +# is the array-level lock, which partition writers do not take. +# * Byte k >= 1 of an array's lock file is the exclusive write lock of +# partition k (file ".farr"). +# * POSIX: classic fcntl record lock F_WRLCK on [k, k + 1), taken with +# non-blocking F_SETLK only. Windows: LockFileEx at offset k, length 1. +# * Locks are exclusive; of an array, only writers lock. The locks are +# advisory, and the operating system releases them when a process exits. +# +# Locks belong to the R session: the C++ registry counts holds per byte, so +# nested writes and several objects for one array share a lock. A write holds +# its locks in a scope, which C++ records by id, so that `on.exit()` releases +# exactly what was taken even after an interrupt. +# +# Locks guard; they never make the work they guard fail. When a lock cannot +# be taken for a reason other than another process holding it, writes warn +# once and proceed without it. + +LOCK_FILE_NAME <- ".filearray.lock" +LOCK_PROTOCOL_VERSION <- 1L + +# Status of FARR_lock_acquire(); keep in sync with src/lock.h +LOCK_STATUS_OK <- 0L +LOCK_STATUS_BUSY <- 1L +LOCK_STATUS_UNSUPPORTED <- 2L + +# Status of FARR_link_noreplace(); keep in sync with src/lock.h +LINK_CREATED <- 0L +LINK_EXISTS <- 1L +LINK_UNSUPPORTED <- 2L + +# Seconds before locking is tried again on a lock file that could not lock, +# so that one transient network error does not end locking for the session +LOCK_RETRY_AFTER <- 60 + +# Temporary files of killed writers older than this (seconds) are removed +STALE_TMP_AGE <- 3600 + +lock_state <- new.env(parent = emptyenv()) +lock_state$scope_id <- 0 +# lock path -> list(time, message) of the last "unsupported" failure +lock_state$unsupported <- new.env(parent = emptyenv()) +# lock path -> TRUE once warned +lock_state$warned <- new.env(parent = emptyenv()) + +lock_file_path <- function(filebase) { + file.path(filebase, LOCK_FILE_NAME) +} + +# C++ takes UTF-8 on Windows (and converts it to UTF-16), native bytes elsewhere +lock_os_path <- function(path) { + if (get_os() == "windows") { + enc2utf8(path) + } else { + enc2native(path) + } +} + +# Whether writes lock automatically +lock_enabled <- function() { + value <- getOption("filearray.lock", NULL) + if (is.null(value)) { + # callr workers inherit environment variables, not options + value <- Sys.getenv("FILEARRAY_LOCK", unset = "TRUE") + } + value <- toupper(trimws(as.character(value)[1])) + !isTRUE(value %in% c("FALSE", "F", "0", "NO", "OFF")) +} + +validate_lock_timeout <- function(timeout, what = "timeout") { + value <- suppressWarnings(as.numeric(timeout)) + if (length(value) != 1L || is.na(value) || value < 0) { + stop(sprintf("`%s` must be a non-negative number of seconds or `Inf`", what), + call. = FALSE) + } + value +} + +# Seconds to wait for locks: `Inf` waits until they are free, `0` tries once. +# A bad option or environment variable warns once and waits without limit: +# it must not stop every write +lock_timeout_default <- function() { + value <- getOption("filearray.lock.timeout", NULL) + what <- "filearray.lock.timeout" + if (is.null(value)) { + value <- Sys.getenv("FILEARRAY_LOCK_TIMEOUT", unset = "") + if (!nzchar(value)) { + return(Inf) + } + what <- "FILEARRAY_LOCK_TIMEOUT" + } + timeout <- suppressWarnings(as.numeric(value)) + if (length(timeout) != 1L || is.na(timeout) || timeout < 0) { + shown <- deparse1(value) + lock_warn_once(sprintf("%s=%s", what, shown), quiet_warning(sprintf(paste( + "`%s` must be a non-negative number of seconds or `Inf`;", + "ignoring %s: locks are waited for without a time limit" + ), what, shown))) + return(Inf) + } + timeout +} + +# Sorted unique partition numbers, each from 1 to `nparts` +validate_lock_parts <- function(parts, nparts = NA) { + if (!length(parts)) { + return(numeric(0)) + } + if (!is.numeric(parts) || anyNA(parts) || any(parts < 1) || + any(parts != round(parts)) || (!is.na(nparts) && any(parts > nparts))) { + msg <- "`parts` must be partition numbers" + if (!is.na(nparts)) { + msg <- sprintf("%s from 1 to %d", msg, as.integer(nparts)) + } + stop(msg, call. = FALSE) + } + sort(unique(as.numeric(parts))) +} + +# Partitions holding the last-margin `slices`. C++ truncates fractional +# indices, so they are truncated here too +partitions_of_slices <- function(x, slices) { + cum <- x$.partition_info[, 3] + sort(unique(findInterval(trunc(slices), cum, left.open = TRUE) + 1)) +} + +lock_unsupported_since <- function(path) { + failure <- lock_state$unsupported[[path]] + if (is.null(failure) || + proc.time()[["elapsed"]] - failure$time >= LOCK_RETRY_AFTER) { + return(NULL) + } + failure +} + +# Records that `path` cannot be locked; callers decide how to report it +lock_mark_unsupported <- function(path, message) { + assign(path, list(time = proc.time()[["elapsed"]], message = message), + envir = lock_state$unsupported) + invisible() +} + +# Evaluates `warning_expr` (a warning) only the first time per session and key +lock_warn_once <- function(key, warning_expr) { + if (!isTRUE(lock_state$warned[[key]])) { + assign(key, TRUE, envir = lock_state$warned) + force(warning_expr) + } + invisible() +} + +# A scope groups the locks of one operation; C++ records them by its id +lock_scope_begin <- function() { + lock_state$scope_id <- lock_state$scope_id + 1 + lock_state$scope_id +} + +lock_scope_end <- function(scope) { + suspendInterrupts(tryCatch( + FARR_lock_release_scope(scope), + error = function(e) { 0 } + )) + invisible() +} + +# Locks `bytes` (sorted) of `lockfile` under `scope`, waiting while other +# processes hold them. Returns a list with `status`: +# * "acquired"; +# * "timeout", with the bytes still `pending`; +# * "unsupported", with a `message`: the lock cannot be taken for a reason +# other than another process holding it (unsupported file system, no +# permission, stale network handle, ...) +lock_path_wait <- function(scope, lockfile, bytes, timeout) { + failure <- lock_unsupported_since(lockfile) + if (!is.null(failure)) { + return(list(status = "unsupported", pending = bytes, message = failure$message)) + } + os_path <- lock_os_path(lockfile) + deadline <- proc.time()[["elapsed"]] + timeout + delay <- 0.001 + repeat { + res <- FARR_lock_acquire(os_path, bytes, scope) + if (res$status == LOCK_STATUS_OK) { + return(list(status = "acquired", pending = numeric(0))) + } + if (res$status != LOCK_STATUS_BUSY) { + lock_mark_unsupported(lockfile, res$message) + return(list(status = "unsupported", pending = bytes, message = res$message)) + } + # the leading bytes are held while waiting for the next one + if (res$acquired > 0) { + bytes <- bytes[-seq_len(res$acquired)] + } + remaining <- deadline - proc.time()[["elapsed"]] + if (remaining <= 0) { + return(list(status = "timeout", pending = bytes)) + } + Sys.sleep(min(delay, remaining)) + delay <- min(delay * 2, 0.1) + } +} + +# lock_path_wait() for partitions of the array at `filebase` +lock_scope_wait <- function(scope, filebase, parts, timeout) { + lock_path_wait(scope, lock_file_path(filebase), parts, timeout) +} + +lock_timeout_error <- function(filebase, pending, timeout) { + msg <- sprintf(paste( + "Timed out after %s second(s) waiting to lock partition %s of file", + "array '%s': another process is writing to it, or holds it with", + "`lock_partition()`. Set `options(filearray.lock.timeout)` to wait longer." + ), format(timeout), format(pending[[1]]), filebase) + structure( + class = c("filearray_lock_timeout", "error", "condition"), + list(message = msg, call = NULL, filebase = filebase, + partitions = pending, timeout = timeout) + ) +} + +# Locks partitions for a write until `lock_scope_end(scope)`. Returns TRUE if +# they are locked; FALSE if automatic locking is off or the file system cannot +# lock. Signals a `filearray_lock_timeout` error on time-out +lock_scope_acquire <- function(scope, x, parts, timeout = NULL) { + if (!length(parts) || !lock_enabled()) { + return(invisible(FALSE)) + } + if (is.null(timeout)) { + timeout <- lock_timeout_default() + } + parts <- sort(unique(as.numeric(parts))) + res <- lock_scope_wait(scope, x$.filebase, parts, timeout) + if (identical(res$status, "timeout")) { + stop(lock_timeout_error(x$.filebase, res$pending, timeout)) + } + if (identical(res$status, "unsupported")) { + lock_warn_once(lock_file_path(x$.filebase), quiet_warning(sprintf(paste( + "Cannot lock partitions of file array '%s' (%s).", + "This session writes to it without partition locks: other", + "processes writing to the same partitions are not excluded." + ), x$.filebase, res$message))) + } + invisible(identical(res$status, "acquired")) +} + +# Locks for `lock_partition()`: TRUE when locked; TRUE with attribute +# `locked = FALSE` (and a warning) if the file system cannot lock; FALSE on +# time-out, with nothing locked +lock_partition_explicit <- function(filebase, parts, timeout) { + if (!length(parts)) { + return(TRUE) + } + scope <- lock_scope_begin() + on.exit(lock_scope_end(scope), add = TRUE) + res <- lock_scope_wait(scope, filebase, parts, timeout) + if (identical(res$status, "timeout")) { + return(FALSE) + } + if (identical(res$status, "unsupported")) { + warning(sprintf("Partitions of file array '%s' are not locked: %s", + filebase, res$message), call. = FALSE) + return(structure(TRUE, locked = FALSE)) + } + suspendInterrupts(FARR_lock_commit_scope(scope)) + TRUE +} + +# Releases `lock_partition()` locks: one per partition, or all of them +lock_release_explicit <- function(filebase, parts = NULL) { + os_path <- lock_os_path(lock_file_path(filebase)) + if (is.null(parts)) { + FARR_lock_release_explicit(os_path, numeric(0), TRUE) + } else { + FARR_lock_release_explicit(os_path, parts, FALSE) + } +} + +# Drops every lock this session holds on an array and closes its lock file +lock_forget <- function(filebase) { + tryCatch( + FARR_lock_forget(lock_os_path(lock_file_path(filebase))), + error = function(e) { 0 } + ) +} + +# Locks this session holds on the array of `x` +lock_status <- function(x) { + FARR_lock_status(lock_os_path(lock_file_path(x$.filebase))) +} + +validate_lockfile <- function(lockfile) { + if (!is.character(lockfile) || length(lockfile) != 1L || + is.na(lockfile) || !nzchar(lockfile)) { + stop("`lockfile` must be the path of a lock file: one non-empty string", + call. = FALSE) + } + # absolute, so that the path names the same file if the working + # directory changes while the lock is held + lockfile <- path.expand(lockfile) + file.path(normalizePath(dirname(lockfile), mustWork = FALSE), basename(lockfile)) +} + +filelock_timeout_error <- function(lockfile, timeout) { + msg <- sprintf(paste( + "Timed out after %s second(s) waiting to lock '%s': another R", + "session holds it. Increase `timeout` to wait longer." + ), format(timeout), lockfile) + structure( + class = c("filearray_lock_timeout", "error", "condition"), + list(message = msg, call = NULL, lockfile = lockfile, timeout = timeout) + ) +} + +filelock_unsupported <- function(lockfile, reason, type = c("warning", "error")) { + type <- match.arg(type) + msg <- sprintf("Cannot lock '%s' (%s)", lockfile, reason) + if (type == "warning") { + msg <- paste0(msg, ": running without the lock") + } + structure( + class = c("filearray_lock_unsupported", type, "condition"), + list(message = msg, call = NULL, lockfile = lockfile, reason = reason) + ) +} + +#' @title Evaluate an expression while holding a file lock +#' @description Runs \code{expr} while holding an exclusive lock on +#' \code{lockfile}, so that only one \R session at a time runs code under +#' that lock; other sessions asking for it wait until it is released. The +#' lock has nothing to do with file arrays: any package or script can use it. +#' @param expr expression to evaluate while the lock is held; it is evaluated +#' in the calling environment +#' @param lockfile path of the lock file, a file used only for locking. It is +#' created if missing (its directory must exist), never written to, and left +#' in place afterwards +#' @param timeout seconds to wait for another session to release the lock: +#' \code{Inf} (default) waits until the lock is free, \code{0} tries once. +#' After the time-out an error of class \code{filearray_lock_timeout} is +#' raised and \code{expr} is not evaluated +#' @param on_unsupported what to do when the lock cannot be taken for a +#' reason other than another session holding it, for example a file system +#' without lock support or a lock file that cannot be created: +#' \code{"warn"} (default) warns once per lock file and session, then +#' evaluates \code{expr} without the lock; \code{"error"} raises an error and +#' does not evaluate \code{expr}; \code{"ignore"} evaluates \code{expr} +#' without the lock and without a warning. The warning and the error have +#' class \code{filearray_lock_unsupported} +#' @return The value of \code{expr} +#' @details The lock is an operating-system lock on the first byte of +#' \code{lockfile} (\code{fcntl} on \verb{POSIX} systems, \code{LockFileEx} +#' on 'Windows'), so it is released when a session ends, even when it +#' crashes. It is the lock that \code{lock()} of the \pkg{filelock} package +#' takes in exclusive mode, so the two exclude each other between sessions. +#' +#' Locks belong to the \R session: a nested \code{with_filelock()} call on the +#' same lock file in the same session runs at once instead of waiting for +#' itself. Code running in the same session is therefore never excluded, even +#' when it serves different users, as several 'Shiny' sessions in one \R +#' process do. +#' +#' Use a dedicated lock file, and do not open it in \code{expr}: on +#' \verb{POSIX} systems closing any connection to the lock file releases the +#' lock, and on 'Windows' the locked byte cannot be read or written. +#' +#' The lock is advisory: it only excludes code that locks the same file. +#' Forked workers (such as those of \code{mclapply} in \pkg{parallel}) do not +#' inherit it. Sessions that nest locks on two files in opposite orders wait for each +#' other until their time-outs. On network file systems the lock works across +#' computers only if the file system supports locks; folders synchronized by +#' cloud services are locked on one computer only. +#' @examples +#' +#' lockfile <- tempfile(fileext = ".lock") +#' log <- tempfile() +#' +#' # one R session at a time appends to the log +#' with_filelock({ +#' cat("one line\n", file = log, append = TRUE) +#' }, lockfile = lockfile) +#' +#' readLines(log) +#' +#' unlink(c(lockfile, log)) +#' +#' @export +with_filelock <- function(expr, lockfile, timeout = Inf, + on_unsupported = c("warn", "error", "ignore")) { + on_unsupported <- match.arg(on_unsupported) + lockfile <- validate_lockfile(lockfile) + timeout <- validate_lock_timeout(timeout) + + scope <- lock_scope_begin() + on.exit(lock_scope_end(scope), add = TRUE) + res <- lock_path_wait(scope, lockfile, 0, timeout) + if (identical(res$status, "timeout")) { + stop(filelock_timeout_error(lockfile, timeout)) + } + if (identical(res$status, "unsupported")) { + switch( + on_unsupported, + "error" = stop(filelock_unsupported(lockfile, res$message, "error")), + "warn" = lock_warn_once( + paste0("with_filelock:", lockfile), + warning(filelock_unsupported(lockfile, res$message, "warning")) + ) + ) + } + + res <- withVisible(expr) + if (res$visible) { + res$value + } else { + invisible(res$value) + } +} diff --git a/R/method-map.R b/R/method-map.R index 37b7e02..179966c 100644 --- a/R/method-map.R +++ b/R/method-map.R @@ -219,8 +219,12 @@ filearray_map_internal <- function(x, fun, y = NULL, .buffer_count = NA_integer_ stop("Input & output arrays have inconsistent lengths. Cannot map functions") } buffer_nelems <- all_lens / .buffer_count - - + + # Hold every partition of `y` while mapping into it; `fun` may write + # `y` itself, since this session's locks are re-entrant + lock_scope <- lock_scope_begin() + on.exit(lock_scope_end(lock_scope), add = TRUE) + lock_scope_acquire(lock_scope, y, y$.partition_info[, 1]) y$initialize_partition() FARR_buffer_map( @@ -301,8 +305,6 @@ fmap_element_wise_internal <- function(x, fun, .y, ..., .input_size = NA) { stop("Dimensions of x[[1]] and .y mismatch") } } - .y$initialize_partition() - if (is.na(.input_size)) { .input_size <- guess_fmap_buffer_size(dim(.y), .y$element_size()) # .input_size <- get_buffer_size() / .y$element_size() @@ -328,6 +330,13 @@ fmap_element_wise_internal <- function(x, fun, .y, ..., .input_size = NA) { if (length(x[[1]]) > length(.y) ) { stop("Inconsistent input and output length") } + + # Hold every partition of `.y` while mapping into it. The inputs were + # initialized above, so this never waits for them while holding `.y` + lock_scope <- lock_scope_begin() + on.exit(lock_scope_end(lock_scope), add = TRUE) + lock_scope_acquire(lock_scope, .y, .y$.partition_info[, 1]) + .y$initialize_partition() FARR_buffer_map( input_filebases = fbases, diff --git a/R/methods-subsetAssign.R b/R/methods-subsetAssign.R index 902af6b..b30d1f6 100644 --- a/R/methods-subsetAssign.R +++ b/R/methods-subsetAssign.R @@ -28,7 +28,7 @@ fa_subsetAssign1 <- function(x, ..., value) { stop("SubsetAssign FileArray `value` length mismatch: `value` length must be either 1 or the same length of the subset.") } target_dim <- dim - x$initialize_partition(x$.partition_info[, 1]) + parts <- x$.partition_info[, 1] } else if (arglen > 1) { if (arglen != length(dim)) { stop("SubsetAssign FileArray dimension mismatch.") @@ -66,25 +66,21 @@ fa_subsetAssign1 <- function(x, ..., value) { stop("SubsetAssign FileArray `value` length mismatch: `value` length must be either 1 or the same length of the subset.") } - # make sure partitions exist - tmp <- locs[[length(locs)]] - sapply(tmp, function(i) { - sel <- x$.partition_info[, 3] <= i - if (any(sel)) { - sel <- max(x$.partition_info[sel, 1]) - if (x$.partition_info[sel, 3] < i) { - sel <- sel + 1 - } - } else { - sel <- 1 - } - x$initialize_partition(sel) - }) + # partitions holding the selected slices of the last margin + parts <- partitions_of_slices(x, locs[[length(locs)]]) } if (prod(target_dim) == 0) { return(invisible(x)) } + + # Hold the partitions written below until this function returns; the + # nested initialize_partition() and fill_partition() reuse these locks + lock_scope <- lock_scope_begin() + on.exit(lock_scope_end(lock_scope), add = TRUE) + lock_scope_acquire(lock_scope, x, parts) + # make sure partitions exist + x$initialize_partition(parts) # decide split_dim buffer_sz <- buf_bytes / x$element_size() diff --git a/R/write.R b/R/write.R index de6daaa..a0814ab 100644 --- a/R/write.R +++ b/R/write.R @@ -35,6 +35,155 @@ ensure_partition <- function( } +# Writes partition `part` of `x` filled with `value` through a temporary file. +# `replace = TRUE` replaces an existing partition; `replace = FALSE` only +# creates a missing one. `locked`: whether this session holds the partition +# lock, so that temporary files of killed writers can be removed +write_partition_file <- function(x, part, value, replace = TRUE, locked = FALSE) { + value <- value[[1]] + type <- x$type() + switch( + type, + "complex" = { + value <- cplxToReal2(as.complex(value)) + }, + "float" = { + value <- realToFloat2(as.double(value)) + }, + { + storage.mode(value) <- type + } + ) + size <- get_elem_size(type) + file <- x$partition_path(part) + if (locked) { + remove_stale_partition_tmp(x, part) + } + + # Write to a temporary file, then publish it as the partition: another + # process that loads this array while the partition is being written + # would reject a file shorter than the 1024-byte header. The name is new + # for every call, so writers never share a temporary file; `load()` lists + # only `.farr` + tmp <- file.path(x$.filebase, sprintf("%d.%s.tmp", as.integer(part) - 1L, + uuid::UUIDgenerate())) + fid <- file(tmp, "wb") + fid_closed <- FALSE + + on.exit({ + if (!fid_closed) { + close(fid) + } + # `tmp` is left if writing or publishing failed, or if the partition + # already existed + if (file.exists(tmp)) { + unlink(tmp) + } + }, add = TRUE) + + if (part <= nrow(x$.partition_info)) { + partition_size <- x$.partition_info[part, 2] + } else { + partition_size <- x$partition_size() + } + dimension <- x$dimension() + dimension[[length(dimension)]] <- partition_size + part_len <- prod(dimension) + buffer_len <- get_buffer_size() / size + if (buffer_len > part_len) { + buffer_len <- part_len + } + + write_header( + fid = fid, + partition = part, + dimension = dimension, + type = type, + size = size + ) + seek(con = fid, where = HEADER_SIZE, rw = "write") + buf <- writeBin( + con = raw(), + object = rep(value, buffer_len), + size = size, + endian = ENDIANNESS + ) + nloop <- floor(part_len / buffer_len) + replicate(nloop, { + writeBin(con = fid, object = buf) + NULL + }) + rest <- part_len - buffer_len * nloop + if ( rest > 0 ) { + writeBin(con = fid, object = buf[seq_len(rest * size)]) + } + seek(con = fid, where = HEADER_SIZE - 8L, rw = "write") + writeBin(con = fid, object = part_len, size = 8L, endian = ENDIANNESS) + + close(fid) + fid_closed <- TRUE + + if (replace) { + # rename() replaces `file` in one step, so other processes see either + # the old partition or the new one, never a partial one + if (!file.rename(tmp, file)) { + stop("Cannot write partition file: ", file) + } + return(invisible(TRUE)) + } + invisible(publish_partition_file(tmp, file)) +} + +# Publishes `tmp` as the partition file `file` unless that exists. A hard link +# is refused by the file system (on NFS, by the server) when `file` exists, +# so a partition another process has created and written is never replaced, +# even when this host still caches it as missing. Returns FALSE if `file` +# already existed +publish_partition_file <- function(tmp, file) { + status <- FARR_link_noreplace(lock_os_path(tmp), lock_os_path(file)) + code <- as.integer(status) + if (identical(code, LINK_CREATED)) { + unlink(tmp) + return(TRUE) + } + if (identical(code, LINK_EXISTS)) { + return(FALSE) + } + if (!identical(code, LINK_UNSUPPORTED)) { + stop("Cannot write partition file: ", file, " (", attr(status, "message"), ")") + } + # No hard links on this file system: rename, unless the partition + # appeared meanwhile + if (file.exists(file)) { + return(FALSE) + } + if (!file.rename(tmp, file)) { + stop("Cannot write partition file: ", file) + } + TRUE +} + +# Removes temporary files that writers of partition `part` left behind when +# they were killed. Writers that do not lock are never this slow, so their +# temporary files are safe +remove_stale_partition_tmp <- function(x, part, age = STALE_TMP_AGE) { + pattern <- sprintf( + "^%d\\.[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}\\.tmp$", + as.integer(part) - 1L + ) + files <- list.files(x$.filebase, pattern = pattern, full.names = TRUE) + if (!length(files)) { + return(invisible()) + } + mtime <- file.mtime(files) + stale <- files[!is.na(mtime) & + difftime(Sys.time(), mtime, units = "secs") > age] + if (length(stale)) { + unlink(stale) + } + invisible() +} + sexp_to_type <- function(sexp) { switch( as.character(sexp), diff --git a/man/FileArray-class.Rd b/man/FileArray-class.Rd index 2750781..86c7654 100644 --- a/man/FileArray-class.Rd +++ b/man/FileArray-class.Rd @@ -24,16 +24,78 @@ to create instances. \item{\code{expand(n)}}{Expand array along the last margin; returns true if expanded; if the \code{dimnames} have been assigned prior to expansion, the last dimension names will be filled with \code{NA}} \item{\code{initialize_partition()}}{Make sure a partition file exists; if not, create one and fill with \code{NA}s or 0 (\code{type='raw'})} \item{\code{load(filebase, mode = c("readwrite", "readonly"))}}{Load file array from existing directory} +\item{\code{lock_partition(parts, timeout = NULL)}}{Lock partitions \code{parts} (default: all) for this \R session, so that other processes cannot write to them; waits and retries until all of them are locked or \code{timeout} seconds have passed (\code{NULL}: option \code{filearray.lock.timeout}, which defaults to waiting without limit; \code{0}: try once). Returns \code{TRUE} once locked, or \code{FALSE} after the time-out, with nothing locked; see section 'Partition locks'} \item{\code{partition_path(part)}}{Get partition file path} \item{\code{partition_size()}}{Get partition size; see \code{\link{filearray}}} \item{\code{set_partition(part, value, ..., strict = TRUE)}}{Set partition value} \item{\code{sexp_type()}}{Get data \code{SEXP} type; see R internal manuals} \item{\code{show()}}{Print information} \item{\code{type()}}{Get data type} +\item{\code{unlock_partition(parts)}}{Release one \code{lock_partition} lock on each of \code{parts}, or every \code{lock_partition} lock this \R session holds on the array if \code{parts} is missing; returns \code{TRUE} if a lock was released} \item{\code{valid()}}{Check if the array is valid.} } } +\section{Partition locks}{ + +Several processes (for example \pkg{parallel}, \pkg{callr} or \pkg{mirai} +workers) can write to one array. Writes lock the partitions they change: +sub-assignment (\code{x[...] <- value}), \code{fill_partition}, +\code{initialize_partition}, \code{set_partition}, \code{\link{fmap}}, +\code{\link{fmap_element_wise}} and \code{\link{filearray_bind}} wait +while another process holds one of their partitions, so concurrent writers +never lose each other's data. Reading takes no locks, so a reader may see a +write that is still in progress. Functions that read but must first create +missing partitions (\code{max}, \code{min}, \code{range} and \code{sum} +with \code{na.rm = FALSE}, inputs of \code{fmap}, \code{filearray_bind}, +evaluating array operations) lock the partitions they create, so they can +wait too. + +To keep partitions to yourself over several operations, for example to +check whether a partition has been computed and write it only if not, lock +them explicitly: +\preformatted{ +if (!x$lock_partition(k, timeout = 60)) stop("partition ", k, " is busy") +on.exit(x$unlock_partition(k), add = TRUE) +if (is.na(x[1, 1, k])) x[, , k] <- compute(k) +} +Locks belong to the \R session, not to the object: every object of the +same array shares them, writes of this session to partitions it holds do +not wait, and each \code{lock_partition} call needs one +\code{unlock_partition} call. Lock everything you need in one call: two +sessions that each hold one partition and wait for the other's wait until +their time-outs. Likewise, do not hold partitions that forked workers +(\code{mclapply}) of the same session write. Non-interactive workers +should set a time-out. + +The locks are byte-range locks on an empty file \code{.filearray.lock} in +the array directory (byte \code{k} locks partition \code{k}). They are +advisory: older versions of this package and other programs ignore them. +The operating system releases them when a process exits, so a crashed or +killed process never leaves a stale lock. The locks guard writes but never +make them fail: when a lock cannot be taken for a reason other than another +process holding it (a file system without lock support such as some network +mounts, no permission to create the lock file, a network error), writes warn +once and proceed unlocked, and \code{lock_partition} warns and returns +\code{TRUE} with attribute \code{locked = FALSE}. Cloud-synchronized +folders lock on one computer only. Copying or archiving the array directory +in the session that holds its locks releases them on \verb{POSIX} systems. +To lock anything else the same way, see \code{\link{with_filelock}}. + +Options (each falls back to the environment variable, which worker +processes inherit): +\describe{ +\item{\code{filearray.lock} (\code{FILEARRAY_LOCK})}{set to \code{FALSE} +to stop writes from locking automatically; \code{lock_partition} still +locks} +\item{\code{filearray.lock.timeout} (\code{FILEARRAY_LOCK_TIMEOUT})}{seconds +a write waits for a locked partition before it fails with an error of class +\code{filearray_lock_timeout}; default \code{Inf}, wait until the +partition is free. A value that is not a non-negative number warns once and +is ignored} +} +} + \seealso{ -\code{\link{filearray}} +\code{\link{filearray}}, \code{\link{with_filelock}} } diff --git a/man/with_filelock.Rd b/man/with_filelock.Rd new file mode 100644 index 0000000..dab6379 --- /dev/null +++ b/man/with_filelock.Rd @@ -0,0 +1,83 @@ +% Generated by roxygen2: do not edit by hand +% Please edit documentation in R/lock.R +\name{with_filelock} +\alias{with_filelock} +\title{Evaluate an expression while holding a file lock} +\usage{ +with_filelock( + expr, + lockfile, + timeout = Inf, + on_unsupported = c("warn", "error", "ignore") +) +} +\arguments{ +\item{expr}{expression to evaluate while the lock is held; it is evaluated +in the calling environment} + +\item{lockfile}{path of the lock file, a file used only for locking. It is +created if missing (its directory must exist), never written to, and left +in place afterwards} + +\item{timeout}{seconds to wait for another session to release the lock: +\code{Inf} (default) waits until the lock is free, \code{0} tries once. +After the time-out an error of class \code{filearray_lock_timeout} is +raised and \code{expr} is not evaluated} + +\item{on_unsupported}{what to do when the lock cannot be taken for a +reason other than another session holding it, for example a file system +without lock support or a lock file that cannot be created: +\code{"warn"} (default) warns once per lock file and session, then +evaluates \code{expr} without the lock; \code{"error"} raises an error and +does not evaluate \code{expr}; \code{"ignore"} evaluates \code{expr} +without the lock and without a warning. The warning and the error have +class \code{filearray_lock_unsupported}} +} +\value{ +The value of \code{expr} +} +\description{ +Runs \code{expr} while holding an exclusive lock on +\code{lockfile}, so that only one \R session at a time runs code under +that lock; other sessions asking for it wait until it is released. The +lock has nothing to do with file arrays: any package or script can use it. +} +\details{ +The lock is an operating-system lock on the first byte of +\code{lockfile} (\code{fcntl} on \verb{POSIX} systems, \code{LockFileEx} +on 'Windows'), so it is released when a session ends, even when it +crashes. It is the lock that \code{lock()} of the \pkg{filelock} package +takes in exclusive mode, so the two exclude each other between sessions. + +Locks belong to the \R session: a nested \code{with_filelock()} call on the +same lock file in the same session runs at once instead of waiting for +itself. Code running in the same session is therefore never excluded, even +when it serves different users, as several 'Shiny' sessions in one \R +process do. + +Use a dedicated lock file, and do not open it in \code{expr}: on +\verb{POSIX} systems closing any connection to the lock file releases the +lock, and on 'Windows' the locked byte cannot be read or written. + +The lock is advisory: it only excludes code that locks the same file. +Forked workers (such as those of \code{mclapply} in \pkg{parallel}) do not +inherit it. Sessions that nest locks on two files in opposite orders wait for each +other until their time-outs. On network file systems the lock works across +computers only if the file system supports locks; folders synchronized by +cloud services are locked on one computer only. +} +\examples{ + +lockfile <- tempfile(fileext = ".lock") +log <- tempfile() + +# one R session at a time appends to the log +with_filelock({ + cat("one line\n", file = log, append = TRUE) +}, lockfile = lockfile) + +readLines(log) + +unlink(c(lockfile, log)) + +} diff --git a/src/RcppExports.cpp b/src/RcppExports.cpp index 6ae7d98..59081db 100644 --- a/src/RcppExports.cpp +++ b/src/RcppExports.cpp @@ -259,6 +259,108 @@ BEGIN_RCPP return rcpp_result_gen; END_RCPP } +// FARR_lock_acquire +List FARR_lock_acquire(const std::string& lock_path, const NumericVector& parts, const double scope_id); +RcppExport SEXP _filearray_FARR_lock_acquire(SEXP lock_pathSEXP, SEXP partsSEXP, SEXP scope_idSEXP) { +BEGIN_RCPP + Rcpp::RObject rcpp_result_gen; + Rcpp::RNGScope rcpp_rngScope_gen; + Rcpp::traits::input_parameter< const std::string& >::type lock_path(lock_pathSEXP); + Rcpp::traits::input_parameter< const NumericVector& >::type parts(partsSEXP); + Rcpp::traits::input_parameter< const double >::type scope_id(scope_idSEXP); + rcpp_result_gen = Rcpp::wrap(FARR_lock_acquire(lock_path, parts, scope_id)); + return rcpp_result_gen; +END_RCPP +} +// FARR_lock_release_scope +double FARR_lock_release_scope(const double scope_id); +RcppExport SEXP _filearray_FARR_lock_release_scope(SEXP scope_idSEXP) { +BEGIN_RCPP + Rcpp::RObject rcpp_result_gen; + Rcpp::RNGScope rcpp_rngScope_gen; + Rcpp::traits::input_parameter< const double >::type scope_id(scope_idSEXP); + rcpp_result_gen = Rcpp::wrap(FARR_lock_release_scope(scope_id)); + return rcpp_result_gen; +END_RCPP +} +// FARR_lock_commit_scope +double FARR_lock_commit_scope(const double scope_id); +RcppExport SEXP _filearray_FARR_lock_commit_scope(SEXP scope_idSEXP) { +BEGIN_RCPP + Rcpp::RObject rcpp_result_gen; + Rcpp::RNGScope rcpp_rngScope_gen; + Rcpp::traits::input_parameter< const double >::type scope_id(scope_idSEXP); + rcpp_result_gen = Rcpp::wrap(FARR_lock_commit_scope(scope_id)); + return rcpp_result_gen; +END_RCPP +} +// FARR_lock_release_explicit +double FARR_lock_release_explicit(const std::string& lock_path, const NumericVector& parts, const bool all); +RcppExport SEXP _filearray_FARR_lock_release_explicit(SEXP lock_pathSEXP, SEXP partsSEXP, SEXP allSEXP) { +BEGIN_RCPP + Rcpp::RObject rcpp_result_gen; + Rcpp::RNGScope rcpp_rngScope_gen; + Rcpp::traits::input_parameter< const std::string& >::type lock_path(lock_pathSEXP); + Rcpp::traits::input_parameter< const NumericVector& >::type parts(partsSEXP); + Rcpp::traits::input_parameter< const bool >::type all(allSEXP); + rcpp_result_gen = Rcpp::wrap(FARR_lock_release_explicit(lock_path, parts, all)); + return rcpp_result_gen; +END_RCPP +} +// FARR_lock_forget +double FARR_lock_forget(const std::string& lock_path); +RcppExport SEXP _filearray_FARR_lock_forget(SEXP lock_pathSEXP) { +BEGIN_RCPP + Rcpp::RObject rcpp_result_gen; + Rcpp::RNGScope rcpp_rngScope_gen; + Rcpp::traits::input_parameter< const std::string& >::type lock_path(lock_pathSEXP); + rcpp_result_gen = Rcpp::wrap(FARR_lock_forget(lock_path)); + return rcpp_result_gen; +END_RCPP +} +// FARR_lock_release_all +double FARR_lock_release_all(); +RcppExport SEXP _filearray_FARR_lock_release_all() { +BEGIN_RCPP + Rcpp::RObject rcpp_result_gen; + Rcpp::RNGScope rcpp_rngScope_gen; + rcpp_result_gen = Rcpp::wrap(FARR_lock_release_all()); + return rcpp_result_gen; +END_RCPP +} +// FARR_link_noreplace +IntegerVector FARR_link_noreplace(const std::string& from, const std::string& to); +RcppExport SEXP _filearray_FARR_link_noreplace(SEXP fromSEXP, SEXP toSEXP) { +BEGIN_RCPP + Rcpp::RObject rcpp_result_gen; + Rcpp::RNGScope rcpp_rngScope_gen; + Rcpp::traits::input_parameter< const std::string& >::type from(fromSEXP); + Rcpp::traits::input_parameter< const std::string& >::type to(toSEXP); + rcpp_result_gen = Rcpp::wrap(FARR_link_noreplace(from, to)); + return rcpp_result_gen; +END_RCPP +} +// FARR_lock_status +List FARR_lock_status(const std::string& lock_path); +RcppExport SEXP _filearray_FARR_lock_status(SEXP lock_pathSEXP) { +BEGIN_RCPP + Rcpp::RObject rcpp_result_gen; + Rcpp::RNGScope rcpp_rngScope_gen; + Rcpp::traits::input_parameter< const std::string& >::type lock_path(lock_pathSEXP); + rcpp_result_gen = Rcpp::wrap(FARR_lock_status(lock_path)); + return rcpp_result_gen; +END_RCPP +} +// FARR_lock_registry +List FARR_lock_registry(); +RcppExport SEXP _filearray_FARR_lock_registry() { +BEGIN_RCPP + Rcpp::RObject rcpp_result_gen; + Rcpp::RNGScope rcpp_rngScope_gen; + rcpp_result_gen = Rcpp::wrap(FARR_lock_registry()); + return rcpp_result_gen; +END_RCPP +} // FARR_buffer_map SEXP FARR_buffer_map(std::vector& input_filebases, const std::string& output_filebase, const Function& map, std::vector& buffer_nelems, int result_nelems); RcppExport SEXP _filearray_FARR_buffer_map(SEXP input_filebasesSEXP, SEXP output_filebaseSEXP, SEXP mapSEXP, SEXP buffer_nelemsSEXP, SEXP result_nelemsSEXP) { @@ -423,6 +525,15 @@ static const R_CallMethodDef CallEntries[] = { {"_filearray_filearray_subset", (DL_FUNC) &_filearray_filearray_subset, 5}, {"_filearray_FARR_subset_sequential", (DL_FUNC) &_filearray_FARR_subset_sequential, 7}, {"_filearray_FARR_subset2", (DL_FUNC) &_filearray_FARR_subset2, 8}, + {"_filearray_FARR_lock_acquire", (DL_FUNC) &_filearray_FARR_lock_acquire, 3}, + {"_filearray_FARR_lock_release_scope", (DL_FUNC) &_filearray_FARR_lock_release_scope, 1}, + {"_filearray_FARR_lock_commit_scope", (DL_FUNC) &_filearray_FARR_lock_commit_scope, 1}, + {"_filearray_FARR_lock_release_explicit", (DL_FUNC) &_filearray_FARR_lock_release_explicit, 3}, + {"_filearray_FARR_lock_forget", (DL_FUNC) &_filearray_FARR_lock_forget, 1}, + {"_filearray_FARR_lock_release_all", (DL_FUNC) &_filearray_FARR_lock_release_all, 0}, + {"_filearray_FARR_link_noreplace", (DL_FUNC) &_filearray_FARR_link_noreplace, 2}, + {"_filearray_FARR_lock_status", (DL_FUNC) &_filearray_FARR_lock_status, 1}, + {"_filearray_FARR_lock_registry", (DL_FUNC) &_filearray_FARR_lock_registry, 0}, {"_filearray_FARR_buffer_map", (DL_FUNC) &_filearray_FARR_buffer_map, 5}, {"_filearray_FARR_buffer_map2", (DL_FUNC) &_filearray_FARR_buffer_map2, 3}, {"_filearray_FARR_buffer_mapreduce", (DL_FUNC) &_filearray_FARR_buffer_mapreduce, 4}, diff --git a/src/conversion.cpp b/src/conversion.cpp index 3a568c2..9daa99f 100644 --- a/src/conversion.cpp +++ b/src/conversion.cpp @@ -290,7 +290,7 @@ SEXP convert_as2(SEXP x, SEXP y, SEXPTYPE type) { return(y); } - if( TYPEOF(y) != type ){ + if( (SEXPTYPE) TYPEOF(y) != type ){ stop("`convert_as2` inconsistent y type"); } @@ -402,7 +402,7 @@ SEXP realToCplx2(SEXP x){ void realToFloat(double* x, float* y, size_t nelem){ - for(R_xlen_t ii = 0; ii < nelem; ii++, x++, y++){ + for(size_t ii = 0; ii < nelem; ii++, x++, y++){ if(*x == NA_REAL){ *y = NA_FLOAT; } else { @@ -412,7 +412,7 @@ void realToFloat(double* x, float* y, size_t nelem){ } void floatToReal(float* x, double* y, size_t nelem){ - for(R_xlen_t ii = 0; ii < nelem; ii++, y++, x++){ + for(size_t ii = 0; ii < nelem; ii++, y++, x++){ if(ISNAN(*x)){ *y = NA_REAL; } else { diff --git a/src/load.cpp b/src/load.cpp index f0f59ba..81f8891 100644 --- a/src/load.cpp +++ b/src/load.cpp @@ -63,7 +63,7 @@ SEXP FARR_subset_sequential( const int64_t from = 0, const int64_t len = 1 ) { - if( TYPEOF(ret) != array_memory_sxptype(array_type) ){ + if( (SEXPTYPE) TYPEOF(ret) != array_memory_sxptype(array_type) ){ stop("Inconsistent `array_type` and return type"); } if( len > Rf_xlength(ret) ){ @@ -339,7 +339,7 @@ struct FARRSubsetter : public TinyParallel::Worker { void operator_mmap(std::size_t begin, std::size_t end) { - for(R_xlen_t ii = begin; ii < end; ii++){ + for(R_xlen_t ii = begin; ii < (R_xlen_t) end; ii++){ int part = partitions[ii]; int64_t skips = 0; @@ -440,7 +440,7 @@ struct FARRSubsetter : public TinyParallel::Worker { void operator_fread(std::size_t begin, std::size_t end) { std::size_t ncores = buf_ptrs.size(); - for(R_xlen_t ii = begin; ii < end; ii++){ + for(R_xlen_t ii = begin; ii < (R_xlen_t) end; ii++){ int part = partitions[ii]; int64_t skips = 0; diff --git a/src/lock.cpp b/src/lock.cpp new file mode 100644 index 0000000..9eaec33 --- /dev/null +++ b/src/lock.cpp @@ -0,0 +1,856 @@ +// Partition locks: registry and operating-system primitives. See lock.h for +// the protocol. No R headers in this file: it includes . + +#include "lock.h" + +#include +#include +#include +#include +#include +#include +#include + +#ifdef _WIN32 +# ifndef NOMINMAX +# define NOMINMAX +# endif +# ifndef WIN32_LEAN_AND_MEAN +# define WIN32_LEAN_AND_MEAN +# endif +# include +#else +# include +# include +# include +# include +#endif + +namespace farr_lock { + +namespace { + +#ifdef _WIN32 +typedef HANDLE os_handle; +const char* const LOCK_CALL = "LockFileEx"; +#else +typedef int os_handle; +const char* const LOCK_CALL = "fcntl(F_SETLK)"; +#endif + +// Identity of a lock file: device and inode (volume serial number and file +// index on Windows), so that every path of one file shares a registry entry +struct Key { + uint64_t dev; + uint64_t ino; + bool operator<(const Key& other) const { + return dev < other.dev || (dev == other.dev && ino < other.ino); + } + bool operator==(const Key& other) const { + return dev == other.dev && ino == other.ino; + } +}; + +struct Counts { + int64_t automatic; + int64_t explicit_count; + Counts() : automatic(0), explicit_count(0) {} + int64_t total() const { return automatic + explicit_count; } +}; + +struct Entry { + os_handle handle; + // Descriptors of the same file opened through another path. On POSIX, + // closing any of them would release the locks held through `handle` + std::vector aux; + // bytes this process holds + std::map bytes; +}; + +struct Scope { + Key key; + // one element per hold + std::vector bytes; +}; + +struct Registry { + std::mutex mutex; + long pid; + std::map entries; + std::map scopes; + std::map > paths; + // POSIX descriptors whose file could not be identified; closing one might + // release locks, so they stay open until release_all() + std::vector orphans; +}; + +// Key for files without a usable identity (some network or FUSE file systems) +Key path_key(const std::string& path) { + Key key; + key.dev = UINT64_MAX; + key.ino = (uint64_t) std::hash()(path); + return key; +} + +#ifdef _WIN32 + +long current_pid() { + return (long) GetCurrentProcessId(); +} + +std::string system_message(const char* call, const std::string& path, int err) { + return std::string(call) + " '" + path + "' failed (Windows error " + + std::to_string(err) + ")"; +} + +std::wstring utf8_to_wide(const std::string& s) { + int n = MultiByteToWideChar(CP_UTF8, 0, s.c_str(), -1, NULL, 0); + if (n <= 0) { + return std::wstring(); + } + std::wstring w((size_t) n, L'\0'); + MultiByteToWideChar(CP_UTF8, 0, s.c_str(), -1, &w[0], n); + w.resize((size_t) n - 1); + return w; +} + +void close_handle(os_handle h) { + if (h != INVALID_HANDLE_VALUE && h != NULL) { + CloseHandle(h); + } +} + +bool handle_key(os_handle h, const std::string& path, Key& key) { + BY_HANDLE_FILE_INFORMATION info; + if (!GetFileInformationByHandle(h, &info)) { + return false; + } + const uint64_t index = ((uint64_t) info.nFileIndexHigh << 32) | + (uint64_t) info.nFileIndexLow; + if (index == 0) { + key = path_key(path); + } else { + key.dev = (uint64_t) info.dwVolumeSerialNumber; + key.ino = index; + } + return true; +} + +// Identity without creating the file; another handle is safe on Windows, +// where locks belong to handles +bool stat_key(const std::string& path, Key& key) { + const std::wstring wpath = utf8_to_wide(path); + if (wpath.empty()) { + return false; + } + HANDLE h = CreateFileW( + wpath.c_str(), FILE_READ_ATTRIBUTES, + FILE_SHARE_READ | FILE_SHARE_WRITE | FILE_SHARE_DELETE, + NULL, OPEN_EXISTING, FILE_ATTRIBUTE_NORMAL, NULL); + if (h == INVALID_HANDLE_VALUE) { + return false; + } + const bool ok = handle_key(h, path, key); + CloseHandle(h); + return ok; +} + +int open_lock_file(const std::string& path, os_handle& h, int& err, + std::string& message) { + const std::wstring wpath = utf8_to_wide(path); + if (wpath.empty()) { + message = "cannot convert '" + path + "' to UTF-16"; + return LOCK_FAILED; + } + DWORD e = 0; + for (int i = 0; i < 4; i++) { + // GENERIC_READ is enough for LockFileEx, and does not conflict with + // programs that open the file sharing only FILE_SHARE_READ + h = CreateFileW( + wpath.c_str(), GENERIC_READ, + FILE_SHARE_READ | FILE_SHARE_WRITE | FILE_SHARE_DELETE, + NULL, OPEN_ALWAYS, FILE_ATTRIBUTE_NORMAL, NULL); + if (h != INVALID_HANDLE_VALUE) { + return LOCK_OK; + } + e = GetLastError(); + if (e != ERROR_SHARING_VIOLATION && e != ERROR_LOCK_VIOLATION) { + break; + } + // antivirus or synchronization software holding the file briefly + if (i < 3) { + Sleep(5); + } + } + err = (int) e; + message = system_message("CreateFileW", path, err); + if (e == ERROR_SHARING_VIOLATION || e == ERROR_LOCK_VIOLATION) { + return LOCK_BUSY; + } + if (e == ERROR_ACCESS_DENIED || e == ERROR_WRITE_PROTECT || + e == ERROR_NOT_SUPPORTED) { + return LOCK_UNSUPPORTED; + } + return LOCK_FAILED; +} + +void set_offset(OVERLAPPED& ov, int64_t byte) { + ZeroMemory(&ov, sizeof(ov)); + ov.Offset = (DWORD) ((uint64_t) byte & 0xFFFFFFFFULL); + ov.OffsetHigh = (DWORD) ((uint64_t) byte >> 32); +} + +int os_trylock(os_handle h, int64_t byte, int& err) { + // one OVERLAPPED per call; the handle is synchronous + OVERLAPPED ov; + set_offset(ov, byte); + if (LockFileEx(h, LOCKFILE_EXCLUSIVE_LOCK | LOCKFILE_FAIL_IMMEDIATELY, + 0, 1, 0, &ov)) { + return LOCK_OK; + } + const DWORD e = GetLastError(); + err = (int) e; + if (e == ERROR_LOCK_VIOLATION) { + return LOCK_BUSY; + } + if (e == ERROR_NOT_SUPPORTED || e == ERROR_INVALID_FUNCTION) { + return LOCK_UNSUPPORTED; + } + return LOCK_FAILED; +} + +void os_unlock(os_handle h, int64_t byte) { + // must match the locked region exactly + OVERLAPPED ov; + set_offset(ov, byte); + UnlockFileEx(h, 0, 1, 0, &ov); +} + +int os_link(const std::string& from, const std::string& to, int& err, + std::string& message) { + const std::wstring wfrom = utf8_to_wide(from); + const std::wstring wto = utf8_to_wide(to); + if (wfrom.empty() || wto.empty()) { + message = "cannot convert '" + to + "' to UTF-16"; + return LINK_FAILED; + } + if (CreateHardLinkW(wto.c_str(), wfrom.c_str(), NULL)) { + return LINK_CREATED; + } + const DWORD e = GetLastError(); + err = (int) e; + if (e == ERROR_ALREADY_EXISTS || e == ERROR_FILE_EXISTS) { + return LINK_EXISTS; + } + message = system_message("CreateHardLinkW", to, err); + // no hard links on this volume (or no permission to make one): the + // caller falls back to a rename, which reports its own errors + if (e == ERROR_INVALID_FUNCTION || e == ERROR_NOT_SUPPORTED || + e == ERROR_NOT_SAME_DEVICE || e == ERROR_ACCESS_DENIED) { + return LINK_UNSUPPORTED; + } +#ifdef ERROR_TOO_MANY_LINKS + if (e == ERROR_TOO_MANY_LINKS) { + return LINK_UNSUPPORTED; + } +#endif + return LINK_FAILED; +} + +#else // POSIX + +long current_pid() { + return (long) getpid(); +} + +std::string system_message(const char* call, const std::string& path, int err) { + return std::string(call) + " '" + path + "' failed: " + std::strerror(err) + + " (errno " + std::to_string(err) + ")"; +} + +void close_handle(os_handle fd) { + if (fd >= 0) { + ::close(fd); + } +} + +Key make_key(const struct stat& st, const std::string& path) { + if (st.st_ino == 0) { + return path_key(path); + } + Key key; + key.dev = (uint64_t) st.st_dev; + key.ino = (uint64_t) st.st_ino; + return key; +} + +// Identity without opening the file +bool stat_key(const std::string& path, Key& key) { + struct stat st; + if (::stat(path.c_str(), &st) != 0) { + return false; + } + key = make_key(st, path); + return true; +} + +bool handle_key(os_handle fd, const std::string& path, Key& key) { + struct stat st; + if (::fstat(fd, &st) != 0) { + return false; + } + key = make_key(st, path); + return true; +} + +// Errors meaning that this file system or directory cannot hold locks +bool errno_unsupported(int err) { + if (err == ENOLCK || err == ENOSYS || err == EINVAL || err == EIO || + err == EACCES || err == EPERM || err == EROFS) { + return true; + } +#ifdef ENOTSUP + if (err == ENOTSUP) { + return true; + } +#endif +#ifdef EOPNOTSUPP + if (err == EOPNOTSUPP) { + return true; + } +#endif + return false; +} + +int open_lock_file(const std::string& path, os_handle& fd, int& err, + std::string& message) { + int flags = O_RDWR | O_CREAT; +#ifdef O_CLOEXEC + flags |= O_CLOEXEC; +#endif + for (int i = 0; i < 3; i++) { + // F_WRLCK needs a descriptor open for writing; 0666 & ~umask, like + // the partition files + fd = ::open(path.c_str(), flags, 0666); + if (fd >= 0 || errno != EINTR) { + break; + } + } + if (fd < 0) { + err = errno; + message = system_message("open", path, err); + return errno_unsupported(err) ? LOCK_UNSUPPORTED : LOCK_FAILED; + } +#ifndef O_CLOEXEC + ::fcntl(fd, F_SETFD, FD_CLOEXEC); +#endif + return LOCK_OK; +} + +void set_range(struct flock& fl, short type, int64_t byte) { + // zeroed, so that platform-specific fields (l_pid, l_sysid) are 0 + std::memset(&fl, 0, sizeof(fl)); + fl.l_type = type; + fl.l_whence = SEEK_SET; + fl.l_start = (off_t) byte; + fl.l_len = 1; +} + +int os_trylock(os_handle fd, int64_t byte, int& err) { + struct flock fl; + set_range(fl, F_WRLCK, byte); + // never F_SETLKW: a blocking wait could not be interrupted + for (int i = 0; i < 3; i++) { + if (::fcntl(fd, F_SETLK, &fl) == 0) { + return LOCK_OK; + } + err = errno; + if (err == EINTR) { + continue; + } + if (err == EAGAIN || err == EACCES) { + return LOCK_BUSY; + } + return errno_unsupported(err) ? LOCK_UNSUPPORTED : LOCK_FAILED; + } + return LOCK_BUSY; +} + +void os_unlock(os_handle fd, int64_t byte) { + struct flock fl; + set_range(fl, F_UNLCK, byte); + // if this fails, closing the descriptor releases the lock + for (int i = 0; i < 3; i++) { + if (::fcntl(fd, F_SETLK, &fl) == 0 || errno != EINTR) { + break; + } + } +} + +int os_link(const std::string& from, const std::string& to, int& err, + std::string& message) { + if (::link(from.c_str(), to.c_str()) == 0) { + return LINK_CREATED; + } + err = errno; + if (err == EEXIST) { + return LINK_EXISTS; + } + message = system_message("link", to, err); + // no hard links on this file system (or no permission to make one): the + // caller falls back to a rename, which reports its own errors + if (err == EPERM || err == EACCES || err == ENOSYS || err == EXDEV || + err == EMLINK) { + return LINK_UNSUPPORTED; + } +#ifdef ENOTSUP + if (err == ENOTSUP) { + return LINK_UNSUPPORTED; + } +#endif +#ifdef EOPNOTSUPP + if (err == EOPNOTSUPP) { + return LINK_UNSUPPORTED; + } +#endif + return LINK_FAILED; +} + +// In a fork() child: close a descriptor inherited from the parent, but only +// if it still refers to the lock file +void close_if_same_file(os_handle fd, const Key& key) { + struct stat st; + if (key.dev != UINT64_MAX && ::fstat(fd, &st) == 0 && + (uint64_t) st.st_dev == key.dev && (uint64_t) st.st_ino == key.ino) { + ::close(fd); + } +} + +#endif + +Registry& registry() { + // never destroyed: the operating system releases the locks at exit + static Registry* r = []() { + Registry* reg = new Registry(); + reg->pid = current_pid(); + return reg; + }(); + return *r; +} + +// A fork() child owns none of its parent's locks; closing its copies of the +// parent's descriptors cannot release the parent's classic fcntl locks +void check_fork(Registry& r) { + const long pid = current_pid(); + if (pid == r.pid) { + return; + } +#ifndef _WIN32 + for (std::map::iterator it = r.entries.begin(); + it != r.entries.end(); ++it) { + close_if_same_file(it->second.handle, it->first); + for (size_t i = 0; i < it->second.aux.size(); i++) { + close_if_same_file(it->second.aux[i], it->first); + } + } +#endif + r.entries.clear(); + r.scopes.clear(); + r.paths.clear(); + r.orphans.clear(); + r.pid = pid; +} + +void forget_key_paths(Registry& r, const Key& key) { + for (std::map >::iterator it = r.paths.begin(); + it != r.paths.end(); ) { + it->second.erase(key); + if (it->second.empty()) { + it = r.paths.erase(it); + } else { + ++it; + } + } +} + +void close_entry(Registry& r, std::map::iterator it) { + const Key key = it->first; + close_handle(it->second.handle); + for (size_t i = 0; i < it->second.aux.size(); i++) { + close_handle(it->second.aux[i]); + } + r.entries.erase(it); + forget_key_paths(r, key); +} + +// Closes a lock file once this process holds nothing in it, so descriptors are +// only open for arrays with locks, and Windows can delete idle arrays +void close_if_idle(Registry& r, const Key& key) { + std::map::iterator it = r.entries.find(key); + if (it != r.entries.end() && it->second.bytes.empty()) { + close_entry(r, it); + } +} + +// Opens (or reuses) the registry entry of a lock file +int open_entry(Registry& r, const std::string& path, Key& key, int& err, + std::string& message) { +#ifndef _WIN32 + // Never open a second descriptor for a file this process holds locks + // in: closing it later would release them + Key known; + if (stat_key(path, known) && r.entries.count(known)) { + key = known; + r.paths[path].insert(key); + return LOCK_OK; + } +#endif + os_handle h; + const int status = open_lock_file(path, h, err, message); + if (status != LOCK_OK) { + return status; + } + Key opened; + if (!handle_key(h, path, opened)) { +#ifdef _WIN32 + err = (int) GetLastError(); + message = system_message("GetFileInformationByHandle", path, err); + CloseHandle(h); +#else + err = errno; + message = system_message("fstat", path, err); + r.orphans.push_back(h); +#endif + return LOCK_FAILED; + } + std::map::iterator it = r.entries.find(opened); + if (it == r.entries.end()) { + Entry entry; + entry.handle = h; + r.entries[opened] = entry; + } else { +#ifdef _WIN32 + // Windows locks belong to handles: closing another one is safe + CloseHandle(h); +#else + // the path changed to a file this process already has open + it->second.aux.push_back(h); +#endif + } + r.paths[path].insert(opened); + key = opened; + return LOCK_OK; +} + +// Drops one hold of `byte`, and unlocks the byte when none is left +bool release_hold(Entry& entry, int64_t byte, bool explicit_hold) { + std::map::iterator it = entry.bytes.find(byte); + if (it == entry.bytes.end()) { + return false; + } + Counts& counts = it->second; + if (explicit_hold) { + if (counts.explicit_count <= 0) { + return false; + } + counts.explicit_count--; + } else { + if (counts.automatic <= 0) { + return false; + } + counts.automatic--; + } + if (counts.total() == 0) { + os_unlock(entry.handle, byte); + entry.bytes.erase(it); + } + return true; +} + +// Registry keys of a lock path: the identities it had when locked, and the one +// it has now +std::set resolve(Registry& r, const std::string& path) { + std::set keys; + std::map >::iterator it = r.paths.find(path); + if (it != r.paths.end()) { + keys = it->second; + } + Key key; + if (stat_key(path, key)) { + keys.insert(key); + } + return keys; +} + +} // namespace + +AcquireResult acquire(const std::string& lock_path, + const std::vector& bytes, int64_t scope_id) { + AcquireResult result; + result.status = LOCK_OK; + result.acquired = 0; + result.sys_error = 0; + + Registry& r = registry(); + std::lock_guard guard(r.mutex); + check_fork(r); + if (bytes.empty()) { + return result; + } + + Key key; + const int status = open_entry(r, lock_path, key, result.sys_error, + result.message); + if (status != LOCK_OK) { + result.status = status; + return result; + } + + std::map::iterator sit = r.scopes.find(scope_id); + if (sit != r.scopes.end() && !(sit->second.key == key)) { + close_if_idle(r, key); + result.status = LOCK_FAILED; + result.message = "a lock scope cannot hold partitions of two arrays"; + return result; + } + Scope& scope = r.scopes[scope_id]; + scope.key = key; + Entry& entry = r.entries[key]; + const size_t first_new = scope.bytes.size(); + + for (size_t i = 0; i < bytes.size(); i++) { + const int64_t byte = bytes[i]; + std::map::iterator held = entry.bytes.find(byte); + if (held == entry.bytes.end()) { + int err = 0; + const int s = os_trylock(entry.handle, byte, err); + if (s == LOCK_BUSY) { + // leading bytes stay held while the caller waits for this one + result.status = LOCK_BUSY; + break; + } + if (s != LOCK_OK) { + // undo this call; earlier calls of the scope keep their holds + for (size_t j = scope.bytes.size(); j > first_new; j--) { + release_hold(entry, scope.bytes[j - 1], false); + } + scope.bytes.resize(first_new); + result.status = s; + result.acquired = 0; + result.sys_error = err; + result.message = system_message(LOCK_CALL, lock_path, err); + break; + } + held = entry.bytes.insert(std::make_pair(byte, Counts())).first; + } + held->second.automatic++; + scope.bytes.push_back(byte); + result.acquired++; + } + + if (scope.bytes.empty()) { + r.scopes.erase(scope_id); + } + close_if_idle(r, key); + return result; +} + +int64_t release_scope(int64_t scope_id) { + Registry& r = registry(); + std::lock_guard guard(r.mutex); + check_fork(r); + + std::map::iterator sit = r.scopes.find(scope_id); + if (sit == r.scopes.end()) { + return 0; + } + const Scope scope = sit->second; + r.scopes.erase(sit); + + int64_t released = 0; + std::map::iterator it = r.entries.find(scope.key); + if (it != r.entries.end()) { + for (size_t j = scope.bytes.size(); j > 0; j--) { + if (release_hold(it->second, scope.bytes[j - 1], false)) { + released++; + } + } + close_if_idle(r, scope.key); + } + return released; +} + +int64_t commit_scope(int64_t scope_id) { + Registry& r = registry(); + std::lock_guard guard(r.mutex); + check_fork(r); + + std::map::iterator sit = r.scopes.find(scope_id); + if (sit == r.scopes.end()) { + return 0; + } + const Scope scope = sit->second; + r.scopes.erase(sit); + + int64_t committed = 0; + std::map::iterator it = r.entries.find(scope.key); + if (it == r.entries.end()) { + return 0; + } + for (size_t j = 0; j < scope.bytes.size(); j++) { + std::map::iterator held = + it->second.bytes.find(scope.bytes[j]); + if (held != it->second.bytes.end() && held->second.automatic > 0) { + held->second.automatic--; + held->second.explicit_count++; + committed++; + } + } + return committed; +} + +int64_t release_explicit(const std::string& lock_path, + const std::vector& bytes, bool all) { + Registry& r = registry(); + std::lock_guard guard(r.mutex); + check_fork(r); + + int64_t released = 0; + const std::set keys = resolve(r, lock_path); + for (std::set::const_iterator k = keys.begin(); k != keys.end(); ++k) { + std::map::iterator it = r.entries.find(*k); + if (it == r.entries.end()) { + continue; + } + Entry& entry = it->second; + if (all) { + std::vector held; + for (std::map::iterator b = entry.bytes.begin(); + b != entry.bytes.end(); ++b) { + if (b->second.explicit_count > 0) { + held.push_back(b->first); + } + } + for (size_t i = 0; i < held.size(); i++) { + while (release_hold(entry, held[i], true)) { + released++; + } + } + } else { + for (size_t i = 0; i < bytes.size(); i++) { + if (release_hold(entry, bytes[i], true)) { + released++; + } + } + } + close_if_idle(r, *k); + } + return released; +} + +int64_t forget(const std::string& lock_path) { + Registry& r = registry(); + std::lock_guard guard(r.mutex); + check_fork(r); + + int64_t dropped = 0; + const std::set keys = resolve(r, lock_path); + for (std::set::const_iterator k = keys.begin(); k != keys.end(); ++k) { + std::map::iterator it = r.entries.find(*k); + if (it == r.entries.end()) { + continue; + } + for (std::map::iterator b = it->second.bytes.begin(); + b != it->second.bytes.end(); ++b) { + os_unlock(it->second.handle, b->first); + dropped++; + } + it->second.bytes.clear(); + close_entry(r, it); + for (std::map::iterator s = r.scopes.begin(); + s != r.scopes.end(); ) { + if (s->second.key == *k) { + s = r.scopes.erase(s); + } else { + ++s; + } + } + } + return dropped; +} + +int64_t release_all() { + Registry& r = registry(); + std::lock_guard guard(r.mutex); + check_fork(r); + + int64_t dropped = 0; + for (std::map::iterator it = r.entries.begin(); + it != r.entries.end(); ++it) { + for (std::map::iterator b = it->second.bytes.begin(); + b != it->second.bytes.end(); ++b) { + os_unlock(it->second.handle, b->first); + dropped++; + } + close_handle(it->second.handle); + for (size_t i = 0; i < it->second.aux.size(); i++) { + close_handle(it->second.aux[i]); + } + } + for (size_t i = 0; i < r.orphans.size(); i++) { + close_handle(r.orphans[i]); + } + r.entries.clear(); + r.scopes.clear(); + r.paths.clear(); + r.orphans.clear(); + return dropped; +} + +int link_noreplace(const std::string& from, const std::string& to, + int& sys_error, std::string& message) { + sys_error = 0; + return os_link(from, to, sys_error, message); +} + +bool status(const std::string& lock_path, std::vector& out) { + Registry& r = registry(); + std::lock_guard guard(r.mutex); + check_fork(r); + + bool open = false; + const std::set keys = resolve(r, lock_path); + for (std::set::const_iterator k = keys.begin(); k != keys.end(); ++k) { + std::map::iterator it = r.entries.find(*k); + if (it == r.entries.end()) { + continue; + } + open = true; + for (std::map::iterator b = it->second.bytes.begin(); + b != it->second.bytes.end(); ++b) { + HeldByte held; + held.byte = b->first; + held.automatic = b->second.automatic; + held.explicit_count = b->second.explicit_count; + out.push_back(held); + } + } + std::sort(out.begin(), out.end(), + [](const HeldByte& a, const HeldByte& b) { return a.byte < b.byte; }); + return open; +} + +RegistryInfo info() { + Registry& r = registry(); + std::lock_guard guard(r.mutex); + check_fork(r); + + RegistryInfo result; + result.entries = (int64_t) r.entries.size(); + result.scopes = (int64_t) r.scopes.size(); + result.paths = (int64_t) r.paths.size(); + result.pid = (int64_t) r.pid; + return result; +} + +} // namespace farr_lock diff --git a/src/lock.h b/src/lock.h new file mode 100644 index 0000000..62f1551 --- /dev/null +++ b/src/lock.h @@ -0,0 +1,101 @@ +#ifndef FARR_LOCK_H +#define FARR_LOCK_H + +/* + * File locks, protocol v1. This is a contract between filearray versions; + * R/lock.R documents the same: + * + * - A lock file is an empty file, created by the first process that locks. + * Nothing writes to it, it is never renamed or removed, and no other code + * opens it. Each array has one: "/.filearray.lock". + * - Byte 0 is the lock that with_filelock() takes. On an array's lock file + * it is the array-level lock, which partition writers do not take. + * - Byte k >= 1 of an array's lock file is the exclusive write lock of + * partition k (file ".farr"). + * - POSIX: classic fcntl record lock F_WRLCK on [k, k + 1), taken with + * non-blocking F_SETLK only. Windows: LockFileEx at offset k, length 1. + * - Locks are exclusive; of an array, only writers lock. The locks are + * advisory, and the operating system releases them when a process exits. + * + * Locks belong to the process. A registry counts the holds per byte, so + * nested writes and several objects for one array share a lock. Only the R + * main thread may call these functions. + * + * No R headers here: lock.cpp includes . + */ + +#include +#include +#include + +namespace farr_lock { + +// Status of acquire(); keep in sync with R/lock.R +const int LOCK_OK = 0; +const int LOCK_BUSY = 1; +const int LOCK_UNSUPPORTED = 2; +const int LOCK_FAILED = 3; + +// Status of link_noreplace(); keep in sync with R/lock.R +const int LINK_CREATED = 0; +const int LINK_EXISTS = 1; +const int LINK_UNSUPPORTED = 2; +const int LINK_FAILED = 3; + +struct AcquireResult { + int status; + // number of leading bytes of the request that are held after the call + int64_t acquired; + int sys_error; + std::string message; +}; + +struct HeldByte { + int64_t byte; + int64_t automatic; + int64_t explicit_count; +}; + +struct RegistryInfo { + int64_t entries; + int64_t scopes; + int64_t paths; + int64_t pid; +}; + +// Locks `bytes` (strictly increasing, >= 0) in order and stops at the first +// byte that another process holds. Every hold taken is recorded under +// `scope_id`; bytes this process already holds are re-entered without a +// system call +AcquireResult acquire(const std::string& lock_path, + const std::vector& bytes, int64_t scope_id); + +// Releases the holds recorded under a scope; returns how many +int64_t release_scope(int64_t scope_id); + +// Turns the holds of a scope into explicit holds; returns how many +int64_t commit_scope(int64_t scope_id); + +// Releases one explicit hold per byte, or every explicit hold when `all`; +// returns how many. Holds of scopes are never released here +int64_t release_explicit(const std::string& lock_path, + const std::vector& bytes, bool all); + +// Drops every hold on a lock file and closes it; returns how many bytes +int64_t forget(const std::string& lock_path); + +// Drops every hold of this process; returns how many bytes +int64_t release_all(); + +// Creates `to` as a hard link to `from`, unless `to` exists +int link_noreplace(const std::string& from, const std::string& to, + int& sys_error, std::string& message); + +// Holds on a lock file, sorted by byte; false if this process holds none +bool status(const std::string& lock_path, std::vector& out); + +RegistryInfo info(); + +} // namespace farr_lock + +#endif // FARR_LOCK_H diff --git a/src/lockExports.cpp b/src/lockExports.cpp new file mode 100644 index 0000000..021d0f6 --- /dev/null +++ b/src/lockExports.cpp @@ -0,0 +1,122 @@ +// R interface of the partition locks (lock.h); used by R/lock.R + +#include "common.h" +#include "lock.h" + +#include + +using namespace Rcpp; + +// Bytes to lock (partition numbers, or 0 for with_filelock()), validated +// before anything is locked +static std::vector lock_bytes(const NumericVector& parts) { + std::vector bytes; + bytes.reserve(parts.size()); + double previous = -1; + for (R_xlen_t i = 0; i < parts.size(); i++) { + const double v = parts[i]; + if (!R_FINITE(v) || v < 0 || v != std::floor(v) || + v > 9007199254740992.0 || v <= previous) { + Rcpp::stop("Lock bytes must be strictly increasing non-negative integers"); + } + previous = v; + bytes.push_back((int64_t) v); + } + return bytes; +} + +static int64_t lock_scope_id(const double scope_id) { + if (!R_FINITE(scope_id) || scope_id < 1 || scope_id != std::floor(scope_id)) { + Rcpp::stop("Invalid partition lock scope"); + } + return (int64_t) scope_id; +} + +// [[Rcpp::export]] +List FARR_lock_acquire(const std::string& lock_path, const NumericVector& parts, + const double scope_id) { + const std::vector bytes = lock_bytes(parts); + const int64_t scope = lock_scope_id(scope_id); + const farr_lock::AcquireResult res = farr_lock::acquire(lock_path, bytes, scope); + return List::create( + _["status"] = res.status, + _["acquired"] = (double) res.acquired, + _["sys_error"] = res.sys_error, + _["message"] = res.message + ); +} + +// [[Rcpp::export]] +double FARR_lock_release_scope(const double scope_id) { + return (double) farr_lock::release_scope(lock_scope_id(scope_id)); +} + +// [[Rcpp::export]] +double FARR_lock_commit_scope(const double scope_id) { + return (double) farr_lock::commit_scope(lock_scope_id(scope_id)); +} + +// [[Rcpp::export]] +double FARR_lock_release_explicit(const std::string& lock_path, + const NumericVector& parts, + const bool all = false) { + std::vector bytes; + if (!all) { + bytes = lock_bytes(parts); + } + return (double) farr_lock::release_explicit(lock_path, bytes, all); +} + +// [[Rcpp::export]] +double FARR_lock_forget(const std::string& lock_path) { + return (double) farr_lock::forget(lock_path); +} + +// [[Rcpp::export]] +double FARR_lock_release_all() { + return (double) farr_lock::release_all(); +} + +// [[Rcpp::export]] +IntegerVector FARR_link_noreplace(const std::string& from, const std::string& to) { + int sys_error = 0; + std::string message; + const int status = farr_lock::link_noreplace(from, to, sys_error, message); + IntegerVector re = IntegerVector::create(status); + if (!message.empty()) { + re.attr("message") = message; + } + return re; +} + +// [[Rcpp::export]] +List FARR_lock_status(const std::string& lock_path) { + std::vector held; + const bool open = farr_lock::status(lock_path, held); + const R_xlen_t n = (R_xlen_t) held.size(); + NumericVector partition(n); + IntegerVector automatic(n); + IntegerVector explicit_count(n); + for (R_xlen_t i = 0; i < n; i++) { + partition[i] = (double) held[i].byte; + automatic[i] = (int) held[i].automatic; + explicit_count[i] = (int) held[i].explicit_count; + } + return List::create( + _["open"] = open, + _["partition"] = partition, + _["automatic"] = automatic, + _["explicit"] = explicit_count + ); +} + +// [[Rcpp::export]] +List FARR_lock_registry() { + const farr_lock::RegistryInfo info = farr_lock::info(); + return List::create( + _["entries"] = (double) info.entries, + _["scopes"] = (double) info.scopes, + _["paths"] = (double) info.paths, + _["pid"] = (double) info.pid + ); +} diff --git a/src/map.cpp b/src/map.cpp index 29b529e..3e9b558 100644 --- a/src/map.cpp +++ b/src/map.cpp @@ -207,6 +207,11 @@ SEXP FARR_buffer_map( current_pos_save += expected_res_nelem; } + } catch (const Rcpp::LongjumpException& e) { + // An R error or interrupt in `map` (including a partition lock + // time-out): let R continue the jump, which also resets the + // protection stack + std::rethrow_exception(std::current_exception()); } catch(std::exception &ex){ UNPROTECT(2 + narrays); forward_exception_to_r(ex); @@ -312,6 +317,10 @@ SEXP FARR_buffer_map2( try{ SET_VECTOR_ELT(ret, iter, Shield(map(argbuffers))); + } catch (const Rcpp::LongjumpException& e) { + // An R error or interrupt in `map`: let R continue the jump with + // the original condition, which also resets the protection stack + std::rethrow_exception(std::current_exception()); } catch(std::exception &ex){ UNPROTECT(2 + narrays); forward_exception_to_r(ex); diff --git a/src/mapreduce.cpp b/src/mapreduce.cpp index f4597d4..efcc383 100644 --- a/src/mapreduce.cpp +++ b/src/mapreduce.cpp @@ -129,12 +129,16 @@ SEXP FARR_buffer_mapreduce( for(R_xlen_t part = 0; part < nparts; part++){ partition_path = fbase + std::to_string(part) + ".farr"; + // pclptr holds the cumulative number of last-margin slices up to + // each partition: this partition has pclptr[part] - pclptr[part - 1] + // slices, and its first element follows those of the partitions + // before it if(part == 0){ count = count2; - psize = (*pclptr + part); + psize = pclptr[0]; } else { - psize = (*pclptr + part) - (*pclptr + (part-1)); - count = count2 + plen * (*pclptr + (part-1)); + psize = pclptr[part] - pclptr[part - 1]; + count = count2 + plen * pclptr[part - 1]; } try { diff --git a/src/save.cpp b/src/save.cpp index 1902e07..de5128d 100644 --- a/src/save.cpp +++ b/src/save.cpp @@ -361,7 +361,7 @@ struct FARRAssigner : public TinyParallel::Worker { int elem_size = sizeof(T); - for(R_xlen_t iter = begin; iter < end; iter++){ + for(R_xlen_t iter = begin; iter < (R_xlen_t) end; iter++){ if( has_error >= 0 ){ continue; } diff --git a/src/utils.cpp b/src/utils.cpp index f74754d..b89a6e0 100644 --- a/src/utils.cpp +++ b/src/utils.cpp @@ -73,18 +73,20 @@ int guess_splitdim(SEXP dim, int elem_size, size_t buffer_bytes){ } void set_buffer(SEXP dim, int elem_size, size_t buffer_bytes, int split_dim){ - int buf_bytes = elem_size; + double buf_bytes = elem_size; for(int ii = 0; ii < split_dim; ii++){ - buf_bytes *= (int) (*(REAL(dim) + ii)); - if( buf_bytes > buffer_bytes ){ - buf_bytes = buffer_bytes; + buf_bytes *= *(REAL(dim) + ii); + if( buf_bytes > (double) buffer_bytes ){ + buf_bytes = (double) buffer_bytes; break; } } - if(buf_bytes == NA_INTEGER || buf_bytes <= 16){ + if( ISNAN(buf_bytes) || buf_bytes <= 16 ){ buf_bytes = 65536; + } else if (buf_bytes > 1073741824) { + buf_bytes = 1073741824; } - set_buffer_size(buf_bytes); + set_buffer_size((int) buf_bytes); } SEXPTYPE file_buffer_sxptype(SEXPTYPE array_type) { diff --git a/tests/testthat/helper-lock.R b/tests/testthat/helper-lock.R new file mode 100644 index 0000000..5a3006e --- /dev/null +++ b/tests/testthat/helper-lock.R @@ -0,0 +1,87 @@ +# Helpers for the partition-lock tests (test-lock.R, test-lock-process.R) + +new_lock_array <- function(dim = c(2, 3, 4), partition_size = 1) { + filearray_create(tempfile(), dim, partition_size = partition_size) +} + +lock_file <- function(x) { + file.path(x$.filebase, LOCK_FILE_NAME) +} + +tmp_files <- function(x) { + list.files(x$.filebase, pattern = "[.]tmp$", all.files = TRUE) +} + +trace_header_writes <- function(tracer) { + invisible(suppressMessages(trace( + "write_header", + tracer = tracer, + where = asNamespace("filearray"), + print = FALSE + ))) +} + +# Untracing a function that is not traced fails in an installed package +untrace_header_writes <- function() { + ns <- asNamespace("filearray") + if (inherits(get("write_header", envir = ns), "functionWithTrace")) { + suppressMessages(untrace("write_header", where = ns)) + } + invisible() +} + +# Sets an environment variable; returns a function that restores it +set_env_var <- function(name, value) { + old <- Sys.getenv(name, unset = NA) + do.call(Sys.setenv, structure(list(value), names = name)) + function() { + if (is.na(old)) { + Sys.unsetenv(name) + } else { + do.call(Sys.setenv, structure(list(old), names = name)) + } + } +} + +# Functions sent to workers must not carry the test environment along +in_worker <- function(f) { + environment(f) <- globalenv() + f +} + +# Starts `n` PSOCK workers that load a filearray build with partition locks +# (and the function named by `need`, if any), or skips the test +lock_test_cluster <- function(n = 1, need = NULL) { + testthat::skip_on_cran() + # `parallel` comes with R but is not declared in DESCRIPTION + parallel <- asNamespace("parallel") + # `R CMD check` sets R_TESTS to a relative path, which the base Rprofile + # sources at start-up; workers started from this directory would fail + restore <- set_env_var("R_TESTS", "") + on.exit(restore(), add = TRUE) + cl <- tryCatch( + parallel$makePSOCKcluster(n, timeout = 120), + error = function(e) { NULL } + ) + if (is.null(cl)) { + testthat::skip("Cannot start PSOCK workers") + } + setup <- in_worker(function(libs, version, need) { + .libPaths(libs) + ns <- tryCatch(asNamespace("filearray"), error = function(e) { NULL }) + # a regression should fail the test, not hang it + options(filearray.lock.timeout = 30) + !is.null(ns) && + identical(get0("LOCK_PROTOCOL_VERSION", envir = ns, inherits = FALSE), version) && + (is.null(need) || exists(need, envir = ns, inherits = FALSE)) + }) + ok <- tryCatch( + all(unlist(parallel$clusterCall(cl, setup, .libPaths(), LOCK_PROTOCOL_VERSION, need))), + error = function(e) { FALSE } + ) + if (!ok) { + try(parallel$stopCluster(cl), silent = TRUE) + testthat::skip("Workers load a filearray build without partition locks; install this build first") + } + cl +} diff --git a/tests/testthat/test-filelock.R b/tests/testthat/test-filelock.R new file mode 100644 index 0000000..6fb2661 --- /dev/null +++ b/tests/testthat/test-filelock.R @@ -0,0 +1,192 @@ +# with_filelock(): an exclusive lock shared by all R sessions + +test_that("with_filelock() returns the value of expr", { + lf <- tempfile(fileext = ".lock") + on.exit(unlink(lf), add = TRUE) + + expect_equal(with_filelock(1 + 1, lf), 2) + expect_invisible(with_filelock(invisible(3), lf)) + expect_visible(with_filelock(4, lf)) +}) + +test_that("with_filelock() holds byte 0 of the lock file while expr runs", { + lf <- tempfile(fileext = ".lock") + on.exit(unlink(lf), add = TRUE) + + held <- with_filelock(FARR_lock_status(lock_os_path(lf)), lf) + expect_true(held$open) + expect_equal(held$partition, 0) + expect_equal(held$automatic, 1L) + expect_equal(held$explicit, 0L) + + registry <- FARR_lock_registry() + expect_equal(registry$entries, 0) + expect_equal(registry$scopes, 0) + # the lock file stays, empty + expect_true(file.exists(lf)) + expect_equal(file.size(lf), 0) +}) + +test_that("with_filelock() releases the lock when expr fails", { + lf <- tempfile(fileext = ".lock") + on.exit(unlink(lf), add = TRUE) + + expect_error(with_filelock(stop("boom"), lf), "boom") + registry <- FARR_lock_registry() + expect_equal(registry$entries, 0) + expect_equal(registry$scopes, 0) +}) + +test_that("nested with_filelock() on one lock file runs at once", { + lf <- tempfile(fileext = ".lock") + on.exit(unlink(lf), add = TRUE) + + # an inner call that waited for the outer lock would time out + expect_equal(with_filelock(with_filelock(1, lf, timeout = 0), lf), 1) +}) + +test_that("with_filelock() validates its arguments", { + lf <- tempfile(fileext = ".lock") + on.exit(unlink(lf), add = TRUE) + + expect_error(with_filelock(1, NA), "lockfile") + expect_error(with_filelock(1, c("a.lock", "b.lock")), "lockfile") + expect_error(with_filelock(1, 1), "lockfile") + expect_error(with_filelock(1, ""), "lockfile") + expect_error(with_filelock(1, lf, timeout = -1), "timeout") + expect_error(with_filelock(1, lf, timeout = NA), "timeout") + expect_error(with_filelock(1, lf, on_unsupported = "maybe")) + # nothing is created before the arguments are checked + expect_false(file.exists(lf)) +}) + +test_that("on_unsupported decides what happens when the lock cannot be taken", { + # a lock file cannot be created in a directory that does not exist + lf <- file.path(tempfile(), "x.lock") + ran <- 0 + + expect_warning( + expect_equal(with_filelock({ + ran <- ran + 1 + "done" + }, lf), "done"), + class = "filearray_lock_unsupported" + ) + # once per lock file and session + expect_warning(with_filelock({ + ran <- ran + 1 + }, lf), NA) + expect_equal(ran, 2) + + expect_error(with_filelock({ + ran <- ran + 1 + }, lf, on_unsupported = "error"), class = "filearray_lock_unsupported") + expect_equal(ran, 2) + + expect_warning(with_filelock({ + ran <- ran + 1 + }, lf, on_unsupported = "ignore"), NA) + expect_equal(ran, 3) +}) + +test_that("another session waits for with_filelock()", { + skip_on_cran() + parallel <- asNamespace("parallel") + cl <- lock_test_cluster(2, need = "with_filelock") + lf <- tempfile(fileext = ".lock") + marker <- tempfile() + on.exit({ + try(parallel$stopCluster(cl), silent = TRUE) + unlink(c(lf, marker)) + }, add = TRUE) + + task <- in_worker(function(role, lf, marker) { + if (role == 1) { + filearray::with_filelock({ + file.create(marker) + Sys.sleep(1.5) + }, lf) + return(NA_real_) + } + deadline <- proc.time()[["elapsed"]] + 20 + while (!file.exists(marker) && proc.time()[["elapsed"]] < deadline) { + Sys.sleep(0.01) + } + start <- proc.time()[["elapsed"]] + filearray::with_filelock(TRUE, lf) + proc.time()[["elapsed"]] - start + }) + res <- parallel$clusterApply(cl, 1:2, task, lf = lf, marker = marker) + + expect_gt(res[[2]], 0.5) +}) + +test_that("with_filelock() gives up after its time-out without running expr", { + skip_on_cran() + parallel <- asNamespace("parallel") + cl <- lock_test_cluster(2, need = "with_filelock") + lf <- tempfile(fileext = ".lock") + marker <- tempfile() + on.exit({ + try(parallel$stopCluster(cl), silent = TRUE) + unlink(c(lf, marker)) + }, add = TRUE) + + task <- in_worker(function(role, lf, marker) { + if (role == 1) { + filearray::with_filelock({ + file.create(marker) + Sys.sleep(2) + }, lf) + return(NULL) + } + deadline <- proc.time()[["elapsed"]] + 20 + while (!file.exists(marker) && proc.time()[["elapsed"]] < deadline) { + Sys.sleep(0.01) + } + ran <- FALSE + cls <- tryCatch({ + filearray::with_filelock({ + ran <- TRUE + }, lf, timeout = 0.3) + "no error" + }, filearray_lock_timeout = function(e) { + "filearray_lock_timeout" + }) + list(class = cls, ran = ran) + }) + res <- parallel$clusterApply(cl, 1:2, task, lf = lf, marker = marker) + + expect_identical(res[[2]], list(class = "filearray_lock_timeout", ran = FALSE)) +}) + +test_that("with_filelock() lets one session at a time run", { + skip_on_cran() + parallel <- asNamespace("parallel") + cl <- lock_test_cluster(2, need = "with_filelock") + lf <- tempfile(fileext = ".lock") + log <- tempfile() + on.exit({ + try(parallel$stopCluster(cl), silent = TRUE) + unlink(c(lf, log)) + }, add = TRUE) + + task <- in_worker(function(k, lf, log) { + for (i in 1:3) { + filearray::with_filelock({ + cat(sprintf("start %d\n", k), file = log, append = TRUE) + Sys.sleep(0.2) + cat(sprintf("end %d\n", k), file = log, append = TRUE) + }, lf) + } + TRUE + }) + parallel$clusterApply(cl, 1:2, task, lf = lf, log = log) + + lines <- readLines(log) + expect_length(lines, 12) + # every start is directly followed by the end of the same session + starts <- lines[c(TRUE, FALSE)] + ends <- lines[c(FALSE, TRUE)] + expect_equal(sub("start", "end", starts), ends) +}) diff --git a/tests/testthat/test-lock-process.R b/tests/testthat/test-lock-process.R new file mode 100644 index 0000000..2377122 --- /dev/null +++ b/tests/testthat/test-lock-process.R @@ -0,0 +1,213 @@ +# Partition locks between processes. PSOCK workers load the installed +# filearray: install this build first (R CMD INSTALL), or these tests skip. +# `parallel` and `tools` come with R but are not declared in DESCRIPTION, so +# they are reached with asNamespace(), only in tests that skip on CRAN + +test_that("another process cannot lock a partition this session holds", { + skip_on_cran() + parallel <- asNamespace("parallel") + cl <- lock_test_cluster(1) + x <- new_lock_array() + on.exit({ + try(parallel$stopCluster(cl), silent = TRUE) + x$unlock_partition() + x$delete(force = TRUE) + }, add = TRUE) + + probe <- in_worker(function(fb, part, timeout) { + y <- filearray::filearray_load(fb) + start <- proc.time()[["elapsed"]] + locked <- isTRUE(y$lock_partition(part, timeout = timeout)) + waited <- proc.time()[["elapsed"]] - start + if (locked) { + y$unlock_partition(part) + } + list(locked = locked, waited = waited) + }) + + x$lock_partition(2) + res <- parallel$clusterCall(cl, probe, x$.filebase, 2, 0.5)[[1]] + expect_false(res$locked) + expect_gte(res$waited, 0.4) + res <- parallel$clusterCall(cl, probe, x$.filebase, 3, 0)[[1]] + expect_true(res$locked) + + x$unlock_partition(2) + res <- parallel$clusterCall(cl, probe, x$.filebase, 2, 0)[[1]] + expect_true(res$locked) +}) + +test_that("a write waits while another process holds the partition", { + skip_on_cran() + parallel <- asNamespace("parallel") + cl <- lock_test_cluster(2) + x <- new_lock_array() + marker <- tempfile() + on.exit({ + try(parallel$stopCluster(cl), silent = TRUE) + unlink(marker) + x$delete(force = TRUE) + }, add = TRUE) + + task <- in_worker(function(role, fb, marker) { + y <- filearray::filearray_load(fb) + if (role == 1) { + y$lock_partition(1) + file.create(marker) + Sys.sleep(1.5) + y[, , 1] <- 1 + y$unlock_partition(1) + return(NA_real_) + } + deadline <- proc.time()[["elapsed"]] + 20 + while (!file.exists(marker) && proc.time()[["elapsed"]] < deadline) { + Sys.sleep(0.01) + } + start <- proc.time()[["elapsed"]] + y[, , 1] <- 2 + proc.time()[["elapsed"]] - start + }) + res <- parallel$clusterApply(cl, 1:2, task, fb = x$.filebase, marker = marker) + + expect_gt(res[[2]], 0.5) + expect_equal(x[, , 1], matrix(2, 2, 3)) +}) + +test_that("processes creating the same partition keep each other's data", { + skip_on_cran() + parallel <- asNamespace("parallel") + cl <- lock_test_cluster(2) + # one partition holding both slices, created by whichever process is first + x <- new_lock_array(c(2, 3, 2), partition_size = 2) + on.exit({ + try(parallel$stopCluster(cl), silent = TRUE) + x$delete(force = TRUE) + }, add = TRUE) + + task <- in_worker(function(k, fb) { + # a slow disk: each partition header takes a second to write + suppressMessages(trace( + "write_header", tracer = quote(Sys.sleep(1)), + where = asNamespace("filearray"), print = FALSE + )) + if (k == 2) { + Sys.sleep(0.3) + } + y <- filearray::filearray_load(fb) + y[, , k] <- k + TRUE + }) + parallel$clusterApply(cl, 1:2, task, fb = x$.filebase) + + expect_equal(x[, , 1], matrix(1, 2, 3)) + expect_equal(x[, , 2], matrix(2, 2, 3)) +}) + +test_that("a write that times out raises an error and holds nothing", { + skip_on_cran() + parallel <- asNamespace("parallel") + cl <- lock_test_cluster(1) + x <- new_lock_array() + old <- options(filearray.lock.timeout = 0.3) + on.exit({ + options(old) + try(parallel$stopCluster(cl), silent = TRUE) + x$delete(force = TRUE) + }, add = TRUE) + + hold <- in_worker(function(fb) { + filearray::filearray_load(fb)$lock_partition(1, timeout = 0) + }) + release <- in_worker(function(fb) { + filearray::filearray_load(fb)$unlock_partition(1) + }) + + x[] <- 0 + expect_true(parallel$clusterCall(cl, hold, x$.filebase)[[1]]) + + start <- proc.time()[["elapsed"]] + expect_error(x[, , 1] <- 1, class = "filearray_lock_timeout") + expect_gte(proc.time()[["elapsed"]] - start, 0.25) + expect_equal(x[, , 1], matrix(0, 2, 3)) + expect_equal(FARR_lock_registry()$entries, 0) + + # other partitions are not affected + x[, , 2] <- 2 + expect_equal(x[, , 2], matrix(2, 2, 3)) + + # set_partition() locks before it reads the partition it rewrites + suppressMessages(trace( + "load_partition", tracer = quote(stop("read before locking")), + where = asNamespace("filearray"), print = FALSE + )) + on.exit(suppressMessages(untrace( + "load_partition", where = asNamespace("filearray") + )), add = TRUE) + expect_error(x$set_partition(1, 7, 1, 1, 1), class = "filearray_lock_timeout") + + # fmap() locks its output before it writes to it + src <- new_lock_array() + on.exit(src$delete(force = TRUE), add = TRUE) + src[] <- 5 + expect_error(fmap(src, function(input) { input[[1]] }, .y = x), + class = "filearray_lock_timeout") + expect_equal(x[, , 1], matrix(0, 2, 3)) + expect_equal(x[, , 2], matrix(2, 2, 3)) + + expect_false(x$lock_partition(1, timeout = 0)) + parallel$clusterCall(cl, release, x$.filebase) + expect_true(x$lock_partition(1, timeout = 0)) + x$unlock_partition() +}) + +test_that("the locks of a killed process are released", { + skip_on_cran() + parallel <- asNamespace("parallel") + tools <- asNamespace("tools") + cl <- lock_test_cluster(1) + x <- new_lock_array() + on.exit({ + try(parallel$stopCluster(cl), silent = TRUE) + x$unlock_partition() + x$delete(force = TRUE) + }, add = TRUE) + + hold <- in_worker(function(fb) { + filearray::filearray_load(fb)$lock_partition(1, timeout = 0) + }) + pid <- parallel$clusterCall(cl, Sys.getpid)[[1]] + # a killed R process cannot remove its own session temporary directory + worker_tmp <- parallel$clusterCall(cl, tempdir)[[1]] + on.exit(unlink(worker_tmp, recursive = TRUE), add = TRUE) + + expect_true(parallel$clusterCall(cl, hold, x$.filebase)[[1]]) + expect_false(x$lock_partition(1, timeout = 0)) + + signal <- if (.Platform$OS.type == "windows") tools$SIGTERM else tools$SIGKILL + tools$pskill(pid, signal) + # Windows may release the locks of a terminated process with a delay + expect_true(x$lock_partition(1, timeout = 10)) +}) + +test_that("forked children do not inherit the parent's locks", { + skip_on_cran() + skip_on_os("windows") + parallel <- asNamespace("parallel") + x <- new_lock_array() + on.exit({ + x$unlock_partition() + x$delete(force = TRUE) + }, add = TRUE) + + x$lock_partition(1) + job <- parallel$mcparallel(list( + busy = x$lock_partition(1, timeout = 0), + free = isTRUE(x$lock_partition(2, timeout = 0)) + )) + expect_identical(parallel$mccollect(job)[[1]], list(busy = FALSE, free = TRUE)) + + # the first child's exit did not release the parent's lock + job <- parallel$mcparallel(x$lock_partition(1, timeout = 0)) + expect_false(parallel$mccollect(job)[[1]]) + expect_equal(lock_status(x)$explicit, 1L) +}) diff --git a/tests/testthat/test-lock.R b/tests/testthat/test-lock.R new file mode 100644 index 0000000..78ef891 --- /dev/null +++ b/tests/testthat/test-lock.R @@ -0,0 +1,411 @@ +# Partition locks within one R session. Tests that need a second process +# are in test-lock-process.R + +test_that("partitions_of_slices() maps last-margin slices to partitions", { + x <- new_lock_array(c(2, 3, 9), partition_size = 2) + on.exit(x$delete(force = TRUE), add = TRUE) + + # partitions hold slices 1-2, 3-4, 5-6, 7-8 and 9 + expect_equal(partitions_of_slices(x, c(9, 1, 2, 3)), c(1, 2, 5)) + # C++ truncates fractional indices, so 2.5 is slice 2 and 4.9 is slice 4 + expect_equal(partitions_of_slices(x, 2.5), 1) + expect_equal(partitions_of_slices(x, c(4.9, 5)), c(2, 3)) +}) + +test_that("lock_partition() holds are counted per session", { + x <- new_lock_array() + on.exit(x$delete(force = TRUE), add = TRUE) + + expect_true(x$lock_partition(c(3, 1))) + st <- lock_status(x) + expect_true(st$open) + expect_equal(st$partition, c(1, 3)) + expect_equal(st$explicit, c(1L, 1L)) + expect_equal(st$automatic, c(0L, 0L)) + + x$lock_partition(1) + expect_equal(lock_status(x)$explicit, c(2L, 1L)) + + # writing a partition this session holds re-enters its lock + x[, , 1] <- 1:6 + expect_equal(x[, , 1], matrix(as.double(1:6), 2, 3)) + st <- lock_status(x) + expect_equal(st$explicit, c(2L, 1L)) + expect_equal(st$automatic, c(0L, 0L)) + + expect_true(x$unlock_partition(1)) + expect_equal(lock_status(x)$explicit, c(1L, 1L)) + + expect_true(x$unlock_partition()) + expect_false(lock_status(x)$open) + expect_false(x$unlock_partition(2)) +}) + +test_that("unlock_partition() cannot release a lock a write is using", { + x <- new_lock_array() + on.exit(x$delete(force = TRUE), add = TRUE) + + scope <- lock_scope_begin() + lock_scope_acquire(scope, x, 2) + expect_equal(lock_status(x)$automatic, 1L) + + expect_false(x$unlock_partition(2)) + expect_equal(lock_status(x)$automatic, 1L) + + lock_scope_end(scope) + expect_false(lock_status(x)$open) +}) + +test_that("writes leave no partition locks behind", { + x <- new_lock_array(c(2, 3, 4), partition_size = 2) + others <- list() + on.exit({ + x$delete(force = TRUE) + for (o in others) { + o$delete(force = TRUE) + } + }, add = TRUE) + + x[] <- 1:24 + # slice 4 <- 101:106, slice 1 <- 107:112, slice 3 <- 113:118 + x[, , c(4, 1, 3)] <- 101:118 + # partition 2 holds slices 3 and 4 + x$fill_partition(2, 0) + # element [1, 1, 1] of partition 1 is x[1, 1, 1] + x$set_partition(1, 7, 1, 1, 1) + expect_equal(x[1, 1, ], c(7, 7, 0, 0)) + expect_equal(x[2, 3, 1], 112) + + others$fmap <- fmap(x, function(input) { + input[[1]] + 1 + }, .buffer_count = 2) + others$element_wise <- fmap_element_wise(x, function(input) { + input[[1]] * 2 + }) + others$as_filearray <- as_filearray(array(1:24, c(2, 3, 4))) + others$bind <- filearray_bind(x, x) + expect_equal(others$fmap[1, 1, ], c(8, 8, 1, 1)) + expect_equal(others$element_wise[2, 3, 1], 224) + expect_equal(others$bind[1, 1, ], c(7, 7, 0, 0, 7, 7, 0, 0)) + + registry <- FARR_lock_registry() + expect_equal(registry$entries, 0) + expect_equal(registry$scopes, 0) + expect_true(file.exists(lock_file(x))) + expect_equal(file.size(lock_file(x)), 0) + expect_false(LOCK_FILE_NAME %in% list.files(x$.filebase)) +}) + +test_that("fmap callbacks can write to the output array", { + x <- new_lock_array(c(2, 3, 2)) + out <- new_lock_array(c(2, 3, 2)) + on.exit({ + x$delete(force = TRUE) + out$delete(force = TRUE) + }, add = TRUE) + x[] <- 1:12 + # fmap holds every partition of `out`; a write that had to wait for + # those locks would time out immediately + old <- options(filearray.lock.timeout = 0) + on.exit(options(old), add = TRUE) + + calls <- 0 + expect_warning(fmap(x, function(input) { + calls <<- calls + 1 + if (calls == 2) { + # partition 1 was written after the first call + out[1, 1, 1] <- -99 + } + input[[1]] * 2 + }, .y = out, .buffer_count = 2), NA) + expect_equal(calls, 2) + expect_equal(out[], array(c(-99, seq(4, 24, by = 2)), c(2, 3, 2))) +}) + +test_that("fmap() stops when the mapping function fails", { + x <- new_lock_array(c(2, 3, 2)) + out <- new_lock_array(c(2, 3, 2)) + on.exit({ + x$delete(force = TRUE) + out$delete(force = TRUE) + }, add = TRUE) + x[] <- 1:12 + + expect_error(fmap(x, function(input) { + stop("boom") + }, .y = out), "boom") + expect_equal(FARR_lock_registry()$scopes, 0) +}) + +test_that("a failed write releases its partition locks", { + x <- new_lock_array() + on.exit({ + untrace_header_writes() + x$delete(force = TRUE) + }, add = TRUE) + + trace_header_writes(quote(stop("disk full"))) + expect_error(x[, , 1] <- 1, "disk full") + expect_error(x$fill_partition(2, 1), "disk full") + untrace_header_writes() + + registry <- FARR_lock_registry() + expect_equal(registry$entries, 0) + expect_equal(registry$scopes, 0) + expect_length(tmp_files(x), 0) +}) + +test_that("fill_partition() writes through a new temporary file each time", { + x <- new_lock_array() + on.exit({ + untrace_header_writes() + x$delete(force = TRUE) + }, add = TRUE) + + # the temporary file exists while its header is written + seen <- new.env() + seen$names <- character(0) + fb <- x$.filebase + trace_header_writes(bquote({ + assign("names", c(get("names", envir = .(seen)), + list.files(.(fb), pattern = "[.]tmp$")), + envir = .(seen)) + })) + x$fill_partition(1, 1) + x$fill_partition(1, 2) + untrace_header_writes() + + expect_length(seen$names, 2) + expect_match(seen$names, "^0\\.[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}\\.tmp$") + expect_false(identical(seen$names[[1]], seen$names[[2]])) + expect_length(tmp_files(x), 0) + expect_equal(x[, , 1], matrix(2, 2, 3)) +}) + +test_that("a locked fill removes stale temporary files of killed writers", { + x <- new_lock_array() + on.exit(x$delete(force = TRUE), add = TRUE) + + stale <- file.path(x$.filebase, "0.0f0e5f0c-6f4a-4f2e-9d55-2c7b6c1d2e3f.tmp") + recent <- file.path(x$.filebase, "0.1a2b3c4d-5e6f-4a1b-8c9d-0e1f2a3b4c5d.tmp") + stale_other <- file.path(x$.filebase, "1.0f0e5f0c-6f4a-4f2e-9d55-2c7b6c1d2e3f.tmp") + file.create(c(stale, recent, stale_other)) + Sys.setFileTime(c(stale, stale_other), Sys.time() - 7200) + + x$fill_partition(1, 5) + + expect_false(file.exists(stale)) + expect_true(file.exists(recent)) + # partition 1's fill only cleans up partition 1 + expect_true(file.exists(stale_other)) +}) + +test_that("initialize_partition() never replaces a partition that exists", { + # Another process creates partition 1 while this one writes its NA-filled + # copy: the partition the other process created must survive + src <- new_lock_array() + src[] <- 42 + x <- new_lock_array() + on.exit({ + untrace_header_writes() + src$delete(force = TRUE) + x$delete(force = TRUE) + }, add = TRUE) + + created <- new.env() + created$done <- FALSE + trace_header_writes(bquote({ + if (!get("done", envir = .(created))) { + assign("done", TRUE, envir = .(created)) + file.copy(.(src$partition_path(1)), .(x$partition_path(1))) + } + })) + x$initialize_partition(1) + untrace_header_writes() + + expect_true(created$done) + expect_equal(x[, , 1], matrix(42, 2, 3)) + expect_length(tmp_files(x), 0) +}) + +test_that("without hard links, missing partitions are still created safely", { + skip_if_not_installed("testthat", "3.1.7") + local_mocked_bindings(FARR_link_noreplace = function(from, to) { + 2L + }) + src <- new_lock_array() + src[] <- 42 + x <- new_lock_array() + on.exit({ + untrace_header_writes() + src$delete(force = TRUE) + x$delete(force = TRUE) + }, add = TRUE) + + x$initialize_partition(1) + expect_true(all(is.na(x[, , 1]))) + + # a partition that appears meanwhile is kept + created <- new.env() + created$done <- FALSE + trace_header_writes(bquote({ + if (!get("done", envir = .(created))) { + assign("done", TRUE, envir = .(created)) + file.copy(.(src$partition_path(2)), .(x$partition_path(2))) + } + })) + x$initialize_partition(2) + untrace_header_writes() + + expect_true(created$done) + expect_equal(x[, , 2], matrix(42, 2, 3)) + expect_length(tmp_files(x), 0) +}) + +test_that("automatic locking can be turned off", { + x <- new_lock_array() + old <- options(filearray.lock = FALSE) + on.exit({ + options(old) + x$delete(force = TRUE) + }, add = TRUE) + + x[] <- 1 + expect_false(file.exists(lock_file(x))) + + options(filearray.lock = NULL) + restore_env <- set_env_var("FILEARRAY_LOCK", "false") + on.exit(restore_env(), add = TRUE) + x[] <- 2 + expect_false(file.exists(lock_file(x))) + + # explicit locks do not depend on the switch + expect_true(x$lock_partition(1)) + expect_true(file.exists(lock_file(x))) + expect_equal(lock_status(x)$explicit, 1L) + x$unlock_partition() +}) + +test_that("lock arguments are validated", { + x <- new_lock_array() + on.exit(x$delete(force = TRUE), add = TRUE) + + expect_error(x$lock_partition(0), "partition") + expect_error(x$lock_partition(NA), "partition") + expect_error(x$lock_partition(5), "partition") + expect_error(x$lock_partition("a"), "partition") + expect_error(x$lock_partition(1, timeout = -1), "timeout") + expect_error(x$lock_partition(1, timeout = NA), "timeout") + expect_error(x$lock_partition(1, timeout = c(1, 2)), "timeout") + expect_false(lock_status(x)$open) +}) + +test_that("a bad time-out option warns instead of stopping writes", { + x <- new_lock_array() + old <- options(filearray.lock.timeout = "soon", filearray.quiet = FALSE) + on.exit({ + options(old) + x$delete(force = TRUE) + }, add = TRUE) + + expect_warning(x[] <- 1, "filearray.lock.timeout") + expect_equal(x[1, 1, 1], 1) + + options(filearray.lock.timeout = NULL) + restore_env <- set_env_var("FILEARRAY_LOCK_TIMEOUT", "-5") + on.exit(restore_env(), add = TRUE) + expect_warning(x[] <- 2, "FILEARRAY_LOCK_TIMEOUT") + expect_equal(x[1, 1, 1], 2) +}) + +test_that("a lock file that cannot be opened does not stop writes", { + x <- new_lock_array() + old <- options(filearray.quiet = FALSE) + on.exit({ + options(old) + x$delete(force = TRUE) + }, add = TRUE) + + # a directory where the lock file belongs cannot be opened as a file + dir.create(lock_file(x)) + + expect_warning(x[, , 1] <- 1:6, "Cannot lock partitions") + expect_equal(x[, , 1], matrix(as.double(1:6), 2, 3)) + + expect_warning(res <- x$lock_partition(2), "not locked") + expect_true(res) + expect_identical(attr(res, "locked"), FALSE) +}) + +test_that("lock_partition() works on read-only arrays", { + x <- new_lock_array() + on.exit(x$delete(force = TRUE), add = TRUE) + + y <- filearray_load(x$.filebase, mode = "readonly") + expect_true(y$lock_partition(1)) + expect_equal(lock_status(x)$explicit, 1L) + expect_true(y$unlock_partition(1)) + expect_false(lock_status(x)$open) +}) + +test_that("delete() releases the session's locks on the array", { + x <- new_lock_array() + fb <- x$.filebase + on.exit(unlink(fb, recursive = TRUE), add = TRUE) + + x$lock_partition(1:2) + x$delete() + + expect_false(dir.exists(fb)) + expect_equal(FARR_lock_registry()$entries, 0) +}) + +test_that("objects and path aliases of one array share its locks", { + skip_on_os("windows") + x <- new_lock_array() + alias <- tempfile() + on.exit({ + unlink(alias) + x$delete(force = TRUE) + }, add = TRUE) + file.symlink(x$.filebase, alias) + + y <- filearray_load(x$.filebase) + z <- filearray_load(alias) + # load() resolves the link; keep the alias so the lock path differs + z$.filebase <- alias + + expect_true(z$lock_partition(2)) + expect_equal(FARR_lock_registry()$entries, 1) + expect_equal(lock_status(x)$explicit, 1L) + expect_true(y$unlock_partition(2)) + expect_false(lock_status(x)$open) +}) + +test_that("an array directory without write access warns and still writes", { + skip_on_os("windows") + x <- new_lock_array() + old <- options(filearray.lock = FALSE, filearray.quiet = FALSE) + on.exit({ + options(old) + Sys.chmod(x$.filebase, "755") + x$delete(force = TRUE) + }, add = TRUE) + + # create every partition, but not the lock file + x[] <- 1 + options(filearray.lock = TRUE) + Sys.chmod(x$.filebase, "555") + skip_if(file.access(x$.filebase, 2) == 0, "the directory stays writable (root)") + + expect_warning(x[, , 1] <- 2, "Cannot lock partitions") + expect_warning(x[, , 1] <- 3, NA) + expect_equal(x[, , 1], matrix(3, 2, 3)) + + # explicit locks report every time that nothing is locked + expect_warning(res <- x$lock_partition(1), "not locked") + expect_true(res) + expect_identical(attr(res, "locked"), FALSE) + expect_warning(x$lock_partition(1), "not locked") + expect_false(file.exists(lock_file(x))) +}) diff --git a/tests/testthat/test-map.R b/tests/testthat/test-map.R index 0fbc2b8..583741b 100644 --- a/tests/testthat/test-map.R +++ b/tests/testthat/test-map.R @@ -214,3 +214,82 @@ test_that("fwhich", { x$delete() }) + +test_that("mapreduce visits every slice of partitions with several slices", { + # 9 slices of 2 elements in partitions of 3 slices: elements 1-6, 7-12, 13-18 + x <- filearray_create(temp_path(check = TRUE), c(2, 9), partition_size = 3) + on.exit(x$delete(force = TRUE), add = TRUE) + x[] <- 1:18 + + chunks <- mapreduce(x, map = function(data, size, first_index) { + c(first_index = first_index, size = size, first = data[[1]]) + }, reduce = function(mapped) { + do.call(rbind, mapped) + }) + expect_equal(unname(chunks[, "first_index"]), c(1, 7, 13)) + expect_equal(unname(chunks[, "size"]), c(6, 6, 6)) + expect_equal(unname(chunks[, "first"]), c(1, 7, 13)) + + expect_equal(sum(x), 171) + expect_equal(sum(x, na.rm = TRUE), 171) + expect_equal(max(x), 18) + expect_equal(min(x), 1) + expect_equal(range(x), c(1, 18)) + expect_equal(fwhich(x, c(5, 17)), c(5, 17)) + + # a shorter last partition: elements 1-6, 7-12, 13-16 + y <- filearray_create(temp_path(check = TRUE), c(2, 8), partition_size = 3) + on.exit(y$delete(force = TRUE), add = TRUE) + y[] <- 1:16 + + chunks <- mapreduce(y, map = function(data, size, first_index) { + c(first_index = first_index, size = size) + }, reduce = function(mapped) { + do.call(rbind, mapped) + }) + expect_equal(unname(chunks[, "first_index"]), c(1, 7, 13)) + expect_equal(unname(chunks[, "size"]), c(6, 6, 4)) + expect_equal(sum(y), 136) + expect_equal(max(y), 16) + expect_equal(fwhich(y, 15), 15) + + # mapreduce skips a missing partition; the next one keeps its position + z <- filearray_create(temp_path(check = TRUE), c(2, 8), partition_size = 3) + on.exit(z$delete(force = TRUE), add = TRUE) + z[, 1:3] <- 1:6 + z[, 7:8] <- 13:16 + + chunks <- mapreduce(z, map = function(data, size, first_index) { + c(first_index = first_index, size = size) + }, reduce = function(mapped) { + do.call(rbind, mapped) + }) + expect_equal(unname(chunks[, "first_index"]), c(1, 13)) + expect_equal(unname(chunks[, "size"]), c(6, 4)) + expect_equal(fwhich(z, 15), 15) + expect_equal(sum(z, na.rm = TRUE), 79) +}) + +test_that("fmap2 passes on errors raised by the mapping function", { + x <- filearray_create(temp_path(check = TRUE), c(2, 3, 2), partition_size = 1) + on.exit(x$delete(force = TRUE), add = TRUE) + x[] <- 1:12 + + expect_error(fmap2(x, function(input) { + stop("boom from fmap2") + }), "boom from fmap2") + + # the condition itself reaches the caller, not only its message + cond <- structure( + class = c("fmap2_test_error", "error", "condition"), + list(message = "custom condition", call = NULL) + ) + expect_error(fmap2(x, function(input) { + stop(cond) + }), class = "fmap2_test_error") + + # and fmap2 still works afterwards + expect_equal(fmap2(x, function(input) { + sum(input[[1]]) + }, .buffer_count = 2), c(21, 57)) +}) diff --git a/tests/testthat/test-partition-create.R b/tests/testthat/test-partition-create.R index 8427191..a9f5dad 100644 --- a/tests/testthat/test-partition-create.R +++ b/tests/testthat/test-partition-create.R @@ -41,7 +41,8 @@ test_that("load() never lists a partition that is still being written", { x$fill_partition(1, 5) expect_equal(loads, list(TRUE, TRUE)) - expect_setequal(list.files(fb, all.files = TRUE, no.. = TRUE), + # apart from the lock file + expect_setequal(setdiff(list.files(fb, all.files = TRUE, no.. = TRUE), LOCK_FILE_NAME), c("meta", "0.farr")) }) @@ -70,7 +71,8 @@ test_that("fill_partition() stops when the partition file cannot be replaced", { expect_error(suppressWarnings(x$fill_partition(1, 5)), "0.farr", fixed = TRUE) # and no temporary file is left behind - expect_setequal(list.files(fb, all.files = TRUE, no.. = TRUE), + # apart from the lock file + expect_setequal(setdiff(list.files(fb, all.files = TRUE, no.. = TRUE), LOCK_FILE_NAME), c("meta", "0.farr")) })