2026-08-05-tier-2-triggers.md 80 KB

Tier 2 Triggers Implementation Plan

For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (- [ ]) syntax for tracking.

Goal: Ship the four Tier 2 entries: real HTTP responses for webhook-triggered workflows, two polling triggers, and a manual approval step that pauses an execution until a person answers it.

Architecture: Three small pieces plus one substantial one. fs.readdir is a single addition to the QuickJS bridge. Webhook responses are a marker a node returns, recorded by the engine and honoured by the webhook controller. Resumable executions add a Waiting status, a _pause marker that stops the topological walk, and a resume path that replays the stored node_results rather than re-running them. The two polling triggers are plain JavaScript over storage and the new fs.readdir.

Tech Stack: C++20 (gRPC, Protobuf, nlohmann/json, cpp-httplib, QuickJS), JavaScript node modules, Python 3 verification harness.

Global Constraints

  • Node files: 4-space indentation, single quotes in schema literals, no // comments inside schema literals, @description on one line - the header parser captures only the first line.
  • Nodes throw on failure. Never return { success: false }.
  • configSchema default: values are not applied by the runner. Defend every config read in code.
  • The engine evaluates {{...}} in node config before execute() and preserves a single expression's native type. Nodes must not re-interpolate.
  • Use smartbotic.utils.getFieldValue(data, path) for paths, never a private copy.
  • Guard a configurable output field name against the node's own reserved keys, and only where a collision is actually possible.
  • C++20, built with -Wall -Wextra -Wpedantic. No new warnings.
  • Never an em-dash or en-dash in any file, commit message, comment or log string.
  • Commits are GPG-signed. If a commit fails with a gpg pinentry timeout, STOP and report rather than committing unsigned.

Reference: the running system

Services: database 9004, webserver 8090 (HTTP) and 9012 (gRPC), runner 9011, web UI 3000.

Rebuild and restart, in this order and as separate commands - the runner loads its node set once over the webserver's gRPC port and silently holds an empty registry if it starts first:

cmake --build build -j$(nproc)
pkill -f smartbotic-webserver
PID=$(pgrep -f smartbotic-runner); [ -n "$PID" ] && kill -9 $PID
sleep 2
nohup ./build/smartbotic-webserver >/tmp/webserver.log 2>&1 & disown
sleep 6
nohup ./build/smartbotic-runner >/tmp/runner.log 2>&1 & disown
sleep 6
ss -ltnp | grep -E '8090|9011|9012'

A stale runner has survived plain pkill in every task so far. Always check the port is free before relaunching, and kill by PID if it is not.

Verification is python3 scripts/verify-node.py tests/nodes/<case>.json. The node schema parser is regex-based and fails silently: if a node is missing from GET /api/v1/nodes after migrating, its schema literal did not parse - check /tmp/webserver.log for "Failed to parse configSchema".


Task 1: Harness support for acting mid-execution and checking HTTP

Three later tasks need assertions the harness cannot express: a webhook's status code and headers, and an execution that pauses, is answered, and then continues. Building that first keeps those tasks from quietly settling for weaker checks.

Files:

  • Modify: scripts/verify-node.py

Interfaces:

  • Produces: two new optional top-level keys in a case file.

    • "http": {method, path, body, expectStatus, expectHeaders, expectBodyContains} - instead of executing the workflow through /execute, the harness calls path on the webserver directly and asserts on the response. Used for webhook cases.
    • "resume": {waitForStatus, tokenFrom, payload} - after starting the execution, poll until the execution reaches waitForStatus, read the token from the node output named by tokenFrom, POST it to the resume endpoint with payload, then continue polling to completion before running the normal expect assertions.
  • [ ] Step 1: Add the HTTP case mode

In scripts/verify-node.py, after the workflow is created and before the normal execute path, branch on the presence of an http key. Insert this function above main:

def run_http_case(case, token, workflow_id):
    """A webhook case asserts on the HTTP response, not on node outputs."""
    spec = case["http"]
    url = BASE.replace("/api/v1", "") + spec["path"].replace("{workflowId}", workflow_id)
    data = json.dumps(spec.get("body", {})).encode()
    req = urllib.request.Request(url, data=data, method=spec.get("method", "POST"))
    req.add_header("Content-Type", "application/json")

    failures = []
    try:
        with urllib.request.urlopen(req, timeout=45) as res:
            status, headers, raw = res.status, dict(res.headers), res.read().decode()
    except urllib.error.HTTPError as e:
        status, headers, raw = e.code, dict(e.headers), e.read().decode()

    print("http: %s %s -> %s" % (spec.get("method", "POST"), spec["path"], status))
    print("  body: %s" % raw[:200])

    if "expectStatus" in spec and status != spec["expectStatus"]:
        failures.append("status: expected %s, got %s" % (spec["expectStatus"], status))

    for name, want in spec.get("expectHeaders", {}).items():
        got = headers.get(name)
        if got is None:
            failures.append("header %s: missing" % name)
        elif want not in got:
            failures.append("header %s: expected to contain %r, got %r" % (name, want, got))

    for fragment in spec.get("expectBodyContains", []):
        if fragment not in raw:
            failures.append("body: expected to contain %r" % fragment)

    return failures
  • Step 2: Add the resume case mode

Insert above main:

def run_resume_case(case, token, execution_id):
    """Wait for the execution to pause, answer it, then let it finish."""
    spec = case["resume"]
    want_status = spec.get("waitForStatus", "waiting")

    execution = {}
    for _ in range(60):
        time.sleep(0.5)
        execution = call("GET", "/executions/%s" % execution_id, token)
        if execution.get("status") == want_status:
            break
        if execution.get("status") in ("completed", "failed", "cancelled"):
            raise SystemExit(
                "execution %s reached %s without pausing" % (execution_id, execution.get("status")))
    else:
        raise SystemExit("execution %s never reached %s" % (execution_id, want_status))

    by_id = {n["nodeId"]: n for n in execution.get("nodeExecutions", [])}
    source = by_id.get(spec["tokenFrom"])
    if not source:
        raise SystemExit("resume: node %s did not run" % spec["tokenFrom"])
    resume_token = (source.get("output") or {}).get("token")
    if not resume_token:
        raise SystemExit("resume: node %s produced no token" % spec["tokenFrom"])

    body = dict(spec.get("payload", {}))
    body["token"] = resume_token
    call("POST", "/executions/%s/resume" % execution_id, token, body)
    print("resumed %s with token %s" % (execution_id, resume_token))
  • Step 3: Wire both into main

In main, immediately after the workflow is created inside the try block, add the HTTP branch before the existing execute call:

        if "http" in case:
            failures = run_http_case(case, token, workflow_id)
            print("\nFAIL" if failures else "\nPASS")
            for f in failures:
                print("  " + f)
            return 1 if failures else 0

Then, after execution_id is obtained and before the existing polling loop, add:

        if "resume" in case:
            run_resume_case(case, token, execution_id)

The existing polling loop then runs unchanged and sees the resumed execution through to completion.

  • Step 4: Confirm nothing regressed

Run:

