9p: add b/ batch namespace for ephemeral one-shot jobs

Introduces b/ as a first-class 9P namespace alongside s/. Writing to
b/new creates one or more ephemeral agent jobs (parallel=N suffix -0..-N)
each running the same prompt in an independent agentCore with no
persistent session state.

Per-job files: spec, status, result, usage, ctxsz.
BatchStore/BatchJobStore mirror SessionStore/SessionFileStore.
rm -r b/{id} cancels and removes the job.

Also removes s/{id}/reply: b/{id}/result is the purpose-built
replacement for one-shot prompt->reply use cases.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
Levi Neely 2026-04-16 19:22:39 +02:00
parent cb00afb8bd
commit 793cf6b4fa
3 changed files with 540 additions and 21 deletions

442
internal/p9/batchstore.go Normal file
View File

@ -0,0 +1,442 @@
package p9
import (
"context"
"fmt"
"os"
"strconv"
"strings"
"sync"
"ollie/pkg/agent"
"ollie/pkg/backend"
"ollie/pkg/config"
"ollie/pkg/tools"
"ollie/pkg/tools/execute"
)
const batchNewTemplate = "name=\ncwd=\nagent=\nbackend=\nmodel=\noutput=\nparallel=1\n---\nWrite your prompt here.\n"
// batchJob holds state for one ephemeral batch agent run.
type batchJob struct {
mu sync.RWMutex
id string
spec string // verbatim input to b/new
prompt string
cwd string
agentName string
backend string
model string
output string
status string // "running" | "done" | "failed: ..."
result string
usage string
ctxsz string
cancel context.CancelFunc
}
// batchSpec is the parsed result of a b/new write.
type batchSpec struct {
name string
cwd string
agentName string
backend string
model string
output string
parallel int
prompt string
}
// BatchStore implements Store for the /b/ directory.
// Entries are job IDs (directories) plus the synthetic files "new" and "idx".
// Put("new", data) creates jobs; Delete(id) cancels and removes one.
type BatchStore struct {
srv *Server
mu sync.RWMutex
jobs map[string]*batchJob
}
func NewBatchStore(srv *Server) *BatchStore {
return &BatchStore{srv: srv, jobs: make(map[string]*batchJob)}
}
func (s *BatchStore) List() ([]os.DirEntry, error) {
entries := []os.DirEntry{
syntheticEntry("new", 0666),
syntheticEntry("idx", 0444),
}
s.mu.RLock()
for id := range s.jobs {
entries = append(entries, syntheticDirEntry(id, 0555))
}
s.mu.RUnlock()
return entries, nil
}
func (s *BatchStore) Stat(name string) (os.FileInfo, error) {
switch name {
case "new":
return &syntheticFileInfo{name: "new", mode: 0666, size: int64(len(batchNewTemplate))}, nil
case "idx":
return &syntheticFileInfo{name: "idx", mode: 0444, size: int64(len(s.index()))}, nil
}
s.mu.RLock()
_, ok := s.jobs[name]
s.mu.RUnlock()
if ok {
return &syntheticFileInfo{name: name, mode: 0555, isDir: true}, nil
}
return nil, fmt.Errorf("%s: not found", name)
}
func (s *BatchStore) Get(name string) ([]byte, error) {
switch name {
case "new":
return []byte(batchNewTemplate), nil
case "idx":
return s.index(), nil
}
return nil, fmt.Errorf("%s: not a readable file", name)
}
func (s *BatchStore) Put(name string, data []byte) error {
if name != "new" {
return fmt.Errorf("%s: not writable", name)
}
return s.handleNewBatch(strings.TrimSpace(string(data)))
}
func (s *BatchStore) Delete(name string) error {
s.mu.Lock()
job, ok := s.jobs[name]
if ok {
delete(s.jobs, name)
}
s.mu.Unlock()
if !ok {
return fmt.Errorf("batch job not found: %s", name)
}
job.mu.RLock()
cancel := job.cancel
job.mu.RUnlock()
if cancel != nil {
cancel()
}
plog.Info("removed batch job %s", name)
return nil
}
func (s *BatchStore) Create(name string) error {
return fmt.Errorf("create not supported for batch jobs")
}
func (s *BatchStore) Rename(oldName, newName string) error {
return fmt.Errorf("rename not supported for batch jobs")
}
// Shutdown cancels all running jobs.
func (s *BatchStore) Shutdown() {
s.mu.Lock()
ids := make([]string, 0, len(s.jobs))
for id := range s.jobs {
ids = append(ids, id)
}
s.mu.Unlock()
for _, id := range ids {
s.Delete(id) //nolint:errcheck
}
}
// job looks up a job by ID (nil if not found).
func (s *BatchStore) job(id string) *batchJob {
s.mu.RLock()
defer s.mu.RUnlock()
return s.jobs[id]
}
// index returns the live b/idx content: id\tstatus\tcwd\tagent per line.
func (s *BatchStore) index() []byte {
var sb strings.Builder
s.mu.RLock()
defer s.mu.RUnlock()
for id, job := range s.jobs {
job.mu.RLock()
fmt.Fprintf(&sb, "%s\t%s\t%s\t%s\n", id, job.status, job.cwd, job.agentName)
job.mu.RUnlock()
}
return []byte(sb.String())
}
// handleNewBatch parses the spec, creates parallel batchJobs, and starts them.
func (s *BatchStore) handleNewBatch(input string) error {
spec, err := parseBatchSpec(input)
if err != nil {
return err
}
baseID := spec.name
if baseID == "" {
baseID = agent.NewSessionID()
}
for i := 0; i < spec.parallel; i++ {
id := fmt.Sprintf("%s-%d", baseID, i)
s.mu.Lock()
if _, exists := s.jobs[id]; exists {
s.mu.Unlock()
return fmt.Errorf("batch job already exists: %s", id)
}
job := &batchJob{
id: id,
spec: input,
prompt: spec.prompt,
cwd: spec.cwd,
agentName: spec.agentName,
backend: spec.backend,
model: spec.model,
output: spec.output,
status: "running",
}
s.jobs[id] = job
s.mu.Unlock()
ctx, cancel := context.WithCancel(context.Background())
job.mu.Lock()
job.cancel = cancel
job.mu.Unlock()
go s.runJob(ctx, job)
}
plog.Info("new batch base=%s parallel=%d", baseID, spec.parallel)
return nil
}
// runJob executes one batch job and updates its status and result.
func (s *BatchStore) runJob(ctx context.Context, job *batchJob) {
result, err := s.executeJob(ctx, job)
job.mu.Lock()
defer job.mu.Unlock()
if err != nil {
if ctx.Err() != nil {
job.status = "failed: cancelled"
} else {
job.status = "failed: " + err.Error()
}
plog.Info("batch job %s failed: %v", job.id, err)
return
}
job.result = result
job.status = "done"
plog.Info("batch job %s done", job.id)
}
// executeJob creates an ephemeral agentCore, submits the prompt, and returns
// the assistant reply. The core is closed before returning.
func (s *BatchStore) executeJob(ctx context.Context, job *batchJob) (string, error) {
job.mu.RLock()
backendName := job.backend
modelName := job.model
agentName := job.agentName
cwd := job.cwd
prompt := job.prompt
output := job.output
jobID := job.id
job.mu.RUnlock()
var (
be backend.Backend
err error
)
if backendName != "" {
old := os.Getenv("OLLIE_BACKEND")
os.Setenv("OLLIE_BACKEND", backendName) //nolint:errcheck
be, err = backend.New()
os.Setenv("OLLIE_BACKEND", old) //nolint:errcheck
} else {
be, err = backend.New()
}
if err != nil {
return "", fmt.Errorf("backend: %w", err)
}
if modelName != "" {
be.SetModel(modelName)
}
cfgPath := agent.AgentConfigPath(s.srv.agentsDir, agentName)
cfg, _ := config.Load(cfgPath)
newDisp := tools.NewDispatcherFunc(map[string]func() tools.Server{
"execute": execute.Decl(cwd),
})
env := agent.BuildAgentEnv(cfg, newDisp(), cwd)
core := agent.NewAgentCore(agent.AgentCoreConfig{
Backend: be,
AgentName: agentName,
AgentsDir: s.srv.agentsDir,
SessionsDir: "", // ephemeral: no persistence
SessionID: jobID,
CWD: cwd,
Env: env,
NewDispatcher: newDisp,
})
defer core.Close()
if output == "json" {
prompt += "\n\nRespond with valid JSON only, no prose."
}
var replyBuf strings.Builder
core.Submit(ctx, prompt, func(ev agent.Event) {
if ev.Role == "assistant" {
replyBuf.WriteString(ev.Content)
}
})
usage := core.Usage()
ctxsz := core.CtxSz()
job.mu.Lock()
job.usage = usage
job.ctxsz = ctxsz
job.mu.Unlock()
if ctx.Err() != nil {
return "", ctx.Err()
}
return replyBuf.String(), nil
}
// parseBatchSpec splits input on the first \n---\n, parses key=value headers
// above it, and treats everything below as the prompt body.
func parseBatchSpec(input string) (batchSpec, error) {
spec := batchSpec{
agentName: "default",
parallel: 1,
}
const sep = "\n---\n"
idx := strings.Index(input, sep)
if idx < 0 {
return spec, fmt.Errorf("missing --- delimiter between headers and prompt")
}
for _, line := range strings.Split(input[:idx], "\n") {
line = strings.TrimSpace(line)
if line == "" {
continue
}
k, v, ok := strings.Cut(line, "=")
if !ok {
continue
}
k = strings.TrimSpace(k)
v = strings.TrimSpace(v)
if v == "" {
continue
}
switch k {
case "name":
spec.name = v
case "cwd":
spec.cwd = v
case "agent":
spec.agentName = v
case "backend":
spec.backend = v
case "model":
spec.model = v
case "output":
spec.output = v
case "parallel":
if n, err := strconv.Atoi(v); err == nil && n > 0 {
spec.parallel = n
}
}
}
spec.prompt = strings.TrimSpace(input[idx+len(sep):])
if spec.cwd == "" {
return spec, fmt.Errorf("cwd is required")
}
if spec.prompt == "" {
return spec, fmt.Errorf("prompt is required")
}
return spec, nil
}
// ---------------------------------------------------------------------------
// BatchJobStore — per-job file access, mirrors SessionFileStore.
// ---------------------------------------------------------------------------
var batchJobFiles = []struct {
name string
mode os.FileMode
}{
{"spec", 0444},
{"status", 0444},
{"result", 0444},
{"usage", 0444},
{"ctxsz", 0444},
}
// BatchJobStore provides Stat/List/Get for the files within a single batch job
// directory (/b/{id}/*). The file set is fixed.
type BatchJobStore struct {
job *batchJob
}
func (s *BatchStore) JobStore(id string) (*BatchJobStore, bool) {
job := s.job(id)
if job == nil {
return nil, false
}
return &BatchJobStore{job: job}, true
}
func (js *BatchJobStore) List() ([]os.DirEntry, error) {
entries := make([]os.DirEntry, len(batchJobFiles))
for i, f := range batchJobFiles {
entries[i] = syntheticEntry(f.name, f.mode)
}
return entries, nil
}
func (js *BatchJobStore) Stat(name string) (os.FileInfo, error) {
for _, f := range batchJobFiles {
if f.name == name {
return &syntheticFileInfo{name: name, mode: f.mode, size: int64(len(js.content(name)))}, nil
}
}
return nil, fmt.Errorf("%s: not found", name)
}
func (js *BatchJobStore) Get(name string) ([]byte, error) {
for _, f := range batchJobFiles {
if f.name == name {
return []byte(js.content(name)), nil
}
}
return nil, fmt.Errorf("%s: not found", name)
}
func (js *BatchJobStore) content(name string) string {
js.job.mu.RLock()
defer js.job.mu.RUnlock()
switch name {
case "spec":
return js.job.spec
case "status":
return js.job.status + "\n"
case "result":
return js.job.result
case "usage":
return js.job.usage + "\n"
case "ctxsz":
return js.job.ctxsz + "\n"
}
return ""
}

