Преглед на файлове

Merge branch 'executions-live-and-cancel'

fszontagh преди 1 месец
родител
ревизия
46097f1d60

+ 10 - 2
nodes/core/if-condition.js

@@ -221,6 +221,14 @@ async function execute(config, input, context) {
   // This allows access to loop variables (currentItem, currentIndex) at the top level
   // as well as nested data from previous nodes
   const data = input;
+  // Conditions are evaluated against the WHOLE input, so loop variables such
+  // as the item variable stay reachable at the top level. What gets forwarded
+  // on the branch is the payload only - input.data - because forwarding the
+  // wrapper too made every downstream path a level deeper than the one the
+  // author wrote, and a loop reading "data.emails" silently found nothing.
+  const payload = (input && typeof input === 'object' && !Array.isArray(input) && 'data' in input)
+    ? input.data
+    : input;
   const conditions = config.conditions || [];
   const combineWith = config.combineWith || 'and';
 
@@ -234,7 +242,7 @@ async function execute(config, input, context) {
       result: true,
       matchedConditions: [],
       _activeBranch: 'true',
-      true: data  // Pass data only on the active branch
+      true: payload  // Pass the payload only, on the active branch
     };
   }
 
@@ -272,7 +280,7 @@ async function execute(config, input, context) {
     result: finalResult,
     matchedConditions: matchedConditions,
     _activeBranch: activeBranch,
-    [activeBranch]: data  // Data is only on the active branch
+    [activeBranch]: payload  // The payload is only on the active branch
   };
 }
 

+ 10 - 2
nodes/core/switch.js

@@ -112,6 +112,14 @@ function ruleMatches(fieldValue, rule) {
 async function execute(config, input, context) {
     const rules = Array.isArray(config.rules) ? config.rules : [];
     const data = input;
+  // Conditions are evaluated against the WHOLE input, so loop variables such
+  // as the item variable stay reachable at the top level. What gets forwarded
+  // on the branch is the payload only - input.data - because forwarding the
+  // wrapper too made every downstream path a level deeper than the one the
+  // author wrote, and a loop reading "data.emails" silently found nothing.
+  const payload = (input && typeof input === 'object' && !Array.isArray(input) && 'data' in input)
+    ? input.data
+    : input;
     const fieldValue = smartbotic.utils.getFieldValue(data, config.field);
 
     for (let i = 0; i < rules.length; i++) {
@@ -122,7 +130,7 @@ async function execute(config, input, context) {
                 matchedRule: i,
                 matchedLabel: rules[i].label || branch,
                 _activeBranch: branch,
-                [branch]: data
+                [branch]: payload
             };
         }
     }
@@ -132,7 +140,7 @@ async function execute(config, input, context) {
         matchedRule: -1,
         matchedLabel: '',
         _activeBranch: 'fallback',
-        fallback: data
+        fallback: payload
     };
 }
 

+ 30 - 0
src/runner/workflow_engine.cpp

@@ -337,6 +337,19 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
     }
     result.workflow_snapshot = snapshot;
 
