diff --git a/package.json b/package.json index c359dc9..a2e96c2 100644 --- a/package.json +++ b/package.json @@ -31,6 +31,7 @@ "@atcute/client": "^4.2.1", "@atcute/crypto": "^2.3.0", "@atcute/did-plc": "^0.3.2", + "@atcute/firehose": "^0.1.0", "@atcute/identity": "^1.1.3", "@atcute/identity-resolver": "^1.2.2", "@atcute/lexicon-doc": "^2.1.1", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 33aee15..4a8c8c9 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -32,6 +32,9 @@ importers: '@atcute/did-plc': specifier: ^0.3.2 version: 0.3.2 + '@atcute/firehose': + specifier: ^0.1.0 + version: 0.1.0 '@atcute/identity': specifier: ^1.1.3 version: 1.1.3 @@ -165,6 +168,9 @@ packages: '@atcute/did-plc@0.3.2': resolution: {integrity: sha512-zOqk5mcZJa+xnfpYFN8aIgRA5uQdTDeiDvQeWWIZKOslFBJTtWWL912FLI8r5JWzIc7SgYEp+SbpXmAB6t26EA==} + '@atcute/firehose@0.1.0': + resolution: {integrity: sha512-xBEKdi6rkODpCIIRpXtXhhcuQ1vTbufDykAM2kA6bmWEpuwI4acoUCg9zbiUWqd21SMMOAthu9Eh72i2VYrD7A==} + '@atcute/identity-resolver@1.2.2': resolution: {integrity: sha512-eUh/UH4bFvuXS0X7epYCeJC/kj4rbBXfSRumLEH4smMVwNOgTo7cL/0Srty+P/qVPoZEyXdfEbS0PHJyzoXmHw==} peerDependencies: @@ -721,6 +727,12 @@ packages: '@marijn/find-cluster-break@1.0.2': resolution: {integrity: sha512-l0h88YhZFyKdXIFNfSWpyjStDjGHwZ/U7iobcK1cQQD8sejsONdQtTVU+1wVN1PBw40PiiHB1vA5S7VTfQiP9g==} + '@mary-ext/event-iterator@1.0.0': + resolution: {integrity: sha512-l6gCPsWJ8aRCe/s7/oCmero70kDHgIK5m4uJvYgwEYTqVxoBOIXbKr5tnkLqUHEg6mNduB4IWvms3h70Hp9ADQ==} + + '@mary-ext/simple-event-emitter@1.0.1': + resolution: {integrity: sha512-9+VvZisxZ/gSg+JJH7hmXaA8Qj42Qjz3O58RSB+INYc8iLA0icATZxHB9vKbj59ojDGZjO3hCKzMXocx3L0H8w==} + '@noble/secp256k1@3.0.0': resolution: {integrity: sha512-NJBaR352KyIvj3t6sgT/+7xrNyF9Xk9QlLSIqUGVUYlsnDTAUqY8LOmwpcgEx4AMJXRITQ5XEVHD+mMaPfr3mg==} @@ -1110,6 +1122,9 @@ packages: esm-env@1.2.2: resolution: {integrity: sha512-Epxrv+Nr/CaL4ZcFGPJIYLWFom+YeV1DqMLHJoEd9SYRxNbaFruBwfEX/kkHUJf55j2+TUbmDcmuilbP1TmXHA==} + event-target-polyfill@0.0.4: + resolution: {integrity: sha512-Gs6RLjzlLRdT8X9ZipJdIZI/Y6/HhRLyq9RdDlCsnpxr/+Nn6bU2EFGuC94GjxqhM+Nmij2Vcq98yoHrU8uNFQ==} + fdir@6.5.0: resolution: {integrity: sha512-tIbYtZbucOs0BRGqPJkshJUYdL+SDH7dVM8gjy+ERp3WAUjLEFJE+02kanyHtwjWOnwrKYBiwAmM0p4kLJAnXg==} engines: {node: '>=12.0.0'} @@ -1287,6 +1302,14 @@ packages: parse5@7.3.0: resolution: {integrity: sha512-IInvU7fabl34qmi9gY8XOVxhYyMyuH2xUNpb2q8/Y+7552KlejkRvqvD19nMoUW/uQGGbqNpA6Tufu5FL5BZgw==} + partysocket@1.1.16: + resolution: {integrity: sha512-d7xFv+ZC7x0p/DAHWJ5FhxQhimIx+ucyZY+kxL0cKddLBmK9c4p2tEA/L+dOOrWm6EYrRwrBjKQV0uSzOY9x1w==} + peerDependencies: + react: '>=17' + peerDependenciesMeta: + react: + optional: true + pathe@2.0.3: resolution: {integrity: sha512-WUjGcAqP1gQacoQe+OBJsFA7Ld4DyXuUIjZ5cc75cLHvJ7dtNsTugphxIADwspS+AraAUePCKrSVtPLFj/F88w==} @@ -1440,6 +1463,10 @@ packages: engines: {node: '>=18.0.0'} hasBin: true + type-fest@4.41.0: + resolution: {integrity: sha512-TeTSQ6H5YHvpqVwBRcnLDCBnDOHWYu7IvGbHT6N8AOymcr9PJGjc1GTtiWZTYg0NCgYwvnYWEkVChQAr9bjfwA==} + engines: {node: '>=16'} + typescript@5.9.3: resolution: {integrity: sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw==} engines: {node: '>=14.17'} @@ -1593,6 +1620,18 @@ snapshots: '@atcute/util-fetch': 1.0.5 '@badrap/valita': 0.4.6 + '@atcute/firehose@0.1.0': + dependencies: + '@atcute/cbor': 2.3.2 + '@atcute/lexicons': 1.2.9 + '@atcute/uint8array': 1.1.1 + '@mary-ext/event-iterator': 1.0.0 + '@mary-ext/simple-event-emitter': 1.0.1 + partysocket: 1.1.16 + type-fest: 4.41.0 + transitivePeerDependencies: + - react + '@atcute/identity-resolver@1.2.2(@atcute/identity@1.1.3)': dependencies: '@atcute/identity': 1.1.3 @@ -2119,6 +2158,12 @@ snapshots: '@marijn/find-cluster-break@1.0.2': {} + '@mary-ext/event-iterator@1.0.0': + dependencies: + yocto-queue: 1.2.2 + + '@mary-ext/simple-event-emitter@1.0.1': {} + '@noble/secp256k1@3.0.0': {} '@rollup/rollup-android-arm-eabi@4.59.0': @@ -2487,6 +2532,8 @@ snapshots: esm-env@1.2.2: {} + event-target-polyfill@0.0.4: {} + fdir@6.5.0(picomatch@4.0.3): optionalDependencies: picomatch: 4.0.3 @@ -2613,6 +2660,10 @@ snapshots: dependencies: entities: 6.0.1 + partysocket@1.1.16: + dependencies: + event-target-polyfill: 0.0.4 + pathe@2.0.3: {} picocolors@1.1.1: {} @@ -2736,6 +2787,8 @@ snapshots: fsevents: 2.3.3 optional: true + type-fest@4.41.0: {} + typescript@5.9.3: {} ufo@1.6.3: {} diff --git a/src/views/stream/index.tsx b/src/views/stream/index.tsx index dcca962..581bddb 100644 --- a/src/views/stream/index.tsx +++ b/src/views/stream/index.tsx @@ -1,4 +1,5 @@ -import { Firehose } from "@skyware/firehose"; +import { ComAtprotoSyncSubscribeRepos } from "@atcute/atproto"; +import { FirehoseSubscription } from "@atcute/firehose"; import { Title } from "@solidjs/meta"; import { A, useLocation, useSearchParams } from "@solidjs/router"; import { createSignal, For, onCleanup, onMount, Show } from "solid-js"; @@ -117,7 +118,7 @@ export const StreamView = () => { const [currentTime, setCurrentTime] = createSignal(Date.now()); let socket: WebSocket; - let firehose: Firehose; + let firehoseIterator: AsyncIterator; let formRef!: HTMLFormElement; let pendingRecords: any[] = []; let rafId: number | null = null; @@ -159,7 +160,7 @@ export const StreamView = () => { const disconnect = () => { if (!config().useFirehoseLib) socket?.close(); - else firehose?.close(); + else firehoseIterator?.return?.(); if (rafId !== null) { cancelAnimationFrame(rafId); @@ -256,43 +257,69 @@ export const StreamView = () => { disconnect(); }); } else { - const cursor = formData.get("cursor")?.toString(); - firehose = new Firehose({ - relay: url, - cursor: cursor, - autoReconnect: false, + const cursorParam = formData.get("cursor")?.toString(); + const cursor = cursorParam ? parseInt(cursorParam, 10) : undefined; + let reconnectCursor: number; + const firehose = new FirehoseSubscription({ + service: url, + nsid: ComAtprotoSyncSubscribeRepos.mainSchema, + params: () => ({ cursor: reconnectCursor ?? cursor }), + onConnectionOpen() { + setNotice(""); + setConnected(true); + }, + onConnectionClose(ev) { + reconnectCursor = records().at(-1)?.seq; + console.log(ev); + if (!ev.wasClean) { + setNotice("Reconnecting..."); + } + }, + onConnectionError(err) { + console.error(err); + setNotice(`Connection error: ${err.message}`); + disconnect(); + }, }); - firehose.on("error", (err) => { - console.error(err); - const message = err instanceof Error ? err.message : "Unknown error"; - setNotice(`Connection error: ${message}`); - disconnect(); - }); - firehose.on("commit", (commit) => { - for (const op of commit.ops) { - addRecord({ - $type: commit.$type, - repo: commit.repo, - seq: commit.seq, - time: commit.time, - rev: commit.rev, - since: commit.since, - op: op, - }); + + firehoseIterator = firehose[Symbol.asyncIterator](); + + while (true) { + const { value: message, done } = await firehoseIterator.next(); + if (done) break; + + switch (message.$type) { + case "com.atproto.sync.subscribeRepos#commit": { + for (const op of message.ops) { + addRecord({ + $type: message.$type, + repo: message.repo, + seq: message.seq, + time: message.time, + rev: message.rev, + since: message.since, + op: op, + }); + } + break; + } + case "com.atproto.sync.subscribeRepos#identity": + case "com.atproto.sync.subscribeRepos#account": { + addRecord(message); + break; + } + case "com.atproto.sync.subscribeRepos#sync": { + addRecord({ + $type: message.$type, + did: message.did, + rev: message.rev, + seq: message.seq, + time: message.time, + }); + break; + } } - }); - firehose.on("identity", (identity) => addRecord(identity)); - firehose.on("account", (account) => addRecord(account)); - firehose.on("sync", (sync) => { - addRecord({ - $type: sync.$type, - did: sync.did, - rev: sync.rev, - seq: sync.seq, - time: sync.time, - }); - }); - firehose.start(); + } } }; @@ -310,7 +337,7 @@ export const StreamView = () => { onCleanup(() => { socket?.close(); - firehose?.close(); + firehoseIterator?.return?.(); if (rafId !== null) cancelAnimationFrame(rafId); if (statsIntervalId !== null) clearInterval(statsIntervalId); if (statsUpdateIntervalId !== null) clearInterval(statsUpdateIntervalId);