virtfs: rename Request to Rdwr
Read, Write, and Rdwr are the three atomic 9P operations: - Read: non-blocking read - Write: non-blocking write (fire-and-forget) - Rdwr: atomic write-then-read (blocking, produces result) BlockOnce and Stream are special cases of Read. Rdwr is its own primitive — not a variant of either.
This commit is contained in:
parent
8b6b08b54f
commit
ab4668cb3f
|
|
@ -97,7 +97,7 @@ func buildTreeSpec(cfg *Config) virtfs.FsNodeDecl {
|
|||
),
|
||||
virtfs.FileNode("ctl", 0666,
|
||||
virtfs.Doc("Server control. Write: 'invalidate', 'kill'"),
|
||||
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
return dispatch(map[string]func([]string) ([]byte, error){
|
||||
"invalidate": func(_ []string) ([]byte, error) {
|
||||
if mc != nil {
|
||||
|
|
@ -127,7 +127,7 @@ func buildTreeSpec(cfg *Config) virtfs.FsNodeDecl {
|
|||
),
|
||||
virtfs.FileNode("generate", 0666,
|
||||
virtfs.Doc("One-shot LLM generation"),
|
||||
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
var req backend.GenerateRequest
|
||||
if err := json.Unmarshal(data, &req); err != nil {
|
||||
req.Prompt = strings.TrimSpace(string(data))
|
||||
|
|
@ -158,7 +158,7 @@ func buildTreeSpec(cfg *Config) virtfs.FsNodeDecl {
|
|||
virtfs.Doc("Session management"),
|
||||
virtfs.FileNode("new", 0666,
|
||||
virtfs.Doc("Create session"),
|
||||
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
args := strings.Fields(string(data))
|
||||
name, remote := "", ""
|
||||
for _, arg := range args {
|
||||
|
|
@ -246,7 +246,7 @@ func buildSessionChildren(
|
|||
}),
|
||||
),
|
||||
virtfs.FileNode("ctl", 0666,
|
||||
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
return dispatch(map[string]func([]string) ([]byte, error){
|
||||
"kill": func(_ []string) ([]byte, error) { return []byte("ok\n"), removeFn() },
|
||||
".": func(_ []string) ([]byte, error) { return []byte("ok\n"), removeFn() },
|
||||
|
|
@ -288,7 +288,7 @@ func buildSessionChildren(
|
|||
),
|
||||
virtfs.DirNode("agent",
|
||||
virtfs.FileNode("new", 0666,
|
||||
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
ag, err := session.CreateAgent(s.Name(), strings.Fields(string(data)))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create agent: %w", err)
|
||||
|
|
@ -509,7 +509,7 @@ func buildAgentChildren(a *agent.Agent, s *session.Session) []virtfs.FsNodeDecl
|
|||
}),
|
||||
),
|
||||
virtfs.FileNode("ctl", 0666,
|
||||
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
return dispatch(map[string]func([]string) ([]byte, error){
|
||||
"kill": func(_ []string) ([]byte, error) {
|
||||
s.RemoveAgent(a.ID())
|
||||
|
|
|
|||
|
|
@ -596,7 +596,7 @@ func write(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
return errFcall(fc, "bad fid")
|
||||
}
|
||||
|
||||
if f.entry != nil && f.entry.RequestRespMode() {
|
||||
if f.entry != nil && f.entry.RdwrMode() {
|
||||
f.readCache = nil // clear cached response for next read
|
||||
cs.mu.Unlock()
|
||||
if err := f.entry.Write(fc.Data); err != nil {
|
||||
|
|
@ -681,7 +681,7 @@ func clunk(root *fs.Tree, cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
data = make([]byte, len(f.writeBuf))
|
||||
copy(data, f.writeBuf)
|
||||
if f.entry != nil {
|
||||
isRequest = f.entry.RequestRespMode()
|
||||
isRequest = f.entry.RdwrMode()
|
||||
}
|
||||
}
|
||||
delete(cs.fids, fc.Fid)
|
||||
|
|
|
|||
|
|
@ -118,7 +118,7 @@ func Spec(srv *Server) virtfs.FsNodeDecl {
|
|||
),
|
||||
virtfs.FileNode("tools", 0666,
|
||||
virtfs.Doc("Tools: write agent ID, read tool list (rdwr)"),
|
||||
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
agentID := strings.TrimSpace(string(data))
|
||||
if agentID == "" {
|
||||
return nil, fmt.Errorf("agent ID required")
|
||||
|
|
@ -179,14 +179,14 @@ func Spec(srv *Server) virtfs.FsNodeDecl {
|
|||
virtfs.DirNode("proc",
|
||||
virtfs.FileNode("list", 0666,
|
||||
virtfs.Doc("List processes: write agent ID (or empty for all), read pid\\tstate\\ttool"),
|
||||
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
agentID := strings.TrimSpace(string(data))
|
||||
return []byte(srv.Fs.ListProcsForAgent(agentID)), nil
|
||||
}),
|
||||
),
|
||||
virtfs.FileNode("new", 0666,
|
||||
virtfs.Doc("Execute tool: write token + tool + args, read result (blocking)"),
|
||||
virtfs.Request(func(ctx context.Context, data []byte) ([]byte, error) {
|
||||
virtfs.Rdwr(func(ctx context.Context, data []byte) ([]byte, error) {
|
||||
result, err := srv.Fs.HandleProcNew(ctx, string(data))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -196,7 +196,7 @@ func Spec(srv *Server) virtfs.FsNodeDecl {
|
|||
),
|
||||
virtfs.FileNode("new.bg", 0666,
|
||||
virtfs.Doc("Execute tool in background: write token + tool + args, read pid"),
|
||||
virtfs.Request(func(ctx context.Context, data []byte) ([]byte, error) {
|
||||
virtfs.Rdwr(func(ctx context.Context, data []byte) ([]byte, error) {
|
||||
result, err := srv.Fs.HandleProcNewBg(ctx, string(data))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
|
|||
|
|
@ -372,7 +372,7 @@ func handleRead(cs *connState, fc *plan9.Fcall, tree *virtfs.Tree) *plan9.Fcall
|
|||
}
|
||||
|
||||
// For rdwr files, check if we've already read
|
||||
if f.entry.RequestRespMode() && f.readDone {
|
||||
if f.entry.RdwrMode() && f.readDone {
|
||||
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Data: nil}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -18,7 +18,7 @@ func BuildTree(spec FsNodeDecl) *Tree {
|
|||
func validate(d *FsNodeDecl, path string) {
|
||||
hasChildren := len(d.Children) > 0
|
||||
hasBindings := d.Bindings != nil
|
||||
hasFileHandlers := d.Read != nil || d.Write != nil || d.BlockOnce != nil || d.Stream != nil || d.Request != nil
|
||||
hasFileHandlers := d.Read != nil || d.Write != nil || d.BlockOnce != nil || d.Stream != nil || d.Rdwr != nil
|
||||
|
||||
isDir := hasChildren
|
||||
|
||||
|
|
@ -39,11 +39,11 @@ func validate(d *FsNodeDecl, path string) {
|
|||
if d.Stream != nil {
|
||||
blocking++
|
||||
}
|
||||
if d.Request != nil {
|
||||
if d.Rdwr != nil {
|
||||
blocking++
|
||||
}
|
||||
if blocking > 1 {
|
||||
panic(fmt.Sprintf("BuildTree: %q: at most one of BlockOnce, Stream, Request may be set", fullPath))
|
||||
panic(fmt.Sprintf("BuildTree: %q: at most one of BlockOnce, Stream, Rdwr may be set", fullPath))
|
||||
}
|
||||
|
||||
// Template nodes must have Bindings set.
|
||||
|
|
@ -145,7 +145,7 @@ func listDir(d *FsNodeDecl, uid, gid string) ([]os.DirEntry, error) {
|
|||
if m == 0 {
|
||||
m = 0755
|
||||
}
|
||||
hasFileHandlers := c.Read != nil || c.Write != nil || c.BlockOnce != nil || c.Stream != nil || c.Request != nil
|
||||
hasFileHandlers := c.Read != nil || c.Write != nil || c.BlockOnce != nil || c.Stream != nil || c.Rdwr != nil
|
||||
isChildDir := len(c.Children) > 0 || (c.Bindings != nil && !hasFileHandlers)
|
||||
if isChildDir {
|
||||
entries = append(entries, dirEntry(c.Name, m))
|
||||
|
|
@ -186,7 +186,7 @@ func statDir(d *FsNodeDecl, uid, gid, name string) (os.FileInfo, error) {
|
|||
}
|
||||
cu := resolveUID(child, uid)
|
||||
cg := resolveGID(child, gid)
|
||||
hasFileHandlers := child.Read != nil || child.Write != nil || child.BlockOnce != nil || child.Stream != nil || child.Request != nil
|
||||
hasFileHandlers := child.Read != nil || child.Write != nil || child.BlockOnce != nil || child.Stream != nil || child.Rdwr != nil
|
||||
isChildDir := len(child.Children) > 0 || (child.Bindings != nil && !hasFileHandlers)
|
||||
if isChildDir {
|
||||
return &SyntheticFileInfo{Name_: name, Mode_: cm | os.ModeDir, IsDir_: true, UID_: cu, GID_: cg}, nil
|
||||
|
|
@ -247,7 +247,7 @@ func openLeaf(d *FsNodeDecl, uid, gid string, mode os.FileMode, name string) (Fi
|
|||
}
|
||||
|
||||
var blockFn func(context.Context, string) ([]byte, string, error)
|
||||
var isBlocking, isStreaming, isRequest bool
|
||||
var isBlocking, isStreaming, isRdwr bool
|
||||
switch {
|
||||
case d.BlockOnce != nil:
|
||||
blockFn = d.BlockOnce
|
||||
|
|
@ -255,8 +255,8 @@ func openLeaf(d *FsNodeDecl, uid, gid string, mode os.FileMode, name string) (Fi
|
|||
case d.Stream != nil:
|
||||
blockFn = d.Stream
|
||||
isStreaming = true
|
||||
case d.Request != nil:
|
||||
fn := d.Request
|
||||
case d.Rdwr != nil:
|
||||
fn := d.Rdwr
|
||||
var result []byte
|
||||
var reqCtx context.Context
|
||||
var reqCancel context.CancelFunc
|
||||
|
|
@ -285,7 +285,7 @@ func openLeaf(d *FsNodeDecl, uid, gid string, mode os.FileMode, name string) (Fi
|
|||
}
|
||||
return nil
|
||||
}
|
||||
isRequest = true
|
||||
isRdwr = true
|
||||
return &fileConfig{
|
||||
StatFn: func() (os.FileInfo, error) {
|
||||
if d.Stat != nil {
|
||||
|
|
@ -300,7 +300,7 @@ func openLeaf(d *FsNodeDecl, uid, gid string, mode os.FileMode, name string) (Fi
|
|||
BlockingReadFn: blockFn,
|
||||
BlockingRead_: isBlocking,
|
||||
StreamMode_: isStreaming,
|
||||
RequestResp_: isRequest,
|
||||
Rdwr_: isRdwr,
|
||||
}, nil
|
||||
default:
|
||||
blockFn = notBlocking
|
||||
|
|
@ -319,7 +319,7 @@ func openLeaf(d *FsNodeDecl, uid, gid string, mode os.FileMode, name string) (Fi
|
|||
BlockingReadFn: blockFn,
|
||||
BlockingRead_: isBlocking,
|
||||
StreamMode_: isStreaming,
|
||||
RequestResp_: isRequest,
|
||||
Rdwr_: isRdwr,
|
||||
}, nil
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -63,8 +63,8 @@ type FsNodeDecl struct {
|
|||
BlockOnce func(ctx context.Context, base string) ([]byte, string, error)
|
||||
Stream func(ctx context.Context, base string) ([]byte, string, error)
|
||||
|
||||
// Request — rdwr: write triggers result, read returns it once.
|
||||
Request func(ctx context.Context, data []byte) ([]byte, error)
|
||||
// Rdwr — atomic write-then-read: write triggers computation, read returns result once.
|
||||
Rdwr func(ctx context.Context, data []byte) ([]byte, error)
|
||||
|
||||
// Tremove + Trename.
|
||||
Remove func() error
|
||||
|
|
@ -187,9 +187,10 @@ func BlockOnceRaw(fn func(context.Context, string) ([]byte, string, error)) Node
|
|||
return func(d *FsNodeDecl) { d.BlockOnce = fn }
|
||||
}
|
||||
|
||||
// Request sets a context-aware request-response handler (rdwr pattern).
|
||||
func Request(fn func(context.Context, []byte) ([]byte, error)) NodeOption {
|
||||
return func(d *FsNodeDecl) { d.Request = fn }
|
||||
// Rdwr sets an atomic write-then-read handler.
|
||||
// Write triggers computation; the subsequent read returns the result once.
|
||||
func Rdwr(fn func(context.Context, []byte) ([]byte, error)) NodeOption {
|
||||
return func(d *FsNodeDecl) { d.Rdwr = fn }
|
||||
}
|
||||
|
||||
// RemoveNode sets the removal handler for a node.
|
||||
|
|
|
|||
|
|
@ -31,12 +31,12 @@ func TestNodeOptions(t *testing.T) {
|
|||
}
|
||||
})
|
||||
|
||||
t.Run("Request", func(t *testing.T) {
|
||||
spec := FileNode("test", 0666, Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
t.Run("Rdwr", func(t *testing.T) {
|
||||
spec := FileNode("test", 0666, Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
return nil, nil
|
||||
}))
|
||||
if spec.Request == nil {
|
||||
t.Error("Request not set")
|
||||
if spec.Rdwr == nil {
|
||||
t.Error("Rdwr not set")
|
||||
}
|
||||
})
|
||||
|
||||
|
|
@ -161,7 +161,7 @@ func TestFileConfigMethods(t *testing.T) {
|
|||
|
||||
func TestFileConfigModes(t *testing.T) {
|
||||
fc := &fileConfig{}
|
||||
if fc.BlockingReadMode() || fc.StreamMode() || fc.RequestRespMode() {
|
||||
if fc.BlockingReadMode() || fc.StreamMode() || fc.RdwrMode() {
|
||||
t.Error("default modes should be false")
|
||||
}
|
||||
|
||||
|
|
@ -171,8 +171,8 @@ func TestFileConfigModes(t *testing.T) {
|
|||
if !(&fileConfig{StreamMode_: true}).StreamMode() {
|
||||
t.Error("StreamMode should be true")
|
||||
}
|
||||
if !(&fileConfig{RequestResp_: true}).RequestRespMode() {
|
||||
t.Error("RequestRespMode should be true")
|
||||
if !(&fileConfig{Rdwr_: true}).RdwrMode() {
|
||||
t.Error("RdwrMode should be true")
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -279,9 +279,9 @@ func TestOpenLeafStream(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestOpenLeafRequest(t *testing.T) {
|
||||
func TestOpenLeafRdwr(t *testing.T) {
|
||||
spec := DirNode("/",
|
||||
FileNode("request", 0666, Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
FileNode("request", 0666, Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
return append([]byte("echo:"), data...), nil
|
||||
})),
|
||||
)
|
||||
|
|
@ -290,8 +290,8 @@ func TestOpenLeafRequest(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
if !f.RequestRespMode() {
|
||||
t.Error("expected RequestRespMode")
|
||||
if !f.RdwrMode() {
|
||||
t.Error("expected RdwrMode")
|
||||
}
|
||||
if err := f.Write([]byte("test")); err != nil {
|
||||
t.Fatalf("Write: %v", err)
|
||||
|
|
@ -305,9 +305,9 @@ func TestOpenLeafRequest(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestOpenLeafRequestError(t *testing.T) {
|
||||
func TestOpenLeafRdwrError(t *testing.T) {
|
||||
spec := DirNode("/",
|
||||
FileNode("request", 0666, Request(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
FileNode("request", 0666, Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
||||
return nil, errors.New("request error")
|
||||
})),
|
||||
)
|
||||
|
|
@ -321,11 +321,11 @@ func TestOpenLeafRequestError(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestOpenLeafRequestWithFallbackRead(t *testing.T) {
|
||||
func TestOpenLeafRdwrWithFallbackRead(t *testing.T) {
|
||||
spec := DirNode("/",
|
||||
FileNode("request", 0666,
|
||||
Read(func() ([]byte, error) { return []byte("fallback"), nil }),
|
||||
Request(func(_ context.Context, data []byte) ([]byte, error) { return []byte("response"), nil }),
|
||||
Rdwr(func(_ context.Context, data []byte) ([]byte, error) { return []byte("response"), nil }),
|
||||
),
|
||||
)
|
||||
tree := BuildTree(spec)
|
||||
|
|
@ -342,11 +342,11 @@ func TestOpenLeafRequestWithFallbackRead(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestOpenLeafRequestCancelOnClose(t *testing.T) {
|
||||
func TestOpenLeafRdwrCancelOnClose(t *testing.T) {
|
||||
canceled := make(chan struct{})
|
||||
started := make(chan struct{})
|
||||
spec := DirNode("/",
|
||||
FileNode("request", 0666, Request(func(ctx context.Context, data []byte) ([]byte, error) {
|
||||
FileNode("request", 0666, Rdwr(func(ctx context.Context, data []byte) ([]byte, error) {
|
||||
close(started)
|
||||
<-ctx.Done()
|
||||
close(canceled)
|
||||
|
|
|
|||
|
|
@ -25,10 +25,10 @@ type File interface {
|
|||
// Read-side contracts (mutually exclusive — at most one true):
|
||||
// BlockingReadMode — blocks up to 5s for next value, returns it or empty, EOF on 2nd read
|
||||
// StreamMode — blocks indefinitely, streams data incrementally, FD lives until EOF
|
||||
// RequestRespMode — write blocks until result produced, read returns result once (rdwr)
|
||||
// RdwrMode — write blocks until result produced, read returns result once
|
||||
BlockingReadMode() bool
|
||||
StreamMode() bool
|
||||
RequestRespMode() bool
|
||||
RdwrMode() bool
|
||||
}
|
||||
|
||||
// fileConfig implements File via function pointers.
|
||||
|
|
@ -40,7 +40,7 @@ type fileConfig struct {
|
|||
BlockingReadFn func(context.Context, string) ([]byte, string, error)
|
||||
BlockingRead_ bool // one-shot waiter: block 5s, return once, EOF on 2nd read
|
||||
StreamMode_ bool // indefinite stream: block forever, keep FD open
|
||||
RequestResp_ bool // rdwr: write blocks on write, read returns result once
|
||||
Rdwr_ bool // rdwr: write blocks on write, read returns result once
|
||||
}
|
||||
|
||||
func (e *fileConfig) Stat() (os.FileInfo, error) { return e.StatFn() }
|
||||
|
|
@ -57,7 +57,7 @@ func (e *fileConfig) BlockingRead(ctx context.Context, base string) ([]byte, str
|
|||
}
|
||||
func (e *fileConfig) BlockingReadMode() bool { return e.BlockingRead_ }
|
||||
func (e *fileConfig) StreamMode() bool { return e.StreamMode_ }
|
||||
func (e *fileConfig) RequestRespMode() bool { return e.RequestResp_ }
|
||||
func (e *fileConfig) RdwrMode() bool { return e.Rdwr_ }
|
||||
|
||||
// SyntheticFileInfo implements os.FileInfo for entries with no backing file.
|
||||
type SyntheticFileInfo struct {
|
||||
|
|
|
|||
|
|
@ -256,7 +256,7 @@ func TestGenerateHelpModes(t *testing.T) {
|
|||
spec := DirNode("/",
|
||||
FileNode("ro", 0444, Doc("read only file")),
|
||||
FileNode("rw", 0644, Doc("read write file"), Write(func(_ []byte) error { return nil })),
|
||||
FileNode("ctl", 0200, Doc("control file"), Request(func(_ context.Context, _ []byte) ([]byte, error) { return nil, nil })),
|
||||
FileNode("ctl", 0200, Doc("control file"), Rdwr(func(_ context.Context, _ []byte) ([]byte, error) { return nil, nil })),
|
||||
)
|
||||
help := GenerateHelp(spec)
|
||||
if help == "" {
|
||||
|
|
|
|||
Loading…
Reference in New Issue