Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Comment thread
BenCodez marked this conversation as resolved.
running.set(true);
responseState.cancel();
responseState = new ResponseState();
Expand All @@ -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;
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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());
Expand All @@ -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); }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
Loading