diff --git a/sync/events.ts b/sync/events.ts index a7dcefa..831f49f 100644 --- a/sync/events.ts +++ b/sync/events.ts @@ -4,8 +4,23 @@ import type { RepoRecord } from "@atp/lexicon"; import type { BlockMap } from "@atp/repo"; import type { AtUri } from "@atp/syntax"; +/** Broad sync event type for all sync events */ export type Event = CommitEvt | SyncEvt | IdentityEvt | AccountEvt; +/** + * Metadata for a {@link CommitEvt} + * @prop seq + * Event Sequence Number + * @see {@link https://atproto.com/specs/event-stream#sequence-numbers} + * @prop time Time of the Commit Event + * @prop commit CID of the Commit + * @prop blocks CAR "slice" for the corresponding repo diff + * @prop rev Repo revision identifier as a TID + * @prop uri AT URI of the record committed + * @prop did DID of the repository + * @prop collection Collection (lexicon) of the record + * @prop rkey Record Key of the record + */ export type CommitMeta = { seq: number; time: string; @@ -18,24 +33,40 @@ export type CommitMeta = { rkey: string; }; +/** {@link Event} for all commit events */ export type CommitEvt = Create | Update | Delete; +/** {@link CommitEvt} for record creation */ export type Create = CommitMeta & { event: "create"; record: RepoRecord; cid: CID; }; +/** {@link CommitEvt} for record updates/edits */ export type Update = CommitMeta & { event: "update"; record: RepoRecord; cid: CID; }; +/** {@link CommitEvt} for record deletions */ export type Delete = CommitMeta & { event: "delete"; }; +/** + * {@link Event} for repository sync events + * @prop seq + * Event Sequence Number + * @see {@link https://atproto.com/specs/event-stream#sequence-numbers} + * @prop time Time of sync event + * @prop event Type of event + * @prop did Repository of event + * @prop cid CID of event + * @prop rev Repository revision identifier as a TID + * @prop blocks CAR "slice" for the corresponding repo diff + */ export type SyncEvt = { seq: number; time: string; @@ -46,6 +77,17 @@ export type SyncEvt = { blocks: BlockMap; }; +/** + * {@link Event} for identity change events + * @prop seq + * Event Sequence Number + * @see {@link https://atproto.com/specs/event-stream#sequence-numbers} + * @prop time Time of sync event + * @prop event Type of event + * @prop did Repository of event + * @prop handle Handle corresponding to DID + * @prop didDocument DID Document corresponding to DID + */ export type IdentityEvt = { seq: number; time: string; @@ -55,6 +97,16 @@ export type IdentityEvt = { didDocument?: DidDocument; }; +/** + * @prop seq + * Event Sequence Number + * @see {@link https://atproto.com/specs/event-stream#sequence-numbers} + * @prop time Time of sync event + * @prop event Type of event + * @prop did Repository of event + * @prop active Whether account has been activated or is deactivated + * @prop status Current Account Status of the repository + */ export type AccountEvt = { seq: number; time: string; @@ -64,6 +116,7 @@ export type AccountEvt = { status?: AccountStatus; }; +/** Upstream status of an account */ export type AccountStatus = | "takendown" | "suspended" diff --git a/sync/firehose/index.ts b/sync/firehose/index.ts index 3ec9795..bf46ad5 100644 --- a/sync/firehose/index.ts +++ b/sync/firehose/index.ts @@ -327,6 +327,14 @@ export class Firehose { } } +/** + * Parse a {@link Commit} object while authenticating the commit + * @param idResolver Identity resolver for DIDs and handles + * @param evt Commit event object to parse + * @param matchCollection Lexicon collection to match record to + * @param forceKeyRefresh Whether to force a refresh when resolving AT Protocol Key + * @returns A parsed authenticated commit + */ export const parseCommitAuthenticated = async ( idResolver: IdResolver, evt: Commit, @@ -372,6 +380,12 @@ export const parseCommitAuthenticated = async ( }); }; +/** + * Parse a {@link Commit} object without authenticating the commit + * @param evt Commit event object to parse + * @param matchCollection Lexicon collection to match record to + * @returns A parsed commit + */ export const parseCommitUnauthenticated = ( evt: Commit, matchCollection?: ((col: string) => boolean) | null, @@ -439,6 +453,10 @@ const formatCommitOps = async ( return evts; }; +/** + * Parse {@link Sync} object to a sync event + * @param evt Sync event to parse + */ export const parseSync = async (evt: Sync): Promise => { const car = await readCarWithRoot(evt.blocks); @@ -453,6 +471,12 @@ export const parseSync = async (evt: Sync): Promise => { }; }; +/** + * Parse and authenticate an identity event + * @param idResolver DID and handle resolver for authentication + * @param evt Identity event to parse + * @param unauthenticated If true authentication is skipped + */ export const parseIdentity = async ( idResolver: IdResolver, evt: Identity, @@ -486,6 +510,10 @@ const verifyHandle = async ( return res === did ? handle : undefined; }; +/** + * Parse an account event + * @param evt Account event to parse + */ export const parseAccount = (evt: Account): AccountEvt | undefined => { if (evt.status && !isValidStatus(evt.status)) return; return { diff --git a/sync/runner/consecutive-list.ts b/sync/runner/consecutive-list.ts index 219854c..cdc25aa 100644 --- a/sync/runner/consecutive-list.ts +++ b/sync/runner/consecutive-list.ts @@ -1,8 +1,10 @@ /** * Add items to a list, and mark those items as * completed. Upon item completion, get list of consecutive - * items completed at the head of the list. Example: + * items completed at the head of the list. * + * @example Get consecultive item list + * ```typescript * const consecutive = new ConsecutiveList() * const item1 = consecutive.push(1) * const item2 = consecutive.push(2) @@ -10,6 +12,7 @@ * item2.complete() // [] * item1.complete() // [1, 2] * item3.complete() // [3] + * ``` */ export class ConsecutiveList { list: ConsecutiveItem[] = []; @@ -29,6 +32,7 @@ export class ConsecutiveList { } } +/** Process being run consecutively in a {@link ConsecutiveList} */ export class ConsecutiveItem { isComplete = false; constructor( diff --git a/sync/runner/memory-runner.ts b/sync/runner/memory-runner.ts index 686deeb..b4f9325 100644 --- a/sync/runner/memory-runner.ts +++ b/sync/runner/memory-runner.ts @@ -2,6 +2,13 @@ import PQueue from "p-queue"; import { ConsecutiveList } from "./consecutive-list.ts"; import type { EventRunner } from "./types.ts"; +/** + * Options for {@link MemoryRunner} + * @param setCursor Method to save the current cursor + * @param concurrency Maximum amount of concurrent events being processed + * @param startCursor Starting Cursor for filling in downtime + * @param setCursorInterval Interval on which to run setCursor + */ export type MemoryRunnerOptions = { setCursor?: (cursor: number) => Promise; concurrency?: number; diff --git a/sync/runner/types.ts b/sync/runner/types.ts index 762ea11..df54b4a 100644 --- a/sync/runner/types.ts +++ b/sync/runner/types.ts @@ -1,3 +1,7 @@ +/** + * Generic event runner interface + * for event tracking and processing + */ export interface EventRunner { getCursor(): Awaited; trackEvent(