Confirm metadata operation journal writes

This commit is contained in:
Logic
2026-08-10 13:53:28 +08:00
parent 9ba2f5cd17
commit bb3aef3ecf
2 changed files with 325 additions and 0 deletions
@@ -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<MigrationOperationSnapshot> 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<MigrationOperationSnapshot> 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<MigrationOperationSnapshot> 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<MigrationOperationSnapshot> 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<MigrationOperationSnapshot> 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<MigrationOperationSnapshot> 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<MigrationOperationSnapshot> intended, byte[] encoded) {
List<MigrationOperationSnapshot> 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<MigrationOperationSnapshot> snapshots) {
while (snapshots.stream().filter(MigrationOperationSnapshot::terminal).count() > HISTORY_LIMIT) {
int oldestTerminal = -1;
@@ -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;
}
}