Merge pull request 'fix(lifecycle): retain deployment diagnostics' (#39) from codex/block-12-deployment-diagnostics into main
CI / validate (push) Successful in 26m0s

Reviewed-on: #39
Reviewed-by: tony <1+tony@noreply.localhost>
This commit was merged in pull request #39.
This commit is contained in:
2026-08-15 23:04:06 +02:00
29 changed files with 1367 additions and 68 deletions
+3
View File
@@ -2,6 +2,9 @@
.cache
dist
data
servers
backups
.playwright-cli
secrets
*.db
*.db-shm
+2 -1
View File
@@ -1,6 +1,6 @@
schema_version: 1
id: palworld-official
version: 1.1.0
version: 1.1.1
source:
type: official
@@ -33,6 +33,7 @@ container:
tag: v1.0.2.101103
entrypoint:
- /pal/helper.sh
user_mode: image
arguments:
- -port=8211
- -useperfthreads
+25 -13
View File
@@ -55,14 +55,14 @@ func run(logger *slog.Logger) error {
return err
}
defer func() { _ = db.Close() }()
repository := sqlite.NewRepository(db)
if err := catalog.InitializeOfficial(templatesRoot, catalogdata.Files); err != nil {
return err
}
scan, err := catalog.ScanDir(templatesRoot)
scan, err := synchronizeCatalog(ctx, repository, templatesRoot)
if err != nil {
return err
}
repository := sqlite.NewRepository(db)
keyFile := environment("DOGAMA_MASTER_KEY_FILE", "secrets/master_key")
key, keyErr := loadMasterKey(keyFile)
if keyErr != nil {
@@ -71,9 +71,6 @@ func run(logger *slog.Logger) error {
if keyErr = repository.SetSecretKey(key); keyErr != nil {
return keyErr
}
if err := repository.Replace(ctx, scan.Valid); err != nil {
return err
}
logger.Info("local catalog synchronized", "event", "catalog.synchronized", "template_count", len(scan.Valid), "invalid_template_count", len(scan.Errors))
var handler http.Handler
var lifecycle *instance.LifecycleService
@@ -121,14 +118,7 @@ func run(logger *slog.Logger) error {
cancel()
}
handler, err = web.NewHandlerCompleteWithCatalogDeploymentAndRuntime(auth.New(db), repository, lifecycle, backupService, importService, auditService, notificationService, func(ctx context.Context) (catalog.ScanResult, error) {
result, scanErr := catalog.ScanDir(templatesRoot)
if scanErr != nil {
return result, scanErr
}
if syncErr := repository.Replace(ctx, result.Valid); syncErr != nil {
return result, syncErr
}
return result, nil
return synchronizeCatalog(ctx, repository, templatesRoot)
}, serversRoot, instance.NewModuleService(repository, repository, modulesRoot), logger)
if err != nil {
return err
@@ -168,6 +158,28 @@ func run(logger *slog.Logger) error {
}
}
// synchronizeCatalog indexes administrator-owned local templates first, then
// makes the currently shipped immutable snapshots active for new deployments.
// Replace retains older versions, so existing instances and diagnostics remain
// pinned to their original snapshot.
func synchronizeCatalog(ctx context.Context, repository *sqlite.Repository, templatesRoot string) (catalog.ScanResult, error) {
result, err := catalog.ScanDir(templatesRoot)
if err != nil {
return result, err
}
if err := repository.Replace(ctx, result.Valid); err != nil {
return result, err
}
current, err := catalog.LoadFS(catalogdata.Files, ".")
if err != nil {
return result, err
}
if err := repository.Sync(ctx, current); err != nil {
return result, err
}
return result, nil
}
func loadMasterKey(path string) ([]byte, error) {
if _, err := internalsecrets.Ensure(path, 0o600); err != nil {
return nil, errors.New("initialize encryption key file")
+3 -1
View File
@@ -73,7 +73,7 @@ Read this compact operational baseline before starting a milestone. Open detaile
- Notification delivery attempts are capped at five with exponential minute-scale backoff and never determine the originating operation result.
- Audit retention defaults to 30 days and 10,000 entries; zero explicitly selects unlimited retention/count within documented bounds.
- At least one active global administrator is always retained; deactivation revokes that user's sessions atomically.
- The local template directory (`/var/lib/dogama/templates`, under the application data bind mount) is the catalog source of truth. Bundled templates are copied only when their destination files are absent; an administrator Scan validates each directory independently and refreshes the available SQLite index without network fetches.
- The local template directory (`/var/lib/dogama/templates`, under the application data bind mount) is the catalog source of truth for administrator-owned customizations. Bundled templates are copied only when their destination files are absent; after every local scan, the current bundled immutable snapshots are synchronized into SQLite and selected for new deployments while older snapshots remain available for existing instances, audit and diagnostics.
- Administrators can persist bounded HTTP(S) template-repository definitions for future use. They are configuration only: remote retrieval, authentication, synchronization and automatic updates are deliberately unavailable, and Catalog Scan remains local-only.
## Known limitations and debt
@@ -99,3 +99,5 @@ Read this compact operational baseline before starting a milestone. Open detaile
- Update this file after every merged milestone or durable architectural change; keep it compact and remove stale statements.
- Web access and i18n are stored in the initial SQLite schema. HTTP defaults to working session/CSRF cookies without Secure and without HSTS; HTTPS enforcement is explicit.
- Lifecycle failures retain a bounded diagnostic record separate from Audit. The restricted agent captures Docker inspection state (including exit/OOM/timestamps/health) and a bounded log tail after a failed start; the application persists this record under the operation ID and writes a stable `DGM-*` error category. Existing SQLite stores receive the additive diagnostic table at open time.
- Approved template helper assets are immutable mode `0555`, so an image-defined non-root user can execute a bind-mounted entrypoint while retaining no write access.
+6 -2
View File
@@ -49,7 +49,7 @@ func writeImmutableAsset(path string, content []byte) error {
if existingDigest != wantedDigest {
return errors.New("approved asset content conflict")
}
return nil
return os.Chmod(path, 0o555)
} else if !errors.Is(err, os.ErrNotExist) {
return errors.New("inspect approved asset path")
}
@@ -59,7 +59,11 @@ func writeImmutableAsset(path string, content []byte) error {
}
temporaryPath := temporary.Name()
defer func() { _ = os.Remove(temporaryPath) }()
if err = temporary.Chmod(0o500); err == nil {
// Template entrypoint assets are bind-mounted into images that can run as a
// non-root image user. The host-side agent owns these files, so the image
// user must be able to read and execute an approved helper without gaining
// write access.
if err = temporary.Chmod(0o555); err == nil {
_, err = temporary.Write(content)
}
if err == nil {
+28
View File
@@ -0,0 +1,28 @@
package agent
import (
"os"
"path/filepath"
"testing"
)
func TestWriteImmutableAssetIsExecutableByImageUser(t *testing.T) {
path := filepath.Join(t.TempDir(), "helper.sh")
content := []byte("#!/bin/sh\nexit 0\n")
if err := writeImmutableAsset(path, content); err != nil {
t.Fatal(err)
}
if err := os.Chmod(path, 0o500); err != nil {
t.Fatal(err)
}
if err := writeImmutableAsset(path, content); err != nil {
t.Fatal(err)
}
info, err := os.Stat(path)
if err != nil {
t.Fatal(err)
}
if got, want := info.Mode().Perm(), os.FileMode(0o555); got != want {
t.Fatalf("asset mode = %04o, want %04o so an image user can execute the bind-mounted entrypoint", got, want)
}
}
+87 -12
View File
@@ -10,6 +10,7 @@ import (
"net"
"net/http"
"net/url"
"os"
"path/filepath"
"sort"
"strconv"
@@ -31,15 +32,25 @@ type DockerRuntime interface {
Restart(context.Context, string, int) error
Delete(context.Context, string) error
Inspect(context.Context, string) (DockerInspection, error)
Logs(context.Context, string, int) (string, error)
Stats(context.Context, string) (agentwire.InstanceStats, error)
}
type DockerInspection struct {
ContainerID string
Running bool
Health string
ExitCode int
Labels map[string]string
ContainerID string
Running bool
Restarting bool
Paused bool
OOMKilled bool
Dead bool
Status string
Error string
StartedAt string
FinishedAt string
RestartCount int
Health string
ExitCode int
Labels map[string]string
}
type AssetMount struct {
@@ -155,16 +166,24 @@ func (d *dockerRuntime) Create(ctx context.Context, plan agentwire.DeploymentPla
bindings[key] = []portBinding{{HostIP: "0.0.0.0", HostPort: strconv.Itoa(port.HostPort)}}
}
}
binds := make([]string, 0, len(plan.Mounts))
binds := make([]string, 0, len(plan.Mounts)+len(assets))
for _, mount := range plan.Mounts {
hostPath, hostErr := d.hostBindPath(ctx, mount.HostPath)
if hostErr != nil {
return "", hostErr
}
mode := "rw"
if mount.ReadOnly {
mode = "ro"
}
binds = append(binds, mount.HostPath+":"+mount.ContainerPath+":"+mode)
binds = append(binds, hostPath+":"+mount.ContainerPath+":"+mode)
}
for _, asset := range assets {
binds = append(binds, asset.HostPath+":"+asset.ContainerPath+":ro")
hostPath, hostErr := d.hostBindPath(ctx, asset.HostPath)
if hostErr != nil {
return "", hostErr
}
binds = append(binds, hostPath+":"+asset.ContainerPath+":ro")
}
pidsLimit := int64(512)
payload := struct {
@@ -224,6 +243,42 @@ func (d *dockerRuntime) Create(ctx context.Context, plan agentwire.DeploymentPla
return created.ID, nil
}
// hostBindPath translates an agent-container path to the exact host source
// recorded for this agent's own bind mount. Docker receives paths in the host
// namespace, not the agent container namespace.
func (d *dockerRuntime) hostBindPath(ctx context.Context, containerPath string) (string, error) {
self, err := os.Hostname()
if err != nil || self == "" {
return "", errors.New("agent container identity is unavailable")
}
response, err := d.call(ctx, http.MethodGet, dockerAPIVersion+"/containers/"+url.PathEscape(self)+"/json", nil, "", 128<<10)
if err != nil || response.status != http.StatusOK {
return "", errors.New("agent bind mount mapping is unavailable")
}
var inspection struct {
Mounts []struct {
Source string
Destination string
}
}
if json.Unmarshal(response.body, &inspection) != nil {
return "", errors.New("agent bind mount mapping is unavailable")
}
clean := filepath.Clean(containerPath)
for _, mount := range inspection.Mounts {
destination := filepath.Clean(mount.Destination)
if !filepath.IsAbs(mount.Source) || (clean != destination && !strings.HasPrefix(clean, destination+string(filepath.Separator))) {
continue
}
relative, relativeErr := filepath.Rel(destination, clean)
if relativeErr != nil || relative == ".." || strings.HasPrefix(relative, ".."+string(filepath.Separator)) {
return "", errors.New("agent bind mount mapping is unavailable")
}
return filepath.Join(mount.Source, relative), nil
}
return "", errors.New("agent bind mount mapping is unavailable")
}
func mergeDockerLabels(custom, technical map[string]string) map[string]string {
result := make(map[string]string, len(custom)+len(technical))
for key, value := range custom {
@@ -291,9 +346,18 @@ func (d *dockerRuntime) Inspect(ctx context.Context, id string) (DockerInspectio
Labels map[string]string `json:"Labels"`
} `json:"Config"`
State struct {
Running bool `json:"Running"`
ExitCode int `json:"ExitCode"`
Health *struct {
Running bool `json:"Running"`
Restarting bool `json:"Restarting"`
Paused bool `json:"Paused"`
OOMKilled bool `json:"OOMKilled"`
Dead bool `json:"Dead"`
Status string `json:"Status"`
Error string `json:"Error"`
StartedAt string `json:"StartedAt"`
FinishedAt string `json:"FinishedAt"`
RestartCount int `json:"RestartCount"`
ExitCode int `json:"ExitCode"`
Health *struct {
Status string `json:"Status"`
} `json:"Health"`
} `json:"State"`
@@ -305,7 +369,18 @@ func (d *dockerRuntime) Inspect(ctx context.Context, id string) (DockerInspectio
if payload.State.Health != nil {
health = payload.State.Health.Status
}
return DockerInspection{ContainerID: payload.ID, Running: payload.State.Running, Health: health, ExitCode: payload.State.ExitCode, Labels: payload.Config.Labels}, nil
return DockerInspection{ContainerID: payload.ID, Running: payload.State.Running, Restarting: payload.State.Restarting, Paused: payload.State.Paused, OOMKilled: payload.State.OOMKilled, Dead: payload.State.Dead, Status: payload.State.Status, Error: payload.State.Error, StartedAt: payload.State.StartedAt, FinishedAt: payload.State.FinishedAt, RestartCount: payload.State.RestartCount, Health: health, ExitCode: payload.State.ExitCode, Labels: payload.Config.Labels}, nil
}
// Logs returns a bounded tail. Docker multiplexes stream frames only when TTY
// is enabled; game templates do not enable TTY, so the raw tail is still useful
// even if Docker returns no stdout/stderr at all.
func (d *dockerRuntime) Logs(ctx context.Context, id string, tail int) (string, error) {
response, err := d.call(ctx, http.MethodGet, dockerAPIVersion+"/containers/"+url.PathEscape(id)+"/logs?stdout=true&stderr=true&tail="+strconv.Itoa(tail), nil, "", 64<<10)
if err != nil || response.status != http.StatusOK {
return "", errors.New("container log retrieval failed")
}
return strings.TrimSpace(string(response.body)), nil
}
func (d *dockerRuntime) Stats(ctx context.Context, id string) (agentwire.InstanceStats, error) {
+27
View File
@@ -6,6 +6,7 @@ import (
"io"
"net"
"net/http"
"os"
"path/filepath"
"strings"
"sync"
@@ -87,6 +88,14 @@ func TestDockerPingerUsesConfiguredUnixSocket(t *testing.T) {
}
func TestDockerRuntimeCreatesFixedSecurityBaseline(t *testing.T) {
hostRoot := t.TempDir()
if err := os.MkdirAll(filepath.Join(hostRoot, ".dogama"), 0o755); err != nil {
t.Fatal(err)
}
hostAsset := filepath.Join(hostRoot, ".dogama", "helper")
if err := os.WriteFile(hostAsset, []byte("#!/bin/sh\n"), 0o555); err != nil {
t.Fatal(err)
}
socket := filepath.Join(t.TempDir(), "docker.sock")
listener, err := net.Listen("unix", socket)
if err != nil {
@@ -95,6 +104,8 @@ func TestDockerRuntimeCreatesFixedSecurityBaseline(t *testing.T) {
createdBodies := make(chan []byte, 2)
server := &http.Server{Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/v1.41/containers/" + hostname(t) + "/json":
_ = json.NewEncoder(w).Encode(map[string]any{"Mounts": []map[string]string{{"Source": hostRoot, "Destination": "/srv/games"}}})
case "/v1.41/images/create":
_, _ = w.Write([]byte("{}\n"))
case "/v1.41/containers/json":
@@ -136,6 +147,7 @@ func TestDockerRuntimeCreatesFixedSecurityBaseline(t *testing.T) {
User string `json:"User"`
Labels map[string]string `json:"Labels"`
HostConfig struct {
Binds []string
NetworkMode string `json:"NetworkMode"`
CapDrop []string `json:"CapDrop"`
SecurityOpt []string `json:"SecurityOpt"`
@@ -157,6 +169,21 @@ func TestDockerRuntimeCreatesFixedSecurityBaseline(t *testing.T) {
if payload.Labels["dashboard.name"] != "Summer" || payload.User != "1000:1001" {
t.Fatalf("custom Docker configuration = %#v user=%q", payload.Labels, payload.User)
}
if got := strings.Join(payload.HostConfig.Binds, " "); !strings.Contains(got, filepath.Join(hostRoot, "saved")+":/game/saved:rw") || !strings.Contains(got, hostAsset+":/pal/helper.sh:ro") {
t.Fatalf("host bind mappings = %q", got)
}
if info, err := os.Stat(hostAsset); err != nil || !info.Mode().IsRegular() {
t.Fatalf("asset became non-file after host bind translation: info=%#v error=%v", info, err)
}
}
func hostname(t *testing.T) string {
t.Helper()
value, err := os.Hostname()
if err != nil {
t.Fatal(err)
}
return value
}
func TestDockerPingerDoesNotExposeConnectionDetails(t *testing.T) {
+10 -4
View File
@@ -46,8 +46,11 @@ type ApprovedAsset struct {
func (p *PlanPolicy) Assets(plan agentwire.DeploymentPlan) ([]ApprovedAsset, error) {
snapshot, ok := p.snapshots[plan.TemplateID+"@"+plan.TemplateVersion]
if !ok || snapshot.Digest != plan.TemplateDigest {
return nil, errors.New("unknown template snapshot")
if !ok {
return nil, fmt.Errorf("template snapshot not registered: %s@%s", plan.TemplateID, plan.TemplateVersion)
}
if snapshot.Digest != plan.TemplateDigest {
return nil, fmt.Errorf("template snapshot digest mismatch: received=%s expected=%s", plan.TemplateDigest, snapshot.Digest)
}
result := make([]ApprovedAsset, 0, len(snapshot.Template.Container.Assets))
for _, asset := range snapshot.Template.Container.Assets {
@@ -68,8 +71,11 @@ func (p *PlanPolicy) Validate(plan agentwire.DeploymentPlan) error {
return errors.New("invalid deployment plan")
}
snapshot, ok := p.snapshots[plan.TemplateID+"@"+plan.TemplateVersion]
if !ok || snapshot.Digest != plan.TemplateDigest {
return errors.New("unknown template snapshot")
if !ok {
return fmt.Errorf("template snapshot not registered: %s@%s", plan.TemplateID, plan.TemplateVersion)
}
if snapshot.Digest != plan.TemplateDigest {
return fmt.Errorf("template snapshot digest mismatch: received=%s expected=%s", plan.TemplateDigest, snapshot.Digest)
}
template := snapshot.Template
imagePrefix := template.Container.Image + ":"
+45
View File
@@ -0,0 +1,45 @@
package agent
import (
"path/filepath"
"testing"
catalogdata "git.zaynet.fr/DoGaMa/DoGaMa-serv/catalog"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/catalog"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/instance"
)
func TestEmbeddedPalworldSnapshotMatchesApplicationDeploymentPlan(t *testing.T) {
snapshots, err := catalog.LoadFS(catalogdata.Files, ".")
if err != nil || len(snapshots) != 1 {
t.Fatalf("embedded snapshots = %#v, error = %v", snapshots, err)
}
snapshot := snapshots[0]
if snapshot.Template.ID != "palworld-official" || snapshot.Template.Version != "1.1.1" {
t.Fatalf("embedded Palworld snapshot = %s@%s", snapshot.Template.ID, snapshot.Template.Version)
}
preview, err := instance.BuildPreview(snapshot, instance.PreviewRequest{
DisplayName: "Snapshot consistency", Slug: "snapshot-consistency",
HostPorts: map[string]int{"game": 38211},
MountPaths: map[string]string{"saved": filepath.Join(t.TempDir(), "saved")},
DataOrigin: "new", BackupRetention: 7,
})
if err != nil {
t.Fatal(err)
}
plan, err := preview.DeploymentPlan("abcdefghijklmnopqrstuvwx")
if err != nil {
t.Fatal(err)
}
policy, err := NewPlanPolicy(snapshots, catalogdata.Files)
if err != nil {
t.Fatal(err)
}
known, ok := policy.snapshots[plan.TemplateID+"@"+plan.TemplateVersion]
if !ok || known.Digest != plan.TemplateDigest {
t.Fatalf("agent snapshot=%#v plan=%s@%s digest=%s", known, plan.TemplateID, plan.TemplateVersion, plan.TemplateDigest)
}
if err := policy.Validate(plan); err != nil {
t.Fatalf("application plan rejected by matching embedded snapshot: %v", err)
}
}
+83 -9
View File
@@ -8,6 +8,8 @@ import (
"io"
"log/slog"
"net/http"
"regexp"
"strings"
"syscall"
"time"
@@ -16,6 +18,8 @@ import (
const maxDiskPaths = 16
var diagnosticSecretPattern = regexp.MustCompile(`(?i)(password|token|secret|api[_-]?key)\s*[:=]\s*[^\s,;]+`)
type service struct {
paths *PathPolicy
registry *Registry
@@ -137,8 +141,13 @@ func (s *service) checkPorts(w http.ResponseWriter, r *http.Request) {
func (s *service) createInstance(w http.ResponseWriter, r *http.Request) {
var plan agentwire.DeploymentPlan
if decodeJSON(r.Body, &plan) != nil || s.plans.Validate(plan) != nil {
writeProblem(w, http.StatusUnprocessableEntity, "invalid_plan", "The deployment plan is invalid.")
if err := decodeJSON(r.Body, &plan); err != nil {
writeDiagnosticProblem(w, http.StatusUnprocessableEntity, "invalid_plan", "The deployment plan is invalid.", map[string]any{"phase": "preparation", "operation_error": redactDiagnostic(err.Error())})
return
}
if err := s.plans.Validate(plan); err != nil {
s.logger.Error("instance plan rejected", "event", "instance.create.plan_rejected", "instance_id", plan.InstanceID, "template_id", plan.TemplateID, "template_version", plan.TemplateVersion, "template_digest", plan.TemplateDigest, "error", err)
writeDiagnosticProblem(w, http.StatusUnprocessableEntity, "invalid_plan", "The deployment plan is invalid.", map[string]any{"phase": "preparation", "template_id": plan.TemplateID, "template_version": plan.TemplateVersion, "template_digest": plan.TemplateDigest, "operation_error": redactDiagnostic(err.Error())})
return
}
for index := range plan.Mounts {
@@ -179,7 +188,8 @@ func (s *service) createInstance(w http.ResponseWriter, r *http.Request) {
}
containerID, err := s.docker.Create(r.Context(), plan, assets)
if err != nil {
writeProblem(w, http.StatusBadGateway, "container_create_failed", "The container could not be created.")
s.logger.Error("instance create failed", "event", "instance.create.failed", "instance_id", plan.InstanceID, "error", err)
writeDiagnosticProblem(w, http.StatusBadGateway, "container_create_failed", "The container could not be created.", map[string]any{"operation_error": redactDiagnostic(err.Error())})
return
}
entry := RegisteredInstance{InstanceID: plan.InstanceID, ContainerID: containerID, PlanDigest: plan.PlanDigest}
@@ -288,10 +298,25 @@ func (s *service) startInstance(w http.ResponseWriter, r *http.Request) {
}
if !state.Running {
if err := s.docker.Start(r.Context(), entry.ContainerID); err != nil {
writeProblem(w, http.StatusBadGateway, "start_failed", "The registered container could not be started.")
s.logger.Error("instance start failed", "event", "instance.start.failed", "instance_id", entry.InstanceID, "container_id", entry.ContainerID, "error", err)
writeDiagnosticProblem(w, http.StatusBadGateway, "start_failed", "The registered container could not be started.", s.containerDiagnostic(r.Context(), entry.ContainerID, err))
return
}
state.Running, state.Health = true, "starting"
inspection, inspectErr := s.docker.Inspect(r.Context(), entry.ContainerID)
if inspectErr != nil || (!inspection.Running && (inspection.Status != "" || inspection.FinishedAt != "")) {
s.logger.Error("instance exited after start", "event", "container.exited", "instance_id", entry.InstanceID, "container_id", entry.ContainerID, "error", err)
writeDiagnosticProblem(w, http.StatusBadGateway, "start_exited", "The registered container exited immediately after start.", s.containerDiagnostic(r.Context(), entry.ContainerID, inspectErr))
return
}
if inspection.Running {
state, err = s.boundState(r.Context(), entry)
if err != nil {
writeDiagnosticProblem(w, http.StatusBadGateway, "start_inspect_failed", "The started container could not be inspected.", s.containerDiagnostic(r.Context(), entry.ContainerID, err))
return
}
} else {
state.Running, state.Health = true, "starting"
}
}
writeJSON(w, http.StatusOK, state)
}
@@ -396,13 +421,57 @@ func (s *service) boundState(ctx context.Context, entry RegisteredInstance) (age
if !inspection.Running {
health = "stopped"
}
return agentwire.InstanceState{InstanceID: entry.InstanceID, ContainerID: entry.ContainerID, PlanDigest: entry.PlanDigest, Running: inspection.Running, Ready: inspection.Running && inspection.Health == "healthy", Health: health, ExitCode: inspection.ExitCode}, nil
return agentwire.InstanceState{InstanceID: entry.InstanceID, ContainerID: entry.ContainerID, PlanDigest: entry.PlanDigest, Running: inspection.Running, Ready: inspection.Running && (inspection.Health == "healthy" || inspection.Health == "none"), Health: health, ExitCode: inspection.ExitCode}, nil
}
func (s *service) bindingProblem(w http.ResponseWriter) {
writeProblem(w, http.StatusConflict, "registration_mismatch", "The registered container binding is invalid.")
}
// containerDiagnostic deliberately contains Docker state rather than request
// payloads: it is safe to return over the authenticated private agent link and
// remains useful when stdout/stderr is empty.
func (s *service) containerDiagnostic(ctx context.Context, containerID string, cause error) map[string]any {
diagnostic := map[string]any{"container_id": containerID}
if cause != nil {
diagnostic["operation_error"] = redactDiagnostic(cause.Error())
}
inspection, err := s.docker.Inspect(ctx, containerID)
if err != nil {
diagnostic["inspect_error"] = redactDiagnostic(err.Error())
return diagnostic
}
diagnostic["state"] = inspection.Status
diagnostic["running"] = inspection.Running
diagnostic["restarting"] = inspection.Restarting
diagnostic["paused"] = inspection.Paused
diagnostic["oom_killed"] = inspection.OOMKilled
diagnostic["dead"] = inspection.Dead
diagnostic["exit_code"] = inspection.ExitCode
diagnostic["error"] = redactDiagnostic(inspection.Error)
diagnostic["started_at"] = inspection.StartedAt
diagnostic["finished_at"] = inspection.FinishedAt
diagnostic["health"] = inspection.Health
diagnostic["restart_count"] = inspection.RestartCount
if logs, logsErr := s.docker.Logs(ctx, containerID, 100); logsErr == nil {
diagnostic["logs_tail"] = redactDiagnostic(logs)
} else {
diagnostic["logs_error"] = redactDiagnostic(logsErr.Error())
}
return diagnostic
}
func redactDiagnostic(value string) string {
return diagnosticSecretPattern.ReplaceAllStringFunc(value, func(match string) string {
separator := "="
if strings.Contains(match, ":") {
separator = ":"
}
parts := strings.SplitN(match, separator, 2)
return parts[0] + separator + "[REDACTED]"
})
}
func (s *service) headers(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Cache-Control", "no-store")
@@ -431,8 +500,13 @@ func writeJSON(w http.ResponseWriter, status int, value any) {
}
func writeProblem(w http.ResponseWriter, status int, code, message string) {
writeDiagnosticProblem(w, status, code, message, nil)
}
func writeDiagnosticProblem(w http.ResponseWriter, status int, code, message string, details map[string]any) {
writeJSON(w, status, struct {
Code string `json:"code"`
Message string `json:"message"`
}{Code: code, Message: message})
Code string `json:"code"`
Message string `json:"message"`
Details map[string]any `json:"details,omitempty"`
}{Code: code, Message: message, Details: details})
}
+61
View File
@@ -6,11 +6,13 @@ import (
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"io"
"log/slog"
"net/http"
"net/http/httptest"
"path/filepath"
"strings"
"testing"
"time"
@@ -24,6 +26,7 @@ import (
type fakeDocker struct {
err error
inspection agent.DockerInspection
logs string
}
type fakeDisk struct{}
@@ -53,6 +56,7 @@ func (d fakeDocker) Inspect(_ context.Context, id string) (agent.DockerInspectio
}
return inspection, nil
}
func (d fakeDocker) Logs(context.Context, string, int) (string, error) { return d.logs, d.err }
func (d fakeDocker) Stats(context.Context, string) (agentwire.InstanceStats, error) {
return agentwire.InstanceStats{MemoryBytes: 42}, d.err
}
@@ -161,6 +165,63 @@ func TestAgentCreatesOnlyValidatedBoundInstances(t *testing.T) {
}
}
func TestAgentTreatsRunningContainerWithoutDockerHealthcheckAsReady(t *testing.T) {
root := t.TempDir()
secret := bytes.Repeat([]byte{0x31}, 32)
instanceID := "abcdefghijklmnopqrstuvwx"
plan := testPlan(t, instanceID, root)
docker := fakeDocker{inspection: agent.DockerInspection{
ContainerID: "container-" + instanceID, Running: true, Status: "running", Health: "none",
Labels: map[string]string{"io.dogama.managed": "true", "io.dogama.instance-id": instanceID, "io.dogama.plan-digest": plan.PlanDigest},
}}
server := httptest.NewServer(newTestHandler(t, root, secret, docker))
t.Cleanup(server.Close)
client, err := agentclient.New(server.URL, secret, server.Client())
if err != nil {
t.Fatal(err)
}
if _, err := client.CreateInstance(context.Background(), plan); err != nil {
t.Fatal(err)
}
state, err := client.InspectInstance(context.Background(), instanceID)
if err != nil || !state.Running || !state.Ready || state.Health != "none" {
t.Fatalf("state=%#v error=%v", state, err)
}
}
func TestAgentStartExitReturnsDockerDiagnosticWithoutLogs(t *testing.T) {
root := t.TempDir()
secret := bytes.Repeat([]byte{0x31}, 32)
instanceID := "abcdefghijklmnopqrstuvwx"
plan := testPlan(t, instanceID, root)
docker := fakeDocker{inspection: agent.DockerInspection{
ContainerID: "container-" + instanceID, Status: "exited", ExitCode: 23,
StartedAt: "2026-01-01T00:00:00Z", FinishedAt: "2026-01-01T00:00:01Z",
OOMKilled: true, Health: "unhealthy", Error: "permission denied password=never-leak",
Labels: map[string]string{"io.dogama.managed": "true", "io.dogama.instance-id": instanceID, "io.dogama.plan-digest": plan.PlanDigest},
}}
server := httptest.NewServer(newTestHandler(t, root, secret, docker))
t.Cleanup(server.Close)
client, err := agentclient.New(server.URL, secret, server.Client())
if err != nil {
t.Fatal(err)
}
if _, err := client.CreateInstance(context.Background(), plan); err != nil {
t.Fatal(err)
}
_, err = client.StartInstance(context.Background(), instanceID)
var problem *agentclient.ProblemError
if !errors.As(err, &problem) || problem.Code != "start_exited" {
t.Fatalf("start error = %#v", err)
}
if problem.Details["exit_code"] != float64(23) || problem.Details["oom_killed"] != true || problem.Details["logs_tail"] != "" {
t.Fatalf("diagnostic = %#v", problem.Details)
}
if strings.Contains(fmt.Sprint(problem.Details), "never-leak") {
t.Fatalf("secret leaked in diagnostic: %#v", problem.Details)
}
}
func TestAgentRejectsPlanSubstitutionAndEscapingMount(t *testing.T) {
root := t.TempDir()
secret := bytes.Repeat([]byte{0x31}, 32)
+17
View File
@@ -48,12 +48,29 @@ type ProblemError struct {
Status int
Code string
Message string
Details map[string]any `json:"details,omitempty"`
}
func (e *ProblemError) Error() string {
return fmt.Sprintf("agent request failed: %s (HTTP %d)", e.Code, e.Status)
}
// DiagnosticDetails lets the lifecycle service persist the agent's bounded
// technical cause without exposing it to unprivileged HTTP clients.
func (e *ProblemError) DiagnosticDetails() map[string]any {
details := make(map[string]any, len(e.Details)+2)
for key, value := range e.Details {
details[key] = value
}
if e.Code != "" {
details["agent_code"] = e.Code
}
if e.Message != "" {
details["agent_message"] = e.Message
}
return details
}
// New constructs a client. The base URL must not contain credentials, a query
// or a path beyond an optional trailing slash.
func New(baseURL string, secret []byte, httpClient *http.Client) (*Client, error) {
+1
View File
@@ -66,6 +66,7 @@ type Template struct {
Image string `json:"image"`
Tag string `json:"tag"`
Entrypoint []string `json:"entrypoint,omitempty"`
UserMode string `json:"user_mode,omitempty"`
Arguments []string `json:"arguments,omitempty"`
Environment map[string]string `json:"environment,omitempty"`
StopTimeoutSeconds int `json:"stop_timeout_seconds"`
+1 -1
View File
@@ -50,7 +50,7 @@ func TestStageValidZIPAndRejectTraversal(t *testing.T) {
if err != nil {
t.Fatal(err)
}
policy := importexport.Policy{TemplateID: "palworld-official", TemplateVersion: "1.1.0", AcceptedFormats: []string{"zip"}, MaxExpandedBytes: 1 << 20, RequiredPaths: []string{"Level.sav", "Players"}}
policy := importexport.Policy{TemplateID: snapshots[0].Template.ID, TemplateVersion: snapshots[0].Template.Version, AcceptedFormats: []string{"zip"}, MaxExpandedBytes: 1 << 20, RequiredPaths: []string{"Level.sav", "Players"}}
valid := zipBytes(t, map[string]string{"Save/Level.sav": "world", "Save/Players/player.sav": "player"})
result, err := service.Stage(ctx, actor.ID, "zip", bytes.NewReader(valid), policy)
if err != nil {
+231 -2
View File
@@ -4,9 +4,11 @@ import (
"context"
"crypto/rand"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"sync"
"time"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/agentwire"
)
@@ -17,6 +19,8 @@ var (
ErrInvalidState = errors.New("invalid instance lifecycle state")
)
const startupVerificationDelay = 3 * time.Second
type StoredInstance struct {
ID string
Preview Preview
@@ -38,6 +42,59 @@ type OperationResult struct {
AgentState agentwire.InstanceState `json:"agent_state,omitempty"`
}
// Diagnostic is a bounded, redacted technical record for an operation. It is
// intentionally separate from audit: audit records who did something;
// diagnostics retain why it failed.
type Diagnostic struct {
OperationID string
InstanceID string
Step string
ErrorCode string
Details string
CreatedAt string
}
type OperationHistory struct {
OperationID string `json:"operation_id"`
Kind string `json:"kind"`
State string `json:"state"`
Phase string `json:"phase"`
ErrorCode string `json:"error_code,omitempty"`
CreatedAt string `json:"created_at"`
CompletedAt string `json:"completed_at,omitempty"`
Diagnostic *Diagnostic `json:"diagnostic,omitempty"`
}
// OperationProgress is the non-technical, persisted view consumed by the UI.
// Diagnostics deliberately remain a separate administrator-only resource.
type OperationProgress struct {
OperationID string `json:"operation_id"`
OperationType string `json:"operation_type"`
Game string `json:"game"`
InstanceID string `json:"instance_id,omitempty"`
GlobalStatus string `json:"global_status"`
CurrentStep string `json:"current_step,omitempty"`
Steps []OperationStep `json:"steps"`
ErrorCode string `json:"error_code,omitempty"`
StartedAt string `json:"started_at"`
FinishedAt string `json:"finished_at,omitempty"`
}
type OperationStep struct {
ID string `json:"id"`
Status string `json:"status"`
}
type ProgressRepository interface {
SetOperationStep(context.Context, string, string, string) error
GetOperationProgress(context.Context, string) (OperationProgress, error)
}
type DiagnosticRepository interface {
RecordDiagnostic(context.Context, Diagnostic) error
ListDiagnostics(context.Context, string, int) ([]Diagnostic, error)
ListOperationHistory(context.Context, string, int) ([]OperationHistory, error)
}
type LifecycleRepository interface {
GetInstance(context.Context, string) (StoredInstance, error)
ListLifecycleInstances(context.Context) ([]StoredInstance, error)
@@ -105,6 +162,9 @@ func (s *LifecycleService) Install(ctx context.Context, instanceID string) (Oper
}
state, err := s.agent.CreateInstance(ctx, plan)
if err != nil {
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, "installation", "failed")
}
return s.fail(ctx, operationID, instanceID, "agent_create_failed", err)
}
if err := s.repository.FinishOperation(ctx, operationID, "stopped", "stopped", state.ContainerID, plan.PlanDigest, false, ""); err != nil {
@@ -114,6 +174,130 @@ func (s *LifecycleService) Install(ctx context.Context, instanceID string) (Oper
})
}
// BeginInstall allocates the existing operation identifier before any long
// running work. Callers may safely return it to HTTP clients and run InstallOperation later.
func (s *LifecycleService) BeginInstall(ctx context.Context, instanceID string) (OperationResult, error) {
return s.exclusive(instanceID, func() (OperationResult, error) {
operationID, err := operationToken()
if err != nil {
return OperationResult{}, err
}
current, err := s.repository.BeginOperation(ctx, operationID, instanceID, "install", "installing")
if err != nil {
return OperationResult{}, err
}
return resultFrom(current, operationID), nil
})
}
// InstallOperation continues an operation already created by BeginInstall.
func (s *LifecycleService) InstallOperation(ctx context.Context, instanceID, operationID string) (OperationResult, error) {
return s.exclusive(instanceID, func() (OperationResult, error) {
current, err := s.repository.GetInstance(ctx, instanceID)
if err != nil {
return OperationResult{}, err
}
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, "installation", "running")
}
plan, err := deploymentPlan(ctx, s.repository, current)
if err != nil {
return s.fail(ctx, operationID, instanceID, "invalid_plan", err)
}
state, err := s.agent.CreateInstance(ctx, plan)
if err != nil {
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, "installation", "failed")
}
return s.fail(ctx, operationID, instanceID, "agent_create_failed", err)
}
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, "installation", "success")
}
return OperationResult{OperationID: operationID, InstanceID: instanceID, State: "stopped", Observed: "stopped", ContainerID: state.ContainerID, AgentState: state}, nil
})
}
// CompleteInstall records the final state after the deployment worker has run
// imports, configuration and the start action under the same operation id.
func (s *LifecycleService) CompleteInstall(ctx context.Context, result OperationResult) error {
current, err := s.repository.GetInstance(ctx, result.InstanceID)
if err != nil {
return err
}
return s.repository.FinishOperation(ctx, result.OperationID, result.State, result.Observed, result.ContainerID, current.PlanDigest, true, "")
}
// StartInstall starts the container without allocating a second operation.
func (s *LifecycleService) StartInstall(ctx context.Context, instanceID, operationID string) (OperationResult, error) {
return s.exclusive(instanceID, func() (OperationResult, error) {
if _, err := s.repository.GetInstance(ctx, instanceID); err != nil {
return OperationResult{}, err
}
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, "starting", "running")
}
state, err := s.agent.StartInstance(ctx, instanceID)
if err != nil {
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, "starting", "failed")
}
return s.fail(ctx, operationID, instanceID, "agent_start_failed", err)
}
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, "starting", "success")
_ = progress.SetOperationStep(ctx, operationID, "verification", "running")
}
verified := state
if !state.Ready {
var verifyErr error
verified, verifyErr = s.verifyStarted(ctx, instanceID)
if verifyErr != nil {
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, "verification", "failed")
}
return s.fail(ctx, operationID, instanceID, "verification_failed", verifyErr)
}
}
lifecycle, observed := stateToLifecycle(verified)
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, "verification", "success")
}
return OperationResult{OperationID: operationID, InstanceID: instanceID, State: lifecycle, Observed: observed, ContainerID: verified.ContainerID, AgentState: verified}, nil
})
}
func (s *LifecycleService) verifyStarted(ctx context.Context, instanceID string) (agentwire.InstanceState, error) {
timer := time.NewTimer(startupVerificationDelay)
defer timer.Stop()
select {
case <-ctx.Done():
return agentwire.InstanceState{}, ctx.Err()
case <-timer.C:
}
state, err := s.agent.InspectInstance(ctx, instanceID)
if err != nil {
return agentwire.InstanceState{}, err
}
if !state.Running || state.Health == "unhealthy" {
return agentwire.InstanceState{}, errors.New("container did not remain running during startup verification")
}
return state, nil
}
func (s *LifecycleService) FailInstall(ctx context.Context, operationID, instanceID, step string, cause error) (OperationResult, error) {
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, step, "failed")
}
return s.fail(ctx, operationID, instanceID, step, cause)
}
func (s *LifecycleService) MarkStep(ctx context.Context, operationID, step, status string) {
if p, ok := s.repository.(ProgressRepository); ok {
_ = p.SetOperationStep(ctx, operationID, step, status)
}
}
func deploymentPlan(ctx context.Context, repository LifecycleRepository, current StoredInstance) (agentwire.DeploymentPlan, error) {
if secrets, ok := repository.(SecretRepository); ok && (len(current.Preview.ResolvedConfiguration.SecretEnvironment) != 0 || len(current.Preview.ResolvedConfiguration.SecretArguments) != 0) {
values, err := secrets.LoadInstanceSecrets(ctx, current.ID)
@@ -335,10 +519,55 @@ func (s *LifecycleService) inspectAndPersist(ctx context.Context, current Stored
}
func (s *LifecycleService) fail(ctx context.Context, operationID, instanceID, code string, cause error) (OperationResult, error) {
if err := s.repository.FailOperation(ctx, operationID, "error", code); err != nil {
stableCode := stableErrorCode(code)
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, code, "failed")
}
if diagnostics, ok := s.repository.(DiagnosticRepository); ok {
_ = diagnostics.RecordDiagnostic(ctx, Diagnostic{OperationID: operationID, InstanceID: instanceID, Step: code, ErrorCode: stableCode, Details: diagnosticDetails(operationID, code, stableCode, cause)})
}
if err := s.repository.FailOperation(ctx, operationID, "error", stableCode); err != nil {
return OperationResult{}, err
}
return OperationResult{OperationID: operationID, InstanceID: instanceID, State: "error", Observed: "unknown"}, fmt.Errorf("%s: %w", code, cause)
return OperationResult{OperationID: operationID, InstanceID: instanceID, State: "error", Observed: "unknown"}, fmt.Errorf("%s: %w", stableCode, cause)
}
func diagnosticDetails(operationID, step, stableCode string, cause error) string {
details := map[string]any{
"operation_id": operationID,
"step": step,
"code": stableCode,
}
type detailed interface{ DiagnosticDetails() map[string]any }
var value detailed
if errors.As(cause, &value) {
for key, detail := range value.DiagnosticDetails() {
details[key] = detail
}
} else if cause != nil {
details["operation_error"] = cause.Error()
}
if _, exists := details["phase"]; !exists {
details["phase"] = step
}
encoded, err := json.Marshal(details)
if err != nil {
return ""
}
return string(encoded)
}
func stableErrorCode(step string) string {
switch step {
case "agent_start_failed", "agent_restart_failed", "verification_failed":
return "DGM-START-001"
case "agent_create_failed":
return "DGM-DEPLOY-001"
case "invalid_plan":
return "DGM-DEPLOY-002"
default:
return "DGM-DOCKER-001"
}
}
func (s *LifecycleService) exclusive(instanceID string, action func() (OperationResult, error)) (OperationResult, error) {
+116
View File
@@ -8,6 +8,7 @@ import (
"testing"
catalogdata "git.zaynet.fr/DoGaMa/DoGaMa-serv/catalog"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/agentclient"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/agentwire"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/catalog"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/instance"
@@ -20,6 +21,29 @@ type lifecycleAgent struct {
created int
}
type failingStartAgent struct{ lifecycleAgent }
type rejectedPlanAgent struct{ lifecycleAgent }
func (a *failingStartAgent) StartInstance(context.Context, string) (agentwire.InstanceState, error) {
return agentwire.InstanceState{}, &agentclient.ProblemError{Status: 502, Code: "start_exited", Details: map[string]any{"state": "exited", "exit_code": 42, "started_at": "2026-01-01T00:00:00Z", "finished_at": "2026-01-01T00:00:01Z", "logs_tail": "", "error": "permission denied password=[REDACTED]"}}
}
func (a *rejectedPlanAgent) CreateInstance(context.Context, agentwire.DeploymentPlan) (agentwire.InstanceState, error) {
return agentwire.InstanceState{}, &agentclient.ProblemError{
Status: 422,
Code: "invalid_plan",
Message: "The deployment plan is invalid.",
Details: map[string]any{
"phase": "preparation",
"template_id": "palworld-official",
"template_version": "1.1.1",
"template_digest": "digest-not-secret",
"operation_error": "template snapshot not registered: palworld-official@1.1.1",
},
}
}
func (a *lifecycleAgent) CreateInstance(_ context.Context, plan agentwire.DeploymentPlan) (agentwire.InstanceState, error) {
a.mu.Lock()
defer a.mu.Unlock()
@@ -162,3 +186,95 @@ func TestLifecycleInstallStartStopAndSafeContainerDeletion(t *testing.T) {
t.Fatalf("instance intent/player paths were removed: %v", err)
}
}
func TestManualStartPersistsAgentDiagnosticAndStableCode(t *testing.T) {
ctx := context.Background()
db, err := sqlite.Open(ctx, filepath.Join(t.TempDir(), "dogama.db"))
if err != nil {
t.Fatal(err)
}
defer db.Close()
repository := sqlite.NewRepository(db)
snapshots, err := catalog.LoadFS(catalogdata.Files, ".")
if err != nil {
t.Fatal(err)
}
if err := repository.Sync(ctx, snapshots); err != nil {
t.Fatal(err)
}
preview, err := instance.BuildPreview(snapshots[0], instance.PreviewRequest{DisplayName: "Failure", Slug: "failure", HostPorts: map[string]int{"game": 38212}, MountPaths: map[string]string{"saved": filepath.Join(t.TempDir(), "saved")}, DataOrigin: "new", BackupRetention: 7})
if err != nil {
t.Fatal(err)
}
id := "abcdefghijklmnopqrstuvwx"
if err := repository.CreateDraft(ctx, instance.Draft{ID: id, Preview: preview}); err != nil {
t.Fatal(err)
}
agent := &failingStartAgent{}
service := instance.NewLifecycleService(repository, agent)
if _, err := service.Install(ctx, id); err != nil {
t.Fatal(err)
}
result, startErr := service.Start(ctx, id)
if startErr == nil || !strings.Contains(startErr.Error(), "DGM-START-001") || result.OperationID == "" {
t.Fatalf("result=%#v err=%v", result, startErr)
}
diagnostics, err := repository.ListDiagnostics(ctx, id, 10)
if err != nil {
t.Fatal(err)
}
if len(diagnostics) != 1 || diagnostics[0].OperationID != result.OperationID || diagnostics[0].ErrorCode != "DGM-START-001" {
t.Fatalf("diagnostics=%#v", diagnostics)
}
if !strings.Contains(diagnostics[0].Details, `"exit_code":42`) || strings.Contains(diagnostics[0].Details, "password=") && !strings.Contains(diagnostics[0].Details, "[REDACTED]") {
t.Fatalf("details=%s", diagnostics[0].Details)
}
if strings.Contains(startErr.Error(), "The instance operation could not be completed") {
t.Fatalf("original cause was lost: %v", startErr)
}
}
func TestInstallPlanRejectedPersistsDiagnostic(t *testing.T) {
ctx := context.Background()
db, err := sqlite.Open(ctx, filepath.Join(t.TempDir(), "dogama.db"))
if err != nil {
t.Fatal(err)
}
defer db.Close()
repository := sqlite.NewRepository(db)
snapshots, err := catalog.LoadFS(catalogdata.Files, ".")
if err != nil {
t.Fatal(err)
}
if err := repository.Sync(ctx, snapshots); err != nil {
t.Fatal(err)
}
preview, err := instance.BuildPreview(snapshots[0], instance.PreviewRequest{DisplayName: "Rejected", Slug: "rejected", HostPorts: map[string]int{"game": 38213}, MountPaths: map[string]string{"saved": filepath.Join(t.TempDir(), "saved")}, DataOrigin: "new", BackupRetention: 7})
if err != nil {
t.Fatal(err)
}
const id = "abcdefghijklmnopqrstuvwx"
if err := repository.CreateDraft(ctx, instance.Draft{ID: id, Preview: preview}); err != nil {
t.Fatal(err)
}
result, installErr := instance.NewLifecycleService(repository, &rejectedPlanAgent{}).Install(ctx, id)
if installErr == nil || !strings.Contains(installErr.Error(), "DGM-DEPLOY-001") || result.OperationID == "" {
t.Fatalf("result=%#v err=%v", result, installErr)
}
diagnostics, err := repository.ListDiagnostics(ctx, id, 10)
if err != nil || len(diagnostics) != 1 {
t.Fatalf("diagnostics=%#v err=%v", diagnostics, err)
}
diagnostic := diagnostics[0]
if diagnostic.OperationID != result.OperationID || diagnostic.ErrorCode != "DGM-DEPLOY-001" || diagnostic.Details == "" {
t.Fatalf("diagnostic=%#v", diagnostic)
}
for _, required := range []string{`"operation_id":"` + result.OperationID + `"`, `"phase":"preparation"`, `"agent_code":"invalid_plan"`, `"template_version":"1.1.1"`, `template snapshot not registered`} {
if !strings.Contains(diagnostic.Details, required) {
t.Fatalf("diagnostic details missing %q: %s", required, diagnostic.Details)
}
}
if strings.Contains(diagnostic.Details, "password=") || strings.Contains(diagnostic.Details, "secret=") {
t.Fatalf("diagnostic leaked a secret: %s", diagnostic.Details)
}
}
+3
View File
@@ -160,6 +160,9 @@ func BuildPreview(snapshot catalog.Snapshot, request PreviewRequest) (Preview, e
}
if request.DockerUser.Mode == "" {
request.DockerUser.Mode = DockerUserDoGaMa
if snapshot.Template.Container.UserMode == DockerUserImage {
request.DockerUser.Mode = DockerUserImage
}
}
if err := ValidateDockerUser(request.DockerUser); err != nil {
return Preview{}, err
+27
View File
@@ -43,6 +43,33 @@ func TestBuildPreviewIsDeterministicAndRedactsSecrets(t *testing.T) {
}
}
func TestBuildPreviewUsesTemplateImageUserUnlessAdministratorSelectsOne(t *testing.T) {
snapshots, err := catalog.LoadFS(catalogdata.Files, ".")
if err != nil {
t.Fatal(err)
}
request := instance.PreviewRequest{
DisplayName: "Image user", Slug: "image-user", HostPorts: map[string]int{"game": 38211},
MountPaths: map[string]string{"saved": filepath.Join(t.TempDir(), "saved")}, DataOrigin: "new", BackupRetention: 7,
}
preview, err := instance.BuildPreview(snapshots[0], request)
if err != nil {
t.Fatal(err)
}
if preview.DockerUser.Mode != instance.DockerUserImage || preview.DockerUserValue != "" {
t.Fatalf("image-mode preview = %#v", preview.DockerUser)
}
plan, err := preview.DeploymentPlan("abcdefghijklmnopqrstuvwx")
if err != nil || plan.User != "" {
t.Fatalf("image-mode plan user=%q error=%v", plan.User, err)
}
request.DockerUser.Mode = instance.DockerUserDoGaMa
explicit, err := instance.BuildPreview(snapshots[0], request)
if err != nil || explicit.DockerUser.Mode != instance.DockerUserDoGaMa || explicit.DockerUserValue == "" {
t.Fatalf("explicit user preview=%#v error=%v", explicit.DockerUser, err)
}
}
func TestBuildPreviewRejectsPrivatePortPublicationAndLowResources(t *testing.T) {
snapshots, err := catalog.LoadFS(catalogdata.Files, ".")
if err != nil {
+8 -2
View File
@@ -275,6 +275,8 @@ func (s *Service) recipients(ctx context.Context, event Event) []string {
}
func categoryFor(typ string) string {
switch {
case strings.HasSuffix(typ, ".failed"):
return "server_error"
case strings.HasPrefix(typ, "start."):
return "server_start"
case strings.HasPrefix(typ, "stop."):
@@ -285,8 +287,6 @@ func categoryFor(typ string) string {
return "restore"
case strings.HasPrefix(typ, "update."):
return "update"
case strings.HasSuffix(typ, ".failed"):
return "server_error"
case strings.HasPrefix(typ, "installation_request."):
return "administration"
}
@@ -584,6 +584,12 @@ func renderText(e Event) string {
if e.Action != "" {
parts = append(parts, "Action: "+e.Action)
}
if e.OperationID != "" {
parts = append(parts, "Operation: "+e.OperationID)
}
if !e.Timestamp.IsZero() {
parts = append(parts, "Date: "+e.Timestamp.UTC().Format(time.RFC3339))
}
if e.Message != "" {
parts = append(parts, e.Message)
}
+84
View File
@@ -1,11 +1,14 @@
package notification_test
import (
"bufio"
"bytes"
"context"
"net"
"path/filepath"
"strings"
"testing"
"time"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/notification"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/persistence/sqlite"
@@ -39,6 +42,87 @@ func TestChannelSecretsAreEncryptedAndWriteOnly(t *testing.T) {
}
}
func TestSMTPFailureNotificationIncludesLifecycleContext(t *testing.T) {
ctx := context.Background()
db, err := sqlite.Open(ctx, filepath.Join(t.TempDir(), "dogama.db"))
if err != nil {
t.Fatal(err)
}
defer db.Close()
if _, err := db.ExecContext(ctx, `INSERT INTO users(id,username,email,password_hash,global_role,created_at) VALUES('admin','admin','admin@example.test','x','admin','2026-01-01T00:00:00Z')`); err != nil {
t.Fatal(err)
}
s, _ := notification.New(db, bytes.Repeat([]byte{5}, 32))
if err := s.SetPreferences(ctx, "admin", map[string]bool{"server_error": true}); err != nil {
t.Fatal(err)
}
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
defer listener.Close()
mail := make(chan string, 1)
go func() {
c, _ := listener.Accept()
if c == nil {
return
}
defer c.Close()
_, _ = c.Write([]byte("220 test\r\n"))
scanner := bufio.NewScanner(c)
var body strings.Builder
data := false
for scanner.Scan() {
line := scanner.Text()
if data {
if line == "." {
mail <- body.String()
_, _ = c.Write([]byte("250 queued\r\n"))
data = false
} else {
body.WriteString(line + "\n")
}
continue
}
switch {
case strings.HasPrefix(line, "EHLO"), strings.HasPrefix(line, "HELO"), strings.HasPrefix(line, "MAIL FROM"), strings.HasPrefix(line, "RCPT TO"):
_, _ = c.Write([]byte("250 ok\r\n"))
case line == "DATA":
data = true
_, _ = c.Write([]byte("354 data\r\n"))
case line == "QUIT":
_, _ = c.Write([]byte("221 bye\r\n"))
return
}
}
}()
host, port, _ := net.SplitHostPort(listener.Addr().String())
if _, err := s.Upsert(ctx, "", notification.Input{Name: "smtp", Type: "email", Enabled: true, Events: []string{"start.failed"}, Config: map[string]string{"host": host, "port": port, "from": "dogama@example.test", "tls_mode": "none"}}); err != nil {
t.Fatal(err)
}
e := notification.Event{Type: "start.failed", Title: "Instance failed", Message: "Code: DGM-START-001", Game: "Palworld", InstanceName: "Broken", Action: "start", OperationID: "operation-42", Timestamp: time.Now().UTC()}
if err := s.Queue(ctx, e); err != nil {
t.Fatal(err)
}
if err := s.RunDue(ctx); err != nil {
t.Fatal(err)
}
var deliveryStatus, deliveryError string
if err := db.QueryRowContext(ctx, `SELECT status,last_error_code FROM notification_deliveries`).Scan(&deliveryStatus, &deliveryError); err != nil {
t.Fatal(err)
}
select {
case body := <-mail:
for _, want := range []string{"Palworld", "Broken", "operation-42", "DGM-START-001", "Date:"} {
if !strings.Contains(body, want) {
t.Fatalf("mail missing %q: %s", want, body)
}
}
case <-time.After(time.Second):
t.Fatalf("mail not sent; delivery=%s error=%s", deliveryStatus, deliveryError)
}
}
func TestDeliveryBlocksPrivateWebhookAndRetriesWithRedactedError(t *testing.T) {
ctx := context.Background()
db, err := sqlite.Open(ctx, filepath.Join(t.TempDir(), "dogama.db"))
+99
View File
@@ -282,6 +282,13 @@ func (r *Repository) BeginOperation(ctx context.Context, operationID, instanceID
if _, err := tx.ExecContext(ctx, `INSERT INTO instance_operations(id, instance_id, kind, state, phase, created_at, updated_at) VALUES (?, ?, ?, 'running', 'dispatch', ?, ?)`, operationID, instanceID, kind, now, now); err != nil {
return instance.StoredInstance{}, fmt.Errorf("insert instance operation: %w", err)
}
if kind == "install" {
for position, step := range []string{"validation", "preparation", "installation", "import", "configuration", "starting", "verification"} {
if _, err := tx.ExecContext(ctx, `INSERT INTO operation_steps(operation_id,step_id,position,status) VALUES(?,?,?,?)`, operationID, step, position, "pending"); err != nil {
return instance.StoredInstance{}, fmt.Errorf("create operation steps: %w", err)
}
}
}
if _, err := tx.ExecContext(ctx, `UPDATE instances SET lifecycle_state=?, last_error_code=NULL, updated_at=? WHERE id=?`, lifecycleState, now, instanceID); err != nil {
return instance.StoredInstance{}, fmt.Errorf("mark instance operation: %w", err)
}
@@ -292,6 +299,43 @@ func (r *Repository) BeginOperation(ctx context.Context, operationID, instanceID
return current, nil
}
func (r *Repository) SetOperationStep(ctx context.Context, operationID, step, status string) error {
_, err := r.db.ExecContext(ctx, `UPDATE operation_steps SET status=? WHERE operation_id=? AND step_id=?`, status, operationID, step)
return err
}
func (r *Repository) GetOperationProgress(ctx context.Context, operationID string) (instance.OperationProgress, error) {
var v instance.OperationProgress
var finished sql.NullString
err := r.db.QueryRowContext(ctx, `SELECT o.id,o.kind,i.id,json_extract(i.preview_json,'$.game.name'),o.state,o.phase,COALESCE(o.error_code,''),o.created_at,o.completed_at FROM instance_operations o JOIN instances i ON i.id=o.instance_id WHERE o.id=?`, operationID).Scan(&v.OperationID, &v.OperationType, &v.InstanceID, &v.Game, &v.GlobalStatus, &v.CurrentStep, &v.ErrorCode, &v.StartedAt, &finished)
if err != nil {
if errors.Is(err, sql.ErrNoRows) {
return v, instance.ErrInstanceNotFound
}
return v, err
}
v.FinishedAt = finished.String
rows, err := r.db.QueryContext(ctx, `SELECT step_id,status FROM operation_steps WHERE operation_id=? ORDER BY position`, operationID)
if err != nil {
return v, err
}
defer rows.Close()
for rows.Next() {
var step instance.OperationStep
if err := rows.Scan(&step.ID, &step.Status); err != nil {
return v, err
}
v.Steps = append(v.Steps, step)
if step.Status == "running" {
v.CurrentStep = step.ID
}
}
if v.GlobalStatus == "succeeded" {
v.GlobalStatus = "success"
}
return v, rows.Err()
}
func (r *Repository) FinishOperation(ctx context.Context, operationID, lifecycleState, observedState, containerID, planDigest string, desiredRunning bool, errorCode string) error {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
@@ -350,6 +394,61 @@ func (r *Repository) FailOperation(ctx context.Context, operationID, lifecycleSt
return nil
}
func (r *Repository) RecordDiagnostic(ctx context.Context, value instance.Diagnostic) error {
_, err := r.db.ExecContext(ctx, `INSERT INTO operation_diagnostics(operation_id, instance_id, step, error_code, details, created_at) VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT(operation_id) DO UPDATE SET step=excluded.step, error_code=excluded.error_code, details=excluded.details, created_at=excluded.created_at`, value.OperationID, value.InstanceID, value.Step, value.ErrorCode, value.Details, r.now().UTC().Format(time.RFC3339Nano))
if err != nil {
return fmt.Errorf("record operation diagnostic: %w", err)
}
return nil
}
func (r *Repository) ListDiagnostics(ctx context.Context, instanceID string, limit int) ([]instance.Diagnostic, error) {
if limit < 1 || limit > 50 {
limit = 20
}
rows, err := r.db.QueryContext(ctx, `SELECT operation_id, instance_id, step, error_code, details, created_at FROM operation_diagnostics WHERE instance_id=? ORDER BY created_at DESC LIMIT ?`, instanceID, limit)
if err != nil {
return nil, fmt.Errorf("list operation diagnostics: %w", err)
}
defer rows.Close()
var result []instance.Diagnostic
for rows.Next() {
var value instance.Diagnostic
if err := rows.Scan(&value.OperationID, &value.InstanceID, &value.Step, &value.ErrorCode, &value.Details, &value.CreatedAt); err != nil {
return nil, fmt.Errorf("scan operation diagnostic: %w", err)
}
result = append(result, value)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate operation diagnostics: %w", err)
}
return result, nil
}
func (r *Repository) ListOperationHistory(ctx context.Context, instanceID string, limit int) ([]instance.OperationHistory, error) {
if limit < 1 || limit > 50 {
limit = 20
}
rows, err := r.db.QueryContext(ctx, `SELECT o.id,o.kind,o.state,o.phase,COALESCE(o.error_code,''),o.created_at,COALESCE(o.completed_at,''),COALESCE(d.step,''),COALESCE(d.details,'') FROM instance_operations o LEFT JOIN operation_diagnostics d ON d.operation_id=o.id WHERE o.instance_id=? ORDER BY o.created_at DESC LIMIT ?`, instanceID, limit)
if err != nil {
return nil, fmt.Errorf("list operation history: %w", err)
}
defer rows.Close()
var out []instance.OperationHistory
for rows.Next() {
var v instance.OperationHistory
var step, details string
if err := rows.Scan(&v.OperationID, &v.Kind, &v.State, &v.Phase, &v.ErrorCode, &v.CreatedAt, &v.CompletedAt, &step, &details); err != nil {
return nil, fmt.Errorf("scan operation history: %w", err)
}
if details != "" {
v.Diagnostic = &instance.Diagnostic{OperationID: v.OperationID, InstanceID: instanceID, Step: step, ErrorCode: v.ErrorCode, Details: details, CreatedAt: v.CreatedAt}
}
out = append(out, v)
}
return out, rows.Err()
}
func (r *Repository) UpdateObservation(ctx context.Context, instanceID, lifecycleState, observedState, containerID string, desiredRunning bool, errorCode string) error {
desired := 0
if desiredRunning {
@@ -63,12 +63,23 @@ func TestCatalogSyncIsImmutableAndDraftPinsSnapshot(t *testing.T) {
if _, err := repository.BeginOperation(ctx, "operation-one", "opaque-instance-id", "install", "installing"); err != nil {
t.Fatal(err)
}
if err := repository.SetOperationStep(ctx, "operation-one", "installation", "running"); err != nil {
t.Fatal(err)
}
progress, err := repository.GetOperationProgress(ctx, "operation-one")
if err != nil || progress.GlobalStatus != "running" || progress.CurrentStep != "installation" || len(progress.Steps) != 7 {
t.Fatalf("progress=%#v err=%v", progress, err)
}
if _, err := repository.BeginOperation(ctx, "operation-two", "opaque-instance-id", "start", "starting"); !errors.Is(err, instance.ErrOperationConflict) {
t.Fatalf("parallel operation error = %v", err)
}
if err := repository.FailOperation(ctx, "operation-one", "error", "test_failure"); err != nil {
t.Fatal(err)
}
progress, err = repository.GetOperationProgress(ctx, "operation-one")
if err != nil || progress.GlobalStatus != "failed" || progress.ErrorCode != "test_failure" || progress.FinishedAt == "" {
t.Fatalf("failed progress=%#v err=%v", progress, err)
}
if _, err := repository.BeginOperation(ctx, "operation-three", "opaque-instance-id", "start", "starting"); err != nil {
t.Fatal(err)
}
@@ -121,3 +132,42 @@ func TestReplaceMakesLocalCatalogDiskStateAuthoritative(t *testing.T) {
t.Fatalf("removed list = %#v, %v", list, err)
}
}
func TestSyncActivatesCurrentBundledSnapshotWithoutDiscardingHistory(t *testing.T) {
ctx := context.Background()
db, err := sqlite.Open(ctx, filepath.Join(t.TempDir(), "dogama.db"))
if err != nil {
t.Fatal(err)
}
defer db.Close()
repository := sqlite.NewRepository(db)
current, err := catalog.LoadFS(catalogdata.Files, ".")
if err != nil || len(current) != 1 {
t.Fatalf("current catalog = %#v, %v", current, err)
}
body, err := catalogdata.Files.ReadFile("palworld/template.yaml")
if err != nil {
t.Fatal(err)
}
oldBody := strings.Replace(string(body), "version: 1.1.1", "version: 1.1.0", 1)
old, err := catalog.Validate([]byte(oldBody), "palworld", catalogdata.Files)
if err != nil {
t.Fatal(err)
}
if err := repository.Replace(ctx, []catalog.Snapshot{old}); err != nil {
t.Fatal(err)
}
if err := repository.Sync(ctx, current); err != nil {
t.Fatal(err)
}
summaries, err := repository.List(ctx)
if err != nil || len(summaries) != 1 {
t.Fatalf("summaries = %#v, error = %v", summaries, err)
}
if summaries[0].Version != current[0].Template.Version || summaries[0].Digest != current[0].Digest {
t.Fatalf("active summary = %#v, want %s@%s", summaries[0], current[0].Template.ID, current[0].Template.Version)
}
if preserved, err := repository.Get(ctx, old.Template.ID, old.Template.Version); err != nil || preserved.Digest != old.Digest {
t.Fatalf("historical snapshot = %#v, error = %v", preserved, err)
}
}
+19 -1
View File
@@ -57,10 +57,28 @@ func initialize(ctx context.Context, db *sql.DB) error {
return fmt.Errorf("inspect sqlite schema: %w", err)
}
if existing == 1 {
return nil
return ensureDiagnosticSchema(ctx, db)
}
if _, err := db.ExecContext(ctx, schema); err != nil {
return fmt.Errorf("initialize sqlite schema: %w", err)
}
return ensureDiagnosticSchema(ctx, db)
}
// This additive migration is safe for existing V1 databases and keeps the
// original embedded schema immutable for fresh installs.
func ensureDiagnosticSchema(ctx context.Context, db *sql.DB) error {
_, err := db.ExecContext(ctx, `CREATE TABLE IF NOT EXISTS operation_diagnostics (
operation_id TEXT PRIMARY KEY REFERENCES instance_operations(id) ON DELETE CASCADE,
instance_id TEXT NOT NULL REFERENCES instances(id) ON DELETE CASCADE,
step TEXT NOT NULL,
error_code TEXT NOT NULL,
details TEXT NOT NULL,
created_at TEXT NOT NULL
); CREATE INDEX IF NOT EXISTS operation_diagnostics_instance_idx ON operation_diagnostics(instance_id, created_at DESC);
CREATE TABLE IF NOT EXISTS operation_steps (operation_id TEXT NOT NULL REFERENCES instance_operations(id) ON DELETE CASCADE, step_id TEXT NOT NULL, position INTEGER NOT NULL, status TEXT NOT NULL, PRIMARY KEY(operation_id, step_id));`)
if err != nil {
return fmt.Errorf("migrate diagnostic schema: %w", err)
}
return nil
}
+126 -15
View File
@@ -301,6 +301,8 @@ func newHandlerServicesWithCatalog(authService *auth.Service, repository reposit
mux.HandleFunc("PUT /api/v1/instances/{id}/memberships/{userID}/permissions/{permission}", s.permissionOverrideSet)
mux.HandleFunc("DELETE /api/v1/instances/{id}/memberships/{userID}/permissions/{permission}", s.permissionOverrideDelete)
if lifecycle != nil {
mux.HandleFunc("GET /api/v1/operations/{operationID}", s.operationProgress)
mux.HandleFunc("GET /api/v1/instances/{id}/diagnostics", s.instanceDiagnostics)
mux.HandleFunc("GET /api/v1/instances/{id}", s.instanceInspect)
mux.HandleFunc("GET /api/v1/instances/{id}/stats", s.instanceStats)
mux.HandleFunc("POST /api/v1/instances/{id}/install", s.instanceInstall)
@@ -911,6 +913,61 @@ func (s *server) instanceStats(w http.ResponseWriter, r *http.Request) {
s.apiJSON(w, http.StatusOK, stats)
}
// instanceDiagnostics is deliberately administrator-only. The operation
// history contains Docker details that must never be made available merely to
// an instance member.
func (s *server) instanceDiagnostics(w http.ResponseWriter, r *http.Request) {
user, ok := s.requireAPIUser(w, r, false)
if !ok {
return
}
if user.Role != "admin" {
s.apiProblem(w, http.StatusForbidden, "diagnostics_forbidden", "Diagnostics are restricted to administrators.")
return
}
repository, ok := s.repository.(instance.DiagnosticRepository)
if !ok {
s.apiProblem(w, http.StatusServiceUnavailable, "diagnostics_unavailable", "Diagnostics are unavailable.")
return
}
if _, err := s.repository.GetInstance(r.Context(), r.PathValue("id")); err != nil {
s.lifecycleProblem(w, err)
return
}
history, err := repository.ListOperationHistory(r.Context(), r.PathValue("id"), 20)
if err != nil {
s.apiProblem(w, http.StatusInternalServerError, "diagnostics_unavailable", "Diagnostics are unavailable.")
return
}
s.apiJSON(w, http.StatusOK, map[string]any{"operations": history})
}
func (s *server) operationProgress(w http.ResponseWriter, r *http.Request) {
user, ok := s.requireAPIUser(w, r, false)
if !ok {
return
}
repository, ok := s.repository.(instance.ProgressRepository)
if !ok {
s.apiProblem(w, http.StatusServiceUnavailable, "operations_unavailable", "Operations are unavailable.")
return
}
progress, err := repository.GetOperationProgress(r.Context(), r.PathValue("operationID"))
if errors.Is(err, instance.ErrInstanceNotFound) {
s.apiProblem(w, http.StatusNotFound, "operation_not_found", "The operation does not exist.")
return
}
if err != nil {
s.apiProblem(w, http.StatusInternalServerError, "operations_unavailable", "Operations are unavailable.")
return
}
if s.permissions.Require(r.Context(), user, progress.InstanceID, authorization.PermissionInstanceView) != nil {
s.apiProblem(w, http.StatusForbidden, "permission_denied", "Permission denied.")
return
}
s.apiJSON(w, http.StatusOK, progress)
}
func (s *server) instanceInstall(w http.ResponseWriter, r *http.Request) {
actor, ok := s.requireAPIUser(w, r, true)
if !ok {
@@ -995,6 +1052,11 @@ func (s *server) runLifecycleAction(w http.ResponseWriter, r *http.Request, acti
}
result, err := action(r.Context(), r.PathValue("id"))
if err != nil {
code := "DGM-DOCKER-001"
if strings.HasPrefix(err.Error(), "DGM-") {
code = strings.SplitN(err.Error(), ":", 2)[0]
}
s.queueNotification(r, notification.Event{Type: "start.failed", Title: "Instance lifecycle action failed", Message: "The instance operation failed. Code: " + code, Action: "start", OperationID: result.OperationID, Severity: "error"})
s.lifecycleProblem(w, err)
return
}
@@ -1025,6 +1087,9 @@ func (s *server) instanceDeleteContainer(w http.ResponseWriter, r *http.Request)
func (s *server) lifecycleProblem(w http.ResponseWriter, err error) {
status, code := http.StatusBadGateway, "lifecycle_failed"
if strings.HasPrefix(err.Error(), "DGM-") {
code = strings.SplitN(err.Error(), ":", 2)[0]
}
switch {
case errors.Is(err, instance.ErrInstanceNotFound):
status, code = http.StatusNotFound, "instance_not_found"
@@ -1034,6 +1099,7 @@ func (s *server) lifecycleProblem(w http.ResponseWriter, err error) {
status, code = http.StatusConflict, "invalid_instance_state"
}
s.apiProblem(w, status, code, "The instance operation could not be completed.")
s.logger.Error("instance lifecycle failed", "event", "instance.lifecycle.failed", "error_code", code, "error", err)
}
func (s *server) buildAPIPreview(w http.ResponseWriter, r *http.Request) (previewAPIRequest, instance.Preview, bool) {
@@ -1812,12 +1878,46 @@ func (s *server) deploymentSubmit(w http.ResponseWriter, r *http.Request) {
return
}
}
if _, err := s.lifecycle.Install(r.Context(), id); err != nil {
data.Deployment.Error = "Installation failed. The instance is retained in an error state."
operation, err := s.lifecycle.BeginInstall(r.Context(), id)
if err != nil {
data.Deployment.Error = "The deployment could not be queued."
s.render(w, 502, "deployment.html", data)
return
}
// Everything needed by the worker is now persisted (draft, encrypted secrets
// and validated import); it never retains the HTTP request or multipart files.
go s.runDeployment(operation.OperationID, id, preview, name)
if s.audit != nil {
actor, _ := s.currentUser(r)
_ = s.audit.Record(r.Context(), audit.Event{ActorID: actor.ID, InstanceID: id, Action: "instance.deploy", Outcome: "allowed", Summary: map[string]string{"target_name": name}})
}
if r.Header.Get("Accept") == "application/json" {
s.apiJSON(w, http.StatusAccepted, map[string]string{"operation_id": operation.OperationID, "instance_id": id})
return
}
http.Redirect(w, r, "/instances/"+id, http.StatusSeeOther)
}
func (s *server) runDeployment(operationID, id string, preview instance.Preview, name string) {
ctx := context.Background()
fail := func(step string, err error) {
_, _ = s.lifecycle.FailInstall(ctx, operationID, id, step, err)
s.logger.Error("deployment worker failed", "operation_id", operationID, "instance_id", id, "step", step, "error", err)
s.queueDeploymentFailure(operationID, preview.Game.Name, name)
}
defer func() {
if recovered := recover(); recovered != nil {
fail("internal", fmt.Errorf("deployment worker panic"))
}
}()
s.lifecycle.MarkStep(ctx, operationID, "validation", "success")
s.lifecycle.MarkStep(ctx, operationID, "preparation", "success")
if _, err := s.lifecycle.InstallOperation(ctx, id, operationID); err != nil {
s.queueDeploymentFailure(operationID, preview.Game.Name, name)
return
}
if preview.DataOrigin == "import" {
s.lifecycle.MarkStep(ctx, operationID, "import", "running")
mountPath := ""
for _, mount := range preview.Mounts {
if mount.ID == preview.Import.DestinationMount {
@@ -1825,27 +1925,38 @@ func (s *server) deploymentSubmit(w http.ResponseWriter, r *http.Request) {
break
}
}
if mountPath == "" || s.imports == nil || s.imports.ApplyToInstance(r.Context(), preview.Import.ID, id, preview.Template.ID, preview.Template.Version, mountPath, preview.Import.DestinationRelativePath) != nil {
data.Deployment.Error = "Installation succeeded but the backup could not be restored."
s.render(w, 422, "deployment.html", data)
if mountPath == "" || s.imports == nil {
fail("import", errors.New("import destination unavailable"))
return
}
if err := s.imports.ApplyToInstance(ctx, preview.Import.ID, id, preview.Template.ID, preview.Template.Version, mountPath, preview.Import.DestinationRelativePath); err != nil {
fail("import", err)
return
}
s.lifecycle.MarkStep(ctx, operationID, "import", "success")
} else {
s.lifecycle.MarkStep(ctx, operationID, "import", "success")
}
if err := s.applyDeploymentConfiguration(r.Context(), id, preview); err != nil {
data.Deployment.Error = "Installation succeeded but configuration could not be applied."
s.render(w, 502, "deployment.html", data)
s.lifecycle.MarkStep(ctx, operationID, "configuration", "running")
if err := s.applyDeploymentConfiguration(ctx, id, preview); err != nil {
fail("configuration", err)
return
}
if _, err := s.lifecycle.Start(r.Context(), id); err != nil {
data.Deployment.Error = "Installation succeeded but the server could not start."
s.render(w, 502, "deployment.html", data)
s.lifecycle.MarkStep(ctx, operationID, "configuration", "success")
result, err := s.lifecycle.StartInstall(ctx, id, operationID)
if err != nil {
s.queueDeploymentFailure(operationID, preview.Game.Name, name)
return
}
if s.audit != nil {
actor, _ := s.currentUser(r)
_ = s.audit.Record(r.Context(), audit.Event{ActorID: actor.ID, InstanceID: id, Action: "instance.deploy", Outcome: "allowed", Summary: map[string]string{"target_name": name}})
if err := s.lifecycle.CompleteInstall(ctx, result); err != nil {
fail("verification", err)
}
}
func (s *server) queueDeploymentFailure(operationID, game, name string) {
if s.notifications != nil {
_ = s.notifications.Queue(context.Background(), notification.Event{Type: "start.failed", Title: "Instance deployment failed", Message: "The instance deployment failed. Code: DGM-START-001", Game: game, InstanceName: name, Action: "start", OperationID: operationID, Severity: "error"})
}
http.Redirect(w, r, "/", http.StatusSeeOther)
}
// applyDeploymentConfiguration intentionally runs after restore: values chosen
+164 -3
View File
@@ -7,6 +7,7 @@ import (
"encoding/json"
"io"
"log/slog"
"mime/multipart"
"net/http"
"net/http/httptest"
"net/url"
@@ -30,6 +31,37 @@ import (
type webLifecycleAgent struct{}
type blockingLifecycleAgent struct {
entered chan struct{}
release chan struct{}
panicOnCreate bool
}
func (a *blockingLifecycleAgent) CreateInstance(_ context.Context, plan agentwire.DeploymentPlan) (agentwire.InstanceState, error) {
close(a.entered)
<-a.release
if a.panicOnCreate {
panic("controlled worker panic")
}
return agentwire.InstanceState{InstanceID: plan.InstanceID, ContainerID: "controlled-container", PlanDigest: plan.PlanDigest, Health: "stopped"}, nil
}
func (a *blockingLifecycleAgent) InspectInstance(_ context.Context, id string) (agentwire.InstanceState, error) {
return agentwire.InstanceState{InstanceID: id}, nil
}
func (a *blockingLifecycleAgent) StartInstance(_ context.Context, id string) (agentwire.InstanceState, error) {
return agentwire.InstanceState{InstanceID: id, ContainerID: "controlled-container", Running: true, Ready: true, Health: "healthy"}, nil
}
func (a *blockingLifecycleAgent) StopInstance(context.Context, string, int) (agentwire.InstanceState, error) {
return agentwire.InstanceState{}, nil
}
func (a *blockingLifecycleAgent) RestartInstance(context.Context, string, int) (agentwire.InstanceState, error) {
return agentwire.InstanceState{}, nil
}
func (a *blockingLifecycleAgent) DeleteContainer(context.Context, string) error { return nil }
func (a *blockingLifecycleAgent) GetInstanceStats(context.Context, string) (agentwire.InstanceStats, error) {
return agentwire.InstanceStats{}, nil
}
func (webLifecycleAgent) CreateInstance(_ context.Context, plan agentwire.DeploymentPlan) (agentwire.InstanceState, error) {
return agentwire.InstanceState{InstanceID: plan.InstanceID, ContainerID: "container-1", PlanDigest: plan.PlanDigest, Health: "stopped"}, nil
}
@@ -260,7 +292,7 @@ func TestCatalogPreviewAndDraftAPIAuthorization(t *testing.T) {
t.Fatalf("catalog response = %s", catalogResponse.Body.String())
}
payload, _ := json.Marshal(map[string]any{
"template_id": "palworld-official", "template_version": "1.1.0",
"template_id": snapshots[0].Template.ID, "template_version": snapshots[0].Template.Version,
"display_name": "Family Palworld", "slug": "family-palworld",
"host_ports": map[string]int{"game": 8211},
"mount_paths": map[string]string{"saved": "/srv/game-servers/family-palworld/saved"},
@@ -722,7 +754,7 @@ func TestInstanceAuthorizationAndInstallationRequestWorkflow(t *testing.T) {
playerCookie := &http.Cookie{Name: sessionCookie, Value: playerSession.Token}
draftPayload, _ := json.Marshal(map[string]any{
"template_id": "palworld-official", "template_version": "1.1.0",
"template_id": snapshots[0].Template.ID, "template_version": snapshots[0].Template.Version,
"display_name": "Authorization Test", "slug": "authorization-test",
"host_ports": map[string]int{"game": 8211},
"mount_paths": map[string]string{"saved": "/srv/game-servers/authorization-test/saved"},
@@ -761,7 +793,10 @@ func TestInstanceAuthorizationAndInstallationRequestWorkflow(t *testing.T) {
substitution := request(t, handler, http.MethodGet, "/api/v1/instances/not-the-member-instance", []*http.Cookie{playerCookie})
assertStatus(t, substitution, http.StatusForbidden)
requestPayload := []byte(`{"template_id":"palworld-official","template_version":"1.1.0","suggested_name":"Friends","player_estimate":8,"desired_schedule":"evenings","mods_requested":true,"message":"Private group"}`)
requestPayload, err := json.Marshal(map[string]any{"template_id": snapshots[0].Template.ID, "template_version": snapshots[0].Template.Version, "suggested_name": "Friends", "player_estimate": 8, "desired_schedule": "evenings", "mods_requested": true, "message": "Private group"})
if err != nil {
t.Fatal(err)
}
installationRequest := jsonRequest(t, handler, "/api/v1/installation-requests", requestPayload, playerCookie, playerSession.CSRFToken)
assertStatus(t, installationRequest, http.StatusCreated)
var createdRequest struct {
@@ -1109,3 +1144,129 @@ func TestDeploymentFormRequiresAdminAndRendersTemplateFields(t *testing.T) {
missing := request(t, handler, http.MethodGet, "/catalog/missing/deploy", adminCookies)
assertStatus(t, missing, http.StatusNotFound)
}
func TestDeploymentHTTPAsyncProgressAndRBAC(t *testing.T) { testDeploymentHTTPAsync(t, false) }
func TestDeploymentWorkerPanicFailsOperation(t *testing.T) { testDeploymentHTTPAsync(t, true) }
func testDeploymentHTTPAsync(t *testing.T, panicWorker bool) {
ctx := context.Background()
db, err := sqlite.Open(ctx, filepath.Join(t.TempDir(), "dogama.db"))
if err != nil {
t.Fatal(err)
}
defer db.Close()
repository := sqlite.NewRepository(db)
if err := repository.SetSecretKey(bytes.Repeat([]byte{1}, 32)); err != nil {
t.Fatal(err)
}
snapshots, err := catalog.LoadFS(catalogdata.Files, ".")
if err != nil {
t.Fatal(err)
}
if err = repository.Sync(ctx, snapshots); err != nil {
t.Fatal(err)
}
authService := auth.New(db)
if err = authService.BootstrapAdmin(ctx, "admin", "correct horse battery staple"); err != nil {
t.Fatal(err)
}
admin, err := authService.Login(ctx, "admin", "correct horse battery staple", "192.0.2.1:1234")
if err != nil {
t.Fatal(err)
}
if _, err = authService.CreateUser(ctx, "viewer", "another correct battery staple", "user"); err != nil {
t.Fatal(err)
}
viewer, err := authService.Login(ctx, "viewer", "another correct battery staple", "192.0.2.2:1234")
if err != nil {
t.Fatal(err)
}
agent := &blockingLifecycleAgent{entered: make(chan struct{}), release: make(chan struct{}), panicOnCreate: panicWorker}
handler, err := NewHandlerCompleteWithCatalogAndDeployment(authService, repository, instance.NewLifecycleService(repository, agent), nil, nil, nil, nil, nil, filepath.Join(t.TempDir(), "servers"), slog.New(slog.NewTextHandler(io.Discard, nil)))
if err != nil {
t.Fatal(err)
}
var body bytes.Buffer
form := multipart.NewWriter(&body)
for key, value := range map[string]string{"csrf_token": admin.CSRFToken, "display_name": "Async test", "description": "safe", "config_server_name": "Async", "config_server_description": "safe", "config_max_players": "16", "config_admin_password": "top-secret-value", "config_rest_api_enabled": "true", "config_rest_api_port": "8212"} {
_ = form.WriteField(key, value)
}
_ = form.Close()
req := httptest.NewRequest(http.MethodPost, "/catalog/palworld-official/deploy", &body)
req.Header.Set("Content-Type", form.FormDataContentType())
req.Header.Set("Accept", "application/json")
req.AddCookie(&http.Cookie{Name: sessionCookie, Value: admin.Token})
req.AddCookie(&http.Cookie{Name: csrfCookie, Value: admin.CSRFToken})
response := httptest.NewRecorder()
handler.ServeHTTP(response, req)
assertStatus(t, response, http.StatusAccepted)
var accepted map[string]string
if err := json.Unmarshal(response.Body.Bytes(), &accepted); err != nil || accepted["operation_id"] == "" {
t.Fatalf("accepted=%s err=%v", response.Body.String(), err)
}
// The HTTP response has returned while the worker is deterministically blocked.
<-agent.entered
progressReq := httptest.NewRequest(http.MethodGet, "/api/v1/operations/"+accepted["operation_id"], nil)
progressReq.AddCookie(&http.Cookie{Name: sessionCookie, Value: admin.Token})
progress := httptest.NewRecorder()
handler.ServeHTTP(progress, progressReq)
assertStatus(t, progress, http.StatusOK)
if strings.Contains(progress.Body.String(), "top-secret-value") || !strings.Contains(progress.Body.String(), "installation") {
t.Fatalf("unexpected progress body: %s", progress.Body.String())
}
users, _ := authService.ListUsers(ctx)
var viewerID, adminID string
for _, u := range users {
if u.Username == "viewer" {
viewerID = u.ID
}
if u.Username == "admin" {
adminID = u.ID
}
}
if err := repository.SetMembership(ctx, adminID, accepted["instance_id"], viewerID, "user"); err != nil {
t.Fatal(err)
}
viewerProgress := httptest.NewRequest(http.MethodGet, "/api/v1/operations/"+accepted["operation_id"], nil)
viewerProgress.AddCookie(&http.Cookie{Name: sessionCookie, Value: viewer.Token})
viewerResponse := httptest.NewRecorder()
handler.ServeHTTP(viewerResponse, viewerProgress)
assertStatus(t, viewerResponse, http.StatusOK)
if strings.Contains(viewerResponse.Body.String(), "top-secret-value") {
t.Fatal("secret leaked to member progress")
}
diagnostic := httptest.NewRequest(http.MethodGet, "/api/v1/instances/"+accepted["instance_id"]+"/diagnostics", nil)
diagnostic.AddCookie(&http.Cookie{Name: sessionCookie, Value: viewer.Token})
diagnosticResponse := httptest.NewRecorder()
handler.ServeHTTP(diagnosticResponse, diagnostic)
assertStatus(t, diagnosticResponse, http.StatusForbidden)
unknown := httptest.NewRequest(http.MethodGet, "/api/v1/operations/missing", nil)
unknown.AddCookie(&http.Cookie{Name: sessionCookie, Value: admin.Token})
unknownResponse := httptest.NewRecorder()
handler.ServeHTTP(unknownResponse, unknown)
assertStatus(t, unknownResponse, http.StatusNotFound)
close(agent.release)
deadline := time.Now().Add(2 * time.Second)
for {
value, err := repository.GetOperationProgress(ctx, accepted["operation_id"])
if err == nil && ((!panicWorker && value.GlobalStatus == "success") || (panicWorker && value.GlobalStatus == "failed")) {
if panicWorker && (value.ErrorCode == "" || value.OperationID != accepted["operation_id"]) {
t.Fatalf("panic result is not actionable: %#v", value)
}
break
}
if time.Now().After(deadline) {
t.Fatalf("operation did not finish: %#v %v", value, err)
}
time.Sleep(10 * time.Millisecond)
}
adminDiagnostic := httptest.NewRequest(http.MethodGet, "/api/v1/instances/"+accepted["instance_id"]+"/diagnostics", nil)
adminDiagnostic.AddCookie(&http.Cookie{Name: sessionCookie, Value: admin.Token})
adminDiagnosticResponse := httptest.NewRecorder()
handler.ServeHTTP(adminDiagnosticResponse, adminDiagnostic)
assertStatus(t, adminDiagnosticResponse, http.StatusOK)
if panicWorker && (strings.Contains(adminDiagnosticResponse.Body.String(), "top-secret-value") || !strings.Contains(adminDiagnosticResponse.Body.String(), accepted["operation_id"])) {
t.Fatalf("panic diagnostic leaked secret or lost operation id: %s", adminDiagnosticResponse.Body.String())
}
}
+39 -2
View File
@@ -31,10 +31,47 @@ if (instanceSearch) {
const deploymentForm = document.querySelector("[data-deployment-form]");
if (deploymentForm) {
deploymentForm.addEventListener("submit", () => {
deploymentForm.addEventListener("submit", (event) => {
event.preventDefault();
const submit = deploymentForm.querySelector("[data-deployment-submit]");
if (submit) { submit.disabled = true; submit.textContent = "…"; }
}, { once: true });
let progress = document.querySelector("[data-deployment-progress]");
if (!progress) { progress = document.createElement("section"); progress.className = "panel deployment-progress"; progress.dataset.deploymentProgress = ""; deploymentForm.after(progress); }
progress.replaceChildren(); const heading = document.createElement("h2"); heading.textContent = "Deployment progress"; const status = document.createElement("p"); status.dataset.operationStatus = ""; const list = document.createElement("ol"); list.dataset.operationSteps = ""; const error = document.createElement("p"); error.className = "error"; error.hidden = true; const link = document.createElement("a"); link.hidden = true; progress.append(heading, status, list, error, link);
fetch(deploymentForm.action || window.location.href, { method: "POST", body: new FormData(deploymentForm), credentials: "same-origin", headers: { Accept: "application/json" } }).then((response) => response.ok ? response.json() : Promise.reject(new Error("Deployment could not be queued."))).then(({ operation_id, instance_id }) => {
progress.hidden = false; deploymentForm.hidden = true;
const poll = () => fetch(`/api/v1/operations/${encodeURIComponent(operation_id)}`, { credentials: "same-origin" }).then((r) => r.ok ? r.json() : Promise.reject(new Error("Progress unavailable."))).then((operation) => {
status.textContent = operation.global_status; list.replaceChildren(...operation.steps.map((step) => { const item = document.createElement("li"); item.textContent = `${step.status === "success" ? "✓" : step.status === "failed" ? "✗" : step.status === "running" ? "▶" : "○"} ${step.id}`; return item; }));
if (operation.global_status === "failed") { error.hidden = false; error.textContent = `Deployment failed${operation.error_code ? `. Code: ${operation.error_code}` : "."}`; return; }
if (operation.global_status === "success") { link.hidden = false; link.href = `/instances/${encodeURIComponent(instance_id)}`; link.textContent = "Open instance"; return; }
window.setTimeout(poll, 800);
}).catch((reason) => { error.hidden = false; error.textContent = reason.message; }); poll();
}).catch((reason) => { if (submit) { submit.disabled = false; submit.textContent = "Go"; } alert(reason.message); });
});
}
// The backend, not a timer, is the source for this history. A 403 simply
// means the signed-in user is not an administrator and leaves no technical
// data in the DOM.
const detailMatch = window.location.pathname.match(/^\/instances\/([^/]+)$/);
if (detailMatch) {
fetch(`/api/v1/instances/${encodeURIComponent(detailMatch[1])}/diagnostics`, { credentials: "same-origin" })
.then((response) => response.ok ? response.json() : null)
.then((data) => {
if (!data || !data.operations || !data.operations.length) return;
const section = document.createElement("section");
section.className = "panel";
section.innerHTML = "<h2>Operation diagnostics</h2><p class=muted>Administrator-only technical history.</p>";
data.operations.forEach((operation) => {
const item = document.createElement("details");
const summary = document.createElement("summary");
summary.textContent = `${operation.created_at} · ${operation.kind} · ${operation.state}${operation.error_code ? ` · ${operation.error_code}` : ""} · ${operation.operation_id}`;
item.append(summary);
if (operation.diagnostic) { const detail = document.createElement("pre"); detail.textContent = operation.diagnostic.details; item.append(detail); }
section.append(item);
});
document.querySelector(".instance-detail-lower")?.append(section);
}).catch(() => {});
}
const catalogSearch = document.querySelector("#catalog-search");
+1
View File
@@ -59,6 +59,7 @@
"image": { "type": "string", "pattern": "^[a-zA-Z0-9._/-]+$", "maxLength": 300 },
"tag": { "type": "string", "pattern": "^[a-zA-Z0-9._-]+$", "maxLength": 128 },
"entrypoint": { "type": "array", "items": { "type": "string", "maxLength": 500 }, "maxItems": 8 },
"user_mode": { "type": "string", "enum": ["dogama", "image"] },
"arguments": { "type": "array", "items": { "type": "string", "maxLength": 500 }, "maxItems": 64 },
"environment": { "type": "object", "maxProperties": 64, "additionalProperties": { "type": "string", "maxLength": 4096 }, "propertyNames": { "pattern": "^[A-Za-z_][A-Za-z0-9_]*$" } },
"assets": {
+1
View File
@@ -40,6 +40,7 @@ func TestV1BootstrapAuthenticationAndHTTPBoundary(t *testing.T) {
"DOGAMA_LISTEN_ADDRESS="+address,
"DOGAMA_DATABASE_PATH="+filepath.Join(t.TempDir(), "dogama.db"),
"DOGAMA_MASTER_KEY_FILE="+filepath.Join(t.TempDir(), "master_key"),
"DOGAMA_TEMPLATES_ROOT="+filepath.Join(t.TempDir(), "templates"),
"DOGAMA_IMPORTS_ROOT="+filepath.Join(t.TempDir(), "imports"),
"DOGAMA_SERVERS_ROOT="+filepath.Join(t.TempDir(), "servers"),
)