Readers and pagination
A reader is an async generator that fetches provider data and yields bounded chunks shaped for its destination table.
Requests and provider clients
Section titled “Requests and provider clients”Start with one HTTP request inside context.attempt. Pass its cancellation signal to fetch, and convert unsuccessful responses with HttpError.fromResponse. For a paginated endpoint, use paginate({ context, fetchPage }); it calls context.attempt for each page.
Authentication, token refresh, and required scopes belong to the provider client. Keep credentials in environment variables or the application’s existing secret mechanism. Direct requests outside context.attempt bypass request-level retry policy and fetch permits.
A reusable page client
Section titled “A reusable page client”This example assumes a provider contract: GET /tickets accepts updated_since, updated_before, and cursor, and returns { data: Ticket[], next_cursor: string | null }. Replace the parameters and response parsing to match the provider’s documentation. HELPDESK_API_URL is the API base URL, including a trailing slash; HELPDESK_TOKEN is its bearer token.
Create src/sources/helpdesk-client.ts:
import { HttpError } from '@chkit/plugin-ingest'
export type Ticket = { id: string; subject: string; updated_at: string }export type TicketPage = { data: Ticket[]; next_cursor: string | null }
export async function fetchTicketPage( range: { from: Date; to: Date }, cursor: string | undefined, signal: AbortSignal,): Promise<TicketPage> { const base = process.env.HELPDESK_API_URL const token = process.env.HELPDESK_TOKEN if (!base || !token) throw new Error('Set HELPDESK_API_URL and HELPDESK_TOKEN') const url = new URL('tickets', base) url.searchParams.set('updated_since', range.from.toISOString()) url.searchParams.set('updated_before', range.to.toISOString()) if (cursor !== undefined) url.searchParams.set('cursor', cursor) const response = await fetch(url, { signal, headers: { Authorization: `Bearer ${token}` }, }) if (!response.ok) throw await HttpError.fromResponse(response) return await response.json() as TicketPage}This helper performs one request. Call it through context.attempt or from paginate’s fetchPage callback, as shown in Incremental syncs. Validate untrusted response shapes at the source boundary in production; a TypeScript assertion does not validate JSON.
Pages and checkpoint metadata
Section titled “Pages and checkpoint metadata”paginate yields the complete { items, next, metadata? } returned by fetchPage, including empty and terminal pages. Read rows from page.items. initial and next support offset, string-token, and compound-object continuations; next: undefined ends pagination. Repeated continuations, including a return to initial, fail with IngestConfigError before the invalid page reaches the reader.
next controls the next request. metadata carries provider details such as a candidate checkpoint, replacement sync token, or page identity. The helper preserves metadata without interpreting or persisting it. With cursorState, explicitly map a safe, complete checkpoint from the page into SourceChunk.state; the executor commits it after the rows are saved. Empty pages can carry meaningful progress through the same path. See the complete provider cursor example.
A page cursor can be temporary, tied to one snapshot, or expire before the next run. Only persist it with cursorState when the provider guarantees it is valid for later executions. Otherwise page through a timestamp window and commit progress after that whole window succeeds.
For a sync-token API, a nonterminal page can retain the input sync token and save its next-page position in metadata. Its terminal page has no next, but can carry a replacement sync token. Metadata and pagination therefore have independent types: paginate<TItem, TCursor, TMetadata>. Protocol-specific token recovery remains the reader’s responsibility; use context.attempt for requests outside paginate and avoid wrapping its fetchPage in another attempt.
Bound work and separate accounts
Section titled “Bound work and separate accounts”Yield pages as they arrive; do not collect a large source into one array. Source page size controls response memory; batchSize controls loading and is not a hard limit on a yielded chunk. See Loading and batching.
Use stable stream IDs for independently resumable accounts or resources, for example helpdesk.account-42.tickets. The stream ID owns the checkpoint. If several accounts share a destination table, include account identity in the record key too; a provider-local ticket ID alone may collide.
Parent records and child collections
Section titled “Parent records and child collections”For a document that needs a root object and its children, assemble one complete root in the reader:
- Fetch a page of roots, such as tickets.
- For each ticket, fetch every required comments page through
paginateorcontext.attempt. - Attach the comments to the ticket and yield that object with
rawRows, or map the assembled object into typed columns.
The default loader writes the assembled row. There is no automatic nested loader or parent/child registration. If a required child request fails, let the reader fail before yielding that root; do not publish a partial object as complete. Bound the assembled object size and report any provider-imposed truncation.
For separately queried children, load a ticket_comments table with its own stream, stable comment IDs, and a ticket_id column. A child stream can discover its own roots when the provider only exposes parent-scoped endpoints. It must not assume another stream has already loaded the parents: pipeline execution does not order dependencies. Choose the row model alongside stored shape.
Use an SDK when it helps
Section titled “Use an SDK when it helps”Wrap a provider SDK or reusable client
Use a provider SDK when it handles signing, authentication, or protocol details the integration needs. Call it inside context.attempt, forward cancellation where supported, and classify SDK-specific errors through the stream’s classifyError. Avoid nested retries in the SDK and chkit unless their combined behavior is deliberate.
Pass FetchContext to reusable clients that need only attempt and signal; they do not need to depend on checkpoint types.
Related pages
Section titled “Related pages”- Incremental syncs: choose the durable resume boundary.
- Scheduling and recovery: retry policies and provider error classification.
- Test a source: check pagination and failure behavior with fixtures.