How the AI-Infra-Guard Agent Runtime Works: WebSocket-Based Distributed Scanning
The AI-Infra-Guard agent runtime is a Go-based plugin architecture that maintains a persistent WebSocket connection to a central server, executes scanning tasks through registered handlers, and streams results back via an asynchronous message pipeline.
The Tencent/AI-Infra-Guard project implements a distributed security scanning system where lightweight agents perform infrastructure assessments. Understanding how the agent runtime executes tasks helps developers extend the platform with custom scanning capabilities or integrate the agent into existing workflows.
Core Architecture Components
The runtime consists of three tightly integrated subsystems defined across the common/agent/ package.
The Agent Struct and Connection State
At the heart of the runtime lies the Agent struct defined in common/agent/agent.go. This core object maintains:
- Identity metadata via
AgentInfo(ID, hostname, IP, version, capabilities) - WebSocket client connection to the central server
- Active task registry storing in-flight
TaskContextinstances - Outbound message channel (
sendChan) for decoupled I/O - Lifecycle context for graceful cancellation
The struct exposes NewAgent() for initialization and Start()/Stop() for lifecycle management, while internally coordinating goroutines for concurrent send and receive operations.
Message Protocol and Type System
All server-agent communication uses JSON-encoded messages defined in common/agent/types.go. The protocol supports:
- Registration messages (
AgentMsgTypeRegister) for agent identification - Task assignment (
ServerMsgTypeTaskAssign) for work distribution - Progress updates (
AgentMsgTypeTaskResult,AgentMsgTypeToolUsed) for real-time feedback - Status notifications for completion or failure states
These definitions include strongly-typed Go structs that map directly to WebSocket payloads, ensuring type safety across the wire format.
Task Interface and Handler Registration
The runtime uses a plugin model where capabilities implement the TaskInterface interface. This contract requires two methods:
GetName() string– Returns the task type identifier (e.g.,TaskTypeTestDemo)Execute(ctx context.Context, req TaskRequest, cb TaskCallbacks) error– Contains the scanning logic
Agents register implementations via RegisterTaskFunc() in agent.go, storing handlers in an internal map keyed by task type. When the server assigns work, the runtime looks up the matching handler in agent_task.go and invokes it.
Runtime Execution Flow
The agent follows a strict lifecycle from connection establishment to task completion.
Initialization and WebSocket Handshake
When Start() is invoked, the runtime executes a three-phase bootstrap:
- Connection establishment –
connect()dials the WebSocket endpoint specified inAgentConfig.ServerURL - Registration –
register()transmits the agent'sAgentInfoto announce capabilities - Goroutine spawning –
handleSend()andhandleReceive()launch as separate goroutines to manage concurrent bidirectional communication
This architecture decouples network I/O from task execution, preventing slow scanning operations from blocking message transmission.
Task Reception and Dispatch
Incoming messages flow through handleReceive(), which unmarshals JSON and routes to processMessage(). Upon receiving ServerMsgTypeTaskAssign:
- The runtime creates a
TaskContextwrapping theTaskRequest - It queries the registered handler map for a matching
TaskInterface - If found, it spawns a dedicated goroutine to execute
TaskInterface.Execute() - The original
TaskRequestand aTaskCallbacksstruct are passed to the handler
The TaskCallbacks struct provides dependency-injected functions that wrap internal methods like SendTaskResult() and SendToolUsed(), automatically piping data into the sendChan for transmission.
Asynchronous Result Reporting
Task handlers communicate with the server through callback methods rather than direct socket writes:
ResultCallback(data map[string]interface{})– Reports final scan results or intermediate dataToolUsedCallback(stepID, status, message string, tools []Tool)– Announces tool invocations during multi-step scansNewPlanStepCallback(stepID, description string)– Indicates progression through scanning phasesToolUseLogCallback(...)– Streams real-time log output from external tools
These callbacks construct protocol-compliant messages and push them onto sendChan. The handleSend() goroutine continuously drains this channel, writes JSON to the WebSocket, and guarantees ordered delivery even under high concurrency.
Implementing Custom Scanning Capabilities
Developers extend the runtime by implementing the TaskInterface and registering instances with the agent.
Creating a Task Implementation
Here is a minimal implementation of a custom scanning task located in common/agent/agent_test.go and executable patterns:
package main
import (
"context"
"log"
"github.com/Tencent/AI-Infra-Guard/common/agent"
)
// DemoTask implements the TaskInterface expected by the Agent runtime.
type DemoTask struct{}
func (d DemoTask) GetName() string {
return agent.TaskTypeTestDemo
}
func (d DemoTask) Execute(ctx context.Context,
req agent.TaskRequest,
cb agent.TaskCallbacks) error {
// Scanning logic here
cb.ResultCallback(map[string]interface{}{"msg": "demo completed"})
return nil
}
func main() {
cfg := agent.AgentConfig{
ServerURL: "ws://127.0.0.1:8088/ws",
Info: agent.AgentInfo{
ID: "demo-agent",
HostName: "demo-host",
IP: "127.0.0.1",
Version: "v0.1",
},
}
a := agent.NewAgent(cfg)
a.RegisterTaskFunc(DemoTask{})
if err := a.Start(); err != nil {
log.Fatalf("agent start failed: %v", err)
}
select {} // Run until interrupted
}
Streaming Progress and Tool Usage
Complex scans requiring multiple steps or external tool invocations use the callback API for granular reporting:
func (d DemoTask) Execute(ctx context.Context,
req agent.TaskRequest,
cb agent.TaskCallbacks) error {
// Report planning phase
cb.NewPlanStepCallback("step-1", "Initialize Demo Environment")
// Report tool execution
cb.ToolUsedCallback("step-1", "status-1", "Running vulnerability scanner",
[]agent.Tool{agent.CreateTool("tool-1", "nmap",
agent.ToolStatusDoing, "port scan", "run", "192.168.1.1", "")})
// Final results
cb.ResultCallback(map[string]interface{}{
"status": "completed",
"findings": []string{"port 80 open", "port 443 open"},
})
return nil
}
These callbacks translate to AgentMsgTypePlanStep, AgentMsgTypeToolUsed, and AgentMsgTypeTaskResult messages on the wire, allowing the central server to render real-time progress indicators.
Graceful Shutdown and Resource Cleanup
The runtime supports clean termination through context cancellation. Calling Stop() triggers:
- Context cancellation – Signals all active task goroutines to terminate via the cancellable context passed to
Execute() - Task abortion – Iterates through active
TaskContextinstances, updating their status to failed - Channel closure – Drains
sendChanand closes the WebSocket connection - Goroutine termination –
handleSend()andhandleReceive()exit their loops upon detecting the closed connection or cancelled context
This ensures no orphaned goroutines remain when the agent process exits, and the server receives proper disconnect notifications.
Summary
- The agent runtime in
common/agent/agent.gomanages WebSocket connections, task lifecycles, and message serialization for distributed scanning operations. - Task handlers implement the
TaskInterfaceand register viaRegisterTaskFunc(), enabling plugin-style extensibility without modifying core runtime code. - Asynchronous messaging through
sendChanand dedicated goroutines (handleSend,handleReceive) decouples network I/O from CPU-intensive scanning tasks. - Callback-driven reporting via
TaskCallbacksprovides type-safe methods for streaming results, tool logs, and execution plans back to the central server. - Graceful shutdown leverages Go's context package to cancel in-flight operations and release resources cleanly.
Frequently Asked Questions
How does the agent handle multiple concurrent tasks?
The runtime spawns a new goroutine for each incoming ServerMsgTypeTaskAssign message in processMessage(). Each goroutine receives its own TaskContext and operates independently, allowing parallel execution of distinct scanning tasks. The sendChan pipeline serializes outbound messages from all concurrent tasks to maintain ordered delivery to the server.
What message types does the WebSocket protocol support?
According to common/agent/types.go, the protocol defines constants for agent-to-server messages (AgentMsgTypeRegister, AgentMsgTypeTaskResult, AgentMsgTypeToolUsed) and server-to-agent commands (ServerMsgTypeTaskAssign, ServerMsgTypeStopTask). These map to Go structs handling JSON serialization of task requests, results, and intermediate progress updates.
Can I implement the agent in languages other than Go?
While the core runtime is implemented in Go in common/agent/, the repository includes agent-scan/agent_scan/core/agent.py, demonstrating Python-based agent interaction. However, the canonical implementation uses Go's common/agent/ package for production deployments due to its efficient goroutine model and compiled performance characteristics.
Where is the server-side WebSocket handler located?
The server-side counterpart that receives agent connections and routes task assignments resides in common/websocket/agent.go. This file implements the WebSocket upgrade handling and message routing logic that pairs with the client-side runtime in common/agent/agent.go.
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 →