Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
38 changes: 38 additions & 0 deletions internal/agent/canvas/runner.go
Original file line number Diff line number Diff line change
Expand Up @@ -163,6 +163,8 @@ type Runner struct {
mu sync.Mutex
interruptIDs map[string]string // key = canvasID + "|" + sessionID; value = eino interrupt id
runCancels map[string]chan struct{}
taskCancels map[string]chan struct{}
taskCanvases map[string]string
}

// NewRunner returns a fresh Runner with the in-memory interrupt-id
Expand All @@ -172,6 +174,8 @@ func NewRunner() *Runner {
return &Runner{
interruptIDs: make(map[string]string),
runCancels: make(map[string]chan struct{}),
taskCancels: make(map[string]chan struct{}),
taskCanvases: make(map[string]string),
}
}

Expand Down Expand Up @@ -259,6 +263,10 @@ func (r *Runner) Run(
if taskID == "" {
taskID = utility.GenerateToken()
}
r.mu.Lock()
r.taskCancels[taskID] = cancel
r.taskCanvases[taskID] = canvasID
r.mu.Unlock()

// Inject the output channel + metadata so the RunFunc can emit
// events during execution (workflow_started, node_started,
Expand All @@ -275,6 +283,10 @@ func (r *Runner) Run(
if r.runCancels[canvasID] == cancel {
delete(r.runCancels, canvasID)
}
if r.taskCancels[taskID] == cancel {
delete(r.taskCancels, taskID)
delete(r.taskCanvases, taskID)
}
r.mu.Unlock()
}()
// Panic sentinel (temporary diagnostic — see plan):
Expand Down Expand Up @@ -372,6 +384,32 @@ func (r *Runner) Cancel(canvasID string) {
}
}

// CancelTask signals an active run identified by the task_id emitted in its
// events. It returns the owning canvas id when the task is still active.
func (r *Runner) CancelTask(taskID string) (string, bool) {
r.mu.Lock()
cancel, ok := r.taskCancels[taskID]
canvasID := r.taskCanvases[taskID]
r.mu.Unlock()
if !ok {
return "", false
}
select {
case <-cancel:
default:
close(cancel)
}
return canvasID, true
}

// TaskCanvas returns the canvas owning an active task without cancelling it.
func (r *Runner) TaskCanvas(taskID string) (string, bool) {
r.mu.Lock()
defer r.mu.Unlock()
canvasID, ok := r.taskCanvases[taskID]
return canvasID, ok
}

// Peek reports whether a paused interrupt id is held for the given
// (canvasID, sessionID). It is intended for tests and diagnostics;
// the real runner does not need it at run time.
Expand Down
16 changes: 16 additions & 0 deletions internal/handler/agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -493,6 +493,22 @@ func (h *AgentHandler) CancelAgent(c *gin.Context) {
common.SuccessWithData(c, true, "success")
}

// CancelTask cancels an active agent run by the task_id emitted in its events.
// Unknown and already-completed task ids are treated as successful no-ops.
func (h *AgentHandler) CancelTask(c *gin.Context) {
user, code, msg := GetUser(c)
if code != common.CodeSuccess {
common.ResponseWithCodeData(c, code, nil, msg)
return
}
if err := h.agentService.CancelTask(c.Request.Context(), user.ID, c.Param("task_id")); err != nil {
ec, em := mapAgentError(err)
common.ResponseWithCodeData(c, ec, nil, em)
return
}
common.SuccessWithData(c, true, "success")
}

