migrate runtime call sites to lex clients and sdk actions

Phase 3 tasks 3-7 (parallel wave): composer/post pipeline on
pdsClient/appviewClient with structural+instance blob guards and a golden
CID fixture test; chat Convo/EventBus/queries on the dedicated chat client;
preferences sugar to SDK actions on the PDS client; remaining state/queries
producers (usePostThread unspecced flip, video scoped-token clients,
notifications, starter packs, lists) to client.call; UI runtime sweep
(AtUri from @atproto/syntax, moderation fns from @bsky.app/sdk/moderation,
SDK RichText, ozone reason tokens, guard rewrites via #/types/bsky).

Intermediate checkpoint (hooks skipped): ~114 typecheck errors remain in
cross-boundary consumer files, resolved by the type-only codemod next.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Samuel Newman
2026-07-16 21:04:44 +03:00
parent b1a4e1cd16
commit 300a50b69a
296 changed files with 4507 additions and 3903 deletions
+121
View File
@@ -0,0 +1,121 @@
/*
* The jest suite ships global manual mocks for `multiformats/cid` and
* `multiformats/hashes/hasher` (in root `__mocks__/`) so unrelated tests don't
* pull in real crypto. This test is precisely about the real CID hashing, so we
* opt back into the actual implementations here.
*/
jest.unmock('multiformats/cid')
jest.unmock('multiformats/hashes/hasher')
import {BlobRef} from '@atproto/api'
import {CID} from 'multiformats/cid'
import {computeCid} from '#/lib/api/computeCid'
import {type app, type com} from '#/lexicons'
/*
* Golden-CID regression test for the composer post pipeline (design section F).
*
* `computeCid` hashes a post record in the client so a thread's later posts can
* reference earlier posts by CID before the server assigns them. The hash is
* byte-sensitive: any drift in how records (especially blobs) are serialized to
* DAG-CBOR silently produces the wrong CID and breaks reply chains with NO type
* error. These golden values were captured from the PRE-migration `computeCid`
* (the `instanceof BlobRef` path) and MUST remain byte-identical after the guard
* is changed to the structural `isBlobRef` shape check.
*
* The blob CID below is a fixed, deterministic CIDv1/raw/sha256 used purely as a
* stable fixture - it is not derived from any real upload.
*/
const BLOB_CID = 'bafkreieq5jui4j25lacwomsqgjeswwl3y5zcdrresptwgmfylxo2depppq'
/**
* Build a post record with an image embed whose blob is the given value. Used to
* prove that a `@atproto/api` `BlobRef` class instance (the shape the not-yet
* -migrated video path still yields) and a plain-JSON lex blob (the shape lex
* `uploadBlob` returns) hash to the SAME CID.
*/
function postWithImageBlob(blob: unknown): app.bsky.feed.post.Main {
return {
$type: 'app.bsky.feed.post',
createdAt: '2024-01-01T00:00:00.001Z',
text: 'post with image',
embed: {
$type: 'app.bsky.embed.images',
images: [
{
image: blob,
alt: 'alt text',
aspectRatio: {width: 100, height: 200},
},
],
},
} as app.bsky.feed.post.Main
}
describe('computeCid', () => {
it('case 1: plain post record with no blob', async () => {
const record: app.bsky.feed.post.Main = {
$type: 'app.bsky.feed.post',
createdAt: '2024-01-01T00:00:00.000Z',
text: 'hello world',
}
expect(await computeCid(record)).toBe(
'bafyreieawtmh7hwfrqpamqkodza5r62bbfhsepe2iyustgxhgbhi6b2lfi',
)
})
it('case 2: record whose embed carries a BlobRef class instance', async () => {
const blob = new BlobRef(CID.parse(BLOB_CID), 'image/jpeg', 12345)
expect(await computeCid(postWithImageBlob(blob))).toBe(
'bafyreiem7g6vja66nebr7he4fshfnlyndyldbvle2n265oixscmepjcbii',
)
})
it('case 2b: a plain-JSON blob object hashes identically to the class instance', async () => {
/*
* This is the post-migration shape: lex `uploadBlob` returns a plain object
* `{$type: 'blob', ref, mimeType, size}` (with `ref` a parsed CID), not a
* `BlobRef` class instance. The structural `isBlobRef` guard must treat it
* exactly like the class instance so the CID is unchanged. Under the
* pre-change `instanceof` code this case already matches because the plain
* object walks through `prepareForHashing` unchanged and DAG-CBOR encodes
* its CID `ref` the same way `.ipld()` does.
*/
const blob = {
$type: 'blob' as const,
ref: CID.parse(BLOB_CID),
mimeType: 'image/jpeg',
size: 12345,
}
expect(await computeCid(postWithImageBlob(blob))).toBe(
'bafyreiem7g6vja66nebr7he4fshfnlyndyldbvle2n265oixscmepjcbii',
)
})
it('case 3: three-post thread chains reply StrongRef CIDs', async () => {
const did = 'did:plc:abc123'
const base = new Date('2024-01-01T00:00:00.000Z')
const golden = [
'bafyreig62rxs34h5rvznfrracwkjlfgad5b25qxglp2hcziqdfas2nw2ee',
'bafyreicxcj2tq5jrh5jcaczg3eli5cvxitgzu7kpu3fm5v3njq2byjxirq',
'bafyreigvaswuhlpd7dllja2xrqswhqbruyv2kar7mvbn7gdm5ldzu6vkti',
]
let reply: app.bsky.feed.post.Main['reply'] | undefined
for (let i = 0; i < 3; i++) {
const now = new Date(base.getTime() + i)
const uri = `at://${did}/app.bsky.feed.post/rkey${i}`
const record = {
$type: 'app.bsky.feed.post',
createdAt: now.toISOString(),
text: `post ${i}`,
reply,
} as app.bsky.feed.post.Main
const cid = await computeCid(record)
expect(cid).toBe(golden[i])
const ref = {cid, uri} as com.atproto.repo.strongRef.Main
reply = {root: reply?.root ?? ref, parent: ref}
}
})
})
+158
View File
@@ -0,0 +1,158 @@
import {sha256} from 'js-sha256'
import {CID} from 'multiformats/cid'
import * as Hasher from 'multiformats/hashes/hasher'
import {app} from '#/lexicons'
/*
* Client-side CID computation for post records, extracted from the post
* pipeline so it can be unit-tested in isolation (importing the pipeline pulls
* in the native gallery/media chain). See `computeCid.test.ts` for the golden
* -CID regression fixtures that gate any change to this serialization.
*/
// The built-in hashing functions from multiformats (`multiformats/hashes/sha2`)
// are meant for Node.js, this is the cross-platform equivalent.
const mf_sha256 = Hasher.from({
name: 'sha2-256',
code: 0x12,
encode: input => {
const digest = sha256.arrayBuffer(input)
return new Uint8Array(digest)
},
})
export async function computeCid(
record: app.bsky.feed.post.Main,
): Promise<string> {
/*
* Lazily loaded since it's only needed when posting a thread, and its
* `cborg` dependency is ~190KB that would otherwise be in the initial
* web bundle.
*/
const dcbor = await importDagCbor()
// IMPORTANT: `prepareObject` prepares the record to be hashed by removing
// fields with undefined value, and converting BlobRef instances to the
// right IPLD representation.
const prepared = prepareForHashing(record)
// 1. Encode the record into DAG-CBOR format
const encoded = dcbor.encode(prepared)
// 2. Hash the record in SHA-256 (code 0x12)
const digest = await mf_sha256.digest(encoded)
// 3. Create a CIDv1, specifying DAG-CBOR as content (code 0x71)
const cid = CID.createV1(0x71, digest)
// 4. Get the Base32 representation of the CID (`b` prefix)
return cid.toString()
}
/**
* True for a plain-JSON lexicon blob, the shape lex `uploadBlob` now returns
* (`{$type: 'blob', ref, mimeType, size}` with `ref` a parsed CID). Replaces
* the old `instanceof BlobRef` check, since lex blobs are plain objects, not
* class instances (design section F).
*/
function isBlobRef(v: unknown): boolean {
if (v == null || typeof v !== 'object') return false
const o = v as Record<string, unknown>
return o.$type === 'blob' && 'ref' in o && 'mimeType' in o
}
/**
* True for a legacy `@atproto/api` `BlobRef` class instance. During the
* migration the video embed path still yields these (its blob comes from the
* not-yet-migrated `app.bsky.video.getJobStatus` bridge call), so we must keep
* handling them here even though the composer's own uploads are now plain lex
* blobs. A class instance is duck-typed by its `ipld()` method plus the
* `ref`/`mimeType` fields; it has NO `$type` and a non-plain prototype, so it
* would otherwise slip past both `isBlobRef` and `isPlainObject` and be encoded
* wrong - silently breaking video reply-chain CIDs.
*/
function isBlobRefInstance(
v: unknown,
): v is {ipld: () => unknown; ref: unknown; mimeType: unknown} {
if (v == null || typeof v !== 'object') return false
const o = v as Record<string, unknown>
return typeof o.ipld === 'function' && 'ref' in o && 'mimeType' in o
}
// Returns a transformed version of the object for use in DAG-CBOR.
// eslint-disable-next-line @typescript-eslint/no-explicit-any
function prepareForHashing(v: any): any {
/*
* A plain-JSON lex blob is already in the right IPLD shape (its `ref` is a
* parsed CID that DAG-CBOR encodes as a CID link), so pass it through
* untouched.
*/
if (isBlobRef(v)) {
return v
}
/*
* A legacy `BlobRef` class instance must be converted via `ipld()` to the
* plain `{$type, ref, mimeType, size}` object; encoding the instance directly
* would emit its internal `original` field and omit `$type`, producing the
* wrong CID. `ipld()` returns exactly what `isBlobRef` accepts above.
*/
if (isBlobRefInstance(v)) {
return v.ipld()
}
// Walk through arrays
if (Array.isArray(v)) {
let pure = true
const mapped = v.map(value => {
if (value !== (value = prepareForHashing(value))) {
pure = false
}
return value
})
return pure ? v : mapped
}
// Walk through plain objects
if (isPlainObject(v)) {
const rec = v as Record<string, unknown>
const obj: Record<string, unknown> = {}
let pure = true
for (const key in rec) {
let value = rec[key]
// `value` is undefined
if (value === undefined) {
pure = false
continue
}
// `prepareObject` returned a value that's different from what we had before
if (value !== (value = prepareForHashing(value))) {
pure = false
}
obj[key] = value
}
// Return as is if we haven't needed to tamper with anything
return pure ? v : obj
}
return v
}
// eslint-disable-next-line @typescript-eslint/no-explicit-any
function isPlainObject(v: any): boolean {
if (typeof v !== 'object' || v === null) {
return false
}
const proto = Object.getPrototypeOf(v)
return proto === Object.prototype || proto === null
}
/**
* Load `@ipld/dag-cbor` on demand. The dynamic `import()` lets web bundlers
* emit it (and its ~190KB `cborg` dependency) as a separate chunk that only
* loads when posting a thread. Under jest (which runs without
* `--experimental-vm-modules`) dynamic import throws, so we fall back to a
* lazy `require`, which resolves through the test moduleNameMapper.
*/
function importDagCbor(): Promise<typeof import('@ipld/dag-cbor')> {
if (process.env.NODE_ENV === 'test') {
// eslint-disable-next-line @typescript-eslint/no-require-imports
return Promise.resolve(require('@ipld/dag-cbor'))
}
return import('@ipld/dag-cbor')
}
+24 -28
View File
@@ -1,23 +1,21 @@
import {
AppBskyFeedDefs,
type AppBskyFeedGetAuthorFeed as GetAuthorFeed,
} from '@atproto/api'
import {type Client} from '@atproto/lex-client'
import {type SessionAgent} from '#/state/session'
import {app} from '#/lexicons'
import * as bsky from '#/types/bsky'
import {type FeedAPI, type FeedAPIResponse} from './types'
export class AuthorFeedAPI implements FeedAPI {
agent: SessionAgent
_params: GetAuthorFeed.QueryParams
client: Client
_params: app.bsky.feed.getAuthorFeed.$Params
constructor({
agent,
client,
feedParams,
}: {
agent: SessionAgent
feedParams: GetAuthorFeed.QueryParams
client: Client
feedParams: app.bsky.feed.getAuthorFeed.$Params
}) {
this.agent = agent
this.client = client
this._params = feedParams
}
@@ -27,12 +25,12 @@ export class AuthorFeedAPI implements FeedAPI {
return params
}
async peekLatest(): Promise<AppBskyFeedDefs.FeedViewPost> {
const res = await this.agent.getAuthorFeed({
async peekLatest(): Promise<app.bsky.feed.defs.FeedViewPost> {
const res = await this.client.call(app.bsky.feed.getAuthorFeed, {
...this.params,
limit: 1,
})
return res.data.feed[0]
return res.feed[0]
}
async fetch({
@@ -42,28 +40,26 @@ export class AuthorFeedAPI implements FeedAPI {
cursor: string | undefined
limit: number
}): Promise<FeedAPIResponse> {
const res = await this.agent.getAuthorFeed({
const res = await this.client.call(app.bsky.feed.getAuthorFeed, {
...this.params,
cursor,
limit,
})
if (res.success) {
return {
cursor: res.data.cursor,
feed: this._filter(res.data.feed),
}
}
return {
feed: [],
cursor: res.cursor,
feed: this._filter(res.feed),
}
}
_filter(feed: AppBskyFeedDefs.FeedViewPost[]) {
_filter(feed: app.bsky.feed.defs.FeedViewPost[]) {
if (this.params.filter === 'posts_and_author_threads') {
return feed.filter(post => {
const isReply = post.reply
const isRepost = AppBskyFeedDefs.isReasonRepost(post.reason)
const isPin = AppBskyFeedDefs.isReasonPin(post.reason)
const isRepost = bsky.isType(
app.bsky.feed.defs.reasonRepost,
post.reason,
)
const isPin = bsky.isType(app.bsky.feed.defs.reasonPin, post.reason)
if (!isReply) return true
if (isRepost || isPin) return true
return isReply && isAuthorReplyChain(this.params.actor, post, feed)
@@ -76,15 +72,15 @@ export class AuthorFeedAPI implements FeedAPI {
function isAuthorReplyChain(
actor: string,
post: AppBskyFeedDefs.FeedViewPost,
posts: AppBskyFeedDefs.FeedViewPost[],
post: app.bsky.feed.defs.FeedViewPost,
posts: app.bsky.feed.defs.FeedViewPost[],
): boolean {
// current post is by a different user (shouldn't happen)
if (post.post.author.did !== actor) return false
const replyParent = post.reply?.parent
if (AppBskyFeedDefs.isPostView(replyParent)) {
if (bsky.isType(app.bsky.feed.defs.postView, replyParent)) {
// reply parent is by a different user
if (replyParent.author.did !== actor) return false
+67 -52
View File
@@ -1,47 +1,53 @@
import {
type AppBskyFeedDefs,
type AppBskyFeedGetFeed as GetCustomFeed,
AtpAgent,
jsonStringToLex,
} from '@atproto/api'
import {lexParse} from '@atproto/lex'
import {Client} from '@atproto/lex-client'
import {
getAppLanguageAsContentLanguage,
getContentLanguages,
} from '#/state/preferences/languages'
import {type SessionAgent} from '#/state/session'
import {app} from '#/lexicons'
import {type FeedAPI, type FeedAPIResponse} from './types'
import {createBskyTopicsHeader, isBlueskyOwnedFeed} from './utils'
/**
* Input params for {@link CustomFeedAPI}. The generated `$Params` type reflects
* post-parse output, where `limit` (which has a lexicon default) is required;
* callers supply only `feed` and let `limit`/`cursor` come from `fetch`.
*/
type CustomFeedParams = {feed: app.bsky.feed.getFeed.$Params['feed']} & Partial<
Omit<app.bsky.feed.getFeed.$Params, 'feed'>
>
export class CustomFeedAPI implements FeedAPI {
agent: SessionAgent
params: GetCustomFeed.QueryParams
client: Client
params: CustomFeedParams
userInterests?: string
constructor({
agent,
client,
feedParams,
userInterests,
}: {
agent: SessionAgent
feedParams: GetCustomFeed.QueryParams
client: Client
feedParams: CustomFeedParams
userInterests?: string
}) {
this.agent = agent
this.client = client
this.params = feedParams
this.userInterests = userInterests
}
async peekLatest(): Promise<AppBskyFeedDefs.FeedViewPost> {
async peekLatest(): Promise<app.bsky.feed.defs.FeedViewPost> {
const contentLangs = getContentLanguages().join(',')
const res = await this.agent.app.bsky.feed.getFeed(
const res = await this.client.call(
app.bsky.feed.getFeed,
{
...this.params,
limit: 1,
},
{headers: {'Accept-Language': contentLangs}},
)
return res.data.feed[0]
return res.feed[0]
}
async fetch({
@@ -52,41 +58,49 @@ export class CustomFeedAPI implements FeedAPI {
limit: number
}): Promise<FeedAPIResponse> {
const contentLangs = getContentLanguages().join(',')
const agent = this.agent
const isBlueskyOwned = isBlueskyOwnedFeed(this.params.feed)
const res = agent.did
? await this.agent.app.bsky.feed.getFeed(
{
...this.params,
cursor,
limit,
let feed: app.bsky.feed.defs.FeedViewPost[]
let resCursor: string | undefined
if (this.client.did) {
const res = await this.client.call(
app.bsky.feed.getFeed,
{
...this.params,
cursor,
limit,
},
{
headers: {
...(isBlueskyOwned
? createBskyTopicsHeader(this.userInterests)
: {}),
'Accept-Language': contentLangs,
},
{
headers: {
...(isBlueskyOwned
? createBskyTopicsHeader(this.userInterests)
: {}),
'Accept-Language': contentLangs,
},
},
)
: await loggedOutFetch({...this.params, cursor, limit})
if (res.success) {
// NOTE
// some custom feeds fail to enforce the pagination limit
// so we manually truncate here
// -prf
if (res.data.feed.length > limit) {
res.data.feed = res.data.feed.slice(0, limit)
}
return {
cursor: res.data.feed.length ? res.data.cursor : undefined,
feed: res.data.feed,
},
)
feed = res.feed
resCursor = res.cursor
} else {
const res = await loggedOutFetch({...this.params, cursor, limit})
if (!res.success) {
return {feed: []}
}
feed = res.data.feed
resCursor = res.data.cursor
}
// NOTE
// some custom feeds fail to enforce the pagination limit
// so we manually truncate here
// -prf
if (feed.length > limit) {
feed = feed.slice(0, limit)
}
return {
feed: [],
cursor: feed.length ? resCursor : undefined,
feed,
}
}
}
@@ -106,15 +120,16 @@ async function loggedOutFetch({
feed: string
limit: number
cursor?: string
}) {
}): Promise<{success: boolean; data: app.bsky.feed.getFeed.$OutputBody}> {
let contentLangs = getAppLanguageAsContentLanguage()
/**
* Copied from our root `Agent` class
* @see https://github.com/bluesky-social/atproto/blob/60df3fc652b00cdff71dd9235d98a7a4bb828f05/packages/api/src/agent.ts#L120
/*
* Copied from our root `Agent` class. The global (`;redact`-suffixed) app
* labelers are kept on the lex `Client` static (synced in
* `#/state/session/moderation`), replacing the old `AtpAgent.appLabelers`.
*/
const labelersHeader = {
'atproto-accept-labelers': AtpAgent.appLabelers
'atproto-accept-labelers': Client.appLabelers
.map(l => `${l};redact`)
.join(', '),
}
@@ -130,7 +145,7 @@ async function loggedOutFetch({
},
)
let data = res.ok
? (jsonStringToLex(await res.text()) as GetCustomFeed.OutputSchema)
? (lexParse(await res.text()) as app.bsky.feed.getFeed.$OutputBody)
: null
if (data?.feed?.length) {
return {
@@ -147,7 +162,7 @@ async function loggedOutFetch({
{method: 'GET', headers: {'Accept-Language': '', ...labelersHeader}},
)
data = res.ok
? (jsonStringToLex(await res.text()) as GetCustomFeed.OutputSchema)
? (lexParse(await res.text()) as app.bsky.feed.getFeed.$OutputBody)
: null
if (data?.feed?.length) {
return {
+15 -8
View File
@@ -1,21 +1,28 @@
import {type AppBskyFeedDefs} from '@atproto/api'
import {type Client} from '@atproto/lex-client'
import {DEMO_FEED} from '#/lib/demo'
import {type SessionAgent} from '#/state/session'
import {type app} from '#/lexicons'
import {toLex} from '#/types/bsky'
import {type FeedAPI, type FeedAPIResponse} from './types'
export class DemoFeedAPI implements FeedAPI {
agent: SessionAgent
client: Client
constructor({agent}: {agent: SessionAgent}) {
this.agent = agent
constructor({client}: {client: Client}) {
this.client = client
}
async peekLatest(): Promise<AppBskyFeedDefs.FeedViewPost> {
return DEMO_FEED.feed[0]
async peekLatest(): Promise<app.bsky.feed.defs.FeedViewPost> {
/*
* TODO(phase4): drop this cast once `#/lib/demo` (DEMO_FEED) sources its
* feed items from `#/lexicons`. Its records are still typed against the old
* `@atproto/api` FeedViewPost, which does not assign to the branded lexicon
* type this method must return.
*/
return toLex(DEMO_FEED.feed[0])
}
async fetch(): Promise<FeedAPIResponse> {
return DEMO_FEED
return toLex(DEMO_FEED)
}
}
+11 -16
View File
@@ -1,20 +1,20 @@
import {type AppBskyFeedDefs} from '@atproto/api'
import {type Client} from '@atproto/lex-client'
import {type SessionAgent} from '#/state/session'
import {app} from '#/lexicons'
import {type FeedAPI, type FeedAPIResponse} from './types'
export class FollowingFeedAPI implements FeedAPI {
agent: SessionAgent
client: Client
constructor({agent}: {agent: SessionAgent}) {
this.agent = agent
constructor({client}: {client: Client}) {
this.client = client
}
async peekLatest(): Promise<AppBskyFeedDefs.FeedViewPost> {
const res = await this.agent.getTimeline({
async peekLatest(): Promise<app.bsky.feed.defs.FeedViewPost> {
const res = await this.client.call(app.bsky.feed.getTimeline, {
limit: 1,
})
return res.data.feed[0]
return res.feed[0]
}
async fetch({
@@ -24,18 +24,13 @@ export class FollowingFeedAPI implements FeedAPI {
cursor: string | undefined
limit: number
}): Promise<FeedAPIResponse> {
const res = await this.agent.getTimeline({
const res = await this.client.call(app.bsky.feed.getTimeline, {
cursor,
limit,
})
if (res.success) {
return {
cursor: res.data.cursor,
feed: res.data.feed,
}
}
return {
feed: [],
cursor: res.cursor,
feed: res.feed,
}
}
}
+21 -16
View File
@@ -1,7 +1,8 @@
import {type AppBskyFeedDefs} from '@atproto/api'
import {type Client} from '@atproto/lex-client'
import {type AtUriString} from '@atproto/syntax'
import {PROD_DEFAULT_FEED} from '#/lib/constants'
import {type SessionAgent} from '#/state/session'
import {type app} from '#/lexicons'
import {CustomFeedAPI} from './custom'
import {FollowingFeedAPI} from './following'
import {type FeedAPI, type FeedAPIResponse} from './types'
@@ -14,7 +15,11 @@ import {type FeedAPI, type FeedAPIResponse} from './types'
// we use this fallback marker post to drive this instead. see Feed.tsx
// for the usage.
// -prf
export const FALLBACK_MARKER_POST: AppBskyFeedDefs.FeedViewPost = {
/*
* A sentinel post whose fields intentionally violate the branded lexicon
* formats (`uri`, `did`, `indexedAt`), so it is asserted into the view type.
*/
export const FALLBACK_MARKER_POST: app.bsky.feed.defs.FeedViewPost = {
post: {
uri: 'fallback-marker-post',
cid: 'fake',
@@ -25,10 +30,10 @@ export const FALLBACK_MARKER_POST: AppBskyFeedDefs.FeedViewPost = {
},
indexedAt: new Date().toISOString(),
},
}
} as unknown as app.bsky.feed.defs.FeedViewPost
export class HomeFeedAPI implements FeedAPI {
agent: SessionAgent
client: Client
following: FollowingFeedAPI
discover: CustomFeedAPI
usingDiscover = false
@@ -37,32 +42,32 @@ export class HomeFeedAPI implements FeedAPI {
constructor({
userInterests,
agent,
client,
}: {
userInterests?: string
agent: SessionAgent
client: Client
}) {
this.agent = agent
this.following = new FollowingFeedAPI({agent})
this.client = client
this.following = new FollowingFeedAPI({client})
this.discover = new CustomFeedAPI({
agent,
feedParams: {feed: PROD_DEFAULT_FEED('whats-hot')},
client,
feedParams: {feed: PROD_DEFAULT_FEED('whats-hot') as AtUriString},
})
this.userInterests = userInterests
}
reset() {
this.following = new FollowingFeedAPI({agent: this.agent})
this.following = new FollowingFeedAPI({client: this.client})
this.discover = new CustomFeedAPI({
agent: this.agent,
feedParams: {feed: PROD_DEFAULT_FEED('whats-hot')},
client: this.client,
feedParams: {feed: PROD_DEFAULT_FEED('whats-hot') as AtUriString},
userInterests: this.userInterests,
})
this.usingDiscover = false
this.itemCursor = 0
}
async peekLatest(): Promise<AppBskyFeedDefs.FeedViewPost> {
async peekLatest(): Promise<app.bsky.feed.defs.FeedViewPost> {
if (this.usingDiscover) {
return this.discover.peekLatest()
}
@@ -81,7 +86,7 @@ export class HomeFeedAPI implements FeedAPI {
}
let returnCursor
let posts: AppBskyFeedDefs.FeedViewPost[] = []
let posts: app.bsky.feed.defs.FeedViewPost[] = []
if (!this.usingDiscover) {
const res = await this.following.fetch({cursor, limit})
+16 -24
View File
@@ -1,32 +1,29 @@
import {
type AppBskyFeedDefs,
type AppBskyFeedGetActorLikes as GetActorLikes,
} from '@atproto/api'
import {type Client} from '@atproto/lex-client'
import {type SessionAgent} from '#/state/session'
import {app} from '#/lexicons'
import {type FeedAPI, type FeedAPIResponse} from './types'
export class LikesFeedAPI implements FeedAPI {
agent: SessionAgent
params: GetActorLikes.QueryParams
client: Client
params: app.bsky.feed.getActorLikes.$Params
constructor({
agent,
client,
feedParams,
}: {
agent: SessionAgent
feedParams: GetActorLikes.QueryParams
client: Client
feedParams: app.bsky.feed.getActorLikes.$Params
}) {
this.agent = agent
this.client = client
this.params = feedParams
}
async peekLatest(): Promise<AppBskyFeedDefs.FeedViewPost> {
const res = await this.agent.getActorLikes({
async peekLatest(): Promise<app.bsky.feed.defs.FeedViewPost> {
const res = await this.client.call(app.bsky.feed.getActorLikes, {
...this.params,
limit: 1,
})
return res.data.feed[0]
return res.feed[0]
}
async fetch({
@@ -36,21 +33,16 @@ export class LikesFeedAPI implements FeedAPI {
cursor: string | undefined
limit: number
}): Promise<FeedAPIResponse> {
const res = await this.agent.getActorLikes({
const res = await this.client.call(app.bsky.feed.getActorLikes, {
...this.params,
cursor,
limit,
})
if (res.success) {
// HACKFIX: the API incorrectly returns a cursor when there are no items -sfn
const isEmptyPage = res.data.feed.length === 0
return {
cursor: isEmptyPage ? undefined : res.data.cursor,
feed: res.data.feed,
}
}
// HACKFIX: the API incorrectly returns a cursor when there are no items -sfn
const isEmptyPage = res.feed.length === 0
return {
feed: [],
cursor: isEmptyPage ? undefined : res.cursor,
feed: res.feed,
}
}
}
+14 -22
View File
@@ -1,32 +1,29 @@
import {
type Agent,
type AppBskyFeedDefs,
type AppBskyFeedGetListFeed as GetListFeed,
} from '@atproto/api'
import {type Client} from '@atproto/lex-client'
import {app} from '#/lexicons'
import {type FeedAPI, type FeedAPIResponse} from './types'
export class ListFeedAPI implements FeedAPI {
agent: Agent
params: GetListFeed.QueryParams
client: Client
params: app.bsky.feed.getListFeed.$Params
constructor({
agent,
client,
feedParams,
}: {
agent: Agent
feedParams: GetListFeed.QueryParams
client: Client
feedParams: app.bsky.feed.getListFeed.$Params
}) {
this.agent = agent
this.client = client
this.params = feedParams
}
async peekLatest(): Promise<AppBskyFeedDefs.FeedViewPost> {
const res = await this.agent.app.bsky.feed.getListFeed({
async peekLatest(): Promise<app.bsky.feed.defs.FeedViewPost> {
const res = await this.client.call(app.bsky.feed.getListFeed, {
...this.params,
limit: 1,
})
return res.data.feed[0]
return res.feed[0]
}
async fetch({
@@ -36,19 +33,14 @@ export class ListFeedAPI implements FeedAPI {
cursor: string | undefined
limit: number
}): Promise<FeedAPIResponse> {
const res = await this.agent.app.bsky.feed.getListFeed({
const res = await this.client.call(app.bsky.feed.getListFeed, {
...this.params,
cursor,
limit,
})
if (res.success) {
return {
cursor: res.data.cursor,
feed: res.data.feed,
}
}
return {
feed: [],
cursor: res.cursor,
feed: res.feed,
}
}
}
+70 -44
View File
@@ -1,4 +1,5 @@
import {type AppBskyFeedDefs, type AppBskyFeedGetTimeline} from '@atproto/api'
import {type Client} from '@atproto/lex-client'
import {type AtUriString} from '@atproto/syntax'
import shuffle from 'lodash.shuffle'
import {bundleAsync} from '#/lib/async/bundle'
@@ -6,7 +7,8 @@ import {timeout} from '#/lib/async/timeout'
import {feedUriToHref} from '#/lib/strings/url-helpers'
import {getContentLanguages} from '#/state/preferences/languages'
import {type FeedParams} from '#/state/queries/post-feed'
import {type SessionAgent} from '#/state/session'
import {app} from '#/lexicons'
import {toLex} from '#/types/bsky'
import {FeedTuner} from '../feed-manip'
import {type FeedTunerFn} from '../feed-manip'
import {
@@ -19,9 +21,21 @@ import {createBskyTopicsHeader, isBlueskyOwnedFeed} from './utils'
const REQUEST_WAIT_MS = 500 // 500ms
const POST_AGE_CUTOFF = 60e3 * 60 * 24 // 24hours
/**
* Internal result shape for a single feed page fetch. Lex `client.call`
* returns the response body directly (throwing on error), so we no longer
* carry the old `{success, headers, data}` wrapper - `success` here just
* distinguishes an empty/errored fetch from a populated one.
*/
type FeedPage = {
success: boolean
cursor?: string
feed: app.bsky.feed.defs.FeedViewPost[]
}
export class MergeFeedAPI implements FeedAPI {
userInterests?: string
agent: SessionAgent
client: Client
params: FeedParams
feedTuners: FeedTunerFn[]
following: MergeFeedSource_Following
@@ -31,29 +45,29 @@ export class MergeFeedAPI implements FeedAPI {
sampleCursor = 0
constructor({
agent,
client,
feedParams,
feedTuners,
userInterests,
}: {
agent: SessionAgent
client: Client
feedParams: FeedParams
feedTuners: FeedTunerFn[]
userInterests?: string
}) {
this.agent = agent
this.client = client
this.params = feedParams
this.feedTuners = feedTuners
this.userInterests = userInterests
this.following = new MergeFeedSource_Following({
agent: this.agent,
client: this.client,
feedTuners: this.feedTuners,
})
}
reset() {
this.following = new MergeFeedSource_Following({
agent: this.agent,
client: this.client,
feedTuners: this.feedTuners,
})
this.customFeeds = []
@@ -65,7 +79,7 @@ export class MergeFeedAPI implements FeedAPI {
this.params.mergeFeedSources.map(
feedUri =>
new MergeFeedSource_Custom({
agent: this.agent,
client: this.client,
feedUri,
feedTuners: this.feedTuners,
userInterests: this.userInterests,
@@ -77,11 +91,11 @@ export class MergeFeedAPI implements FeedAPI {
}
}
async peekLatest(): Promise<AppBskyFeedDefs.FeedViewPost> {
const res = await this.agent.getTimeline({
async peekLatest(): Promise<app.bsky.feed.defs.FeedViewPost> {
const res = await this.client.call(app.bsky.feed.getTimeline, {
limit: 1,
})
return res.data.feed[0]
return res.feed[0]
}
async fetch({
@@ -124,7 +138,7 @@ export class MergeFeedAPI implements FeedAPI {
await Promise.all(promises)
// assemble a response by sampling from feeds with content
const posts: AppBskyFeedDefs.FeedViewPost[] = []
const posts: app.bsky.feed.defs.FeedViewPost[] = []
while (posts.length < limit) {
let slice = this.sampleItem()
if (slice[0]) {
@@ -172,21 +186,21 @@ export class MergeFeedAPI implements FeedAPI {
}
class MergeFeedSource {
agent: SessionAgent
client: Client
feedTuners: FeedTunerFn[]
sourceInfo: ReasonFeedSource | undefined
cursor: string | undefined = undefined
queue: AppBskyFeedDefs.FeedViewPost[] = []
queue: app.bsky.feed.defs.FeedViewPost[] = []
hasMore = true
constructor({
agent,
client,
feedTuners,
}: {
agent: SessionAgent
client: Client
feedTuners: FeedTunerFn[]
}) {
this.agent = agent
this.client = client
this.feedTuners = feedTuners
}
@@ -198,7 +212,7 @@ class MergeFeedSource {
return this.hasMore && this.queue.length === 0
}
take(n: number): AppBskyFeedDefs.FeedViewPost[] {
take(n: number): app.bsky.feed.defs.FeedViewPost[] {
return this.queue.splice(0, n)
}
@@ -209,9 +223,9 @@ class MergeFeedSource {
_fetchNextInner = bundleAsync(async (n: number) => {
const res = await this._getFeed(this.cursor, n)
if (res.success) {
this.cursor = res.data.cursor
if (res.data.feed.length) {
this.queue = this.queue.concat(res.data.feed)
this.cursor = res.cursor
if (res.feed.length) {
this.queue = this.queue.concat(res.feed)
} else {
this.hasMore = false
}
@@ -223,7 +237,7 @@ class MergeFeedSource {
protected _getFeed(
_cursor: string | undefined,
_limit: number,
): Promise<AppBskyFeedGetTimeline.Response> {
): Promise<FeedPage> {
throw new Error('Must be overridden')
}
}
@@ -238,39 +252,51 @@ class MergeFeedSource_Following extends MergeFeedSource {
protected async _getFeed(
cursor: string | undefined,
limit: number,
): Promise<AppBskyFeedGetTimeline.Response> {
const res = await this.agent.getTimeline({cursor, limit})
): Promise<FeedPage> {
const res = await this.client.call(app.bsky.feed.getTimeline, {
cursor,
limit,
})
// run the tuner pre-emptively to ensure better mixing
const slices = this.tuner.tune(res.data.feed, {
const slices = this.tuner.tune(res.feed, {
dryRun: false,
})
res.data.feed = slices.map(slice => slice._feedPost)
return res
return {
success: true,
cursor: res.cursor,
/*
* TODO(phase4): drop the toLex once `#/lib/api/feed-manip` (FeedTuner)
* flips its FeedViewPost source from `@atproto/api` to `#/lexicons`. Its
* `_feedPost` is still the old-world view type, which does not assign to
* the branded lexicon type this page carries.
*/
feed: slices.map(slice => toLex(slice._feedPost)),
}
}
}
class MergeFeedSource_Custom extends MergeFeedSource {
agent: SessionAgent
client: Client
minDate: Date
feedUri: string
userInterests?: string
constructor({
agent,
client,
feedUri,
feedTuners,
userInterests,
}: {
agent: SessionAgent
client: Client
feedUri: string
feedTuners: FeedTunerFn[]
userInterests?: string
}) {
super({
agent,
client,
feedTuners,
})
this.agent = agent
this.client = client
this.feedUri = feedUri
this.userInterests = userInterests
this.sourceInfo = {
@@ -284,15 +310,16 @@ class MergeFeedSource_Custom extends MergeFeedSource {
protected async _getFeed(
cursor: string | undefined,
limit: number,
): Promise<AppBskyFeedGetTimeline.Response> {
): Promise<FeedPage> {
try {
const contentLangs = getContentLanguages().join(',')
const isBlueskyOwned = isBlueskyOwnedFeed(this.feedUri)
const res = await this.agent.app.bsky.feed.getFeed(
const res = await this.client.call(
app.bsky.feed.getFeed,
{
cursor,
limit,
feed: this.feedUri,
feed: this.feedUri as AtUriString,
},
{
headers: {
@@ -303,26 +330,25 @@ class MergeFeedSource_Custom extends MergeFeedSource {
},
},
)
let feed = res.feed
// NOTE
// some custom feeds fail to enforce the pagination limit
// so we manually truncate here
// -prf
if (limit && res.data.feed.length > limit) {
res.data.feed = res.data.feed.slice(0, limit)
if (limit && feed.length > limit) {
feed = feed.slice(0, limit)
}
// filter out older posts
res.data.feed = res.data.feed.filter(
post => new Date(post.post.indexedAt) > this.minDate,
)
feed = feed.filter(post => new Date(post.post.indexedAt) > this.minDate)
// attach source info
for (const post of res.data.feed) {
for (const post of feed) {
// @ts-ignore
post.__source = this.sourceInfo
}
return res
return {success: true, cursor: res.cursor, feed}
} catch {
// dont bubble custom-feed errors
return {success: false, headers: {}, data: {feed: []}}
return {success: false, feed: []}
}
}
}
+13 -21
View File
@@ -1,25 +1,22 @@
import {
type Agent,
type AppBskyFeedDefs,
type AppBskyFeedGetPosts,
} from '@atproto/api'
import {type Client} from '@atproto/lex-client'
import {logger} from '#/logger'
import {app} from '#/lexicons'
import {type FeedAPI, type FeedAPIResponse} from './types'
export class PostListFeedAPI implements FeedAPI {
agent: Agent
params: AppBskyFeedGetPosts.QueryParams
peek: AppBskyFeedDefs.FeedViewPost | null = null
client: Client
params: app.bsky.feed.getPosts.$Params
peek: app.bsky.feed.defs.FeedViewPost | null = null
constructor({
agent,
client,
feedParams,
}: {
agent: Agent
feedParams: AppBskyFeedGetPosts.QueryParams
client: Client
feedParams: app.bsky.feed.getPosts.$Params
}) {
this.agent = agent
this.client = client
if (feedParams.uris.length > 25) {
logger.warn(
`Too many URIs provided - expected 25, got ${feedParams.uris.length}`,
@@ -30,23 +27,18 @@ export class PostListFeedAPI implements FeedAPI {
}
}
async peekLatest(): Promise<AppBskyFeedDefs.FeedViewPost> {
async peekLatest(): Promise<app.bsky.feed.defs.FeedViewPost> {
if (this.peek) return this.peek
throw new Error('Has not fetched yet')
}
async fetch({}: {}): Promise<FeedAPIResponse> {
const res = await this.agent.app.bsky.feed.getPosts({
const res = await this.client.call(app.bsky.feed.getPosts, {
...this.params,
})
if (res.success) {
this.peek = {post: res.data.posts[0]}
return {
feed: res.data.posts.map(post => ({post})),
}
}
this.peek = {post: res.posts[0]}
return {
feed: [],
feed: res.posts.map(post => ({post})),
}
}
}
+3 -3
View File
@@ -1,12 +1,12 @@
import {type AppBskyFeedDefs} from '@atproto/api'
import {type app} from '#/lexicons'
export interface FeedAPIResponse {
cursor?: string
feed: AppBskyFeedDefs.FeedViewPost[]
feed: app.bsky.feed.defs.FeedViewPost[]
}
export interface FeedAPI {
peekLatest(): Promise<AppBskyFeedDefs.FeedViewPost>
peekLatest(): Promise<app.bsky.feed.defs.FeedViewPost>
fetch({
cursor,
limit,
+1 -1
View File
@@ -1,4 +1,4 @@
import {AtUri} from '@atproto/api'
import {AtUri} from '@atproto/syntax'
import {BSKY_FEED_OWNER_DIDS} from '#/lib/constants'
import {type UsePreferencesQueryResponse} from '#/state/queries/preferences'
+129 -190
View File
@@ -1,25 +1,14 @@
import {
type $Typed,
type AppBskyEmbedExternal,
type AppBskyEmbedGallery,
type AppBskyEmbedImages,
type AppBskyEmbedRecord,
type AppBskyEmbedRecordWithMedia,
type AppBskyEmbedVideo,
AppBskyFeedPost,
BlobRef,
ChatBskyGroupDefs,
type ComAtprotoLabelDefs,
type ComAtprotoRepoApplyWrites,
type ComAtprotoRepoStrongRef,
RichText,
} from '@atproto/api'
import {TID} from '@atproto/common-web'
import {type $Typed} from '@atproto/lex'
import {type Client} from '@atproto/lex-client'
import {
type AtUriString,
toDatetimeString,
type UriString,
} from '@atproto/syntax'
import {RichText} from '@bsky.app/sdk/richtext'
import {t} from '@lingui/core/macro'
import {type QueryClient} from '@tanstack/react-query'
import {sha256} from 'js-sha256'
import {CID} from 'multiformats/cid'
import * as Hasher from 'multiformats/hashes/hasher'
import {IMAGE_SIZE_CONFIG_POSTS} from '#/lib/constants'
import {isNetworkError} from '#/lib/strings/errors'
@@ -34,18 +23,33 @@ import {
createThreadgateRecord,
threadgateAllowUISettingToAllowRecordValue,
} from '#/state/queries/threadgate'
import {type SessionAgent} from '#/state/session'
import {
type EmbedDraft,
type PostDraft,
type ThreadDraft,
} from '#/view/com/composer/state/composer'
import {app, chat, com} from '#/lexicons'
import * as bsky from '#/types/bsky'
import {createGIFDescription} from '../gif-alt-text'
import {computeCid} from './computeCid'
import {type ResolveClients} from './resolve'
import {uploadBlob} from './upload-blob'
export {uploadBlob}
/**
* The lex clients the post pipeline needs. `pdsClient` handles writes to the
* user's repo (applyWrites, uploadBlob) - proxied to their PDS, never the
* appview. `appviewClient` handles the `app.bsky.*` reads (getPosts, when
* resolving a reply's root). `resolveClients` is threaded to the link/gif
* resolvers, which need the appview + chat + bridge agent (design section H).
*/
export type PostClients = {
pdsClient: Client
appviewClient: Client
resolveClients: ResolveClients
}
interface PostOpts {
thread: ThreadDraft
replyTo?: string
@@ -54,20 +58,21 @@ interface PostOpts {
}
export async function post(
agent: SessionAgent,
clients: PostClients,
queryClient: QueryClient,
opts: PostOpts,
) {
const {pdsClient, appviewClient, resolveClients} = clients
const thread = opts.thread
opts.onStateChange?.(t`Processing...`)
let replyPromise:
| Promise<AppBskyFeedPost.Record['reply']>
| AppBskyFeedPost.Record['reply']
| Promise<app.bsky.feed.post.Main['reply']>
| app.bsky.feed.post.Main['reply']
| undefined
if (opts.replyTo) {
// Not awaited to avoid waterfalls.
replyPromise = resolveReply(agent, opts.replyTo)
replyPromise = resolveReply(appviewClient, opts.replyTo)
}
// add top 3 languages from user preferences if langs is provided
@@ -76,8 +81,8 @@ export async function post(
langs = opts.langs.slice(0, 3)
}
const did = agent.assertDid
const writes: $Typed<ComAtprotoRepoApplyWrites.Create>[] = []
const did = pdsClient.assertDid
const writes: com.atproto.repo.applyWrites.$InputBody['writes'] = []
const uris: string[] = []
let now = new Date()
@@ -86,15 +91,23 @@ export async function post(
for (let i = 0; i < thread.posts.length; i++) {
const draft = thread.posts[i]
// Not awaited to avoid waterfalls.
const rtPromise = resolveRT(agent, draft.richtext)
/*
* Not awaited to avoid waterfalls. `draft.richtext` is still an
* `@atproto/api` RichText (composer state is migrated by Task 7); the SDK
* RichText is structurally the same transplant, so bridge it here.
* TODO(phase4): drop toLex once composer state migrates to SDK RichText.
*/
const rtPromise = resolveRT(
resolveClients.appview,
bsky.toLex<RichText>(draft.richtext),
)
const embedPromise = resolveEmbed(
agent,
clients,
queryClient,
draft,
opts.onStateChange,
)
let labels: $Typed<ComAtprotoLabelDefs.SelfLabels> | undefined
let labels: $Typed<com.atproto.label.defs.SelfLabels> | undefined
if (draft.labels.length) {
labels = {
$type: 'com.atproto.label.defs#selfLabels',
@@ -107,17 +120,17 @@ export async function post(
now.setMilliseconds(now.getMilliseconds() + 1)
tid = TID.next(tid)
const rkey = tid.toString()
const uri = `at://${did}/app.bsky.feed.post/${rkey}`
const uri = `at://${did}/app.bsky.feed.post/${rkey}` as AtUriString
uris.push(uri)
const rt = await rtPromise
const embed = await embedPromise
const reply = await replyPromise
const record: AppBskyFeedPost.Record = {
const record: app.bsky.feed.post.Main = {
// IMPORTANT: $type has to exist, CID is calculated with the `$type` field
// present and will produce the wrong CID if you omit it.
$type: 'app.bsky.feed.post',
createdAt: now.toISOString(),
createdAt: toDatetimeString(now),
text: rt.text,
facets: rt.facets,
reply,
@@ -125,11 +138,17 @@ export async function post(
langs,
labels,
}
/*
* `value` is typed `LexMap` (a loose index-signature map). A generated
* record type is a valid LexMap at runtime but a strict interface is not
* assignable to an index-signature type, so widen with `toLex` at this
* boundary. Not a brand cast - the record is already fully typed.
*/
writes.push({
$type: 'com.atproto.repo.applyWrites#create',
collection: 'app.bsky.feed.post',
rkey: rkey,
value: record,
value: bsky.toLex(record),
})
if (i === 0 && thread.threadgate.some(tg => tg.type !== 'everybody')) {
@@ -137,11 +156,14 @@ export async function post(
$type: 'com.atproto.repo.applyWrites#create',
collection: 'app.bsky.feed.threadgate',
rkey: rkey,
value: createThreadgateRecord({
createdAt: now.toISOString(),
post: uri,
allow: threadgateAllowUISettingToAllowRecordValue(thread.threadgate),
}),
value: bsky.toLex(
createThreadgateRecord({
post: uri,
allow: threadgateAllowUISettingToAllowRecordValue(
thread.threadgate,
),
}),
),
})
}
@@ -153,12 +175,12 @@ export async function post(
$type: 'com.atproto.repo.applyWrites#create',
collection: 'app.bsky.feed.postgate',
rkey: rkey,
value: {
value: bsky.toLex({
...thread.postgate,
$type: 'app.bsky.feed.postgate',
createdAt: now.toISOString(),
post: uri,
},
}),
})
}
@@ -174,8 +196,8 @@ export async function post(
}
try {
await agent.com.atproto.repo.applyWrites({
repo: agent.assertDid,
await pdsClient.call(com.atproto.repo.applyWrites, {
repo: did,
writes: writes,
validate: true,
})
@@ -196,14 +218,18 @@ export async function post(
return {uris}
}
async function resolveRT(agent: SessionAgent, richtext: RichText) {
async function resolveRT(appviewClient: Client, richtext: RichText) {
const trimmedText = richtext.text
// Trim leading whitespace-only lines (but don't break ASCII art).
.replace(/^(\s*\n)+/, '')
// Trim any trailing whitespace.
.trimEnd()
let rt = new RichText({text: trimmedText}, {cleanNewlines: true})
await rt.detectFacets(agent)
/*
* Facet detection resolves handles via `com.atproto.identity.resolveHandle`,
* which the appview client serves (design section B).
*/
await rt.detectFacets(appviewClient)
rt = shortenLinks(rt)
rt = stripInvalidMentions(rt)
@@ -216,9 +242,9 @@ export class ReplyDeletedError extends Error {
}
}
async function resolveReply(agent: SessionAgent, replyTo: string) {
const {data} = await agent.app.bsky.feed.getPosts({
uris: [replyTo],
async function resolveReply(appviewClient: Client, replyTo: string) {
const data = await appviewClient.call(app.bsky.feed.getPosts, {
uris: [replyTo as AtUriString],
})
const parentPost = data.posts[0]
if (!parentPost) {
@@ -229,14 +255,9 @@ async function resolveReply(agent: SessionAgent, replyTo: string) {
uri: parentPost.uri,
cid: parentPost.cid,
}
let rootRef = parentRef
let rootRef: com.atproto.repo.strongRef.Main = parentRef
if (
bsky.dangerousIsType<AppBskyFeedPost.Record>(
parentPost.record,
AppBskyFeedPost.isRecord,
)
) {
if (bsky.isType(app.bsky.feed.post, parentPost.record)) {
if (parentPost.record.reply) {
rootRef = parentPost.record.reply.root
}
@@ -249,23 +270,15 @@ async function resolveReply(agent: SessionAgent, replyTo: string) {
}
async function resolveEmbed(
agent: SessionAgent,
clients: PostClients,
queryClient: QueryClient,
draft: PostDraft,
onStateChange: ((state: string) => void) | undefined,
): Promise<
| $Typed<AppBskyEmbedImages.Main>
| $Typed<AppBskyEmbedGallery.Main>
| $Typed<AppBskyEmbedVideo.Main>
| $Typed<AppBskyEmbedExternal.Main>
| $Typed<AppBskyEmbedRecord.Main>
| $Typed<AppBskyEmbedRecordWithMedia.Main>
| undefined
> {
): Promise<app.bsky.feed.post.Main['embed']> {
if (draft.embed.quote) {
const [resolvedMedia, resolvedQuote] = await Promise.all([
resolveMedia(agent, queryClient, draft.embed, onStateChange),
resolveRecord(agent, queryClient, draft.embed.quote.uri),
resolveMedia(clients, queryClient, draft.embed, onStateChange),
resolveRecord(clients, queryClient, draft.embed.quote.uri),
])
if (resolvedMedia) {
return {
@@ -283,7 +296,7 @@ async function resolveEmbed(
}
}
const resolvedMedia = await resolveMedia(
agent,
clients,
queryClient,
draft.embed,
onStateChange,
@@ -294,7 +307,7 @@ async function resolveEmbed(
if (draft.embed.link) {
const resolvedLink = await fetchResolveLinkQuery(
queryClient,
agent,
clients.resolveClients,
draft.embed.link.uri,
)
if (resolvedLink.type === 'record') {
@@ -308,24 +321,25 @@ async function resolveEmbed(
}
async function resolveMedia(
agent: SessionAgent,
clients: PostClients,
queryClient: QueryClient,
embedDraft: EmbedDraft,
onStateChange: ((state: string) => void) | undefined,
): Promise<
| $Typed<AppBskyEmbedExternal.Main>
| $Typed<AppBskyEmbedImages.Main>
| $Typed<AppBskyEmbedGallery.Main>
| $Typed<AppBskyEmbedVideo.Main>
| $Typed<app.bsky.embed.external.Main>
| $Typed<app.bsky.embed.images.Main>
| $Typed<app.bsky.embed.gallery.Main>
| $Typed<app.bsky.embed.video.Main>
| undefined
> {
const {pdsClient, resolveClients} = clients
if (embedDraft.media?.type === 'images') {
const imagesDraft = embedDraft.media.images
logger.debug(`Uploading images`, {
count: imagesDraft.length,
})
onStateChange?.(t`Uploading images...`)
const images: AppBskyEmbedImages.Image[] = await Promise.all(
const images: app.bsky.embed.images.Image[] = await Promise.all(
imagesDraft.map(async (image, i) => {
logger.debug(`Compressing image #${i}`)
const {path, width, height, mime} = await compressImage(
@@ -333,9 +347,9 @@ async function resolveMedia(
IMAGE_SIZE_CONFIG_POSTS,
)
logger.debug(`Uploading image #${i}`)
const res = await uploadBlob(agent, path, mime)
const res = await uploadBlob(pdsClient, path, mime)
return {
image: res.data.blob,
image: res.blob,
alt: image.alt,
aspectRatio: {width, height},
}
@@ -352,7 +366,7 @@ async function resolveMedia(
count: imagesDraft.length,
})
onStateChange?.(t`Uploading images...`)
const items: $Typed<AppBskyEmbedGallery.Image>[] = await Promise.all(
const items: $Typed<app.bsky.embed.gallery.Image>[] = await Promise.all(
imagesDraft.map(async (image, i) => {
logger.debug(`Compressing image #${i}`)
const {path, width, height, mime} = await compressImage(
@@ -360,10 +374,10 @@ async function resolveMedia(
IMAGE_SIZE_CONFIG_POSTS,
)
logger.debug(`Uploading image #${i}`)
const res = await uploadBlob(agent, path, mime)
const res = await uploadBlob(pdsClient, path, mime)
return {
$type: 'app.bsky.embed.gallery#image' as const,
image: res.data.blob,
image: res.blob,
alt: image.alt,
aspectRatio: {width, height},
}
@@ -383,10 +397,10 @@ async function resolveMedia(
videoDraft.captions
.filter(caption => caption.lang !== '')
.map(async caption => {
const {data} = await agent.uploadBlob(caption.file, {
const res = await pdsClient.uploadBlob(caption.file, {
encoding: 'text/vtt',
})
return {lang: caption.lang, file: data.blob}
return {lang: caption.lang, file: res.body.blob}
}),
)
@@ -406,7 +420,13 @@ async function resolveMedia(
return {
$type: 'app.bsky.embed.video',
video: videoDraft.pendingPublish.blobRef,
/*
* The video blob is a legacy `@atproto/api` BlobRef from the not-yet
* -migrated video pipeline (getJobStatus, in composer state/video). Its
* structural shape matches the lexicon blob field; the CID hasher handles
* both class instances and plain lex blobs (see computeCid).
*/
video: bsky.toLex(videoDraft.pendingPublish.blobRef),
alt: videoDraft.altText || undefined,
captions: captions.length === 0 ? undefined : captions,
aspectRatio,
@@ -416,22 +436,18 @@ async function resolveMedia(
}
if (embedDraft.media?.type === 'gif') {
const gifDraft = embedDraft.media
const resolvedGif = await fetchResolveGifQuery(
queryClient,
agent,
gifDraft.gif,
)
let blob: BlobRef | undefined
const resolvedGif = await fetchResolveGifQuery(queryClient, gifDraft.gif)
let blob: app.bsky.embed.external.External['thumb']
if (resolvedGif.thumb) {
onStateChange?.(t`Uploading link thumbnail...`)
const {path, mime} = resolvedGif.thumb.source
const response = await uploadBlob(agent, path, mime)
blob = response.data.blob
const response = await uploadBlob(pdsClient, path, mime)
blob = response.blob
}
return {
$type: 'app.bsky.embed.external',
external: {
uri: resolvedGif.uri,
uri: resolvedGif.uri as UriString,
title: resolvedGif.title,
description: createGIFDescription(resolvedGif.title, gifDraft.alt),
thumb: blob,
@@ -441,36 +457,42 @@ async function resolveMedia(
if (embedDraft.link) {
const resolvedLink = await fetchResolveLinkQuery(
queryClient,
agent,
resolveClients,
embedDraft.link.uri,
)
if (resolvedLink.type === 'external') {
let blob: BlobRef | undefined
let blob: app.bsky.embed.external.External['thumb']
if (resolvedLink.thumb) {
onStateChange?.(t`Uploading link thumbnail...`)
const {path, mime} = resolvedLink.thumb.source
const response = await uploadBlob(agent, path, mime)
blob = response.data.blob
const response = await uploadBlob(pdsClient, path, mime)
blob = response.blob
}
return {
$type: 'app.bsky.embed.external',
external: {
uri: resolvedLink.uri,
uri: resolvedLink.uri as UriString,
title: resolvedLink.title,
description: resolvedLink.description,
thumb: blob,
associatedRefs: resolvedLink.associatedRefs,
/*
* associatedRefs comes from getLinkMeta, still typed with the old
* `@atproto/api` StrongRef (link-meta.ts is out of scope). The shape
* is identical; bridge with toLex.
* TODO(phase4): drop toLex once link-meta migrates.
*/
associatedRefs: bsky.toLex(resolvedLink.associatedRefs),
},
}
}
if (
resolvedLink.type === 'chat-invite' &&
ChatBskyGroupDefs.isJoinLinkPreviewView(resolvedLink.view)
bsky.isType(chat.bsky.group.defs.joinLinkPreviewView, resolvedLink.view)
) {
return {
$type: 'app.bsky.embed.external',
external: {
uri: resolvedLink.uri,
uri: resolvedLink.uri as UriString,
title: resolvedLink.view.name,
description: `${resolvedLink.view.memberCount}/${resolvedLink.view.memberLimit}`,
},
@@ -481,100 +503,17 @@ async function resolveMedia(
}
async function resolveRecord(
agent: SessionAgent,
clients: PostClients,
queryClient: QueryClient,
uri: string,
): Promise<ComAtprotoRepoStrongRef.Main> {
const resolvedLink = await fetchResolveLinkQuery(queryClient, agent, uri)
): Promise<com.atproto.repo.strongRef.Main> {
const resolvedLink = await fetchResolveLinkQuery(
queryClient,
clients.resolveClients,
uri,
)
if (resolvedLink.type !== 'record') {
throw Error(t`Expected uri to resolve to a record`)
}
return resolvedLink.record
}
// The built-in hashing functions from multiformats (`multiformats/hashes/sha2`)
// are meant for Node.js, this is the cross-platform equivalent.
const mf_sha256 = Hasher.from({
name: 'sha2-256',
code: 0x12,
encode: input => {
const digest = sha256.arrayBuffer(input)
return new Uint8Array(digest)
},
})
async function computeCid(record: AppBskyFeedPost.Record): Promise<string> {
/*
* Lazily loaded since it's only needed when posting a thread, and its
* `cborg` dependency is ~190KB that would otherwise be in the initial
* web bundle.
*/
const dcbor = await import('@ipld/dag-cbor')
// IMPORTANT: `prepareObject` prepares the record to be hashed by removing
// fields with undefined value, and converting BlobRef instances to the
// right IPLD representation.
const prepared = prepareForHashing(record)
// 1. Encode the record into DAG-CBOR format
const encoded = dcbor.encode(prepared)
// 2. Hash the record in SHA-256 (code 0x12)
const digest = await mf_sha256.digest(encoded)
// 3. Create a CIDv1, specifying DAG-CBOR as content (code 0x71)
const cid = CID.createV1(0x71, digest)
// 4. Get the Base32 representation of the CID (`b` prefix)
return cid.toString()
}
// Returns a transformed version of the object for use in DAG-CBOR.
// eslint-disable-next-line @typescript-eslint/no-explicit-any
function prepareForHashing(v: any): any {
// IMPORTANT: BlobRef#ipld() returns the correct object we need for hashing,
// the API client will convert this for you but we're hashing in the client,
// so we need it *now*.
if (v instanceof BlobRef) {
return v.ipld()
}
// Walk through arrays
if (Array.isArray(v)) {
let pure = true
const mapped = v.map(value => {
if (value !== (value = prepareForHashing(value))) {
pure = false
}
return value
})
return pure ? v : mapped
}
// Walk through plain objects
if (isPlainObject(v)) {
const obj: Record<string, unknown> = {}
let pure = true
for (const key in v) {
// eslint-disable-next-line @typescript-eslint/no-unsafe-member-access
let value = v[key]
// `value` is undefined
if (value === undefined) {
pure = false
continue
}
// `prepareObject` returned a value that's different from what we had before
if (value !== (value = prepareForHashing(value))) {
pure = false
}
obj[key] = value
}
// Return as is if we haven't needed to tamper with anything
return pure ? v : obj
}
return v
}
// eslint-disable-next-line @typescript-eslint/no-explicit-any
function isPlainObject(v: any): boolean {
if (typeof v !== 'object' || v === null) {
return false
}
const proto = Object.getPrototypeOf(v)
return proto === Object.prototype || proto === null
}
+63 -49
View File
@@ -1,11 +1,7 @@
import {
type AppBskyFeedDefs,
type AppBskyGraphDefs,
type ComAtprotoRepoStrongRef,
} from '@atproto/api'
import {AtUri} from '@atproto/api'
import {type Client} from '@atproto/lex-client'
import {AtUri, type AtUriString, type HandleString} from '@atproto/syntax'
import {DM_SERVICE_HEADERS, IMAGE_SIZE_CONFIG_2K_1MB} from '#/lib/constants'
import {IMAGE_SIZE_CONFIG_2K_1MB} from '#/lib/constants'
import {getLinkMeta, type LinkMeta} from '#/lib/link-meta/link-meta'
import {resolveShortLink} from '#/lib/link-meta/resolve-short-link'
import {downloadAndResize} from '#/lib/media/manip'
@@ -29,8 +25,21 @@ import {createComposerImage} from '#/state/gallery'
import {type ChatInvitePreview} from '#/state/queries/join-links'
import {type SessionAgent} from '#/state/session'
import {type Gif} from '#/features/gifPicker/types'
import {app, chat, com} from '#/lexicons'
import {createGIFDescription} from '../gif-alt-text'
/**
* Clients the link resolver needs. `appview` serves the `app.bsky.*` feed/graph
* reads plus handle resolution; `chat` serves the group join-link previews.
* `agent` is still threaded through to {@link getLinkMeta}, which has not yet
* migrated off the bridge - it only reads the (unused) service URL from it.
*/
export type ResolveClients = {
appview: Client
chat: Client
agent: SessionAgent
}
type ResolvedExternalLink = {
type: 'external'
uri: string
@@ -47,30 +56,30 @@ type ResolvedExternalLink = {
type ResolvedPostRecord = {
type: 'record'
record: ComAtprotoRepoStrongRef.Main
record: com.atproto.repo.strongRef.Main
kind: 'post'
view: AppBskyFeedDefs.PostView
view: app.bsky.feed.defs.PostView
}
type ResolvedFeedRecord = {
type: 'record'
record: ComAtprotoRepoStrongRef.Main
record: com.atproto.repo.strongRef.Main
kind: 'feed'
view: AppBskyFeedDefs.GeneratorView
view: app.bsky.feed.defs.GeneratorView
}
type ResolvedListRecord = {
type: 'record'
record: ComAtprotoRepoStrongRef.Main
record: com.atproto.repo.strongRef.Main
kind: 'list'
view: AppBskyGraphDefs.ListView
view: app.bsky.graph.defs.ListView
}
type ResolvedStarterPackRecord = {
type: 'record'
record: ComAtprotoRepoStrongRef.Main
record: com.atproto.repo.strongRef.Main
kind: 'starter-pack'
view: AppBskyGraphDefs.StarterPackView
view: app.bsky.graph.defs.StarterPackView
}
type ResolvedChatInvite = {
@@ -95,9 +104,10 @@ export class EmbeddingDisabledError extends Error {
}
export async function resolveLink(
agent: SessionAgent,
clients: ResolveClients,
uri: string,
): Promise<ResolvedLink> {
const {appview, chat: chatClient} = clients
if (isShortLink(uri)) {
uri = await resolveShortLink(uri)
}
@@ -124,15 +134,15 @@ export async function resolveLink(
const [_0, handleOrDid, _1, rkey] = uri.split('/').filter(Boolean)
const did = await fetchDid(handleOrDid)
const feed = makeRecordUri(did, 'app.bsky.feed.generator', rkey)
const res = await agent.app.bsky.feed.getFeedGenerator({feed})
const {view} = await appview.call(app.bsky.feed.getFeedGenerator, {feed})
return {
type: 'record',
record: {
uri: res.data.view.uri,
cid: res.data.view.cid,
uri: view.uri,
cid: view.cid,
},
kind: 'feed',
view: res.data.view,
view,
}
}
if (isBskyListUrl(uri)) {
@@ -140,28 +150,27 @@ export async function resolveLink(
const [_0, handleOrDid, _1, rkey] = uri.split('/').filter(Boolean)
const did = await fetchDid(handleOrDid)
const list = makeRecordUri(did, 'app.bsky.graph.list', rkey)
const res = await agent.app.bsky.graph.getList({list})
const res = await appview.call(app.bsky.graph.getList, {list})
return {
type: 'record',
record: {
uri: res.data.list.uri,
cid: res.data.list.cid,
uri: res.list.uri,
cid: res.list.cid,
},
kind: 'list',
view: res.data.list,
view: res.list,
}
}
const chatInviteCode = getChatInviteCodeFromUrl(uri)
if (chatInviteCode) {
const res = await agent.chat.bsky.group.getJoinLinkPreviews(
{codes: [chatInviteCode]},
{headers: DM_SERVICE_HEADERS},
)
const res = await chatClient.call(chat.bsky.group.getJoinLinkPreviews, {
codes: [chatInviteCode],
})
return {
type: 'chat-invite',
uri,
code: chatInviteCode,
view: res.data.joinLinkPreviews[0],
view: res.joinLinkPreviews[0],
}
}
if (isBskyStartUrl(uri) || isBskyStarterPackUrl(uri)) {
@@ -173,34 +182,35 @@ export async function resolveLink(
}
const did = await fetchDid(parsed.name)
const starterPack = createStarterPackUri({did, rkey: parsed.rkey})
const res = await agent.app.bsky.graph.getStarterPack({starterPack})
const res = await appview.call(app.bsky.graph.getStarterPack, {
starterPack: starterPack as AtUriString,
})
return {
type: 'record',
record: {
uri: res.data.starterPack.uri,
cid: res.data.starterPack.cid,
uri: res.starterPack.uri,
cid: res.starterPack.cid,
},
kind: 'starter-pack',
view: res.data.starterPack,
view: res.starterPack,
}
}
return resolveExternal(agent, uri)
return resolveExternal(clients, uri)
// Forked from useGetPost. TODO: move into RQ.
async function getPost({uri}: {uri: string}) {
const urip = new AtUri(uri)
if (!urip.host.startsWith('did:')) {
const res = await agent.resolveHandle({
handle: urip.host,
const {did} = await appview.call(com.atproto.identity.resolveHandle, {
handle: urip.host as HandleString,
})
// @ts-expect-error TODO new-sdk-migration
urip.host = res.data.did
urip.host = did
}
const res = await agent.getPosts({
const res = await appview.call(app.bsky.feed.getPosts, {
uris: [urip.toString()],
})
if (res.success && res.data.posts[0]) {
return res.data.posts[0]
if (res.posts[0]) {
return res.posts[0]
}
throw new Error('getPost: post not found')
}
@@ -209,17 +219,16 @@ export async function resolveLink(
async function fetchDid(handleOrDid: string) {
let identifier = handleOrDid
if (!identifier.startsWith('did:')) {
const res = await agent.resolveHandle({handle: identifier})
identifier = res.data.did
const {did} = await appview.call(com.atproto.identity.resolveHandle, {
handle: identifier as HandleString,
})
identifier = did
}
return identifier
}
}
export async function resolveGif(
agent: SessionAgent,
gif: Gif,
): Promise<ResolvedExternalLink> {
export async function resolveGif(gif: Gif): Promise<ResolvedExternalLink> {
const gifUrl = gif.media_formats.gif.url
const params = new URLSearchParams()
params.set('hh', String(gif.media_formats.gif.dims[1]))
@@ -259,10 +268,15 @@ function getFileSlug(url: string | undefined): string | undefined {
}
async function resolveExternal(
agent: SessionAgent,
clients: ResolveClients,
uri: string,
): Promise<ResolvedExternalLink> {
const result = await getLinkMeta(agent, uri)
/*
* getLinkMeta still takes the bridge agent (not yet migrated); it only reads
* a service URL from it that LINK_META_PROXY ignores. Keep threading the
* agent here until link-meta migrates.
*/
const result = await getLinkMeta(clients.agent, uri)
return {
type: 'external',
uri: result.url,
+32 -8
View File
@@ -1,38 +1,62 @@
import {copyAsync} from 'expo-file-system/legacy'
import {type Agent, type ComAtprotoRepoUploadBlob} from '@atproto/api'
import {type BlobRef} from '@atproto/lex'
import {type Client, type EncodingString} from '@atproto/lex-client'
import {safeDeleteAsync} from '#/lib/media/manip'
/**
* @param encoding Allows overriding the blob's type
* The blob-upload response body: `{blob}`. lex `Client.uploadBlob` returns the
* full XRPC response, so callers read `res.body.blob` (the parsed blob ref).
*/
type UploadBlobResult = {blob: BlobRef}
/**
* @param encoding Allows overriding the blob's type. Passed as the lex upload
* option (NEVER a content-type header - lex-client throws if the encoding is
* set via headers).
*/
export async function uploadBlob(
agent: Agent,
client: Client,
input: string | Blob,
encoding?: string,
): Promise<ComAtprotoRepoUploadBlob.Response> {
): Promise<UploadBlobResult> {
if (typeof input === 'string' && input.startsWith('file:')) {
const blob = await asBlob(input)
return agent.uploadBlob(blob, {encoding})
return uploadBlobResult(client, blob, encoding)
}
if (typeof input === 'string' && input.startsWith('/')) {
const blob = await asBlob(`file://${input}`)
return agent.uploadBlob(blob, {encoding})
return uploadBlobResult(client, blob, encoding)
}
if (typeof input === 'string' && input.startsWith('data:')) {
const blob = await fetch(input).then(r => r.blob())
return agent.uploadBlob(blob, {encoding})
return uploadBlobResult(client, blob, encoding)
}
if (input instanceof Blob) {
return agent.uploadBlob(input, {encoding})
return uploadBlobResult(client, input, encoding)
}
throw new TypeError(`Invalid uploadBlob input: ${typeof input}`)
}
async function uploadBlobResult(
client: Client,
blob: Blob,
encoding?: string,
): Promise<UploadBlobResult> {
/*
* The lex encoding option is a branded mime string (`${string}/${string}`);
* callers pass a plain mime string, so assert the brand here.
*/
const res = await client.uploadBlob(blob, {
encoding: encoding as EncodingString | undefined,
})
return {blob: res.body.blob}
}
async function asBlob(uri: string): Promise<Blob> {
return withSafeFile(uri, async safeUri => {
// Note
+22 -7
View File
@@ -1,28 +1,43 @@
import {type Agent, type ComAtprotoRepoUploadBlob} from '@atproto/api'
import {type BlobRef} from '@atproto/lex'
import {type Client, type EncodingString} from '@atproto/lex-client'
/**
* The blob-upload response body: `{blob}`. lex `Client.uploadBlob` returns the
* full XRPC response, so callers read `res.body.blob` (the parsed blob ref).
*/
type UploadBlobResult = {blob: BlobRef}
/**
* @note It is recommended, on web, to use the `file` instance of the file
* selector input element, rather than a `data:` URL, to avoid
* loading the file into memory. `File` extends `Blob` "file" instances can
* be passed directly to this function.
*
* @param encoding Passed as the lex upload option (NEVER a content-type header
* - lex-client throws if the encoding is set via headers).
*/
export async function uploadBlob(
agent: Agent,
client: Client,
input: string | Blob,
encoding?: string,
): Promise<ComAtprotoRepoUploadBlob.Response> {
): Promise<UploadBlobResult> {
/*
* The lex encoding option is a branded mime string (`${string}/${string}`);
* callers pass a plain mime string, so assert the brand here.
*/
const enc = encoding as EncodingString | undefined
if (
typeof input === 'string' &&
(input.startsWith('data:') || input.startsWith('blob:'))
) {
const blob = await fetch(input).then(r => r.blob())
return agent.uploadBlob(blob, {encoding})
const res = await client.uploadBlob(blob, {encoding: enc})
return {blob: res.body.blob}
}
if (input instanceof Blob) {
return agent.uploadBlob(input, {
encoding,
})
const res = await client.uploadBlob(input, {encoding: enc})
return {blob: res.body.blob}
}
throw new TypeError(`Invalid uploadBlob input: ${typeof input}`)