R/pipeflow.R

Defines functions error pip_collect pip_clone pip_add pip_new .pip_bind .pip_insert .pip_append .pip_get_reachable_nodes .pip_steps_to_rows .pip_steps_to_nodes .pip_step_exists .pip_reindex .pip_is_indexed .pip_filter_nodes .pip_filter .pip_select_rows .pip_step_in_view .pip_view_data .pip_view_rows .pip_pipenv .pip_update_downstream .pip_halt .pip_restart .pip_run_row .pip_execute_step_call .as_pipeflow_partitioned .extract_depends .extract_fun_params .wrap_self .wrap_pipenv .assert_state .assert_pip .assert_pip_or_view .assert_logger .assert_exec_mode .is_pipeflow_partitioned .is_pipeflow_view .is_pipeflow .new_step .empty_pipeline

Documented in pip_add pip_clone pip_collect pip_new

# ---------------------
# Pipeline construction
# ---------------------
.empty_pipeline <- function() {
    data.table::data.table(
        step = character(0),
        fun = list(),
        params = list(),
        out = list(),
        state = character(0),
        tags = list(),
        locked = logical(0),
        exec = character(0), # execution mode
        time = as.POSIXct(character(0)),
        depends = list(), # names of parameters referencing other steps
        unbound = list(), # names of unbound parameters
        nodeId = integer() # the unique node ID in the internal DAG
    )
}

.new_step <- function(
    step,
    fun,
    params,
    depends,
    nodeId,
    tags = character(0),
    exec = "auto"
) {
    list(
        step = step,
        fun = list(fun),
        params = list(params),
        out = list(NULL),
        state = .step_states[["new"]][["name"]],
        tags = list(tags),
        locked = FALSE,
        exec = exec,
        time = Sys.time(),
        depends = list(depends),
        unbound = list(setdiff(names(params), names(depends))),
        nodeId = nodeId
    )
}


# ---------------
# Type predicates
# ---------------
.is_pipeflow <- function(x) {
    inherits(x, "pipeflow")
}

# A view is a pipeflow whose `rows` field is not NULL.
.is_pipeflow_view <- function(x) {
    inherits(x, "pipeflow") && !is.null(x[["view"]])
}

.is_pipeflow_partitioned <- function(x) {
    inherits(x, "pipeflow_partitioned")
}


# -------
# Asserts
# -------
.assert_exec_mode <- function(exec) {
    if (!.is_single(exec, "character") || is.na(exec)) {
        stop("exec must be a single string")
    }
    allowed <- c("auto", "split", "reduce", "plain")
    if (!(exec %in% allowed)) {
        stop("exec must be one of: ", toString(allowed))
    }
    invisible(exec)
}

.assert_logger <- function(lgr) {
    if (!is.function(lgr)) {
        stop("lgr must be a function")
    }
    if (!all(c("level", "msg") %in% names(formals(lgr)))) {
        stop("lgr must be a function with arguments 'level' and 'msg'")
    }
    invisible(lgr)
}

.assert_pip_or_view <- function(x) {
    if (!.is_pipeflow(x)) {
        stop_no_call("x must be a pipeflow pip or view")
    }
    invisible(x)
}

# Structural operations require a full pipeline, not a view.
.assert_pip <- function(x) {
    if (!.is_pipeflow(x)) {
        stop_no_call("x must be a pipeflow pip")
    }
    if (.is_pipeflow_view(x)) {
        stop_no_call("x must be a full pipeline, not a view")
    }
    invisible(x)
}

.assert_state <- function(state) {
    if (!.is_single(state, "character") || is.na(state)) {
        stop("state must be a single string")
    }
    allowed <- names(.step_states)
    if (!(state %in% allowed)) {
        stop("state must be one of: ", toString(allowed))
    }
    invisible(state)
}

# -------
# Wrapper
# -------

# Outer pipeline env wrapper that allows to create views as copied objects with
# different names and view specifications, while sharing (i.e. pointing to) the
# same underlying pipeline environment.
.wrap_pipenv <- function(pipenv, name, view = NULL) {
    structure(
        list(name = name, view = view, pipenv = pipenv),
        class = "pipeflow"
    )
}

# Wrap a step function so that `.self` is available in its body. The wrapper
# gets its own environment holding the pipeline reference, so the original
# function object is never mutated.
.wrap_self <- function(fun, self) {
    env <- new.env(parent = environment(fun))
    env[[".self"]] <- self
    eval(call("function", formals(fun), body(fun)), envir = env)
}

# ------------------------------
# Parameter & dependency parsing
# ------------------------------
.extract_fun_params <- function(fun) {
    args <- formals(fun)

    # Remove potential "..." argument
    hasDots <- "..." %in% names(args)
    if (hasDots) {
        args <- args[!names(args) %in% "..."]
    }

    # First verify that all args have default values
    is_missing_default <- function(x) {
        identical(x, quote(expr = ))
    }
    undef <- names(Filter(args, f = is_missing_default))
    if (length(undef) > 0L) {
        stop_no_call(
            paste0("'", undef, "'", collapse = ", "),
            ifelse(length(undef) > 1L, " have ", " has "),
            "no default value"
        )
    }

    # Make sure default values are returned as resolved values
    lapply(args, \(x) eval(x, envir = environment(fun)))
}

.extract_depends <- function(
    params,
    steps,
    toPos = as.integer(length(steps))
) {
    if (!is.list(params)) {
        stop_no_call("params must be a list")
    }
    if (!is.character(steps)) {
        stop_no_call("steps must be a character vector")
    }
    if (toPos < 1) {
        stop_no_call("toPos (", toPos, ") must be at least 1")
    }
    if (toPos > length(steps)) {
        stop_no_call(
            sprintf(
                "toPos (%d) exceeds number of steps (%d)",
                toPos,
                length(steps)
            )
        )
    }

    # References to other steps are marked using a formula and can be either
    # referencing earlier steps (e.g. x = ~step1) or using positional indices
    # by pointing backwards a certain number of steps (e.g. x = ~-2)
    depends <- formula_deps(params)
    if (length(depends) == 0) {
        return(character(0))
    }

    # Convert relative positional indices (e.g. x = ~-2) to step names.
    iRelPos <- which(depends |> startsWith("-"))
    if (length(iRelPos) > 0) {
        stepNumbers <- toPos + as.integer(depends[iRelPos])
        if (any(stepNumbers < 1)) {
            firstBad <- which(stepNumbers < 1)[1]
            badIdx <- depends[iRelPos[firstBad]]
            info <- "relative index %s = ~%s points outside pipeline"
            stop_no_call(sprintf(info, names(badIdx), badIdx))
        }
        depends[iRelPos] <- steps[as.integer(stepNumbers)]
    }

    unlist(depends)
}

# --------------
# Step execution
# --------------
.as_pipeflow_partitioned <- function(x) {
    if (!is.list(x)) {
        stop("split mode requires step output to be a list")
    }

    nm <- names(x)
    if (is.null(nm) || anyNA(nm) || any(!nzchar(nm))) {
        stop("split output must be a named list with non-empty keys")
    }

    if (anyDuplicated(nm) > 0L) {
        stop("split output keys must be unique")
    }

    class(x) <- c(class(x), "pipeflow_partitioned")
    x
}

.pip_execute_step_call <- function(fun, args, exec) {
    # A partitioned argument is a list of per-key values produced by a
    # step with exec = "split" and tagged with class "pipeflow_partitioned".
    partIdx <- which(vapply(
        args,
        FUN = .is_pipeflow_partitioned,
        FUN.VALUE = logical(1)
    ))

    # Split mode: run the function once on the full (unpartitioned) inputs and
    # mark the result as partitioned so downstream steps can map over it.
    if (exec == "split") {
        out <- do.call(fun, args = args)
        return(.as_pipeflow_partitioned(out))
    }

    # Plain mode: only executes a single call, so partitioned inputs
    # (which would require mapping) are not allowed.
    if (exec == "plain" && length(partIdx) > 0L) {
        stop("plain mode does not accept partitioned inputs")
    }

    # Reduce mode: combines partitioned inputs in a single call,
    # so it needs at least one partitioned input to be meaningful.
    if (exec == "reduce" && length(partIdx) == 0L) {
        stop("reduce mode requires at least one partitioned input")
    }

    # Single-call, which happens in three scenarios:
    # 1) there are no partitioned inputs, which is the standard case in a
    #    standard pipeline that runs without any split/reduce steps
    # 2) plain mode was explicitly requested (to ensure non-split mode)
    # 3) reduce mode was requested to re-combine partitioned inputs
    if (length(partIdx) == 0L || exec == "plain" || exec == "reduce") {
        return(do.call(fun, args = args))
    }

    # Auto mode with partitioned inputs: map the function over the partition
    # keys. Keys are taken from the first partitioned argument...
    keys <- names(args[[partIdx[[1]]]])
    for (k in partIdx[-1]) {
        kk <- names(args[[k]])
        if (!identical(kk, keys)) {
            stop("partitioned arguments must share identical keys")
        }
    }

    # ...and each key is computed separately by slicing all partitioned
    # arguments down to that key before calling the function.
    out <- stats::setNames(vector(mode = "list", length = length(keys)), keys)
    for (key in keys) {
        keyArgs <- args
        for (idx in partIdx) {
            keyArgs[[idx]] <- args[[idx]][[key]]
        }

        # Errors are re-raised with the key name so a failing partition can
        # be located without inspecting the whole output.
        out[[key]] <- tryCatch(
            expr = do.call(fun, args = keyArgs),
            error = function(e) {
                stop_no_call("key '", key, "': ", e$message)
            }
        )
    }

    .as_pipeflow_partitioned(out)
}

.pip_run_row <- function(x, i, lgr) {
    if (!.is_pipeflow(x)) {
        stop("x must be a pipeflow pip")
    }
    if (!.pip_is_indexed(x)) {
        .pip_reindex(x)
    }

    dat <- .pip_pipenv(x)[["data"]]
    fun <- dat[["fun"]][[i]]
    args <- dat[["params"]][[i]]
    depends <- dat[["depends"]][[i]]
    exec <- dat[["exec"]][[i]]

    # If calculation depends on results of earlier steps, get them from
    # respective referenced output slots of the pipeline.
    if (length(depends) > 0) {
        refsOut <- .pip_filter(x, on = "step", values = depends)[["out"]]
        args[names(depends)] <- refsOut
    }

    step <- dat[["step"]][[i]]

    # Keep `.self` pointing at the pipeline object being run, so steps still
    # reference the correct pipeline after cloning, subsetting or replacing.
    environment(fun)[[".self"]] <- x

    out <- withCallingHandlers(
        .pip_execute_step_call(fun = fun, args = args, exec = exec),
        error = function(e) {
            data.table::set(
                dat,
                i = i,
                j = "state",
                value = .step_states[["failed"]][["name"]]
            )
            lgr(level = "error", msg = e$message)
            stop_no_call(e$message)
        },
        warning = function(w) {
            lgr(level = "warn", msg = w$message)
        },
        message = function(m) {
            lgr(level = "info", msg = m$message)
        }
    )

    # Re-read pipeline after execution to handle scenarios where the pipeline
    # modified itself at runtime. Since we update by step name, if the current
    # step does not exist anymore, we simply skip the update.
    dat <- .pip_pipenv(x)[["data"]]
    rowNow <- data.table::chmatch(step, dat[["step"]])
    stepStillExists <- !is.na(rowNow)
    if (stepStillExists) {
        data.table::set(
            dat,
            i = rowNow,
            j = c("out", "time", "state"),
            value = list(
                list(out),
                Sys.time(),
                .step_states[["done"]][["name"]]
            )
        )
    }

    out
}

