diff --git a/.dockerignore b/.dockerignore index 0d27dc8..c219151 100644 --- a/.dockerignore +++ b/.dockerignore @@ -2,6 +2,9 @@ .cache dist data +servers +backups +.playwright-cli secrets *.db *.db-shm diff --git a/catalog/palworld/template.yaml b/catalog/palworld/template.yaml index 4d27c09..b336d01 100644 --- a/catalog/palworld/template.yaml +++ b/catalog/palworld/template.yaml @@ -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 diff --git a/cmd/dogama/main.go b/cmd/dogama/main.go index 05f5819..a038f14 100644 --- a/cmd/dogama/main.go +++ b/cmd/dogama/main.go @@ -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") diff --git a/docs/PROJECT-STATE.md b/docs/PROJECT-STATE.md index eb2dbbb..9b844e7 100644 --- a/docs/PROJECT-STATE.md +++ b/docs/PROJECT-STATE.md @@ -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 @@ -100,3 +100,4 @@ Read this compact operational baseline before starting a milestone. Open detaile - 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. diff --git a/internal/agent/agent_assets.go b/internal/agent/agent_assets.go index 7f6bcea..9a67201 100644 --- a/internal/agent/agent_assets.go +++ b/internal/agent/agent_assets.go @@ -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 { diff --git a/internal/agent/agent_assets_test.go b/internal/agent/agent_assets_test.go new file mode 100644 index 0000000..6634b8b --- /dev/null +++ b/internal/agent/agent_assets_test.go @@ -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) + } +} diff --git a/internal/agent/agent_docker.go b/internal/agent/agent_docker.go index b74d8ad..948ea61 100644 --- a/internal/agent/agent_docker.go +++ b/internal/agent/agent_docker.go @@ -10,6 +10,7 @@ import ( "net" "net/http" "net/url" + "os" "path/filepath" "sort" "strconv" @@ -165,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 { @@ -234,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 { diff --git a/internal/agent/agent_docker_test.go b/internal/agent/agent_docker_test.go index dd13742..aa86d7b 100644 --- a/internal/agent/agent_docker_test.go +++ b/internal/agent/agent_docker_test.go @@ -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) { diff --git a/internal/agent/agent_plan.go b/internal/agent/agent_plan.go index eb6eab9..ae0289c 100644 --- a/internal/agent/agent_plan.go +++ b/internal/agent/agent_plan.go @@ -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 + ":" diff --git a/internal/agent/agent_plan_test.go b/internal/agent/agent_plan_test.go new file mode 100644 index 0000000..5a3028f --- /dev/null +++ b/internal/agent/agent_plan_test.go @@ -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) + } +} diff --git a/internal/agent/agent_server.go b/internal/agent/agent_server.go index 81ee8f0..1aefb3a 100644 --- a/internal/agent/agent_server.go +++ b/internal/agent/agent_server.go @@ -141,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 { @@ -183,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} @@ -415,7 +421,7 @@ 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) { diff --git a/internal/agent/agent_server_test.go b/internal/agent/agent_server_test.go index a2c0ad5..4a39047 100644 --- a/internal/agent/agent_server_test.go +++ b/internal/agent/agent_server_test.go @@ -165,6 +165,30 @@ 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) diff --git a/internal/agentclient/client.go b/internal/agentclient/client.go index 193823d..44c98fd 100644 --- a/internal/agentclient/client.go +++ b/internal/agentclient/client.go @@ -56,8 +56,20 @@ func (e *ProblemError) Error() string { } // DiagnosticDetails lets the lifecycle service persist the agent's bounded -// Docker inspection without exposing it to unprivileged HTTP clients. -func (e *ProblemError) DiagnosticDetails() map[string]any { return e.Details } +// 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. diff --git a/internal/catalog/template.go b/internal/catalog/template.go index 7e94463..7b149e5 100644 --- a/internal/catalog/template.go +++ b/internal/catalog/template.go @@ -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"` diff --git a/internal/instance/lifecycle.go b/internal/instance/lifecycle.go index d07ce04..98729d2 100644 --- a/internal/instance/lifecycle.go +++ b/internal/instance/lifecycle.go @@ -8,6 +8,7 @@ import ( "errors" "fmt" "sync" + "time" "git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/agentwire" ) @@ -18,6 +19,8 @@ var ( ErrInvalidState = errors.New("invalid instance lifecycle state") ) +const startupVerificationDelay = 3 * time.Second + type StoredInstance struct { ID string Preview Preview @@ -244,12 +247,41 @@ func (s *LifecycleService) StartInstall(ctx context.Context, instanceID, operati lifecycle, observed := stateToLifecycle(state) if progress, ok := s.repository.(ProgressRepository); ok { _ = progress.SetOperationStep(ctx, operationID, "starting", "success") + _ = progress.SetOperationStep(ctx, operationID, "verification", "running") + } + 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: state.ContainerID, AgentState: state}, nil + 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") @@ -489,7 +521,7 @@ func (s *LifecycleService) fail(ctx context.Context, operationID, instanceID, co _ = 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(cause)}) + _ = 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 @@ -497,24 +529,34 @@ func (s *LifecycleService) fail(ctx context.Context, operationID, instanceID, co return OperationResult{OperationID: operationID, InstanceID: instanceID, State: "error", Observed: "unknown"}, fmt.Errorf("%s: %w", stableCode, cause) } -func diagnosticDetails(cause error) string { - if cause == nil { - return "" +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) { - encoded, err := json.Marshal(value.DiagnosticDetails()) - if err == nil { - return string(encoded) + for key, detail := range value.DiagnosticDetails() { + details[key] = detail } + } else if cause != nil { + details["operation_error"] = cause.Error() } - return 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": + case "agent_start_failed", "agent_restart_failed", "verification_failed": return "DGM-START-001" case "agent_create_failed": return "DGM-DEPLOY-001" diff --git a/internal/instance/lifecycle_test.go b/internal/instance/lifecycle_test.go index 93a5559..087a891 100644 --- a/internal/instance/lifecycle_test.go +++ b/internal/instance/lifecycle_test.go @@ -23,10 +23,27 @@ type lifecycleAgent struct { 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() @@ -216,3 +233,48 @@ func TestManualStartPersistsAgentDiagnosticAndStableCode(t *testing.T) { 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) + } +} diff --git a/internal/instance/preview.go b/internal/instance/preview.go index 2437302..0b583ad 100644 --- a/internal/instance/preview.go +++ b/internal/instance/preview.go @@ -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 diff --git a/internal/instance/preview_test.go b/internal/instance/preview_test.go index 8d65917..d998114 100644 --- a/internal/instance/preview_test.go +++ b/internal/instance/preview_test.go @@ -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 { diff --git a/internal/persistence/sqlite/catalog_test.go b/internal/persistence/sqlite/catalog_test.go index cb008de..e9764b8 100644 --- a/internal/persistence/sqlite/catalog_test.go +++ b/internal/persistence/sqlite/catalog_test.go @@ -63,9 +63,13 @@ 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) } + 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 != 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) } @@ -73,7 +77,9 @@ func TestCatalogSyncIsImmutableAndDraftPinsSnapshot(t *testing.T) { 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 != 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) } @@ -126,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) + } +} diff --git a/specs/template.schema.json b/specs/template.schema.json index 1e8696e..7a4f58d 100644 --- a/specs/template.schema.json +++ b/specs/template.schema.json @@ -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": {