From 560d9d4125ae29c32c212ad56ae20d984058aa5d Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Wed, 15 Apr 2026 08:59:37 +0300 Subject: [PATCH] feat(harness): docker isolation backend + canonical synapbus-agent image MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit New `internal/harness/docker` package: per-run ephemeral container backend that runs each agent in `docker run --rm`, bind-mounts the materialized workdir at /workspace, and captures stdout/stderr/exit code/result.json the same way the subprocess backend does. Inspired by scion's pkg/runtime/docker.go: shell out to the docker CLI (zero new Go deps, zero CGO), per-task ephemeral containers with no warm pool, host-side scratch dir bind-mounted in. Default security posture (overridable per agent): --rm --cap-drop=ALL --security-opt=no-new-privileges --read-only with tmpfs /tmp --pids-limit=512 --user= --network=bridge (configurable; --network=none for air-gap) --memory / --cpus from agent config --add-host host.docker.internal:host-gateway on Linux The backend reuses subprocess.AgentConfig for gemini_md/claude_md/ mcp_servers/skills materialization so existing example configs work unchanged. Per-agent docker tunables go under a new `docker` block in harness_config_json: image, memory, cpus, network, extra_mounts, cap_add, read_only_root, user, entrypoint, command, extra_args. MCP host rewrite: `.gemini/settings.json` URLs of the form http://127.0.0.1:/mcp are rewritten to http://host.docker.internal:/mcp at materialization time so the in-container Gemini CLI can reach the SynapBus MCP server on the host without code changes in the example wrappers. Wired into the reactor and Registry resolver: - Registry.Resolve picks "docker" when harness_config_json contains a `"docker"` block, taking precedence over local_command so explicit isolation never silently downgrades. - reactor.agentBackendKind() returns backendDocker for the same case. - evaluateTrigger's harness-backend gate accepts backendDocker alongside subprocess + webhook. - main.go registers docker.Harness with the harness registry, passing the SynapBus listen port so the URL rewrite uses the correct host port. Smoke tests in docker_test.go (skipped when no docker daemon): - TestExecute_Hello: env injection + bind-mount writeback + message.json + result.json + stdout capture using alpine:3.20 - TestExecute_NoImage: rejects agents missing docker.image - TestExecute_TimeoutCancel: wall-clock budget kills the container New canonical agent image at image-build/synapbus-agent/: - Debian bookworm-slim base - Node 22 + @google/gemini-cli + @anthropic-ai/claude-code - jq, sqlite3, curl, git, python3, tini (PID 1 for signal forwarding) - Non-root agent user uid/gid 1000 - ENTRYPOINT tini, CMD /workspace/wrapper.sh No SynapBus binary inside the image — agents reach the host MCP server over the network at host.docker.internal:. Pre-existing reactor test failures (TestReactorNoK8sImage, TestReactorDepthExceeded, TestReactorBudgetExhausted, TestReactorCooldownSkipped, TestReactorSequentialExecution) verified to exist on f319290 unchanged — not introduced by this commit. Co-Authored-By: Claude Opus 4.6 (1M context) --- cmd/synapbus/main.go | 13 + image-build/README.md | 73 ++++ image-build/synapbus-agent/Dockerfile | 69 +++ internal/harness/docker/config.go | 159 +++++++ internal/harness/docker/docker.go | 567 +++++++++++++++++++++++++ internal/harness/docker/docker_test.go | 219 ++++++++++ internal/harness/registry.go | 9 + internal/reactor/reactor.go | 10 +- 8 files changed, 1118 insertions(+), 1 deletion(-) create mode 100644 image-build/README.md create mode 100644 image-build/synapbus-agent/Dockerfile create mode 100644 internal/harness/docker/config.go create mode 100644 internal/harness/docker/docker.go create mode 100644 internal/harness/docker/docker_test.go diff --git a/cmd/synapbus/main.go b/cmd/synapbus/main.go index b70eac0..3a07771 100644 --- a/cmd/synapbus/main.go +++ b/cmd/synapbus/main.go @@ -48,6 +48,7 @@ import ( "github.com/synapbus/synapbus/internal/messaging" prommetrics "github.com/synapbus/synapbus/internal/metrics" "github.com/synapbus/synapbus/internal/harness" + "github.com/synapbus/synapbus/internal/harness/docker" "github.com/synapbus/synapbus/internal/harness/k8sjob" "github.com/synapbus/synapbus/internal/harness/runs" "github.com/synapbus/synapbus/internal/harness/subprocess" @@ -521,6 +522,18 @@ func runServe(cmd *cobra.Command, args []string) error { KeepWorkdirOnSuccess: keepWorkdir, }, slog.Default())) harnessRegistry.Register(webhook.New(webhook.Config{}, slog.Default())) + // Docker isolation backend — agents whose harness_config_json has a + // `docker.image` block run inside ephemeral containers. Same per-run + // workdir convention as subprocess; the workdir is bind-mounted at + // /workspace so wrappers and config files (.gemini/settings.json, + // CLAUDE.md, message.json) reach the container unchanged. The MCP + // host gets rewritten from 127.0.0.1 to host.docker.internal so the + // in-container Gemini/Claude CLI can reach the SynapBus MCP server. + harnessRegistry.Register(docker.New(docker.Config{ + BaseDir: filepath.Join(dataDir, "harness", "docker"), + KeepWorkdirOnSuccess: keepWorkdir, + HostMCPPort: port, + }, slog.Default())) harnessRunsStore := runs.New(db.DB, slog.Default()) harnessRegistry.Observer = harnessRunsStore reactorEngine.SetHarnessRegistry(harnessRegistry) diff --git a/image-build/README.md b/image-build/README.md new file mode 100644 index 0000000..30da4ad --- /dev/null +++ b/image-build/README.md @@ -0,0 +1,73 @@ +# SynapBus container images + +The `docker` harness backend (`internal/harness/docker/`) runs each agent +inside an ephemeral container. This directory holds the canonical agent +image SynapBus's bundled examples reference. + +## synapbus-agent + +The default image. Debian bookworm-slim base with: + +- `gemini` CLI (`@google/gemini-cli`) +- `claude` CLI (`@anthropic-ai/claude-code`) +- `tini` as PID 1 (signal forwarding + zombie reaping) +- Standard tooling the example wrappers use: `jq`, `sqlite3`, `curl`, `git`, `python3` +- Non-root `agent` user (uid 1000, gid 1000) matching the typical host user + +No SynapBus binary lives in the image. Agents reach the SynapBus MCP +server on the host at `host.docker.internal:` — the harness +rewrites `.gemini/settings.json` URLs from `127.0.0.1` to the gateway +hostname automatically. + +### Build + +Local single-arch: + +```bash +docker build -t synapbus-agent:latest image-build/synapbus-agent +``` + +Multi-arch via buildx (recommended for sharing the image): + +```bash +docker buildx build \ + --platform linux/amd64,linux/arm64 \ + -t synapbus-agent:latest \ + --load \ + image-build/synapbus-agent +``` + +Pin specific CLI versions with build args: + +```bash +docker build \ + --build-arg GEMINI_CLI_VERSION=0.37.1 \ + --build-arg CLAUDE_CODE_VERSION=1.0.0 \ + -t synapbus-agent:0.37.1 \ + image-build/synapbus-agent +``` + +### Wire an agent to use it + +In `harness_config_json` add a `docker` block: + +```json +{ + "gemini_md": "...", + "mcp_servers": [...], + "env": {...}, + "docker": { + "image": "synapbus-agent:latest", + "memory": "1g", + "cpus": "1.0", + "network": "bridge" + } +} +``` + +The reactor will pick the docker backend automatically when it sees the +`docker.image` field. Default security posture: `--cap-drop=ALL`, +`--security-opt=no-new-privileges`, `--read-only` root with tmpfs +`/tmp`, `--pids-limit=512`, `--user=`. Override via the +typed fields in the `docker` block (`memory`, `cpus`, `cap_add`, +`extra_mounts`, `read_only_root`, `user`). diff --git a/image-build/synapbus-agent/Dockerfile b/image-build/synapbus-agent/Dockerfile new file mode 100644 index 0000000..203bdb4 --- /dev/null +++ b/image-build/synapbus-agent/Dockerfile @@ -0,0 +1,69 @@ +# synapbus-agent: the canonical container image for SynapBus's docker +# harness backend. One image, all the agent CLIs the bundled examples +# invoke (gemini, claude — add codex/opencode here when needed). No +# SynapBus binary inside — agents reach the host's MCP server over the +# network at host.docker.internal:. +# +# Build with multi-arch buildx: +# docker buildx build \ +# --platform linux/amd64,linux/arm64 \ +# -t synapbus-agent:latest \ +# --load image-build/synapbus-agent +# +# Or for local dev (single-arch matching your host): +# docker build -t synapbus-agent:latest image-build/synapbus-agent + +FROM debian:bookworm-slim + +ARG NODE_MAJOR=22 +ARG GEMINI_CLI_VERSION=latest +ARG CLAUDE_CODE_VERSION=latest + +ENV DEBIAN_FRONTEND=noninteractive \ + LANG=C.UTF-8 \ + LC_ALL=C.UTF-8 + +# Base tooling. Most demos shell out to one of these from wrapper.sh. +RUN apt-get update && apt-get install -y --no-install-recommends \ + ca-certificates \ + curl \ + git \ + gnupg \ + jq \ + sqlite3 \ + python3 \ + python3-pip \ + tini \ + && rm -rf /var/lib/apt/lists/* + +# Node.js for the agent CLIs (gemini, claude). NodeSource keeps a +# pinned major version so the image is reproducible-ish across builds. +RUN curl -fsSL https://deb.nodesource.com/setup_${NODE_MAJOR}.x | bash - \ + && apt-get update && apt-get install -y --no-install-recommends nodejs \ + && rm -rf /var/lib/apt/lists/* \ + && npm config set update-notifier false + +# Agent CLIs. Install globally so any user inside the container can +# call them. Pinned versions are accepted via build args above. +RUN npm install -g \ + @google/gemini-cli@${GEMINI_CLI_VERSION} \ + @anthropic-ai/claude-code@${CLAUDE_CODE_VERSION} + +# Non-root user with UID/GID 1000 — matches the typical host user on +# Linux dev machines and lets `docker run --user 1000:1000` (which the +# harness sets by default) write into bind-mounted workdirs without +# permission errors. +RUN groupadd -g 1000 agent && useradd -u 1000 -g 1000 -m -s /bin/bash agent + +# Use tini as PID 1 so: +# * SIGTERM from `docker stop` reaches our wrapper.sh +# * Zombie node/python child processes get reaped properly +# Wrappers can override the entrypoint via harness_config_json.docker. +ENTRYPOINT ["/usr/bin/tini", "--"] + +# Default: run /workspace/wrapper.sh (the convention every example +# follows). Overridden via harness_config_json.docker.command when +# you want a different entry script. +WORKDIR /workspace +USER 1000:1000 +CMD ["/workspace/wrapper.sh"] diff --git a/internal/harness/docker/config.go b/internal/harness/docker/config.go new file mode 100644 index 0000000..3e8d201 --- /dev/null +++ b/internal/harness/docker/config.go @@ -0,0 +1,159 @@ +// Package docker is the container-isolation implementation of +// harness.Harness. Each Execute call materializes the agent's per-run +// workdir on the host (CLAUDE.md / GEMINI.md / .mcp.json / message.json +// — same layout as the subprocess backend), then runs an ephemeral +// `docker run --rm` with that workdir bind-mounted at /workspace. +// +// Inspired by scion's pkg/runtime/docker.go: per-task ephemeral +// containers, host-side scratch dir, secret/env injection via -e flags, +// shell-out to the docker CLI (no Docker SDK dependency, zero CGO, +// trivial cross-compile). +package docker + +import ( + "encoding/json" + "fmt" +) + +// Config tunes the docker harness at process-startup time. Per-agent +// overrides go into the agent's harness_config_json (parsed by +// ParseDockerConfig below). +type Config struct { + // BaseDir is the parent directory under which a per-run workdir is + // created on the host. The host writes config files here and bind- + // mounts the directory at /workspace inside the container. + BaseDir string + + // LogsCap bounds the number of bytes kept in ExecResult.Logs. The + // full stdout/stderr stream is written to stdout.log / stderr.log + // inside the workdir for forensics. + LogsCap int + + // KeepWorkdirOnSuccess leaves the workdir behind even for zero-exit + // runs. Useful when debugging MCP traces or container exit codes. + KeepWorkdirOnSuccess bool + + // HostGatewayName is the hostname the agent inside the container + // uses to reach the SynapBus MCP server on the host. Defaults to + // "host.docker.internal" which works on Docker Desktop (mac/win) + // natively and on Linux when --add-host=host.docker.internal: + // host-gateway is supplied (we add it automatically). + HostGatewayName string + + // HostMCPPort is the port SynapBus listens on. Used to rewrite the + // .gemini/settings.json materialized by the harness so MCP URLs + // like http://127.0.0.1:18090/mcp become + // http://host.docker.internal:18090/mcp inside the container. + // When 0 the harness leaves the URL alone (the agent prompt may + // reference an externally addressable URL already). + HostMCPPort int + + // DockerBin is the path to the docker CLI. Empty = "docker" from + // PATH. Override for podman or a wrapper script. + DockerBin string +} + +// AgentConfig is the per-agent docker block parsed from +// harness_config_json. Fields named alongside the existing subprocess +// AgentConfig so the same JSON file can carry both backends: +// +// { +// "gemini_md": "...", +// "mcp_servers": [...], +// "env": {...}, +// "docker": { +// "image": "synapbus-agent:latest", +// "memory": "1g", +// "cpus": "1.0", +// "network": "bridge", +// "extra_mounts": [{"source": "/host/path", "target": "/in/container", "read_only": true}], +// "cap_add": [], +// "extra_args": [] +// } +// } +type AgentConfig struct { + // Image is the container image to run. Required. May be a local tag + // ("synapbus-agent:latest") or a fully-qualified registry path. + Image string `json:"image"` + + // Memory is the --memory limit, e.g. "1g", "512m". Empty = no limit. + Memory string `json:"memory,omitempty"` + + // CPUs is the --cpus quota, e.g. "1.0", "0.5". Empty = no limit. + CPUs string `json:"cpus,omitempty"` + + // PIDsLimit is --pids-limit. Defaults to 512 when zero. + PIDsLimit int `json:"pids_limit,omitempty"` + + // Network is the --network mode. Empty defaults to "bridge". Use + // "none" for fully air-gapped runs. + Network string `json:"network,omitempty"` + + // ExtraMounts is a list of additional host bind-mounts. The + // per-run workdir is always mounted at /workspace; this is for + // extra read-only resources like CA bundles or shared caches. + ExtraMounts []ExtraMount `json:"extra_mounts,omitempty"` + + // CapAdd is the list of Linux capabilities to grant on top of the + // default --cap-drop=ALL. Most agents need none. + CapAdd []string `json:"cap_add,omitempty"` + + // ReadOnlyRoot makes the container's root filesystem read-only. + // Defaults to true. The harness always tmpfs-mounts /tmp so the + // agent has a writable scratch dir. + ReadOnlyRoot *bool `json:"read_only_root,omitempty"` + + // User is the --user flag value, e.g. "1000:1000". Empty leaves the + // container's default user. Set explicitly when the host has + // permission constraints on the bind-mounted workdir. + User string `json:"user,omitempty"` + + // Entrypoint overrides the image ENTRYPOINT. Empty leaves it alone. + // The harness always passes /workspace/wrapper.sh as the first arg + // after entrypoint, so the image's ENTRYPOINT must accept a script + // path (e.g. ["/usr/bin/dumb-init", "--"] then args become argv[1:]). + Entrypoint []string `json:"entrypoint,omitempty"` + + // Command overrides what the harness passes after the entrypoint. + // Defaults to ["/workspace/wrapper.sh"] — the convention every + // example in this repo follows. + Command []string `json:"command,omitempty"` + + // ExtraArgs are passed verbatim to `docker run` between the + // security flags and the image name. Use sparingly; prefer the + // typed fields above. + ExtraArgs []string `json:"extra_args,omitempty"` +} + +// ExtraMount describes one additional bind-mount. +type ExtraMount struct { + Source string `json:"source"` + Target string `json:"target"` + ReadOnly bool `json:"read_only,omitempty"` +} + +// ParseDockerConfig pulls the `docker` block out of the agent's +// harness_config_json. Tolerates an absent block (returns zero value +// + ErrNoDockerConfig so callers can decide whether to error or +// fall through). +func ParseDockerConfig(raw string) (AgentConfig, error) { + var cfg AgentConfig + if raw == "" { + return cfg, ErrNoDockerConfig + } + var envelope struct { + Docker *AgentConfig `json:"docker"` + } + if err := json.Unmarshal([]byte(raw), &envelope); err != nil { + return cfg, fmt.Errorf("docker: parse harness_config_json: %w", err) + } + if envelope.Docker == nil { + return cfg, ErrNoDockerConfig + } + return *envelope.Docker, nil +} + +// ErrNoDockerConfig signals that the agent's harness_config_json had no +// `docker` block. The Registry uses this to fall through to a different +// backend rather than failing the dispatch. +var ErrNoDockerConfig = fmt.Errorf("docker: agent has no docker config block") diff --git a/internal/harness/docker/docker.go b/internal/harness/docker/docker.go new file mode 100644 index 0000000..bd57415 --- /dev/null +++ b/internal/harness/docker/docker.go @@ -0,0 +1,567 @@ +package docker + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "log/slog" + "os" + "os/exec" + "path/filepath" + "runtime" + "sort" + "strings" + "time" + + "github.com/synapbus/synapbus/internal/agents" + "github.com/synapbus/synapbus/internal/harness" + "github.com/synapbus/synapbus/internal/harness/subprocess" + "github.com/synapbus/synapbus/internal/messaging" +) + +// Harness runs agents inside ephemeral Docker containers. One container +// per Execute call, bind-mounted workdir, --rm cleanup, no warm pool. +type Harness struct { + cfg Config + logger *slog.Logger +} + +// New builds a docker harness with sensible defaults. +func New(cfg Config, logger *slog.Logger) *Harness { + if cfg.LogsCap <= 0 { + cfg.LogsCap = 64 * 1024 + } + if cfg.HostGatewayName == "" { + cfg.HostGatewayName = "host.docker.internal" + } + if cfg.DockerBin == "" { + cfg.DockerBin = "docker" + } + if logger == nil { + logger = slog.Default() + } + return &Harness{ + cfg: cfg, + logger: logger.With("harness", "docker"), + } +} + +// Name is the registered backend identifier. Match it in agent rows via +// harness_name = "docker". +func (h *Harness) Name() string { return "docker" } + +// Capabilities mirrors subprocess: same workdir convention, same MCP +// support, same OTel env-var injection. Skills aren't materialized +// today (subprocess doesn't either). +func (h *Harness) Capabilities() harness.Capabilities { + return harness.Capabilities{ + SystemPrompt: true, + SessionResume: true, + Skills: false, + OTelNative: true, + MaxConcurrency: 4, + } +} + +// TestEnvironment runs `docker version --format {{.Server.Version}}` +// and fails fast if the daemon isn't reachable. +func (h *Harness) TestEnvironment(ctx context.Context) error { + cmd := exec.CommandContext(ctx, h.cfg.DockerBin, "version", "--format", "{{.Server.Version}}") + out, err := cmd.CombinedOutput() + if err != nil { + return fmt.Errorf("docker: daemon unreachable: %w (output: %s)", err, strings.TrimSpace(string(out))) + } + return nil +} + +// Provision is a no-op. Image pulls happen lazily on the first Execute +// (docker run will pull missing images automatically). +func (h *Harness) Provision(ctx context.Context, agent *agents.Agent) error { return nil } + +// Cancel asks the docker daemon to kill the container for runID. +// Best-effort: returns nil even if no container exists. +func (h *Harness) Cancel(ctx context.Context, runID string) error { + name := containerName(runID) + _ = exec.CommandContext(ctx, h.cfg.DockerBin, "kill", name).Run() + return nil +} + +// Execute is the hot path. Materializes the workdir, runs `docker run +// --rm` synchronously, captures exit + stdout/stderr. +func (h *Harness) Execute(ctx context.Context, req *harness.ExecRequest) (*harness.ExecResult, error) { + if req == nil { + return nil, errors.New("docker: nil ExecRequest") + } + if req.Agent == nil { + return nil, errors.New("docker: ExecRequest.Agent is required") + } + + dockerCfg, err := ParseDockerConfig(req.Agent.HarnessConfigJSON) + if err != nil { + return nil, err + } + if dockerCfg.Image == "" { + return nil, errors.New("docker: agent's harness_config_json.docker.image is required") + } + + subCfg, err := subprocess.ParseAgentConfig(req.Agent.HarnessConfigJSON) + if err != nil { + return nil, err + } + + workdir, err := h.makeWorkdir(req.RunID) + if err != nil { + return nil, err + } + + if req.Message != nil { + if err := writeMessageFile(workdir, req.Message); err != nil { + return nil, err + } + } + + if err := subprocess.MaterialiseAgentConfig(workdir, subCfg); err != nil { + return nil, err + } + + // Rewrite MCP host in .gemini/settings.json so the agent inside the + // container can reach SynapBus on the host. The host writes + // 127.0.0.1: by default; the container sees that loopback as + // itself, not the host. + if h.cfg.HostMCPPort > 0 { + if err := rewriteGeminiMCPHost(workdir, h.cfg.HostGatewayName, h.cfg.HostMCPPort); err != nil { + h.logger.Warn("rewrite gemini MCP host failed", + "workdir", workdir, "error", err) + } + } + + runCtx := ctx + if req.Budget.MaxWallClock > 0 { + var cancel context.CancelFunc + runCtx, cancel = context.WithTimeout(ctx, req.Budget.MaxWallClock) + defer cancel() + } + + args, err := h.buildRunArgs(req, dockerCfg, subCfg, workdir) + if err != nil { + return nil, err + } + + h.logger.Info("docker launching", + "run_id", req.RunID, + "agent", req.AgentName, + "image", dockerCfg.Image, + "workdir", workdir, + "network", argOr(dockerCfg.Network, "bridge"), + ) + + cmd := exec.CommandContext(runCtx, h.cfg.DockerBin, args...) + var stdout, stderr bytes.Buffer + cmd.Stdout = io.MultiWriter(&stdout, fileWriter(workdir, "stdout.log")) + cmd.Stderr = io.MultiWriter(&stderr, fileWriter(workdir, "stderr.log")) + + startedAt := time.Now() + runErr := cmd.Run() + duration := time.Since(startedAt) + + exitCode := 0 + if runErr != nil { + var exitErr *exec.ExitError + if errors.As(runErr, &exitErr) { + exitCode = exitErr.ExitCode() + } else { + exitCode = 1 + } + } + + var resultJSON json.RawMessage + if raw, readErr := os.ReadFile(filepath.Join(workdir, "result.json")); readErr == nil && len(raw) > 0 { + if json.Valid(raw) { + resultJSON = raw + } + } + + promptText := readFileSafe(filepath.Join(workdir, "prompt.txt")) + responseText := readFileSafe(filepath.Join(workdir, "response.txt")) + logs := mergeLogs(&stdout, &stderr, h.cfg.LogsCap) + + if exitCode == 0 && !h.cfg.KeepWorkdirOnSuccess { + _ = os.RemoveAll(workdir) + } + + h.logger.Info("docker finished", + "run_id", req.RunID, + "agent", req.AgentName, + "exit", exitCode, + "duration_ms", duration.Milliseconds(), + ) + + result := &harness.ExecResult{ + ExitCode: exitCode, + Logs: logs, + ResultJSON: resultJSON, + Prompt: promptText, + Response: responseText, + } + + if runErr != nil && errors.Is(runCtx.Err(), context.DeadlineExceeded) { + _ = h.Cancel(context.Background(), req.RunID) + return result, fmt.Errorf("docker: wall-clock budget %s exceeded", req.Budget.MaxWallClock) + } + if runErr != nil && errors.Is(runCtx.Err(), context.Canceled) { + _ = h.Cancel(context.Background(), req.RunID) + return result, runCtx.Err() + } + return result, nil +} + +// buildRunArgs constructs the full `docker run` argv. Defaults are +// security-conservative: --rm, --cap-drop=ALL, no-new-privileges, pids +// limit, read-only root with tmpfs /tmp, no privilege escalation, no +// host networking. Per-agent config layers on top. +func (h *Harness) buildRunArgs( + req *harness.ExecRequest, + dockerCfg AgentConfig, + subCfg subprocess.AgentConfig, + workdir string, +) ([]string, error) { + name := containerName(req.RunID) + args := []string{ + "run", + "--rm", + "--name", name, + "--workdir", "/workspace", + "--mount", fmt.Sprintf("type=bind,source=%s,target=/workspace", workdir), + "--security-opt", "no-new-privileges", + "--cap-drop", "ALL", + "--pids-limit", fmt.Sprintf("%d", pidsLimit(dockerCfg.PIDsLimit)), + } + + // Read-only root + tmpfs scratch unless explicitly disabled. + if dockerCfg.ReadOnlyRoot == nil || *dockerCfg.ReadOnlyRoot { + args = append(args, "--read-only", "--tmpfs", "/tmp:rw,size=64m") + } + + if dockerCfg.Memory != "" { + args = append(args, "--memory", dockerCfg.Memory) + // Match memory-swap to memory so swap doesn't silently double + // the effective limit. -1 would mean unlimited; equal disables. + args = append(args, "--memory-swap", dockerCfg.Memory) + } + if dockerCfg.CPUs != "" { + args = append(args, "--cpus", dockerCfg.CPUs) + } + + network := dockerCfg.Network + if network == "" { + network = "bridge" + } + args = append(args, "--network", network) + + // Add host.docker.internal pointer on Linux so the agent can reach + // the SynapBus MCP server at the same hostname as on Docker Desktop. + // Skipped for --network=host (not needed) and --network=none + // (would fail the gateway lookup). + if network != "host" && network != "none" && runtime.GOOS == "linux" { + args = append(args, "--add-host", h.cfg.HostGatewayName+":host-gateway") + } + + // User namespacing — host UID/GID injection so files written into + // the bind-mounted workdir end up owned by the SynapBus user. If + // the agent overrides User explicitly use that. + if dockerCfg.User != "" { + args = append(args, "--user", dockerCfg.User) + } else { + args = append(args, "--user", currentUserSpec()) + } + + for _, c := range dockerCfg.CapAdd { + if c == "" { + continue + } + args = append(args, "--cap-add", c) + } + + for _, m := range dockerCfg.ExtraMounts { + if m.Source == "" || m.Target == "" { + continue + } + spec := fmt.Sprintf("type=bind,source=%s,target=%s", m.Source, m.Target) + if m.ReadOnly { + spec += ",readonly" + } + args = append(args, "--mount", spec) + } + + // Environment variables — caller-provided + harness-injected. Pass + // through as -e KEY=VALUE; sort for deterministic output. + envMap := buildEnvMap(req, subCfg) + keys := make([]string, 0, len(envMap)) + for k := range envMap { + keys = append(keys, k) + } + sort.Strings(keys) + for _, k := range keys { + args = append(args, "--env", k+"="+envMap[k]) + } + + args = append(args, dockerCfg.ExtraArgs...) + + if len(dockerCfg.Entrypoint) > 0 { + args = append(args, "--entrypoint", dockerCfg.Entrypoint[0]) + } + + args = append(args, dockerCfg.Image) + + // Args after the image become the container's CMD. If Entrypoint is + // set, prepend the rest of its argv first, then the Command. + if len(dockerCfg.Entrypoint) > 1 { + args = append(args, dockerCfg.Entrypoint[1:]...) + } + if len(dockerCfg.Command) > 0 { + args = append(args, dockerCfg.Command...) + } else { + args = append(args, "/workspace/wrapper.sh") + } + + return args, nil +} + +// makeWorkdir creates BaseDir// with mode 0755. The +// directory is removed on successful completion unless KeepWorkdirOn +// Success is set. +func (h *Harness) makeWorkdir(runID string) (string, error) { + base := h.cfg.BaseDir + if base == "" { + base = filepath.Join(os.TempDir(), "synapbus-docker") + } + if err := os.MkdirAll(base, 0o755); err != nil { + return "", fmt.Errorf("docker: mkdir base: %w", err) + } + name := sanitizeRunDir(runID) + if name == "" { + name = fmt.Sprintf("run-%d", time.Now().UnixNano()) + } + wd := filepath.Join(base, name) + if err := os.MkdirAll(wd, 0o755); err != nil { + return "", fmt.Errorf("docker: mkdir workdir: %w", err) + } + return wd, nil +} + +// containerName turns a run id into a docker-safe container name. +func containerName(runID string) string { + clean := sanitizeRunDir(runID) + if clean == "" { + clean = fmt.Sprintf("run%d", time.Now().UnixNano()) + } + return "synapbus-" + clean +} + +func sanitizeRunDir(runID string) string { + if runID == "" { + return "" + } + var b strings.Builder + for _, r := range runID { + switch { + case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z', r >= '0' && r <= '9': + b.WriteRune(r) + case r == '-' || r == '_': + b.WriteRune(r) + default: + b.WriteByte('-') + } + } + out := b.String() + if len(out) > 64 { + out = out[:64] + } + return out +} + +func writeMessageFile(workdir string, msg *messaging.Message) error { + raw, err := json.Marshal(msg) + if err != nil { + return fmt.Errorf("docker: marshal message: %w", err) + } + return os.WriteFile(filepath.Join(workdir, "message.json"), raw, 0o644) +} + +// rewriteGeminiMCPHost reads .gemini/settings.json (written by the +// subprocess MaterialiseAgentConfig step) and rewrites every mcpServers +// URL whose host is 127.0.0.1 / localhost / 0.0.0.0 to the host +// gateway. The container can't reach the host on loopback. +func rewriteGeminiMCPHost(workdir, gateway string, port int) error { + path := filepath.Join(workdir, ".gemini", "settings.json") + raw, err := os.ReadFile(path) + if err != nil { + if errors.Is(err, os.ErrNotExist) { + return nil + } + return err + } + var settings struct { + MCPServers map[string]map[string]json.RawMessage `json:"mcpServers"` + } + if err := json.Unmarshal(raw, &settings); err != nil { + return fmt.Errorf("docker: parse gemini settings: %w", err) + } + if len(settings.MCPServers) == 0 { + return nil + } + dirty := false + for name, server := range settings.MCPServers { + urlRaw, ok := server["url"] + if !ok { + continue + } + var url string + if err := json.Unmarshal(urlRaw, &url); err != nil { + continue + } + newURL := rewriteLoopback(url, gateway, port) + if newURL == url { + continue + } + fixed, _ := json.Marshal(newURL) + settings.MCPServers[name]["url"] = fixed + dirty = true + } + if !dirty { + return nil + } + out, err := json.MarshalIndent(settings, "", " ") + if err != nil { + return err + } + return os.WriteFile(path, out, 0o644) +} + +func rewriteLoopback(url, gateway string, port int) string { + for _, host := range []string{"127.0.0.1", "localhost", "0.0.0.0"} { + needle := "//" + host + if i := strings.Index(url, needle); i >= 0 { + rest := url[i+len(needle):] + // Replace the port on the URL with the harness-known host + // port so a misconfigured agent cant accidentally point at + // a different listener. + if strings.HasPrefix(rest, ":") { + if slash := strings.IndexByte(rest, '/'); slash >= 0 { + rest = rest[slash:] + } else { + rest = "" + } + } + return url[:i] + "//" + gateway + ":" + fmt.Sprintf("%d", port) + rest + } + } + return url +} + +// buildEnvMap mirrors subprocess.buildEnv but does NOT inherit the +// parent process's environment. Containers start clean — only what we +// explicitly forward gets in. Order: agent k8s_env_json → harness +// config env → caller overrides → SYNAPBUS_* run context. +func buildEnvMap(req *harness.ExecRequest, cfg subprocess.AgentConfig) map[string]string { + env := map[string]string{} + + if req.Agent != nil && req.Agent.K8sEnvJSON != "" { + var m map[string]json.RawMessage + if err := json.Unmarshal([]byte(req.Agent.K8sEnvJSON), &m); err == nil { + for k, v := range m { + var s string + if err := json.Unmarshal(v, &s); err == nil { + env[k] = s + continue + } + env[k] = strings.Trim(string(v), "\"") + } + } + } + + for k, v := range cfg.Env { + env[k] = v + } + + for k, v := range req.Env { + env[k] = v + } + + env["SYNAPBUS_RUN_ID"] = req.RunID + env["SYNAPBUS_AGENT"] = req.AgentName + env["SYNAPBUS_WORKDIR"] = "/workspace" + if req.Message != nil { + env["SYNAPBUS_MESSAGE_ID"] = fmt.Sprintf("%d", req.Message.ID) + env["SYNAPBUS_FROM_AGENT"] = req.Message.FromAgent + } + return env +} + +// currentUserSpec returns "uid:gid" for the host user so files written +// inside the bind-mount land with sane ownership instead of root. Only +// meaningful on Linux; on Docker Desktop (mac) the bind-mount layer +// handles ownership translation transparently but passing --user is +// still a defence-in-depth measure. +func currentUserSpec() string { + uid := os.Getuid() + gid := os.Getgid() + if uid <= 0 { + // fall through to image default if we can't discover (e.g. + // running on Windows or under a daemon without /proc). + return "" + } + return fmt.Sprintf("%d:%d", uid, gid) +} + +func argOr(s, fallback string) string { + if s == "" { + return fallback + } + return s +} + +func pidsLimit(cfg int) int { + if cfg <= 0 { + return 512 + } + return cfg +} + +func fileWriter(workdir, name string) io.Writer { + f, err := os.OpenFile(filepath.Join(workdir, name), os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o644) + if err != nil { + return io.Discard + } + return f +} + +func readFileSafe(path string) string { + raw, err := os.ReadFile(path) + if err != nil { + return "" + } + return string(raw) +} + +func mergeLogs(out, errb *bytes.Buffer, cap int) string { + var b strings.Builder + if out.Len() > 0 { + b.WriteString(out.String()) + } + if errb.Len() > 0 { + if b.Len() > 0 { + b.WriteString("\n") + } + b.WriteString("-- stderr --\n") + b.WriteString(errb.String()) + } + s := b.String() + if cap > 0 && len(s) > cap { + s = "... [truncated " + fmt.Sprintf("%d", len(s)-cap) + " bytes] ...\n" + s[len(s)-cap:] + } + return s +} diff --git a/internal/harness/docker/docker_test.go b/internal/harness/docker/docker_test.go new file mode 100644 index 0000000..4fe151d --- /dev/null +++ b/internal/harness/docker/docker_test.go @@ -0,0 +1,219 @@ +package docker_test + +import ( + "context" + "encoding/json" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/synapbus/synapbus/internal/agents" + "github.com/synapbus/synapbus/internal/harness" + "github.com/synapbus/synapbus/internal/harness/docker" + "github.com/synapbus/synapbus/internal/messaging" +) + +// TestExecute_Hello is the smoke test for the docker backend. Skipped +// when the docker daemon isn't reachable so CI without docker won't +// fail. Builds nothing — uses `alpine:3.20` which is small and +// universally available. +func TestExecute_Hello(t *testing.T) { + if err := exec.Command("docker", "version", "--format", "{{.Server.Version}}").Run(); err != nil { + t.Skip("docker daemon not available, skipping") + } + + base := t.TempDir() + h := docker.New(docker.Config{ + BaseDir: base, + KeepWorkdirOnSuccess: true, + HostMCPPort: 0, // skip URL rewrite for this test + }, nil) + + // Minimal agent with a docker block. wrapper.sh writes a marker + // file to /workspace and prints a known string so we can assert + // both bind-mount writeback and stdout capture. + cfgJSON, err := json.Marshal(map[string]any{ + "env": map[string]string{ + "GREETING": "hello-from-container", + }, + "docker": map[string]any{ + "image": "alpine:3.20", + "command": []string{"sh", "/workspace/wrapper.sh"}, + "network": "none", // air-gapped; we don't need MCP for this test + "memory": "128m", + "cpus": "0.5", + "pids_limit": 32, + }, + }) + if err != nil { + t.Fatal(err) + } + + agent := &agents.Agent{ + ID: 1, + Name: "smoke-test-agent", + HarnessConfigJSON: string(cfgJSON), + } + + req := &harness.ExecRequest{ + RunID: "smoke-test-1", + AgentName: agent.Name, + Agent: agent, + Message: &messaging.Message{ + ID: 42, + FromAgent: "tester", + ToAgent: agent.Name, + Body: "hello", + }, + Budget: harness.Budget{ + MaxWallClock: 60 * time.Second, + }, + } + + // We need a wrapper.sh staged BEFORE Execute creates the container. + // In production the harness materializes one via the gemini_md / + // claude_md fields, but for the test we drop a tiny shell script + // directly into the run workdir. + runDir := filepath.Join(base, "smoke-test-1") + if err := os.MkdirAll(runDir, 0o755); err != nil { + t.Fatal(err) + } + wrapper := `#!/bin/sh +set -eu +echo "wrapper running as uid=$(id -u) gid=$(id -g) cwd=$(pwd)" +echo "GREETING=$GREETING" +echo "SYNAPBUS_RUN_ID=$SYNAPBUS_RUN_ID" +echo "SYNAPBUS_FROM_AGENT=$SYNAPBUS_FROM_AGENT" +[ -f /workspace/message.json ] && echo "message.json present" +echo '{"ok":true,"phase":"smoke"}' > /workspace/result.json +echo "wrapper done" +` + if err := os.WriteFile(filepath.Join(runDir, "wrapper.sh"), []byte(wrapper), 0o755); err != nil { + t.Fatal(err) + } + + res, err := h.Execute(context.Background(), req) + if err != nil { + t.Fatalf("Execute returned error: %v\nlogs:\n%s", err, func() string { + if res != nil { + return res.Logs + } + return "" + }()) + } + if res == nil { + t.Fatal("nil result") + } + if res.ExitCode != 0 { + t.Fatalf("expected exit 0, got %d\nlogs:\n%s", res.ExitCode, res.Logs) + } + + // Check stdout capture. + if !strings.Contains(res.Logs, "wrapper done") { + t.Errorf("stdout missing 'wrapper done':\n%s", res.Logs) + } + if !strings.Contains(res.Logs, "GREETING=hello-from-container") { + t.Errorf("env injection missing GREETING:\n%s", res.Logs) + } + if !strings.Contains(res.Logs, "SYNAPBUS_RUN_ID=smoke-test-1") { + t.Errorf("SYNAPBUS_RUN_ID not propagated:\n%s", res.Logs) + } + if !strings.Contains(res.Logs, "message.json present") { + t.Errorf("message.json not bind-mounted:\n%s", res.Logs) + } + + // Check result.json bind-mount writeback. + if len(res.ResultJSON) == 0 { + t.Error("result.json was not captured from bind-mount") + } else { + var parsed map[string]any + if err := json.Unmarshal(res.ResultJSON, &parsed); err != nil { + t.Errorf("result.json invalid: %v", err) + } else if parsed["ok"] != true { + t.Errorf("result.json content unexpected: %v", parsed) + } + } + + // Workdir should still exist (KeepWorkdirOnSuccess=true). Verify + // the result.json the container wrote actually landed on the host. + hostResult, err := os.ReadFile(filepath.Join(runDir, "result.json")) + if err != nil { + t.Errorf("result.json not on host post-run: %v", err) + } else if !strings.Contains(string(hostResult), "smoke") { + t.Errorf("host result.json content unexpected: %s", hostResult) + } +} + +// TestExecute_NoImage verifies the backend rejects agents whose +// harness_config_json lacks docker.image rather than silently picking +// some default. +func TestExecute_NoImage(t *testing.T) { + h := docker.New(docker.Config{}, nil) + agent := &agents.Agent{ + Name: "no-image", + HarnessConfigJSON: `{"docker":{}}`, + } + req := &harness.ExecRequest{ + RunID: "noimg", + AgentName: agent.Name, + Agent: agent, + } + _, err := h.Execute(context.Background(), req) + if err == nil { + t.Fatal("expected error for missing image, got nil") + } + if !strings.Contains(err.Error(), "image is required") { + t.Errorf("unexpected error: %v", err) + } +} + +// TestExecute_TimeoutCancel verifies the wall-clock budget kills a +// long-running container. +func TestExecute_TimeoutCancel(t *testing.T) { + if err := exec.Command("docker", "version", "--format", "{{.Server.Version}}").Run(); err != nil { + t.Skip("docker daemon not available, skipping") + } + + base := t.TempDir() + h := docker.New(docker.Config{BaseDir: base, KeepWorkdirOnSuccess: true}, nil) + + cfgJSON, _ := json.Marshal(map[string]any{ + "docker": map[string]any{ + "image": "alpine:3.20", + "command": []string{"sh", "/workspace/wrapper.sh"}, + "network": "none", + "memory": "64m", + }, + }) + agent := &agents.Agent{ + Name: "slow", + HarnessConfigJSON: string(cfgJSON), + } + runDir := filepath.Join(base, "slow") + _ = os.MkdirAll(runDir, 0o755) + _ = os.WriteFile(filepath.Join(runDir, "wrapper.sh"), + []byte("#!/bin/sh\nsleep 30\n"), 0o755) + + req := &harness.ExecRequest{ + RunID: "slow", + AgentName: agent.Name, + Agent: agent, + Budget: harness.Budget{MaxWallClock: 2 * time.Second}, + } + start := time.Now() + res, err := h.Execute(context.Background(), req) + elapsed := time.Since(start) + + if err == nil { + t.Fatalf("expected timeout error, got nil (exit=%d)", res.ExitCode) + } + if !strings.Contains(err.Error(), "budget") { + t.Errorf("unexpected error: %v", err) + } + if elapsed > 10*time.Second { + t.Errorf("timeout took too long: %s", elapsed) + } +} diff --git a/internal/harness/registry.go b/internal/harness/registry.go index 2c21e81..899c639 100644 --- a/internal/harness/registry.go +++ b/internal/harness/registry.go @@ -106,6 +106,15 @@ func (r *Registry) Resolve(agent *agents.Agent) (Harness, error) { return h, nil } } + // Docker block in harness_config_json takes precedence over a plain + // local_command — same precedence the reactor uses, so explicit + // isolation never silently downgrades to subprocess. + if agent != nil && agent.HarnessConfigJSON != "" && + strings.Contains(agent.HarnessConfigJSON, "\"docker\"") { + if h, ok := r.byName["docker"]; ok { + return h, nil + } + } if agent != nil && agent.LocalCommand != "" { if h, ok := r.byName["subprocess"]; ok { return h, nil diff --git a/internal/reactor/reactor.go b/internal/reactor/reactor.go index 2c3382e..ba03915 100644 --- a/internal/reactor/reactor.go +++ b/internal/reactor/reactor.go @@ -27,6 +27,7 @@ import ( const ( backendK8s = "k8s" backendSubprocess = "subprocess" + backendDocker = "docker" backendWebhook = "webhook" backendNone = "" ) @@ -197,7 +198,7 @@ func (r *Reactor) evaluateTrigger(ctx context.Context, agentName string, event d r.recordSkippedRun(ctx, agentName, event, StatusFailed, "K8s runner not available") return nil } - case backendSubprocess, backendWebhook: + case backendSubprocess, backendDocker, backendWebhook: // 4. Harness-specific precondition: registry must be wired and // the resolver must pick the backend we think we should get. if r.registry == nil { @@ -286,6 +287,8 @@ func (r *Reactor) agentBackendKind(agent *agents.Agent) string { return backendK8s case "subprocess": return backendSubprocess + case "docker": + return backendDocker case "webhook": return backendWebhook } @@ -293,6 +296,11 @@ func (r *Reactor) agentBackendKind(agent *agents.Agent) string { if agent.K8sImage != "" { return backendK8s } + // Docker block in harness_config_json wins over local_command — an + // agent that defines both is asking to run isolated. + if agent.HarnessConfigJSON != "" && strings.Contains(agent.HarnessConfigJSON, "\"docker\"") { + return backendDocker + } if agent.LocalCommand != "" { return backendSubprocess }