for f in tests/nodes/*.json; do python3 scripts/verify-node.py "$f" >/dev/null || echo "FAILED: $f"; done; echo done

Expected: no FAILED: lines. Every existing case lacks both new keys, so all take the unchanged path.

  • [ ] Step 5: Commit

    git add scripts/verify-node.py
    git commit -m "test: let the harness assert on HTTP responses and answer a paused execution"
    

Task 2: fs.readdir, and a workflow id a node can see

Files:

  • Modify: src/runner/engine/script_engine.cpp (after the fs.stat registration, before JS_SetPropertyStr(ctx, smartbotic, "fs", fs);, and the context object around line 4393)
  • Modify: src/runner/workflow_engine.cpp:905 (the ctx filled before a node runs)
  • Create: tests/nodes/fs-readdir.json

Interfaces:

  • Produces: smartbotic.fs.readdir(path) returning {success: true, entries: [{name, path, size, modifiedAt, isDirectory}]} or {success: false, error}. modifiedAt is milliseconds since the epoch, matching fs.stat's mtime. Not recursive.
  • Produces: context.workflowId, alongside the context.executionId and context.nodeId a node already receives. Tasks 3 and 4 key their cursors on it.

  • [ ] Step 1: Register the function

Insert after the closing }, "stat", 1));:

    // fs.readdir(path) - List a directory, one level only. Recursing here would
    // be an unbounded walk driven by whatever happens to be on disk, so a node
    // that wants a tree walks it itself.
    JS_SetPropertyStr(ctx, fs, "readdir", JS_NewCFunction(ctx, [](JSContext* ctx, JSValue this_val, int argc, JSValue* argv) -> JSValue {
        if (argc < 1) {
            return JS_ThrowTypeError(ctx, "fs.readdir requires path argument");
        }

        const char* path = JS_ToCString(ctx, argv[0]);
        if (!path) {
            return JS_ThrowTypeError(ctx, "path must be a string");
        }
        std::string path_str(path);
        JS_FreeCString(ctx, path);

        try {
            if (!std::filesystem::exists(path_str)) {
                JSValue response = JS_NewObject(ctx);
                JS_SetPropertyStr(ctx, response, "success", JS_FALSE);
                JS_SetPropertyStr(ctx, response, "error",
                                  JS_NewString(ctx, ("Directory not found: " + path_str).c_str()));
                return response;
            }
            if (!std::filesystem::is_directory(path_str)) {
                JSValue response = JS_NewObject(ctx);
                JS_SetPropertyStr(ctx, response, "success", JS_FALSE);
                JS_SetPropertyStr(ctx, response, "error",
                                  JS_NewString(ctx, ("Not a directory: " + path_str).c_str()));
                return response;
            }

            JSValue entries = JS_NewArray(ctx);
            uint32_t index = 0;

            for (const auto& entry : std::filesystem::directory_iterator(path_str)) {
                JSValue item = JS_NewObject(ctx);
                JS_SetPropertyStr(ctx, item, "name",
                                  JS_NewString(ctx, entry.path().filename().string().c_str()));
                JS_SetPropertyStr(ctx, item, "path",
                                  JS_NewString(ctx, entry.path().string().c_str()));

                const bool is_dir = entry.is_directory();
                JS_SetPropertyStr(ctx, item, "isDirectory", is_dir ? JS_TRUE : JS_FALSE);

                int64_t size = 0;
                int64_t mtime_ms = 0;
                // A file removed between listing and stat is ordinary on a
                // directory being written to, and is reported with zeroes
                // rather than failing the whole listing.
                try {
                    if (entry.is_regular_file()) {
                        size = static_cast<int64_t>(entry.file_size());
                    }
                    auto mtime = entry.last_write_time();
                    mtime_ms = std::chrono::duration_cast<std::chrono::milliseconds>(
                        mtime.time_since_epoch()).count();
                } catch (const std::exception&) {
                }

                JS_SetPropertyStr(ctx, item, "size", JS_NewInt64(ctx, size));
                JS_SetPropertyStr(ctx, item, "modifiedAt", JS_NewInt64(ctx, mtime_ms));

                JS_SetPropertyUint32(ctx, entries, index++, item);
            }

            JSValue response = JS_NewObject(ctx);
            JS_SetPropertyStr(ctx, response, "success", JS_TRUE);
            JS_SetPropertyStr(ctx, response, "entries", entries);
            return response;
        } catch (const std::exception& e) {
            JSValue response = JS_NewObject(ctx);
            JS_SetPropertyStr(ctx, response, "success", JS_FALSE);
            JS_SetPropertyStr(ctx, response, "error", JS_NewString(ctx, e.what()));
            return response;
        }
    }, "readdir", 1));
  • Step 2: Expose the workflow id to nodes

A node currently receives only executionId and nodeId on its context. The two polling triggers in Tasks 3 and 4 key their cursor on the workflow, and without it every workflow's cursor would collide under one key - two workflows watching different directories would consume each other's changes.

ScriptContext already declares workflow_id (script_engine.hpp:147) and executeNode already receives the Workflow, so nothing new has to be plumbed through. It is simply never set and never exposed.

In workflow_engine.cpp, in executeNode beside the other ctx assignments around line 903:

    ctx.workflow_id = workflow.id;

In script_engine.cpp, in the context object built for the execute function around line 4393:

    nlohmann::json ctx_json = {
        {"executionId", context.execution_id},
        {"nodeId", context.node_id},
        {"workflowId", context.workflow_id}
    };
  • Step 3: Build and restart

Run the rebuild and restart sequence from the Reference section. Expected: clean build, all three ports listening.

  • Step 4: Write the fixture

tests/nodes/fs-readdir.json - a code node creates a directory with two files, lists it, then cleans up:

{
  "name": "verify-fs-readdir",
  "nodes": [
    {"id": "n1", "name": "Trigger", "type": "click-trigger", "position": {"x": 0, "y": 0}, "config": {}},
    {"id": "n2", "name": "List", "type": "code", "position": {"x": 0, "y": 100},
     "config": {"code": "const dir = '/tmp/sb-readdir-test';\nsmartbotic.fs.mkdir(dir);\nsmartbotic.fs.writeFile(dir + '/one.txt', smartbotic.utils.base64Encode('hello'));\nsmartbotic.fs.writeFile(dir + '/two.txt', smartbotic.utils.base64Encode('worldwide'));\nsmartbotic.fs.mkdir(dir + '/sub');\n\nconst listing = smartbotic.fs.readdir(dir);\nconst missing = smartbotic.fs.readdir('/tmp/sb-readdir-does-not-exist');\nconst notDir = smartbotic.fs.readdir(dir + '/one.txt');\n\nconst names = listing.entries.map(function (e) { return e.name; }).sort();\nconst one = listing.entries.filter(function (e) { return e.name === 'one.txt'; })[0];\nconst sub = listing.entries.filter(function (e) { return e.name === 'sub'; })[0];\n\nsmartbotic.fs.unlink(dir + '/one.txt');\nsmartbotic.fs.unlink(dir + '/two.txt');\n\nreturn {\n    ok: listing.success,\n    names: names,\n    oneSize: one.size,\n    oneHasMtime: one.modifiedAt > 1600000000000,\n    subIsDirectory: sub.isDirectory,\n    oneIsDirectory: one.isDirectory,\n    missingRejected: missing.success === false,\n    notDirRejected: notDir.success === false,\n    hasWorkflowId: typeof context.workflowId === 'string' && context.workflowId.length > 0\n};"}}
  ],
  "connections": [
    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "n2", "targetInput": "data"}
  ],
  "expect": {
    "n2": {"status": "completed", "output": {"result": {
      "ok": true,
      "names": ["one.txt", "sub", "two.txt"],
      "oneSize": 5,
      "oneHasMtime": true,
      "subIsDirectory": true,
      "oneIsDirectory": false,
      "missingRejected": true,
      "notDirRejected": true
    }}}
  }
}

oneSize is 5 because the file holds hello. That pins the size actually being read rather than defaulted to zero.

  • Step 5: Run it

Run: python3 scripts/verify-node.py tests/nodes/fs-readdir.json Expected: PASS.

  • [ ] Step 6: Commit

    git add src/runner/engine/script_engine.cpp src/runner/workflow_engine.cpp tests/nodes/fs-readdir.json
    git commit -m "feat: fs.readdir, and a workflow id a node can see"
    

Task 3: File Watch trigger

Files:

  • Create: nodes/triggers/file-watch.js
  • Create: tests/nodes/file-watch.json

Interfaces:

  • Consumes: smartbotic.fs.readdir(path) from Task 2.
  • Produces: output {files: [...], count, isFirstRun}. Cursor documents live in the _watch_cursors collection, keyed <workflowId>:<nodeId>.

  • [ ] Step 1: Write the node

    /**
    * @node file-watch
    * @name File Watch
    * @category triggers
    * @version 1.0.0
    * @description Report files added or changed in a directory since the last run, paired with a schedule trigger for its cadence
    * @icon folder-search
    */
    
    const configSchema = {
    type: 'object',
    properties: {
        directory: {
            type: 'string',
            title: 'Directory',
            description: 'Absolute path to watch. Not recursive'
        },
        pattern: {
            type: 'string',
            title: 'Name Pattern',
            description: 'Optional regular expression a file name must match, such as \\.csv$'
        },
        includeDirectories: {
            type: 'boolean',
            title: 'Include Directories',
            description: 'Report subdirectories as well as files',
            default: false
        },
        emitOnFirstRun: {
            type: 'boolean',
            title: 'Report Everything On First Run',
            description: 'On the very first run there is nothing to compare against. Off records what is there and reports nothing, so adding this to a live workflow does not fire for every file already present',
            default: false
        },
        cursorCollection: {
            type: 'string',
            title: 'Cursor Collection',
            description: 'Collection holding the last-seen state',
            default: '_watch_cursors'
        }
    },
    required: ['directory']
    };
    
    const inputSchema = {
    type: 'object',
    properties: {
        data: { type: 'any' }
    }
    };
    
    const outputSchema = {
    type: 'object',
    properties: {
        files: { type: 'array', description: 'Files new or changed since the last run' },
        count: { type: 'number' },
        isFirstRun: { type: 'boolean', description: 'True when no cursor existed yet' }
    }
    };
    
    async function execute(config, input, context) {
    const directory = config.directory;
    if (!directory) {
        throw new Error('File Watch: a directory is required');
    }
    
    const collection = config.cursorCollection || '_watch_cursors';
    const workflowId = (context && context.workflowId) || 'unknown';
    const nodeId = (context && context.nodeId) || 'file-watch';
    const cursorId = workflowId + ':' + nodeId;
    
    const listing = smartbotic.fs.readdir(directory);
    if (!listing || listing.success !== true) {
        throw new Error('File Watch: could not read ' + directory + ': ' +
            ((listing && listing.error) || 'unknown error'));
    }
    
    let matcher = null;
    if (config.pattern) {
        try {
            matcher = new RegExp(config.pattern);
        } catch (e) {
            throw new Error('File Watch: "' + config.pattern + '" is not a valid pattern: ' + e.message);
        }
    }
    
    const current = {};
    const candidates = [];
    for (const entry of listing.entries) {
        if (entry.isDirectory && config.includeDirectories !== true) {
            continue;
        }
        if (matcher && !matcher.test(entry.name)) {
            continue;
        }
        // Size and modification time together, because a file rewritten within
        // the same second at the same length is not a change worth waking a
        // workflow for, and a timestamp alone misses a rewrite that preserves
        // mtime granularity.
        current[entry.name] = entry.modifiedAt + ':' + entry.size;
        candidates.push(entry);
    }
    
    const stored = smartbotic.storage.get(collection, cursorId);
    const isFirstRun = !stored || stored.found !== true;
    const previous = (!isFirstRun && stored.document && stored.document.seen) || {};
    
    let changed = [];
    if (isFirstRun && config.emitOnFirstRun !== true) {
        smartbotic.log.info('File Watch: first run on ' + directory + ', recorded ' +
            candidates.length + ' entries without reporting them');
    } else {
        changed = candidates.filter(function (entry) {
            return previous[entry.name] !== current[entry.name];
        });
    }
    
    const cursor = { seen: current, updatedAt: Date.now(), directory: directory };
    if (isFirstRun) {
        smartbotic.storage.insert(collection, cursor, cursorId);
    } else {
        smartbotic.storage.update(collection, cursorId, cursor);
    }
    
    smartbotic.log.info('File Watch: ' + changed.length + ' new or changed in ' + directory);
    
    return {
        files: changed,
        count: changed.length,
        isFirstRun: isFirstRun
    };
    }
    
    module.exports = { configSchema, inputSchema, outputSchema, execute };
    
  • [ ] Step 2: Check the context fields exist

