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

> Discover how AI-Infra-Guard schedules and executes tasks. Learn about the TaskManager and AgentManager, database persistence, and real-time progress streaming.

- Repository: [Tencent/AI-Infra-Guard](https://github.com/tencent/AI-Infra-Guard)
- Tags: internals
- Published: 2026-08-22

---

**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`](https://github.com/Tencent/AI-Infra-Guard/blob/main/common/websocket/task_manager.go) and `AgentManager` in [`common/websocket/agent.go`](https://github.com/Tencent/AI-Infra-Guard/blob/main/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:

- **`TaskManager`** ([`common/websocket/task_manager.go`](https://github.com/Tencent/AI-Infra-Guard/blob/main/common/websocket/task_manager.go)) — owns the task lifecycle: creation, dispatch, enrichment, cleanup, and termination.
- **`AgentManager`** ([`common/websocket/agent.go`](https://github.com/Tencent/AI-Infra-Guard/blob/main/common/websocket/agent.go)) — maintains the registry of active agent WebSocket connections and provides selection logic.

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`.

```go
// 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`](https://github.com/Tencent/AI-Infra-Guard/blob/main/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`](https://github.com/Tencent/AI-Infra-Guard/blob/main/agent.go) line 12). Inactive or disconnected agents are excluded automatically.

### Round-Robin Policy

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

```go
// 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`:

```go
// 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:

```go
// 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`](https://github.com/Tencent/AI-Infra-Guard/blob/main/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`](https://github.com/Tencent/AI-Infra-Guard/blob/main/common/websocket/task_manager.go) | [TaskManager implementation](https://github.com/Tencent/AI-Infra-Guard/blob/main/common/websocket/task_manager.go) |
| Agent WebSocket handling & registry | [`common/websocket/agent.go`](https://github.com/Tencent/AI-Infra-Guard/blob/main/common/websocket/agent.go) | [AgentManager & AgentConnection](https://github.com/Tencent/AI-Infra-Guard/blob/main/common/websocket/agent.go) |
| Database models for tasks, events, models | [`pkg/database/task.go`](https://github.com/Tencent/AI-Infra-Guard/blob/main/pkg/database/task.go) | [TaskStore definitions](https://github.com/Tencent/AI-Infra-Guard/blob/main/pkg/database/task.go) |
| SSE manager for front-end streaming | [`common/websocket/sse_manager.go`](https://github.com/Tencent/AI-Infra-Guard/blob/main/common/websocket/sse_manager.go) | [SSEManager implementation](https://github.com/Tencent/AI-Infra-Guard/blob/main/common/websocket/sse_manager.go) |
| API entry point (CLI/Web) | [`cmd/cli/main.go`](https://github.com/Tencent/AI-Infra-Guard/blob/main/cmd/cli/main.go) | [CLI server routes](https://github.com/Tencent/AI-Infra-Guard/blob/main/cmd/cli/main.go) |

## 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.