Skip to content

Commit 697fdaa

Browse files
committed
[do] Broker moved to transport
1 parent 12cf5fb commit 697fdaa

7 files changed

Lines changed: 108 additions & 97 deletions

File tree

src/durable-object/broker.ts

Lines changed: 38 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,6 @@ import {
1010
getClientId,
1111
getPathId,
1212
} from './common.ts';
13-
import {SyncletDurableObject} from './synclet.ts';
1413

1514
export const createDurableObjectBrokerTransport: typeof createDurableObjectBrokerTransportDecl =
1615
({path, brokerPaths, ...options}) => {
@@ -27,62 +26,48 @@ export const createDurableObjectBrokerTransport: typeof createDurableObjectBroke
2726

2827
const sendPacket = async (_packet: string): Promise<void> => {};
2928

30-
const fetch = async (_request: Request): Promise<Response | undefined> => {
31-
return undefined;
29+
const fetch = async (
30+
ctx: DurableObjectState,
31+
request: Request,
32+
): Promise<Response | undefined> => {
33+
const pathId = getPathId(request);
34+
return ifNotNull(
35+
getClientId(request),
36+
(clientId) => {
37+
const [webSocket, client] = objValues(new WebSocketPair());
38+
ctx.acceptWebSocket(client, [clientId, pathId]);
39+
return createResponse(101, webSocket);
40+
},
41+
createUpgradeRequiredResponse,
42+
) as Response;
3243
};
3344

45+
const webSocketMessage = async (
46+
ctx: DurableObjectState,
47+
ws: WebSocket,
48+
message: ArrayBuffer | string,
49+
): Promise<boolean | undefined> =>
50+
ifNotUndefined(ctx.getTags(ws)[0], (clientId) => {
51+
const packet = message.toString();
52+
const splitAt = packet.indexOf(SPACE);
53+
if (splitAt !== -1) {
54+
const to = slice(packet, 0, splitAt);
55+
const remainder = slice(packet, splitAt + 1);
56+
const forwardedPacket = clientId + SPACE + remainder;
57+
if (to === ASTERISK) {
58+
arrayForEach(ctx.getWebSockets(), (otherClient) =>
59+
otherClient !== ws ? otherClient.send(forwardedPacket) : 0,
60+
);
61+
} else if (to != clientId) {
62+
ctx.getWebSockets(to)[0]?.send(forwardedPacket);
63+
}
64+
}
65+
return true;
66+
});
67+
3468
return createDurableObjectTransport(
3569
{connect, disconnect, sendPacket},
36-
{fetch},
70+
{fetch, webSocketMessage},
3771
options,
3872
);
3973
};
40-
41-
export class BrokerOnlyDurableObject<
42-
Env = unknown,
43-
> extends SyncletDurableObject<Env> {
44-
constructor(ctx: DurableObjectState, env: Env) {
45-
super(ctx, env);
46-
}
47-
48-
async fetch(request: Request): Promise<Response> {
49-
const pathId = getPathId(request);
50-
return ifNotNull(
51-
getClientId(request),
52-
(clientId) => {
53-
const [webSocket, client] = objValues(new WebSocketPair());
54-
this.ctx.acceptWebSocket(client, [clientId, pathId]);
55-
return createResponse(101, webSocket);
56-
},
57-
createUpgradeRequiredResponse,
58-
) as Response;
59-
}
60-
61-
webSocketMessage(client: WebSocket, message: ArrayBuffer | string) {
62-
ifNotUndefined(this.ctx.getTags(client)[0], (clientId) =>
63-
this.#handleMessage(clientId, message.toString(), client),
64-
);
65-
}
66-
67-
webSocketClose(_client: WebSocket) {}
68-
69-
#handleMessage(id: string, packet: string, fromClient?: WebSocket) {
70-
const splitAt = packet.indexOf(SPACE);
71-
if (splitAt !== -1) {
72-
const to = slice(packet, 0, splitAt);
73-
const remainder = slice(packet, splitAt + 1);
74-
const forwardedPacket = id + SPACE + remainder;
75-
if (to === ASTERISK) {
76-
arrayForEach(this.#getClients(), (otherClient) =>
77-
otherClient !== fromClient ? otherClient.send(forwardedPacket) : 0,
78-
);
79-
} else if (to != id) {
80-
this.#getClients(to)[0]?.send(forwardedPacket);
81-
}
82-
}
83-
}
84-
85-
#getClients(tag?: string) {
86-
return this.ctx.getWebSockets(tag);
87-
}
88-
}

