batch: rename status→state, add statewait blocking read
- Rename batchJob.status → .state; update all references
- Add done chan struct{} to batchJob, closed on terminal state
- Add statewait file to b/<id>/: blocks until state leaves "running"
- Add CurrentWaitValue/Wait methods to BatchJobStore
- Extend Topen handler to snapshot waitBase for b/ wait files
- Extend Tread handler for b/ to serve statewait with 5s timeout/retry
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
80651d3b98
commit
b9509a0961
|
|
@ -30,14 +30,15 @@ ollie/
|
|||
|
||||
b/ dir: batch jobs and job scripts
|
||||
new r/w: read: KV template; write: submit a job spec
|
||||
idx read: job index (id, status, cwd, agent — one per line)
|
||||
idx read: job index (id, state, cwd, agent — one per line)
|
||||
job exec: core batch job runner (submit, wait, print result)
|
||||
q exec: foreground one-shot query; thin wrapper around job
|
||||
sched exec: submit background job; prints b/ path; wrapper around job
|
||||
cleanup exec: remove all jobs with status "done"
|
||||
cleanup exec: remove all jobs with state "done"
|
||||
<job-id>/ rm -r to cancel and remove
|
||||
spec read: original job spec as submitted
|
||||
status read: running | done | failed: <reason>
|
||||
state read: running | done | failed: <reason>
|
||||
statewait read: blocks until state changes; returns new state
|
||||
result read: assistant reply (populated when done)
|
||||
usage read: token counts
|
||||
ctxsz read: context size
|
||||
|
|
|
|||
|
|
@ -28,11 +28,12 @@ type batchJob struct {
|
|||
backend string
|
||||
model string
|
||||
output string
|
||||
status string // "running" | "done" | "failed: ..."
|
||||
state string // "running" | "done" | "failed: ..."
|
||||
result string
|
||||
usage string
|
||||
ctxsz string
|
||||
cancel context.CancelFunc
|
||||
done chan struct{} // closed when state reaches a terminal value
|
||||
}
|
||||
|
||||
// batchSpec is the parsed result of a b/new write.
|
||||
|
|
@ -154,14 +155,14 @@ func (s *BatchStore) job(id string) *batchJob {
|
|||
return s.jobs[id]
|
||||
}
|
||||
|
||||
// index returns the live b/idx content: id\tstatus\tcwd\tagent per line.
|
||||
// index returns the live b/idx content: id\tstate\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)
|
||||
fmt.Fprintf(&sb, "%s\t%s\t%s\t%s\n", id, job.state, job.cwd, job.agentName)
|
||||
job.mu.RUnlock()
|
||||
}
|
||||
return []byte(sb.String())
|
||||
|
|
@ -196,7 +197,8 @@ func (s *BatchStore) handleNewBatch(input string) error {
|
|||
backend: spec.backend,
|
||||
model: spec.model,
|
||||
output: spec.output,
|
||||
status: "running",
|
||||
state: "running",
|
||||
done: make(chan struct{}),
|
||||
}
|
||||
s.jobs[id] = job
|
||||
s.mu.Unlock()
|
||||
|
|
@ -213,24 +215,26 @@ func (s *BatchStore) handleNewBatch(input string) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// runJob executes one batch job and updates its status and result.
|
||||
// runJob executes one batch job and updates its state 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"
|
||||
job.state = "failed: cancelled"
|
||||
} else {
|
||||
job.status = "failed: " + err.Error()
|
||||
job.state = "failed: " + err.Error()
|
||||
}
|
||||
job.result = job.status
|
||||
job.result = job.state
|
||||
plog.Info("batch job %s failed: %v", job.id, err)
|
||||
close(job.done)
|
||||
return
|
||||
}
|
||||
job.result = result
|
||||
job.status = "done"
|
||||
job.state = "done"
|
||||
plog.Info("batch job %s done", job.id)
|
||||
close(job.done)
|
||||
}
|
||||
|
||||
// executeJob creates an ephemeral agentCore, submits the prompt, and returns
|
||||
|
|
@ -383,11 +387,12 @@ var batchJobFiles = []struct {
|
|||
name string
|
||||
mode os.FileMode
|
||||
}{
|
||||
{"spec", 0444},
|
||||
{"status", 0444},
|
||||
{"result", 0444},
|
||||
{"usage", 0444},
|
||||
{"ctxsz", 0444},
|
||||
{"spec", 0444},
|
||||
{"state", 0444},
|
||||
{"statewait", 0444},
|
||||
{"result", 0444},
|
||||
{"usage", 0444},
|
||||
{"ctxsz", 0444},
|
||||
}
|
||||
|
||||
// BatchJobStore provides Stat/List/Get for the files within a single batch job
|
||||
|
|
@ -436,8 +441,8 @@ func (js *BatchJobStore) content(name string) string {
|
|||
switch name {
|
||||
case "spec":
|
||||
return js.job.spec
|
||||
case "status":
|
||||
return js.job.status + "\n"
|
||||
case "state":
|
||||
return js.job.state + "\n"
|
||||
case "result":
|
||||
return js.job.result
|
||||
case "usage":
|
||||
|
|
@ -447,3 +452,36 @@ func (js *BatchJobStore) content(name string) string {
|
|||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// CurrentWaitValue returns the current value for the named *wait file.
|
||||
func (js *BatchJobStore) CurrentWaitValue(name string) string {
|
||||
if name == "statewait" {
|
||||
js.job.mu.RLock()
|
||||
defer js.job.mu.RUnlock()
|
||||
return js.job.state
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// Wait blocks until the job leaves "running" state, then returns the new state.
|
||||
// Returns nil content (empty read) on context cancellation.
|
||||
func (js *BatchJobStore) Wait(ctx context.Context, name, base string) ([]byte, error) {
|
||||
if name != "statewait" {
|
||||
return nil, fmt.Errorf("%s: not a wait file", name)
|
||||
}
|
||||
js.job.mu.RLock()
|
||||
current := js.job.state
|
||||
js.job.mu.RUnlock()
|
||||
if current != base {
|
||||
return []byte(current + "\n"), nil
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return nil, nil
|
||||
case <-js.job.done:
|
||||
}
|
||||
js.job.mu.RLock()
|
||||
current = js.job.state
|
||||
js.job.mu.RUnlock()
|
||||
return []byte(current + "\n"), nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -539,8 +539,15 @@ func (s *Server) open(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
if strings.HasSuffix(pathBase(f.path), "wait") {
|
||||
parts := strings.SplitN(strings.TrimPrefix(f.path, "/"), "/", 3)
|
||||
if len(parts) == 3 {
|
||||
if store, ok := s.sessionFileStore(parts[1]); ok {
|
||||
f.waitBase = store.CurrentWaitValue(parts[2])
|
||||
switch parts[0] {
|
||||
case "s":
|
||||
if store, ok := s.sessionFileStore(parts[1]); ok {
|
||||
f.waitBase = store.CurrentWaitValue(parts[2])
|
||||
}
|
||||
case "b":
|
||||
if js, ok := s.batchStore.JobStore(parts[1]); ok {
|
||||
f.waitBase = js.CurrentWaitValue(parts[2])
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -636,6 +643,32 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall, ctx context.Context) *plan
|
|||
if !ok {
|
||||
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
|
||||
}
|
||||
if parts[2] == "statewait" {
|
||||
if fc.Offset > 0 {
|
||||
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
|
||||
}
|
||||
cs.mu.RLock()
|
||||
f, fidOK := cs.fids[fc.Fid]
|
||||
var base string
|
||||
if fidOK {
|
||||
base = f.waitBase
|
||||
}
|
||||
cs.mu.RUnlock()
|
||||
waitCtx, waitCancel := context.WithTimeout(ctx, 5*time.Second)
|
||||
defer waitCancel()
|
||||
content, err := js.Wait(waitCtx, parts[2], base)
|
||||
if err != nil {
|
||||
return errFcall(fc, err.Error())
|
||||
}
|
||||
if content != nil {
|
||||
cs.mu.Lock()
|
||||
if f, ok := cs.fids[fc.Fid]; ok {
|
||||
f.waitBase = strings.TrimSuffix(string(content), "\n")
|
||||
}
|
||||
cs.mu.Unlock()
|
||||
}
|
||||
return s.readSlice(fc, content)
|
||||
}
|
||||
content, err := js.Get(parts[2])
|
||||
if err != nil {
|
||||
return errFcall(fc, err.Error())
|
||||
|
|
|
|||
Reference in New Issue