diff --git a/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpBackendTransportConnector.java b/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpBackendTransportConnector.java index 322088b4..b8a66acd 100644 --- a/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpBackendTransportConnector.java +++ b/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpBackendTransportConnector.java @@ -75,6 +75,9 @@ public final class HttpBackendTransportConnector implements AutoCloseable { private volatile Thread callbackWorker; // Guarded by lifecycle. A flush drains the accepted queue but must not admit a new send. private boolean flushingOutgoing; + // Guarded by lifecycle. Staged connectors poll normally but must not cross the + // durable inbound fence until their application handler has been published. + private boolean incomingActive = true; // Guarded by state. Keep admission and queue insertion in the same critical section. private boolean sendAdmissionOpen; private long sequence; @@ -189,12 +192,34 @@ private static HttpClientCredentialStore.ClientCredential enrollLocked( } public void start() { + start(true); + } + /** + * Starts normal transport while holding inbound callbacks before their first + * durable journal transition. Call {@link #activateIncoming()} only after the + * application handler backed by this connector has been published. + */ + public void startPaused() { + start(false); + } + private void start(boolean activateIncoming) { synchronized (lifecycle) { - if (closing.get() || flushingOutgoing || running.get()) return; + if (closing.get() || flushingOutgoing) return; + if (running.get()) { + // start() is also the idempotent lifecycle entry point used by callers + // that publish their handler after starting a staged connector. Reopen + // the barrier and wake callbacks already waiting on it. + if (activateIncoming && !incomingActive) { + incomingActive = true; + lifecycle.notifyAll(); + } + return; + } // A timed-out flush can leave the interrupted poller winding down. It still // owns the single long-poll slot until it exits, so a restart must wait for a // later start call rather than creating an overlapping poller. if (poller != null && poller.isAlive()) return; + incomingActive = activateIncoming; running.set(true); responseState.cancel(); responseState = new ResponseState(); @@ -204,6 +229,14 @@ public void start() { poller.start(); } } + /** Opens the one-way inbound callback barrier for a staged running connector. */ + public void activateIncoming() { + synchronized (lifecycle) { + if (closing.get() || !running.get() || incomingActive) return; + incomingActive = true; + lifecycle.notifyAll(); + } + } /** Waits for one authenticated, protocol-valid transport response. */ public boolean awaitFirstResponse(long deadlineNanos) throws InterruptedException { ResponseState expected = responseState; @@ -329,6 +362,7 @@ public boolean flushOutgoing(long deadlineNanos) { current = poller; if (current != null) current.interrupt(); } + lifecycle.notifyAll(); } if (!alreadyClosing) synchronized (state) { sendAdmissionOpen = false; } if (alreadyClosing) { @@ -471,6 +505,10 @@ void dispatch(HttpTransportProtocol.Delivery delivery) { Runnable callback = () -> { boolean success = false; try { + if (!awaitIncomingActivation()) { + completeIncoming(delivery.id(), false); + return; + } if (inboundDeliveries != null) { if (inboundDeliveries.state(delivery.id()) == null) inboundDeliveries.reserve(delivery.id()); inboundDeliveries.markRunning(delivery.id()); @@ -489,6 +527,15 @@ void dispatch(HttpTransportProtocol.Delivery delivery) { }; if (!executeOrdered(callbackExecutor, callback)) completeIncoming(delivery.id(), false); } + private boolean awaitIncomingActivation() { + boolean interrupted = false; + synchronized (lifecycle) { + while (!incomingActive && !closing.get()) try { lifecycle.wait(); } + catch (InterruptedException stopRequested) { interrupted = true; break; } + if (interrupted) Thread.currentThread().interrupt(); + return !interrupted && incomingActive && !closing.get(); + } + } void completeIncoming(String id, boolean success) { synchronized (state) { processing.remove(id); if (success) { received.add(id); while (received.size() > HttpTransportProtocol.MAX_QUEUE) received.remove(received.iterator().next()); queueAck(id); } diff --git a/SimpleAPI/src/test/java/com/bencodez/simpleapi/servercomm/http/HttpTransportRuntimeTest.java b/SimpleAPI/src/test/java/com/bencodez/simpleapi/servercomm/http/HttpTransportRuntimeTest.java index f6ee809e..e2185563 100644 --- a/SimpleAPI/src/test/java/com/bencodez/simpleapi/servercomm/http/HttpTransportRuntimeTest.java +++ b/SimpleAPI/src/test/java/com/bencodez/simpleapi/servercomm/http/HttpTransportRuntimeTest.java @@ -1421,6 +1421,85 @@ void backendCallbackQueueBackpressuresWithoutBreakingFifo() throws Exception { } finally { releaseFirst.countDown(); } } + @Test + void pausedBackendDefersCallbackAndDurableFenceUntilActivation() throws Exception { + HttpTlsIdentity identity = HttpTlsIdentity.loadOrCreate(directory.resolve("paused-proxy"), "localhost"); + HttpTlsIdentity.IssuedClientCertificate issued = identity.issueClientCertificate("lobby-1"); + Path clientDirectory = directory.resolve("paused-client"); + HttpConnectionCode code = new HttpConnectionCode("lobby-1", java.net.URI.create("https://localhost:8443/"), + identity.serverCertificatePin(), identity.caCertificatePin(), Instant.now().plusSeconds(300), "D".repeat(43)); + HttpClientCredentialStore.saveEnrolled(clientDirectory, code, issued); + CountDownLatch callback = new CountDownLatch(1); + String id = java.util.UUID.randomUUID().toString(); + try (HttpBackendTransportConnector connector = new HttpBackendTransportConnector(clientDirectory, + ignored -> callback.countDown())) { + connector.startPaused(); + connector.dispatch(new HttpTransportProtocol.Delivery(id, JsonEnvelope.builder("vote").build())); + assertFalse(callback.await(150, TimeUnit.MILLISECONDS)); + var inboundField = HttpBackendTransportConnector.class.getDeclaredField("inboundDeliveries"); + inboundField.setAccessible(true); + HttpInboundDeliveryStore store = (HttpInboundDeliveryStore) inboundField.get(connector); + assertEquals(null, store.state(id), "publication staging must not create a durable replay fence"); + + connector.activateIncoming(); + assertTrue(callback.await(2, TimeUnit.SECONDS)); + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(2); + while (store.state(id) != HttpInboundDeliveryStore.State.COMPLETED && System.nanoTime() < deadline) + Thread.sleep(5); + assertEquals(HttpInboundDeliveryStore.State.COMPLETED, store.state(id)); + } + } + + @Test + void restartingPausedBackendWithStartReleasesWaitingCallbacks() throws Exception { + HttpTlsIdentity identity = HttpTlsIdentity.loadOrCreate(directory.resolve("paused-start-proxy"), "localhost"); + HttpTlsIdentity.IssuedClientCertificate issued = identity.issueClientCertificate("lobby-1"); + Path clientDirectory = directory.resolve("paused-start-client"); + HttpConnectionCode code = new HttpConnectionCode("lobby-1", java.net.URI.create("https://localhost:8443/"), + identity.serverCertificatePin(), identity.caCertificatePin(), Instant.now().plusSeconds(300), "F".repeat(43)); + HttpClientCredentialStore.saveEnrolled(clientDirectory, code, issued); + CountDownLatch callback = new CountDownLatch(1); + try (HttpBackendTransportConnector connector = new HttpBackendTransportConnector(clientDirectory, + ignored -> callback.countDown())) { + connector.startPaused(); + connector.dispatch(new HttpTransportProtocol.Delivery(java.util.UUID.randomUUID().toString(), + JsonEnvelope.builder("vote").build())); + assertFalse(callback.await(150, TimeUnit.MILLISECONDS)); + + // start() remains idempotent for the poller but must reopen a barrier + // left by startPaused(), including callbacks already waiting on it. + connector.start(); + assertTrue(callback.await(2, TimeUnit.SECONDS)); + } + } + + @Test + void closingPausedBackendLeavesDeliveryRecoverable() throws Exception { + HttpTlsIdentity identity = HttpTlsIdentity.loadOrCreate(directory.resolve("paused-close-proxy"), "localhost"); + HttpTlsIdentity.IssuedClientCertificate issued = identity.issueClientCertificate("lobby-1"); + Path clientDirectory = directory.resolve("paused-close-client"); + HttpConnectionCode code = new HttpConnectionCode("lobby-1", java.net.URI.create("https://localhost:8443/"), + identity.serverCertificatePin(), identity.caCertificatePin(), Instant.now().plusSeconds(300), "E".repeat(43)); + HttpClientCredentialStore.saveEnrolled(clientDirectory, code, issued); + AtomicInteger callbacks = new AtomicInteger(); + String id = java.util.UUID.randomUUID().toString(); + HttpTransportProtocol.Delivery delivery = new HttpTransportProtocol.Delivery(id, JsonEnvelope.builder("vote").build()); + HttpBackendTransportConnector staged = new HttpBackendTransportConnector(clientDirectory, + ignored -> callbacks.incrementAndGet()); + staged.startPaused(); + staged.dispatch(delivery); + Thread.sleep(150); + staged.close(); + assertEquals(0, callbacks.get()); + + try (HttpBackendTransportConnector restarted = new HttpBackendTransportConnector(clientDirectory, + ignored -> callbacks.incrementAndGet())) { + assertEquals(java.util.List.of(delivery), restarted.accept(java.util.List.of(delivery)), + "abandoned staging must leave the proxy delivery replayable"); + assertTrue(restarted.drainAcknowledgements().isEmpty()); + } + } + @Test void durableBackendFencePreventsCallbackReplayAfterRestartBeforeAck() throws Exception { HttpTlsIdentity identity = HttpTlsIdentity.loadOrCreate(directory.resolve("fence-proxy"), "localhost");