fix: background process lifecycle, streaming output, and connection deadlocks
- Timeout: timeout=0 means no deadline (was defaulting to 30s)
- Signal: send to process group (-pgid) not just process; SIGTERM no longer
cancels context (only SIGKILL does); cmd.Cancel sends SIGTERM with 5s WaitDelay
- Streaming: background procs stream output in real-time via procWriter;
shell tool no longer buffers all output into a bash variable
- Proc tree: olliesrv exposes proc/{id}/out, proc/{id}/ctl, proc/{id}/status
as proper 9P directory (was broken flat file)
- Connection: proc handlers dial fresh toolsrv conn per request via
Session.DialToolServer() to avoid deadlocking the agent's blocked conn
- Stat format: key=value (exited=true, exit_code=N, id=N) matching client parser
- GC: procs auto-removed 10min after LastRead (exited procs only)
- Rename: PID -> ID throughout (synthetic, not OS PID)
- Ctl commands: term (SIGTERM), kill (SIGKILL), signal <n>, dismiss
- System prompt: correct ollie-9p commands for proc management
This commit is contained in:
parent
12babf2f92
commit
0d1eeaa46b
|
|
@ -94,7 +94,7 @@ func (bt *bgTracker) CollectInterrupts(conn *toolsrv.Conn) string {
|
|||
// Tail last 20 lines
|
||||
lines := tailLines(output, 20)
|
||||
|
||||
sb.WriteString("\n\n<system-proc-interrupt pid=\"")
|
||||
sb.WriteString("\n\n<system-proc-interrupt id=\"")
|
||||
sb.WriteString(fmt.Sprintf("%d", p.PID))
|
||||
sb.WriteString("\" cmd=\"")
|
||||
sb.WriteString(escapeAttr(p.Cmd))
|
||||
|
|
|
|||
|
|
@ -578,7 +578,7 @@ func execOne(rt *Runtime, ctx TurnCtx, call backend.ToolCall, resultCache *sync.
|
|||
if ctx.BgTracker != nil {
|
||||
ctx.BgTracker.Add(pid, call.Name, cmdDesc)
|
||||
}
|
||||
result = fmt.Sprintf("<system-proc-background>\npid=%d cmd=%q\n</system-proc-background>", pid, cmdDesc)
|
||||
result = fmt.Sprintf("<system-proc-background>\nid=%d cmd=%q\n</system-proc-background>", pid, cmdDesc)
|
||||
}
|
||||
}
|
||||
emit(ctx, Event{Role: "tool", Name: call.Name, Content: result, OutputFormat: toolOutputFormat(rt, call.Name)})
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ import (
|
|||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
"syscall"
|
||||
|
||||
"ollie/cmd/olliesrv/internal/agent"
|
||||
"ollie/paths"
|
||||
|
|
@ -495,19 +496,28 @@ func streamChat(al *AgentLog, cctx context.Context, base string) ([]byte, string
|
|||
|
||||
// processBindings returns bindings for detached processes.
|
||||
func processBindings(ctx HandlerCtx) ([]Binding, error) {
|
||||
if ctx.Agent == nil {
|
||||
if ctx.Session == nil {
|
||||
return nil, nil
|
||||
}
|
||||
procs := ctx.Agent.ListDetached()
|
||||
out := make([]Binding, len(procs))
|
||||
for i, p := range procs {
|
||||
pid := p.PID // capture
|
||||
out[i] = Binding{
|
||||
Name: strconv.Itoa(p.PID),
|
||||
Applier: func(c HandlerCtx) HandlerCtx {
|
||||
c.Data = pid
|
||||
return c
|
||||
},
|
||||
conn := ctx.Session.DialToolServer()
|
||||
if conn == nil {
|
||||
return nil, nil
|
||||
}
|
||||
defer conn.Close()
|
||||
raw := conn.ListDetachedRaw()
|
||||
out := make([]Binding, 0, len(raw))
|
||||
for _, r := range raw {
|
||||
if m, ok := r.(map[string]any); ok {
|
||||
if pid, ok := m["pid"].(int); ok {
|
||||
id := pid // capture
|
||||
out = append(out, Binding{
|
||||
Name: strconv.Itoa(id),
|
||||
Applier: func(c HandlerCtx) HandlerCtx {
|
||||
c.Data = id
|
||||
return c
|
||||
},
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
|
|
@ -515,17 +525,72 @@ func processBindings(ctx HandlerCtx) ([]Binding, error) {
|
|||
|
||||
func readProcess(ctx HandlerCtx) ([]byte, error) {
|
||||
pid := ctx.Data.(int)
|
||||
output, err := ctx.Agent.GetDetachedOutput(pid)
|
||||
conn := ctx.Session.DialToolServer()
|
||||
if conn == nil {
|
||||
return nil, fmt.Errorf("no tool server")
|
||||
}
|
||||
defer conn.Close()
|
||||
output, err := conn.GetDetachedOutput(pid)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return []byte(output), nil
|
||||
}
|
||||
|
||||
func removeProcess(ctx HandlerCtx) error {
|
||||
func readProcessStatus(ctx HandlerCtx) ([]byte, error) {
|
||||
pid := ctx.Data.(int)
|
||||
if !ctx.Agent.DismissDetached(pid) {
|
||||
return fmt.Errorf("process %d not found or still running", pid)
|
||||
conn := ctx.Session.DialToolServer()
|
||||
if conn == nil {
|
||||
return nil, fmt.Errorf("no tool server")
|
||||
}
|
||||
defer conn.Close()
|
||||
raw := conn.ListDetachedRaw()
|
||||
for _, r := range raw {
|
||||
if m, ok := r.(map[string]any); ok {
|
||||
if id, ok := m["pid"].(int); ok && id == pid {
|
||||
if exited, ok := m["exited"].(bool); ok && exited {
|
||||
exitCode := 0
|
||||
if code, ok := m["exit_code"].(int); ok {
|
||||
exitCode = code
|
||||
}
|
||||
return []byte(fmt.Sprintf("exited %d\n", exitCode)), nil
|
||||
}
|
||||
return []byte("running\n"), nil
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil, fmt.Errorf("process %d not found", pid)
|
||||
}
|
||||
|
||||
func writeProcessCtl(ctx HandlerCtx, data []byte) error {
|
||||
pid := ctx.Data.(int)
|
||||
conn := ctx.Session.DialToolServer()
|
||||
if conn == nil {
|
||||
return fmt.Errorf("no tool server")
|
||||
}
|
||||
defer conn.Close()
|
||||
cmd := strings.TrimSpace(string(data))
|
||||
switch {
|
||||
case cmd == "term":
|
||||
return conn.SignalDetached(pid, syscall.SIGTERM)
|
||||
case cmd == "kill":
|
||||
return conn.SignalDetached(pid, syscall.SIGKILL)
|
||||
case cmd == "dismiss":
|
||||
if !conn.DismissDetached(pid) {
|
||||
return fmt.Errorf("process %d not found", pid)
|
||||
}
|
||||
return nil
|
||||
case strings.HasPrefix(cmd, "signal "):
|
||||
parts := strings.Fields(cmd)
|
||||
if len(parts) < 2 {
|
||||
return fmt.Errorf("signal requires signal number")
|
||||
}
|
||||
sig, err := strconv.Atoi(parts[1])
|
||||
if err != nil {
|
||||
return fmt.Errorf("invalid signal: %s", parts[1])
|
||||
}
|
||||
return conn.SignalDetached(pid, syscall.Signal(sig))
|
||||
default:
|
||||
return fmt.Errorf("unknown command: %s (valid: term, kill, signal <n>, dismiss)", cmd)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -69,6 +69,10 @@ func (s *SessionNode) ToolsConn() *toolsrv.Conn {
|
|||
return s.Core.ToolsConn()
|
||||
}
|
||||
|
||||
func (s *SessionNode) DialToolServer() *toolsrv.Conn {
|
||||
return s.Core.DialToolServer()
|
||||
}
|
||||
|
||||
// --- AgentLog management ---
|
||||
|
||||
func (s *SessionNode) AgentLog() *AgentLog {
|
||||
|
|
|
|||
|
|
@ -294,10 +294,21 @@ FileNode("generate", 0666,
|
|||
Doc("Detached background processes started by this agent."),
|
||||
|
||||
Each("{pid}", processBindings,
|
||||
Doc("Process output buffer. Read: output. Remove: dismiss after completion."),
|
||||
Doc("Background process. Directory contains out, ctl, status."),
|
||||
Bind(BindData),
|
||||
Read(readProcess),
|
||||
Remove(removeProcess),
|
||||
|
||||
FileNode("out", 0444,
|
||||
Doc("Process output buffer."),
|
||||
Read(readProcess),
|
||||
),
|
||||
FileNode("ctl", 0222,
|
||||
Doc("Process control. Write: term, kill, signal <n>, dismiss."),
|
||||
Write(writeProcessCtl),
|
||||
),
|
||||
FileNode("status", 0444,
|
||||
Doc("Process status: running or exited <code>."),
|
||||
Read(readProcessStatus),
|
||||
),
|
||||
),
|
||||
),
|
||||
),
|
||||
|
|
|
|||
|
|
@ -50,13 +50,13 @@ When a tool or shell command fails due to sandbox restrictions, you may retry wi
|
|||
|
||||
Run long-running commands in the background by passing `"background": true` to shell: `{"cmd": "go test ./...", "background": true}`.
|
||||
|
||||
The result is a `<system-proc-background>` tag containing the PID. You do NOT need to poll for output — background process updates are automatically injected into your context as `<system-proc-interrupt>` blocks alongside tool results. These include status (running/exited), exit code, and the last 20 lines of output.
|
||||
The result is a `<system-proc-background>` tag containing the process ID. You do NOT need to poll for output — background process updates are automatically injected into your context as `<system-proc-interrupt>` blocks alongside tool results. These include status (running/exited), exit code, and the last 20 lines of output.
|
||||
|
||||
**Rules**:
|
||||
- Use background for long operations (builds, test suites, deployments, log tailing). Do NOT background short commands (< 5s).
|
||||
- React to `<system-proc-interrupt>` naturally. If a build fails, fix it. If output is irrelevant, kill the process.
|
||||
- Kill background processes when done: `{"cmd": "kill <pid>"}`.
|
||||
- You can still read full output explicitly via `process_output` if needed.
|
||||
- Stop background processes when done: `echo term | ollie-9p write session/$OLLIE_SESSION_ID/agent/$OLLIE_UNAME/proc/<id>/ctl` (use `kill` instead of `term` to force-kill)
|
||||
- Read full output: `ollie-9p read session/$OLLIE_SESSION_ID/agent/$OLLIE_UNAME/proc/<id>/out`
|
||||
|
||||
# Security
|
||||
|
||||
|
|
@ -102,6 +102,9 @@ Your world model is a 9P filesystem. Your session ID is `${OLLIE_SESSION_ID}`. U
|
|||
| `stats` | read | Usage, cost, context size |
|
||||
| `tools` | r/w | Write: load tool by name. Read: list loaded tools. |
|
||||
| `proc/` | dir | Detached background processes |
|
||||
| `proc/{id}/out` | read | Process output buffer |
|
||||
| `proc/{id}/ctl` | write | Process control: term, kill, signal \<n\>, dismiss |
|
||||
| `proc/{id}/status` | read | Process status: running or exited \<code\> |
|
||||
|
||||
## Operations
|
||||
|
||||
|
|
|
|||
|
|
@ -84,6 +84,22 @@ func (s *Session) ToolsConn() *toolsrv.Conn {
|
|||
return s.toolsConn
|
||||
}
|
||||
|
||||
// DialToolServer dials a fresh toolsrv connection via the keeper.
|
||||
// Caller is responsible for closing the returned conn.
|
||||
func (s *Session) DialToolServer() *toolsrv.Conn {
|
||||
s.mu.RLock()
|
||||
k := s.Keeper
|
||||
s.mu.RUnlock()
|
||||
if k == nil {
|
||||
return nil
|
||||
}
|
||||
conn, err := k.Dial()
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
return conn
|
||||
}
|
||||
|
||||
// SetToolsConn sets the tool server connection.
|
||||
func (s *Session) SetToolsConn(conn *toolsrv.Conn) {
|
||||
s.mu.Lock()
|
||||
|
|
|
|||
|
|
@ -27,6 +27,8 @@ type Config struct {
|
|||
Env map[string]string
|
||||
Yolo bool
|
||||
Timeout int
|
||||
Output io.Writer // if set, stream stdout/stderr here in real-time
|
||||
Started chan *os.Process // if set, receives the os.Process after Start()
|
||||
}
|
||||
|
||||
// StreamFunc returns the streaming output function from context, if any.
|
||||
|
|
@ -55,7 +57,8 @@ func ExecuteTool(ctx context.Context, info toolsrv.ToolInfo, args json.RawMessag
|
|||
// Extract dispatch-level flags from args
|
||||
bypassed := false
|
||||
timeout := cfg.Timeout
|
||||
if timeout <= 0 {
|
||||
timeoutExplicit := timeout > 0
|
||||
if !timeoutExplicit {
|
||||
timeout = 30
|
||||
}
|
||||
sandboxName := "default"
|
||||
|
|
@ -72,6 +75,7 @@ func ExecuteTool(ctx context.Context, info toolsrv.ToolInfo, args json.RawMessag
|
|||
delete(argMap, "bypass")
|
||||
}
|
||||
if t, ok := argMap["timeout"]; ok {
|
||||
timeoutExplicit = true
|
||||
switch v := t.(type) {
|
||||
case float64:
|
||||
timeout = int(v)
|
||||
|
|
@ -116,7 +120,7 @@ func ExecuteTool(ctx context.Context, info toolsrv.ToolInfo, args json.RawMessag
|
|||
bypassCode := fmt.Sprintf("cat <<'OLLIE_EOF' | %s\n%s\nOLLIE_EOF", toolPath, stdinData)
|
||||
result, err = executeBypassDirect(ctx, bypassCode, cwd, cfg.Env, timeout, false)
|
||||
} else {
|
||||
result, err = executeSandboxed(ctx, toolPath, stdinData, cwd, cfg.Env, timeout, sandboxName, cfg.Yolo)
|
||||
result, err = executeSandboxed(ctx, toolPath, stdinData, cwd, cfg.Env, timeout, sandboxName, cfg.Yolo, cfg.Output, cfg.Started)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
|
|
@ -131,11 +135,7 @@ func ExecuteTool(ctx context.Context, info toolsrv.ToolInfo, args json.RawMessag
|
|||
}
|
||||
|
||||
// executeSandboxed runs a tool script inside the sandbox.
|
||||
func executeSandboxed(ctx context.Context, toolPath, stdinData, cwd string, envExtra map[string]string, timeout int, sandboxName string, yolo bool) (string, error) {
|
||||
if timeout <= 0 {
|
||||
timeout = 30
|
||||
}
|
||||
|
||||
func executeSandboxed(ctx context.Context, toolPath, stdinData, cwd string, envExtra map[string]string, timeout int, sandboxName string, yolo bool, streamOut io.Writer, started chan *os.Process) (string, error) {
|
||||
var sandboxCfg *sandbox.Config
|
||||
if !yolo {
|
||||
cfgPath := filepath.Join(paths.CfgDir(), "sandbox", sandboxName+".yaml")
|
||||
|
|
@ -151,8 +151,11 @@ func executeSandboxed(ctx context.Context, toolPath, stdinData, cwd string, envE
|
|||
}
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(ctx, time.Duration(timeout)*time.Second)
|
||||
defer cancel()
|
||||
if timeout > 0 {
|
||||
var cancel context.CancelFunc
|
||||
ctx, cancel = context.WithTimeout(ctx, time.Duration(timeout)*time.Second)
|
||||
defer cancel()
|
||||
}
|
||||
|
||||
interpreter := []string{"bash", "-c", toolPath}
|
||||
|
||||
|
|
@ -197,15 +200,19 @@ func executeSandboxed(ctx context.Context, toolPath, stdinData, cwd string, envE
|
|||
cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true}
|
||||
cmd.Cancel = func() error {
|
||||
if cmd.Process != nil {
|
||||
syscall.Kill(-cmd.Process.Pid, syscall.SIGKILL)
|
||||
syscall.Kill(-cmd.Process.Pid, syscall.SIGTERM)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
cmd.WaitDelay = time.Second
|
||||
cmd.WaitDelay = 5 * time.Second
|
||||
|
||||
var outputBuf bytes.Buffer
|
||||
var w io.Writer = &outputBuf
|
||||
if streamOut != nil {
|
||||
w = io.MultiWriter(&outputBuf, streamOut)
|
||||
}
|
||||
lw := &limitedWriter{
|
||||
w: &outputBuf,
|
||||
w: w,
|
||||
limit: 10 * 1024 * 1024,
|
||||
stream: StreamFunc(ctx),
|
||||
}
|
||||
|
|
@ -213,9 +220,16 @@ func executeSandboxed(ctx context.Context, toolPath, stdinData, cwd string, envE
|
|||
cmd.Stderr = lw
|
||||
|
||||
if err := cmd.Start(); err != nil {
|
||||
if started != nil {
|
||||
close(started)
|
||||
}
|
||||
return "", fmt.Errorf("execution failed: %w", err)
|
||||
}
|
||||
|
||||
if started != nil {
|
||||
started <- cmd.Process
|
||||
}
|
||||
|
||||
err := cmd.Wait()
|
||||
output := outputBuf.Bytes()
|
||||
if lw.truncated {
|
||||
|
|
@ -244,8 +258,11 @@ func executeBypassDirect(ctx context.Context, cmd, cwd string, envExtra map[stri
|
|||
}
|
||||
sockPath := filepath.Join(xdg, "ollie", "bypass.sock")
|
||||
|
||||
ctx, cancel := context.WithTimeout(ctx, time.Duration(timeout)*time.Second)
|
||||
defer cancel()
|
||||
if timeout > 0 {
|
||||
var cancel context.CancelFunc
|
||||
ctx, cancel = context.WithTimeout(ctx, time.Duration(timeout)*time.Second)
|
||||
defer cancel()
|
||||
}
|
||||
|
||||
// Build environment map
|
||||
envMap := make(map[string]string)
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ import (
|
|||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"runtime"
|
||||
"strconv"
|
||||
|
|
@ -33,7 +34,7 @@ type State struct {
|
|||
// Process management
|
||||
procMu sync.Mutex
|
||||
procs map[int]*Proc
|
||||
nextPID int
|
||||
nextID int
|
||||
procLimit int
|
||||
|
||||
// Callbacks
|
||||
|
|
@ -42,13 +43,14 @@ type State struct {
|
|||
|
||||
// Proc represents a running or completed process.
|
||||
type Proc struct {
|
||||
PID int
|
||||
ID int
|
||||
Tool string
|
||||
Args map[string]string
|
||||
StartTime time.Time
|
||||
EndTime time.Time
|
||||
ExitCode int
|
||||
Exited bool
|
||||
LastRead time.Time // last time output was read; GC TTL counts from here
|
||||
|
||||
mu sync.Mutex
|
||||
output bytes.Buffer
|
||||
|
|
@ -63,7 +65,7 @@ func NewState(cwd string) *State {
|
|||
cwd: cwd,
|
||||
env: make(map[string]string),
|
||||
procs: make(map[int]*Proc),
|
||||
nextPID: 1,
|
||||
nextID: 1,
|
||||
procLimit: 32,
|
||||
}
|
||||
}
|
||||
|
|
@ -177,12 +179,12 @@ func (st *State) UnloadTool(name string) error {
|
|||
|
||||
// --- Process Management ---
|
||||
|
||||
// allocPID allocates a new process ID.
|
||||
func (st *State) allocPID() int {
|
||||
// allocID allocates a new process ID.
|
||||
func (st *State) allocID() int {
|
||||
st.procMu.Lock()
|
||||
defer st.procMu.Unlock()
|
||||
pid := st.nextPID
|
||||
st.nextPID++
|
||||
pid := st.nextID
|
||||
st.nextID++
|
||||
return pid
|
||||
}
|
||||
|
||||
|
|
@ -229,7 +231,7 @@ func (st *State) NewProc(ctx context.Context, payload string, background bool) (
|
|||
// Create proc
|
||||
procCtx, cancel := context.WithCancel(ctx)
|
||||
proc := &Proc{
|
||||
PID: st.allocPID(),
|
||||
ID: st.allocID(),
|
||||
Tool: toolName,
|
||||
Args: args,
|
||||
StartTime: time.Now(),
|
||||
|
|
@ -238,26 +240,43 @@ func (st *State) NewProc(ctx context.Context, payload string, background bool) (
|
|||
}
|
||||
|
||||
st.procMu.Lock()
|
||||
st.procs[proc.PID] = proc
|
||||
st.procs[proc.ID] = proc
|
||||
st.procMu.Unlock()
|
||||
|
||||
// For background procs, stream output directly into the proc buffer.
|
||||
// For foreground, use nil (local buffer in exec path).
|
||||
var outputWriter io.Writer
|
||||
var startedCh chan *os.Process
|
||||
if background {
|
||||
outputWriter = &procWriter{proc: proc}
|
||||
startedCh = make(chan *os.Process, 1)
|
||||
}
|
||||
|
||||
// Execute in goroutine
|
||||
go func() {
|
||||
defer close(proc.done)
|
||||
defer cancel()
|
||||
|
||||
out, exitCode := st.executeTool(procCtx, info, args, cwd, envCopy, yolo)
|
||||
out, exitCode := st.executeTool(procCtx, info, args, cwd, envCopy, yolo, outputWriter, startedCh)
|
||||
|
||||
proc.mu.Lock()
|
||||
proc.output.WriteString(out)
|
||||
if !background {
|
||||
proc.output.WriteString(out)
|
||||
}
|
||||
proc.ExitCode = exitCode
|
||||
proc.Exited = true
|
||||
proc.EndTime = time.Now()
|
||||
proc.mu.Unlock()
|
||||
}()
|
||||
|
||||
// For background procs, wait for the process to start and store its reference.
|
||||
if background {
|
||||
return fmt.Sprintf("%d", proc.PID), proc.PID, nil
|
||||
if p, ok := <-startedCh; ok && p != nil {
|
||||
proc.mu.Lock()
|
||||
proc.proc = p
|
||||
proc.mu.Unlock()
|
||||
}
|
||||
return fmt.Sprintf("%d", proc.ID), proc.ID, nil
|
||||
}
|
||||
|
||||
// Wait for completion
|
||||
|
|
@ -270,24 +289,26 @@ func (st *State) NewProc(ctx context.Context, payload string, background bool) (
|
|||
|
||||
// Auto-cleanup for foreground procs
|
||||
st.procMu.Lock()
|
||||
delete(st.procs, proc.PID)
|
||||
delete(st.procs, proc.ID)
|
||||
st.procMu.Unlock()
|
||||
|
||||
if exitCode != 0 {
|
||||
return result, proc.PID, fmt.Errorf("exit %d", exitCode)
|
||||
return result, proc.ID, fmt.Errorf("exit %d", exitCode)
|
||||
}
|
||||
return result, proc.PID, nil
|
||||
return result, proc.ID, nil
|
||||
}
|
||||
|
||||
// executeTool runs a tool and returns output + exit code.
|
||||
func (st *State) executeTool(ctx context.Context, info toolsrv.ToolInfo, args map[string]string, cwd string, envExtra map[string]string, yolo bool) (string, int) {
|
||||
func (st *State) executeTool(ctx context.Context, info toolsrv.ToolInfo, args map[string]string, cwd string, envExtra map[string]string, yolo bool, output io.Writer, started chan *os.Process) (string, int) {
|
||||
// Convert args map to JSON for the execution path
|
||||
jsonArgs := argsToJSON(args)
|
||||
|
||||
cfg := exec.Config{
|
||||
CWD: cwd,
|
||||
Env: envExtra,
|
||||
Yolo: yolo,
|
||||
CWD: cwd,
|
||||
Env: envExtra,
|
||||
Yolo: yolo,
|
||||
Output: output,
|
||||
Started: started,
|
||||
}
|
||||
|
||||
result, err := exec.ExecuteTool(ctx, info, jsonArgs, cfg)
|
||||
|
|
@ -326,6 +347,40 @@ func (st *State) DismissProc(pid int) bool {
|
|||
return false
|
||||
}
|
||||
|
||||
const procGCTTL = 10 * time.Minute
|
||||
|
||||
// StartProcGC runs a background goroutine that removes exited procs
|
||||
// whose output has been read and whose LastRead is older than procGCTTL.
|
||||
func (st *State) StartProcGC(ctx context.Context) {
|
||||
go func() {
|
||||
ticker := time.NewTicker(1 * time.Minute)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
st.gcProcs()
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
func (st *State) gcProcs() {
|
||||
st.procMu.Lock()
|
||||
defer st.procMu.Unlock()
|
||||
now := time.Now()
|
||||
for pid, p := range st.procs {
|
||||
p.mu.Lock()
|
||||
exited := p.Exited
|
||||
lastRead := p.LastRead
|
||||
p.mu.Unlock()
|
||||
if exited && !lastRead.IsZero() && now.Sub(lastRead) >= procGCTTL {
|
||||
delete(st.procs, pid)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// SignalProc sends a signal to a process.
|
||||
func (st *State) SignalProc(pid int, sig syscall.Signal) error {
|
||||
proc := st.GetProc(pid)
|
||||
|
|
@ -333,17 +388,17 @@ func (st *State) SignalProc(pid int, sig syscall.Signal) error {
|
|||
return fmt.Errorf("process not found: %d", pid)
|
||||
}
|
||||
|
||||
// Cancel the context (for graceful stop)
|
||||
if sig == syscall.SIGTERM || sig == syscall.SIGINT {
|
||||
proc.cancel()
|
||||
}
|
||||
|
||||
// Also signal the underlying process if available
|
||||
// Signal the process group (we set Setpgid: true)
|
||||
proc.mu.Lock()
|
||||
p := proc.proc
|
||||
proc.mu.Unlock()
|
||||
if p != nil {
|
||||
p.Signal(sig)
|
||||
syscall.Kill(-p.Pid, sig)
|
||||
}
|
||||
|
||||
// Cancel context on SIGKILL (hard stop)
|
||||
if sig == syscall.SIGKILL {
|
||||
proc.cancel()
|
||||
}
|
||||
|
||||
return nil
|
||||
|
|
@ -441,6 +496,10 @@ func (st *State) HandleProcCtl(pid int, input string) error {
|
|||
}
|
||||
|
||||
switch parts[0] {
|
||||
case "term":
|
||||
return st.SignalProc(pid, syscall.SIGTERM)
|
||||
case "kill":
|
||||
return st.SignalProc(pid, syscall.SIGKILL)
|
||||
case "signal":
|
||||
if len(parts) < 2 {
|
||||
return fmt.Errorf("signal requires signal number")
|
||||
|
|
@ -458,12 +517,26 @@ func (st *State) HandleProcCtl(pid int, input string) error {
|
|||
}
|
||||
}
|
||||
|
||||
// --- procWriter ---
|
||||
|
||||
// procWriter is a thread-safe writer that streams into a Proc's output buffer.
|
||||
type procWriter struct {
|
||||
proc *Proc
|
||||
}
|
||||
|
||||
func (pw *procWriter) Write(p []byte) (int, error) {
|
||||
pw.proc.mu.Lock()
|
||||
defer pw.proc.mu.Unlock()
|
||||
return pw.proc.output.Write(p)
|
||||
}
|
||||
|
||||
// --- Proc Methods ---
|
||||
|
||||
// Output returns the current output buffer.
|
||||
// Output returns the current output buffer and updates LastRead.
|
||||
func (p *Proc) Output() string {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
p.LastRead = time.Now()
|
||||
return p.output.String()
|
||||
}
|
||||
|
||||
|
|
@ -475,7 +548,7 @@ func (p *Proc) Wait() int {
|
|||
return p.ExitCode
|
||||
}
|
||||
|
||||
// Stat returns a status string.
|
||||
// Stat returns a status string in key=value format.
|
||||
func (p *Proc) Stat() string {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
|
|
@ -483,9 +556,9 @@ func (p *Proc) Stat() string {
|
|||
rt := time.Since(p.StartTime)
|
||||
if p.Exited {
|
||||
rt = p.EndTime.Sub(p.StartTime)
|
||||
return fmt.Sprintf("exited %d\nruntime=%v\ntool=%s\n", p.ExitCode, rt, p.Tool)
|
||||
return fmt.Sprintf("id=%d\nexited=true\nexit_code=%d\nruntime=%v\ntool=%s\n", p.ID, p.ExitCode, rt, p.Tool)
|
||||
}
|
||||
return fmt.Sprintf("running\nruntime=%v\ntool=%s\n", rt, p.Tool)
|
||||
return fmt.Sprintf("id=%d\nexited=false\nruntime=%v\ntool=%s\n", p.ID, rt, p.Tool)
|
||||
}
|
||||
|
||||
// --- Helpers ---
|
||||
|
|
|
|||
|
|
@ -126,7 +126,7 @@ func TestState_HandleProcCtl(t *testing.T) {
|
|||
|
||||
func TestProc_Stat(t *testing.T) {
|
||||
proc := &Proc{
|
||||
PID: 1,
|
||||
ID: 1,
|
||||
Tool: "shell",
|
||||
StartTime: time.Now(),
|
||||
done: make(chan struct{}),
|
||||
|
|
@ -134,8 +134,8 @@ func TestProc_Stat(t *testing.T) {
|
|||
|
||||
// Running state
|
||||
stat := proc.Stat()
|
||||
if !strings.Contains(stat, "running") {
|
||||
t.Errorf("Stat() for running proc should contain 'running', got %q", stat)
|
||||
if !strings.Contains(stat, "exited=false") {
|
||||
t.Errorf("Stat() for running proc should contain 'exited=false', got %q", stat)
|
||||
}
|
||||
|
||||
// Mark as exited
|
||||
|
|
@ -144,8 +144,8 @@ func TestProc_Stat(t *testing.T) {
|
|||
proc.EndTime = time.Now()
|
||||
|
||||
stat = proc.Stat()
|
||||
if !strings.Contains(stat, "exited 0") {
|
||||
t.Errorf("Stat() for exited proc should contain 'exited 0', got %q", stat)
|
||||
if !strings.Contains(stat, "exited=true") {
|
||||
t.Errorf("Stat() for exited proc should contain 'exited=true', got %q", stat)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -191,3 +191,56 @@ func TestParsePayload(t *testing.T) {
|
|||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestGcProcs(t *testing.T) {
|
||||
st := NewState("/tmp")
|
||||
|
||||
// Add a proc that's exited and was read long ago
|
||||
st.procMu.Lock()
|
||||
st.procs[1] = &Proc{
|
||||
ID: 1,
|
||||
Exited: true,
|
||||
EndTime: time.Now().Add(-20 * time.Minute),
|
||||
LastRead: time.Now().Add(-15 * time.Minute),
|
||||
done: make(chan struct{}),
|
||||
}
|
||||
// Add a proc that's exited but never read
|
||||
st.procs[2] = &Proc{
|
||||
ID: 2,
|
||||
Exited: true,
|
||||
done: make(chan struct{}),
|
||||
}
|
||||
// Add a proc that's exited and recently read
|
||||
st.procs[3] = &Proc{
|
||||
ID: 3,
|
||||
Exited: true,
|
||||
LastRead: time.Now(),
|
||||
done: make(chan struct{}),
|
||||
}
|
||||
// Add a still-running proc
|
||||
st.procs[4] = &Proc{
|
||||
ID: 4,
|
||||
Exited: false,
|
||||
LastRead: time.Now().Add(-20 * time.Minute),
|
||||
done: make(chan struct{}),
|
||||
}
|
||||
st.procMu.Unlock()
|
||||
|
||||
st.gcProcs()
|
||||
|
||||
st.procMu.Lock()
|
||||
defer st.procMu.Unlock()
|
||||
|
||||
if _, ok := st.procs[1]; ok {
|
||||
t.Error("proc 1 should have been GC'd (exited, read >10min ago)")
|
||||
}
|
||||
if _, ok := st.procs[2]; !ok {
|
||||
t.Error("proc 2 should NOT be GC'd (never read)")
|
||||
}
|
||||
if _, ok := st.procs[3]; !ok {
|
||||
t.Error("proc 3 should NOT be GC'd (recently read)")
|
||||
}
|
||||
if _, ok := st.procs[4]; !ok {
|
||||
t.Error("proc 4 should NOT be GC'd (still running)")
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -165,7 +165,7 @@ func handleProcCtlWrite(ctx Ctx, data []byte) error {
|
|||
if ctx.Proc == nil {
|
||||
return fmt.Errorf("no process context")
|
||||
}
|
||||
return ctx.Server.Fs.HandleProcCtl(ctx.Proc.PID, strings.TrimSpace(string(data)))
|
||||
return ctx.Server.Fs.HandleProcCtl(ctx.Proc.ID, strings.TrimSpace(string(data)))
|
||||
}
|
||||
|
||||
// Mode constants for convenience.
|
||||
|
|
|
|||
|
|
@ -110,6 +110,9 @@ func main() {
|
|||
go idleMonitor(runCtx, cancel, *idleTimeout)
|
||||
}
|
||||
|
||||
// Start proc GC
|
||||
srv.Fs.StartProcGC(runCtx)
|
||||
|
||||
// Remove stale socket if it exists
|
||||
os.Remove(*listenPath)
|
||||
|
||||
|
|
|
|||
|
|
@ -30,10 +30,8 @@ fi
|
|||
|
||||
# Execute the command. The sandbox (landlock) is already applied by toolsrv.
|
||||
# We use eval to handle pipes, redirects, and compound commands.
|
||||
# Capture output and exit code, always report both.
|
||||
output=$(eval "$cmd" 2>&1) && rc=$? || rc=$?
|
||||
if [ -n "$output" ]; then
|
||||
printf '%s\n' "$output"
|
||||
fi
|
||||
# Output streams directly to stdout/stderr for real-time capture.
|
||||
eval "$cmd" 2>&1
|
||||
rc=$?
|
||||
echo "exit: $rc"
|
||||
exit $rc
|
||||
|
|
@ -379,7 +379,7 @@ func (c *Conn) readProcInfo(name string) map[string]any {
|
|||
key := line[:idx]
|
||||
val := line[idx+1:]
|
||||
switch key {
|
||||
case "pid":
|
||||
case "id":
|
||||
var pid int
|
||||
fmt.Sscanf(val, "%d", &pid)
|
||||
info["pid"] = pid
|
||||
|
|
|
|||
Loading…
Reference in New Issue