From 46073a5729e7d43ae93dffd3ba5732803c93d44d Mon Sep 17 00:00:00 2001 From: Graham Barber Date: Sun, 9 Nov 2025 14:11:57 -0800 Subject: [PATCH] feat(consumer): list items from PDS --- packages/consumer/mod.ts | 41 ++++++++++++++++++++++++++++++++++++---- 1 file changed, 37 insertions(+), 4 deletions(-) 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. -- 2.51.2