package nativebackendexecutor import ( "context" "crypto/sha256" "encoding/hex" "errors" "io" "log/slog" "net/http" "os" "path/filepath" "sync" "testing" "time" "yms-daemon/internal/filestore" "yms-daemon/internal/healthcheck" "yms-daemon/internal/systemd" "yms-daemon/internal/transaction" ) // TestExecutorInstallsJarStartsExactUnitAndReachesSwitching 验证成功路径:安装 JAR、启动精确单元、 // 槽位软链接指向新版本、事务推进到 switching,且重复运行不会再次启动单元。 func TestExecutorInstallsJarStartsExactUnitAndReachesSwitching(t *testing.T) { ctx := context.Background() executor, store, releaseStore, units, request, previousTarget := testNativeExecutor(t) record := createNativeTransaction(t, store, "success") if err := executor.Run(ctx, record.ID, request); err != nil { t.Fatalf("run native backend executor: %v", err) } current, err := store.Transaction(ctx, record.ID) if err != nil || current.State != transaction.StateSwitching { t.Fatalf("unexpected transaction after native preparation: record=%+v err=%v", current, err) } installed, found, err := releaseStore.Inspect(request.ReleasePath, request.ArtifactIdentity) if err != nil || !found { t.Fatalf("inspect installed backend JAR: file=%+v found=%t err=%v", installed, found, err) } actualTarget, err := os.Readlink(request.SlotJarPath) if err != nil || actualTarget != installed.Path || actualTarget == previousTarget { t.Fatalf("unexpected slot link: target=%q installed=%q previous=%q err=%v", actualTarget, installed.Path, previousTarget, err) } units.mu.Lock() startCalls := units.startCalls stopCalls := units.stopCalls unitActiveState := units.unit.ActiveState startedName := units.startedName units.mu.Unlock() if startCalls != 1 || stopCalls != 0 || unitActiveState != activeState || startedName != request.UnitName { t.Fatalf("unexpected systemd calls: start=%d stop=%d state=%s name=%s", startCalls, stopCalls, unitActiveState, startedName) } pending, err := store.PendingSteps(ctx, record.ID) if err != nil || len(pending) != 0 { t.Fatalf("unexpected pending native steps: steps=%+v err=%v", pending, err) } if err := executor.Run(ctx, record.ID, request); err != nil { t.Fatalf("repeat executor at switching state: %v", err) } units.mu.Lock() repeatedStartCalls := units.startCalls units.mu.Unlock() if repeatedStartCalls != startCalls { t.Fatalf("switching state repeated systemd start: before=%d after=%d", startCalls, repeatedStartCalls) } } // TestExecutorRollsBackSlotAndStopsUnitWhenHealthFails 验证健康检查失败时:事务回滚、槽位软链接恢复、 // 单元被停止且状态回到非活跃。 func TestExecutorRollsBackSlotAndStopsUnitWhenHealthFails(t *testing.T) { ctx := context.Background() executor, store, _, units, request, previousTarget := testNativeExecutor(t) executor.checker = &fakeActuatorChecker{ waitErr: errors.New("Actuator rejected native backend"), checkReady: false, } record := createNativeTransaction(t, store, "health-failure") err := executor.Run(ctx, record.ID, request) if err == nil { t.Fatal("expected native backend health failure") } current, readErr := store.Transaction(ctx, record.ID) if readErr != nil || current.State != transaction.StateRolledBack { t.Fatalf("unexpected compensated transaction: record=%+v err=%v", current, readErr) } actualTarget, linkErr := os.Readlink(request.SlotJarPath) if linkErr != nil || actualTarget != previousTarget { t.Fatalf("native slot was not restored: target=%q previous=%q err=%v", actualTarget, previousTarget, linkErr) } units.mu.Lock() startCalls := units.startCalls stopCalls := units.stopCalls unitActiveState := units.unit.ActiveState units.mu.Unlock() if startCalls != 1 || stopCalls != 1 || unitActiveState != inactiveState { t.Fatalf("unexpected compensated systemd state: start=%d stop=%d state=%s", startCalls, stopCalls, unitActiveState) } } // TestExecutorResumesPersistedRollback 验证从持久化的 rolling back 状态恢复补偿:恢复槽位并停止单元。 func TestExecutorResumesPersistedRollback(t *testing.T) { ctx := context.Background() executor, store, _, units, request, previousTarget := testNativeExecutor(t) record := createNativeTransaction(t, store, "resume-rollback") transitionNativeToPrepared(t, store, record.ID) if _, err := executor.prepare(ctx, record.ID, request); err != nil { t.Fatalf("prepare native backend before rollback interruption: %v", err) } if _, err := store.Transition(ctx, record.ID, transaction.StateStarting, "test starting"); err != nil { t.Fatalf("transition native transaction to starting: %v", err) } if err := units.Start(ctx, request.UnitName); err != nil { t.Fatalf("start native backend before rollback interruption: %v", err) } if _, err := store.Transition(ctx, record.ID, transaction.StateRollingBack, "test interrupted rollback"); err != nil { t.Fatalf("persist interrupted rollback state: %v", err) } if err := executor.Run(ctx, record.ID, request); err != nil { t.Fatalf("resume native backend rollback: %v", err) } current, err := store.Transaction(ctx, record.ID) if err != nil || current.State != transaction.StateRolledBack { t.Fatalf("unexpected resumed rollback state: record=%+v err=%v", current, err) } actualTarget, err := os.Readlink(request.SlotJarPath) if err != nil || actualTarget != previousTarget { t.Fatalf("resumed rollback did not restore slot: target=%q previous=%q err=%v", actualTarget, previousTarget, err) } units.mu.Lock() stopCalls := units.stopCalls unitActiveState := units.unit.ActiveState units.mu.Unlock() if stopCalls != 1 || unitActiveState != inactiveState { t.Fatalf("resumed rollback did not stop unit: stop=%d state=%s", stopCalls, unitActiveState) } } // TestExecutorRollbackRemovesFirstDeploymentSlotLink 验证首次部署(无先前目标)失败时,回滚会移除槽位软链接。 func TestExecutorRollbackRemovesFirstDeploymentSlotLink(t *testing.T) { ctx := context.Background() executor, store, _, _, request, _ := testNativeExecutor(t) if err := os.Remove(request.SlotJarPath); err != nil { t.Fatalf("remove seeded slot link: %v", err) } request.PreviousSlotTarget = "" executor.checker = &fakeActuatorChecker{ waitErr: errors.New("Actuator rejected first native backend deployment"), checkReady: false, } record := createNativeTransaction(t, store, "first-deployment-rollback") if err := executor.Run(ctx, record.ID, request); err == nil { t.Fatal("expected first native backend deployment health failure") } current, err := store.Transaction(ctx, record.ID) if err != nil || current.State != transaction.StateRolledBack { t.Fatalf("unexpected first deployment rollback state: record=%+v err=%v", current, err) } if _, err := os.Lstat(request.SlotJarPath); !errors.Is(err, os.ErrNotExist) { t.Fatalf("first deployment rollback retained slot link: %v", err) } } // TestExecutorRecoversRecordedSlotIntentWithoutChangingRequest 验证崩溃恢复:已记录的槽位意图不依赖请求变更, // 执行器能直接复用并完成后续步骤。 func TestExecutorRecoversRecordedSlotIntentWithoutChangingRequest(t *testing.T) { ctx := context.Background() executor, store, releaseStore, units, request, _ := testNativeExecutor(t) record := createNativeTransaction(t, store, "recover-slot") transitionNativeToPrepared(t, store, record.ID) installOperation := &installJarOperation{ store: releaseStore, sourcePath: request.ArtifactPath, releasePath: request.ReleasePath, identity: request.ArtifactIdentity, } if _, err := executor.coordinator.ExecuteStep(ctx, record.ID, installIntent(request), installOperation); err != nil { t.Fatalf("install backend JAR before simulated crash: %v", err) } installed, found, err := releaseStore.Inspect(request.ReleasePath, request.ArtifactIdentity) if err != nil || !found { t.Fatalf("inspect backend JAR before simulated crash: file=%+v found=%t err=%v", installed, found, err) } if _, _, err := store.RecordStepIntent(ctx, record.ID, bindIntent(request, installed.Path)); err != nil { t.Fatalf("record slot intent before simulated crash: %v", err) } if err := executor.Run(ctx, record.ID, request); err != nil { t.Fatalf("recover native backend executor: %v", err) } actualTarget, err := os.Readlink(request.SlotJarPath) if err != nil || actualTarget != installed.Path { t.Fatalf("unexpected recovered slot link: target=%q installed=%q err=%v", actualTarget, installed.Path, err) } units.mu.Lock() startCalls := units.startCalls units.mu.Unlock() if startCalls != 1 { t.Fatalf("unexpected recovered systemd start count: %d", startCalls) } } // TestExecutorRejectsChangedRecoveryIntentAndPreservesStartingState 验证请求变更导致步骤冲突时被拒绝,且事务停留在 starting 状态。 func TestExecutorRejectsChangedRecoveryIntentAndPreservesStartingState(t *testing.T) { ctx := context.Background() executor, store, _, _, request, _ := testNativeExecutor(t) record := createNativeTransaction(t, store, "changed-request") transitionNativeToPrepared(t, store, record.ID) if _, err := executor.prepare(ctx, record.ID, request); err != nil { t.Fatalf("prepare original native backend request: %v", err) } if _, err := store.Transition(ctx, record.ID, transaction.StateStarting, "test starting"); err != nil { t.Fatalf("transition native backend to starting: %v", err) } changed := request changed.ReleasePath = "different-release.jar" err := executor.Run(ctx, record.ID, changed) if !errors.Is(err, transaction.ErrStepConflict) { t.Fatalf("expected persisted native 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) } } // TestExecutorRejectsActiveUnitBeforeChangingFiles 验证单元已活跃时校验失败:事务进入 failed,且不改动发布存储与槽位软链接。 func TestExecutorRejectsActiveUnitBeforeChangingFiles(t *testing.T) { ctx := context.Background() executor, store, releaseStore, units, request, previousTarget := testNativeExecutor(t) units.mu.Lock() units.unit.ActiveState = activeState units.unit.SubState = "running" units.mu.Unlock() record := createNativeTransaction(t, store, "active-unit") if err := executor.Run(ctx, record.ID, request); err == nil { t.Fatal("expected active unit validation failure") } current, err := store.Transaction(ctx, record.ID) if err != nil || current.State != transaction.StateFailed { t.Fatalf("unexpected active-unit transaction state: record=%+v err=%v", current, err) } if _, found, err := releaseStore.Inspect(request.ReleasePath, request.ArtifactIdentity); err != nil || found { t.Fatalf("validation failure changed release store: found=%t err=%v", found, err) } actualTarget, err := os.Readlink(request.SlotJarPath) if err != nil || actualTarget != previousTarget { t.Fatalf("validation failure changed slot link: target=%q previous=%q err=%v", actualTarget, previousTarget, err) } } // TestUnitStopOperationAcceptsSystemdFailedAsStopped 验证 unitStopOperation 把 systemd 的 failed 状态视为已停止。 func TestUnitStopOperationAcceptsSystemdFailedAsStopped(t *testing.T) { units := &fakeUnitManager{unit: systemd.Unit{ Name: "yms-backend@8080.service", LoadState: "loaded", ActiveState: failedState, SubState: "failed", }} operation := &unitStopOperation{units: units, name: units.unit.Name} inspection, err := operation.Inspect(context.Background()) if err != nil { t.Fatalf("inspect stopped native backend unit: %v", err) } if inspection.Status != transaction.InspectionApplied { t.Fatalf("unexpected stopped native backend unit inspection: %+v", inspection) } } // testNativeExecutor 构造一套完整且可复用的原生后端测试夹具。 // t 用于失败报告与清理注册。返回执行器、事务存储、发布存储、伪造的单元管理器、请求与先前槽位目标。 func testNativeExecutor(t *testing.T) (*Executor, *transaction.Store, *filestore.Store, *fakeUnitManager, Request, string) { 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) } releaseStore, err := filestore.New(t.TempDir()) if err != nil { t.Fatalf("create native release store: %v", err) } artifactContent := []byte("native backend jar content") artifactPath := filepath.Join(t.TempDir(), "glory-soft-yms.jar") if err := os.WriteFile(artifactPath, artifactContent, 0o640); err != nil { t.Fatalf("write native backend artifact: %v", err) } previousTarget := filepath.Join(t.TempDir(), "previous-backend.jar") if err := os.WriteFile(previousTarget, []byte("previous backend jar"), 0o640); err != nil { t.Fatalf("write previous backend JAR: %v", err) } slotDirectory := t.TempDir() slotJarPath := filepath.Join(slotDirectory, "backend-green.jar") if err := os.Symlink(previousTarget, slotJarPath); err != nil { t.Fatalf("create previous native backend slot link: %v", err) } request := Request{ ArtifactPath: artifactPath, ArtifactIdentity: testIdentity(artifactContent), ReleasePath: "glory-soft-yms-20260815.jar", SlotJarPath: slotJarPath, PreviousSlotTarget: previousTarget, UnitName: "yms-green.service", Port: 8081, HealthEndpoint: "http://127.0.0.1:8081/yms/actuator/health", } units := &fakeUnitManager{unit: systemd.Unit{ Name: request.UnitName, LoadState: "loaded", ActiveState: inactiveState, SubState: "dead", }} executor, err := New(store, coordinator, releaseStore, units, &http.Client{}) if err != nil { t.Fatalf("create native backend executor: %v", err) } executor.checker = &fakeActuatorChecker{ waitReport: healthcheck.ActuatorReport{Status: "UP", Components: map[string]string{"db": "UP"}}, checkReport: healthcheck.ActuatorReport{Status: "UP", Components: map[string]string{"db": "UP"}}, checkReady: true, } return executor, store, releaseStore, units, request, previousTarget } // createNativeTransaction 在测试事务存储中创建一条原生后端事务记录。 // t 用于失败报告;store 是事务存储;suffix 用于构造唯一的事务 ID。返回创建的事务记录。 func createNativeTransaction(t *testing.T, store *transaction.Store, suffix string) transaction.Transaction { t.Helper() record, _, err := store.CreateTransaction(context.Background(), transaction.CreateRequest{ ID: "native-backend-" + suffix, IdempotencyKey: "native-backend-request-" + suffix, Source: "test", Service: "backend", }) if err != nil { t.Fatalf("create native backend transaction: %v", err) } return record } // transitionNativeToPrepared 把事务依次推进到 validating 与 prepared 状态,供测试预置前置步骤。 // t 用于失败报告;store 是事务存储;transactionID 定位事务。 func transitionNativeToPrepared(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 native transaction to validating: %v", err) } if _, err := store.Transition(ctx, transactionID, transaction.StatePrepared, "test prepared"); err != nil { t.Fatalf("transition native transaction to prepared: %v", err) } } // testIdentity 根据内容计算文件身份(大小与 SHA-256),用于测试夹具与校验。 // content 文件内容。返回对应的 filestore.Identity。 func testIdentity(content []byte) filestore.Identity { digest := sha256.Sum256(content) return filestore.Identity{Size: int64(len(content)), SHA256: hex.EncodeToString(digest[:])} } // fakeUnitManager systemd.Manager 的测试替身,记录单元状态与启动/停止调用次数。 type fakeUnitManager struct { mu sync.Mutex unit systemd.Unit startErr error stopErr error startCalls int stopCalls int startedName string } // Inspect 返回单元状态:名称不匹配时返回 ErrUnitNotFound。 // 参数 ctx 与 name 用于取消与定位,name 决定返回哪个单元。返回单元状态与错误。 func (m *fakeUnitManager) Inspect(_ context.Context, name string) (systemd.Unit, error) { m.mu.Lock() defer m.mu.Unlock() if name != m.unit.Name { return systemd.Unit{}, systemd.ErrUnitNotFound } return m.unit, nil } // Start 模拟启动单元:记录调用与名称,成功置为活跃,失败置为 failed 并返回 startErr。 // 参数 ctx 与 name 用于取消与定位,name 决定记录的名称。返回值是启动错误(若配置了 startErr)。 func (m *fakeUnitManager) Start(_ context.Context, name string) error { m.mu.Lock() defer m.mu.Unlock() m.startCalls++ m.startedName = name if m.startErr != nil { m.unit.ActiveState = failedState m.unit.SubState = "failed" return m.startErr } m.unit.ActiveState = activeState m.unit.SubState = "running" return nil } // Stop 模拟停止单元:记录调用,成功置为非活跃,失败返回 stopErr。 // 参数 ctx 与 name 用于取消与定位。返回值是停止错误(若配置了 stopErr)。 func (m *fakeUnitManager) Stop(context.Context, string) error { m.mu.Lock() defer m.mu.Unlock() m.stopCalls++ if m.stopErr != nil { return m.stopErr } m.unit.ActiveState = inactiveState m.unit.SubState = "dead" return nil } // fakeActuatorChecker actuatorChecker 的测试替身,返回可配置的健康报告与错误。 type fakeActuatorChecker struct { waitReport healthcheck.ActuatorReport waitErr error checkReport healthcheck.ActuatorReport checkReady bool checkErr error } // Wait 模拟等待健康检查:先执行运行探针,未运行返回 ErrWorkloadStopped,否则返回 waitReport 与 waitErr。 // ctx 用于取消;running 是运行探针。返回值是健康报告与错误。 func (c *fakeActuatorChecker) Wait(ctx context.Context, _ string, _ time.Duration, running healthcheck.RunningProbe) (healthcheck.ActuatorReport, error) { isRunning, err := running(ctx) if err != nil { return healthcheck.ActuatorReport{}, err } if !isRunning { return healthcheck.ActuatorReport{}, healthcheck.ErrWorkloadStopped } return c.waitReport, c.waitErr } // Check 模拟单次健康检查:先执行运行探针,未运行返回 ErrWorkloadStopped,否则返回 checkReport、checkReady 与 checkErr。 // ctx 用于取消;running 是运行探针。返回值是健康报告、就绪标志与错误。 func (c *fakeActuatorChecker) Check(ctx context.Context, _ string, running healthcheck.RunningProbe) (healthcheck.ActuatorReport, bool, error) { isRunning, err := running(ctx) if err != nil { return healthcheck.ActuatorReport{}, false, err } if !isRunning { return healthcheck.ActuatorReport{}, false, healthcheck.ErrWorkloadStopped } return c.checkReport, c.checkReady, c.checkErr } // 以下编译期断言确保测试替身实现了相应接口。 var _ systemd.Manager = (*fakeUnitManager)(nil) var _ actuatorChecker = (*fakeActuatorChecker)(nil)