diff --git a/cmd/ollie-9p/main.go b/cmd/ollie-9p/main.go index e7ac640..cd15426 100644 --- a/cmd/ollie-9p/main.go +++ b/cmd/ollie-9p/main.go @@ -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() diff --git a/experiments/streamrdwr/main.go b/experiments/streamrdwr/main.go new file mode 100644 index 0000000..3592081 --- /dev/null +++ b/experiments/streamrdwr/main.go @@ -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) + } +} diff --git a/experiments/streamrdwr/streamrdwr b/experiments/streamrdwr/streamrdwr new file mode 100755 index 0000000..282026b Binary files /dev/null and b/experiments/streamrdwr/streamrdwr differ