ollie/cmd/olliesrv/internal/fs/spec.go

761 lines
21 KiB
Go

package fs
// spec.go — the single source of truth for the olliesrv 9P namespace.
// Every handler is inline. No archaeological expeditions needed.
import (
"context"
"encoding/json"
"fmt"
"os"
"os/user"
"strconv"
"strings"
"syscall"
"ollie/cmd/olliesrv/internal/agent"
"ollie/cmd/olliesrv/internal/backend"
"ollie/cmd/olliesrv/internal/session"
"ollie/paths"
"ollie/virtfs"
)
// Type aliases for virtfs types used throughout the package.
type (
Tree = virtfs.Tree
File = virtfs.File
SyntheticFileInfo = virtfs.SyntheticFileInfo
)
// Structural permissions.
const (
PermChildDir = virtfs.PermChildDir
PermIdx = virtfs.PermIdx
PermMkdir = virtfs.PermMkdir
PermMkdirPrivate = virtfs.PermMkdirPrivate
)
var treeSpecUID, treeSpecGID string
// buildTreeSpec constructs the full 9P namespace.
func buildTreeSpec(cfg *Config) virtfs.FsNodeDecl {
mc := cfg.ModelCache
u, _ := user.Current()
treeSpecUID = "ollie"
treeSpecGID = "agent"
if u != nil {
treeSpecUID = u.Username
}
// Event subscription for eventwait (lives for server lifetime).
eventRead, eventSignal := session.EventValue(session.SubscribeEvents(context.Background(), "*"))
return virtfs.DirNode("/",
virtfs.UID(treeSpecUID),
virtfs.GID(treeSpecGID),
virtfs.FileNode("backends", 0444,
virtfs.Doc("Available backend names, one per line"),
virtfs.Read(func() ([]byte, error) {
return []byte(strings.Join(backend.Backends(), "\n") + "\n"), nil
}),
),
virtfs.FileNode("help", 0444,
virtfs.Doc("Filesystem reference"),
virtfs.Read(func() ([]byte, error) {
return []byte(helpFn()), nil
}),
),
virtfs.FileNode("models", 0444,
virtfs.Doc("Available models. Format: backend<tab>model per line."),
virtfs.Read(func() ([]byte, error) {
if mc != nil {
return mc.Get(), nil
}
return []byte("(no model cache)\n"), nil
}),
),
virtfs.FileNode("agents", 0444,
virtfs.Doc("Available agent profiles, one per line"),
virtfs.Read(func() ([]byte, error) {
var sb strings.Builder
for _, dir := range agent.AgentsDirs() {
entries, err := os.ReadDir(dir)
if err != nil {
continue
}
for _, e := range entries {
if !e.IsDir() && strings.HasSuffix(e.Name(), ".json") {
sb.WriteString(strings.TrimSuffix(e.Name(), ".json"))
sb.WriteByte('\n')
}
}
}
return []byte(sb.String()), nil
}),
),
virtfs.FileNode("ctl", 0666,
virtfs.Doc("Server control. Write: 'invalidate', 'kill'"),
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
return dispatch(map[string]func([]string) ([]byte, error){
"invalidate": func(_ []string) ([]byte, error) {
if mc != nil {
mc.Invalidate()
}
if cfg.Invalidate != nil {
cfg.Invalidate()
}
return []byte("ok\n"), nil
},
"kill": func(_ []string) ([]byte, error) {
if cfg.Shutdown != nil {
cfg.Shutdown()
}
return []byte("ok\n"), nil
},
}, data)
}),
),
virtfs.FileNode("eventwait", 0444,
virtfs.Doc("Blocking read for server events"),
virtfs.Read(func() ([]byte, error) {
data, _, _ := eventRead()
return data, nil
}),
virtfs.BlockOnce(eventRead, eventSignal),
),
virtfs.FileNode("generate", 0666,
virtfs.Doc("One-shot LLM generation"),
virtfs.Request(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))
}
result, err := backend.Generate(cfg.Ctx, req)
if err != nil {
return nil, err
}
return []byte(result + "\n"), nil
}),
),
virtfs.FileNode("aliases", 0444,
virtfs.Doc("Alias table. Format: id<tab>path per line."),
virtfs.Read(func() ([]byte, error) {
var sb strings.Builder
for name, sess := range session.Sessions() {
fmt.Fprintf(&sb, "%s\tsession/%s\n", sess.ID, name)
for _, ag := range sess.Agents() {
fmt.Fprintf(&sb, "%s\tsession/%s/agent/%s\n", ag.ID(), name, ag.Name())
}
}
return []byte(sb.String()), nil
}),
),
// session/
virtfs.DirNode("session",
virtfs.Doc("Session management"),
virtfs.FileNode("new", 0666,
virtfs.Doc("Create session"),
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
args := strings.Fields(string(data))
name, remote := "", ""
for _, arg := range args {
if k, v, ok := strings.Cut(arg, "="); ok {
switch k {
case "name":
name = v
case "remote":
remote = v
}
}
}
sess, err := session.CreateEmpty(name, remote)
if err != nil {
return nil, err
}
return []byte(sess.Name() + "\n"), nil
}),
),
virtfs.FileNode("idx", 0444,
virtfs.Doc("Session index"),
virtfs.Read(func() ([]byte, error) { return session.BuildIndex(), nil }),
),
virtfs.Each("{sname}", func() ([]virtfs.FsNodeDecl, error) {
return buildSessionEntries()
}),
),
)
}
// buildSessionEntries produces one FsNodeDecl per session.
func buildSessionEntries() ([]virtfs.FsNodeDecl, error) {
sessions := session.Sessions()
var out []virtfs.FsNodeDecl
for name, sess := range sessions {
s := sess
n := name
removeFn := func() error {
session.Kill(n)
return nil
}
renameFn := func(newName string) error {
return session.Rename(n, newName)
}
out = append(out, virtfs.FsNodeDecl{
Name: s.Name(),
Aliases: []string{s.ID},
Remove: removeFn,
Rename: renameFn,
Children: buildSessionChildren(s, removeFn, renameFn),
})
}
return out, nil
}
// buildSessionChildren returns the file nodes for a single session.
func buildSessionChildren(
s *session.Session,
removeFn func() error,
renameFn func(string) error,
) []virtfs.FsNodeDecl {
return []virtfs.FsNodeDecl{
virtfs.FileNode("env", 0444,
virtfs.Read(func() ([]byte, error) {
var sb strings.Builder
fmt.Fprintf(&sb, "OLLIE_SESSION_ID=%s\n", s.ID)
for _, e := range os.Environ() {
if strings.HasPrefix(e, "OLLIE_") && !strings.HasPrefix(e, "OLLIE_SESSION_ID=") {
sb.WriteString(e)
sb.WriteByte('\n')
}
}
return []byte(sb.String()), nil
}),
),
virtfs.FileNode("paused", 0444,
virtfs.Read(func() ([]byte, error) {
if s.IsPaused() {
return []byte("true\n"), nil
}
return []byte("false\n"), nil
}),
),
virtfs.FileNode("ctl", 0666,
virtfs.Request(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() },
"save": func(_ []string) ([]byte, error) {
s.Save()
return []byte("ok\n"), nil
},
"invalidate": func(_ []string) ([]byte, error) {
s.InvalidateModelsCache()
return []byte("ok\n"), nil
},
"pause": func(_ []string) ([]byte, error) {
if err := s.Pause(); err != nil {
return nil, err
}
return []byte("ok\n"), nil
},
"resume": func(_ []string) ([]byte, error) {
if err := s.Resume(); err != nil {
return nil, err
}
return []byte("ok\n"), nil
},
}, data)
}),
),
virtfs.FileNode("name", 0666,
virtfs.Read(func() ([]byte, error) { return []byte(s.Name() + "\n"), nil }),
virtfs.Write(func(data []byte) error {
newName := strings.TrimSpace(string(data))
if newName == "" || newName == s.Name() {
return nil
}
return renameFn(newName)
}),
),
virtfs.FileNode("id", 0444,
virtfs.Read(func() ([]byte, error) { return []byte(s.ID + "\n"), nil }),
),
virtfs.DirNode("agent",
virtfs.FileNode("new", 0666,
virtfs.Request(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)
}
wireAgentEvents(s.ID, ag)
return []byte(ag.ID() + "\n"), nil
}),
),
virtfs.FileNode("idx", 0444,
virtfs.Read(func() ([]byte, error) { return session.BuildAgentIndex(s), nil }),
),
virtfs.Each("{aname}", func() ([]virtfs.FsNodeDecl, error) {
return buildAgentEntries(s)
}),
),
}
}
// buildAgentEntries produces one FsNodeDecl per agent in a session.
func buildAgentEntries(s *session.Session) ([]virtfs.FsNodeDecl, error) {
agents := s.Agents()
var out []virtfs.FsNodeDecl
for _, ag := range agents {
a := ag
wireAgentEvents(s.ID, a)
out = append(out, virtfs.FsNodeDecl{
Name: a.Name(),
Aliases: []string{a.ID()},
Remove: func() error {
s.RemoveAgent(a.ID())
session.PublishEvent("session."+s.ID+".agent."+a.ID()+".kill", "")
return nil
},
Children: buildAgentChildren(a, s),
})
}
return out, nil
}
// buildAgentChildren returns the file nodes for a single agent.
func buildAgentChildren(a *agent.Agent, s *session.Session) []virtfs.FsNodeDecl {
chatStat := func() os.FileInfo {
return &SyntheticFileInfo{Name_: "chat", Mode_: 0444, Size_: 64 * 1024}
}
return []virtfs.FsNodeDecl{
virtfs.FileNode("prompt", 0666,
virtfs.Write(func(data []byte) error {
input := strings.TrimSpace(string(data))
if input == "" {
return nil
}
if input == "/invalidate" {
s.InvalidateModelsCache()
return nil
}
go func() {
a.Submit(s.Ctx, input)
a.EnsureTrailingNewline()
}()
return nil
}),
),
virtfs.FileNode("fifo", 0666,
virtfs.Write(func(data []byte) error {
input := strings.TrimSpace(string(data))
if input == "" {
return nil
}
a.Queue(input)
// If the agent is idle, trigger processing.
if !a.IsRunning() {
go func() {
if next, ok := a.PopQueue(); ok {
a.Submit(s.Ctx, next)
a.EnsureTrailingNewline()
}
}()
}
return nil
}),
virtfs.Read(func() ([]byte, error) {
item, ok := a.PopQueue()
if !ok {
return nil, nil
}
return []byte(item), nil
}),
),
virtfs.FileNode("feed", 0666,
virtfs.Write(func(data []byte) error {
if len(data) == 0 {
return nil
}
for _, b := range data {
if b != 0 {
a.FeedWrite(data)
return nil
}
}
return nil
}),
virtfs.BlockOnce(func() ([]byte, string, error) {
return a.Feed.Get(), a.Feed.Hash(), nil
}, a.SignalCh),
),
virtfs.FileNode("chat", 0444,
virtfs.StatOverride(chatStat),
virtfs.Read(func() ([]byte, error) {
a.ChatMu().RLock()
defer a.ChatMu().RUnlock()
return stripMarkers(a.ChatLog()), nil
}),
virtfs.Stream(func(base string) ([]byte, string, error) {
data, nextBase, err := a.ChatRead(base)
if err != nil || len(data) == 0 {
return data, nextBase, err
}
stripped := stripMarkers(data)
if len(stripped) > 0 {
return stripped, nextBase, nil
}
return nil, nextBase, nil
}, a.ChatSignal),
),
virtfs.FileNode("chat.raw", 0444,
virtfs.StatOverride(chatStat),
virtfs.Read(func() ([]byte, error) {
a.ChatMu().RLock()
defer a.ChatMu().RUnlock()
log := a.ChatLog()
data := make([]byte, len(log))
copy(data, log)
return data, nil
}),
virtfs.Stream(a.ChatRead, a.ChatSignal),
),
virtfs.FileNode("statewait", 0444,
virtfs.Read(func() ([]byte, error) {
return []byte(a.State() + "\n"), nil
}),
virtfs.BlockOnce(func() ([]byte, string, error) {
st := a.State()
return []byte(st + "\n"), st, nil
}, a.SignalCh),
),
virtfs.FileNode("log", 0444,
virtfs.Read(func() ([]byte, error) {
const maxWindow = 64 * 1024
a.ChatMu().RLock()
log := a.ChatLog()
start := 0
if len(log) > maxWindow {
start = len(log) - maxWindow
}
data := make([]byte, len(log)-start)
copy(data, log[start:])
a.ChatMu().RUnlock()
return data, nil
}),
),
virtfs.FileNode("plan", 0666,
virtfs.Read(func() ([]byte, error) {
return a.Plan(), nil
}),
virtfs.Write(func(data []byte) error {
a.SetPlan(data)
return nil
}),
),
virtfs.FileNode("cfg", 0666,
virtfs.Read(func() ([]byte, error) {
p := a.GenParams()
var sb strings.Builder
fmt.Fprintf(&sb, "name=%s\n", s.ID)
fmt.Fprintf(&sb, "backend=%s\n", a.BackendName())
fmt.Fprintf(&sb, "model=%s\n", a.ModelName())
fmt.Fprintf(&sb, "id=%s\n", a.ID())
fmt.Fprintf(&sb, "profile=%s\n", a.Profile())
fmt.Fprintf(&sb, "displayName=%s\n", a.Name())
fmt.Fprintf(&sb, "cwd=%s\n", a.Cwd())
fmt.Fprintf(&sb, "remote=%s\n", s.Remote)
fmt.Fprintf(&sb, "maxTokens=%d\n", p.MaxTokens)
if p.Temperature != nil {
fmt.Fprintf(&sb, "temperature=%g\n", *p.Temperature)
}
if p.TopP != nil {
fmt.Fprintf(&sb, "topP=%g\n", *p.TopP)
}
if p.TopK != nil {
fmt.Fprintf(&sb, "topK=%d\n", *p.TopK)
}
if len(p.Stop) > 0 {
fmt.Fprintf(&sb, "stop=%s\n", strings.Join(p.Stop, ","))
}
return []byte(sb.String()), nil
}),
virtfs.Write(func(data []byte) error {
input := strings.TrimSpace(string(data))
if input == "" {
return nil
}
parts := strings.SplitN(input, "=", 2)
if len(parts) != 2 {
return fmt.Errorf("invalid cfg format (expected key=value)")
}
switch parts[0] {
case "name":
if strings.TrimSpace(parts[1]) != "" {
a.SetName(strings.TrimSpace(parts[1]))
}
default:
return fmt.Errorf("unknown cfg key: %s", parts[0])
}
return nil
}),
),
virtfs.FileNode("ctl", 0666,
virtfs.Request(func(_ context.Context, data []byte) ([]byte, error) {
return dispatch(map[string]func([]string) ([]byte, error){
"kill": func(_ []string) ([]byte, error) {
s.RemoveAgent(a.ID())
session.PublishEvent("session."+s.ID+".agent."+a.ID()+".kill", "")
return []byte("ok\n"), nil
},
"stop": func(_ []string) ([]byte, error) {
a.Interrupt(agent.ErrInterrupted)
return []byte("ok\n"), nil
},
"compact": func(_ []string) ([]byte, error) {
if err := a.Compact(s.Ctx); err != nil {
return nil, err
}
return []byte("ok\n"), nil
},
"clear": func(_ []string) ([]byte, error) {
if err := a.Clear(); err != nil {
return nil, err
}
return []byte("ok\n"), nil
},
"inject": func(args []string) ([]byte, error) {
text := strings.Join(args, " ")
if text == "" {
return nil, fmt.Errorf("inject requires text")
}
a.InjectRewrite(text)
return []byte("ok\n"), nil
},
"i": func(args []string) ([]byte, error) {
text := strings.Join(args, " ")
if text == "" {
return nil, fmt.Errorf("inject requires text")
}
a.InjectRewrite(text)
return []byte("ok\n"), nil
},
"agent": func(args []string) ([]byte, error) {
if len(args) == 0 {
return []byte(a.Profile() + "\n"), nil
}
if err := a.SwitchProfile(args[0]); err != nil {
return nil, err
}
return []byte(args[0] + "\n"), nil
},
"model": func(args []string) ([]byte, error) {
if len(args) == 0 {
if be := a.Backend(); be != nil {
return []byte(be.Model() + "\n"), nil
}
return nil, nil
}
if be := a.Backend(); be != nil {
be.SetModel(strings.Join(args, " "))
}
return []byte(strings.Join(args, " ") + "\n"), nil
},
"models": func(_ []string) ([]byte, error) {
return []byte(s.CachedListModels() + "\n"), nil
},
"tools": func(_ []string) ([]byte, error) {
ts := a.ToolServer()
if ts == nil {
return []byte("(no tool server)\n"), nil
}
loaded, err := ts.ListTools()
if err != nil {
return nil, err
}
var sb strings.Builder
for _, ti := range loaded {
if ti.Description != "" {
fmt.Fprintf(&sb, "%-20s %s\n", ti.Name, ti.Description)
} else {
sb.WriteString(ti.Name + "\n")
}
}
return []byte(sb.String()), nil
},
"tool_load": func(args []string) ([]byte, error) {
if len(args) == 0 {
return nil, fmt.Errorf("tool_load requires a tool name")
}
ts := a.ToolServer()
if ts == nil {
return nil, fmt.Errorf("no tool server")
}
if err := ts.LoadTool(args[0]); err != nil {
return nil, err
}
if infos, err := ts.ListTools(); err == nil {
a.SetToolsPreamble(agent.RenderTools(infos))
}
return []byte(args[0] + "\n"), nil
},
"tools_all": func(args []string) ([]byte, error) {
ts := a.ToolServer()
if ts == nil {
return nil, fmt.Errorf("no tool server")
}
return ts.ListAllTools()
},
"tool_unload": func(args []string) ([]byte, error) {
if len(args) == 0 {
return nil, fmt.Errorf("tool_unload requires a tool name")
}
ts := a.ToolServer()
if ts == nil {
return nil, fmt.Errorf("no tool server")
}
if err := ts.UnloadTool(args[0]); err != nil {
return nil, err
}
if infos, err := ts.ListTools(); err == nil {
a.SetToolsPreamble(agent.RenderTools(infos))
}
return []byte(args[0] + "\n"), nil
},
"backend": func(args []string) ([]byte, error) {
if len(args) == 0 {
if be := a.Backend(); be != nil {
return []byte(be.Name() + "\n"), nil
}
return nil, nil
}
return nil, fmt.Errorf("backend switching not supported via ctl; use cfg")
},
"name": func(args []string) ([]byte, error) {
if len(args) == 0 {
return []byte(a.Name() + "\n"), nil
}
a.SetName(strings.Join(args, " "))
return []byte(strings.Join(args, " ") + "\n"), nil
},
"cwd": func(args []string) ([]byte, error) {
if len(args) == 0 {
return []byte(a.Cwd() + "\n"), nil
}
dir := paths.ExpandHome(strings.Join(args, " "))
if _, err := os.Stat(dir); err != nil {
return nil, fmt.Errorf("cwd: %w", err)
}
a.SetCWD(dir)
return []byte(dir + "\n"), nil
},
"proc": func(args []string) ([]byte, error) {
ts := a.ToolServer()
if ts == nil {
return nil, fmt.Errorf("no tool server")
}
if len(args) == 0 || args[0] == "top" {
data, err := ts.ListProcs()
if err != nil {
return nil, err
}
if len(data) == 0 {
return []byte("(no background processes)\n"), nil
}
return data, nil
}
subcmd := args[0]
switch subcmd {
case "term":
if len(args) < 2 {
return nil, fmt.Errorf("proc term requires pid")
}
pid, err := strconv.Atoi(args[1])
if err != nil {
return nil, fmt.Errorf("invalid pid: %s", args[1])
}
if err := ts.SignalDetached(pid, syscall.SIGTERM); err != nil {
return nil, err
}
return []byte("ok\n"), nil
case "kill":
if len(args) < 2 {
return nil, fmt.Errorf("proc kill requires pid")
}
pid, err := strconv.Atoi(args[1])
if err != nil {
return nil, fmt.Errorf("invalid pid: %s", args[1])
}
if err := ts.SignalDetached(pid, syscall.SIGKILL); err != nil {
return nil, err
}
return []byte("ok\n"), nil
case "out":
if len(args) < 2 {
return nil, fmt.Errorf("proc out requires pid")
}
pid, err := strconv.Atoi(args[1])
if err != nil {
return nil, fmt.Errorf("invalid pid: %s", args[1])
}
out, err := ts.GetDetachedOutput(pid)
if err != nil {
return nil, err
}
return []byte(out), nil
case "dismiss":
if len(args) < 2 {
return nil, fmt.Errorf("proc dismiss requires pid")
}
pid, err := strconv.Atoi(args[1])
if err != nil {
return nil, fmt.Errorf("invalid pid: %s", args[1])
}
if !ts.DismissDetached(pid) {
return nil, fmt.Errorf("process %d not found", pid)
}
return []byte("ok\n"), nil
default:
return nil, fmt.Errorf("unknown proc subcommand: %s (use: top, term <pid>, kill <pid>, out <pid>, dismiss <pid>)", subcmd)
}
},
"systemprompt": func(_ []string) ([]byte, error) {
return []byte(a.SystemPrompt() + "\n"), nil
},
"help": func(_ []string) ([]byte, error) {
return []byte("stop compact clear inject agent model models tools tool_load tool_unload cwd proc name backend systemprompt help\n"), nil
},
}, data)
}),
),
virtfs.FileNode("stats", 0444,
virtfs.Read(func() ([]byte, error) {
return []byte("usage=" + a.UsageStr() + "\ncost=" + a.CostStr() + "\nctxsz=" + a.CtxSz() + "\n"), nil
}),
),
virtfs.FileNode("name", 0666,
virtfs.Read(func() ([]byte, error) { return []byte(a.Name() + "\n"), nil }),
virtfs.Write(func(data []byte) error {
newName := strings.TrimSpace(string(data))
if newName == "" {
return nil
}
a.SetName(newName)
session.PersistSession(s.Name())
return nil
}),
),
virtfs.FileNode("id", 0444,
virtfs.Read(func() ([]byte, error) { return []byte(a.ID() + "\n"), nil }),
),
}
}