# How to Implement Observer Endpoints for Retrieval Monitoring in OpenViking

> Learn to implement observer endpoints for retrieval monitoring in OpenViking by injecting TrafficMonitor or wrapping Snapshot in an HTTP handler. Get real-time statistics now.

- Repository: [Volcengine/OpenViking](https://github.com/volcengine/OpenViking)
- Tags: how-to-guide
- Published: 2026-03-08

---

**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`](https://github.com/volcengine/OpenViking/blob/main/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

```go
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`](https://github.com/volcengine/OpenViking/blob/main/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`.

```go
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`](https://github.com/volcengine/OpenViking/blob/main/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.

```go
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:

```go
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`](https://github.com/volcengine/OpenViking/blob/main/main.go) under a dedicated path:

```go
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`](https://github.com/volcengine/OpenViking/blob/main/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:

```yaml
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:

- [`third_party/agfs/agfs-server/pkg/handlers/traffic_monitor_test.go`](https://github.com/volcengine/OpenViking/blob/main/third_party/agfs/agfs-server/pkg/handlers/traffic_monitor_test.go) – verifies counter increments
- [`third_party/agfs/agfs-server/pkg/plugins/serverinfofs/serverinfofs_test.go`](https://github.com/volcengine/OpenViking/blob/main/third_party/agfs/agfs-server/pkg/plugins/serverinfofs/serverinfofs_test.go) – validates JSON payload structure

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

## Summary

- **Instantiate** `TrafficMonitor` in [`main.go`](https://github.com/volcengine/OpenViking/blob/main/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`](https://github.com/volcengine/OpenViking/blob/main/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`](https://github.com/volcengine/OpenViking/blob/main/traffic_monitor_test.go) and [`serverinfofs_test.go`](https://github.com/volcengine/OpenViking/blob/main/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.