# -------------
# State updates
# -------------
.pip_restart <- function(pipenv, force = TRUE, times = 1L) {
    if (!.is_single(force, "logical")) {
        stop("force must be a single logical value")
    }
    if (!.is_single(times, "numeric") || is.na(times) || times < 1L) {
        stop("times must be a single integer value >= 1")
    }

    count <- pipenv[[".restart_count"]]
    if (count >= times) {
        pipenv[[".restart_count"]] <- 0L
        return(invisible())
    }

    pipenv[[".run_state"]][] <- "restart"
    pipenv[[".restart_count"]] <- count + 1L
    pipenv[[".restart_force"]] <- force
    invisible()
}

.pip_halt <- function(pipenv) {
    pipenv[[".run_state"]][] <- "halted"
    invisible()
}

.pip_update_downstream <- function(x, steps, what, value) {
    nodes <- .pip_get_reachable_nodes(x, steps)
    .pip_pipenv(x)[["data"]][list(nodes), (what) := value, on = "nodeId"]

    invisible(x)
}


# ----------------------------------
# Pipeline data access and filtering
# ----------------------------------

# Convenience helper function to retrieve the shared inner environment.
.pip_pipenv <- function(x) {
    .subset2(x, "pipenv")
}

# The rows covered by `x`: all pipeline rows for a full pipeline, or the
# view's `rows` selector for a view.
.pip_view_rows <- function(x) {
    view <- .subset2(x, "view")
    if (is.null(view)) {
        seq_len(nrow(.pip_pipenv(x)[["data"]]))
    } else {
        as.integer(view)
    }
}

.pip_view_data <- function(x) {
    rows <- .pip_view_rows(x)
    .pip_pipenv(x)[["data"]][rows, ]
}

# TRUE if `step` is covered by the rows of `x`. Always TRUE for a full
# pipeline (where all rows are covered) and FALSE for a view that does not
# cover the step.
.pip_step_in_view <- function(x, step) {
    absRow <- data.table::chmatch(step, .pip_pipenv(x)[["data"]][["step"]])
    if (is.na(absRow)) {
        return(FALSE)
    }
    absRow %in% .pip_view_rows(x)
}

# Resolve the `i` selector of `[.pipeflow` to absolute row indices. `i_expr`
# is the unevaluated `i` expression; it is evaluated in the context of the
# covered rows, i.e. the data of a view or the full step table of a pipeline.
# `enclos` is used as the enclosing environment for boolean filter expressions.
.pip_select_rows <- function(x, i_expr, enclos, keep_order = FALSE) {
    dat <- .pip_pipenv(x)[["data"]]
    view <- .is_pipeflow_view(x)
    rows <- .pip_view_rows(x)
    sub <- if (view) dat[rows] else dat

    value <- eval(i_expr, envir = sub, enclos = enclos)

    if (is.logical(value)) {
        if (length(value) == 1L) {
            value <- rep(value, nrow(sub))
        }
        if (length(value) != nrow(sub)) {
            stop(sprintf(
                "logical filter has length %d but the %s has %d rows",
                length(value),
                if (view) "view" else "pipeline",
                nrow(sub)
            ))
        }
        if (anyNA(value)) {
            stop("logical filter must not contain NA")
        }
        return(rows[which(value)])
    }

    if (is.numeric(value)) {
        if (anyNA(value)) {
            stop("row indices in 'i' must not contain NA")
        }
        idx <- as.integer(value)
        if (any(idx < 0L)) {
            # Negative indices exclude the corresponding rows like in base R
            if (any(idx > 0L)) {
                stop("only 0's may be mixed with negative subscripts")
            }
            idx <- seq_len(nrow(sub))[idx]
        }
        if (!keep_order) {
            idx <- sort(unique(idx))
        }
        bad <- idx[idx < 1L | idx > nrow(sub)]
        if (length(bad) > 0L) {
            stop("Invalid row indices in 'i': ", toString(bad))
        }
        return(rows[idx])
    }

    if (is.character(value)) {
        if (anyNA(value)) {
            stop("step names must not contain NA")
        }
        if (!all(nzchar(value))) {
            stop("step names must be non-empty strings")
        }
        m <- data.table::chmatch(value, dat[["step"]])
        if (anyNA(m)) {
            unknown <- unique(value[is.na(m)])
            stop("Unknown step names: ", toString(unknown), call. = FALSE)
        }
        if (!keep_order) {
            m <- sort(unique(m))
        }
        outside <- m[!(m %in% rows)]
        if (length(outside) > 0L) {
            stop(
                "step '",
                dat[["step"]][[outside[[1L]]]],
                "' is not part of the view"
            )
        }
        return(m)
    }

    stop(sprintf(
        "`i` must evaluate to row indices, step names, or a logical filter, ",
        "not %s",
        typeof(value)
    ))
}

.pip_filter <- function(x, on, values) {
    .pip_pipenv(x)[["data"]][list(values), on = on]
}

.pip_filter_nodes <- function(x, nodes) {
    .pip_pipenv(x)[["data"]][list(nodes), on = "nodeId"]
}


# --------
# Indexing
# -------.
.pip_is_indexed <- function(x) {
    !is.null(data.table::indices(.pip_pipenv(x)[["data"]]))
}

.pip_reindex <- function(x) {
    data.table::setindexv(.pip_pipenv(x)[["data"]], list("step", "nodeId"))
}


# ---------------------------
# Step lookup & DAG traversal
# ---------------------------
.pip_step_exists <- function(x, step) {
    exists(
        step,
        where = .pip_pipenv(x)[[".steps_to_nodes"]],
        inherits = FALSE
    )
}

.pip_steps_to_nodes <- function(x, steps) {
    mget(
        steps,
        envir = .pip_pipenv(x)[[".steps_to_nodes"]],
        ifnotfound = NA_integer_,
        inherits = FALSE
    )
}

.pip_steps_to_rows <- function(x, steps) {
    data <- .pip_pipenv(x)[["data"]]

    if (anyNA(steps)) {
        stop("step names must not contain NA", call. = FALSE)
    }
    if (!all(nzchar(steps))) {
        stop("step names must be non-empty strings", call. = FALSE)
    }

    i <- data.table::chmatch(steps, data[["step"]])
    if (anyNA(i)) {
        unknown <- unique(steps[is.na(i)])
        stop("Unknown step names: ", toString(unknown), call. = FALSE)
    }
    as.integer(i)
}

.pip_get_reachable_nodes <- function(x, steps, downstream = TRUE) {
    env <- .pip_pipenv(x)
    known <- intersect(steps, names(env[[".steps_to_nodes"]]))
    if (length(known) == 0L) {
        return(integer(0))
    }

    start_ids <- as.integer(mget(
        known,
        envir = env[[".steps_to_nodes"]],
        ifnotfound = NA_integer_,
        inherits = FALSE
    ))
    start_ids <- start_ids[!is.na(start_ids)]
    if (length(start_ids) == 0L) {
        return(integer(0))
    }

    if (downstream) {
        dag_get_reachable_nodes_down(env[[".dag"]], start_ids)
    } else {
        dag_get_reachable_nodes_up(env[[".dag"]], start_ids)
    }
}


# -----------------
# Pipeline addition
# -----------------

.pip_append <- function(x, step, fun, tags, exec = "auto", params = list()) {
    funParams <- .extract_fun_params(fun)
    params[names(funParams)] <- funParams

    # Provide `.self` to the step via a dedicated wrapper environment.
    fun <- .wrap_self(fun, x)

    # Determine and verify potential links to existing steps
    env <- .pip_pipenv(x)
    steps <- c(env[["data"]][["step"]], step)
    depends <- .extract_depends(params = params, steps = steps)
    refNodes <- mget(
        depends,
        envir = env[[".steps_to_nodes"]],
        ifnotfound = NA_integer_,
        inherits = FALSE
    )
    if (anyNA(refNodes)) {
        notFound <- Filter(is.na, refNodes)
        stop_no_call(
            "while adding step '",
            step,
            "' - cannot reference unknown steps: ",
            paste0("'", names(notFound), "'", collapse = ", ")
        )
    }

    # Update DAG
    d <- env[[".dag"]]
    nodeId <- as.integer(dag_add_node(d))
    if (length(refNodes) > 0) {
        dag_add_edges_to(d, from = as.integer(refNodes), to = nodeId)
    }

    # Create and append step
    newStep <- .new_step(
        step = step,
        fun = fun,
        params = params,
        depends = depends,
        tags = tags,
        exec = exec,
        nodeId = nodeId
    )

    env[["data"]] <- data.table::rbindlist(
        list(env[["data"]], newStep),
        use.names = TRUE
    )
    env[[".steps_to_nodes"]][[step]] <- nodeId
    x
}

# Insert a step at position `pos` (0 <= pos < number of steps) of the
# pipeline, i.e. after the first `pos` steps. The new row is inserted into
# the data table (one O(n) rbindlist) and the node is added to the DAG at
# the corresponding position of its topological order, leaving all existing
# steps (and their runtime state) untouched.
.pip_insert <- function(x, step, fun, tags, exec, params, pos) {
    funParams <- .extract_fun_params(fun)
    params[names(funParams)] <- funParams

    # Provide `.self` to the step via a dedicated wrapper environment.
    fun <- .wrap_self(fun, x)

    env <- .pip_pipenv(x)
    data <- env[["data"]]

    # Dependencies may only point to steps that precede the insertion point
    # (mirroring the previous step-by-step rebuild, this keeps the row order
    # a valid topological order).
    prefixSteps <- data[["step"]][seq_len(pos)]
    steps <- c(prefixSteps, step)
    depends <- .extract_depends(params = params, steps = steps)
    if (length(depends) > 0L) {
        forbidden <- depends[!(depends %in% prefixSteps)]
        if (length(forbidden) > 0L) {
            stop_no_call(
                "while adding step '",
                step,
                "' - cannot reference unknown steps: ",
                paste0("'", unname(forbidden), "'", collapse = ", ")
            )
        }
    }

    refNodes <- mget(
        depends,
        envir = env[[".steps_to_nodes"]],
        ifnotfound = NA_integer_,
        inherits = FALSE
    )
    if (anyNA(refNodes)) {
        notFound <- Filter(is.na, refNodes)
        stop_no_call(
            "while adding step '",
            step,
            "' - cannot reference unknown steps: ",
            paste0("'", names(notFound), "'", collapse = ", ")
        )
    }

    # Update DAG: add the node at the insertion position of its order
    d <- env[[".dag"]]
    nodeId <- as.integer(dag_add_node_at(d, pos))
    if (length(refNodes) > 0) {
        dag_add_edges_to(d, from = as.integer(refNodes), to = nodeId)
    }

    # Create the new step and insert its row after the first `pos` steps
    newStep <- .new_step(
        step = step,
        fun = fun,
        params = params,
        depends = depends,
        tags = tags,
        exec = exec,
        nodeId = nodeId
    )
    n <- nrow(data)
    env[["data"]] <- data.table::rbindlist(
        list(
            data[seq_len(pos)],
            newStep,
            data[seq_len(n - pos) + pos]
        ),
        use.names = TRUE
    )
    env[[".steps_to_nodes"]][[step]] <- nodeId
    x
}


