How Task Scheduling and Execution Works in AI-Infra-Guard

AI-Infra-Guard schedules and executes tasks through a round-robin WebSocket dispatcher driven by two core components — TaskManager in common/websocket/task_manager.go and AgentManager in common/websocket/agent.go — which persist tasks to a database, select active agents via an atomic counter, and stream real-time progress through Server-Sent Events.

AI-Infra-Guard, an open-source infrastructure-scanning platform from Tencent, orchestrates its security scanning workloads through a deliberately lightweight scheduler built on persistent WebSocket connections. Rather than relying on a heavyweight queue system, the project uses a combination of database-backed task state, atomic round-robin agent selection, and SSE-driven front-end updates. This article walks through the complete task lifecycle, referencing exact source files and function names to give you a ground-level view of the architecture.

Task Scheduling Architecture Overview

The scheduling system lives entirely within the common/websocket/ package, with two primary governors:

Tasks flow through a linear pipeline: create → persist → select agent → enrich payload → dispatch over WebSocket → event updates → final state. Each stage is deliberately simple so the system remains deterministic and easy to debug.

Task Creation and Persistence

When a client — either the web UI or an API caller — invokes the AddTask endpoint, the request is wrapped in a TaskCreateRequest struct and passed to TaskManager.AddTask.

// 1️⃣ Create a new task (API handler)
func createTaskHandler(c *gin.Context) {
    var req TaskCreateRequest
    if err := c.ShouldBindJSON(&req); err != nil { ... }

    // taskMgr is a *TaskManager injected into the server
    if err := taskMgr.AddTask(&req, generateTraceID()); err != nil { ... }

    c.JSON(http.StatusOK, gin.H{"message": "task scheduled"})
}

The AddTask flow performs three critical steps before dispatch:

  1. Persistence — The task is written to the database via TaskStore.CreateSession with a status of doing.
  2. Title generation — A readable title is auto-generated through generateTaskTitle (referenced at task_manager.go line 96).
  3. SSE availability check — The system confirms an SSE connection exists (SSEManager.HasConnection) so the front-end can receive live updates before the task is even dispatched.

This ordering matters: if the front-end has disconnected, logs and progress events would be lost, so the manager gates dispatch on SSE readiness.

Agent Selection Using Round-Robin

After persistence, TaskManager.dispatchTask is invoked. This is the heart of the scheduler.

Getting Available Agents

First, the manager calls AgentManager.GetAvailableAgents, which iterates over the map of active WebSocket connections and returns only those whose isActive flag is true (see GetAvailableAgents in agent.go line 12). Inactive or disconnected agents are excluded automatically.

Round-Robin Policy

A global atomic counter — dispatchCounter — implements the round-robin policy:

// 2️⃣ Scheduler core – round-robin dispatch
func (tm *TaskManager) dispatchTask(sessionId, traceID string) error {
    // fetch the in-memory task
    task, _ := tm.GetTask(sessionId)

    // obtain currently active agents
    agents := tm.agentManager.GetAvailableAgents()
    if len(agents) == 0 { return fmt.Errorf("no agents") }

    // round-robin selection
    idx := atomic.AddUint64(&tm.dispatchCounter, 1) - 1
    chosen := agents[idx%uint64(len(agents))]

    // enrich parameters, construct WSMessage, and send
    msg := WSMessage{Type: WSMsgTypeTaskAssign, Content: ...}
    chosen.conn.SetWriteDeadline(time.Now().Add(writeWait))
    return chosen.conn.WriteJSON(msg)
}

The selection formula idx % len(availableAgents) ensures that successive tasks cycle through all active agents evenly. Because the counter is atomic, concurrent task submissions do not race — each dispatch receives a distinct index even when multiple goroutines call dispatchTask simultaneously.

Task Payload Enrichment

Before the assignment message is constructed, TaskManager enriches the request parameters:

  • Model resolution — Any model_id found in req.Params is fetched from ModelStore, and the resolved model metadata is injected into the task payload as enhancedParams.
  • Agent configuration — If an agent_id is supplied, the corresponding YAML configuration file is read and attached to the task as agent_data.

This enrichment happens server-side so that agents never need to query the database themselves — they receive a fully self-contained payload over the socket.

Dispatching Tasks Over WebSocket

Once the WSMessage is fully prepared, it is written to the chosen agent's WebSocket connection using AgentConnection.conn.WriteJSON. Two safeguards are in place:

  1. Write deadlines — SetWriteDeadline(time.Now().Add(writeWait)) prevents a hung agent from blocking the scheduler.
  2. Failure cleanup — If the write fails, the task is removed from both in-memory state and the database via cleanupFailedTask.

On the agent side, the assignment arrives as a task_assign event, which triggers scanning or monitoring work. The agent then begins streaming events back to the server.

Event Handling and State Updates

