feat(deploy): run deployments asynchronously with progress
CI / validate (pull_request) Failing after 6m57s

This commit is contained in:
2026-08-15 19:52:39 +02:00
parent 06f99750ea
commit 265136fc30
6 changed files with 272 additions and 18 deletions
+119
View File
@@ -62,6 +62,30 @@ type OperationHistory struct {
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)
@@ -135,6 +159,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 {
@@ -144,6 +171,95 @@ 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 {
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)
}
lifecycle, observed := stateToLifecycle(state)
if progress, ok := s.repository.(ProgressRepository); ok {
_ = progress.SetOperationStep(ctx, operationID, "starting", "success")
_ = progress.SetOperationStep(ctx, operationID, "verification", "success")
}
return OperationResult{OperationID: operationID, InstanceID: instanceID, State: lifecycle, Observed: observed, ContainerID: state.ContainerID, AgentState: 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)
@@ -366,6 +482,9 @@ func (s *LifecycleService) inspectAndPersist(ctx context.Context, current Stored
func (s *LifecycleService) fail(ctx context.Context, operationID, instanceID, code string, cause error) (OperationResult, error) {
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(cause)})
}
+44
View File
@@ -282,6 +282,13 @@ func (r *Repository) BeginOperation(ctx context.Context, operationID, instanceID
if _, err := tx.ExecContext(ctx, `INSERT INTO instance_operations(id, instance_id, kind, state, phase, created_at, updated_at) VALUES (?, ?, ?, 'running', 'dispatch', ?, ?)`, operationID, instanceID, kind, now, now); err != nil {
return instance.StoredInstance{}, fmt.Errorf("insert instance operation: %w", err)
}
if kind == "install" {
for position, step := range []string{"validation", "preparation", "installation", "import", "configuration", "starting", "verification"} {
if _, err := tx.ExecContext(ctx, `INSERT INTO operation_steps(operation_id,step_id,position,status) VALUES(?,?,?,?)`, operationID, step, position, "pending"); err != nil {
return instance.StoredInstance{}, fmt.Errorf("create operation steps: %w", err)
}
}
}
if _, err := tx.ExecContext(ctx, `UPDATE instances SET lifecycle_state=?, last_error_code=NULL, updated_at=? WHERE id=?`, lifecycleState, now, instanceID); err != nil {
return instance.StoredInstance{}, fmt.Errorf("mark instance operation: %w", err)
}
@@ -292,6 +299,43 @@ func (r *Repository) BeginOperation(ctx context.Context, operationID, instanceID
return current, nil
}
func (r *Repository) SetOperationStep(ctx context.Context, operationID, step, status string) error {
_, err := r.db.ExecContext(ctx, `UPDATE operation_steps SET status=? WHERE operation_id=? AND step_id=?`, status, operationID, step)
return err
}
func (r *Repository) GetOperationProgress(ctx context.Context, operationID string) (instance.OperationProgress, error) {
var v instance.OperationProgress
var finished sql.NullString
err := r.db.QueryRowContext(ctx, `SELECT o.id,o.kind,i.id,json_extract(i.preview_json,'$.game.name'),o.state,o.phase,COALESCE(o.error_code,''),o.created_at,o.completed_at FROM instance_operations o JOIN instances i ON i.id=o.instance_id WHERE o.id=?`, operationID).Scan(&v.OperationID, &v.OperationType, &v.InstanceID, &v.Game, &v.GlobalStatus, &v.CurrentStep, &v.ErrorCode, &v.StartedAt, &finished)
if err != nil {
if errors.Is(err, sql.ErrNoRows) {
return v, instance.ErrInstanceNotFound
}
return v, err
}
v.FinishedAt = finished.String
rows, err := r.db.QueryContext(ctx, `SELECT step_id,status FROM operation_steps WHERE operation_id=? ORDER BY position`, operationID)
if err != nil {
return v, err
}
defer rows.Close()
for rows.Next() {
var step instance.OperationStep
if err := rows.Scan(&step.ID, &step.Status); err != nil {
return v, err
}
v.Steps = append(v.Steps, step)
if step.Status == "running" {
v.CurrentStep = step.ID
}
}
if v.GlobalStatus == "succeeded" {
v.GlobalStatus = "success"
}
return v, rows.Err()
}
func (r *Repository) FinishOperation(ctx context.Context, operationID, lifecycleState, observedState, containerID, planDigest string, desiredRunning bool, errorCode string) error {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
@@ -63,12 +63,17 @@ 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)
}
+2 -1
View File
@@ -75,7 +75,8 @@ func ensureDiagnosticSchema(ctx context.Context, db *sql.DB) error {
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 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)
}
+87 -15
View File
@@ -301,6 +301,7 @@ 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)
@@ -941,6 +942,32 @@ func (s *server) instanceDiagnostics(w http.ResponseWriter, r *http.Request) {
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 {
@@ -1851,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 {
@@ -1864,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
+15 -2
View File
@@ -31,10 +31,23 @@ 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