Files
dockwatch/internal/stacks/stacks.go
T
jbergner ad54651558
release-tag / release-image (push) Successful in 2m13s
v9.3.3
2026-08-31 22:34:50 +02:00

1474 lines
40 KiB
Go

package stacks
import (
"bufio"
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"os"
"os/exec"
"path/filepath"
"regexp"
"sort"
"strings"
"sync"
"time"
"github.com/creack/pty"
"github.com/gorilla/websocket"
)
var validName = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$`)
var validSecretName = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$`)
var validServiceName = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$`)
type Stack struct {
Name string `json:"name"`
Compose string `json:"compose,omitempty"`
Env string `json:"env,omitempty"`
Secrets []SecretFile `json:"secrets,omitempty"`
EnvFiles []SecretFile `json:"env_files,omitempty"`
Configs []SecretFile `json:"configs,omitempty"`
Services []ServiceInfo `json:"services,omitempty"`
Status string `json:"status"`
Error string `json:"error,omitempty"`
}
type SecretFile struct {
Name string `json:"name"`
Content string `json:"content,omitempty"`
Size int64 `json:"size,omitempty"`
}
type ServiceInfo struct {
ID string `json:"id,omitempty"`
Name string `json:"name,omitempty"`
Service string `json:"service,omitempty"`
State string `json:"state,omitempty"`
Status string `json:"status,omitempty"`
Image string `json:"image,omitempty"`
Command string `json:"command,omitempty"`
Ports string `json:"ports,omitempty"`
}
type ExecInput struct {
Service string `json:"service"`
Command string `json:"command"`
}
type SaveInput struct {
Compose string `json:"compose"`
Env string `json:"env"`
Secrets []SecretFile `json:"secrets"`
EnvFiles []SecretFile `json:"env_files"`
Configs []SecretFile `json:"configs"`
}
type Service struct {
root string
hostRoot string
allowHostUserManagement bool
allowHostPermissionManagement bool
hostIdentityMu sync.Mutex
hostMappingMu sync.Mutex
hostMapping string
hostMappingAt time.Time
locks sync.Map
}
func New(root string) (*Service, error) {
if err := os.MkdirAll(root, 0750); err != nil {
return nil, err
}
return &Service{root: root}, nil
}
func (s *Service) Root() string { return s.root }
// ConfigureHostAccess enables optional host identity inspection. HostRoot is expected
// to point at a deliberately mounted host root (for example /host). User creation is
// disabled unless allowManagement is explicitly true.
func (s *Service) ConfigureHostAccess(hostRoot string, allowManagement bool) {
hostRoot = strings.TrimSpace(hostRoot)
if hostRoot != "" {
hostRoot = filepath.Clean(hostRoot)
}
s.hostRoot = hostRoot
s.allowHostUserManagement = allowManagement
}
// ConfigureHostPermissionManagement enables explicit host bind-mount ownership/mode repairs.
// This is intentionally independent from host user creation so operators can opt in to one
// class of host mutation without granting the other.
func (s *Service) ConfigureHostPermissionManagement(allow bool) {
s.allowHostPermissionManagement = allow
}
func (s *Service) lockFor(name string) *sync.Mutex {
v, _ := s.locks.LoadOrStore(name, &sync.Mutex{})
return v.(*sync.Mutex)
}
func (s *Service) stackDir(name string) (string, error) {
if !validName.MatchString(name) {
return "", errors.New("invalid stack name")
}
dir := filepath.Join(s.root, name)
info, err := os.Lstat(dir)
if err == nil {
if info.Mode()&os.ModeSymlink != 0 || !info.IsDir() {
return "", errors.New("stack path must be a real directory, not a symlink or file")
}
} else if !os.IsNotExist(err) {
return "", err
}
return dir, nil
}
func (s *Service) path(name string) (string, error) {
dir, err := s.stackDir(name)
if err != nil {
return "", err
}
return filepath.Join(dir, "compose.yaml"), nil
}
func (s *Service) envPath(name string) (string, error) {
dir, err := s.stackDir(name)
if err != nil {
return "", err
}
return filepath.Join(dir, ".env"), nil
}
func (s *Service) secretsDir(name string) (string, error) {
dir, err := s.stackDir(name)
if err != nil {
return "", err
}
return filepath.Join(dir, "secrets"), nil
}
func (s *Service) List(ctx context.Context) ([]Stack, error) {
ents, err := os.ReadDir(s.root)
if err != nil {
return nil, err
}
out := []Stack{}
for _, e := range ents {
if !e.IsDir() || !validName.MatchString(e.Name()) {
continue
}
p, _ := s.path(e.Name())
if _, err := os.Stat(p); err != nil {
continue
}
st := Stack{Name: e.Name(), Status: "unknown"}
if services, err := s.PS(ctx, e.Name()); err != nil {
st.Error = err.Error()
} else {
st.Services = services
st.Status = summarizeStatus(services)
}
out = append(out, st)
}
sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name })
return out, nil
}
func (s *Service) Get(ctx context.Context, name string) (Stack, error) {
p, err := s.path(name)
if err != nil {
return Stack{}, err
}
b, err := os.ReadFile(p)
if err != nil {
return Stack{}, err
}
st := Stack{Name: name, Compose: string(b), Status: "unknown"}
if eb, err := os.ReadFile(filepath.Join(filepath.Dir(p), ".env")); err == nil {
st.Env = string(eb)
}
st.Secrets, _ = s.ReadSecrets(name, true)
st.EnvFiles, _ = s.readManagedFiles(name, "envs", true)
st.Configs, _ = s.readManagedFiles(name, "configs", true)
if services, err := s.PS(ctx, name); err == nil {
st.Services = services
st.Status = summarizeStatus(services)
} else {
st.Error = err.Error()
}
return st, nil
}
func (s *Service) Save(ctx context.Context, name string, in SaveInput) error {
mu := s.lockFor(name)
mu.Lock()
defer mu.Unlock()
if len(in.Compose) == 0 || len(in.Compose) > 2<<20 {
return errors.New("compose file must be 1 byte to 2 MiB")
}
p, err := s.path(name)
if err != nil {
return err
}
finalDir := filepath.Dir(p)
if err := os.MkdirAll(s.root, 0750); err != nil {
return err
}
stageDir, err := os.MkdirTemp(s.root, ".dockwatch-stage-"+name+"-")
if err != nil {
return err
}
defer os.RemoveAll(stageDir)
composePath := filepath.Join(stageDir, "compose.yaml")
if err := os.WriteFile(composePath, []byte(in.Compose), 0640); err != nil {
return err
}
if strings.TrimSpace(in.Env) != "" {
if len(in.Env) > 512<<10 {
return errors.New(".env too large")
}
if err := os.WriteFile(filepath.Join(stageDir, ".env"), []byte(in.Env), 0640); err != nil {
return err
}
}
if err := writeFilesToDir(filepath.Join(stageDir, "secrets"), in.Secrets, 0600); err != nil {
return err
}
if err := writeFilesToDir(filepath.Join(stageDir, "envs"), in.EnvFiles, 0640); err != nil {
return err
}
if err := writeFilesToDir(filepath.Join(stageDir, "configs"), in.Configs, 0640); err != nil {
return err
}
if err := s.validateStaged(ctx, name, stageDir, composePath); err != nil {
return err
}
if err := os.MkdirAll(finalDir, 0750); err != nil {
return err
}
// Snapshot only the files Dockwatch owns. A save is applied as one logical
// operation and is rolled back when any copy/remove step fails. Unrelated
// bind-mount data beside compose.yaml is never part of the snapshot.
backupDir, err := os.MkdirTemp(s.root, ".dockwatch-backup-"+name+"-")
if err != nil {
return err
}
defer os.RemoveAll(backupDir)
if err := snapshotManaged(finalDir, backupDir); err != nil {
return fmt.Errorf("snapshot current stack: %w", err)
}
rollback := func(cause error) error {
if restoreErr := restoreManaged(finalDir, backupDir); restoreErr != nil {
return fmt.Errorf("%w (rollback failed: %v)", cause, restoreErr)
}
return cause
}
// Replace managed files only after validation succeeded.
if err := copyFile(composePath, filepath.Join(finalDir, "compose.yaml"), 0640); err != nil {
return rollback(err)
}
stageEnv := filepath.Join(stageDir, ".env")
if _, err := os.Stat(stageEnv); err == nil {
if err := copyFile(stageEnv, filepath.Join(finalDir, ".env"), 0640); err != nil {
return rollback(err)
}
} else if os.IsNotExist(err) {
if err := os.Remove(filepath.Join(finalDir, ".env")); err != nil && !os.IsNotExist(err) {
return rollback(err)
}
} else {
return rollback(err)
}
if err := replaceManagedDir(stageDir, finalDir, "secrets", in.Secrets, 0600); err != nil {
return rollback(err)
}
if err := replaceManagedDir(stageDir, finalDir, "envs", in.EnvFiles, 0640); err != nil {
return rollback(err)
}
if err := replaceManagedDir(stageDir, finalDir, "configs", in.Configs, 0640); err != nil {
return rollback(err)
}
return nil
}
func (s *Service) ValidateProject(ctx context.Context, name, projectDir, composeFile string) error {
if !validName.MatchString(name) {
return errors.New("invalid stack name")
}
projectDir = filepath.Clean(projectDir)
composeFile = filepath.Clean(composeFile)
if projectDir == "." || filepath.IsAbs(composeFile) || strings.HasPrefix(composeFile, "..") {
return errors.New("invalid compose project path")
}
return s.validateStaged(ctx, name, projectDir, filepath.Join(projectDir, composeFile))
}
func (s *Service) validateStaged(ctx context.Context, name, projectDir, composeFile string) error {
cctx, cancel := context.WithTimeout(ctx, 25*time.Second)
defer cancel()
args := []string{"compose", "--project-name", name, "--project-directory", projectDir, "-f", composeFile}
if _, err := os.Stat(filepath.Join(projectDir, ".env")); err == nil {
args = append(args, "--env-file", filepath.Join(projectDir, ".env"))
}
args = append(args, "config", "--quiet")
cmd := exec.CommandContext(cctx, "docker", args...)
cmd.Dir = projectDir
out, err := cmd.CombinedOutput()
if err == nil {
return nil
}
msg := strings.TrimSpace(string(out))
if msg == "" {
msg = err.Error()
}
if errors.Is(cctx.Err(), context.DeadlineExceeded) {
msg = "validation timed out after 25s"
}
return fmt.Errorf("compose validation failed: %s", msg)
}
var managedStackEntries = []string{"compose.yaml", ".env", "secrets", "envs", "configs"}
func snapshotManaged(srcDir, backupDir string) error {
for _, name := range managedStackEntries {
src := filepath.Join(srcDir, name)
if _, err := os.Lstat(src); os.IsNotExist(err) {
continue
} else if err != nil {
return err
}
if err := copyTree(src, filepath.Join(backupDir, name)); err != nil {
return err
}
}
return nil
}
func restoreManaged(dstDir, backupDir string) error {
for _, name := range managedStackEntries {
if err := os.RemoveAll(filepath.Join(dstDir, name)); err != nil {
return err
}
}
for _, name := range managedStackEntries {
src := filepath.Join(backupDir, name)
if _, err := os.Lstat(src); os.IsNotExist(err) {
continue
} else if err != nil {
return err
}
if err := copyTree(src, filepath.Join(dstDir, name)); err != nil {
return err
}
}
return nil
}
func copyTree(src, dst string) error {
info, err := os.Lstat(src)
if err != nil {
return err
}
if info.Mode()&os.ModeSymlink != 0 {
return fmt.Errorf("refusing symlink in managed stack data: %s", src)
}
if !info.IsDir() {
return copyFile(src, dst, info.Mode().Perm())
}
if err := os.MkdirAll(dst, info.Mode().Perm()); err != nil {
return err
}
entries, err := os.ReadDir(src)
if err != nil {
return err
}
for _, entry := range entries {
if err := copyTree(filepath.Join(src, entry.Name()), filepath.Join(dst, entry.Name())); err != nil {
return err
}
}
return nil
}
func copyFile(src, dst string, mode os.FileMode) error {
b, err := os.ReadFile(src)
if err != nil {
return err
}
if err := os.MkdirAll(filepath.Dir(dst), 0750); err != nil {
return err
}
f, err := os.CreateTemp(filepath.Dir(dst), ".dockwatch-write-*")
if err != nil {
return err
}
tmp := f.Name()
defer os.Remove(tmp)
if err := f.Chmod(mode); err != nil {
_ = f.Close()
return err
}
if _, err := f.Write(b); err != nil {
_ = f.Close()
return err
}
if err := f.Sync(); err != nil {
_ = f.Close()
return err
}
if err := f.Close(); err != nil {
return err
}
return os.Rename(tmp, dst)
}
func copyDir(src, dst string, mode os.FileMode) error {
if err := os.MkdirAll(dst, 0750); err != nil {
return err
}
ents, err := os.ReadDir(src)
if err != nil {
return err
}
for _, e := range ents {
if e.IsDir() {
continue
}
if err := copyFile(filepath.Join(src, e.Name()), filepath.Join(dst, e.Name()), mode); err != nil {
return err
}
}
return nil
}
func writeFilesToDir(dir string, files []SecretFile, mode os.FileMode) error {
if len(files) == 0 {
return nil
}
if err := os.MkdirAll(dir, 0750); err != nil {
return err
}
seen := map[string]bool{}
for _, f := range files {
f.Name = strings.TrimSpace(f.Name)
if !validSecretName.MatchString(f.Name) {
return fmt.Errorf("invalid managed file name %q", f.Name)
}
if seen[f.Name] {
return fmt.Errorf("duplicate managed file name %q", f.Name)
}
seen[f.Name] = true
if len(f.Content) > 512<<10 {
return fmt.Errorf("managed file %s too large", f.Name)
}
if err := os.WriteFile(filepath.Join(dir, f.Name), []byte(f.Content), mode); err != nil {
return err
}
}
return nil
}
func replaceManagedDir(stageDir, finalDir, sub string, files []SecretFile, mode os.FileMode) error {
dst := filepath.Join(finalDir, sub)
if len(files) == 0 {
return os.RemoveAll(dst)
}
if err := os.RemoveAll(dst); err != nil {
return err
}
return copyDir(filepath.Join(stageDir, sub), dst, mode)
}
// SaveCompose preserves backwards compatibility with the first scaffold.
func (s *Service) SaveCompose(ctx context.Context, name, compose string) error {
return s.Save(ctx, name, SaveInput{Compose: compose})
}
func (s *Service) writeEnv(name, env string) error {
p, err := s.envPath(name)
if err != nil {
return err
}
if strings.TrimSpace(env) == "" {
if err := os.Remove(p); err != nil && !os.IsNotExist(err) {
return err
}
return nil
}
if len(env) > 512<<10 {
return errors.New(".env too large")
}
return os.WriteFile(p, []byte(env), 0640)
}
func (s *Service) writeSecrets(name string, secrets []SecretFile) error {
dir, err := s.secretsDir(name)
if err != nil {
return err
}
if len(secrets) == 0 {
return nil
}
if err := os.MkdirAll(dir, 0750); err != nil {
return err
}
for _, sec := range secrets {
sec.Name = strings.TrimSpace(sec.Name)
if !validSecretName.MatchString(sec.Name) {
return fmt.Errorf("invalid secret name %q", sec.Name)
}
if len(sec.Content) > 512<<10 {
return fmt.Errorf("secret %s too large", sec.Name)
}
if err := os.WriteFile(filepath.Join(dir, sec.Name), []byte(sec.Content), 0600); err != nil {
return err
}
}
return nil
}
func (s *Service) ReadSecrets(name string, includeContent bool) ([]SecretFile, error) {
dir, err := s.secretsDir(name)
if err != nil {
return nil, err
}
ents, err := os.ReadDir(dir)
if os.IsNotExist(err) {
return []SecretFile{}, nil
}
if err != nil {
return nil, err
}
out := []SecretFile{}
for _, e := range ents {
if e.IsDir() || !validSecretName.MatchString(e.Name()) {
continue
}
info, _ := e.Info()
sf := SecretFile{Name: e.Name()}
if info != nil {
sf.Size = info.Size()
}
if includeContent {
b, err := os.ReadFile(filepath.Join(dir, e.Name()))
if err == nil {
sf.Content = string(b)
}
}
out = append(out, sf)
}
sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name })
return out, nil
}
func (s *Service) readManagedFiles(name, sub string, includeContent bool) ([]SecretFile, error) {
if !validName.MatchString(name) {
return nil, errors.New("invalid stack name")
}
dir := filepath.Join(s.root, name, sub)
ents, err := os.ReadDir(dir)
if os.IsNotExist(err) {
return []SecretFile{}, nil
}
if err != nil {
return nil, err
}
out := []SecretFile{}
for _, e := range ents {
if e.IsDir() || !validSecretName.MatchString(e.Name()) {
continue
}
info, _ := e.Info()
f := SecretFile{Name: e.Name()}
if info != nil {
f.Size = info.Size()
}
if includeContent {
if b, er := os.ReadFile(filepath.Join(dir, e.Name())); er == nil {
f.Content = string(b)
}
}
out = append(out, f)
}
sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name })
return out, nil
}
func (s *Service) Action(ctx context.Context, name, action string) (string, error) {
mu := s.lockFor(name)
mu.Lock()
defer mu.Unlock()
return s.actionUnlocked(ctx, name, action)
}
func (s *Service) actionUnlocked(ctx context.Context, name, action string) (string, error) {
switch action {
case "up":
return s.runString(ctx, name, "up", "-d", "--remove-orphans")
case "down":
return s.runString(ctx, name, "down")
case "restart":
return s.runString(ctx, name, "restart")
case "stop":
return s.runString(ctx, name, "stop")
case "start":
return s.runString(ctx, name, "start")
case "pull":
return s.runString(ctx, name, "pull")
case "update":
a, e1 := s.runString(ctx, name, "pull")
b, e2 := s.runString(ctx, name, "up", "-d", "--remove-orphans")
if e1 != nil {
return a + "\n" + b, e1
}
return a + "\n" + b, e2
case "recreate":
return s.runString(ctx, name, "up", "-d", "--force-recreate", "--remove-orphans")
default:
return "", errors.New("unsupported action")
}
}
func (s *Service) Logs(ctx context.Context, name string, tail int) (string, error) {
if tail < 1 {
tail = 200
}
if tail > 5000 {
tail = 5000
}
return s.runString(ctx, name, "logs", "--no-color", "--tail", fmt.Sprint(tail))
}
func (s *Service) StreamLogs(ctx context.Context, name string, tail int, w http.ResponseWriter) error {
if tail < 1 {
tail = 200
}
if tail > 2000 {
tail = 2000
}
p, err := s.path(name)
if err != nil {
return err
}
if _, err := os.Stat(p); err != nil {
return err
}
args := append(s.baseArgs(name), "logs", "--follow", "--no-color", "--tail", fmt.Sprint(tail))
cmd := exec.CommandContext(ctx, "docker", args...)
stdout, err := cmd.StdoutPipe()
if err != nil {
return err
}
stderr, err := cmd.StderrPipe()
if err != nil {
return err
}
if err := cmd.Start(); err != nil {
return err
}
defer func() {
if cmd.Process != nil {
_ = cmd.Process.Kill()
}
}()
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache, no-store")
w.Header().Set("X-Accel-Buffering", "no")
fl, _ := w.(http.Flusher)
lines := make(chan string, 256)
errCh := make(chan error, 3)
var wg sync.WaitGroup
scan := func(r io.Reader) {
defer wg.Done()
sc := bufio.NewScanner(r)
sc.Buffer(make([]byte, 64<<10), 1<<20)
for sc.Scan() {
select {
case lines <- sc.Text():
case <-ctx.Done():
return
}
}
if err := sc.Err(); err != nil {
select {
case errCh <- err:
default:
}
}
}
wg.Add(2)
go scan(stdout)
go scan(stderr)
go func() { wg.Wait(); close(lines) }()
enc := json.NewEncoder(w)
for {
select {
case <-ctx.Done():
_ = cmd.Wait()
return nil
case err := <-errCh:
if err != nil {
_ = cmd.Wait()
return err
}
case line, ok := <-lines:
if !ok {
err := cmd.Wait()
if err != nil && ctx.Err() == nil {
return err
}
return nil
}
fmt.Fprint(w, "data: ")
_ = enc.Encode(line)
fmt.Fprint(w, "\n")
if fl != nil {
fl.Flush()
}
}
}
}
func (s *Service) Exec(ctx context.Context, name string, in ExecInput) (string, error) {
in.Service = strings.TrimSpace(in.Service)
in.Command = strings.TrimSpace(in.Command)
if !validServiceName.MatchString(in.Service) || in.Command == "" {
return "", errors.New("valid service and command required")
}
if len(in.Command) > 4000 {
return "", errors.New("command too long")
}
return s.runString(ctx, name, "exec", "-T", in.Service, "sh", "-lc", in.Command)
}
func (s *Service) Delete(ctx context.Context, name string, down, purge bool) error {
mu := s.lockFor(name)
mu.Lock()
defer mu.Unlock()
p, err := s.path(name)
if err != nil {
return err
}
dir := filepath.Dir(p)
if down {
if _, err := s.actionUnlocked(ctx, name, "down"); err != nil {
return err
}
}
if purge {
return os.RemoveAll(dir)
}
// Safe delete removes only files managed by Dockwatch. Arbitrary bind-mount
// data living next to compose.yaml is never recursively deleted by default.
for _, entry := range []string{"compose.yaml", ".env", "secrets", "envs", "configs"} {
if err := os.RemoveAll(filepath.Join(dir, entry)); err != nil {
return err
}
}
// Git-backed stacks keep only Dockwatch's own manifest under .dockwatch.
// Never remove the whole directory because a repository/user may keep other
// metadata there.
manifest := filepath.Join(dir, ".dockwatch", "git-manifest.json")
if err := os.Remove(manifest); err != nil && !os.IsNotExist(err) {
return err
}
_ = os.Remove(filepath.Dir(manifest)) // succeeds only when empty
if ents, err := os.ReadDir(dir); err == nil && len(ents) == 0 {
_ = os.Remove(dir)
}
return nil
}
func (s *Service) PS(ctx context.Context, name string) ([]ServiceInfo, error) {
out, err := s.run(ctx, name, "ps", "--format", "json")
if err != nil {
return nil, err
}
var arr []map[string]any
dec := json.NewDecoder(bytes.NewReader(out))
if err := dec.Decode(&arr); err != nil {
// Some docker compose versions output one JSON object per line.
var services []ServiceInfo
sc := bufio.NewScanner(bytes.NewReader(out))
for sc.Scan() {
var m map[string]any
if json.Unmarshal(sc.Bytes(), &m) == nil {
services = append(services, mapService(m))
}
}
if len(services) > 0 {
return services, nil
}
return nil, err
}
services := make([]ServiceInfo, 0, len(arr))
for _, m := range arr {
services = append(services, mapService(m))
}
return services, nil
}
func mapService(m map[string]any) ServiceInfo {
return ServiceInfo{ID: str(m, "ID", "Id"), Name: str(m, "Name"), Service: str(m, "Service"), State: strings.ToLower(str(m, "State")), Status: str(m, "Status"), Image: str(m, "Image"), Command: str(m, "Command"), Ports: str(m, "Publishers", "Ports")}
}
func str(m map[string]any, keys ...string) string {
for _, k := range keys {
if v, ok := m[k]; ok && v != nil {
switch x := v.(type) {
case string:
return x
default:
b, _ := json.Marshal(x)
return string(b)
}
}
}
return ""
}
func summarizeStatus(sv []ServiceInfo) string {
if len(sv) == 0 {
return "stopped"
}
running := 0
bad := 0
for _, s := range sv {
st := strings.ToLower(s.State + " " + s.Status)
if strings.Contains(st, "running") || strings.Contains(st, "up") {
running++
}
if strings.Contains(st, "exit") || strings.Contains(st, "dead") || strings.Contains(st, "error") {
bad++
}
}
if bad > 0 {
return "degraded"
}
if running == len(sv) {
return "running"
}
if running > 0 {
return "partial"
}
return "stopped"
}
func (s *Service) runString(ctx context.Context, name string, args ...string) (string, error) {
out, err := s.run(ctx, name, args...)
return string(out), err
}
func (s *Service) run(ctx context.Context, name string, args ...string) ([]byte, error) {
p, err := s.path(name)
if err != nil {
return nil, err
}
if _, err := os.Stat(p); err != nil {
return nil, err
}
cctx, cancel := context.WithTimeout(ctx, 5*time.Minute)
defer cancel()
cmd := exec.CommandContext(cctx, "docker", append(s.baseArgs(name), args...)...)
var b bytes.Buffer
cmd.Stdout = &b
cmd.Stderr = &b
err = cmd.Run()
if err != nil {
return b.Bytes(), fmt.Errorf("docker compose: %w: %s", err, strings.TrimSpace(b.String()))
}
return b.Bytes(), nil
}
func (s *Service) baseArgs(name string) []string {
p, _ := s.path(name)
args := []string{"compose", "--project-name", name, "-f", p}
if ep, err := s.envPath(name); err == nil {
if _, stat := os.Stat(ep); stat == nil {
args = append(args, "--env-file", ep)
}
}
return args
}
// DockerInventory returns lightweight Docker CLI inventory records. It intentionally
// keeps Docker-specific fields as strings so different Engine versions remain compatible.
func (s *Service) DockerInventory(ctx context.Context, kind string) ([]map[string]string, error) {
var args []string
switch kind {
case "containers":
args = []string{"ps", "-a", "--format", "{{json .}}"}
case "images":
args = []string{"image", "ls", "--format", "{{json .}}"}
case "volumes":
args = []string{"volume", "ls", "--format", "{{json .}}"}
case "networks":
args = []string{"network", "ls", "--format", "{{json .}}"}
default:
return nil, errors.New("unsupported inventory kind")
}
cctx, cancel := context.WithTimeout(ctx, 20*time.Second)
defer cancel()
out, err := exec.CommandContext(cctx, "docker", args...).CombinedOutput()
if err != nil {
return nil, fmt.Errorf("docker inventory: %s", strings.TrimSpace(string(out)))
}
items := []map[string]string{}
sc := bufio.NewScanner(bytes.NewReader(out))
for sc.Scan() {
var v map[string]string
if json.Unmarshal(sc.Bytes(), &v) == nil {
items = append(items, v)
}
}
return items, sc.Err()
}
type DockerActionInput struct {
Name string `json:"name"`
Registry string `json:"registry"`
Username string `json:"username"`
Password string `json:"password"`
ID string `json:"id"`
Driver string `json:"driver"`
Force bool `json:"force"`
Internal bool `json:"internal"`
Attachable bool `json:"attachable"`
Labels map[string]string `json:"labels"`
}
func safeDockerPositional(v, label string) (string, error) {
v = strings.TrimSpace(v)
if v == "" {
return "", fmt.Errorf("%s required", label)
}
if strings.HasPrefix(v, "-") || strings.ContainsAny(v, "\r\n\x00") || len(v) > 4096 {
return "", fmt.Errorf("invalid %s", label)
}
return v, nil
}
func (s *Service) DockerAction(ctx context.Context, kind, action string, in DockerActionInput) (string, error) {
name := strings.TrimSpace(in.Name)
id := strings.TrimSpace(in.ID)
var args []string
switch kind {
case "containers":
target := id
if target == "" {
target = name
}
var err error
if target, err = safeDockerPositional(target, "container id/name"); err != nil {
return "", err
}
switch action {
case "start", "stop", "restart":
args = []string{action, target}
case "remove":
args = []string{"rm"}
if in.Force {
args = append(args, "-f")
}
args = append(args, target)
default:
return "", errors.New("unsupported container action")
}
case "images":
switch action {
case "login":
registry, err := safeDockerPositional(in.Registry, "registry")
user := strings.TrimSpace(in.Username)
if err != nil || user == "" || strings.ContainsAny(user, "\r\n\x00") || in.Password == "" {
return "", errors.New("valid registry, username and password required")
}
cctx, cancel := context.WithTimeout(ctx, 45*time.Second)
defer cancel()
cmd := exec.CommandContext(cctx, "docker", "login", registry, "--username", user, "--password-stdin")
cmd.Stdin = strings.NewReader(in.Password)
out, err := cmd.CombinedOutput()
msg := strings.TrimSpace(string(out))
if err != nil {
return msg, fmt.Errorf("docker login: %s", fallbackOutput(out, err))
}
return msg, nil
case "logout":
registry, err := safeDockerPositional(in.Registry, "registry")
if err != nil {
return "", err
}
return runDocker(ctx, "logout", registry)
case "pull":
var err error
if name, err = safeDockerPositional(name, "image reference"); err != nil {
return "", err
}
args = []string{"pull", name}
case "remove":
target := id
if target == "" {
target = name
}
var err error
if target, err = safeDockerPositional(target, "image id/reference"); err != nil {
return "", err
}
args = []string{"image", "rm"}
if in.Force {
args = append(args, "-f")
}
args = append(args, target)
case "prune":
args = []string{"image", "prune", "-f"}
default:
return "", errors.New("unsupported image action")
}
case "volumes":
switch action {
case "create":
var err error
if name, err = safeDockerPositional(name, "volume name"); err != nil {
return "", err
}
args = []string{"volume", "create"}
if in.Driver != "" {
args = append(args, "--driver", in.Driver)
}
args = appendLabels(args, in.Labels)
args = append(args, name)
case "remove":
target := name
if target == "" {
target = id
}
var err error
if target, err = safeDockerPositional(target, "volume name"); err != nil {
return "", err
}
args = []string{"volume", "rm"}
if in.Force {
args = append(args, "-f")
}
args = append(args, target)
case "prune":
args = []string{"volume", "prune", "-f"}
default:
return "", errors.New("unsupported volume action")
}
case "networks":
switch action {
case "create":
var err error
if name, err = safeDockerPositional(name, "network name"); err != nil {
return "", err
}
args = []string{"network", "create"}
if in.Driver != "" {
args = append(args, "--driver", in.Driver)
}
if in.Internal {
args = append(args, "--internal")
}
if in.Attachable {
args = append(args, "--attachable")
}
args = appendLabels(args, in.Labels)
args = append(args, name)
case "remove":
target := name
if target == "" {
target = id
}
var err error
if target, err = safeDockerPositional(target, "network name"); err != nil {
return "", err
}
args = []string{"network", "rm", target}
case "prune":
args = []string{"network", "prune", "-f"}
default:
return "", errors.New("unsupported network action")
}
default:
return "", errors.New("unsupported docker resource kind")
}
return runDocker(ctx, args...)
}
func appendLabels(args []string, labels map[string]string) []string {
keys := make([]string, 0, len(labels))
for k := range labels {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
v := labels[k]
if v == "" {
args = append(args, "--label", k)
} else {
args = append(args, "--label", k+"="+v)
}
}
return args
}
func runDocker(ctx context.Context, args ...string) (string, error) {
cctx, cancel := context.WithTimeout(ctx, 5*time.Minute)
defer cancel()
out, err := exec.CommandContext(cctx, "docker", args...).CombinedOutput()
msg := strings.TrimSpace(string(out))
if err != nil {
if msg == "" {
msg = err.Error()
}
return msg, fmt.Errorf("docker %s: %s", strings.Join(args, " "), msg)
}
return msg, nil
}
func (s *Service) DockerInspect(ctx context.Context, kind, id string) (map[string]any, error) {
var err error
if id, err = safeDockerPositional(id, "docker resource id/name"); err != nil {
return nil, err
}
var args []string
switch kind {
case "containers":
args = []string{"inspect", id}
case "images":
args = []string{"image", "inspect", id}
case "volumes":
args = []string{"volume", "inspect", id}
case "networks":
args = []string{"network", "inspect", id}
default:
return nil, errors.New("unsupported inspect kind")
}
raw, err := runDocker(ctx, args...)
if err != nil {
return nil, err
}
var arr []map[string]any
if err := json.Unmarshal([]byte(raw), &arr); err != nil || len(arr) == 0 {
return nil, errors.New("invalid docker inspect response")
}
out := map[string]any{"inspect": arr[0]}
if kind == "containers" {
statsRaw, _ := runDocker(ctx, "stats", "--no-stream", "--format", "{{json .}}", id)
stats := map[string]string{}
_ = json.Unmarshal([]byte(statsRaw), &stats)
out["stats"] = stats
}
return out, nil
}
type GraphNode struct {
ID string `json:"id"`
Kind string `json:"kind"`
Label string `json:"label"`
Image string `json:"image,omitempty"`
Running bool `json:"running,omitempty"`
}
type GraphEdge struct {
From string `json:"from"`
To string `json:"to"`
Kind string `json:"kind"`
}
type ServiceGraph struct {
Nodes []GraphNode `json:"nodes"`
Edges []GraphEdge `json:"edges"`
}
// Graph is based on Docker Compose's normalized JSON output, so profiles,
// interpolation and long/short syntax are interpreted by Compose itself.
func (s *Service) Graph(ctx context.Context, name string) (ServiceGraph, error) {
out, err := s.run(ctx, name, "config", "--format", "json")
if err != nil {
return ServiceGraph{}, err
}
var cfg struct {
Services map[string]struct {
Image string `json:"image"`
DependsOn map[string]any `json:"depends_on"`
Networks map[string]any `json:"networks"`
Volumes []struct {
Source string `json:"source"`
Target string `json:"target"`
Type string `json:"type"`
} `json:"volumes"`
} `json:"services"`
Networks map[string]any `json:"networks"`
Volumes map[string]any `json:"volumes"`
}
if err := json.Unmarshal(out, &cfg); err != nil {
return ServiceGraph{}, fmt.Errorf("decode compose config: %w", err)
}
ps, _ := s.PS(ctx, name)
running := map[string]bool{}
for _, p := range ps {
if strings.Contains(strings.ToLower(p.State+" "+p.Status), "running") || strings.Contains(strings.ToLower(p.Status), "up") {
running[p.Service] = true
}
}
g := ServiceGraph{}
serviceNames := make([]string, 0, len(cfg.Services))
for n := range cfg.Services {
serviceNames = append(serviceNames, n)
}
sort.Strings(serviceNames)
for _, n := range serviceNames {
sv := cfg.Services[n]
g.Nodes = append(g.Nodes, GraphNode{ID: "service:" + n, Kind: "service", Label: n, Image: sv.Image, Running: running[n]})
deps := make([]string, 0, len(sv.DependsOn))
for d := range sv.DependsOn {
deps = append(deps, d)
}
sort.Strings(deps)
for _, d := range deps {
g.Edges = append(g.Edges, GraphEdge{From: "service:" + n, To: "service:" + d, Kind: "depends_on"})
}
for netName := range sv.Networks {
g.Edges = append(g.Edges, GraphEdge{From: "service:" + n, To: "network:" + netName, Kind: "network"})
}
for _, v := range sv.Volumes {
if v.Type == "volume" && v.Source != "" {
g.Edges = append(g.Edges, GraphEdge{From: "service:" + n, To: "volume:" + v.Source, Kind: "volume"})
}
}
}
nets := make([]string, 0, len(cfg.Networks))
for n := range cfg.Networks {
nets = append(nets, n)
}
sort.Strings(nets)
for _, n := range nets {
g.Nodes = append(g.Nodes, GraphNode{ID: "network:" + n, Kind: "network", Label: n})
}
vols := make([]string, 0, len(cfg.Volumes))
for n := range cfg.Volumes {
vols = append(vols, n)
}
sort.Strings(vols)
for _, n := range vols {
g.Nodes = append(g.Nodes, GraphNode{ID: "volume:" + n, Kind: "volume", Label: n})
}
return g, nil
}
type ImageUpdate struct {
Service string `json:"service"`
Image string `json:"image"`
LocalDigest string `json:"local_digest,omitempty"`
RemoteDigest string `json:"remote_digest,omitempty"`
Update bool `json:"update"`
Error string `json:"error,omitempty"`
CheckedAt int64 `json:"checked_at"`
}
func (s *Service) ImageUpdates(ctx context.Context, name string) ([]ImageUpdate, error) {
out, err := s.run(ctx, name, "config", "--format", "json")
if err != nil {
return nil, err
}
var cfg struct {
Services map[string]struct {
Image string `json:"image"`
} `json:"services"`
}
if err = json.Unmarshal(out, &cfg); err != nil {
return nil, err
}
names := make([]string, 0, len(cfg.Services))
for n := range cfg.Services {
names = append(names, n)
}
sort.Strings(names)
res := make([]ImageUpdate, 0, len(names))
for _, svc := range names {
img := strings.TrimSpace(cfg.Services[svc].Image)
r := ImageUpdate{Service: svc, Image: img, CheckedAt: time.Now().Unix()}
if img == "" {
r.Error = "service has no image"
res = append(res, r)
continue
}
if strings.Contains(img, "@sha256:") {
r.LocalDigest = strings.SplitN(img, "@", 2)[1]
r.RemoteDigest = r.LocalDigest
res = append(res, r)
continue
}
r.LocalDigest, _ = localImageDigest(ctx, img)
r.RemoteDigest, err = remoteImageDigest(ctx, img)
if err != nil {
r.Error = err.Error()
} else if r.LocalDigest != "" && r.RemoteDigest != "" {
r.Update = r.LocalDigest != r.RemoteDigest
}
res = append(res, r)
}
return res, nil
}
func localImageDigest(ctx context.Context, image string) (string, error) {
cctx, cancel := context.WithTimeout(ctx, 20*time.Second)
defer cancel()
out, err := exec.CommandContext(cctx, "docker", "image", "inspect", "--format", "{{json .RepoDigests}}", image).CombinedOutput()
if err != nil {
return "", fmt.Errorf("local inspect: %s", strings.TrimSpace(string(out)))
}
var ds []string
if json.Unmarshal(bytes.TrimSpace(out), &ds) != nil || len(ds) == 0 {
return "", nil
}
for _, d := range ds {
if i := strings.LastIndex(d, "@sha256:"); i >= 0 {
return strings.TrimPrefix(d[i+1:], "@"), nil
}
}
return "", nil
}
func remoteImageDigest(ctx context.Context, image string) (string, error) {
cctx, cancel := context.WithTimeout(ctx, 45*time.Second)
defer cancel()
// buildx imagetools reports the top-level registry digest. This is important
// for multi-arch tags because a platform descriptor digest is not the same
// value as the RepoDigest stored by Docker for the manifest index.
if out, err := exec.CommandContext(cctx, "docker", "buildx", "imagetools", "inspect", image).CombinedOutput(); err == nil {
re := regexp.MustCompile(`(?m)^Digest:\s*(sha256:[a-fA-F0-9]{64})\s*$`)
if m := re.FindSubmatch(out); len(m) == 2 {
return string(m[1]), nil
}
}
out, err := exec.CommandContext(cctx, "docker", "manifest", "inspect", "--verbose", image).CombinedOutput()
if err != nil {
return "", fmt.Errorf("manifest inspect: %s", fallbackOutput(out, err))
}
var v any
if e := json.Unmarshal(out, &v); e != nil {
return "", fmt.Errorf("manifest decode: %w", e)
}
if d := descriptorDigest(v); d != "" {
return d, nil
}
return "", errors.New("registry manifest did not expose a digest")
}
func descriptorDigest(v any) string {
switch x := v.(type) {
case map[string]any:
if d, ok := x["Descriptor"].(map[string]any); ok {
if s, ok := d["digest"].(string); ok {
return s
}
}
if s, ok := x["digest"].(string); ok && strings.HasPrefix(s, "sha256:") {
return s
}
for _, k := range []string{"descriptor", "manifests"} {
if z, ok := x[k]; ok {
if d := descriptorDigest(z); d != "" {
return d
}
}
}
case []any:
for _, z := range x {
if d := descriptorDigest(z); d != "" {
return d
}
}
}
return ""
}
func fallbackOutput(out []byte, err error) string {
v := strings.TrimSpace(string(out))
if v == "" {
return err.Error()
}
return v
}
type TerminalMessage struct {
Type string `json:"type"`
Data string `json:"data,omitempty"`
Cols uint16 `json:"cols,omitempty"`
Rows uint16 `json:"rows,omitempty"`
Code int `json:"code,omitempty"`
}
func allowedShell(v string) string {
v = strings.TrimSpace(v)
switch v {
case "bash", "/bin/bash", "sh", "/bin/sh", "ash", "/bin/ash", "zsh", "/bin/zsh":
return v
default:
return "sh"
}
}
// Terminal runs docker compose exec inside a real PTY and bridges the PTY to
// a WebSocket. Client messages are JSON {type:input,data:"..."} or
// {type:resize,cols:120,rows:40}. Server output uses {type:output,data:"..."}.
func (s *Service) Terminal(ctx context.Context, name, service, shell string, ws *websocket.Conn) error {
service = strings.TrimSpace(service)
if !validServiceName.MatchString(service) {
return errors.New("valid service required")
}
p, err := s.path(name)
if err != nil {
return err
}
if _, err := os.Stat(p); err != nil {
return err
}
args := append(s.baseArgs(name), "exec", service, allowedShell(shell))
cmd := exec.CommandContext(ctx, "docker", args...)
ptmx, err := pty.Start(cmd)
if err != nil {
return fmt.Errorf("start terminal: %w", err)
}
defer func() {
_ = ptmx.Close()
if cmd.Process != nil {
_ = cmd.Process.Kill()
}
_ = cmd.Wait()
}()
_ = pty.Setsize(ptmx, &pty.Winsize{Cols: 120, Rows: 32})
writeMu := sync.Mutex{}
writeJSON := func(v TerminalMessage) error { writeMu.Lock(); defer writeMu.Unlock(); return ws.WriteJSON(v) }
done := make(chan error, 2)
go func() {
buf := make([]byte, 16<<10)
for {
n, e := ptmx.Read(buf)
if n > 0 {
if werr := writeJSON(TerminalMessage{Type: "output", Data: string(buf[:n])}); werr != nil {
done <- werr
return
}
}
if e != nil {
done <- e
return
}
}
}()
go func() {
for {
var m TerminalMessage
if e := ws.ReadJSON(&m); e != nil {
done <- e
return
}
switch m.Type {
case "input":
if _, e := io.WriteString(ptmx, m.Data); e != nil {
done <- e
return
}
case "resize":
if m.Cols > 0 && m.Rows > 0 {
_ = pty.Setsize(ptmx, &pty.Winsize{Cols: m.Cols, Rows: m.Rows})
}
}
}
}()
select {
case <-ctx.Done():
return nil
case e := <-done:
_ = writeJSON(TerminalMessage{Type: "exit"})
return e
}
}