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
5 changes: 5 additions & 0 deletions .changeset/rpc-promise-from-promise.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"capnweb": minor
---

`RpcPromise` can now be constructed from a `Promise`: pipelined calls queue in order until it settles, so you can publish a capability that doesn't exist yet.
22 changes: 22 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -267,6 +267,28 @@ let profile = await api.getUserProfile(user.id);

Whenever an `RpcPromise` is passed in the parameters to an RPC, or returned as part of the result, the promise will be replaced with its resolution before delivery to the receiving application. So, you can use an `RpcPromise<T>` anywhere where a `T` is required!

#### Constructing `RpcPromise` from a `Promise`

You can construct an `RpcPromise<T>` directly from a regular `Promise<T>`, allowing you to perform promise pipelining on a regular local promise. Pipelined calls will wait until the inner promise resolves, then will be delivered, in-order, to the resolution. This is useful when you plan to obtain some stub in the future, but you want to allow code to start queuing calls on it immediately.

Wrapping a `Promise<T>` in this way is semantically identical to creating a local-loopback RPC and then invoking it. That is:

```ts
// this...
let rpcPromise = new RpcPromise(myPromise);

// is semantically the same as this...
let rpcFunc = new RpcStub(() => myPromise);
let rpcPromise = rpcFunc();
```

In other words, this means:
* The result of the promise must be serializable.
* If the promise resolution contains `RpcTarget`s or `Function`s, the `RpcPromise`'s resolution will replace them with stubs.
* Ownership of any stubs in the Promise result is transferred away. If you want to keep your own copies, you need to `dup()` them.
* If the promise rejects, the rejection propagates to all pipelined calls.
* etc.

### The magic `map()` method

Every RPC promise has a special method `.map()` which can be used to remotely transform a value, without pulling it back locally. Here's an example:
Expand Down
247 changes: 246 additions & 1 deletion __tests__/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@

