Browse Source

feat: a node that switches workflows on and off, without the REST API

Re-enabling a workflow from a workflow meant logging in to this instance's own
REST API: a credential holding the admin password, a base64 decoder in a code
node because QuickJS has no atob, a POST to /auth/login and another to
/activate. It worked, and it was wrong - a password stored to talk to ourselves,
and on the first attempt the decoded password landed in a node's output, which
is kept on the execution record for anyone who can read one.

The Workflow Control node does it directly: status, activate, deactivate.

It goes to the webserver over the gRPC channel the runner already holds for
credentials, on the same port, because the webserver owns the scheduler.
Writing "active" into the workflows collection would set a flag nothing acts
on - trigger registration happens in updateScheduledTriggers, and a workflow
enabled by a bare database write sits there looking on and never runs. The
proof is in the log: a node-driven activate is followed by "Registered workflow
... for scheduled execution", which a database write could never produce.

Both paths now share WorkflowController::setActive, so the REST endpoint and
the node cannot drift.

Three decisions worth keeping:

Same project only, for reading as well as switching. A workflow able to switch
off anything on the instance would be a way around every project boundary the
rest of the system keeps, and unlike a person it cannot be asked whether it
meant to.

A workflow cannot control itself. Switching itself off mid-run would leave the
run finishing under a workflow that is no longer registered, and one that
disables itself has no way back until somebody notices.

Activating clears the failure counter, as the REST endpoint does. A workflow
re-enabled at 3 of 3 failures would otherwise switch off again on its very next
failure, however long it had been running happily in between.

Verified against the live instance: status reads the state, deactivate and
activate change it and the scheduler follows, activating twice reports
changed=false rather than pretending, and a workflow in another project is
refused by name for both reading and switching. The re-enable workflow is
rebuilt on it and its password credential is deleted.

Node suite: 90 passed, 0 failed, 2 skipped - both sd.cpp cases needing a model
loaded on that server, which has unloaded since the last run.
fszontagh 1 tháng trước cách đây
mục cha
commit
3ced2cb73e

+ 5 - 1
CMakeLists.txt

@@ -141,6 +141,8 @@ ${CMAKE_CURRENT_BINARY_DIR}/proto/workflow.pb.cc
     ${CMAKE_CURRENT_BINARY_DIR}/proto/runner.grpc.pb.cc
     ${CMAKE_CURRENT_BINARY_DIR}/proto/credentials.pb.cc
     ${CMAKE_CURRENT_BINARY_DIR}/proto/credentials.grpc.pb.cc
