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 b8a66ac..92e86dc 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 @@ -22,6 +22,7 @@ import java.util.LinkedHashMap; import java.util.LinkedHashSet; import java.util.List; +import java.util.Objects; import java.util.Set; import java.util.UUID; import java.util.concurrent.ArrayBlockingQueue; @@ -50,6 +51,7 @@ public final class HttpBackendTransportConnector implements AutoCloseable { private volatile HttpClientCredentialStore.HttpClientProfile profile; private final String serverId; private final Consumer onEnvelope; + private final HttpEnvelopeWireCodec wireCodec; private volatile HttpClient client; private volatile HttpClientCredentialStore.ClientCredential credential; private final Path credentialDirectory; @@ -105,14 +107,23 @@ public final class HttpBackendTransportConnector implements AutoCloseable { private HttpBackendTransportConnector(HttpClientCredentialStore.EnrolledClient enrolled, Consumer onEnvelope, Path credentialDirectory) throws Exception { - this(enrolled == null ? null : enrolled.profile(), enrolled == null ? null : enrolled.credential(), onEnvelope, credentialDirectory); + this(enrolled == null ? null : enrolled.profile(), enrolled == null ? null : enrolled.credential(), onEnvelope, + credentialDirectory, HttpEnvelopeWireCodec.identity()); } private HttpBackendTransportConnector(HttpClientCredentialStore.HttpClientProfile profile, HttpClientCredentialStore.ClientCredential credential, Consumer onEnvelope, Path credentialDirectory) throws Exception { - if (profile == null || credential == null || onEnvelope == null) throw new IllegalArgumentException("HTTP backend transport configuration is invalid"); + this(profile, credential, onEnvelope, credentialDirectory, HttpEnvelopeWireCodec.identity()); + } + + private HttpBackendTransportConnector(HttpClientCredentialStore.HttpClientProfile profile, + HttpClientCredentialStore.ClientCredential credential, Consumer onEnvelope, Path credentialDirectory, + HttpEnvelopeWireCodec wireCodec) throws Exception { + if (profile == null || credential == null || onEnvelope == null || wireCodec == null) + throw new IllegalArgumentException("HTTP backend transport configuration is invalid"); if (!matchesCredential(profile, credential)) throw new IllegalArgumentException("HTTP client certificate does not match transport profile"); this.profile = profile; this.serverId = profile.serverId(); this.onEnvelope = onEnvelope; + this.wireCodec = HttpEnvelopeWireCodec.serialized(wireCodec); this.credential = credential; this.credentialDirectory = credentialDirectory; HttpInboundDeliveryStore loadedInbound = null, loadedAcknowledgements = null; @@ -165,6 +176,21 @@ public HttpBackendTransportConnector(Path credentials, Consumer on this(HttpClientCredentialStore.loadEnrolled(credentials), onEnvelope, credentials); } + /** + * Starts normal transport with a codec applied only at the HTTP wire boundary. Queued semantic + * envelopes remain independent of the codec's current key or policy. + */ + public HttpBackendTransportConnector(Path credentials, Consumer onEnvelope, + HttpEnvelopeWireCodec wireCodec) throws Exception { + this(HttpClientCredentialStore.loadEnrolled(credentials), onEnvelope, credentials, wireCodec); + } + + private HttpBackendTransportConnector(HttpClientCredentialStore.EnrolledClient enrolled, + Consumer onEnvelope, Path credentialDirectory, HttpEnvelopeWireCodec wireCodec) throws Exception { + this(enrolled == null ? null : enrolled.profile(), enrolled == null ? null : enrolled.credential(), onEnvelope, + credentialDirectory, wireCodec); + } + /** Performs enrollment network I/O; call this from a connector/setup worker, never a platform main thread. */ public static HttpClientCredentialStore.ClientCredential enroll(HttpConnectionCode code, String serverId, Path credentials) throws Exception { if (credentials == null) throw new IllegalArgumentException("Enrollment configuration is invalid"); @@ -252,8 +278,13 @@ private boolean responseStateIsActive(ResponseState expected) { synchronized (li */ public boolean send(JsonEnvelope envelope) { if (envelope == null) return false; - try { HttpTransportProtocol.validateEnvelope(envelope); } - catch (IllegalArgumentException invalid) { return false; } + try { + HttpTransportProtocol.validateEnvelope(envelope); + JsonEnvelope encoded = wireCodec.encode(envelope); + if (encoded == null) throw new IllegalArgumentException("HTTP wire codec returned no envelope"); + HttpTransportProtocol.validateEnvelope(encoded); + } + catch (RuntimeException invalid) { return false; } synchronized (state) { if (!sendAdmissionOpen || !running.get()) return false; if (outgoing.size() >= HttpTransportProtocol.MAX_QUEUE) return false; @@ -279,8 +310,16 @@ private boolean pollOnce(Duration timeout, boolean requireRunning, boolean accep acks = first(acknowledgements); ackConfirmations = first(acknowledgementConfirmations); requestSequence = sequence++; + List encoded = new java.util.ArrayList<>( + Math.min(outgoing.size(), HttpTransportProtocol.MAX_BATCH)); + for (HttpTransportProtocol.Delivery delivery : outgoing.values()) { + if (encoded.size() == HttpTransportProtocol.MAX_BATCH) break; + JsonEnvelope envelope = Objects.requireNonNull(wireCodec.encode(delivery.envelope()), "encoded envelope"); + HttpTransportProtocol.validateEnvelope(envelope); + encoded.add(new HttpTransportProtocol.Delivery(delivery.id(), envelope)); + } messages = HttpTransportProtocol.fittingMessages(serverId, session, requestSequence, acks, - ackConfirmations, outgoing.values()); + ackConfirmations, encoded); for (int index = 0; index < acks.size(); index++) acknowledgements.removeFirst(); for (int index = 0; index < ackConfirmations.size(); index++) acknowledgementConfirmations.removeFirst(); } @@ -292,6 +331,7 @@ private boolean pollOnce(Duration timeout, boolean requireRunning, boolean accep HttpTransportProtocol.Packet packet = HttpTransportProtocol.parsePacket(response.body()); if (!serverId.equals(packet.server()) || !session.equals(packet.session()) || packet.sequence() != requestSequence) return false; if (!acks.equals(packet.ackConfirmations())) return false; + packet = decodePacket(packet); confirmAcknowledgements(packet.ackConfirmations()); acknowledgementsConfirmed = true; if (!confirmSentAcknowledgementConfirmations(ackConfirmations)) return false; @@ -312,6 +352,16 @@ private boolean pollOnce(Duration timeout, boolean requireRunning, boolean accep if (!confirmationsConfirmed) requeue(acknowledgementConfirmations, ackConfirmations); } } + private HttpTransportProtocol.Packet decodePacket(HttpTransportProtocol.Packet packet) { + List decoded = new java.util.ArrayList<>(packet.messages().size()); + for (HttpTransportProtocol.Delivery delivery : packet.messages()) { + JsonEnvelope envelope = Objects.requireNonNull(wireCodec.decode(delivery.envelope()), "decoded envelope"); + HttpTransportProtocol.validateEnvelope(envelope); + decoded.add(new HttpTransportProtocol.Delivery(delivery.id(), envelope)); + } + return new HttpTransportProtocol.Packet(packet.server(), packet.session(), packet.sequence(), packet.acks(), + packet.ackConfirmations(), decoded); + } /** Stops normal polling and gives already-queued outbound messages a bounded final delivery attempt. */ public boolean flushOutgoing(long deadlineNanos) { Thread current; diff --git a/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpEnvelopeWireCodec.java b/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpEnvelopeWireCodec.java new file mode 100644 index 0000000..4c096db --- /dev/null +++ b/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpEnvelopeWireCodec.java @@ -0,0 +1,62 @@ +package com.bencodez.simpleapi.servercomm.http; + +import java.util.Objects; + +import com.bencodez.simpleapi.servercomm.codec.JsonEnvelope; + +/** + * Transforms semantic HTTP transport envelopes at the wire boundary. + * + *

The HTTP delivery queues always retain the untransformed semantic envelope. Implementations + * may therefore change keys or policy across a restart without making already accepted deliveries + * unreadable. A codec must be non-blocking and keep a stable policy for the lifetime of one + * transport instance, produce envelopes within the HTTP transport size limit, and reject invalid + * wire envelopes by throwing an unchecked exception.

+ */ +public interface HttpEnvelopeWireCodec { + JsonEnvelope encode(JsonEnvelope envelope); + + JsonEnvelope decode(JsonEnvelope envelope); + + static HttpEnvelopeWireCodec identity() { + return Identity.INSTANCE; + } + + /** Serializes a codec shared by concurrent HTTP sessions. */ + static HttpEnvelopeWireCodec serialized(HttpEnvelopeWireCodec codec) { + Objects.requireNonNull(codec, "wireCodec"); + return codec == Identity.INSTANCE ? codec : new Serialized(codec); + } + + final class Identity implements HttpEnvelopeWireCodec { + private static final Identity INSTANCE = new Identity(); + + private Identity() { } + + @Override public JsonEnvelope encode(JsonEnvelope envelope) { return envelope; } + + @Override public JsonEnvelope decode(JsonEnvelope envelope) { return envelope; } + } + + final class Serialized implements HttpEnvelopeWireCodec { + private final HttpEnvelopeWireCodec delegate; + + private Serialized(HttpEnvelopeWireCodec delegate) { + this.delegate = delegate; + } + + @Override + public JsonEnvelope encode(JsonEnvelope envelope) { + synchronized (delegate) { + return delegate.encode(envelope); + } + } + + @Override + public JsonEnvelope decode(JsonEnvelope envelope) { + synchronized (delegate) { + return delegate.decode(envelope); + } + } + } +} diff --git a/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpProxyTransportServer.java b/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpProxyTransportServer.java index ce319c5..8ee39a2 100644 --- a/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpProxyTransportServer.java +++ b/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/http/HttpProxyTransportServer.java @@ -84,6 +84,7 @@ public final class HttpProxyTransportServer implements AutoCloseable { private final Path durableIncomingRoot; private final Consumer onEnvelope; private final DeliveryAcknowledgement onAcknowledged; + private final HttpEnvelopeWireCodec wireCodec; private final LongSupplier nanoTime; private volatile boolean closed; private boolean closeFinalizing, closeFinalized; @@ -91,7 +92,8 @@ public final class HttpProxyTransportServer implements AutoCloseable { /** In-memory constructor for tests; production callers must supply a durable state directory. */ HttpProxyTransportServer(InetSocketAddress bind, HttpTlsIdentity identity, HttpEnrollmentAuthority authority, Consumer onEnvelope) throws Exception { - this(bind, identity, authority, null, onEnvelope, (serverId, deliveryId) -> { }, System::nanoTime); + this(bind, identity, authority, null, onEnvelope, (serverId, deliveryId) -> { }, + HttpEnvelopeWireCodec.identity(), System::nanoTime); } public HttpProxyTransportServer(InetSocketAddress bind, HttpTlsIdentity identity, HttpEnrollmentAuthority authority, @@ -103,17 +105,37 @@ public HttpProxyTransportServer(InetSocketAddress bind, HttpTlsIdentity identity Path outgoingDirectory, Consumer onEnvelope, DeliveryAcknowledgement onAcknowledged) throws Exception { this(bind, identity, authority, Objects.requireNonNull(outgoingDirectory, "outgoingDirectory is required"), - onEnvelope, onAcknowledged, System::nanoTime); + onEnvelope, onAcknowledged, HttpEnvelopeWireCodec.identity(), System::nanoTime); + } + + /** + * Creates a durable transport whose queues retain semantic envelopes while the supplied codec + * is applied independently to every HTTP transmission and receipt. + */ + public HttpProxyTransportServer(InetSocketAddress bind, HttpTlsIdentity identity, HttpEnrollmentAuthority authority, + Path outgoingDirectory, Consumer onEnvelope, + DeliveryAcknowledgement onAcknowledged, HttpEnvelopeWireCodec wireCodec) throws Exception { + this(bind, identity, authority, Objects.requireNonNull(outgoingDirectory, "outgoingDirectory is required"), + onEnvelope, onAcknowledged, wireCodec, System::nanoTime); } HttpProxyTransportServer(InetSocketAddress bind, HttpTlsIdentity identity, HttpEnrollmentAuthority authority, Path outgoingDirectory, Consumer onEnvelope, DeliveryAcknowledgement onAcknowledged, LongSupplier nanoTime) throws Exception { + this(bind, identity, authority, outgoingDirectory, onEnvelope, onAcknowledged, + HttpEnvelopeWireCodec.identity(), nanoTime); + } + + HttpProxyTransportServer(InetSocketAddress bind, HttpTlsIdentity identity, HttpEnrollmentAuthority authority, + Path outgoingDirectory, Consumer onEnvelope, + DeliveryAcknowledgement onAcknowledged, HttpEnvelopeWireCodec wireCodec, LongSupplier nanoTime) throws Exception { if (bind == null || identity == null || authority == null || onEnvelope == null || onAcknowledged == null) throw new IllegalArgumentException("HTTP transport configuration is required"); - if (nanoTime == null) throw new IllegalArgumentException("HTTP transport clock is required"); + if (wireCodec == null || nanoTime == null) throw new IllegalArgumentException("HTTP transport codec and clock are required"); this.identity = identity; this.authority = authority; this.onEnvelope = onEnvelope; - this.onAcknowledged = onAcknowledged; this.nanoTime = nanoTime; + this.onAcknowledged = onAcknowledged; + this.wireCodec = HttpEnvelopeWireCodec.serialized(wireCodec); + this.nanoTime = nanoTime; durableOutgoing = outgoingDirectory == null ? null : new DurableOutgoingQueue(outgoingDirectory); HttpsServer createdServer = null; ThreadPoolExecutor createdListener = null, createdHandler = null; @@ -138,6 +160,7 @@ public HttpProxyTransportServer(InetSocketAddress bind, HttpTlsIdentity identity } if (durableOutgoing != null) for (Map.Entry> pending : durableOutgoing.load().entrySet()) { + for (HttpTransportProtocol.Delivery delivery : pending.getValue()) validateForWire(delivery.envelope()); BackendState state = backendState(pending.getKey()); state.restore(pending.getValue()); } @@ -297,8 +320,9 @@ private boolean send(String serverId, String deliveryId, JsonEnvelope envelope, serverId = HttpTlsIdentity.canonicalServerId(serverId); HttpTransportProtocol.validId(deliveryId); HttpTransportProtocol.validateEnvelope(envelope); + validateForWire(envelope); } - catch (IllegalArgumentException invalid) { return false; } + catch (RuntimeException invalid) { return false; } BackendState backend; final String canonicalServerId = serverId; try { backend = backendState(canonicalServerId); } @@ -310,6 +334,11 @@ private boolean send(String serverId, String deliveryId, JsonEnvelope envelope, throw indeterminate; } } + private void validateForWire(JsonEnvelope envelope) { + JsonEnvelope encoded = wireCodec.encode(envelope); + if (encoded == null) throw new IllegalArgumentException("HTTP wire codec returned no envelope"); + HttpTransportProtocol.validateEnvelope(encoded); + } @Override public void close() { boolean callbackWorker = Thread.currentThread() == handlerWorker.get(); @@ -400,15 +429,17 @@ private void transport(HttpsExchange exchange) throws IOException { if (bodyError != 0) { reply(exchange, bodyError, new byte[0]); return; } if (!admission.tryAcquire()) { reply(exchange, 429, new byte[0]); return; } try { - HttpTransportProtocol.Packet packet = HttpTransportProtocol.parsePacket(read(exchange.getRequestBody(), HttpTransportProtocol.MAX_BODY_BYTES)); + HttpTransportProtocol.Packet packet = HttpTransportProtocol.parsePacket( + read(exchange.getRequestBody(), HttpTransportProtocol.MAX_BODY_BYTES)); X509Certificate certificate = peerCertificate(exchange); if (certificate == null || !authority.authenticate(packet.server(), certificate)) { reply(exchange, 401, new byte[0]); return; } + packet = decodePacket(packet); BackendState backend; backend = backendState(packet.server()); if (!backend.beginPoll(packet.session())) { reply(exchange, 409, new byte[0]); return; } try { handlePacket(packet, backend); - Response response = backend.await(packet.server(), packet.session(), packet.sequence(), packet.acks()); + Response response = backend.await(packet.server(), packet.session(), packet.sequence(), packet.acks(), wireCodec); reply(exchange, 200, HttpTransportProtocol.response(packet.server(), packet.session(), packet.sequence(), response.acks(), packet.acks(), response.messages())); } finally { backend.endPoll(); } @@ -429,6 +460,16 @@ private void handlePacket(HttpTransportProtocol.Packet packet, BackendState back synchronized (backend) { accepted = backend.acceptIncoming(packet.messages()); } for (HttpTransportProtocol.Delivery delivery : accepted) dispatch(packet.server(), backend, delivery); } + private HttpTransportProtocol.Packet decodePacket(HttpTransportProtocol.Packet packet) { + List decoded = new java.util.ArrayList<>(packet.messages().size()); + for (HttpTransportProtocol.Delivery delivery : packet.messages()) { + JsonEnvelope envelope = Objects.requireNonNull(wireCodec.decode(delivery.envelope()), "decoded envelope"); + HttpTransportProtocol.validateEnvelope(envelope); + decoded.add(new HttpTransportProtocol.Delivery(delivery.id(), envelope)); + } + return new HttpTransportProtocol.Packet(packet.server(), packet.session(), packet.sequence(), packet.acks(), + packet.ackConfirmations(), decoded); + } private void dispatch(String serverId, BackendState backend, HttpTransportProtocol.Delivery delivery) { Runnable callback = () -> { boolean success = false; @@ -718,6 +759,11 @@ synchronized Response await(String serverId, String requestedSession, long reque } synchronized Response await(String serverId, String requestedSession, long requestedSequence, Collection ackConfirmations) { + return await(serverId, requestedSession, requestedSequence, ackConfirmations, + HttpEnvelopeWireCodec.identity()); + } + synchronized Response await(String serverId, String requestedSession, long requestedSequence, + Collection ackConfirmations, HttpEnvelopeWireCodec wireCodec) { long deadline = System.nanoTime() + LONG_POLL.toNanos(); while (acknowledgements.isEmpty() && !hasUndelivered()) { long retryRemaining = nanosUntilRedelivery(nanoTime.getAsLong()); @@ -726,15 +772,26 @@ synchronized Response await(String serverId, String requestedSession, long reque long wait = Math.min(requestRemaining, retryRemaining); try { TimeUnit.NANOSECONDS.timedWait(this, wait); } catch (InterruptedException interrupted) { Thread.currentThread().interrupt(); break; } } - List acks = new java.util.ArrayList<>(); while (!acknowledgements.isEmpty() && acks.size() < HttpTransportProtocol.MAX_BATCH) acks.add(acknowledgements.remove()); + List acks = new java.util.ArrayList<>(); + for (String acknowledgement : acknowledgements) { + if (acks.size() == HttpTransportProtocol.MAX_BATCH) break; + acks.add(acknowledgement); + } List candidates = new java.util.ArrayList<>(); long now = nanoTime.getAsLong(); if (hasUndelivered() || redeliveryDue(now)) for (HttpTransportProtocol.Delivery delivery : outgoing.values()) { if (!deliveredAtNanos.containsKey(delivery.id()) || redeliveryDue(delivery.id(), now)) candidates.add(delivery); if (candidates.size() == HttpTransportProtocol.MAX_BATCH) break; } + List encoded = new java.util.ArrayList<>(candidates.size()); + for (HttpTransportProtocol.Delivery delivery : candidates) { + JsonEnvelope envelope = Objects.requireNonNull(wireCodec.encode(delivery.envelope()), "encoded envelope"); + HttpTransportProtocol.validateEnvelope(envelope); + encoded.add(new HttpTransportProtocol.Delivery(delivery.id(), envelope)); + } List messages = HttpTransportProtocol.fittingMessages(serverId, requestedSession, - requestedSequence, acks, ackConfirmations, candidates); + requestedSequence, acks, ackConfirmations, encoded); + for (int index = 0; index < acks.size(); index++) acknowledgements.remove(); long deliveredAt = nanoTime.getAsLong(); for (HttpTransportProtocol.Delivery delivery : messages) deliveredAtNanos.put(delivery.id(), deliveredAt); return new Response(acks, messages); diff --git a/SimpleAPI/src/test/java/com/bencodez/simpleapi/servercomm/http/HttpEnvelopeWireCodecTest.java b/SimpleAPI/src/test/java/com/bencodez/simpleapi/servercomm/http/HttpEnvelopeWireCodecTest.java new file mode 100644 index 0000000..1ce89eb --- /dev/null +++ b/SimpleAPI/src/test/java/com/bencodez/simpleapi/servercomm/http/HttpEnvelopeWireCodecTest.java @@ -0,0 +1,245 @@ +package com.bencodez.simpleapi.servercomm.http; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import com.bencodez.simpleapi.servercomm.codec.JsonEnvelope; +import java.net.InetSocketAddress; +import java.nio.file.Files; +import java.nio.file.Path; +import java.time.Duration; +import java.util.LinkedHashMap; +import java.util.UUID; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +class HttpEnvelopeWireCodecTest { + @TempDir Path directory; + + @Test + void serializedWrappersCoordinateConcurrentAccessToOneCodec() throws Exception { + CountDownLatch firstEntered = new CountDownLatch(1); + CountDownLatch releaseFirst = new CountDownLatch(1); + AtomicInteger active = new AtomicInteger(); + AtomicInteger maximumActive = new AtomicInteger(); + HttpEnvelopeWireCodec delegate = new HttpEnvelopeWireCodec() { + private JsonEnvelope apply(JsonEnvelope envelope) { + int current = active.incrementAndGet(); + maximumActive.accumulateAndGet(current, Math::max); + try { + if (current == 1) { + firstEntered.countDown(); + assertTrue(releaseFirst.await(5, TimeUnit.SECONDS)); + } + return envelope; + } catch (InterruptedException interrupted) { + Thread.currentThread().interrupt(); + throw new IllegalStateException(interrupted); + } finally { + active.decrementAndGet(); + } + } + + @Override public JsonEnvelope encode(JsonEnvelope envelope) { return apply(envelope); } + @Override public JsonEnvelope decode(JsonEnvelope envelope) { return apply(envelope); } + }; + HttpEnvelopeWireCodec first = HttpEnvelopeWireCodec.serialized(delegate); + HttpEnvelopeWireCodec second = HttpEnvelopeWireCodec.serialized(delegate); + JsonEnvelope envelope = JsonEnvelope.builder("test").build(); + var executor = Executors.newFixedThreadPool(2); + try { + var firstCall = executor.submit(() -> first.encode(envelope)); + assertTrue(firstEntered.await(5, TimeUnit.SECONDS)); + CountDownLatch secondAttempted = new CountDownLatch(1); + var secondCall = executor.submit(() -> { + secondAttempted.countDown(); + return second.decode(envelope); + }); + assertTrue(secondAttempted.await(5, TimeUnit.SECONDS)); + assertThrows(java.util.concurrent.TimeoutException.class, + () -> secondCall.get(100, TimeUnit.MILLISECONDS)); + assertEquals(1, maximumActive.get()); + releaseFirst.countDown(); + assertEquals(envelope, firstCall.get(5, TimeUnit.SECONDS)); + assertEquals(envelope, secondCall.get(5, TimeUnit.SECONDS)); + assertEquals(1, maximumActive.get()); + } finally { + releaseFirst.countDown(); + executor.shutdownNow(); + } + } + + @Test + void durableProxyQueueKeepsSemanticEnvelopeAcrossCodecChange() throws Exception { + Path outgoing = directory.resolve("outgoing"); + HttpTlsIdentity identity = HttpTlsIdentity.loadOrCreate(directory.resolve("proxy"), "localhost"); + HttpEnrollmentAuthority authority = new HttpEnrollmentAuthority(identity, directory.resolve("authority")); + String deliveryId = UUID.randomUUID().toString(); + JsonEnvelope semantic = JsonEnvelope.builder("vote").put("player", "Example").build(); + + try (HttpProxyTransportServer first = server(identity, authority, outgoing, codec("old"), ignored -> { })) { + assertTrue(first.send("backend-a", deliveryId, semantic)); + } + + Path persisted; + try (var files = Files.walk(outgoing)) { + persisted = files.filter(Files::isRegularFile).findFirst().orElseThrow(); + } + HttpTransportProtocol.Delivery stored = HttpTransportProtocol.parseStoredDelivery(Files.readAllBytes(persisted)); + assertEquals("vote", stored.envelope().getSubChannel()); + assertEquals("Example", stored.envelope().getFields().get("player")); + assertFalse(stored.envelope().getFields().containsKey("wire-key")); + + HttpEnvelopeWireCodec replacement = codec("new"); + try (HttpProxyTransportServer restarted = server(identity, authority, outgoing, replacement, ignored -> { })) { + HttpProxyTransportServer.Response response = restarted.backendStateForTest("backend-a") + .await("backend-a", UUID.randomUUID().toString(), 1L, java.util.List.of(), replacement); + assertEquals(1, response.messages().size()); + JsonEnvelope wire = response.messages().iterator().next().envelope(); + assertEquals("new", wire.getFields().get("wire-key")); + assertEquals("Example", replacement.decode(wire).getFields().get("player")); + } + } + + @Test + void codecAppliesAtBothAuthenticatedHttpWireBoundaries() throws Exception { + HttpTlsIdentity identity = HttpTlsIdentity.loadOrCreate(directory.resolve("integration-proxy"), "localhost"); + HttpEnrollmentAuthority authority = new HttpEnrollmentAuthority(identity, directory.resolve("integration-authority")); + HttpEnvelopeWireCodec codec = codec("shared"); + CountDownLatch proxyReceived = new CountDownLatch(1), backendReceived = new CountDownLatch(1); + AtomicReference proxyEnvelope = new AtomicReference<>(), backendEnvelope = new AtomicReference<>(); + try (HttpProxyTransportServer server = server(identity, authority, directory.resolve("integration-outgoing"), codec, + received -> { proxyEnvelope.set(received.envelope()); proxyReceived.countDown(); })) { + server.start(); + HttpConnectionCode code = authority.createConnectionCode("backend-a", server.endpoint("localhost"), + Duration.ofMinutes(5)); + Path credentials = directory.resolve("integration-client"); + HttpBackendTransportConnector.enroll(code, "backend-a", credentials); + try (HttpBackendTransportConnector connector = new HttpBackendTransportConnector(credentials, + envelope -> { backendEnvelope.set(envelope); backendReceived.countDown(); }, codec)) { + connector.start(); + assertTrue(connector.awaitFirstResponse(System.nanoTime() + TimeUnit.SECONDS.toNanos(8))); + assertTrue(connector.send(JsonEnvelope.builder("backend-to-proxy").put("value", "one").build())); + assertTrue(proxyReceived.await(8, TimeUnit.SECONDS)); + assertEquals("one", proxyEnvelope.get().getFields().get("value")); + assertFalse(proxyEnvelope.get().getFields().containsKey("wire-key")); + + assertTrue(server.send("backend-a", JsonEnvelope.builder("proxy-to-backend").put("value", "two").build())); + assertTrue(backendReceived.await(8, TimeUnit.SECONDS)); + assertEquals("two", backendEnvelope.get().getFields().get("value")); + assertFalse(backendEnvelope.get().getFields().containsKey("wire-key")); + } + } + } + + @Test + void wireExpansionIsRejectedBeforeItCanBlockAQueue() throws Exception { + Path outgoing = directory.resolve("bounded-outgoing"); + HttpTlsIdentity identity = HttpTlsIdentity.loadOrCreate(directory.resolve("bounded-proxy"), "localhost"); + HttpEnrollmentAuthority authority = new HttpEnrollmentAuthority(identity, directory.resolve("bounded-authority")); + JsonEnvelope large = JsonEnvelope.builder("large").put("value", "x".repeat(46_000)).build(); + HttpEnvelopeWireCodec expanding = expandingCodec(); + + try (HttpProxyTransportServer server = server(identity, authority, outgoing, expanding, ignored -> { })) { + assertFalse(server.send("backend-a", UUID.randomUUID().toString(), large)); + assertTrue(server.send("backend-a", UUID.randomUUID().toString(), + JsonEnvelope.builder("small").build())); + } + + Path restartOutgoing = directory.resolve("restart-bounded-outgoing"); + try (HttpProxyTransportServer first = server(identity, authority, restartOutgoing, + HttpEnvelopeWireCodec.identity(), ignored -> { })) { + assertTrue(first.send("backend-a", UUID.randomUUID().toString(), large)); + } + assertThrows(IllegalArgumentException.class, + () -> server(identity, authority, restartOutgoing, expanding, ignored -> { })); + try (HttpProxyTransportServer recovered = server(identity, authority, restartOutgoing, + HttpEnvelopeWireCodec.identity(), ignored -> { })) { + assertTrue(recovered.hasPendingDeliveries()); + } + + try (HttpProxyTransportServer server = server(identity, authority, directory.resolve("backend-bounded-outgoing"), + HttpEnvelopeWireCodec.identity(), ignored -> { })) { + server.start(); + HttpConnectionCode code = authority.createConnectionCode("backend-b", server.endpoint("localhost"), + Duration.ofMinutes(5)); + Path credentials = directory.resolve("backend-bounded-client"); + HttpBackendTransportConnector.enroll(code, "backend-b", credentials); + try (HttpBackendTransportConnector connector = new HttpBackendTransportConnector(credentials, + ignored -> { }, expanding)) { + connector.start(); + assertTrue(connector.awaitFirstResponse(System.nanoTime() + TimeUnit.SECONDS.toNanos(8))); + assertFalse(connector.send(large)); + assertTrue(connector.send(JsonEnvelope.builder("small").build())); + } + } + } + + @Test + void uncheckedCodecRejectionReturnsFalseAtBothSendBoundaries() throws Exception { + HttpEnvelopeWireCodec rejecting = new HttpEnvelopeWireCodec() { + @Override public JsonEnvelope encode(JsonEnvelope envelope) { + throw new IllegalStateException("codec rejected envelope"); + } + + @Override public JsonEnvelope decode(JsonEnvelope envelope) { return envelope; } + }; + HttpTlsIdentity identity = HttpTlsIdentity.loadOrCreate(directory.resolve("rejecting-proxy"), "localhost"); + HttpEnrollmentAuthority authority = new HttpEnrollmentAuthority(identity, directory.resolve("rejecting-authority")); + try (HttpProxyTransportServer server = server(identity, authority, directory.resolve("rejecting-outgoing"), + rejecting, ignored -> { })) { + assertFalse(server.send("backend-a", UUID.randomUUID().toString(), JsonEnvelope.builder("test").build())); + server.start(); + HttpConnectionCode code = authority.createConnectionCode("backend-a", server.endpoint("localhost"), + Duration.ofMinutes(5)); + Path credentials = directory.resolve("rejecting-client"); + HttpBackendTransportConnector.enroll(code, "backend-a", credentials); + try (HttpBackendTransportConnector connector = new HttpBackendTransportConnector(credentials, + ignored -> { }, rejecting)) { + connector.start(); + assertTrue(connector.awaitFirstResponse(System.nanoTime() + TimeUnit.SECONDS.toNanos(8))); + assertFalse(connector.send(JsonEnvelope.builder("test").build())); + } + } + } + + private static HttpProxyTransportServer server(HttpTlsIdentity identity, HttpEnrollmentAuthority authority, + Path outgoing, HttpEnvelopeWireCodec codec, + java.util.function.Consumer consumer) throws Exception { + return new HttpProxyTransportServer(new InetSocketAddress("localhost", 0), identity, authority, outgoing, + consumer, (serverId, deliveryId) -> { }, codec); + } + + private static HttpEnvelopeWireCodec codec(String key) { + return new HttpEnvelopeWireCodec() { + @Override public JsonEnvelope encode(JsonEnvelope envelope) { + return envelope.toBuilder().put("wire-key", key).build(); + } + + @Override public JsonEnvelope decode(JsonEnvelope envelope) { + if (!key.equals(envelope.getFields().get("wire-key"))) + throw new IllegalArgumentException("wire key mismatch"); + LinkedHashMap fields = new LinkedHashMap<>(envelope.getFields()); + fields.remove("wire-key"); + return new JsonEnvelope(envelope.getSubChannel(), envelope.getSchema(), fields); + } + }; + } + + private static HttpEnvelopeWireCodec expandingCodec() { + return new HttpEnvelopeWireCodec() { + @Override public JsonEnvelope encode(JsonEnvelope envelope) { + return envelope.toBuilder().put("overhead", "y".repeat(4_000)).build(); + } + + @Override public JsonEnvelope decode(JsonEnvelope envelope) { return envelope; } + }; + } +}