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:
TaskManager(common/websocket/task_manager.go) — owns the task lifecycle: creation, dispatch, enrichment, cleanup, and termination.AgentManager(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.
// 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:
- Persistence — The task is written to the database via
TaskStore.CreateSessionwith a status ofdoing. - Title generation — A readable title is auto-generated through
generateTaskTitle(referenced attask_manager.goline 96). - 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_idfound inreq.Paramsis fetched fromModelStore, and the resolved model metadata is injected into the task payload asenhancedParams. - Agent configuration — If an
agent_idis supplied, the corresponding YAML configuration file is read and attached to the task asagent_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:
- Write deadlines —
SetWriteDeadline(time.Now().Add(writeWait))prevents a hung agent from blocking the scheduler. - 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:
- Database update — On a
resultUpdateevent,TaskStore.UpdateSessionStatuschanges the task's status todone. - Resource reclamation —
cleanupTaskremoves the in-memory task and releases associated resources. - 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
terminatedand 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →