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 TaskContext instances
  • 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:

  1. Connection establishment – connect() dials the WebSocket endpoint specified in AgentConfig.ServerURL
  2. Registration – register() transmits the agent's AgentInfo to announce capabilities
  3. Goroutine spawning – handleSend() and handleReceive() 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:

  1. The runtime creates a TaskContext wrapping the TaskRequest
  2. It queries the registered handler map for a matching TaskInterface
  3. If found, it spawns a dedicated goroutine to execute TaskInterface.Execute()
  4. The original TaskRequest and a TaskCallbacks struct 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 data
  • ToolUsedCallback(stepID, status, message string, tools []Tool) – Announces tool invocations during multi-step scans
  • NewPlanStepCallback(stepID, description string) – Indicates progression through scanning phases
  • ToolUseLogCallback(...) – 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:

  1. Context cancellation – Signals all active task goroutines to terminate via the cancellable context passed to Execute()
  2. Task abortion – Iterates through active TaskContext instances, updating their status to failed
  3. Channel closure – Drains sendChan and closes the WebSocket connection
  4. Goroutine termination – handleSend() and handleReceive() 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.go manages WebSocket connections, task lifecycles, and message serialization for distributed scanning operations.
  • Task handlers implement the TaskInterface and register via RegisterTaskFunc(), enabling plugin-style extensibility without modifying core runtime code.
  • Asynchronous messaging through sendChan and dedicated goroutines (handleSend, handleReceive) decouples network I/O from CPU-intensive scanning tasks.
  • Callback-driven reporting via TaskCallbacks provides 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:

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 →