Agents report progress back through a continuous stream of events: liveStatus, planUpdate, and resultUpdate. Each event is routed through AgentConnection.handleAgentEvent, which deserializes and forwards it to TaskManager.HandleAgentEvent:

// 3️⃣ Agent receives the assignment and starts processing
func (ac *AgentConnection) handleAgentEvent(am *AgentManager, content interface{}, eventType string) {
    // deserialize and validate
    var ev TaskEventMessage
    json.Unmarshal(...)

    // forward to the TaskManager for DB updates / SSE pushes
    am.taskManager.HandleAgentEvent(ev.SessionID, eventType, ev.Event)
}

The full update flow works like this:

  1. Database update — On a resultUpdate event, TaskStore.UpdateSessionStatus changes the task's status to done.
  2. Resource reclamation — cleanupTask removes the in-memory task and releases associated resources.
  3. SSE push — Simultaneously, the event is streamed to the front-end through SSEManager, which maintains endpoint channels for each session.

This event-driven model means the UI never polls — progress bars, live logs, and plan details update in real time.

Terminating a Running Task

Users can stop a task mid-execution via TaskManager.TerminateTask. The flow contracts what happens during normal dispatch:

// 4️⃣ Terminate a running task
func (tm *TaskManager) TerminateTask(sessionId, username, traceID string) error {
    // permission check omitted for brevity
    tm.notifyAgentToTerminate(assignedAgent, sessionId, traceID)
    tm.taskStore.UpdateSessionStatus(sessionId, TaskStatusTerminated)
    tm.sendTerminationEvent(sessionId, traceID)
    go tm.cleanupTask(sessionId)
    return nil
}

The critical difference between termination and natural completion is the state transition: the database status is set to terminated instead of done. Everything else — the SSE event, the in-memory cleanup, the agent notification — follows the same deterministic path.

SSE Notifications: A Real-Time Dashboard

Every lifecycle event, from creation to assignment to termination, is persisted to the database via TaskStore.StoreEvent and simultaneously pushed to the front-end through SSEManager.

This dual persistence creates an audit trail:

  • The database stores a full event history that can be replayed or audited later.
  • The SSE channel gives you real-time progress bars, logs, and plan status updates.

The sse_manager.go file provides the connection registry and event broadcasting logic, tying the agent-to-server WebSocket channel to the server-to-browser SSE channel.

Key Implementation Files

Component Important File Link
Task lifecycle (creation, dispatch, cleanup) common/websocket/task_manager.go TaskManager implementation
Agent WebSocket handling & registry common/websocket/agent.go AgentManager & AgentConnection
Database models for tasks, events, models pkg/database/task.go TaskStore definitions
SSE manager for front-end streaming common/websocket/sse_manager.go SSEManager implementation
API entry point (CLI/Web) cmd/cli/main.go CLI server routes

Summary

  • AI-Infra-Guard schedules tasks through a round-robin dispatcher powered by an atomic counter, with each task assigned to the next active WebSocket agent.
  • The task lifecycle is persistence-first: tasks are saved to the database (via TaskStore.CreateSession) before any agent is contacted.
  • Every task is enriched server-side with model metadata and agent configuration, making payloads self-contained for the chosen agent.
  • Dispatch happens over persistent WebSocket connections with write deadlines that prevent hung agents from blocking the pipeline.
  • SSE pushes every event to the front-end in real time — status changes, plan updates, and results appear without polling.
  • Termination follows a dedicated path that updates the DB to terminated and reuses the same cleanup machinery as natural completion.

Frequently Asked Questions

How does AI-Infra-Guard select an agent for a task?

The scheduler calls AgentManager.GetAvailableAgents, which filters the connection registry to return only agents with isActive set to true. A global atomic counter (dispatchCounter) is then applied to this slice with the formula agents[idx % len(agents)], implementing a round-robin policy that evenly distributes work across all available agents.

What happens if an agent disconnects mid-task?

When an agent's WebSocket connection breaks, its isActive flag is set to false, so future dispatch runs will not select it. For tasks already assigned, the agent is expected to complete its current work and the server-side uses the resultUpdate event to finalize the task state; any task left incomplete triggers the standard cleanup flow during the event receiver path.

Can tasks be cancelled after they've started?

Yes. A user logs on can call TaskManager.TerminateTask(sessionId, username, traceID), which sends a terminate message to the assigned agent, updates the database status to terminated (via TaskStore.UpdateSessionStatus), pushes an SSE termination event to the frontend, and triggers cleanupTask asynchronously.

How does the front-end receive real-time progress?

Every lifecycle event from the task is written to the database by TaskStore.StoreEvent and simultaneously pushed through SSEManager. The front-end maintains an SSE connection, so updates like liveStatus, planUpdate, and resultUpdate appear as streaming JSON events without any page reload or polling.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →