From 1933260831d888d8e352874a0a13fa7e1af0909a Mon Sep 17 00:00:00 2001 From: Mohsen Zainalpour Date: Tue, 29 Sep 2026 18:31:57 -0400 Subject: [PATCH] fix(core): report an unknown events() port as ImposterNotFound The engine answers 404 both when it has no /events route and when the requested port has no imposter, and the body cannot tell them apart. A 404 with a port set now checks GET /imposters/{port} first, so a wrong port surfaces as the same ImposterNotFound every other per-port call throws instead of claiming the engine cannot stream. Bare {"error": ...} bodies from engines through 0.18.1 now yield their message too. Fixes #254 --- docs/design/sdk-api.md | 3 +- docs/events.md | 5 ++ .../rift/conformance/EventStreamIT.java | 15 +++++ .../java/io/github/achirdlabs/rift/Rift.java | 3 + .../rift/transport/RemoteTransport.java | 26 +++++++- .../rift/transport/EventStreamTest.java | 60 +++++++++++++++++++ 6 files changed, 109 insertions(+), 3 deletions(-) diff --git a/docs/design/sdk-api.md b/docs/design/sdk-api.md index 179d27d..4cb6209 100644 --- a/docs/design/sdk-api.md +++ b/docs/design/sdk-api.md @@ -857,7 +857,8 @@ Reconnect policy is the caller's: a retry schedule belongs to whatever drives th Capability probe: `events()` throws `UnsupportedOperationException` when streaming is unavailable — an engine too old to serve `/events` (404). The caller's move is to poll, which is the supported -baseline, not a degraded mode. +baseline, not a degraded mode. A 404 for a `port(...)` that names no imposter is not that: the +SDK confirms with `GET /imposters/{port}` and throws `ImposterNotFound`, as every per-port call does. Transport coverage is total, embedded included: the FFI transport delegates `events()` to the same lazily-started in-process admin server it already delegates `replaceAllImposters` to, and that diff --git a/docs/events.md b/docs/events.md index 46788f0..a0c78ec 100644 --- a/docs/events.md +++ b/docs/events.md @@ -157,3 +157,8 @@ What still refuses is an engine too old to serve `/events`: that throws `Unsuppo rather than returning an empty stream that would look like "nothing is happening". Same principle the filtered cursor reads follow — a connection that cannot answer the question says so, instead of returning a plausible-looking answer that is wrong. + +A `port(...)` that names no imposter is a different mistake and gets a different answer: +`ImposterNotFound`, the same error every other per-port call throws. The engine answers both cases +with a 404, so when a port was asked for the SDK checks `GET /imposters/{port}` before deciding which +one it is. diff --git a/rift-java-conformance/src/test/java/io/github/achirdlabs/rift/conformance/EventStreamIT.java b/rift-java-conformance/src/test/java/io/github/achirdlabs/rift/conformance/EventStreamIT.java index 646e787..b1fa704 100644 --- a/rift-java-conformance/src/test/java/io/github/achirdlabs/rift/conformance/EventStreamIT.java +++ b/rift-java-conformance/src/test/java/io/github/achirdlabs/rift/conformance/EventStreamIT.java @@ -8,6 +8,7 @@ import io.github.achirdlabs.rift.Rift; import io.github.achirdlabs.rift.RiftEvent; import io.github.achirdlabs.rift.error.EngineUnavailable; +import io.github.achirdlabs.rift.error.ImposterNotFound; import org.junit.jupiter.api.DynamicTest; import org.junit.jupiter.api.TestFactory; @@ -84,6 +85,20 @@ Stream theStreamPushesWhatTheJournalRecords() { }); } + @TestFactory + Stream aStreamForAPortWithNoImposterIsImposterNotFound() { + // The engine refuses with a 404 — the same status an engine without /events gives — so this + // pins that the SDK reports the wrong port, not a missing capability (#254). + return gated("an unknown port is ImposterNotFound, not 'cannot stream'", () -> { + try (Rift rift = engine()) { + int nobody = 1; // privileged, so never an imposter this suite created + ImposterNotFound e = assertThrows(ImposterNotFound.class, + () -> rift.events(EventStreamOptions.builder().port(nobody).build()).close()); + assertEquals(nobody, e.port()); + } + }); + } + @TestFactory Stream aBrokenStreamIsRecoveredByPollingFromTheLastIndexSeen() { return gated("the reconcile loop, end to end", () -> { diff --git a/rift-java-core/src/main/java/io/github/achirdlabs/rift/Rift.java b/rift-java-core/src/main/java/io/github/achirdlabs/rift/Rift.java index 10f7316..d404d07 100644 --- a/rift-java-core/src/main/java/io/github/achirdlabs/rift/Rift.java +++ b/rift-java-core/src/main/java/io/github/achirdlabs/rift/Rift.java @@ -118,6 +118,9 @@ static boolean isEmbeddedAvailable() { * old to serve {@code /events}. That means poll instead, * which is a supported baseline rather than a degraded * mode. + * @throws io.github.achirdlabs.rift.error.ImposterNotFound if {@link EventStreamOptions#port()} + * names no imposter — the same error every other per-port + * call reports, not a claim that the engine cannot stream. * @see EventStream */ EventStream events(EventStreamOptions options); diff --git a/rift-java-core/src/main/java/io/github/achirdlabs/rift/transport/RemoteTransport.java b/rift-java-core/src/main/java/io/github/achirdlabs/rift/transport/RemoteTransport.java index 7a71f26..26c6706 100644 --- a/rift-java-core/src/main/java/io/github/achirdlabs/rift/transport/RemoteTransport.java +++ b/rift-java-core/src/main/java/io/github/achirdlabs/rift/transport/RemoteTransport.java @@ -380,9 +380,17 @@ public EventStream events(EventStreamOptions options) { int status = response.statusCode(); if (status == 404) { + drainQuietly(response.body()); + // A 404 has two meanings, and its body cannot tell them apart: an engine too old to route + // /events answers with the same "no such resource" envelope an engine that has the stream + // uses for a port naming no imposter. When a port was asked for, the imposter route + // decides it — the same GET every other per-port call makes, so a missing imposter + // surfaces as the same ImposterNotFound they throw (#254). + if (options.port().isPresent()) { + requireImposter(options.port().getAsInt()); + } // An engine too old to serve /events. Same answer as a transport that cannot stream at // all, because the caller's move is the same: poll. - drainQuietly(response.body()); throw new UnsupportedOperationException( "this rift engine has no admin event stream (404 " + path + "); poll recordedSince(...) instead"); } @@ -394,6 +402,11 @@ public EventStream events(EventStreamOptions options) { return new SseEventStream(response.body(), URI.create(base + path), options.idleTimeout()); } + /** Throws the per-port error {@code GET /imposters/{port}} maps to, if the imposter is not there. */ + private void requireImposter(int port) { + executeVoid("GET", "/imposters/" + port, null, OptionalInt.of(port)); + } + private HttpResponse sendStreaming(String path) throws IOException, InterruptedException { if (closed.get()) { throw new IllegalStateException("this RiftTransport is closed"); @@ -573,7 +586,11 @@ private static RuntimeException mapError(int status, String body, OptionalInt po return new EngineError(status, message); } - /** Extracts {@code errors[0].message} from a {@code {"errors":[{"code","message"}]}} error body, falling back to the raw body if it isn't that shape. */ + /** + * Extracts {@code errors[0].message} from a {@code {"errors":[{"code","message"}]}} error body, or + * {@code error} from a bare {@code {"error":"..."}} one, falling back to the raw body if it is + * neither shape. + */ private static String extractErrorMessage(String body) { try { if (JsonValue.parse(body) instanceof JsonObject obj @@ -583,6 +600,11 @@ private static String extractErrorMessage(String body) { && first.get("message") instanceof JsonString message) { return message.value(); } + // Engines through 0.18.1 answer the /events refusals with a bare {"error": "..."}. + if (JsonValue.parse(body) instanceof JsonObject obj + && obj.get("error") instanceof JsonString message) { + return message.value(); + } } catch (RuntimeException ignored) { // Not the {"errors": [...]} shape (or not JSON at all) — fall through to the raw body. } diff --git a/rift-java-core/src/test/java/io/github/achirdlabs/rift/transport/EventStreamTest.java b/rift-java-core/src/test/java/io/github/achirdlabs/rift/transport/EventStreamTest.java index e1feee4..2955638 100644 --- a/rift-java-core/src/test/java/io/github/achirdlabs/rift/transport/EventStreamTest.java +++ b/rift-java-core/src/test/java/io/github/achirdlabs/rift/transport/EventStreamTest.java @@ -10,6 +10,8 @@ import io.github.achirdlabs.rift.error.CommunicationError; import io.github.achirdlabs.rift.error.EngineError; import io.github.achirdlabs.rift.error.EngineUnavailable; +import io.github.achirdlabs.rift.error.ImposterNotFound; +import io.github.achirdlabs.rift.error.InvalidDefinition; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; @@ -271,6 +273,64 @@ void a404MeansThisEngineCannotStreamAtAll() { } } + private static EventStreamOptions onPort(int port) { + return EventStreamOptions.builder().port(port).build(); + } + + private static final String IMPOSTER_NOT_FOUND = + "{\"errors\":[{\"code\":\"404\",\"type\":\"no such resource\",\"message\":\"Imposter not found on port 4545\"}]}"; + + @Test + void anUnknownPortIsImposterNotFoundNotAMissingStream() { + // rift <= 0.18.1: a bare {"error"} body from /events, the canonical envelope from the imposter route. + try (FakeAdminServer s = new FakeAdminServer()) { + s.respond("GET /events", 404, "{\"error\":\"no imposter on port 4545\"}"); + s.respond("GET /imposters/4545", 404, IMPOSTER_NOT_FOUND); + try (Rift rift = connect(s)) { + ImposterNotFound e = assertThrows(ImposterNotFound.class, () -> rift.events(onPort(4545))); + assertEquals(4545, e.port()); + assertEquals("Imposter not found on port 4545", e.getMessage()); + } + } + } + + @Test + void anUnknownPortIsImposterNotFoundUnderTheEnvelopeToo() { + // rift#1226: /events answers the same envelope an engine without the route would, so only the + // imposter lookup can tell them apart. + try (FakeAdminServer s = new FakeAdminServer()) { + s.respond("GET /events", 404, IMPOSTER_NOT_FOUND); + s.respond("GET /imposters/4545", 404, IMPOSTER_NOT_FOUND); + try (Rift rift = connect(s)) { + ImposterNotFound e = assertThrows(ImposterNotFound.class, () -> rift.events(onPort(4545))); + assertEquals(4545, e.port()); + } + } + } + + @Test + void a404ForAPortThatExistsStillMeansThisEngineCannotStream() { + try (FakeAdminServer s = new FakeAdminServer()) { + s.respond("GET /events", 404, "{\"errors\":[{\"code\":\"404\",\"message\":\"Not Found\"}]}"); + s.respond("GET /imposters/4545", 200, "{\"protocol\":\"http\",\"port\":4545}"); + try (Rift rift = connect(s)) { + assertThrows(UnsupportedOperationException.class, () -> rift.events(onPort(4545))); + } + } + } + + @Test + void aBareErrorBodyStillCarriesItsMessage() { + // rift <= 0.18.1 refuses a bad filter with {"error": "..."}; the message must not be the raw JSON. + try (FakeAdminServer s = new FakeAdminServer()) { + s.respond("GET /events", 400, "{\"error\":\"unknown types value 'bogus' (expected requests|lifecycle)\"}"); + try (Rift rift = connect(s)) { + InvalidDefinition e = assertThrows(InvalidDefinition.class, () -> rift.events(opts())); + assertEquals("unknown types value 'bogus' (expected requests|lifecycle)", e.getMessage()); + } + } + } + @Test void aRejectedConnectFailsUpFrontNotMidIteration() { try (FakeAdminServer s = new FakeAdminServer()) {