import { expect, it, describe, inject } from "vitest"
import { deserialize, serialize, RpcSession, type RpcSessionOptions, RpcTransport,
type RpcTransportWithCustomEncoding, RpcTarget, RpcStub, newWebSocketRpcSession,
type RpcTransportWithCustomEncoding, RpcTarget, RpcStub, RpcPromise, newWebSocketRpcSession,
newMessagePortRpcSession,
newHttpBatchRpcSession} from "../src/index.js"
import { swapByteOrder } from "../src/serialize.js"
Expand Down Expand Up @@ -549,6 +549,7 @@ class TestTransport implements RpcTransport {
private waiter?: () => void;
private aborter?: (err: any) => void;
public log = false;
public sentLog: string[] = [];
private fenced = false;

send(message: string): void {
Expand All @@ -557,6 +558,7 @@ class TestTransport implements RpcTransport {
message = message.replaceAll("$remove$", "");

if (this.log) console.log(`${this.name}: ${message}`);
this.sentLog.push(message);
this.partner!.queue.push(message);
if (this.partner!.waiter && !this.partner!.fenced) {
this.partner!.waiter();
Expand Down Expand Up @@ -2378,6 +2380,249 @@ describe("RpcImportHook argument disposal", () => {

// =======================================================================================

describe("constructing RpcPromise from a promise", () => {
it("pipelines through a pending promise without pulling the resolution", async () => {
await using harness = new TestHarness(new TestTarget());

let {promise, resolve} = Promise.withResolvers<RpcStub<TestTarget>>();
using stub = new RpcPromise<TestTarget>(promise);

using counter = stub.makeCounter(1);
let result = counter.increment(2);
resolve(harness.stub.dup());
expect(await result).toBe(3);

// Only the final result was pulled: neither the promise's resolution nor the intermediate
// counter was transmitted.
let sent = harness.clientTransport.sentLog;
expect(sent.some(msg => msg.startsWith('["push"'))).toBe(true);
expect(sent.filter(msg => msg.startsWith('["pull"'))).toHaveLength(1);
});

it("queues calls made before resolution and delivers them in order", async () => {
let calls: number[] = [];
class Recorder extends RpcTarget {
record(i: number) { calls.push(i); return i; }
}

let {promise, resolve} = Promise.withResolvers<Recorder>();
using stub = new RpcPromise<Recorder>(promise);

let results = [stub.record(1), stub.record(2), stub.record(3)];
expect(calls).toStrictEqual([]);

resolve(new Recorder());
expect(await Promise.all(results)).toStrictEqual([1, 2, 3]);
expect(calls).toStrictEqual([1, 2, 3]);
});

it("accepts a promise for a target, a remote stub, or a plain value", async () => {
await using harness = new TestHarness(new TestTarget());

using target = new RpcPromise<Counter>(Promise.resolve(new Counter(1)));
expect(await target.increment()).toBe(2);

using remote = new RpcPromise<TestTarget>(Promise.resolve(harness.stub.dup()));
expect(await remote.square(3)).toBe(9);

using value = new RpcPromise<{foo: number}>(Promise.resolve({foo: 123}));
expect(await value.foo).toBe(123);
});

it("awaiting the RpcPromise yields the resolution", async () => {
using plain = new RpcPromise<{foo: number}>(Promise.resolve({foo: 123}));
expect(await plain).toStrictEqual({foo: 123});

await using harness = new TestHarness(new TestTarget());
using remote = new RpcPromise<TestTarget>(Promise.resolve(harness.stub.dup()));
let resolved = await remote;
expect(await resolved.square(4)).toBe(16);
});

it("reports rejection to queued calls, await, and onRpcBroken", async () => {
let error = new Error("nope");
using stub = new RpcPromise<Counter>(Promise.reject(error));

let broken: any[] = [];
stub.onRpcBroken(err => { broken.push(err); });

await expect(() => stub.increment()).rejects.toThrow("nope");
await expect(Promise.resolve(stub)).rejects.toThrow("nope");
expect(broken).toStrictEqual([error]);
});

it("does not report an unhandled rejection for an unused promise", async () => {
new RpcPromise<Counter>(Promise.reject(new Error("ignored")));
await pumpMicrotasks();
});

it("does not report an unhandled rejection for a disposed, unawaited queued call", async () => {
using stub = new RpcPromise<Counter>(Promise.reject(new Error("ignored")));
using result = stub.increment(); // never awaited; disposal alone must observe the error
await pumpMicrotasks();
});

it("does not report an unhandled rejection for a disposed, unawaited map() result", async () => {
using stub = new RpcPromise<number>(Promise.reject(new Error("ignored")));
using result = stub.map(i => i); // never awaited; disposal alone must observe the error
await pumpMicrotasks();
});

it("disposes the eventual target when disposed before resolution", async () => {
let disposed = false;
class Disposable extends RpcTarget {
[Symbol.dispose]() { disposed = true; }
}

let {promise, resolve} = Promise.withResolvers<Disposable>();
let stub = new RpcPromise<Disposable>(promise);
stub[Symbol.dispose]();

resolve(new Disposable());
await pumpMicrotasks();
expect(disposed).toBe(true);
});

it("delivers a call initiated before disposal", async () => {
let disposed = false;
class DisposableCounter extends Counter {
[Symbol.dispose]() { disposed = true; }
}

let stub = new RpcPromise<DisposableCounter>(Promise.resolve(new DisposableCounter(1)));
await pumpMicrotasks();

let result = stub.increment(2);
stub[Symbol.dispose]();

expect(disposed).toBe(false);
expect(await result).toBe(3);
expect(disposed).toBe(true);
});

it("keeps an adopted RpcPromise lazy", async () => {
await using harness = new TestHarness(new TestTarget());

using counter = new RpcPromise<Counter>(harness.stub.makeCounter(1));
expect(await counter.increment(2)).toBe(3);

let sent = harness.clientTransport.sentLog;
expect(sent.filter(msg => msg.startsWith('["pull"'))).toHaveLength(1);
});

it("transmits nothing when constructed from a remote property promise", async () => {
await using harness = new TestHarness(new TestTarget());

using counter = harness.stub.makeCounter(5);
await pumpMicrotasks(); // let the makeCounter push flush

let source = counter.value;
let wrapped = new RpcPromise<number>(source);

// Construction shares the source's hook and path; the get() producing an independent hook
// happens lazily on first await, so nothing goes over the wire yet.
let sentBefore = harness.clientTransport.sentLog.length;
await pumpMicrotasks();
expect(harness.clientTransport.sentLog.length).toBe(sentBefore);

expect(await wrapped).toBe(5);

// The source property promise remains usable.
expect(await source).toBe(5);
});

it("does not invoke a local getter when constructed from a property promise", async () => {
let reads = 0;
class Gettable extends RpcTarget {
get prop() { ++reads; return 42; }
}

using stub = new RpcStub(new Gettable());
let source = stub.prop;
let wrapped = new RpcPromise<number>(source);

await pumpMicrotasks();
expect(reads).toBe(0);

expect(await wrapped).toBe(42);
expect(await source).toBe(42);
});

it("consumes the source when adopting an existing RpcPromise", async () => {
await using harness = new TestHarness(new TestTarget());

let source = harness.stub.makeCounter(1);
using wrapper = new RpcPromise<Counter>(source);

// The source was neutered: using it now reports the standard disposed error, and disposing
// it is a harmless no-op that doesn't affect the wrapper.
await expect(source.increment(1)).rejects.toThrow(
"Attempted to use RPC stub after it has been disposed.");
source[Symbol.dispose]();

expect(await wrapper.increment(2)).toBe(3);
});

it("stays lazy when a deferred promise is resolved with dup()", async () => {
await using harness = new TestHarness(new TestTarget());

using counter = harness.stub.makeCounter(1);

// Resolving with the RpcPromise itself would let the native promise machinery assimilate it
// as a thenable, pulling the resolution. dup() returns a non-thenable stub, which the
// resolution adopts, keeping calls pipelined.
let {promise, resolve} = Promise.withResolvers<RpcStub<Counter>>();
using stub = new RpcPromise<Counter>(promise);

let result = stub.increment(2);
resolve(counter.dup());
expect(await result).toBe(3);

let sent = harness.clientTransport.sentLog;
expect(sent.filter(msg => msg.startsWith('["pull"'))).toHaveLength(1);
});

it("preserves brokenness of a bare stub it was constructed from", async () => {
await using harness = new TestHarness(new TestTarget());
using stub = new RpcPromise<TestTarget>(<any>harness.stub.dup());

let errors: any[] = [];
stub.onRpcBroken(error => { errors.push(error); });

harness.clientTransport.forceReceiveError(new Error("test disconnect"));
await pumpMicrotasks();
expect(errors).toStrictEqual([new Error("test disconnect")]);
});

it("keeps disposal idempotent when constructed from a bare stub", async () => {
let disposals = 0;
class Disposable extends RpcTarget {
[Symbol.dispose]() { ++disposals; }
}

let inner = new RpcStub(new Disposable());
let outer = new RpcPromise<Disposable>(<any>inner);
inner[Symbol.dispose]();
outer[Symbol.dispose]();

await pumpMicrotasks();
expect(disposals).toBe(1);
});

it("resolves when awaited after construction from a bare local stub", async () => {
// Regression test: the constructor previously adopted a bare stub's hook directly, producing
// a promise whose pipelined calls worked but whose await rejected, because non-promise hooks
// don't implement pull().
using stub = new RpcPromise<Counter>(<any>new RpcStub(new Counter(1)));

expect(await stub.increment(2)).toBe(3);
let resolved = await stub;
expect(await resolved.increment(3)).toBe(6);
});
});

// =======================================================================================

describe("HTTP requests", () => {
it("can perform a batch HTTP request", async () => {
let cap = newHttpBatchRpcSession<TestTarget>(`http://${inject("testServerHost")}`);
Expand Down
60 changes: 59 additions & 1 deletion __tests__/workerd.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
/// <reference types="@cloudflare/workers-types" />
import { expect, it, describe } from "vitest";
import { RpcStub as NativeRpcStub, RpcTarget as NativeRpcTarget, env, DurableObject } from "cloudflare:workers";
import { newHttpBatchRpcSession, newWebSocketRpcSession, RpcStub, RpcTarget } from "../src/index-workers.js";
import { newHttpBatchRpcSession, newWebSocketRpcSession, RpcStub, RpcPromise, RpcTarget } from "../src/index-workers.js";
import { v, wrapServerTarget, type ServiceValidator } from "../packages/capnweb-validate/src/internal/core.js";
import { Counter, TestTarget } from "./test-util.js";

Expand Down Expand Up @@ -44,6 +44,10 @@ class CounterFactory extends RpcTarget {
return new NativeRpcStub(new NativeCounter());
}

getBroken(): NativeRpcStub<NativeCounter> {
throw new RangeError("test error");
}

getNativeEmbedded() {
return {stub: new NativeRpcStub(new NativeCounter())};
}
Expand All @@ -57,6 +61,12 @@ class CounterFactory extends RpcTarget {
}
}

async function pumpMicrotasks() {
for (let i = 0; i < 16; i++) {
await Promise.resolve();
}
}

describe("workerd compatibility", () => {
it("allows native RpcStubs to be created using userspace RpcTargets", async () => {
let stub = new NativeRpcStub(new JsCounter());
Expand Down Expand Up @@ -126,6 +136,54 @@ describe("workerd compatibility", () => {
}
})

it("can wrap a native promise in a userspace promise", async () => {
// Wrapping in RpcPromise (rather than RpcStub) exercises the rpc-thenable adoption path,
// which pipelines calls on the native thenable without eagerly awaiting it.
let factory = new NativeRpcStub(new CounterFactory());
let stub = new RpcPromise(factory.getNative());
expect(await stub.increment()).toBe(1);
expect(await stub.increment()).toBe(2);

expect(await stub.value).toBe(2);
})

it("can dup a userspace promise wrapping a native promise", async () => {
let factory = new NativeRpcStub(new CounterFactory());
let promise = new RpcPromise(factory.getNative());

// dup() routes through get([]), which must produce an independent hook aliasing the same
// underlying native promise.
let dup = promise.dup();
expect(await dup.increment()).toBe(1);
expect(await promise.increment()).toBe(2);
})

it("can pass a wrapped native promise as an RPC argument", async () => {
class CounterUser extends RpcTarget {
useCounter(counter: RpcStub<NativeCounter>) {
return counter.increment(5);
}
}

let factory = new NativeRpcStub(new CounterFactory());
let user = new RpcStub(new CounterUser());
let arg = new RpcPromise(factory.getNative());
expect(await user.useCounter(<any>arg)).toBe(5);
})

it("reports brokenness when a wrapped native promise rejects", async () => {
let factory = new NativeRpcStub(new CounterFactory());
let promise = new RpcPromise(<any>factory.getBroken());

let errors: any[] = [];
promise.onRpcBroken(err => { errors.push(err); });

await expect(Promise.resolve(promise)).rejects.toThrow("test error");
await pumpMicrotasks();
expect(errors.length).toBe(1);
expect(errors[0].message).toBe("test error");
})

it("can pipeline on a native stub returned from a userspace call", async () => {
{
let factory = new RpcStub(new CounterFactory());
Expand Down
Loading
Loading