370 lines
9.9 KiB
Go
370 lines
9.9 KiB
Go
package toolsrv
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"strings"
|
|
"sync"
|
|
"syscall"
|
|
"time"
|
|
|
|
"ollie/paths"
|
|
)
|
|
|
|
const (
|
|
failureWindow = 1 * time.Minute
|
|
maxFailures = 5
|
|
blockDuration = 5 * time.Minute
|
|
)
|
|
|
|
// RateLimitedError is returned by checkRateLimit when the shell is blocked
|
|
// due to too many validation failures. The agent loop uses this to detect
|
|
// the condition and inject a clear signal rather than counting it as a
|
|
// normal tool error.
|
|
type RateLimitedError struct {
|
|
Remaining time.Duration
|
|
}
|
|
|
|
func (e *RateLimitedError) Error() string {
|
|
return fmt.Sprintf("rate limited: too many validation failures, blocked for %v", e.Remaining)
|
|
}
|
|
|
|
// Server runs code in a sandboxed environment.
|
|
type Server struct {
|
|
// cwd is the working directory for sandboxed commands. If empty,
|
|
// the process working directory is used.
|
|
wdMu sync.RWMutex
|
|
cwd string
|
|
|
|
// envExtra holds per-session environment variables injected via SetEnv.
|
|
envMu sync.RWMutex
|
|
envExtra map[string]string
|
|
|
|
// Hooks for lifecycle events (harnesses like 9P can inject mount logic here)
|
|
OnClose func()
|
|
OnEnvSet func(key, value string)
|
|
|
|
// Yolo skips the landrun sandbox for all execution.
|
|
Yolo bool
|
|
|
|
toolRegistry *Registry
|
|
sessionID string
|
|
|
|
// OnInjection is called when a skill is loaded and its content
|
|
// should be injected into the agent's context.
|
|
OnInjection func(content string)
|
|
|
|
// OnToolsChanged is called when the set of loaded tools changes.
|
|
// The receiver should re-fetch the tool list to update its state.
|
|
OnToolsChanged func()
|
|
|
|
// rate limiting state (per-Server)
|
|
rateLimitMu sync.Mutex
|
|
validationFailures int
|
|
lastFailure time.Time
|
|
blockedUntil time.Time
|
|
|
|
// Detached process management
|
|
detachMu sync.Mutex
|
|
detachCh chan struct{} // signal to detach the currently running process
|
|
detached []*DetachedProcess
|
|
OnDetach func(pid int, cmd string) // hook: called when a process is detached
|
|
OnExit func(pid int, exitCode int) // hook: called when a detached process exits
|
|
}
|
|
|
|
// Option configures a Server.
|
|
type Option func(*Server)
|
|
|
|
// WithYolo skips the landrun sandbox.
|
|
func WithYolo() Option { return func(s *Server) { s.Yolo = true } }
|
|
|
|
|
|
|
|
// WithToolRegistry attaches a tool registry and session ID to the Server.
|
|
func WithToolRegistry(r *Registry, sessionID string) Option {
|
|
return func(s *Server) {
|
|
s.toolRegistry = r
|
|
s.sessionID = sessionID
|
|
}
|
|
}
|
|
|
|
// LoadTool loads a tool by name into this server's tool registry.
|
|
// Returns an error if no registry is configured or the tool cannot be found.
|
|
// If OnToolsChanged is set, it is called to signal that the available tools
|
|
// have changed. The caller is responsible for re-fetching the current list.
|
|
func (s *Server) LoadTool(name string) error {
|
|
if s.toolRegistry == nil {
|
|
return fmt.Errorf("no tool registry configured")
|
|
}
|
|
if s.sessionID == "" {
|
|
return fmt.Errorf("no session ID configured")
|
|
}
|
|
if err := s.toolRegistry.Load(s.sessionID, name); err != nil {
|
|
return err
|
|
}
|
|
if s.OnToolsChanged != nil {
|
|
s.OnToolsChanged()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
|
|
|
|
// ListTools implements Server, returning tools loaded in the session's
|
|
// tool registry. With zero built-in tools, only dynamically loaded tools
|
|
// appear here.
|
|
func (s *Server) ListTools() ([]ToolInfo, error) {
|
|
var all []ToolInfo
|
|
|
|
if s.toolRegistry != nil && s.sessionID != "" {
|
|
all = append(all, s.toolRegistry.Loaded(s.sessionID)...)
|
|
}
|
|
return all, nil
|
|
}
|
|
|
|
// CallTool implements Runner.
|
|
// With zero built-in tools, all tools must be loaded in the registry first.
|
|
func (s *Server) CallTool(ctx context.Context, tool string, args json.RawMessage) (json.RawMessage, error) {
|
|
if s.toolRegistry != nil && s.sessionID != "" {
|
|
if _, promoted := s.toolRegistry.Lookup(s.sessionID, tool); promoted {
|
|
return s.callPromotedTool(ctx, tool, args)
|
|
}
|
|
}
|
|
return nil, fmt.Errorf("unknown tool: %s (not loaded)", tool)
|
|
}
|
|
|
|
// New creates a new Server with the given working directory.
|
|
func New(cwd string) *Server { return &Server{cwd: paths.ExpandHome(cwd)} }
|
|
|
|
// SetCWD updates the working directory used for subsequent command executions.
|
|
func (s *Server) SetCWD(dir string) {
|
|
s.wdMu.Lock()
|
|
s.cwd = paths.ExpandHome(dir)
|
|
s.wdMu.Unlock()
|
|
}
|
|
|
|
// SetEnv adds a session-scoped environment variable injected into all
|
|
// subsequent subprocess invocations for this session.
|
|
func (s *Server) SetEnv(key, value string) {
|
|
s.envMu.Lock()
|
|
if s.envExtra == nil {
|
|
s.envExtra = make(map[string]string)
|
|
}
|
|
s.envExtra[key] = value
|
|
s.envMu.Unlock()
|
|
if s.OnEnvSet != nil {
|
|
s.OnEnvSet(key, value)
|
|
}
|
|
}
|
|
|
|
// SetToolRegistry attaches a session-local tool registry.
|
|
func (s *Server) SetToolRegistry(r *Registry, sessionID string) {
|
|
s.toolRegistry = r
|
|
s.sessionID = sessionID
|
|
}
|
|
|
|
// callPromotedTool executes a tool promoted via the registry by running the
|
|
// script file inside the sandbox, piping the JSON args to stdin.
|
|
func (s *Server) callPromotedTool(ctx context.Context, tool string, args json.RawMessage) (json.RawMessage, error) {
|
|
// Resolve script path.
|
|
if strings.Contains(tool, "/") || strings.Contains(tool, "..") {
|
|
return nil, fmt.Errorf("invalid tool name")
|
|
}
|
|
path, err := ResolveTool(tool)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// Extract dispatch-level flags (not passed to tool script).
|
|
bypassed := false
|
|
timeout := 30
|
|
sandboxName := "default"
|
|
var argMap map[string]interface{}
|
|
if err := json.Unmarshal(args, &argMap); err == nil {
|
|
if e, ok := argMap["bypass"]; ok {
|
|
switch v := e.(type) {
|
|
case bool:
|
|
bypassed = v
|
|
case string:
|
|
bypassed = v == "true" || v == "1"
|
|
}
|
|
delete(argMap, "bypass")
|
|
}
|
|
if t, ok := argMap["timeout"]; ok {
|
|
switch v := t.(type) {
|
|
case float64:
|
|
timeout = int(v)
|
|
case string:
|
|
var n int
|
|
if _, err := fmt.Sscanf(v, "%d", &n); err == nil {
|
|
timeout = n
|
|
}
|
|
}
|
|
delete(argMap, "timeout")
|
|
}
|
|
if s, ok := argMap["sandbox"]; ok {
|
|
if v, ok := s.(string); ok && v != "" {
|
|
sandboxName = v
|
|
}
|
|
delete(argMap, "sandbox")
|
|
}
|
|
args, _ = json.Marshal(argMap)
|
|
}
|
|
|
|
// Check if tool requires sudo (from .meta).
|
|
needsSudo := false
|
|
if m, err := LoadMetaFile(tool); err == nil && m != nil {
|
|
resolved := m.Resolve()
|
|
if resolved != nil && resolved.Sudo {
|
|
needsSudo = true
|
|
bypassed = true // sudo implies bypass
|
|
}
|
|
}
|
|
|
|
// The tool script receives JSON args on stdin.
|
|
code := path
|
|
stdinData := string(args)
|
|
|
|
var result string
|
|
if needsSudo {
|
|
s.wdMu.RLock()
|
|
workDir := s.cwd
|
|
s.wdMu.RUnlock()
|
|
sudoCode := fmt.Sprintf("cat <<'OLLIE_EOF' | %s\n%s\nOLLIE_EOF", code, stdinData)
|
|
result, err = s.executeBypassSudo(ctx, sudoCode, workDir, timeout)
|
|
} else if bypassed {
|
|
s.wdMu.RLock()
|
|
workDir := s.cwd
|
|
s.wdMu.RUnlock()
|
|
// Broker protocol has no stdin support; pipe JSON via heredoc.
|
|
bypassCode := fmt.Sprintf("cat <<'OLLIE_EOF' | %s\n%s\nOLLIE_EOF", code, stdinData)
|
|
result, err = s.executeBypass(ctx, bypassCode, workDir, timeout)
|
|
} else {
|
|
result, err = s.executeWithStdin(ctx, code, "bash", timeout, sandboxName, false, stdinData)
|
|
}
|
|
if err != nil {
|
|
return json.Marshal(ToolResult{
|
|
IsError: true,
|
|
Content: []ToolResultContent{{Type: "text", Text: result + ": " + err.Error()}},
|
|
})
|
|
}
|
|
return json.Marshal(ToolResult{
|
|
Content: []ToolResultContent{{Type: "text", Text: result}},
|
|
})
|
|
}
|
|
|
|
// Close is called when the session ends. Calls OnClose hook if registered.
|
|
func (s *Server) Close() {
|
|
s.cleanupDetached()
|
|
if s.OnClose != nil {
|
|
s.OnClose()
|
|
}
|
|
}
|
|
|
|
// Detach signals the currently running process to be detached from the agent.
|
|
// The process continues running; its output is captured in a ring buffer.
|
|
// Returns false if no process is currently running.
|
|
func (s *Server) Detach() bool {
|
|
s.detachMu.Lock()
|
|
ch := s.detachCh
|
|
s.detachMu.Unlock()
|
|
if ch == nil {
|
|
return false
|
|
}
|
|
select {
|
|
case ch <- struct{}{}:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
// ListDetachedRaw returns detached process info as []any (each element is map[string]any)
|
|
// for consumption by packages that can't import this package directly.
|
|
func (s *Server) ListDetachedRaw() []any {
|
|
s.detachMu.Lock()
|
|
defer s.detachMu.Unlock()
|
|
out := make([]any, len(s.detached))
|
|
for i, p := range s.detached {
|
|
info := p.Info()
|
|
out[i] = map[string]any{
|
|
"pid": info.PID,
|
|
"command": info.Command,
|
|
"started": info.Started,
|
|
"exited": info.Exited,
|
|
"exit_code": info.ExitCode,
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// SignalDetached sends a signal to a detached process by PID.
|
|
func (s *Server) SignalDetached(pid int, sig syscall.Signal) error {
|
|
s.detachMu.Lock()
|
|
defer s.detachMu.Unlock()
|
|
for _, p := range s.detached {
|
|
if p.PID == pid {
|
|
p.Mu.Lock()
|
|
defer p.Mu.Unlock()
|
|
if p.Exited {
|
|
return fmt.Errorf("process %d already exited", pid)
|
|
}
|
|
return syscall.Kill(-pid, sig)
|
|
}
|
|
}
|
|
return fmt.Errorf("no detached process with pid %d", pid)
|
|
}
|
|
|
|
// GetDetachedOutput returns the ring buffer contents for a detached process.
|
|
func (s *Server) GetDetachedOutput(pid int) (string, error) {
|
|
s.detachMu.Lock()
|
|
defer s.detachMu.Unlock()
|
|
for _, p := range s.detached {
|
|
if p.PID == pid {
|
|
return p.Output(), nil
|
|
}
|
|
}
|
|
return "", fmt.Errorf("no detached process with pid %d", pid)
|
|
}
|
|
|
|
// DismissDetached removes an exited process from the list.
|
|
func (s *Server) DismissDetached(pid int) bool {
|
|
s.detachMu.Lock()
|
|
defer s.detachMu.Unlock()
|
|
for i, p := range s.detached {
|
|
if p.PID == pid && p.Exited {
|
|
s.detached = append(s.detached[:i], s.detached[i+1:]...)
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// cleanupDetached sends SIGTERM to all running detached processes.
|
|
func (s *Server) cleanupDetached() {
|
|
for _, p := range s.detached {
|
|
p.Mu.Lock()
|
|
if !p.Exited {
|
|
syscall.Kill(-p.PID, syscall.SIGTERM)
|
|
}
|
|
p.Mu.Unlock()
|
|
}
|
|
}
|
|
|
|
// SetOnToolsChanged sets the callback for tool directory changes.
|
|
func (s *Server) SetOnToolsChanged(fn func()) {
|
|
s.OnToolsChanged = fn
|
|
}
|
|
|
|
// ToolRegistryRevision returns the current registry revision for this server's session.
|
|
// Returns 0 if no registry is configured.
|
|
func (s *Server) ToolRegistryRevision() uint64 {
|
|
if s.toolRegistry == nil || s.sessionID == "" {
|
|
return 0
|
|
}
|
|
return s.toolRegistry.Revision(s.sessionID)
|
|
}
|
|
|
|
|