The node reads context.workflowId and context.nodeId. Confirm both are actually provided:

grep -n "workflowId\|nodeId" src/runner/engine/script_engine.cpp | grep -i "context" | head

If either is absent, the cursor id would collapse to unknown:file-watch and two workflows watching different directories would share one cursor - a real bug. If they are missing, report it rather than working around it; the fix belongs in the engine and needs its own review.

  • Step 3: Write the fixture

tests/nodes/file-watch.json - first run records without reporting, a second run after adding a file reports exactly that file. Both runs happen in one workflow, with code nodes doing the filesystem work between them:

{
  "name": "verify-file-watch",
  "nodes": [
    {"id": "n1", "name": "Trigger", "type": "click-trigger", "position": {"x": 0, "y": 0}, "config": {}},
    {"id": "setup", "name": "Setup", "type": "code", "position": {"x": 0, "y": 100},
     "config": {"code": "const dir = '/tmp/sb-file-watch-test';\nconst listing = smartbotic.fs.readdir(dir);\nif (listing.success) {\n    for (const e of listing.entries) { smartbotic.fs.unlink(e.path); }\n}\nsmartbotic.fs.mkdir(dir);\nsmartbotic.fs.writeFile(dir + '/existing.txt', smartbotic.utils.base64Encode('old'));\nsmartbotic.storage.delete('_watch_cursors', 'verify-file-watch:first');\nsmartbotic.storage.delete('_watch_cursors', 'verify-file-watch:second');\nreturn { dir: dir };"}},
    {"id": "first", "name": "First Run", "type": "file-watch", "position": {"x": 0, "y": 200},
     "config": {"directory": "/tmp/sb-file-watch-test", "emitOnFirstRun": false}},
    {"id": "add", "name": "Add A File", "type": "code", "position": {"x": 0, "y": 300},
     "config": {"code": "smartbotic.fs.writeFile('/tmp/sb-file-watch-test/fresh.txt', smartbotic.utils.base64Encode('new'));\nreturn { added: 'fresh.txt' };"}},
    {"id": "second", "name": "Second Run", "type": "file-watch", "position": {"x": 0, "y": 400},
     "config": {"directory": "/tmp/sb-file-watch-test", "emitOnFirstRun": false}}
  ],
  "connections": [
    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "setup", "targetInput": "data"},
    {"sourceNodeId": "setup", "sourceOutput": "main", "targetNodeId": "first", "targetInput": "data"},
    {"sourceNodeId": "first", "sourceOutput": "main", "targetNodeId": "add", "targetInput": "data"},
    {"sourceNodeId": "add", "sourceOutput": "main", "targetNodeId": "second", "targetInput": "data"}
  ],
  "expect": {
    "first": {"status": "completed", "output": {"count": 0, "isFirstRun": true}},
    "second": {"status": "completed", "output": {"count": 1, "isFirstRun": false, "files": [{"name": "fresh.txt"}]}}
  }
}

Both file-watch nodes share a workflow but have different node ids, so they hold separate cursors - which is why first sees a first run and second does too unless the cursor id includes the node id. If second reports isFirstRun: true, the cursor key is not node-scoped and that is the bug to report.

Note the fixture deletes its cursors in setup so it is repeatable. The node ids in those delete calls must match the fixture's actual node ids; if you rename them, update the setup code.

  • Step 4: Run it

Run: python3 scripts/verify-node.py tests/nodes/file-watch.json Expected: PASS. Run it twice in a row - it must pass both times, which is what the cursor cleanup in setup is for.

  • [ ] Step 5: Commit

    git add nodes/triggers/file-watch.js tests/nodes/file-watch.json
    git commit -m "feat: a File Watch trigger, reporting what changed in a directory"
    

Task 4: Database Change trigger

Files:

  • Create: nodes/triggers/database-change.js
  • Create: tests/nodes/database-change.json

Interfaces:

  • Produces: output {documents: [...], count, isFirstRun}. Cursor documents live in the same _watch_cursors collection, keyed <workflowId>:<nodeId>.

  • [ ] Step 1: Write the node

    /**
    * @node database-change
    * @name Database Change
    * @category triggers
    * @version 1.0.0
    * @description Report documents added or changed in a collection since the last run, paired with a schedule trigger for its cadence
    * @icon database
    */
    
    const configSchema = {
    type: 'object',
    properties: {
        collection: {
            type: 'string',
            title: 'Collection',
            description: 'Collection to watch'
        },
        timestampField: {
            type: 'string',
            title: 'Timestamp Field',
            description: 'Field holding the last-modified time, in milliseconds. Defaults to the _updatedAt the database maintains itself',
            default: '_updatedAt'
        },
        filter: {
            type: 'object',
            title: 'Filter',
            description: 'Optional query restricting which documents are watched',
            additionalProperties: true
        },
        maxDocuments: {
            type: 'number',
            title: 'Max Documents',
            description: 'Most documents to report in one run, so a large backlog does not arrive as one enormous payload',
            default: 100
        },
        emitOnFirstRun: {
            type: 'boolean',
            title: 'Report Everything On First Run',
            description: 'Off records the current high-water mark and reports nothing, so adding this to a live workflow does not fire for every document already there',
            default: false
        },
        cursorCollection: {
            type: 'string',
            title: 'Cursor Collection',
            default: '_watch_cursors'
        }
    },
    required: ['collection']
    };
    
    const inputSchema = {
    type: 'object',
    properties: {
        data: { type: 'any' }
    }
    };
    
    const outputSchema = {
    type: 'object',
    properties: {
        documents: { type: 'array', description: 'Documents new or changed since the last run' },
        count: { type: 'number' },
        isFirstRun: { type: 'boolean' },
        cursor: { type: 'number', description: 'High-water timestamp stored for the next run' }
    }
    };
    
    async function execute(config, input, context) {
    const collection = config.collection;
    if (!collection) {
        throw new Error('Database Change: a collection is required');
    }
    
    const timestampField = config.timestampField || '_updatedAt';
    const cursorCollection = config.cursorCollection || '_watch_cursors';
    const maxDocuments = Number(config.maxDocuments) || 100;
    const workflowId = (context && context.workflowId) || 'unknown';
    const nodeId = (context && context.nodeId) || 'database-change';
    const cursorId = workflowId + ':' + nodeId;
    
    const stored = smartbotic.storage.get(cursorCollection, cursorId);
    const isFirstRun = !stored || stored.found !== true;
    const since = (!isFirstRun && stored.document && Number(stored.document.since)) || 0;
    
    const query = config.filter && typeof config.filter === 'object' ? config.filter : {};
    const result = smartbotic.storage.query(collection, query);
    if (!result || result.success !== true) {
        throw new Error('Database Change: could not query ' + collection + ': ' +
            ((result && result.error) || 'unknown error'));
    }
    
    const documents = result.documents || [];
    
    // The database stores _updatedAt in nanoseconds while everything a node
    // sees is milliseconds, so the value is normalised rather than compared
    // against a cursor in different units.
    function stampOf(document) {
        const raw = Number(smartbotic.utils.getFieldValue(document, timestampField));
        if (!raw || isNaN(raw)) {
            return 0;
        }
        return raw > 1e15 ? Math.floor(raw / 1000000) : raw;
    }
    
    let highWater = since;
    for (const document of documents) {
        const stamp = stampOf(document);
        if (stamp > highWater) {
            highWater = stamp;
        }
    }
    
    let changed = [];
    if (isFirstRun && config.emitOnFirstRun !== true) {
        smartbotic.log.info('Database Change: first run on ' + collection + ', recorded the mark at ' +
            highWater + ' without reporting ' + documents.length + ' documents');
    } else {
        changed = documents
            .filter(function (document) { return stampOf(document) > since; })
            .sort(function (left, right) { return stampOf(left) - stampOf(right); });
        if (changed.length > maxDocuments) {
            smartbotic.log.warn('Database Change: ' + changed.length + ' documents changed, reporting the oldest ' +
                maxDocuments + '. The rest arrive on the next run');
            changed = changed.slice(0, maxDocuments);
            // The mark only advances as far as what was actually reported, or
            // the remainder would be skipped rather than deferred.
            highWater = stampOf(changed[changed.length - 1]);
        }
    }
    
    const cursor = { since: highWater, updatedAt: Date.now(), collection: collection };
    if (isFirstRun) {
        smartbotic.storage.insert(cursorCollection, cursor, cursorId);
    } else {
        smartbotic.storage.update(cursorCollection, cursorId, cursor);
    }
    
    smartbotic.log.info('Database Change: ' + changed.length + ' documents from ' + collection);
    
    return {
        documents: changed,
        count: changed.length,
        isFirstRun: isFirstRun,
        cursor: highWater
    };
    }
    
    module.exports = { configSchema, inputSchema, outputSchema, execute };
    
  • [ ] Step 2: Confirm the timestamp unit

