feat: backend container executor implement
This commit is contained in:
@@ -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
|
||||
)
|
||||
|
||||
@@ -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=
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
@@ -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)
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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"`
|
||||
}
|
||||
@@ -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)),
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user