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
44 changes: 44 additions & 0 deletions packages/client-runtime/src/rpc/session.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1112,6 +1112,50 @@ describe("RpcSessionFactory", () => {
}).pipe(Effect.scoped, Effect.provide(TestClock.layer())),
);

it.effect("keeps reading replies after closing a stream with a full buffer", () =>
Effect.gen(function* () {
const { factory, sockets } = yield* makeFactory();
const session = yield* factory.connect(PREPARED);
const readyFiber = yield* Effect.forkChild(session.ready);
const socket = yield* awaitSocket(sockets);
socket.open();
yield* completeInitialConfig(socket);
yield* Fiber.join(readyFiber);

// The consumer takes one event and stops pulling, so the stream buffer fills.
const consuming = yield* Deferred.make<void>();
const streamFiber = yield* session.client[WS_METHODS.subscribeServerConfig]({}).pipe(
Stream.runForEach(() =>
Deferred.succeed(consuming, undefined).pipe(Effect.andThen(Effect.never)),
),
Effect.forkChild,
);
const streamRequest = yield* awaitRequest(socket, 1);
const snapshot = { version: 1, type: "snapshot", config: ENCODED_SERVER_CONFIG };
socket.serverMessage(
encodeJson({
_tag: "Chunk",
requestId: streamRequest.id,
values: Array.from({ length: 64 }, () => snapshot),
}),
);
yield* Deferred.await(consuming);
yield* Fiber.interrupt(streamFiber);

const probeFiber = yield* Effect.forkChild(session.probe);
const probeRequest = yield* awaitRequest(socket, 2);
expect(probeRequest).toMatchObject({ tag: WS_METHODS.serverProbe });
socket.serverMessage(
encodeJson({
_tag: "Exit",
requestId: probeRequest.id,
exit: { _tag: "Success", value: {} },
}),
);
yield* Fiber.join(probeFiber);
}).pipe(Effect.scoped),
);