The node assumes _updatedAt may be nanoseconds. Confirm against a real document:

TOKEN=$(curl -s http://localhost:8090/api/v1/auth/login -H "Content-Type: application/json" \
  -d '{"username":"admin","password":"admin"}' | jq -r '.accessToken')
curl -s "http://localhost:8090/api/v1/executions?limit=1" -H "Authorization: Bearer $TOKEN" \
  | jq '.executions[0]._updated_at, .executions[0]._created_at'

Record the magnitude in your report. If those values are around 1.7e18 the normalisation is needed as written; if they are around 1.7e12 they are already milliseconds and the branch simply never fires. Either way do not remove it - a mix across collections is exactly the case it exists for.

  • Step 3: Write the fixture

tests/nodes/database-change.json - first run records the mark, a document is inserted, second run reports exactly it:

{
  "name": "verify-database-change",
  "nodes": [
    {"id": "n1", "name": "Trigger", "type": "click-trigger", "position": {"x": 0, "y": 0}, "config": {}},
    {"id": "setup", "name": "Setup", "type": "code", "position": {"x": 0, "y": 100},
     "config": {"code": "smartbotic.storage.delete('_watch_cursors', 'verify-database-change:first');\nsmartbotic.storage.delete('_watch_cursors', 'verify-database-change:second');\nsmartbotic.storage.insert('_dbchange_test', { name: 'before', at: Date.now() }, 'before-1');\nreturn { ready: true };"}},
    {"id": "first", "name": "First Run", "type": "database-change", "position": {"x": 0, "y": 200},
     "config": {"collection": "_dbchange_test", "emitOnFirstRun": false}},
    {"id": "add", "name": "Insert One", "type": "code", "position": {"x": 0, "y": 300},
     "config": {"code": "smartbotic.utils.sleep(1100);\nsmartbotic.storage.insert('_dbchange_test', { name: 'after', at: Date.now() }, 'after-1');\nreturn { inserted: 'after-1' };"}},
    {"id": "second", "name": "Second Run", "type": "database-change", "position": {"x": 0, "y": 400},
     "config": {"collection": "_dbchange_test", "emitOnFirstRun": false}},
    {"id": "cleanup", "name": "Cleanup", "type": "code", "position": {"x": 0, "y": 500},
     "config": {"code": "smartbotic.storage.delete('_dbchange_test', 'before-1');\nsmartbotic.storage.delete('_dbchange_test', 'after-1');\nreturn { cleaned: true };"}}
  ],
  "connections": [
    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "setup", "targetInput": "data"},
    {"sourceNodeId": "setup", "sourceOutput": "main", "targetNodeId": "first", "targetInput": "data"},
    {"sourceNodeId": "first", "sourceOutput": "main", "targetNodeId": "add", "targetInput": "data"},
    {"sourceNodeId": "add", "sourceOutput": "main", "targetNodeId": "second", "targetInput": "data"},
    {"sourceNodeId": "second", "sourceOutput": "main", "targetNodeId": "cleanup", "targetInput": "data"}
  ],
  "expect": {
    "first": {"status": "completed", "output": {"count": 0, "isFirstRun": true}},
    "second": {"status": "completed", "output": {"count": 1, "isFirstRun": false, "documents": [{"name": "after"}]}}
  }
}

The sleep(1100) exists so the inserted document's timestamp is strictly greater than the recorded mark even if the database's timestamp resolution is coarse. Without it the test is a race and would fail intermittently, which is worse than failing outright.

  • Step 4: Run it

Run: python3 scripts/verify-node.py tests/nodes/database-change.json Expected: PASS, twice in a row.

If second reports count: 2, the first run's mark was not recorded or the comparison is inclusive rather than strict. If it reports count: 0, the inserted document's timestamp field is not what timestampField names - check what the stored document actually contains before changing the node.

  • [ ] Step 5: Commit

    git add nodes/triggers/database-change.js tests/nodes/database-change.json
    git commit -m "feat: a Database Change trigger, reporting documents changed since the last run"
    

Task 5: Webhook responses in the engine and controller

Files:

  • Modify: src/runner/workflow_engine.hpp (the ExecutionResult struct, around line 90-106)
  • Modify: src/runner/workflow_engine.cpp (toJson, and the node-completed path around line 660)
  • Modify: src/runner/runner_service.cpp:125 (where set_result is called)
  • Modify: src/webserver/api/webhook_controller.cpp:172-182 (the response section)

Interfaces:

  • Produces: a node output containing _webhookResponse is recorded on ExecutionResult::webhook_response and serialized under webhookResponse. The webhook controller honours {status, headers, body} from it.

  • [ ] Step 1: Add the field

In workflow_engine.hpp, in ExecutionResult after nlohmann::json workflow_snapshot;:

    nlohmann::json webhook_response;   // Set by a respond-to-webhook node, if any

In workflow_engine.cpp's ExecutionResult::toJson, beside the other members:

    if (!webhook_response.is_null()) {
        j["webhookResponse"] = webhook_response;
    }
  • Step 2: Record the marker when any node sets it

In the main walk, immediately after a node's result is stored and its status is known to be Completed, add:

            // Any node may declare the HTTP response, not just the last one to
            // run, so appending a node to a workflow cannot silently change
            // what its webhook returns.
            if (node_result.status == NodeStatus::Completed &&
                node_result.output.contains("_webhookResponse")) {
                if (!result.webhook_response.is_null()) {
                    LOG_WARN("Node {} overrides a webhook response already set by an earlier node",
                             node_id);
                }
                result.webhook_response = node_result.output["_webhookResponse"];
            }

Place it directly before the existing // Check for loop node block, so it runs for every completed node.

  • Step 3: Put it on the wire

In runner_service.cpp, where the response is filled after a successful execution, the result JSON already goes out through set_result(result.value().final_output.dump()). Change that to send the webhook response when there is one:

        if (!result.value().webhook_response.is_null()) {
            nlohmann::json envelope;
            envelope["_webhookResponse"] = result.value().webhook_response;
            response->set_result(envelope.dump());
        } else {
            response->set_result(result.value().final_output.dump());
        }
  • Step 4: Honour it in the controller

In webhook_controller.cpp, replace the result block (currently parsing grpc_res.result() and calling sendJson) with:

    if (!grpc_res.result().empty()) {
        nlohmann::json parsed;
        bool parsed_ok = true;
        try {
            parsed = nlohmann::json::parse(grpc_res.result());
        } catch (...) {
            parsed_ok = false;
        }

        if (parsed_ok && parsed.is_object() && parsed.contains("_webhookResponse")) {
            const auto& spec = parsed["_webhookResponse"];

            int status_code = spec.value("status", 200);
            if (status_code < 100 || status_code > 599) {
                LOG_WARN("Webhook response asked for status {}, which is not a valid HTTP status; sending 500",
                         status_code);
                status_code = 500;
            }

            std::string content_type = "application/json";
            if (spec.contains("headers") && spec["headers"].is_object()) {
                for (auto it = spec["headers"].begin(); it != spec["headers"].end(); ++it) {
                    if (!it.value().is_string()) {
                        continue;
                    }
                    // Content-Type reaches httplib through set_content rather
                    // than as a header, and setting both sends it twice.
                    std::string name = it.key();
                    std::string lowered;
                    for (char c : name) {
                        lowered += static_cast<char>(std::tolower(static_cast<unsigned char>(c)));
                    }
                    if (lowered == "content-type") {
                        content_type = it.value().get<std::string>();
                    } else {
                        res.set_header(name.c_str(), it.value().get<std::string>().c_str());
                    }
                }
            }

            std::string body;
            if (spec.contains("body")) {
                const auto& value = spec["body"];
                body = value.is_string() ? value.get<std::string>() : value.dump();
            }

            res.status = status_code;
            res.set_content(body, content_type.c_str());
            return;
        }

        if (parsed_ok) {
            sendJson(res, parsed);
        } else {
            res.set_content(grpc_res.result(), "text/plain");
        }
    } else {
        sendJson(res, {{"executionId", grpc_res.execution_id()}, {"status", "completed"}});
    }

Add #include <cctype> at the top of the file if it is not already present.

  • Step 5: Build and restart

Run the rebuild and restart sequence. Expected: clean build, all ports listening.

  • Step 6: Confirm nothing regressed

Run:

