towerops/AGENT_PROTOCOL.md

9.2 KiB

Agent WebSocket Protocol

This document describes the binary WebSocket protocol used for communication between remote SNMP polling agents and the Towerops server.

Overview

The agent uses a single persistent WebSocket connection to:

  • Receive SNMP query jobs from the server
  • Submit raw SNMP query results back to the server
  • Send heartbeat updates
  • Report errors

All messages use Protocol Buffers (protobuf) binary encoding for efficiency.

Connection

WebSocket URL

ws://towerops.net/socket/agent?token=<agent_token>
wss://towerops.net/socket/agent?token=<agent_token>  # Production (TLS)

Authentication

  • Token is passed as query parameter during WebSocket handshake
  • Server validates token and establishes session
  • Invalid tokens result in connection rejection

Connection Lifecycle

1. Agent connects with token
2. Server authenticates and assigns to channel "agent:<token_id>"
3. Server immediately sends "jobs" message with initial job list
4. Agent executes jobs and sends "result" messages
5. Server pushes new jobs as devices are added/changed
6. Agent sends periodic "heartbeat" messages
7. Connection remains open indefinitely

Message Types

All messages are binary protobuf-encoded. The WebSocket frame format is:

{
  event: "jobs" | "result" | "heartbeat" | "error",
  payload: <binary protobuf>
}

Server → Agent

jobs - Job List

Sent when agent first connects or when jobs change.

message AgentJobList {
  repeated AgentJob jobs = 1;
}

message AgentJob {
  string job_id = 1;           // e.g., "discover:abc123" or "poll:def456"
  JobType job_type = 2;        // DISCOVER or POLL
  string equipment_id = 3;     // Equipment UUID
  SnmpDevice device = 4;       // SNMP connection details
  repeated SnmpQuery queries = 5;  // What to query
}

message SnmpDevice {
  string ip = 1;           // 10.250.1.26
  string community = 2;    // SNMP community string
  string version = 3;      // "1", "2c", or "3"
  uint32 port = 4;         // 161
}

message SnmpQuery {
  QueryType query_type = 1;    // GET or WALK
  repeated string oids = 2;    // OIDs to query
}

enum JobType {
  DISCOVER = 0;  // Full device discovery
  POLL = 1;      // Poll known sensors/interfaces
}

enum QueryType {
  GET = 0;   // SNMP GET operation
  WALK = 1;  // SNMP WALK operation
}

Example DISCOVER job:

job_id: "discover:123e4567-e89b-12d3-a456-426614174000"
job_type: DISCOVER
equipment_id: "123e4567-e89b-12d3-a456-426614174000"
device: {
  ip: "10.250.1.26"
  community: "public"
  version: "1"
  port: 161
}
queries: [
  { query_type: GET, oids: ["1.3.6.1.2.1.1.1.0", "1.3.6.1.2.1.1.2.0", ...] },
  { query_type: WALK, oids: ["1.3.6.1.2.1.2.2.1"] },
  { query_type: WALK, oids: ["1.3.6.1.4.1.41112"] }
]

Example POLL job:

job_id: "poll:123e4567-e89b-12d3-a456-426614174000"
job_type: POLL
equipment_id: "123e4567-e89b-12d3-a456-426614174000"
device: {
  ip: "10.250.1.26"
  community: "public"
  version: "1"
  port: 161
}
queries: [
  { query_type: GET, oids: [
    "1.3.6.1.4.1.41112.1.3.2.1.3.1",  // Sensor 1
    "1.3.6.1.4.1.41112.1.3.2.1.4.1",  // Sensor 2
    "1.3.6.1.2.1.2.2.1.10.3",         // Interface 3 stats
    ...
  ]}
]

Agent → Server

result - SNMP Query Results

Sent after executing SNMP queries for a job.

message SnmpResult {
  string equipment_id = 1;
  JobType job_type = 2;
  map<string, string> oid_values = 3;  // OID → value mapping
  int64 timestamp = 4;  // Unix timestamp in seconds
}

Example:

equipment_id: "123e4567-e89b-12d3-a456-426614174000"
job_type: DISCOVER
oid_values: {
  "1.3.6.1.2.1.1.1.0": "Linux AF11 2.6.33 ...",
  "1.3.6.1.2.1.1.2.0": "1.3.6.1.4.1.41112",
  "1.3.6.1.2.1.1.5.0": "AF11-Tower1",
  "1.3.6.1.2.1.2.2.1.2.1": "lo",
  "1.3.6.1.2.1.2.2.1.2.3": "eth0",
  "1.3.6.1.4.1.41112.1.3.2.1.3.1": "5725",  // TX frequency
  ...
}
timestamp: 1705363200

heartbeat - Agent Status

Sent periodically (e.g., every 60 seconds) to indicate agent is alive and report metadata.

message AgentHeartbeat {
  string version = 1;        // Agent version (e.g., "1.0.0")
  string hostname = 2;       // Agent hostname
  uint64 uptime_seconds = 3; // Agent uptime
  string ip_address = 4;     // Agent's IP address
}

