From 7529fccbcc5c82847f37ff1d564c4b57094b69c3 Mon Sep 17 00:00:00 2001 From: QuickMythril Date: Mon, 20 Jul 2026 06:52:07 -0400 Subject: [PATCH 1/2] qdn: make full-node relaying unconditional --- QORTIUM-CHANGELOG.md | 11 ++++++++++ preview/settings-preview-seed-netcup.json | 1 - preview/settings-preview-seed.json | 1 - preview/settings-preview.json | 1 - .../api/resource/ArbitraryResource.java | 6 +++-- .../ArbitraryDataFileListManager.java | 22 ++++++------------- .../arbitrary/ArbitraryDataFileManager.java | 11 ++-------- .../arbitrary/ArbitraryMetadataManager.java | 6 ++--- .../java/org/qortium/settings/Settings.java | 8 ------- .../org/qortium/test/api/AdminApiTests.java | 11 ++++++++++ .../qortium/test/api/ArbitraryApiTests.java | 12 ++++++++++ 11 files changed, 50 insertions(+), 40 deletions(-) diff --git a/QORTIUM-CHANGELOG.md b/QORTIUM-CHANGELOG.md index 21a8d738f..b57bf94d2 100644 --- a/QORTIUM-CHANGELOG.md +++ b/QORTIUM-CHANGELOG.md @@ -34,6 +34,17 @@ own chain. ## Change Entries +### 2026-07-20 - qdn: make full-node relaying unconditional + +Removes the inherited operator switch that allowed a full node to opt out of +relaying QDN discovery, metadata, and chunks for other peers. Qortium full nodes +now always participate in QDN relay while retaining separate storage-policy, +block-list, capacity, hop/time, and queue controls. Existing settings files that +still contain `relayModeEnabled` remain loadable, but the value is ignored; the +setting is no longer writable or shipped in Previewnet templates. The legacy +`/arbitrary/relaymode` endpoint remains temporarily for compatibility and +always returns `true`. + ### 2026-07-19 - fix(qdn): honor identifiers on status query requests Makes the shorter QDN status endpoint honor an `identifier` query diff --git a/preview/settings-preview-seed-netcup.json b/preview/settings-preview-seed-netcup.json index 8e9e1dd3d..731e2a30f 100644 --- a/preview/settings-preview-seed-netcup.json +++ b/preview/settings-preview-seed-netcup.json @@ -91,7 +91,6 @@ "recoveryModeTimeout": 0, "autoUpdateMode": "OFF", "qdnEnabled": true, - "relayModeEnabled": true, "qdnAuthBypassEnabled": true, "gatewayEnabled": true, "gatewayPort": 8090, diff --git a/preview/settings-preview-seed.json b/preview/settings-preview-seed.json index f50d741b0..303248d6f 100644 --- a/preview/settings-preview-seed.json +++ b/preview/settings-preview-seed.json @@ -91,7 +91,6 @@ "recoveryModeTimeout": 0, "autoUpdateMode": "OFF", "qdnEnabled": true, - "relayModeEnabled": true, "qdnAuthBypassEnabled": true, "gatewayEnabled": true, "gatewayPort": 8090, diff --git a/preview/settings-preview.json b/preview/settings-preview.json index 6f6f2ff55..006c26ab5 100644 --- a/preview/settings-preview.json +++ b/preview/settings-preview.json @@ -90,7 +90,6 @@ "recoveryModeTimeout": 0, "autoUpdateMode": "NOTIFY", "qdnEnabled": true, - "relayModeEnabled": true, "qdnAuthBypassEnabled": true, "publicQdnPublishMaxSize": 104857600, "publicApiWriteMaxBodySize": 262144, diff --git a/src/main/java/org/qortium/api/resource/ArbitraryResource.java b/src/main/java/org/qortium/api/resource/ArbitraryResource.java index 057c5f016..4698617be 100644 --- a/src/main/java/org/qortium/api/resource/ArbitraryResource.java +++ b/src/main/java/org/qortium/api/resource/ArbitraryResource.java @@ -620,7 +620,9 @@ public String createArbitrary(ArbitraryTransactionData transactionData) { @GET @Path("/relaymode") @Operation( - summary = "Returns whether relay mode is enabled or not", + summary = "Returns true because QDN relay is always enabled", + description = "Deprecated compatibility endpoint. Qortium full nodes always relay QDN data.", + deprecated = true, responses = { @ApiResponse( content = @Content(mediaType = MediaType.TEXT_PLAIN, schema = @Schema(type = "boolean")) @@ -631,7 +633,7 @@ public String createArbitrary(ArbitraryTransactionData transactionData) { public boolean getRelayMode(@HeaderParam(Security.API_KEY_HEADER) String apiKey) { Security.checkApiCallAllowed(request); - return Settings.getInstance().isRelayModeEnabled(); + return true; } @GET diff --git a/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManager.java b/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManager.java index ac12dc343..3b193c32c 100644 --- a/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManager.java +++ b/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManager.java @@ -88,12 +88,9 @@ public class ArbitraryDataFileListManager { /** Maximum number of hops that a file list relay request is allowed to make */ public static int RELAY_REQUEST_MAX_HOPS = 4; // was 4, this is no longer ArbData, only metaData/Lists public static int SEARCH_DEPTH_MAX_HOPS = 6; // Only used to determine if we should forward or terminate a search for QDN data - private final Boolean isRelayAvailable; - private ArbitraryDataFileListManager() { getArbitraryDataFileListMessageScheduler.scheduleAtFixedRate(this::processNetworkGetArbitraryDataFileListMessage, 60 * 1000, 500, TimeUnit.MILLISECONDS); arbitraryDataFileListMessageScheduler.scheduleAtFixedRate(this::processNetworkArbitraryDataFileListMessage, 60 * 1000, 500, TimeUnit.MILLISECONDS); - this.isRelayAvailable = Settings.getInstance().isRelayModeEnabled(); } @@ -662,7 +659,8 @@ private void processNetworkArbitraryDataFileListMessage() { // Determine how many peers to process based on request type List peersToProcess; - if (isRelayRequest != null && isRelayRequest && this.isRelayAvailable) { + boolean relayRequest = Boolean.TRUE.equals(isRelayRequest); + if (relayRequest) { // For relay requests, only process first peer to avoid duplicate forwards peersToProcess = Collections.singletonList(peerMessages.get(0)); @@ -680,7 +678,7 @@ private void processNetworkArbitraryDataFileListMessage() { ArbitraryDataFileListMessage arbitraryDataFileListMessage = (ArbitraryDataFileListMessage) message; // Process direct download responses (not relay forwarding) - if (!isRelayRequest || !this.isRelayAvailable) { + if (!relayRequest) { Long now = NTP.getTime(); @@ -769,7 +767,7 @@ private void processNetworkArbitraryDataFileListMessage() { // Forwarding - We are not the original requestor, just in the middle LOGGER.trace("Status of isRelayRequest {}", isRelayRequest); - if (isRelayRequest && this.isRelayAvailable) { + if (relayRequest) { Triple request = requestBySignature58.get(signature58); Peer requestingPeer = request.getB(); if (requestingPeer != null) { @@ -1064,7 +1062,7 @@ private void processNetworkGetArbitraryDataFileListMessage() { String nodeId = NetworkData.getInstance().getOurNodeId(); arbitraryDataFileListMessage = new ArbitraryDataFileListMessage(signature, - hashes, NTP.getTime(), 0, ourAddress, nodeId, this.isRelayAvailable, this.getIsDirectConnectable()); + hashes, NTP.getTime(), 0, ourAddress, nodeId, true, this.getIsDirectConnectable()); arbitraryDataFileListMessage.setId(message.getId()); @@ -1092,11 +1090,8 @@ private void processNetworkGetArbitraryDataFileListMessage() { } } - LOGGER.trace("We don't have hashes or all hashes - Checking Relay Mode"); - // We may need to forward this request on - - //if (this.isRelayAvailable ) { = We do not want to block the relay of file find requests, this prevents Direct-Connect finding - // In relay mode - so ask our other peers if they have it + LOGGER.trace("We don't have hashes or all hashes - forwarding the QDN search"); + // Qortium full nodes always relay QDN discovery requests. GetArbitraryDataFileListMessage getArbitraryDataFileListMessage = (GetArbitraryDataFileListMessage) message; @@ -1135,9 +1130,6 @@ private void processNetworkGetArbitraryDataFileListMessage() { } else { LOGGER.trace("Relay Request has timed out"); } - //} else { - // LOGGER.debug("Relay (fetch-reserve) is disabled"); - //} } } catch (Exception e) { LOGGER.error(e.getMessage(), e); diff --git a/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileManager.java b/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileManager.java index 532684707..aaccb882e 100644 --- a/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileManager.java +++ b/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileManager.java @@ -1666,9 +1666,6 @@ private void processDataFile(Peer peer, byte[] hash, byte[] sig, Message origina ArbitraryRelayInfo relayInfo = this.getOptimalRelayInfoEntryForHash(hash58); if (relayInfo != null) { removeFromRelayMap(relayInfo); - if (!Settings.getInstance().isRelayModeEnabled()) { - LOGGER.info("Relay info exists for hash {} but relay mode is disabled", hash58); - } LOGGER.trace("Selected optimal relay info for hash {}: peer={}, hops={}", hash58, relayInfo.getPeerData() != null ? relayInfo.getPeerData().getAddress() : "null", relayInfo.getRequestHops()); @@ -1748,7 +1745,7 @@ private void processDataFile(Peer peer, byte[] hash, byte[] sig, Message origina LOGGER.trace("Resolved relayPeer: {}, socketOpen: {}", relayPeer, socketOpen); - if (!socketOpen && Settings.getInstance().isRelayModeEnabled()) { + if (!socketOpen) { // Socket is closed or peer not found - try to force connect LOGGER.info("Socket closed or peer not found for relay to {}. Attempting reconnect...", relayPeerAddress); @@ -1758,7 +1755,7 @@ private void processDataFile(Peer peer, byte[] hash, byte[] sig, Message origina LOGGER.warn("Skipping relay for hash {} because socket is closed. Forcing reconnect for future requests.", hash58); } - else if (socketOpen && Settings.getInstance().isRelayModeEnabled()) { + else { // Track that this peer requested this hash from us, using composite key @@ -1800,10 +1797,6 @@ else if (socketOpen && Settings.getInstance().isRelayModeEnabled()) { } } } - else { - LOGGER.warn("Cannot relay for hash {}: relayPeer={}, socketOpen={}, relayModeEnabled={}", - hash58, relayPeer, socketOpen, Settings.getInstance().isRelayModeEnabled()); - } } } else { diff --git a/src/main/java/org/qortium/controller/arbitrary/ArbitraryMetadataManager.java b/src/main/java/org/qortium/controller/arbitrary/ArbitraryMetadataManager.java index fc3af8233..5e97a9cd5 100644 --- a/src/main/java/org/qortium/controller/arbitrary/ArbitraryMetadataManager.java +++ b/src/main/java/org/qortium/controller/arbitrary/ArbitraryMetadataManager.java @@ -485,7 +485,7 @@ public void onNetworkArbitraryMetadataMessage(Peer peer, Message message) { } // Forwarding - if (isRelayRequest && Settings.getInstance().isRelayModeEnabled()) { + if (isRelayRequest) { if (!isBlocked) { Peer requestingPeer = request.getB(); if (requestingPeer != null) { @@ -637,8 +637,8 @@ private void processNetworkGetArbitraryMetadataMessage() { // We may need to forward this request on boolean isBlocked = (transactionDataList == null || ListUtils.isQdnBlocked(transactionData.getService(), transactionData.getName(), transactionData.getIdentifier())); - if (Settings.getInstance().isRelayModeEnabled() && !isBlocked) { - // In relay mode - so ask our other peers if they have it + if (!isBlocked) { + // Ask our other peers because Qortium full nodes always relay QDN metadata. PeerMessage peerMessage = peerMessageBySignature58.get(signature58); GetArbitraryMetadataMessage getArbitraryMetadataMessage = (GetArbitraryMetadataMessage) peerMessage.message; diff --git a/src/main/java/org/qortium/settings/Settings.java b/src/main/java/org/qortium/settings/Settings.java index ce089aec9..eee178e72 100644 --- a/src/main/java/org/qortium/settings/Settings.java +++ b/src/main/java/org/qortium/settings/Settings.java @@ -513,9 +513,6 @@ public enum Transport { /** Storage policy to indicate which data should be hosted */ private String storagePolicy = "FOLLOWED_OR_VIEWED"; - /** Whether to allow data outside of the storage policy to be relayed between other peers */ - private boolean relayModeEnabled = true; - /** Whether, after publishing/holding our own resource, we proactively PUSH it out to a few * reachable (outbound, push-capable) peers so the data reaches the network even when this node * is NAT'd and cannot be pulled from. Receivers still gate acceptance on their own storage policy. */ @@ -1127,7 +1124,6 @@ private static Map buildWritableSettings() { settings.put("publicQdnPublishMaxSize", new WritableSetting(WritableSettingType.LONG, false)); settings.put("publicQdnPublishChunkMaxSize", new WritableSetting(WritableSettingType.LONG, false)); settings.put("publicQdnPublishChunkSessionLimit", new WritableSetting(WritableSettingType.INTEGER, false)); - settings.put("relayModeEnabled", new WritableSetting(WritableSettingType.BOOLEAN, false)); settings.put("qdnPushOnPublishEnabled", new WritableSetting(WritableSettingType.BOOLEAN, false)); settings.put("publicDataEnabled", new WritableSetting(WritableSettingType.BOOLEAN, false)); settings.put("privateDataEnabled", new WritableSetting(WritableSettingType.BOOLEAN, false)); @@ -2726,10 +2722,6 @@ public StoragePolicy getStoragePolicy() { return StoragePolicy.valueOf(this.storagePolicy); } - public boolean isRelayModeEnabled() { - return this.relayModeEnabled; - } - public boolean isQdnPushOnPublishEnabled() { return this.qdnPushOnPublishEnabled; } diff --git a/src/test/java/org/qortium/test/api/AdminApiTests.java b/src/test/java/org/qortium/test/api/AdminApiTests.java index 3134e856f..5a85f012b 100644 --- a/src/test/java/org/qortium/test/api/AdminApiTests.java +++ b/src/test/java/org/qortium/test/api/AdminApiTests.java @@ -86,6 +86,16 @@ public void testUpdateSettingsRejectsDisallowedSetting() throws Exception { assertEquals(StoragePolicy.NONE, Settings.getInstance().getStoragePolicy()); } + @Test + public void testLegacyRelayModeSettingCannotDisableRelay() throws Exception { + Path settingsPath = createWritableApiSettings("{\"relayModeEnabled\":false}"); + Settings.fileInstance(settingsPath.toString()); + + assertFalse(this.adminResource.settingsMetadata().writable.containsKey("relayModeEnabled")); + assertApiError(ApiError.INVALID_CRITERIA, + () -> this.adminResource.updateSettings(ApiCommon.TEST_API_KEY, "{\"relayModeEnabled\":false}")); + } + @Test public void testSettingsMetadata() throws Exception { Path settingsPath = createWritableApiSettings("{\"storagePolicy\":\"NONE\"}"); @@ -148,6 +158,7 @@ public void testSettingsMetadata() throws Exception { assertEquals(false, metadata.writable.get("publicQdnPublishChunkMaxSize").restartRequired); assertEquals("INTEGER", metadata.writable.get("publicQdnPublishChunkSessionLimit").type); assertEquals(false, metadata.writable.get("publicQdnPublishChunkSessionLimit").restartRequired); + assertFalse(metadata.writable.containsKey("relayModeEnabled")); assertTrue(metadata.pendingRestart.isEmpty()); } diff --git a/src/test/java/org/qortium/test/api/ArbitraryApiTests.java b/src/test/java/org/qortium/test/api/ArbitraryApiTests.java index 3077a47cf..b042e83f8 100644 --- a/src/test/java/org/qortium/test/api/ArbitraryApiTests.java +++ b/src/test/java/org/qortium/test/api/ArbitraryApiTests.java @@ -132,6 +132,18 @@ public void testDefaultStatusEndpointHonorsIdentifierQueryParameter() throws Exc } } + @Test + public void testRelayModeCompatibilityEndpointIsAlwaysEnabled() { + ApiCommon.installTestApiKey(); + try { + ArbitraryResource resource = (ArbitraryResource) ApiCommon + .buildResource(ArbitraryResource.class, ApiCommon.TEST_API_KEY); + assertTrue(resource.getRelayMode(ApiCommon.TEST_API_KEY)); + } finally { + ApiCommon.clearTestApiKey(); + } + } + @Test public void testPreviewPathWorksWithoutName() throws Exception { ApiCommon.installTestApiKey(); From baca61412c93defb69a037a932bc72c5559a59cc Mon Sep 17 00:00:00 2001 From: QuickMythril Date: Mon, 20 Jul 2026 07:34:12 -0400 Subject: [PATCH 2/2] Bound relay cache and cover forwarding behavior --- QORTIUM-CHANGELOG.md | 8 +- .../api/resource/ArbitraryResource.java | 2 +- .../ArbitraryDataFileListManager.java | 18 +- .../arbitrary/ArbitraryDataFileManager.java | 241 +++++++++++++----- .../ArbitraryDataStorageManager.java | 46 +++- .../arbitrary/ArbitraryMetadataManager.java | 2 +- .../ArbitraryDataFileListManagerTests.java | 179 +++++++++++++ .../ArbitraryDataFileManagerTests.java | 85 ++++++ .../ArbitraryDataStorageCapacityTests.java | 15 ++ 9 files changed, 523 insertions(+), 73 deletions(-) create mode 100644 src/test/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManagerTests.java diff --git a/QORTIUM-CHANGELOG.md b/QORTIUM-CHANGELOG.md index b57bf94d2..0e362f47c 100644 --- a/QORTIUM-CHANGELOG.md +++ b/QORTIUM-CHANGELOG.md @@ -34,11 +34,11 @@ own chain. ## Change Entries -### 2026-07-20 - qdn: make full-node relaying unconditional +### 2026-07-20 - qdn: make QDN relaying unconditional -Removes the inherited operator switch that allowed a full node to opt out of -relaying QDN discovery, metadata, and chunks for other peers. Qortium full nodes -now always participate in QDN relay while retaining separate storage-policy, +Removes the inherited operator switch that allowed a QDN-enabled node to opt out +of relaying discovery, metadata, and chunks for other peers. QDN-enabled Qortium +nodes now always participate in relay while retaining separate storage-policy, block-list, capacity, hop/time, and queue controls. Existing settings files that still contain `relayModeEnabled` remain loadable, but the value is ignored; the setting is no longer writable or shipped in Previewnet templates. The legacy diff --git a/src/main/java/org/qortium/api/resource/ArbitraryResource.java b/src/main/java/org/qortium/api/resource/ArbitraryResource.java index 4698617be..802b998ef 100644 --- a/src/main/java/org/qortium/api/resource/ArbitraryResource.java +++ b/src/main/java/org/qortium/api/resource/ArbitraryResource.java @@ -621,7 +621,7 @@ public String createArbitrary(ArbitraryTransactionData transactionData) { @Path("/relaymode") @Operation( summary = "Returns true because QDN relay is always enabled", - description = "Deprecated compatibility endpoint. Qortium full nodes always relay QDN data.", + description = "Deprecated compatibility endpoint. QDN-enabled Qortium nodes always relay QDN data.", deprecated = true, responses = { @ApiResponse( diff --git a/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManager.java b/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManager.java index 3b193c32c..b84f140cc 100644 --- a/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManager.java +++ b/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManager.java @@ -548,14 +548,18 @@ public void onNetworkArbitraryDataFileListMessage(Peer peer, Message message) { * Received a ArbitraryDataList Message */ private void processNetworkArbitraryDataFileListMessage() { + List messagesToProcess; + synchronized (arbitraryDataFileListMessageLock) { + messagesToProcess = new ArrayList<>(arbitraryDataFileListMessageList); + arbitraryDataFileListMessageList.clear(); + } - try { - List messagesToProcess; - synchronized (arbitraryDataFileListMessageLock) { - messagesToProcess = new ArrayList<>(arbitraryDataFileListMessageList); - arbitraryDataFileListMessageList.clear(); - } + processNetworkArbitraryDataFileListMessages(messagesToProcess); + } + /** Synchronous processing seam used by the scheduler and relay-behavior tests. */ + void processNetworkArbitraryDataFileListMessages(List messagesToProcess) { + try { if (messagesToProcess.isEmpty()) return; // Store ALL peer messages per signature (not just the last one) @@ -1091,7 +1095,7 @@ private void processNetworkGetArbitraryDataFileListMessage() { } LOGGER.trace("We don't have hashes or all hashes - forwarding the QDN search"); - // Qortium full nodes always relay QDN discovery requests. + // QDN-enabled Qortium nodes always relay discovery requests. GetArbitraryDataFileListMessage getArbitraryDataFileListMessage = (GetArbitraryDataFileListMessage) message; diff --git a/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileManager.java b/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileManager.java index aaccb882e..4f383b222 100644 --- a/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileManager.java +++ b/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataFileManager.java @@ -20,12 +20,14 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.stream.Collectors; import org.apache.commons.io.FileUtils; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import org.qortium.arbitrary.ArbitraryDataFile; +import org.qortium.arbitrary.ArbitraryDataFolderSizeEstimator; import org.qortium.controller.Controller; import org.qortium.data.arbitrary.ArbitraryDirectConnectionInfo; import org.qortium.data.arbitrary.ArbitraryFileListResponseInfo; @@ -156,6 +158,8 @@ private static class MetadataHashEntry { private static final long RELAY_CACHE_MIN_FILE_AGE_MS = 5 * 60 * 1000L; // 5 minutes minimum age before deletion private Path relayCacheDir; private final AtomicInteger relayCacheFileCount = new AtomicInteger(0); + private final AtomicLong relayCacheSize = new AtomicLong(0L); + private final Object relayCacheLock = new Object(); /** * Creates a composite key for tracking relay requests by hash and source peer @@ -168,27 +172,41 @@ private String createRelayKey(String hash58, Peer sourcePeer) { } /** - * Initializes the relay cache directory in the system temp directory + * Initializes the relay cache directory inside the QDN data directory. */ private void initializeRelayCache() { try { this.relayCacheDir = Paths.get(Settings.getInstance().getDataPath() + File.separator + RELAY_CACHE_DIR_NAME); Files.createDirectories(relayCacheDir); + + // A prior process may have stopped between staging and atomic replacement. + // Remove those incomplete files before establishing the exact startup size. + deleteStaleRelayCacheStagingFiles(false); // Count existing files on startup File[] existingFiles = relayCacheDir.toFile().listFiles(); if (existingFiles != null) { int fileCount = 0; + long cacheSize = 0L; for (File file : existingFiles) { if (file.isFile() && file.getName().endsWith(".tmp")) { fileCount++; + cacheSize += file.length(); } } this.relayCacheFileCount.set(fileCount); + this.relayCacheSize.set(cacheSize); LOGGER.debug("Initialized relay cache directory: {} ({} existing files)", relayCacheDir, fileCount); } else { + this.relayCacheSize.set(0L); LOGGER.debug("Initialized relay cache directory: {}", relayCacheDir); } + + // Establish exact QDN usage/capacity before the first relay-cache write. + // If this cannot be calculated, admission remains fail-closed at zero bytes. + Long now = NTP.getTime(); + ArbitraryDataStorageManager.getInstance().recalculateDataDirectorySize( + now != null ? now : System.currentTimeMillis()); cleanupRelayCache(); // Run cleanup in case disk conditions changed while node was offline } catch (IOException e) { LOGGER.error("Failed to initialize relay cache directory: {}", e.getMessage()); @@ -218,43 +236,147 @@ private Path getRelayCachePath(String hash58) { * @param hash58 The hash of the chunk * @param data The chunk data */ - private void saveToRelayCache(String hash58, byte[] data) { + boolean saveToRelayCache(String hash58, byte[] data) { if (relayCacheDir == null || data == null) { - return; + return false; } Path cachePath = getRelayCachePath(hash58); if (cachePath == null) { LOGGER.debug("Invalid relay cache path for hash {}", hash58); - return; + return false; } - boolean isNewFile = false; - try { - isNewFile = !Files.exists(cachePath); - if (isNewFile) { - int currentCount = relayCacheFileCount.incrementAndGet(); - if (currentCount > RELAY_CACHE_CLEANUP_TRIGGER) { - LOGGER.trace("Relay cache has {} files (threshold: {}), triggering cleanup", - currentCount, RELAY_CACHE_CLEANUP_TRIGGER); - cleanupRelayCache(); + synchronized (relayCacheLock) { + Path stagingPath = null; + try { + boolean isNewFile = !Files.exists(cachePath); + long existingSize = isNewFile ? 0L : Files.size(cachePath); + long sizeDelta = (long) data.length - existingSize; + + if (!hasRelayCacheWriteCapacity(sizeDelta, data.length)) { + return false; + } + + // Write to a sibling first so an interrupted write can never corrupt the + // existing cache entry or leave a partial file outside capacity accounting. + stagingPath = Files.createTempFile(relayCacheDir, ".relay-", ".part"); + try (java.io.OutputStream out = Files.newOutputStream(stagingPath, + java.nio.file.StandardOpenOption.TRUNCATE_EXISTING, + java.nio.file.StandardOpenOption.WRITE)) { + out.write(data); + out.flush(); + } + + // Re-read the destination and recheck logical headroom immediately before + // the atomic replacement in case capacity settings or external state changed. + isNewFile = !Files.exists(cachePath); + existingSize = isNewFile ? 0L : Files.size(cachePath); + sizeDelta = (long) data.length - existingSize; + if (!hasRelayCacheWriteCapacity(sizeDelta, 0L)) { + return false; + } + + Files.move(stagingPath, cachePath, + StandardCopyOption.ATOMIC_MOVE, + StandardCopyOption.REPLACE_EXISTING); + stagingPath = null; + + if (sizeDelta > 0L) { + ArbitraryDataFolderSizeEstimator.getInstance().add(sizeDelta); + relayCacheSize.addAndGet(sizeDelta); + } else if (sizeDelta < 0L) { + ArbitraryDataFolderSizeEstimator.getInstance().subtract(-sizeDelta); + relayCacheSize.addAndGet(sizeDelta); + } + + if (isNewFile) { + int currentCount = relayCacheFileCount.incrementAndGet(); + if (currentCount > RELAY_CACHE_CLEANUP_TRIGGER) { + LOGGER.trace("Relay cache has {} files (threshold: {}), triggering cleanup", + currentCount, RELAY_CACHE_CLEANUP_TRIGGER); + cleanupRelayCache(); + } + } + return true; + } catch (IOException e) { + LOGGER.warn("Failed to save to relay cache for hash {}: {}", hash58, e.getMessage()); + return false; + } finally { + if (stagingPath != null) { + try { + Files.deleteIfExists(stagingPath); + } catch (IOException e) { + LOGGER.warn("Failed to remove relay cache staging file {}: {}", + stagingPath.getFileName(), e.getMessage()); + Long now = NTP.getTime(); + ArbitraryDataStorageManager.getInstance().recalculateDataDirectorySize( + now != null ? now : System.currentTimeMillis()); + } } } + } + } - // Use streaming write to avoid holding byte[] reference in memory during I/O - try (java.io.OutputStream out = Files.newOutputStream(cachePath, - java.nio.file.StandardOpenOption.CREATE, - java.nio.file.StandardOpenOption.TRUNCATE_EXISTING, - java.nio.file.StandardOpenOption.WRITE)) { - out.write(data); - out.flush(); + private void deleteStaleRelayCacheStagingFiles(boolean updateEstimator) { + long deletedSize = 0L; + try (DirectoryStream stagingFiles = Files.newDirectoryStream(relayCacheDir, "*.part")) { + for (Path stagingFile : stagingFiles) { + try { + long size = Files.size(stagingFile); + if (Files.deleteIfExists(stagingFile)) { + deletedSize += size; + } + } catch (IOException e) { + LOGGER.warn("Failed to remove stale relay cache staging file {}: {}", + stagingFile.getFileName(), e.getMessage()); + } } } catch (IOException e) { - LOGGER.warn("Failed to save to relay cache for hash {}: {}", hash58, e.getMessage()); - if (isNewFile) { - relayCacheFileCount.decrementAndGet(); // Rollback counter on failure - } + LOGGER.warn("Failed to scan relay cache staging files: {}", e.getMessage()); + } + + if (updateEstimator && deletedSize > 0L) { + ArbitraryDataFolderSizeEstimator.getInstance().subtract(deletedSize); + } + } + + private boolean hasRelayCacheWriteCapacity(long logicalGrowth, long temporaryDiskBytes) { + long usableSpace = relayCacheDir.toFile().getUsableSpace(); + if (temporaryDiskBytes > usableSpace) { + LOGGER.debug("Skipping relay cache write: need {} temporary bytes, filesystem headroom {} bytes", + temporaryDiskBytes, usableSpace); + return false; + } + + if (logicalGrowth <= 0L) { + return true; + } + + long configuredHeadroom = ArbitraryDataStorageManager.getInstance() + .getRemainingStorageCapacityAtFullThreshold(); + long softCacheHeadroom = Math.max(0L, + calculateRelayCacheLimit(relayCacheSize.get()) - relayCacheSize.get()); + long writableHeadroom = Math.max(0L, Math.min(configuredHeadroom, softCacheHeadroom)); + if (logicalGrowth > writableHeadroom) { + LOGGER.debug("Skipping relay cache write: need {} logical bytes, headroom {} bytes", + logicalGrowth, writableHeadroom); + return false; + } + return true; + } + + private long calculateRelayCacheLimit(long currentCacheSize) { + ArbitraryDataStorageManager storageManager = ArbitraryDataStorageManager.getInstance(); + Long fullThresholdCapacity = storageManager.getStorageCapacityAtFullThreshold(); + if (fullThresholdCapacity == null) { + return 0L; } + + long estimatedTotalSize = ArbitraryDataFolderSizeEstimator.getInstance().get(); + long nonCacheSize = Math.max(0L, estimatedTotalSize - currentCacheSize); + long nonCacheHeadroom = Math.max(0L, fullThresholdCapacity - nonCacheSize); + return nonCacheHeadroom / 10L; } /** @@ -290,10 +412,15 @@ private void cleanupRelayCache() { if (relayCacheDir == null || !Files.exists(relayCacheDir)) { return; } - + + synchronized (relayCacheLock) { try { + deleteStaleRelayCacheStagingFiles(true); + File[] files = relayCacheDir.toFile().listFiles(); if (files == null || files.length == 0) { + relayCacheFileCount.set(0); + relayCacheSize.set(0L); return; } @@ -316,43 +443,27 @@ private void cleanupRelayCache() { } if (fileInfos.isEmpty()) { + relayCacheSize.set(totalSize); return; } // Sort by creation time (oldest first) fileInfos.sort(Comparator.comparingLong(fi -> fi.creationTime)); - // Calculate max allowed size based on QDN storage headroom - long maxAllowedSize; - ArbitraryDataStorageManager storageManager = ArbitraryDataStorageManager.getInstance(); - Long qdnStorageCapacity = storageManager.getStorageCapacity(); - long qdnUsedSpace = storageManager.getTotalDirectorySize(); - - if (qdnStorageCapacity != null && qdnUsedSpace > 0) { - // Calculate space before QDN hits its cleanup threshold (90%) - long qdnCleanupThreshold = (long)(qdnStorageCapacity * ArbitraryDataStorageManager.DELETION_THRESHOLD); - long qdnHeadroom = qdnCleanupThreshold - qdnUsedSpace; - - if (qdnHeadroom < 0) { - // QDN is already over threshold, use minimal relay cache - maxAllowedSize = 500L * 1024 * 1024; // 500MB minimum - LOGGER.debug("QDN over storage threshold, relay cache limited to 500MB"); - } else { - // Use up to 10% of the headroom, with min/max bounds - long calculatedSize = (long)(qdnHeadroom * 0.10); - maxAllowedSize = Math.max(500L * 1024 * 1024, // Min 500MB - calculatedSize); // 10% of free QDN space - RELAY_CACHE_CLEANUP_TRIGGER = (int)(maxAllowedSize / (512L * 1024)); // 500KB avg per file - LOGGER.debug("Relay cache limit: {} MB (based on {}% of {} MB headroom)", - maxAllowedSize / (1024 * 1024), - (int)(0.10 * 100), - qdnHeadroom / (1024 * 1024)); - } - } else { - // Fallback: conservative fixed size if QDN storage not calculated yet - maxAllowedSize = 2L * 1024 * 1024 * 1024; // 2GB - LOGGER.debug("Relay cache limit: 2GB (fallback - QDN storage not calculated)"); - } + // Calculate a soft cache ceiling from exact non-cache QDN usage. There is + // no positive fallback/minimum: the hard capacity boundary wins. + long maxAllowedSize = calculateRelayCacheLimit(totalSize); + Long fullThresholdCapacity = ArbitraryDataStorageManager.getInstance().getStorageCapacityAtFullThreshold(); + long estimatedTotalSize = ArbitraryDataFolderSizeEstimator.getInstance().get(); + long nonCacheSize = Math.max(0L, estimatedTotalSize - totalSize); + long nonCacheHeadroom = fullThresholdCapacity == null + ? 0L + : Math.max(0L, fullThresholdCapacity - nonCacheSize); + RELAY_CACHE_CLEANUP_TRIGGER = maxAllowedSize == 0L + ? 0 + : (int) Math.min(Integer.MAX_VALUE, Math.max(1L, maxAllowedSize / (512L * 1024L))); + LOGGER.debug("Relay cache limit: {} MB (10% of {} MB non-cache QDN headroom)", + maxAllowedSize / (1024 * 1024), nonCacheHeadroom / (1024 * 1024)); int deletedCount = 0; long deletedSize = 0; @@ -395,6 +506,10 @@ private void cleanupRelayCache() { } } relayCacheFileCount.set(actualCount); + relayCacheSize.set(totalSize); + if (deletedSize > 0L) { + ArbitraryDataFolderSizeEstimator.getInstance().subtract(deletedSize); + } if (deletedCount > 0) { String youngFilesMsg = skippedYoungFiles > 0 ? String.format(" (skipped %d young files)", skippedYoungFiles) : ""; @@ -410,6 +525,7 @@ private void cleanupRelayCache() { } catch (Exception e) { LOGGER.error("Error during relay cache cleanup: {}", e.getMessage(), e); } + } } /** @@ -474,10 +590,12 @@ public boolean eraseRelayCache() { return false; } + synchronized (relayCacheLock) { try { // Get the count before cleaning int deletedCount = relayCacheFileCount.get(); final AtomicInteger failedCount = new AtomicInteger(0); + final AtomicLong deletedSize = new AtomicLong(0L); // Recursive delete using Files.walkFileTree for better error handling Files.walkFileTree(relayCacheDir, new SimpleFileVisitor() { @@ -485,6 +603,7 @@ public boolean eraseRelayCache() { public FileVisitResult visitFile(Path file, BasicFileAttributes attrs) throws IOException { try { Files.delete(file); + deletedSize.addAndGet(attrs.size()); return FileVisitResult.CONTINUE; } catch (IOException e) { // File might be locked or have permission issues - log and continue @@ -522,6 +641,10 @@ public FileVisitResult visitFileFailed(Path file, IOException exc) throws IOExce // Update file count int failed = failedCount.get(); relayCacheFileCount.set(failed); + if (deletedSize.get() > 0L) { + ArbitraryDataFolderSizeEstimator.getInstance().subtract(deletedSize.get()); + relayCacheSize.updateAndGet(size -> Math.max(0L, size - deletedSize.get())); + } if (failed > 0) { LOGGER.warn("Erased relay cache: deleted {} files, {} failed", deletedCount - failed, failed); @@ -534,6 +657,7 @@ public FileVisitResult visitFileFailed(Path file, IOException exc) throws IOExce LOGGER.error("Error erasing relay cache: {}", e.getMessage(), e); return false; } + } } @@ -1012,7 +1136,7 @@ public void receivedArbitraryDataFile(Peer peer, ArbitraryDataFile adf) { } // Save to relay cache for persistence (uses streaming write to avoid holding reference) - saveToRelayCache(hash58, chunkData); + boolean cachedForRelay = saveToRelayCache(hash58, chunkData); // Capture primitives for MessageFactory lambdas byte[] hashCopy = hash; @@ -1113,7 +1237,8 @@ public void receivedArbitraryDataFile(Peer peer, ArbitraryDataFile adf) { } } // Relay chunks are now cached to temp dir - no need to skip further processing - LOGGER.trace("Completed relay forwarding for chunk {} - saved to relay cache, using zero-copy for {} forwards", hash58, pendingRequests.size()); + LOGGER.trace("Completed relay forwarding for chunk {} - cache persisted={}, using zero-copy for {} forwards", + hash58, cachedForRelay, pendingRequests.size()); return; } diff --git a/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataStorageManager.java b/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataStorageManager.java index 8bc62e942..43b554ec9 100644 --- a/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataStorageManager.java +++ b/src/main/java/org/qortium/controller/arbitrary/ArbitraryDataStorageManager.java @@ -48,7 +48,7 @@ public enum StoragePolicy { private static ArbitraryDataStorageManager instance; private volatile boolean isStopping = false; - private Long storageCapacity = null; + private volatile Long storageCapacity = null; private long totalDirectorySize = 0L; private long lastDirectorySizeCheck = 0; @@ -430,7 +430,7 @@ public boolean shouldCalculateDirectorySize(Long now) { return false; } - public void getDataDirectorySize(Long now) { + public synchronized void getDataDirectorySize(Long now) { if (now == null) { return; } @@ -526,6 +526,48 @@ public long getTotalDirectorySize() { return this.totalDirectorySize; } + /** + * Forces an exact directory scan and refreshes the usable/configured capacity. + * Relay-cache admission uses this at startup so it never relies on an + * uninitialized or stale fallback allowance. + */ + public synchronized void recalculateDataDirectorySize(Long now) { + if (now == null) { + return; + } + + this.recalculate.set(true); + this.getDataDirectorySize(now); + } + + /** + * Returns the bytes that can still be written before QDN reaches its normal + * storage-full threshold. The live estimator is used so cache and permanent + * writes are accounted for between full directory scans. + */ + public long getRemainingStorageCapacityAtFullThreshold() { + Long thresholdCapacity = this.getStorageCapacityAtFullThreshold(); + if (thresholdCapacity == null) { + return 0L; + } + + long estimatedSize = ArbitraryDataFolderSizeEstimator.getInstance().get(); + return Math.max(0L, thresholdCapacity - estimatedSize); + } + + public Long getStorageCapacityAtFullThreshold() { + if (this.storageCapacity == null) { + return null; + } + + long effectiveCapacity = this.storageCapacity; + Long configuredCapacity = Settings.getInstance().getMaxStorageCapacity(); + if (configuredCapacity != null) { + effectiveCapacity = Math.min(effectiveCapacity, configuredCapacity); + } + return (long) (effectiveCapacity * STORAGE_FULL_THRESHOLD); + } + public boolean isStorageSpaceAvailable(double threshold) { if (!this.isStorageCapacityCalculated()) { return false; diff --git a/src/main/java/org/qortium/controller/arbitrary/ArbitraryMetadataManager.java b/src/main/java/org/qortium/controller/arbitrary/ArbitraryMetadataManager.java index 5e97a9cd5..ae31ad3d8 100644 --- a/src/main/java/org/qortium/controller/arbitrary/ArbitraryMetadataManager.java +++ b/src/main/java/org/qortium/controller/arbitrary/ArbitraryMetadataManager.java @@ -638,7 +638,7 @@ private void processNetworkGetArbitraryMetadataMessage() { // We may need to forward this request on boolean isBlocked = (transactionDataList == null || ListUtils.isQdnBlocked(transactionData.getService(), transactionData.getName(), transactionData.getIdentifier())); if (!isBlocked) { - // Ask our other peers because Qortium full nodes always relay QDN metadata. + // Ask our other peers because QDN-enabled Qortium nodes always relay metadata. PeerMessage peerMessage = peerMessageBySignature58.get(signature58); GetArbitraryMetadataMessage getArbitraryMetadataMessage = (GetArbitraryMetadataMessage) peerMessage.message; diff --git a/src/test/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManagerTests.java b/src/test/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManagerTests.java new file mode 100644 index 000000000..cf401abbb --- /dev/null +++ b/src/test/java/org/qortium/controller/arbitrary/ArbitraryDataFileListManagerTests.java @@ -0,0 +1,179 @@ +package org.qortium.controller.arbitrary; + +import org.apache.commons.lang3.reflect.FieldUtils; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.qortium.account.PrivateKeyAccount; +import org.qortium.arbitrary.ArbitraryDataFile; +import org.qortium.arbitrary.misc.Service; +import org.qortium.controller.arbitrary.ArbitraryDataStorageManager.StoragePolicy; +import org.qortium.data.arbitrary.ArbitraryFileListResponseInfo; +import org.qortium.data.network.PeerData; +import org.qortium.data.transaction.ArbitraryTransactionData; +import org.qortium.data.transaction.RegisterNameTransactionData; +import org.qortium.network.Peer; +import org.qortium.network.PeerAddress; +import org.qortium.network.message.ArbitraryDataFileListMessage; +import org.qortium.network.message.Message; +import org.qortium.repository.Repository; +import org.qortium.repository.RepositoryManager; +import org.qortium.settings.Settings; +import org.qortium.test.common.ArbitraryUtils; +import org.qortium.test.common.Common; +import org.qortium.test.common.TransactionUtils; +import org.qortium.test.common.transaction.TestTransaction; +import org.qortium.transaction.RegisterNameTransaction; +import org.qortium.utils.Base58; +import org.qortium.utils.NTP; +import org.qortium.utils.Triple; + +import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; + +import static org.junit.Assert.assertArrayEquals; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; + +public class ArbitraryDataFileListManagerTests extends Common { + + private ArbitraryDataFileListManager fileListManager; + private ArbitraryDataFileManager fileManager; + + @Before + public void beforeTest() throws Exception { + Common.useDefaultSettings(); + FieldUtils.writeField(ArbitraryDataManager.getInstance(), "powDifficultyOverride", 1, true); + this.fileListManager = ArbitraryDataFileListManager.getInstance(); + this.fileManager = ArbitraryDataFileManager.getInstance(); + clearRelayState(); + } + + @After + public void afterTest() throws Exception { + clearRelayState(); + Common.useDefaultSettings(); + } + + @Test + public void testLegacyFalseSettingStillForwardsRelayResponseWithoutLocalDownload() throws Exception { + byte[] signature; + try (Repository repository = RepositoryManager.getRepository()) { + PrivateKeyAccount alice = Common.getTestAccount(repository, "alice"); + String name = "RELAY-TEST"; + + RegisterNameTransactionData registerName = new RegisterNameTransactionData( + TestTransaction.generateBase(alice), name, ""); + registerName.setFee(new RegisterNameTransaction(null, null).getUnitFee(registerName.getTimestamp())); + TransactionUtils.signAndMint(repository, registerName, alice); + + Path dataPath = ArbitraryUtils.generateRandomDataPath(128); + ArbitraryDataFile dataFile = ArbitraryUtils.createAndMintTxn(repository, + Base58.encode(alice.getPublicKey()), dataPath, name, null, + ArbitraryTransactionData.Method.PUT, Service.ARBITRARY_DATA, alice, 64); + signature = dataFile.getSignature(); + } + + Path legacySettings = Files.createTempFile("legacy-relay-setting", ".json"); + Files.write(legacySettings, "{\"relayModeEnabled\":false}\n".getBytes(StandardCharsets.UTF_8)); + Settings.fileInstance(legacySettings.toString()); + assertEquals(StoragePolicy.FOLLOWED_OR_VIEWED, Settings.getInstance().getStoragePolicy()); + + byte[] advertisedHash = new byte[32]; + Arrays.fill(advertisedHash, (byte) 7); + List hashes = List.of(advertisedHash); + long now = NTP.getTime(); + int requestId = 41001; + + CapturingPeer requester = new CapturingPeer("127.0.0.1:9100"); + CapturingPeer firstHolder = new CapturingPeer("127.0.0.1:9101"); + CapturingPeer secondHolder = new CapturingPeer("127.0.0.1:9102"); + String signature58 = Base58.encode(signature); + this.fileListManager.arbitraryDataFileListRequests.put(requestId, + new Triple<>(signature58, requester, now)); + + ArbitraryDataFileListMessage firstResponse = incomingFileList(requestId, signature, hashes, + now, 1, "127.0.0.1:9101", "holder-1"); + ArbitraryDataFileListMessage secondResponse = incomingFileList(requestId, signature, hashes, + now, 1, "127.0.0.1:9102", "holder-2"); + + this.fileListManager.processNetworkArbitraryDataFileListMessages(List.of( + new PeerMessage(firstHolder, firstResponse), + new PeerMessage(secondHolder, secondResponse))); + + assertEquals(1, requester.sentMessages.size()); + ArbitraryDataFileListMessage forwarded = parseFileList(requester.sentMessages.get(0)); + assertEquals(requestId, forwarded.getId()); + assertArrayEquals(signature, forwarded.getSignature()); + assertEquals(1, forwarded.getHashes().size()); + assertArrayEquals(advertisedHash, forwarded.getHashes().get(0)); + assertEquals(Integer.valueOf(2), forwarded.getRequestHops()); + + assertEquals(1, this.fileManager.arbitraryRelayMap.size()); + assertEquals(firstHolder.getPeerData(), this.fileManager.arbitraryRelayMap.get(0).getPeerData()); + assertEquals(0, localResponseCount()); + + int directRequestId = 41002; + this.fileListManager.arbitraryDataFileListRequests.put(directRequestId, + new Triple<>(signature58, null, now)); + ArbitraryDataFileListMessage directResponse = incomingFileList(directRequestId, signature, hashes, + now, 1, "127.0.0.1:9101", "holder-1"); + this.fileListManager.processNetworkArbitraryDataFileListMessages( + List.of(new PeerMessage(firstHolder, directResponse))); + + assertEquals(1, requester.sentMessages.size()); + assertEquals(1, localResponseCount()); + } + + private static ArbitraryDataFileListMessage incomingFileList(int id, byte[] signature, List hashes, + long requestTime, int requestHops, + String peerAddress, String nodeId) throws Exception { + ArbitraryDataFileListMessage outgoing = new ArbitraryDataFileListMessage(signature, hashes, + requestTime, requestHops, peerAddress, nodeId, true, false); + outgoing.setId(id); + return parseFileList(outgoing); + } + + private static ArbitraryDataFileListMessage parseFileList(Message message) throws Exception { + Message parsed = Message.fromByteBuffer(ByteBuffer.wrap(message.toBytes())); + assertNotNull(parsed); + return (ArbitraryDataFileListMessage) parsed; + } + + @SuppressWarnings("unchecked") + private int localResponseCount() throws IllegalAccessException { + List responses = (List) + FieldUtils.readField(this.fileManager, "arbitraryDataFileHashResponses", true); + return responses.size(); + } + + @SuppressWarnings("unchecked") + private void clearRelayState() throws IllegalAccessException { + this.fileListManager.arbitraryDataFileListRequests.clear(); + this.fileManager.arbitraryRelayMap.clear(); + List responses = (List) + FieldUtils.readField(this.fileManager, "arbitraryDataFileHashResponses", true); + responses.clear(); + } + + private static class CapturingPeer extends Peer { + private final List sentMessages = Collections.synchronizedList(new ArrayList<>()); + + private CapturingPeer(String address) { + super(new PeerData(new PeerAddress(address)), Peer.NETWORKDATA); + } + + @Override + public boolean sendMessage(Message message) { + this.sentMessages.add(message); + return true; + } + } +} diff --git a/src/test/java/org/qortium/controller/arbitrary/ArbitraryDataFileManagerTests.java b/src/test/java/org/qortium/controller/arbitrary/ArbitraryDataFileManagerTests.java index 3867a458e..05ff6acdc 100644 --- a/src/test/java/org/qortium/controller/arbitrary/ArbitraryDataFileManagerTests.java +++ b/src/test/java/org/qortium/controller/arbitrary/ArbitraryDataFileManagerTests.java @@ -1,10 +1,18 @@ package org.qortium.controller.arbitrary; +import org.apache.commons.io.FileUtils; +import org.apache.commons.lang3.reflect.FieldUtils; import org.junit.Before; import org.junit.Test; +import org.qortium.arbitrary.ArbitraryDataFolderSizeEstimator; import org.qortium.repository.DataException; import org.qortium.test.common.Common; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.stream.Stream; + +import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; @@ -59,4 +67,81 @@ public void testInFlightRequestCountIsScopedBySignature() { manager.clearChunkReceived("hashB1", signatureB); } } + + @Test + public void testRelayCacheWriteRespectsHeadroomAndTracksOverwriteDeltas() throws Exception { + ArbitraryDataFileManager manager = ArbitraryDataFileManager.getInstance(); + ArbitraryDataStorageManager storageManager = ArbitraryDataStorageManager.getInstance(); + Path testRoot = Files.createTempDirectory("relay-cache-capacity"); + Path relayCacheDir = testRoot.resolve("relay-cache"); + Files.createDirectories(relayCacheDir); + + Object originalRelayCacheDir = FieldUtils.readField(manager, "relayCacheDir", true); + java.util.concurrent.atomic.AtomicInteger fileCount = + (java.util.concurrent.atomic.AtomicInteger) FieldUtils.readField(manager, "relayCacheFileCount", true); + int originalFileCount = fileCount.get(); + java.util.concurrent.atomic.AtomicLong cacheSize = + (java.util.concurrent.atomic.AtomicLong) FieldUtils.readField(manager, "relayCacheSize", true); + long originalCacheSize = cacheSize.get(); + int originalCleanupTrigger = (Integer) FieldUtils.readField(manager, "RELAY_CACHE_CLEANUP_TRIGGER", true); + Object originalStorageCapacity = FieldUtils.readField(storageManager, "storageCapacity", true); + long originalEstimate = ArbitraryDataFolderSizeEstimator.getInstance().get(); + + try { + FieldUtils.writeField(manager, "relayCacheDir", relayCacheDir, true); + fileCount.set(0); + cacheSize.set(0L); + FieldUtils.writeField(storageManager, "storageCapacity", 1_000L, true); + ArbitraryDataFolderSizeEstimator.getInstance().set(800L); + + assertFalse(manager.saveToRelayCache("blocked", new byte[] { 1 })); + assertFalse(Files.exists(relayCacheDir.resolve("blocked.tmp"))); + assertEquals(800L, ArbitraryDataFolderSizeEstimator.getInstance().get()); + + ArbitraryDataFolderSizeEstimator.getInstance().set(700L); + assertTrue(manager.saveToRelayCache("allowed", new byte[] { 1, 2, 3, 4, 5 })); + assertEquals(705L, ArbitraryDataFolderSizeEstimator.getInstance().get()); + + assertTrue(manager.saveToRelayCache("allowed", new byte[] { 1, 2, 3, 4, 5, 6, 7, 8 })); + assertEquals(708L, ArbitraryDataFolderSizeEstimator.getInstance().get()); + + assertFalse(manager.saveToRelayCache("allowed", new byte[20])); + assertArrayEquals(new byte[] { 1, 2, 3, 4, 5, 6, 7, 8 }, + Files.readAllBytes(relayCacheDir.resolve("allowed.tmp"))); + assertEquals(708L, ArbitraryDataFolderSizeEstimator.getInstance().get()); + + assertTrue(manager.saveToRelayCache("allowed", new byte[] { 1, 2, 3, 4 })); + assertEquals(704L, ArbitraryDataFolderSizeEstimator.getInstance().get()); + + // A failed atomic replacement must preserve the destination and remove staging files. + Files.createDirectory(relayCacheDir.resolve("collision.tmp")); + assertFalse(manager.saveToRelayCache("collision", new byte[] { 9 })); + assertTrue(Files.isDirectory(relayCacheDir.resolve("collision.tmp"))); + try (Stream paths = Files.list(relayCacheDir)) { + assertFalse(paths.anyMatch(path -> path.getFileName().toString().endsWith(".part"))); + } + assertEquals(704L, ArbitraryDataFolderSizeEstimator.getInstance().get()); + + Path staleStagingFile = relayCacheDir.resolve("crash-orphan.part"); + Files.write(staleStagingFile, new byte[] { 6, 7, 8 }); + ArbitraryDataFolderSizeEstimator.getInstance().set(707L); + java.lang.reflect.Method cleanupMethod = ArbitraryDataFileManager.class + .getDeclaredMethod("cleanupRelayCache"); + cleanupMethod.setAccessible(true); + cleanupMethod.invoke(manager); + assertFalse(Files.exists(staleStagingFile)); + assertEquals(704L, ArbitraryDataFolderSizeEstimator.getInstance().get()); + + assertTrue(manager.eraseRelayCache()); + assertEquals(700L, ArbitraryDataFolderSizeEstimator.getInstance().get()); + } finally { + FieldUtils.writeField(manager, "relayCacheDir", originalRelayCacheDir, true); + fileCount.set(originalFileCount); + cacheSize.set(originalCacheSize); + FieldUtils.writeField(manager, "RELAY_CACHE_CLEANUP_TRIGGER", originalCleanupTrigger, true); + FieldUtils.writeField(storageManager, "storageCapacity", originalStorageCapacity, true); + ArbitraryDataFolderSizeEstimator.getInstance().set(originalEstimate); + FileUtils.deleteQuietly(testRoot.toFile()); + } + } } diff --git a/src/test/java/org/qortium/test/arbitrary/ArbitraryDataStorageCapacityTests.java b/src/test/java/org/qortium/test/arbitrary/ArbitraryDataStorageCapacityTests.java index eca2d667d..67ccb4bc4 100644 --- a/src/test/java/org/qortium/test/arbitrary/ArbitraryDataStorageCapacityTests.java +++ b/src/test/java/org/qortium/test/arbitrary/ArbitraryDataStorageCapacityTests.java @@ -86,6 +86,21 @@ public void testCalculateTotalStorageCapacity() { assertTrue(storageManager.shouldCalculateDirectorySize(now)); } + @Test + public void testRemainingStorageCapacityAtFullThresholdIsClamped() throws IllegalAccessException { + ArbitraryDataStorageManager storageManager = ArbitraryDataStorageManager.getInstance(); + FieldUtils.writeField(storageManager, "storageCapacity", 1_000L, true); + + ArbitraryDataFolderSizeEstimator.getInstance().set(700L); + assertEquals(100L, storageManager.getRemainingStorageCapacityAtFullThreshold()); + + ArbitraryDataFolderSizeEstimator.getInstance().set(800L); + assertEquals(0L, storageManager.getRemainingStorageCapacityAtFullThreshold()); + + ArbitraryDataFolderSizeEstimator.getInstance().set(900L); + assertEquals(0L, storageManager.getRemainingStorageCapacityAtFullThreshold()); + } + @Test public void testCalculateStorageCapacityPerName() { ArbitraryDataStorageManager storageManager = ArbitraryDataStorageManager.getInstance();