for f in tests/nodes/*.json; do python3 scripts/verify-node.py "$f" >/dev/null || echo "FAILED: $f"; done; echo done

Expected: no failures. No node sets _webhookResponse yet, so every path is the unchanged one.

  • [ ] Step 7: Commit

    git add src/runner/workflow_engine.hpp src/runner/workflow_engine.cpp src/runner/runner_service.cpp src/webserver/api/webhook_controller.cpp
    git commit -m "feat: let a node declare the HTTP response a webhook returns"
    

Task 6: Respond to Webhook node

Files:

  • Create: nodes/core/respond-to-webhook.js
  • Create: tests/nodes/respond-to-webhook.json

Interfaces:

  • Consumes: the _webhookResponse handling from Task 5, and the harness http case mode from Task 1.

  • [ ] Step 1: Write the node

    /**
    * @node respond-to-webhook
    * @name Respond to Webhook
    * @category core
    * @version 1.0.0
    * @description Set the status, headers and body a webhook-triggered workflow returns to its caller
    * @icon reply
    */
    
    const configSchema = {
    type: 'object',
    properties: {
        status: {
            type: 'number',
            title: 'Status Code',
            description: 'HTTP status to return, such as 200, 201, 400 or 404',
            default: 200
        },
        bodySource: {
            type: 'string',
            title: 'Body From',
            description: 'json sends the text below parsed as JSON, text sends it as it is, and input sends a value taken from the incoming data',
            enum: ['json', 'text', 'input'],
            default: 'json'
        },
        body: {
            type: 'string',
            title: 'Body',
            description: 'The response body, for the json and text sources. Supports {{variable}} interpolation',
            format: 'textarea'
        },
        bodyField: {
            type: 'string',
            title: 'Body Field',
            description: 'Path to the value to send, for the input source, such as data.result',
            default: 'data'
        },
        contentType: {
            type: 'string',
            title: 'Content Type',
            description: 'Content-Type header. Leave empty to send application/json',
            default: 'application/json'
        },
        headers: {
            type: 'object',
            title: 'Extra Headers',
            description: 'Additional response headers',
            additionalProperties: { type: 'string' }
        }
    }
    };
    
    const inputSchema = {
    type: 'object',
    properties: {
        data: { type: 'any' }
    }
    };
    
    const outputSchema = {
    type: 'object',
    properties: {
        status: { type: 'number', description: 'Status this node asked for' },
        respondedWith: { type: 'string', description: 'Content type sent' }
    }
    };
    
    async function execute(config, input, context) {
    const status = config.status === undefined || config.status === null || config.status === ''
        ? 200
        : Number(config.status);
    if (isNaN(status) || status < 100 || status > 599) {
        throw new Error('Respond to Webhook: "' + config.status +
            '" is not an HTTP status code. Use a number from 100 to 599');
    }
    
    const source = config.bodySource || 'json';
    const contentType = config.contentType || 'application/json';
    
    let body;
    if (source === 'input') {
        body = smartbotic.utils.getFieldValue(input, config.bodyField || 'data');
    } else if (source === 'text') {
        body = config.body === undefined || config.body === null ? '' : String(config.body);
    } else {
        const raw = config.body === undefined || config.body === null ? '' : String(config.body);
        if (raw.length === 0) {
            body = {};
        } else {
            try {
                body = JSON.parse(raw);
            } catch (e) {
                throw new Error('Respond to Webhook: the body is not valid JSON: ' + e.message +
                    '. Use the text source to send it as it is');
            }
        }
    }
    
    const headers = {};
    if (config.headers && typeof config.headers === 'object') {
        for (const name of Object.keys(config.headers)) {
            headers[name] = String(config.headers[name]);
        }
    }
    headers['Content-Type'] = contentType;
    
    smartbotic.log.info('Respond to Webhook: ' + status + ' as ' + contentType);
    
    return {
        status: status,
        respondedWith: contentType,
        _webhookResponse: {
            status: status,
            headers: headers,
            body: body
        }
    };
    }
    
    module.exports = { configSchema, inputSchema, outputSchema, execute };
    
  • [ ] Step 2: Find the webhook path

The fixture calls a webhook URL rather than /execute. Find the route and how a workflow's webhook path is formed:

grep -n "Post\|Get\|webhook" src/webserver/api/webhook_controller.cpp | grep -i "server\." | head

Record the exact path shape in your report - the fixture's http.path must match it. It is likely /webhook/<workflowId> or similar, and the fixture below uses {workflowId}, which the harness substitutes.

  • Step 3: Write the fixture

tests/nodes/respond-to-webhook.json. Replace the path with the real shape found in Step 2. The workflow needs a webhook trigger rather than a click trigger, since it is reached over HTTP:

{
  "name": "verify-respond-to-webhook",
  "nodes": [
    {"id": "n1", "name": "Webhook", "type": "post-trigger", "position": {"x": 0, "y": 0}, "config": {}},
    {"id": "n2", "name": "Respond", "type": "respond-to-webhook", "position": {"x": 0, "y": 100},
     "config": {"status": 201, "bodySource": "json", "body": "{\"created\": true, \"id\": \"abc\"}",
                "contentType": "application/json", "headers": {"X-Smartbotic-Test": "tier2"}}},
    {"id": "n3", "name": "After", "type": "code", "position": {"x": 0, "y": 200},
     "config": {"code": "return { ranAfterResponding: true };"}}
  ],
  "connections": [
    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "n2", "targetInput": "data"},
    {"sourceNodeId": "n2", "sourceOutput": "main", "targetNodeId": "n3", "targetInput": "data"}
  ],
  "http": {
    "method": "POST",
    "path": "/webhook/{workflowId}",
    "body": {"hello": "world"},
    "expectStatus": 201,
    "expectHeaders": {"Content-Type": "application/json", "X-Smartbotic-Test": "tier2"},
    "expectBodyContains": ["\"created\":true", "\"id\":\"abc\""]
  }
}

The n3 node after the responder is deliberate: it proves the response survives a node running after it, which is the whole point of recording the marker from any node rather than the last one.

expectBodyContains uses the compact form nlohmann emits, without spaces. If the assertion fails on spacing, read the actual body the harness prints and match it rather than loosening the check to something that would pass on any body.

  • Step 4: Run it

Run: python3 scripts/verify-node.py tests/nodes/respond-to-webhook.json Expected: PASS, with the harness printing the status and body it received.

If the status is 200 rather than 201, the marker did not reach the controller - check webhookResponse in GET /api/v1/executions/<id> to see whether the engine recorded it, which tells you which side of the gRPC hop lost it.

  • [ ] Step 5: Commit

    git add nodes/core/respond-to-webhook.js tests/nodes/respond-to-webhook.json
    git commit -m "feat: a Respond to Webhook node, so the HTTP triggers can back a real API"
    

Task 7: A Waiting status and a pause that stops the walk

Files:

  • Modify: src/runner/workflow_engine.hpp (ExecutionStatus enum around line 79, ExecutionResult struct)
  • Modify: src/runner/workflow_engine.cpp (executionStatusToString, the main walk, storeExecution)

Interfaces:

  • Produces: ExecutionStatus::Waiting, serialized as "waiting". ExecutionResult gains paused_node_id, pause_token and pause_expires_at, serialized as pausedNodeId, pauseToken and pauseExpiresAt. A node output containing _pause stops the walk.

  • [ ] Step 1: Add the status and the fields

In workflow_engine.hpp, extend the enum - appended at the end so no existing value's meaning shifts:

enum class ExecutionStatus {
    Pending,
    Running,
    Completed,
    Failed,
    Cancelled,
    Waiting
};

In ExecutionResult, after webhook_response:

    std::string paused_node_id;        // Node that asked to pause, when Waiting
    std::string pause_token;           // Must be presented to resume
    int64_t pause_expires_at = 0;      // Milliseconds since the epoch, 0 for never
  • Step 2: Teach the string conversion

In workflow_engine.cpp, find executionStatusToString and add the case:

        case ExecutionStatus::Waiting: return "waiting";

In ExecutionResult::toJson, beside the other members:

    if (!paused_node_id.empty()) {
        j["pausedNodeId"] = paused_node_id;
        j["pauseToken"] = pause_token;
        j["pauseExpiresAt"] = pause_expires_at;
    }
  • Step 3: Stop the walk on a pause marker

In the main walk, directly after the _webhookResponse block added in Task 5, add:

            // A node asking to pause ends this pass. The execution is stored as
            // Waiting with everything computed so far, and a later resume picks
            // it up from here rather than starting again.
            if (node_result.status == NodeStatus::Completed &&
                node_result.output.contains("_pause")) {
                const auto& pause = node_result.output["_pause"];

                result.status = ExecutionStatus::Waiting;
                result.paused_node_id = node_id;
                result.pause_token = pause.value("token", "");
                result.pause_expires_at = pause.value("expiresAt", static_cast<int64_t>(0));

                // The marker is engine plumbing. What is stored is the request a
                // person has to answer, not the mechanism that carried it.
                auto& stored_node = result.node_results[node_id];
                stored_node.output.erase("_pause");

                if (callback) {
                    callback("execution.waiting", {
                        {"executionId", result.execution_id},
                        {"nodeId", node_id},
                        {"expiresAt", result.pause_expires_at}
                    });
                }

                LOG_INFO("Execution {} paused at node {}", result.execution_id, node_id);
                break;
            }

break leaves the walk over execution_order. Confirm the code after the loop does not unconditionally set the status to Completed - it currently does so only if (result.status == ExecutionStatus::Running) around line 666, which Waiting correctly fails, but verify that is still true where you are inserting.

  • Step 4: Store a paused execution for as long as it may live

In storeExecution, replace the fixed TTL with one derived from the pause:

