// Package compose executes Docker Compose through a fixed executable and argument arrays. package compose import ( "context" "crypto/rand" "encoding/hex" "errors" "fmt" "io" "os" "os/exec" "strings" "sync" "time" "github.com/aritmolab/thothii/tools/tht/internal/config" ) const defaultCaptureBytes = 4 * 1024 * 1024 const maximumCaptureBytes = 64 * 1024 * 1024 const cleanupCaptureBytes = 4 * 1024 const containerCleanupBound = 2 * time.Second // ErrOutputLimit reports that a child exceeded one of its capture limits. var ErrOutputLimit = errors.New("Docker command output limit exceeded") // ErrProcessReap reports that a cancelled Docker CLI could not be reaped in its final bound. var ErrProcessReap = errors.New("Docker command process could not be reaped") // ErrContainerCleanup reports that Docker could not confirm removal of a named one-shot container. var ErrContainerCleanup = errors.New("Docker one-shot container cleanup failed") // CaptureLimits bounds each captured stream while the child is running. type CaptureLimits struct { StdoutBytes int StderrBytes int } // Result is the captured output and process exit code for one Docker invocation. type Result struct { Stdout string Stderr string ExitCode int } // Runner is the shell-free Docker command boundary used by the host CLI. type Runner interface { Run(context.Context, []string, io.Reader) (Result, error) } type boundedRunner interface { RunBounded(context.Context, []string, io.Reader, CaptureLimits) (Result, error) } type boundedStreamingRunner interface { RunBoundedStreaming(context.Context, []string, io.Reader, CaptureLimits, func([]byte)) (Result, error) } // execRunner executes the Docker CLI. It never invokes a shell. type execRunner struct { binary string environment []string } // NewRunnerWithEnvironment uses an explicit process environment, for document-owned Compose // configuration that must not inherit unrelated operator shell parameter overrides. func NewRunnerWithEnvironment(binary string, environment []string) Runner { if binary == "" { binary = "docker" } return execRunner{binary: binary, environment: append([]string{}, environment...)} } // NewRunner returns a runner for binary. An empty binary selects docker from PATH. func NewRunner(binary string) Runner { if binary == "" { binary = "docker" } return execRunner{binary: binary} } // Run invokes Docker with the supplied argument array and optional standard input. func (r execRunner) Run(ctx context.Context, args []string, stdin io.Reader) (Result, error) { return r.RunBounded(ctx, args, stdin, CaptureLimits{ StdoutBytes: defaultCaptureBytes, StderrBytes: defaultCaptureBytes, }) } // RunBounded invokes Docker while enforcing both stream limits during capture. func (r execRunner) RunBounded(ctx context.Context, args []string, stdin io.Reader, limits CaptureLimits) (Result, error) { return r.runBounded(ctx, args, stdin, limits, true, nil) } // RunBoundedStreaming retains bounded stderr capture while synchronously observing only the // retained prefix. It is used for the one validated interactive device-flow prompt. func (r execRunner) RunBoundedStreaming(ctx context.Context, args []string, stdin io.Reader, limits CaptureLimits, stderrObserver func([]byte)) (Result, error) { return r.runBounded(ctx, args, stdin, limits, true, stderrObserver) } func (r execRunner) runBounded(ctx context.Context, args []string, stdin io.Reader, limits CaptureLimits, manageOneShot bool, stderrObserver func([]byte)) (Result, error) { if err := validCaptureLimits(limits); err != nil { return Result{}, err } if err := ctx.Err(); err != nil { return Result{}, err } preparedArgs := append([]string(nil), args...) containerName := "" if manageOneShot { var err error preparedArgs, containerName, err = prepareOneShot(preparedArgs) if err != nil { return Result{}, err } } command := exec.Command(r.binary, preparedArgs...) if r.environment != nil { command.Env = r.environment } configureProcess(command) command.Stdin = stdin overflow := make(chan struct{}, 1) stdout := newCappedBuffer(limits.StdoutBytes, overflow, nil) stderr := newCappedBuffer(limits.StderrBytes, overflow, stderrObserver) command.Stdout = stdout command.Stderr = stderr process, err := startProcessForRunner(command) if err != nil { return startFailure(err) } done := make(chan error, 1) go func() { done <- command.Wait() }() var processErr error var lifecycleErr error interrupted := false select { case processErr = <-done: case <-ctx.Done(): interrupted = true processErr = ctx.Err() if err := terminateProcessForRunner(process, done); err != nil { lifecycleErr = errors.Join(lifecycleErr, err) } case <-overflow: interrupted = true processErr = ErrOutputLimit if err := terminateProcessForRunner(process, done); err != nil { lifecycleErr = errors.Join(lifecycleErr, err) } } if err := process.close(); err != nil { lifecycleErr = errors.Join(lifecycleErr, ErrProcessReap, err) } if containerName != "" { if err := r.cleanupOneShotContainer(containerName); err != nil { lifecycleErr = errors.Join(lifecycleErr, ErrContainerCleanup) } } result := Result{Stdout: stdout.String(), Stderr: stderr.String()} if command.ProcessState != nil { result.ExitCode = command.ProcessState.ExitCode() } if stdout.Overflowed() || stderr.Overflowed() { return result, joinLifecycleError(ErrOutputLimit, lifecycleErr) } if interrupted { return result, joinLifecycleError(processErr, lifecycleErr) } if processErr == nil { return result, lifecycleErr } var exitError *exec.ExitError if errors.As(processErr, &exitError) { result.ExitCode = exitError.ExitCode() return result, joinLifecycleError(processErr, lifecycleErr) } return result, joinLifecycleError(processErr, lifecycleErr) } func joinLifecycleError(processErr, lifecycleErr error) error { if lifecycleErr == nil { return processErr } return errors.Join(processErr, lifecycleErr) } // NewOneShotContainerName returns a Docker-safe, cross-process unique name with a bounded prefix. func NewOneShotContainerName(prefix string) (string, error) { if len(prefix) < 1 || len(prefix) > 96 || !safeContainerName(prefix) { return "", errors.New("Docker one-shot container name is invalid") } random := make([]byte, 12) if _, err := rand.Read(random); err != nil { return "", errors.New("Docker one-shot container name is unavailable") } return prefix + "-" + hex.EncodeToString(random), nil } func prepareOneShot(args []string) ([]string, string, error) { for runIndex, arg := range args { if arg != "run" || !(runIndex == 0 || args[0] == "compose") { continue } _, name, valid := oneShotRunOptions(args, runIndex) if !valid { continue } if name != "" { return args, name, nil } return args, "", nil } return args, "", nil } func oneShotRunOptions(args []string, runIndex int) (rmIndex int, name string, valid bool) { rmIndex = -1 valueOptions := map[string]struct{}{ "--name": {}, "--entrypoint": {}, "--network": {}, "--mount": {}, "--env": {}, "-e": {}, "--user": {}, "-u": {}, "--volume": {}, "-v": {}, "--workdir": {}, "-w": {}, "--label": {}, "-l": {}, "--pull": {}, "--cap-add": {}, "--cap-drop": {}, } for index := runIndex + 1; index < len(args); index++ { arg := args[index] if arg == "--rm" { rmIndex = index continue } if _, expectsValue := valueOptions[arg]; expectsValue { if index+1 >= len(args) { return -1, "", false } if arg == "--name" { if !safeContainerName(args[index+1]) { return -1, "", false } name = args[index+1] } index++ continue } if strings.HasPrefix(arg, "--name=") { name = strings.TrimPrefix(arg, "--name=") if !safeContainerName(name) { return -1, "", false } continue } if strings.HasPrefix(arg, "-") { continue } return rmIndex, name, rmIndex >= 0 } return -1, "", false } func safeContainerName(value string) bool { if len(value) < 1 || len(value) > 128 { return false } for index, character := range value { letter := character >= 'a' && character <= 'z' || character >= 'A' && character <= 'Z' digit := character >= '0' && character <= '9' if letter || digit || index > 0 && (character == '-' || character == '_' || character == '.') { continue } return false } return true } func (r execRunner) cleanupOneShotContainer(name string) error { if !safeContainerName(name) { return ErrContainerCleanup } ctx, cancel := context.WithTimeout(context.Background(), containerCleanupBound) defer cancel() exists, err := r.oneShotContainerExists(ctx, name) if err != nil || !exists { if err != nil { return ErrContainerCleanup } return nil } limits := CaptureLimits{StdoutBytes: cleanupCaptureBytes, StderrBytes: cleanupCaptureBytes} if _, err := r.runBounded(ctx, []string{"container", "rm", "-f", name}, nil, limits, false, nil); err != nil { return ErrContainerCleanup } exists, err = r.oneShotContainerExists(ctx, name) if err != nil || exists { return ErrContainerCleanup } return nil } func (r execRunner) oneShotContainerExists(ctx context.Context, name string) (bool, error) { limits := CaptureLimits{StdoutBytes: cleanupCaptureBytes, StderrBytes: cleanupCaptureBytes} result, err := r.runBounded(ctx, []string{ "container", "ls", "--all", "--quiet", "--filter", "name=^/" + name + "$", }, nil, limits, false, nil) if err != nil { return false, ErrContainerCleanup } return strings.TrimSpace(result.Stdout) != "", nil } // RunBounded uses the production runner's during-capture limits while retaining compatibility // with injected runners, whose already-bounded test results are checked before use. func RunBounded(runner Runner, ctx context.Context, args []string, stdin io.Reader, limits CaptureLimits) (Result, error) { if err := validCaptureLimits(limits); err != nil { return Result{}, err } if bounded, ok := runner.(boundedRunner); ok { return bounded.RunBounded(ctx, args, stdin, limits) } result, err := runner.Run(ctx, args, stdin) if len(result.Stdout) > limits.StdoutBytes || len(result.Stderr) > limits.StderrBytes { return Result{ Stdout: boundedString(result.Stdout, limits.StdoutBytes), Stderr: boundedString(result.Stderr, limits.StderrBytes), ExitCode: result.ExitCode, }, ErrOutputLimit } return result, err } // RunBoundedStreaming uses the production runner's during-capture observer. Compatibility // runners are observed only after their already-bounded result returns. func RunBoundedStreaming(runner Runner, ctx context.Context, args []string, stdin io.Reader, limits CaptureLimits, stderrObserver func([]byte)) (Result, error) { if err := validCaptureLimits(limits); err != nil { return Result{}, err } if streaming, ok := runner.(boundedStreamingRunner); ok { return streaming.RunBoundedStreaming(ctx, args, stdin, limits, stderrObserver) } result, err := RunBounded(runner, ctx, args, stdin, limits) if stderrObserver != nil && result.Stderr != "" { stderrObserver([]byte(result.Stderr)) } return result, err } func validCaptureLimits(limits CaptureLimits) error { if limits.StdoutBytes < 1 || limits.StderrBytes < 1 || limits.StdoutBytes > maximumCaptureBytes || limits.StderrBytes > maximumCaptureBytes { return errors.New("Docker command capture limits are invalid") } return nil } func boundedString(value string, maximum int) string { if len(value) <= maximum { return value } return value[:maximum] } func startFailure(err error) (Result, error) { result := Result{} if errors.Is(err, exec.ErrNotFound) || errors.Is(err, os.ErrNotExist) { result.ExitCode = 127 return result, fmt.Errorf("%w: %w", exec.ErrNotFound, err) } return result, err } type cappedBuffer struct { mu sync.Mutex contents []byte maximum int overflow chan<- struct{} observer func([]byte) exceeded bool } func newCappedBuffer(maximum int, overflow chan<- struct{}, observer func([]byte)) *cappedBuffer { capacity := maximum if capacity > 4096 { capacity = 4096 } return &cappedBuffer{contents: make([]byte, 0, capacity), maximum: maximum, overflow: overflow, observer: observer} } func (b *cappedBuffer) Write(value []byte) (int, error) { b.mu.Lock() remaining := b.maximum - len(b.contents) kept := 0 if remaining > 0 { kept = len(value) if kept > remaining { kept = remaining } b.contents = append(b.contents, value[:kept]...) } if len(value) > remaining && !b.exceeded { b.exceeded = true select { case b.overflow <- struct{}{}: default: } } var observed []byte if kept > 0 && b.observer != nil { observed = append([]byte(nil), value[:kept]...) } b.mu.Unlock() if len(observed) > 0 { b.observer(observed) } return len(value), nil } func (b *cappedBuffer) String() string { b.mu.Lock() defer b.mu.Unlock() return string(b.contents) } func (b *cappedBuffer) Overflowed() bool { b.mu.Lock() defer b.mu.Unlock() return b.exceeded } // InstallationRunner applies an installation's validated Compose arguments to commands that // explicitly start with compose. Direct Docker image commands remain host-side. type InstallationRunner struct { Installation config.Installation Runner Runner } // SessionInventoryScope supplies the installation's intended visibility to lifecycle callers. func (r InstallationRunner) SessionInventoryScope() string { if r.Installation.Profile == "local" { return "mine" } return "all" } // Run transforms only Compose invocations. It never performs shell interpolation. func (r InstallationRunner) Run(ctx context.Context, args []string, stdin io.Reader) (Result, error) { if len(args) > 0 && args[0] == "compose" { return r.Runner.Run(ctx, r.Installation.ComposeArgs(args[1:]...), stdin) } return r.Runner.Run(ctx, args, stdin) } // RunBounded preserves bounded capture when Compose argument injection is wrapped per installation. func (r InstallationRunner) RunBounded(ctx context.Context, args []string, stdin io.Reader, limits CaptureLimits) (Result, error) { if len(args) > 0 && args[0] == "compose" { return RunBounded(r.Runner, ctx, r.Installation.ComposeArgs(args[1:]...), stdin, limits) } return RunBounded(r.Runner, ctx, args, stdin, limits) } // RunBoundedStreaming preserves the stderr observer across installation argument injection. func (r InstallationRunner) RunBoundedStreaming(ctx context.Context, args []string, stdin io.Reader, limits CaptureLimits, stderrObserver func([]byte)) (Result, error) { if len(args) > 0 && args[0] == "compose" { return RunBoundedStreaming(r.Runner, ctx, r.Installation.ComposeArgs(args[1:]...), stdin, limits, stderrObserver) } return RunBoundedStreaming(r.Runner, ctx, args, stdin, limits, stderrObserver) }