+    // Write the record before the walk starts, not only when it ends.
+    //
+    // Until this, an execution existed in the database only once it had
+    // finished, which made a running execution invisible to everything outside
+    // this process: the executions list could not show it, its elapsed time
+    // could not be read, and cancelling it returned 404 because there was no
+    // record to look up. A run that never finishes - a hung HTTP call, a killed
+    // runner - left no trace at all.
+    //
+    // storeExecution already distinguishes insert from update, so the write at
+    // the end of the walk updates this row rather than colliding with it.
+    storeExecution(result);
+
     ++active_count_;
 
     {
@@ -1874,10 +1887,27 @@ bool WorkflowEngine::executeLoopBody(
     // loop and the per-item loop can unwind.
     bool stopped_in_body = false;
 
+
     // Execute body for each item
     for (size_t i = 0; i < ctx.items.size(); ++i) {
         ctx.current_index = i;
 
+        // The main walk checks for cancellation between nodes; without the same
+        // check here a loop was uncancellable, and a loop is exactly where a run
+        // spends its time - ten thousand items, or a body that waits on a slow
+        // HTTP call, would ignore Cancel until the whole loop finished.
+        {
+            std::lock_guard<std::mutex> lock(mutex_);
+            if (cancelled_executions_.contains(result.execution_id)) {
+                result.status = ExecutionStatus::Cancelled;
+                result.error = "Execution cancelled";
+                // Breaking here leaves the outer walk to notice: its own
+                // between-nodes check sees the same cancellation set and ends
+                // the run, so the nodes after the loop do not execute.
+                break;
+            }
+        }
+
         if (callback) {
             callback("loop.iteration.start", {
                 {"executionId", result.execution_id},

+ 66 - 7
src/webserver/api/execution_controller.cpp

@@ -121,17 +121,76 @@ void ExecutionController::cancelExecution(const httplib::Request& req, httplib::
                                            const auth::AuthContext& ctx) {
     std::string id = req.matches[1];
 
-    // TODO: Send cancel request to runner via gRPC
-    // For now, just update status in database
-
-    auto result = storage_.update("executions", id,
-                                  {{"status", "cancelled"}}, 0, true);
-    if (result.failed()) {
+    auto existing = storage_.get("executions", id);
+    if (existing.failed()) {
         sendError(res, "Execution not found", 404);
         return;
     }
 
-    sendJson(res, {{"success", true}, {"status", "cancelled"}});
+    const auto& record = existing.value();
+    const std::string status = record.value("status", "");
+
+    // A finished execution cannot be cancelled, and saying so beats rewriting a
+    // completed run's status to "cancelled" and losing what actually happened.
+    if (status == "completed" || status == "failed" || status == "cancelled") {
+        sendError(res, "Execution already finished as " + status, 409);
+        return;
+    }
+
+    // The cancel has to reach the runner that owns this execution. Its engine
+    // holds the cancellation set that the walk consults between nodes; any
+    // other runner has never heard of this execution id. Picking one through
+    // the load balancer, as the resume path does, would report success while
+    // the real execution carried on.
+    const std::string runner_id = record.value("runnerId", "");
+    auto runner = runner_id.empty() ? std::nullopt : load_balancer_.getRunner(runner_id);
+
+    if (!runner) {
+        // No runner to ask: it was never dispatched, or the runner is gone. The
+        // record is the only thing left to correct, and leaving it Running for
+        // ever is worse than marking it cancelled.
+        auto updated = storage_.update("executions", id, {{"status", "cancelled"}}, 0, true);
+        if (updated.failed()) {
+            sendError(res, "Could not cancel: " + updated.error().message(), 500);
+            return;
+        }
+        LOG_WARN("Cancelled execution {} in the database only - runner '{}' is not registered",
+                 id, runner_id);
+        sendJson(res, {{"success", true}, {"status", "cancelled"},
+                       {"note", "The owning runner is not registered, so the record was marked "
+                                "cancelled without reaching a runner"}});
+        return;
+    }
+
+    // Qualified with the leading "::" for the same reason as the resume path.
+    auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
+    auto stub = proto::RunnerService::NewStub(channel);
+
+    proto::CancelExecutionRequest grpc_req;
+    grpc_req.set_execution_id(id);
+
+    proto::Empty grpc_res;
+    ::grpc::ClientContext grpc_ctx;
+    grpc_ctx.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(15));
+
+    auto grpc_status = stub->CancelExecution(&grpc_ctx, grpc_req, &grpc_res);
+    if (!grpc_status.ok()) {
+        sendError(res, "Could not reach the runner to cancel: " + grpc_status.error_message(), 502);
+        return;
+    }
+
+    // Deliberately not writing the status here. Cancellation is honoured at the
+    // next node boundary, and the runner writes the finished record itself - a
+    // status written now would be overwritten moments later, and would claim
+    // the run had stopped while a node was still running.
+    ws_server_.broadcast("executions." + id + ".cancelling", {
+        {"executionId", id},
+        {"requestedBy", ctx.user_id}
+    });
+
+    LOG_INFO("Cancellation requested for execution {} on runner {}", id, runner_id);
+    sendJson(res, {{"success", true}, {"status", "cancelling"},
+                   {"note", "The runner stops at the next node boundary"}});
 }
 
 void ExecutionController::retryExecution(const httplib::Request& req, httplib::Response& res,

+ 10 - 0
src/webserver/runners/load_balancer.hpp

@@ -33,6 +33,16 @@ public:
     // Returns runner address or nullopt if no runner available
     std::optional<Runner> selectRunner(const std::string& required_node_type = "");
 
+    // Look up one specific runner rather than choosing one.
+    //
+    // Some requests are about an execution that already belongs to a runner -
+    // cancelling, for instance - and only that runner knows about it. Those
+    // must not go through selectRunner, which would happily hand back a
+    // different runner that has never heard of the execution.
+    std::optional<Runner> getRunner(const std::string& runner_id) {
+        return registry_.getRunner(runner_id);
+    }
+
     // Get current strategy
     LoadBalancingStrategy getStrategy() const { return config_.strategy; }
 

+ 24 - 0
tests/nodes/if-branch-payload-depth.json

@@ -0,0 +1,24 @@
+{
+  "name": "verify-if-branch-payload-depth",
+  "nodes": [
+    {"id": "n1", "name": "Trigger", "type": "click-trigger", "position": {"x": 0, "y": 0}, "config": {}},
+    {"id": "src", "name": "Src", "type": "code", "position": {"x": 0, "y": 100},
+     "config": {"code": "return { count: 2, emails: [{ id: 'a' }, { id: 'b' }] };"}},
+    {"id": "gate", "name": "Gate", "type": "if-condition", "position": {"x": 0, "y": 200},
+     "config": {"combineWith": "and", "conditions": [{"field": "data.result.count", "operator": "greater_than", "value": "0"}]}},
+    {"id": "lp", "name": "Loop", "type": "loop", "position": {"x": 0, "y": 300},
+     "config": {"inputField": "data.result.emails", "itemVariableName": "email", "outputField": "res", "continueOnError": true}},
+    {"id": "body", "name": "Body", "type": "code", "position": {"x": 0, "y": 400},
+     "config": {"code": "return { seen: (input.email && input.email.id) || 'NOT-AN-EMAIL' };"}}
+  ],
+  "connections": [
+    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "src", "targetInput": "data"},
+    {"sourceNodeId": "src", "sourceOutput": "main", "targetNodeId": "gate", "targetInput": "data"},
+    {"sourceNodeId": "gate", "sourceOutput": "true", "targetNodeId": "lp", "targetInput": "data"},
+    {"sourceNodeId": "lp", "sourceOutput": "loop", "targetNodeId": "body", "targetInput": "data"}
+  ],
+  "expect": {
+    "lp": {"status": "completed"},
+    "body": {"status": "completed", "output": {"result": {"seen": "b"}}}
+  }
+}

+ 71 - 8
webui/src/pages/ExecutionsPage.tsx

@@ -1,10 +1,10 @@
-import { useState } from 'react'
+import { useState, useEffect } from 'react'
 import { useNavigate } from 'react-router-dom'
-import { useQuery } from '@tanstack/react-query'
+import { useQuery, useMutation, useQueryClient } from '@tanstack/react-query'
 import { executionsApi } from '../api/workflows'
 import { useExecutionListUpdates } from '../hooks/useExecutionListUpdates'
 import { formatDistanceToNow } from 'date-fns'
-import { CheckCircle, XCircle, Clock, Ban, RefreshCw, ChevronLeft, ChevronRight, AlertCircle, Octagon } from 'lucide-react'
+import { CheckCircle, XCircle, Clock, Ban, RefreshCw, ChevronLeft, ChevronRight, AlertCircle, Octagon, Square, Hourglass } from 'lucide-react'
 import clsx from 'clsx'
 
 // Helper to safely format timestamps that might be invalid or too large
@@ -19,10 +19,30 @@ function safeFormatDistanceToNow(timestamp: number): string {
   }
 }
 
+// An execution that has not finished has finishedAt 0, so its duration has to
+// be measured against now. "Running..." told you nothing about whether a run had
+// been going for two seconds or two hours.
+function formatElapsed(startedAt: number, finishedAt: number, now: number): string {
+  if (!startedAt) return '-'
+  const end = finishedAt && finishedAt > 0 ? finishedAt : now
+  const seconds = (end - startedAt) / 1000
+  if (seconds < 0) return '-'
+  if (seconds < 60) return `${seconds.toFixed(2)}s`
+  const minutes = Math.floor(seconds / 60)
+  const rest = Math.floor(seconds % 60)
+  if (minutes < 60) return `${minutes}m ${rest}s`
+  const hours = Math.floor(minutes / 60)
+  return `${hours}h ${minutes % 60}m`
+}
+
+// Statuses a run can still leave on its own, so a Stop is meaningful.
+const IN_FLIGHT = ['running', 'pending', 'waiting', 'resuming']
+
 const PAGE_SIZE = 20
 
 export default function ExecutionsPage() {
   const navigate = useNavigate()
+  const queryClient = useQueryClient()
   const [page, setPage] = useState(1)
 
   // The workflow editor already renders an execution across the canvas with each
@@ -43,6 +63,24 @@ export default function ExecutionsPage() {
   // Subscribe to WebSocket for real-time updates instead of polling
   useExecutionListUpdates()
 
+  // A ticking clock only while something is actually in flight. Without the
+  // guard this would re-render the page once a second for ever on a list of
+  // finished runs.
+  const [now, setNow] = useState(() => Date.now())
+  const hasInFlight = (data?.executions || []).some((e: any) => IN_FLIGHT.includes(e.status))
+  useEffect(() => {
+    if (!hasInFlight) return
+    const timer = setInterval(() => setNow(Date.now()), 1000)
+    return () => clearInterval(timer)
+  }, [hasInFlight])
+
+  const cancelMutation = useMutation({
+    mutationFn: (id: string) => executionsApi.cancel(id),
+    onSuccess: () => {
+      queryClient.invalidateQueries({ queryKey: ['executions'] })
+    },
+  })
+
   const executions = data?.executions || []
   const total = data?.total || 0
   const totalPages = Math.ceil(total / PAGE_SIZE)
@@ -58,6 +96,11 @@ export default function ExecutionsPage() {
         return <RefreshCw className="w-5 h-5 text-blue-500 animate-spin" />
       case 'cancelled':
         return <Ban className="w-5 h-5 text-gray-500" />
+      case 'waiting':
+        // Distinct from running on purpose: a waiting run is not working, it
+        // needs an answer, and stopping it is a different act from stopping
+        // something mid-flight.
+        return <Hourglass className="w-5 h-5 text-amber-500" />
       default:
         return <Clock className="w-5 h-5 text-gray-400" />
     }
@@ -73,6 +116,8 @@ export default function ExecutionsPage() {
         return 'bg-blue-100 dark:bg-blue-900/30 text-blue-700 dark:text-blue-400'
       case 'cancelled':
         return 'bg-gray-100 dark:bg-slate-700 text-gray-700 dark:text-gray-400'
+      case 'waiting':
+        return 'bg-amber-100 dark:bg-amber-900/30 text-amber-700 dark:text-amber-400'
       default:
         return 'bg-gray-100 dark:bg-slate-700 text-gray-600 dark:text-gray-400'
     }
@@ -123,6 +168,9 @@ export default function ExecutionsPage() {
                   <th className="px-4 py-3 text-left text-sm font-medium text-gray-600 dark:text-gray-300">
                     Duration
                   </th>
+                  <th className="px-4 py-3 text-right text-sm font-medium text-gray-600 dark:text-gray-300">
+                    Actions
+                  </th>
                 </tr>
               </thead>
               <tbody className="divide-y divide-gray-200 dark:divide-slate-700">
@@ -186,11 +234,26 @@ export default function ExecutionsPage() {
                         : '-'}
                     </td>
                     <td className="px-4 py-3 align-top text-gray-500 dark:text-gray-400 text-sm">
-                      {execution.finishedAt && execution.startedAt
-                        ? `${((execution.finishedAt - execution.startedAt) / 1000).toFixed(2)}s`
-                        : execution.status === 'running'
-                        ? 'Running...'
-                        : '-'}
+                      <span className={clsx(IN_FLIGHT.includes(execution.status) && 'tabular-nums')}>
+                        {formatElapsed(execution.startedAt, execution.finishedAt, now)}
+                      </span>
+                    </td>
+                    <td className="px-4 py-3 align-top text-right">
+                      {IN_FLIGHT.includes(execution.status) && (
+                        <button
+                          onClick={(e) => {
+                            // The row opens the execution; the button must not.
+                            e.stopPropagation()
+                            cancelMutation.mutate(execution.id)
+                          }}
+                          disabled={cancelMutation.isPending}
+                          title="Stop this execution. The runner stops at the next node boundary"
+                          className="inline-flex items-center gap-1.5 px-2 py-1 text-xs text-red-600 dark:text-red-400 hover:bg-red-50 dark:hover:bg-red-900/20 rounded-lg disabled:opacity-50"
+                        >
+                          <Square className="w-3 h-3" />
+                          Stop
+                        </button>
+                      )}
                     </td>
                   </tr>
                 ))}