# Bind two pipelines together by concatenating their steps.
.pip_bind <- function(x, y) {
    out <- pip_clone(x, name = paste0(x[["name"]], "-", y[["name"]]))
    yy <- pip_clone(y)
    yyDat <- .pip_pipenv(yy)[["data"]]

    # Resolve all name clashes directly on the cloned source pipeline.
    reserved <- .pip_pipenv(out)[["data"]][["step"]]

    `%chin%` <- data.table::`%chin%`
    for (k in seq_len(nrow(yyDat))) {
        step <- yyDat[["step"]][[k]]
        if (step %chin% reserved) {
            to <- step
            i <- 2L
            allSteps <- yyDat[["step"]]
            while (to %chin% reserved || to %chin% allSteps) {
                to <- paste0(step, i)
                i <- i + 1L
            }
            pip_rename(yy, from = step, to = to)
        }
        reserved <- c(reserved, yyDat[["step"]][[k]])
    }

    # Append all (potentially renamed) steps of y in one pass: transform each
    # row (re-create the formula dependencies after renaming and re-bind
    # `.self` to the target), allocate a fresh node id, wire the DAG edges,
    # and collect the rows for a single rbindlist.
    outEnv <- .pip_pipenv(out)
    d <- outEnv[[".dag"]]
    stepsToNodes <- outEnv[[".steps_to_nodes"]]
    rows <- vector("list", nrow(yyDat))
    for (k in seq_len(nrow(yyDat))) {
        step <- yyDat[["step"]][[k]]
        fun <- yyDat[["fun"]][[k]]
        tags <- yyDat[["tags"]][[k]]
        exec <- yyDat[["exec"]][[k]]
        params <- yyDat[["params"]][[k]]
        depends <- yyDat[["depends"]][[k]]

        # Re-create the formula dependencies in the target pipeline's context
        # so that they point to the (renamed) steps of y.
        for (arg in intersect(names(depends), names(params))) {
            params[[arg]] <- stats::as.formula(paste("~", depends[[arg]]))
        }
        # Fold current parameter values into the defaults of the existing
        # function arguments.
        fml <- formals(fun)
        for (nm in intersect(names(params), setdiff(names(fml), "..."))) {
            fml[[nm]] <- params[[nm]]
        }
        formals(fun) <- fml

        fun <- .wrap_self(fun, out)

        # Update DAG: allocate the node, register the step, add edges to its
        # upstream steps (which are either in `out` or already appended).
        nodeId <- as.integer(dag_add_node(d))
        stepsToNodes[[step]] <- nodeId
        if (length(depends) > 0L) {
            refNodes <- as.integer(unlist(mget(
                unname(depends),
                envir = stepsToNodes,
                inherits = FALSE
            )))
            dag_add_edges_to(d, from = refNodes, to = nodeId)
        }

        row <- .new_step(step, fun, params, depends, nodeId, tags, exec)
        rows[[k]] <- row
    }

    # rbind all appended rows and preserve runtime state of the source steps
    appRows <- data.table::rbindlist(rows, use.names = TRUE)
    data.table::set(
        appRows,
        j = c("out", "state", "time", "locked"),
        value = list(
            yyDat[["out"]],
            yyDat[["state"]],
            yyDat[["time"]],
            yyDat[["locked"]]
        )
    )
    outEnv[["data"]] <- data.table::rbindlist(
        list(outEnv[["data"]], appRows),
        use.names = TRUE
    )
    out
}

# ---------------------------
# Exported pipeline functions
# ---------------------------

#' Create a pipeline
#'
#' Creates a new, empty pipeline. Add steps with [pip_add()] and execute
#' them with [pip_run()].
#'
#' @param name The name of the pipeline used for display and logging.
#' @return A pipeflow pipeline object.
#'
#' @examples
#' p <- pip_new("demo") |>
#'     pip_add("numbers", \(n = 5) seq_len(n)) |>
#'     pip_add("squared", \(x = ~numbers) x^2) |>
#'     pip_add("total",   \(x = ~squared) sum(x))
#' p
#' str(p)
#' p[["name"]]  # "demo"
#' p[["view"]]  # initially NULL
#'
#' # Inner pipeline environment (for advanced usage)
#' ls(p[["pipenv"]])                # shows "data"
#' ls(p[["pipenv"]], all = TRUE)    # also shows hidden variables
#' @export
pip_new <- function(name = "pipe") {
    if (!.is_single(name, "character")) {
        stop("name must be a single string")
    }
    if (is.na(name)) {
        stop("name must not be NA")
    }

    # Main pipeline components live in an inner environment that is shared by
    # reference with all views of this pipeline.
    hash_map <- function() new.env(parent = emptyenv())
    env <- hash_map()
    env[["data"]] <- .empty_pipeline()
    env[[".dag"]] <- dag_new()
    env[[".steps_to_nodes"]] <- hash_map()

    # Pipeline states
    env[[".run_state"]] <- factor(
        "ready",
        levels = c("ready", "restart", "running", "halted", "failed")
    )
    env[[".last_run"]] <- NULL

    # Restart tracking
    env[[".restart_count"]] <- 0L
    env[[".restart_force"]] <- TRUE

    structure(
        list(name = name, view = NULL, pipenv = env),
        class = "pipeflow"
    )
}


#' Add a step
#'
#' Adds a named step to the pipeline. Each step is a function whose parameters
#' either hold constant defaults or reference the output of a prior step using
#' formula notation (`~step_name`). Dependencies are validated when the step
#' is added.
#'
#' @param x A pipeflow pipeline object.
#' @param step Unique step name.
#' @param fun Function to execute for the step. Each function parameter must
#' have a default value. Default values that are simple constants are resolved
#' immediately. Default values that are formulas like `~other_step` are
#' treated as dependencies to those steps and resolved to the respective output
#' values at runtime once the step is executed.
#' @param tags Optional character vector of tags belonging to the step.
#' Can also be adjusted later using `[pip_tag()]`.
#' @param after Optional position after which the new step should be inserted
#' (defaults to last position). Can be a step name or an integer index. If
#' set to 0, the new step will be inserted at the beginning of the pipeline.
#' @param params Optional named list of parameter values, which will be merged
#' with the defaults of `fun` (if overlapping names, the default values in `fun`
#' take precedence). There are two use cases for `params`:
#' 1. Provide param values programmatically when adding steps at runtime
#' 2. Provide extra param values to "mark" dependencies that are defined in
#'   pipelines nested in a step, which ensures that the step (and with that
#'   the pipeline in the step) is re-executed when one of the respective
#'   param values change.
#' @param exec Execution mode for this step. One of "auto", "split",
#' "reduce" or "plain".
#' Using execution mode `exec = split`, the output of the step is marked as
#' partitioned output. In this mode, any step that depends on the split step
#' (directly or indirectly) will have its output automatically mapped
#' partition-wise during step execution. The `reduce` mode expects
#' partitioned input and passes it through without mapping, while `plain`
#' mode only accepts non-partitioned input and always intends to execute
#' a single call. In summary:
#' * auto: map if partitioned input appears, otherwise single call
#' * split: single call, then mark output as partitioned
#' * reduce: single call, but only valid with partitioned input
#' * plain: single call, only valid with non-partitioned input
#'
#' @details
#' If `after` was specified, the new step will be inserted after the given
#' step or position. Be aware that in contrast to adding a step at the end,
#' inserting a step in the middle is a rather expensive operation as it
#' requires re-wiring parts of the internal pipeline structure, especially
#' if the new step is inserted at an early position.
#'
#' @return The updated pipeline, invisibly.
#' @examples
#' # --- Tags, and view filtering ---
#' p <- pip_new("analysis") |>
#'   pip_add("load", \(n = 5) seq_len(n), tags = c("io", "raw")) |>
#'   pip_add("clean", \(x = ~load) x * 2, tags = c("io", "process")) |>
#'   pip_add("fit", \(x = ~clean) sum(x), tags = c("model", "core", "daily")) |>
#'   pip_add("report", \(x = ~fit) paste("result:", x), tags = "report")
#'
#' pip_run(p)
#' p
#'
#' # Filter by tag using pip_view — keeps steps with any matching tag
#' pip_view(p, tags = "daily")
#' pip_view(p, tags = "core")
#' pip_view(p, tags = c("raw", "report"))
#'
#' # --- Split / reduce execution modes ---
#' q <- pip_new("split-demo") |>
#'   pip_add("data", \(x = iris) x) |>
#'   pip_add("split", \(x = ~data) split(x, x$Species),
#'     exec = "split"
#'   ) |>
#'   pip_add("stats", \(x = ~split) summary(x)) |>
#'   pip_add("combine", \(x = ~stats) do.call(rbind, x),
#'     exec = "reduce"
#'   )
#'
#' pip_run(q)
#' q[["stats", "out"]]   # partitioned list — one summary per species
#' q[["combine", "out"]] # combined table
#'
#' # --- Insert a step at a specific position with 'after' ---
#' p2 <- pip_new("insert-demo") |>
#'   pip_add("load", \(x = 1) x) |>
#'   pip_add("fit", \(x = ~load) x + 1)
#' p2
#' pip_add(p2, "clean", \(x = ~load) x * 2, after = "load")
#' p2  # "load", "clean", "fit" — inserted after "load"
#'
#' # after = 0 inserts the new step at the very beginning
#' pip_add(p2, "preload", \(x = 1) x, after = 0)
#' p2
#'
#' # --- Provide parameter values programmatically with 'params' ---
#' p3 <- pip_new("params-demo") |>
#'   pip_add("load", \(x = 1, ...) c(x, ...), params = list(size = 42))
#' pip_run(p3)
#' p3
#' p3[["load", "out"]] # c(1, size = 42) — extra params passed via `...`
#' @export
pip_add <- function(
    x,
    step,
    fun,
    tags = character(0),
    after = length(x),
    params = list(),
    exec = "auto"
) {
    .assert_pip(x)
    if (!.is_single(step, "character")) {
        stop("step must be a single string")
    }
    if (is.na(step)) {
        stop("step must not be NA")
    }
    if (!nzchar(step)) {
        stop("step must be a non-empty string")
    }
    if (.pip_step_exists(x, step)) {
        stop("step '", step, "' already exists in the pipeline")
    }
    if (!is.function(fun)) {
        stop("fun must be a function")
    }
    if (".self" %in% names(formals(fun))) {
        stop_no_call(
            "'.self' is a reserved parameter and cannot be used as a step name"
        )
    }
    if (".self" %in% names(params)) {
        stop_no_call(
            "'.self' is a reserved parameter name and cannot be used in params"
        )
    }
    .assert_exec_mode(exec)

    n <- length(x)

    if (is.character(after)) {
        if (!.is_single(after, "character") || is.na(after) || !nzchar(after)) {
            stop("after must be a non-empty step name or integer index")
        }
        if (!.pip_step_exists(x, after)) {
            stop("step '", after, "' does not exist")
        }
        # Most of the time the new step is added at the end, so we check that
        # first to avoid the more expensive match() call in the common case.
        pos <- n
        last <- .pip_pipenv(x)[["data"]][["step"]][n]
        if (after != last) {
            pos <- data.table::chmatch(
                after,
                .pip_pipenv(x)[["data"]][["step"]]
            )
        }
    } else if (is.numeric(after)) {
        if (length(after) != 1 || is.na(after)) {
            stop("after must be a non-empty step name or integer index")
        }
        if (!is.finite(after) || after != as.integer(after)) {
            stop("after index must be a whole number")
        }
        pos <- as.integer(after)
        if (pos < 0L || pos > n) {
            stop("after index must be between 0 and ", n)
        }
    } else {
        stop("after must be a non-empty step name or integer index")
    }

    if (pos == n) {
        .pip_append(
            x,
            step = step,
            fun = fun,
            tags = tags,
            exec = exec,
            params = params
        )
    } else {
        .pip_insert(
            x,
            step = step,
            fun = fun,
            tags = tags,
            exec = exec,
            params = params,
            pos = pos
        )
    }
    invisible(x)
}