void WorkflowEngine::storeExecution(const ExecutionResult& result) {
    // A finished execution is a log entry and ages out after a week. One that is
    // waiting for a person is work still owed an answer, and having it expire
    // under them loses the run silently, so it lives until its own deadline
    // plus a day of slack for a late approver.
    int64_t ttl_ms = 7 * 24 * 60 * 60 * 1000;
    if (result.status == ExecutionStatus::Waiting && result.pause_expires_at > 0) {
        const int64_t remaining = result.pause_expires_at - common::TimeUtils::nowMs();
        const int64_t grace = 24 * 60 * 60 * 1000;
        if (remaining + grace > ttl_ms) {
            ttl_ms = remaining + grace;
        }
    }

    auto insert_result = storage_.insert("executions", result.toJson(), result.execution_id, ttl_ms);

    if (insert_result.failed()) {
        LOG_ERROR("Failed to store execution {}: {}", result.execution_id, insert_result.error().message());
    } else {
        LOG_DEBUG("Stored execution {} successfully", result.execution_id);
    }
}

Check how TimeUtils::nowMs() is referenced elsewhere in this file and match it - the file has using namespace common; in some translation units, in which case drop the common:: prefix.

  • Step 5: Build and re-verify

Run the rebuild and restart sequence, then:

for f in tests/nodes/*.json; do python3 scripts/verify-node.py "$f" >/dev/null || echo "FAILED: $f"; done; echo done

Expected: no failures. Nothing emits _pause yet.

  • [ ] Step 6: Commit

    git add src/runner/workflow_engine.hpp src/runner/workflow_engine.cpp
    git commit -m "feat: an execution can pause, and a paused one outlives the seven day log TTL"
    

Task 8: Resuming a paused execution in the engine

Files:

  • Modify: src/runner/workflow_engine.hpp (a resume declaration)
  • Modify: src/runner/workflow_engine.cpp (execute gains a resume path, plus the new resume method)

Interfaces:

  • Consumes: ExecutionStatus::Waiting and the pause fields from Task 7.
  • Produces: common::Result<ExecutionResult> WorkflowEngine::resume(const std::string& execution_id, const std::string& token, const nlohmann::json& payload, ExecutionCallback callback = nullptr).

  • [ ] Step 1: Declare it

In workflow_engine.hpp, after the execute declaration:

    // Continue an execution that stopped at a pause marker. The workflow is
    // rebuilt from the snapshot stored with the execution rather than from the
    // workflow as it stands now, because it may have been edited while the
    // approval was waiting and finishing a run against a different workflow
    // than it started under is worse than refusing.
    common::Result<ExecutionResult> resume(const std::string& execution_id,
                                           const std::string& token,
                                           const nlohmann::json& payload,
                                           ExecutionCallback callback = nullptr);
  • Step 2: Let the walk skip work already done

In execute, the walk executes every node in execution_order. For a resume, nodes already recorded must not run again. Immediately after input is collected and before the _skip check, add:

            // On a resume the results of everything that ran before the pause
            // are already seeded, and they are not run a second time.
            {
                auto seeded = result.node_results.find(node_id);
                if (seeded != result.node_results.end() && seeded->second.from_resume) {
                    LOG_DEBUG("Node {} already ran before the pause, keeping its result", node_id);
                    continue;
                }
            }

Add the flag to NodeExecutionResult in workflow_engine.hpp, beside from_cache:

    bool from_resume = false;   // Seeded from a stored execution rather than run now
  • Step 3: Write the resume method

Add to workflow_engine.cpp, after execute:

common::Result<ExecutionResult> WorkflowEngine::resume(const std::string& execution_id,
                                                       const std::string& token,
                                                       const nlohmann::json& payload,
                                                       ExecutionCallback callback) {
    auto stored = storage_.get("executions", execution_id);
    if (stored.failed()) {
        return common::Error(common::ErrorCode::NotFound, "No execution " + execution_id);
    }

    const nlohmann::json& record = stored.value();

    if (record.value("status", "") != "waiting") {
        return common::Error(common::ErrorCode::InvalidArgument,
            "Execution " + execution_id + " is " + record.value("status", "unknown") +
            ", not waiting for an answer");
    }

    const std::string expected_token = record.value("pauseToken", "");
    if (expected_token.empty() || expected_token != token) {
        return common::Error(common::ErrorCode::PermissionDenied,
            "The token does not match the one this execution is waiting on");
    }

    const int64_t expires_at = record.value("pauseExpiresAt", static_cast<int64_t>(0));
    if (expires_at > 0 && TimeUtils::nowMs() > expires_at) {
        return common::Error(common::ErrorCode::InvalidArgument,
            "This approval expired at " + std::to_string(expires_at) + " and can no longer be answered");
    }

    const std::string paused_node_id = record.value("pausedNodeId", "");
    if (paused_node_id.empty()) {
        return common::Error(common::ErrorCode::Internal,
            "Execution " + execution_id + " is waiting but records no paused node");
    }

    if (!record.contains("workflowSnapshot")) {
        return common::Error(common::ErrorCode::Internal,
            "Execution " + execution_id + " has no workflow snapshot to resume against");
    }

    Workflow workflow = Workflow::fromJson(record["workflowSnapshot"]);

    // Seed everything that ran before the pause, including the paused node
    // itself, whose output becomes the answer that was given.
    nlohmann::json seed = nlohmann::json::object();
    if (record.contains("nodeExecutions") && record["nodeExecutions"].is_array()) {
        for (const auto& entry : record["nodeExecutions"]) {
            const std::string node_id = entry.value("nodeId", "");
            if (node_id.empty()) {
                continue;
            }
            seed[node_id] = entry.value("output", nlohmann::json::object());
        }
    }

    nlohmann::json answer = payload.is_object() ? payload : nlohmann::json::object();
    answer["answeredAt"] = TimeUtils::nowMs();
    seed[paused_node_id] = answer;

    nlohmann::json trigger_data = record.value("triggerData", nlohmann::json::object());
    trigger_data["_resumeSeed"] = seed;
    trigger_data["_resumeExecutionId"] = execution_id;

    LOG_INFO("Resuming execution {} from node {} with {} seeded results",
             execution_id, paused_node_id, seed.size());

    return execute(workflow, record.value("triggerType", "resume"), trigger_data, callback);
}
  • Step 4: Consume the seed in execute

In execute, beside the existing _cachedOutputs extraction around line 212, add:

    // A resume carries the results of everything that ran before the pause, and
    // reuses the original execution id so the run reads as one execution rather
    // than two.
    nlohmann::json resume_seed;
    std::string resume_execution_id;
    if (actual_trigger_data.contains("_resumeSeed")) {
        resume_seed = actual_trigger_data["_resumeSeed"];
        actual_trigger_data.erase("_resumeSeed");
    }
    if (actual_trigger_data.contains("_resumeExecutionId")) {
        resume_execution_id = actual_trigger_data["_resumeExecutionId"].get<std::string>();
        actual_trigger_data.erase("_resumeExecutionId");
    }

Then, after result.execution_id is assigned, override it and seed the results:

    if (!resume_execution_id.empty()) {
        result.execution_id = resume_execution_id;
    }
    for (auto it = resume_seed.begin(); it != resume_seed.end(); ++it) {
        NodeExecutionResult seeded;
        seeded.node_id = it.key();
        seeded.status = NodeStatus::Completed;
        seeded.output = it.value();
        seeded.started_at = TimeUtils::nowMs();
        seeded.finished_at = seeded.started_at;
        seeded.from_resume = true;
        result.node_results[it.key()] = seeded;
    }
  • Step 5: Build and re-verify

Run the rebuild and restart sequence, then the full suite. Expected: no failures - nothing calls resume yet, and no trigger data carries a seed.

  • [ ] Step 6: Commit

    git add src/runner/workflow_engine.hpp src/runner/workflow_engine.cpp
    git commit -m "feat: resume a paused execution from its stored results and snapshot"
    

Task 9: The resume RPC and REST endpoints

Files:

  • Modify: proto/runner.proto (a ResumeExecution RPC and its request message)
  • Modify: src/runner/runner_service.cpp (the handler) and src/runner/runner_service.hpp
  • Modify: src/webserver/api/execution_controller.cpp (two routes) and its header

Interfaces:

  • Consumes: WorkflowEngine::resume from Task 8.
  • Produces: POST /api/v1/executions/{id}/resume taking {token, approved, data}, and GET /api/v1/executions/pending listing executions in waiting.

  • [ ] Step 1: Add the RPC

In proto/runner.proto, inside service RunnerService after CancelExecution:

    // Continue an execution that paused for an answer
    rpc ResumeExecution(ResumeExecutionRequest) returns (ExecuteWorkflowResponse);

And beside the other request messages:

message ResumeExecutionRequest {
    string execution_id = 1;
    string token = 2;
    string payload = 3;  // JSON given by whoever answered
}
  • Step 2: Implement the handler

Declare it in runner_service.hpp beside CancelExecution, matching that method's signature style, then implement in runner_service.cpp:

::grpc::Status RunnerServiceImpl::ResumeExecution(
    ::grpc::ServerContext* context,
    const proto::ResumeExecutionRequest* request,
    proto::ExecuteWorkflowResponse* response) {

    nlohmann::json payload = nlohmann::json::object();
    if (!request->payload().empty()) {
        try {
            payload = nlohmann::json::parse(request->payload());
        } catch (const std::exception& e) {
            return ::grpc::Status(::grpc::StatusCode::INVALID_ARGUMENT,
                                  std::string("payload is not JSON: ") + e.what());
        }
    }

    auto result = engine_->resume(request->execution_id(), request->token(), payload);
    if (result.failed()) {
        return ::grpc::Status(::grpc::StatusCode::FAILED_PRECONDITION, result.error().message());
    }

    response->set_execution_id(result.value().execution_id);
    response->set_status(toProtoStatus(result.value().status));
    response->set_result(result.value().final_output.dump());
    return ::grpc::Status::OK;
}

Match how the existing handlers reach the engine and convert the status - copy the pattern from ExecuteWorkflow in the same file rather than inventing names. If there is no toProtoStatus helper, use whatever ExecuteWorkflow uses.

  • Step 3: Add the REST routes

In execution_controller.cpp's route registration, beside the existing cancel and retry routes:

    server.Post(R"(/api/v1/executions/([^/]+)/resume)", [this](const httplib::Request& req, httplib::Response& res) {
        requireAuth(req, res, [&](const auth::AuthContext& ctx) { resumeExecution(req, res, ctx); });
    });

    server.Get("/api/v1/executions/pending", [this](const httplib::Request& req, httplib::Response& res) {
        requireAuth(req, res, [&](const auth::AuthContext& ctx) { listPending(req, res, ctx); });
    });

Match the surrounding registrations exactly - copy the shape of the existing cancel route rather than the sketch above if they differ.

Route order matters: /api/v1/executions/pending must be registered BEFORE the existing /api/v1/executions/([^/]+) GET route, or the regex swallows pending as an execution id and the listing returns a 404 for an execution named "pending".

  • [ ] Step 4: Implement the handlers

    void ExecutionController::resumeExecution(const httplib::Request& req, httplib::Response& res,
                                          const auth::AuthContext& ctx) {
    std::string execution_id = req.matches[1];
    
    nlohmann::json body;
    try {
        body = req.body.empty() ? nlohmann::json::object() : nlohmann::json::parse(req.body);
    } catch (...) {
        sendError(res, "Body is not valid JSON", 400);
        return;
    }
    
    const std::string token = body.value("token", "");
    if (token.empty()) {
        sendError(res, "A token is required to answer a paused execution", 400);
        return;
    }
    
    nlohmann::json payload = body.value("data", nlohmann::json::object());
    payload["approved"] = body.value("approved", true);
    payload["answeredBy"] = ctx.user_id;
    
    auto runner = load_balancer_.selectRunner();
    if (!runner) {
        sendError(res, "No runners available", 503);
        return;
    }
    
    auto channel = grpc::CreateChannel(runner->address, grpc::InsecureChannelCredentials());
    auto stub = proto::RunnerService::NewStub(channel);
    
    proto::ResumeExecutionRequest grpc_req;
    grpc_req.set_execution_id(execution_id);
    grpc_req.set_token(token);
    grpc_req.set_payload(payload.dump());
    
    proto::ExecuteWorkflowResponse grpc_res;
    grpc::ClientContext grpc_ctx;
    grpc_ctx.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(60));
    
    auto status = stub->ResumeExecution(&grpc_ctx, grpc_req, &grpc_res);
    if (!status.ok()) {
        sendError(res, "Could not resume: " + status.error_message(), 400);
        return;
    }
    
    ws_server_.broadcast("executions." + execution_id + ".resumed", {
        {"executionId", execution_id},
        {"answeredBy", ctx.user_id}
    });
    
    sendJson(res, {{"executionId", execution_id}, {"status", "resumed"}});
    }
    
    void ExecutionController::listPending(const httplib::Request& req, httplib::Response& res,
                                      const auth::AuthContext& ctx) {
    auto result = storage_.query("executions", {{"status", "waiting"}});
    if (result.failed()) {
        sendError(res, "Could not list pending approvals", 500);
        return;
    }
    
    nlohmann::json pending = nlohmann::json::array();
    for (const auto& record : result.value()) {
        pending.push_back({
            {"executionId", record.value("_id", "")},
            {"workflowId", record.value("workflowId", "")},
            {"workflowName", record.value("workflowName", "")},
            {"pausedNodeId", record.value("pausedNodeId", "")},
            {"pauseExpiresAt", record.value("pauseExpiresAt", static_cast<int64_t>(0))},
            {"startedAt", record.value("startedAt", static_cast<int64_t>(0))}
        });
    }
    
    sendJson(res, {{"pending", pending}, {"total", pending.size()}});
    }
    

The listing deliberately does not include the pause token. Anyone who can list pending approvals would otherwise be able to answer all of them, and listing is a weaker permission than approving. The token reaches the approver through the node's own output, on the execution detail.

Declare both in the header beside the existing handlers, and match how storage_, load_balancer_ and ws_server_ are reached in this class - copy from the neighbouring methods.

  • Step 5: Build and re-verify

Run the rebuild and restart sequence, then the full suite. Expected: no failures.

Check the routes exist:

TOKEN=$(curl -s http://localhost:8090/api/v1/auth/login -H "Content-Type: application/json" \
  -d '{"username":"admin","password":"admin"}' | jq -r '.accessToken')
curl -s http://localhost:8090/api/v1/executions/pending -H "Authorization: Bearer $TOKEN"

Expected: {"pending":[],"total":0} rather than a 404. A 404 means the route ordering problem from Step 3.

  • [ ] Step 6: Commit

    git add proto/runner.proto src/runner/runner_service.hpp src/runner/runner_service.cpp src/webserver/api/execution_controller.hpp src/webserver/api/execution_controller.cpp
    git commit -m "feat: an endpoint to answer a paused execution, and one to list them"
    

Task 10: Wait for Approval node, and the end-to-end resume

Files:

  • Create: nodes/core/wait-for-approval.js
  • Create: tests/nodes/wait-for-approval.json
  • Create: tests/nodes/wait-for-approval-loop.json

Interfaces:

  • Consumes: the pause handling from Task 7, the resume path from Tasks 8 and 9, and the harness resume case mode from Task 1.

  • [ ] Step 1: Write the node

    /**
    * @node wait-for-approval
    * @name Wait for Approval
    * @category flow-control
    * @version 1.0.0
    * @description Pause the execution until someone answers, then continue with what they sent
    * @icon user-check
    */
    
    const configSchema = {
    type: 'object',
    properties: {
        reason: {
            type: 'string',
            title: 'Reason',
            description: 'What the approver is being asked to decide. Supports {{variable}} interpolation',
            format: 'textarea'
        },
        expiresIn: {
            type: 'number',
            title: 'Expires In (hours)',
            description: 'How long the request stays answerable. 0 means it never expires, though the execution is still capped at 30 days',
            default: 24
        },
        fields: {
            type: 'array',
            title: 'Fields To Collect',
            description: 'Extra values the approver is asked for, sent back on the answer',
            items: {
                type: 'object',
                properties: {
                    name: { type: 'string', title: 'Name' },
                    title: { type: 'string', title: 'Label' },
                    type: { type: 'string', title: 'Type', enum: ['string', 'number', 'boolean'], default: 'string' }
                }
            }
        }
    }
    };
    
    const inputSchema = {
    type: 'object',
    properties: {
        data: { type: 'any' }
    }
    };
    
    const outputSchema = {
    type: 'object',
    properties: {
        token: { type: 'string', description: 'Token that must be presented to answer this request' },
        reason: { type: 'string' },
        expiresAt: { type: 'number', description: 'Milliseconds since the epoch, 0 when it never expires' },
        approved: { type: 'boolean', description: 'Set on the answer, after the execution resumes' },
        answeredBy: { type: 'string', description: 'Set on the answer' },
        answeredAt: { type: 'number', description: 'Set on the answer' }
    }
    };
    
    const MAX_EXPIRY_HOURS = 30 * 24;
    
    async function execute(config, input, context) {
    // Loop iterations keep their state in engine locals that the stored
    // execution does not describe, so a pause inside one could be recorded but
    // never resumed. Failing here is far better than accepting an approval that
    // can never be answered.
    if (input && (input.index !== undefined || input.currentIndex !== undefined)) {
        throw new Error('Wait for Approval: this node cannot be used inside a Loop body, ' +
            'because a paused loop iteration cannot be resumed. Collect the items first, ' +
            'approve once, then loop');
    }
    
    const hours = config.expiresIn === undefined || config.expiresIn === null || config.expiresIn === ''
        ? 24
        : Number(config.expiresIn);
    if (isNaN(hours) || hours < 0) {
        throw new Error('Wait for Approval: expiresIn must be a number of hours, got "' +
            config.expiresIn + '"');
    }
    if (hours > MAX_EXPIRY_HOURS) {
        throw new Error('Wait for Approval: expiresIn is capped at ' + MAX_EXPIRY_HOURS +
            ' hours. An approval nobody answers should eventually stop occupying the queue');
    }
    
    const token = smartbotic.utils.uuid();
    const expiresAt = hours === 0 ? 0 : Date.now() + Math.round(hours * 60 * 60 * 1000);
    const reason = config.reason ? String(config.reason) : 'Approval required';
    const fields = Array.isArray(config.fields) ? config.fields : [];
    
    smartbotic.log.info('Wait for Approval: pausing for "' + reason + '", token ' + token);
    
    return {
        token: token,
        reason: reason,
        expiresAt: expiresAt,
        fields: fields,
        _pause: {
            token: token,
            reason: reason,
            expiresAt: expiresAt,
            fields: fields
        }
    };
    }
    
    module.exports = { configSchema, inputSchema, outputSchema, execute };
    
  • [ ] Step 2: Write the end-to-end fixture

