diff --git a/.zed/settings.json b/.zed/settings.json new file mode 100644 index 0000000..1ac3de3 --- /dev/null +++ b/.zed/settings.json @@ -0,0 +1,13 @@ +{ + "lsp": { + "vtsls": { + "settings": { + "typescript": { + "suggestionActions": { + "enabled": false + } + } + } + } + } +} diff --git a/src/components/FollowingItem.svelte b/src/components/FollowingItem.svelte index 4f93032..237a228 100644 --- a/src/components/FollowingItem.svelte +++ b/src/components/FollowingItem.svelte @@ -1,5 +1,4 @@ @@ -8,11 +7,11 @@ import { getRelativeTime } from '$lib/date'; import { generateColorForDid } from '$lib/accounts'; import type { Did } from '@atcute/lexicons'; - import type { AtprotoDid } from '@atcute/lexicons/syntax'; import type { calculateFollowedUserStats, Sort } from '$lib/following'; - import type { AtpClient } from '$lib/at/client'; + import { resolveDidDoc, type AtpClient } from '$lib/at/client'; import { SvelteMap } from 'svelte/reactivity'; - import { clients, getClient, router } from '$lib/state.svelte'; + import { router } from '$lib/state.svelte'; + import { map } from '$lib/result'; interface Props { style: string; @@ -35,48 +34,20 @@ const c = profileCache.get(targetDid)!; displayName = c.displayName; handle = c.handle; - } else { - const existingClient = clients.get(targetDid as AtprotoDid); - if (existingClient?.user?.handle) { - handle = existingClient.user.handle; - } else { - handle = 'handle.invalid'; - displayName = undefined; - } } try { - // Optimization: Check clients map first to avoid async overhead if possible - // but we need to ensure we have the profile data, not just client existence. - const userClient = await getClient(targetDid as AtprotoDid); - - // Check if the component has been recycled for a different user while we were awaiting + const [profileRes, handleRes] = await Promise.all([ + client.getProfile(), + resolveDidDoc(targetDid).then((r) => map(r, (doc) => doc.handle)) + ]); if (did !== targetDid) return; + if (profileRes.ok) displayName = profileRes.value.displayName; + if (handleRes.ok) handle = handleRes.value; - let newHandle = handle; - let newDisplayName = displayName; - - if (userClient.user?.handle) { - newHandle = userClient.user.handle; - handle = newHandle; - } else { - newHandle = targetDid; - handle = newHandle; - } - - const profileRes = await userClient.getProfile(); - - if (did !== targetDid) return; - - if (profileRes.ok) { - newDisplayName = profileRes.value.displayName; - displayName = newDisplayName; - } - - // Update cache profileCache.set(targetDid, { - handle: newHandle, - displayName: newDisplayName + handle, + displayName }); } catch (e) { if (did !== targetDid) return; @@ -85,7 +56,6 @@ } }; - // Re-run whenever `did` changes $effect(() => { loadProfile(did); }); diff --git a/src/components/TimelineView.svelte b/src/components/TimelineView.svelte index bc62583..1e19413 100644 --- a/src/components/TimelineView.svelte +++ b/src/components/TimelineView.svelte @@ -11,7 +11,7 @@ fetchTimeline, allPosts, timelines, - fetchInteractionsUntil + fetchInteractionsToTimelineEnd } from '$lib/state.svelte'; import Icon from '@iconify/svelte'; import { buildThreads, filterThreads, type ThreadPost } from '$lib/thread'; @@ -53,7 +53,6 @@ const loaderState = new LoaderState(); let scrollContainer = $state(); let loading = $state(false); - let fetchMoreInteractions: boolean | undefined = $state(false); let loadError = $state(''); const loadMore = async () => { @@ -63,9 +62,9 @@ loaderState.status = 'LOADING'; try { - await fetchTimeline(did as AtprotoDid, 7, showReplies); - // interaction fetching is done lazily so we dont block loading posts - fetchMoreInteractions = true; + await fetchTimeline(client, did as AtprotoDid, 7, showReplies); + // only fetch interactions if logged in (because if not who is the interactor) + if (client.user) await fetchInteractionsToTimelineEnd(client, did); loaderState.loaded(); } catch (error) { loadError = `${error}`; @@ -87,11 +86,12 @@ const cursor = did ? postCursors.get(did as AtprotoDid) : undefined; if (!cursor?.end) loadMore(); } - if (client && did && fetchMoreInteractions) { - // set to false so it doesnt attempt to fetch again while its already fetching - fetchMoreInteractions = false; - fetchInteractionsUntil(client, did).then(() => (fetchMoreInteractions = undefined)); - } + }); + // we want to load interactions when changing logged in user on timelines + // only on timelines that arent logged in users, because those are already + // loaded by loadMore + $effect(() => { + if (client && did && client.user?.did !== did) fetchInteractionsToTimelineEnd(client, did); }); diff --git a/src/lib/at/client.ts b/src/lib/at/client.ts index 0612f50..0cbcfd6 100644 --- a/src/lib/at/client.ts +++ b/src/lib/at/client.ts @@ -8,7 +8,6 @@ import { Client as AtcuteClient, simpleFetchHandler } from '@atcute/client'; import { safeParse, type Blob as AtpBlob, type Handle, type InferOutput } from '@atcute/lexicons'; import { isDid, - parseCanonicalResourceUri, parseResourceUri, type ActorIdentifier, type AtprotoDid, @@ -36,7 +35,7 @@ import { WebSocket } from '@soffinal/websocket'; import type { Notification } from './stardust'; import type { OAuthUserAgent } from '@atcute/oauth-browser-client'; import { timestampFromCursor, toCanonicalUri, toResourceUri } from '$lib'; -import { constellationUrl, slingshotUrl, spacedustUrl } from '.'; +import { constellationUrl, httpToDidWeb, slingshotUrl, spacedustUrl } from '.'; export type RecordOutput = { uri: ResourceUri; cid: Cid | undefined; record: Output }; @@ -90,32 +89,19 @@ export const xhrPost = ( if (onProgress && xhr.upload) { xhr.upload.onprogress = (event: ProgressEvent) => { - if (event.lengthComputable) { - onProgress(event.loaded, event.total); - } + if (event.lengthComputable) onProgress(event.loaded, event.total); }; } - Object.keys(headers).forEach((key) => { - xhr.setRequestHeader(key, headers[key]); - }); + Object.keys(headers).forEach((key) => xhr.setRequestHeader(key, headers[key])); xhr.onload = () => { - if (xhr.status >= 200 && xhr.status < 300) { - resolve(ok(JSON.parse(xhr.responseText))); - } else { - resolve(err(JSON.parse(xhr.responseText))); - } - }; - - xhr.onerror = () => { - resolve(err({ error: 'xhr_error', message: 'network error' })); - }; - - xhr.onabort = () => { - resolve(err({ error: 'xhr_error', message: 'upload aborted' })); + if (xhr.status >= 200 && xhr.status < 300) resolve(ok(JSON.parse(xhr.responseText))); + else resolve(err(JSON.parse(xhr.responseText))); }; + xhr.onerror = () => resolve(err({ error: 'xhr_error', message: 'network error' })); + xhr.onabort = () => resolve(err({ error: 'xhr_error', message: 'upload aborted' })); xhr.send(body); }); }; @@ -202,16 +188,23 @@ export class AtpClient { } async listRecords( + ident: ActorIdentifier, collection: Collection, cursor?: string, limit: number = 100 ): Promise< Result, string> > { - if (!this.atcute || !this.user) return err('not authenticated'); - const res = await this.atcute.get('com.atproto.repo.listRecords', { + if (!this.atcute) return err('not authenticated'); + const docRes = await resolveDidDoc(ident); + if (!docRes.ok) return docRes; + const atp = + this.user?.did === docRes.value.did + ? this.atcute + : new AtcuteClient({ handler: simpleFetchHandler({ service: docRes.value.pds }) }); + const res = await atp.get('com.atproto.repo.listRecords', { params: { - repo: this.user.did, + repo: docRes.value.did, collection, cursor, limit, @@ -227,6 +220,7 @@ export class AtpClient { } async listRecordsUntil( + ident: ActorIdentifier, collection: Collection, cursor?: string, timestamp: number = -1 @@ -238,7 +232,7 @@ export class AtpClient { let end = false; while (!end) { - const res = await this.listRecords(collection, data.cursor); + const res = await this.listRecords(ident, collection, data.cursor); if (!res.ok) return res; data.cursor = res.value.cursor; data.records.push(...res.value.records); @@ -255,7 +249,7 @@ export class AtpClient { end = true; } else { console.info( - `${this.user?.did}: continuing to fetch ${collection}, on ${cursorTimestamp} until ${timestamp}` + `${ident}: continuing to fetch ${collection}, on ${cursorTimestamp} until ${timestamp}` ); } } @@ -264,34 +258,22 @@ export class AtpClient { return ok(data); } - async getBacklinksUri( - uri: ResourceUri, - source: BacklinksSource - ): Promise> { - const parsedResourceUri = expect(parseCanonicalResourceUri(uri)); - return await this.getBacklinks( - parsedResourceUri.repo, - parsedResourceUri.collection, - parsedResourceUri.rkey, - source - ); - } - async getBacklinks( - repo: ActorIdentifier, - collection: Nsid, - rkey: RecordKey, + subject: ResourceUri, source: BacklinksSource, + filterBy?: Did[], limit?: number ): Promise> { + const { repo, collection, rkey } = expect(parseResourceUri(subject)); const did = await resolveHandle(repo); if (!did.ok) return err(`cant resolve handle: ${did.error}`); const timeout = new Promise((resolve) => setTimeout(() => resolve(null), 2000)); const query = fetchMicrocosm(constellationUrl, BacklinksQuery, { - subject: toCanonicalUri({ did: did.value, collection, rkey }), + subject: collection ? toCanonicalUri({ did: did.value, collection, rkey: rkey! }) : did.value, source, - limit: limit || 100 + limit: limit || 100, + did: filterBy }); const results = await Promise.race([query, timeout]); @@ -303,10 +285,7 @@ export class AtpClient { async getServiceAuth(lxm: keyof XRPCProcedures, exp: number): Promise> { if (!this.atcute || !this.user) return err('not authenticated'); const serviceAuthUrl = new URL(`${this.user.pds}xrpc/com.atproto.server.getServiceAuth`); - serviceAuthUrl.searchParams.append( - 'aud', - this.user.pds.replace('https://', 'did:web:').slice(0, -1) - ); + serviceAuthUrl.searchParams.append('aud', httpToDidWeb(this.user.pds)); serviceAuthUrl.searchParams.append('lxm', 'com.atproto.repo.uploadBlob'); serviceAuthUrl.searchParams.append('exp', exp.toString()); // 30 minutes @@ -412,41 +391,32 @@ export class AtpClient { } } -export const newPublicClient = async (ident: ActorIdentifier): Promise => { - const atp = new AtpClient(); - const didDoc = await resolveDidDoc(ident); - if (!didDoc.ok) { - console.error('failed to resolve did doc', didDoc.error); - return atp; - } - atp.atcute = new AtcuteClient({ handler: simpleFetchHandler({ service: didDoc.value.pds }) }); - atp.user = didDoc.value; - return atp; -}; - -// Wrappers that use the cache - -export const resolveHandle = async ( - identifier: ActorIdentifier -): Promise> => { - if (isDid(identifier)) return ok(identifier as AtprotoDid); - - try { - const did = await cache.resolveHandle(identifier); - return ok(did); - } catch (e) { - return err(String(e)); - } +// export const newPublicClient = async (ident: ActorIdentifier) => { +// const atp = new AtpClient(); +// const didDoc = await resolveDidDoc(ident); +// if (!didDoc.ok) { +// console.error('failed to resolve did doc', didDoc.error); +// return atp; +// } +// atp.atcute = new AtcuteClient({ handler: simpleFetchHandler({ service: didDoc.value.pds }) }); +// atp.user = didDoc.value; +// return atp; +// }; + +export const resolveHandle = (identifier: ActorIdentifier) => { + if (isDid(identifier)) return Promise.resolve(ok(identifier as AtprotoDid)); + + return cache + .resolveHandle(identifier) + .then((did) => ok(did)) + .catch((e) => err(String(e))); }; -export const resolveDidDoc = async (ident: ActorIdentifier): Promise> => { - try { - const doc = await cache.resolveDidDoc(ident); - return ok(doc); - } catch (e) { - return err(String(e)); - } -}; +export const resolveDidDoc = (ident: ActorIdentifier) => + cache + .resolveDidDoc(ident) + .then((doc) => ok(doc)) + .catch((e) => err(String(e))); type NotificationsStreamEncoder = WebSocket.Encoder; export type NotificationsStream = WebSocket; @@ -485,7 +455,9 @@ const fetchMicrocosm = async < ): Promise> => { if (!schema.output || schema.output.type === 'blob') return err('schema must be blob'); api.pathname = `/xrpc/${schema.nsid}`; - api.search = params ? `?${new URLSearchParams(params)}` : ''; + api.search = params + ? `?${new URLSearchParams(Object.entries(params).flatMap(([k, v]) => (v === undefined ? [] : [[k, String(v)]])))}` + : ''; try { const body = await fetchJson(api, init); if (!body.ok) return err(body.error); diff --git a/src/lib/at/constellation.ts b/src/lib/at/constellation.ts index 9291aac..bf1d69c 100644 --- a/src/lib/at/constellation.ts +++ b/src/lib/at/constellation.ts @@ -9,7 +9,7 @@ export const BacklinkSchema = v.object({ }); export const BacklinksQuery = v.query('blue.microcosm.links.getBacklinks', { params: v.object({ - subject: v.resourceUriString(), + subject: v.string(), source: v.string(), did: v.optional(v.array(v.didString())), limit: v.optional(v.integer()) diff --git a/src/lib/at/fetch.ts b/src/lib/at/fetch.ts index 0285d36..eab11d0 100644 --- a/src/lib/at/fetch.ts +++ b/src/lib/at/fetch.ts @@ -17,12 +17,13 @@ export type PostWithBacklinks = PostWithUri & { }; export const fetchPosts = async ( + subject: Did, client: AtpClient, cursor?: string, limit?: number, withBacklinks: boolean = true ): Promise> => { - const recordsList = await client.listRecords('app.bsky.feed.post', cursor, limit); + const recordsList = await client.listRecords(subject, 'app.bsky.feed.post', cursor, limit); if (!recordsList.ok) return err(`can't retrieve posts: ${recordsList.error}`); cursor = recordsList.value.cursor; const records = recordsList.value.records; @@ -41,7 +42,7 @@ export const fetchPosts = async ( try { const allBacklinks = await Promise.all( records.map(async (r): Promise => { - const result = await client.getBacklinksUri(r.uri, replySource); + const result = await client.getBacklinks(r.uri, replySource); if (!result.ok) throw `cant fetch replies: ${result.error}`; const replies = result.value; return { @@ -120,7 +121,7 @@ export const hydratePosts = async ( if (repo === postRepo) return; // get chains that are the same author until we exhaust them - const backlinks = await client.getBacklinksUri(post.uri, replySource); + const backlinks = await client.getBacklinks(post.uri, replySource); if (!backlinks.ok) return; const promises = []; diff --git a/src/lib/at/index.ts b/src/lib/at/index.ts index 8fcb7a0..7863fcb 100644 --- a/src/lib/at/index.ts +++ b/src/lib/at/index.ts @@ -1,6 +1,9 @@ import { settings } from '$lib/settings'; +import type { Did } from '@atcute/lexicons'; import { get } from 'svelte/store'; export const slingshotUrl: URL = new URL(get(settings).endpoints.slingshot); export const spacedustUrl: URL = new URL(get(settings).endpoints.spacedust); export const constellationUrl: URL = new URL(get(settings).endpoints.constellation); + +export const httpToDidWeb = (url: string): Did => `did:web:${new URL(url).hostname}`; diff --git a/src/lib/index.ts b/src/lib/index.ts index 2424036..0f7b29b 100644 --- a/src/lib/index.ts +++ b/src/lib/index.ts @@ -28,6 +28,8 @@ export const extractDidFromUri = (uri: string): Did | null => { export const likeSource: BacklinksSource = 'app.bsky.feed.like:subject.uri'; export const repostSource: BacklinksSource = 'app.bsky.feed.repost:subject.uri'; export const replySource: BacklinksSource = 'app.bsky.feed.post:reply.parent.uri'; +export const replyRootSource: BacklinksSource = 'app.bsky.feed.post:reply.root.uri'; +export const blockSource: BacklinksSource = 'app.bsky.graph.block:subject'; export const timestampFromCursor = (cursor: string | undefined) => { if (!cursor) return undefined; diff --git a/src/lib/state.svelte.ts b/src/lib/state.svelte.ts index 4cb513d..83b4f63 100644 --- a/src/lib/state.svelte.ts +++ b/src/lib/state.svelte.ts @@ -1,30 +1,32 @@ import { writable } from 'svelte/store'; -import { - AtpClient, - newPublicClient, - type NotificationsStream, - type NotificationsStreamEvent -} from './at/client'; +import { AtpClient, type NotificationsStream, type NotificationsStreamEvent } from './at/client'; import { SvelteMap, SvelteDate, SvelteSet } from 'svelte/reactivity'; -import type { Did, Handle, InferOutput, Nsid, RecordKey, ResourceUri } from '@atcute/lexicons'; +import type { Did, Handle, Nsid, RecordKey, ResourceUri } from '@atcute/lexicons'; import { fetchPosts, hydratePosts, type PostWithUri } from './at/fetch'; import { parseCanonicalResourceUri, type AtprotoDid } from '@atcute/lexicons/syntax'; -import { AppBskyActorProfile, AppBskyFeedPost, type AppBskyGraphFollow } from '@atcute/bluesky'; -import type { ComAtprotoRepoListRecords } from '@atcute/atproto'; +import { + AppBskyActorProfile, + AppBskyFeedPost, + AppBskyGraphBlock, + type AppBskyGraphFollow +} from '@atcute/bluesky'; import type { JetstreamSubscription, JetstreamEvent } from '@atcute/jetstream'; import { expect, ok } from './result'; import type { Backlink, BacklinksSource } from './at/constellation'; import { now as tidNow } from '@atcute/tid'; import type { Records } from '@atcute/lexicons/ambient'; import { + blockSource, extractDidFromUri, likeSource, + replyRootSource, replySource, repostSource, timestampFromCursor, toCanonicalUri } from '$lib'; import { Router } from './router.svelte'; +import type { Account } from './accounts'; export const notificationStream = writable(null); export const jetstream = writable(null); @@ -134,17 +136,15 @@ export const backlinksCursors = new SvelteMap< >(); export const fetchLinksUntil = async ( + subject: Did, client: AtpClient, backlinkSource: BacklinksSource, timestamp: number = -1 ) => { - const did = client.user?.did; - if (!did) return; - - let cursorMap = backlinksCursors.get(did); + let cursorMap = backlinksCursors.get(subject); if (!cursorMap) { cursorMap = new SvelteMap(); - backlinksCursors.set(did, cursorMap); + backlinksCursors.set(subject, cursorMap); } const [_collection, source] = backlinkSource.split(':'); @@ -155,8 +155,8 @@ export const fetchLinksUntil = async ( const cursorTimestamp = timestampFromCursor(cursor); if (cursorTimestamp && cursorTimestamp <= timestamp) return; - console.log(`${did}: fetchLinksUntil`, backlinkSource, cursor, timestamp); - const result = await client.listRecordsUntil(collection, cursor, timestamp); + console.log(`${subject}: fetchLinksUntil`, backlinkSource, cursor, timestamp); + const result = await client.listRecordsUntil(subject, collection, cursor, timestamp); if (!result.ok) { console.error('failed to fetch links until', result.error); @@ -237,10 +237,6 @@ export const pulsingPostId = writable(null); export const viewClient = new AtpClient(); export const clients = new SvelteMap(); -export const getClient = async (did: Did): Promise => { - if (!clients.has(did)) clients.set(did, await newPublicClient(did)); - return clients.get(did)!; -}; export const follows = new SvelteMap>(); @@ -248,37 +244,74 @@ export const addFollows = ( did: Did, followMap: Iterable<[ResourceUri, AppBskyGraphFollow.Main]> ) => { - if (!follows.has(did)) { - follows.set(did, new SvelteMap(followMap)); + let map = follows.get(did)!; + if (!map) { + map = new SvelteMap(followMap); + follows.set(did, map); return; } - const map = follows.get(did)!; for (const [uri, record] of followMap) map.set(uri, record); }; -export const fetchFollows = async (did: AtprotoDid) => { - const client = await getClient(did); - const res = await client.listRecordsUntil('app.bsky.graph.follow'); - if (!res.ok) return; +export const fetchFollows = async ( + account: Account +): Promise> => { + const client = clients.get(account.did)!; + const res = await client.listRecordsUntil(account.did, 'app.bsky.graph.follow'); + if (!res.ok) { + console.error("can't fetch follows:", res.error); + return [].values(); + } addFollows( - did, + account.did, res.value.records.map((follow) => [follow.uri, follow.value as AppBskyGraphFollow.Main]) ); + return res.value.records.values().map((follow) => follow.value as AppBskyGraphFollow.Main); }; // this fetches up to three days of posts and interactions for using in following list -export const fetchForInteractions = async (did: AtprotoDid) => { +export const fetchForInteractions = async (client: AtpClient, subject: Did) => { const threeDaysAgo = (Date.now() - 3 * 24 * 60 * 60 * 1000) * 1000; - const client = await getClient(did); - const res = await client.listRecordsUntil('app.bsky.feed.post', undefined, threeDaysAgo); + const res = await client.listRecordsUntil(subject, 'app.bsky.feed.post', undefined, threeDaysAgo); if (!res.ok) return; - addPostsRaw(did, res.value); + const postsWithUri = res.value.records.map( + (post) => + ({ cid: post.cid, uri: post.uri, record: post.value as AppBskyFeedPost.Main }) as PostWithUri + ); + addPosts(postsWithUri); const cursorTimestamp = timestampFromCursor(res.value.cursor) ?? -1; const timestamp = Math.min(cursorTimestamp, threeDaysAgo); - console.log(`${did}: fetchForInteractions`, res.value.cursor, timestamp); - await Promise.all([repostSource].map((s) => fetchLinksUntil(client, s, timestamp))); + console.log(`${subject}: fetchForInteractions`, res.value.cursor, timestamp); + await Promise.all([repostSource].map((s) => fetchLinksUntil(subject, client, s, timestamp))); +}; + +// if did is in set, we have fetched blocks for them already (against logged in users) +export const blockFlags = new SvelteMap>(); + +export const fetchBlocked = async (client: AtpClient, subject: Did, blocker: Did) => { + const subjectUri = `at://${subject}` as ResourceUri; + const res = await client.getBacklinks(subjectUri, blockSource, [blocker], 1); + if (!res.ok) return; + if (res.value.total > 0) addBacklinks(subjectUri, blockSource, res.value.records); +}; + +export const fetchBlocks = async (account: Account) => { + const client = clients.get(account.did)!; + const res = await client.listRecordsUntil(account.did, 'app.bsky.graph.block'); + if (!res.ok) return; + for (const block of res.value.records) { + const record = block.value as AppBskyGraphBlock.Main; + const parsedUri = expect(parseCanonicalResourceUri(block.uri)); + addBacklinks(`at://${record.subject}`, blockSource, [ + { + did: parsedUri.repo, + collection: parsedUri.collection, + rkey: parsedUri.rkey + } + ]); + } }; export const allPosts = new SvelteMap>(); @@ -292,17 +325,6 @@ const hydrateCacheFn: Parameters[3] = (did, rkey) => { return cached ? ok(cached) : undefined; }; -export const addPostsRaw = ( - did: AtprotoDid, - newPosts: InferOutput -) => { - const postsWithUri = newPosts.records.map( - (post) => - ({ cid: post.cid, uri: post.uri, record: post.value as AppBskyFeedPost.Main }) as PostWithUri - ); - addPosts(postsWithUri); -}; - export const addPosts = (newPosts: Iterable) => { for (const post of newPosts) { const parsedUri = expect(parseCanonicalResourceUri(post.uri)); @@ -313,13 +335,13 @@ export const addPosts = (newPosts: Iterable) => { } posts.set(post.uri, post); if (post.record.reply) { - addBacklinks(post.record.reply.parent.uri, replySource, [ - { - did: parsedUri.repo, - collection: parsedUri.collection, - rkey: parsedUri.rkey - } - ]); + const link = { + did: parsedUri.repo, + collection: parsedUri.collection, + rkey: parsedUri.rkey + }; + addBacklinks(post.record.reply.parent.uri, replySource, [link]); + addBacklinks(post.record.reply.root.uri, replyRootSource, [link]); // update reply index const parentDid = extractDidFromUri(post.record.reply.parent.uri); @@ -363,34 +385,60 @@ export const addTimeline = (did: Did, uris: Iterable) => { }; export const fetchTimeline = async ( - did: AtprotoDid, + client: AtpClient, + subject: AtprotoDid, limit: number = 6, withBacklinks: boolean = true ) => { - const targetClient = await getClient(did); - - const cursor = postCursors.get(did); + const cursor = postCursors.get(subject); if (cursor && cursor.end) return; - const accPosts = await fetchPosts(targetClient, cursor?.value, limit, withBacklinks); - if (!accPosts.ok) throw `cant fetch posts ${did}: ${accPosts.error}`; + const accPosts = await fetchPosts(subject, client, cursor?.value, limit, withBacklinks); + if (!accPosts.ok) throw `cant fetch posts ${subject}: ${accPosts.error}`; // if the cursor is undefined, we've reached the end of the timeline - postCursors.set(did, { value: accPosts.value.cursor, end: !accPosts.value.cursor }); - const hydrated = await hydratePosts(targetClient, did, accPosts.value.posts, hydrateCacheFn); - if (!hydrated.ok) throw `cant hydrate posts ${did}: ${hydrated.error}`; + const newCursor = { value: accPosts.value.cursor, end: !accPosts.value.cursor }; + postCursors.set(subject, newCursor); + const hydrated = await hydratePosts(client, subject, accPosts.value.posts, hydrateCacheFn); + if (!hydrated.ok) throw `cant hydrate posts ${subject}: ${hydrated.error}`; addPosts(hydrated.value.values()); - addTimeline(did, hydrated.value.keys()); + addTimeline(subject, hydrated.value.keys()); + + // we only need to check blocks if the user is the subject (ie. logged in) + if (client.user?.did === subject) { + // check if any of the post authors block the user + // eslint-disable-next-line svelte/prefer-svelte-reactivity + let distinctDids = new Set(hydrated.value.keys().map((uri) => extractDidFromUri(uri)!)); + distinctDids.delete(subject); // dont need to check if user blocks themselves + const alreadyFetched = blockFlags.get(subject); + if (alreadyFetched) distinctDids = distinctDids.difference(alreadyFetched); + if (distinctDids.size > 0) + await Promise.all(distinctDids.values().map((did) => fetchBlocked(client, subject, did))); + } - console.log(`${did}: fetchTimeline`, accPosts.value.cursor); + console.log(`${subject}: fetchTimeline`, accPosts.value.cursor); + return newCursor; }; -export const fetchInteractionsUntil = async (client: AtpClient, did: Did) => { +export const fetchInteractionsToTimelineEnd = async (client: AtpClient, did: Did) => { const cursor = postCursors.get(did); if (!cursor) return; const timestamp = timestampFromCursor(cursor.value); - await Promise.all([likeSource, repostSource].map((s) => fetchLinksUntil(client, s, timestamp))); + await Promise.all( + [likeSource, repostSource].map((s) => fetchLinksUntil(did, client, s, timestamp)) + ); +}; + +export const fetchInitial = async (account: Account) => { + const client = clients.get(account.did)!; + await Promise.all([ + fetchBlocks(account), + fetchForInteractions(client, account.did), + fetchFollows(account).then((follows) => + Promise.all(follows.map((follow) => fetchForInteractions(client, follow.subject)) ?? []) + ) + ]); }; export const handleJetstreamEvent = async (event: JetstreamEvent) => { @@ -407,7 +455,7 @@ export const handleJetstreamEvent = async (event: JetstreamEvent) => { cid: commit.cid } ]; - const client = await getClient(did); + const client = clients.get(did) ?? viewClient; const hydrated = await hydratePosts(client, did, posts, hydrateCacheFn); if (!hydrated.ok) { console.error(`cant hydrate posts ${did}: ${hydrated.error}`); @@ -416,7 +464,15 @@ export const handleJetstreamEvent = async (event: JetstreamEvent) => { addPosts(hydrated.value.values()); addTimeline(did, hydrated.value.keys()); } else if (commit.operation === 'delete') { - allPosts.get(did)?.delete(uri); + const post = allPosts.get(did)?.get(uri); + if (post) { + allPosts.get(did)?.delete(uri); + // remove from timeline + timelines.get(did)?.delete(uri); + // remove reply from index + const subjectDid = extractDidFromUri(post.record.reply?.parent.uri ?? ''); + if (subjectDid) replyIndex.get(subjectDid)?.delete(uri); + } } } }; @@ -424,7 +480,11 @@ export const handleJetstreamEvent = async (event: JetstreamEvent) => { const handlePostNotification = async (event: NotificationsStreamEvent & { type: 'message' }) => { const parsedSubjectUri = expect(parseCanonicalResourceUri(event.data.link.subject)); const did = parsedSubjectUri.repo as AtprotoDid; - const client = await getClient(did); + const client = clients.get(did); + if (!client) { + console.error(`${did}: cant handle post notification, client not found !?`); + return; + } const subjectPost = await client.getRecord( AppBskyFeedPost.mainSchema, did, diff --git a/src/routes/[...catchall]/+page.svelte b/src/routes/[...catchall]/+page.svelte index 4980387..f0e4f01 100644 --- a/src/routes/[...catchall]/+page.svelte +++ b/src/routes/[...catchall]/+page.svelte @@ -12,8 +12,6 @@ import { clients, postCursors, - fetchForInteractions, - fetchFollows, follows, notificationStream, viewClient, @@ -22,7 +20,8 @@ handleNotification, addPosts, addTimeline, - router + router, + fetchInitial } from '$lib/state.svelte'; import { get } from 'svelte/store'; import Icon from '@iconify/svelte'; @@ -113,7 +112,8 @@ 'app.bsky.feed.post:embed.record.uri', 'app.bsky.feed.repost:subject.uri', 'app.bsky.feed.like:subject.uri', - 'app.bsky.graph.follow:subject' + 'app.bsky.graph.follow:subject', + 'app.bsky.graph.block:subject' ) ); }); @@ -144,16 +144,7 @@ } if (!$accounts.some((account) => account.did === selectedDid)) selectedDid = $accounts[0].did; // console.log('onMount selectedDid', selectedDid); - Promise.all($accounts.map(loginAccount)).then(() => { - $accounts.forEach((account) => { - fetchFollows(account.did).then(() => - follows - .get(account.did) - ?.forEach((follow) => fetchForInteractions(follow.subject as AtprotoDid)) - ); - fetchForInteractions(account.did); - }); - }); + Promise.all($accounts.map(loginAccount)).then(() => $accounts.forEach(fetchInitial)); } else { selectedDid = null; } @@ -163,11 +154,11 @@ $effect(() => { const wantedDids: Did[] = ['did:web:guestbook.gaze.systems']; - - for (const followMap of follows.values()) - for (const follow of followMap.values()) wantedDids.push(follow.subject); - for (const account of $accounts) wantedDids.push(account.did); - + const followDids = follows + .values() + .flatMap((followMap) => followMap.values().map((follow) => follow.subject)); + const accountDids = $accounts.values().map((account) => account.did); + wantedDids.push(...followDids, ...accountDids); // console.log('updating jetstream options:', wantedDids); $jetstream?.updateOptions({ wantedDids }); });