diff --git a/packages/consumer/mod.ts b/packages/consumer/mod.ts index 70bb5d9..46285ee 100644 --- a/packages/consumer/mod.ts +++ b/packages/consumer/mod.ts @@ -1,9 +1,13 @@ import { produceRequirements } from "@cistern/shared"; import { generateKeys } from "@cistern/crypto"; import { generateRandomName } from "@puregarlic/randimal"; +import { parse } from "@atcute/lexicons"; import type { Did } from "@atcute/lexicons/syntax"; import type { Client, CredentialManager } from "@atcute/client"; -import type { AppCisternLexiconPubkey } from "@cistern/lexicon"; +import { + AppCisternLexiconItem, + type AppCisternLexiconPubkey, +} from "@cistern/lexicon"; import type { ConsumerOptions, ConsumerParams, LocalKeyPair } from "./types.ts"; import type {} from "@atcute/atproto"; @@ -80,10 +84,39 @@ export class Consumer { } /** - * Returns an async iterator that returns pages of the user's items from their PDS. - * @todo List items from repo + * Asynchronously iterate through items in the user's PDS */ - async listItems() {} + async *listItems(): AsyncIterator< + AppCisternLexiconItem.Main, + void, + undefined + > { + let cursor: string | undefined; + + while (true) { + const res = await this.rpc.get("com.atproto.repo.listRecords", { + params: { + collection: "app.cistern.lexicon.item", + repo: this.did, + cursor, + }, + }); + + if (!res.ok) { + throw new Error( + `failed to list items: ${res.status} ${res.data.error}`, + ); + } + + if (res.data.cursor) cursor = res.data.cursor; + + for (const record of res.data.records) { + yield parse(AppCisternLexiconItem.mainSchema, record.value); + } + + if (!cursor) return; + } + } /** * Subscribes to the Jetstreams for the user's items.