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

1283 lines
38 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/exec"
"os/user"
"sort"
"strconv"
"strings"
"syscall"
"time"
"ollie/cmd/olliesrv/internal/agent"
"ollie/cmd/olliesrv/internal/backend"
"ollie/cmd/olliesrv/internal/metrics"
"ollie/cmd/olliesrv/internal/session"
"ollie/util"
"ollie/virtfs"
)
// Type aliases for virtfs types used throughout the package.
type (
Tree = virtfs.Tree
File = virtfs.File
SyntheticFileInfo = virtfs.SyntheticFileInfo
)
var treeSpecUID, treeSpecGID string
const asyncWorkLimit = 64
var asyncWorkSlots = make(chan struct{}, asyncWorkLimit)
func startAsync(ctx context.Context, fn func()) bool {
select {
case asyncWorkSlots <- struct{}{}:
case <-ctx.Done():
return false
default:
return false
}
go func() {
defer func() { <-asyncWorkSlots }()
fn()
}()
return true
}
// runWorkflow executes a workflow script in the background.
// If variant is non-empty and not "default", the matching conf sidecar is
// sourced as environment variables before the script runs.
// If workflow is empty or "none", nothing happens.
func runWorkflow(s *session.Session, workflow, variant string) {
if workflow == "" || workflow == "none" {
return
}
startAsync(s.Ctx, func() {
workflowDir := util.CfgDir() + "/workflows/"
scriptPath := workflowDir + workflow
if _, err := os.Stat(scriptPath); err != nil {
s.SetGoalStatus("error: workflow not found: " + workflow)
return
}
cwd := s.Cwd()
cmd := exec.CommandContext(s.Ctx, scriptPath)
cmd.Env = append(os.Environ(),
"OLLIE_SESSION_ID="+s.ID,
"OLLIE_SESSION_NAME="+s.Name(),
"OLLIE_CWD="+cwd,
)
// Source variant conf if specified.
if variant != "" && variant != "default" {
confPath := workflowDir + workflow + "-" + variant + ".conf"
confData, err := os.ReadFile(confPath)
if err != nil {
s.SetGoalStatus("error: variant conf not found: " + workflow + "-" + variant + ".conf")
return
}
for _, line := range strings.Split(string(confData), "\n") {
line = strings.TrimSpace(line)
if line == "" || line[0] == '#' {
continue
}
if _, _, ok := strings.Cut(line, "="); ok {
cmd.Env = append(cmd.Env, line)
}
}
}
cmd.Dir = cwd
out, err := cmd.CombinedOutput()
if err != nil {
msg := strings.TrimSpace(string(out))
if msg == "" {
msg = err.Error()
}
s.SetGoalStatus("error: " + msg)
}
// Script exited — don't auto-set complete; the conductor agent handles that.
})
}
// 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
}
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[<tab>in<tab>out<tab>cache_read<tab>cache_write]. Pricing per 1M tokens USD; (e) suffix marks estimates."),
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("workflows", 0444,
virtfs.Doc("Available workflows with variants. Format: name\\tdefault,variant1,variant2"),
virtfs.Read(func() ([]byte, error) {
dir := util.CfgDir() + "/workflows"
entries, err := os.ReadDir(dir)
if err != nil {
return nil, nil
}
// Collect executable workflow scripts (no extension).
workflows := make(map[string][]string)
for _, e := range entries {
if e.IsDir() {
continue
}
name := e.Name()
if strings.Contains(name, ".") {
continue // skip .conf files
}
workflows[name] = []string{"default"}
}
// Find variant confs: {name}-{variant}.conf
for _, e := range entries {
name := e.Name()
if !strings.HasSuffix(name, ".conf") {
continue
}
base := strings.TrimSuffix(name, ".conf")
// Find the longest workflow name that is a prefix.
for wf := range workflows {
if strings.HasPrefix(base, wf+"-") {
variant := strings.TrimPrefix(base, wf+"-")
if variant != "" {
workflows[wf] = append(workflows[wf], variant)
}
}
}
}
var sb strings.Builder
// "none" is always the first entry — it means no workflow runs.
sb.WriteString("none\tdefault\n")
names := make([]string, 0, len(workflows))
for wf := range workflows {
names = append(names, wf)
}
sort.Strings(names)
for _, wf := range names {
sb.WriteString(wf)
sb.WriteByte('\t')
sb.WriteString(strings.Join(workflows[wf], ","))
sb.WriteByte('\n')
}
return []byte(sb.String()), nil
}),
),
virtfs.FileNode("ctl", 0666,
virtfs.Doc("Server control. Write: 'invalidate', 'kill'"),
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
return dispatch([]ctlCmd{
{"invalidate", "clear the models cache", func(_ []string) ([]byte, error) {
if mc != nil {
mc.Invalidate()
}
if cfg.Invalidate != nil {
cfg.Invalidate()
}
return []byte("ok\n"), nil
}},
{"kill", "shut down the server", func(_ []string) ([]byte, error) {
if cfg.Shutdown != nil {
cfg.Shutdown()
}
return []byte("ok\n"), nil
}},
}, data)
}),
),
virtfs.FileNode("event", 0666,
virtfs.Doc("Event stream. Read: all events. Write filter then read: filtered events."),
// Stream mode marker — actual handling is in server.go Tread handler
virtfs.StreamRaw(func(ctx context.Context, base string) ([]byte, string, error) {
// Never called — server.go intercepts event file reads
<-ctx.Done()
return nil, base, nil
}),
),
virtfs.FileNode("event.pub", 0220,
virtfs.Doc("Publish event: write topic<tab>payload"),
virtfs.Write(func(data []byte) error {
line := strings.TrimSpace(string(data))
topic, payload, _ := strings.Cut(line, "\t")
if topic == "" {
return fmt.Errorf("empty topic")
}
session.PublishEvent(topic, payload)
return nil
}),
),
virtfs.FileNode("generate", 0666,
virtfs.Doc("One-shot LLM generation"),
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))
}
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
}),
),
virtfs.FileNode("stats", 0666,
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
result, err := metrics.QueryFormat(string(data))
if err != nil {
return nil, err
}
return []byte(result), nil
}),
),
virtfs.FileNode("metrics", 0444,
virtfs.Read(func() ([]byte, error) {
a, err := metrics.AggregateAll()
if err != nil {
return nil, err
}
return []byte(metrics.Format(a)), nil
}),
),
virtfs.FileNode("metrics.by-backend-model", 0444,
virtfs.Read(func() ([]byte, error) {
groups, err := metrics.GroupedAggregate("", "")
if err != nil {
return nil, err
}
return []byte(metrics.FormatGroups(groups)), nil
}),
),
// session/
virtfs.DirNode("session",
virtfs.Doc("Session management"),
virtfs.FileNode("new", 0666,
virtfs.Doc("Create session"),
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
args := strings.Fields(string(data))
name, remote, workflow, variant, cwd := "", "", "", "", ""
yolo := false
for _, arg := range args {
if k, v, ok := strings.Cut(arg, "="); ok {
switch k {
case "name":
name = v
case "remote":
remote = v
case "workflow":
workflow = v
case "variant":
variant = v
case "cwd":
cwd = v
case "yolo":
yolo = v == "true"
}
}
}
if cwd == "" {
return nil, fmt.Errorf("cwd is required")
}
sess, created, err := session.CreateEmpty(name, remote, yolo)
if err != nil {
return nil, err
}
// Get-or-create: only a freshly created session takes the
// provided configuration. An existing session is returned
// untouched.
if created {
sess.SetCwd(cwd)
if workflow != "" {
sess.SetWorkflow(workflow)
}
if variant != "" {
sess.SetVariant(variant)
}
}
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("goal", 0666,
virtfs.Doc("Session goal text. Writing triggers the workflow if not already running."),
virtfs.GID("agent"),
virtfs.Read(func() ([]byte, error) {
text, _ := s.Goal()
if text == "" {
return nil, nil
}
return []byte(text + "\n"), nil
}),
virtfs.Write(func(data []byte) error {
input := strings.TrimSpace(string(data))
if input == "" {
s.ClearGoal()
return nil
}
_, status := s.Goal()
s.SetGoal(input)
// Trigger workflow only if not already running.
if status == "" || status == "complete" || status == "blocked" || strings.HasPrefix(status, "error") {
runWorkflow(s, s.Workflow(), s.Variant())
}
return nil
}),
),
virtfs.FileNode("goalstatus", 0666,
virtfs.Doc("Goal status. Read/write."),
virtfs.GID("agent"),
virtfs.Read(func() ([]byte, error) {
_, status := s.Goal()
return []byte(status + "\n"), nil
}),
virtfs.Write(func(data []byte) error {
status := strings.TrimSpace(string(data))
s.SetGoalStatus(status)
return nil
}),
),
virtfs.FileNode("goalwait", 0444,
virtfs.Doc("Blocks until goal status changes."),
virtfs.BlockOnce(func() ([]byte, string, error) {
_, status := s.Goal()
return []byte(status + "\n"), status, nil
}, s.GoalSignal),
),
virtfs.FileNode("bypass", 0666,
virtfs.Doc("Bypass request handling. Read: pending request JSON or empty. Write: 'id approve' or 'id deny'."),
virtfs.Read(func() ([]byte, error) {
req := s.BypassPending()
if req == nil {
return nil, nil
}
data, err := json.Marshal(req)
if err != nil {
return nil, err
}
return append(data, '\n'), nil
}),
virtfs.Write(func(data []byte) error {
input := strings.TrimSpace(string(data))
parts := strings.SplitN(input, " ", 2)
if len(parts) != 2 {
return fmt.Errorf("expected 'id approve' or 'id deny'")
}
id, action := parts[0], parts[1]
approved := action == "approve"
return s.ResolveBypass(id, approved)
}),
),
virtfs.FileNode("stats", 0444,
virtfs.Read(func() ([]byte, error) {
a, err := metrics.AggregateSession(s.ID)
if err != nil {
return nil, err
}
return []byte(metrics.Format(a)), nil
}),
),
virtfs.FileNode("metrics.by-backend-model", 0444,
virtfs.Read(func() ([]byte, error) {
groups, err := metrics.GroupedAggregate(s.ID, "")
if err != nil {
return nil, err
}
return []byte(metrics.FormatGroups(groups)), nil
}),
),
virtfs.FileNode("ctl", 0666,
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
return dispatch([]ctlCmd{
{"kill", "destroy this session", func(_ []string) ([]byte, error) { return []byte("ok\n"), removeFn() }},
{"save", "persist session state to disk", func(_ []string) ([]byte, error) {
s.Save()
return []byte("ok\n"), nil
}},
{"invalidate", "clear the models cache", func(_ []string) ([]byte, error) {
s.InvalidateModelsCache()
return []byte("ok\n"), nil
}},
{"pause", "pause all agents in this session", func(_ []string) ([]byte, error) {
if err := s.Pause(); err != nil {
return nil, err
}
return []byte("ok\n"), nil
}},
{"resume", "resume all agents in this session", func(_ []string) ([]byte, error) {
if err := s.Resume(); err != nil {
return nil, err
}
return []byte("ok\n"), nil
}},
{"run", "run a workflow: run [workflow] [variant]", func(args []string) ([]byte, error) {
workflow := s.Workflow()
variant := s.Variant()
if len(args) > 0 {
workflow = args[0]
}
if len(args) > 1 {
variant = args[1]
}
runWorkflow(s, workflow, variant)
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("cfg", 0666,
virtfs.Doc("Session configuration. Read: key=value lines. Write: single key=value to update."),
virtfs.Read(func() ([]byte, error) {
var sb strings.Builder
fmt.Fprintf(&sb, "name=%s\n", s.Name())
fmt.Fprintf(&sb, "cwd=%s\n", s.Cwd())
fmt.Fprintf(&sb, "remote=%s\n", s.Remote)
fmt.Fprintf(&sb, "workflow=%s\n", s.Workflow())
fmt.Fprintf(&sb, "variant=%s\n", s.Variant())
fmt.Fprintf(&sb, "yolo=%t\n", s.Yolo)
return []byte(sb.String()), nil
}),
virtfs.Write(func(data []byte) error {
input := strings.TrimSpace(string(data))
if input == "" {
return nil
}
for _, line := range strings.Split(input, "\n") {
line = strings.TrimSpace(line)
if line == "" {
continue
}
k, v, ok := strings.Cut(line, "=")
if !ok {
return fmt.Errorf("invalid cfg format (expected key=value)")
}
switch k {
case "workflow":
s.SetWorkflow(v)
case "variant":
s.SetVariant(v)
case "cwd":
if v != "" {
s.SetCwd(v)
}
default:
return fmt.Errorf("unknown session cfg key: %s", k)
}
}
return nil
}),
),
virtfs.FileNode("id", 0444,
virtfs.Read(func() ([]byte, error) { return []byte(s.ID + "\n"), nil }),
),
virtfs.DirNode("agent",
virtfs.FileNode("new", 0666,
virtfs.Doc("Create agent. With prompt=, runs as sub-agent: blocks until done, returns reply."),
virtfs.Rdwr(func(ctx context.Context, data []byte) ([]byte, error) {
req := parseAgentNewRequest(data)
req.Params.ParentID = req.ParentID
if req.Prompt == "" {
// Persistent agent creation (no sub-agent mode).
ag, err := session.CreateAgentWithParams(s.Name(), req.Params)
if err != nil {
return nil, fmt.Errorf("create agent: %w", err)
}
wireAgentEvents(s.ID, ag)
return []byte(ag.ID() + "\n"), nil
}
// Sub-agent mode: enforce limits, create, run, destroy.
maxDepth := req.MaxDepth
if maxDepth <= 0 {
maxDepth = 1
}
maxParallel := 5 // default
if req.MaxParallel != nil {
maxParallel = *req.MaxParallel
}
if maxParallel == 0 {
return nil, fmt.Errorf("sub-agent spawning is disabled (max_parallel=0)")
}
timeout := req.Timeout
if timeout <= 0 {
timeout = 600
}
var parent *agent.Agent
if req.ParentID != "" {
parent = s.FindAgent(req.ParentID)
if parent == nil {
return nil, fmt.Errorf("parent agent %q not found", req.ParentID)
}
}
if parent != nil {
// Depth check.
if parent.Depth()+1 > maxDepth {
return nil, fmt.Errorf("sub-agent depth limit exceeded (max %d)", maxDepth)
}
// Parallelism check.
if maxParallel > 0 && parent.ActiveChildren() >= int32(maxParallel) {
return nil, fmt.Errorf("sub-agent parallelism limit exceeded (max %d)", maxParallel)
}
// Inherit context (optionally truncated at fork_at).
msgs := parent.Messages()
if req.ForkAt > 0 {
msgs = messagesUpToTurn(msgs, req.ForkAt)
}
// Sanitize to remove dangling tool calls
msgs = backend.SanitizeMessages(msgs)
req.Params.History = agent.RestoreHistoryFromMessages(msgs)
parent.IncChildren()
defer parent.DecChildren()
}
ag, err := session.CreateAgentWithParams(s.Name(), req.Params)
if err != nil {
return nil, fmt.Errorf("create agent: %w", err)
}
wireAgentEvents(s.ID, ag)
if parent != nil {
ag.SetDepth(parent.Depth() + 1)
}
ag.SetEnv("OLLIE_SUBAGENT_DEPTH", fmt.Sprintf("%d", ag.Depth()))
// Run with timeout.
subCtx, cancel := context.WithTimeout(ctx, time.Duration(timeout)*time.Second)
defer cancel()
session.PublishEvent("session."+s.ID+".agent."+ag.ID()+".new", "")
subPrompt := "You are a sub-agent. Your task is below. " +
"Execute it completely using tools — investigate, implement, verify. " +
"Do NOT respond with a plan or intentions. Do NOT say what you will do. " +
"Call tools. Do the work. Your final text response must summarize what you ACCOMPLISHED.\n\n" + req.Prompt
ag.Submit(subCtx, subPrompt)
ag.EnsureTrailingNewline()
reply := ag.Reply()
s.RemoveAgent(ag.ID())
session.PublishEvent("session."+s.ID+".agent."+ag.ID()+".kill", "")
if reply == "" {
return []byte("(no reply)\n"), nil
}
return []byte(reply), 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()},
Mode: 0755, // owner rwx, group/world rx
UID: a.ID(),
GID: "agent",
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_: 0440, Size_: 64 * 1024}
}
stats := func(_ []string) ([]byte, error) {
u := a.Usage()
var sb strings.Builder
fmt.Fprintf(&sb, "usage=%s\n", a.UsageStr())
fmt.Fprintf(&sb, "cost=%s\n", a.CostStr())
fmt.Fprintf(&sb, "ctxsz=%s\n", a.CtxSz())
if u != nil {
fmt.Fprintf(&sb, "cachedInputTokens=%d\n", u.TotalCachedInputTokens)
fmt.Fprintf(&sb, "cacheCreationTokens=%d\n", u.TotalCacheCreationTokens)
fmt.Fprintf(&sb, "cacheHitRatio=%g\n", u.CacheHitRatio)
}
return []byte(sb.String()), nil
}
inject := func(args []string) ([]byte, error) {
text := strings.Join(args, " ")
if text == "" {
return nil, fmt.Errorf("inject requires text")
}
if a.IsRunning() {
a.InjectRewrite(text)
} else {
startAsync(s.Ctx, func() {
a.Submit(s.Ctx, text)
a.EnsureTrailingNewline()
})
}
return []byte("ok\n"), nil
}
return []virtfs.FsNodeDecl{
virtfs.FileNode("prompt", 0220, // write-only, owner+group
virtfs.Write(func(data []byte) error {
input := strings.TrimSpace(string(data))
if input == "" {
return nil
}
if input == "/invalidate" {
s.InvalidateModelsCache()
return nil
}
startAsync(s.Ctx, func() {
a.Submit(s.Ctx, input)
a.EnsureTrailingNewline()
})
return nil
}),
),
virtfs.FileNode("fifo", 0660, // owner + group (frontend needs access)
virtfs.Write(func(data []byte) error {
input := strings.TrimSpace(string(data))
if input == "" {
return nil
}
if err := a.Queue(input); err != nil {
return err
}
// If the agent is idle, trigger processing.
if !a.IsRunning() {
if !startAsync(s.Ctx, func() {
if next, ok := a.PopQueue(); ok {
a.Submit(s.Ctx, next)
a.EnsureTrailingNewline()
}
}) {
return fmt.Errorf("async work limit reached")
}
}
return nil
}),
virtfs.Read(func() ([]byte, error) {
item, ok := a.PopQueue()
if !ok {
return nil, nil
}
return []byte(item), nil
}),
),
virtfs.FileNode("chat", 0440, // read-only, owner + group
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", 0440, // read-only, owner + group
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("state", 0444, // world-readable (non-sensitive)
virtfs.Doc("Current agent state (idle, calling, thinking, paused)"),
virtfs.Read(func() ([]byte, error) {
return []byte(a.State() + "\n"), nil
}),
),
virtfs.FileNode("status", 0444,
virtfs.Doc("Human-readable status: current activity and elapsed time"),
virtfs.Read(func() ([]byte, error) {
return []byte(a.Status() + "\n"), nil
}),
),
virtfs.FileNode("log", 0440, // read-only, owner + group
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", 0660, // owner + group (frontend needs access)
virtfs.Read(func() ([]byte, error) {
return a.Plan(), nil
}),
virtfs.Write(func(data []byte) error {
a.SetPlan(data)
return nil
}),
),
virtfs.FileNode("cfg", 0640, // owner rw, group r
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, "sessionCwd=%s\n", a.SessionCwd())
fmt.Fprintf(&sb, "cwdOverride=%s\n", a.CwdOverride())
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]))
}
case "cwd", "cwdOverride":
// Per-agent working directory override. An empty value
// clears the override so the agent inherits the session cwd.
a.SetCwdOverride(strings.TrimSpace(parts[1]))
default:
return fmt.Errorf("unknown cfg key: %s", parts[0])
}
return nil
}),
),
virtfs.Each("peer", func() ([]virtfs.FsNodeDecl, error) {
var out []virtfs.FsNodeDecl
for _, peerName := range a.Peers() {
name := peerName
out = append(out, virtfs.FsNodeDecl{
Name: name,
Mode: 0222,
Write: func(data []byte) error {
target := s.FindAgent(name)
if target == nil {
return fmt.Errorf("peer %q not found", name)
}
input := strings.TrimSpace(string(data))
if input == "" {
return nil
}
startAsync(s.Ctx, func() {
target.Submit(s.Ctx, input)
target.EnsureTrailingNewline()
})
return nil
},
})
}
return out, nil
}),
virtfs.FileNode("ctl", 0660, // owner + group (frontend needs access)
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
return dispatch([]ctlCmd{
{"kill", "destroy this agent", func(_ []string) ([]byte, error) {
s.RemoveAgent(a.ID())
session.PublishEvent("session."+s.ID+".agent."+a.ID()+".kill", "")
return []byte("ok\n"), nil
}},
{"stop", "interrupt the current turn", func(_ []string) ([]byte, error) {
a.Interrupt(agent.ErrInterrupted)
return []byte("ok\n"), nil
}},
{"detach", "detach the foreground process to background", func(_ []string) ([]byte, error) {
ts := a.ToolServer()
if ts == nil {
return nil, fmt.Errorf("no tool server")
}
if !ts.Detach() {
return nil, fmt.Errorf("no foreground process to detach")
}
return []byte("ok\n"), nil
}},
{"compact", "compact context now", func(_ []string) ([]byte, error) {
if err := a.Compact(s.Ctx); err != nil {
return nil, err
}
return []byte("ok\n"), nil
}},
{"compactionmodel", "get/set the compaction model", func(args []string) ([]byte, error) {
model := strings.TrimSpace(strings.Join(args, " "))
if model == "" {
return []byte(a.CompactionModel() + "\n"), nil
}
a.SetCompactionModel(model)
return []byte(model + "\n"), nil
}},
{"clear", "clear chat history", func(_ []string) ([]byte, error) {
if err := a.Clear(); err != nil {
return nil, err
}
return []byte("ok\n"), nil
}},
{"i", "inject text into the running turn", inject},
{"agent", "get/switch agent profile: agent [name]", func(args []string) ([]byte, error) {
if len(args) == 0 {
return []byte(a.Profile() + "\n"), nil
}
cfg, err := a.SwitchProfile(args[0])
if err != nil {
return nil, err
}
// Reload tools for the new profile
if ts := a.ToolServer(); ts != nil {
ts.ClearTools()
session.LoadTools(cfg, ts, s.ID, a.ID(), nil)
}
return []byte(args[0] + "\n"), nil
}},
{"model", "get/set model: model [name]", func(args []string) ([]byte, error) {
if len(args) == 0 {
if be := a.Backend(); be != nil {
return []byte(be.Model() + "\n"), nil
}
return nil, nil
}
name := strings.Join(args, " ")
if be := a.Backend(); be != nil {
be.SetModel(name)
return []byte(modelPricingBlock(be, name)), nil
}
return []byte(name + "\n"), nil
}},
{"models", "list available models", func(_ []string) ([]byte, error) {
return []byte(s.CachedListModels() + "\n"), nil
}},
{"tools", "list loaded 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", "load a tool: tool_load <name>", 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", "list all available tools", func(args []string) ([]byte, error) {
ts := a.ToolServer()
if ts == nil {
return nil, fmt.Errorf("no tool server")
}
return ts.ListAllTools()
}},
{"tool_unload", "unload a tool: tool_unload <name>", 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", "get/set backend: backend [name]", func(args []string) ([]byte, error) {
if len(args) == 0 {
if be := a.Backend(); be != nil {
return []byte(be.Name() + "\n"), nil
}
return nil, nil
}
name := strings.Join(args, " ")
if err := a.SwitchBackend(name); err != nil {
return nil, err
}
return []byte(name + "\n"), nil
}},
{"name", "get/set agent name: name [value]", 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", "print, set, or clear the per-agent working directory override: cwd [<dir>|-]", func(args []string) ([]byte, error) {
if len(args) == 0 {
return []byte(a.Cwd() + "\n"), nil
}
dir := strings.TrimSpace(strings.Join(args, " "))
if dir == "-" {
dir = "" // clear override, inherit session cwd
}
a.SetCwdOverride(dir)
return []byte(a.Cwd() + "\n"), nil
}},
{"proc", "manage background procs: proc [top|idx|term|kill|out|dismiss <pid>]", 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 "idx":
data, err := ts.ListProcsIdx()
if err != nil {
return nil, err
}
return data, nil
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, idx, term <pid>, kill <pid>, out <pid>, dismiss <pid>)", subcmd)
}
}},
{"systemprompt", "print the rendered system prompt", func(_ []string) ([]byte, error) {
return []byte(a.SystemPrompt() + "\n"), nil
}},
{"peeradd", "link a peer agent: peeradd <name>", func(args []string) ([]byte, error) {
if len(args) == 0 {
return nil, fmt.Errorf("peeradd requires agent name")
}
target := s.FindAgent(args[0])
if target == nil {
return nil, fmt.Errorf("agent %q not found in session", args[0])
}
a.AddPeer(target.Name())
target.AddPeer(a.Name())
go s.Save()
return []byte("ok\n"), nil
}},
{"peerdel", "unlink a peer agent: peerdel <name>", func(args []string) ([]byte, error) {
if len(args) == 0 {
return nil, fmt.Errorf("peerdel requires agent name")
}
a.RemovePeer(args[0])
if target := s.FindAgent(args[0]); target != nil {
target.RemovePeer(a.Name())
}
go s.Save()
return []byte("ok\n"), nil
}},
{"peers", "list linked peer agents", func(_ []string) ([]byte, error) {
peers := a.Peers()
if len(peers) == 0 {
return []byte("(no peers)\n"), nil
}
return []byte(strings.Join(peers, "\n") + "\n"), nil
}},
{"stats", "print usage and cost statistics", stats},
}, data)
}),
),
virtfs.FileNode("stats", 0440, // owner + group readable
virtfs.Read(func() ([]byte, error) {
a, err := metrics.AggregateAgent(s.ID, a.ID())
if err != nil {
return nil, err
}
return []byte(metrics.Format(a)), nil
}),
),
virtfs.FileNode("metrics.by-backend-model", 0440, // owner + group readable
virtfs.Read(func() ([]byte, error) {
groups, err := metrics.GroupedAggregate(s.ID, a.ID())
if err != nil {
return nil, err
}
return []byte(metrics.FormatGroups(groups)), nil
}),
),
virtfs.FileNode("name", 0640, // owner rw, group r
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, // world-readable (non-sensitive)
virtfs.Read(func() ([]byte, error) { return []byte(a.ID() + "\n"), nil }),
),
}
}