it.effect("reaches ready when a newer server sends unknown config members", () =>
Effect.gen(function* () {
const { factory, sockets } = yield* makeFactory();
Expand Down
27 changes: 14 additions & 13 deletions patches/effect@4.0.0-rc.115.patch
Original file line number Diff line number Diff line change
Expand Up @@ -157,7 +157,7 @@ index c586c60..49886d8 100644
});
});
const onStreamRequest = Effect.fnUntraced(function* (rpc, middleware, payload, headers, streamBufferSize, context) {
@@ -143,11 +150,17 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
@@ -143,11 +150,18 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
});
const fiber = Fiber.getCurrent();
const id = generateRequestId();
Expand All @@ -173,11 +173,12 @@ index c586c60..49886d8 100644
+ if (!entry) return Effect.void;
entries.delete(id);
- return sendInterrupt(id, Exit.isFailure(exit) ? Array.from(Cause.interruptors(exit.cause)) : [], context);
+ return sendInterrupt(id, Exit.isFailure(exit) ? Array.from(Cause.interruptors(exit.cause)) : [], context, entry.rpc._tag);
+ // Shutting down releases a socket reader blocked offering into this full buffer.
+ return Queue.shutdown(entry.queue).pipe(Effect.andThen(sendInterrupt(id, Exit.isFailure(exit) ? Array.from(Cause.interruptors(exit.cause)) : [], context, entry.rpc._tag)));
});
const queue = yield* Queue.bounded(streamBufferSize);
entries.set(id, {
@@ -155,9 +168,10 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
@@ -155,9 +169,10 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
rpc,
queue,
scope,
Expand All @@ -190,7 +191,7 @@ index c586c60..49886d8 100644
message,
context,
discard: false
@@ -172,7 +186,7 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
@@ -172,7 +187,7 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
sampled: span.sampled
} : {}),
headers: Headers.merge(fiber.getRef(CurrentHeaders), headers)
Expand All @@ -199,7 +200,7 @@ index c586c60..49886d8 100644
captureStackTrace: false
}) : identity, Effect.catchCause(error => Queue.failCause(queue, error)), Effect.interruptible, Effect.forkIn(scope, {
startImmediately: true
@@ -202,7 +216,10 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
@@ -202,7 +217,10 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
});
};
};
Expand All @@ -211,7 +212,7 @@ index c586c60..49886d8 100644
const parentFiber = Fiber.getCurrent();
const fiber = options.onFromClient({
message: {
@@ -216,7 +233,7 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
@@ -216,7 +234,7 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
fiber.addObserver(() => {
resume(Effect.void);
});
Expand All @@ -220,7 +221,7 @@ index c586c60..49886d8 100644
const write = message => {
switch (message._tag) {
case "Chunk":
@@ -224,7 +241,13 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
@@ -224,7 +242,13 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
const requestId = message.requestId;
const entry = entries.get(requestId);
if (!entry || entry._tag !== "Queue") return Effect.void;
Expand All @@ -235,7 +236,7 @@ index c586c60..49886d8 100644
message: {
_tag: "Ack",
requestId: message.requestId
@@ -239,11 +262,18 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
@@ -239,11 +263,18 @@ export const makeNoSerialization = /*#__PURE__*/Effect.fnUntraced(function* (gro
const entry = entries.get(requestId);
if (!entry) return Effect.void;
entries.delete(requestId);
Expand All @@ -257,7 +258,7 @@ index c586c60..49886d8 100644
}
case "Defect":
{
@@ -598,7 +628,7 @@ export const makeProtocolSocket = options => Protocol.make(Effect.fnUntraced(fun
@@ -598,7 +629,7 @@ export const makeProtocolSocket = options => Protocol.make(Effect.fnUntraced(fun
// `parser` is replaced on every connect, and a stateful serialization
// encodes against the connection it is writing to, so the ping is encoded
// when it is sent rather than once up front.
Expand All @@ -266,7 +267,7 @@ index c586c60..49886d8 100644
let currentError;
const broadcast = response => Effect.forEach(clientIds, clientId => writeResponse(clientId, response));
const broadcastError = error => {
@@ -618,8 +648,7 @@ export const makeProtocolSocket = options => Protocol.make(Effect.fnUntraced(fun
@@ -618,8 +649,7 @@ export const makeProtocolSocket = options => Protocol.make(Effect.fnUntraced(fun
body: () => {
const response = responses[i++];
if (response._tag === "Pong") {
Expand All @@ -276,7 +277,7 @@ index c586c60..49886d8 100644
}
if (Object.hasOwn(response, "requestId")) {
const requestId = response.requestId;
@@ -664,12 +693,12 @@ export const makeProtocolSocket = options => Protocol.make(Effect.fnUntraced(fun
@@ -664,12 +694,12 @@ export const makeProtocolSocket = options => Protocol.make(Effect.fnUntraced(fun
yield* processData(frames[i]);
}
}
Expand All @@ -291,7 +292,7 @@ index c586c60..49886d8 100644
}).pipe(Option.isSome(hooks) ? Effect.ensuring(hooks.value.onDisconnect) : identity, Effect.tapCause(cause => {
const error = Cause.findError(cause);
const hasError = Result.isSuccess(error);
@@ -710,20 +739,28 @@ export const makeProtocolSocket = options => Protocol.make(Effect.fnUntraced(fun
@@ -710,20 +740,28 @@ export const makeProtocolSocket = options => Protocol.make(Effect.fnUntraced(fun
};
}));
const defaultRetryPolicy = /*#__PURE__*/Schedule.min([/*#__PURE__*/Schedule.exponential(500, 1.5), /*#__PURE__*/Schedule.spaced(5000)]);
Expand Down Expand Up @@ -325,7 +326,7 @@ index c586c60..49886d8 100644
}).pipe(Effect.delay("5 seconds"), Effect.ignore, Effect.forever, Effect.interruptible, Effect.forkScoped);
return {
timeout: latch.await,
@@ -874,6 +911,11 @@ export const makeProtocolWorker = options => Protocol.make(Effect.fnUntraced(fun
@@ -874,6 +912,11 @@ export const makeProtocolWorker = options => Protocol.make(Effect.fnUntraced(fun
* @since 4.0.0
*/
export const layerProtocolWorker = /*#__PURE__*/flow(makeProtocolWorker, /*#__PURE__*/Layer.effect(Protocol));
Expand Down
Loading
Loading