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
3 changes: 2 additions & 1 deletion docs/design/sdk-api.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions docs/events.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -84,6 +85,20 @@ Stream<DynamicTest> theStreamPushesWhatTheJournalRecords() {
});
}

@TestFactory
Stream<DynamicTest> 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<DynamicTest> aBrokenStreamIsRecoveredByPollingFromTheLastIndexSeen() {
return gated("the reconcile loop, end to end", () -> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
Expand All @@ -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<InputStream> sendStreaming(String path) throws IOException, InterruptedException {
if (closed.get()) {
throw new IllegalStateException("this RiftTransport is closed");
Expand Down Expand Up @@ -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
Expand All @@ -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.
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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()) {
Expand Down
Loading