Implement and extend PostHog Data warehouse import sources. Use when adding a new source under products/warehouse_sources/backend/temporal/data_imports/sources, adding datasets/endpoints to an existing source, or adding incremental sync, resumable imports, webhook ingestion, pagination, credentials validation, and source tests.
71
88%
Does it follow best practices?
Run evals on this skill
Adds up to 20 points to the overall score
View guide
Passed
No findings from the security scan
Use this skill when building or updating Data warehouse sources in products/warehouse_sources/backend/temporal/data_imports/sources/.
Before coding, read:
products/warehouse_sources/backend/temporal/data_imports/sources/source.template (the top-of-file TODOs are the bootstrap checklist; still verify target files against current source implementations, since the template can drift)products/warehouse_sources/backend/temporal/data_imports/sources/README.mdproducts/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md — inventory of every registered source with its communication method (HTTP / vendor SDK / gRPC / DB protocol / webhook) and tracked-transport state. Skim this first to see how similar sources are wired and what state today's source you're touching is in. Keep it in sync — see "Updating SOURCES.md" below.products/warehouse_sources/backend/temporal/data_imports/sources/common/base.py — base classes (SimpleSource, ResumableSource, WebhookSource) and the FieldType unionproducts/warehouse_sources/backend/temporal/data_imports/sources/common/resumable.py — ResumableSourceManagerproducts/warehouse_sources/backend/temporal/data_imports/sources/common/webhook_s3.py — WebhookSourceManagerchargebee/ — the canonical reference for a new REST source. It uses the shared rest_source framework (declarative RESTAPIConfig + rest_api_resource, framework auth + paginators, tracked+retrying transport) and is resumable — proof the framework covers the dominant "paginate a list endpoint and yield, resumably" shape. Read it first, alongside "Prefer the shared REST framework" below. Read klaviyo/ or github/ only as a bespoke-transport fallback: they hand-roll their client for edge cases (custom query-string encoding, multi-level fan-out, JSON:API reshaping) that most sources don't have — don't copy that boilerplate into a source that doesn't need it. For dependent-resource fan-out (parent→child with type: "resolve"), also read products/warehouse_sources/backend/temporal/data_imports/sources/common/rest_source/__init__.py and config_setup.py (e.g. process_parent_data_item, make_parent_key_name).products/warehouse_sources/backend/temporal/data_imports/sources/stripe/source.py as the reference implementation.Every new source must inherit from one (or a combination) of these:
SimpleSource[Config] — default for straightforward pull-based APIs where each run fully iterates the endpoint.ResumableSource[Config, ResumableData] — preferred for any new API-backed source whose underlying API supports resumption (cursor/link-header pagination, time windows, offset tokens, or any other deterministic way to pick back up where we left off). If the API gives us a next-page token, a Link header, or a stable time filter, use ResumableSource. This lets Temporal resume after heartbeat timeouts without restarting from scratch. The manager persists state to Redis (24h TTL).WebhookSource[Config] — only when the source can push events to us (e.g. Stripe webhook endpoints). Typically combined with ResumableSource so the initial backfill is resumable and subsequent deltas come via webhook.Combine by multiple inheritance when both apply, e.g.:
class StripeSource(
ResumableSource[StripeSourceConfig, StripeResumeConfig],
WebhookSource[StripeSourceConfig],
OAuthMixin,
):
...Rule of thumb:
SimpleSource.ResumableSource.WebhookSource on top of whichever pull base fits.Databases and file-transfer sources (SFTP, S3) stay on SimpleSource unless there's a clear reason otherwise.
Most REST sources should be built on the shared rest_source framework
(common/rest_source/), not a hand-rolled client. It already provides — so you write none of it:
RESTClient defaults to make_tracked_session() and retries
429 + transient 5xx honoring Retry-After. No tenacity, no RetryableError, no fetch loop.rest_source/paginators.py, chosen by string/dict in the config, not hand-written):
single_page, header_link, json_response (next-URL in body), cursor, offset, page_number.rest_source/auth.py): bearer, api_key (header/query/cookie), http_basic, oauth2
(customer-owned client-credentials/refresh). Each redacts its own secrets — no _get_headers builder.data_selector, response actions, resume (resume_hook /
initial_paginator_state), and parent/child fan-out (fanout.build_dependent_resource).chargebee/ is the canonical example (declarative endpoints + framework auth + resume). zendesk/
shows multi-endpoint + data_selector; attio/ shows cursor pagination.
When hand-rolling is justified (read klaviyo/ then): the API needs query strings the framework
can't produce (literal brackets/operators, e.g. filter=greater-than(...), page[size]);
multi-level (2+ deep) fan-out; or per-item reshaping the data_selector can't express (e.g.
flattening JSON:API attributes into the row root). Single-level fan-out and per-item maps are
supported declaratively — don't hand-roll for those. If you must hand-roll, still ride
make_tracked_session() and do not add a second status-code retry layer (see "Retry and throttling").
Follow this order. Each step maps to TODOs in source.template.
Survey the source. Pick the endpoints a user will actually want. Cross-reference:
Bootstrap the source. Copy the template and wire up the enum/type references:
mkdir -p products/warehouse_sources/backend/temporal/data_imports/sources/{SOURCE_NAME}
cp products/warehouse_sources/backend/temporal/data_imports/sources/source.template products/warehouse_sources/backend/temporal/data_imports/sources/{SOURCE_NAME}/source.pyThen update the two hand-edited files (the template still lists posthog/schema.py too, but that file is regenerated by pnpm run schema:build in step 12 — don't maintain it by hand):
ExternalDataSourceType at products/warehouse_sources/backend/types.py — follow the existing convention in that file: ALL_CAPS with no underscores between words (e.g. ACTIVECAMPAIGN, APPLESEARCHADS), value is PascalCaseexternalDataSources at frontend/src/queries/schema/schema-general.ts — PascalCase, identical to the ExternalDataSourceType value (e.g. 'ActiveCampaign', 'GoogleAds', 'CustomerIO'). NOT kebab-case. (The only kebab-case identifier in the flow is the optional featureFlag="dwh-{source_name}".)Pick the base class (see above) and rename the class / source_type return.
Define get_source_config — name, category (required — see "Source category & keywords"), label, caption, docsUrl, iconPath, fields, and optional keywords. Use appropriate field types (see below). Also set the vendor API version metadata class attributes — see "Vendor API version metadata".
Register the source — add an import line to products/warehouse_sources/backend/temporal/data_imports/sources/__init__.py and include it in __all__. (The @SourceRegistry.register decorator on the class handles runtime registration.)
Run the config generator: pnpm run generate:source-configs. Confirm the new config class appears in products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/<your_source>.py (one generated module per source; the package __init__.py is hand-written and never regenerated). Do not edit generated modules by hand. Every time you change get_source_config.fields, re-run the generator.
Swap the generic Config type in source.py for the generated {Source}SourceConfig class.
Implement: validate_credentials, get_schemas, source_for_pipeline (plus get_resumable_source_manager / get_webhook_source_manager as needed).
Split transport logic. Put API client, paginator, row normalization, and SourceResponse assembly in {source}.py. Keep endpoint catalog/incremental fields/primary keys/partition defaults in settings.py.
Add icon. Place at frontend/public/services/{source}.png — PNG is the repo convention (~800 png vs ~58 svg, and source.template defaults to .png). SVG is accepted but not the norm; set iconPath to match whichever extension you commit. If the logo isn't already committed, fetch from Logo.dev — ask the user for the Logo.dev API key; do not hardcode one. Logo.dev's image API returns PNG (not SVG). Keep file size reasonable.
Run migrations. DEBUG=1 python manage.py makemigrations && DEBUG=1 ./bin/migrate (only needed if a new enum value triggers a Django migration).
Rebuild schema types: pnpm run schema:build. This updates posthog/schema.py from schema-general.ts and makes the source appear in frontend dropdowns. Re-run whenever schema-general.ts changes.
Release status — a finished source has no unreleasedSource flag. The default for the deliverable this skill produces is no unreleasedSource — a completed, working source ships visible and connectable. You don't need anyone's sign-off to ship it released; that's just the finished state. The scaffolded stub ships with unreleasedSource=True pre-set, so deleting that line is part of finishing the source — go ahead and remove it. (Why it matters: unreleasedSource=True hides the connector from users entirely — the frontend filters out every source where it's truthy; see DataWarehouseQueryVariant.tsx, InlineSourceSetup.tsx, and the "coming soon / Notify me" path in nonHogFunctionTemplatesLogic.tsx.)
Deleting that line is mandatory, and it is not gated on anything you can't do in your environment. In particular, "I couldn't curl the live API" or "I couldn't verify against a real account" is NOT a reason to keep the flag — that is exactly what releaseStatus=ReleaseStatus.ALPHA is for (a soft "new, lightly tested" label on a visible source). The only time unreleasedSource=True legitimately stays is when the source physically cannot sync yet because it is being landed across several PRs and the implementing code isn't all there. A source with working get_schemas / source_for_pipeline and passing tests is finished — the flag comes out. Never write a test that asserts unreleasedSource is True — that locks the bug in and is what kept 166 finished sources hidden until they had to be released in bulk.
So a newly finished, tested source ships with:
unreleasedSource (visible and connectable),releaseStatus=ReleaseStatus.ALPHA for a new source that hasn't been extensively tested (ReleaseStatus.BETA once rough edges are ironed out; ReleaseStatus.GA, or omit releaseStatus entirely, for general availability) — a soft label on a visible source, not a gate,featureFlag="dwh-{source_name}" (kebab-case) only if you want a controlled rollout to flagged users instead of releasing to everyone.Whenever you set releaseStatus, use the ReleaseStatus enum from posthog.schema — never a bare string literal. Add ReleaseStatus to your existing from posthog.schema import (...) block.
Document the source. Write or update the user-facing doc on posthog.com following the
/documenting-warehouse-sources skill (template, shared snippets, <SourceParameters /> +
<SourceTables />). Ensure docsUrl in get_source_config matches the doc filename
(kebab-case), and — if get_schemas is a static endpoint catalog — set
lists_tables_without_credentials = True (see below) so the doc's Supported tables section
renders. A finished source ships with a consistent doc, not a stub.
Delete the template TODO comments before PR.
For API-backed sources, use this split:
source.py: source registration, source form fields, schema list, credential validation, resumable/webhook manager wiring, pipeline handoff.settings.py: endpoint catalog, incremental fields, primary key, partition defaults.{source}.py: API client/auth, paginator, request params, row normalization, and SourceResponse.This keeps endpoint behavior declarative and easy to extend.
The warehouse_sources presentation layer (products/warehouse_sources/backend/presentation/views/external_data_source.py, external_data_schema.py) must stay source-agnostic.
Do not add if source_type == ExternalDataSourceType.X / source.is_direct_<engine> branches there — a CI guard (.github/scripts/check-dwh-source-agnostic.py) blocks new ones.
When a source needs behaviour the API must invoke, expose it on the source instead:
_BaseSource with a safe default (like supports_column_selection, connection_host_fields, has_managed_hogql_schema), and let the API branch on the flag.isinstance(source, <Capability>).DataWarehouseTable, maps columns) is keyed on the engine, not the source type — dispatch on source.direct_engine through the engine adapter/registry (posthog/hogql/direct_sql/ for query concerns, the data_warehouse engine registry for materialization), never source_type.Keep source-domain semantics (how to talk to the engine, how it names things, whether filters push down) on the source; the warehouse-domain work it drives (DataWarehouseTable rows, managed viewsets, hog functions) stays in data_warehouse, keyed off what the source or adapter returns.
Source capabilities never import data_warehouse types.
See products/data_warehouse/backend/presentation/README.md.
For REST sources that mix top-level and fan-out endpoints, keep endpoint metadata in settings.py and route in {source}.py with this priority:
After a table syncs, a background activity (workflow_activities/enrich_table_semantics.py) writes
WarehouseColumnAnnotation rows describing each table/column, surfaced to the AI agent. For
fixed-schema sources (SaaS APIs) the schema is the same for everyone, so document it once from the
official API docs instead of paying an LLM to re-derive it per team. These curated descriptions are
authoritative — they're applied directly (description_source="canonical") and never sent to the LLM.
Add a canonical_descriptions.py accompanying the source (sibling of source.py / settings.py):
# products/warehouse_sources/backend/temporal/data_imports/sources/{source}/canonical_descriptions.py
from products.warehouse_sources.backend.temporal.data_imports.sources.common.canonical_descriptions import CanonicalDescriptions
CANONICAL_DESCRIPTIONS: CanonicalDescriptions = {
"Charge": { # key = ExternalDataSchema.name (the endpoint name from get_schemas / ENDPOINTS)
"description": "A single attempt to move money into your account by charging a payment source.",
"docs_url": "https://stripe.com/docs/api/charges", # passed to the LLM for columns not covered here
"columns": { # column name -> one-line description, taken from the official API docs
"id": "Unique identifier for the charge.",
"amount": "Amount intended to be collected, in the smallest currency unit (e.g. cents).",
},
},
}Then override the hook on the source class with a lazy import of the sibling file:
def get_canonical_descriptions(self) -> CanonicalDescriptions:
from products.warehouse_sources.backend.temporal.data_imports.sources.{source}.canonical_descriptions import CANONICAL_DESCRIPTIONS
return CANONICAL_DESCRIPTIONSRules:
get_schemas returns (matches ENDPOINTS), not the
prefixed warehouse table name.description falls back to the LLM, which is given the
source name, endpoint, docs_url, and column data types.{}.source.py/settings.py transport logic — this is purely additive metadata.The posthog.com docs render a Supported tables section via a <SourceTables /> component fed by the
public_source_configs API, which calls get_documented_tables() on each source. The base
implementation lists tables from get_schemas (merged with canonical_descriptions) only when the
source opts in:
class MySource(SimpleSource[MySourceConfig]):
lists_tables_without_credentials = True # static endpoint catalog — safe for public docsSet this to True only when get_schemas iterates a static endpoint catalog with no I/O — no
network, no DB, no credentials (the common fixed-schema SaaS pattern: for endpoint in ENDPOINTS). The
endpoint builds a placeholder config and calls get_schemas with no real credentials, so a source that
connects to discover schemas (SQL, file storage, MongoDB, ad platforms that list accounts) must leave
this False (the default) — otherwise it would try to connect to an empty host, hang, or close the DB
session. When False, the docs render a generic "discovered from your account" note instead.
The richer the table list, the better the docs — so pair this with canonical_descriptions.py
(table/column descriptions). Verify the rendered output via the API:
GET /api/public_source_configs → your source → tables.
Every source must set category on its SourceConfig — it groups the source in the new-source wizard
catalog (a category rail + tile grid). A test (tests/test_source_categories.py) fails if any registered
source has no category, so this is non-optional. Import the enum from posthog.schema:
from posthog.schema import DataWarehouseSourceCategory
...
return SourceConfig(
name=SchemaExternalDataSourceType.STRIPE,
category=DataWarehouseSourceCategory.PAYMENTS___BILLING,
keywords=["billing", "subscriptions"],
...
)Pick the single closest bucket. The enum members (note the triple underscore where the label has " & "):
DATABASES — OLTP/OLAP databases, warehouses, data streams (Postgres, Snowflake, BigQuery, Kafka, …)FILE_STORAGE — object/file stores & file transfer (S3, Azure Blob, GCS, Google Drive, SFTP, …)ADVERTISING — ad platforms & mobile attribution (Google Ads, Meta Ads, Reddit Ads, Adjust, …)MARKETING___EMAIL — email/SMS/marketing automation (Klaviyo, Mailchimp, Braze, SendGrid, …)CRM — CRM & sales intelligence (HubSpot, Salesforce, Attio, Pipedrive, ZoomInfo, …)SALES — sales engagement/enablement, contracts (Salesloft, Outreach, Gong, DocuSign, …)CUSTOMER_SUPPORT — helpdesk/support/CX (Zendesk, Intercom, Freshdesk, Front, …)PAYMENTS___BILLING — payment processors & subscription billing (Stripe, Chargebee, PayPal, …)FINANCE___ACCOUNTING — accounting/ERP/expense/spend (QuickBooks, Xero, NetSuite, SAP ERP, …)ANALYTICS — product/web/marketing analytics & experimentation (Amplitude, Mixpanel, GA, …)ENGINEERING___MONITORING — dev tooling, CI, error/uptime monitoring, feature flags, identity/auth (GitHub, Datadog, Sentry, LaunchDarkly, Auth0, …)PRODUCTIVITY — project mgmt, docs, forms, scheduling (Notion, Airtable, Jira, Linear, Typeform, …)HR___RECRUITING — HRIS/ATS/payroll/people (Ashby, Greenhouse, BambooHR, Workday, Gusto, …)COMMUNICATION — messaging/meetings/telephony/social (Slack, Zoom, Microsoft Teams, Twilio, …)E_COMMERCE — online store/commerce (Shopify, WooCommerce, BigCommerce, …)The category list is the source of truth in frontend/src/queries/schema/schema-general.ts
(dataWarehouseSourceCategories); pnpm run schema:build regenerates the Python DataWarehouseSourceCategory
enum. Adding a new category means editing that array and rebuilding — don't invent ad-hoc strings.
keywords is an optional list of lowercase search aliases — only add when the source has a common acronym or
alternate spelling a user might type (e.g. ["ga4", "ga"], ["sql server"], ["facebook ads"]). Skip it when
the name already obviously matches; don't add noise.
Every source declares three class attributes (on the source class body, alongside lists_tables_without_credentials)
describing the vendor's API version.
The framework (common/base.py) records the version each ExternalDataSource runs against so old pins keep working
and deprecations can be surfaced;
sources/tests/test_source_versions.py enforces the invariants below across every registered source, so a new
source that gets these wrong fails CI.
Two cases:
The vendor exposes a real, pinnable API version — a URL path segment (/v3/, /2/), a required version
header value (a dated 2022-11-28), a dated query/version param, or a named release. Declare all three:
class MySource(SimpleSource[MySourceConfig]):
supported_versions = ("v3",) # opaque vendor labels — never parsed or ordered
default_version = "v3" # stamped onto newly created sources; must be in supported_versions
api_docs_url = "https://vendor.example/docs/api" # API reference or changelog page (https, not the marketing site)Pin the version the source's own code actually calls (the base URL path, a version header, or a version
constant in settings.py / {source}.py) — not the vendor's newest version. Examples already in the tree:
GitHub ("2022-11-28",) (dated header), HubSpot ("v3",) (path), Klaviyo ("2024-10-15",) (dated revision).
The vendor has no meaningful API versioning — set only api_docs_url; leave supported_versions /
default_version at the framework default (("v1",), the UNVERSIONED_API_VERSION sentinel). A bare /v1/
that has never changed and isn't a documented version choice is this case.
Rules:
default_version must equal the single entry in supported_versions, and api_docs_url must be https://.self.resolve_api_version(inputs.api_version)), which already falls back to default_version./warehouse-source-new-version skill — not this one.Defined in get_source_config.fields. All field types live in posthog/schema.py and are unioned as FieldType in products/warehouse_sources/backend/temporal/data_imports/sources/common/base.py.
SourceFieldInputConfig — basic input (text, email, number, password, textarea). Rendered as <LemonInput />.SourceFieldSwitchGroupConfig — toggle that reveals a sub-group of fields. Use for optional feature blocks.SourceFieldSelectConfig — dropdown. Options can carry sub-fields shown when selected (use for alternative auth methods — e.g. API key vs OAuth).SourceFieldOauthConfig — OAuth via Integration model. See OAuth section.SourceFieldFileUploadConfig — file upload (JSON). Use keys=["..."] allow-list or "*".SourceFieldSSHTunnelConfig — renders SSH tunnel sub-fields; adds ssh_tunnel: SSHTunnel to the config with helpers.Guidelines:
SourceFieldSelectConfig with child fields per option.SourceFieldSwitchGroupConfig.SourceFieldInputConfigType.PASSWORD. The serializer derives sensitive vs nonsensitive keys automatically from the field definitions — you do not need to maintain an allow-list elsewhere.source_for_pipelineReturn a SourceResponse directly. Do not use dlt_source_to_source_response for new sources — DLT is being removed.
Prefer yielding data in the shape the API returns it. No custom dataclasses, no heavy parsing. Yield either dict, list[dict] (preferred when possible), or a pyarrow.Table. The pipeline buffers and batches for you.
Default to yielding raw dict / list[dict] and let the pipeline batch for you. The pipeline already runs a Batcher (pipelines/pipeline/pipeline.py) at 5000-row / 200 MiB thresholds, so the common case needs no batcher of its own. Reach for pyarrow.Table only when you already have arrow-shaped data (e.g. a ClickHouse adapter). A source may instantiate its own Batcher with smaller thresholds (e.g. chunk_size=2000, chunk_size_bytes=100 * 1024 * 1024, as klaviyo and ~70 other sources do) when it deliberately wants a tighter memory footprint for large/wide rows — that's a valid choice, not the default. What to avoid is a second full-size batcher, which just double-buffers with no win.
For pyarrow tables, cap in-memory rows at ~200 MiB or ~5000 rows. Use helpers like table_from_iterator() / table_from_py_list() from products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline/utils.py.
URL construction: use urllib.parse.urlencode for query strings. Don't use requests.Request(...).prepare().url — PreparedRequest.url is typed Optional[str] and the typical workaround (prepared.url or f"...") carries an unreachable fallback. urlencode is shorter, dependency-free, and produces identical output for ASCII-safe params.
@dataclasses.dataclass
class MyResumeConfig:
next_url: str # or cursor, offset, time window — whatever the API uses
class MySource(ResumableSource[MySourceConfig, MyResumeConfig]):
def get_resumable_source_manager(self, inputs: SourceInputs) -> ResumableSourceManager[MyResumeConfig]:
return ResumableSourceManager[MyResumeConfig](inputs, MyResumeConfig)
def source_for_pipeline(
self,
config: MySourceConfig,
resumable_source_manager: ResumableSourceManager[MyResumeConfig],
inputs: SourceInputs,
) -> SourceResponse:
return my_source(..., resumable_source_manager=resumable_source_manager)In the transport function:
resume = manager.load_state() if manager.can_resume() else None
url = resume.next_url if resume else initial_url
while True:
data = fetch_page(url)
# yield batch
next_url = data.get("links", {}).get("next")
if not next_url:
break
manager.save_state(MyResumeConfig(next_url=next_url))
url = next_url # advance before the next fetch, otherwise we loop on the same pageSave state after yielding each batch, not before — so if we crash we re-yield the last batch (merge dedupes on primary key) rather than skipping it.
webhook_template returning a HogFunctionTemplateDC that transforms incoming webhook payloads.webhook_resource_map mapping our schema name → external object type.create_webhook, delete_webhook, get_external_webhook_info if the API allows programmatic webhook management. Otherwise return a failed result and provide a webhookSetupCaption explaining manual setup.webhookFields to SourceConfig for post-setup inputs (e.g. signing secret).source_for_pipeline, call self.get_webhook_source_manager(inputs) and pass its iterator alongside the pull iterator so a single sync pulls historical + webhook-delivered rows.SourceSchema.supports_webhooks=True only for endpoints where webhooks are actually viable (usually incremental/append-only ones).table_transformer. WebhookSourceManager.get_items() takes an optional table_transformer: Callable[[pa.Table], pa.Table] applied after the raw webhook payloads are deserialized into row dicts. Delta merge only de-dupes across syncs (on primary_keys), not within a single source batch — so when one batch can carry multiple events for the same object (e.g. customer.created then customer.updated), pass a transformer that keeps only the latest version per id. Reference: _webhook_table_transformer in stripe/stripe.py, wired via webhook_source_manager.get_items(table_transformer=_webhook_table_transformer) in stripe_source. It groups rows by object.id, keeps the one with the greatest event created timestamp, and rebuilds the table shaped like the underlying object (ready to merge on primary_keys=["id"]).SQL DB sources (Postgres, MSSQL, Snowflake, Redshift today) can import tables from every namespace (schema) in one connection: a blank namespace field discovers tables across all non-system namespaces, the wizard groups them by namespace, and sync writes one warehouse table per namespace.table. Reference implementation: postgres/postgres.py + PostgresImplementation; the shared seam lives in common/sql/.
The capability marker is the source's schema field being optional (required=False) in get_source_config — is_multi_schema_capable_sql_source() (products/data_warehouse/backend/sql_warehouse_migration.py) keys off it, so flipping the field optional is what turns on the viewset migration behavior. Treat None / "" / whitespace as "all namespaces" (normalize_namespace in common/sql/location.py) and never emit WHERE table_schema = ''.
Checklist for bringing a SQL source to multi-schema parity:
required=False on the schema field, rerun pnpm run generate:source-configs. Keep database required: the database/catalog stays fixed per connection.get_columns, get_primary_keys, index/row-count/foreign-key helpers: when the namespace is blank, drop the WHERE table_schema = <ns> predicate (excluding system namespaces like information_schema, pg_catalog, sys) and return qualified display names (namespace.table). Keep the single-namespace fast path when the field is set.get_source_metadata — return SourceMetadata(catalog_by_table, schema_by_table, table_name_by_table) keyed by the qualified display name. SQLSource.get_schemas stamps it onto each SourceSchema, and reconcile_schema_metadata persists it into ExternalDataSchema.sync_type_config["schema_metadata"].build_pipeline — resolve (schema, table_name, response_name) with resolve_source_location (common/sql/location.py): per-row metadata → dotted-name self-heal → config namespace. Run SQL against the resolved schema + unqualified table; set SourceResponse.name = response_name (dwh_storage_key or schema.name, normalized) — never the bare table name, or the row's Delta path moves and orphans synced data.(schema, table). Missing one degrades silently (no partitioning / wrong stats).quote("a.b") yields one wrong identifier. Split into (schema, table) first and use quote_qualified (common/sql/identifiers.py).analytics.users; S3/Delta subdir analytics_users (normalized response_name); HogQL table {prefix}_analytics_users. The one stored exception is dwh_storage_key, which pins a migrated legacy row to its original Delta path.sql_warehouse_migration.py renames rows to qualified form and stamps dwh_storage_key, preserving synced data with no re-sync. Don't add source_type == "..." branches to the shared layer.Discovery cost: validate_credentials and database_schema run discovery with no name filter, so a blank namespace on a catalog with hundreds of schemas must not issue per-table queries per namespace — batch the listing queries or cap enumeration (see Snowflake's SHOW PRIMARY KEYS handling).
Every HTTP call from products/warehouse_sources/backend/temporal/data_imports/sources/** must go through make_tracked_session() (from
products.warehouse_sources.backend.temporal.data_imports.sources.common.http). The tracked session attaches team_id, source_type,
external_data_source_id, external_data_schema_id, and external_data_job_id to every outbound request's
log line and OTel metric, and participates in opt-in sample capture.
requests usage: make_tracked_session(headers=..., retry=...) returns a requests.Session. Use
session.get/post/... instead of the module-level requests.get/... shortcuts.redact_values=(api_key, token, ...) to
make_tracked_session(...) so the tracked transport masks those literal values in logged URLs, headers,
and sampled bodies — important for keys that ride in a query param or an odd header name. rest_source
auth classes do this automatically (each implements secret_values()); only raw-session sources need to
pass redact_values themselves.rest_source.RESTClient: it defaults to a tracked session
automatically; no change needed.RequestsClient(session=...),
gspread authorize(credentials, session=...), BigQuery via AuthorizedSession + TrackedHTTPAdapter),
inject one. Reference patterns live in stripe/stripe.py, google_sheets/google_sheets.py, and
bigquery/bigquery.py.bingads, linkedin-api's RestliClient), add a
# nosemgrep: data-imports-http-transport-... pragma with a one-line reason and record the source as
⚠️ Vendor SDK in SOURCES.md.CI enforces this via .semgrep/rules/security/data-imports-http-transport.yaml. The rule bans direct requests.Session(),
requests.<verb>(...), and httpx.Client/AsyncClient/<verb> inside sources/**. Type-only imports
(from requests import Response, from requests.exceptions import HTTPError) remain allowed.
gRPC calls from sources/** ride client interceptors from
products.warehouse_sources.backend.temporal.data_imports.sources.common.grpc, which attach the same JobContext labels to logs and
OTel metrics (data_import_grpc_*) and feed opt-in sample capture (protobuf → scrubbed JSON). Two seams:
interceptors= list (google-ads GoogleAdsClient.get_service(...)), pass
interceptors=tracked_interceptors(host) on every get_service call — google-ads rebuilds the channel
per call, so the interceptors must be re-supplied each time. Reference: google_ads/google_ads.py.channel= / transport= (BigQuery Storage Read API), build the credential-bearing
channel, wrap it with make_tracked_channel(channel, host=...), then hand it to the transport. Reference:
bigquery/bigquery.py:bigquery_storage_read_client.CI enforces this via .semgrep/rules/security/data-imports-grpc-transport.yaml, which bans raw grpc.*_channel(...)
and direct BigQueryReadClient(...) / GoogleAdsClient(...) construction inside sources/** (outside the
common/grpc/ package and the two reference source files). Operators arm sample capture with
python manage.py warehouse_sources_capture_grpc_samples enable ....
products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md is the inventory of every registered source, its
communication method, and whether its outbound traffic is tracked. Update it as part of the same PR
whenever you:
⚠️ Vendor SDK to ✅.requests to a vendor SDK. Update both the comm method and tracked-transport columns.Keep the entries alphabetical within each table. The scaffolded list is one source per line (one bullet each, also alphabetical) so adding or removing a source only touches its own line and avoids conflicts with concurrent PRs — don't collapse it back into a comma-separated paragraph. If you add a partially-tracked source, also append a short "Notes on partially-tracked sources" entry explaining what blocks tracking (e.g. a vendor SDK with no session/interceptor seam).
common/base.py exposes class-level flags most API sources leave at their defaults, but which matter
when they apply:
supports_column_selection (default False; SQLSource sets True) — whether the source honors
enabled_columns via SELECT projection.supports_row_filters (default False) — must be True for a source to apply saved row_filters;
otherwise filters are ignored.has_managed_hogql_schema (default False; True for revenue-analytics sources like Stripe,
Paddle, Zendesk) — a fixed field set powers a managed HogQL schema, which disables column selection.
A new billing/revenue source that ignores this can break revenue analytics.cleanup_cdc_resources_on_deletion — best-effort teardown hook for sources that provision CDC
resources; override only if your source creates such resources.For vendor API versioning (supported_versions, default_version, api_docs_url,
deprecated_versions, resolve_api_version()), use the dedicated warehouse-source-new-version
skill — don't hand-roll version handling.
@SourceRegistry.register.SimpleSource[GeneratedConfig] unless resumable/webhook behavior is required.table_format="delta" in endpoint resources.primary_keys are endpoint-specific (declare in settings.py, not always id). Use composite keys when no single field is unique. The key must be unique across the whole table, not per parent: fan-out child endpoints aggregate rows from every parent, so include the parent identifier in the key (e.g. ["form_id", "token"]) unless the API explicitly documents global uniqueness. Non-unique keys seed duplicate rows in the Delta table, and every later merge multi-matches them — merges get slower each sync until the pod OOMs.partition_mode="datetime" with a stable datetime field.partition_count and partition_size.created_at, dateCreated, firstSeen. Never use updated_at or lastSeen.get_non_retryable_errors() for known permanent failures (401/403, invalid/expired credentials, missing scopes).supports_incremental=True when the API exposes a server-side timestamp filter (<field>_gte, since, modified_after, etc.). A "client-side cursor" that fetches every page and skips already-seen rows in Python is not incremental — every run still hits every page, so the API cost of an "incremental" sync ends up identical to a full refresh. If the API has no server filter, ship full refresh only.db_incremental_field_last_value.inputs.incremental_field — that's the user's chosen cursor field from the schema settings. INCREMENTAL_FIELDS per-endpoint is the menu of advertised options; don't reach into INCREMENTAL_FIELDS[endpoint][0] to pick a default and silently override the user's selection.?sorting=created_at (or whatever) globally. Verify each list endpoint's allowed sort values against the API spec and with a curl smoke-test against the live API — APIs frequently document one set of options and silently reject another, or use a different timestamp column on certain resources.?sorting= explicitly on a stable monotonic field when paginating. For incremental sources, the request sort must match SourceResponse.sort_mode ("asc" typically; "desc" only when forced by the API — see stripe/stripe.py, github/settings.py) so the pipeline's cursor watermark advances correctly. For full-refresh sources, an explicit sort prevents page-boundary skips/duplicates if the API's implicit default is unstable or shifts as rows are inserted during the sync.sort_mode must match the order rows actually arrive in — verify it, don't assume it. The pipeline trusts sort_mode="asc" to checkpoint the incremental watermark after every batch and to allow safe mid-sync worker shutdowns; declaring asc while the API returns newest-first corrupts the watermark and breaks resume semantics. Check the API's default sort (it applies when you can't pass sort), and remember cursor pagination often rejects or ignores sort params entirely.sort_mode="desc" only if the endpoint truly cannot return ascending. For descending sources, handle db_incremental_field_earliest_value to scroll earlier rows before newer ones (see Stripe).db_incremental_field_last_value (see typeform/typeform.py:TypeformResponsesPaginator) — otherwise every incremental sync re-fetches and re-merges each parent's full history, which is both an API-cost bug and a per-sync memory amplifier.Before finalizing endpoint logic, verify from docs and with curl against the live API (not just docs — APIs frequently silently ignore unknown params or document outdated enums):
{"data": [...]}).before/after tokens), confirm whether the API allows sort and time-window params alongside it — many reject or ignore them, which dictates both your sort_mode and how pagination terminates on incremental syncs.sorting= values does each list endpoint accept? Some APIs vary the allowed enum per resource. Confirm with curl that the value you intend to pass returns 200, and probe with a future-date cutoff to confirm whether timestamp filters are honored or silently ignored.<field>_gte / since / modified_after actually filter, or does the API accept it and ignore it? Test by passing a future date and checking whether the row drops out.If undocumented, keep parsing/merge logic conservative and add a short code comment noting the uncertainty.
products/warehouse_sources/backend/temporal/data_imports/sources/<source>/api_inventory.md).settings.py (path, primary_key, incremental_fields, partition_key, sort_mode).limit, required filters).Link headers — check both rel="next" and any results flag.update_request to avoid duplicate query params.make_tracked_session() and rest_source.RESTClient already retry 429 + transient 5xx at the transport layer, honoring Retry-After. Do NOT wrap a second tenacity @retry around your fetch for status codes — it compounds (e.g. 3 transport attempts × 5 tenacity attempts) and is the single most-copied mistake across existing sources. If you use the framework or the tracked session, status-code retries are already handled; write none.tenacity for a condition the transport does not cover — e.g. an app-level "still processing" body that isn't a retryable HTTP status. If you do, disable transport retries so they don't compound: make_tracked_session(retry=Retry(total=0)).429 — the transport already honors Retry-After. Keep any custom retry bounded and deterministic (stop_after_attempt), with clear terminal behavior.The backoff above is the right control when the customer owns the credential — their own PAT / API key / OAuth token on their own third-party account, which is nearly every source.
PostHog can't overspend a budget it doesn't own, so honoring 429 / Retry-After at the source is enough.
The exception is a credential PostHog owns and shares across processes — today that's the PostHog GitHub App installation token (many PostHog subsystems draw from one per-installation budget at once).
There, reactive backoff isn't enough: without coordination, concurrent PostHog callers collectively blow past the shared limit before any 429 comes back.
Those calls must route through posthog/egress/ — a Redis-backed shared budget plus telemetry, gated by construction — never hand-rolled requests.
The GitHub source is the reference: it keys the limiter on the GitHub App installation id (the budget owner in GitHub's own id space, not a PostHog DB row), and the customer-PAT path skips the limiter token-blind.
Raw calls to api.github.com are blocked by the github-api-calls-go-through-egress semgrep rule, so a GitHub-shaped source lands on the egress path by construction.
Deciding question is never "is this a warehouse source?" — it's "who owns the token, and could concurrent PostHog processes trample each other on it?"
Fan-out = iterate a parent resource, then query child endpoints per parent.
Prefer dependent resources for single-hop fan-out. Use rest_api_resources with a parent and child that declares type: "resolve" for the parent field. Shared infra (rest_source/__init__.py, config_setup.process_parent_data_item) paginates the parent and calls the child per parent row. Use include_from_parent so child rows carry parent fields (injected as _<parent>_<field> via make_parent_key_name).
Make fan-out declarative. Add a fan-out config object in settings.py (e.g. DependentEndpointConfig) with parent_name, resolve_param, resolve_field, include_from_parent, optional parent field renames, and optional parent endpoint params. Route single-hop fan-out through a shared helper (e.g. common/rest_source/fanout.py:build_dependent_resource).
Parent field rename mapping belongs in the helper. Callers should not branch on whether renames exist.
Per-endpoint pagination/selectors — build_dependent_resource supports endpoint overrides (parent_endpoint_extra, child_endpoint_extra for paginator / data_selector, page_size_param for non-limit size params).
Path pre-formatting: process_parent_data_item only does str.format() with the resolved param. Pre-format static placeholders with .replace() before passing to the resource config, so only the resolved placeholder remains.
Custom iterator only when fan-out is 2+ levels deep. Reuse the same pagination/retry helpers as elsewhere.
Before implementing OAuth, check if the integration already exists — search posthog/models/integration.py loosely for the service name before concluding it's new.
If new:
Env vars. Add to posthog/settings/integrations.py:
YOUR_SOURCE_CLIENT_ID = get_from_env("YOUR_SOURCE_CLIENT_ID", "")
YOUR_SOURCE_CLIENT_SECRET = get_from_env("YOUR_SOURCE_CLIENT_SECRET", "")Integration kind. In posthog/models/integration.py:
IntegrationKind enum.OauthIntegration.supported_kinds.elif kind == "your-source": return OauthConfig(...) branch in oauth_config_for_kind().
Raise NotImplementedError("<Source> app not configured") when the env vars are empty — that's the
fail-closed message, so code and charts can ship before the secret values exist.id_path / name_path from a claim (sub), mirroring the reddit/bing
branches in integration_from_oauth_response.Register the client + deploy the credentials. Registering the OAuth client with the provider,
the redirect URIs (US/EU/dev/localhost), the charts PR (wiring the env vars into both
posthog-django-shared-secrets for the web app and the worker's secret_env_app_specific store),
and writing the values into AWS Secrets Manager via the PostHog/secrets CLI — plus which of these
an agent can vs. must not automate — are all in
references/oauth-app-deployment.md.
Override get_non_retryable_errors() to mark errors that should permanently fail instead of retrying:
def get_non_retryable_errors(self) -> dict[str, str | None]:
return {
"401 Client Error: Unauthorized for url: https://api.example.com": "Your API key is invalid or expired. Please generate a new key and reconnect.",
"403 Client Error: Forbidden for url: https://api.example.com": "Your API key does not have the required permissions. Please check the key permissions and try again.",
}Common cases: 401 Unauthorized, 403 Forbidden, invalid/expired tokens, OAuth tokens needing re-auth.
validate_credentialsCalled with schema_name=None at source-create (one cheap probe to confirm the token is genuine) and with schema_name=<name> from the per-schema incremental_fields action (confirm scope for that specific endpoint).
If the API distinguishes 401 (bad token) from 403 (valid token, missing scope), accept 403 at source-create — users may legitimately only grant scopes for the endpoints they want to sync. Re-raise 403 only when schema_name is set. Sync-time 403s are handled separately by get_non_retryable_errors().
For per-table scope status in the schema picker, override get_endpoint_permissions(config, team_id, endpoints) -> {name: None | reason}: probe each endpoint and return None when reachable or a short reason when not. The database_schema action surfaces it as each table's permission_error, so the user sees which tables need extra scopes and deselects them — it must never block source-create. The base default reports everything reachable.
When you do surface a missing scope, name it — providers usually state it (Required access: `read_x` access scope.), so parse that into your own message instead of dumping the raw exception or collapsing it to a bare table list. Probe whatever field the sync query needs (not just id) so the per-table check reflects what syncing that table actually requires. Keep the probe narrow: only a real denial is a missing scope — a throttle, 5xx, or network blip is not, so route those through the retryable path rather than bucketing every exception as "missing permission".
If the API issues OAuth scopes or per-resource access tokens, declare every scope the source actually calls so users know what to grant — don't make them grant the full set defensively.
requiredScopes on SourceFieldOauthConfig (space-separated string, matches the OAuth scope parameter format). The frontend diffs it against the integration's granted scopes and warns the user with a Reconnect action when any are missing.caption instead. Captions render through LemonMarkdown, so backticks, bold, and links work.If your source stores a secret (API token, password) and sends it to a host that the user configures in a
non-host field, declare that field on the source class's connection_host_fields property (from the
base source in common/base.py):
@property
def connection_host_fields(self) -> list[str]:
# `okta_domain` is where the stored API token is sent; retargeting it must re-require the token.
return ["okta_domain"]The update serializer reads this list and forces the editor to re-enter the source's secrets whenever one of
these fields changes. Without it, an org member could PATCH the host field to a server they control while the
preserved (omitted) secret is reused — exfiltrating the credential. host and the SSH-tunnel target are
already handled separately, so only sources whose connection target lives in a differently named field (e.g.
Okta's okta_domain) need to override this. The default is [] (no extra fields).
Pair this with _is_host_safe (common/mixins.py) at both
source-create and sync time to block hosts resolving to internal/private IPs.
From products/warehouse_sources/backend/temporal/data_imports/sources/common/mixins.py:
SSHTunnelMixin — with_ssh_tunnel() context plus make_ssh_tunnel_func() for deferred tunnel opening.OAuthMixin — get_oauth_integration() to pull Integration from the DB.ValidateDatabaseHostMixin — is_database_host_valid() to block internal VPC IPs (unless SSH tunnel is used).source.template defaults to .png; ~800 png vs ~58 svg). SVG is accepted but not the norm. Keep file size reasonable.frontend/public/services/ and reference as /static/services/{name}.png (or .svg) in iconPath — match the extension you commit.Add at least two test modules:
tests/test_<source>_source.py (source-class level):
source_typeget_source_config fields and labelsget_schemas outputsvalidate_credentials success/failuresource_for_pipeline argument plumbingget_resumable_source_manager returns a manager bound to the right data classcreate_webhook / delete_webhook / get_external_webhook_info behavior, webhook_resource_map correctness, webhook_template presencetests/test_<source>.py (transport level):
rest_api_resources, pass rows with _<parent>_<field> keys to exercise parent-field injection and rename behaviorsettings.pyPrefer behavior tests over config-shape tests. Avoid brittle assertions on internal config dict structure unless they protect a known regression that cannot be asserted via output behavior.
Use parameterized tests for status codes and edge cases. Lean toward over-covering.
Bootstrapping:
- [ ] Enum added to products/warehouse_sources/backend/types.py (ALL_CAPS, no underscores between words)
- [ ] Entry added to frontend/src/queries/schema/schema-general.ts (PascalCase, matching the ExternalDataSourceType value — NOT kebab-case) — `pnpm run schema:build` regenerates posthog/schema.py from this; don't hand-edit posthog/schema.py
- [ ] Source imported in products/warehouse_sources/backend/temporal/data_imports/sources/__init__.py + __all__
- [ ] Class inherits from SimpleSource / ResumableSource / WebhookSource (or combo) — see "Picking the right base class"
Source implementation:
- [ ] Set category on get_source_config (required — DataWarehouseSourceCategory; groups the source in the wizard catalog)
- [ ] Add keywords if the source has a common acronym / alternate spelling (optional, lowercase)
- [ ] Set api_docs_url (https, vendor API docs/changelog); add supported_versions + default_version if the vendor
exposes a real version token — pin what the code actually calls (see "Vendor API version metadata")
- [ ] Define source fields in get_source_config
- [ ] Implement validate_credentials
- [ ] Implement get_schemas
- [ ] Add endpoint settings (settings.py)
- [ ] Implement transport + paginator ({source}.py)
- [ ] Return SourceResponse with correct primary_keys, partitioning, sort_mode
(keys unique table-wide — parent id in fan-out child keys; sort_mode verified against actual response order;
incremental cursor pagination stops at the watermark)
- [ ] Implement get_resumable_source_manager if ResumableSource
- [ ] Implement webhook methods if WebhookSource
- [ ] Add get_non_retryable_errors for auth/permission errors
- [ ] (Fixed-schema sources) Add canonical_descriptions.py from the API docs + override get_canonical_descriptions
- [ ] (Fixed-schema sources, static get_schemas only) Set lists_tables_without_credentials = True so public docs render the table catalog
Tooling & assets:
- [ ] Icon in frontend/public/services/ (PNG is the convention; ask user for Logo.dev key if needed)
- [ ] Run `pnpm run generate:source-configs`
- [ ] Swap generic Config for generated {Source}SourceConfig in source.py
- [ ] Run `pnpm run schema:build`
- [ ] Django migrations run if enum value requires it
Release status (a finished source has NO unreleasedSource flag — it hides the source from users entirely):
- [ ] REQUIRED: delete `unreleasedSource=True` from the finished source (the scaffolded stub ships with it).
Not being able to curl the live API is NOT a reason to keep it — use releaseStatus=ALPHA.
Keep it ONLY when the code genuinely can't sync yet (landed across multiple PRs).
- [ ] No test asserts `unreleasedSource is True` (that anti-pattern locks the source hidden)
- [ ] When set, releaseStatus uses the `ReleaseStatus` enum, never a string literal
- [ ] releaseStatus=ReleaseStatus.ALPHA for a new source not yet extensively tested
(ReleaseStatus.BETA later; ReleaseStatus.GA or omit for GA)
- [ ] featureFlag="dwh-{source_name}" ONLY if you want a controlled rollout instead of releasing to all
Tests & handoff:
- [ ] Source tests (test_<source>_source.py)
- [ ] Transport tests (test_<source>.py)
- [ ] User-facing doc written/updated per /documenting-warehouse-sources (docsUrl matches filename; `audit_source_docs` passes)
- [ ] `ruff check . --fix` and `ruff format .`
- [ ] List any new env vars (OAuth client IDs/secrets, etc) in the PR / handoffAfter changing source fields, re-run pnpm run generate:source-configs and pnpm run schema:build, then the targeted tests for the new source. Run ruff check . --fix and ruff format . on modified Python files.
sources/__init__.py, or schema:build not rerun.test_source_categories failing: the source's get_source_config is missing category — set it to the closest DataWarehouseSourceCategory bucket.generate:source-configs after updating fields.sort_mode="asc" declared on an API that returns newest-first: the watermark checkpoints to ≈now after the first batch and mid-sync shutdowns lose data ordering guarantees.get_non_retryable_errors.validate_credentials(schema_name=None) probes every resource's scope instead of just the token, so one missing scope — often on a table the user won't sync — blocks the whole source. Probe only the token at create; report per-table scope via get_endpoint_permissions.save_state after yielding a batch; or saved before yield and a crash causes data loss.is_webhook=False, or initial_sync_complete=False.KeyError: pre-format static path placeholders (see Fan-out).Endpoint, ClientConfig, IncrementalConfig) to keep static checks precise.updated_at instead of created_at; partitions rewrite on every sync.d1dd198
If you maintain this skill, you can claim it as your own. Once claimed, you can manage eval scenarios, bundle related skills, attach documentation or rules, and ensure cross-agent compatibility.