#' Clone a pipeline
#'
#' Creates an independent copy of the pipeline. Changes to the cloned
#' pipeline do not affect the original pipeline, and vice versa.
#'
#' @param x A pipeflow pipeline object.
#' @param name Optional name for the cloned pipeline. If `NULL`, the
#' original name is used.
#'
#' @return A cloned pipeflow pipeline object.
#' @examples
#' p <- pip_new("original") |>
#'   pip_add("s1", \(x = 1) x) |>
#'   pip_add("s2", \(x = ~s1) x + 1)
#'
#' # Clone produces a fully independent copy
#' cp <- pip_clone(p, name = "copy")
#' pip_add(cp, "s3", \(x = ~s2) x * 10)
#'
#' # As a result, the clone has the new step ...
#' cp
#'
#' # ... while the original is left unchanged
#' p
#' @export
pip_clone <- function(x, name = NULL) {
    .assert_pip(x)
    if (!is.null(name) && (!.is_single(name, "character") || is.na(name))) {
        stop("name must be a single non-NA string")
    }

    newName <- if (is.null(name)) x[["name"]] else name
    out <- pip_new(name = newName)
    env <- .pip_pipenv(out)

    env[[".dag"]] <- dag_clone(.pip_pipenv(x)[[".dag"]])
    dat <- data.table::copy(.pip_pipenv(x)[["data"]])

    # Clone steps to nodes mapping
    stepsToNodes <- env[[".steps_to_nodes"]]
    for (k in seq_len(nrow(dat))) {
        step <- dat[["step"]][[k]]
        nodeId <- dat[["nodeId"]][[k]]
        stepsToNodes[[step]] <- nodeId
    }
    env[["data"]] <- dat

    out
}


#' Collect step outputs
#'
#' Returns the outputs of the steps in a pipeline or view. With
#' `by = "step"` (the default) the result is a named list of step outputs.
#' With any other column, the outputs are grouped by the values of that
#' column (typically `"tags"`).
#'
#' The result always takes one of two shapes:
#'
#' * **Flat**: a named list whose elements are the step outputs directly,
#'   e.g. `list(s1 = 1, s2 = 2)`.
#' * **Grouped**: a named list whose elements are themselves named lists of
#'   step outputs, one per group, e.g.
#'   `list(io = list(s1 = 1, s2 = 2), model = list(s3 = 4))`.
#'
#' With `simplify = TRUE` (the default) the result is *flat* whenever every
#' group contains exactly one step, which, for example, is always the case
#' for `by = "step"` as step names are unique.
#' If any group contains more than one step (or `simplify = FALSE`) the
#' result is *grouped*. For list columns such as `tags`, a step with
#' several entries contributes its output to every
#' corresponding group, and steps without an entry (e.g. untagged steps)
#' are omitted.
#' You can use [pip_view()] to further narrow the selection before collecting.
#'
#' @param x A pipeflow pip or view.
#' @param by Single step-table column name to group by.
#' @param as.table If TRUE, return a `data.table` instead of a named list.
#' @param simplify If TRUE (default), if the list of collected outputs
#' contains exactly one step per group, the result is flattened by one level,
#' otherwise it is returned as a grouped list.
#' @return A named list of outputs
#' @examples
#' p <- pip_new() |>
#'   pip_add("load", \(x = 1) x, tags = "io") |>
#'   pip_add("clean", \(x = ~load) x + 1, tags = "io") |>
#'   pip_add("model", \(x = ~clean) x * 2, tags = "model")
#' pip_run(p)
#'
#' # By default, a flat named list with one entry per step
#' pip_collect(p)
#'
#' # The same output as a data.table
#' pip_collect(p, as.table = TRUE)
#'
#' # Group the outputs by tag ...
#' pip_collect(p, by = "tags")
#'
#' # ... which is equivalent to
#' list(
#'   io = pip_view(p, tags = "io") |> pip_collect(),
#'   model = pip_view(p, tags = "model") |> pip_collect()
#' )
#'
#' # Grouped table output
#' pip_collect(p, by = "tags", as.table = TRUE)
#'
#' # Keep single-step groups nested
#' pip_collect(p, simplify = FALSE)
#'
#' # Collect output from a view
#' v <- p[step %in% c("clean", "model"), ]
#' pip_collect(v)
#' pip_collect(v, as.table = TRUE)
#' @export
pip_collect <- function(x, by = "step", as.table = FALSE, simplify = TRUE) {
    .assert_pip_or_view(x)
    if (!.is_single(as.table, "logical") || is.na(as.table)) {
        stop("as.table must be a single logical value")
    }
    if (!.is_single(simplify, "logical") || is.na(simplify)) {
        stop("simplify must be a single logical value")
    }
    if (is.null(by)) {
        by <- "step"
    }
    if (!.is_single(by, "character") || is.na(by) || !nzchar(by)) {
        stop("by must be a single column name")
    }

    dat <- .pip_view_data(x)

    if (!(by %in% colnames(dat))) {
        stop("unknown column: ", by)
    }
    # fmt: skip
    allowedBy <- c(
        "step", "state", "tags", "time", "locked", "exec",
        "depends", "unbound", "nodeId"
    )
    if (!by %in% allowedBy) {
        stop(sprintf(
            "'%s' cannot be used as a grouping column - must be one of %s",
            by,
            toString(allowedBy)
        ))
    }

    byCol <- dat[[by]]
    steps <- dat[["step"]]
    outs <- dat[["out"]]

    emptyResult <- function() {
        if (!as.table) {
            return(list())
        }
        key0 <- if (is.list(byCol)) character(0) else byCol[0]
        tbl <- data.table::data.table(grp = key0, out = vector("list", 0L))
        data.table::setnames(tbl, new = c(by, "out"))
        tbl
    }

    if (nrow(dat) == 0L) {
        return(emptyResult())
    }

    tab <- if (is.list(byCol)) {
        # For list columns (like tags) the table is expanded by duplicating
        # rows with multiple entries so that each row contains one entry.
        # For example, tags = c("a", "b") will produce two rows, one for
        # tags = "a" and one for tags = "b".
        keysAll <- byCol
        lens <- lengths(keysAll)
        has <- which(lens > 0L)
        idx <- rep(has, lens[has])
        grp <- if (length(has) > 0L) {
            unlist(keysAll[has], use.names = FALSE)
        } else {
            character(0)
        }
        data.table::data.table(grp = grp, step = steps[idx], out = outs[idx])
    } else {
        # Scalar columns (like step) can stay as-is.
        data.table::data.table(grp = byCol, step = steps, out = outs)
    }

    tab <- tab[!is.na(grp)]
    if (nrow(tab) == 0L) {
        return(emptyResult())
    }

    # Delegate the actual grouping to data.table
    . <- out <- step <- NULL # silence R CMD check
    res <- tab[, .(collect = list(stats::setNames(out, step))), by = grp]

    # Keep the first-appearance order of the groups.
    groupNames <- res[["grp"]]
    first <- match(unique(tab[["grp"]]), groupNames)
    res <- res[first]
    collected <- res[["collect"]]

    # If all groups have a single element ...
    if (simplify && all(lengths(collected) == 1L)) {
        # ... list(grpA = list(x = 1), grpB = list(y = 2)) can be flattened
        # to  list(grpA = 1, grpB = 2)
        collected <- unlist1(collected)
    }

    if (as.table) {
        tbl <- data.table::data.table(grp = groupNames, out = collected)
        data.table::setnames(tbl, new = c(by, "out"))
        tbl
    } else {
        stats::setNames(collected, nm = groupNames)
    }
}


#' @rdname pip_collect
#' @section Lifecycle: Deprecated
#'
#' `pip_collect_out()` is a legacy alias for [pip_collect()]. It raises a
#' deprecation warning and will be removed in a future release.
#' @export
pip_collect_out <- function(
    x,
    by = "step",
    as.table = FALSE,
    simplify = TRUE
) {
    .Deprecated(
        new = "pip_collect",
        old = "pip_collect_out",
        package = "pipeflow"
    )
    pip_collect(x, by = by, as.table = as.table, simplify = simplify)
}


#' Access the underlying step table
#'
#' This is a convenience wrapper for accessing the internal `data.table`,
#' which normally is reachable via `p[["pipenv"]][["data"]]`.
#'
#' @details The internal `data.table` is holding the pipeline steps, one row
#' per step. Unless you know what you are doing, this table should not be
#' modified directly, as this can corrupt the pipeline including any views
#' that are derived from it. If you want to experiment, consider cloning the
#' pipeline first with [pip_clone()].
#'
#' @param x A pipeflow pipeline or view.
#' @return The underlying step table as a `data.table`.
#' @export
#' @examples
#' p <- pip_new() |>
#'   pip_add("load", \(n = 5) seq_len(n)) |>
#'   pip_add("model", \(x = ~load) sum(x))
#'
#' pip_data(p)
pip_data <- function(x) {
    .assert_pip_or_view(x)
    .pip_pipenv(x)[["data"]]
}