+    ${CMAKE_CURRENT_BINARY_DIR}/proto/workflow_control.pb.cc
+    ${CMAKE_CURRENT_BINARY_DIR}/proto/workflow_control.grpc.pb.cc
 )
 target_include_directories(smartbotic_proto PUBLIC
     ${CMAKE_CURRENT_BINARY_DIR}
@@ -158,7 +160,7 @@ set(PROTO_DIR ${CMAKE_CURRENT_SOURCE_DIR}/proto)
 set(PROTO_OUT_DIR ${CMAKE_CURRENT_BINARY_DIR}/proto)
 file(MAKE_DIRECTORY ${PROTO_OUT_DIR})
 
-set(PROTO_FILES common workflow runner credentials)
+set(PROTO_FILES common workflow runner credentials workflow_control)
 foreach(PROTO ${PROTO_FILES})
     add_custom_command(
         OUTPUT
@@ -197,6 +199,7 @@ add_executable(smartbotic-webserver
     src/webserver/scheduler/database_watcher.cpp
     src/webserver/grpc/node_sync_service.cpp
     src/webserver/grpc/credential_service.cpp
+    src/webserver/grpc/workflow_control_service.cpp
     src/webserver/api/auth_controller.cpp
     src/webserver/api/credential_controller.cpp
     src/webserver/api/certificate_controller.cpp
@@ -248,6 +251,7 @@ set(RUNNER_SOURCES
     src/runner/node_registry.cpp
     src/runner/error_handler.cpp
     src/runner/runner_service.cpp
+    src/runner/workflow_control_client.cpp
     src/runner/collection_permissions.cpp
     src/runner/engine/script_engine.cpp
     src/runner/imap/imap_client.cpp

+ 116 - 0
nodes/core/workflow-control.js

@@ -0,0 +1,116 @@
+/**
+ * @node workflow-control
+ * @name Workflow Control
+ * @category flow-control
+ * @version 1.0.0
+ * @description Switch another workflow on or off, or read whether it is running
+ * @icon power
+ */
+
+const configSchema = {
+    type: 'object',
+    properties: {
+        operation: {
+            type: 'string',
+            title: 'Operation',
+            description: 'Activate switches the workflow on and registers its triggers. Deactivate switches it off. Status only reads, and changes nothing',
+            enum: ['status', 'activate', 'deactivate'],
+            enumLabels: ['Read its state', 'Switch it on', 'Switch it off'],
+            default: 'status'
+        },
+        workflowId: {
+            type: 'string',
+            title: 'Workflow',
+            description: 'The workflow to control. It has to be in the same project as this one',
+            dynamicOptions: {
+                source: 'workflows',
+                labelField: 'name',
+                valueField: 'id'
+            }
+        },
+        failIfRefused: {
+            type: 'boolean',
+            title: 'Fail If Refused',
+            description: 'Fail the node when the workflow cannot be controlled - it is in another project, or no longer exists. Turn off to carry on and decide in the workflow, reading success from the output',
+            default: true
+        }
+    },
+    required: ['operation', 'workflowId']
+};
+
+const inputSchema = {
+    type: 'object',
+    properties: {
+        data: { type: 'any' }
+    }
+};
+
+const outputSchema = {
+    type: 'object',
+    properties: {
+        success: { type: 'boolean', description: 'False when the request was refused, and only reachable when Fail If Refused is off' },
+        error: { type: 'string', description: 'Why it was refused' },
+        name: { type: 'string', description: 'The controlled workflow\'s name' },
+        active: { type: 'boolean', description: 'Whether it is on - after the change, for activate and deactivate' },
+        changed: { type: 'boolean', description: 'False when it was already in that state, so nothing happened' },
+        consecutiveFailures: { type: 'number', description: 'How many times in a row it has failed. Status only' },
+        deactivateAfterFailures: { type: 'number', description: 'How many consecutive failures switch it off, or 0 when it never does. Status only' },
+        deactivatedReason: { type: 'string', description: 'Why it switched itself off, when it did. Status only' },
+        publishedVersion: { type: 'number', description: 'Status only' }
+    }
+};
+
+async function execute(config, input, context) {
+    const workflowId = String(config.workflowId || '').trim();
+    if (!workflowId) {
+        throw new Error('Workflow Control: choose the workflow to control');
+    }
+
+    const operation = config.operation || 'status';
+    const failIfRefused = config.failIfRefused !== false;
+
+    // The webserver decides whether this is allowed and does the switching -
+    // it owns the scheduler, and a workflow enabled by writing the flag
+    // straight into the database would look on and never run.
+    let result;
+    if (operation === 'status') {
+        result = smartbotic.workflows.get(workflowId);
+    } else {
+        result = smartbotic.workflows.setActive(workflowId, operation === 'activate');
+    }
+
+    if (!result.success) {
+        if (failIfRefused) {
+            throw new Error('Workflow Control: ' + (result.error || 'the request was refused'));
+        }
+        smartbotic.log.warn('Workflow Control: ' + (result.error || 'refused'));
+        return { success: false, error: result.error || 'refused' };
+    }
+
+    if (operation === 'status') {
+        smartbotic.log.info('Workflow Control: "' + result.name + '" is ' +
+                            (result.active ? 'on' : 'off'));
+        return {
+            success: true,
+            error: '',
+            name: result.name,
+            active: result.active,
+            consecutiveFailures: result.consecutiveFailures,
+            deactivateAfterFailures: result.deactivateAfterFailures,
+            deactivatedReason: result.deactivatedReason,
+            publishedVersion: result.publishedVersion
+        };
+    }
+
+    smartbotic.log.info('Workflow Control: switched ' + workflowId + ' ' +
+                        (result.active ? 'on' : 'off') +
+                        (result.changed ? '' : ' (it was already)'));
+    return {
+        success: true,
+        error: '',
+        active: result.active,
+        changed: result.changed
+    };
+}
+
+module.exports = { configSchema, inputSchema, outputSchema, execute };

+ 54 - 0
proto/workflow_control.proto

@@ -0,0 +1,54 @@
+syntax = "proto3";
+
+package smartbotic.proto;
+
+// Lets a running workflow switch another workflow on or off, and read its
+// state, without going through the REST API.
+//
+// It lives on the webserver rather than in the runner because the webserver
+// owns the scheduler. Writing "active" straight into the workflows collection
+// would set a flag nothing acts on: registering and unregistering triggers
+// happens here, and a workflow switched on by a database write would sit there
+// looking enabled and never run.
+//
+// Served on the same port as CredentialService, and for the same reason - the
+// runner already holds a channel to it, and a request says which workflow is
+// asking so the webserver can decide whether to allow it.
+service WorkflowControlService {
+    // Switch a workflow on or off. Registers or unregisters its triggers as
+    // part of the same call.
+    rpc SetWorkflowActive(SetWorkflowActiveRequest) returns (SetWorkflowActiveResponse);
+
+    // What a workflow's state is: whether it is on, and how close it is to
+    // being switched off by consecutive failures.
+    rpc GetWorkflowState(GetWorkflowStateRequest) returns (GetWorkflowStateResponse);
+}
+
+message SetWorkflowActiveRequest {
+    string workflow_id = 1;             // the workflow to switch
+    bool active = 2;
+    string requesting_workflow_id = 3;  // who is asking, for access control
+}
+
+message SetWorkflowActiveResponse {
+    bool success = 1;
+    bool active = 2;                    // the state afterwards
+    bool changed = 3;                   // false when it was already in that state
+    string error = 4;
+}
+
+message GetWorkflowStateRequest {
+    string workflow_id = 1;
+    string requesting_workflow_id = 2;
+}
+
+message GetWorkflowStateResponse {
+    bool success = 1;
+    string name = 2;
+    bool active = 3;
+    int32 consecutive_failures = 4;
+    int32 deactivate_after_failures = 5;  // 0 when the workflow never switches itself off
+    string deactivated_reason = 6;        // why it switched itself off, when it did
+    int64 published_version = 7;
+    string error = 8;
+}

+ 93 - 0
src/runner/engine/script_engine.cpp

@@ -4435,6 +4435,98 @@ static void setupSmtpAPI(JSContext* ctx, JSValue smartbotic) {
     JS_SetPropertyStr(ctx, smartbotic, "smtp", smtp_obj);
 }
 
+// workflows.setActive(id, active) - switch another workflow on or off
+static JSValue js_workflows_set_active(JSContext* ctx, JSValue this_val, int argc, JSValue* argv) {
+    const ScriptContext* script_ctx = static_cast<const ScriptContext*>(JS_GetContextOpaque(ctx));
+    JSValue response = JS_NewObject(ctx);
+
+    if (!script_ctx || !script_ctx->workflow_set_active) {
+        JS_SetPropertyStr(ctx, response, "success", JS_FALSE);
+        JS_SetPropertyStr(ctx, response, "error",
+                          JS_NewString(ctx, "Workflow control is not available on this runner"));
+        return response;
+    }
+    if (argc < 1 || !JS_IsString(argv[0])) {
+        JS_SetPropertyStr(ctx, response, "success", JS_FALSE);
+        JS_SetPropertyStr(ctx, response, "error",
+                          JS_NewString(ctx, "workflows.setActive needs a workflow id"));
+        return response;
+    }
+
+    const char* id = JS_ToCString(ctx, argv[0]);
+    const std::string workflow_id = id ? id : "";
+    if (id) JS_FreeCString(ctx, id);
+    const bool active = argc > 1 ? static_cast<bool>(JS_ToBool(ctx, argv[1])) : true;
+
+    auto result = script_ctx->workflow_set_active(workflow_id, active);
+    if (result.failed()) {
+        JS_SetPropertyStr(ctx, response, "success", JS_FALSE);
+        JS_SetPropertyStr(ctx, response, "error", JS_NewString(ctx, result.error().message().c_str()));
+        return response;
+    }
+
+    JS_SetPropertyStr(ctx, response, "success", JS_TRUE);
+    JS_SetPropertyStr(ctx, response, "active", active ? JS_TRUE : JS_FALSE);
+    // False when the workflow was already in that state - nothing happened,
+    // which is worth being able to tell apart from having done something.
+    JS_SetPropertyStr(ctx, response, "changed", result.value() ? JS_TRUE : JS_FALSE);
+    return response;
+}
+
+// workflows.get(id) - read another workflow's state
+static JSValue js_workflows_get(JSContext* ctx, JSValue this_val, int argc, JSValue* argv) {
+    const ScriptContext* script_ctx = static_cast<const ScriptContext*>(JS_GetContextOpaque(ctx));
+    JSValue response = JS_NewObject(ctx);
+
+    if (!script_ctx || !script_ctx->workflow_get_state) {
+        JS_SetPropertyStr(ctx, response, "success", JS_FALSE);
+        JS_SetPropertyStr(ctx, response, "error",
+                          JS_NewString(ctx, "Workflow control is not available on this runner"));
+        return response;
+    }
+    if (argc < 1 || !JS_IsString(argv[0])) {
+        JS_SetPropertyStr(ctx, response, "success", JS_FALSE);
+        JS_SetPropertyStr(ctx, response, "error",
+                          JS_NewString(ctx, "workflows.get needs a workflow id"));
+        return response;
+    }
+
+    const char* id = JS_ToCString(ctx, argv[0]);
+    const std::string workflow_id = id ? id : "";
+    if (id) JS_FreeCString(ctx, id);
+
+    auto result = script_ctx->workflow_get_state(workflow_id);
+    if (result.failed()) {
+        JS_SetPropertyStr(ctx, response, "success", JS_FALSE);
+        JS_SetPropertyStr(ctx, response, "error", JS_NewString(ctx, result.error().message().c_str()));
+        return response;
+    }
+
+    const auto& state = result.value();
+    JS_SetPropertyStr(ctx, response, "success", JS_TRUE);
+    JS_SetPropertyStr(ctx, response, "name", JS_NewString(ctx, state.name.c_str()));
+    JS_SetPropertyStr(ctx, response, "active", state.active ? JS_TRUE : JS_FALSE);
+    JS_SetPropertyStr(ctx, response, "consecutiveFailures",
+                      JS_NewInt32(ctx, state.consecutive_failures));
+    JS_SetPropertyStr(ctx, response, "deactivateAfterFailures",
+                      JS_NewInt32(ctx, state.deactivate_after_failures));
+    JS_SetPropertyStr(ctx, response, "deactivatedReason",
+                      JS_NewString(ctx, state.deactivated_reason.c_str()));
+    JS_SetPropertyStr(ctx, response, "publishedVersion",
+                      JS_NewInt64(ctx, state.published_version));
+    return response;
+}
+
+// Setup workflow control API on smartbotic namespace
+static void setupWorkflowsAPI(JSContext* ctx, JSValue smartbotic) {
+    JSValue workflows = JS_NewObject(ctx);
+    JS_SetPropertyStr(ctx, workflows, "setActive",
+                      JS_NewCFunction(ctx, js_workflows_set_active, "setActive", 2));
+    JS_SetPropertyStr(ctx, workflows, "get",
+                      JS_NewCFunction(ctx, js_workflows_get, "get", 1));
+    JS_SetPropertyStr(ctx, smartbotic, "workflows", workflows);
+}
+
 // Setup credentials API on smartbotic namespace
 static void setupCredentialsAPI(JSContext* ctx, JSValue smartbotic) {
     JSValue credentials = JS_NewObject(ctx);
@@ -4759,6 +4851,7 @@ ScriptResult ScriptEngine::execute(const std::string& script, const ScriptContex
     if (JS_IsObject(smartbotic)) {
         setupStorageAPI(context_, smartbotic);
         setupCredentialsAPI(context_, smartbotic);
+        setupWorkflowsAPI(context_, smartbotic);
         setupImapAPI(context_, smartbotic);
         setupSmtpAPI(context_, smartbotic);
 #ifdef MYSQL_SUPPORT

+ 13 - 0
src/runner/engine/script_engine.hpp

@@ -192,6 +192,19 @@ struct ScriptContext {
     std::function<common::Result<ScriptFileInfo>(const std::string&)> storage_get_file_info;
     std::function<common::Result<void>(const std::string&)> storage_delete_file;
 
+    // Switching another workflow on or off, and reading its state. Goes to the
+    // webserver, which owns the scheduler - see WorkflowControlClient.
+    struct WorkflowStateInfo {
+        std::string name;
+        bool active = false;
+        int consecutive_failures = 0;
+        int deactivate_after_failures = 0;
+        std::string deactivated_reason;
+        int64_t published_version = 0;
+    };
+    std::function<common::Result<bool>(const std::string&, bool)> workflow_set_active;
+    std::function<common::Result<WorkflowStateInfo>(const std::string&)> workflow_get_state;
+
     // Credential operations
     std::function<common::Result<CredentialAuth>(const std::string&)> credentials_get;
     std::function<common::Result<ImapCredential>(const std::string&)> credentials_get_imap;

+ 28 - 0
src/runner/runner_service.cpp

@@ -7,6 +7,7 @@
 #include <sys/resource.h>
 #include <fstream>
 #include <curl/curl.h>
+#include "workflow_control_client.hpp"
 #ifdef MYSQL_SUPPORT
 #include "runner/mysql/mysql_client.hpp"
 #endif
@@ -495,6 +496,33 @@ RunnerService::RunnerService(const RunnerServiceConfig& config)
             return auth;
         });
 
+    // Switching workflows on and off. Shares the webserver address the
+    // credential client already uses - it is the same server, and a second
+    // setting would be another thing to keep in step.
+    workflow_control_client_ =
+        std::make_unique<WorkflowControlClient>(config_.credential_service_address);
+    engine_->setWorkflowSetActiveCallback(
+        [this](const std::string& target_id, bool active, const std::string& requesting)
+            -> common::Result<bool> {
+            return workflow_control_client_->setActive(target_id, active, requesting);
+        });
+    engine_->setWorkflowGetStateCallback(
+        [this](const std::string& target_id, const std::string& requesting)
+            -> common::Result<engine::ScriptContext::WorkflowStateInfo> {
+            auto state = workflow_control_client_->getState(target_id, requesting);
+            if (state.failed()) {
+                return state.error();
+            }
+            engine::ScriptContext::WorkflowStateInfo info;
+            info.name = state.value().name;
+            info.active = state.value().active;
+            info.consecutive_failures = state.value().consecutive_failures;
+            info.deactivate_after_failures = state.value().deactivate_after_failures;
+            info.deactivated_reason = state.value().deactivated_reason;
+            info.published_version = state.value().published_version;
+            return info;
+        });
+
     // The mTLS identity a workflow presents, fetched over the same mediated
     // service the other secrets use.
     engine_->setClientCertificateCallback(

+ 3 - 0
src/runner/runner_service.hpp

@@ -15,6 +15,8 @@
 
 namespace smartbotic::runner {
 
+class WorkflowControlClient;
+
 // Forward declaration
 struct RunnerMetrics;
 
@@ -124,6 +126,7 @@ private:
 
     std::unique_ptr<storage::StorageClient> storage_;
     std::unique_ptr<credentials::CredentialClient> credential_client_;
+    std::unique_ptr<WorkflowControlClient> workflow_control_client_;
     std::unique_ptr<NodeRegistry> registry_;
     std::unique_ptr<WorkflowEngine> engine_;
     std::unique_ptr<RunnerServiceImpl> service_impl_;

+ 72 - 0
src/runner/workflow_control_client.cpp

@@ -0,0 +1,72 @@
+#include "workflow_control_client.hpp"
+
+#include <chrono>
+
+#include "logging/logger.hpp"
+
+namespace smartbotic::runner {
+
+using common::Error;
+using common::ErrorCode;
+using common::Result;
+
+WorkflowControlClient::WorkflowControlClient(const std::string& webserver_address, int timeout_ms)
+    : timeout_ms_(timeout_ms) {
+    auto channel = ::grpc::CreateChannel(webserver_address, ::grpc::InsecureChannelCredentials());
+    stub_ = proto::WorkflowControlService::NewStub(channel);
+}
+
+Result<bool> WorkflowControlClient::setActive(const std::string& workflow_id, bool active,
+                                              const std::string& requesting_workflow_id) {
+    proto::SetWorkflowActiveRequest request;
+    request.set_workflow_id(workflow_id);
+    request.set_active(active);
+    request.set_requesting_workflow_id(requesting_workflow_id);
+
+    proto::SetWorkflowActiveResponse response;
+    ::grpc::ClientContext context;
+    context.set_deadline(std::chrono::system_clock::now() + std::chrono::milliseconds(timeout_ms_));
+
+    auto status = stub_->SetWorkflowActive(&context, request, &response);
+    if (!status.ok()) {
+        return Error(ErrorCode::Unavailable,
+                     "Could not reach the workflow control service: " + status.error_message());
+    }
+    if (!response.success()) {
+        // The refusal is the server's, and it explains itself - passed through
+        // rather than replaced, so the workflow author reads why.
+        return Error(ErrorCode::PermissionDenied, response.error());
+    }
+    return response.changed();
+}
+
+Result<WorkflowState> WorkflowControlClient::getState(const std::string& workflow_id,
+                                                      const std::string& requesting_workflow_id) {
+    proto::GetWorkflowStateRequest request;
+    request.set_workflow_id(workflow_id);
+    request.set_requesting_workflow_id(requesting_workflow_id);
+
+    proto::GetWorkflowStateResponse response;
+    ::grpc::ClientContext context;
+    context.set_deadline(std::chrono::system_clock::now() + std::chrono::milliseconds(timeout_ms_));
+
+    auto status = stub_->GetWorkflowState(&context, request, &response);
+    if (!status.ok()) {
+        return Error(ErrorCode::Unavailable,
+                     "Could not reach the workflow control service: " + status.error_message());
+    }
+    if (!response.success()) {
+        return Error(ErrorCode::PermissionDenied, response.error());
+    }
+
+    WorkflowState state;
+    state.name = response.name();
+    state.active = response.active();
+    state.consecutive_failures = response.consecutive_failures();
+    state.deactivate_after_failures = response.deactivate_after_failures();
+    state.deactivated_reason = response.deactivated_reason();
+    state.published_version = response.published_version();
+    return state;
+}
+
+} // namespace smartbotic::runner

+ 43 - 0
src/runner/workflow_control_client.hpp

@@ -0,0 +1,43 @@
+#pragma once
+
+#include <memory>
+#include <string>
+
+#include <grpcpp/grpcpp.h>
+
+#include "common/error.hpp"
+#include "proto/workflow_control.grpc.pb.h"
+
+namespace smartbotic::runner {
+
+// What a workflow is allowed to know about another workflow.
+struct WorkflowState {
+    std::string name;
+    bool active = false;
+    int consecutive_failures = 0;
+    int deactivate_after_failures = 0;
+    std::string deactivated_reason;
+    int64_t published_version = 0;
+};
+
+// Switches workflows on and off through the webserver, which owns the
+// scheduler. Deliberately not a database write: "active" is a flag the
+// scheduler acts on when told, and a workflow enabled by writing that field
+// would look on and never run.
+class WorkflowControlClient {
+public:
+    WorkflowControlClient(const std::string& webserver_address, int timeout_ms = 30000);
+
+    // Returns whether the state actually changed.
+    common::Result<bool> setActive(const std::string& workflow_id, bool active,
+                                   const std::string& requesting_workflow_id);
+
+    common::Result<WorkflowState> getState(const std::string& workflow_id,
+                                           const std::string& requesting_workflow_id);
+
+private:
+    std::unique_ptr<proto::WorkflowControlService::Stub> stub_;
+    int timeout_ms_;
+};
+
+} // namespace smartbotic::runner

+ 19 - 0
src/runner/workflow_engine.cpp

@@ -2029,6 +2029,25 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
         return imap_credential_callback_(credential_id, workflow.id);
     };
 
+    // Workflow control. The workflow being run is named here rather than
+    // taken from the node, so a node cannot claim to be a different one.
+    ctx.workflow_set_active = [this, &workflow](const std::string& target_id, bool active)
+        -> common::Result<bool> {
+        if (!workflow_set_active_callback_) {
+            return common::Error(common::ErrorCode::Unavailable,
+                                 "Workflow control is not configured on this runner");
+        }
+        return workflow_set_active_callback_(target_id, active, workflow.id);
+    };
+    ctx.workflow_get_state = [this, &workflow](const std::string& target_id)
+        -> common::Result<engine::ScriptContext::WorkflowStateInfo> {
+        if (!workflow_get_state_callback_) {
+            return common::Error(common::ErrorCode::Unavailable,
+                                 "Workflow control is not configured on this runner");
+        }
+        return workflow_get_state_callback_(target_id, workflow.id);
+    };
+
     // SMTP Credential API callback
     ctx.credentials_get_smtp = [this, &workflow](const std::string& credential_id)
         -> common::Result<engine::SmtpCredential> {

+ 11 - 0
src/runner/workflow_engine.hpp

@@ -224,6 +224,13 @@ using SmtpCredentialCallback = std::function<common::Result<engine::SmtpCredenti
 using ImapCredentialCallback = std::function<common::Result<engine::ImapCredential>(
     const std::string& credential_id, const std::string& workflow_id)>;
 
+// Switching another workflow on or off, and reading its state. The workflow
+// doing the asking is passed along so the webserver can decide whether it may.
+using WorkflowSetActiveCallback = std::function<common::Result<bool>(
+    const std::string& target_id, bool active, const std::string& requesting_workflow_id)>;
+using WorkflowGetStateCallback = std::function<common::Result<engine::ScriptContext::WorkflowStateInfo>(
+    const std::string& target_id, const std::string& requesting_workflow_id)>;
+
 // The mTLS identity a workflow presents, fetched through the credential
 // service like every other secret.
 using ClientCertificateCallback = std::function<common::Result<engine::TlsClientIdentity>(
@@ -298,6 +305,8 @@ public:
     // Set IMAP credential callback (called by runner service to provide IMAP credential access)
     void setImapCredentialCallback(ImapCredentialCallback callback) { imap_credential_callback_ = callback; }
     void setClientCertificateCallback(ClientCertificateCallback callback) { client_certificate_callback_ = callback; }
+    void setWorkflowSetActiveCallback(WorkflowSetActiveCallback callback) { workflow_set_active_callback_ = callback; }
+    void setWorkflowGetStateCallback(WorkflowGetStateCallback callback) { workflow_get_state_callback_ = callback; }
     void setSmtpCredentialCallback(SmtpCredentialCallback callback) { smtp_credential_callback_ = callback; }
 
     // Set MySQL credential callback (called by runner service to provide MySQL credential access)
@@ -546,6 +555,8 @@ private:
     CredentialAuthCallback credential_auth_callback_;
     ImapCredentialCallback imap_credential_callback_;
     ClientCertificateCallback client_certificate_callback_;
+    WorkflowSetActiveCallback workflow_set_active_callback_;
+    WorkflowGetStateCallback workflow_get_state_callback_;
     SmtpCredentialCallback smtp_credential_callback_;
     MysqlCredentialCallback mysql_credential_callback_;
     MysqlQueryCallback mysql_query_callback_;

+ 30 - 0
src/webserver/api/workflow_controller.cpp

@@ -1145,6 +1145,36 @@ void WorkflowController::transferWorkflowOwner(const httplib::Request& req, http
     }
 }
 
+common::Result<bool> WorkflowController::setActive(const std::string& workflow_id, bool active) {
+    auto before = storage_.get("workflows", workflow_id);
+    if (before.failed()) {
+        return common::Error(common::ErrorCode::NotFound, "No workflow " + workflow_id);
+    }
+    const bool was_active = before.value().value("active", false);
+
+    nlohmann::json patch = {{"active", active}, {"updatedAt", TimeUtils::nowMs()}};
+    if (active) {
+        // Switching on clears the failure counter and the reason, the same way
+        // the REST endpoint does. Leaving a count of three behind would switch
+        // the workflow off again on its very next failure, however long it had
+        // been running happily in between.
+        patch["consecutiveFailures"] = 0;
+        patch["deactivatedReason"] = "";
+    }
+
+    auto result = storage_.update("workflows", workflow_id, patch, 0, true);
+    if (result.failed()) {
+        return common::Error(common::ErrorCode::Internal, result.error().message());
+    }
+
+    updateScheduledTriggers(workflow_id, active);
+    ws_server_.broadcast(active ? "workflows.activated" : "workflows.deactivated",
+                         {{"id", workflow_id}});
+
+    LOG_INFO("Workflow {} switched {} by a workflow", workflow_id, active ? "on" : "off");
+    return was_active != active;
+}
+
 void WorkflowController::updateScheduledTriggers(const std::string& workflow_id, bool activate) {
     if (!activate) {
         scheduler_.unregisterWorkflow(workflow_id);

+ 12 - 0
src/webserver/api/workflow_controller.hpp

@@ -106,6 +106,18 @@ private:
     DispatchPool dispatch_pool_;
 
     // Helper to check for scheduled triggers and register/unregister with scheduler
+public:
+    // Switch a workflow on or off from inside the process, doing everything the
+    // REST endpoint does - including registering or unregistering its triggers.
+    //
+    // Public so the gRPC control service can use the same path. Writing
+    // "active" into the collection instead would set a flag nothing acts on:
+    // trigger registration happens in updateScheduledTriggers below, and a
+    // workflow enabled by a bare database write sits there looking on and never
+    // runs. Returns whether the state actually changed.
+    common::Result<bool> setActive(const std::string& workflow_id, bool active);
+
+private:
     void updateScheduledTriggers(const std::string& workflow_id, bool activate);
     int getScheduledInterval(const nlohmann::json& workflow);
 

+ 13 - 0
src/webserver/grpc/credential_service.cpp

@@ -204,6 +204,16 @@ CredentialServer::~CredentialServer() {
     stop();
 }
 
+void CredentialServer::addService(::grpc::Service* service) {
+    if (running_.load()) {
+        LOG_ERROR("A service was added after the gRPC server started; it will not be reachable");
+        return;
+    }
+    if (service) {
+        extra_services_.push_back(service);
+    }
+}
+
 void CredentialServer::start() {
     if (running_.load()) {
         return;
@@ -217,6 +227,9 @@ void CredentialServer::start() {
         ::grpc::ServerBuilder builder;
         builder.AddListeningPort(server_address, ::grpc::InsecureServerCredentials());
         builder.RegisterService(service_impl_.get());
+        for (auto* extra : extra_services_) {
+            builder.RegisterService(extra);
+        }
 
         server_ = builder.BuildAndStart();
         LOG_INFO("Credential gRPC server listening on {}", server_address);

+ 14 - 1
src/webserver/grpc/credential_service.hpp

@@ -1,5 +1,7 @@
 #pragma once
 
+#include <vector>
+
 #include <memory>
 #include <thread>
 #include <atomic>
@@ -53,7 +55,12 @@ private:
     credentials::CredentialStore& credential_store_;
 };
 
-// Credential service server wrapper
+// The gRPC endpoint the runner calls back on.
+//
+// It serves more than credentials now - workflow control shares the port,
+// because the runner already holds a channel to it and a second port would be
+// another thing to configure, open and get wrong. Extra services are registered
+// through addService before start().
 class CredentialServer {
 public:
     CredentialServer(credentials::CredentialStore& credential_store, int port);
@@ -65,8 +72,14 @@ public:
     CredentialServiceImpl& service() { return *service_impl_; }
     int port() const { return port_; }
 
+    // Register another service on the same port. Must be called before start();
+    // gRPC builds its service registry when the server is built, so anything
+    // added afterwards would be silently unreachable.
+    void addService(::grpc::Service* service);
+
 private:
     int port_;
+    std::vector<::grpc::Service*> extra_services_;
     std::unique_ptr<CredentialServiceImpl> service_impl_;
     std::unique_ptr<::grpc::Server> server_;
     std::thread server_thread_;

+ 127 - 0
src/webserver/grpc/workflow_control_service.cpp

@@ -0,0 +1,127 @@
+#include "workflow_control_service.hpp"
+
+#include "logging/logger.hpp"
+
+namespace smartbotic::webserver::grpc {
+
+WorkflowControlServiceImpl::WorkflowControlServiceImpl(storage::StorageClient& storage,
+                                                       SetActive set_active)
+    : storage_(storage), set_active_(std::move(set_active)) {}
+
+bool WorkflowControlServiceImpl::mayControl(const std::string& requesting_workflow_id,
+                                            const std::string& target_id,
+                                            std::string& error_out) {
+    if (target_id.empty()) {
+        error_out = "No workflow was named to control";
+        return false;
+    }
+    if (requesting_workflow_id.empty()) {
+        // Every real call carries it; an empty one means the runner did not
+        // know which workflow was asking, and an unattributed request is not
+        // one to trust with switching things off.
+        error_out = "The request did not say which workflow was asking";
+        return false;
+    }
+
+    // A workflow controlling itself is refused rather than allowed to switch
+    // itself off mid-run, which would leave a run finishing under a workflow
+    // that is no longer registered - and a workflow that disables itself has no
+    // way back short of somebody noticing.
+    if (requesting_workflow_id == target_id) {
+        error_out = "A workflow cannot switch itself on or off";
+        return false;
+    }
+
+    auto asking = storage_.get("workflows", requesting_workflow_id);
+    if (asking.failed()) {
+        error_out = "The asking workflow no longer exists";
+        return false;
+    }
+    auto target = storage_.get("workflows", target_id);
+    if (target.failed()) {
+        error_out = "No workflow " + target_id;
+        return false;
+    }
+
+    const std::string asking_project = asking.value().value("projectId", "");
+    const std::string target_project = target.value().value("projectId", "");
+    if (asking_project.empty() || asking_project != target_project) {
+        error_out = "\"" + target.value().value("name", target_id) + "\" is in another project, "
+                    "so this workflow cannot control it";
+        return false;
+    }
+    return true;
+}
+
+::grpc::Status WorkflowControlServiceImpl::SetWorkflowActive(
+    ::grpc::ServerContext*,
+    const proto::SetWorkflowActiveRequest* request,
+    proto::SetWorkflowActiveResponse* response) {
+
+    std::string refusal;
+    if (!mayControl(request->requesting_workflow_id(), request->workflow_id(), refusal)) {
+        response->set_success(false);
+        response->set_error(refusal);
+        LOG_WARN("Workflow {} refused control of {}: {}", request->requesting_workflow_id(),
+                 request->workflow_id(), refusal);
+        return ::grpc::Status::OK;
+    }
+
+    if (!set_active_) {
+        response->set_success(false);
+        response->set_error("Workflow control is not available on this server");
+        return ::grpc::Status::OK;
+    }
+
+    auto changed = set_active_(request->workflow_id(), request->active());
+    if (changed.failed()) {
+        response->set_success(false);
+        response->set_error(changed.error().message());
+        return ::grpc::Status::OK;
+    }
+
+    response->set_success(true);
+    response->set_active(request->active());
+    response->set_changed(changed.value());
+    LOG_INFO("Workflow {} switched {} {} by workflow {}", request->workflow_id(),
+             request->active() ? "on" : "off",
+             changed.value() ? "" : "(it was already)", request->requesting_workflow_id());
+    return ::grpc::Status::OK;
+}
+
+::grpc::Status WorkflowControlServiceImpl::GetWorkflowState(
+    ::grpc::ServerContext*,
+    const proto::GetWorkflowStateRequest* request,
+    proto::GetWorkflowStateResponse* response) {
+
+    std::string refusal;
+    // Reading is gated the same way as switching. What a workflow is called and
+    // how close it is to switching itself off is not much, but it is still
+    // another project's business.
+    if (!mayControl(request->requesting_workflow_id(), request->workflow_id(), refusal)) {
+        response->set_success(false);
+        response->set_error(refusal);
+        return ::grpc::Status::OK;
+    }
+
+    auto stored = storage_.get("workflows", request->workflow_id());
+    if (stored.failed()) {
+        response->set_success(false);
+        response->set_error("No workflow " + request->workflow_id());
+        return ::grpc::Status::OK;
+    }
+
+    const auto& doc = stored.value();
+    const auto settings = doc.value("settings", nlohmann::json::object());
+
+    response->set_success(true);
+    response->set_name(doc.value("name", ""));
+    response->set_active(doc.value("active", false));
+    response->set_consecutive_failures(doc.value("consecutiveFailures", 0));
+    response->set_deactivate_after_failures(settings.value("deactivateAfterFailures", 0));
+    response->set_deactivated_reason(doc.value("deactivatedReason", ""));
+    response->set_published_version(doc.value("publishedVersion", int64_t{0}));
+    return ::grpc::Status::OK;
+}
+
+} // namespace smartbotic::webserver::grpc

+ 51 - 0
src/webserver/grpc/workflow_control_service.hpp

@@ -0,0 +1,51 @@
+#pragma once
+
+#include <functional>
+#include <string>
+
+#include <grpcpp/grpcpp.h>
+
+#include "common/error.hpp"
+#include "proto/workflow_control.grpc.pb.h"
+#include "storage/storage_client.hpp"
+
+namespace smartbotic::webserver::grpc {
+
+// Lets a running workflow switch another workflow on or off.
+//
+// The switching itself is done by the callback, which the service is handed at
+// construction and which goes through WorkflowController - the one place that
+// registers and unregisters triggers. This service does the deciding, not the
+// doing: whether the caller is allowed to touch the workflow it named.
+class WorkflowControlServiceImpl final : public proto::WorkflowControlService::Service {
+public:
+    // Returns whether the state actually changed, or an error.
+    using SetActive = std::function<common::Result<bool>(const std::string& workflow_id,
+                                                          bool active)>;
+
+    WorkflowControlServiceImpl(storage::StorageClient& storage, SetActive set_active);
+
+    ::grpc::Status SetWorkflowActive(::grpc::ServerContext* context,
+                                     const proto::SetWorkflowActiveRequest* request,
+                                     proto::SetWorkflowActiveResponse* response) override;
+
+    ::grpc::Status GetWorkflowState(::grpc::ServerContext* context,
+                                    const proto::GetWorkflowStateRequest* request,
+                                    proto::GetWorkflowStateResponse* response) override;
+
+private:
+    // Whether the asking workflow may touch the named one.
+    //
+    // Same project, and nothing else. A workflow that could switch off anything
+    // on the instance would be a way around every project boundary the rest of
+    // the system keeps - and unlike a person, it cannot be asked whether it
+    // meant to. The error says which rule was hit, because "not allowed" with
+    // no reason is the kind of thing people work around rather than fix.
+    bool mayControl(const std::string& requesting_workflow_id, const std::string& target_id,
+                    std::string& error_out);
+
+    storage::StorageClient& storage_;
+    SetActive set_active_;
+};
+
+} // namespace smartbotic::webserver::grpc

+ 19 - 0
src/webserver/webserver_service.cpp

@@ -17,6 +17,7 @@
 #include "nodes/node_store.hpp"
 #include "grpc/node_sync_service.hpp"
 #include "grpc/credential_service.hpp"
+#include "grpc/workflow_control_service.hpp"
 #include "credentials/credential_store.hpp"
 #include "scheduler/workflow_scheduler.hpp"
 #include "common/time_utils.hpp"
@@ -361,6 +362,24 @@ void WebServerService::start() {
     node_sync_server_->start();
 
     // Start Credential gRPC server for runners
+    // Registered before the server starts - gRPC builds its service registry
+    // when the server is built, so a service added later would be silently
+    // unreachable.
+    //
+    // The callback resolves workflow_ctrl_ when it is called, not now: the
+    // gRPC server starts before the REST controllers are built, and capturing
+    // the controller here would capture a null.
+    workflow_control_service_ = std::make_unique<grpc::WorkflowControlServiceImpl>(
+        *storage_,
+        [this](const std::string& workflow_id, bool active) -> common::Result<bool> {
+            if (!workflow_ctrl_) {
+                return common::Error(common::ErrorCode::Unavailable,
+                                     "The server is still starting up");
+            }
+            return workflow_ctrl_->setActive(workflow_id, active);
+        });
+    credential_server_->addService(workflow_control_service_.get());
+
     credential_server_->start();
 
     // Start WebSocket server

+ 2 - 0
src/webserver/webserver_service.hpp

@@ -27,6 +27,7 @@ namespace smartbotic::webserver::nodes {
 namespace smartbotic::webserver::grpc {
     class NodeSyncServer;
     class CredentialServer;
+    class WorkflowControlServiceImpl;
 }
 
 namespace smartbotic::credentials {
@@ -160,6 +161,7 @@ private:
     std::unique_ptr<WebSocketServer> ws_server_;
     std::unique_ptr<grpc::NodeSyncServer> node_sync_server_;
     std::unique_ptr<grpc::CredentialServer> credential_server_;
+    std::unique_ptr<grpc::WorkflowControlServiceImpl> workflow_control_service_;
 
     // Scheduler
     std::unique_ptr<WorkflowScheduler> scheduler_;