diff --git a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/setup/workflow/FileMigrationOperationStore.java b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/setup/workflow/FileMigrationOperationStore.java index ca05277d90..171a5d7c0a 100644 --- a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/setup/workflow/FileMigrationOperationStore.java +++ b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/setup/workflow/FileMigrationOperationStore.java @@ -67,6 +67,32 @@ public final class FileMigrationOperationStore implements MigrationOperationStor }); } + /** Creates or confirms one fully equal PENDING snapshot under the store lock. */ + MigrationOperationSnapshot createOrConfirm(MigrationOperationSnapshot snapshot) { + Objects.requireNonNull(snapshot, "snapshot"); + if (snapshot.state() != MigrationOperationState.PENDING) { + throw failure(SetupErrorCode.INVALID_REQUEST); + } + return locked(() -> { + List snapshots = read(); + for (MigrationOperationSnapshot current : snapshots) { + if (current.operationId().equals(snapshot.operationId())) { + if (current.equals(snapshot)) { + writeAndConfirm(snapshots); + return snapshot; + } + throw failure(SetupErrorCode.OPERATION_CONFLICT); + } + if (!current.terminal()) { + throw failure(SetupErrorCode.OPERATION_CONFLICT); + } + } + snapshots.add(snapshot); + writeAndConfirm(snapshots); + return snapshot; + }); + } + @Override public Optional find(String operationId) { requireSafeId(operationId); @@ -88,6 +114,15 @@ public final class FileMigrationOperationStore implements MigrationOperationStor return locked(() -> transition(read(), operationId, expectedState, replacement)); } + /** Transitions or confirms one fully equal replacement under the store lock. */ + MigrationOperationSnapshot compareAndTransitionOrConfirm( + String operationId, MigrationOperationState expectedState, MigrationOperationSnapshot replacement) { + requireSafeId(operationId); + Objects.requireNonNull(expectedState, "expectedState"); + Objects.requireNonNull(replacement, "replacement"); + return locked(() -> transitionOrConfirm(read(), operationId, expectedState, replacement)); + } + private MigrationOperationSnapshot transition( List snapshots, String operationId, MigrationOperationState expectedState, MigrationOperationSnapshot replacement) { @@ -107,6 +142,29 @@ public final class FileMigrationOperationStore implements MigrationOperationStor throw failure(SetupErrorCode.OPERATION_NOT_FOUND); } + private MigrationOperationSnapshot transitionOrConfirm( + List snapshots, String operationId, + MigrationOperationState expectedState, MigrationOperationSnapshot replacement) { + for (int index = 0; index < snapshots.size(); index++) { + MigrationOperationSnapshot current = snapshots.get(index); + if (current.operationId().equals(operationId)) { + if (current.equals(replacement)) { + writeAndConfirm(snapshots); + return replacement; + } + if (current.state() != expectedState) { + throw failure(SetupErrorCode.OPERATION_CONFLICT); + } + transitionPolicy.requireAllowed(current, replacement); + snapshots.set(index, replacement); + trim(snapshots); + writeAndConfirm(snapshots); + return replacement; + } + } + throw failure(SetupErrorCode.OPERATION_NOT_FOUND); + } + private List read() { if (!Files.exists(operationFile, LinkOption.NOFOLLOW_LINKS)) { return new ArrayList<>(); @@ -139,6 +197,32 @@ public final class FileMigrationOperationStore implements MigrationOperationStor } } + private void writeAndConfirm(List snapshots) { + collectionPolicy.validate(snapshots); + byte[] encoded = codec.encode(snapshots); + try { + publisher.publish(operationFile, encoded); + } catch (CommittedSetupFileDurabilityException uncertain) { + confirmAndRepublish(snapshots, encoded); + } catch (IOException failure) { + throw failure(SetupErrorCode.CONFIG_WRITE_FAILED); + } finally { + Arrays.fill(encoded, (byte) 0); + } + } + + private void confirmAndRepublish(List intended, byte[] encoded) { + List persisted = read(); + if (!persisted.equals(intended)) { + throw failure(SetupErrorCode.CONFIG_RECOVERY_REQUIRED); + } + try { + publisher.publish(operationFile, encoded); + } catch (IOException failure) { + throw failure(SetupErrorCode.CONFIG_RECOVERY_REQUIRED); + } + } + private void trim(List snapshots) { while (snapshots.stream().filter(MigrationOperationSnapshot::terminal).count() > HISTORY_LIMIT) { int oldestTerminal = -1; diff --git a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/setup/workflow/FileMigrationOperationStoreExactTest.java b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/setup/workflow/FileMigrationOperationStoreExactTest.java new file mode 100644 index 0000000000..3355162d7f --- /dev/null +++ b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/setup/workflow/FileMigrationOperationStoreExactTest.java @@ -0,0 +1,241 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0. + */ + +package org.apache.hertzbeat.manager.setup.workflow; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.time.Instant; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import org.apache.hertzbeat.manager.setup.api.DeploymentApiContract.MigrationOperationState; +import org.apache.hertzbeat.manager.setup.api.DeploymentApiContract.MigrationStage; +import org.apache.hertzbeat.manager.setup.api.DeploymentApiContract.MigrationTarget; +import org.apache.hertzbeat.manager.setup.api.DeploymentApiContract.VerificationState; +import org.apache.hertzbeat.manager.setup.api.SetupApiContract.ApplyMode; +import org.apache.hertzbeat.manager.setup.api.SetupApiContract.SetupErrorCode; +import org.apache.hertzbeat.manager.setup.security.CommittedSetupFileDurabilityException; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +class FileMigrationOperationStoreExactTest { + + private static final String IDENTITY = "a".repeat(64); + private static final String GENERATION = "candidate-generation"; + private static final Instant CREATED = Instant.parse("2026-08-10T01:00:00Z"); + + @TempDir + private Path root; + + @Test + void exactCreateAndTransitionRetriesAreIdempotent() { + FileMigrationOperationStore store = new FileMigrationOperationStore(root); + MigrationOperationSnapshot pending = pending("operation-a", IDENTITY, GENERATION); + MigrationOperationSnapshot running = running(pending); + + assertThat(store.createOrConfirm(pending)).isEqualTo(pending); + assertThat(store.createOrConfirm(pending)).isEqualTo(pending); + assertThat(store.compareAndTransitionOrConfirm( + pending.operationId(), MigrationOperationState.PENDING, running)).isEqualTo(running); + assertThat(store.compareAndTransitionOrConfirm( + pending.operationId(), MigrationOperationState.PENDING, running)).isEqualTo(running); + } + + @Test + void rejectsDifferingIdentityAndGeneration() { + FileMigrationOperationStore store = new FileMigrationOperationStore(root); + MigrationOperationSnapshot pending = pending("operation-a", IDENTITY, GENERATION); + store.createOrConfirm(pending); + + assertStoreError(SetupErrorCode.OPERATION_CONFLICT, + () -> store.createOrConfirm(pending("operation-a", "b".repeat(64), GENERATION))); + assertStoreError(SetupErrorCode.OPERATION_CONFLICT, + () -> store.createOrConfirm(pending("operation-a", IDENTITY, "other-generation"))); + } + + @Test + void differentActiveOperationAndAdvancedNonExactStateRemainConflicts() { + FileMigrationOperationStore store = new FileMigrationOperationStore(root); + MigrationOperationSnapshot pending = pending("operation-a", IDENTITY, GENERATION); + store.createOrConfirm(pending); + + assertStoreError(SetupErrorCode.OPERATION_CONFLICT, + () -> store.createOrConfirm(pending("operation-b", IDENTITY, GENERATION))); + + MigrationOperationSnapshot running = running(pending); + store.compareAndTransitionOrConfirm(pending.operationId(), MigrationOperationState.PENDING, running); + MigrationOperationSnapshot later = new MigrationOperationSnapshot( + running.operationId(), running.state(), running.target(), running.applyMode(), running.stage(), + 20, running.createdAt(), running.startedAt(), running.completedAt(), running.verificationState(), + running.errorCode(), running.rollbackOrigin(), running.nextPollAfterMillis(), + running.activationAvailable(), running.restartRequired(), running.externalApplyRequired(), + running.targetIdentityHash(), running.managedCandidateGeneration()); + assertStoreError(SetupErrorCode.OPERATION_CONFLICT, () -> store.compareAndTransitionOrConfirm( + pending.operationId(), MigrationOperationState.PENDING, later)); + } + + @Test + void committedCreateIsConfirmedByAuthoritativeReadBack() { + MigrationOperationSnapshot pending = pending("operation-a", IDENTITY, GENERATION); + MigrationOperationFilePublisher committed = new MigrationOperationFilePublisher(root); + AtomicBoolean first = new AtomicBoolean(true); + FileMigrationOperationStore uncertain = new FileMigrationOperationStore(root, (target, content) -> { + committed.publish(target, content); + if (first.getAndSet(false)) { + throw new CommittedSetupFileDurabilityException(); + } + }); + + assertThat(uncertain.createOrConfirm(pending)).isEqualTo(pending); + } + + @Test + void committedTransitionIsConfirmedByAuthoritativeReadBack() { + MigrationOperationSnapshot pending = pending("operation-a", IDENTITY, GENERATION); + new FileMigrationOperationStore(root).create(pending); + MigrationOperationSnapshot running = running(pending); + MigrationOperationFilePublisher committed = new MigrationOperationFilePublisher(root); + AtomicBoolean first = new AtomicBoolean(true); + FileMigrationOperationStore uncertain = new FileMigrationOperationStore(root, (target, content) -> { + committed.publish(target, content); + if (first.getAndSet(false)) { + throw new CommittedSetupFileDurabilityException(); + } + }); + + assertThat(uncertain.compareAndTransitionOrConfirm( + pending.operationId(), MigrationOperationState.PENDING, running)).isEqualTo(running); + } + + @Test + void exactReplayMustConfirmDurabilityBeforeReturningSuccess() { + MigrationOperationSnapshot pending = pending("operation-a", IDENTITY, GENERATION); + MigrationOperationFilePublisher committed = new MigrationOperationFilePublisher(root); + AtomicInteger publications = new AtomicInteger(); + FileMigrationOperationStore uncertain = new FileMigrationOperationStore(root, (target, content) -> { + committed.publish(target, content); + if (publications.incrementAndGet() <= 4) { + throw new CommittedSetupFileDurabilityException(); + } + }); + + assertStoreError(SetupErrorCode.CONFIG_RECOVERY_REQUIRED, + () -> uncertain.createOrConfirm(pending)); + assertStoreError(SetupErrorCode.CONFIG_RECOVERY_REQUIRED, + () -> uncertain.createOrConfirm(pending)); + assertThat(uncertain.createOrConfirm(pending)).isEqualTo(pending); + assertThat(publications).hasValue(5); + } + + @Test + void exactTransitionReplayMustConfirmDurabilityBeforeReturningSuccess() { + MigrationOperationSnapshot pending = pending("operation-a", IDENTITY, GENERATION); + MigrationOperationSnapshot running = running(pending); + new FileMigrationOperationStore(root).create(pending); + MigrationOperationFilePublisher committed = new MigrationOperationFilePublisher(root); + AtomicInteger publications = new AtomicInteger(); + FileMigrationOperationStore uncertain = new FileMigrationOperationStore(root, (target, content) -> { + committed.publish(target, content); + if (publications.incrementAndGet() <= 4) { + throw new CommittedSetupFileDurabilityException(); + } + }); + + assertStoreError(SetupErrorCode.CONFIG_RECOVERY_REQUIRED, + () -> uncertain.compareAndTransitionOrConfirm( + pending.operationId(), MigrationOperationState.PENDING, running)); + assertStoreError(SetupErrorCode.CONFIG_RECOVERY_REQUIRED, + () -> uncertain.compareAndTransitionOrConfirm( + pending.operationId(), MigrationOperationState.PENDING, running)); + assertThat(uncertain.compareAndTransitionOrConfirm( + pending.operationId(), MigrationOperationState.PENDING, running)).isEqualTo(running); + assertThat(publications).hasValue(5); + } + + @Test + void uncertainCreateMissingOrCorruptFailsClosed() { + MigrationOperationSnapshot pending = pending("operation-a", IDENTITY, GENERATION); + FileMigrationOperationStore missing = new FileMigrationOperationStore(root, (target, content) -> { + throw new CommittedSetupFileDurabilityException(); + }); + assertStoreError(SetupErrorCode.CONFIG_RECOVERY_REQUIRED, + () -> missing.createOrConfirm(pending)); + + FileMigrationOperationStore corrupt = new FileMigrationOperationStore(root, (target, content) -> { + Files.createDirectories(target.getParent()); + Files.writeString(target, "schema=99\n", StandardCharsets.UTF_8); + throw new CommittedSetupFileDurabilityException(); + }); + assertStoreError(SetupErrorCode.CONFIG_RECOVERY_REQUIRED, + () -> corrupt.createOrConfirm(pending)); + } + + @Test + void uncertainTransitionMissingOrCorruptFailsClosed() { + MigrationOperationSnapshot pending = pending("operation-a", IDENTITY, GENERATION); + MigrationOperationSnapshot running = running(pending); + new FileMigrationOperationStore(root).create(pending); + FileMigrationOperationStore missing = new FileMigrationOperationStore(root, (target, content) -> { + Files.delete(target); + throw new CommittedSetupFileDurabilityException(); + }); + assertStoreError(SetupErrorCode.CONFIG_RECOVERY_REQUIRED, + () -> missing.compareAndTransitionOrConfirm( + pending.operationId(), MigrationOperationState.PENDING, running)); + + new FileMigrationOperationStore(root).create(pending); + FileMigrationOperationStore corrupt = new FileMigrationOperationStore(root, (target, content) -> { + Files.writeString(target, "schema=99\n", StandardCharsets.UTF_8); + throw new CommittedSetupFileDurabilityException(); + }); + assertStoreError(SetupErrorCode.CONFIG_RECOVERY_REQUIRED, + () -> corrupt.compareAndTransitionOrConfirm( + pending.operationId(), MigrationOperationState.PENDING, running)); + } + + @Test + void exactMethodSurfaceContainsNoMigrationPayloadOrCredentialFields() { + FileMigrationOperationStore store = new FileMigrationOperationStore(root); + assertThat(store.toString()) + .doesNotContain("jdbc:", "password", "username", IDENTITY, GENERATION); + } + + private static MigrationOperationSnapshot pending(String operation, String identity, String generation) { + return new MigrationOperationSnapshot(operation, MigrationOperationState.PENDING, MigrationTarget.MYSQL, + ApplyMode.MANAGED_WRITE, MigrationStage.QUEUED, 0, CREATED, null, null, + VerificationState.PENDING, null, null, 1000, false, false, false, + identity, generation); + } + + private static MigrationOperationSnapshot running(MigrationOperationSnapshot pending) { + return new MigrationOperationSnapshot( + pending.operationId(), MigrationOperationState.RUNNING, pending.target(), pending.applyMode(), + MigrationStage.COPYING, 10, pending.createdAt(), pending.createdAt().plusSeconds(1), null, + VerificationState.PENDING, null, null, 1000, false, false, false, + pending.targetIdentityHash(), pending.managedCandidateGeneration()); + } + + private static void assertStoreError(SetupErrorCode code, ThrowingAction action) { + assertThatThrownBy(action::run) + .isInstanceOfSatisfying(MigrationOperationStoreException.class, + failure -> assertThat(failure.errorCode()).isEqualTo(code)) + .hasNoCause() + .hasMessageNotContaining("jdbc") + .hasMessageNotContaining("password") + .hasMessageNotContaining("username") + .hasMessageNotContaining("schema=99"); + } + + @FunctionalInterface + private interface ThrowingAction { + void run() throws Exception; + } +}