#' @name pip_set_params
#' @rdname pip_set_params
#' @export
pip_get_params <- function(x) {
    .assert_pip_or_view(x)
    dat <- .pip_view_data(x)

    params <- mapply(
        par = dat[["params"]],
        unbound = dat[["unbound"]],
        FUN = \(par, unbound) par[unbound],
        SIMPLIFY = FALSE
    ) |>
        Filter(f = \(x) length(x) > 0)

    parNames <- unlist(sapply(params, FUN = names))
    parValues <- stats::setNames(unlist1(params), parNames)
    as.list(parValues[!duplicated(names(parValues))])
}


#' Build pipeline graph data
#'
#' Builds graph data (nodes and edges) describing the pipeline's step
#' structure, suitable for visualisation with [visNetwork::visNetwork()].
#'
#' @details
#' Node shapes reflect execution mode:
#' * `auto`/`plain`: `hexagon`
#' * `reduce`: `dot`
#' * `split`: `star`
#'
#' @param x A pipeflow pip or view.
#' @param include_upstream Logical. Only relevant for views. If `TRUE`, add
#' all upstream dependencies of selected steps.
#'
#' @return A named list with two `data.frame`s: `nodes` and `edges`.
#' @export
#' @examples
#' p <- pip_new()
#' pip_add(p, "load", \(x = 1) x, tags = "io")
#' pip_add(p, "clean", \(x = ~load) x + 1, tags = "io")
#' pip_add(p, "fit", \(x = ~clean) x * 2, tags = "model")
#'
#' graph <- pip_graph(p)
#' graph$nodes # data.frame: id, label, shape, color
#' graph$edges # data.frame: from, to, arrows
#'
#' # For a view, include_upstream = TRUE adds upstream deps to the graph
#' v <- pip_view(p, step = "fit")
#' pip_graph(v, include_upstream = TRUE)
#'
#' if (require("visNetwork", quietly = TRUE)) {
#'   do.call(what = visNetwork::visNetwork, args = graph)
#' }
pip_graph <- function(x, include_upstream = FALSE) {
    .assert_pip_or_view(x)
    if (!.is_single(include_upstream, "logical")) {
        stop("include_upstream must be a single logical value")
    }

    isView <- .is_pipeflow_view(x)
    env <- .pip_pipenv(x)
    dat <- env[["data"]]
    dag <- env[[".dag"]]

    rows <- .pip_view_rows(x)
    rows <- sort(unique(rows))

    if (isView && include_upstream && length(rows) > 0L) {
        startNodes <- dat[["nodeId"]][rows]
        keepNodes <- dag_get_reachable_nodes_up(
            dag,
            as.integer(unique(startNodes))
        )
        rows <- which(dat[["nodeId"]] %in% keepNodes)
        rows <- sort(unique(as.integer(rows)))
    }

    sub <- dat[rows]

    # Nodes
    colors <- vapply(
        sub[["state"]],
        FUN = \(st) .step_states[[st]][["color"]],
        FUN.VALUE = character(1)
    )

    ids <- as.integer(sub[["nodeId"]])
    shape <- rep("hexagon", nrow(sub))
    if ("exec" %in% names(sub)) {
        shape[sub[["exec"]] %in% "split"] <- "star"
        shape[sub[["exec"]] %in% "reduce"] <- "dot"
    }
    nodes <- data.frame(
        id = ids,
        label = sub[["step"]],
        shape = shape,
        color = colors
    )

    # Edges from direct dependencies only (no transitive links).
    stepToId <- stats::setNames(ids, sub[["step"]])
    edgeRows <- lapply(seq_len(nrow(sub)), FUN = \(k) {
        depSteps <- unname(sub[["depends"]][[k]])
        if (length(depSteps) == 0L) {
            return(NULL)
        }

        fromIds <- as.integer(stepToId[depSteps])
        fromIds <- fromIds[!is.na(fromIds)]
        if (length(fromIds) == 0L) {
            return(NULL)
        }

        data.frame(
            from = fromIds,
            to = rep.int(ids[[k]], length(fromIds)),
            stringsAsFactors = FALSE
        )
    })
    edgeRows <- Filter(f = Negate(is.null), x = edgeRows)

    edges <- if (length(edgeRows) == 0L) {
        data.frame(
            from = integer(0),
            to = integer(0),
            arrows = character(0),
            stringsAsFactors = FALSE
        )
    } else {
        edges <- do.call(what = rbind, args = edgeRows)
        edges <- unique(edges)
        cbind(edges, "arrows" = "to")
    }

    list(nodes = nodes, edges = edges)
}


#' @rdname pip_graph
#' @section Lifecycle: Deprecated
#'
#' `pip_get_graph()` is a legacy alias for [pip_graph()]. It raises a
#' deprecation warning and will be removed in a future release.
#' @export
pip_get_graph <- function(x, include_upstream = FALSE) {
    .Deprecated(new = "pip_graph", old = "pip_get_graph", package = "pipeflow")
    pip_graph(x, include_upstream = include_upstream)
}


#' Remove a step
#'
#' Removes a pipeline step by its name. If other steps depend on it,
#' an error is given and the removal of the step is blocked, unless
#' `force` was set to `TRUE`, which will remove the selected step
#' together with all its downstream dependent steps.
#'
#' @note
#' If called on a view, the step to be removed must be part of the view,
#' since removing steps drops rows from the pipeline, the row positions
#' of the view will be shifted after the removal operation.
#' For views, you therefore need to assign the result back to the view:
#' `v <- pip_remove(v, ...)`.
#' Views that are not reassigned (or other views on the same pipeline),
#' keep their old row positions and therefore can become invalid after a
#' step removal.
#'
#' @param x A pipeflow pip or view
#' @param step `string` the name of the step to be removed.
#' @param force `logical` if `TRUE` the step is removed together
#' with all its downstream dependencies.
#' @return The updated pipeline, invisibly. For a view, the view with its
#' remapped selector.
#' @examples
#' p <- pip_new() |>
#'   pip_add("load", \(x = 1) x) |>
#'   pip_add("transform", \(x = ~load) x * 2) |>
#'   pip_add("model", \(x = ~transform) x + 10)
#'
#' # Removing a leaf step (nothing depends on it) works directly
#' pip_remove(p, "model")
#' p                        # "load", "transform"
#'
#' # Trying to remove a step that others depend on raises an error:
#' # pip_remove(p, "load")  # Error!
#'
#' # If a view is passed, the step must be part of the view.
#' v <- pip_view(p, step = "transform")
#' try(pip_remove(v, "load"))  # Error: "load" is not part of the view
#' v <- pip_remove(v, "transform")
#' v[["step"]]              # view is remapped and stays valid
#' p                        # "load"
#'
#' # force = TRUE removes the step and all its downstream dependents
#' pip_remove(p, "load", force = TRUE)
#' p                        # pipeline is now empty
#' @export
pip_remove <- function(x, step, force = FALSE) {
    .assert_pip_or_view(x)
    if (!.is_single(step, "character")) {
        stop("step must be a single string")
    }
    if (is.na(step)) {
        stop("step must not be NA")
    }
    if (!.pip_step_exists(x, step)) {
        stop("step '", step, "' does not exist")
    }
    if (!.pip_step_in_view(x, step)) {
        stop("step '", step, "' is not part of the view")
    }
    if (!is.logical(force) || length(force) != 1L || is.na(force)) {
        stop("force must be a single logical value")
    }

    pipenv <- .pip_pipenv(x)
    dat <- pipenv[["data"]]
    `%chin%` <- data.table::`%chin%`

    # For views, remember the covered steps so that the view can be remapped.
    isView <- .is_pipeflow_view(x)
    if (isView) {
        viewSteps <- dat[["step"]][.pip_view_rows(x)]
    }

    directDeps <- dat[["step"]][
        vapply(
            dat[["depends"]],
            FUN = \(dep) step %chin% dep,
            FUN.VALUE = logical(1)
        )
    ]

    if (length(directDeps) > 0L && !force) {
        stepsString <- paste0("'", directDeps, "'", collapse = ", ")
        stop(
            "cannot remove step '",
            step,
            "' because the following steps depend on it: ",
            stepsString
        )
    }

    stepsToRemove <- step
    if (force) {
        downNodes <- .pip_get_reachable_nodes(x, step)
        downNodes <- as.integer(downNodes)
        stepNode <- as.integer(.pip_steps_to_nodes(x, step)[[1]])

        downDeps <- dat[["step"]][
            dat[["nodeId"]] %in% setdiff(downNodes, stepNode)
        ]
        if (length(downDeps) > 0L) {
            stepsString <- paste0("'", downDeps, "'", collapse = ", ")
            message(
                "Removing step '",
                step,
                "' and its downstream dependencies: ",
                stepsString
            )
        }

        stepsToRemove <- dat[["step"]][dat[["nodeId"]] %in% downNodes]
    }

    nodesToRemove <- as.integer(unname(unlist(
        .pip_steps_to_nodes(x, stepsToRemove)
    )))

    # Remove DAG nodes first to keep node references stable during filtering.
    for (nid in rev(nodesToRemove)) {
        ok <- dag_remove_node(pipenv[[".dag"]], nid, force = force)
        if (!ok) {
            stop("failed to remove node ", nid, " from DAG")
        }
    }

    dag_tidy_up(pipenv[[".dag"]])
    keep <- !(dat[["step"]] %chin% stepsToRemove)
    pipenv[["data"]] <- dat[keep]
    suppressWarnings(
        rm(
            list = stepsToRemove,
            envir = pipenv[[".steps_to_nodes"]],
            inherits = FALSE
        )
    )

    data.table::setindexv(pipenv[["data"]], list("step", "nodeId"))

    # Remap the view to the shifted (and potentially dropped) row positions.
    if (isView) {
        remaining <- pipenv[["data"]][["step"]]
        newRows <- as.integer(which(remaining %chin% viewSteps))
        x <- .wrap_pipenv(pipenv, name = x[["name"]], view = newRows)
    }
    invisible(x)
}


#' Rename a step
#'
#' Renames the selected step and updates dependency references in
#' downstream steps.
#' @param x A pipeflow pip or view
#' @param from Existing step name
#' @param to New step name
#' @return The updated pipeline, invisibly.
#' @examples
#' p <- pip_new() |>
#'   pip_add("s1", \(x = 1) x) |>
#'   pip_add("s2", \(x = ~s1) x + 1)           # "s2" depends on "s1"
#'
#' # Downstream dependency references are updated automatically
#' pip_rename(p, from = "s1", to = "load_data")
#' p
#'
#' # Trying to rename to an existing step name raises an error:
#' try(pip_rename(p, "load_data", to = "s2"))  # step 's2' already exists!
#'
#' # If a view is passed, the step must be part of the view
#' v <- pip_view(p, step = c("load_data", "s2"))
#' pip_rename(v, from = "load_data", to = "input")
#' p[["step"]]                                 # "input", "s2"

