Merge "index: recover from crashes between NoteDb write and index write"
diff --git a/Documentation/config-gerrit.txt b/Documentation/config-gerrit.txt
index b7bb053..36c9077 100644
--- a/Documentation/config-gerrit.txt
+++ b/Documentation/config-gerrit.txt
@@ -4215,6 +4215,40 @@
link:#schedule-configuration-examples[Schedule examples] can be found
in the link:#schedule-configuration[Schedule Configuration] section.
+[[index.staleChangeRecovery]]index.staleChangeRecovery::
++
+Whether to enable automatic recovery of change index updates that were not
+completed due to a user interrupt or a crash.
++
+When enabled, Gerrit writes a per-change intent file under
+`$site_path/data/pending-index/<pid>_<startTime>/<threadId>/<hash>` before each
+NoteDb update and removes it once the index write succeeds. On startup, any
+intent files left behind by a previously crashed process are recovered
+immediately. A background scanner then periodically picks up intent files whose
+writer thread within the current process is no longer alive, reindexing the
+affected changes automatically.
++
+Disabling this also disables both the startup recovery and the background
+scanner, so any index inconsistencies caused by a crash must be resolved with a
+manual link:cmd-index-changes.html[reindex].
++
+> **NOTE**: The stale change recovery is enabled only when `index.commitWithin`
+> is set to zero and `index.indexChangesAsync` is false. By default, Lucene flush
+> to disk is deferred until the `commitWithin` interval elapses, making the stale
+> change recovery ineffective.
++
+Defaults to `false`.
+
+[[index.staleChangeRecoveryInterval]]index.staleChangeRecoveryInterval::
++
+How often the background scanner checks for dead-thread intent files within the
+current process. Recovery of intent files from a previously crashed process
+happens once at startup and is not affected by this interval.
++
+Only used when link:#index.staleChangeRecovery[index.staleChangeRecovery] is `true`.
++
+Defaults to `5m`.
+
==== Lucene configuration
Open and closed changes are indexed in separate indexes named
diff --git a/java/com/google/gerrit/server/index/IndexModule.java b/java/com/google/gerrit/server/index/IndexModule.java
index 6b7e87b..7b22f3a 100644
--- a/java/com/google/gerrit/server/index/IndexModule.java
+++ b/java/com/google/gerrit/server/index/IndexModule.java
@@ -47,6 +47,7 @@
import com.google.gerrit.server.index.change.ChangeIndexRewriter;
import com.google.gerrit.server.index.change.ChangeIndexer;
import com.google.gerrit.server.index.change.ChangeSchemaDefinitions;
+import com.google.gerrit.server.index.change.PendingIndexUpdateScanner;
import com.google.gerrit.server.index.change.StalenessChecker;
import com.google.gerrit.server.index.group.GroupIndexCollection;
import com.google.gerrit.server.index.group.GroupIndexDefinition;
@@ -130,6 +131,7 @@
factory(ChangeIndexer.Factory.class);
factory(StalenessChecker.Factory.class);
factory(AllChangesIndexer.Factory.class);
+ install(new PendingIndexUpdateScanner.Module());
bind(GroupIndexRewriter.class);
// GroupIndexCollection is already bound very high up in SchemaModule.
diff --git a/java/com/google/gerrit/server/index/change/PendingIndexUpdate.java b/java/com/google/gerrit/server/index/change/PendingIndexUpdate.java
new file mode 100644
index 0000000..9d63ff5
--- /dev/null
+++ b/java/com/google/gerrit/server/index/change/PendingIndexUpdate.java
@@ -0,0 +1,179 @@
+// Copyright (C) 2026 The Android Open Source Project
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package com.google.gerrit.server.index.change;
+
+import com.google.common.flogger.FluentLogger;
+import com.google.common.hash.Hashing;
+import com.google.gerrit.entities.Change;
+import com.google.gerrit.entities.Project;
+import com.google.gerrit.server.config.GerritServerConfig;
+import com.google.gerrit.server.config.SitePaths;
+import com.google.gerrit.server.project.NoSuchChangeException;
+import com.google.gson.Gson;
+import com.google.gson.JsonSyntaxException;
+import com.google.inject.Inject;
+import com.google.inject.Singleton;
+import java.io.IOException;
+import java.lang.management.ManagementFactory;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.DirectoryNotEmptyException;
+import java.nio.file.Files;
+import java.nio.file.NoSuchFileException;
+import java.nio.file.Path;
+import java.nio.file.StandardCopyOption;
+import org.eclipse.jgit.lib.Config;
+
+/**
+ * Manages the change-index write-ahead intent files under {@code $site_dir/data/pending-index/}.
+ *
+ * <p>Each intent is a file at {@code <data_dir>/<pid>_<start_time>/<threadId>/sha(project, change)}
+ * with the JSON content of {@link Intent}.
+ */
+@Singleton
+public final class PendingIndexUpdate {
+ record Intent(String project, int changeId, String operation) {}
+
+ private static final FluentLogger logger = FluentLogger.forEnclosingClass();
+ private static final String PROCESS_MARKER =
+ ProcessHandle.current().pid() + "_" + ManagementFactory.getRuntimeMXBean().getStartTime();
+ private static final Gson GSON = new Gson();
+ private final ChangeIndexer indexer;
+ private final boolean enabled;
+ final Path intentDir;
+ final Path buildingDir;
+ final Path runningDir;
+
+ @Inject
+ public PendingIndexUpdate(
+ SitePaths sitePaths, ChangeIndexer indexer, @GerritServerConfig Config cfg) {
+ intentDir = sitePaths.data_dir.resolve("pending-index");
+ buildingDir = intentDir.resolve("building");
+ runningDir = intentDir.resolve(PROCESS_MARKER);
+ this.indexer = indexer;
+ this.enabled = computeEnabled(cfg);
+ }
+
+ /** Returns {@code true} if stale change recovery is active for this process. */
+ public boolean isEnabled() {
+ return enabled;
+ }
+
+ private static boolean computeEnabled(Config cfg) {
+ if (!cfg.getBoolean("index", null, "staleChangeRecovery", false)) {
+ return false;
+ }
+ if (cfg.getBoolean("index", null, "indexChangesAsync", false)) {
+ logger.atWarning().log(
+ "index.staleChangeRecovery has no effect when index.indexChangesAsync is true;"
+ + " stale change recovery is disabled");
+ return false;
+ }
+ for (String subsection : new String[] {"changes", "changes_open", "changes_closed"}) {
+ long commitWithin = cfg.getLong("index", subsection, "commitWithin", 0L);
+ if (commitWithin != 0) {
+ logger.atWarning().log(
+ "index.staleChangeRecovery has no effect when index.%s.commitWithin is non-zero;"
+ + " stale change recovery is disabled",
+ subsection);
+ return false;
+ }
+ }
+ return true;
+ }
+
+ /** Returns the per-thread intent directory for {@code threadId}. */
+ public Path threadDir(long threadId) {
+ return runningDir.resolve(String.valueOf(threadId));
+ }
+
+ public String filename(Project.NameKey project, Change.Id changeId) {
+ return Hashing.sha256()
+ .hashString("%s_%s".formatted(project, changeId), StandardCharsets.UTF_8)
+ .toString();
+ }
+
+ public void cleanIfEmpty(Path dir) {
+ try {
+ Files.delete(dir);
+ } catch (NoSuchFileException | DirectoryNotEmptyException ignored) {
+ // Already gone or not empty.
+ } catch (IOException e) {
+ logger.atWarning().withCause(e).log("Failed to delete directory %s", dir);
+ }
+ }
+
+ /** Writes an intent file for the given change under the thread's pending directory. */
+ public void write(long threadId, Project.NameKey project, Change.Id changeId, boolean delete)
+ throws IOException {
+ Files.createDirectories(buildingDir);
+ Path tmp =
+ Files.writeString(
+ Files.createTempFile(buildingDir, null, null),
+ GSON.toJson(new Intent(project.get(), changeId.get(), delete ? "delete" : "index")));
+
+ Path dir = threadDir(threadId);
+ Files.createDirectories(dir);
+ Files.move(tmp, dir.resolve(filename(project, changeId)), StandardCopyOption.ATOMIC_MOVE);
+ }
+
+ /** Deletes the intent file for {@code changeId} under the thread's pending directory. */
+ public void delete(long threadId, Project.NameKey project, Change.Id changeId) {
+ try {
+ Path threadDir = threadDir(threadId);
+ Files.deleteIfExists(threadDir.resolve(filename(project, changeId)));
+ cleanIfEmpty(threadDir);
+ } catch (IOException e) {
+ logger.atWarning().withCause(e).log(
+ "Failed to delete pending index intent for change %s in thread %d", changeId, threadId);
+ }
+ }
+
+ /** Reads the intent file, applies the index operation, then deletes the file. */
+ public void recover(Path file) throws IOException {
+ Intent intent;
+ try {
+ intent = GSON.fromJson(Files.readString(file), Intent.class);
+ } catch (JsonSyntaxException e) {
+ logger.atWarning().withCause(e).log(
+ "Malformed pending index intent, deleting %s", file.getFileName());
+ Files.deleteIfExists(file);
+ return;
+ }
+ if (intent == null
+ || intent.project() == null
+ || intent.operation() == null
+ || intent.changeId() <= 0) {
+ logger.atWarning().log("Malformed pending index intent, deleting %s", file.getFileName());
+ Files.deleteIfExists(file);
+ return;
+ }
+ Project.NameKey project = Project.nameKey(intent.project());
+ try {
+ switch (intent.operation()) {
+ case "delete" -> indexer.delete(project, Change.id(intent.changeId()));
+ case "index" -> indexer.index(project, Change.id(intent.changeId()));
+ default ->
+ logger.atSevere().log(
+ "Unknown operation '%s' in pending index intent: %s", intent.operation(), intent);
+ }
+ } catch (NoSuchChangeException e) {
+ // Ignore silently. change got deleted after intent.
+ } catch (Exception e) {
+ // catch all indexing exceptions to not propagate further.
+ logger.atSevere().withCause(e).log("Exception while recovering index intent: %s", intent);
+ }
+ Files.deleteIfExists(file);
+ }
+}
diff --git a/java/com/google/gerrit/server/index/change/PendingIndexUpdateScanner.java b/java/com/google/gerrit/server/index/change/PendingIndexUpdateScanner.java
new file mode 100644
index 0000000..0a1cd9b
--- /dev/null
+++ b/java/com/google/gerrit/server/index/change/PendingIndexUpdateScanner.java
@@ -0,0 +1,177 @@
+// Copyright (C) 2026 The Android Open Source Project
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package com.google.gerrit.server.index.change;
+
+import com.google.common.flogger.FluentLogger;
+import com.google.common.io.MoreFiles;
+import com.google.gerrit.extensions.events.LifecycleListener;
+import com.google.gerrit.lifecycle.LifecycleModule;
+import com.google.gerrit.server.config.GerritServerConfig;
+import com.google.gerrit.server.git.WorkQueue;
+import com.google.inject.Inject;
+import com.google.inject.Singleton;
+import java.io.IOException;
+import java.lang.management.ManagementFactory;
+import java.nio.file.DirectoryStream;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.time.Duration;
+import java.util.concurrent.TimeUnit;
+import org.eclipse.jgit.lib.Config;
+
+/** Background scanner that recovers change index updates missed due to a crash/interrupt. */
+@Singleton
+public final class PendingIndexUpdateScanner implements Runnable, LifecycleListener {
+ public static class Module extends LifecycleModule {
+ @Override
+ protected void configure() {
+ listener().to(PendingIndexUpdateScanner.class);
+ }
+ }
+
+ private static final FluentLogger logger = FluentLogger.forEnclosingClass();
+ private static final Duration DEFAULT_SCAN_INTERVAL = Duration.ofMinutes(5);
+
+ private final PendingIndexUpdate pendingIndexUpdate;
+ private final WorkQueue workQueue;
+ private final Duration scanInterval;
+
+ @Inject
+ PendingIndexUpdateScanner(
+ PendingIndexUpdate pendingIndexUpdate, WorkQueue workQueue, @GerritServerConfig Config cfg) {
+ this.pendingIndexUpdate = pendingIndexUpdate;
+ this.workQueue = workQueue;
+ this.scanInterval =
+ Duration.ofMillis(
+ cfg.getTimeUnit(
+ "index",
+ null,
+ "staleChangeRecoveryInterval",
+ DEFAULT_SCAN_INTERVAL.toMillis(),
+ TimeUnit.MILLISECONDS));
+ }
+
+ @Override
+ public void start() {
+ if (!pendingIndexUpdate.isEnabled()) {
+ return;
+ }
+
+ // Remove the in-process intents from previous crash.
+ try {
+ Path buildingDir = pendingIndexUpdate.buildingDir;
+ if (Files.exists(buildingDir)) {
+ MoreFiles.deleteRecursively(buildingDir);
+ }
+ } catch (IOException e) {
+ logger.atWarning().withCause(e).log("Unable to clean up building index directory");
+ }
+
+ // recover intents from previous crash.
+ var unused =
+ workQueue
+ .getDefaultQueue()
+ .submit(
+ () -> {
+ Path intentDir = pendingIndexUpdate.intentDir;
+ if (!Files.exists(intentDir)) {
+ // fresh install or feature newly enabled
+ return;
+ }
+
+ try (DirectoryStream<Path> pidDirs = Files.newDirectoryStream(intentDir)) {
+ for (Path pidDir : pidDirs) {
+ if (pendingIndexUpdate.runningDir.equals(pidDir)
+ || pendingIndexUpdate.buildingDir.equals(pidDir)) {
+ continue;
+ }
+
+ processPidDir(pidDir, true);
+ pendingIndexUpdate.cleanIfEmpty(pidDir);
+ }
+ } catch (Exception e) {
+ logger.atSevere().withCause(e).log(
+ "Unable to recover index intents from previous run");
+ }
+ });
+
+ unused =
+ workQueue
+ .getDefaultQueue()
+ .scheduleWithFixedDelay(
+ this, scanInterval.toMillis(), scanInterval.toMillis(), TimeUnit.MILLISECONDS);
+ }
+
+ @Override
+ public void stop() {}
+
+ @Override
+ public void run() {
+ try {
+ Path runningDir = pendingIndexUpdate.runningDir;
+ if (!Files.isDirectory(runningDir)) {
+ // no intents written yet.
+ return;
+ }
+
+ processPidDir(runningDir, false);
+ } catch (Exception e) {
+ // catch all to not disrupt next run.
+ logger.atSevere().withCause(e).log("Error in pending index intent run");
+ }
+ }
+
+ private void processPidDir(Path pidDir, boolean skipDeadCheck) {
+ try (DirectoryStream<Path> threadDirs = Files.newDirectoryStream(pidDir)) {
+ for (Path threadDir : threadDirs) {
+ long threadId;
+ try {
+ threadId = Long.parseLong(threadDir.getFileName().toString());
+ } catch (NumberFormatException e) {
+ logger.atWarning().log(
+ "Unexpected entry in pending index dir: %s; skipping", threadDir.getFileName());
+ MoreFiles.deleteRecursively(threadDir);
+ continue;
+ }
+ if (skipDeadCheck || isThreadDead(threadId)) {
+ processDeadThreadDir(threadDir);
+ }
+ }
+ } catch (IOException e) {
+ logger.atSevere().withCause(e).log("Failed to run pending index update scan");
+ }
+ }
+
+ private static boolean isThreadDead(long threadId) {
+ return ManagementFactory.getThreadMXBean().getThreadInfo(threadId) == null;
+ }
+
+ private void processDeadThreadDir(Path threadDir) {
+ try (DirectoryStream<Path> intents = Files.newDirectoryStream(threadDir)) {
+ for (Path intent : intents) {
+ try {
+ pendingIndexUpdate.recover(intent);
+ } catch (IOException e) {
+ logger.atWarning().withCause(e).log(
+ "Failed to recover pending index intent %s", intent.getFileName());
+ }
+ }
+ } catch (IOException e) {
+ logger.atWarning().withCause(e).log(
+ "Failed to recover pending index updates for %s", threadDir);
+ }
+ pendingIndexUpdate.cleanIfEmpty(threadDir);
+ }
+}
diff --git a/java/com/google/gerrit/server/update/BatchUpdate.java b/java/com/google/gerrit/server/update/BatchUpdate.java
index 76bf5ef..5e60aa9 100644
--- a/java/com/google/gerrit/server/update/BatchUpdate.java
+++ b/java/com/google/gerrit/server/update/BatchUpdate.java
@@ -63,6 +63,7 @@
import com.google.gerrit.server.git.GitRepositoryManager;
import com.google.gerrit.server.git.validators.OnSubmitValidators;
import com.google.gerrit.server.index.change.ChangeIndexer;
+import com.google.gerrit.server.index.change.PendingIndexUpdate;
import com.google.gerrit.server.logging.Metadata;
import com.google.gerrit.server.logging.RequestId;
import com.google.gerrit.server.logging.TraceContext;
@@ -668,6 +669,23 @@
}
return indexFutures.build();
}
+
+ void writeIndexIntents(PendingIndexUpdate pendingIndexUpdate, long threadId)
+ throws IOException {
+ for (Map.Entry<Change.Id, ChangeResult> e : results.entrySet()) {
+ if (e.getValue() == ChangeResult.SKIPPED) {
+ continue;
+ }
+ pendingIndexUpdate.write(
+ threadId, project, e.getKey(), e.getValue() == ChangeResult.DELETED);
+ }
+ }
+
+ void deleteIndexIntents(PendingIndexUpdate pendingIndexUpdate, long threadId) {
+ for (Map.Entry<Change.Id, ChangeResult> e : results.entrySet()) {
+ pendingIndexUpdate.delete(threadId, project, e.getKey());
+ }
+ }
}
ChangesHandle executeChangeOps(
diff --git a/java/com/google/gerrit/server/update/BatchUpdates.java b/java/com/google/gerrit/server/update/BatchUpdates.java
index aa727f1..a4060bc 100644
--- a/java/com/google/gerrit/server/update/BatchUpdates.java
+++ b/java/com/google/gerrit/server/update/BatchUpdates.java
@@ -31,6 +31,7 @@
import com.google.gerrit.extensions.restapi.ResourceConflictException;
import com.google.gerrit.extensions.restapi.ResourceNotFoundException;
import com.google.gerrit.extensions.restapi.RestApiException;
+import com.google.gerrit.server.index.change.PendingIndexUpdate;
import com.google.gerrit.server.notedb.LimitExceededException;
import com.google.gerrit.server.project.InvalidChangeOperationException;
import com.google.gerrit.server.project.NoSuchChangeException;
@@ -83,10 +84,12 @@
}
private final ChangeData.Factory changeDataFactory;
+ private final PendingIndexUpdate pendingIndexUpdate;
@Inject
- BatchUpdates(ChangeData.Factory changeDataFactory) {
+ BatchUpdates(ChangeData.Factory changeDataFactory, PendingIndexUpdate pendingIndexUpdate) {
this.changeDataFactory = changeDataFactory;
+ this.pendingIndexUpdate = pendingIndexUpdate;
}
@CanIgnoreReturnValue
@@ -100,9 +103,10 @@
checkDifferentProject(updates);
+ List<ListenableFuture<ChangeData>> indexFutures = new ArrayList<>();
+ List<ChangesHandle> changesHandles = new ArrayList<>(updates.size());
+ long threadId = Thread.currentThread().threadId();
try {
- List<ListenableFuture<ChangeData>> indexFutures = new ArrayList<>();
- List<ChangesHandle> changesHandles = new ArrayList<>(updates.size());
try {
for (BatchUpdate u : updates) {
u.executeUpdateRepo();
@@ -111,6 +115,11 @@
for (BatchUpdate u : updates) {
changesHandles.add(u.executeChangeOps(listeners, dryrun));
}
+ if (!dryrun && pendingIndexUpdate.isEnabled()) {
+ for (ChangesHandle h : changesHandles) {
+ h.writeIndexIntents(pendingIndexUpdate, threadId);
+ }
+ }
for (ChangesHandle h : changesHandles) {
h.execute();
if (h.requiresReindex()) {
@@ -137,6 +146,12 @@
updates.forEach(BatchUpdate::fireRefChangeEvents);
if (!dryrun) {
+ if (pendingIndexUpdate.isEnabled()) {
+ for (ChangesHandle h : changesHandles) {
+ h.deleteIndexIntents(pendingIndexUpdate, threadId);
+ }
+ }
+
for (BatchUpdate u : updates) {
u.executePostOps(changeDatas);
}
diff --git a/javatests/com/google/gerrit/acceptance/server/index/change/PendingIndexUpdateIT.java b/javatests/com/google/gerrit/acceptance/server/index/change/PendingIndexUpdateIT.java
new file mode 100644
index 0000000..78a5cd9
--- /dev/null
+++ b/javatests/com/google/gerrit/acceptance/server/index/change/PendingIndexUpdateIT.java
@@ -0,0 +1,79 @@
+// Copyright (C) 2026 The Android Open Source Project
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package com.google.gerrit.acceptance.server.index.change;
+
+import static com.google.common.truth.Truth.assertThat;
+import static com.google.gerrit.acceptance.WaitUtil.waitUntil;
+
+import com.google.gerrit.acceptance.AbstractDaemonTest;
+import com.google.gerrit.acceptance.PushOneCommit;
+import com.google.gerrit.acceptance.config.GerritConfig;
+import com.google.gerrit.entities.Change;
+import com.google.gerrit.server.index.change.PendingIndexUpdate;
+import com.google.inject.Inject;
+import java.time.Duration;
+import org.junit.Test;
+
+public class PendingIndexUpdateIT extends AbstractDaemonTest {
+ private static final long DEAD_THREAD_ID = Long.MAX_VALUE;
+ @Inject private PendingIndexUpdate pendingIndexUpdate;
+
+ @Test
+ @GerritConfig(name = "index.staleChangeRecovery", value = "true")
+ @GerritConfig(name = "index.staleChangeRecoveryInterval", value = "1s")
+ @GerritConfig(name = "index.changes.commitWithin", value = "0")
+ public void scannerRecoversMissedIndexWrite() throws Exception {
+ PushOneCommit.Result r = createChange();
+ Change.Id changeId = r.getChange().getId();
+
+ // Simulate a crash: the change is in NoteDb but the index write was missed.
+ indexer.delete(project, changeId);
+ pendingIndexUpdate.write(DEAD_THREAD_ID, project, changeId, /* delete= */ false);
+ assertThat(gApi.changes().query("change:" + changeId).get()).isEmpty();
+
+ waitUntil(
+ () -> {
+ try {
+ return gApi.changes().query("change:" + changeId).get().size() == 1;
+ } catch (Exception e) {
+ return false;
+ }
+ },
+ Duration.ofSeconds(5));
+ assertThat(gApi.changes().query("change:" + changeId).get()).hasSize(1);
+ }
+
+ @Test
+ @GerritConfig(name = "index.staleChangeRecovery", value = "true")
+ @GerritConfig(name = "index.staleChangeRecoveryInterval", value = "1s")
+ @GerritConfig(name = "index.changes.commitWithin", value = "0")
+ public void scannerRecoversMissedIndexDelete() throws Exception {
+ PushOneCommit.Result r = createChange();
+ Change.Id changeId = r.getChange().getId();
+ pendingIndexUpdate.write(DEAD_THREAD_ID, project, changeId, /* delete= */ true);
+
+ assertThat(gApi.changes().query("change:" + changeId).get()).hasSize(1);
+ waitUntil(
+ () -> {
+ try {
+ return gApi.changes().query("change:" + changeId).get().isEmpty();
+ } catch (Exception e) {
+ return false;
+ }
+ },
+ Duration.ofSeconds(5));
+ assertThat(gApi.changes().query("change:" + changeId).get()).isEmpty();
+ }
+}
diff --git a/javatests/com/google/gerrit/server/index/change/PendingIndexUpdateScannerTest.java b/javatests/com/google/gerrit/server/index/change/PendingIndexUpdateScannerTest.java
new file mode 100644
index 0000000..4a23d46
--- /dev/null
+++ b/javatests/com/google/gerrit/server/index/change/PendingIndexUpdateScannerTest.java
@@ -0,0 +1,169 @@
+// Copyright (C) 2026 The Android Open Source Project
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package com.google.gerrit.server.index.change;
+
+import static com.google.common.truth.Truth.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.lenient;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import com.google.gerrit.entities.Change;
+import com.google.gerrit.entities.Project;
+import com.google.gerrit.server.config.SitePaths;
+import com.google.gerrit.server.git.WorkQueue;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.concurrent.ScheduledExecutorService;
+import org.eclipse.jgit.lib.Config;
+import org.junit.Before;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+import org.mockito.Mock;
+import org.mockito.junit.MockitoJUnit;
+import org.mockito.junit.MockitoRule;
+
+public class PendingIndexUpdateScannerTest {
+ private static final long DEAD_THREAD_ID = Long.MAX_VALUE;
+ private static final Project.NameKey PROJECT = Project.nameKey("test-project");
+ private static final Change.Id CHANGE_ID = Change.id(42);
+
+ @Rule public final MockitoRule mockito = MockitoJUnit.rule();
+ @Rule public final TemporaryFolder tempDir = new TemporaryFolder();
+
+ @Mock private ChangeIndexer indexer;
+ @Mock private WorkQueue workQueue;
+ @Mock private ScheduledExecutorService fakeQueue;
+
+ private SitePaths sitePaths;
+ private PendingIndexUpdate pendingIndexUpdate;
+ private PendingIndexUpdateScanner scanner;
+
+ @Before
+ public void setUp() throws Exception {
+ lenient().doNothing().when(indexer).index(any(), any());
+ lenient().doNothing().when(indexer).delete(any(), any());
+ // Run submitted tasks synchronously so startup recovery completes inline.
+ doAnswer(
+ inv -> {
+ ((Runnable) inv.getArgument(0)).run();
+ return null;
+ })
+ .when(fakeQueue)
+ .submit(any(Runnable.class));
+ when(workQueue.getDefaultQueue()).thenReturn(fakeQueue);
+ sitePaths = new SitePaths(tempDir.getRoot().toPath());
+ pendingIndexUpdate = new PendingIndexUpdate(sitePaths, indexer, recoveryConfig());
+ scanner = new PendingIndexUpdateScanner(pendingIndexUpdate, workQueue, recoveryConfig());
+ }
+
+ @Test
+ public void scannerIndexesStaleFile() throws Exception {
+ pendingIndexUpdate.write(DEAD_THREAD_ID, PROJECT, CHANGE_ID, /* delete= */ false);
+
+ scanner.run();
+
+ verify(indexer).index(PROJECT, CHANGE_ID);
+ assertThat(intentFile(DEAD_THREAD_ID, PROJECT, CHANGE_ID).toFile().exists()).isFalse();
+ }
+
+ @Test
+ public void scannerDeletesChangeWhenOperationIsDelete() throws Exception {
+ pendingIndexUpdate.write(DEAD_THREAD_ID, PROJECT, CHANGE_ID, /* delete= */ true);
+
+ scanner.run();
+
+ verify(indexer).delete(PROJECT, CHANGE_ID);
+ assertThat(intentFile(DEAD_THREAD_ID, PROJECT, CHANGE_ID).toFile().exists()).isFalse();
+ }
+
+ @Test
+ public void scannerSkipsIntentsForLiveThread() throws Exception {
+ long liveThreadId = Thread.currentThread().threadId();
+ pendingIndexUpdate.write(liveThreadId, PROJECT, CHANGE_ID, /* delete= */ false);
+
+ scanner.run();
+
+ verify(indexer, never()).index(any(), any());
+ assertThat(intentFile(liveThreadId, PROJECT, CHANGE_ID).toFile().exists()).isTrue();
+ }
+
+ @Test
+ public void scannerDeletesMalformedFile() throws Exception {
+ Path file = intentFile(DEAD_THREAD_ID, PROJECT, CHANGE_ID);
+ Files.createDirectories(file.getParent());
+ Files.writeString(file, "not-a-valid-blob");
+
+ scanner.run();
+
+ verify(indexer, never()).index(any(), any());
+ verify(indexer, never()).delete(any(), any());
+ assertThat(file.toFile().exists()).isFalse();
+ }
+
+ @Test
+ public void scannerDoesNothingWhenNoPendingFiles() throws Exception {
+ scanner.run();
+
+ verify(indexer, never()).index(any(), any());
+ verify(indexer, never()).delete(any(), any());
+ }
+
+ @Test
+ public void startRecoversPreviousProcessIntents() throws Exception {
+ // Write an intent under a foreign process marker dir, simulating a previous crash.
+ Path intentDir = sitePaths.data_dir.resolve("pending-index");
+ Path prevThreadDir = intentDir.resolve("99999_1234567890000").resolve("1");
+ Files.createDirectories(prevThreadDir);
+ Files.writeString(
+ prevThreadDir.resolve(pendingIndexUpdate.filename(PROJECT, CHANGE_ID)),
+ "{\"project\":\"test-project\",\"changeId\":42,\"operation\":\"index\"}");
+
+ scanner = new PendingIndexUpdateScanner(pendingIndexUpdate, workQueue, recoveryConfig());
+ scanner.start();
+
+ verify(indexer).index(PROJECT, CHANGE_ID);
+ assertThat(prevThreadDir.toFile().exists()).isFalse();
+ }
+
+ @Test
+ public void startCleansBuildingDir() throws Exception {
+ // Leave an orphaned temp file in buildingDir as if a crash happened mid-write.
+ Path buildingDir = sitePaths.data_dir.resolve("pending-index").resolve("building");
+ Files.createDirectories(buildingDir);
+ Path orphan = Files.createTempFile(buildingDir, null, null);
+
+ scanner = new PendingIndexUpdateScanner(pendingIndexUpdate, workQueue, recoveryConfig());
+ scanner.start();
+
+ assertThat(orphan.toFile().exists()).isFalse();
+ }
+
+ private static Config recoveryConfig() {
+ Config cfg = new Config();
+ cfg.setBoolean("index", null, "staleChangeRecovery", true);
+ cfg.setInt("index", "changes", "commitWithin", 0);
+ return cfg;
+ }
+
+ private Path intentFile(long threadId, Project.NameKey project, Change.Id changeId) {
+ return pendingIndexUpdate
+ .threadDir(threadId)
+ .resolve(pendingIndexUpdate.filename(project, changeId));
+ }
+}
diff --git a/javatests/com/google/gerrit/server/update/BatchUpdateIndexIntentTest.java b/javatests/com/google/gerrit/server/update/BatchUpdateIndexIntentTest.java
new file mode 100644
index 0000000..2cbff2b
--- /dev/null
+++ b/javatests/com/google/gerrit/server/update/BatchUpdateIndexIntentTest.java
@@ -0,0 +1,144 @@
+// Copyright (C) 2026 The Android Open Source Project
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package com.google.gerrit.server.update;
+
+import static com.google.common.truth.Truth.assertThat;
+import static com.google.gerrit.testing.TestActionRefUpdateContext.openTestRefUpdateContext;
+
+import com.google.gerrit.entities.Change;
+import com.google.gerrit.entities.Project;
+import com.google.gerrit.server.CurrentUser;
+import com.google.gerrit.server.Sequences;
+import com.google.gerrit.server.change.ChangeInserter;
+import com.google.gerrit.server.config.SitePaths;
+import com.google.gerrit.server.git.GitRepositoryManager;
+import com.google.gerrit.server.index.change.PendingIndexUpdate;
+import com.google.gerrit.server.update.context.RefUpdateContext;
+import com.google.gerrit.server.util.time.TimeUtil;
+import com.google.gerrit.testing.InMemoryTestEnvironment;
+import com.google.inject.Inject;
+import com.google.inject.Provider;
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.concurrent.atomic.AtomicBoolean;
+import org.eclipse.jgit.junit.TestRepository;
+import org.eclipse.jgit.lib.Config;
+import org.eclipse.jgit.lib.Repository;
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Rule;
+import org.junit.Test;
+
+/** Tests for the pending-index lifecycle in {@link BatchUpdate} and {@link BatchUpdates}. */
+public class BatchUpdateIndexIntentTest {
+ @Rule
+ public InMemoryTestEnvironment testEnvironment =
+ new InMemoryTestEnvironment(
+ () -> {
+ Config cfg = new Config();
+ cfg.setString("index", null, "type", "fake");
+ cfg.setBoolean("index", null, "staleChangeRecovery", true);
+ return cfg;
+ });
+
+ @Inject private BatchUpdate.Factory batchUpdateFactory;
+ @Inject private ChangeInserter.Factory changeInserterFactory;
+ @Inject private GitRepositoryManager repoManager;
+ @Inject private Provider<CurrentUser> user;
+ @Inject private Sequences sequences;
+ @Inject private SitePaths sitePaths;
+ @Inject private PendingIndexUpdate pendingIndexUpdate;
+
+ private Project.NameKey project;
+ private TestRepository<Repository> repo;
+ private RefUpdateContext testRefUpdateContext;
+
+ @Before
+ public void setUp() throws Exception {
+ project = Project.nameKey("test");
+ repo = new TestRepository<>(repoManager.createRepository(project));
+ testRefUpdateContext = openTestRefUpdateContext();
+ }
+
+ @After
+ public void tearDown() {
+ testRefUpdateContext.close();
+ }
+
+ @Test
+ public void pendingIndexIntentFilePresentDuringUpdate() throws Exception {
+ Change.Id id = createChange();
+ AtomicBoolean intentFound = new AtomicBoolean(false);
+
+ BatchUpdateListener listener =
+ new BatchUpdateListener() {
+ @Override
+ public void afterUpdateRefs() throws Exception {
+ intentFound.set(hasPendingIntentFile(id));
+ }
+ };
+
+ try (BatchUpdate bu = batchUpdateFactory.create(project, user.get(), TimeUtil.now())) {
+ bu.addOp(id, addMessageOp("Pending intent test"));
+ bu.execute(listener);
+ }
+
+ assertThat(intentFound.get()).isTrue();
+ }
+
+ @Test
+ public void pendingIndexIntentFilesRemovedAfterSuccessfulUpdate() throws Exception {
+ Change.Id id = createChange();
+
+ try (BatchUpdate bu = batchUpdateFactory.create(project, user.get(), TimeUtil.now())) {
+ bu.addOp(id, addMessageOp("Cleanup test"));
+ bu.execute();
+ }
+
+ assertThat(hasPendingIntentFile(id)).isFalse();
+ }
+
+ private boolean hasPendingIntentFile(Change.Id id) throws IOException {
+ Path intentDir = sitePaths.data_dir.resolve("pending-index");
+ String expectedFilename = pendingIndexUpdate.filename(project, id);
+ try (var stream = Files.walk(intentDir)) {
+ return stream
+ .filter(Files::isRegularFile)
+ .anyMatch(p -> p.getFileName().toString().equals(expectedFilename));
+ }
+ }
+
+ private Change.Id createChange() throws Exception {
+ Change.Id id = Change.id(sequences.nextChangeId());
+ try (BatchUpdate bu = batchUpdateFactory.create(project, user.get(), TimeUtil.now())) {
+ bu.insertChange(
+ changeInserterFactory.create(
+ id, repo.commit().message("Change").insertChangeId().create(), "refs/heads/master"));
+ bu.execute();
+ }
+ return id;
+ }
+
+ private static BatchUpdateOp addMessageOp(String message) {
+ return new BatchUpdateOp() {
+ @Override
+ public boolean updateChange(ChangeContext ctx) {
+ ctx.getUpdate(ctx.getChange().currentPatchSetId()).setChangeMessage(message);
+ return true;
+ }
+ };
+ }
+}