Files
jbergner 45ca18b74e
release-tag / release-image (push) Failing after 1m20s
init
2026-08-31 17:09:21 +02:00

766 lines
21 KiB
Go

package gitops
import (
"context"
"crypto/aes"
"crypto/cipher"
"crypto/hmac"
"crypto/rand"
"crypto/sha256"
"database/sql"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"os/exec"
"path/filepath"
"regexp"
"sort"
"strings"
"sync"
"time"
"git.send.nrw/sendnrw/dockwatch/internal/nodes"
"git.send.nrw/sendnrw/dockwatch/internal/stacks"
)
var stackNameRx = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$`)
type Source struct {
ID int64 `json:"id"`
NodeID *int64 `json:"node_id,omitempty"`
StackName string `json:"stack_name"`
RepoURL string `json:"repo_url"`
Branch string `json:"branch"`
Workdir string `json:"workdir"`
ComposeFile string `json:"compose_file"`
AutoDeploy bool `json:"auto_deploy"`
LastCommit string `json:"last_commit"`
LastSyncAt *int64 `json:"last_sync_at,omitempty"`
LastError string `json:"last_error"`
CreatedAt int64 `json:"created_at"`
UpdatedAt int64 `json:"updated_at"`
}
type Input struct {
NodeID *int64 `json:"node_id"`
StackName string `json:"stack_name"`
RepoURL string `json:"repo_url"`
Branch string `json:"branch"`
Workdir string `json:"workdir"`
ComposeFile string `json:"compose_file"`
AutoDeploy bool `json:"auto_deploy"`
}
type Service struct {
db *sql.DB
key []byte
stacks *stacks.Service
nodes *nodes.Manager
locks sync.Map
}
func New(db *sql.DB, key []byte, ss *stacks.Service, nm *nodes.Manager) *Service {
return &Service{db: db, key: key, stacks: ss, nodes: nm}
}
func normalize(in *Input) error {
in.StackName = strings.TrimSpace(in.StackName)
in.RepoURL = strings.TrimSpace(in.RepoURL)
in.Branch = strings.TrimSpace(in.Branch)
in.Workdir = filepath.Clean(strings.TrimSpace(in.Workdir))
in.ComposeFile = filepath.Clean(strings.TrimSpace(in.ComposeFile))
if in.StackName == "" || in.RepoURL == "" {
return errors.New("stack_name and repo_url required")
}
if !stackNameRx.MatchString(in.StackName) || strings.ContainsAny(in.RepoURL, "\r\n") || strings.HasPrefix(in.RepoURL, "-") {
return errors.New("invalid stack name or repository URL")
}
if len(in.RepoURL) > 4096 || len(in.Branch) > 255 || strings.ContainsAny(in.Branch, "\r\n") {
return errors.New("git source fields too long")
}
if in.Branch == "" {
in.Branch = "main"
}
if in.Workdir == "." || in.Workdir == "" {
in.Workdir = "."
}
if strings.HasPrefix(in.Workdir, "..") || filepath.IsAbs(in.Workdir) {
return errors.New("invalid workdir")
}
if in.ComposeFile == "." || in.ComposeFile == "" {
in.ComposeFile = "compose.yaml"
}
if strings.HasPrefix(in.ComposeFile, "..") || filepath.IsAbs(in.ComposeFile) {
return errors.New("invalid compose_file")
}
return nil
}
func (s *Service) List(ctx context.Context) ([]Source, error) {
rows, e := s.db.QueryContext(ctx, `SELECT id,node_id,stack_name,repo_url,branch,workdir,compose_file,auto_deploy,last_commit,last_sync_at,last_error,created_at,updated_at FROM git_sources ORDER BY stack_name`)
if e != nil {
return nil, e
}
defer rows.Close()
out := []Source{}
for rows.Next() {
var x Source
var sync, node sql.NullInt64
if e := rows.Scan(&x.ID, &node, &x.StackName, &x.RepoURL, &x.Branch, &x.Workdir, &x.ComposeFile, &x.AutoDeploy, &x.LastCommit, &sync, &x.LastError, &x.CreatedAt, &x.UpdatedAt); e != nil {
return nil, e
}
if node.Valid {
v := node.Int64
x.NodeID = &v
}
if sync.Valid {
v := sync.Int64
x.LastSyncAt = &v
}
out = append(out, x)
}
return out, rows.Err()
}
func (s *Service) Get(ctx context.Context, id int64) (Source, error) {
var x Source
var sync, node sql.NullInt64
e := s.db.QueryRowContext(ctx, `SELECT id,node_id,stack_name,repo_url,branch,workdir,compose_file,auto_deploy,last_commit,last_sync_at,last_error,created_at,updated_at FROM git_sources WHERE id=?`, id).Scan(&x.ID, &node, &x.StackName, &x.RepoURL, &x.Branch, &x.Workdir, &x.ComposeFile, &x.AutoDeploy, &x.LastCommit, &sync, &x.LastError, &x.CreatedAt, &x.UpdatedAt)
if node.Valid {
v := node.Int64
x.NodeID = &v
}
if sync.Valid {
v := sync.Int64
x.LastSyncAt = &v
}
return x, e
}
func (s *Service) Create(ctx context.Context, in Input) (Source, string, error) {
if e := normalize(&in); e != nil {
return Source{}, "", e
}
secret := make([]byte, 32)
if _, e := rand.Read(secret); e != nil {
return Source{}, "", e
}
sec := hex.EncodeToString(secret)
enc, e := s.encrypt([]byte(sec))
if e != nil {
return Source{}, "", e
}
now := time.Now().Unix()
r, e := s.db.ExecContext(ctx, `INSERT INTO git_sources(node_id,stack_name,repo_url,branch,workdir,compose_file,auto_deploy,webhook_secret_enc,last_commit,last_error,created_at,updated_at) VALUES(?,?,?,?,?,?,?,?,'','',?,?)`, in.NodeID, in.StackName, in.RepoURL, in.Branch, in.Workdir, in.ComposeFile, in.AutoDeploy, enc, now, now)
if e != nil {
return Source{}, "", e
}
id, _ := r.LastInsertId()
x, e := s.Get(ctx, id)
return x, sec, e
}
func (s *Service) Update(ctx context.Context, id int64, in Input) (Source, error) {
if e := normalize(&in); e != nil {
return Source{}, e
}
now := time.Now().Unix()
r, e := s.db.ExecContext(ctx, `UPDATE git_sources SET node_id=?,stack_name=?,repo_url=?,branch=?,workdir=?,compose_file=?,auto_deploy=?,updated_at=? WHERE id=?`, in.NodeID, in.StackName, in.RepoURL, in.Branch, in.Workdir, in.ComposeFile, in.AutoDeploy, now, id)
if e != nil {
return Source{}, e
}
n, _ := r.RowsAffected()
if n == 0 {
return Source{}, sql.ErrNoRows
}
return s.Get(ctx, id)
}
func (s *Service) Delete(ctx context.Context, id int64) error {
_, e := s.db.ExecContext(ctx, `DELETE FROM git_sources WHERE id=?`, id)
return e
}
func (s *Service) RotateSecret(ctx context.Context, id int64) (string, error) {
secret := make([]byte, 32)
if _, e := rand.Read(secret); e != nil {
return "", e
}
sec := hex.EncodeToString(secret)
enc, e := s.encrypt([]byte(sec))
if e != nil {
return "", e
}
_, e = s.db.ExecContext(ctx, `UPDATE git_sources SET webhook_secret_enc=?,updated_at=? WHERE id=?`, enc, time.Now().Unix(), id)
return sec, e
}
func (s *Service) VerifyWebhook(ctx context.Context, id int64, body []byte, signature, token string) error {
var enc []byte
if e := s.db.QueryRowContext(ctx, `SELECT webhook_secret_enc FROM git_sources WHERE id=?`, id).Scan(&enc); e != nil {
return e
}
plain, e := s.decrypt(enc)
if e != nil {
return e
}
secret := string(plain)
if token != "" && hmac.Equal([]byte(token), []byte(secret)) {
return nil
}
signature = strings.TrimPrefix(signature, "sha256=")
if signature != "" {
mac := hmac.New(sha256.New, []byte(secret))
_, _ = mac.Write(body)
want := hex.EncodeToString(mac.Sum(nil))
if hmac.Equal([]byte(strings.ToLower(signature)), []byte(want)) {
return nil
}
}
return errors.New("invalid webhook signature")
}
func (s *Service) lockFor(id int64) *sync.Mutex {
v, _ := s.locks.LoadOrStore(id, &sync.Mutex{})
return v.(*sync.Mutex)
}
func (s *Service) Sync(ctx context.Context, id int64) (Source, error) {
mu := s.lockFor(id)
mu.Lock()
defer mu.Unlock()
x, e := s.Get(ctx, id)
if e != nil {
return x, e
}
in := Input{NodeID: x.NodeID, StackName: x.StackName, RepoURL: x.RepoURL, Branch: x.Branch, Workdir: x.Workdir, ComposeFile: x.ComposeFile, AutoDeploy: x.AutoDeploy}
var commit string
if x.NodeID != nil {
if s.nodes == nil {
return s.syncFailed(ctx, x, errors.New("node manager unavailable"))
}
b, _, err := s.nodes.Do(ctx, *x.NodeID, "POST", "/agent/v1/git/sync", in)
if err != nil {
return s.syncFailed(ctx, x, err)
}
var resp struct {
Commit string `json:"commit"`
}
if err = json.Unmarshal(b, &resp); err != nil {
return s.syncFailed(ctx, x, err)
}
commit = resp.Commit
} else {
commit, e = s.SyncTransient(ctx, in)
if e != nil {
return s.syncFailed(ctx, x, e)
}
}
now := time.Now().Unix()
_, _ = s.db.ExecContext(ctx, `UPDATE git_sources SET last_commit=?,last_sync_at=?,last_error='',updated_at=? WHERE id=?`, commit, now, now, id)
return s.Get(ctx, id)
}
// SyncTransient performs a Git-backed stack synchronization on the current
// Docker environment. Agents expose this operation to the master without
// persisting Git source metadata locally.
func (s *Service) SyncTransient(ctx context.Context, in Input) (string, error) {
if e := normalize(&in); e != nil {
return "", e
}
tmp, e := os.MkdirTemp("", "dockwatch-git-*")
if e != nil {
return "", e
}
defer os.RemoveAll(tmp)
cloneDir := filepath.Join(tmp, "repo")
cctx, cancel := context.WithTimeout(ctx, 3*time.Minute)
defer cancel()
cmd := exec.CommandContext(cctx, "git", "clone", "--depth", "1", "--branch", in.Branch, "--single-branch", in.RepoURL, cloneDir)
out, e := cmd.CombinedOutput()
if e != nil {
return "", fmt.Errorf("git clone: %w: %s", e, strings.TrimSpace(string(out)))
}
commitRaw, e := exec.CommandContext(cctx, "git", "-C", cloneDir, "rev-parse", "HEAD").Output()
if e != nil {
return "", e
}
src := filepath.Join(cloneDir, in.Workdir)
if fi, e := os.Stat(src); e != nil || !fi.IsDir() {
return "", errors.New("git workdir not found")
}
if _, e := os.Stat(filepath.Join(src, in.ComposeFile)); e != nil {
return "", fmt.Errorf("compose file not found: %s", in.ComposeFile)
}
dst := filepath.Join(s.stacks.Root(), in.StackName)
stage := filepath.Join(tmp, "stage")
if e := copyDir(src, stage); e != nil {
return "", e
}
if filepath.Clean(in.ComposeFile) != "compose.yaml" {
b, e := os.ReadFile(filepath.Join(stage, in.ComposeFile))
if e != nil {
return "", e
}
if e = os.WriteFile(filepath.Join(stage, "compose.yaml"), b, 0640); e != nil {
return "", e
}
}
if e := s.stacks.ValidateProject(ctx, in.StackName, stage, "compose.yaml"); e != nil {
return "", e
}
if e := syncManagedTree(stage, dst); e != nil {
return "", e
}
if in.AutoDeploy {
if _, e = s.stacks.Action(ctx, in.StackName, "up"); e != nil {
return "", e
}
}
return strings.TrimSpace(string(commitRaw)), nil
}
func (s *Service) syncFailed(ctx context.Context, x Source, e error) (Source, error) {
_, _ = s.db.ExecContext(ctx, `UPDATE git_sources SET last_error=?,updated_at=? WHERE id=?`, e.Error(), time.Now().Unix(), x.ID)
x.LastError = e.Error()
return x, e
}
const gitManifestPath = ".dockwatch/git-manifest.json"
type gitManifest struct {
Files []string `json:"files"`
}
// syncManagedTree applies a Git checkout without treating the stack directory as
// disposable storage. Only files previously managed by Git and files present in
// the new checkout are changed. Unrelated files (for example bind-mount data)
// survive a sync. The touched files are snapshotted so a failed apply can be
// rolled back without copying or deleting the whole stack directory.
func syncManagedTree(stage, dst string) error {
files, err := collectManagedFiles(stage)
if err != nil {
return err
}
if err := ensureSafeRoot(dst); err != nil {
return err
}
old, err := readGitManifest(dst)
if err != nil {
return err
}
backup, err := os.MkdirTemp("", "dockwatch-git-rollback-*")
if err != nil {
return err
}
defer os.RemoveAll(backup)
affected := map[string]struct{}{}
for _, rel := range old {
affected[rel] = struct{}{}
}
for _, rel := range files {
affected[rel] = struct{}{}
}
manifestAbs := filepath.Join(dst, filepath.FromSlash(gitManifestPath))
manifestBackup := filepath.Join(backup, "manifest.json")
manifestExisted := false
if info, err := os.Lstat(manifestAbs); err == nil {
if !info.Mode().IsRegular() {
return errors.New("Git manifest path is not a regular file")
}
if err := copyRegularFile(manifestAbs, manifestBackup, info.Mode().Perm()); err != nil {
return err
}
manifestExisted = true
} else if !os.IsNotExist(err) {
return err
}
backedUp := map[string]bool{}
for rel := range affected {
target, err := safeManagedPath(dst, rel)
if err != nil {
return err
}
if err := ensureSafeParent(dst, filepath.Dir(target)); err != nil {
return err
}
info, err := os.Lstat(target)
if os.IsNotExist(err) {
continue
}
if err != nil {
return err
}
if info.Mode()&os.ModeSymlink != 0 {
return fmt.Errorf("refusing to replace symlink in Git-managed path %q", rel)
}
if info.Mode().IsRegular() {
bp := filepath.Join(backup, filepath.FromSlash(rel))
if err := copyRegularFile(target, bp, info.Mode().Perm()); err != nil {
return err
}
backedUp[rel] = true
}
}
rollback := func() {
for rel := range affected {
target, err := safeManagedPath(dst, rel)
if err != nil {
continue
}
if info, err := os.Lstat(target); err == nil && info.Mode().IsRegular() {
_ = os.Remove(target)
}
if backedUp[rel] {
bp := filepath.Join(backup, filepath.FromSlash(rel))
if info, err := os.Stat(bp); err == nil {
_ = copyRegularFile(bp, target, info.Mode().Perm())
}
}
}
_ = os.Remove(manifestAbs)
if manifestExisted {
if info, err := os.Stat(manifestBackup); err == nil {
_ = copyRegularFile(manifestBackup, manifestAbs, info.Mode().Perm())
}
}
}
newSet := make(map[string]struct{}, len(files))
for _, rel := range files {
newSet[rel] = struct{}{}
}
// Remove files that disappeared from Git, but never recursively remove a
// directory. This deliberately leaves unrelated data untouched.
for _, rel := range old {
if _, ok := newSet[rel]; ok {
continue
}
target, err := safeManagedPath(dst, rel)
if err != nil {
rollback()
return err
}
if info, err := os.Lstat(target); err == nil {
if info.Mode()&os.ModeSymlink != 0 || (!info.Mode().IsRegular() && !info.IsDir()) {
rollback()
return fmt.Errorf("refusing to remove non-regular Git-managed path %q", rel)
}
if info.Mode().IsRegular() {
if err := os.Remove(target); err != nil {
rollback()
return err
}
}
} else if !os.IsNotExist(err) {
rollback()
return err
}
}
for _, rel := range files {
source, err := safeManagedPath(stage, rel)
if err != nil {
rollback()
return err
}
target, err := safeManagedPath(dst, rel)
if err != nil {
rollback()
return err
}
if err := ensureSafeParent(dst, filepath.Dir(target)); err != nil {
rollback()
return err
}
if info, err := os.Lstat(target); err == nil && info.IsDir() {
// A file replacing a directory is safe only when that directory is
// empty. os.Remove intentionally refuses non-empty directories.
if err := os.Remove(target); err != nil {
rollback()
return fmt.Errorf("Git file %q conflicts with existing directory containing unmanaged data: %w", rel, err)
}
} else if err == nil && info.Mode()&os.ModeSymlink != 0 {
rollback()
return fmt.Errorf("refusing to replace symlink in Git-managed path %q", rel)
} else if err != nil && !os.IsNotExist(err) {
rollback()
return err
}
info, err := os.Stat(source)
if err != nil {
rollback()
return err
}
if err := copyRegularFile(source, target, info.Mode().Perm()); err != nil {
rollback()
return err
}
}
if err := writeGitManifest(dst, files); err != nil {
rollback()
return err
}
return nil
}
func collectManagedFiles(root string) ([]string, error) {
out := []string{}
err := filepath.Walk(root, func(path string, info os.FileInfo, err error) error {
if err != nil {
return err
}
if path == root {
return nil
}
rel, err := filepath.Rel(root, path)
if err != nil {
return err
}
rel = filepath.ToSlash(rel)
if rel == gitManifestPath {
return nil
}
if info.Mode()&os.ModeSymlink != 0 {
return fmt.Errorf("Git stack contains unsupported symlink %q", rel)
}
if info.IsDir() {
return nil
}
if info.Mode().IsRegular() {
out = append(out, rel)
}
return nil
})
if err != nil {
return nil, err
}
sort.Strings(out)
return out, nil
}
func readGitManifest(dst string) ([]string, error) {
path := filepath.Join(dst, filepath.FromSlash(gitManifestPath))
if err := ensureSafeParent(dst, filepath.Dir(path)); err != nil {
return nil, err
}
b, err := os.ReadFile(path)
if os.IsNotExist(err) {
return []string{}, nil
}
if err != nil {
return nil, err
}
var m gitManifest
if err := json.Unmarshal(b, &m); err != nil {
return nil, fmt.Errorf("invalid Git-managed file manifest: %w", err)
}
out := make([]string, 0, len(m.Files))
seen := map[string]bool{}
for _, rel := range m.Files {
rel = filepath.ToSlash(filepath.Clean(filepath.FromSlash(rel)))
if rel == "." || rel == gitManifestPath || strings.HasPrefix(rel, "../") || filepath.IsAbs(filepath.FromSlash(rel)) || seen[rel] {
continue
}
seen[rel] = true
out = append(out, rel)
}
return out, nil
}
func writeGitManifest(dst string, files []string) error {
path := filepath.Join(dst, filepath.FromSlash(gitManifestPath))
if err := ensureSafeParent(dst, filepath.Dir(path)); err != nil {
return err
}
b, err := json.MarshalIndent(gitManifest{Files: files}, "", " ")
if err != nil {
return err
}
f, err := os.CreateTemp(filepath.Dir(path), ".dockwatch-git-manifest-*")
if err != nil {
return err
}
tmp := f.Name()
defer os.Remove(tmp)
if err := f.Chmod(0640); err != nil {
_ = f.Close()
return err
}
if _, err := f.Write(append(b, '\n')); 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, path)
}
func safeManagedPath(root, rel string) (string, error) {
rel = filepath.Clean(filepath.FromSlash(rel))
if rel == "." || filepath.IsAbs(rel) || rel == ".." || strings.HasPrefix(rel, ".."+string(os.PathSeparator)) {
return "", errors.New("invalid Git-managed path")
}
return filepath.Join(root, rel), nil
}
func ensureSafeRoot(root string) error {
info, err := os.Lstat(root)
if os.IsNotExist(err) {
return os.MkdirAll(root, 0750)
}
if err != nil {
return err
}
if info.Mode()&os.ModeSymlink != 0 || !info.IsDir() {
return errors.New("Git stack destination must be a real directory")
}
return nil
}
func ensureSafeParent(root, parent string) error {
rel, err := filepath.Rel(root, parent)
if err != nil || rel == ".." || strings.HasPrefix(rel, ".."+string(os.PathSeparator)) {
return errors.New("path escapes stack root")
}
cur := root
if err := ensureSafeRoot(root); err != nil {
return err
}
if rel == "." {
return nil
}
for _, part := range strings.Split(rel, string(os.PathSeparator)) {
cur = filepath.Join(cur, part)
info, err := os.Lstat(cur)
if os.IsNotExist(err) {
if err := os.Mkdir(cur, 0750); err != nil && !os.IsExist(err) {
return err
}
continue
}
if err != nil {
return err
}
if info.Mode()&os.ModeSymlink != 0 || !info.IsDir() {
return fmt.Errorf("unsafe parent path %q", cur)
}
}
return nil
}
func copyRegularFile(src, dst string, mode os.FileMode) error {
if err := os.MkdirAll(filepath.Dir(dst), 0750); err != nil {
return err
}
in, err := os.Open(src)
if err != nil {
return err
}
tmp, err := os.CreateTemp(filepath.Dir(dst), ".dockwatch-git-write-*")
if err != nil {
_ = in.Close()
return err
}
tmpName := tmp.Name()
defer os.Remove(tmpName)
if err := tmp.Chmod(mode); err != nil {
_ = in.Close()
_ = tmp.Close()
return err
}
_, copyErr := io.Copy(tmp, in)
inErr := in.Close()
if copyErr == nil {
copyErr = tmp.Sync()
}
outErr := tmp.Close()
if copyErr != nil {
return copyErr
}
if inErr != nil {
return inErr
}
if outErr != nil {
return outErr
}
return os.Rename(tmpName, dst)
}
func copyDir(src, dst string) error {
return filepath.Walk(src, func(path string, info os.FileInfo, e error) error {
if e != nil {
return e
}
rel, e := filepath.Rel(src, path)
if e != nil {
return e
}
if rel == ".git" || strings.HasPrefix(rel, ".git"+string(os.PathSeparator)) {
if info.IsDir() {
return filepath.SkipDir
}
return nil
}
target := filepath.Join(dst, rel)
if info.IsDir() {
return os.MkdirAll(target, info.Mode().Perm())
}
if !info.Mode().IsRegular() {
return nil
}
if e := os.MkdirAll(filepath.Dir(target), 0750); e != nil {
return e
}
in, e := os.Open(path)
if e != nil {
return e
}
out, e := os.OpenFile(target, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, info.Mode().Perm())
if e != nil {
_ = in.Close()
return e
}
_, copyErr := io.Copy(out, in)
inErr := in.Close()
outErr := out.Close()
if copyErr != nil {
return copyErr
}
if inErr != nil {
return inErr
}
return outErr
})
}
func (s *Service) encrypt(p []byte) ([]byte, error) {
b, e := aes.NewCipher(s.key)
if e != nil {
return nil, e
}
g, e := cipher.NewGCM(b)
if e != nil {
return nil, e
}
nonce := make([]byte, g.NonceSize())
if _, e = rand.Read(nonce); e != nil {
return nil, e
}
return g.Seal(nonce, nonce, p, nil), nil
}
func (s *Service) decrypt(v []byte) ([]byte, error) {
b, e := aes.NewCipher(s.key)
if e != nil {
return nil, e
}
g, e := cipher.NewGCM(b)
if e != nil {
return nil, e
}
if len(v) < g.NonceSize() {
return nil, errors.New("invalid encrypted secret")
}
return g.Open(nil, v[:g.NonceSize()], v[g.NonceSize():], nil)
}