#' v2 <- pip_view(p, step = "s2")
#' try(pip_rename(v2, from = "input", to = "data"))
#' @export
pip_rename <- function(x, from, to) {
    .assert_pip_or_view(x)

    if (!.is_single(from, "character")) {
        stop("from must be a single string")
    }
    if (is.na(from)) {
        stop("from must not be NA")
    }
    if (!nzchar(from)) {
        stop("from must be a non-empty string")
    }

    if (!.is_single(to, "character")) {
        stop("to must be a single string")
    }
    if (is.na(to)) {
        stop("to must not be NA")
    }
    if (!nzchar(to)) {
        stop("to must be a non-empty string")
    }

    if (!.pip_step_exists(x, from)) {
        stop("step '", from, "' does not exist")
    }
    if (!.pip_step_in_view(x, from)) {
        stop("step '", from, "' is not part of the view")
    }
    if (identical(from, to)) {
        return(invisible(x))
    }
    if (.pip_step_exists(x, to)) {
        stop("step '", to, "' already exists")
    }

    env <- .pip_pipenv(x)
    dat <- env[["data"]]
    `%chin%` <- data.table::`%chin%`
    newSteps <- dat[["step"]]
    newSteps[newSteps %chin% from] <- to

    newDepends <- lapply(
        dat[["depends"]],
        FUN = \(dep) {
            if (length(dep) == 0L) {
                return(dep)
            }
            dep[dep %chin% from] <- to
            dep
        }
    )

    # Keep the formula references in params in sync with the renamed steps.
    # The formulas are informational (step resolution uses `depends`), but
    # they are user visible (e.g. when printing the params column).
    newParams <- dat[["params"]]
    for (i in seq_along(newParams)) {
        dep <- newDepends[[i]]
        if (length(dep) == 0L) {
            next
        }
        for (arg in names(dep)) {
            fml <- newParams[[i]][[arg]]
            if (inherits(fml, "formula")) {
                # Only swap the referenced name to keep the formula's
                # environment (as.formula() would attach this frame).
                fml[[2L]] <- as.name(dep[[arg]])
                newParams[[i]][[arg]] <- fml
            }
        }
    }

    data.table::set(dat, j = "step", value = newSteps)
    data.table::set(dat, j = "depends", value = newDepends)
    data.table::set(dat, j = "params", value = list(newParams))

    stepsToNodes <- env[[".steps_to_nodes"]]
    nodeId <- stepsToNodes[[from]]
    stepsToNodes[[to]] <- nodeId
    rm(list = from, envir = stepsToNodes, inherits = FALSE)

    data.table::setindexv(dat, list("step", "nodeId"))
    invisible(x)
}


#' Replace a step
#'
#' Replaces a step's function while keeping it in the same position in the
#' pipeline. Downstream steps are automatically marked as outdated and will
#' re-run on the next [pip_run()].
#'
#' @inheritParams pip_add
#' @param x A pipeflow pipeline or view object.
#' @return The updated pipeline, invisibly.
#' @examples
#' p <- pip_new() |>
#'     pip_add("load", \(n = 5) seq_len(n)) |>
#'     pip_add("double", \(x = ~load) x * 2)
#' pip_run(p)
#' p
#'
#' # Replace "load" — downstream steps are automatically marked "outdated"
#' pip_replace(p, "load", \(n = 3) seq_len(n))
#' p
#'
#' # Re-run to bring everything up to date
#' pip_run(p)
#' p
#'
#' # If a view is passed, the step must be part of the view
#' v <- pip_view(p, step = "double")
#' pip_replace(v, "double", \(x = ~load) x * 3)
#' p[["state"]]                                # "double" is "new" again
#'
#' try(pip_replace(v, "load", \(n = 2) seq_len(n)))
#' @export
pip_replace <- function(
    x,
    step,
    fun,
    tags = character(0),
    params = list(),
    exec = "auto"
) {
    .assert_pip_or_view(x)
    if (!.is_single(step, "character")) {
        stop("step must be a single string")
    }
    if (is.na(step)) {
        stop("step must not be NA")
    }
    if (!nzchar(step)) {
        stop("step must be a non-empty string")
    }
    if (!.pip_step_exists(x, step)) {
        stop("step '", step, "' does not exist")
    }
    if (!.pip_step_in_view(x, step)) {
        stop("step '", step, "' is not part of the view")
    }
    if (!is.function(fun)) {
        stop("fun must be a function")
    }
    if (".self" %in% names(formals(fun))) {
        stop_no_call(
            "'.self' is a reserved parameter name and must not be declared ",
            "in step '",
            step,
            "' - it is provided automatically"
        )
    }
    if (".self" %in% names(params)) {
        stop_no_call(
            "'.self' is a reserved parameter name and must not be set via ",
            "params - it is provided automatically"
        )
    }
    .assert_exec_mode(exec)

    env <- .pip_pipenv(x)
    data <- env[["data"]]
    iStep <- data.table::chmatch(step, data[["step"]])

    funParams <- .extract_fun_params(fun)
    params[names(funParams)] <- funParams

    # Provide `.self` to the step via a dedicated wrapper environment.
    fun <- .wrap_self(fun, x)

    # Dependencies of the replacement may only point to earlier steps
    prefixSteps <- data[["step"]][seq_len(iStep - 1L)]
    steps <- c(prefixSteps, step)
    depends <- .extract_depends(params = params, steps = steps)
    if (length(depends) > 0L) {
        bad <- depends[!(depends %in% prefixSteps)]
        if (length(bad) > 0L) {
            stop_no_call(
                "while replacing step '",
                step,
                "' - cannot reference unknown steps: ",
                paste0("'", bad, "'", collapse = ", ")
            )
        }
    }

    refNodes <- mget(
        depends,
        envir = env[[".steps_to_nodes"]],
        ifnotfound = NA_integer_,
        inherits = FALSE
    )
    if (anyNA(refNodes)) {
        notFound <- Filter(is.na, refNodes)
        stop_no_call(
            "while adding step '",
            step,
            "' - cannot reference unknown steps: ",
            paste0("'", notFound, "'", collapse = ", ")
        )
    }

    # Update the incoming DAG edges of the replaced step (its node and all
    # outgoing edges stay the same).
    d <- env[[".dag"]]
    nodeId <- env[[".steps_to_nodes"]][[step]]
    oldDeps <- data[["depends"]][[iStep]]
    oldRefs <- if (length(oldDeps) > 0L) {
        as.integer(unlist(mget(
            unname(oldDeps),
            envir = env[[".steps_to_nodes"]],
            inherits = FALSE
        )))
    } else {
        integer()
    }
    newRefs <- as.integer(refNodes)

    for (from in setdiff(oldRefs, newRefs)) {
        dag_remove_edge(d, from = from, to = nodeId, force = TRUE)
    }
    toAdd <- setdiff(newRefs, oldRefs)
    if (length(toAdd) > 0L) {
        dag_add_edges_to(d, from = toAdd, to = nodeId)
    }

    # Reset the step row: new function, params, tags and exec, fresh runtime
    # state (like a freshly added step).
    values <- list(
        "fun" = list(fun),
        "params" = list(params),
        "out" = list(NULL),
        "state" = .step_states[["new"]][["name"]],
        "tags" = list(tags),
        "locked" = FALSE,
        "exec" = exec,
        "time" = Sys.time(),
        "depends" = list(depends),
        "unbound" = list(setdiff(names(params), names(depends)))
    )
    data.table::set(data, i = iStep, j = names(values), value = values)

    # Mark downstream dependent steps as outdated, but keep the replaced
    # step itself as "new".
    downNodes <- .pip_get_reachable_nodes(x, step)
    downNodes <- unique(setdiff(
        as.integer(unlist(downNodes)),
        as.integer(nodeId)
    ))
    if (length(downNodes) > 0L) {
        rowsDown <- data[
            list(downNodes),
            which = TRUE,
            on = "nodeId"
        ]
        if (length(rowsDown) > 0L) {
            data.table::set(
                data,
                i = rowsDown,
                j = "state",
                value = .step_states[["outdated"]][["name"]]
            )
        }
    }

    invisible(x)
}


