diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md b/products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md index ecc11c226937..04829123d04c 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md @@ -586,6 +586,7 @@ the row lists both. | rapid7_insightvm | HTTP | requests | ✅ | | raygun | HTTP | requests | ✅ | | razorpay | HTTP | requests + `rest_source.RESTClient` | ✅ | +| recall_ai | HTTP | requests + `rest_source.RESTClient` | ✅ | | recharge | HTTP | requests | ✅ | | recreation | HTTP | requests + `rest_source.RESTClient` | ✅ | | recruitee | HTTP | requests | ✅ | @@ -1316,7 +1317,6 @@ doesn't conflict with concurrent PRs. - raygun - rb2b - rd_station_marketing -- recall_ai - reddit - redis - redpanda_cloud diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/recallai.py b/products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/recallai.py index 6d5f48f5a4fd..48436d96090a 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/recallai.py +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/recallai.py @@ -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") diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/canonical_descriptions.py b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/canonical_descriptions.py new file mode 100644 index 000000000000..436b3feba9fe --- /dev/null +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/canonical_descriptions.py @@ -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.", + }, + }, +} diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/recall_ai.py b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/recall_ai.py new file mode 100644 index 000000000000..672d4448783a --- /dev/null +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/recall_ai.py @@ -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"}, + # 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)], + } + + 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, + ) diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/settings.py b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/settings.py new file mode 100644 index 000000000000..dc9b431c676c --- /dev/null +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/settings.py @@ -0,0 +1,88 @@ +from posthog.dataclasses import frozen + +from products.warehouse_sources.backend.types import IncrementalField, IncrementalFieldType + + +@frozen +class RecallAIEndpointConfig: + path: str + # Stable datetime used for Delta partitioning; never an update-tracking column. + partition_key: str + # Server-side lower-bound filter param, and the response field it filters on. None + # means the endpoint has no usable server-side filter and stays full-refresh only. + incremental_param: str | None = None + incremental_field: str | None = None + # Response fields dropped before rows reach the warehouse (credentials the API echoes back). + scrub_fields: tuple[str, ...] = () + extra_params: tuple[tuple[str, str], ...] = () + + +RECALL_AI_ENDPOINTS: dict[str, RecallAIEndpointConfig] = { + # The bot list only filters on join_at (date granularity), and join_at is the scheduled + # join time: a bot scheduled for next week advances a join_at watermark past bots created + # later with earlier join times, so rows would be skipped. Full refresh only. + "bots": RecallAIEndpointConfig( + path="/api/v1/bot/", + partition_key="join_at", + # The bot list paginates by page number unless use_cursor is set; every other list + # endpoint is cursor-only. Opting in keeps all endpoints on the next-URL paginator. + extra_params=(("use_cursor", "true"),), + ), + "recordings": RecallAIEndpointConfig( + path="/api/v1/recording/", + partition_key="created_at", + incremental_param="created_at_after", + incremental_field="created_at", + ), + "transcripts": RecallAIEndpointConfig( + path="/api/v1/transcript/", + partition_key="created_at", + incremental_param="created_at_after", + incremental_field="created_at", + ), + "participant_events": RecallAIEndpointConfig( + path="/api/v1/participant_events/", + partition_key="created_at", + incremental_param="created_at_after", + incremental_field="created_at", + ), + "meeting_metadata": RecallAIEndpointConfig( + path="/api/v1/meeting_metadata/", + partition_key="created_at", + incremental_param="created_at_after", + incremental_field="created_at", + ), + # Calendars only filter on created_at, which misses status changes on existing rows + # (connected -> disconnected). The table is small (one row per connected calendar), + # so full refresh keeps it correct. The API echoes the customer's OAuth app secrets + # back in list responses; those fields never reach the warehouse. + "calendars": RecallAIEndpointConfig( + path="/api/v2/calendars/", + partition_key="created_at", + scrub_fields=("oauth_client_secret", "oauth_refresh_token"), + ), + "calendar_events": RecallAIEndpointConfig( + path="/api/v2/calendar-events/", + partition_key="created_at", + incremental_param="updated_at__gte", + incremental_field="updated_at", + ), +} + +ENDPOINTS = tuple(RECALL_AI_ENDPOINTS.keys()) + +INCREMENTAL_FIELDS: dict[str, list[IncrementalField]] = { + name: ( + [ + { + "label": config.incremental_field, + "type": IncrementalFieldType.DateTime, + "field": config.incremental_field, + "field_type": IncrementalFieldType.DateTime, + } + ] + if config.incremental_param and config.incremental_field + else [] + ) + for name, config in RECALL_AI_ENDPOINTS.items() +} diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/source.py b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/source.py index 5cc86d39eb96..03127c9b384f 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/source.py +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/source.py @@ -1,28 +1,157 @@ -from typing import cast +from typing import Optional, cast -from products.warehouse_sources.backend.facade.source_config import DataWarehouseSourceCategory, SourceConfig -from products.warehouse_sources.backend.temporal.data_imports.sources.common.base import FieldType, SimpleSource +from products.warehouse_sources.backend.facade.source_config import ( + DataWarehouseSourceCategory, + ReleaseStatus, + SourceConfig, + SourceFieldInputConfig, + SourceFieldInputConfigType, + SourceFieldSelectConfig, + SourceFieldSelectConfigOption, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.common.base import FieldType, ResumableSource +from products.warehouse_sources.backend.temporal.data_imports.sources.common.canonical_descriptions import ( + CanonicalDescriptions, +) from products.warehouse_sources.backend.temporal.data_imports.sources.common.registry import SourceRegistry +from products.warehouse_sources.backend.temporal.data_imports.sources.common.resumable import ResumableSourceManager +from products.warehouse_sources.backend.temporal.data_imports.sources.common.schema import ( + SourceSchema, + build_endpoint_schemas, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.common.typings import SourceInputs, SourceResponse from products.warehouse_sources.backend.temporal.data_imports.sources.generated_configs.recallai import ( RecallAISourceConfig, ) +from products.warehouse_sources.backend.temporal.data_imports.sources.recall_ai.recall_ai import ( + REGIONS, + RecallAIResumeConfig, + recall_ai_source, + validate_credentials as validate_recall_ai_credentials, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.recall_ai.settings import ( + ENDPOINTS, + INCREMENTAL_FIELDS, +) from products.warehouse_sources.backend.types import ExternalDataSourceType @SourceRegistry.register -class RecallAISource(SimpleSource[RecallAISourceConfig]): +class RecallAISource(ResumableSource[RecallAISourceConfig, RecallAIResumeConfig]): + lists_tables_without_credentials = True # static endpoint catalog, safe for public docs + # Recall.ai versions per resource (/api/v1/ for bots and artifacts, /api/v2/ for + # calendars) with no workspace-wide version to pin, so the framework's unversioned + # default applies. + api_docs_url = "https://docs.recall.ai/reference" + @property def source_type(self) -> ExternalDataSourceType: return ExternalDataSourceType.RECALLAI + def get_non_retryable_errors(self) -> dict[str, str | None]: + return { + "401 Client Error: Unauthorized for url": "Your Recall.ai API key is invalid, or it belongs to a different region. Check the key and the region, then reconnect.", + "403 Client Error: Forbidden for url": "Your Recall.ai API key doesn't have access to this data. Check the key in your Recall.ai dashboard.", + } + + def get_canonical_descriptions(self) -> CanonicalDescriptions: + from products.warehouse_sources.backend.temporal.data_imports.sources.recall_ai.canonical_descriptions import ( + CANONICAL_DESCRIPTIONS, + ) + + return CANONICAL_DESCRIPTIONS + + def get_schemas( + self, + config: RecallAISourceConfig, + team_id: int, + with_counts: bool = False, + names: list[str] | None = None, + force_refresh: bool = False, + api_version: str | None = None, + ) -> list[SourceSchema]: + return build_endpoint_schemas(ENDPOINTS, INCREMENTAL_FIELDS, names) + + def validate_credentials( + self, + config: RecallAISourceConfig, + team_id: int, + schema_name: Optional[str] = None, + api_version: str | None = None, + ) -> tuple[bool, str | None]: + if config.region not in REGIONS: + return ( + False, + f"Choose one of the supported Recall.ai regions: {', '.join(REGIONS)}.", + ) + + is_valid, _status = validate_recall_ai_credentials(config.api_key, config.region) + if is_valid: + return True, None + + return ( + False, + "Couldn't connect to Recall.ai. Check that your API key is valid and that you selected the region it was created in.", + ) + + def get_resumable_source_manager(self, inputs: SourceInputs) -> ResumableSourceManager[RecallAIResumeConfig]: + return ResumableSourceManager[RecallAIResumeConfig](inputs, RecallAIResumeConfig) + + def source_for_pipeline( + self, + config: RecallAISourceConfig, + resumable_source_manager: ResumableSourceManager[RecallAIResumeConfig], + inputs: SourceInputs, + ) -> SourceResponse: + return recall_ai_source( + api_key=config.api_key, + region=config.region, + endpoint=inputs.schema_name, + team_id=inputs.team_id, + job_id=inputs.job_id, + resumable_source_manager=resumable_source_manager, + should_use_incremental_field=inputs.should_use_incremental_field, + db_incremental_field_last_value=inputs.db_incremental_field_last_value + if inputs.should_use_incremental_field + else None, + ) + @property def get_source_config(self) -> SourceConfig: return SourceConfig( name=ExternalDataSourceType.RECALLAI, category=DataWarehouseSourceCategory.COMMUNICATION, label="Recall.ai", + caption="Sync meeting bots, recordings, transcripts, participant events, and calendar data from Recall.ai.", + docsUrl="https://posthog.com/docs/cdp/sources/recall-ai", iconPath="/static/services/recall_ai.png", keywords=["meetings", "call recording", "transcripts", "bots"], - fields=cast(list[FieldType], []), - unreleasedSource=True, + fields=cast( + list[FieldType], + [ + SourceFieldInputConfig( + name="api_key", + label="API key", + type=SourceFieldInputConfigType.PASSWORD, + required=True, + placeholder="", + secret=True, + caption="Create an API key in your Recall.ai dashboard under API keys.", + ), + SourceFieldSelectConfig( + name="region", + label="Region", + required=True, + defaultValue="us-east-1", + caption="API keys only work in the region they were created in. Find yours in your Recall.ai dashboard URL.", + options=[ + SourceFieldSelectConfigOption(label="US East (us-east-1)", value="us-east-1"), + SourceFieldSelectConfigOption(label="US West (us-west-2)", value="us-west-2"), + SourceFieldSelectConfigOption(label="EU (eu-central-1)", value="eu-central-1"), + SourceFieldSelectConfigOption(label="Japan (ap-northeast-1)", value="ap-northeast-1"), + ], + ), + ], + ), + releaseStatus=ReleaseStatus.ALPHA, ) diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/tests/test_recall_ai.py b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/tests/test_recall_ai.py new file mode 100644 index 000000000000..c12d41a3442a --- /dev/null +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/tests/test_recall_ai.py @@ -0,0 +1,234 @@ +import json +from collections.abc import Iterable +from datetime import UTC, date, datetime +from typing import Any, cast + +import pytest +from unittest.mock import MagicMock, patch + +from requests import Response + +from products.warehouse_sources.backend.temporal.data_imports.sources.common.resumable import ResumableSourceManager +from products.warehouse_sources.backend.temporal.data_imports.sources.recall_ai.recall_ai import ( + RecallAIResumeConfig, + _to_iso8601, + base_url_for_region, + recall_ai_source, +) + +BASE_URL = "https://us-east-1.recall.ai" + + +def _make_http_response(body: dict[str, Any], status_code: int = 200) -> Response: + resp = Response() + resp.status_code = status_code + resp._content = json.dumps(body).encode() + resp.headers["Content-Type"] = "application/json" + resp.url = f"{BASE_URL}/api/v1/recording/" + return resp + + +def _drive( + endpoint: str, + manager: MagicMock, + responses: list[Response], + *, + should_use_incremental_field: bool = False, + db_incremental_field_last_value: Any = None, +) -> tuple[list[str], list[dict[str, Any]], list[dict[str, Any]]]: + """Drive ``recall_ai_source`` with a mocked HTTP session. + + Returns ``(sent_urls, sent_params, rows)``. Params are captured as shallow copies at + send-time because the paginator mutates the Request object in place between pages. + """ + sent_urls: list[str] = [] + sent_params: list[dict[str, Any]] = [] + response_iter = iter(responses) + + def fake_send(request: Any, *_args: Any, **_kwargs: Any) -> Response: + sent_urls.append(request.url) + sent_params.append(dict(request.params or {})) + return next(response_iter) + + with patch( + "products.warehouse_sources.backend.temporal.data_imports.sources.common.rest_source.rest_client.make_tracked_session" + ) as MockSession: + mock_session = MockSession.return_value + mock_session.headers = {} + mock_session.prepare_request.side_effect = lambda req: req + mock_session.send.side_effect = fake_send + + source_response = recall_ai_source( + api_key="test-key", + region="us-east-1", + endpoint=endpoint, + team_id=123, + job_id="test_job", + resumable_source_manager=manager, + db_incremental_field_last_value=db_incremental_field_last_value, + should_use_incremental_field=should_use_incremental_field, + ) + rows = [row for page in cast("Iterable[Any]", source_response.items()) for row in page] + return sent_urls, sent_params, rows + + +def _manager(can_resume: bool = False) -> MagicMock: + manager = MagicMock(spec=ResumableSourceManager) + manager.can_resume.return_value = can_resume + return manager + + +class TestRecallAIRequestShaping: + INCREMENTAL_PARAMS = { + "recordings": "created_at_after", + "transcripts": "created_at_after", + "participant_events": "created_at_after", + "meeting_metadata": "created_at_after", + "calendar_events": "updated_at__gte", + } + + @pytest.mark.parametrize("endpoint", sorted(INCREMENTAL_PARAMS)) + def test_incremental_sync_with_watermark_sends_iso_filter(self, endpoint: str) -> None: + _, sent_params, _ = _drive( + endpoint, + _manager(), + [_make_http_response({"next": None, "previous": None, "results": [{"id": "r1"}]})], + should_use_incremental_field=True, + db_incremental_field_last_value=datetime(2026, 1, 15, 10, 30, 0, tzinfo=UTC), + ) + + assert sent_params[0][self.INCREMENTAL_PARAMS[endpoint]] == "2026-01-15T10:30:00Z" + + @pytest.mark.parametrize("endpoint", sorted(INCREMENTAL_PARAMS)) + @pytest.mark.parametrize( + ("should_use_incremental_field", "last_value"), + [ + # First incremental sync has no watermark yet and must go out unfiltered. + (True, None), + # Full refresh must ignore a stale watermark. + (False, datetime(2026, 1, 15, tzinfo=UTC)), + ], + ) + def test_filter_omitted_without_active_watermark( + self, endpoint: str, should_use_incremental_field: bool, last_value: Any + ) -> None: + _, sent_params, _ = _drive( + endpoint, + _manager(), + [_make_http_response({"next": None, "previous": None, "results": [{"id": "r1"}]})], + should_use_incremental_field=should_use_incremental_field, + db_incremental_field_last_value=last_value, + ) + + assert self.INCREMENTAL_PARAMS[endpoint] not in sent_params[0] + + def test_bots_request_uses_cursor_pagination_and_no_time_filter(self) -> None: + # Without use_cursor the bot list paginates by page number, which skips or repeats + # rows when bots are created mid-walk. Bots also have no server-side created-at + # filter, so no watermark param may ever be attached. + _, sent_params, _ = _drive( + "bots", + _manager(), + [_make_http_response({"next": None, "previous": None, "results": [{"id": "b1"}]})], + should_use_incremental_field=True, + db_incremental_field_last_value=datetime(2026, 1, 15, tzinfo=UTC), + ) + + assert sent_params[0]["use_cursor"] == "true" + assert "join_at_after" not in sent_params[0] + assert "created_at_after" not in sent_params[0] + + +class TestRecallAIPaginationAndResume: + def test_fresh_run_follows_next_urls_and_saves_state_per_page(self) -> None: + manager = _manager() + page2 = f"{BASE_URL}/api/v1/recording/?cursor=page2" + page3 = f"{BASE_URL}/api/v1/recording/?cursor=page3" + responses = [ + _make_http_response({"next": page2, "previous": None, "results": [{"id": "r1"}]}), + _make_http_response({"next": page3, "previous": None, "results": [{"id": "r2"}]}), + _make_http_response({"next": None, "previous": page2, "results": [{"id": "r3"}]}), + ] + + sent_urls, _, rows = _drive("recordings", manager, responses) + + assert sent_urls == [f"{BASE_URL}/api/v1/recording/", page2, page3] + assert [row["id"] for row in rows] == ["r1", "r2", "r3"] + saved = [call.args[0] for call in manager.save_state.call_args_list] + assert saved == [ + RecallAIResumeConfig(next_url=page2), + RecallAIResumeConfig(next_url=page3), + ] + + def test_terminal_single_page_does_not_save_state(self) -> None: + manager = _manager() + + _drive( + "recordings", + manager, + [_make_http_response({"next": None, "previous": None, "results": [{"id": "r1"}]})], + ) + + manager.save_state.assert_not_called() + manager.load_state.assert_not_called() + + def test_resume_starts_at_saved_next_url(self) -> None: + manager = _manager(can_resume=True) + saved_url = f"{BASE_URL}/api/v1/recording/?cursor=saved" + manager.load_state.return_value = RecallAIResumeConfig(next_url=saved_url) + + sent_urls, _, _ = _drive( + "recordings", + manager, + [_make_http_response({"next": None, "previous": None, "results": [{"id": "r9"}]})], + ) + + assert sent_urls == [saved_url] + manager.load_state.assert_called_once() + + +class TestRecallAICalendarScrubbing: + def test_oauth_secrets_never_reach_the_warehouse(self) -> None: + calendar = { + "id": "cal-1", + "platform": "google_calendar", + "oauth_client_id": "client-id", + "oauth_client_secret": "super-secret", + "oauth_refresh_token": "refresh-token", + "status": "connected", + } + + _, _, rows = _drive( + "calendars", + _manager(), + [_make_http_response({"next": None, "previous": None, "results": [calendar]})], + ) + + assert rows == [ + { + "id": "cal-1", + "platform": "google_calendar", + "oauth_client_id": "client-id", + "status": "connected", + } + ] + + +class TestRecallAIHelpers: + @pytest.mark.parametrize( + ("value", "expected"), + [ + (datetime(2026, 1, 15, 10, 30, 45, tzinfo=UTC), "2026-01-15T10:30:45Z"), + # Naive datetimes from the warehouse watermark are UTC. + (datetime(2026, 1, 15, 10, 30, 45), "2026-01-15T10:30:45Z"), + (date(2026, 1, 15), "2026-01-15T00:00:00Z"), + ("2026-01-15T10:30:45Z", "2026-01-15T10:30:45Z"), + (None, None), + ], + ) + def test_to_iso8601(self, value: Any, expected: str | None) -> None: + assert _to_iso8601(value) == expected + + def test_unknown_region_raises_before_any_request(self) -> None: + with pytest.raises(ValueError, match="Unknown Recall.ai region"): + base_url_for_region("us-central-99") diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/tests/test_recall_ai_source.py b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/tests/test_recall_ai_source.py new file mode 100644 index 000000000000..8d35f847d065 --- /dev/null +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/recall_ai/tests/test_recall_ai_source.py @@ -0,0 +1,50 @@ +import pytest +from unittest.mock import patch + +from products.warehouse_sources.backend.temporal.data_imports.sources.generated_configs.recallai import ( + RecallAISourceConfig, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.recall_ai.source import RecallAISource + +_VALIDATE = ( + "products.warehouse_sources.backend.temporal.data_imports.sources.recall_ai.source.validate_recall_ai_credentials" +) + + +class TestRecallAIValidateCredentials: + @pytest.mark.parametrize( + "region", + [ + "us-east-1.recall.ai", + "https://us-east-1.recall.ai", + "us-central-1", + "", + ], + ) + def test_rejects_unknown_region_without_calling_the_api(self, region: str) -> None: + config = RecallAISourceConfig(api_key="key", region=region) # type: ignore[arg-type] + with patch(_VALIDATE) as mock_validate: + is_valid, message = RecallAISource().validate_credentials(config, team_id=1) + assert is_valid is False + assert message is not None and "us-east-1" in message + mock_validate.assert_not_called() + + @pytest.mark.parametrize( + ("probe_result", "expected_valid"), + [ + ((True, 200), True), + ((False, 401), False), + ], + ) + def test_maps_probe_result_to_user_message(self, probe_result: tuple[bool, int], expected_valid: bool) -> None: + config = RecallAISourceConfig(api_key="key", region="eu-central-1") + with patch(_VALIDATE, return_value=probe_result) as mock_validate: + is_valid, message = RecallAISource().validate_credentials(config, team_id=1) + + assert is_valid is expected_valid + mock_validate.assert_called_once_with("key", "eu-central-1") + if expected_valid: + assert message is None + else: + # The API key and the region must match, so a failed probe points at both. + assert message is not None and "region" in message