ollie-9p: add rdwrs command for streaming rdwr
rdwrs writes stdin then streams reads (does not wait for EOF). Use for event subscription: echo filter | rdwrs event Also adds prototype in experiments/streamrdwr demonstrating: - Write topic filter, stream matching events - Glob-style filtering: * (single segment), > (rest)
This commit is contained in:
parent
6ea2bc052e
commit
9a66944859
|
|
@ -37,7 +37,7 @@ var (
|
|||
|
||||
func usage() {
|
||||
fmt.Fprintf(os.Stderr, "usage: ollie-9p [-a addr] cmd args...\n")
|
||||
fmt.Fprintf(os.Stderr, "commands: read write rdwr stat ls create rm mkdir mv\n")
|
||||
fmt.Fprintf(os.Stderr, "commands: read write rdwr rdwrs stat ls create rm mkdir mv\n")
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
|
|
@ -77,6 +77,11 @@ func main() {
|
|||
fatalf("usage: ollie-9p rdwr path")
|
||||
}
|
||||
cmdRdwr(args[0])
|
||||
case "rdwrs":
|
||||
if len(args) != 1 {
|
||||
fatalf("usage: ollie-9p rdwrs path")
|
||||
}
|
||||
cmdRdwrs(args[0])
|
||||
case "stat":
|
||||
if len(args) != 1 {
|
||||
fatalf("usage: ollie-9p stat path")
|
||||
|
|
@ -220,6 +225,43 @@ func cmdRdwr(path string) {
|
|||
}
|
||||
}
|
||||
|
||||
func cmdRdwrs(path string) {
|
||||
fsys := dial()
|
||||
defer fsys.Close()
|
||||
|
||||
data, err := io.ReadAll(os.Stdin)
|
||||
if err != nil {
|
||||
fatalf("read stdin: %v", err)
|
||||
}
|
||||
fid, err := fsys.Open(path, plan9.ORDWR)
|
||||
if err != nil {
|
||||
fatalf("open %s: %v", path, err)
|
||||
}
|
||||
defer fid.Close()
|
||||
|
||||
if _, err := fid.Write(data); err != nil {
|
||||
fatalf("write %s: %v", path, err)
|
||||
}
|
||||
if _, err := fid.Seek(0, 0); err != nil {
|
||||
fatalf("seek %s: %v", path, err)
|
||||
}
|
||||
|
||||
// Stream reads until EOF or error
|
||||
buf := make([]byte, 8192)
|
||||
for {
|
||||
n, err := fid.Read(buf)
|
||||
if n > 0 {
|
||||
os.Stdout.Write(buf[:n])
|
||||
}
|
||||
if err != nil {
|
||||
if err != io.EOF {
|
||||
fatalf("read %s: %v", path, err)
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func cmdStat(path string) {
|
||||
fsys := dial()
|
||||
defer fsys.Close()
|
||||
|
|
|
|||
|
|
@ -0,0 +1,295 @@
|
|||
// Prototype: streaming rdwr for event subscription
|
||||
//
|
||||
// Write topic filter, then stream matching events.
|
||||
// Example: echo "session.*.state" | ollie-9p -sock /tmp/test.sock rdwr event
|
||||
package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log"
|
||||
"net"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"9fans.net/go/plan9"
|
||||
)
|
||||
|
||||
// Simple pub/sub
|
||||
type Event struct {
|
||||
Topic string
|
||||
Payload string
|
||||
}
|
||||
|
||||
var (
|
||||
subsMu sync.RWMutex
|
||||
subs = make(map[chan Event]string) // channel -> topic filter
|
||||
)
|
||||
|
||||
func subscribe(filter string) chan Event {
|
||||
ch := make(chan Event, 64)
|
||||
subsMu.Lock()
|
||||
subs[ch] = filter
|
||||
subsMu.Unlock()
|
||||
return ch
|
||||
}
|
||||
|
||||
func unsubscribe(ch chan Event) {
|
||||
subsMu.Lock()
|
||||
delete(subs, ch)
|
||||
subsMu.Unlock()
|
||||
close(ch)
|
||||
}
|
||||
|
||||
func publish(topic, payload string) {
|
||||
ev := Event{Topic: topic, Payload: payload}
|
||||
subsMu.RLock()
|
||||
defer subsMu.RUnlock()
|
||||
for ch, filter := range subs {
|
||||
if matchTopic(filter, topic) {
|
||||
select {
|
||||
case ch <- ev:
|
||||
default:
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func matchTopic(filter, topic string) bool {
|
||||
if filter == "*" {
|
||||
return true
|
||||
}
|
||||
// Glob-style matching with * for single segment and > for rest
|
||||
// e.g. "session.*.agent.*.state" matches "session.abc.agent.xyz.state"
|
||||
filterParts := strings.Split(filter, ".")
|
||||
topicParts := strings.Split(topic, ".")
|
||||
|
||||
fi, ti := 0, 0
|
||||
for fi < len(filterParts) && ti < len(topicParts) {
|
||||
fp := filterParts[fi]
|
||||
if fp == ">" {
|
||||
// > matches rest
|
||||
return true
|
||||
}
|
||||
if fp == "*" {
|
||||
// * matches single segment
|
||||
fi++
|
||||
ti++
|
||||
continue
|
||||
}
|
||||
if fp != topicParts[ti] {
|
||||
return false
|
||||
}
|
||||
fi++
|
||||
ti++
|
||||
}
|
||||
return fi == len(filterParts) && ti == len(topicParts)
|
||||
}
|
||||
|
||||
// Per-fid state for streaming rdwr
|
||||
type fidState struct {
|
||||
mu sync.Mutex
|
||||
filter string // written topic filter
|
||||
events chan Event // subscription channel
|
||||
buf []byte // pending read data
|
||||
started bool // subscription started
|
||||
}
|
||||
|
||||
var (
|
||||
fidsMu sync.RWMutex
|
||||
fids = make(map[uint32]*fidState)
|
||||
)
|
||||
|
||||
func getFid(fid uint32) *fidState {
|
||||
fidsMu.RLock()
|
||||
f := fids[fid]
|
||||
fidsMu.RUnlock()
|
||||
return f
|
||||
}
|
||||
|
||||
func createFid(fid uint32) *fidState {
|
||||
f := &fidState{}
|
||||
fidsMu.Lock()
|
||||
fids[fid] = f
|
||||
fidsMu.Unlock()
|
||||
return f
|
||||
}
|
||||
|
||||
func removeFid(fid uint32) {
|
||||
fidsMu.Lock()
|
||||
f := fids[fid]
|
||||
delete(fids, fid)
|
||||
fidsMu.Unlock()
|
||||
if f != nil && f.events != nil {
|
||||
unsubscribe(f.events)
|
||||
}
|
||||
}
|
||||
|
||||
// Minimal 9P server
|
||||
func serve(conn net.Conn) {
|
||||
defer conn.Close()
|
||||
|
||||
for {
|
||||
fc, err := plan9.ReadFcall(conn)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
var resp *plan9.Fcall
|
||||
|
||||
switch fc.Type {
|
||||
case plan9.Tversion:
|
||||
resp = &plan9.Fcall{Type: plan9.Rversion, Tag: fc.Tag, Msize: fc.Msize, Version: "9P2000"}
|
||||
|
||||
case plan9.Tattach:
|
||||
resp = &plan9.Fcall{Type: plan9.Rattach, Tag: fc.Tag, Qid: plan9.Qid{Type: plan9.QTDIR}}
|
||||
|
||||
case plan9.Twalk:
|
||||
if len(fc.Wname) == 0 {
|
||||
// Clone fid
|
||||
if fc.Fid != fc.Newfid {
|
||||
createFid(fc.Newfid)
|
||||
}
|
||||
resp = &plan9.Fcall{Type: plan9.Rwalk, Tag: fc.Tag}
|
||||
} else if len(fc.Wname) == 1 && fc.Wname[0] == "event" {
|
||||
createFid(fc.Newfid)
|
||||
resp = &plan9.Fcall{Type: plan9.Rwalk, Tag: fc.Tag, Wqid: []plan9.Qid{{Type: plan9.QTFILE}}}
|
||||
} else {
|
||||
resp = &plan9.Fcall{Type: plan9.Rerror, Tag: fc.Tag, Ename: "not found"}
|
||||
}
|
||||
|
||||
case plan9.Topen:
|
||||
resp = &plan9.Fcall{Type: plan9.Ropen, Tag: fc.Tag, Qid: plan9.Qid{Type: plan9.QTFILE}}
|
||||
|
||||
case plan9.Tread:
|
||||
f := getFid(fc.Fid)
|
||||
if f == nil {
|
||||
resp = &plan9.Fcall{Type: plan9.Rerror, Tag: fc.Tag, Ename: "bad fid"}
|
||||
break
|
||||
}
|
||||
|
||||
f.mu.Lock()
|
||||
// Start subscription on first read
|
||||
// If no filter was written, default to "*" (all events)
|
||||
if !f.started {
|
||||
if f.filter == "" {
|
||||
f.filter = "*"
|
||||
}
|
||||
f.events = subscribe(f.filter)
|
||||
f.started = true
|
||||
}
|
||||
|
||||
// Try to get data
|
||||
if len(f.buf) == 0 && f.events != nil {
|
||||
f.mu.Unlock()
|
||||
// Block waiting for event
|
||||
select {
|
||||
case ev, ok := <-f.events:
|
||||
if !ok {
|
||||
resp = &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
|
||||
break
|
||||
}
|
||||
line := ev.Topic
|
||||
if ev.Payload != "" {
|
||||
line += " " + ev.Payload
|
||||
}
|
||||
f.mu.Lock()
|
||||
f.buf = []byte(line + "\n")
|
||||
}
|
||||
}
|
||||
|
||||
if resp == nil {
|
||||
// Return buffered data
|
||||
count := int(fc.Count)
|
||||
if count > len(f.buf) {
|
||||
count = len(f.buf)
|
||||
}
|
||||
data := f.buf[:count]
|
||||
f.buf = f.buf[count:]
|
||||
f.mu.Unlock()
|
||||
resp = &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: uint32(len(data)), Data: data}
|
||||
}
|
||||
|
||||
case plan9.Twrite:
|
||||
f := getFid(fc.Fid)
|
||||
if f == nil {
|
||||
resp = &plan9.Fcall{Type: plan9.Rerror, Tag: fc.Tag, Ename: "bad fid"}
|
||||
break
|
||||
}
|
||||
f.mu.Lock()
|
||||
f.filter = strings.TrimSpace(string(fc.Data))
|
||||
if f.filter == "" {
|
||||
f.filter = "*"
|
||||
}
|
||||
f.mu.Unlock()
|
||||
resp = &plan9.Fcall{Type: plan9.Rwrite, Tag: fc.Tag, Count: uint32(len(fc.Data))}
|
||||
|
||||
case plan9.Tclunk:
|
||||
removeFid(fc.Fid)
|
||||
resp = &plan9.Fcall{Type: plan9.Rclunk, Tag: fc.Tag}
|
||||
|
||||
case plan9.Tstat:
|
||||
d := plan9.Dir{
|
||||
Name: "event",
|
||||
Mode: 0666,
|
||||
Qid: plan9.Qid{Type: plan9.QTFILE},
|
||||
}
|
||||
stat, _ := d.Bytes()
|
||||
resp = &plan9.Fcall{Type: plan9.Rstat, Tag: fc.Tag, Stat: stat}
|
||||
|
||||
default:
|
||||
resp = &plan9.Fcall{Type: plan9.Rerror, Tag: fc.Tag, Ename: "not implemented"}
|
||||
}
|
||||
|
||||
if err := plan9.WriteFcall(conn, resp); err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func main() {
|
||||
sockPath := "/tmp/streamrdwr.sock"
|
||||
os.Remove(sockPath)
|
||||
ln, err := net.Listen("unix", sockPath)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
defer os.Remove(sockPath)
|
||||
|
||||
fmt.Printf("Listening on %s\n\n", sockPath)
|
||||
fmt.Println("Test commands:")
|
||||
fmt.Printf(" # All events:\n")
|
||||
fmt.Printf(" ollie-9p -sock %s read event\n\n", sockPath)
|
||||
fmt.Printf(" # Filtered events (rdwr pattern):\n")
|
||||
fmt.Printf(" echo '*' | ollie-9p -sock %s rdwr event\n", sockPath)
|
||||
fmt.Printf(" echo 'session.*.state' | ollie-9p -sock %s rdwr event\n\n", sockPath)
|
||||
|
||||
// Generate test events
|
||||
go func() {
|
||||
i := 0
|
||||
for {
|
||||
time.Sleep(500 * time.Millisecond)
|
||||
i++
|
||||
states := []string{"idle", "calling: shell", "thinking", "paused"}
|
||||
state := states[i%len(states)]
|
||||
publish("session.abc.agent.xyz.state", state)
|
||||
if i%3 == 0 {
|
||||
publish("session.abc.agent.xyz.new", "")
|
||||
}
|
||||
if i%5 == 0 {
|
||||
publish("session.abc.bypass.request", "id123\tcmd\tcwd")
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
fmt.Println("Publishing test events every 500ms...")
|
||||
for {
|
||||
conn, err := ln.Accept()
|
||||
if err != nil {
|
||||
log.Printf("accept: %v", err)
|
||||
continue
|
||||
}
|
||||
go serve(conn)
|
||||
}
|
||||
}
|
||||
Binary file not shown.
Loading…
Reference in New Issue