src/durable-object/common.ts

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,14 +10,27 @@ const PATH_REGEX = /\/([^?]*)/;
1010

1111
export const createDurableObjectTransport = (
1212
{connect, disconnect, sendPacket}: TransportImplementations,
13-
{fetch}: {fetch: (request: Request) => Promise<Response | undefined>},
13+
{
14+
fetch,
15+
webSocketMessage,
16+
}: {
17+
fetch: (
18+
ctx: DurableObjectState,
19+
request: Request,
20+
) => Promise<Response | undefined>;
21+
webSocketMessage?: (
22+
ctx: DurableObjectState,
23+
client: WebSocket,
24+
message: ArrayBuffer | string,
25+
) => Promise<boolean | undefined>;
26+
},
1427
{durableObject, ...options}: DurableObjectTransportOptions,
1528
) => {
1629
const getDurableObject = () => durableObject;
1730

1831
return createTransport({connect, disconnect, sendPacket}, options, {
1932
_brand2: 'DurableObjectTransport',
20-
__: [fetch],
33+
__: [fetch, webSocketMessage],
2134
getDurableObject,
2235
}) as DurableObjectTransport;
2336
};

src/durable-object/index.ts

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,4 @@
1-
export {
2-
BrokerOnlyDurableObject,
3-
createDurableObjectBrokerTransport,
4-
} from './broker.ts';
1+
export {createDurableObjectBrokerTransport} from './broker.ts';
52
export {
63
createDurableObjectSqliteDataConnector,
74
createDurableObjectSqliteMetaConnector,

src/durable-object/synclet.ts

Lines changed: 17 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -25,8 +25,6 @@ import {
2525
} from './common.ts';
2626
import {ProtectedDurableObjectTransport} from './types.ts';
2727

28-
const FETCH = 0;
29-
3028
export abstract class SyncletDurableObject<
3129
Env = unknown,
3230
Depth extends number = number,
@@ -63,15 +61,30 @@ export abstract class SyncletDurableObject<
6361
}
6462

