222 lines
8.2 KiB
Go
222 lines
8.2 KiB
Go
package sqlite
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"time"
|
|
|
|
"git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/instance"
|
|
)
|
|
|
|
const globalLabelsKey = "game_container_labels"
|
|
const gameContainerUIDKey = "game_container_uid"
|
|
const gameContainerGIDKey = "game_container_gid"
|
|
|
|
func (r *Repository) GetGameContainerRuntimeIdentity(ctx context.Context) (instance.RuntimeIdentity, error) {
|
|
identity := instance.RuntimeIdentity{UID: instance.DefaultGameContainerUID, GID: instance.DefaultGameContainerGID}
|
|
for _, setting := range []struct {
|
|
key string
|
|
value *uint32
|
|
}{{gameContainerUIDKey, &identity.UID}, {gameContainerGIDKey, &identity.GID}} {
|
|
var body string
|
|
err := r.db.QueryRowContext(ctx, `SELECT value_json FROM system_settings WHERE key=?`, setting.key).Scan(&body)
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
continue
|
|
}
|
|
if err != nil {
|
|
return instance.RuntimeIdentity{}, fmt.Errorf("load game-container runtime identity: %w", err)
|
|
}
|
|
var value uint64
|
|
if err := json.Unmarshal([]byte(body), &value); err != nil || value > ^uint64(0)>>32 {
|
|
return instance.RuntimeIdentity{}, errors.New("stored game-container runtime identity is invalid")
|
|
}
|
|
*setting.value = uint32(value)
|
|
}
|
|
return identity, nil
|
|
}
|
|
|
|
func (r *Repository) SetGameContainerRuntimeIdentity(ctx context.Context, identity instance.RuntimeIdentity) error {
|
|
if err := identity.Validate(); err != nil {
|
|
return err
|
|
}
|
|
tx, err := r.db.BeginTx(ctx, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer func() { _ = tx.Rollback() }()
|
|
now := r.now().UTC().Format(time.RFC3339Nano)
|
|
for _, setting := range []struct {
|
|
key string
|
|
value uint32
|
|
}{{gameContainerUIDKey, identity.UID}, {gameContainerGIDKey, identity.GID}} {
|
|
body := string(mustJSON(setting.value))
|
|
if _, err := tx.ExecContext(ctx, `INSERT INTO system_settings(key,value_json,revision,updated_at) VALUES(?,?,1,?) ON CONFLICT(key) DO UPDATE SET value_json=excluded.value_json, revision=system_settings.revision+1, updated_at=excluded.updated_at`, setting.key, body, now); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return tx.Commit()
|
|
}
|
|
|
|
func (r *Repository) GetGlobalLabels(ctx context.Context) (map[string]string, error) {
|
|
var body string
|
|
err := r.db.QueryRowContext(ctx, `SELECT value_json FROM system_settings WHERE key=?`, globalLabelsKey).Scan(&body)
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return map[string]string{}, nil
|
|
}
|
|
if err != nil {
|
|
return nil, fmt.Errorf("load global game-container labels: %w", err)
|
|
}
|
|
var labels map[string]string
|
|
if err := json.Unmarshal([]byte(body), &labels); err != nil {
|
|
return nil, fmt.Errorf("decode global game-container labels: %w", err)
|
|
}
|
|
return labels, nil
|
|
}
|
|
|
|
func (r *Repository) SetGlobalLabels(ctx context.Context, labels map[string]string, pending bool) (int, int, error) {
|
|
tx, err := r.db.BeginTx(ctx, nil)
|
|
if err != nil {
|
|
return 0, 0, err
|
|
}
|
|
defer func() { _ = tx.Rollback() }()
|
|
now := r.now().UTC().Format(time.RFC3339Nano)
|
|
body := string(mustJSON(labels))
|
|
if _, err := tx.ExecContext(ctx, `INSERT INTO system_settings(key,value_json,revision,updated_at) VALUES(?,?,1,?) ON CONFLICT(key) DO UPDATE SET value_json=excluded.value_json, revision=system_settings.revision+1, updated_at=excluded.updated_at`, globalLabelsKey, body, now); err != nil {
|
|
return 0, 0, err
|
|
}
|
|
rows, err := tx.QueryContext(ctx, `SELECT id, preview_json, desired_running FROM instances WHERE deleted_at IS NULL AND lifecycle_state!='draft'`)
|
|
if err != nil {
|
|
return 0, 0, err
|
|
}
|
|
type update struct {
|
|
id, body string
|
|
running bool
|
|
}
|
|
var updates []update
|
|
for rows.Next() {
|
|
var id, previewBody string
|
|
var running int
|
|
if err := rows.Scan(&id, &previewBody, &running); err != nil {
|
|
rows.Close()
|
|
return 0, 0, err
|
|
}
|
|
var preview instance.Preview
|
|
if json.Unmarshal([]byte(previewBody), &preview) != nil {
|
|
rows.Close()
|
|
return 0, 0, errors.New("decode instance configuration")
|
|
}
|
|
preview.GlobalLabels = labels
|
|
updates = append(updates, update{id: id, body: string(mustJSON(preview)), running: running != 0})
|
|
}
|
|
rows.Close()
|
|
running := 0
|
|
pendingValue := 0
|
|
if pending {
|
|
pendingValue = 1
|
|
}
|
|
for _, update := range updates {
|
|
if update.running {
|
|
running++
|
|
}
|
|
if _, err := tx.ExecContext(ctx, `UPDATE instances SET preview_json=?, container_config_pending=?, revision=revision+1, updated_at=? WHERE id=?`, update.body, pendingValue, now, update.id); err != nil {
|
|
return 0, 0, err
|
|
}
|
|
if _, err := tx.ExecContext(ctx, `INSERT INTO configuration_revisions(instance_id, revision, redacted_snapshot, reason, created_at) SELECT id, revision, preview_json, 'global_labels', ? FROM instances WHERE id=?`, now, update.id); err != nil {
|
|
return 0, 0, err
|
|
}
|
|
if _, err := tx.ExecContext(ctx, `DELETE FROM configuration_revisions WHERE instance_id=? AND revision NOT IN (SELECT revision FROM configuration_revisions WHERE instance_id=? ORDER BY revision DESC LIMIT 10)`, update.id, update.id); err != nil {
|
|
return 0, 0, err
|
|
}
|
|
}
|
|
if err := tx.Commit(); err != nil {
|
|
return 0, 0, err
|
|
}
|
|
return len(updates), running, nil
|
|
}
|
|
|
|
func (r *Repository) SaveInstanceConfiguration(ctx context.Context, id string, preview instance.Preview, pending bool, reason, actorID string) error {
|
|
body := string(mustJSON(preview))
|
|
labels := string(mustJSON(preview.CustomLabels))
|
|
pendingValue := 0
|
|
if pending {
|
|
pendingValue = 1
|
|
}
|
|
tx, err := r.db.BeginTx(ctx, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer func() { _ = tx.Rollback() }()
|
|
now := r.now().UTC().Format(time.RFC3339Nano)
|
|
result, err := tx.ExecContext(ctx, `UPDATE instances SET preview_json=?, custom_labels_json=?, image_tag_mode=?, image_tag=?, container_config_pending=?, revision=revision+1, updated_at=? WHERE id=? AND deleted_at IS NULL`, body, labels, preview.ImageTag.Mode, preview.ImageTag.Tag, pendingValue, now, id)
|
|
if err != nil {
|
|
return fmt.Errorf("save instance container configuration: %w", err)
|
|
}
|
|
changed, _ := result.RowsAffected()
|
|
if changed != 1 {
|
|
return instance.ErrInstanceNotFound
|
|
}
|
|
var revision int
|
|
if err := tx.QueryRowContext(ctx, `SELECT revision FROM instances WHERE id=?`, id).Scan(&revision); err != nil {
|
|
return err
|
|
}
|
|
if _, err := tx.ExecContext(ctx, `INSERT INTO configuration_revisions(instance_id, revision, redacted_snapshot, reason, created_by, created_at) VALUES(?,?,?,?,?,?)`, id, revision, body, reason, nullable(actorID), now); err != nil {
|
|
return err
|
|
}
|
|
if _, err := tx.ExecContext(ctx, `DELETE FROM configuration_revisions WHERE instance_id=? AND revision NOT IN (SELECT revision FROM configuration_revisions WHERE instance_id=? ORDER BY revision DESC LIMIT 10)`, id, id); err != nil {
|
|
return err
|
|
}
|
|
return tx.Commit()
|
|
}
|
|
|
|
func (r *Repository) ListConfigurationRevisions(ctx context.Context, id string) ([]instance.ConfigurationRevision, error) {
|
|
rows, err := r.db.QueryContext(ctx, `SELECT revision, redacted_snapshot, reason, COALESCE(created_by,''), created_at FROM configuration_revisions WHERE instance_id=? ORDER BY revision DESC`, id)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var result []instance.ConfigurationRevision
|
|
for rows.Next() {
|
|
var value instance.ConfigurationRevision
|
|
var body, created string
|
|
value.InstanceID = id
|
|
if err := rows.Scan(&value.Revision, &body, &value.Reason, &value.CreatedBy, &created); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := json.Unmarshal([]byte(body), &value.Snapshot); err != nil {
|
|
return nil, err
|
|
}
|
|
value.CreatedAt, _ = time.Parse(time.RFC3339Nano, created)
|
|
result = append(result, value)
|
|
}
|
|
return result, rows.Err()
|
|
}
|
|
|
|
func (r *Repository) GetConfigurationRevision(ctx context.Context, id string, revision int) (instance.ConfigurationRevision, error) {
|
|
values, err := r.ListConfigurationRevisions(ctx, id)
|
|
if err != nil {
|
|
return instance.ConfigurationRevision{}, err
|
|
}
|
|
for _, value := range values {
|
|
if value.Revision == revision {
|
|
return value, nil
|
|
}
|
|
}
|
|
return instance.ConfigurationRevision{}, instance.ErrInstanceNotFound
|
|
}
|
|
|
|
func (r *Repository) ClearContainerConfigPending(ctx context.Context, id, planDigest string) error {
|
|
_, err := r.db.ExecContext(ctx, `UPDATE instances SET container_config_pending=0, plan_digest=?, updated_at=? WHERE id=?`, planDigest, r.now().UTC().Format(time.RFC3339Nano), id)
|
|
return err
|
|
}
|
|
|
|
func mustJSON(value any) []byte {
|
|
body, err := json.Marshal(value)
|
|
if err != nil {
|
|
panic(err)
|
|
}
|
|
return body
|
|
}
|