doc: document streaming rdwr pattern for filtered subscriptions
This commit is contained in:
parent
772d77940c
commit
59f162a57e
|
|
@ -0,0 +1,182 @@
|
|||
# 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:
|
||||
1. A subscription API (connect, subscribe to topic, receive messages)
|
||||
2. Multiple protocol operations
|
||||
3. 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
|
||||
|
||||
1. **Single file, dual mode** — Same path serves both use cases
|
||||
2. **Per-fid state** — Each open file descriptor maintains its own filter and subscription
|
||||
3. **Lazy subscription** — Subscription is created on first read after write, not on write
|
||||
4. **Clean teardown** — Clunk cancels the subscription automatically
|
||||
5. **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:
|
||||
|
||||
```go
|
||||
type fid struct {
|
||||
// ... standard fid fields ...
|
||||
|
||||
// Event stream subscription
|
||||
eventFilter string
|
||||
eventCh <-chan Event
|
||||
eventCancel context.CancelFunc
|
||||
}
|
||||
```
|
||||
|
||||
### Write Handler
|
||||
|
||||
```go
|
||||
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
|
||||
|
||||
```go
|
||||
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
|
||||
|
||||
```go
|
||||
if f.eventCancel != nil {
|
||||
f.eventCancel()
|
||||
}
|
||||
```
|
||||
|
||||
## Client Usage
|
||||
|
||||
### Shell (ollie-9p)
|
||||
|
||||
```bash
|
||||
# All events
|
||||
ollie-9p read event
|
||||
|
||||
# Filtered events (rdwrs = streaming rdwr)
|
||||
echo "session.*.agent.*.state" | ollie-9p rdwrs event
|
||||
```
|
||||
|
||||
### Go
|
||||
|
||||
```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
|
||||
|
||||
1. **No replay** — Missed events are lost; no persistent queue
|
||||
2. **Single filter per fid** — Rewriting filter resets subscription
|
||||
3. **Server-side filtering** — CPU cost scales with event rate × subscribers
|
||||
|
||||
For high-throughput scenarios, consider client-side filtering on the unfiltered stream.
|
||||
Loading…
Reference in New Issue