From 7a5bde07f6c8923703a6c4d2b65356c3048ebc41 Mon Sep 17 00:00:00 2001 From: Florian <45694132+flo-bit@users.noreply.github.com> Date: Fri, 24 Apr 2026 11:18:06 +0200 Subject: [PATCH] commit --- .../src/lib/rooms/channels-context.ts | 21 ++ .../src/lib/rooms/realtime.svelte.ts | 7 +- .../src/lib/rooms/rooms.remote.ts | 3 +- .../src/lib/rooms/server.ts | 15 +- .../src/lib/rooms/watch.svelte.ts | 5 +- .../routes/c/[communityDid]/+layout.server.ts | 64 +---- .../routes/c/[communityDid]/+layout.svelte | 69 ++++- .../src/routes/c/[communityDid]/+page.svelte | 12 +- .../[communityDid]/[channelKey]/+page.svelte | 11 +- src/core/realtime/ticket.ts | 6 +- src/core/router/collection.ts | 255 +++++++++++++----- src/core/router/index.ts | 8 +- 12 files changed, 330 insertions(+), 146 deletions(-) create mode 100644 examples/sveltekit-group-chat/src/lib/rooms/channels-context.ts diff --git a/examples/sveltekit-group-chat/src/lib/rooms/channels-context.ts b/examples/sveltekit-group-chat/src/lib/rooms/channels-context.ts new file mode 100644 index 0000000..fe350f8 --- /dev/null +++ b/examples/sveltekit-group-chat/src/lib/rooms/channels-context.ts @@ -0,0 +1,21 @@ +/** Context key + type for the community channel list. + * + * The community layout (`+layout.svelte`) maintains a live channel list via + * a cross-space `createWatchQuery` on `tools.atmo.chat.channel`. Child pages + * read the list through Svelte context so they don't each run their own + * watch query. The getter ensures reads stay reactive. */ + +export interface ChannelMeta { + spaceUri: string; + key: string; + name: string; + topic?: string; + visibility: 'public' | 'private'; + createdAt: string; +} + +export interface ChannelsContext { + readonly list: readonly ChannelMeta[]; +} + +export const CHANNELS_CTX = Symbol('community-channels'); diff --git a/examples/sveltekit-group-chat/src/lib/rooms/realtime.svelte.ts b/examples/sveltekit-group-chat/src/lib/rooms/realtime.svelte.ts index 70bd04a..db74c8b 100644 --- a/examples/sveltekit-group-chat/src/lib/rooms/realtime.svelte.ts +++ b/examples/sveltekit-group-chat/src/lib/rooms/realtime.svelte.ts @@ -138,17 +138,16 @@ export function connectCommunityRealtime(communityDid: string): () => void { bumpUnread(p.space, rec.createdAt); } } - } else if (p.collection === 'tools.atmo.chat.channel') { - void invalidateAll(); } + // Channel record events are handled by the layout's channel + // watch query now — no server-loader reinvoke needed. } else if (kind === 'record.deleted') { const p = ev.payload as { uri: string; did: string; collection: string; rkey: string; space?: string }; if (!p.space) return; if (p.collection === 'tools.atmo.chat.message') { channelMessages.remove(p.space, p.rkey); - } else if (p.collection === 'tools.atmo.chat.channel') { - void invalidateAll(); } + // Channel deletions: handled by the layout's channel watch query. } else if (kind === 'member.added' || kind === 'member.removed') { void invalidateAll(); } diff --git a/examples/sveltekit-group-chat/src/lib/rooms/rooms.remote.ts b/examples/sveltekit-group-chat/src/lib/rooms/rooms.remote.ts index 2f8402e..44ae8ce 100644 --- a/examples/sveltekit-group-chat/src/lib/rooms/rooms.remote.ts +++ b/examples/sveltekit-group-chat/src/lib/rooms/rooms.remote.ts @@ -189,7 +189,8 @@ export const postMessage = command(PostMessageInput, async (input) => { // --------------------------------------------------------------------------- const MintWatchTicketInput = v.object({ - spaceUri: v.pipe(v.string(), v.minLength(1)), + spaceUri: v.optional(v.pipe(v.string(), v.minLength(1))), + actor: v.optional(v.pipe(v.string(), v.minLength(1))), watchRecordsNsid: v.pipe(v.string(), v.minLength(1)), limit: v.optional(v.pipe(v.number(), v.integer(), v.minValue(1), v.maxValue(200))) }); diff --git a/examples/sveltekit-group-chat/src/lib/rooms/server.ts b/examples/sveltekit-group-chat/src/lib/rooms/server.ts index c81e4ac..0d42c98 100644 --- a/examples/sveltekit-group-chat/src/lib/rooms/server.ts +++ b/examples/sveltekit-group-chat/src/lib/rooms/server.ts @@ -198,10 +198,21 @@ export async function getRealtimeTicket( * snapshot on the ticketed connection. */ export async function mintWatchTicket( ctx: AuthedCallContext, - input: { watchRecordsNsid: string; spaceUri: string; limit?: number } + input: { + watchRecordsNsid: string; + /** Per-space watch. Mutually exclusive with `actor`. */ + spaceUri?: string; + /** Cross-space watch — actor must currently be a community DID. */ + actor?: string; + limit?: number; + } ): Promise<{ ticket: string; expiresAt: number }> { + if (!input.spaceUri && !input.actor) { + throw new Error('mintWatchTicket: spaceUri or actor required'); + } const url = new URL(`http://localhost/xrpc/${input.watchRecordsNsid}`); - url.searchParams.set('spaceUri', input.spaceUri); + if (input.spaceUri) url.searchParams.set('spaceUri', input.spaceUri); + if (input.actor) url.searchParams.set('actor', input.actor); url.searchParams.set('mode', 'ws'); if (input.limit) url.searchParams.set('limit', String(input.limit)); const req = new Request(url, { headers: { accept: 'application/json' } }); diff --git a/examples/sveltekit-group-chat/src/lib/rooms/watch.svelte.ts b/examples/sveltekit-group-chat/src/lib/rooms/watch.svelte.ts index 34057ea..85c2ddb 100644 --- a/examples/sveltekit-group-chat/src/lib/rooms/watch.svelte.ts +++ b/examples/sveltekit-group-chat/src/lib/rooms/watch.svelte.ts @@ -91,9 +91,12 @@ export interface TicketMintContext { export type TicketMinter = (ctx: TicketMintContext) => Promise; let ticketMinter: TicketMinter | null = async ({ endpoint, params }) => { + const spaceUri = params.spaceUri ? String(params.spaceUri) : undefined; + const actor = params.actor ? String(params.actor) : undefined; const res = await mintWatchTicketCmd({ watchRecordsNsid: `${endpoint}.watchRecords`, - spaceUri: String(params.spaceUri ?? ''), + ...(spaceUri ? { spaceUri } : {}), + ...(actor ? { actor } : {}), limit: typeof params.limit === 'number' ? params.limit : 50 }); return res.ticket; diff --git a/examples/sveltekit-group-chat/src/routes/c/[communityDid]/+layout.server.ts b/examples/sveltekit-group-chat/src/routes/c/[communityDid]/+layout.server.ts index ff8f27b..e97550f 100644 --- a/examples/sveltekit-group-chat/src/routes/c/[communityDid]/+layout.server.ts +++ b/examples/sveltekit-group-chat/src/routes/c/[communityDid]/+layout.server.ts @@ -1,8 +1,7 @@ import type { LayoutServerLoad } from './$types'; import { redirect } from '@sveltejs/kit'; -import type { Client } from '@atcute/client'; import { authedFetch } from '$lib/rooms/server'; -import { parseSpaceUri, buildAdminUri } from '$lib/rooms/uri'; +import { buildAdminUri } from '$lib/rooms/uri'; interface ServerMeta { communityDid: string; @@ -12,25 +11,12 @@ interface ServerMeta { membersUri: string; } -interface ChannelMeta { - spaceUri: string; - key: string; - name: string; - topic?: string; - visibility: 'public' | 'private'; - createdAt: string; -} - export const load: LayoutServerLoad = async ({ locals, params, platform }) => { - if (!locals.did || !locals.client) { + if (!locals.did) { throw redirect(302, '/'); } const communityDid = decodeURIComponent(params.communityDid); - const ctx = { - env: platform!.env, - client: locals.client as Client, - did: locals.did as string - }; + const ctx = { env: platform!.env, did: locals.did as string }; // --- fetch server record ------------------------------------------------- let server: ServerMeta | null = null; @@ -73,44 +59,10 @@ export const load: LayoutServerLoad = async ({ locals, params, platform }) => { // fallthrough — server stays null, UI shows "server" fallback } - // --- fetch channels ------------------------------------------------------ - const channels: ChannelMeta[] = []; - try { - const data = await authedFetch<{ - records: Array<{ - did: string; - rkey: string; - record: { - communityDid?: string; - name?: string; - topic?: string; - visibility?: 'public' | 'private'; - createdAt?: string; - }; - space?: string; - }>; - }>(ctx, 'tools.atmo.chat.channel.listRecords', { - query: { actor: communityDid, limit: '100' } - }); - for (const r of data.records) { - if (!r.space || r.did !== communityDid || r.rkey !== 'self') continue; - if (r.record?.communityDid !== communityDid) continue; - if (!r.record.name || !r.record.visibility || !r.record.createdAt) continue; - const parsed = parseSpaceUri(r.space); - if (!parsed || parsed.key.startsWith('$') || parsed.key === 'members') continue; - channels.push({ - spaceUri: r.space, - key: parsed.key, - name: r.record.name, - topic: r.record.topic, - visibility: r.record.visibility, - createdAt: r.record.createdAt - }); - } - channels.sort((a, b) => (a.createdAt < b.createdAt ? -1 : 1)); - } catch { - // fallthrough - } + // Channels are no longer fetched here — the layout now derives them from + // a live `createWatchQuery` against `tools.atmo.chat.channel` scoped by + // actor=. New/renamed/deleted channels reflect instantly + // without an `invalidateAll()` → server-loader roundtrip. // --- caller's access level on $admin ----------------------------------- let isAdmin = false; @@ -126,5 +78,5 @@ export const load: LayoutServerLoad = async ({ locals, params, platform }) => { // stay false } - return { communityDid, server, channels, isAdmin, myDid: locals.did }; + return { communityDid, server, isAdmin, myDid: locals.did }; }; diff --git a/examples/sveltekit-group-chat/src/routes/c/[communityDid]/+layout.svelte b/examples/sveltekit-group-chat/src/routes/c/[communityDid]/+layout.svelte index 27980bc..7bdbf64 100644 --- a/examples/sveltekit-group-chat/src/routes/c/[communityDid]/+layout.svelte +++ b/examples/sveltekit-group-chat/src/routes/c/[communityDid]/+layout.svelte @@ -1,13 +1,17 @@
- {#if data.channels.length === 0} + {#if channelsCtx.list.length === 0}

No channels yet.

{#if data.isAdmin} diff --git a/examples/sveltekit-group-chat/src/routes/c/[communityDid]/[channelKey]/+page.svelte b/examples/sveltekit-group-chat/src/routes/c/[communityDid]/[channelKey]/+page.svelte index 447da11..13da1b9 100644 --- a/examples/sveltekit-group-chat/src/routes/c/[communityDid]/[channelKey]/+page.svelte +++ b/examples/sveltekit-group-chat/src/routes/c/[communityDid]/[channelKey]/+page.svelte @@ -6,19 +6,20 @@ import { setCurrentChannel } from '$lib/rooms/realtime.svelte'; import { markLastRead } from '$lib/rooms/unread.svelte'; import { displayName } from '$lib/rooms/profiles.svelte'; + import { getContext } from 'svelte'; import { createWatchQuery } from '$lib/rooms/watch.svelte'; import { setConnectionStatus, resetConnectionStatus } from '$lib/rooms/connection.svelte'; import { nextTid } from '@atmo-dev/contrail'; + import { CHANNELS_CTX, type ChannelsContext } from '$lib/rooms/channels-context'; let { data } = $props(); - let channel = $derived(data.channels.find((c) => c.key === data.channelKey)); + const channelsCtx = getContext(CHANNELS_CTX); + + let channel = $derived(channelsCtx.list.find((c) => c.key === data.channelKey)); let channelName = $derived(channel?.name ?? data.channelKey); - // Live message feed. `$derived` recreates the query when spaceUri changes; - // the old instance is auto-torn down via `createSubscriber` once no - // component reads it. Initial ticket is pre-minted during SSR to save one - // roundtrip; reconnects call the configured ticket minter transparently. + // Live message feed. `$derived` recreates the query when spaceUri changes let messagesQuery = $derived( createWatchQuery({ endpoint: 'tools.atmo.chat.message', diff --git a/src/core/realtime/ticket.ts b/src/core/realtime/ticket.ts index 7e0fd6a..1c6e3de 100644 --- a/src/core/realtime/ticket.ts +++ b/src/core/realtime/ticket.ts @@ -26,7 +26,11 @@ export interface TicketPayload { export interface TicketQuerySpec { collection: string; - spaceUri: string; + /** Exactly one of `spaceUri` or `actor` is set. `spaceUri` = per-space + * watch; `actor` = cross-space watch for records authored by this DID + * (the ticket's `topics` list carries the expanded delivery topics). */ + spaceUri?: string; + actor?: string; hydrate?: Record; } diff --git a/src/core/router/collection.ts b/src/core/router/collection.ts index 1918aa8..d54f6d8 100644 --- a/src/core/router/collection.ts +++ b/src/core/router/collection.ts @@ -23,19 +23,41 @@ import type { SpacesContext } from "."; import type { Nsid } from "@atcute/lexicons"; import type { RealtimeEvent } from "../realtime/types"; import { sseResponse } from "../realtime/sse"; -import { spaceTopic } from "../realtime/types"; +import { spaceTopic, communityTopic, parseSpaceTopic } from "../realtime/types"; import type { SubscriberQuerySpec } from "../realtime/durable-object"; import { DurableObjectPubSub } from "../realtime/durable-object"; import { TicketSigner, type TicketQuerySpec } from "../realtime/ticket"; +import { resolveTopicForCaller } from "../realtime/resolve"; +import { mergeAsyncIterables } from "../realtime/merge"; +import type { CommunityAdapter } from "../community/adapter"; import { getRelationField, getNestedValue } from "../types"; +/** Scope of a watch stream. + * - `space`: single permissioned space — one `space:` topic. + * - `actor`: records authored by `actor` across multiple spaces — the + * resolver expanded these to a per-caller subset of space topics (plus + * `actor:` for public records). Events outside `allowedSpaces` + * are filtered out. */ +type WatchScope = + | { kind: "space"; spaceUri: string } + | { + kind: "actor"; + actor: string; + /** Concrete pubsub topics to subscribe to (from resolveTopicForCaller). */ + topics: string[]; + /** Space URIs the caller can see. Events with `space` outside this + * set are dropped. Undefined `space` on an event (public record) + * is allowed only when `actor` topic is in `topics`. */ + allowedSpaces: Set; + }; + /** Shared implementation of the watchRecords snapshot+live loop. Called by * both transport branches (SSE and Worker-terminated WS). The caller owns * the actual socket/stream and provides a `send(kind, data)` closure. */ async function runQueryStream(opts: { send: (kind: string, data: unknown) => void; abort: AbortController; - spaceUri: string; + scope: WatchScope; callerDid: string | undefined; params: URLSearchParams; db: Database; @@ -53,7 +75,7 @@ async function runQueryStream(opts: { const { send, abort, - spaceUri, + scope, callerDid, params, db, @@ -66,6 +88,13 @@ async function runQueryStream(opts: { childCollectionMap } = opts; + // Predicate: does this event belong in the caller's scope? + const inScope = (space: string | undefined): boolean => { + if (scope.kind === "space") return space === scope.spaceUri; + if (space == null) return false; // actor mode: require space for now (app topic) + return scope.allowedSpaces.has(space); + }; + const hydrateSpec = parseHydrateParams(params, relations, references); const trackHydration = Object.keys(hydrateSpec.relations).length > 0; const parentUris = new Set(); @@ -80,7 +109,11 @@ async function runQueryStream(opts: { const meta = childCollectionMap.get(event.payload.collection); if (!meta) return; if (!(hydrateSpec.relations as Record)[meta.relName]) return; - if (event.payload.space !== spaceUri) return; + if (!inScope(event.payload.space)) return; + // Actor mode: additionally require the record's author match our actor + // (the caller might share spaces with other authors — we only surface + // records by the actor under watch). + if (scope.kind === "actor" && event.payload.did !== scope.actor) return; if (event.kind === "record.created") { const matched = getNestedValue(event.payload.record, meta.matchField); @@ -108,7 +141,7 @@ async function runQueryStream(opts: { collection: event.payload.collection, cid: event.payload.cid, record: event.payload.record, - _space: spaceUri + _space: event.payload.space } }); } else { @@ -132,7 +165,8 @@ async function runQueryStream(opts: { return; } if (event.kind !== "record.created" && event.kind !== "record.deleted") return; - if (event.payload.space !== spaceUri) return; + if (!inScope(event.payload.space)) return; + if (scope.kind === "actor" && event.payload.did !== scope.actor) return; if (event.payload.collection !== colNsid) { handleChildEvent(event); @@ -154,7 +188,7 @@ async function runQueryStream(opts: { record: event.payload.record, time_us: nowUs, indexed_at: event.ts, - _space: spaceUri + _space: event.payload.space } }); } else { @@ -167,8 +201,15 @@ async function runQueryStream(opts: { } }; - const topic = spaceTopic(spaceUri); - const iter = pubsub.subscribe(topic, abort.signal); + // Subscribe: one topic for space-scoped, merge across all topics for + // actor-scoped. `mergeAsyncIterables` exists for exactly this case. + let iter: AsyncIterable; + if (scope.kind === "space") { + iter = pubsub.subscribe(spaceTopic(scope.spaceUri), abort.signal); + } else { + const sources = scope.topics.map((t) => pubsub.subscribe(t, abort.signal)); + iter = mergeAsyncIterables(sources, abort.signal); + } const buffered: RealtimeEvent[] = []; let snapshotDone = false; @@ -186,8 +227,15 @@ async function runQueryStream(opts: { })(); try { - send("snapshot.start", { spaceUri, collection: colNsid }); - const result = await runPipeline(db, config, collection, params, undefined, [spaceUri]); + send( + "snapshot.start", + scope.kind === "space" + ? { spaceUri: scope.spaceUri, collection: colNsid } + : { actor: scope.actor, collection: colNsid } + ); + const snapshotSpaces = + scope.kind === "space" ? [scope.spaceUri] : Array.from(scope.allowedSpaces); + const result = await runPipeline(db, config, collection, params, undefined, snapshotSpaces); for (const record of result.records) { if (abort.signal.aborted) break; if (typeof record.uri === "string") parentUris.add(record.uri); @@ -209,7 +257,15 @@ async function runQueryStream(opts: { } } } - send("snapshot.record", { record }); + // Normalize `space` → `_space` so snapshot records carry the same + // field name as live `record.created` payloads. Clients then have a + // single key to read regardless of origin. + const normalized = record as Record; + if (normalized.space !== undefined && normalized._space === undefined) { + normalized._space = normalized.space; + delete normalized.space; + } + send("snapshot.record", { record: normalized }); } send("snapshot.end", { cursor: result.cursor }); snapshotDone = true; @@ -397,10 +453,14 @@ export function registerCollectionRoutes( db: Database, config: ContrailConfig, spacesCtx?: SpacesContext | null, - options: { pubsub?: import("../realtime/types").PubSub | null } = {} + options: { + pubsub?: import("../realtime/types").PubSub | null; + community?: CommunityAdapter | null; + } = {} ): void { const ns = config.namespace; const pubsub = options.pubsub ?? null; + const community = options.community ?? null; /** When a per-collection endpoint receives `?spaceUri=...`, verify the JWT, * resolve membership, run the space ACL, and return the caller DID if allowed. @@ -561,67 +621,128 @@ export function registerCollectionRoutes( app.get(`/xrpc/${ns}.${collection}.watchRecords`, async (c) => { const params = new URL(c.req.url).searchParams; const spaceUri = params.get("spaceUri"); - if (!spaceUri) { + const actorParam = params.get("actor"); + + if (!spaceUri && !actorParam) { return c.json( - { error: "InvalidRequest", message: "spaceUri required (cross-space watch is deferred)" }, + { error: "InvalidRequest", message: "spaceUri or actor required" }, 400 ); } - // Try ticket-auth first: if a valid watch ticket scoped to this - // spaceUri/collection is present, use it and skip the JWT gate. + // Resolve the caller and their scope. Two parallel paths: + // - space-scoped: single `space:` topic, per-space ACL gate. + // - actor-scoped: caller's reachable spaces in the actor's + // community (v1 only supports community DIDs as the actor). + // Events are delivered via N `space:` topics and filtered + // to `did === actor`. let callerDid: string | undefined; - let querySpec: SubscriberQuerySpec; - let ticketSpec: SubscriberQuerySpec | null = null; + let scope: WatchScope; + let scopeTopics: string[]; // for ticket signing + let ticketSpec: TicketQuerySpec | null = null; const providedTicket = params.get("ticket"); if (providedTicket && ticketSigner) { const payload = await ticketSigner.verify(providedTicket); - if (payload?.querySpec) { + if (payload?.querySpec && payload.querySpec.collection === colNsid) { const ts = payload.querySpec; - if ( - ts.collection === colNsid && - ts.spaceUri === spaceUri && - payload.topics.includes(spaceTopic(spaceUri)) - ) { + if (spaceUri && ts.spaceUri === spaceUri) { + if (payload.topics.includes(spaceTopic(spaceUri))) { + callerDid = payload.did; + ticketSpec = { + collection: ts.collection, + spaceUri: ts.spaceUri, + ...(ts.hydrate ? { hydrate: ts.hydrate } : {}) + }; + } + } else if (actorParam && ts.actor === actorParam) { callerDid = payload.did; ticketSpec = { collection: ts.collection, - spaceUri: ts.spaceUri, + actor: ts.actor, ...(ts.hydrate ? { hydrate: ts.hydrate } : {}) }; } } } - if (!ticketSpec) { - const gated = await gateSpaceAccess(c, spaceUri, "read"); - if (gated instanceof Response) return gated; - callerDid = "callerDid" in gated ? gated.callerDid : undefined; + const hydrateSpec = parseHydrateParams(params, relations, references); + const hydrateForSpec = Object.keys(hydrateSpec.relations).length > 0 + ? Object.fromEntries( + Object.entries(hydrateSpec.relations).map(([relName]) => { + const rel = relations[relName]!; + const childNsid = + nsidForShortName(config, rel.collection) ?? rel.collection; + return [ + relName, + { childCollection: childNsid, matchField: getRelationField(rel) } + ]; + }) + ) + : undefined; + + if (spaceUri) { + if (!ticketSpec) { + const gated = await gateSpaceAccess(c, spaceUri, "read"); + if (gated instanceof Response) return gated; + callerDid = "callerDid" in gated ? gated.callerDid : undefined; + } + scope = { kind: "space", spaceUri }; + scopeTopics = [spaceTopic(spaceUri)]; + } else { + // Actor-scoped path — v1 only supports community DIDs. + const actor = actorParam!; + if (!community || !spacesCtx) { + return c.json( + { error: "NotSupported", reason: "community-module-disabled" }, + 400 + ); + } + const isCommunity = !!(await community.getCommunity(actor)); + if (!isCommunity) { + return c.json( + { error: "InvalidRequest", reason: "actor-must-be-community-did", message: "cross-space watch currently only supports community DIDs as actor" }, + 400 + ); + } + + if (!ticketSpec) { + // Verify the caller via the same JWT/in-process path used for + // per-space queries, then resolve the community topic to the + // caller's accessible space topics. + const nsidLxm = new URL(c.req.url).pathname.match(/\/xrpc\/([^?]+)/)?.[1] as Nsid | null; + const auth = await verifyServiceAuthRequest(spacesCtx.verifier, c.req.raw, nsidLxm); + if (!auth) { + return c.json( + { error: "AuthRequired", message: "service-auth JWT or in-process principal required" }, + 401 + ); + } + callerDid = auth.issuer; + } + const resolved = await resolveTopicForCaller(communityTopic(actor), callerDid!, { + spaces: spacesCtx.adapter, + community + }); + if (!resolved.ok) { + const status = + resolved.error === "NotFound" ? 404 : + resolved.error === "Forbidden" ? 403 : 400; + return c.json({ error: resolved.error, reason: resolved.reason }, status); + } + const allowedSpaces = new Set(); + for (const t of resolved.topics) { + const uri = parseSpaceTopic(t); + if (uri) allowedSpaces.add(uri); + } + scope = { kind: "actor", actor, topics: resolved.topics, allowedSpaces }; + scopeTopics = resolved.topics; } - // Build the query spec the DO will filter events against. Prefer - // the ticket's spec when present (guarantees parity with what the - // client asked for at handshake time, no param drift). - const hydrateSpec = parseHydrateParams(params, relations, references); - querySpec = ticketSpec ?? { + const querySpec: TicketQuerySpec = ticketSpec ?? { collection: colNsid, - spaceUri, - ...(Object.keys(hydrateSpec.relations).length > 0 - ? { - hydrate: Object.fromEntries( - Object.entries(hydrateSpec.relations).map(([relName]) => { - const rel = relations[relName]!; - const childNsid = - nsidForShortName(config, rel.collection) ?? rel.collection; - return [ - relName, - { childCollection: childNsid, matchField: getRelationField(rel) } - ]; - }) - ) - } - : {}) + ...(spaceUri ? { spaceUri } : { actor: actorParam! }), + ...(hydrateForSpec ? { hydrate: hydrateForSpec } : {}) }; // Upgrade-to-WS path — forward directly to the DO with the spec, @@ -634,29 +755,27 @@ export function registerCollectionRoutes( if (isWsMode && !isUpgrade) { // Handshake: return snapshot + a ticket the client uses to - // upgrade. Ticket carries the (did, topic, querySpec) signed + // upgrade. Ticket carries the (did, topics, querySpec) signed // so the WS-upgrade route skips any other auth. try { - // Capture a server-side timestamp BEFORE running the snapshot. - // Any event published after this moment will have ts > sinceTs - // and be replayed by the DO on WS connect — so the client - // never misses events during the snapshot→WS gap. const sinceTs = Date.now(); + const snapshotSpaces = + scope.kind === "space" ? [scope.spaceUri] : Array.from(scope.allowedSpaces); const result = await runPipeline( db, config, collection, params, undefined, - [spaceUri] + snapshotSpaces ); let ticket: string | undefined; if (ticketSigner && callerDid) { ticket = await ticketSigner.sign({ - topics: [spaceTopic(spaceUri)], + topics: scopeTopics, did: callerDid, ttlMs: ticketTtl, - querySpec: querySpec as TicketQuerySpec + querySpec }); } const wsUrl = (() => { @@ -683,15 +802,23 @@ export function registerCollectionRoutes( } } - if (isUpgrade && pubsub instanceof DurableObjectPubSub) { + if (isUpgrade && pubsub instanceof DurableObjectPubSub && scope.kind === "space") { // Forward the WS upgrade to the DO. The DO owns the socket from // here and hibernates when idle. Replays any events buffered // since the handshake `sinceTs` so the client closes the gap. + // + // Actor-scoped queries fall through to the worker-terminated + // path below — the DO binding is single-topic today; extending + // it to fan out over N topics is future work. const sinceTsParam = params.get("sinceTs"); const sinceTs = sinceTsParam ? Number(sinceTsParam) : 0; - return pubsub.forwardSubscribe(spaceTopic(spaceUri), c.req.raw, { + return pubsub.forwardSubscribe(spaceTopic(scope.spaceUri), c.req.raw, { did: callerDid, - querySpec, + querySpec: { + collection: querySpec.collection, + spaceUri: scope.spaceUri, + ...(querySpec.hydrate ? { hydrate: querySpec.hydrate } : {}) + }, sinceTs: Number.isFinite(sinceTs) ? sinceTs : 0 }); } @@ -733,7 +860,7 @@ export function registerCollectionRoutes( void runQueryStream({ send: sendWs, abort: ac, - spaceUri, + scope, callerDid, params, db, @@ -802,7 +929,7 @@ export function registerCollectionRoutes( void runQueryStream({ send, abort: ac, - spaceUri, + scope, callerDid, params, db, diff --git a/src/core/router/index.ts b/src/core/router/index.ts index 8a1b732..bf714fd 100644 --- a/src/core/router/index.ts +++ b/src/core/router/index.ts @@ -114,7 +114,13 @@ export function createApp( } registerAdminRoutes(app, db, config); - registerCollectionRoutes(app, db, config, spacesCtx, { pubsub: realtimePubsub }); + const communityAdapterForCollection = config.community + ? new CommunityAdapter(spacesDb) + : null; + registerCollectionRoutes(app, db, config, spacesCtx, { + pubsub: realtimePubsub, + community: communityAdapterForCollection, + }); registerFeedRoutes(app, db, config); registerNotifyRoute(app, db, config); const communityAdapterForSpaces = config.community && spacesCtx -- 2.51.2