tests/nodes/wait-for-approval.json - the workflow pauses, the harness answers it, and the node after the approval runs with the answer:

{
  "name": "verify-wait-for-approval",
  "nodes": [
    {"id": "n1", "name": "Trigger", "type": "click-trigger", "position": {"x": 0, "y": 0}, "config": {}},
    {"id": "approve", "name": "Approve", "type": "wait-for-approval", "position": {"x": 0, "y": 100},
     "config": {"reason": "Approve the test", "expiresIn": 1}},
    {"id": "after", "name": "After Approval", "type": "code", "position": {"x": 0, "y": 200},
     "config": {"code": "return { sawApproval: input.data.approved === true, note: input.data.note };"}}
  ],
  "connections": [
    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "approve", "targetInput": "data"},
    {"sourceNodeId": "approve", "sourceOutput": "main", "targetNodeId": "after", "targetInput": "data"}
  ],
  "resume": {
    "waitForStatus": "waiting",
    "tokenFrom": "approve",
    "payload": {"approved": true, "data": {"note": "looks fine"}}
  },
  "expect": {
    "after": {"status": "completed", "output": {"result": {"sawApproval": true, "note": "looks fine"}}}
  }
}

This is the load-bearing test of the whole task. It proves the execution paused, that the after node did not run during the first pass, that the answer reached it, and that the resumed run is the same execution.

  • Step 3: Write the loop-rejection fixture

