diff --git a/build.xml b/build.xml
index bee7d97decc2f642bcc4d6b8655d6f2f00148cc8..f9d0fb735bd9507e627b54ac1d955824c655c0be 100644
--- a/build.xml
+++ b/build.xml
@@ -12,7 +12,7 @@
-
+
diff --git a/src/build b/src/build
index 665ab024daba01df34cdf6ada5f1c5666102c351..fc78d3d7be9c14352c2a5311b7f1edc4b5f215f8 160000
--- a/src/build
+++ b/src/build
@@ -1 +1 @@
-Subproject commit 665ab024daba01df34cdf6ada5f1c5666102c351
+Subproject commit fc78d3d7be9c14352c2a5311b7f1edc4b5f215f8
diff --git a/src/main/java/org/torproject/metrics/onionoo/docs/DocumentStore.java b/src/main/java/org/torproject/metrics/onionoo/docs/DocumentStore.java
index b1404224a4f985dea09b9650ed32107e84f62389..57373b507eb2b8176e8a7b4bcf9c6a24e8d92da2 100644
--- a/src/main/java/org/torproject/metrics/onionoo/docs/DocumentStore.java
+++ b/src/main/java/org/torproject/metrics/onionoo/docs/DocumentStore.java
@@ -22,8 +22,13 @@ import java.io.FileOutputStream;
import java.io.FileReader;
import java.io.IOException;
import java.io.InputStream;
+import java.io.OutputStream;
+import java.nio.channels.FileChannel;
import java.nio.charset.StandardCharsets;
+import java.nio.file.AtomicMoveNotSupportedException;
import java.nio.file.Files;
+import java.nio.file.StandardCopyOption;
+import java.nio.file.StandardOpenOption;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
@@ -346,11 +351,7 @@ public class DocumentStore {
}
}
}
- File documentTempFile = new File(
- documentFile.getAbsolutePath() + ".tmp");
- writeToFile(documentTempFile, documentString);
- documentFile.delete();
- documentTempFile.renameTo(documentFile);
+ this.writeToFileAtomic(documentFile, documentString);
this.storedFiles++;
this.storedBytes += documentString.length();
} catch (IOException e) {
@@ -766,7 +767,7 @@ public class DocumentStore {
String documentString = sb.toString();
try {
summaryFile.getParentFile().mkdirs();
- writeToFile(summaryFile, documentString);
+ this.writeToFileAtomic(summaryFile, documentString);
this.lastModifiedNodeStatuses = summaryFile.lastModified();
this.updatedNodeStatuses.clear();
this.storedFiles++;
@@ -777,12 +778,42 @@ public class DocumentStore {
}
}
- private static void writeToFile(File file, String content)
- throws IOException {
- try (BufferedOutputStream bos = new BufferedOutputStream(
- new FileOutputStream(file))) {
+ /** Atomically replaces {@code file} with {@code content}: writes the bytes
+ * to a sibling {@code .tmp}, fsyncs the file descriptor, then renames the
+ * tmp into place with {@link StandardCopyOption#ATOMIC_MOVE} and best-effort
+ * fsyncs the directory. Guarantees that any reader sees either the old or
+ * the new contents in full, never an empty or truncated file, even if the
+ * JVM exits mid-write. */
+ void writeToFileAtomic(File file, String content) throws IOException {
+ File parent = file.getParentFile();
+ File tmp = new File(parent, file.getName() + ".tmp");
+ try (OutputStream raw = this.openTempFileOutputStream(tmp);
+ BufferedOutputStream bos = new BufferedOutputStream(raw)) {
bos.write(content.getBytes(StandardCharsets.UTF_8));
+ bos.flush();
+ if (raw instanceof FileOutputStream) {
+ ((FileOutputStream) raw).getFD().sync();
+ }
+ }
+ try {
+ Files.move(tmp.toPath(), file.toPath(),
+ StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE);
+ } catch (AtomicMoveNotSupportedException e) {
+ Files.move(tmp.toPath(), file.toPath(),
+ StandardCopyOption.REPLACE_EXISTING);
}
+ try (FileChannel dir = FileChannel.open(parent.toPath(),
+ StandardOpenOption.READ)) {
+ dir.force(true);
+ } catch (IOException ignored) {
+ /* Directory fsync is best-effort; not supported on all filesystems. */
+ }
+ }
+
+ /** Hook for tests to inject a failing {@code OutputStream}. Production
+ * callers get a plain {@link FileOutputStream}. */
+ protected OutputStream openTempFileOutputStream(File tmp) throws IOException {
+ return new FileOutputStream(tmp);
}
private void writeSummaryDocuments() {
@@ -810,7 +841,7 @@ public class DocumentStore {
File summaryFile = new File(this.outDir, "summary");
try {
summaryFile.getParentFile().mkdirs();
- writeToFile(summaryFile, documentString);
+ this.writeToFileAtomic(summaryFile, documentString);
this.lastModifiedSummaryDocuments = summaryFile.lastModified();
this.updatedSummaryDocuments.clear();
this.storedFiles++;
diff --git a/src/main/java/org/torproject/metrics/onionoo/updater/NodeDetailsStatusUpdater.java b/src/main/java/org/torproject/metrics/onionoo/updater/NodeDetailsStatusUpdater.java
index 6516bf2e98336f6bafd02bc4871fa1d53b6f6fb9..c1f1625e30bac75ceb99e117bde94474a912b57b 100644
--- a/src/main/java/org/torproject/metrics/onionoo/updater/NodeDetailsStatusUpdater.java
+++ b/src/main/java/org/torproject/metrics/onionoo/updater/NodeDetailsStatusUpdater.java
@@ -14,11 +14,16 @@ import org.torproject.descriptor.Microdescriptor;
import org.torproject.descriptor.NetworkStatusEntry;
import org.torproject.descriptor.RelayNetworkStatusConsensus;
import org.torproject.descriptor.ServerDescriptor;
+import org.torproject.metrics.onionoo.docs.BandwidthStatus;
+import org.torproject.metrics.onionoo.docs.ClientsStatus;
import org.torproject.metrics.onionoo.docs.DateTimeHelper;
import org.torproject.metrics.onionoo.docs.DetailsStatus;
import org.torproject.metrics.onionoo.docs.DocumentStore;
import org.torproject.metrics.onionoo.docs.DocumentStoreFactory;
import org.torproject.metrics.onionoo.docs.NodeStatus;
+import org.torproject.metrics.onionoo.docs.SummaryDocument;
+import org.torproject.metrics.onionoo.docs.UptimeStatus;
+import org.torproject.metrics.onionoo.docs.WeightsStatus;
import org.torproject.metrics.onionoo.util.FormattingUtils;
import org.slf4j.Logger;
@@ -567,6 +572,48 @@ public class NodeDetailsStatusUpdater implements DescriptorListener,
logger.info("Finished reverse domain name lookups");
this.updateNodeDetailsStatuses();
logger.info("Updated node and details statuses");
+ this.removeStaleNodes();
+ logger.info("Removed stale nodes");
+ }
+
+ /** Removes per-fingerprint status and document files (and the
+ * corresponding node status row in the monolithic summary) for any
+ * relay or bridge whose last appearance is older than one week before
+ * the most recent consensus or bridge network status processed. This
+ * matches the {@code /summary} retention policy already enforced by
+ * {@code SummaryDocumentWriter} and prevents indefinite accumulation of
+ * stale per-fingerprint files (the asymmetry that turned the #40028
+ * one-time corruption into a permanent visible symptom). */
+ private void removeStaleNodes() {
+ long lastValidAfter = Math.max(this.relaysLastValidAfterMillis,
+ this.bridgesLastPublishedMillis);
+ if (lastValidAfter <= 0L) {
+ return;
+ }
+ long cutoff = lastValidAfter - DateTimeHelper.ONE_WEEK;
+ SortedSet fingerprints = new TreeSet<>(
+ this.documentStore.list(NodeStatus.class));
+ int removed = 0;
+ for (String fp : fingerprints) {
+ NodeStatus ns = this.documentStore.retrieve(
+ NodeStatus.class, true, fp);
+ if (ns != null && ns.getLastSeenMillis() >= cutoff) {
+ continue;
+ }
+ this.documentStore.remove(NodeStatus.class, fp);
+ this.documentStore.remove(SummaryDocument.class, fp);
+ this.documentStore.remove(DetailsStatus.class, fp);
+ this.documentStore.remove(BandwidthStatus.class, fp);
+ this.documentStore.remove(WeightsStatus.class, fp);
+ this.documentStore.remove(ClientsStatus.class, fp);
+ this.documentStore.remove(UptimeStatus.class, fp);
+ removed++;
+ }
+ if (removed > 0) {
+ logger.info("Removed {} stale node statuses and their per-fingerprint "
+ + "documents (cutoff: {}).", removed,
+ DateTimeHelper.format(cutoff));
+ }
}
/* Step 2: read node statuses from disk. */
@@ -1191,7 +1238,10 @@ public class NodeDetailsStatusUpdater implements DescriptorListener,
detailsStatus.setAddress(nodeStatus.getAddress());
detailsStatus.setOrAddressesAndPorts(
nodeStatus.getOrAddressesAndPorts());
- detailsStatus.setFirstSeenMillis(nodeStatus.getFirstSeenMillis());
+ long nsFirstSeen = nodeStatus.getFirstSeenMillis();
+ if (nsFirstSeen > 0L) {
+ detailsStatus.setFirstSeenMillis(nsFirstSeen);
+ }
detailsStatus.setLastSeenMillis(nodeStatus.getLastSeenMillis());
detailsStatus.setOrPort(nodeStatus.getOrPort());
detailsStatus.setDirPort(nodeStatus.getDirPort());
diff --git a/src/test/java/org/torproject/metrics/onionoo/docs/DocumentStoreTest.java b/src/test/java/org/torproject/metrics/onionoo/docs/DocumentStoreTest.java
new file mode 100644
index 0000000000000000000000000000000000000000..84040dce49286f8b8a13f7cc158d91a0d287bcf6
--- /dev/null
+++ b/src/test/java/org/torproject/metrics/onionoo/docs/DocumentStoreTest.java
@@ -0,0 +1,85 @@
+/* Copyright 2026 The Tor Project
+ * See LICENSE for licensing information */
+
+package org.torproject.metrics.onionoo.docs;
+
+import static org.junit.Assert.assertArrayEquals;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.fail;
+
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+
+import java.io.File;
+import java.io.IOException;
+import java.io.OutputStream;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+
+public class DocumentStoreTest {
+
+ @Rule
+ public TemporaryFolder folder = new TemporaryFolder();
+
+ @Test
+ public void writeToFileAtomicRoundTrip() throws IOException {
+ DocumentStore store = new DocumentStore();
+ File target = new File(folder.getRoot(), "summary");
+ store.writeToFileAtomic(target, "v1\n");
+ assertArrayEquals("v1\n".getBytes(StandardCharsets.UTF_8),
+ Files.readAllBytes(target.toPath()));
+ store.writeToFileAtomic(target, "v2\n");
+ assertArrayEquals("v2\n".getBytes(StandardCharsets.UTF_8),
+ Files.readAllBytes(target.toPath()));
+ }
+
+ @Test
+ public void writeToFileAtomicLeavesNoTmpOnSuccess() throws IOException {
+ DocumentStore store = new DocumentStore();
+ File target = new File(folder.getRoot(), "summary");
+ store.writeToFileAtomic(target, "v1\n");
+ File tmp = new File(folder.getRoot(), "summary.tmp");
+ assertEquals("No .tmp should remain after a successful atomic write.",
+ false, tmp.exists());
+ }
+
+ @Test
+ public void writeToFileAtomicPreservesPreviousOnWriteFailure()
+ throws IOException {
+ File target = new File(folder.getRoot(), "summary");
+ DocumentStore good = new DocumentStore();
+ good.writeToFileAtomic(target, "v1-original\n");
+ byte[] originalBytes = Files.readAllBytes(target.toPath());
+
+ DocumentStore failing = new FailingDocumentStore();
+ try {
+ failing.writeToFileAtomic(target, "v2-new-content\n");
+ fail("Expected IOException to propagate from failing stream.");
+ } catch (IOException expected) {
+ /* Expected. */
+ }
+
+ assertArrayEquals(
+ "Original file contents must survive a mid-write failure.",
+ originalBytes, Files.readAllBytes(target.toPath()));
+ }
+
+ private static class FailingDocumentStore extends DocumentStore {
+ @Override
+ protected OutputStream openTempFileOutputStream(File tmp)
+ throws IOException {
+ return new OutputStream() {
+ @Override
+ public void write(int singleByte) throws IOException {
+ throw new IOException("simulated write failure");
+ }
+
+ @Override
+ public void write(byte[] bytes, int off, int len) throws IOException {
+ throw new IOException("simulated write failure");
+ }
+ };
+ }
+ }
+}
diff --git a/src/test/java/org/torproject/metrics/onionoo/updater/NodeDetailsStatusUpdaterTest.java b/src/test/java/org/torproject/metrics/onionoo/updater/NodeDetailsStatusUpdaterTest.java
index c950b54f3628a1d325418977b587e052b1255342..909c30449da6be4c5b95df4e1f4a98eeb6608e9f 100644
--- a/src/test/java/org/torproject/metrics/onionoo/updater/NodeDetailsStatusUpdaterTest.java
+++ b/src/test/java/org/torproject/metrics/onionoo/updater/NodeDetailsStatusUpdaterTest.java
@@ -7,6 +7,8 @@ import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
+import static org.torproject.metrics.onionoo.docs.DateTimeHelper.ONE_DAY;
+import static org.torproject.metrics.onionoo.docs.DateTimeHelper.ONE_WEEK;
import org.torproject.descriptor.BridgePoolAssignment;
import org.torproject.descriptor.Descriptor;
@@ -14,10 +16,15 @@ import org.torproject.descriptor.DescriptorParser;
import org.torproject.descriptor.DescriptorSourceFactory;
import org.torproject.descriptor.Microdescriptor;
import org.torproject.descriptor.ServerDescriptor;
+import org.torproject.metrics.onionoo.docs.BandwidthStatus;
+import org.torproject.metrics.onionoo.docs.ClientsStatus;
+import org.torproject.metrics.onionoo.docs.DateTimeHelper;
import org.torproject.metrics.onionoo.docs.DetailsStatus;
import org.torproject.metrics.onionoo.docs.DocumentStoreFactory;
import org.torproject.metrics.onionoo.docs.DummyDocumentStore;
import org.torproject.metrics.onionoo.docs.NodeStatus;
+import org.torproject.metrics.onionoo.docs.UptimeStatus;
+import org.torproject.metrics.onionoo.docs.WeightsStatus;
import org.junit.Before;
import org.junit.Test;
@@ -451,6 +458,169 @@ public class NodeDetailsStatusUpdaterTest {
assertTrue(nsC.getExtendedFamily().contains(FP_B));
}
+ /**
+ * Populate this.knownNodes on a fresh updater and invoke the
+ * private updateNodeDetailsStatuses().
+ */
+ private NodeDetailsStatusUpdater invokeUpdateNodeDetailsStatuses(
+ String fingerprint, NodeStatus nodeStatus) throws Exception {
+ NodeDetailsStatusUpdater ndsu = new NodeDetailsStatusUpdater(null, null);
+ SortedMap knownNodes = new TreeMap<>();
+ knownNodes.put(fingerprint, nodeStatus);
+ setField(ndsu, "knownNodes", knownNodes);
+
+ Method method = NodeDetailsStatusUpdater.class.getDeclaredMethod(
+ "updateNodeDetailsStatuses");
+ method.setAccessible(true);
+ method.invoke(ndsu);
+ return ndsu;
+ }
+
+ /**
+ * A minimal NodeStatus builder.
+ */
+ private NodeStatus newRelayNodeStatus(String fingerprint,
+ long firstSeenMillis) {
+ NodeStatus ns = new NodeStatus(fingerprint);
+ ns.setRelay(true);
+ ns.setRelayFlags(new TreeSet<>());
+ ns.setFirstSeenMillis(firstSeenMillis);
+ ns.setLastSeenMillis(firstSeenMillis);
+ ns.setOrAddressesAndPorts(new TreeSet<>());
+ return ns;
+ }
+
+ @Test
+ public void testFirstSeenMillisOverwriteWhenNonZero() throws Exception {
+ /**
+ * Non-zero NodeStatus.firstSeenMillis
+ * overwrites whatever DetailsStatus had.
+ */
+ DetailsStatus existing = new DetailsStatus();
+ existing.setFirstSeenMillis(9999L);
+ this.docStore.addDocument(existing, FP);
+
+ NodeStatus ns = newRelayNodeStatus(FP, 1234567L);
+ invokeUpdateNodeDetailsStatuses(FP, ns);
+
+ DetailsStatus stored = this.docStore.getDocument(DetailsStatus.class, FP);
+ assertNotNull(stored);
+ assertEquals(1234567L, stored.getFirstSeenMillis());
+ }
+
+ @Test
+ public void testFirstSeenMillisPreservedWhenNodeStatusIsZero()
+ throws Exception {
+ /**
+ * NodeStatus.firstSeenMillis = 0L must NOT clobber
+ * a historically-correct DetailsStatus.firstSeenMillis.
+ */
+ DetailsStatus existing = new DetailsStatus();
+ existing.setFirstSeenMillis(9999L);
+ this.docStore.addDocument(existing, FP);
+
+ NodeStatus ns = newRelayNodeStatus(FP, 0L);
+ invokeUpdateNodeDetailsStatuses(FP, ns);
+
+ DetailsStatus stored = this.docStore.getDocument(DetailsStatus.class, FP);
+ assertNotNull(stored);
+ assertEquals("DetailsStatus.firstSeenMillis must be preserved when "
+ + "NodeStatus carries 0L.", 9999L, stored.getFirstSeenMillis());
+ }
+
+ @Test
+ public void testFirstSeenMillisStaysZeroWhenNoPriorDetailsStatus()
+ throws Exception {
+ /* If there's nothing on disk and NodeStatus carries 0L, the freshly-
+ * created DetailsStatus also starts at 0L.
+ */
+ NodeStatus ns = newRelayNodeStatus(FP, 0L);
+ invokeUpdateNodeDetailsStatuses(FP, ns);
+
+ DetailsStatus stored = this.docStore.getDocument(DetailsStatus.class, FP);
+ assertNotNull(stored);
+ assertEquals(0L, stored.getFirstSeenMillis());
+ }
+
+ /** Invoke the private removeStaleNodes() against the supplied updater
+ * after setting relaysLastValidAfterMillis to {@code lastValidAfter}. */
+ private void invokeRemoveStaleNodes(NodeDetailsStatusUpdater ndsu,
+ long lastValidAfter) throws Exception {
+ setField(ndsu, "relaysLastValidAfterMillis", lastValidAfter);
+ Method method = NodeDetailsStatusUpdater.class.getDeclaredMethod(
+ "removeStaleNodes");
+ method.setAccessible(true);
+ method.invoke(ndsu);
+ }
+
+ @Test
+ public void testRemoveStaleNodesPrunesStaleAndKeepsFresh() throws Exception {
+ long now = DateTimeHelper.parse("2026-05-19 12:00:00");
+ long fresh = now - ONE_WEEK / 2;
+ long stale = now - 2 * ONE_WEEK;
+
+ NodeStatus nsFresh = newRelayNodeStatus(FP_A, fresh);
+ NodeStatus nsStale = newRelayNodeStatus(FP_B, stale);
+ this.docStore.addDocument(nsFresh, FP_A);
+ this.docStore.addDocument(nsStale, FP_B);
+ DetailsStatus dsStale = new DetailsStatus();
+ this.docStore.addDocument(dsStale, FP_B);
+
+ NodeDetailsStatusUpdater ndsu =
+ new NodeDetailsStatusUpdater(null, null);
+ invokeRemoveStaleNodes(ndsu, now);
+
+ assertNotNull("Fresh NodeStatus survives.",
+ this.docStore.getDocument(NodeStatus.class, FP_A));
+ assertNull("Stale NodeStatus is removed.",
+ this.docStore.getDocument(NodeStatus.class, FP_B));
+ assertNull("Stale DetailsStatus is removed.",
+ this.docStore.getDocument(DetailsStatus.class, FP_B));
+ }
+
+ @Test
+ public void testRemoveStaleNodesIsNoOpWhenNoConsensusYet()
+ throws Exception {
+ NodeStatus ns = newRelayNodeStatus(FP_A, 0L);
+ this.docStore.addDocument(ns, FP_A);
+
+ NodeDetailsStatusUpdater ndsu =
+ new NodeDetailsStatusUpdater(null, null);
+ /* relaysLastValidAfterMillis stays at its default (-1); cutoff must be
+ * skipped to avoid wiping everything when we haven't yet seen a
+ * consensus this run. */
+ invokeRemoveStaleNodes(ndsu, -1L);
+
+ assertNotNull("With no consensus seen, removeStaleNodes is a no-op.",
+ this.docStore.getDocument(NodeStatus.class, FP_A));
+ }
+
+ @Test
+ public void testRemoveStaleNodesRemovesAllPerFingerprintFiles()
+ throws Exception {
+ long now = DateTimeHelper.parse("2026-05-19 12:00:00");
+ long stale = now - ONE_WEEK - ONE_DAY;
+
+ NodeStatus nsStale = newRelayNodeStatus(FP_A, stale);
+ this.docStore.addDocument(nsStale, FP_A);
+ this.docStore.addDocument(new DetailsStatus(), FP_A);
+ this.docStore.addDocument(new BandwidthStatus(), FP_A);
+ this.docStore.addDocument(new WeightsStatus(), FP_A);
+ this.docStore.addDocument(new ClientsStatus(), FP_A);
+ this.docStore.addDocument(new UptimeStatus(), FP_A);
+
+ NodeDetailsStatusUpdater ndsu =
+ new NodeDetailsStatusUpdater(null, null);
+ invokeRemoveStaleNodes(ndsu, now);
+
+ assertNull(this.docStore.getDocument(NodeStatus.class, FP_A));
+ assertNull(this.docStore.getDocument(DetailsStatus.class, FP_A));
+ assertNull(this.docStore.getDocument(BandwidthStatus.class, FP_A));
+ assertNull(this.docStore.getDocument(WeightsStatus.class, FP_A));
+ assertNull(this.docStore.getDocument(ClientsStatus.class, FP_A));
+ assertNull(this.docStore.getDocument(UptimeStatus.class, FP_A));
+ }
+
@Test
public void testFamilyIdsByFingerprintPopulatedFromMicrodescriptor()
throws Exception {