How to Implement Observer Endpoints for Retrieval Monitoring in OpenViking

To implement observer endpoints for retrieval monitoring in OpenViking, inject the TrafficMonitor into the ServerInfoFSPlugin to expose real-time statistics via the virtual file /monitor/traffic, or wrap the Snapshot() method in an HTTP handler for REST-based monitoring.

OpenViking (the AGFS file-system server) includes a built-in traffic monitoring system that tracks bytes read, bytes written, and request counts across all request handlers. By leveraging the TrafficMonitor struct and the virtual file system plugin architecture, you can create observer endpoints that expose retrieval statistics without modifying core handler logic. This approach allows both file-based and HTTP-based monitoring interfaces to consume the same underlying metrics from a single source of truth.

Understanding the TrafficMonitor Architecture

Core Monitoring Component

The retrieval statistics aggregation lives in third_party/agfs/agfs-server/pkg/handlers/traffic_monitor.go. This file implements the TrafficMonitor struct, a thread-safe component that maintains counters for total bytes read, total bytes written, request counts, and timestamps.

The monitor exposes three critical methods:

  • AddRead(n int64) – increments the read counter and request count
  • AddWrite(n int64) – increments the write counter
  • Snapshot() – returns a copy of current statistics suitable for serving to observers
type TrafficMonitor struct {
    mu           sync.RWMutex
    totalRead    int64
    totalWrite   int64
    requestCount int64
    startTime    time.Time
}

func (t *TrafficMonitor) AddRead(n int64) {
    t.mu.Lock()
    t.totalRead += n
    t.requestCount++
    t.mu.Unlock()
}

Every request handler calls these methods after I/O operations complete, ensuring real-time accuracy.

Virtual File System Plugin

The ServerInfoFSPlugin in third_party/agfs/agfs-server/pkg/plugins/serverinfofs/serverinfofs.go exposes the monitor through the virtual file system. This plugin holds a reference to the TrafficMonitor via SetTrafficMonitor and serves JSON statistics when clients read the special path /monitor/traffic.

When a read operation targets this path, the plugin fetches the latest snapshot and returns a JSON payload containing fields such as total_read_bytes, total_write_bytes, and request_count.

func (p *ServerInfoFSPlugin) Read(ctx context.Context, path string) ([]byte, error) {
    if path == "/monitor/traffic" {
        snap := p.monitor.Snapshot()
        return json.Marshal(snap) // returns JSON like {"total_read_bytes":123,...}
    }
    // other virtual files …
}

Wiring the Monitor into the Server

The server initialization in third_party/agfs/agfs-server/cmd/server/main.go handles instantiation and dependency injection. Around line 225, the code creates a single monitor instance and passes it to both the plugin layer and the handler stack.

monitor := trafficmonitor.NewTrafficMonitor()           // <-- core monitor

// inject into ServerInfoFS plug‑in
serverInfoPlugin := serverinfofs.NewServerInfoFSPlugin()
serverInfoPlugin.SetTrafficMonitor(monitor)

// later, when building the handler stack
handler := handlers.NewHandlerStack(monitor) // handlers will call monitor.AddRead/Write

This wiring ensures that all request handlers update the same monitor instance that the observer endpoints query.

Implementing Observer Endpoints

OpenViking supports two primary patterns for exposing retrieval statistics to observers: virtual file system integration and direct HTTP endpoints.

Virtual File System Approach

Mount the ServerInfoFSPlugin into the AGFS namespace by enabling it in the server configuration. Any client capable of reading from the virtual file system—including HTTP gateways, gRPC layers, or FUSE clients—can retrieve statistics by reading /monitor/traffic. This approach requires no additional HTTP routes and leverages the existing file system infrastructure.

HTTP REST Endpoint Approach

For classic REST monitoring, add a thin wrapper around the monitor's Snapshot() method in the HTTP handler layer:

func trafficHandler(monitor *trafficmonitor.TrafficMonitor) http.HandlerFunc {
    return func(w http.ResponseWriter, r *http.Request) {
        snap := monitor.Snapshot()
        json.NewEncoder(w).Encode(snap) // {total_read_bytes:…, total_write_bytes:…, …}
    }
}

Register this handler in main.go under a dedicated path:

http.HandleFunc("/v1/monitor/traffic", trafficHandler(monitor))

Both the virtual file system and HTTP endpoints consume the same TrafficMonitor instance, guaranteeing consistent metrics across interfaces.

Creating Custom Observer Plugins

When you need specialized monitoring views—such as per-client quotas or per-directory traffic statistics—implement a custom plugin:

  1. Create a new plugin under pkg/plugins/ (e.g., retrievalfs)
  2. Inject the TrafficMonitor during construction, mirroring the ServerInfoFSPlugin pattern
  3. Implement custom aggregation inside the plugin's Read method (e.g., filter snapshots by path prefix)
  4. Mount the plugin in the server configuration file

The plumbing for monitor injection already exists in main.go, so custom plugins only need to implement the observation logic.

Configuration and Testing

The TrafficMonitor requires no special configuration; it activates automatically when instantiated. However, you can enable or disable the observer endpoint by configuring the serverinfofs plugin:

plugins:
  - name: serverinfofs
    enabled: true   # turn on/off the observer endpoint

Add custom plugins to the same plugins array in the YAML configuration.

OpenViking includes unit tests for validating monitor accuracy and plugin output:

Run go test ./... to ensure your observer implementation correctly aggregates and exposes retrieval statistics.

Summary

  • Instantiate TrafficMonitor in main.go using trafficmonitor.NewTrafficMonitor()
  • Inject the monitor into ServerInfoFSPlugin via SetTrafficMonitor() to enable the virtual file endpoint at /monitor/traffic
  • Query real-time statistics using Snapshot(), which returns current read/write bytes and request counts
  • Extend monitoring by wrapping Snapshot() in HTTP handlers or creating custom plugins that consume the same monitor instance
  • Configure observer availability through the serverinfofs plugin settings in the server configuration file

Frequently Asked Questions

How does the TrafficMonitor ensure thread safety across concurrent requests?

The TrafficMonitor uses a sync.RWMutex to protect all internal counters. Methods like AddRead() and AddWrite() acquire write locks when incrementing statistics, while Snapshot() acquires a read lock when copying current values. This design allows concurrent request handlers to update metrics safely without blocking each other during read-heavy observation queries.

Can I expose retrieval statistics through both the virtual file system and HTTP simultaneously?

Yes. Since the TrafficMonitor is a shared singleton injected into multiple components, you can mount the ServerInfoFSPlugin for virtual file access while also registering an HTTP handler that calls monitor.Snapshot() directly. Both interfaces consume the same underlying data structure, ensuring consistent metrics across all observer endpoints.

What fields are included in the JSON output from the /monitor/traffic virtual file?

The JSON payload contains total_read_bytes, total_write_bytes, request_count, and start_time fields by default. The exact structure matches the fields defined in the TrafficMonitor struct within traffic_monitor.go. You can extend this by modifying the Snapshot() method return type or adding custom aggregation logic in your observer plugin.

How do I test a custom observer plugin during development?

Use the existing test patterns in traffic_monitor_test.go and serverinfofs_test.go as templates. Instantiate a TrafficMonitor, call AddRead() and AddWrite() with known values, then verify that your plugin's Read() method returns the expected JSON structure. Run go test ./pkg/plugins/yourplugin/ to validate that your observer correctly interprets the monitor's snapshot data.

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 →