fix(lifecycle): retain deployment diagnostics #39
@@ -2,6 +2,9 @@
|
||||
.cache
|
||||
dist
|
||||
data
|
||||
servers
|
||||
backups
|
||||
.playwright-cli
|
||||
secrets
|
||||
*.db
|
||||
*.db-shm
|
||||
|
||||
@@ -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
@@ -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")
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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) {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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 + ":"
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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})
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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"`
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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"))
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
@@ -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
@@ -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())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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"),
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user