feat(harness): docker isolation backend + canonical synapbus-agent image

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=<host uid:gid>
  --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:<port>/mcp are rewritten to
http://host.docker.internal:<port>/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:<port>.

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) <noreply@anthropic.com>
This commit is contained in:
Algis Dumbris
2026-04-15 08:59:37 +03:00
co-authored by Claude Opus 4.6
parent f319290ef9
commit 560d9d4125
8 changed files with 1118 additions and 1 deletions
+13
View File
@@ -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)
+73
View File
@@ -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:<port>` — 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=<host uid:gid>`. Override via the
typed fields in the `docker` block (`memory`, `cpus`, `cap_add`,
`extra_mounts`, `read_only_root`, `user`).
+69
View File
@@ -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:<port>.
#
# 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"]
+159
View File
@@ -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")
+567
View File
@@ -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:<port> 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/<runID-sanitised>/ 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
}
+219
View File
@@ -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)
}
}
+9
View File
@@ -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
+9 -1
View File
@@ -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
}