Example:

version: "1.0.0"
hostname: "datacenter-agent-1"
uptime_seconds: 86400
ip_address: "192.168.1.100"

error - Job Execution Error

Sent when a job fails (e.g., SNMP timeout, device unreachable).

message AgentError {
  string equipment_id = 1;
  string job_id = 2;
  string message = 3;
  int64 timestamp = 4;
}

Example:

equipment_id: "123e4567-e89b-12d3-a456-426614174000"
job_id: "poll:123e4567-e89b-12d3-a456-426614174000"
message: "SNMP timeout after 5 seconds"
timestamp: 1705363200

Agent Implementation Guide

Minimal Rust Agent Architecture

use tokio::net::TcpStream;
use tokio_tungstenite::{connect_async, WebSocketStream};
use futures::{SinkExt, StreamExt};
use prost::Message;

// 1. Connect to WebSocket
let url = format!("ws://towerops.net/socket/agent?token={}", token);
let (ws_stream, _) = connect_async(url).await?;
let (mut write, mut read) = ws_stream.split();

// 2. Receive jobs
while let Some(msg) = read.next().await {
    let msg = msg?;
    if msg.is_binary() {
        let frame: PhoenixMessage = serde_json::from_slice(&msg.into_data())?;

        if frame.event == "jobs" {
            let job_list = AgentJobList::decode(frame.payload)?;
            for job in job_list.jobs {
                tokio::spawn(execute_job(job, write.clone()));
            }
        }
    }
}

// 3. Execute SNMP job
async fn execute_job(job: AgentJob, mut sink: SplitSink) {
    let mut oid_values = HashMap::new();

    for query in job.queries {
        match query.query_type {
            QueryType::Get => {
                for oid in query.oids {
                    let value = snmp_get(&job.device, &oid).await?;
                    oid_values.insert(oid, value);
                }
            }
            QueryType::Walk => {
                for base_oid in query.oids {
                    let results = snmp_walk(&job.device, &base_oid).await?;
                    oid_values.extend(results);
                }
            }
        }
    }

    // 4. Send results back
    let result = SnmpResult {
        equipment_id: job.equipment_id,
        job_type: job.job_type,
        oid_values,
        timestamp: SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs() as i64,
    };

    let payload = result.encode_to_vec();
    let frame = PhoenixMessage {
        event: "result".to_string(),
        payload,
    };

    sink.send(Message::Binary(serde_json::to_vec(&frame)?)).await?;
}

// 5. Raw UDP SNMP implementation
async fn snmp_get(device: &SnmpDevice, oid: &str) -> Result<String> {
    // Construct SNMP GET PDU (BER encoded)
    // Send via UDP to device.ip:device.port
    // Parse response and extract value
    // Return value as string
}

async fn snmp_walk(device: &SnmpDevice, base_oid: &str) -> Result<HashMap<String, String>> {
    // Construct SNMP GETNEXT PDUs (BER encoded)
    // Send via UDP, follow GETNEXT chain
    // Parse responses and extract OID/value pairs
    // Return map of OID → value
}

Phoenix Channel Message Format

Phoenix channels use JSON-wrapped binary for messages:

{
  "event": "jobs",
  "payload": <binary protobuf>,
  "ref": null,
  "topic": "agent:abc123"
}

The agent library should handle this framing automatically, or you can use phoenix_client_rs crate.

Benefits Over REST API

Feature REST API WebSocket Channel
Connection New TCP handshake per request Persistent connection
Overhead HTTP headers (~500 bytes) Phoenix frame (~50 bytes)
Latency Round-trip per poll Server push (0 latency)
Encoding JSON (~3x larger) Protobuf (binary)
Heartbeat Explicit POST requests Implicit (connection alive)
Complexity Multiple endpoints Single channel

Size comparison for a typical job:

  • REST JSON: ~2.5 KB
  • WebSocket Protobuf: ~800 bytes
  • Savings: 68%

Security Notes

  • Use wss:// (WebSocket over TLS) in production
  • Token is validated once at connection time
  • Invalid tokens reject connection immediately
  • No need for per-message authentication

Error Handling

Connection Errors

  • Token invalid: Connection rejected, status 403
  • Network error: Reconnect with exponential backoff
  • Server restart: Reconnect, re-authenticate

Job Errors

  • SNMP timeout: Send error message with details
  • Invalid OID: Skip OID, log warning
  • Device unreachable: Send error message

Testing

Manual Testing with websocat

# Install websocat
brew install websocat

# Connect and listen
websocat "ws://localhost:4000/socket/agent?token=YOUR_TOKEN"

# Server will send jobs message immediately

Unit Testing

Mock the WebSocket connection and verify:

  1. Agent sends correct protobuf on "result" event
  2. Agent parses "jobs" message correctly
  3. SNMP queries execute with correct parameters

Protocol Version

Current version: 1.0

Future versions may add:

  • Job priority queuing
  • Partial result streaming
  • Compression (gzip protobuf)
  • Multi-tenant isolation