#' Run a pipeline
#'
#' Executes all pending steps in order. On repeated runs, steps that are
#' already `"done"` are skipped and only steps that are still `"new"` or were
#' marked `"outdated"` (because one of their dependencies changed) are
#' executed. Use `force = TRUE` to re-execute every step.
#'
#' @param x A pipeflow pip or view
#' @param lgr A logging function of the form `function(level, msg, ...)`.
#' To suppress logging, you can set `lgr = NULL`.
#' @param force Logical indicating if all steps should be forced to run,
#' regardless of whether they are outdated or not.
#' @param progress Optional callback of the form
#' `function(value, detail)` called before each step.
#' @return The updated pipeline or view, invisibly.
#'
#' @details
#'
#' ## Step states
#'
#' A step can take the following states:
#'
#' * new: the step has been added but not yet executed, or was reset via
#'   [pip_reset()].
#' * outdated: a step can get outdated in the following scenarios:
#'   - one of the step's parameters were changed (e.g. via [pip_set_params()])
#'   - a step it depends on was re-executed or replaced
#'   - the last run did not reach it either on purpose (see section 'Running
#'     views below) or because a run was aborted early due to failed step
#' * done: the step was executed successfully with its current inputs
#' * failed: the step raised an error during its last execution
#'
#' A "done" step is skipped unless `force = TRUE` was set. For all other
#' states, the step will be re-executed in the next [pip_run()].
#'
#' ## Running views
#' When `x` is a view, the requested rows are run together with their
#' upstream dependencies, so the steps covered by the view are brought up to
#' date even if their inputs come from steps outside the view. The rest of
#' the pipeline is not executed; downstream steps that were not processed
#' are marked `"outdated"`.
#'
#' ## Runtime errors
#' If a step fails with an error, the failing step's state is set to
#' `"failed"`, the run is aborted (no further steps are executed), and the
#' pipeline run state is set to `"failed"`. Steps that have not been
#' executed are marked `"outdated"`, so a subsequent run retries them.
#'
#' ## Runtime control flow via restart and halt
#'
#' A running pipeline can be interrupted via the `restart()` and `halt()`
#' virtual methods. They are intended for advanced, self-modifying pipelines
#' and are most often called from within a step function via the `.self`
#' argument.
#'
#' - `.self$restart(force = TRUE, times = 1L)`: aborts the current run after the
#'   current step has finished and restarts it from the first step.
#'   The default parameters are `force = TRUE` and `times = 1L`, that is, the
#'   above call is the same as just invoking .self$restart().
#'   To skip steps that are already in state `"done"`, set `force = FALSE`.
#'   The `times` parameter limits the number of consecutive restarts within a
#'   single `pip_run()` call.
#'   If a view is being run, a restart covers the view steps together with
#'   their upstream dependencies.
#' - `p$halt()`: aborts the current run after the current step has finished.
#'   This is a *controlled halt* and deliberately distinct from [base::stop()]:
#'   no error is raised and the pipeline is not marked as `"failed"`. The run
#'   simply ends, the steps that have not been executed are marked
#'   `"outdated"`, and a subsequent [pip_run()] will continue where the run
#'   left off.
#'
#' In both cases steps that have not been executed until the restart or halt
#' happens are marked as `"outdated"`.
#'
#' @seealso `vignette("v06-self-modify-pipeline", package = "pipeflow")`
#'   for an advanced example of dynamic pipelines.
#' @examples
#' p <- pip_new() |>
#' pip_add("load", \(n = 3) seq_len(n)) |>
#'   pip_add("prep", \(x = ~load, weight = 1) x * weight) |>
#'   pip_add("square", \(x = ~prep) x^2) |>
#'   pip_add("total", \(x = ~square) sum(x))
#'
#' pip_run(p)
#' p
#'
#' pip_set_params(p, list(weight = 2))
#' p
#'
#' # Already-done steps are skipped on a second run
#' pip_run(p) # first step skipped
#'
#' # lgr = NULL suppresses log output
#' pip_run(p, lgr = NULL)
#'
#' # force = TRUE re-executes every step regardless of state
#' pip_run(p, force = TRUE)
#' p
#'
#' # Run only a subset of steps via a view;
#' # upstream dependencies are automatically included
#' v <- pip_view(p, step = "total")
#' pip_run(v)
#'
#' # Halt or restart pipeline at runtime (for advanced usage)
#' p <- pip_new("restart") |>
#'   pip_add("load", \(n = 3) seq_len(n)) |>
#'   pip_add("check", \(x = ~load) {
#'     if (length(x) > 10L) .self$halt()
#'   }) |>
#'   pip_add("model", \(x = ~load) {
#'     if (length(x) == 3L) {
#'       .self$set_params(list(n = 5))
#'       .self$restart()
#'     }
#'     x * 2
#'   })
#'
#' pip_run(p)
#' p
#'
#' pip_set_params(p, list(n = 15)) # now halt() in 'check' step is triggered
#' pip_run(p)
#' @export
pip_run <- function(
    x,
    lgr = pipeflow_lgr,
    force = FALSE,
    progress = NULL
) {
    .assert_pip_or_view(x)
    if (!.is_single(force, "logical")) {
        stop("force must be a single logical value")
    }
    if (!is.null(progress) && !is.function(progress)) {
        stop("progress must be a function")
    }
    if (is.null(lgr)) {
        lgr <- function(level, msg, ...) {}
    } else {
        .assert_logger(lgr)
    }
    log_info <- function(msg) lgr(level = "info", msg = msg)

    isView <- .is_pipeflow_view(x)
    pipenv <- .pip_pipenv(x)
    pipname <- x[["name"]]
    dat <- pipenv[["data"]]
    rowsToRun <- seq_len(nrow(dat))

    if (isView) {
        requested <- .pip_view_rows(x)
        reqSteps <- dat[["step"]][requested]
        upNodes <- .pip_get_reachable_nodes(
            x,
            reqSteps,
            downstream = FALSE
        )
        upRows <- as.integer(dat[list(upNodes), which = TRUE, on = "nodeId"])
        rowsToRun <- as.integer(sort(unique(c(requested, upRows))))
        upstreamRows <- setdiff(rowsToRun, requested)
        names(rowsToRun)[match(requested, rowsToRun)] <- "view"
        names(rowsToRun)[match(upstreamRows, rowsToRun)] <- "upstream"
    }
    processedSteps <- character()
    restartDelegated <- FALSE
    on.exit({
        # At the end, mark all downstream dependent steps as outdated that
        # were *not* processed, which can happen in two different ways:
        # a) when running a view that does not cover the entire pipeline or
        # b) the run was aborted in the middle (due to an error or manual halt).
        # When the run was restarted, a nested pip_run() has already handled
        # the whole pipeline (including marking), so nothing to do here.
        if (!restartDelegated) {
            processedNodes <- as.integer(.pip_steps_to_nodes(x, processedSteps))
            outdatedNodes <- .pip_get_reachable_nodes(x, processedSteps) |>
                unlist() |>
                unique() |>
                setdiff(processedNodes)
            if (length(outdatedNodes) > 0L) {
                iOut <- dat[list(outdatedNodes), which = TRUE, on = "nodeId"]
                if (length(iOut) > 0L) {
                    data.table::set(
                        dat,
                        i = iOut,
                        j = "state",
                        value = "outdated"
                    )
                }
            }
        }
    })

    state <- pipenv[[".run_state"]]
    action <- if (state == "restart") "Restarting" else "Starting"
    log_info(sprintf("%s run of %s '%s'", action, data.class(x), pipname))
    pipenv[[".run_state"]][] <- "running"
    tryCatch(
        {
            for (i in seq_along(rowsToRun)) {
                row <- rowsToRun[[i]]
                step <- dat[["step"]][[row]]
                processedSteps <- c(processedSteps, step)
                if (!is.null(progress)) {
                    progress(value = i, detail = step)
                }
                msg <- if (isView) {
                    marker <- names(rowsToRun)[[i]]
                    sprintf(
                        "Step %i/%i [%s] %s",
                        i,
                        length(rowsToRun),
                        marker,
                        step
                    )
                } else {
                    sprintf("Step %i/%i %s", i, length(rowsToRun), step)
                }

                if (dat[["state"]][[row]] == "done" && !force) {
                    log_info(sprintf("%s - skipping done step", msg))
                    next()
                }
                if (dat[["locked"]][[row]]) {
                    log_info(sprintf("%s - skipping locked step", msg))
                    next()
                }

                # Always pass the full pipeline object (not a view) to the
                # step function, so that it can modify itself if needed.
                self <- .wrap_pipenv(pipenv, pipname, view = NULL)

                log_info(msg)
                .pip_run_row(x = self, i = row, lgr = lgr)
                stateAfterStep <- pipenv[[".run_state"]]

                # Check for restart or stop signals
                if (stateAfterStep == "restart") {
                    log_info("Restarting pipeline execution.")
                    doForce <- pipenv[[".restart_force"]]
                    restartDelegated <- TRUE
                    pip_run(
                        x,
                        lgr = lgr,
                        force = doForce,
                        progress = progress
                    )
                    return(invisible(x))
                }

                if (stateAfterStep == "halted") {
                    log_info("Aborting pipeline execution on manual halt.")
                    break
                }
            }

            log_info(
                sprintf("Finished run of %s '%s'", data.class(x), pipname)
            )
            pipenv[[".run_state"]][] <- "ready"
            pipenv[[".last_run"]] <- Sys.time()
            invisible(x)
        },
        error = function(e) {
            pipenv[[".run_state"]][] <- "failed"
            pipenv[[".last_run"]] <- Sys.time()
            stop_no_call(e$message)
        }
    )
}


#' Reset a pipeline to its initial state
#'
#' Resets all unlocked steps of a pipeline (or a subset of steps defined by a
#' view) to state `"new"` and clears their outputs, so a subsequent
#' [pip_run()] re-executes the cleaned steps from scratch.
#' Any pending restart counters are cleared, but parameters, tags, and
#' locked flags are left unchanged.
#'
#' @details Locked steps are skipped: their state and output are preserved.
#' If all selected steps are locked, a warning is issued and nothing is
#' changed.
#'
#' @param x A pipeflow pip or view. If a view is given, only the steps covered
#' by the view are reset.
#'
#' @return The updated pipeline or view, invisibly.
#' @examples
#' p <- pip_new() |>
#'   pip_add("load", \(n = 3) seq_len(n)) |>
#'   pip_add("square", \(x = ~load) x^2)
#'
#' pip_run(p)
#' p[["state"]] # "done", "done"
#'
#' # Locked steps keep their state and output when resetting
#' pip_lock(pip_view(p, step = "square"))
#' pip_reset(p)
#' p
#' p[["state"]] # "new", "done"
#' p[["out"]]   # NULL, (x^2 result)
#'
#' pip_unlock(p)
#' pip_reset(p)
#' p
#' p[["state"]] # "new", "new"
#' p[["out"]]   # NULL, NULL
#' @export
pip_reset <- function(x) {
    .assert_pip_or_view(x)
    env <- .pip_pipenv(x)
    dat <- env[["data"]]

    rows <- .pip_view_rows(x)
    if (length(rows) == 0L) {
        return(invisible(x))
    }
    rowsConsidered <- setdiff(rows, which(dat[["locked"]]))

    if (length(rowsConsidered) == 0L) {
        message("No steps to update: all selected steps are locked")
        return(invisible(x))
    }

    data.table::set(
        dat,
        i = rowsConsidered,
        j = c("out", "state"),
        value = list(
            rep(list(NULL), length(rowsConsidered)),
            rep(.step_states[["new"]][["name"]], length(rowsConsidered))
        )
    )

    env[[".run_state"]][] <- "ready"
    env[[".restart_count"]] <- 0L
    env[[".last_run"]] <- NULL

    invisible(x)
}


