Files
DoGaMa-serv/internal/persistence/sqlite/catalog.go
T
codex 4ed3c8ae58
CI / validate (pull_request) Successful in 26m33s
refactor(catalog): make template artwork fully local
2026-08-26 21:58:30 +02:00

514 lines
23 KiB
Go

package sqlite
import (
"context"
"crypto/aes"
"crypto/cipher"
"crypto/rand"
"database/sql"
"encoding/json"
"errors"
"fmt"
"strings"
"time"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/catalog"
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/instance"
)
// Repository persists catalog snapshots and the desired instance registry.
type Repository struct {
db *sql.DB
now func() time.Time
aead cipher.AEAD
}
// Ping verifies that the SQLite connection backing the current request remains usable.
func (r *Repository) Ping(ctx context.Context) error { return r.db.PingContext(ctx) }
func NewRepository(db *sql.DB) *Repository {
return &Repository{db: db, now: time.Now}
}
// SetSecretKey configures encryption for instance runtime secrets. It is set
// from the service-owned master key and deliberately has no read API for web.
func (r *Repository) SetSecretKey(key []byte) error {
if len(key) != 32 {
return errors.New("instance encryption key must be exactly 32 bytes")
}
block, err := aes.NewCipher(key)
if err != nil {
return err
}
aead, err := cipher.NewGCM(block)
if err != nil {
return err
}
r.aead = aead
return nil
}
func (r *Repository) SaveInstanceSecrets(ctx context.Context, instanceID string, values map[string]string) error {
if len(values) == 0 {
return nil
}
if r.aead == nil {
return errors.New("instance secret storage unavailable")
}
body, err := json.Marshal(values)
if err != nil {
return err
}
nonce := make([]byte, r.aead.NonceSize())
if _, err = rand.Read(nonce); err != nil {
return err
}
encrypted := r.aead.Seal(nonce, nonce, body, []byte(instanceID))
_, err = r.db.ExecContext(ctx, `INSERT INTO instance_secrets(instance_id, encrypted_values, updated_at) VALUES(?,?,?) ON CONFLICT(instance_id) DO UPDATE SET encrypted_values=excluded.encrypted_values, updated_at=excluded.updated_at`, instanceID, encrypted, r.now().UTC().Format(time.RFC3339Nano))
if err != nil {
return fmt.Errorf("save instance secrets: %w", err)
}
return nil
}
func (r *Repository) LoadInstanceSecrets(ctx context.Context, instanceID string) (map[string]string, error) {
if r.aead == nil {
return nil, errors.New("instance secret storage unavailable")
}
var encrypted []byte
if err := r.db.QueryRowContext(ctx, `SELECT encrypted_values FROM instance_secrets WHERE instance_id=?`, instanceID).Scan(&encrypted); err != nil {
if errors.Is(err, sql.ErrNoRows) {
return map[string]string{}, nil
}
return nil, err
}
n := r.aead.NonceSize()
if len(encrypted) < n {
return nil, errors.New("invalid encrypted instance secrets")
}
body, err := r.aead.Open(nil, encrypted[:n], encrypted[n:], []byte(instanceID))
if err != nil {
return nil, errors.New("invalid encrypted instance secrets")
}
out := map[string]string{}
if err := json.Unmarshal(body, &out); err != nil {
return nil, errors.New("invalid encrypted instance secrets")
}
return out, nil
}
func (r *Repository) Sync(ctx context.Context, snapshots []catalog.Snapshot) error {
for _, snapshot := range snapshots {
var digest string
err := r.db.QueryRowContext(ctx, "SELECT digest FROM template_versions WHERE template_id=? AND version=?", snapshot.Template.ID, snapshot.Template.Version).Scan(&digest)
if err == nil && digest != snapshot.Digest {
return fmt.Errorf("%w: %s@%s", catalog.ErrImmutableSnapshot, snapshot.Template.ID, snapshot.Template.Version)
}
if err != nil && !errors.Is(err, sql.ErrNoRows) {
return fmt.Errorf("check template snapshot: %w", err)
}
}
return r.Replace(ctx, snapshots)
}
// Replace synchronizes the current local catalog. It is used only for
// administrator-owned local templates, for which the directory is authoritative.
func (r *Repository) Replace(ctx context.Context, snapshots []catalog.Snapshot) error {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin catalog sync: %w", err)
}
defer func() { _ = tx.Rollback() }()
now := r.now().UTC().Format(time.RFC3339Nano)
ids := make([]string, 0, len(snapshots))
for _, snapshot := range snapshots {
moduleBundle, marshalErr := json.Marshal(snapshot.ModuleFiles)
if marshalErr != nil {
return fmt.Errorf("encode template module bundle: %w", marshalErr)
}
assetBundle, marshalErr := json.Marshal(snapshot.AssetFiles)
if marshalErr != nil {
return fmt.Errorf("encode template asset bundle: %w", marshalErr)
}
ids = append(ids, snapshot.Template.ID)
_, err = tx.ExecContext(ctx, `INSERT INTO templates(id, origin, trust_status, active_version, available, created_at, updated_at)
VALUES (?, ?, ?, ?, 1, ?, ?)
ON CONFLICT(id) DO UPDATE SET origin=excluded.origin, trust_status=excluded.trust_status, active_version=excluded.active_version, available=1, updated_at=excluded.updated_at`,
snapshot.Template.ID, snapshot.Origin, snapshot.Origin, snapshot.Template.Version, now, now)
if err != nil {
return fmt.Errorf("upsert catalog template: %w", err)
}
_, err = tx.ExecContext(ctx, `INSERT INTO template_versions(template_id, version, schema_version, canonical_yaml, digest, game_id, game_name, description, image, asset_bundle, module_bundle, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(template_id, version) DO UPDATE SET schema_version=excluded.schema_version, canonical_yaml=excluded.canonical_yaml, digest=excluded.digest, game_id=excluded.game_id, game_name=excluded.game_name, description=excluded.description, image=excluded.image, asset_bundle=excluded.asset_bundle, module_bundle=excluded.module_bundle`, snapshot.Template.ID, snapshot.Template.Version, snapshot.Template.SchemaVersion, snapshot.CanonicalYAML, snapshot.Digest, snapshot.Template.Game.ID, snapshot.Template.Game.Name, snapshot.Template.Game.Description, snapshot.Template.Game.Artwork.Image, assetBundle, moduleBundle, now)
if err != nil {
return fmt.Errorf("upsert template snapshot: %w", err)
}
}
if len(ids) == 0 {
if _, err = tx.ExecContext(ctx, "UPDATE templates SET available=0, updated_at=?", now); err != nil {
return fmt.Errorf("hide removed templates: %w", err)
}
} else {
placeholders, args := make([]string, len(ids)), make([]any, 0, len(ids)+1)
args = append(args, now)
for i, id := range ids {
placeholders[i], args = "?", append(args, id)
}
if _, err = tx.ExecContext(ctx, "UPDATE templates SET available=0, updated_at=? WHERE id NOT IN ("+strings.Join(placeholders, ",")+")", args...); err != nil {
return fmt.Errorf("hide removed templates: %w", err)
}
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit catalog sync: %w", err)
}
return nil
}
func (r *Repository) List(ctx context.Context) ([]catalog.Summary, error) {
rows, err := r.db.QueryContext(ctx, `SELECT t.id, t.active_version, v.game_id, v.game_name, v.description, v.image, t.trust_status, v.digest
FROM templates t JOIN template_versions v ON v.template_id=t.id AND v.version=t.active_version WHERE t.available=1 ORDER BY v.game_name, t.id`)
if err != nil {
return nil, fmt.Errorf("list catalog: %w", err)
}
defer rows.Close()
var result []catalog.Summary
for rows.Next() {
var summary catalog.Summary
if err := rows.Scan(&summary.ID, &summary.Version, &summary.GameID, &summary.GameName, &summary.Description, &summary.Image, &summary.TrustStatus, &summary.Digest); err != nil {
return nil, fmt.Errorf("scan catalog: %w", err)
}
result = append(result, summary)
}
return result, rows.Err()
}
func (r *Repository) Get(ctx context.Context, id, version string) (catalog.Snapshot, error) {
var canonical, digest, origin string
var assetBundle, moduleBundle []byte
err := r.db.QueryRowContext(ctx, `SELECT v.canonical_yaml, v.digest, t.origin, v.asset_bundle, v.module_bundle FROM template_versions v JOIN templates t ON t.id=v.template_id
WHERE v.template_id=? AND v.version=?`, id, version).Scan(&canonical, &digest, &origin, &assetBundle, &moduleBundle)
if errors.Is(err, sql.ErrNoRows) {
return catalog.Snapshot{}, catalog.ErrTemplateNotFound
}
if err != nil {
return catalog.Snapshot{}, fmt.Errorf("load template snapshot: %w", err)
}
var template catalog.Template
if err := json.Unmarshal([]byte(canonical), &template); err != nil {
return catalog.Snapshot{}, fmt.Errorf("decode stored template snapshot: %w", err)
}
var moduleFiles map[string][]byte
if len(moduleBundle) != 0 && json.Unmarshal(moduleBundle, &moduleFiles) != nil {
return catalog.Snapshot{}, errors.New("decode stored template module bundle")
}
var assetFiles map[string][]byte
if len(assetBundle) != 0 && json.Unmarshal(assetBundle, &assetFiles) != nil {
return catalog.Snapshot{}, errors.New("decode stored template asset bundle")
}
return catalog.Snapshot{Template: template, CanonicalYAML: canonical, Digest: digest, Origin: origin, AssetFiles: assetFiles, ModuleFiles: moduleFiles}, nil
}
func (r *Repository) CreateDraft(ctx context.Context, draft instance.Draft) error {
if draft.ID == "" {
return errors.New("draft instance ID is required")
}
now := r.now().UTC().Format(time.RFC3339Nano)
previewJSON, err := json.Marshal(draft.Preview)
if err != nil {
return fmt.Errorf("encode draft preview: %w", err)
}
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()
_, err = tx.ExecContext(ctx, `INSERT INTO instances(id, slug, display_name, template_id, template_version, template_digest, revision, lifecycle_state, preview_json, plan_digest, custom_labels_json, docker_user_mode, docker_uid, docker_gid, image_tag_mode, image_tag, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, 1, 'draft', ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, draft.ID, draft.Preview.Slug, draft.Preview.DisplayName, draft.Preview.Template.ID, draft.Preview.Template.Version, draft.Preview.Template.Digest, string(previewJSON), draft.Preview.PlanDigest, string(mustJSON(draft.Preview.CustomLabels)), draft.Preview.DockerUser.Mode, draft.Preview.DockerUser.UID, draft.Preview.DockerUser.GID, draft.Preview.ImageTag.Mode, draft.Preview.ImageTag.Tag, now, now)
if err != nil {
return fmt.Errorf("create draft instance: %w", err)
}
if _, err := tx.ExecContext(ctx, `INSERT INTO configuration_revisions(instance_id, revision, redacted_snapshot, reason, created_at) VALUES(?,1,?,'creation',?)`, draft.ID, string(previewJSON), now); err != nil {
return fmt.Errorf("create initial configuration revision: %w", err)
}
return tx.Commit()
}
func (r *Repository) GetInstance(ctx context.Context, id string) (instance.StoredInstance, error) {
return scanInstance(r.db.QueryRowContext(ctx, `SELECT id, preview_json, lifecycle_state, observed_state, COALESCE(container_id, ''), plan_digest, desired_running, container_config_pending FROM instances WHERE id=? AND deleted_at IS NULL`, id))
}
func (r *Repository) ListLifecycleInstances(ctx context.Context) ([]instance.StoredInstance, error) {
rows, err := r.db.QueryContext(ctx, `SELECT id, preview_json, lifecycle_state, observed_state, COALESCE(container_id, ''), plan_digest, desired_running, container_config_pending FROM instances WHERE deleted_at IS NULL AND lifecycle_state != 'draft' ORDER BY id`)
if err != nil {
return nil, fmt.Errorf("list lifecycle instances: %w", err)
}
defer rows.Close()
var result []instance.StoredInstance
for rows.Next() {
value, err := scanInstance(rows)
if err != nil {
return nil, err
}
result = append(result, value)
}
return result, rows.Err()
}
type rowScanner interface{ Scan(...any) error }
func scanInstance(row rowScanner) (instance.StoredInstance, error) {
var value instance.StoredInstance
var previewJSON string
var desired int
var pending int
err := row.Scan(&value.ID, &previewJSON, &value.LifecycleState, &value.ObservedState, &value.ContainerID, &value.PlanDigest, &desired, &pending)
if errors.Is(err, sql.ErrNoRows) {
return instance.StoredInstance{}, instance.ErrInstanceNotFound
}
if err != nil {
return instance.StoredInstance{}, fmt.Errorf("scan instance: %w", err)
}
if err := json.Unmarshal([]byte(previewJSON), &value.Preview); err != nil {
return instance.StoredInstance{}, fmt.Errorf("decode instance preview: %w", err)
}
value.DesiredRunning = desired != 0
value.ContainerConfigPending = pending != 0
return value, nil
}
func (r *Repository) BeginOperation(ctx context.Context, operationID, instanceID, kind, lifecycleState string) (instance.StoredInstance, error) {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return instance.StoredInstance{}, fmt.Errorf("begin instance operation: %w", err)
}
defer func() { _ = tx.Rollback() }()
current, err := scanInstance(tx.QueryRowContext(ctx, `SELECT id, preview_json, lifecycle_state, observed_state, COALESCE(container_id, ''), plan_digest, desired_running, container_config_pending FROM instances WHERE id=? AND deleted_at IS NULL`, instanceID))
if err != nil {
return instance.StoredInstance{}, err
}
var active int
err = tx.QueryRowContext(ctx, `SELECT 1 FROM instance_operations WHERE instance_id=? AND state='running' LIMIT 1`, instanceID).Scan(&active)
if err == nil {
return instance.StoredInstance{}, instance.ErrOperationConflict
}
if !errors.Is(err, sql.ErrNoRows) {
return instance.StoredInstance{}, fmt.Errorf("check active operation: %w", err)
}
now := r.now().UTC().Format(time.RFC3339Nano)
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)
}
if err := tx.Commit(); err != nil {
return instance.StoredInstance{}, fmt.Errorf("commit instance operation: %w", err)
}
current.LifecycleState = lifecycleState
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 {
return fmt.Errorf("begin operation completion: %w", err)
}
defer func() { _ = tx.Rollback() }()
now := r.now().UTC().Format(time.RFC3339Nano)
desired := 0
if desiredRunning {
desired = 1
}
result, err := tx.ExecContext(ctx, `UPDATE instance_operations SET state='succeeded', phase='complete', error_code=NULL, updated_at=?, completed_at=? WHERE id=? AND state='running'`, now, now, operationID)
if err != nil {
return fmt.Errorf("complete instance operation: %w", err)
}
changed, _ := result.RowsAffected()
if changed != 1 {
return instance.ErrOperationConflict
}
var container any
if containerID != "" {
container = containerID
}
_, err = tx.ExecContext(ctx, `UPDATE instances SET lifecycle_state=?, observed_state=?, container_id=?, plan_digest=?, desired_running=?, last_error_code=?, updated_at=? WHERE id=(SELECT instance_id FROM instance_operations WHERE id=?)`, lifecycleState, observedState, container, planDigest, desired, nullable(errorCode), now, operationID)
if err != nil {
return fmt.Errorf("update completed instance: %w", err)
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit operation completion: %w", err)
}
return nil
}
func (r *Repository) FailOperation(ctx context.Context, operationID, lifecycleState, errorCode string) error {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin operation failure: %w", err)
}
defer func() { _ = tx.Rollback() }()
now := r.now().UTC().Format(time.RFC3339Nano)
result, err := tx.ExecContext(ctx, `UPDATE instance_operations SET state='failed', phase='failed', error_code=?, updated_at=?, completed_at=? WHERE id=? AND state='running'`, errorCode, now, now, operationID)
if err != nil {
return fmt.Errorf("fail instance operation: %w", err)
}
changed, _ := result.RowsAffected()
if changed != 1 {
return instance.ErrOperationConflict
}
_, err = tx.ExecContext(ctx, `UPDATE instances SET lifecycle_state=?, observed_state='unknown', last_error_code=?, updated_at=? WHERE id=(SELECT instance_id FROM instance_operations WHERE id=?)`, lifecycleState, errorCode, now, operationID)
if err != nil {
return fmt.Errorf("update failed instance: %w", err)
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit operation failure: %w", err)
}
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 {
desired = 1
}
var container any
if containerID != "" {
container = containerID
}
result, err := r.db.ExecContext(ctx, `UPDATE instances SET lifecycle_state=?, observed_state=?, container_id=?, desired_running=?, last_error_code=?, updated_at=? WHERE id=? AND deleted_at IS NULL`, lifecycleState, observedState, container, desired, nullable(errorCode), r.now().UTC().Format(time.RFC3339Nano), instanceID)
if err != nil {
return fmt.Errorf("update instance observation: %w", err)
}
changed, _ := result.RowsAffected()
if changed != 1 {
return instance.ErrInstanceNotFound
}
return nil
}
func (r *Repository) RecoverInterruptedOperations(ctx context.Context) error {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin interrupted-operation recovery: %w", err)
}
defer func() { _ = tx.Rollback() }()
now := r.now().UTC().Format(time.RFC3339Nano)
if _, err := tx.ExecContext(ctx, `UPDATE instances SET lifecycle_state='intervention_required', last_error_code='operation_interrupted', updated_at=? WHERE id IN (SELECT instance_id FROM instance_operations WHERE state='running')`, now); err != nil {
return fmt.Errorf("mark interrupted instances: %w", err)
}
if _, err := tx.ExecContext(ctx, `UPDATE instance_operations SET state='intervention_required', phase='interrupted', error_code='operation_interrupted', updated_at=?, completed_at=? WHERE state='running'`, now, now); err != nil {
return fmt.Errorf("mark interrupted operations: %w", err)
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit interrupted-operation recovery: %w", err)
}
return nil
}
func nullable(value string) any {
if value == "" {
return nil
}
return value
}