From a7f73a85901a6738956c393318be39b54de7a211 Mon Sep 17 00:00:00 2001 From: Juliet Date: Mon, 19 Jan 2026 21:25:45 +0100 Subject: [PATCH] add stardust streaming --- src/index.tsx | 2 +- src/views/stream/config.ts | 217 ++++++++++++++++++++++ src/views/stream/index.tsx | 368 ++++++++++++++++--------------------- src/views/stream/stats.tsx | 18 +- 4 files changed, 386 insertions(+), 219 deletions(-) create mode 100644 src/views/stream/config.ts diff --git a/src/index.tsx b/src/index.tsx index 0257eeb..d64875d 100644 --- a/src/index.tsx +++ b/src/index.tsx @@ -19,7 +19,7 @@ render( () => ( - + diff --git a/src/views/stream/config.ts b/src/views/stream/config.ts new file mode 100644 index 0000000..711ccea --- /dev/null +++ b/src/views/stream/config.ts @@ -0,0 +1,217 @@ +import { localDateFromTimestamp } from "../../utils/date"; + +export type StreamType = "jetstream" | "firehose" | "spacedust"; + +export type FormField = { + name: string; + label: string; + type: "text" | "textarea" | "checkbox"; + placeholder?: string; + searchParam: string; +}; + +export type RecordInfo = { + type: string; + did?: string; + collection?: string; + rkey?: string; + action?: string; + time?: string; +}; + +export type StreamConfig = { + label: string; + icon: string; + defaultInstance: string; + fields: FormField[]; + useFirehoseLib: boolean; + buildUrl: (instance: string, formData: FormData) => string; + parseRecord: (record: any) => RecordInfo; + showEventTypes: boolean; + collectionsLabel: string; +}; + +export const STREAM_CONFIGS: Record = { + jetstream: { + label: "Jetstream", + icon: "lucide--radio-tower", + defaultInstance: "wss://jetstream1.us-east.bsky.network/subscribe", + useFirehoseLib: false, + showEventTypes: true, + collectionsLabel: "Top Collections", + fields: [ + { + name: "collections", + label: "Collections", + type: "textarea", + placeholder: "Comma-separated list of collections", + searchParam: "collections", + }, + { + name: "dids", + label: "DIDs", + type: "textarea", + placeholder: "Comma-separated list of DIDs", + searchParam: "dids", + }, + { + name: "cursor", + label: "Cursor", + type: "text", + placeholder: "Leave empty for live-tail", + searchParam: "cursor", + }, + { + name: "allEvents", + label: "Show account and identity events", + type: "checkbox", + searchParam: "allEvents", + }, + ], + buildUrl: (instance, formData) => { + let url = instance + "?"; + + const collections = formData.get("collections")?.toString().split(","); + collections?.forEach((c) => { + if (c.trim().length) url += `wantedCollections=${c.trim()}&`; + }); + + const dids = formData.get("dids")?.toString().split(","); + dids?.forEach((d) => { + if (d.trim().length) url += `wantedDids=${d.trim()}&`; + }); + + const cursor = formData.get("cursor")?.toString(); + if (cursor?.length) url += `cursor=${cursor}&`; + + return url.replace(/[&?]$/, ""); + }, + parseRecord: (rec) => { + const collection = rec.commit?.collection || rec.kind; + const rkey = rec.commit?.rkey; + const action = rec.commit?.operation; + const time = rec.time_us ? localDateFromTimestamp(rec.time_us / 1000) : undefined; + return { type: rec.kind, did: rec.did, collection, rkey, action, time }; + }, + }, + + firehose: { + label: "Firehose", + icon: "lucide--rss", + defaultInstance: "wss://bsky.network", + useFirehoseLib: true, + showEventTypes: true, + collectionsLabel: "Top Collections", + fields: [ + { + name: "cursor", + label: "Cursor", + type: "text", + placeholder: "Leave empty for live-tail", + searchParam: "cursor", + }, + ], + buildUrl: (instance, _formData) => { + let url = instance; + url = url.replace("/xrpc/com.atproto.sync.subscribeRepos", ""); + if (!(url.startsWith("wss://") || url.startsWith("ws://"))) { + url = "wss://" + url; + } + return url; + }, + parseRecord: (rec) => { + const type = rec.$type?.split("#").pop() || rec.$type; + const did = rec.repo ?? rec.did; + const pathParts = rec.op?.path?.split("/") || []; + const collection = pathParts[0]; + const rkey = pathParts[1]; + const time = rec.time ? localDateFromTimestamp(Date.parse(rec.time)) : undefined; + return { type, did, collection, rkey, action: rec.op?.action, time }; + }, + }, + + spacedust: { + label: "Spacedust", + icon: "lucide--link", + defaultInstance: "wss://spacedust.microcosm.blue/subscribe", + useFirehoseLib: false, + showEventTypes: false, + collectionsLabel: "Top Sources", + fields: [ + { + name: "sources", + label: "Sources", + type: "textarea", + placeholder: "e.g. app.bsky.graph.follow:subject", + searchParam: "sources", + }, + { + name: "subjectDids", + label: "Subject DIDs", + type: "textarea", + placeholder: "Comma-separated list of DIDs", + searchParam: "subjectDids", + }, + { + name: "subjects", + label: "Subjects", + type: "textarea", + placeholder: "Comma-separated list of AT URIs", + searchParam: "subjects", + }, + { + name: "instant", + label: "Instant mode (bypass 21s delay buffer)", + type: "checkbox", + searchParam: "instant", + }, + ], + buildUrl: (instance, formData) => { + let url = instance + "?"; + + const sources = formData.get("sources")?.toString().split(","); + sources?.forEach((s) => { + if (s.trim().length) url += `wantedSources=${s.trim()}&`; + }); + + const subjectDids = formData.get("subjectDids")?.toString().split(","); + subjectDids?.forEach((d) => { + if (d.trim().length) url += `wantedSubjectDids=${d.trim()}&`; + }); + + const subjects = formData.get("subjects")?.toString().split(","); + subjects?.forEach((s) => { + if (s.trim().length) url += `wantedSubjects=${encodeURIComponent(s.trim())}&`; + }); + + const instant = formData.get("instant")?.toString(); + if (instant === "on") url += `instant=true&`; + + return url.replace(/[&?]$/, ""); + }, + parseRecord: (rec) => { + const source = rec.link?.source; + const sourceRecord = rec.link?.source_record; + const uriParts = sourceRecord?.replace("at://", "").split("/") || []; + const did = uriParts[0]; + const collection = uriParts[1]; + const rkey = uriParts[2]; + return { + type: rec.kind, + did, + collection: source || collection, + rkey, + action: rec.link?.operation, + time: undefined, + }; + }, + }, +}; + +export const STREAM_TYPES = Object.keys(STREAM_CONFIGS) as StreamType[]; + +export const getStreamType = (pathname: string): StreamType => { + if (pathname === "/firehose") return "firehose"; + if (pathname === "/spacedust") return "spacedust"; + return "jetstream"; +}; diff --git a/src/views/stream/index.tsx b/src/views/stream/index.tsx index 180934a..36577a3 100644 --- a/src/views/stream/index.tsx +++ b/src/views/stream/index.tsx @@ -7,43 +7,28 @@ import DidHoverCard from "../../components/hover-card/did"; import { JSONValue } from "../../components/json"; import { TextInput } from "../../components/text-input"; import { addToClipboard } from "../../utils/copy"; -import { localDateFromTimestamp } from "../../utils/date"; +import { getStreamType, STREAM_CONFIGS, STREAM_TYPES, StreamType } from "./config"; import { StreamStats, StreamStatsPanel } from "./stats"; const LIMIT = 20; -type Parameter = { name: string; param: string | string[] | undefined }; -const StreamRecordItem = (props: { record: any; streamType: "jetstream" | "firehose" }) => { - const [expanded, setExpanded] = createSignal(false); - - const getBasicInfo = () => { - const rec = props.record; - if (props.streamType === "jetstream") { - const collection = rec.commit?.collection || rec.kind; - const rkey = rec.commit?.rkey; - const action = rec.commit?.operation; - const time = rec.time_us ? localDateFromTimestamp(rec.time_us / 1000) : undefined; - return { type: rec.kind, did: rec.did, collection, rkey, action, time }; - } else { - const type = rec.$type?.split("#").pop() || rec.$type; - const did = rec.repo ?? rec.did; - const pathParts = rec.op?.path?.split("/") || []; - const collection = pathParts[0]; - const rkey = pathParts[1]; - const time = rec.time ? localDateFromTimestamp(Date.parse(rec.time)) : undefined; - return { type, did, collection, rkey, action: rec.op?.action, time }; - } - }; +const TYPE_COLORS: Record = { + create: "bg-green-100 text-green-700 dark:bg-green-900/30 dark:text-green-300", + update: "bg-orange-100 text-orange-700 dark:bg-orange-900/30 dark:text-orange-300", + delete: "bg-red-100 text-red-700 dark:bg-red-900/30 dark:text-red-300", + identity: "bg-purple-100 text-purple-700 dark:bg-purple-900/30 dark:text-purple-300", + account: "bg-blue-100 text-blue-700 dark:bg-blue-900/30 dark:text-blue-300", + sync: "bg-pink-100 text-pink-700 dark:bg-pink-900/30 dark:text-pink-300", +}; - const info = getBasicInfo(); +const StreamRecordItem = (props: { record: any; streamType: StreamType }) => { + const [expanded, setExpanded] = createSignal(false); + const config = () => STREAM_CONFIGS[props.streamType]; + const info = () => config().parseRecord(props.record); - const typeColors: Record = { - create: "bg-green-100 text-green-700 dark:bg-green-900/30 dark:text-green-300", - update: "bg-orange-100 text-orange-700 dark:bg-orange-900/30 dark:text-orange-300", - delete: "bg-red-100 text-red-700 dark:bg-red-900/30 dark:text-red-300", - identity: "bg-purple-100 text-purple-700 dark:bg-purple-900/30 dark:text-purple-300", - account: "bg-blue-100 text-blue-700 dark:bg-blue-900/30 dark:text-blue-300", - sync: "bg-pink-100 text-pink-700 dark:bg-pink-900/30 dark:text-pink-300", + const displayType = () => { + const i = info(); + return i.type === "commit" || i.type === "link" ? i.action : i.type; }; const copyRecord = (e: MouseEvent) => { @@ -65,27 +50,29 @@ const StreamRecordItem = (props: { record: any; streamType: "jetstream" | "fireh : }
-
+
- {info.type === "commit" ? info.action : info.type} + {displayType()} - - {info.collection} + + + {info().collection} + - - {info.rkey} + + {info().rkey}
- + e.stopPropagation()}> - + - - {info.time} + + {info().time}
@@ -103,7 +90,7 @@ const StreamRecordItem = (props: { record: any; streamType: "jetstream" | "fireh
- +
@@ -111,14 +98,16 @@ const StreamRecordItem = (props: { record: any; streamType: "jetstream" | "fireh ); }; -const StreamView = () => { +export const StreamView = () => { const [searchParams, setSearchParams] = useSearchParams(); - const [parameters, setParameters] = createSignal([]); - const streamType = useLocation().pathname === "/firehose" ? "firehose" : "jetstream"; + const streamType = getStreamType(useLocation().pathname); + const config = () => STREAM_CONFIGS[streamType]; + const [records, setRecords] = createSignal([]); const [connected, setConnected] = createSignal(false); const [paused, setPaused] = createSignal(false); const [notice, setNotice] = createSignal(""); + const [parameters, setParameters] = createSignal<{ name: string; value?: string }[]>([]); const [stats, setStats] = createSignal({ totalEvents: 0, eventsPerSecond: 0, @@ -126,6 +115,7 @@ const StreamView = () => { collections: {}, }); const [currentTime, setCurrentTime] = createSignal(Date.now()); + let socket: WebSocket; let firehose: Firehose; let formRef!: HTMLFormElement; @@ -133,22 +123,25 @@ const StreamView = () => { let rafId: number | null = null; let statsIntervalId: number | null = null; let statsUpdateIntervalId: number | null = null; - let lastSecondEventCount = 0; let currentSecondEventCount = 0; - // Track stats in variables for batching let totalEventsCount = 0; let eventTypesMap: Record = {}; let collectionsMap: Record = {}; const addRecord = (record: any) => { currentSecondEventCount++; - - // Track statistics in variables (batched update) totalEventsCount++; - const eventType = record.kind || record.$type || "unknown"; + + const rawEventType = record.kind || record.$type || "unknown"; + const eventType = rawEventType.includes("#") ? rawEventType.split("#").pop() : rawEventType; eventTypesMap[eventType] = (eventTypesMap[eventType] || 0) + 1; + if (eventType !== "account" && eventType !== "identity") { - const collection = record.commit?.collection || record.op?.path?.split("/")[0] || "unknown"; + const collection = + record.commit?.collection || + record.op?.path?.split("/")[0] || + record.link?.source || + "unknown"; collectionsMap[collection] = (collectionsMap[collection] || 0) + 1; } @@ -165,8 +158,9 @@ const StreamView = () => { }; const disconnect = () => { - if (streamType === "jetstream") socket?.close(); + if (!config().useFirehoseLib) socket?.close(); else firehose?.close(); + if (rafId !== null) { cancelAnimationFrame(rafId); rafId = null; @@ -179,23 +173,17 @@ const StreamView = () => { clearInterval(statsUpdateIntervalId); statsUpdateIntervalId = null; } + pendingRecords = []; totalEventsCount = 0; eventTypesMap = {}; collectionsMap = {}; setConnected(false); setPaused(false); - setStats((prev) => ({ - ...prev, - eventsPerSecond: 0, - })); + setStats((prev) => ({ ...prev, eventsPerSecond: 0 })); }; - const togglePause = () => { - setPaused(!paused()); - }; - - const connectSocket = async (formData: FormData) => { + const connectStream = async (formData: FormData) => { setNotice(""); if (connected()) { disconnect(); @@ -203,54 +191,31 @@ const StreamView = () => { } setRecords([]); - let url = ""; - if (streamType === "jetstream") { - url = - formData.get("instance")?.toString() ?? "wss://jetstream1.us-east.bsky.network/subscribe"; - url = url.concat("?"); - } else { - url = formData.get("instance")?.toString() ?? "wss://bsky.network"; - url = url.replace("/xrpc/com.atproto.sync.subscribeRepos", ""); - if (!(url.startsWith("wss://") || url.startsWith("ws://"))) url = "wss://" + url; - } - - const collections = formData.get("collections")?.toString().split(","); - collections?.forEach((collection) => { - if (collection.length) url = url.concat(`wantedCollections=${collection}&`); - }); - - const dids = formData.get("dids")?.toString().split(","); - dids?.forEach((did) => { - if (did.length) url = url.concat(`wantedDids=${did}&`); - }); - - const cursor = formData.get("cursor")?.toString(); - if (streamType === "jetstream") { - if (cursor?.length) url = url.concat(`cursor=${cursor}`); - if (url.endsWith("&")) url = url.slice(0, -1); - } + const instance = formData.get("instance")?.toString() ?? config().defaultInstance; + const url = config().buildUrl(instance, formData); - setSearchParams({ - instance: formData.get("instance")?.toString(), - collections: formData.get("collections")?.toString(), - dids: formData.get("dids")?.toString(), - cursor: formData.get("cursor")?.toString(), - allEvents: formData.get("allEvents")?.toString(), + // Save all form fields to URL params + const params: Record = { instance }; + config().fields.forEach((field) => { + params[field.searchParam] = formData.get(field.name)?.toString(); }); + setSearchParams(params); + // Build parameters display setParameters([ - { name: "Instance", param: formData.get("instance")?.toString() }, - { name: "Collections", param: formData.get("collections")?.toString() }, - { name: "DIDs", param: formData.get("dids")?.toString() }, - { name: "Cursor", param: formData.get("cursor")?.toString() }, - { name: "All Events", param: formData.get("allEvents")?.toString() }, + { name: "Instance", value: instance }, + ...config() + .fields.filter((f) => f.type !== "checkbox") + .map((f) => ({ name: f.label, value: formData.get(f.name)?.toString() })), + ...config() + .fields.filter((f) => f.type === "checkbox" && formData.get(f.name) === "on") + .map((f) => ({ name: f.label, value: "on" })), ]); setConnected(true); const now = Date.now(); setCurrentTime(now); - // Reset tracking variables totalEventsCount = 0; eventTypesMap = {}; collectionsMap = {}; @@ -272,21 +237,18 @@ const StreamView = () => { })); }, 50); - // Calculate events/sec every second statsIntervalId = window.setInterval(() => { - setStats((prev) => ({ - ...prev, - eventsPerSecond: currentSecondEventCount, - })); - lastSecondEventCount = currentSecondEventCount; + setStats((prev) => ({ ...prev, eventsPerSecond: currentSecondEventCount })); currentSecondEventCount = 0; setCurrentTime(Date.now()); }, 1000); - if (streamType === "jetstream") { + + if (!config().useFirehoseLib) { socket = new WebSocket(url); socket.addEventListener("message", (event) => { const rec = JSON.parse(event.data); - if (searchParams.allEvents === "on" || (rec.kind !== "account" && rec.kind !== "identity")) + const isFilteredEvent = rec.kind === "account" || rec.kind === "identity"; + if (!isFilteredEvent || streamType !== "jetstream" || searchParams.allEvents === "on") addRecord(rec); }); socket.addEventListener("error", () => { @@ -294,6 +256,7 @@ const StreamView = () => { disconnect(); }); } else { + const cursor = formData.get("cursor")?.toString(); firehose = new Firehose({ relay: url, cursor: cursor, @@ -307,7 +270,7 @@ const StreamView = () => { }); firehose.on("commit", (commit) => { for (const op of commit.ops) { - const record = { + addRecord({ $type: commit.$type, repo: commit.repo, seq: commit.seq, @@ -315,164 +278,139 @@ const StreamView = () => { rev: commit.rev, since: commit.since, op: op, - }; - addRecord(record); + }); } }); - firehose.on("identity", (identity) => { - addRecord(identity); - }); - firehose.on("account", (account) => { - addRecord(account); - }); + firehose.on("identity", (identity) => addRecord(identity)); + firehose.on("account", (account) => addRecord(account)); firehose.on("sync", (sync) => { - const event = { + addRecord({ $type: sync.$type, did: sync.did, rev: sync.rev, seq: sync.seq, time: sync.time, - }; - addRecord(event); + }); }); firehose.start(); } }; - onMount(async () => { - const formData = new FormData(); - if (searchParams.instance) formData.append("instance", searchParams.instance.toString()); - if (searchParams.collections) - formData.append("collections", searchParams.collections.toString()); - if (searchParams.dids) formData.append("dids", searchParams.dids.toString()); - if (searchParams.cursor) formData.append("cursor", searchParams.cursor.toString()); - if (searchParams.allEvents) formData.append("allEvents", searchParams.allEvents.toString()); - if (searchParams.instance) connectSocket(formData); + onMount(() => { + if (searchParams.instance) { + const formData = new FormData(); + formData.append("instance", searchParams.instance.toString()); + config().fields.forEach((field) => { + const value = searchParams[field.searchParam]; + if (value) formData.append(field.name, value.toString()); + }); + connectStream(formData); + } }); onCleanup(() => { socket?.close(); - if (rafId !== null) { - cancelAnimationFrame(rafId); - } - if (statsIntervalId !== null) { - clearInterval(statsIntervalId); - } - if (statsUpdateIntervalId !== null) { - clearInterval(statsUpdateIntervalId); - } + firehose?.close(); + if (rafId !== null) cancelAnimationFrame(rafId); + if (statsIntervalId !== null) clearInterval(statsIntervalId); + if (statsUpdateIntervalId !== null) clearInterval(statsUpdateIntervalId); }); return ( <> - {streamType === "firehose" ? "Firehose" : "Jetstream"} - PDSls + {config().label} - PDSls
+ {/* Tab Navigation */} + + {/* Connection Form */} -
+ - -