#' Get or set unbound parameters
#'
#' [pip_get_params()] returns the current default values of all unbound
#' parameters of a pipeline or view, that is, parameters wired to another
#' step's output via `~step_name` are excluded.
#' [pip_set_params()] updates these parameters for the whole pipeline or
#' or view and marks the affected steps and their downstream dependents
#' as outdated.
#'
#' @note Parameters of locked steps are never changed and their state
#' remains unchanged.
#'
#' @param x A pipeflow pip or view.
#' @param params Named list of parameters to set (only used by
#' [pip_set_params()]).
#' @return For [pip_get_params()], a named list of unbound parameter values;
#' if the same parameter name appears in multiple steps, the first
#' occurrence in pipeline order is returned. For [pip_set_params()], the
#' updated pipeline or view, invisibly.
#' @examples
#' p <- pip_new() |>
#'   pip_add("load", \(n = 10) seq_len(n)) |>
#'   pip_add("scale", \(x = ~load, factor = 0.5) x * factor)
#'
#' # See all unbound/adjustable parameters before running
#' pip_get_params(p) # list(n = 10, factor = 0.5)
#' (pip_run(p))
#'
#' # Updating params marks affected steps (and their dependents) outdated
#' pip_set_params(p, params = list(n = 5, factor = 2.0))
#' p
#' (pip_run(p))
#'
#' # Setting a parameter that is not defined in the pipeline yields a warning
#' \donttest{
#' pip_set_params(p, params = list(nope = 1))
#' }
#' @export
pip_set_params <- function(x, params = list()) {
    # Input checking
    .assert_pip_or_view(x)
    if (!is.list(params)) {
        stop("params must be a list")
    }
    parNames <- names(params)
    if (length(params) == 0) {
        return(invisible(x))
    }
    allNamed <- length(parNames) == length(params) && all(nzchar(parNames))
    if (!allNamed) {
        stop("All parameters must be named")
    }

    # Narrow down the considered rows
    env <- .pip_pipenv(x)
    dat <- env[["data"]]
    rows <- .pip_view_rows(x)
    rowsConsidered <- setdiff(rows, which(dat[["locked"]]))

    if (length(rowsConsidered) == 0L) {
        message("No steps to update: all selected steps are locked")
        return(invisible(x))
    }

    # Determine which steps/rows are affected, i.e. have intersecting params
    unbound <- dat[["unbound"]][rowsConsidered] # names of unbound params
    intersects <- lapply(unbound, FUN = intersect, y = parNames)
    hasOverlap <- lengths(intersects) > 0
    namesAffected <- intersects[hasOverlap]
    rowsAffected <- rowsConsidered[hasOverlap]

    # Signal parameters that are not defined in any of the affected steps
    used <- unique(unlist(intersects))
    undefined <- setdiff(parNames, used)
    if (length(undefined) > 0L) {
        warning(
            "Trying to set parameters not defined in the target: ",
            toString(undefined)
        )
    }

    if (any(hasOverlap)) {
        # Update parameters in all affected rows
        for (j in seq_along(rowsAffected)) {
            i <- rowsAffected[[j]]
            names <- namesAffected[[j]]
            rowPars <- dat[["params"]][[i]]
            rowPars[names] <- params[names]
            value <- list(list(rowPars)) # need to wrap in list() for call below
            data.table::set(dat, i = i, j = "params", value = value)
        }

        # Update states of affected steps and their downstream steps
        steps <- dat[["step"]][rowsAffected]
        .pip_update_downstream(
            x,
            steps = steps,
            what = "state",
            value = "outdated"
        )
    }

    invisible(x)
}


#' Add or remove tags
#'
#' Adds tags to, or removes tags from, all steps of a pipeline or a subset
#' of steps defined by a view. Tagged steps can later be selected via
#' [pip_view()].
#' @param x A pipeflow pip or view.
#' @param tags Character vector of tags to add to or remove from each
#' selected step.
#' @return The updated pipeline or view, invisibly.
#' @examples
#' p <- pip_new() |>
#'   pip_add("load", \(x = 1) x) |>
#'   pip_add("fit", \(x = ~load) x + 1)
#'
#' # Tag every step in the pipeline at once
#' pip_tag(p, c("daily", "core"))
#' p
#' p[, "tags"] # both steps have c("daily", "core")
#'
#' # Add an extra tag to only one step via a view
#' p[step == "fit"] |> pip_tag("model")
#' p
#' p[, "tags"] # "fit" also has "model"
#'
#' # Remove "daily" from all steps
#' pip_untag(p, "daily")
#' p[, "tags"]
#' @export
pip_tag <- function(x, tags = character()) {
    .assert_pip_or_view(x)
    if (!is.character(tags)) {
        stop("tags must be a character vector")
    }

    env <- .pip_pipenv(x)
    dat <- env[["data"]]
    rows <- .pip_view_rows(x)

    if (length(rows) == 0L || length(tags) == 0L) {
        return(invisible(x))
    }

    for (i in rows) {
        if (isTRUE(dat[["locked"]][[i]])) {
            next
        }

        oldTags <- dat[["tags"]][[i]]
        newTags <- unique(c(oldTags, tags))
        data.table::set(dat, i = i, j = "tags", value = list(list(newTags)))
    }

    invisible(x)
}


#' @rdname pip_tag
#' @export
pip_untag <- function(x, tags = character()) {
    .assert_pip_or_view(x)
    if (!is.character(tags)) {
        stop("tags must be a character vector")
    }

    env <- .pip_pipenv(x)
    dat <- env[["data"]]
    rows <- .pip_view_rows(x)

    if (length(rows) == 0L || length(tags) == 0L) {
        return(invisible(x))
    }

    for (i in rows) {
        if (isTRUE(dat[["locked"]][[i]])) {
            next
        }

        oldTags <- dat[["tags"]][[i]]
        newTags <- setdiff(oldTags, tags)
        data.table::set(dat, i = i, j = "tags", value = list(list(newTags)))
    }

    invisible(x)
}


#' Lock or unlock steps
#'
#' Locks or unlocks all steps of a pipeline, or a subset of steps defined by
#' a view. Locked steps are skipped during [pip_run()] and are protected
#' against modification from [pip_set_params()], [pip_tag()] or
#' [pip_untag()]. Calling [pip_unlock()] removes the lock again.
#' @param x A pipeflow pip or view.
#' @return The updated pipeline or view, invisibly.
#' @examples
#' p <- pip_new() |>
#'   pip_add("x", \(x = 1) x) |>
#'   pip_add("y", \(y = 2) y) |>
#'   pip_add("sum", \(x = 1, y = 2) x + y)
#' (pip_run(p))
#' p[["sum", "out"]] # 3
#'
#' # Lock "sum" step via a view so it cannot be overwritten
#' pip_set_params(p, params = list(x = 10, y = 20))
#' p[["sum", "params"]] # x = 10, y = 20
#' pip_lock(p["sum", ])
#' (pip_run(p))
#' p[["sum", "out"]] # still 3
#'
#' # Note that locking also prevents any parameter updates
#' pip_set_params(p, params = list(x = 100, y = 200))
#' p[["x", "params"]] # x = 100
#' p[["y", "params"]] # y = 200
#' p[["sum", "params"]] # still x = 10, y = 20
#'
#' # Unlock everything to allow updates again
#' pip_unlock(p)
#' pip_set_params(p, params = list(x = 100, y = 200))
#' (pip_run(p))
#' @export
pip_lock <- function(x) {
    .assert_pip_or_view(x)

    env <- .pip_pipenv(x)
    dat <- env[["data"]]
    rows <- .pip_view_rows(x)

    if (length(rows) == 0L) {
        return(invisible(x))
    }

    data.table::set(dat, i = rows, j = "locked", value = TRUE)
    invisible(x)
}


#' @rdname pip_lock
#' @export
pip_unlock <- function(x) {
    .assert_pip_or_view(x)

    env <- .pip_pipenv(x)
    dat <- env[["data"]]
    rows <- .pip_view_rows(x)

    if (length(rows) == 0L) {
        return(invisible(x))
    }

    data.table::set(dat, i = rows, j = "locked", value = FALSE)
    invisible(x)
}


#' Create a pipeline view
#'
#' Creates a filtered view showing only a selected subset of steps.
#' A view references the underlying pipeline without copying it, so
#' operations like [pip_run()] and [pip_set_params()] applied to a view
#' work directly on the underlying pipeline, but are restricted to the
#' the steps defined by the view.
#'
#' @param x A pipeflow pipeline or view.
#' @param ... Named filters, which can be one or more of `step`, `params`,
#' `state`, `tags`, `exec`, and `depends`. Each filter value is a character
#' vector of values to keep, or - if `fixed` is `FALSE` - a regular
#' expression. The `params` filter matches against the actual parameter
#' names of each step.
#' @param join How individual filters are combined:
#' * `"intersect"` (the default) keeps steps that match *all* filters,
#' * `"union"` keeps steps that match *any* filter. Within a single
#'    filter, multiple values are always treated as alternatives (OR).
#' @param fixed If TRUE, values in `...` are treated as fixed strings,
#' otherwise they are treated as regular expressions.
#'
#' @return A `pipeflow_view` object.
#' @export
#' @examples
#' p <- pip_new() |>
#'   pip_add("load", \(a = 1) a, tags = c("io", "core", "daily")) |>
#'   pip_add("fit", \(b = 2) b + 1, tags = c("model")) |>
#'   pip_add("eval_fit", \(fit = ~fit) fit,
#'     tags = c("model", "daily", "report")
#'   )
#' p
#'
#' # Filter by one or more column values
#' pip_view(p, state = "new")
#' pip_view(p, step = c("load", "fit"))
#'
#' # Filter by tag — keeps steps that have *any* of the given tags
#' pip_view(p, tags = "daily")
#'
#' # Combine filters: step pattern AND state (logical AND)
#' pip_view(p, step = "fit", state = "new")
#'
#' # Combine filters as a union (step OR state)
#' pip_view(p, step = "load", tags = "report", join = "union")
#'
#' # Use a regex pattern
#' pip_view(p, step = "fit$", fixed = FALSE)
#'
#' # Filter by parameter names — steps with any of the given parameters
#' pip_view(p, params = c("a", "fit"))
#'
#' # Views are composable: create a view-of-view for progressive narrowing
#' v1 <- pip_view(p, tags = "daily")
#' print(v1) # load, eval_fit
#' v2 <- pip_view(v1, tags = "report")
#' print(v2) # eval_fit only
pip_view <- function(x, ..., join = c("intersect", "union"), fixed = TRUE) {
    .assert_pip_or_view(x)
    join <- match.arg(join)
    if (!.is_single(fixed, "logical")) {
        stop("fixed must be a single logical value")
    }
    env <- x[["pipenv"]]
    dat <- env[["data"]]

    filters <- list(...)
    validFilters <- c("step", "params", "depends", "state", "tags", "exec")
    unknown <- setdiff(names(filters), validFilters)
    if (length(unknown) > 0) {
        stop(sprintf(
            "Invalid filter name: '%s' - can be one of: %s",
            unknown[[1]],
            paste(validFilters, collapse = ", ")
        ))
    }

    # For view-of-view, filter only within parent view rows and map local
    # matches back to absolute row indices of the underlying pipeline.
    parent_rows <- .pip_view_rows(x)
    sub <- dat[parent_rows]

    # Identity element for the join: TRUE for intersect, FALSE for union.
    keep <- rep(join == "intersect", nrow(sub))

    # Resolve each filter name to the column it filters on.
    for (name in names(filters)) {
        values <- filters[[name]]
        if (!is.character(values)) {
            stop(sprintf(
                "filter '%s' must be a character vector, not %s",
                name,
                typeof(values)
            ))
        }
        col <- if (name == "params") {
            # Special case "params": it matches against the parameter *names*
            lapply(sub[["params"]], names)
        } else {
            sub[[name]]
        }
        matchFun <- if (fixed) {
            function(x) any(x %in% values)
        } else {
            function(x) {
                any(vapply(values, FUN = \(p) any(grepl(p, x = x)), logical(1)))
            }
        }
        hasMatch <- vapply(col, FUN = matchFun, logical(1))
        join_op <- if (join == "intersect") `&` else `|`
        keep <- join_op(keep, hasMatch)
    }

    rows <- parent_rows[which(keep)]
    .wrap_pipenv(env, name = sprintf("%s view", x[["name"]]), view = rows)
}

Try the pipeflow package in your browser

Any scripts or data that you put into this service are public.

pipeflow documentation built on Sept. 28, 2026, 1:06 a.m.