tests/nodes/wait-for-approval-loop.json:

{
  "name": "verify-wait-for-approval-rejects-loop",
  "nodes": [
    {"id": "n1", "name": "Trigger", "type": "click-trigger", "position": {"x": 0, "y": 0}, "config": {}},
    {"id": "items", "name": "Items", "type": "code", "position": {"x": 0, "y": 100},
     "config": {"code": "return { items: [1, 2] };"}},
    {"id": "loop", "name": "Loop", "type": "loop", "position": {"x": 0, "y": 200},
     "config": {"inputField": "data.result.items"}},
    {"id": "approve", "name": "Approve", "type": "wait-for-approval", "position": {"x": 0, "y": 300},
     "config": {"reason": "Should never pause", "expiresIn": 1}}
  ],
  "connections": [
    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "items", "targetInput": "data"},
    {"sourceNodeId": "items", "sourceOutput": "main", "targetNodeId": "loop", "targetInput": "data"},
    {"sourceNodeId": "loop", "sourceOutput": "loop", "targetNodeId": "approve", "targetInput": "data"}
  ],
  "expect": {
    "approve": {"status": "failed", "errorContains": "cannot be used inside a Loop body"}
  }
}

Before running it, check what a loop body node actually receives - the loop node's config in this fixture uses inputField, and the loop variable names default to item and index. If the node's guard does not fire, print the input a body node sees and adjust the guard to match reality rather than deleting the guard.

  • Step 4: Run both

Run:

python3 scripts/verify-node.py tests/nodes/wait-for-approval.json
python3 scripts/verify-node.py tests/nodes/wait-for-approval-loop.json

Expected: PASS for both.

If the first hangs at "never reached waiting", the pause marker is not stopping the walk - check the execution's status in the database. If it reaches waiting but resume fails, the error from the resume endpoint names which precondition rejected it.

  • [ ] Step 5: Commit

    git add nodes/core/wait-for-approval.js tests/nodes/wait-for-approval.json tests/nodes/wait-for-approval-loop.json
    git commit -m "feat: a Wait for Approval node that pauses an execution for a person"
    

Task 11: Documentation and the full suite

Files:

  • Modify: docs/nodes.md
  • Modify: docs/node-roadmap.md

  • [ ] Step 1: Run the entire suite

Run:

pass=0; fail=0; failed=""
for f in tests/nodes/*.json; do
  if python3 scripts/verify-node.py "$f" >/tmp/vn.log 2>&1; then pass=$((pass+1));
  else fail=$((fail+1)); failed="$failed $(basename $f)"; fi
done
echo "PASS=$pass FAIL=$fail${failed:+ FAILED:$failed}"

Expected: zero failures, roughly 36 fixtures. If any fail, STOP and report which and why - a failure here means something regressed in already-reviewed work, and that is the most valuable thing this task can find. Do not fix it yourself and do not skip it.

  • Step 2: Document the new APIs

In docs/nodes.md, in the filesystem section, add readdir beside the existing calls:

// List a directory, one level
const listing = smartbotic.fs.readdir('/var/spool/incoming');
// { success: true, entries: [{ name, path, size, modifiedAt, isDirectory }] }

Then add a new section after the outputs sections:

## Pausing for an Answer

A node pauses its execution by returning a `_pause` marker, the way a branching
node returns `_activeBranch`:

```javascript
return {
    token: token,
    reason: 'Approve the refund',
    expiresAt: Date.now() + 86400000,
    _pause: { token: token, reason: 'Approve the refund', expiresAt: expiresAt }
};
```

The engine stops the walk, stores the execution as `waiting` with everything
computed so far, and returns. Nothing downstream runs. The marker is stripped
from the stored output, so a reader sees the request rather than the mechanism.

Answering it continues the run:

```
POST /api/v1/executions/{id}/resume
{ "token": "...", "approved": true, "data": { "note": "looks fine" } }
```

The paused node's output becomes that payload, and the walk continues from
there. Nodes that already ran are not run again. The workflow is rebuilt from
the snapshot stored with the execution, not from the workflow as it stands now,
because it may have been edited while the approval waited.

`GET /api/v1/executions/pending` lists executions waiting for an answer. It
deliberately omits the token: listing is a weaker permission than approving.

Two limits are deliberate. A pause inside a Loop body cannot be resumed, because
loop iteration state is not part of the stored execution, so a node must refuse
to pause there rather than record something unanswerable. And a webhook cannot
wait for an approval - the HTTP request is still open and its deadline is
35 seconds - so a webhook-triggered workflow that pauses returns immediately
with the execution id.

## Answering a Webhook

A webhook-triggered workflow returns the last node's output as JSON by default.
To control the response, return a `_webhookResponse` marker from any node:

```javascript
return {
    _webhookResponse: {
        status: 201,
        headers: { 'Content-Type': 'application/json' },
        body: { id: created.id }
    }
};
```

Any node may set it, not only the last one to run, so adding a node to the end
of a workflow cannot silently change what its API returns. If several set it,
the last one wins and the runner logs that it happened.
  • Step 3: Update the roadmap

In docs/node-roadmap.md, replace the Tier 2 section with:

## Tier 2 - triggers

Built. `respond-to-webhook`, `database-change`, `file-watch` and
`wait-for-approval` all live in `nodes/`, with fixtures under `tests/nodes/`.

Three platform changes came with them: `fs.readdir`, so a node can see what is
in a directory; a `_webhookResponse` marker any node can return to set the
status, headers and body a webhook replies with; and resumable executions - an
execution can pause on a `_pause` marker and be continued later through
`POST /api/v1/executions/{id}/resume`, rebuilt from the workflow snapshot stored
with it.

Two things the original entries claimed turned out not to hold. The webhook
controller never fired and forgot: it already waited for completion and returned
a body, and what was missing was control over the status code, the headers, and
which node decides the response. And File Watch could not be pure JavaScript,
because the filesystem API had no way to list a directory.

Cron and interval scheduling are already covered by `schedule-trigger`,
including its overlap policy.

Also update the node count in the "What exists today" section: count the files with find nodes -name '*.js' | wc -l and use the real number rather than assuming.

  • [ ] Step 4: Commit

    git add docs/nodes.md docs/node-roadmap.md
    git commit -m "docs: pausing for an answer, answering a webhook, and fs.readdir"
    

Notes for the implementer

A node missing after migrate. The schema parser is regex-based and fails silently. Check /tmp/webserver.log for Failed to parse configSchema. Usual causes are a // comment inside a schema literal or a double quote inside a single-quoted string.

The runner surviving pkill. It has happened in every C++ task so far. Always confirm 9011 is free before relaunching, and kill by PID if it is not. A stale runner holds the old binary and the old node registry, and the resulting failures point nowhere near the cause.

Testing a pause by hand. Start a workflow with an approval node, then:

TOKEN=$(curl -s http://localhost:8090/api/v1/auth/login -H "Content-Type: application/json" \
  -d '{"username":"admin","password":"admin"}' | jq -r '.accessToken')
curl -s http://localhost:8090/api/v1/executions/pending -H "Authorization: Bearer $TOKEN" | jq
curl -s http://localhost:8090/api/v1/executions/<id> -H "Authorization: Bearer $TOKEN" \
  | jq '.status, .pausedNodeId, .nodeExecutions[] | select(.nodeId=="approve") | .output.token'

The token comes from the node's output on the execution detail, not from the pending listing.

If resume returns FAILED_PRECONDITION. The message names which check rejected it: not waiting, token mismatch, expired, no snapshot, or no paused node. Each means something different, which is why they are separate messages rather than one generic failure.