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 @@ -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;
Expand Down Expand Up @@ -50,6 +51,7 @@ public final class HttpBackendTransportConnector implements AutoCloseable {
private volatile HttpClientCredentialStore.HttpClientProfile profile;
private final String serverId;
private final Consumer<JsonEnvelope> onEnvelope;
private final HttpEnvelopeWireCodec wireCodec;
private volatile HttpClient client;
private volatile HttpClientCredentialStore.ClientCredential credential;
private final Path credentialDirectory;
Expand Down Expand Up @@ -105,14 +107,23 @@ public final class HttpBackendTransportConnector implements AutoCloseable {

private HttpBackendTransportConnector(HttpClientCredentialStore.EnrolledClient enrolled, Consumer<JsonEnvelope> 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<JsonEnvelope> 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<JsonEnvelope> 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;
Expand Down Expand Up @@ -165,6 +176,21 @@ public HttpBackendTransportConnector(Path credentials, Consumer<JsonEnvelope> 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<JsonEnvelope> onEnvelope,
HttpEnvelopeWireCodec wireCodec) throws Exception {
this(HttpClientCredentialStore.loadEnrolled(credentials), onEnvelope, credentials, wireCodec);
}

private HttpBackendTransportConnector(HttpClientCredentialStore.EnrolledClient enrolled,
Consumer<JsonEnvelope> 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");
Expand Down Expand Up @@ -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;
Expand All @@ -279,8 +310,16 @@ private boolean pollOnce(Duration timeout, boolean requireRunning, boolean accep
acks = first(acknowledgements);
ackConfirmations = first(acknowledgementConfirmations);
requestSequence = sequence++;
List<HttpTransportProtocol.Delivery> 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();
}
Expand All @@ -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;
Expand All @@ -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<HttpTransportProtocol.Delivery> 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;
Expand Down
Original file line number Diff line number Diff line change
@@ -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.
*
* <p>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.</p>
*/
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);
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -84,14 +84,16 @@ public final class HttpProxyTransportServer implements AutoCloseable {
private final Path durableIncomingRoot;
private final Consumer<ReceivedEnvelope> onEnvelope;
private final DeliveryAcknowledgement onAcknowledged;
private final HttpEnvelopeWireCodec wireCodec;
private final LongSupplier nanoTime;
private volatile boolean closed;
private boolean closeFinalizing, closeFinalized;

/** In-memory constructor for tests; production callers must supply a durable state directory. */
HttpProxyTransportServer(InetSocketAddress bind, HttpTlsIdentity identity, HttpEnrollmentAuthority authority,
Consumer<ReceivedEnvelope> 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,
Expand All @@ -103,17 +105,37 @@ public HttpProxyTransportServer(InetSocketAddress bind, HttpTlsIdentity identity
Path outgoingDirectory, Consumer<ReceivedEnvelope> 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<ReceivedEnvelope> 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<ReceivedEnvelope> 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<ReceivedEnvelope> 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;
Expand All @@ -138,6 +160,7 @@ public HttpProxyTransportServer(InetSocketAddress bind, HttpTlsIdentity identity
}
if (durableOutgoing != null) for (Map.Entry<String, List<HttpTransportProtocol.Delivery>> pending
: durableOutgoing.load().entrySet()) {
for (HttpTransportProtocol.Delivery delivery : pending.getValue()) validateForWire(delivery.envelope());
BackendState state = backendState(pending.getKey());
state.restore(pending.getValue());
}
Expand Down Expand Up @@ -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); }
Expand All @@ -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);
Comment thread
BenCodez marked this conversation as resolved.
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();
Expand Down Expand Up @@ -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(); }
Expand All @@ -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<HttpTransportProtocol.Delivery> 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;
Expand Down Expand Up @@ -718,6 +759,11 @@ synchronized Response await(String serverId, String requestedSession, long reque
}
synchronized Response await(String serverId, String requestedSession, long requestedSequence,
Collection<String> ackConfirmations) {
return await(serverId, requestedSession, requestedSequence, ackConfirmations,
HttpEnvelopeWireCodec.identity());
}
synchronized Response await(String serverId, String requestedSession, long requestedSequence,
Collection<String> ackConfirmations, HttpEnvelopeWireCodec wireCodec) {
long deadline = System.nanoTime() + LONG_POLL.toNanos();
while (acknowledgements.isEmpty() && !hasUndelivered()) {
long retryRemaining = nanosUntilRedelivery(nanoTime.getAsLong());
Expand All @@ -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<String> acks = new java.util.ArrayList<>(); while (!acknowledgements.isEmpty() && acks.size() < HttpTransportProtocol.MAX_BATCH) acks.add(acknowledgements.remove());
List<String> acks = new java.util.ArrayList<>();
for (String acknowledgement : acknowledgements) {
if (acks.size() == HttpTransportProtocol.MAX_BATCH) break;
acks.add(acknowledgement);
}
List<HttpTransportProtocol.Delivery> 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<HttpTransportProtocol.Delivery> encoded = new java.util.ArrayList<>(candidates.size());
for (HttpTransportProtocol.Delivery delivery : candidates) {
JsonEnvelope envelope = Objects.requireNonNull(wireCodec.encode(delivery.envelope()), "encoded envelope");
Comment thread
BenCodez marked this conversation as resolved.
HttpTransportProtocol.validateEnvelope(envelope);
encoded.add(new HttpTransportProtocol.Delivery(delivery.id(), envelope));
}
List<HttpTransportProtocol.Delivery> 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);
Expand Down
Loading
Loading