From fbce91d7979f347f488a1ec6df3bfe3acb3d011a Mon Sep 17 00:00:00 2001 From: Zhan Ziyang Date: Sat, 15 Aug 2026 20:58:01 +0800 Subject: [PATCH] feat: backend container executor implement --- go.mod | 24 +- go.sum | 57 +++ internal/backendexecutor/executor.go | 361 ++++++++++++++++++ internal/backendexecutor/executor_test.go | 437 ++++++++++++++++++++++ internal/backendexecutor/operations.go | 200 ++++++++++ internal/containerengine/engine.go | 80 ++++ internal/containerengine/moby.go | 206 ++++++++++ internal/containerengine/moby_test.go | 31 ++ internal/healthcheck/actuator.go | 150 ++++++++ internal/healthcheck/actuator_test.go | 117 ++++++ 10 files changed, 1662 insertions(+), 1 deletion(-) create mode 100644 internal/backendexecutor/executor.go create mode 100644 internal/backendexecutor/executor_test.go create mode 100644 internal/backendexecutor/operations.go create mode 100644 internal/containerengine/engine.go create mode 100644 internal/containerengine/moby.go create mode 100644 internal/containerengine/moby_test.go create mode 100644 internal/healthcheck/actuator.go create mode 100644 internal/healthcheck/actuator_test.go diff --git a/go.mod b/go.mod index efe66ef..5d70689 100644 --- a/go.mod +++ b/go.mod @@ -2,10 +2,32 @@ module yms-daemon go 1.26 -require github.com/ncruces/go-sqlite3 v0.35.3 +require ( + github.com/containerd/errdefs v1.0.0 + github.com/distribution/reference v0.6.0 + github.com/moby/moby/api v1.55.0 + github.com/moby/moby/client v0.5.1 + github.com/ncruces/go-sqlite3 v0.35.3 + github.com/opencontainers/go-digest v1.0.0 + github.com/opencontainers/image-spec v1.1.1 +) require ( + github.com/Microsoft/go-winio v0.6.2 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/containerd/errdefs/pkg v0.3.0 // indirect + github.com/docker/go-connections v0.8.1 // indirect + github.com/docker/go-units v0.5.0 // indirect + github.com/felixge/httpsnoop v1.1.0 // indirect + github.com/go-logr/logr v1.4.4 // indirect + github.com/go-logr/stdr v1.2.2 // indirect + github.com/moby/docker-image-spec v1.3.1 // indirect github.com/ncruces/go-sqlite3-wasm/v3 v3.2.35304 // indirect github.com/ncruces/julianday v1.0.0 // indirect + go.opentelemetry.io/auto/sdk v1.2.1 // indirect + go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.70.0 // indirect + go.opentelemetry.io/otel v1.45.0 // indirect + go.opentelemetry.io/otel/metric v1.45.0 // indirect + go.opentelemetry.io/otel/trace v1.45.0 // indirect golang.org/x/sys v0.47.0 // indirect ) diff --git a/go.sum b/go.sum index ae032de..b8b1557 100644 --- a/go.sum +++ b/go.sum @@ -1,9 +1,66 @@ +github.com/Microsoft/go-winio v0.6.2 h1:F2VQgta7ecxGYO8k3ZZz3RS8fVIXVxONVUPlNERoyfY= +github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA8Ipt1oGCvU= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/containerd/errdefs v1.0.0 h1:tg5yIfIlQIrxYtu9ajqY42W3lpS19XqdxRQeEwYG8PI= +github.com/containerd/errdefs v1.0.0/go.mod h1:+YBYIdtsnF4Iw6nWZhJcqGSg/dwvV7tyJ/kCkyJ2k+M= +github.com/containerd/errdefs/pkg v0.3.0 h1:9IKJ06FvyNlexW690DXuQNx2KA2cUJXx151Xdx3ZPPE= +github.com/containerd/errdefs/pkg v0.3.0/go.mod h1:NJw6s9HwNuRhnjJhM7pylWwMyAkmCQvQ4GpJHEqRLVk= +github.com/distribution/reference v0.6.0 h1:0IXCQ5g4/QMHHkarYzh5l+u8T3t73zM5QvfrDyIgxBk= +github.com/distribution/reference v0.6.0/go.mod h1:BbU0aIcezP1/5jX/8MP0YiH4SdvB5Y4f/wlDRiLyi3E= +github.com/docker/go-connections v0.7.0 h1:6SsRfJddP22WMrCkj19x9WKjEDTB+ahsdiGYf0mN39c= +github.com/docker/go-connections v0.7.0/go.mod h1:no1qkHdjq7kLMGUXYAduOhYPSJxxvgWBh7ogVvptn3Q= +github.com/docker/go-connections v0.8.1 h1:JibmG5hULs5qXSr/cp/w3Pw5fZuStt4MOHMUExb29/M= +github.com/docker/go-connections v0.8.1/go.mod h1:no1qkHdjq7kLMGUXYAduOhYPSJxxvgWBh7ogVvptn3Q= +github.com/docker/go-units v0.5.0 h1:69rxXcBk27SvSaaxTtLh/8llcHD8vYHT7WSdRZ/jvr4= +github.com/docker/go-units v0.5.0/go.mod h1:fgPhTUdO+D/Jk86RDLlptpiXQzgHJF7gydDDbaIK4Dk= +github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg= +github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U= +github.com/felixge/httpsnoop v1.1.0 h1:3YtUj32ZZkqZtt3sZZsClsymw/QDuVfpNhoA31zeORc= +github.com/felixge/httpsnoop v1.1.0/go.mod h1:Zqxgdd+1Rkcz8euOqdr7lqgCRJztwr5hp9vDSi5UZCE= +github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= +github.com/go-logr/logr v1.4.2 h1:6pFjapn8bFcIbiKo3XT4j/BhANplGihG6tvd+8rYgrY= +github.com/go-logr/logr v1.4.2/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/logr v1.4.4 h1:tG4xh9yMsRCAiodLVTxyrkzSZ9+o0L1Kg/+cPVcbP/8= +github.com/go-logr/logr v1.4.4/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= +github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/moby/docker-image-spec v1.3.1 h1:jMKff3w6PgbfSa69GfNg+zN/XLhfXJGnEx3Nl2EsFP0= +github.com/moby/docker-image-spec v1.3.1/go.mod h1:eKmb5VW8vQEh/BAr2yvVNvuiJuY6UIocYsFu/DxxRpo= +github.com/moby/moby/api v1.55.0 h1:2/sexvQyqIWS8pRSCFddBfpW2qE7vR7FCL+vN8pxwMc= +github.com/moby/moby/api v1.55.0/go.mod h1:+RQ6wluLwtYaTd1WnPLykIDPekkuyD/ROWQClE83pzs= +github.com/moby/moby/client v0.5.1 h1:tYNaJno4c0HXz12y5BiqEDy0rVTYkWzI26lGvnTMiJw= +github.com/moby/moby/client v0.5.1/go.mod h1:odLstlZ6uSnfvAgVxMpvgmb8SUdd+siH2T0GBuxVAlM= github.com/ncruces/go-sqlite3 v0.35.3 h1:Ei07Zv1qfV/vyXzelhFsyS5Oh9TArBZHsmFk14Xv3GY= github.com/ncruces/go-sqlite3 v0.35.3/go.mod h1:i1rhym/NIiB5xeEfzbN+e24Y+i7NGUpf7C2xZ3Dpwks= github.com/ncruces/go-sqlite3-wasm/v3 v3.2.35304 h1:5NoQAewtgKNK3G4bjNPxVoGXu6F6NzLXWCTdD5FFAEY= github.com/ncruces/go-sqlite3-wasm/v3 v3.2.35304/go.mod h1:o8gr9w/50fXA5TDskg6bNUjvqmFfw4KaXth4q+yDSjg= github.com/ncruces/julianday v1.0.0 h1:fH0OKwa7NWvniGQtxdJRxAgkBMolni2BjDHaWTxqt7M= github.com/ncruces/julianday v1.0.0/go.mod h1:Dusn2KvZrrovOMJuOt0TNXL6tB7U2E8kvza5fFc9G7g= +github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U= +github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= +github.com/opencontainers/image-spec v1.1.1 h1:y0fUlFfIZhPF1W537XOLg0/fcx6zcHCJwooC2xJA040= +github.com/opencontainers/image-spec v1.1.1/go.mod h1:qpqAh3Dmcf36wStyyWU+kCeDgrGnAve2nCC8+7h8Q0M= +go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA= +go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A= +go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= +go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.60.0 h1:sbiXRNDSWJOTobXh5HyQKjq6wUC5tNybqjIqDpAY4CU= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.60.0/go.mod h1:69uWxva0WgAA/4bu2Yy70SLDBwZXuQ6PbBpbsa5iZrQ= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.70.0 h1:LMuyCAyfalSjDyjdC65nK6N0zoTT63+E/u95X0JovZI= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.70.0/go.mod h1:085m8qbm4hgc8rZWGDEa4vmyyo2c3nPxUslYUKUIU04= +go.opentelemetry.io/otel v1.35.0 h1:xKWKPxrxB6OtMCbmMY021CqC45J+3Onta9MqjhnusiQ= +go.opentelemetry.io/otel v1.35.0/go.mod h1:UEqy8Zp11hpkUrL73gSlELM0DupHoiq72dR+Zqel/+Y= +go.opentelemetry.io/otel v1.45.0 h1:pdrWmLHofpubmArBv1LgFSv1Z0Ie/ppdZzu+kUN5EeU= +go.opentelemetry.io/otel v1.45.0/go.mod h1:XZxIqPapzEYnhNSScF5DIqXhm/rYi0FzCe2XddAwZfQ= +go.opentelemetry.io/otel/metric v1.35.0 h1:0znxYu2SNyuMSQT4Y9WDWej0VpcsxkuklLa4/siN90M= +go.opentelemetry.io/otel/metric v1.35.0/go.mod h1:nKVFgxBZ2fReX6IlyW28MgZojkoAkJGaE8CpgeAU3oE= +go.opentelemetry.io/otel/metric v1.45.0 h1:7Eg1uH7CJ5cXv9is6tnBe1FI6rj1nwUdbFypRm3br/M= +go.opentelemetry.io/otel/metric v1.45.0/go.mod h1:HAPbm1nd3p1PmFH7v2dR+6BjXxw+Lq4a2+pndMAm08s= +go.opentelemetry.io/otel/trace v1.35.0 h1:dPpEfJu1sDIqruz7BHFG3c7528f6ddfSWfFDVt/xgMs= +go.opentelemetry.io/otel/trace v1.35.0/go.mod h1:WUk7DtFp1Aw2MkvqGdwiXYDZZNvA/1J8o6xRXLrIkyc= +go.opentelemetry.io/otel/trace v1.45.0 h1:l/mP6Uv7oNO7/TblbhpbgMidxhq1uO/rPsikOyVhxag= +go.opentelemetry.io/otel/trace v1.45.0/go.mod h1:qoJJA2xNMnxRrdISU/kLtfUH2wNeQbiv+jhs/CxI8bc= golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs= diff --git a/internal/backendexecutor/executor.go b/internal/backendexecutor/executor.go new file mode 100644 index 0000000..6508d7a --- /dev/null +++ b/internal/backendexecutor/executor.go @@ -0,0 +1,361 @@ +// Package backendexecutor prepares and starts one explicitly named backend container. +// Gateway switching is deliberately outside this package. +package backendexecutor + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "net/http" + "net/url" + "os" + "path/filepath" + "strconv" + "strings" + "time" + + "github.com/distribution/reference" + opencontainersdigest "github.com/opencontainers/go-digest" + + "yms-daemon/internal/containerengine" + "yms-daemon/internal/healthcheck" + "yms-daemon/internal/transaction" +) + +const ( + healthPath = "/yms/actuator/health" + healthTimeout = 120 * time.Second + healthInterval = time.Second + hostNetworkMode = "host" + bindMountType = "bind" + stepLoadImage = "backend.image.load" + stepCreateContainer = "backend.container.create" + stepStartContainer = "backend.container.start" + stepCheckHealth = "backend.container.health" +) + +// Request contains exact values supplied by the update package and local deployment configuration. +// ImageReference is opaque: the executor never extracts meaning from its tag. +type Request struct { + ArchivePath string + ImageReference string + ExpectedImageDigest string + Platform containerengine.Platform + ContainerName string + Port int + PortEnvironmentKey string + ConfigSource string + ConfigTarget string + RestartPolicy containerengine.RestartPolicy + HealthEndpoint string +} + +// Executor drives the persisted transaction up to SWITCHING after the new container is healthy. +type Executor struct { + store *transaction.Store + coordinator *transaction.Coordinator + engine containerengine.Engine + checker *healthcheck.ActuatorChecker +} + +func New(store *transaction.Store, coordinator *transaction.Coordinator, engine containerengine.Engine, httpClient *http.Client) (*Executor, error) { + if store == nil { + return nil, errors.New("transaction store is required") + } + if coordinator == nil { + return nil, errors.New("transaction coordinator is required") + } + if engine == nil { + return nil, errors.New("container engine is required") + } + checker, err := healthcheck.NewActuatorChecker(httpClient, healthInterval) + if err != nil { + return nil, err + } + return &Executor{store: store, coordinator: coordinator, engine: engine, checker: checker}, nil +} + +// Run resumes from the transaction's persisted state. It does not switch gateway traffic. +func (e *Executor) Run(ctx context.Context, transactionID string, request Request) error { + if strings.TrimSpace(transactionID) == "" { + return errors.New("transaction ID is required") + } + return e.coordinator.RunExclusive(ctx, func(ctx context.Context) error { + return e.run(ctx, transactionID, request) + }) +} + +func (e *Executor) run(ctx context.Context, transactionID string, request Request) error { + for { + record, err := e.store.Transaction(ctx, transactionID) + if err != nil { + return err + } + switch record.State { + case transaction.StateCreated: + if _, err := e.store.Transition(ctx, transactionID, transaction.StateValidating, "backend container validation started"); err != nil { + return err + } + case transaction.StateValidating: + if err := e.validate(ctx, request); err != nil { + _, transitionErr := e.store.Transition(ctx, transactionID, transaction.StateFailed, err.Error()) + return errors.Join(err, transitionErr) + } + if _, err := e.store.Transition(ctx, transactionID, transaction.StatePrepared, "backend container inputs validated"); err != nil { + return err + } + case transaction.StatePrepared: + if err := e.prepare(ctx, transactionID, request); err != nil { + return e.failUnlessRecoverable(ctx, transactionID, err) + } + if _, err := e.store.Transition(ctx, transactionID, transaction.StateStarting, "backend container prepared"); err != nil { + return err + } + case transaction.StateStarting: + // 重放 PREPARED 步骤会核对持久化意图,阻止恢复时换入另一组请求参数。 + if err := e.prepare(ctx, transactionID, request); err != nil { + return e.failUnlessRecoverable(ctx, transactionID, err) + } + if err := e.startAndCheck(ctx, transactionID, request); err != nil { + return e.failUnlessRecoverable(ctx, transactionID, err) + } + if _, err := e.store.Transition(ctx, transactionID, transaction.StateSwitching, "backend container is healthy"); err != nil { + return err + } + case transaction.StateSwitching: + return nil + default: + return fmt.Errorf("backend container executor cannot run transaction %s in state %s", transactionID, record.State) + } + } +} + +func (e *Executor) failUnlessRecoverable(ctx context.Context, transactionID string, cause error) error { + var uncertain *transaction.UncertainStepError + if errors.As(cause, &uncertain) || errors.Is(cause, transaction.ErrStepConflict) { + return cause + } + _, transitionErr := e.store.Transition(ctx, transactionID, transaction.StateFailed, cause.Error()) + return errors.Join(cause, transitionErr) +} + +func (e *Executor) validate(ctx context.Context, request Request) error { + if err := validateRequest(request); err != nil { + return err + } + if err := regularFile(request.ArchivePath, "image archive"); err != nil { + return err + } + if err := regularFile(request.ConfigSource, "backend configuration"); err != nil { + return err + } + if err := e.engine.Ping(ctx); err != nil { + return err + } + return nil +} + +func (e *Executor) prepare(ctx context.Context, transactionID string, request Request) error { + loadOperation := &loadImageOperation{ + engine: e.engine, + archivePath: request.ArchivePath, + imageReference: request.ImageReference, + expectedDigest: request.ExpectedImageDigest, + platform: request.Platform, + } + if _, err := e.coordinator.ExecuteStep(ctx, transactionID, loadIntent(request), loadOperation); err != nil { + return err + } + + image, err := e.engine.InspectImage(ctx, request.ImageReference) + if err != nil { + return fmt.Errorf("inspect prepared backend image: %w", err) + } + createOperation := &createContainerOperation{ + engine: e.engine, + expectedImage: image, + spec: containerSpec(request), + } + _, err = e.coordinator.ExecuteStep(ctx, transactionID, createIntent(request, image.ID), createOperation) + return err +} + +func (e *Executor) startAndCheck(ctx context.Context, transactionID string, request Request) error { + startOperation := &startContainerOperation{engine: e.engine, name: request.ContainerName} + if _, err := e.coordinator.ExecuteStep(ctx, transactionID, startIntent(request), startOperation); err != nil { + return err + } + healthOperation := &healthOperation{ + engine: e.engine, + checker: e.checker, + name: request.ContainerName, + endpoint: request.HealthEndpoint, + timeout: healthTimeout, + } + _, err := e.coordinator.ExecuteStep(ctx, transactionID, healthIntent(request), healthOperation) + return err +} + +func validateRequest(request Request) error { + if !filepath.IsAbs(request.ArchivePath) { + return errors.New("image archive path must be absolute") + } + if request.ImageReference == "" || strings.TrimSpace(request.ImageReference) != request.ImageReference { + return errors.New("exact image reference is required") + } + if _, err := opencontainersdigest.Parse(request.ExpectedImageDigest); err != nil { + return fmt.Errorf("invalid expected image digest: %w", err) + } + if request.Platform.OS == "" || request.Platform.Architecture == "" { + return errors.New("explicit image operating system and architecture are required") + } + if request.ContainerName == "" || strings.TrimSpace(request.ContainerName) != request.ContainerName { + return errors.New("exact container name is required") + } + if request.Port != 8080 && request.Port != 8081 { + return fmt.Errorf("backend container port must be 8080 or 8081: %d", request.Port) + } + if request.PortEnvironmentKey == "" || strings.Contains(request.PortEnvironmentKey, "=") || strings.TrimSpace(request.PortEnvironmentKey) != request.PortEnvironmentKey { + return errors.New("exact port environment key is required") + } + if !filepath.IsAbs(request.ConfigSource) || !filepath.IsAbs(request.ConfigTarget) { + return errors.New("backend configuration source and target must be absolute paths") + } + if err := validateRestartPolicy(request.RestartPolicy); err != nil { + return err + } + parsed, err := url.ParseRequestURI(request.HealthEndpoint) + if err != nil || parsed.Scheme != "http" || parsed.Host == "" || parsed.Path != healthPath { + return fmt.Errorf("health endpoint must be an HTTP URL with exact path %s", healthPath) + } + if parsed.Port() != strconv.Itoa(request.Port) { + return fmt.Errorf("health endpoint port must equal backend container port %d", request.Port) + } + return nil +} + +func validateRestartPolicy(policy containerengine.RestartPolicy) error { + switch policy.Name { + case "no", "always", "unless-stopped": + if policy.MaximumRetryCount != 0 { + return fmt.Errorf("restart policy %s does not accept a maximum retry count", policy.Name) + } + case "on-failure": + if policy.MaximumRetryCount < 0 { + return errors.New("on-failure maximum retry count cannot be negative") + } + default: + return fmt.Errorf("unsupported explicit restart policy: %q", policy.Name) + } + return nil +} + +func regularFile(path, description string) error { + info, err := os.Stat(path) + if err != nil { + return fmt.Errorf("inspect %s %s: %w", description, path, err) + } + if !info.Mode().IsRegular() { + return fmt.Errorf("%s is not a regular file: %s", description, path) + } + return nil +} + +func containerSpec(request Request) containerengine.ContainerSpec { + return containerengine.ContainerSpec{ + Name: request.ContainerName, + ImageReference: request.ImageReference, + Platform: request.Platform, + Environment: []string{ + request.PortEnvironmentKey + "=" + strconv.Itoa(request.Port), + }, + NetworkMode: hostNetworkMode, + RestartPolicy: request.RestartPolicy, + Mounts: []containerengine.Mount{{ + Type: bindMountType, + Source: request.ConfigSource, + Target: request.ConfigTarget, + ReadOnly: true, + }}, + } +} + +func loadIntent(request Request) transaction.StepIntent { + return intent(stepLoadImage, "load and verify backend image", struct { + ArchivePath string `json:"archivePath"` + ImageReference string `json:"imageReference"` + ImageDigest string `json:"imageDigest"` + Platform containerengine.Platform `json:"platform"` + }{request.ArchivePath, request.ImageReference, request.ExpectedImageDigest, request.Platform}) +} + +func createIntent(request Request, imageID string) transaction.StepIntent { + return intent(stepCreateContainer, "create inactive backend container", struct { + Spec containerengine.ContainerSpec `json:"spec"` + ImageID string `json:"imageId"` + }{containerSpec(request), imageID}) +} + +func startIntent(request Request) transaction.StepIntent { + return intent(stepStartContainer, "start inactive backend container", struct { + ContainerName string `json:"containerName"` + }{request.ContainerName}) +} + +func healthIntent(request Request) transaction.StepIntent { + return intent(stepCheckHealth, "wait for backend Actuator health", struct { + ContainerName string `json:"containerName"` + Endpoint string `json:"endpoint"` + Timeout time.Duration `json:"timeout"` + }{request.ContainerName, request.HealthEndpoint, healthTimeout}) +} + +func intent(key, name string, value any) transaction.StepIntent { + payload, err := json.Marshal(value) + if err != nil { + panic(fmt.Sprintf("marshal internal step intent: %v", err)) + } + return transaction.StepIntent{Key: key, Name: name, Intent: payload} +} + +func imageMatches(image containerengine.Image, expectedDigest string, expectedPlatform containerengine.Platform) (bool, error) { + expected, err := opencontainersdigest.Parse(expectedDigest) + if err != nil { + return false, fmt.Errorf("parse expected image digest: %w", err) + } + if image.Platform != expectedPlatform { + return false, nil + } + if image.DescriptorDigest != "" { + actual, err := opencontainersdigest.Parse(image.DescriptorDigest) + if err != nil { + return false, fmt.Errorf("parse inspected image descriptor digest: %w", err) + } + if actual == expected { + return true, nil + } + } + for _, repoDigest := range image.RepoDigests { + parsed, err := reference.ParseAnyReference(repoDigest) + if err != nil { + return false, fmt.Errorf("parse inspected repository digest %q: %w", repoDigest, err) + } + digested, ok := parsed.(reference.Digested) + if !ok { + return false, fmt.Errorf("inspected repository digest is not digest-qualified: %q", repoDigest) + } + if digested.Digest() == expected { + return true, nil + } + } + return false, nil +} + +func resultJSON(value any) json.RawMessage { + payload, err := json.Marshal(value) + if err != nil { + panic(fmt.Sprintf("marshal internal step result: %v", err)) + } + return payload +} diff --git a/internal/backendexecutor/executor_test.go b/internal/backendexecutor/executor_test.go new file mode 100644 index 0000000..070a67a --- /dev/null +++ b/internal/backendexecutor/executor_test.go @@ -0,0 +1,437 @@ +package backendexecutor + +import ( + "bytes" + "context" + "errors" + "io" + "log/slog" + "net/http" + "os" + "path/filepath" + "slices" + "sync" + "testing" + + "yms-daemon/internal/containerengine" + "yms-daemon/internal/transaction" +) + +const ( + testDigest = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + testRepository = "harbor.ymswell.asia/ymswell/glory-ymswell" + healthyResponse = `{"status":"UP","components":{"db":{"status":"UP"},"diskSpace":{"status":"UP"},"ping":{"status":"UP"},"redis":{"status":"UP"}}}` +) + +func TestExecutorPreservesOpaqueImageTagsAndReachesSwitching(t *testing.T) { + tags := []string{ + "20260814-093609-d7ed70f0-v1.1.8.1", + "20260814-093609-d7ed70f0", + } + for _, tag := range tags { + tag := tag + t.Run(tag, func(t *testing.T) { + ctx := context.Background() + store, coordinator := testTransactionKernel(t) + imageReference := testRepository + ":" + tag + request := testRequest(t, imageReference) + engine := newFakeEngine(request) + executor := testExecutor(t, store, coordinator, engine, healthyResponse) + record := createTransaction(t, store, "opaque-"+tag) + + if err := executor.Run(ctx, record.ID, request); err != nil { + t.Fatalf("run backend executor: %v", err) + } + current, err := store.Transaction(ctx, record.ID) + if err != nil { + t.Fatalf("read completed preparation: %v", err) + } + if current.State != transaction.StateSwitching { + t.Fatalf("unexpected transaction state: %s", current.State) + } + pending, err := store.PendingSteps(ctx, record.ID) + if err != nil || len(pending) != 0 { + t.Fatalf("unexpected pending steps: steps=%+v err=%v", pending, err) + } + + engine.mu.Lock() + if engine.lastCreateSpec.ImageReference != imageReference { + engine.mu.Unlock() + t.Fatalf("image reference changed: got %q want %q", engine.lastCreateSpec.ImageReference, imageReference) + } + if !slices.Contains(engine.lastCreateSpec.Environment, "SERVER_PORT=8081") { + engine.mu.Unlock() + t.Fatalf("missing explicit port environment: %+v", engine.lastCreateSpec.Environment) + } + initialCalls := engine.callCounts() + engine.mu.Unlock() + + if err := executor.Run(ctx, record.ID, request); err != nil { + t.Fatalf("repeat executor at switching state: %v", err) + } + engine.mu.Lock() + repeatedCalls := engine.callCounts() + engine.mu.Unlock() + if repeatedCalls != initialCalls { + t.Fatalf("switching state repeated engine calls: before=%+v after=%+v", initialCalls, repeatedCalls) + } + }) + } +} + +func TestExecutorRecoversRecordedCreateIntentWithoutRepeatingCreate(t *testing.T) { + ctx := context.Background() + store, coordinator := testTransactionKernel(t) + request := testRequest(t, testRepository+":20260814-093609-d7ed70f0-v1.1.8.1") + engine := newFakeEngine(request) + engine.imageAvailable = true + engine.containers[request.ContainerName] = engine.containerFromSpec(containerSpec(request), false) + executor := testExecutor(t, store, coordinator, engine, healthyResponse) + record := createTransaction(t, store, "recover-create") + transitionToPrepared(t, store, record.ID) + + loadStep := loadIntent(request) + if _, _, err := store.RecordStepIntent(ctx, record.ID, loadStep); err != nil { + t.Fatalf("record completed image intent: %v", err) + } + if _, err := store.CompleteStep(ctx, record.ID, loadStep.Key, transaction.StepSucceeded, []byte(`{"loaded":true}`), ""); err != nil { + t.Fatalf("complete image step: %v", err) + } + createStep := createIntent(request, engine.loadedImage.ID) + if _, _, err := store.RecordStepIntent(ctx, record.ID, createStep); err != nil { + t.Fatalf("record crash-window create intent: %v", err) + } + + if err := executor.Run(ctx, record.ID, request); err != nil { + t.Fatalf("recover backend executor: %v", err) + } + engine.mu.Lock() + calls := engine.callCounts() + engine.mu.Unlock() + if calls.load != 0 || calls.create != 0 || calls.start != 1 { + t.Fatalf("unexpected recovery calls: %+v", calls) + } + current, err := store.Transaction(ctx, record.ID) + if err != nil || current.State != transaction.StateSwitching { + t.Fatalf("unexpected recovered state: record=%+v err=%v", current, err) + } +} + +func TestExecutorMarksValidationFailureTerminal(t *testing.T) { + ctx := context.Background() + store, coordinator := testTransactionKernel(t) + request := testRequest(t, testRepository+":20260814-093609-d7ed70f0") + request.Port = 9090 + engine := newFakeEngine(request) + executor := testExecutor(t, store, coordinator, engine, healthyResponse) + record := createTransaction(t, store, "invalid-port") + + if err := executor.Run(ctx, record.ID, request); err == nil { + t.Fatal("expected validation failure") + } + current, err := store.Transaction(ctx, record.ID) + if err != nil || current.State != transaction.StateFailed { + t.Fatalf("unexpected failed transaction: record=%+v err=%v", current, err) + } +} + +func TestExecutorKeepsPreparedStateForConflictingContainerInspection(t *testing.T) { + ctx := context.Background() + store, coordinator := testTransactionKernel(t) + request := testRequest(t, testRepository+":20260814-093609-d7ed70f0") + engine := newFakeEngine(request) + engine.imageAvailable = true + conflicting := engine.containerFromSpec(containerSpec(request), false) + conflicting.ImageID = "sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" + engine.containers[request.ContainerName] = conflicting + executor := testExecutor(t, store, coordinator, engine, healthyResponse) + record := createTransaction(t, store, "container-conflict") + + err := executor.Run(ctx, record.ID, request) + var uncertain *transaction.UncertainStepError + if !errors.As(err, &uncertain) { + t.Fatalf("expected uncertain container step, got %v", err) + } + current, readErr := store.Transaction(ctx, record.ID) + if readErr != nil || current.State != transaction.StatePrepared { + t.Fatalf("conflict did not preserve prepared state: record=%+v err=%v", current, readErr) + } +} + +func TestExecutorRejectsChangedRecoveryRequestWithoutChangingState(t *testing.T) { + ctx := context.Background() + store, coordinator := testTransactionKernel(t) + request := testRequest(t, testRepository+":20260814-093609-d7ed70f0-v1.1.8.1") + engine := newFakeEngine(request) + engine.imageAvailable = true + executor := testExecutor(t, store, coordinator, engine, healthyResponse) + record := createTransaction(t, store, "changed-recovery-request") + transitionToPrepared(t, store, record.ID) + if err := executor.prepare(ctx, record.ID, request); err != nil { + t.Fatalf("prepare original request: %v", err) + } + if _, err := store.Transition(ctx, record.ID, transaction.StateStarting, "test starting"); err != nil { + t.Fatalf("transition to starting: %v", err) + } + + changed := request + changed.ImageReference = testRepository + ":20260814-093609-d7ed70f0" + err := executor.Run(ctx, record.ID, changed) + if !errors.Is(err, transaction.ErrStepConflict) { + t.Fatalf("expected persisted intent conflict, got %v", err) + } + current, readErr := store.Transaction(ctx, record.ID) + if readErr != nil || current.State != transaction.StateStarting { + t.Fatalf("changed recovery request altered transaction: record=%+v err=%v", current, readErr) + } +} + +func TestImageMatchesRequiresDigestEvidenceAndExactPlatform(t *testing.T) { + t.Parallel() + platform := containerengine.Platform{OS: "linux", Architecture: "arm64"} + image := containerengine.Image{ + ID: "sha256:cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc", + RepoDigests: []string{testRepository + "@" + testDigest}, + Platform: platform, + } + matched, err := imageMatches(image, testDigest, platform) + if err != nil || !matched { + t.Fatalf("expected exact image match: matched=%t err=%v", matched, err) + } + image.RepoDigests = nil + matched, err = imageMatches(image, testDigest, platform) + if err != nil || matched { + t.Fatalf("image ID alone must not satisfy manifest digest: matched=%t err=%v", matched, err) + } + image.DescriptorDigest = testDigest + matched, err = imageMatches(image, testDigest, containerengine.Platform{OS: "linux", Architecture: "amd64"}) + if err != nil || matched { + t.Fatalf("platform mismatch accepted: matched=%t err=%v", matched, err) + } +} + +func testRequest(t *testing.T, imageReference string) Request { + t.Helper() + directory := t.TempDir() + archivePath := filepath.Join(directory, "backend-image.tar") + configPath := filepath.Join(directory, "yms.yaml") + if err := os.WriteFile(archivePath, []byte("image archive"), 0o600); err != nil { + t.Fatalf("write image archive: %v", err) + } + if err := os.WriteFile(configPath, []byte("server: {}\n"), 0o600); err != nil { + t.Fatalf("write backend configuration: %v", err) + } + return Request{ + ArchivePath: archivePath, + ImageReference: imageReference, + ExpectedImageDigest: testDigest, + Platform: containerengine.Platform{OS: "linux", Architecture: "amd64"}, + ContainerName: "explicit-backend-8081", + Port: 8081, + PortEnvironmentKey: "SERVER_PORT", + ConfigSource: configPath, + ConfigTarget: "/app/config/application.yaml", + RestartPolicy: containerengine.RestartPolicy{Name: "unless-stopped"}, + HealthEndpoint: "http://127.0.0.1:8081/yms/actuator/health", + } +} + +func testTransactionKernel(t *testing.T) (*transaction.Store, *transaction.Coordinator) { + t.Helper() + store, err := transaction.OpenStore(context.Background(), filepath.Join(t.TempDir(), "transactions.db")) + if err != nil { + t.Fatalf("open transaction store: %v", err) + } + t.Cleanup(func() { _ = store.Close() }) + coordinator, err := transaction.NewCoordinator(store, slog.New(slog.NewTextHandler(io.Discard, nil))) + if err != nil { + t.Fatalf("create transaction coordinator: %v", err) + } + return store, coordinator +} + +func testExecutor(t *testing.T, store *transaction.Store, coordinator *transaction.Coordinator, engine containerengine.Engine, body string) *Executor { + t.Helper() + client := &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) { + return &http.Response{StatusCode: http.StatusOK, Header: make(http.Header), Body: io.NopCloser(bytes.NewBufferString(body))}, nil + })} + executor, err := New(store, coordinator, engine, client) + if err != nil { + t.Fatalf("create backend executor: %v", err) + } + return executor +} + +func createTransaction(t *testing.T, store *transaction.Store, suffix string) transaction.Transaction { + t.Helper() + record, _, err := store.CreateTransaction(context.Background(), transaction.CreateRequest{ + ID: "backend-" + suffix, + IdempotencyKey: "backend-request-" + suffix, + Source: "test", + Service: "backend", + }) + if err != nil { + t.Fatalf("create backend transaction: %v", err) + } + return record +} + +func transitionToPrepared(t *testing.T, store *transaction.Store, transactionID string) { + t.Helper() + ctx := context.Background() + if _, err := store.Transition(ctx, transactionID, transaction.StateValidating, "test validating"); err != nil { + t.Fatalf("transition to validating: %v", err) + } + if _, err := store.Transition(ctx, transactionID, transaction.StatePrepared, "test prepared"); err != nil { + t.Fatalf("transition to prepared: %v", err) + } +} + +type roundTripFunc func(*http.Request) (*http.Response, error) + +func (f roundTripFunc) RoundTrip(request *http.Request) (*http.Response, error) { + return f(request) +} + +type engineCalls struct { + ping int + load int + inspect int + create int + start int + remove int +} + +type fakeEngine struct { + mu sync.Mutex + request Request + loadedImage containerengine.Image + imageAvailable bool + containers map[string]containerengine.Container + lastCreateSpec containerengine.ContainerSpec + pingCalls int + loadCalls int + inspectCalls int + createCalls int + startCalls int + removeCalls int +} + +func newFakeEngine(request Request) *fakeEngine { + return &fakeEngine{ + request: request, + loadedImage: containerengine.Image{ + ID: "sha256:cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc", + RepoDigests: []string{testRepository + "@" + testDigest}, + Platform: request.Platform, + }, + containers: make(map[string]containerengine.Container), + } +} + +func (e *fakeEngine) Ping(context.Context) error { + e.mu.Lock() + defer e.mu.Unlock() + e.pingCalls++ + return nil +} + +func (e *fakeEngine) LoadImage(_ context.Context, input io.Reader) error { + e.mu.Lock() + defer e.mu.Unlock() + e.loadCalls++ + if _, err := io.ReadAll(input); err != nil { + return err + } + e.imageAvailable = true + return nil +} + +func (e *fakeEngine) InspectImage(context.Context, string) (containerengine.Image, error) { + e.mu.Lock() + defer e.mu.Unlock() + e.inspectCalls++ + if !e.imageAvailable { + return containerengine.Image{}, containerengine.ErrNotFound + } + return e.loadedImage, nil +} + +func (e *fakeEngine) CreateContainer(_ context.Context, spec containerengine.ContainerSpec) (containerengine.Container, error) { + e.mu.Lock() + defer e.mu.Unlock() + e.createCalls++ + e.lastCreateSpec = spec + if _, exists := e.containers[spec.Name]; exists { + return containerengine.Container{}, errors.New("container name already exists") + } + record := e.containerFromSpec(spec, false) + e.containers[spec.Name] = record + return record, nil +} + +func (e *fakeEngine) StartContainer(_ context.Context, name string) error { + e.mu.Lock() + defer e.mu.Unlock() + e.startCalls++ + record, exists := e.containers[name] + if !exists { + return containerengine.ErrNotFound + } + record.Running = true + record.Status = "running" + e.containers[name] = record + return nil +} + +func (e *fakeEngine) InspectContainer(_ context.Context, name string) (containerengine.Container, error) { + e.mu.Lock() + defer e.mu.Unlock() + record, exists := e.containers[name] + if !exists { + return containerengine.Container{}, containerengine.ErrNotFound + } + return record, nil +} + +func (e *fakeEngine) RemoveContainer(_ context.Context, name string, _ bool) error { + e.mu.Lock() + defer e.mu.Unlock() + e.removeCalls++ + if _, exists := e.containers[name]; !exists { + return containerengine.ErrNotFound + } + delete(e.containers, name) + return nil +} + +func (e *fakeEngine) Close() error { return nil } + +func (e *fakeEngine) containerFromSpec(spec containerengine.ContainerSpec, running bool) containerengine.Container { + return containerengine.Container{ + ID: "container-id-" + spec.Name, + Name: spec.Name, + ImageID: e.loadedImage.ID, + ImageReference: spec.ImageReference, + Platform: spec.Platform.OS + "/" + spec.Platform.Architecture, + Running: running, + Status: "created", + Environment: append([]string(nil), spec.Environment...), + NetworkMode: spec.NetworkMode, + RestartPolicy: spec.RestartPolicy, + Mounts: append([]containerengine.Mount(nil), spec.Mounts...), + } +} + +func (e *fakeEngine) callCounts() engineCalls { + return engineCalls{ + ping: e.pingCalls, + load: e.loadCalls, + inspect: e.inspectCalls, + create: e.createCalls, + start: e.startCalls, + remove: e.removeCalls, + } +} + +var _ containerengine.Engine = (*fakeEngine)(nil) diff --git a/internal/backendexecutor/operations.go b/internal/backendexecutor/operations.go new file mode 100644 index 0000000..6265029 --- /dev/null +++ b/internal/backendexecutor/operations.go @@ -0,0 +1,200 @@ +package backendexecutor + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "os" + "slices" + "sync" + "time" + + "yms-daemon/internal/containerengine" + "yms-daemon/internal/healthcheck" + "yms-daemon/internal/transaction" +) + +type loadImageOperation struct { + engine containerengine.Engine + archivePath string + imageReference string + expectedDigest string + platform containerengine.Platform +} + +func (o *loadImageOperation) Apply(ctx context.Context) error { + archive, err := os.Open(o.archivePath) + if err != nil { + return fmt.Errorf("open image archive: %w", err) + } + defer archive.Close() + return o.engine.LoadImage(ctx, archive) +} + +func (o *loadImageOperation) Inspect(ctx context.Context) (transaction.Inspection, error) { + image, err := o.engine.InspectImage(ctx, o.imageReference) + if errors.Is(err, containerengine.ErrNotFound) { + return transaction.Inspection{Status: transaction.InspectionNotApplied}, nil + } + if err != nil { + return transaction.Inspection{}, err + } + matches, err := imageMatches(image, o.expectedDigest, o.platform) + if err != nil { + return transaction.Inspection{}, err + } + result := resultJSON(struct { + ImageID string `json:"imageId"` + RepoDigests []string `json:"repoDigests"` + DescriptorDigest string `json:"descriptorDigest"` + Platform containerengine.Platform `json:"platform"` + }{image.ID, image.RepoDigests, image.DescriptorDigest, image.Platform}) + if !matches { + return transaction.Inspection{Status: transaction.InspectionNotApplied, Result: result}, nil + } + return transaction.Inspection{Status: transaction.InspectionApplied, Result: result}, nil +} + +type createContainerOperation struct { + engine containerengine.Engine + expectedImage containerengine.Image + spec containerengine.ContainerSpec +} + +func (o *createContainerOperation) Apply(ctx context.Context) error { + _, err := o.engine.CreateContainer(ctx, o.spec) + return err +} + +func (o *createContainerOperation) Inspect(ctx context.Context) (transaction.Inspection, error) { + record, err := o.engine.InspectContainer(ctx, o.spec.Name) + if errors.Is(err, containerengine.ErrNotFound) { + return transaction.Inspection{Status: transaction.InspectionNotApplied}, nil + } + if err != nil { + return transaction.Inspection{}, err + } + result := containerResult(record) + if !containerMatches(record, o.expectedImage.ID, o.spec) { + return transaction.Inspection{Status: transaction.InspectionUnknown, Result: result}, nil + } + return transaction.Inspection{Status: transaction.InspectionApplied, Result: result}, nil +} + +type startContainerOperation struct { + engine containerengine.Engine + name string +} + +func (o *startContainerOperation) Apply(ctx context.Context) error { + return o.engine.StartContainer(ctx, o.name) +} + +func (o *startContainerOperation) Inspect(ctx context.Context) (transaction.Inspection, error) { + record, err := o.engine.InspectContainer(ctx, o.name) + if errors.Is(err, containerengine.ErrNotFound) { + return transaction.Inspection{Status: transaction.InspectionNotApplied}, nil + } + if err != nil { + return transaction.Inspection{}, err + } + result := containerResult(record) + if record.Running && !record.Dead { + return transaction.Inspection{Status: transaction.InspectionApplied, Result: result}, nil + } + return transaction.Inspection{Status: transaction.InspectionNotApplied, Result: result}, nil +} + +type healthOperation struct { + engine containerengine.Engine + checker *healthcheck.ActuatorChecker + name string + endpoint string + timeout time.Duration + + mu sync.Mutex + confirmedReport healthcheck.ActuatorReport + confirmed bool +} + +func (o *healthOperation) Apply(ctx context.Context) error { + report, err := o.checker.Wait(ctx, o.endpoint, o.timeout, o.running) + if err != nil { + return err + } + o.mu.Lock() + o.confirmedReport = report + o.confirmed = true + o.mu.Unlock() + return nil +} + +func (o *healthOperation) Inspect(ctx context.Context) (transaction.Inspection, error) { + o.mu.Lock() + if o.confirmed { + report := o.confirmedReport + o.mu.Unlock() + return healthInspection(report, true, nil) + } + o.mu.Unlock() + report, ready, err := o.checker.Check(ctx, o.endpoint, o.running) + return healthInspection(report, ready, err) +} + +func (o *healthOperation) running(ctx context.Context) (bool, error) { + record, err := o.engine.InspectContainer(ctx, o.name) + if errors.Is(err, containerengine.ErrNotFound) { + return false, nil + } + if err != nil { + return false, err + } + return record.Running && !record.Dead, nil +} + +func healthInspection(report healthcheck.ActuatorReport, ready bool, err error) (transaction.Inspection, error) { + result := resultJSON(report) + if errors.Is(err, healthcheck.ErrWorkloadStopped) { + return transaction.Inspection{Status: transaction.InspectionNotApplied, Result: result}, nil + } + if err != nil { + return transaction.Inspection{Status: transaction.InspectionNotApplied, Result: result}, nil + } + if !ready { + return transaction.Inspection{Status: transaction.InspectionNotApplied, Result: result}, nil + } + return transaction.Inspection{Status: transaction.InspectionApplied, Result: result}, nil +} + +func containerMatches(record containerengine.Container, expectedImageID string, spec containerengine.ContainerSpec) bool { + if record.ImageID != expectedImageID || record.NetworkMode != spec.NetworkMode || record.RestartPolicy != spec.RestartPolicy { + return false + } + for _, expected := range spec.Environment { + if !slices.Contains(record.Environment, expected) { + return false + } + } + for _, expected := range spec.Mounts { + if !slices.Contains(record.Mounts, expected) { + return false + } + } + return true +} + +func containerResult(record containerengine.Container) json.RawMessage { + return resultJSON(struct { + ID string `json:"id"` + ImageID string `json:"imageId"` + Running bool `json:"running"` + Dead bool `json:"dead"` + Status string `json:"status"` + }{record.ID, record.ImageID, record.Running, record.Dead, record.Status}) +} + +var _ transaction.Operation = (*loadImageOperation)(nil) +var _ transaction.Operation = (*createContainerOperation)(nil) +var _ transaction.Operation = (*startContainerOperation)(nil) +var _ transaction.Operation = (*healthOperation)(nil) diff --git a/internal/containerengine/engine.go b/internal/containerengine/engine.go new file mode 100644 index 0000000..0be70bf --- /dev/null +++ b/internal/containerengine/engine.go @@ -0,0 +1,80 @@ +// Package containerengine defines the container runtime boundary used by update executors. +package containerengine + +import ( + "context" + "errors" + "io" +) + +var ErrNotFound = errors.New("container engine object not found") + +// Platform is an explicit OCI operating system and CPU platform. +type Platform struct { + OS string + Architecture string + Variant string +} + +// Image is the immutable image information returned by the engine. +type Image struct { + ID string + RepoDigests []string + DescriptorDigest string + Platform Platform +} + +// RestartPolicy is passed to the engine without an implicit default. +type RestartPolicy struct { + Name string + MaximumRetryCount int +} + +// Mount is one explicit container mount. +type Mount struct { + Type string + Source string + Target string + ReadOnly bool +} + +// ContainerSpec contains every property controlled by the backend executor. +type ContainerSpec struct { + Name string + ImageReference string + Platform Platform + Environment []string + Labels map[string]string + NetworkMode string + RestartPolicy RestartPolicy + Mounts []Mount +} + +// Container is the runtime state required for idempotent inspection. +type Container struct { + ID string + Name string + ImageID string + ImageReference string + Platform string + Running bool + Dead bool + Status string + Environment []string + Labels map[string]string + NetworkMode string + RestartPolicy RestartPolicy + Mounts []Mount +} + +// Engine is the smallest container runtime API required by an update executor. +type Engine interface { + Ping(context.Context) error + LoadImage(context.Context, io.Reader) error + InspectImage(context.Context, string) (Image, error) + CreateContainer(context.Context, ContainerSpec) (Container, error) + StartContainer(context.Context, string) error + InspectContainer(context.Context, string) (Container, error) + RemoveContainer(context.Context, string, bool) error + Close() error +} diff --git a/internal/containerengine/moby.go b/internal/containerengine/moby.go new file mode 100644 index 0000000..4f562ef --- /dev/null +++ b/internal/containerengine/moby.go @@ -0,0 +1,206 @@ +package containerengine + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "io" + + cerrdefs "github.com/containerd/errdefs" + "github.com/moby/moby/api/types/container" + "github.com/moby/moby/api/types/jsonstream" + "github.com/moby/moby/api/types/mount" + "github.com/moby/moby/client" + ocispec "github.com/opencontainers/image-spec/specs-go/v1" +) + +// MobyEngine adapts the official Docker Engine Go client. +type MobyEngine struct { + client *client.Client +} + +// NewMobyEngine creates a client from Docker's documented environment variables. +// API negotiation remains enabled, including when DOCKER_HOST selects a non-default socket. +func NewMobyEngine() (*MobyEngine, error) { + apiClient, err := client.New(client.FromEnv) + if err != nil { + return nil, fmt.Errorf("create Docker Engine client: %w", err) + } + return &MobyEngine{client: apiClient}, nil +} + +func (e *MobyEngine) Ping(ctx context.Context) error { + if _, err := e.client.Ping(ctx, client.PingOptions{NegotiateAPIVersion: true}); err != nil { + return fmt.Errorf("ping Docker Engine: %w", err) + } + return nil +} + +func (e *MobyEngine) LoadImage(ctx context.Context, input io.Reader) error { + if input == nil { + return errors.New("image archive reader is required") + } + response, err := e.client.ImageLoad(ctx, input) + if err != nil { + return fmt.Errorf("load image archive: %w", err) + } + defer response.Close() + if err := decodeImageLoadResponse(response); err != nil { + return fmt.Errorf("load image archive response: %w", err) + } + return nil +} + +func decodeImageLoadResponse(input io.Reader) error { + decoder := json.NewDecoder(input) + for { + var message jsonstream.Message + if err := decoder.Decode(&message); err != nil { + if errors.Is(err, io.EOF) { + return nil + } + return fmt.Errorf("decode Docker JSON stream: %w", err) + } + if message.Error != nil { + return message.Error + } + } +} + +func (e *MobyEngine) InspectImage(ctx context.Context, reference string) (Image, error) { + response, err := e.client.ImageInspect(ctx, reference) + if err != nil { + return Image{}, engineError("inspect image", err) + } + descriptorDigest := "" + if response.Descriptor != nil { + descriptorDigest = response.Descriptor.Digest.String() + } + return Image{ + ID: response.ID, + RepoDigests: append([]string(nil), response.RepoDigests...), + DescriptorDigest: descriptorDigest, + Platform: Platform{ + OS: response.Os, + Architecture: response.Architecture, + Variant: response.Variant, + }, + }, nil +} + +func (e *MobyEngine) CreateContainer(ctx context.Context, spec ContainerSpec) (Container, error) { + apiMounts := make([]mount.Mount, 0, len(spec.Mounts)) + for _, item := range spec.Mounts { + apiMounts = append(apiMounts, mount.Mount{ + Type: mount.Type(item.Type), + Source: item.Source, + Target: item.Target, + ReadOnly: item.ReadOnly, + }) + } + result, err := e.client.ContainerCreate(ctx, client.ContainerCreateOptions{ + Config: &container.Config{ + Env: append([]string(nil), spec.Environment...), + Labels: cloneMap(spec.Labels), + }, + HostConfig: &container.HostConfig{ + NetworkMode: container.NetworkMode(spec.NetworkMode), + RestartPolicy: container.RestartPolicy{ + Name: container.RestartPolicyMode(spec.RestartPolicy.Name), + MaximumRetryCount: spec.RestartPolicy.MaximumRetryCount, + }, + Mounts: apiMounts, + }, + Platform: &ocispec.Platform{ + OS: spec.Platform.OS, + Architecture: spec.Platform.Architecture, + Variant: spec.Platform.Variant, + }, + Name: spec.Name, + Image: spec.ImageReference, + }) + if err != nil { + return Container{}, engineError("create container", err) + } + return e.InspectContainer(ctx, result.ID) +} + +func (e *MobyEngine) StartContainer(ctx context.Context, idOrName string) error { + if _, err := e.client.ContainerStart(ctx, idOrName, client.ContainerStartOptions{}); err != nil { + return engineError("start container", err) + } + return nil +} + +func (e *MobyEngine) InspectContainer(ctx context.Context, idOrName string) (Container, error) { + result, err := e.client.ContainerInspect(ctx, idOrName, client.ContainerInspectOptions{}) + if err != nil { + return Container{}, engineError("inspect container", err) + } + response := result.Container + record := Container{ + ID: response.ID, + Name: response.Name, + ImageID: response.Image, + Platform: response.Platform, + } + if response.State != nil { + record.Running = response.State.Running + record.Dead = response.State.Dead + record.Status = string(response.State.Status) + } + if response.Config != nil { + record.ImageReference = response.Config.Image + record.Environment = append([]string(nil), response.Config.Env...) + record.Labels = cloneMap(response.Config.Labels) + } + if response.HostConfig != nil { + record.NetworkMode = string(response.HostConfig.NetworkMode) + record.RestartPolicy = RestartPolicy{ + Name: string(response.HostConfig.RestartPolicy.Name), + MaximumRetryCount: response.HostConfig.RestartPolicy.MaximumRetryCount, + } + } + for _, item := range response.Mounts { + record.Mounts = append(record.Mounts, Mount{ + Type: string(item.Type), + Source: item.Source, + Target: item.Destination, + ReadOnly: !item.RW, + }) + } + return record, nil +} + +func (e *MobyEngine) RemoveContainer(ctx context.Context, idOrName string, force bool) error { + _, err := e.client.ContainerRemove(ctx, idOrName, client.ContainerRemoveOptions{Force: force}) + if err != nil { + return engineError("remove container", err) + } + return nil +} + +func (e *MobyEngine) Close() error { + return e.client.Close() +} + +func engineError(action string, err error) error { + if cerrdefs.IsNotFound(err) { + return fmt.Errorf("%s: %w: %v", action, ErrNotFound, err) + } + return fmt.Errorf("%s: %w", action, err) +} + +func cloneMap(source map[string]string) map[string]string { + if source == nil { + return nil + } + result := make(map[string]string, len(source)) + for key, value := range source { + result[key] = value + } + return result +} + +var _ Engine = (*MobyEngine)(nil) diff --git a/internal/containerengine/moby_test.go b/internal/containerengine/moby_test.go new file mode 100644 index 0000000..864bf81 --- /dev/null +++ b/internal/containerengine/moby_test.go @@ -0,0 +1,31 @@ +package containerengine + +import ( + "strings" + "testing" +) + +func TestDecodeImageLoadResponseConsumesCompleteSuccessStream(t *testing.T) { + t.Parallel() + input := strings.NewReader("{\"stream\":\"Loaded image: repository:20260814-093609-d7ed70f0-v1.1.8.1\\n\"}\n" + + "{\"stream\":\"Loaded image: repository:20260814-093609-d7ed70f0\\n\"}\n") + if err := decodeImageLoadResponse(input); err != nil { + t.Fatalf("decode successful load response: %v", err) + } +} + +func TestDecodeImageLoadResponseReturnsStreamError(t *testing.T) { + t.Parallel() + err := decodeImageLoadResponse(strings.NewReader(`{"errorDetail":{"code":500,"message":"load failed"}}`)) + if err == nil || err.Error() != "load failed" { + t.Fatalf("unexpected load error: %v", err) + } +} + +func TestDecodeImageLoadResponseRejectsMalformedJSON(t *testing.T) { + t.Parallel() + err := decodeImageLoadResponse(strings.NewReader(`{"stream":`)) + if err == nil { + t.Fatalf("expected malformed response error, got %v", err) + } +} diff --git a/internal/healthcheck/actuator.go b/internal/healthcheck/actuator.go new file mode 100644 index 0000000..409dd85 --- /dev/null +++ b/internal/healthcheck/actuator.go @@ -0,0 +1,150 @@ +package healthcheck + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/url" + "time" +) + +const maximumResponseBytes = 1 << 20 + +var ErrWorkloadStopped = errors.New("workload stopped before becoming healthy") + +// RunningProbe 在每次健康采样前核对容器是否仍处于运行状态。 +type RunningProbe func(context.Context) (bool, error) + +// ActuatorReport 是一次成功健康采样的非敏感摘要。 +type ActuatorReport struct { + Status string + Components map[string]string +} + +// ActuatorChecker 等待 Spring Boot Actuator 顶层状态进入 UP。 +// 它只采样健康状态,不负责重启容器。 +type ActuatorChecker struct { + client *http.Client + interval time.Duration +} + +// Check 执行一次健康采样,供事务恢复时核对已经记录意图的健康步骤。 +func (c *ActuatorChecker) Check(ctx context.Context, endpoint string, running RunningProbe) (ActuatorReport, bool, error) { + if running == nil { + return ActuatorReport{}, false, errors.New("running probe is required") + } + if err := validateEndpoint(endpoint); err != nil { + return ActuatorReport{}, false, err + } + return c.sample(ctx, endpoint, running) +} + +// NewActuatorChecker 创建固定间隔的 Actuator 检查器。 +func NewActuatorChecker(client *http.Client, interval time.Duration) (*ActuatorChecker, error) { + if client == nil { + return nil, errors.New("HTTP client is required") + } + if interval <= 0 { + return nil, errors.New("health check interval must be positive") + } + return &ActuatorChecker{client: client, interval: interval}, nil +} + +// Wait 在 timeout 范围内持续采样。容器停止时立即失败;其他未就绪结果保留到截止时间。 +func (c *ActuatorChecker) Wait(ctx context.Context, endpoint string, timeout time.Duration, running RunningProbe) (ActuatorReport, error) { + if timeout <= 0 { + return ActuatorReport{}, errors.New("health check timeout must be positive") + } + if running == nil { + return ActuatorReport{}, errors.New("running probe is required") + } + if err := validateEndpoint(endpoint); err != nil { + return ActuatorReport{}, err + } + + waitContext, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + var lastErr error + for { + report, ready, err := c.sample(waitContext, endpoint, running) + if ready { + return report, nil + } + if errors.Is(err, ErrWorkloadStopped) { + return ActuatorReport{}, err + } + lastErr = err + + timer := time.NewTimer(c.interval) + select { + case <-waitContext.Done(): + timer.Stop() + return ActuatorReport{}, fmt.Errorf("Actuator did not become UP within %s: %w", timeout, errors.Join(waitContext.Err(), lastErr)) + case <-timer.C: + } + } +} + +func validateEndpoint(endpoint string) error { + parsed, err := url.ParseRequestURI(endpoint) + if err != nil || parsed.Scheme != "http" || parsed.Host == "" { + return fmt.Errorf("invalid Actuator HTTP endpoint: %q", endpoint) + } + return nil +} + +func (c *ActuatorChecker) sample(ctx context.Context, endpoint string, running RunningProbe) (ActuatorReport, bool, error) { + isRunning, err := running(ctx) + if err != nil { + return ActuatorReport{}, false, fmt.Errorf("inspect workload before health check: %w", err) + } + if !isRunning { + return ActuatorReport{}, false, ErrWorkloadStopped + } + + request, err := http.NewRequestWithContext(ctx, http.MethodGet, endpoint, nil) + if err != nil { + return ActuatorReport{}, false, fmt.Errorf("create Actuator request: %w", err) + } + request.Header.Set("Accept", "application/json") + response, err := c.client.Do(request) + if err != nil { + return ActuatorReport{}, false, fmt.Errorf("request Actuator health: %w", err) + } + defer response.Body.Close() + body, err := io.ReadAll(io.LimitReader(response.Body, maximumResponseBytes+1)) + if err != nil { + return ActuatorReport{}, false, fmt.Errorf("read Actuator response: %w", err) + } + if len(body) > maximumResponseBytes { + return ActuatorReport{}, false, errors.New("Actuator response exceeds size limit") + } + if response.StatusCode != http.StatusOK { + return ActuatorReport{}, false, fmt.Errorf("Actuator returned HTTP status %d", response.StatusCode) + } + + var payload actuatorPayload + if err := json.Unmarshal(body, &payload); err != nil { + return ActuatorReport{}, false, fmt.Errorf("decode Actuator response: %w", err) + } + report := ActuatorReport{Status: payload.Status, Components: make(map[string]string, len(payload.Components))} + for name, component := range payload.Components { + report.Components[name] = component.Status + } + if payload.Status != "UP" { + return report, false, fmt.Errorf("Actuator status is %q", payload.Status) + } + return report, true, nil +} + +type actuatorPayload struct { + Status string `json:"status"` + Components map[string]actuatorComponent `json:"components"` +} + +type actuatorComponent struct { + Status string `json:"status"` +} diff --git a/internal/healthcheck/actuator_test.go b/internal/healthcheck/actuator_test.go new file mode 100644 index 0000000..c943c8c --- /dev/null +++ b/internal/healthcheck/actuator_test.go @@ -0,0 +1,117 @@ +package healthcheck + +import ( + "bytes" + "context" + "errors" + "io" + "net/http" + "sync/atomic" + "testing" + "time" +) + +const healthyActuatorSample = `{"status":"UP","components":{"db":{"status":"UP","components":{"dorisDataSource":{"status":"UP"},"postgresqlDataSource":{"status":"UP"}}},"diskSpace":{"status":"UP"},"ping":{"status":"UP"},"redis":{"status":"UP"}}}` + +func TestActuatorCheckerAcceptsConfirmedResponseShape(t *testing.T) { + t.Parallel() + client := newHTTPClient(func(request *http.Request) *http.Response { + if request.URL.Path != "/yms/actuator/health" { + t.Errorf("unexpected health path: %s", request.URL.Path) + } + return response(http.StatusOK, healthyActuatorSample) + }) + checker := newTestChecker(t, client) + report, err := checker.Wait(context.Background(), "http://127.0.0.1:8081/yms/actuator/health", time.Second, alwaysRunning) + if err != nil { + t.Fatalf("wait for healthy Actuator: %v", err) + } + if report.Status != "UP" || report.Components["db"] != "UP" || report.Components["redis"] != "UP" { + t.Fatalf("unexpected health report: %+v", report) + } +} + +func TestActuatorCheckerChecksOnceForTransactionRecovery(t *testing.T) { + t.Parallel() + client := newHTTPClient(func(*http.Request) *http.Response { + return response(http.StatusOK, healthyActuatorSample) + }) + checker := newTestChecker(t, client) + report, ready, err := checker.Check(context.Background(), "http://127.0.0.1:8081/yms/actuator/health", alwaysRunning) + if err != nil || !ready || report.Status != "UP" { + t.Fatalf("unexpected one-shot health result: report=%+v ready=%t err=%v", report, ready, err) + } +} + +func TestActuatorCheckerSamplesReadinessWithoutRestarting(t *testing.T) { + t.Parallel() + var requests atomic.Int32 + client := newHTTPClient(func(*http.Request) *http.Response { + if requests.Add(1) < 3 { + return response(http.StatusServiceUnavailable, `{"status":"DOWN"}`) + } + return response(http.StatusOK, healthyActuatorSample) + }) + checker := newTestChecker(t, client) + if _, err := checker.Wait(context.Background(), "http://127.0.0.1:8081/yms/actuator/health", time.Second, alwaysRunning); err != nil { + t.Fatalf("wait for delayed readiness: %v", err) + } + if requests.Load() != 3 { + t.Fatalf("unexpected request count: %d", requests.Load()) + } +} + +func TestActuatorCheckerStopsImmediatelyWhenContainerStops(t *testing.T) { + t.Parallel() + checker := newTestChecker(t, http.DefaultClient) + _, err := checker.Wait(context.Background(), "http://127.0.0.1:8081/yms/actuator/health", time.Second, + func(context.Context) (bool, error) { return false, nil }) + if !errors.Is(err, ErrWorkloadStopped) { + t.Fatalf("expected stopped workload error, got %v", err) + } +} + +func TestActuatorCheckerHonorsOverallTimeout(t *testing.T) { + t.Parallel() + client := newHTTPClient(func(*http.Request) *http.Response { + return response(http.StatusServiceUnavailable, `{"status":"DOWN"}`) + }) + checker := newTestChecker(t, client) + _, err := checker.Wait(context.Background(), "http://127.0.0.1:8081/yms/actuator/health", 25*time.Millisecond, alwaysRunning) + if !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("expected deadline error, got %v", err) + } +} + +func newTestChecker(t *testing.T, client *http.Client) *ActuatorChecker { + t.Helper() + checker, err := NewActuatorChecker(client, 5*time.Millisecond) + if err != nil { + t.Fatalf("create Actuator checker: %v", err) + } + return checker +} + +func alwaysRunning(context.Context) (bool, error) { + return true, nil +} + +type roundTripFunc func(*http.Request) (*http.Response, error) + +func (f roundTripFunc) RoundTrip(request *http.Request) (*http.Response, error) { + return f(request) +} + +func newHTTPClient(handle func(*http.Request) *http.Response) *http.Client { + return &http.Client{Transport: roundTripFunc(func(request *http.Request) (*http.Response, error) { + return handle(request), nil + })} +} + +func response(statusCode int, body string) *http.Response { + return &http.Response{ + StatusCode: statusCode, + Header: make(http.Header), + Body: io.NopCloser(bytes.NewBufferString(body)), + } +}