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 {