// publishAgentRequest is the wire shape for POST /api/v1/agents/:canvas_id/publish.
type publishAgentRequest struct {
Title *string `json:"title,omitempty"`
Expand Down
3 changes: 3 additions & 0 deletions internal/router/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -260,6 +260,9 @@ func (r *Router) Setup(engine *gin.Engine) {
// API v1 route group
v1 := authorized.Group("/api/v1")
{
// Agent stop button compatibility with Python's task API.
v1.POST("/tasks/:task_id/cancel", r.agentHandler.CancelTask)

// Auth routes
auth := v1.Group("/auth")
{
Expand Down
18 changes: 18 additions & 0 deletions internal/service/agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -1881,3 +1881,21 @@ func (s *AgentService) CancelAgent(ctx context.Context, userID, canvasID string)
}
return nil
}

// CancelTask cancels an active agent run by the task_id exposed to clients.
// Unknown or already-finished tasks are intentionally idempotent, matching
// Python's /tasks/<task_id>/cancel endpoint.
func (s *AgentService) CancelTask(ctx context.Context, userID, taskID string) error {
if taskID == "" || s.runner == nil {
return nil
}
canvasID, active := s.runner.TaskCanvas(taskID)
if !active {
return nil
}
if _, err := s.loadCanvasForUser(ctx, userID, canvasID); err != nil {
return err
}
s.runner.CancelTask(taskID)
return nil
}
Comment on lines +1888 to +1901

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔒 Security & Privacy | 🔴 Critical | ⚡ Quick win

Prevent TOCTOU race condition during task cancellation.

There is a Time-Of-Check to Time-Of-Use (TOCTOU) race condition between s.runner.TaskCanvas(taskID) and s.runner.CancelTask(taskID). If an attacker quickly reassigns a taskID (e.g., via a colliding version_id) to a victim's canvas in the short window after the authorization check but before the actual cancellation, they could bypass authorization and cancel the victim's task. Fix this by passing the authorized canvas ID to the runner and atomically validating it inside the runner's lock.

  • internal/service/agent.go#L1888-L1901: pass the authorized canvasID to s.runner.CancelTask to ensure we only cancel the verified canvas.
  • internal/agent/canvas/runner.go#L387-L403: update CancelTask to accept and enforce expectedCanvasID while holding the lock.
🔒️ Proposed fix to enforce canvas ownership during cancellation

In internal/service/agent.go:

 	if _, err := s.loadCanvasForUser(ctx, userID, canvasID); err != nil {
 		return err
 	}
-	s.runner.CancelTask(taskID)
+	s.runner.CancelTask(taskID, canvasID)
 	return nil

In internal/agent/canvas/runner.go:

-// CancelTask signals an active run identified by the task_id emitted in its
-// events. It returns the owning canvas id when the task is still active.
-func (r *Runner) CancelTask(taskID string) (string, bool) {
+// CancelTask signals an active run by task_id, ensuring it belongs to the expected canvas.
+func (r *Runner) CancelTask(taskID, expectedCanvasID string) bool {
 	r.mu.Lock()
 	cancel, ok := r.taskCancels[taskID]
 	canvasID := r.taskCanvases[taskID]
 	r.mu.Unlock()
-	if !ok {
-		return "", false
+	if !ok || canvasID != expectedCanvasID {
+		return false
 	}
 	select {
 	case <-cancel:
 	default:
 		close(cancel)
 	}
-	return canvasID, true
+	return true
 }
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
func (s *AgentService) CancelTask(ctx context.Context, userID, taskID string) error {
if taskID == "" || s.runner == nil {
return nil
}
canvasID, active := s.runner.TaskCanvas(taskID)
if !active {
return nil
}
if _, err := s.loadCanvasForUser(ctx, userID, canvasID); err != nil {
return err
}
s.runner.CancelTask(taskID)
return nil
}
func (s *AgentService) CancelTask(ctx context.Context, userID, taskID string) error {
if taskID == "" || s.runner == nil {
return nil
}
canvasID, active := s.runner.TaskCanvas(taskID)
if !active {
return nil
}
if _, err := s.loadCanvasForUser(ctx, userID, canvasID); err != nil {
return err
}
s.runner.CancelTask(taskID, canvasID)
return nil
}
Suggested change
func (s *AgentService) CancelTask(ctx context.Context, userID, taskID string) error {
if taskID == "" || s.runner == nil {
return nil
}
canvasID, active := s.runner.TaskCanvas(taskID)
if !active {
return nil
}
if _, err := s.loadCanvasForUser(ctx, userID, canvasID); err != nil {
return err
}
s.runner.CancelTask(taskID)
return nil
}
// CancelTask signals an active run by task_id, ensuring it belongs to the expected canvas.
func (r *Runner) CancelTask(taskID, expectedCanvasID string) bool {
r.mu.Lock()
cancel, ok := r.taskCancels[taskID]
canvasID := r.taskCanvases[taskID]
r.mu.Unlock()
if !ok || canvasID != expectedCanvasID {
return false
}
select {
case <-cancel:
default:
close(cancel)
}
return true
}
📍 Affects 2 files
  • internal/service/agent.go#L1888-L1901 (this comment)
  • internal/agent/canvas/runner.go#L387-L403
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@internal/service/agent.go` around lines 1888 - 1901, Prevent the cancellation
TOCTOU race by passing the authorized canvasID from AgentService.CancelTask to
runner.CancelTask in internal/service/agent.go lines 1888-1901. Update
runner.CancelTask in internal/agent/canvas/runner.go lines 387-403 to accept
expectedCanvasID and, while holding the runner lock, verify the task still
belongs to that canvas before cancelling; otherwise leave it unchanged.