View File

@ -18,7 +18,6 @@
// enqueue (write) queue a prompt for later execution
// dequeue (read) pop the next queued prompt
// chat (read) cumulative chat history
// reply (read) assistant text from the most recent turn
// state (read) current agent state
// backend (r/w) active backend name
// agent (r/w) active agent name
@ -86,8 +85,6 @@ type session struct {
// directly from this buffer by offset, so tail -f works via polling.
chatLog []byte
chatVers uint32 // incremented on each append; used as Qid.Vers
// replyVers is incremented on each turn to drive Qid.Vers for polling.
replyVers uint32
// mutableVers tracks changes to tailable mutable-state files
// (state, backend, agent, model, usage, cwd). Bumped on change;
// reported via Qid.Vers in stat so tail -f detects truncation/rewrite.
@ -145,6 +142,7 @@ type Server struct {
toolStore Store
skillStore Store
sessionStore Store
batchStore *BatchStore
}
// New creates a new Server.
@ -153,8 +151,8 @@ func New() *Server {
os.MkdirAll(memDir, 0755) //nolint:errcheck
agentsDir := agent.DefaultAgentsDir()
s := &Server{
sessions: make(map[string]*session),
agentsDir: agentsDir,
sessions: make(map[string]*session),
agentsDir: agentsDir,
sessionsDir: agent.DefaultSessionsDir(),
agentStore: NewFlatDirStore(agentsDir, 0644),
promptStore: NewFlatDirStore(agent.DefaultPromptsDir(), 0444),
@ -164,6 +162,7 @@ func New() *Server {
skillStore: NewSkillStore(),
}
s.sessionStore = &SessionStore{srv: s}
s.batchStore = NewBatchStore(s)
return s
}
@ -314,6 +313,23 @@ func (s *Server) pathType(path string) string {
return "dir"
case len(parts) == 1 && parts[0] == "t":
return "dir"
case len(parts) == 1 && parts[0] == "b":
return "dir"
case len(parts) == 2 && parts[0] == "b":
switch parts[1] {
case "new", "idx":
return "file"
default:
if _, err := s.batchStore.Stat(parts[1]); err == nil {
return "dir"
}
}
case len(parts) == 3 && parts[0] == "b":
if js, ok := s.batchStore.JobStore(parts[1]); ok {
if _, err := js.Stat(parts[2]); err == nil {
return "file"
}
}
case len(parts) == 1 && parts[0] == "backends":
return "file"
case len(parts) == 1 && parts[0] == "help":
@ -518,6 +534,31 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: uint32(len(data)), Data: data}
}
// Batch job files: /b/new, /b/idx, /b/{id}/{file}
if path == "/b/new" || path == "/b/idx" {
plog.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count)
content, err := s.batchStore.Get(pathBase(path))
if err != nil {
return errFcall(fc, err.Error())
}
return s.readSlice(fc, content)
}
if strings.HasPrefix(path, "/b/") {
parts := strings.SplitN(strings.TrimPrefix(path, "/"), "/", 3)
if len(parts) == 3 && parts[0] == "b" {
plog.Debug("Tread batch file path=%q offset=%d count=%d", path, fc.Offset, fc.Count)
js, ok := s.batchStore.JobStore(parts[1])
if !ok {
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
}
content, err := js.Get(parts[2])
if err != nil {
return errFcall(fc, err.Error())
}
return s.readSlice(fc, content)
}
}
// s/new and s/idx are served from the session store.
if path == "/s/new" || path == "/s/idx" {
plog.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count)
@ -894,6 +935,15 @@ func (s *Server) remove(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
err = s.skillStore.Delete(pathBase(path))
case strings.HasPrefix(path, "/t/"):
err = s.toolStore.Delete(pathBase(path))
case strings.HasPrefix(path, "/b/") && path != "/b/new":
// Synthetic batch files are no-ops so "rm -r b/{id}" can proceed to
// remove the directory itself, which triggers job cancellation.
parts := strings.SplitN(strings.TrimPrefix(path, "/b/"), "/", 2)
if len(parts) == 2 {
err = nil // synthetic file; let rm -r continue
} else {
err = s.batchStore.Delete(parts[0])
}
case strings.HasPrefix(path, "/s/") && path != "/s/new":
// If path is /s/{id}/{file}, it's a synthetic session file — no-op so
// that "rm -r s/{id}" can proceed to remove the directory itself, which
@ -922,6 +972,10 @@ func (s *Server) handleWrite(path, input string) error {
return nil
}
if path == "/b/new" {
return s.batchStore.Put("new", []byte(input))
}
if path == "/s/new" {
return s.sessionStore.Put("new", []byte(input))
}
@ -1100,7 +1154,7 @@ func (s *Server) createSession(args []string) error {
return nil
}
// Shutdown kills all active sessions, triggering Close() on each core.
// Shutdown kills all active sessions and batch jobs, triggering Close() on each core.
func (s *Server) Shutdown() {
s.mu.Lock()
ids := make([]string, 0, len(s.sessions))
@ -1111,6 +1165,8 @@ func (s *Server) Shutdown() {
for _, id := range ids {
s.killSession(id)
}
s.batchStore.Shutdown()
}
func (s *Server) killSession(id string) {
@ -1177,6 +1233,7 @@ func (s *Server) readDir(path string, offset uint64, count uint32) []byte {
if path == "/" {
dirs = append(dirs, makeDir("a", "/a", true, plan9.DMDIR|0755))
dirs = append(dirs, makeDir("b", "/b", true, plan9.DMDIR|0755))
dirs = append(dirs, makeDir("backends", "/backends", false, 0444))
dirs = append(dirs, makeDir("help", "/help", false, 0444))
dirs = append(dirs, makeDir("m", "/m", true, plan9.DMDIR|0755))
@ -1241,6 +1298,26 @@ func (s *Server) readDir(path string, offset uint64, count uint32) []byte {
}
dirs = append(dirs, makeDir(e.Name(), "/t/"+e.Name(), false, mode))
}
} else if path == "/b" {
entries, _ := s.batchStore.List()
for _, e := range entries {
isDir := e.IsDir()
info, _ := e.Info()
perm := plan9.Perm(info.Mode() & 0777)
if isDir {
perm = plan9.DMDIR | perm
}
dirs = append(dirs, makeDir(e.Name(), "/b/"+e.Name(), isDir, perm))
}
} else if strings.HasPrefix(path, "/b/") {
// /b/{id} directory listing
if js, ok := s.batchStore.JobStore(pathBase(path)); ok {
entries, _ := js.List()
for _, e := range entries {
info, _ := e.Info()
dirs = append(dirs, makeDir(e.Name(), path+"/"+e.Name(), false, plan9.Perm(info.Mode())))
}
}
} else if path == "/s" {
entries, _ := s.sessionStore.List()
for _, e := range entries {
@ -1321,7 +1398,7 @@ func (s *Server) makeStat(path string) plan9.Dir {
switch base {
case "ctl", "prompt", "enqueue":
mode = 0200
case "chat", "reply", "state", "usage", "ctxsz", "models", "mcp", "dequeue":
case "chat", "state", "usage", "ctxsz", "models", "mcp", "dequeue":
mode = 0444
case "backend", "agent", "model", "cwd":
mode = 0666
@ -1361,7 +1438,7 @@ func (s *Server) makeStat(path string) plan9.Dir {
Muid: "ollie",
}
// For chat/reply and tailable mutable files, report actual size and
// For chat and tailable mutable files, report actual size and
// Qid version so polling tools (tail -f) can detect changes via stat.
// Path format: /s/{sessid}/{file}
if strings.HasPrefix(path, "/s/") {
@ -1377,11 +1454,6 @@ func (s *Server) makeStat(path string) plan9.Dir {
dir.Length = uint64(len(sess.chatLog))
dir.Qid.Vers = sess.chatVers
sess.mu.RUnlock()
case "reply":
sess.mu.RLock()
dir.Length = uint64(len(sess.core.Reply()))
dir.Qid.Vers = sess.replyVers
sess.mu.RUnlock()
default:
sess.mu.RLock()
mf := sess.mutableVers[base]
@ -1448,6 +1520,19 @@ func (s *Server) makeStat(path string) plan9.Dir {
if content, err := s.toolStore.Get(base); err == nil {
dir.Length = uint64(len(content))
}
case path == "/b/new" || path == "/b/idx":
if content, err := s.batchStore.Get(base); err == nil {
dir.Length = uint64(len(content))
}
case strings.HasPrefix(path, "/b/"):
parts := strings.SplitN(strings.TrimPrefix(path, "/b/"), "/", 2)
if len(parts) == 2 {
if js, ok := s.batchStore.JobStore(parts[0]); ok {
if info, err := js.Stat(parts[1]); err == nil {
dir.Length = uint64(info.Size())
}
}
}
}
}

View File

@ -19,7 +19,6 @@ var sessionFileList = []struct {
{"enqueue", 0200},
{"dequeue", 0444},
{"chat", 0444},
{"reply", 0444},
{"state", 0444},
{"backend", 0666},
{"agent", 0666},
@ -62,8 +61,6 @@ func (s *SessionFileStore) Stat(name string) (os.FileInfo, error) {
s.sess.mu.RLock()
size = int64(len(s.sess.chatLog))
s.sess.mu.RUnlock()
case "reply":
size = int64(len(s.sess.core.Reply()))
default:
size = int64(len(s.content(name)))
}
@ -81,8 +78,6 @@ func (s *SessionFileStore) Get(name string) ([]byte, error) {
copy(data, s.sess.chatLog)
s.sess.mu.RUnlock()
return data, nil
case "reply":
return []byte(s.sess.core.Reply()), nil
case "dequeue":
item, ok := s.sess.core.PopQueue()
if !ok {
@ -108,9 +103,6 @@ func (s *SessionFileStore) Put(name string, data []byte) error {
case "prompt":
s.sess.core.Submit(s.sess.ctx, input, s.makePublish())
s.sess.trackMutable()
s.sess.mu.Lock()
s.sess.replyVers++
s.sess.mu.Unlock()
case "enqueue":
s.sess.core.Queue(input)