yao/agent/robot/pool/worker.go
Max bc4787f857 Update executor to support V2 execution model and enhance event handling
- Implement V2 execution model in the standard executor, simplifying task execution to a single call without validation loops.
- Introduce support for resuming suspended executions, allowing for human input during task processing.
- Enhance event handling by pushing task completion and failure events to the event bus for better tracking and integration.
- Update tests to reflect changes in execution flow and ensure robust handling of task statuses and results.
2026-02-25 18:40:48 +08:00

138 lines
3.8 KiB
Go

package pool
import (
"fmt"
"sync"
"time"
"github.com/yaoapp/yao/agent/robot/types"
)
// Worker represents a worker goroutine that processes jobs
type Worker struct {
id int
pool *Pool
stopChan chan struct{}
wg *sync.WaitGroup
}
// newWorker creates a new worker
func newWorker(id int, pool *Pool, wg *sync.WaitGroup) *Worker {
return &Worker{
id: id,
pool: pool,
stopChan: make(chan struct{}),
wg: wg,
}
}
// start starts the worker goroutine
func (w *Worker) start() {
w.wg.Add(1)
go w.run()
}
// stop signals the worker to stop
func (w *Worker) stop() {
close(w.stopChan)
}
// run is the main worker loop
func (w *Worker) run() {
defer w.wg.Done()
ticker := time.NewTicker(100 * time.Millisecond) // poll queue every 100ms
defer ticker.Stop()
for {
select {
case <-w.stopChan:
return
case <-ticker.C:
// Try to get a job from the queue
item := w.pool.queue.Dequeue()
if item == nil {
continue // queue empty, wait for next tick
}
// Execute the job
w.execute(item)
}
}
}
// execute processes a single queue item
func (w *Worker) execute(item *QueueItem) {
// Pre-check if robot can run (non-atomic, just for early rejection)
// The actual atomic check happens inside Executor.Execute() via TryAcquireSlot()
if !item.Robot.CanRun() {
// Robot likely at quota, re-enqueue for later
w.requeue(item, "quota pre-check failed")
return
}
// Mark as running (only when actually executing)
w.pool.incrementRunning()
defer w.pool.decrementRunning()
// Get executor based on mode (uses factory if available, otherwise default)
exec := w.pool.GetExecutor(item.ExecutorMode)
// Execute via Executor interface with pre-generated ID and control
// Note: Executor.ExecuteWithControl() does atomic quota check via TryAcquireSlot()
// The control parameter allows executor to check pause state during execution
execution, err := exec.ExecuteWithControl(item.Ctx, item.Robot, item.Trigger, item.Data, item.ExecID, item.Control)
if err != nil {
// Check if it's a quota error (race condition - another worker got the slot)
if err == types.ErrQuotaExceeded {
w.requeue(item, "quota exceeded (race)")
return
}
// Suspended execution: state is persisted, worker slot released gracefully.
// Do NOT call onComplete — the execution stays in robot.Executions and execController
// so that Resume can find it later (§16.1).
if err == types.ErrExecutionSuspended {
if execution != nil {
fmt.Printf("Worker %d: Execution %s suspended for robot %s (waiting for input)\n",
w.id, execution.ID, item.Robot.MemberID)
}
return
}
fmt.Printf("Worker %d: Execution failed for robot %s: %v\n",
w.id, item.Robot.MemberID, err)
// Notify completion callback with appropriate status
if w.pool.onComplete != nil {
status := types.ExecFailed
if err == types.ErrExecutionCancelled {
status = types.ExecCancelled
}
w.pool.onComplete(item.ExecID, item.Robot.MemberID, status)
}
return
}
if execution != nil {
fmt.Printf("Worker %d: Execution %s completed for robot %s (status: %s)\n",
w.id, execution.ID, item.Robot.MemberID, execution.Status)
// Notify completion callback
if w.pool.onComplete != nil {
w.pool.onComplete(execution.ID, item.Robot.MemberID, execution.Status)
}
}
}
// requeue attempts to put the item back in the queue
func (w *Worker) requeue(item *QueueItem, reason string) {
// Queue length is our system load threshold:
// - If queue has space: task waits for robot quota
// - If queue is full: system is overloaded, drop task
if !w.pool.queue.Enqueue(item) {
// Queue full = system overloaded, drop task (protective discard)
fmt.Printf("Worker %d: Task for robot %s dropped (queue full, %s)\n",
w.id, item.Robot.MemberID, reason)
}
}