1281 lines
38 KiB
Go
1281 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/format"
|
|
"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 requests JSON array. Write: 'id approve' or 'id deny'."),
|
|
virtfs.Read(func() ([]byte, error) {
|
|
reqs := s.BypassPending()
|
|
if len(reqs) == 0 {
|
|
return []byte("[]\n"), nil
|
|
}
|
|
data, err := json.Marshal(reqs)
|
|
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)
|
|
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)
|
|
})
|
|
}
|
|
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)
|
|
})
|
|
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)
|
|
}
|
|
}) {
|
|
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
|
|
}),
|
|
),
|
|
// log.raw — JSONL (one Block per line), for GUI/programmatic access
|
|
virtfs.FileNode("log.raw", 0440,
|
|
virtfs.Doc("JSONL chat log. One JSON block per line. Streaming."),
|
|
virtfs.StatOverride(chatStat),
|
|
virtfs.Read(func() ([]byte, error) {
|
|
a.ChatMu().RLock()
|
|
defer a.ChatMu().RUnlock()
|
|
raw := a.RawLog()
|
|
data := make([]byte, len(raw))
|
|
copy(data, raw)
|
|
return data, nil
|
|
}),
|
|
virtfs.Stream(a.RawLogRead, a.ChatSignal),
|
|
),
|
|
// log — Rendered plain text for TUI/humans (non-blocking read)
|
|
virtfs.FileNode("log", 0440,
|
|
virtfs.Doc("Rendered chat log (plain text). Last 64KB."),
|
|
virtfs.Read(func() ([]byte, error) {
|
|
const maxWindow = 64 * 1024
|
|
a.ChatMu().RLock()
|
|
text := a.TextLog()
|
|
start := 0
|
|
if len(text) > maxWindow {
|
|
start = len(text) - maxWindow
|
|
}
|
|
data := make([]byte, len(text)-start)
|
|
copy(data, text[start:])
|
|
a.ChatMu().RUnlock()
|
|
return data, nil
|
|
}),
|
|
),
|
|
// chat — Rendered plain text stream for TUI (blocking)
|
|
virtfs.FileNode("chat", 0440,
|
|
virtfs.Doc("Rendered chat log stream (plain text). Blocks for new data."),
|
|
virtfs.Stream(a.TextLogRead, a.ChatSignal),
|
|
),
|
|
// block — Read block by ID (rdwr: write ID, read JSON block)
|
|
virtfs.FileNode("block", 0660,
|
|
virtfs.Doc("Get block by ID. Write block ID, read JSON block."),
|
|
virtfs.Rdwr(func(_ context.Context, data []byte) ([]byte, error) {
|
|
blockID := strings.TrimSpace(string(data))
|
|
block, found := a.BlockByID(blockID)
|
|
if !found {
|
|
return nil, fmt.Errorf("block not found: %s", blockID)
|
|
}
|
|
return format.MarshalBlock(block), nil
|
|
}),
|
|
),
|
|
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("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)
|
|
})
|
|
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 }),
|
|
),
|
|
}
|
|
}
|