6563
async fetch(request: Request): Promise<Response> {
66-
for (const durableObjectTransport of this.#durableObjectTransports) {
67-
const response = await durableObjectTransport.__[FETCH](request);
64+
for (const {
65+
__: [fetch],
66+
} of this.#durableObjectTransports) {
67+
const response = await fetch(this.ctx, request);
6868
if (!isUndefined(response)) {
6969
return response;
7070
}
7171
}
7272
return createNotImplementedResponse();
7373
}
7474

75+
async webSocketMessage(
76+
ws: WebSocket,
77+
message: string | ArrayBuffer,
78+
): Promise<void> {
79+
for (const {
80+
__: [, webSocketMessage],
81+
} of this.#durableObjectTransports) {
82+
if (await webSocketMessage(this.ctx, ws, message)) {
83+
break;
84+
}
85+
}
86+
}
87+
7588
getCreateDataConnector?(): DataConnectorType;
7689

7790
getCreateMetaConnector?(): MetaConnectorType;

src/durable-object/types.ts

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,15 @@
11
import {DurableObjectTransport} from '@synclets/@types/durable-object';
22

33
export interface ProtectedDurableObjectTransport extends DurableObjectTransport {
4-
__: [fetch: (request: Request) => Promise<Response | undefined>];
4+
__: [
5+
fetch: (
6+
ctx: DurableObjectState,
7+
request: Request,
8+
) => Promise<Response | undefined>,
9+
webSocketMessage: (
10+
ctx: DurableObjectState,
11+
client: WebSocket,
12+
message: ArrayBuffer | string,
13+
) => Promise<boolean | undefined>,
14+
];
515
}

test/durable-object/broker-only.test.ts

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -63,8 +63,8 @@ test('accept WebSocket upgrade requests', async () => {
6363
});
6464

6565
test('2 clients communicate', async () => {
66-
const a = await createClients(2);
67-
const [[webSocket1, webSocket2], [received1, received2]] = a;
66+
const [[webSocket1, webSocket2], [received1, received2]] =
67+
await createClients(2);
6868

6969
expect(await api('getClientCount')).toEqual(2);
7070

@@ -95,6 +95,9 @@ test('2 clients communicate', async () => {
9595
[client1Id, 'from1To*'],
9696
[client1Id, 'from1To2'],
9797
]);
98+
99+
webSocket1.close();
100+
webSocket2.close();
98101
});
99102

100103
test('3 clients communicate', async () => {
@@ -103,6 +106,8 @@ test('3 clients communicate', async () => {
103106
[received1, received2, received3],
104107
] = await createClients(3);
105108

109+
expect(await api('getClientCount')).toEqual(3);
110+
106111
webSocket1.send('* from1To*');
107112
await pause(10);
108113
webSocket2.send('* from2To*');
@@ -155,4 +160,8 @@ test('3 clients communicate', async () => {
155160
[client1Id, 'from1To3'],
156161
[client2Id, 'from2To3'],
157162
]);
163+
164+
webSocket1.close();
165+
webSocket2.close();
166+
webSocket3.close();
158167
});

test/durable-object/servers.ts

Lines changed: 15 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,4 @@
11
import {
2-
BrokerOnlyDurableObject,
32
createDurableObjectBrokerTransport,
43
createDurableObjectSqliteDataConnector,
54
createDurableObjectSqliteMetaConnector,
@@ -21,45 +20,30 @@ export class TestSyncletDurableObject extends SyncletDurableObject {
2120
async api(method: string, ...args: any[]): Promise<any> {
2221
return await api(this, method, ...args);
2322
}
24-
}
25-
26-
export class TestBrokerOnlyDurableObject extends BrokerOnlyDurableObject {
27-
async api(method: string, ...args: any[]): Promise<any> {
28-
return await api(this, method, ...args);
29-
}
3023

3124
getClientCount() {
3225
return this.ctx.getWebSockets().length;
3326
}
3427
}
3528

36-
export class TestBrokerOnlyDurableObject2 extends SyncletDurableObject {
37-
getCreateComponents() {
38-
return {
39-
transport: createDurableObjectBrokerTransport({durableObject: this}),
40-
};
41-
}
42-
43-
async api(method: string, ...args: any[]): Promise<any> {
44-
return await api(this, method, ...args);
45-
}
46-
47-
getClientCount() {
48-
return this.ctx.getWebSockets().length;
29+
export class TestBrokerOnlyDurableObject extends TestSyncletDurableObject {
30+
getCreateTransport() {
31+
return createDurableObjectBrokerTransport({durableObject: this});
4932
}
5033
}
5134

5235
export class TestConnectorsOnlyDurableObject extends TestSyncletDurableObject {
53-
getCreateComponents() {
54-
return {
55-
dataConnector: createDurableObjectSqliteDataConnector({
56-
depth: 1,
57-
sqlStorage: this.ctx.storage.sql,
58-
}),
59-
metaConnector: createDurableObjectSqliteMetaConnector({
60-
depth: 1,
61-
sqlStorage: this.ctx.storage.sql,
62-
}),
63-
};
36+
getCreateDataConnector() {
37+
return createDurableObjectSqliteDataConnector({
38+
depth: 1,
39+
sqlStorage: this.ctx.storage.sql,
40+
});
41+
}
42+
43+
getCreateMetaConnector() {
44+
return createDurableObjectSqliteMetaConnector({
45+
depth: 1,
46+
sqlStorage: this.ctx.storage.sql,
47+
});
6448
}
6549
}

0 commit comments

Comments
 (0)