-
Notifications
You must be signed in to change notification settings - Fork 3.4k
feat(data-warehouse): implement recall_ai import source #107242
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from all commits
Commits
Show all changes
6 commits
Select commit
Hold shift + click to select a range
b7e6f93
feat(data-warehouse): implement recall_ai import source
Gilbert09 4068ca7
fix(recall_ai): pin pagination to region host, disable capture for ca…
Gilbert09 d55142e
Merge branch 'master' into posthog/recall-ai-import-source
Gilbert09 45f9c10
fix(recall_ai): defensively pin credential probe against redirects
Gilbert09 d56d4b3
Merge branch 'master' into posthog/recall-ai-import-source
Gilbert09 0d59103
Merge branch 'master' into posthog/recall-ai-import-source
Gilbert09 File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
5 changes: 4 additions & 1 deletion
5
...cts/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/recallai.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,9 +1,12 @@ | ||
| # This file is automatically generated from `SourceRegistry.get_all_sources()` | ||
| # Do not edit manually - run `pnpm generate:source-configs` to regenerate. | ||
|
|
||
| from typing import Literal | ||
|
|
||
| from products.warehouse_sources.backend.temporal.data_imports.sources.common import config | ||
|
|
||
|
|
||
| @config.config | ||
| class RecallAISourceConfig(config.Config): | ||
| pass | ||
| api_key: str | ||
| region: Literal["us-east-1", "us-west-2", "eu-central-1", "ap-northeast-1"] = config.value(default="us-east-1") |
106 changes: 106 additions & 0 deletions
106
...rehouse_sources/backend/temporal/data_imports/sources/recall_ai/canonical_descriptions.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,106 @@ | ||
| from products.warehouse_sources.backend.temporal.data_imports.sources.common.canonical_descriptions import ( | ||
| CanonicalDescriptions, | ||
| ) | ||
|
|
||
| CANONICAL_DESCRIPTIONS: CanonicalDescriptions = { | ||
| "bots": { | ||
| "description": "A meeting bot sent to join and record a call. One row per bot, with its meeting URL, scheduled join time, status history, and references to the recordings it produced.", | ||
| "docs_url": "https://docs.recall.ai/reference/bot_list", | ||
| "columns": { | ||
| "id": "Unique identifier for the bot.", | ||
| "meeting_url": "The URL of the meeting the bot joined. Cleared by Recall.ai a few days after the bot has joined a call.", | ||
| "bot_name": "The display name the bot used in the call.", | ||
| "join_at": "The time the bot was scheduled to join the call. For ad-hoc bots, the time the bot was requested to join.", | ||
| "status_changes": "History of the bot's lifecycle states (joining, in call, done, fatal), each with a timestamp.", | ||
| "recordings": "Recordings produced by this bot, referencing rows in the recordings table.", | ||
| "calendar_meetings": "Calendar meetings associated with this bot when it was dispatched via the Calendar V1 integration.", | ||
| "metadata": "Arbitrary key-value metadata attached to the bot at creation.", | ||
| }, | ||
| }, | ||
| "recordings": { | ||
| "description": "A recording captured by a meeting bot or uploaded through the desktop SDK. One row per recording, with its processing status and shortcuts to the media artifacts produced from it.", | ||
| "docs_url": "https://docs.recall.ai/reference/recording_list", | ||
| "columns": { | ||
| "id": "Unique identifier for the recording.", | ||
| "created_at": "When the recording object was created.", | ||
| "started_at": "When capture started, or null if it never started.", | ||
| "completed_at": "When capture finished, or null while still recording.", | ||
| "status": "Processing status of the recording (processing, paused, done, failed).", | ||
| "media_shortcuts": "Latest artifact per media type (video, audio, transcript, participant events) with temporary download URLs.", | ||
| "bot": "The bot that produced this recording, or null for desktop SDK uploads.", | ||
| "desktop_sdk_upload": "The desktop SDK upload that produced this recording, or null for bot recordings.", | ||
| "expires_at": "When Recall.ai deletes the recording's media.", | ||
| "metadata": "Arbitrary key-value metadata attached to the recording.", | ||
| }, | ||
| }, | ||
| "transcripts": { | ||
| "description": "Transcript artifacts generated from recordings. Rows carry the transcript's status and a temporary download URL for the transcript content, not the transcript text itself.", | ||
| "docs_url": "https://docs.recall.ai/reference/transcript_list", | ||
| "columns": { | ||
| "id": "Unique identifier for the transcript artifact.", | ||
| "recording": "The recording this transcript was generated from.", | ||
| "created_at": "When the transcript artifact was created.", | ||
| "status": "Processing status of the transcript (processing, done, failed).", | ||
| "data": "Temporary download URL for the transcript content. URLs expire and must be re-fetched from the API.", | ||
| "provider": "The speech-to-text provider that produced the transcript.", | ||
| "diarization": "Speaker diarization configuration used for the transcript.", | ||
| }, | ||
| }, | ||
| "participant_events": { | ||
| "description": "Participant event artifacts generated from recordings, covering joins, leaves, active speakers, and chat. Rows carry a temporary download URL for the event log, not the events themselves.", | ||
| "docs_url": "https://docs.recall.ai/reference/participant_events_list", | ||
| "columns": { | ||
| "id": "Unique identifier for the participant events artifact.", | ||
| "recording": "The recording this artifact was generated from.", | ||
| "created_at": "When the artifact was created.", | ||
| "status": "Processing status of the artifact (processing, done, failed).", | ||
| "data": "Temporary download URLs for the participant event logs. URLs expire and must be re-fetched from the API.", | ||
| }, | ||
| }, | ||
| "meeting_metadata": { | ||
| "description": "Meeting metadata artifacts generated from recordings, such as the meeting title and platform details. Rows carry a temporary download URL for the metadata payload.", | ||
| "docs_url": "https://docs.recall.ai/reference/meeting_metadata_list", | ||
| "columns": { | ||
| "id": "Unique identifier for the meeting metadata artifact.", | ||
| "recording": "The recording this artifact was generated from.", | ||
| "created_at": "When the artifact was created.", | ||
| "status": "Processing status of the artifact (processing, done, failed).", | ||
| "data": "Temporary download URL for the meeting metadata payload.", | ||
| }, | ||
| }, | ||
| "calendars": { | ||
| "description": "Calendar accounts connected through the Calendar V2 integration. One row per connected Google Calendar or Microsoft Outlook account. OAuth client secrets and refresh tokens are removed before rows reach the warehouse.", | ||
| "docs_url": "https://docs.recall.ai/reference/calendars_list", | ||
| "columns": { | ||
| "id": "Unique identifier for the connected calendar.", | ||
| "platform": "The calendar provider (google_calendar or microsoft_outlook).", | ||
| "oauth_client_id": "OAuth client ID of the app the calendar was connected with.", | ||
| "oauth_email": "Email address of the account that authorized the connection.", | ||
| "platform_email": "Email address of the calendar on the provider's side.", | ||
| "status": "Connection status (connecting, connected, disconnected).", | ||
| "status_changes": "History of connection status changes.", | ||
| "created_at": "When the calendar was connected.", | ||
| "updated_at": "When the calendar record last changed.", | ||
| }, | ||
| }, | ||
| "calendar_events": { | ||
| "description": "Events synced from connected calendars through the Calendar V2 integration. One row per calendar event, including the detected meeting URL and any bots scheduled for it.", | ||
| "docs_url": "https://docs.recall.ai/reference/calendar_events_list", | ||
| "columns": { | ||
| "id": "Unique identifier for the calendar event.", | ||
| "start_time": "When the event starts.", | ||
| "end_time": "When the event ends.", | ||
| "calendar_id": "The connected calendar this event belongs to.", | ||
| "raw": "The raw event payload from the calendar provider.", | ||
| "platform": "The calendar provider the event came from.", | ||
| "platform_id": "The event's identifier on the provider's side.", | ||
| "ical_uid": "The event's iCalendar UID.", | ||
| "meeting_platform": "The meeting platform detected from the event (zoom, google_meet, microsoft_teams, and others).", | ||
| "meeting_url": "The meeting URL detected in the event.", | ||
| "created_at": "When the event was first synced from the provider.", | ||
| "updated_at": "When the event last changed.", | ||
| "is_deleted": "Whether the event has been deleted on the provider's side.", | ||
| "bots": "Bots scheduled to record this event.", | ||
| }, | ||
| }, | ||
| } |
178 changes: 178 additions & 0 deletions
178
products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/recall_ai.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,178 @@ | ||
| from datetime import UTC, date, datetime | ||
| from typing import Any, Optional | ||
|
|
||
| from posthog.dataclasses import frozen | ||
|
|
||
| from products.warehouse_sources.backend.temporal.data_imports.sources.common.http import make_tracked_session | ||
| from products.warehouse_sources.backend.temporal.data_imports.sources.common.rest_source import ( | ||
| RESTAPIConfig, | ||
| rest_api_resource, | ||
| ) | ||
| from products.warehouse_sources.backend.temporal.data_imports.sources.common.rest_source.paginators import ( | ||
| JSONResponsePaginator, | ||
| ) | ||
| from products.warehouse_sources.backend.temporal.data_imports.sources.common.rest_source.typing import EndpointResource | ||
| from products.warehouse_sources.backend.temporal.data_imports.sources.common.resumable import ResumableSourceManager | ||
| from products.warehouse_sources.backend.temporal.data_imports.sources.common.source_helpers import validate_via_probe | ||
| from products.warehouse_sources.backend.temporal.data_imports.sources.common.typings import SourceResponse | ||
| from products.warehouse_sources.backend.temporal.data_imports.sources.recall_ai.settings import RECALL_AI_ENDPOINTS | ||
|
|
||
| # API keys are region-scoped and each region is its own host. The allowlist pins outbound | ||
| # traffic to *.recall.ai even if a stored config carries an unexpected region value. | ||
| REGIONS = ("us-east-1", "us-west-2", "eu-central-1", "ap-northeast-1") | ||
|
|
||
|
|
||
| @frozen | ||
| class RecallAIResumeConfig: | ||
| next_url: str | ||
|
|
||
|
|
||
| def base_url_for_region(region: str) -> str: | ||
| if region not in REGIONS: | ||
| raise ValueError(f"Unknown Recall.ai region: {region}") | ||
| return f"https://{region}.recall.ai" | ||
|
|
||
|
|
||
| def _to_iso8601(value: Any) -> Optional[str]: | ||
| """Format the incremental watermark for Recall.ai's ISO 8601 datetime filters. | ||
| Truncating to whole seconds only widens the window, so a boundary row is re-fetched | ||
| and deduped on merge rather than skipped.""" | ||
| if value is None: | ||
| return None | ||
| if isinstance(value, datetime): | ||
| utc = value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC) | ||
| return utc.strftime("%Y-%m-%dT%H:%M:%SZ") | ||
| if isinstance(value, date): | ||
| return value.strftime("%Y-%m-%dT00:00:00Z") | ||
| return str(value) | ||
|
|
||
|
|
||
| def _scrub_fields(fields: tuple[str, ...]): | ||
| def scrub(item: dict[str, Any]) -> dict[str, Any]: | ||
| for field in fields: | ||
| item.pop(field, None) | ||
| return item | ||
|
|
||
| return scrub | ||
|
|
||
|
|
||
| def get_resource( | ||
| endpoint: str, should_use_incremental_field: bool, db_incremental_field_last_value: Optional[Any] | ||
| ) -> EndpointResource: | ||
| config = RECALL_AI_ENDPOINTS[endpoint] | ||
|
|
||
| params: dict[str, Any] = dict(config.extra_params) | ||
| # Only send the lower-bound filter once a real watermark exists; the first incremental | ||
| # sync goes out unfiltered and still advances the watermark from the synced rows. | ||
| if should_use_incremental_field and db_incremental_field_last_value is not None and config.incremental_param: | ||
| params[config.incremental_param] = { | ||
| "type": "incremental", | ||
| "cursor_path": config.incremental_field, | ||
| "initial_value": None, | ||
| "convert": _to_iso8601, | ||
| } | ||
|
|
||
| return { | ||
| "name": endpoint, | ||
| "table_name": endpoint, | ||
| "write_disposition": { | ||
| "disposition": "merge", | ||
| "strategy": "upsert", | ||
| } | ||
| if should_use_incremental_field | ||
| else "replace", | ||
| "endpoint": { | ||
| "data_selector": "results", | ||
| "path": config.path, | ||
| "params": params, | ||
| # Every list response carries an absolute next-page URL in its body, which also | ||
| # keeps any filter params on later pages. | ||
| "paginator": JSONResponsePaginator(next_url_path="next"), | ||
| }, | ||
| "table_format": "delta", | ||
| } | ||
|
|
||
|
|
||
| def recall_ai_source( | ||
| api_key: str, | ||
| region: str, | ||
| endpoint: str, | ||
| team_id: int, | ||
| job_id: str, | ||
| resumable_source_manager: ResumableSourceManager[RecallAIResumeConfig], | ||
| db_incremental_field_last_value: Optional[Any], | ||
| should_use_incremental_field: bool = False, | ||
| ) -> SourceResponse: | ||
| endpoint_config = RECALL_AI_ENDPOINTS[endpoint] | ||
|
|
||
| config: RESTAPIConfig = { | ||
| "client": { | ||
| "base_url": base_url_for_region(region), | ||
| # Recall.ai requires the literal "Token " prefix, not "Bearer ". | ||
| "auth": { | ||
| "type": "api_key", | ||
| "api_key": f"Token {api_key}", | ||
| "name": "Authorization", | ||
| "location": "header", | ||
| }, | ||
| "headers": {"Accept": "application/json"}, | ||
|
veria-ai[bot] marked this conversation as resolved.
|
||
| # Pin next-page and resume URLs to the region origin: a forged/off-host `next` | ||
| # link would otherwise carry the Authorization header to an attacker-controlled | ||
| # or internal destination. Pair host-pinning with rejecting redirects outright. | ||
| "allowed_hosts": [], | ||
| "allow_redirects": False, | ||
| # Calendar rows carry OAuth client secrets/refresh tokens (scrubbed below before | ||
| # they reach the warehouse), which aren't covered by the generic sample denylist — | ||
| # disable diagnostic HTTP capture for this endpoint so raw responses are never stored. | ||
| **({"capture": False} if endpoint_config.scrub_fields else {}), | ||
| }, | ||
| "resources": [get_resource(endpoint, should_use_incremental_field, db_incremental_field_last_value)], | ||
| } | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
|
|
||
| initial_paginator_state: Optional[dict[str, Any]] = None | ||
| if resumable_source_manager.can_resume(): | ||
| resume_config = resumable_source_manager.load_state() | ||
| if resume_config is not None: | ||
| initial_paginator_state = {"next_url": resume_config.next_url} | ||
|
|
||
| def save_checkpoint(state: Optional[dict[str, Any]]) -> None: | ||
| if state and state.get("next_url"): | ||
| resumable_source_manager.save_state(RecallAIResumeConfig(next_url=str(state["next_url"]))) | ||
|
|
||
| resource = rest_api_resource( | ||
| config, | ||
| team_id, | ||
| job_id, | ||
| db_incremental_field_last_value, | ||
| resume_hook=save_checkpoint, | ||
| initial_paginator_state=initial_paginator_state, | ||
| ) | ||
| if endpoint_config.scrub_fields: | ||
| resource = resource.add_map(_scrub_fields(endpoint_config.scrub_fields)) | ||
|
|
||
| return SourceResponse( | ||
| name=endpoint, | ||
| items=lambda: resource, | ||
| primary_keys=["id"], | ||
| partition_count=1, | ||
| partition_size=1, | ||
| partition_mode="datetime", | ||
| partition_format="month", | ||
| partition_keys=[endpoint_config.partition_key], | ||
| # The API documents no ordering and its cursor pagination accepts no sort param, so | ||
| # assume nothing: "desc" commits the incremental watermark only when the sync | ||
| # completes, which is correct whatever order rows actually arrive in. | ||
| sort_mode="desc", | ||
| ) | ||
|
|
||
|
|
||
| def validate_credentials(api_key: str, region: str) -> tuple[bool, int | None]: | ||
| return validate_via_probe( | ||
| lambda: make_tracked_session(redact_values=(api_key,)), | ||
| f"{base_url_for_region(region)}/api/v1/bot/", | ||
| headers={"Authorization": f"Token {api_key}"}, | ||
| # The probe sends the key via Authorization, which `requests` already strips on a | ||
| # cross-host redirect, but pin it anyway: it's a one-word defense against a future | ||
| # change to this probe (e.g. a custom header) reintroducing the leak silently. | ||
| allow_redirects=False, | ||
| ) | ||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.