Files
ThothII/tools/tht/internal/serverops/operations.go

391 lines
12 KiB
Go

// Package serverops implements bounded, installation-aware server maintenance operations.
package serverops
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"strconv"
"strings"
"github.com/aritmolab/thothii/tools/tht/internal/compose"
"github.com/aritmolab/thothii/tools/tht/internal/config"
)
var (
ErrConfirmationRequired = errors.New("explicit confirmation is required")
ErrUnsafeState = errors.New("server operation refused in the current state")
)
type Runner interface {
Run(context.Context, []string, io.Reader) (compose.Result, error)
}
type Stage string
const (
StageContainerInspection Stage = "container-inspection"
StageMigrationConfig Stage = "migration-config"
StageMigrationVerification Stage = "migration-config-verification"
StageSessionMigration Stage = "session-migration"
StageContainerRemoval Stage = "container-removal"
StageRemovalVerification Stage = "removal-verification"
)
type ExitClass string
const (
ExitClassNonzero ExitClass = "nonzero-exit"
ExitClassUnavailable ExitClass = "unavailable"
ExitClassTimeout ExitClass = "timeout"
ExitClassInvocation ExitClass = "invocation-failure"
)
// OperationError reports only allowlisted operation metadata from Error. Complete subprocess
// detail is exposed separately so the CLI can redact it before applying its display bound.
type OperationError struct {
stage Stage
class ExitClass
detail string
}
func (e *OperationError) Error() string {
return fmt.Sprintf("stage=%s class=%s", e.stage, e.class)
}
func (e *OperationError) Stage() Stage {
return e.stage
}
func (e *OperationError) Class() ExitClass {
return e.class
}
func (e *OperationError) Detail() string {
return e.detail
}
type MigrationStatus struct {
Applied []string `json:"applied"`
Drifted []string `json:"drifted"`
Pending []string `json:"pending"`
}
type Container struct {
ID string `json:"ID"`
Name string `json:"Name"`
Service string `json:"Service"`
State string `json:"State"`
}
type RemovalResult struct {
Targets []Container
Preserved int
}
// MigrateSessions runs only the one-shot migration service and proves the resulting schema state.
func MigrateSessions(ctx context.Context, installation config.Installation, runner Runner, confirmed bool) (MigrationStatus, error) {
if !confirmed {
return MigrationStatus{}, ErrConfirmationRequired
}
if installation.Profile != "server" {
return MigrationStatus{}, fmt.Errorf("%w: session migration requires a server installation", ErrUnsafeState)
}
containers, err := inspectContainers(ctx, installation, runner, StageContainerInspection)
if err != nil {
return MigrationStatus{}, err
}
if err := requireStopped(containers); err != nil {
return MigrationStatus{}, err
}
rendered, err := runCompose(ctx, runner, StageMigrationConfig, installation.ComposeArgs("--profile", "session-migrate", "config", "--format", "json"))
if err != nil {
return MigrationStatus{}, err
}
coreImage, err := selectedCoreImage(rendered.Stdout)
if err != nil {
return MigrationStatus{}, err
}
override, cleanup, err := migrationOverride(installation, coreImage)
if err != nil {
return MigrationStatus{}, err
}
defer cleanup()
configArgs, err := installation.ComposeArgsWithFinalOverride(override, "--profile", "session-migrate", "config", "--format", "json")
if err != nil {
return MigrationStatus{}, err
}
finalConfig, err := runCompose(ctx, runner, StageMigrationVerification, configArgs)
if err != nil {
return MigrationStatus{}, err
}
if err := requireMigrationImage(finalConfig.Stdout, coreImage); err != nil {
return MigrationStatus{}, err
}
runArgs, err := installation.ComposeArgsWithFinalOverride(
override, "--profile", "session-migrate", "run", "--rm", "--no-deps", "--no-TTY", "session-migrate",
)
if err != nil {
return MigrationStatus{}, err
}
result, err := runCompose(ctx, runner, StageSessionMigration, runArgs)
if err != nil {
return MigrationStatus{}, err
}
status, err := parseMigrationStatus(result.Stdout)
if err != nil {
return MigrationStatus{}, err
}
if len(status.Pending) != 0 || len(status.Drifted) != 0 {
return status, fmt.Errorf("%w: session migration did not finish cleanly", ErrUnsafeState)
}
return status, nil
}
// Remove deletes only the exact stopped core/frontend container IDs displayed by the command.
// A nil confirmation performs inspection only; a non-nil confirmation must equal every target ID.
func Remove(ctx context.Context, installation config.Installation, runner Runner, confirmedIDs []string) (RemovalResult, error) {
if installation.Profile != "server" {
return RemovalResult{}, fmt.Errorf("%w: removal requires a server installation", ErrUnsafeState)
}
targets, err := inspectContainers(ctx, installation, runner, StageContainerInspection)
result := RemovalResult{Targets: targets}
if err != nil {
return result, err
}
if err := requireStopped(targets); err != nil {
return result, err
}
if confirmedIDs == nil {
return result, ErrConfirmationRequired
}
if !sameTargetIDs(targets, confirmedIDs) {
return result, fmt.Errorf("%w: confirmed container IDs differ from current targets", ErrUnsafeState)
}
paths, err := installation.PreservationPaths()
if err != nil {
return result, fmt.Errorf("%w: preservation paths could not be verified", ErrUnsafeState)
}
snapshots, err := snapshotPaths(paths)
if err != nil {
return result, err
}
if len(targets) > 0 {
args := []string{"rm"}
for _, target := range targets {
args = append(args, target.ID)
}
if _, err := runDocker(ctx, runner, StageContainerRemoval, args); err != nil {
return result, err
}
}
remaining, err := inspectContainers(ctx, installation, runner, StageRemovalVerification)
if err != nil {
return result, err
}
if len(remaining) != 0 {
return result, fmt.Errorf("%w: installation containers changed during removal", ErrUnsafeState)
}
if err := verifySnapshots(snapshots); err != nil {
return result, err
}
result.Preserved = len(snapshots)
return result, nil
}
func sameTargetIDs(targets []Container, confirmed []string) bool {
if len(targets) != len(confirmed) {
return false
}
wanted := make(map[string]struct{}, len(confirmed))
for _, id := range confirmed {
if strings.TrimSpace(id) == "" {
return false
}
if _, duplicate := wanted[id]; duplicate {
return false
}
wanted[id] = struct{}{}
}
for _, target := range targets {
if _, exists := wanted[target.ID]; !exists {
return false
}
}
return true
}
func inspectContainers(ctx context.Context, installation config.Installation, runner Runner, stage Stage) ([]Container, error) {
result, err := runCompose(ctx, runner, stage, installation.ComposeArgs("ps", "--all", "--format", "json", "core", "frontend"))
if err != nil {
return nil, err
}
var containers []Container
if err := json.Unmarshal([]byte(result.Stdout), &containers); err != nil {
return nil, fmt.Errorf("%w: Compose returned invalid container status", ErrUnsafeState)
}
seen := make(map[string]struct{})
for _, container := range containers {
if (container.Service != "core" && container.Service != "frontend") || container.ID == "" || container.Name == "" {
return nil, fmt.Errorf("%w: Compose returned an unexpected removal target", ErrUnsafeState)
}
if _, exists := seen[container.ID]; exists {
return nil, fmt.Errorf("%w: Compose returned duplicate container IDs", ErrUnsafeState)
}
seen[container.ID] = struct{}{}
}
return containers, nil
}
func requireStopped(containers []Container) error {
for _, container := range containers {
if strings.ToLower(container.State) != "exited" {
return fmt.Errorf("%w: %s is not stopped", ErrUnsafeState, container.Service)
}
}
return nil
}
func selectedCoreImage(document string) (string, error) {
services, err := renderedServices(document)
if err != nil {
return "", err
}
core, exists := services["core"]
if !exists || strings.TrimSpace(core.Image) == "" {
return "", fmt.Errorf("%w: rendered core image is missing", ErrUnsafeState)
}
if _, exists := services["session-migrate"]; !exists {
return "", fmt.Errorf("%w: rendered migration service is missing", ErrUnsafeState)
}
return core.Image, nil
}
type renderedService struct {
Image string `json:"image"`
Build json.RawMessage `json:"build"`
}
func renderedServices(document string) (map[string]renderedService, error) {
var configDocument struct {
Services map[string]renderedService `json:"services"`
}
if err := json.Unmarshal([]byte(document), &configDocument); err != nil {
return nil, fmt.Errorf("%w: Compose returned invalid rendered configuration", ErrUnsafeState)
}
return configDocument.Services, nil
}
func requireMigrationImage(document, coreImage string) error {
services, err := renderedServices(document)
if err != nil {
return err
}
migrator, exists := services["session-migrate"]
if !exists || migrator.Image != coreImage {
return fmt.Errorf("%w: migration image differs from selected core image", ErrUnsafeState)
}
if len(migrator.Build) != 0 && strings.TrimSpace(string(migrator.Build)) != "null" {
return fmt.Errorf("%w: migration service unexpectedly declares a build", ErrUnsafeState)
}
return nil
}
func migrationOverride(installation config.Installation, image string) (string, func(), error) {
control := installation.ControlDirectory()
if err := os.MkdirAll(control, 0o700); err != nil {
return "", func() {}, errors.New("migration control directory could not be created")
}
info, err := os.Lstat(control)
if err != nil || !info.IsDir() || info.Mode()&os.ModeSymlink != 0 {
return "", func() {}, errors.New("migration control directory is unsafe")
}
directory, err := os.MkdirTemp(control, "session-migrate-")
if err != nil {
return "", func() {}, errors.New("migration override directory could not be created")
}
cleanup := func() {
_ = os.Remove(filepath.Join(directory, "override.yaml"))
_ = os.Remove(directory)
}
path := filepath.Join(directory, "override.yaml")
contents := "services:\n session-migrate:\n build: !reset null\n image: " + strconv.Quote(image) + "\n"
if err := os.WriteFile(path, []byte(contents), 0o600); err != nil {
cleanup()
return "", func() {}, errors.New("migration override could not be written")
}
return path, cleanup, nil
}
func parseMigrationStatus(document string) (MigrationStatus, error) {
var status MigrationStatus
decoder := json.NewDecoder(strings.NewReader(document))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&status); err != nil || status.Applied == nil || status.Drifted == nil || status.Pending == nil {
return MigrationStatus{}, fmt.Errorf("%w: migration did not return verified JSON status", ErrUnsafeState)
}
var extra any
if err := decoder.Decode(&extra); !errors.Is(err, io.EOF) {
return MigrationStatus{}, fmt.Errorf("%w: migration returned trailing output", ErrUnsafeState)
}
return status, nil
}
type pathSnapshot struct {
path string
info os.FileInfo
}
func snapshotPaths(paths []string) ([]pathSnapshot, error) {
snapshots := make([]pathSnapshot, 0, len(paths))
for _, path := range paths {
info, err := os.Stat(path)
if err != nil {
return nil, fmt.Errorf("%w: preservation target is unavailable", ErrUnsafeState)
}
snapshots = append(snapshots, pathSnapshot{path: path, info: info})
}
return snapshots, nil
}
func verifySnapshots(snapshots []pathSnapshot) error {
for _, snapshot := range snapshots {
info, err := os.Stat(snapshot.path)
if err != nil || !os.SameFile(snapshot.info, info) {
return fmt.Errorf("%w: a preserved path changed during removal", ErrUnsafeState)
}
}
return nil
}
func runCompose(ctx context.Context, runner Runner, stage Stage, args []string) (compose.Result, error) {
return runDocker(ctx, runner, stage, args)
}
func runDocker(ctx context.Context, runner Runner, stage Stage, args []string) (compose.Result, error) {
result, err := runner.Run(ctx, args, nil)
if err != nil {
class := ExitClassInvocation
switch {
case errors.Is(ctx.Err(), context.DeadlineExceeded):
class = ExitClassTimeout
case result.ExitCode == 127:
class = ExitClassUnavailable
case result.ExitCode != 0:
class = ExitClassNonzero
}
detail := result.Stderr
if strings.TrimSpace(detail) == "" {
detail = err.Error()
}
return result, &OperationError{stage: stage, class: class, detail: detail}
}
return result, nil
}