5.6 KiB
Streaming Rdwr: A 9P Pattern for Filtered Subscriptions
This document describes a novel 9P pattern used by Ollie for filtered event subscriptions. The pattern extends the standard rdwr (read-write) file semantics to support streaming reads after an initial write.
The Problem
Traditional pub/sub systems require:
- A subscription API (connect, subscribe to topic, receive messages)
- Multiple protocol operations
- Separate endpoints for different subscription types
9P provides a simple file interface, but standard file semantics don't naturally support "subscribe to X, then stream matching events."
The Pattern: Streaming Rdwr
A single file supports both unfiltered and filtered streaming through the write-then-read pattern:
┌─────────────────────────────────────────────────────────────┐
│ event file │
├─────────────────────────────────────────────────────────────┤
│ Read-only open: Stream all events │
│ Read-write open: Write filter → Stream matching events │
└─────────────────────────────────────────────────────────────┘
Protocol Flow
Unfiltered (read-only):
Topen fid, OREAD
Tread fid, ... → blocks until event, returns "topic payload\n"
Tread fid, ... → blocks until event, returns "topic payload\n"
...
Tclunk fid
Filtered (read-write):
Topen fid, ORDWR
Twrite fid, "session.*.agent.*.state" → stores filter in per-fid state
Tread fid, ... → blocks until matching event, returns "topic payload\n"
Tread fid, ... → blocks until matching event, returns "topic payload\n"
...
Tclunk fid → cancels subscription, frees state
Key Properties
- Single file, dual mode — Same path serves both use cases
- Per-fid state — Each open file descriptor maintains its own filter and subscription
- Lazy subscription — Subscription is created on first read after write, not on write
- Clean teardown — Clunk cancels the subscription automatically
- Standard 9P semantics — No protocol extensions required
Filter Syntax
Filters use a simple glob-like syntax:
| Pattern | Matches |
|---|---|
* |
Exactly one segment |
> |
One or more segments (must be last) |
literal |
Exact match |
Examples:
* # All events
session.*.agent.*.state # State changes for any agent in any session
session.abc123.> # All events for session abc123
session.*.bypass.> # All bypass events
Implementation
Server State
Per-fid structure includes subscription state:
type fid struct {
// ... standard fid fields ...
// Event stream subscription
eventFilter string
eventCh <-chan Event
eventCancel context.CancelFunc
}
Write Handler
if path == "event" {
filter := strings.TrimSpace(string(data))
f.eventFilter = filter
// Clear existing subscription (recreated on next read)
if f.eventCancel != nil {
f.eventCancel()
}
f.eventCh = nil
return success
}
Read Handler
if path == "event" && f.eventFilter != "" {
if f.eventCh == nil {
// Create filtered subscription on first read
ctx, cancel := context.WithCancel(ctx)
f.eventCh = SubscribeEventsFiltered(ctx, f.eventFilter)
f.eventCancel = cancel
}
// Block until matching event
ev := <-f.eventCh
return formatEvent(ev)
}
// Fall through to default stream handling
Clunk Handler
if f.eventCancel != nil {
f.eventCancel()
}
Client Usage
Shell (ollie-9p)
# All events
ollie-9p read event
# Filtered events (rdwrs = streaming rdwr)
echo "session.*.agent.*.state" | ollie-9p rdwrs event
Go
fid, _ := fsys.Open("event", plan9.ORDWR)
fid.Write([]byte("session.*.agent.*.state"))
fid.Seek(0, 0)
for {
n, _ := fid.Read(buf)
fmt.Print(string(buf[:n]))
}
C++/Qt (GUI)
The GUI uses the unfiltered stream for the global event bus, handling filtering client-side for flexibility.
Comparison with Alternatives
| Approach | Pros | Cons |
|---|---|---|
| Separate files per filter | Simple read semantics | Explosion of files, static filters |
| Query parameter in path | RESTful | Non-standard 9P, complex paths |
| Custom protocol extension | Full control | Breaks 9P compatibility |
| Streaming rdwr | Standard 9P, dynamic filters, single file | Requires per-fid state |
Prior Art
- Plan 9 /dev/cons — Multiplexed per-process I/O through single file
- 9P2000.u extensions — Per-fid state for Unix compatibility
- NATS subscriptions — Topic filtering with wildcards
The streaming rdwr pattern combines these ideas: per-fid state from 9P2000.u, wildcard filtering from NATS, and the "write configuration, read results" pattern common in Plan 9 control files.
Limitations
- No replay — Missed events are lost; no persistent queue
- Single filter per fid — Rewriting filter resets subscription
- Server-side filtering — CPU cost scales with event rate × subscribers
For high-throughput scenarios, consider client-side filtering on the unfiltered stream.