From faced0708caa0dffb886f60b09fd593fa7227084 Mon Sep 17 00:00:00 2001 From: uartnet <140632163+uartnet@users.noreply.github.com> Date: Tue, 21 Jul 2026 22:43:29 +0200 Subject: [PATCH 1/6] feat(edge): support region-aware tunnels --- README.md | 2 + docs/CONFIGURATION.md | 8 + src/main/java/io/rstream/ClientOptions.java | 18 ++ src/main/java/io/rstream/ConfigResolver.java | 179 +++++++++++++++- src/main/java/io/rstream/ControlChannel.java | 6 + .../java/io/rstream/CreateTunnelOptions.java | 8 + src/main/java/io/rstream/Protocol.java | 7 +- .../io/rstream/ResolvedClientOptions.java | 3 + .../java/io/rstream/RstreamApiClient.java | 84 +++++++- src/main/java/io/rstream/RstreamClient.java | 43 +++- .../java/io/rstream/RstreamTransport.java | 13 +- .../java/io/rstream/TunnelProperties.java | 12 +- src/main/proto/rstream.proto | 4 +- .../java/io/rstream/ConfigResolverTest.java | 129 ++++++++++++ src/test/java/io/rstream/ProtocolTest.java | 1 + .../java/io/rstream/RstreamApiClientTest.java | 132 ++++++++++++ .../java/io/rstream/RstreamTransportTest.java | 10 + .../java/io/rstream/RuntimeFakeEngineIT.java | 192 ++++++++++++------ 18 files changed, 771 insertions(+), 80 deletions(-) diff --git a/README.md b/README.md index 446ba84..53e7df5 100644 --- a/README.md +++ b/README.md @@ -93,6 +93,8 @@ SDKs: | `RSTREAM_MTLS_CERT_FILE` | Client certificate file for mTLS authentication. | | `RSTREAM_MTLS_KEY_FILE` | Client private key file for mTLS authentication. | | `RSTREAM_API_URL` | Control plane API URL for managed project discovery. | +| `RSTREAM_REGION` | Authorized region to select for a managed project. | +| `RSTREAM_CONTROL_PLANE_HEADERS` | Additional Control plane request headers encoded as a JSON object. | `RSTREAM_ENGINE_ADDRESS` is also accepted for compatibility with older local SDK workflows. Prefer `RSTREAM_ENGINE` in new code. diff --git a/docs/CONFIGURATION.md b/docs/CONFIGURATION.md index 787147f..0cb1c23 100644 --- a/docs/CONFIGURATION.md +++ b/docs/CONFIGURATION.md @@ -26,12 +26,20 @@ Set `RSTREAM_CONFIG` to use another file. | `RSTREAM_MTLS_CERT_FILE` | mTLS client certificate path. | | `RSTREAM_MTLS_KEY_FILE` | mTLS client key path. | | `RSTREAM_API_URL` | Control plane API URL for managed project discovery. | +| `RSTREAM_REGION` | Authorized region to select for a managed project. | +| `RSTREAM_CONTROL_PLANE_HEADERS` | Additional Control plane request headers encoded as a JSON object. | | `RSTREAM_TUNNEL_TRANSPORT` | `auto`, `tls`, or `quic`. Java maps `auto` to TLS and rejects explicit `quic`. | | `RSTREAM_QUIC_TRANSPORT` | Legacy selector. Prefer `RSTREAM_TUNNEL_TRANSPORT`. | The SDK also accepts `RSTREAM_ENGINE_ADDRESS` for compatibility with older local SDK workflows. Prefer `RSTREAM_ENGINE` in new code. +Region selection requires a managed project endpoint and cannot be combined +with an explicit engine override. Control plane headers may satisfy a separate +deployment access layer. Authentication, forwarding, and hop-by-hop headers are +reserved; malformed values and case-insensitive duplicates are rejected before +network I/O. + ## Config file ```yaml diff --git a/src/main/java/io/rstream/ClientOptions.java b/src/main/java/io/rstream/ClientOptions.java index 5c1d1c7..a7b4a58 100644 --- a/src/main/java/io/rstream/ClientOptions.java +++ b/src/main/java/io/rstream/ClientOptions.java @@ -1,12 +1,14 @@ package io.rstream; import java.time.Duration; +import java.util.Map; /** Options accepted by {@link RstreamClient}. */ public record ClientOptions( String apiUrl, String configPath, String context, + Map controlPlaneHeaders, String engine, boolean heartbeat, Duration heartbeatInterval, @@ -15,11 +17,13 @@ public record ClientOptions( Boolean noToken, String projectEndpoint, boolean readConfigFile, + String region, boolean requireToken, String token, TlsOptions tls, boolean zeroRtt) { public ClientOptions { + controlPlaneHeaders = controlPlaneHeaders == null ? Map.of() : Map.copyOf(controlPlaneHeaders); heartbeatInterval = heartbeatInterval == null ? Duration.ofSeconds(5) : heartbeatInterval; connectTimeout = connectTimeout == null ? Duration.ofSeconds(15) : connectTimeout; operationTimeout = operationTimeout == null ? Duration.ofSeconds(30) : operationTimeout; @@ -43,6 +47,7 @@ public static final class Builder { private String apiUrl; private String configPath; private String context; + private Map controlPlaneHeaders = Map.of(); private String engine; private boolean heartbeat = true; private Duration heartbeatInterval = Duration.ofSeconds(5); @@ -51,6 +56,7 @@ public static final class Builder { private Boolean noToken; private String projectEndpoint; private boolean readConfigFile = true; + private String region; private boolean requireToken; private String token; private TlsOptions tls; @@ -71,6 +77,11 @@ public Builder context(String context) { return this; } + public Builder controlPlaneHeaders(Map controlPlaneHeaders) { + this.controlPlaneHeaders = Map.copyOf(controlPlaneHeaders); + return this; + } + public Builder engine(String engine) { this.engine = engine; return this; @@ -111,6 +122,11 @@ public Builder readConfigFile(boolean readConfigFile) { return this; } + public Builder region(String region) { + this.region = region; + return this; + } + public Builder requireToken(boolean requireToken) { this.requireToken = requireToken; return this; @@ -136,6 +152,7 @@ public ClientOptions build() { apiUrl, configPath, context, + controlPlaneHeaders, engine, heartbeat, heartbeatInterval, @@ -144,6 +161,7 @@ public ClientOptions build() { noToken, projectEndpoint, readConfigFile, + region, requireToken, token, tls, diff --git a/src/main/java/io/rstream/ConfigResolver.java b/src/main/java/io/rstream/ConfigResolver.java index 0ce10c1..719ac1c 100644 --- a/src/main/java/io/rstream/ConfigResolver.java +++ b/src/main/java/io/rstream/ConfigResolver.java @@ -1,5 +1,8 @@ package io.rstream; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; import java.io.IOException; import java.net.URI; import java.net.URISyntaxException; @@ -8,13 +11,35 @@ import java.nio.file.Path; import java.time.Instant; import java.util.Base64; +import java.util.LinkedHashMap; import java.util.List; +import java.util.Locale; import java.util.Map; +import java.util.Set; +import java.util.regex.Pattern; import org.snakeyaml.engine.v2.api.Load; import org.snakeyaml.engine.v2.api.LoadSettings; final class ConfigResolver { static final String DEFAULT_API_URL = "https://rstream.io"; + private static final ObjectMapper JSON = new ObjectMapper(); + private static final Pattern HEADER_NAME = Pattern.compile("^[!#$%&'*+.^_`|~0-9A-Za-z-]+$"); + private static final Pattern REGION = Pattern.compile("^[a-z0-9](?:[a-z0-9._-]{0,62}[a-z0-9])?$"); + private static final Set RESERVED_CONTROL_PLANE_HEADERS = + Set.of( + "authorization", + "connection", + "content-length", + "cookie", + "forwarded", + "host", + "keep-alive", + "proxy-authorization", + "proxy-connection", + "te", + "trailer", + "transfer-encoding", + "upgrade"); private ConfigResolver() {} @@ -49,18 +74,35 @@ static ResolvedClientOptions resolve(ClientOptions options, Map throw new ConfigurationException( "Authentication is required but not configured.", "ERR_RSTREAM_AUTH_REQUIRED"); } + var region = normalizeRegion(firstDefined(options.region(), env.region(), config.region())); + var explicitEngine = normalizeOptional(firstDefined(options.engine(), env.engine())); + var projectEndpoint = + normalizeOptional(firstDefined(options.projectEndpoint(), config.projectEndpoint())); + if (region != null && explicitEngine != null) { + throw new ConfigurationException( + "Region selection cannot be combined with an explicit engine override.", + "ERR_RSTREAM_REGION_ENGINE_CONFLICT"); + } + if (region != null && projectEndpoint == null) { + throw new ConfigurationException( + "Managed project endpoint is required for region selection.", + "ERR_RSTREAM_PROJECT_ENDPOINT_REQUIRED"); + } var engine = - normalizeEngine( - firstDefined(options.engine(), env.engine(), config.contextEngine(), config.engine())); + region == null + ? normalizeEngine(firstDefined(explicitEngine, config.contextEngine(), config.engine())) + : null; return new ResolvedClientOptions( firstDefined(options.apiUrl(), env.apiUrl(), config.apiUrl(), DEFAULT_API_URL), + config.controlPlaneHeaders(), engine, options.heartbeat(), options.heartbeatInterval(), options.connectTimeout(), options.operationTimeout(), options.noToken() != null ? options.noToken() : token == null && !hasClientCertificate(tls), - normalizeOptional(firstDefined(options.projectEndpoint(), config.projectEndpoint())), + projectEndpoint, + region, tls, token, "tls", @@ -80,6 +122,8 @@ private static ResolvedConfig resolveConfig(ClientOptions options, EnvironmentSe if (!options.readConfigFile()) { return new ResolvedConfig( firstDefined(options.apiUrl(), env.apiUrl(), DEFAULT_API_URL), + mergeControlPlaneHeaders(env.controlPlaneHeaders(), options.controlPlaneHeaders()), + null, null, null, null, @@ -136,9 +180,14 @@ && engineOverrideUsesStoredAuth(explicitEngine, context)) { tunnelTransport = transportMode(environment == null ? null : environment.transport()); return new ResolvedConfig( apiUrl, + mergeControlPlaneHeaders( + environment == null ? null : environment.headers(), + env.controlPlaneHeaders(), + options.controlPlaneHeaders()), context == null ? null : context.engine(), explicitEngine, context == null ? null : context.projectEndpoint(), + context == null ? null : context.region(), tls, token, tunnelTransport); @@ -188,6 +237,7 @@ private static ContextConfig contextConfig(Map value) { authConfig(value.get("auth")), normalizeOptional(string(value.get("engine"))), normalizeOptional(string(value.get("projectEndpoint"))), + normalizeOptional(string(value.get("region"))), transportConfig(value.get("transport"))); } @@ -195,6 +245,7 @@ private static EnvironmentConfig environmentConfig(Map value) { return new EnvironmentConfig( normalizeApiUrl(string(value.get("apiUrl"))), authConfig(value.get("auth")), + controlPlaneHeaders(value.get("headers")), transportConfig(value.get("transport"))); } @@ -384,6 +435,99 @@ private static TlsOptions mergeTls(TlsOptions inherited, TlsOptions explicit) { .build(); } + @SafeVarargs + static Map mergeControlPlaneHeaders(Map... sources) { + var merged = new LinkedHashMap(); + for (var source : sources) merged.putAll(normalizeControlPlaneHeaders(source)); + return Map.copyOf(merged); + } + + static Map normalizeControlPlaneHeaders(Map headers) { + if (headers == null || headers.isEmpty()) return Map.of(); + var normalized = new LinkedHashMap(); + for (var entry : headers.entrySet()) { + var rawName = entry.getKey(); + var value = entry.getValue(); + if (rawName == null || value == null) { + throw invalidControlPlaneHeaders(); + } + var name = rawName.trim(); + var lowerName = name.toLowerCase(Locale.ROOT); + if (!HEADER_NAME.matcher(name).matches()) { + throw new ConfigurationException( + "Invalid control plane header name '" + rawName + "'.", "ERR_RSTREAM_INVALID_CONFIG"); + } + if (RESERVED_CONTROL_PLANE_HEADERS.contains(lowerName) + || lowerName.startsWith("x-forwarded-")) { + throw new ConfigurationException( + "Control plane header '" + rawName + "' is reserved.", "ERR_RSTREAM_INVALID_CONFIG"); + } + if (value.indexOf('\r') >= 0 || value.indexOf('\n') >= 0) { + throw new ConfigurationException( + "Control plane header '" + rawName + "' has an invalid value.", + "ERR_RSTREAM_INVALID_CONFIG"); + } + var canonicalName = canonicalHeaderName(name); + if (normalized.containsKey(canonicalName)) { + throw new ConfigurationException( + "Duplicate control plane header '" + canonicalName + "'.", + "ERR_RSTREAM_INVALID_CONFIG"); + } + normalized.put(canonicalName, value); + } + return Map.copyOf(normalized); + } + + private static Map controlPlaneHeaders(Object value) { + if (value == null) return Map.of(); + if (!(value instanceof Map values)) throw invalidControlPlaneHeaders(); + var headers = new LinkedHashMap(); + for (var entry : values.entrySet()) { + if (!(entry.getKey() instanceof String name) + || !(entry.getValue() instanceof String headerValue)) { + throw invalidControlPlaneHeaders(); + } + headers.put(name, headerValue); + } + return normalizeControlPlaneHeaders(headers); + } + + private static Map controlPlaneHeadersFromJson(String value) { + value = normalizeOptional(value); + if (value == null) return Map.of(); + try { + JsonNode object = JSON.readTree(value); + if (object == null || !object.isObject()) throw invalidControlPlaneHeaders(); + var headers = new LinkedHashMap(); + for (var field : object.properties()) { + if (!field.getValue().isTextual()) throw invalidControlPlaneHeaders(); + headers.put(field.getKey(), field.getValue().textValue()); + } + return normalizeControlPlaneHeaders(headers); + } catch (JsonProcessingException error) { + throw new ConfigurationException( + "RSTREAM_CONTROL_PLANE_HEADERS must be a JSON object of string values.", + "ERR_RSTREAM_INVALID_CONFIG", + error); + } + } + + private static String canonicalHeaderName(String name) { + var canonical = new StringBuilder(name.length()); + var upper = true; + for (var index = 0; index < name.length(); index++) { + var character = name.charAt(index); + canonical.append(upper ? Character.toUpperCase(character) : Character.toLowerCase(character)); + upper = character == '-'; + } + return canonical.toString(); + } + + private static ConfigurationException invalidControlPlaneHeaders() { + return new ConfigurationException( + "Control plane headers must be an object of string values.", "ERR_RSTREAM_INVALID_CONFIG"); + } + private static void validateTokenExpiry(String token) { var parts = token.split("\\."); if (parts.length < 2) return; @@ -443,12 +587,23 @@ private static String normalizeOptional(String value) { return normalized.isEmpty() ? null : normalized; } + private static String normalizeRegion(String value) { + var normalized = normalizeOptional(value); + if (normalized == null || normalized.equalsIgnoreCase("auto")) return null; + var region = normalized.toLowerCase(Locale.ROOT); + if (!REGION.matcher(region).matches()) { + throw new ConfigurationException( + "Region can only contain letters, numbers, dots, underscores, or hyphens.", + "ERR_RSTREAM_INVALID_REGION"); + } + return region; + } + private static Map map(Object value) { if (!(value instanceof Map input)) return null; - return input.entrySet().stream() - .collect( - java.util.stream.Collectors.toMap( - entry -> String.valueOf(entry.getKey()), Map.Entry::getValue)); + var output = new LinkedHashMap(); + input.forEach((key, item) -> output.put(String.valueOf(key), item)); + return output; } private static List> records(Object value) { @@ -469,9 +624,11 @@ private record EnvironmentSettings( String apiUrl, String configPath, String context, + Map controlPlaneHeaders, String engine, String mtlsCert, String mtlsKey, + String region, String token, String tunnelTransport, Boolean useQuic) { @@ -480,11 +637,13 @@ static EnvironmentSettings read(Map environment) { normalizeApiUrl(environment.get("RSTREAM_API_URL")), normalizeOptional(environment.get("RSTREAM_CONFIG")), normalizeOptional(environment.get("RSTREAM_CONTEXT")), + controlPlaneHeadersFromJson(environment.get("RSTREAM_CONTROL_PLANE_HEADERS")), normalizeOptional( firstDefined( environment.get("RSTREAM_ENGINE"), environment.get("RSTREAM_ENGINE_ADDRESS"))), normalizeOptional(environment.get("RSTREAM_MTLS_CERT_FILE")), normalizeOptional(environment.get("RSTREAM_MTLS_KEY_FILE")), + normalizeOptional(environment.get("RSTREAM_REGION")), normalizeOptional(environment.get("RSTREAM_AUTHENTICATION_TOKEN")), normalizeOptional(environment.get("RSTREAM_TUNNEL_TRANSPORT")), legacyQuic(environment.get("RSTREAM_QUIC_TRANSPORT"))); @@ -509,9 +668,11 @@ private record ContextConfig( AuthConfig auth, String engine, String projectEndpoint, + String region, TransportConfig transport) {} - private record EnvironmentConfig(String apiUrl, AuthConfig auth, TransportConfig transport) {} + private record EnvironmentConfig( + String apiUrl, AuthConfig auth, Map headers, TransportConfig transport) {} private record AuthConfig(TokenConfig token, MtlsConfig mtls) {} @@ -528,9 +689,11 @@ private record TransportConfig(Map raw) {} private record ResolvedConfig( String apiUrl, + Map controlPlaneHeaders, String contextEngine, String engine, String projectEndpoint, + String region, TlsOptions tls, String token, String tunnelTransport) {} diff --git a/src/main/java/io/rstream/ControlChannel.java b/src/main/java/io/rstream/ControlChannel.java index d75ff28..5134c87 100644 --- a/src/main/java/io/rstream/ControlChannel.java +++ b/src/main/java/io/rstream/ControlChannel.java @@ -312,6 +312,11 @@ private ScheduledFuture operationTimeout( } private static TunnelProperties normalizeBytestreamOptions(CreateTunnelOptions options) { + if (options.allowCrossRegionRouting() != null && options.protocol() != TunnelProtocol.TCP) { + throw new RstreamException( + "Cross-region routing policy requires the TCP protocol.", + "ERR_RSTREAM_INVALID_TUNNEL_OPTIONS"); + } if (options.httpVersion() == HttpVersion.H3) { throw new UnsupportedFeatureException( "HTTP/3 tunnels require datagram support, which rstream-java does not support.", @@ -384,6 +389,7 @@ private static TunnelProperties normalizeBytestreamOptions(CreateTunnelOptions o .hostname(options.hostname()) .port(options.port()) .upstreamTls(options.upstreamTls()) + .allowCrossRegionRouting(options.allowCrossRegionRouting()) .build(); } diff --git a/src/main/java/io/rstream/CreateTunnelOptions.java b/src/main/java/io/rstream/CreateTunnelOptions.java index 1108176..5a48885 100644 --- a/src/main/java/io/rstream/CreateTunnelOptions.java +++ b/src/main/java/io/rstream/CreateTunnelOptions.java @@ -24,6 +24,7 @@ public record CreateTunnelOptions( String hostname, Integer port, Boolean upstreamTls, + Boolean allowCrossRegionRouting, TunnelAuth auth) { public CreateTunnelOptions { labels = labels == null ? Map.of() : Map.copyOf(labels); @@ -61,6 +62,7 @@ public static final class Builder { private String hostname; private Integer port; private Boolean upstreamTls; + private Boolean allowCrossRegionRouting; private TunnelAuth auth; public Builder name(String name) { @@ -158,6 +160,11 @@ public Builder upstreamTls(Boolean upstreamTls) { return this; } + public Builder allowCrossRegionRouting(Boolean allowCrossRegionRouting) { + this.allowCrossRegionRouting = allowCrossRegionRouting; + return this; + } + public Builder auth(TunnelAuth auth) { this.auth = auth; return this; @@ -184,6 +191,7 @@ public CreateTunnelOptions build() { hostname, port, upstreamTls, + allowCrossRegionRouting, auth); } } diff --git a/src/main/java/io/rstream/Protocol.java b/src/main/java/io/rstream/Protocol.java index e5c9723..b7b1368 100644 --- a/src/main/java/io/rstream/Protocol.java +++ b/src/main/java/io/rstream/Protocol.java @@ -141,6 +141,8 @@ static Rstream.TunnelProperties tunnelPropertiesToPb(TunnelProperties properties builder.setUpstreamTls(boolValue(properties.upstreamTls())); if (properties.datagramGuaranteedDelivery() != null) builder.setDatagramGuaranteedDelivery(boolValue(properties.datagramGuaranteedDelivery())); + if (properties.allowCrossRegionRouting() != null) + builder.setAllowCrossRegionRouting(boolValue(properties.allowCrossRegionRouting())); return builder.build(); } @@ -173,8 +175,9 @@ static TunnelProperties tunnelPropertiesFromPb(Rstream.TunnelProperties properti wrapperInt(properties.hasPort(), properties.getPort()), wrapperBool(properties.hasUpstreamTls(), properties.getUpstreamTls()), wrapperBool( - properties.hasDatagramGuaranteedDelivery(), - properties.getDatagramGuaranteedDelivery())); + properties.hasDatagramGuaranteedDelivery(), properties.getDatagramGuaranteedDelivery()), + wrapperBool( + properties.hasAllowCrossRegionRouting(), properties.getAllowCrossRegionRouting())); } static ServerDetails serverDetailsFromPb(Rstream.ServerDetails details) { diff --git a/src/main/java/io/rstream/ResolvedClientOptions.java b/src/main/java/io/rstream/ResolvedClientOptions.java index 7b5477b..758335f 100644 --- a/src/main/java/io/rstream/ResolvedClientOptions.java +++ b/src/main/java/io/rstream/ResolvedClientOptions.java @@ -1,9 +1,11 @@ package io.rstream; import java.time.Duration; +import java.util.Map; record ResolvedClientOptions( String apiUrl, + Map controlPlaneHeaders, String engine, boolean heartbeat, Duration heartbeatInterval, @@ -11,6 +13,7 @@ record ResolvedClientOptions( Duration operationTimeout, boolean noToken, String projectEndpoint, + String region, TlsOptions tls, String token, String tunnelTransport, diff --git a/src/main/java/io/rstream/RstreamApiClient.java b/src/main/java/io/rstream/RstreamApiClient.java index 28905d8..47e3856 100644 --- a/src/main/java/io/rstream/RstreamApiClient.java +++ b/src/main/java/io/rstream/RstreamApiClient.java @@ -10,15 +10,25 @@ import java.net.http.HttpResponse; import java.nio.charset.StandardCharsets; import java.time.Duration; +import java.util.ArrayList; +import java.util.Locale; +import java.util.Map; +import java.util.TreeSet; final class RstreamApiClient { private static final ObjectMapper JSON = new ObjectMapper(); private final String apiUrl; + private final Map controlPlaneHeaders; private final String token; private final HttpClient client; RstreamApiClient(String apiUrl, String token) { + this(apiUrl, token, Map.of()); + } + + RstreamApiClient(String apiUrl, String token, Map controlPlaneHeaders) { this.apiUrl = apiUrl == null ? ConfigResolver.DEFAULT_API_URL : apiUrl.replaceAll("/+$", ""); + this.controlPlaneHeaders = ConfigResolver.normalizeControlPlaneHeaders(controlPlaneHeaders); this.token = token; this.client = HttpClient.newBuilder() @@ -28,6 +38,10 @@ final class RstreamApiClient { } String resolveEngine(String projectEndpoint) { + return resolveEngine(projectEndpoint, null); + } + + String resolveEngine(String projectEndpoint, String region) { var endpoint = projectEndpoint == null ? "" : projectEndpoint.trim(); if (endpoint.isEmpty()) { throw new ConfigurationException( @@ -35,7 +49,7 @@ String resolveEngine(String projectEndpoint) { } var encoded = URLEncoder.encode(endpoint, StandardCharsets.UTF_8).replace("+", "%20"); var project = requestJson("/api/projects/tunnels/resolve/" + encoded); - return engineFromProject(project); + return engineFromProject(project, region); } private JsonNode requestJson(String path) { @@ -45,6 +59,7 @@ private JsonNode requestJson(String path) { } var builder = HttpRequest.newBuilder(URI.create(apiUrl + path)).timeout(Duration.ofSeconds(15)).GET(); + controlPlaneHeaders.forEach(builder::header); if (token != null) builder.header("Authorization", "Bearer " + token); try { var response = @@ -68,7 +83,8 @@ private JsonNode requestJson(String path) { } } - private static String engineFromProject(JsonNode project) { + private static String engineFromProject(JsonNode project, String region) { + if (region != null) return regionalEngineFromProject(project, region); var endpoint = optionalString(project, "endpoint"); var domain = optionalString(project, "domain"); var port = optionalInt(project, "enginePort"); @@ -81,6 +97,70 @@ private static String engineFromProject(JsonNode project) { "ERR_RSTREAM_ENGINE_RESOLUTION"); } + private static String regionalEngineFromProject(JsonNode project, String region) { + var requested = region.toLowerCase(Locale.ROOT); + var projectEndpoint = optionalString(project, "endpoint"); + if (projectEndpoint == null) { + throw new RstreamException( + "Managed tunnels project response is missing its endpoint.", + "ERR_RSTREAM_API_INVALID_RESPONSE"); + } + var endpoints = project.get("regionalEndpoints"); + if (endpoints == null || !endpoints.isArray()) { + throw new RstreamException( + "Control plane response has invalid 'regionalEndpoints'.", + "ERR_RSTREAM_API_INVALID_RESPONSE"); + } + var matches = new ArrayList(); + var available = new TreeSet(); + for (var endpoint : endpoints) { + if (!endpoint.isObject()) { + throw new RstreamException( + "Control plane response has invalid 'regionalEndpoints'.", + "ERR_RSTREAM_API_INVALID_RESPONSE"); + } + var endpointRegion = requiredString(endpoint, "region").toLowerCase(Locale.ROOT); + requiredString(endpoint, "provider"); + requiredString(endpoint, "domain"); + requiredPort(endpoint, "enginePort"); + available.add(endpointRegion); + if (endpointRegion.equals(requested)) matches.add(endpoint); + } + if (matches.isEmpty()) { + var suffix = + available.isEmpty() ? "" : " Available regions: " + String.join(", ", available) + "."; + throw new ConfigurationException( + "Region '" + requested + "' is not available for this project." + suffix, + "ERR_RSTREAM_REGION_UNAVAILABLE"); + } + if (matches.size() > 1) { + throw new ConfigurationException( + "Region '" + requested + "' is ambiguous for this project.", + "ERR_RSTREAM_REGION_AMBIGUOUS"); + } + var selected = matches.get(0); + return ConfigResolver.normalizeEngine( + projectEndpoint + + "." + + requiredString(selected, "domain") + + ":" + + requiredPort(selected, "enginePort")); + } + + private static String requiredString(JsonNode data, String key) { + var value = optionalString(data, key); + if (value != null) return value; + throw new RstreamException( + "Control plane response has invalid '" + key + "'.", "ERR_RSTREAM_API_INVALID_RESPONSE"); + } + + private static int requiredPort(JsonNode data, String key) { + var value = optionalInt(data, key); + if (value != null) return value; + throw new RstreamException( + "Control plane response has invalid '" + key + "'.", "ERR_RSTREAM_API_INVALID_RESPONSE"); + } + private static String optionalString(JsonNode data, String key) { var value = data.get(key); if (value != null && value.isTextual() && !value.asText().isBlank()) return value.asText(); diff --git a/src/main/java/io/rstream/RstreamClient.java b/src/main/java/io/rstream/RstreamClient.java index b7c2267..7648c93 100644 --- a/src/main/java/io/rstream/RstreamClient.java +++ b/src/main/java/io/rstream/RstreamClient.java @@ -180,11 +180,25 @@ private static RstreamException closedError() { private RstreamStream openProxyConnection( String engine, ResolvedClientOptions resolvedOptions, Rstream.ProxyConnReq request) { var token = request.hasSecret() ? request.getSecret().getValue() : null; + var proxyEngine = engine; + if (request.hasProxyEndpoint()) { + if (token == null || token.isBlank()) { + throw new ProtocolException( + "Engine did not provide credentials for the redirected stream.", + "ERR_RSTREAM_PROTOCOL"); + } + proxyEngine = request.getProxyEndpoint().getValue().trim(); + if (proxyEngine.isEmpty()) { + throw new ProtocolException( + "Engine returned an empty proxy endpoint.", "ERR_RSTREAM_PROTOCOL"); + } + } var socket = openStreamSocket( - engine, + proxyEngine, resolvedOptions, - Protocol.proxyRequest(request.getStreamId(), token, resolvedOptions.zeroRtt())); + Protocol.proxyRequest(request.getStreamId(), token, resolvedOptions.zeroRtt()), + proxyEngine.equals(engine)); if (!resolvedOptions.zeroRtt()) { try { var response = Protocol.readMessage(socket.getInputStream()); @@ -207,9 +221,22 @@ private RstreamStream openProxyConnection( private Socket openStreamSocket( String engine, ResolvedClientOptions resolvedOptions, Rstream.Message request) { + return openStreamSocket(engine, resolvedOptions, request, true); + } + + private Socket openStreamSocket( + String engine, + ResolvedClientOptions resolvedOptions, + Rstream.Message request, + boolean useConfiguredServerName) { Socket socket = null; try { - socket = transport.dial(engine, resolvedOptions.tls(), resolvedOptions.connectTimeout()); + socket = + transport.dial( + engine, + resolvedOptions.tls(), + resolvedOptions.connectTimeout(), + useConfiguredServerName); setReadTimeout(socket, resolvedOptions.operationTimeout()); Protocol.writeMessage(socket.getOutputStream(), request); return socket; @@ -233,13 +260,17 @@ private ResolvedClientOptions resolved() { } private String resolveEngine(ResolvedClientOptions resolvedOptions) { - if (resolvedOptions.engine() != null) return resolvedOptions.engine(); + if (resolvedOptions.region() == null && resolvedOptions.engine() != null) + return resolvedOptions.engine(); if (resolvedOptions.projectEndpoint() == null) { throw new RstreamException( "Engine is required but not configured.", "ERR_RSTREAM_ENGINE_REQUIRED"); } - return new RstreamApiClient(resolvedOptions.apiUrl(), resolvedOptions.token()) - .resolveEngine(resolvedOptions.projectEndpoint()); + return new RstreamApiClient( + resolvedOptions.apiUrl(), + resolvedOptions.token(), + resolvedOptions.controlPlaneHeaders()) + .resolveEngine(resolvedOptions.projectEndpoint(), resolvedOptions.region()); } private String resolveToken(ResolvedClientOptions resolvedOptions) { diff --git a/src/main/java/io/rstream/RstreamTransport.java b/src/main/java/io/rstream/RstreamTransport.java index 4385ba7..d1e4735 100644 --- a/src/main/java/io/rstream/RstreamTransport.java +++ b/src/main/java/io/rstream/RstreamTransport.java @@ -8,10 +8,14 @@ final class RstreamTransport { SSLSocket dial(String engine, TlsOptions tls, Duration timeout) throws IOException { + return dial(engine, tls, timeout, true); + } + + SSLSocket dial(String engine, TlsOptions tls, Duration timeout, boolean useConfiguredServerName) + throws IOException { var address = EngineAddress.parse(engine); var context = TlsSupport.context(tls); - var serverName = tls == null || tls.serverName() == null ? "" : tls.serverName().trim(); - var peerHost = serverName.isEmpty() ? address.host() : serverName; + var peerHost = peerHost(address, tls, useConfiguredServerName); var rawSocket = new Socket(); rawSocket.connect( new InetSocketAddress(address.host(), address.port()), timeoutMillis(timeout)); @@ -25,6 +29,11 @@ SSLSocket dial(String engine, TlsOptions tls, Duration timeout) throws IOExcepti return socket; } + static String peerHost(EngineAddress address, TlsOptions tls, boolean useConfiguredServerName) { + var serverName = tls == null || tls.serverName() == null ? "" : tls.serverName().trim(); + return useConfiguredServerName && !serverName.isEmpty() ? serverName : address.host(); + } + private static int timeoutMillis(Duration timeout) { var millis = timeout.toMillis(); if (millis > Integer.MAX_VALUE) return Integer.MAX_VALUE; diff --git a/src/main/java/io/rstream/TunnelProperties.java b/src/main/java/io/rstream/TunnelProperties.java index 9b9c937..199bb7a 100644 --- a/src/main/java/io/rstream/TunnelProperties.java +++ b/src/main/java/io/rstream/TunnelProperties.java @@ -29,7 +29,8 @@ public record TunnelProperties( String hostname, Integer port, Boolean upstreamTls, - Boolean datagramGuaranteedDelivery) { + Boolean datagramGuaranteedDelivery, + Boolean allowCrossRegionRouting) { public TunnelProperties { labels = labels == null ? Map.of() : Map.copyOf(labels); geoIp = geoIp == null ? List.of() : List.copyOf(geoIp); @@ -67,6 +68,7 @@ public static final class Builder { private Integer port; private Boolean upstreamTls; private Boolean datagramGuaranteedDelivery; + private Boolean allowCrossRegionRouting; public Builder id(String id) { this.id = id; @@ -188,6 +190,11 @@ public Builder datagramGuaranteedDelivery(Boolean datagramGuaranteedDelivery) { return this; } + public Builder allowCrossRegionRouting(Boolean allowCrossRegionRouting) { + this.allowCrossRegionRouting = allowCrossRegionRouting; + return this; + } + public TunnelProperties build() { return new TunnelProperties( id, @@ -213,7 +220,8 @@ public TunnelProperties build() { hostname, port, upstreamTls, - datagramGuaranteedDelivery); + datagramGuaranteedDelivery, + allowCrossRegionRouting); } } } diff --git a/src/main/proto/rstream.proto b/src/main/proto/rstream.proto index c6a9c76..12d23db 100644 --- a/src/main/proto/rstream.proto +++ b/src/main/proto/rstream.proto @@ -14,7 +14,7 @@ extend google.protobuf.FieldOptions { string access = 51234; } -option (protocol_version) = "1.4.3"; +option (protocol_version) = "1.4.4"; package rstream.io_rstrm.protobuf; @@ -128,6 +128,7 @@ message TunnelProperties { google.protobuf.UInt32Value port = 23 [(access) = "read-write"]; google.protobuf.BoolValue upstream_tls = 24 [(access) = "read-write"]; google.protobuf.BoolValue datagram_guaranteed_delivery = 25 [(access) = "read-write"]; + google.protobuf.BoolValue allow_cross_region_routing = 26 [(access) = "read-write"]; } // When a client opens a new control channel to the server @@ -210,6 +211,7 @@ message ProxyConnReq { string stream_id = 2; google.protobuf.StringValue secret = 3; IpAddress source_ip = 4; + google.protobuf.StringValue proxy_endpoint = 5; } // Client's response to a 'ConnectionInitReq'. The client can refuse the connection by providing an error. diff --git a/src/test/java/io/rstream/ConfigResolverTest.java b/src/test/java/io/rstream/ConfigResolverTest.java index 676f05d..9eb1166 100644 --- a/src/test/java/io/rstream/ConfigResolverTest.java +++ b/src/test/java/io/rstream/ConfigResolverTest.java @@ -9,6 +9,7 @@ import java.time.Duration; import java.time.Instant; import java.util.Base64; +import java.util.List; import java.util.Map; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; @@ -113,6 +114,59 @@ void selectedContextLoadsEngineProjectTokenAndTls() throws Exception { assertThat(resolved.tls().insecureSkipVerify()).isTrue(); } + @Test + void regionResolutionFollowsOptionEnvironmentAndContextPrecedence() throws Exception { + var configPath = + config( + """ + defaults: + context: + name: global + contexts: + - name: global + engine: project.global.example.test:443 + projectEndpoint: project + region: eu-west-3 + """); + var environment = + ConfigResolver.resolve( + ClientOptions.builder().configPath(configPath.toString()).noToken(true).build(), + Map.of("RSTREAM_REGION", "us-east-1")); + assertThat(environment.region()).isEqualTo("us-east-1"); + assertThat(environment.engine()).isNull(); + var option = + ConfigResolver.resolve( + ClientOptions.builder() + .configPath(configPath.toString()) + .noToken(true) + .region("EU-CENTRAL-1") + .build(), + Map.of("RSTREAM_REGION", "us-east-1")); + assertThat(option.region()).isEqualTo("eu-central-1"); + var context = + ConfigResolver.resolve( + ClientOptions.builder().configPath(configPath.toString()).noToken(true).build(), + Map.of()); + assertThat(context.region()).isEqualTo("eu-west-3"); + } + + @Test + void regionResolutionRejectsExplicitEngineOverride() { + assertThatThrownBy( + () -> + ConfigResolver.resolve( + ClientOptions.builder() + .engine("engine.example.test:443") + .noToken(true) + .projectEndpoint("project") + .readConfigFile(false) + .region("eu-west-3") + .build(), + Map.of())) + .isInstanceOf(ConfigurationException.class) + .hasMessageContaining("explicit engine override"); + } + @Test void environmentContextSelectionLoadsEnvironmentCredentials() throws Exception { var configPath = @@ -314,6 +368,21 @@ void expiredJwtLikeTokenIsRejected() { .hasMessageContaining("expired"); } + @Test + void jwtLikeTokenWithoutNumericExpiryIsAccepted() { + var token = + "x." + + Base64.getUrlEncoder() + .withoutPadding() + .encodeToString("{\"exp\":null}".getBytes(StandardCharsets.UTF_8)) + + ".x"; + var resolved = + ConfigResolver.resolve( + ClientOptions.builder().engine("engine.example.com:443").token(token).build(), + Map.of()); + assertThat(resolved.token()).isEqualTo(token); + } + @Test void tokenAndMtlsCannotBeUsedTogether() { var options = @@ -491,6 +560,66 @@ void requireTokenRejectsUnauthenticatedRuntime() { .hasMessageContaining("Authentication is required"); } + @Test + void controlPlaneHeadersMergeConfigEnvironmentAndOptions() throws Exception { + var configPath = + config( + """ + defaults: + context: + name: local + environments: + - apiUrl: https://rstream.io + headers: + X-Environment: config + X-Shared: config + contexts: + - name: local + apiUrl: https://rstream.io + engine: engine.example.com:443 + """); + var resolved = + ConfigResolver.resolve( + ClientOptions.builder() + .configPath(configPath.toString()) + .controlPlaneHeaders(Map.of("X-Explicit", "option", "X-Shared", "option")) + .noToken(true) + .build(), + Map.of( + "RSTREAM_CONTROL_PLANE_HEADERS", + "{\"X-Runtime\":\"environment\",\"X-Shared\":\"environment\"}")); + assertThat(resolved.controlPlaneHeaders()) + .containsExactlyInAnyOrderEntriesOf( + Map.of( + "X-Environment", + "config", + "X-Explicit", + "option", + "X-Runtime", + "environment", + "X-Shared", + "option")); + } + + @Test + void invalidControlPlaneHeadersAreRejected() { + var options = ClientOptions.builder().readConfigFile(false).noToken(true).build(); + for (var value : + List.of( + "not-json", + "[]", + "{\"X-Test\":1}", + "{\"Authorization\":\"secret\"}", + "{\"X-Forwarded-Host\":\"example.test\"}", + "{\"Bad Header\":\"value\"}", + "{\"X-Test\":\"first\",\"x-test\":\"second\"}")) { + assertThatThrownBy( + () -> ConfigResolver.resolve(options, Map.of("RSTREAM_CONTROL_PLANE_HEADERS", value))) + .isInstanceOf(ConfigurationException.class) + .hasMessageMatching("(?i).*header.*"); + } + } + private Path config(String content) throws Exception { var path = temp.resolve("config.yaml"); Files.writeString(path, content); diff --git a/src/test/java/io/rstream/ProtocolTest.java b/src/test/java/io/rstream/ProtocolTest.java index c500291..3a87ae7 100644 --- a/src/test/java/io/rstream/ProtocolTest.java +++ b/src/test/java/io/rstream/ProtocolTest.java @@ -41,6 +41,7 @@ void tunnelPropertiesRoundTripAllSupportedFields() { "api.example.test", 443, true, + true, true); var decoded = Protocol.tunnelPropertiesFromPb(Protocol.tunnelPropertiesToPb(properties)); assertThat(decoded).isEqualTo(properties); diff --git a/src/test/java/io/rstream/RstreamApiClientTest.java b/src/test/java/io/rstream/RstreamApiClientTest.java index a4f6ab2..b22260f 100644 --- a/src/test/java/io/rstream/RstreamApiClientTest.java +++ b/src/test/java/io/rstream/RstreamApiClientTest.java @@ -9,6 +9,7 @@ import java.net.InetAddress; import java.net.InetSocketAddress; import java.nio.charset.StandardCharsets; +import java.util.Map; import org.junit.jupiter.api.Test; final class RstreamApiClientTest { @@ -28,6 +29,115 @@ void resolvesEngineFromManagedProjectFields() throws Exception { } } + @Test + void resolvesOnlyAuthorizedRegionalEngines() throws Exception { + try (var api = + ApiServer.start( + 200, + """ + { + "endpoint":"project", + "domain":"global.example.test", + "enginePort":443, + "placement":"global", + "regionalEndpoints":[ + { + "provider":"aws", + "region":"eu-west-3", + "domain":"eu.example.test", + "enginePort":8443 + }, + { + "provider":"aws", + "region":"us-east-1", + "domain":"us.example.test", + "enginePort":443 + } + ] + } + """)) { + var client = new RstreamApiClient(api.url(), null); + assertThat(client.resolveEngine("project", "US-EAST-1")) + .isEqualTo("project.us.example.test:443"); + assertThatThrownBy(() -> client.resolveEngine("project", "ap-southeast-1")) + .isInstanceOf(ConfigurationException.class) + .hasMessageContaining("Available regions: eu-west-3, us-east-1"); + } + } + + @Test + void rejectsAmbiguousRegionalEngines() throws Exception { + try (var api = + ApiServer.start( + 200, + """ + { + "endpoint":"project", + "regionalEndpoints":[ + { + "provider":"aws", + "region":"eu-west-3", + "domain":"eu.example.test", + "enginePort":443 + }, + { + "provider":"aws", + "region":"eu-west-3", + "domain":"eu-alt.example.test", + "enginePort":443 + } + ] + } + """)) { + var client = new RstreamApiClient(api.url(), null); + assertThatThrownBy(() -> client.resolveEngine("project", "eu-west-3")) + .isInstanceOf(ConfigurationException.class) + .hasMessageContaining("ambiguous"); + } + } + + @Test + void sendsValidatedControlPlaneHeaders() throws Exception { + try (var api = + ApiServer.start( + 200, + """ + {"endpoint":"project","domain":"t.localhost.rstream.io","enginePort":9443} + """)) { + var client = + new RstreamApiClient( + api.url(), null, Map.of("x-vercel-protection-bypass", "test-secret")); + client.resolveEngine("project"); + assertThat(api.controlPlaneHeader()).isEqualTo("test-secret"); + } + } + + @Test + void rejectsReservedControlPlaneHeaders() { + assertThatThrownBy( + () -> new RstreamApiClient("https://rstream.io", null, Map.of("Host", "evil.test"))) + .isInstanceOf(ConfigurationException.class) + .hasMessageContaining("reserved"); + } + + @Test + void doesNotForwardControlPlaneHeadersAcrossRedirects() throws Exception { + try (var target = + ApiServer.start( + 200, + """ + {"endpoint":"project","domain":"t.localhost.rstream.io","enginePort":9443} + """); + var redirect = ApiServer.redirectTo(target.url())) { + var client = + new RstreamApiClient(redirect.url(), null, Map.of("X-Deployment-Bypass", "test-secret")); + assertThatThrownBy(() -> client.resolveEngine("project")) + .isInstanceOf(RstreamException.class) + .hasMessageContaining("HTTP error 302"); + assertThat(target.requests()).isZero(); + } + } + @Test void resolvesEngineFromFallbackUrl() throws Exception { try (var api = ApiServer.start(200, "{\"url\":\"fallback.example.com:9443\"}")) { @@ -70,21 +180,35 @@ void surfacesInvalidJsonResponses() throws Exception { private static final class ApiServer implements Closeable { private final HttpServer server; private volatile String authorization; + private volatile String controlPlaneHeader; private volatile String path; + private volatile int requests; private ApiServer(HttpServer server) { this.server = server; } static ApiServer start(int status, String body) throws IOException { + return start(status, body, null); + } + + static ApiServer redirectTo(String location) throws IOException { + return start(302, "", location); + } + + private static ApiServer start(int status, String body, String location) throws IOException { var server = HttpServer.create(new InetSocketAddress(InetAddress.getLoopbackAddress(), 0), 0); var api = new ApiServer(server); server.createContext( "/", exchange -> { + api.requests++; api.path = exchange.getRequestURI().getRawPath(); api.authorization = exchange.getRequestHeaders().getFirst("authorization"); + api.controlPlaneHeader = + exchange.getRequestHeaders().getFirst("x-vercel-protection-bypass"); var bytes = body.getBytes(StandardCharsets.UTF_8); + if (location != null) exchange.getResponseHeaders().set("Location", location); exchange.sendResponseHeaders(status, bytes.length); try (var output = exchange.getResponseBody()) { output.write(bytes); @@ -102,10 +226,18 @@ String authorization() { return authorization; } + String controlPlaneHeader() { + return controlPlaneHeader; + } + String path() { return path; } + int requests() { + return requests; + } + @Override public void close() { server.stop(0); diff --git a/src/test/java/io/rstream/RstreamTransportTest.java b/src/test/java/io/rstream/RstreamTransportTest.java index 501d631..4c74cc4 100644 --- a/src/test/java/io/rstream/RstreamTransportTest.java +++ b/src/test/java/io/rstream/RstreamTransportTest.java @@ -1,5 +1,6 @@ package io.rstream; +import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import java.net.InetAddress; @@ -11,6 +12,15 @@ import org.junit.jupiter.api.Test; final class RstreamTransportTest { + @Test + void redirectedEngineUsesItsHostnameForTls() { + var tls = TlsOptions.builder().serverName("owner.example.test").build(); + var owner = EngineAddress.parse("owner.example.test:443"); + var ingress = EngineAddress.parse("ingress.example.test:443"); + assertThat(RstreamTransport.peerHost(owner, tls, true)).isEqualTo("owner.example.test"); + assertThat(RstreamTransport.peerHost(ingress, tls, false)).isEqualTo("ingress.example.test"); + } + @Test void tlsHandshakeUsesConfiguredTimeout() throws Exception { var accepted = new CountDownLatch(1); diff --git a/src/test/java/io/rstream/RuntimeFakeEngineIT.java b/src/test/java/io/rstream/RuntimeFakeEngineIT.java index a1661c3..eaab5b7 100644 --- a/src/test/java/io/rstream/RuntimeFakeEngineIT.java +++ b/src/test/java/io/rstream/RuntimeFakeEngineIT.java @@ -110,6 +110,63 @@ void proxyConnectionIsDeliveredToTunnel() throws Exception { } } + @Test + void proxyConnectionCanDialIngressEngine() throws Exception { + try (var owner = FakeEngine.start(temp); + var ingress = FakeEngine.start(temp); + var client = client(owner, false, "owner-pat"); + var control = client.connect()) { + var tunnel = + control.createTunnel( + CreateTunnelOptions.builder().name("web").protocol(TunnelProtocol.HTTP).build()); + owner.sendProxyConnection(tunnel.id(), "stream_direct", "stream-secret", ingress.address()); + try (var stream = tunnel.accept(Duration.ofSeconds(2))) { + stream.outputStream().write("direct".getBytes(StandardCharsets.UTF_8)); + stream.outputStream().flush(); + assertThat(stream.inputStream().readNBytes(6)) + .isEqualTo("direct".getBytes(StandardCharsets.UTF_8)); + } + var response = owner.proxyConnectionResponses.poll(2, TimeUnit.SECONDS); + assertThat(response.hasError()).isFalse(); + assertThat(owner.proxyRequests.poll(100, TimeUnit.MILLISECONDS)).isNull(); + var proxyRequest = ingress.proxyRequests.poll(2, TimeUnit.SECONDS); + assertThat(proxyRequest.getStreamId()).isEqualTo("stream_direct"); + assertThat(proxyRequest.getClientDetails().getToken().getValue()).isEqualTo("stream-secret"); + } + } + + @Test + void proxyRedirectWithoutStreamSecretIsRejected() throws Exception { + try (var owner = FakeEngine.start(temp); + var ingress = FakeEngine.start(temp); + var client = client(owner); + var control = client.connect()) { + var tunnel = control.createTunnel(CreateTunnelOptions.builder().name("web").build()); + owner.sendProxyConnection(tunnel.id(), "stream_missing_secret", null, ingress.address()); + var response = owner.proxyConnectionResponses.poll(2, TimeUnit.SECONDS); + assertThat(response.hasError()).isTrue(); + assertThat(response.getError().getMessage().getValue()).contains("credentials"); + assertThat(owner.proxyRequests.poll(100, TimeUnit.MILLISECONDS)).isNull(); + assertThat(ingress.proxyRequests.poll(100, TimeUnit.MILLISECONDS)).isNull(); + } + } + + @Test + void proxyRedirectWithEmptyStreamSecretIsRejected() throws Exception { + try (var owner = FakeEngine.start(temp); + var ingress = FakeEngine.start(temp); + var client = client(owner); + var control = client.connect()) { + var tunnel = control.createTunnel(CreateTunnelOptions.builder().name("web").build()); + owner.sendProxyConnection(tunnel.id(), "stream_empty_secret", "", ingress.address()); + var response = owner.proxyConnectionResponses.poll(2, TimeUnit.SECONDS); + assertThat(response.hasError()).isTrue(); + assertThat(response.getError().getMessage().getValue()).contains("credentials"); + assertThat(owner.proxyRequests.poll(100, TimeUnit.MILLISECONDS)).isNull(); + assertThat(ingress.proxyRequests.poll(100, TimeUnit.MILLISECONDS)).isNull(); + } + } + @Test void concurrentTunnelCreationUsesOneControlChannelSafely() throws Exception { var count = 24; @@ -511,6 +568,7 @@ void publishedTCPOptionsAreSent() throws Exception { .name("ssh") .protocol(TunnelProtocol.TCP) .port(10042) + .allowCrossRegionRouting(true) .build()); var request = engine.openTunnelRequests.poll(2, TimeUnit.SECONDS); assertThat(request).isNotNull(); @@ -518,6 +576,7 @@ void publishedTCPOptionsAreSent() throws Exception { assertThat(request.getTunnelProperties().getPublish().getValue()).isTrue(); assertThat(request.getTunnelProperties().getProtocol().getValue()).isEqualTo("tcp"); assertThat(request.getTunnelProperties().getPort().getValue()).isEqualTo(10042); + assertThat(request.getTunnelProperties().getAllowCrossRegionRouting().getValue()).isTrue(); } } @@ -535,6 +594,15 @@ void publishedTCPRejectsIncompatibleOptionsBeforeRequest() throws Exception { .build())) .isInstanceOf(RstreamException.class) .hasMessageContaining("do not accept"); + assertThatThrownBy( + () -> + control.createTunnel( + CreateTunnelOptions.builder() + .protocol(TunnelProtocol.HTTP) + .allowCrossRegionRouting(true) + .build())) + .isInstanceOf(RstreamException.class) + .hasMessageContaining("requires the TCP protocol"); assertThat(engine.openTunnelRequests).isEmpty(); } } @@ -580,15 +648,20 @@ private static RstreamClient timeoutClient(FakeEngine engine, boolean zeroRtt) { } private static RstreamClient client(FakeEngine engine, boolean zeroRtt) { - return RstreamClient.fromEnv( + return client(engine, zeroRtt, null); + } + + private static RstreamClient client(FakeEngine engine, boolean zeroRtt, String token) { + var builder = ClientOptions.builder() .engine(engine.address()) .readConfigFile(false) - .noToken(true) .heartbeat(false) .zeroRtt(zeroRtt) - .tls(TlsOptions.builder().insecureSkipVerify(true).build()) - .build()); + .tls(TlsOptions.builder().insecureSkipVerify(true).build()); + if (token == null) builder.noToken(true); + else builder.token(token); + return RstreamClient.fromEnv(builder.build()); } private static String dialEcho( @@ -666,13 +739,16 @@ Path certificatePath() { } void sendProxyConnection(String tunnelId, String streamId, String secret) { - var request = - Rstream.ProxyConnReq.newBuilder() - .setTunnelId(tunnelId) - .setStreamId(streamId) - .setSecret(StringValue.newBuilder().setValue(secret)) - .build(); - writeControl(Rstream.Message.newBuilder().setProxyConnReq(request).build()); + sendProxyConnection(tunnelId, streamId, secret, null); + } + + void sendProxyConnection( + String tunnelId, String streamId, String secret, String proxyEndpoint) { + var request = Rstream.ProxyConnReq.newBuilder().setTunnelId(tunnelId).setStreamId(streamId); + if (secret != null) request.setSecret(StringValue.newBuilder().setValue(secret)); + if (proxyEndpoint != null) + request.setProxyEndpoint(StringValue.newBuilder().setValue(proxyEndpoint)); + writeControl(Rstream.Message.newBuilder().setProxyConnReq(request.build()).build()); } void closeControlSocket() throws IOException { @@ -879,52 +955,54 @@ private static void echo(SSLSocket socket) throws IOException { private static SSLContext serverContext(Path temp) throws Exception { var keyStorePath = temp.resolve("server.p12"); var keytool = Path.of(System.getProperty("java.home"), "bin", "keytool").toString(); - var process = - new ProcessBuilder( - keytool, - "-genkeypair", - "-alias", - "server", - "-keyalg", - "RSA", - "-storetype", - "PKCS12", - "-keystore", - keyStorePath.toString(), - "-storepass", - "changeit", - "-keypass", - "changeit", - "-dname", - "CN=localhost", - "-validity", - "2", - "-ext", - "SAN=dns:localhost") - .redirectErrorStream(true) - .start(); - if (process.waitFor() != 0) { - throw new IllegalStateException( - new String(process.getInputStream().readAllBytes(), StandardCharsets.UTF_8)); - } - var export = - new ProcessBuilder( - keytool, - "-exportcert", - "-rfc", - "-alias", - "server", - "-keystore", - keyStorePath.toString(), - "-storepass", - "changeit", - "-file", - temp.resolve("server.crt").toString()) - .redirectErrorStream(true) - .start(); - if (export.waitFor() != 0) { - throw new IllegalStateException( - new String(export.getInputStream().readAllBytes(), StandardCharsets.UTF_8)); + if (!java.nio.file.Files.exists(keyStorePath)) { + var process = + new ProcessBuilder( + keytool, + "-genkeypair", + "-alias", + "server", + "-keyalg", + "RSA", + "-storetype", + "PKCS12", + "-keystore", + keyStorePath.toString(), + "-storepass", + "changeit", + "-keypass", + "changeit", + "-dname", + "CN=localhost", + "-validity", + "2", + "-ext", + "SAN=dns:localhost") + .redirectErrorStream(true) + .start(); + if (process.waitFor() != 0) { + throw new IllegalStateException( + new String(process.getInputStream().readAllBytes(), StandardCharsets.UTF_8)); + } + var export = + new ProcessBuilder( + keytool, + "-exportcert", + "-rfc", + "-alias", + "server", + "-keystore", + keyStorePath.toString(), + "-storepass", + "changeit", + "-file", + temp.resolve("server.crt").toString()) + .redirectErrorStream(true) + .start(); + if (export.waitFor() != 0) { + throw new IllegalStateException( + new String(export.getInputStream().readAllBytes(), StandardCharsets.UTF_8)); + } } var store = KeyStore.getInstance("PKCS12"); try (var input = java.nio.file.Files.newInputStream(keyStorePath)) { From f5ea65e2cef39e32829bb72a59e2d486bb4100b7 Mon Sep 17 00:00:00 2001 From: uartnet <140632163+uartnet@users.noreply.github.com> Date: Fri, 24 Jul 2026 19:57:59 +0200 Subject: [PATCH 2/6] fix(edge): allow cross-region routing for every protocol --- src/main/java/io/rstream/ControlChannel.java | 5 ---- .../java/io/rstream/RuntimeFakeEngineIT.java | 26 ++++++++++++------- 2 files changed, 17 insertions(+), 14 deletions(-) diff --git a/src/main/java/io/rstream/ControlChannel.java b/src/main/java/io/rstream/ControlChannel.java index 5134c87..70e54ca 100644 --- a/src/main/java/io/rstream/ControlChannel.java +++ b/src/main/java/io/rstream/ControlChannel.java @@ -312,11 +312,6 @@ private ScheduledFuture operationTimeout( } private static TunnelProperties normalizeBytestreamOptions(CreateTunnelOptions options) { - if (options.allowCrossRegionRouting() != null && options.protocol() != TunnelProtocol.TCP) { - throw new RstreamException( - "Cross-region routing policy requires the TCP protocol.", - "ERR_RSTREAM_INVALID_TUNNEL_OPTIONS"); - } if (options.httpVersion() == HttpVersion.H3) { throw new UnsupportedFeatureException( "HTTP/3 tunnels require datagram support, which rstream-java does not support.", diff --git a/src/test/java/io/rstream/RuntimeFakeEngineIT.java b/src/test/java/io/rstream/RuntimeFakeEngineIT.java index eaab5b7..54bb9cd 100644 --- a/src/test/java/io/rstream/RuntimeFakeEngineIT.java +++ b/src/test/java/io/rstream/RuntimeFakeEngineIT.java @@ -594,19 +594,27 @@ void publishedTCPRejectsIncompatibleOptionsBeforeRequest() throws Exception { .build())) .isInstanceOf(RstreamException.class) .hasMessageContaining("do not accept"); - assertThatThrownBy( - () -> - control.createTunnel( - CreateTunnelOptions.builder() - .protocol(TunnelProtocol.HTTP) - .allowCrossRegionRouting(true) - .build())) - .isInstanceOf(RstreamException.class) - .hasMessageContaining("requires the TCP protocol"); assertThat(engine.openTunnelRequests).isEmpty(); } } + @Test + void crossRegionRoutingPolicyIsSentForHTTP() throws Exception { + try (var engine = FakeEngine.start(temp); + var client = client(engine); + var control = client.connect()) { + control.createTunnel( + CreateTunnelOptions.builder() + .protocol(TunnelProtocol.HTTP) + .allowCrossRegionRouting(true) + .build()); + var request = engine.openTunnelRequests.poll(2, TimeUnit.SECONDS); + assertThat(request).isNotNull(); + assertThat(request.getTunnelProperties().getProtocol().getValue()).isEqualTo("http"); + assertThat(request.getTunnelProperties().getAllowCrossRegionRouting().getValue()).isTrue(); + } + } + @Test void privateTunnelPublicExposureOptionsAreRejectedBeforeRequest() throws Exception { try (var engine = FakeEngine.start(temp); From 1087b17d0694525f873274eb2773c2a1f01e3b51 Mon Sep 17 00:00:00 2001 From: uartnet <140632163+uartnet@users.noreply.github.com> Date: Fri, 24 Jul 2026 19:58:07 +0200 Subject: [PATCH 3/6] build(deps): update Java dependencies --- examples/forward-local-port/pom.xml | 2 +- examples/micronaut-published-http/pom.xml | 10 +++++----- examples/micronaut-webhook-receiver/pom.xml | 12 +++++------ examples/private-dial/pom.xml | 2 +- examples/published-http-server/pom.xml | 2 +- examples/quarkus-published-http/pom.xml | 2 +- examples/spring-boot-published-http/pom.xml | 2 +- examples/spring-boot-webhook-receiver/pom.xml | 2 +- examples/vertx-published-http/pom.xml | 6 +++--- examples/webhook-receiver/pom.xml | 2 +- pom.xml | 20 +++++++++---------- 11 files changed, 31 insertions(+), 31 deletions(-) diff --git a/examples/forward-local-port/pom.xml b/examples/forward-local-port/pom.xml index 502da1b..b4097a1 100644 --- a/examples/forward-local-port/pom.xml +++ b/examples/forward-local-port/pom.xml @@ -20,7 +20,7 @@ org.codehaus.mojo exec-maven-plugin - 3.6.2 + 3.6.3 io.rstream.examples.ForwardLocalPort diff --git a/examples/micronaut-published-http/pom.xml b/examples/micronaut-published-http/pom.xml index d8b5592..a3b393d 100644 --- a/examples/micronaut-published-http/pom.xml +++ b/examples/micronaut-published-http/pom.xml @@ -8,8 +8,8 @@ 17 ${java.version} - 4.10.18 - 4.10.10 + 5.1.10 + 5.0.6 UTF-8 @@ -36,14 +36,14 @@ com.fasterxml.jackson.core jackson-databind - 2.20.1 + 2.22.1 maven-compiler-plugin - 3.14.1 + 3.15.0 @@ -71,7 +71,7 @@ org.codehaus.mojo exec-maven-plugin - 3.6.2 + 3.6.3 io.rstream.examples.micronaut.PublishedHttpApplication diff --git a/examples/micronaut-webhook-receiver/pom.xml b/examples/micronaut-webhook-receiver/pom.xml index d477bba..69f038e 100644 --- a/examples/micronaut-webhook-receiver/pom.xml +++ b/examples/micronaut-webhook-receiver/pom.xml @@ -8,10 +8,10 @@ 17 ${java.version} - 4.10.18 - 4.10.10 - 4.10.2 - 2.16.2 + 5.1.10 + 5.0.6 + 5.0.2 + 3.1.0 UTF-8 @@ -49,7 +49,7 @@ maven-compiler-plugin - 3.14.1 + 3.15.0 @@ -82,7 +82,7 @@ org.codehaus.mojo exec-maven-plugin - 3.6.2 + 3.6.3 io.rstream.examples.micronaut.WebhookReceiverApplication diff --git a/examples/private-dial/pom.xml b/examples/private-dial/pom.xml index 2732f93..94d23b9 100644 --- a/examples/private-dial/pom.xml +++ b/examples/private-dial/pom.xml @@ -20,7 +20,7 @@ org.codehaus.mojo exec-maven-plugin - 3.6.2 + 3.6.3 io.rstream.examples.PrivateDial diff --git a/examples/published-http-server/pom.xml b/examples/published-http-server/pom.xml index a01292c..8e9de5c 100644 --- a/examples/published-http-server/pom.xml +++ b/examples/published-http-server/pom.xml @@ -20,7 +20,7 @@ org.codehaus.mojo exec-maven-plugin - 3.6.2 + 3.6.3 io.rstream.examples.PublishedHttpServer diff --git a/examples/quarkus-published-http/pom.xml b/examples/quarkus-published-http/pom.xml index b1e5a73..6be8a48 100644 --- a/examples/quarkus-published-http/pom.xml +++ b/examples/quarkus-published-http/pom.xml @@ -6,7 +6,7 @@ 0.3.0 rstream Quarkus published HTTP example - 3.14.1 + 3.15.0 17 ${java.version} UTF-8 diff --git a/examples/spring-boot-published-http/pom.xml b/examples/spring-boot-published-http/pom.xml index 7bed59a..4a9b125 100644 --- a/examples/spring-boot-published-http/pom.xml +++ b/examples/spring-boot-published-http/pom.xml @@ -4,7 +4,7 @@ org.springframework.boot spring-boot-starter-parent - 3.5.14 + 4.1.0 io.rstream.examples diff --git a/examples/spring-boot-webhook-receiver/pom.xml b/examples/spring-boot-webhook-receiver/pom.xml index 3c2501e..b775bc3 100644 --- a/examples/spring-boot-webhook-receiver/pom.xml +++ b/examples/spring-boot-webhook-receiver/pom.xml @@ -4,7 +4,7 @@ org.springframework.boot spring-boot-starter-parent - 3.5.14 + 4.1.0 io.rstream.examples diff --git a/examples/vertx-published-http/pom.xml b/examples/vertx-published-http/pom.xml index ab4a901..0a8e6be 100644 --- a/examples/vertx-published-http/pom.xml +++ b/examples/vertx-published-http/pom.xml @@ -9,7 +9,7 @@ 17 ${java.version} UTF-8 - 5.1.1 + 5.1.5 @@ -33,7 +33,7 @@ maven-compiler-plugin - 3.14.1 + 3.15.0 -Xlint:all @@ -45,7 +45,7 @@ org.codehaus.mojo exec-maven-plugin - 3.6.2 + 3.6.3 io.rstream.examples.vertx.PublishedHttpVerticle diff --git a/examples/webhook-receiver/pom.xml b/examples/webhook-receiver/pom.xml index a61695c..b16c058 100644 --- a/examples/webhook-receiver/pom.xml +++ b/examples/webhook-receiver/pom.xml @@ -20,7 +20,7 @@ org.codehaus.mojo exec-maven-plugin - 3.6.2 + 3.6.3 io.rstream.examples.WebhookReceiver diff --git a/pom.xml b/pom.xml index 54e48fc..51ea6fb 100644 --- a/pom.xml +++ b/pom.xml @@ -33,10 +33,10 @@ 17 UTF-8 ${java.version} - 4.33.1 - 2.20.1 - 5.14.1 - 3.27.6 + 4.35.1 + 2.22.1 + 6.1.2 + 3.27.7 @@ -47,7 +47,7 @@ org.snakeyaml snakeyaml-engine - 2.10 + 3.0.1 com.fasterxml.jackson.core @@ -78,7 +78,7 @@ org.apache.maven.plugins maven-enforcer-plugin - 3.6.2 + 3.6.3 enforce @@ -117,7 +117,7 @@ com.diffplug.spotless spotless-maven-plugin - 3.1.0 + 3.8.0 @@ -133,7 +133,7 @@ org.apache.maven.plugins maven-compiler-plugin - 3.14.1 + 3.15.0 true @@ -164,7 +164,7 @@ org.jacoco jacoco-maven-plugin - 0.8.14 + 0.8.15 rstream/io_rstrm/protobuf/** @@ -208,7 +208,7 @@ org.apache.maven.plugins maven-source-plugin - 3.3.1 + 3.4.0 attach-sources From a3365280fb075143e0cb92bdaa403f3f5506db1a Mon Sep 17 00:00:00 2001 From: uartnet <140632163+uartnet@users.noreply.github.com> Date: Wed, 29 Jul 2026 14:10:06 +0200 Subject: [PATCH 4/6] docs: organize internal documentation --- CONTRIBUTING.md | 2 +- README.md | 8 ++++---- TODO | 1 - docs/{CONFIGURATION.md => 001-configuration.md} | 0 docs/{TUNNELS.md => 002-tunnels.md} | 0 docs/{WEBHOOKS.md => 003-webhooks.md} | 0 docs/{TESTING.md => 004-testing.md} | 2 +- docs/{TEST_MATRIX.md => 005-test-matrix.md} | 0 docs/{GITHUB_SETUP.md => 006-github-setup.md} | 0 docs/README.md | 10 ++++++++++ 10 files changed, 16 insertions(+), 7 deletions(-) delete mode 100644 TODO rename docs/{CONFIGURATION.md => 001-configuration.md} (100%) rename docs/{TUNNELS.md => 002-tunnels.md} (100%) rename docs/{WEBHOOKS.md => 003-webhooks.md} (100%) rename docs/{TESTING.md => 004-testing.md} (96%) rename docs/{TEST_MATRIX.md => 005-test-matrix.md} (100%) rename docs/{GITHUB_SETUP.md => 006-github-setup.md} (100%) create mode 100644 docs/README.md diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 68716a2..b6d7d85 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -50,7 +50,7 @@ Use the smallest test that exercises the behavior: - fake-engine integration tests for runtime control-channel and stream behavior; - opt-in e2e tests for real engines and managed environments. -See [docs/TESTING.md](docs/TESTING.md) for the e2e matrix. +See [docs/004-testing.md](docs/004-testing.md) for the e2e matrix. ## Generated code diff --git a/README.md b/README.md index 53e7df5..b0bb1c3 100644 --- a/README.md +++ b/README.md @@ -99,7 +99,7 @@ SDKs: `RSTREAM_ENGINE_ADDRESS` is also accepted for compatibility with older local SDK workflows. Prefer `RSTREAM_ENGINE` in new code. -See [docs/CONFIGURATION.md](docs/CONFIGURATION.md) for supported YAML fields and +See [docs/001-configuration.md](docs/001-configuration.md) for supported YAML fields and error behavior. ## Published HTTP tunnel @@ -179,7 +179,7 @@ if (event.type().wireValue().equals("tunnel.created")) { `event.id()` is suitable for idempotency. Keep the raw request body unchanged when verifying the signature. -See [docs/WEBHOOKS.md](docs/WEBHOOKS.md) for the payload shape and headers. +See [docs/003-webhooks.md](docs/003-webhooks.md) for the payload shape and headers. ## Examples @@ -209,7 +209,7 @@ Real-engine tests are opt-in: RSTREAM_JAVA_E2E=1 mvn verify ``` -See [docs/TESTING.md](docs/TESTING.md) for local-engine and managed-environment +See [docs/004-testing.md](docs/004-testing.md) for local-engine and managed-environment test commands. ## Repository setup and release @@ -219,7 +219,7 @@ secret for normal pull request checks. Release automation uses release-please an requires the maintainer-managed `RELEASE_PLEASE_TOKEN` secret plus the `CI_ALLOWED_ACTOR` repository variable. -See [docs/GITHUB_SETUP.md](docs/GITHUB_SETUP.md) before creating or publishing +See [docs/006-github-setup.md](docs/006-github-setup.md) before creating or publishing the repository. ## License diff --git a/TODO b/TODO deleted file mode 100644 index 4640904..0000000 --- a/TODO +++ /dev/null @@ -1 +0,0 @@ -# TODO diff --git a/docs/CONFIGURATION.md b/docs/001-configuration.md similarity index 100% rename from docs/CONFIGURATION.md rename to docs/001-configuration.md diff --git a/docs/TUNNELS.md b/docs/002-tunnels.md similarity index 100% rename from docs/TUNNELS.md rename to docs/002-tunnels.md diff --git a/docs/WEBHOOKS.md b/docs/003-webhooks.md similarity index 100% rename from docs/WEBHOOKS.md rename to docs/003-webhooks.md diff --git a/docs/TESTING.md b/docs/004-testing.md similarity index 96% rename from docs/TESTING.md rename to docs/004-testing.md index 3c97833..85f1697 100644 --- a/docs/TESTING.md +++ b/docs/004-testing.md @@ -18,7 +18,7 @@ The default suite includes: JaCoCo coverage is enforced on maintained SDK code. Generated Protobuf classes are excluded from the coverage threshold. -See [TEST_MATRIX.md](TEST_MATRIX.md) for the pre-release coverage matrix. +See [005-test-matrix.md](005-test-matrix.md) for the pre-release coverage matrix. ## Real-engine e2e diff --git a/docs/TEST_MATRIX.md b/docs/005-test-matrix.md similarity index 100% rename from docs/TEST_MATRIX.md rename to docs/005-test-matrix.md diff --git a/docs/GITHUB_SETUP.md b/docs/006-github-setup.md similarity index 100% rename from docs/GITHUB_SETUP.md rename to docs/006-github-setup.md diff --git a/docs/README.md b/docs/README.md new file mode 100644 index 0000000..745c9ad --- /dev/null +++ b/docs/README.md @@ -0,0 +1,10 @@ +# Java SDK documentation + +Read these documents in order: + +1. [Configuration](001-configuration.md) documents SDK configuration. +2. [Tunnels](002-tunnels.md) documents tunnel APIs and behavior. +3. [Webhooks](003-webhooks.md) documents webhook support. +4. [Testing](004-testing.md) describes local and CI validation. +5. [Test matrix](005-test-matrix.md) defines supported runtime combinations. +6. [GitHub setup](006-github-setup.md) documents repository CI configuration. From 1f4399c242fa8af135b43f116322f18e5ee2fc80 Mon Sep 17 00:00:00 2001 From: uartnet <140632163+uartnet@users.noreply.github.com> Date: Thu, 30 Jul 2026 08:21:18 +0200 Subject: [PATCH 5/6] test(api): align project routing metadata --- src/test/java/io/rstream/RstreamApiClientTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/test/java/io/rstream/RstreamApiClientTest.java b/src/test/java/io/rstream/RstreamApiClientTest.java index b22260f..fe430f9 100644 --- a/src/test/java/io/rstream/RstreamApiClientTest.java +++ b/src/test/java/io/rstream/RstreamApiClientTest.java @@ -39,7 +39,7 @@ void resolvesOnlyAuthorizedRegionalEngines() throws Exception { "endpoint":"project", "domain":"global.example.test", "enginePort":443, - "placement":"global", + "routing":"global", "regionalEndpoints":[ { "provider":"aws", From 2a0a1694f22d6a321bb4a8ce6fa1d0340f62c304 Mon Sep 17 00:00:00 2001 From: uartnet <140632163+uartnet@users.noreply.github.com> Date: Sat, 1 Aug 2026 21:49:32 +0200 Subject: [PATCH 6/6] fix(examples): preserve the Java 17 baseline --- examples/micronaut-published-http/pom.xml | 4 ++-- examples/micronaut-webhook-receiver/pom.xml | 8 ++++---- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/examples/micronaut-published-http/pom.xml b/examples/micronaut-published-http/pom.xml index a3b393d..752c47b 100644 --- a/examples/micronaut-published-http/pom.xml +++ b/examples/micronaut-published-http/pom.xml @@ -8,8 +8,8 @@ 17 ${java.version} - 5.1.10 - 5.0.6 + 4.10.18 + 4.10.10 UTF-8 diff --git a/examples/micronaut-webhook-receiver/pom.xml b/examples/micronaut-webhook-receiver/pom.xml index 69f038e..64fd358 100644 --- a/examples/micronaut-webhook-receiver/pom.xml +++ b/examples/micronaut-webhook-receiver/pom.xml @@ -8,10 +8,10 @@ 17 ${java.version} - 5.1.10 - 5.0.6 - 5.0.2 - 3.1.0 + 4.10.18 + 4.10.10 + 4.10.2 + 2.16.2 UTF-8