{
  "cells": [
    {
      "cell_type": "markdown",
      "id": "7291fe19-e630-4945-be60-595dcdf9140e",
      "metadata": {},
      "source": [
        "# Salesforce to OCI Streaming"
      ]
    },
    {
      "cell_type": "code",
      "execution_count": null,
      "id": "454170c0-8083-491a-ba8c-f58a771e5776",
      "metadata": {
        "command_metadata": {
          "end_time": 1781876903592.94,
          "start_time": 1781875944828.1987
        },
        "execution": {
          "iopub.status.busy": "2026-06-19T13:48:40.461Z"
        },
        "result_type": "result",
        "trusted": true,
        "type": "python"
      },
      "outputs": [],
      "source": [
        "%pip install -r /Workspace/Shared/Eloi/salesforce/requirements.txt"
      ]
    },
    {
      "cell_type": "code",
      "execution_count": null,
      "id": "a14cf1b3-7a77-4014-a438-28a46e705693",
      "metadata": {
        "command_metadata": {
          "end_time": 1781879843453.6794,
          "start_time": 1781879841215.9138
        },
        "execution": {
          "iopub.status.busy": "2026-06-19T14:37:22.399Z"
        },
        "result_type": "result",
        "trusted": true,
        "type": "python"
      },
      "outputs": [],
      "source": [
        "\"\"\"Embedded sf_to_oracle bootstrap for Oracle AI Data Platform Workbench.\n",
        "\n",
        "Paste this file into the first notebook cell, or run it as a cell with\n",
        "``%run`` if the file itself is uploaded beside the notebook.\n",
        "\n",
        "It intentionally avoids Path.cwd(), directory discovery, sys.path mutation,\n",
        "importlib file loaders, and direct reads from /Workspace/.../src.  The\n",
        "sf_to_oracle project modules are registered from in-memory source strings.\n",
        "\"\"\"\n",
        "\n",
        "from __future__ import annotations\n",
        "\n",
        "import sys\n",
        "import types\n",
        "\n",
        "\n",
        "MODULE_SOURCES: dict[str, str] = {\n",
        "    \"sf_to_oracle\": r'''\n",
        "__all__ = [\"__version__\"]\n",
        "\n",
        "__version__ = \"0.1.0\"\n",
        "''',\n",
        "    \"sf_to_oracle.config\": r'''\n",
        "''',\n",
        "    \"sf_to_oracle.config.kafka_properties\": r'''\n",
        "import re\n",
        "from pathlib import Path\n",
        "\n",
        "\n",
        "JAAS_USERNAME_RE = re.compile(r'username=\"([^\"]+)\"')\n",
        "JAAS_PASSWORD_RE = re.compile(r'password=\"([^\"]+)\"')\n",
        "\n",
        "\n",
        "def parse_client_properties(path: Path) -> dict[str, str]:\n",
        "    properties: dict[str, str] = {}\n",
        "    if not path.exists():\n",
        "        return properties\n",
        "\n",
        "    for raw_line in path.read_text(encoding=\"utf-8\").splitlines():\n",
        "        line = raw_line.strip()\n",
        "        if not line or line.startswith(\"#\"):\n",
        "            continue\n",
        "        if \"=\" not in line:\n",
        "            continue\n",
        "        key, value = line.split(\"=\", 1)\n",
        "        properties[key.strip()] = value.strip()\n",
        "\n",
        "    jaas_config = properties.get(\"sasl.jaas.config\", \"\")\n",
        "    username_match = JAAS_USERNAME_RE.search(jaas_config)\n",
        "    password_match = JAAS_PASSWORD_RE.search(jaas_config)\n",
        "    if username_match:\n",
        "        properties[\"sasl.username\"] = username_match.group(1)\n",
        "    if password_match:\n",
        "        properties[\"sasl.password\"] = password_match.group(1)\n",
        "    return properties\n",
        "''',\n",
        "    \"sf_to_oracle.config.settings\": r'''\n",
        "from functools import lru_cache\n",
        "from pathlib import Path\n",
        "\n",
        "from pydantic import Field\n",
        "from pydantic_settings import BaseSettings, SettingsConfigDict\n",
        "\n",
        "from sf_to_oracle.config.kafka_properties import parse_client_properties\n",
        "\n",
        "\n",
        "class Settings(BaseSettings):\n",
        "    model_config = SettingsConfigDict(env_file=\".env\", env_file_encoding=\"utf-8\", extra=\"ignore\")\n",
        "\n",
        "    app_env: str = \"local\"\n",
        "    log_level: str = \"INFO\"\n",
        "    oci_config_file: Path = Path(\"~/.oci/config\")\n",
        "    oci_config_profile: str = \"ATEAMCPT\"\n",
        "\n",
        "    oci_stream_bootstrap_servers: str | None = None\n",
        "    oci_stream_topic: str = \"salesforce-cdc\"\n",
        "    oci_stream_client_properties: Path | None = Path(\"config/client.properties\")\n",
        "    oci_stream_username: str | None = None\n",
        "    oci_stream_password_secret_ocid: str | None = None\n",
        "    oci_stream_password: str | None = None\n",
        "    oci_stream_security_protocol: str = \"SASL_SSL\"\n",
        "    oci_stream_sasl_mechanism: str = \"SCRAM-SHA-512\"\n",
        "    oci_stream_group_id: str = \"sf-oracle-ai-platform\"\n",
        "\n",
        "    oci_vault_sf_client_id_secret_ocid: str | None = None\n",
        "    oci_vault_sf_client_secret_secret_ocid: str | None = None\n",
        "\n",
        "    salesforce_login_url: str = \"https://oracle2.my.salesforce.com\"\n",
        "    salesforce_cdc_channels: str = Field(\n",
        "        default=\"/data/AccountChangeEvent,/data/ContactChangeEvent,/data/OpportunityChangeEvent,/data/CaseChangeEvent\"\n",
        "    )\n",
        "\n",
        "    @property\n",
        "    def expanded_oci_config_file(self) -> str:\n",
        "        return str(self.oci_config_file.expanduser())\n",
        "\n",
        "    @property\n",
        "    def cdc_channels(self) -> list[str]:\n",
        "        return [channel.strip() for channel in self.salesforce_cdc_channels.split(\",\") if channel.strip()]\n",
        "\n",
        "    @property\n",
        "    def kafka_properties(self) -> dict[str, str]:\n",
        "        file_properties = parse_client_properties(self.oci_stream_client_properties) if self.oci_stream_client_properties else {}\n",
        "        return {\n",
        "            \"security.protocol\": file_properties.get(\"security.protocol\", self.oci_stream_security_protocol),\n",
        "            \"sasl.mechanism\": file_properties.get(\"sasl.mechanism\", self.oci_stream_sasl_mechanism),\n",
        "            \"sasl.username\": file_properties.get(\"sasl.username\", self.oci_stream_username or \"\"),\n",
        "            \"sasl.password\": file_properties.get(\"sasl.password\", self.oci_stream_password or \"\"),\n",
        "        }\n",
        "\n",
        "\n",
        "@lru_cache\n",
        "def get_settings() -> Settings:\n",
        "    return Settings()\n",
        "''',\n",
        "    \"sf_to_oracle.config.vault\": r'''\n",
        "import base64\n",
        "import os\n",
        "from dataclasses import dataclass\n",
        "\n",
        "\n",
        "class SecretProvider:\n",
        "    def get_secret(self, secret_ocid: str) -> str:\n",
        "        raise NotImplementedError\n",
        "\n",
        "\n",
        "@dataclass(frozen=True)\n",
        "class EnvironmentSecretProvider(SecretProvider):\n",
        "    \"\"\"Environment-backed provider used by tests and local development.\"\"\"\n",
        "\n",
        "    prefix: str = \"SECRET_\"\n",
        "\n",
        "    def get_secret(self, secret_ocid: str) -> str:\n",
        "        key = f\"{self.prefix}{secret_ocid}\".upper().replace(\"-\", \"_\")\n",
        "        return os.environ.get(key, \"\")\n",
        "\n",
        "\n",
        "class OciVaultSecretProvider(SecretProvider):\n",
        "    def __init__(self, config_file: str | None = None, profile: str | None = None):\n",
        "        try:\n",
        "            import oci\n",
        "        except ImportError as exc:\n",
        "            raise RuntimeError(\"Install cloud dependencies with `pip install -e '.[cloud]'`.\") from exc\n",
        "\n",
        "        self._oci = oci\n",
        "        config = (\n",
        "            oci.config.from_file(file_location=config_file, profile_name=profile)\n",
        "            if config_file or profile\n",
        "            else oci.config.from_file()\n",
        "        )\n",
        "        self._client = oci.secrets.SecretsClient(config)\n",
        "\n",
        "    def get_secret(self, secret_ocid: str) -> str:\n",
        "        bundle = self._client.get_secret_bundle(secret_ocid).data\n",
        "        encoded = bundle.secret_bundle_content.content\n",
        "        return base64.b64decode(encoded).decode(\"utf-8\")\n",
        "''',\n",
        "    \"sf_to_oracle.schemas\": r'''\n",
        "from sf_to_oracle.schemas.cdc import (\n",
        "    AccountPayload,\n",
        "    CasePayload,\n",
        "    ChangeType,\n",
        "    ContactPayload,\n",
        "    CdcEvent,\n",
        "    EntityName,\n",
        "    OpportunityPayload,\n",
        ")\n",
        "\n",
        "__all__ = [\n",
        "    \"AccountPayload\",\n",
        "    \"CasePayload\",\n",
        "    \"ChangeType\",\n",
        "    \"ContactPayload\",\n",
        "    \"CdcEvent\",\n",
        "    \"EntityName\",\n",
        "    \"OpportunityPayload\",\n",
        "]\n",
        "''',\n",
        "    \"sf_to_oracle.schemas.cdc\": r'''\n",
        "from datetime import datetime, timezone\n",
        "from enum import StrEnum\n",
        "from typing import Any, Literal\n",
        "from uuid import uuid4\n",
        "\n",
        "from pydantic import BaseModel, ConfigDict, Field\n",
        "\n",
        "\n",
        "class EntityName(StrEnum):\n",
        "    ACCOUNT = \"Account\"\n",
        "    CONTACT = \"Contact\"\n",
        "    OPPORTUNITY = \"Opportunity\"\n",
        "    CASE = \"Case\"\n",
        "\n",
        "\n",
        "class ChangeType(StrEnum):\n",
        "    CREATE = \"CREATE\"\n",
        "    UPDATE = \"UPDATE\"\n",
        "    DELETE = \"DELETE\"\n",
        "    UNDELETE = \"UNDELETE\"\n",
        "\n",
        "\n",
        "class AccountPayload(BaseModel):\n",
        "    id: str\n",
        "    name: str\n",
        "    industry: str | None = None\n",
        "    annual_revenue: float | None = None\n",
        "    owner_id: str | None = None\n",
        "    billing_country: str | None = None\n",
        "\n",
        "\n",
        "class ContactPayload(BaseModel):\n",
        "    id: str\n",
        "    account_id: str\n",
        "    first_name: str | None = None\n",
        "    last_name: str | None = None\n",
        "    email: str | None = None\n",
        "    title: str | None = None\n",
        "\n",
        "\n",
        "class OpportunityPayload(BaseModel):\n",
        "    id: str\n",
        "    account_id: str\n",
        "    name: str\n",
        "    stage_name: str\n",
        "    amount: float | None = None\n",
        "    close_date: str | None = None\n",
        "\n",
        "\n",
        "class CasePayload(BaseModel):\n",
        "    id: str\n",
        "    account_id: str\n",
        "    case_number: str\n",
        "    subject: str\n",
        "    status: str\n",
        "    priority: str | None = None\n",
        "\n",
        "\n",
        "Payload = AccountPayload | ContactPayload | OpportunityPayload | CasePayload | dict[str, Any]\n",
        "\n",
        "\n",
        "class CdcEvent(BaseModel):\n",
        "    model_config = ConfigDict(use_enum_values=True)\n",
        "\n",
        "    event_id: str = Field(default_factory=lambda: str(uuid4()))\n",
        "    entity: EntityName\n",
        "    change_type: ChangeType = ChangeType.UPDATE\n",
        "    commit_timestamp: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))\n",
        "    record_id: str\n",
        "    replay_id: str | None = None\n",
        "    payload: Payload\n",
        "    source: Literal[\"salesforce_cdc\"] = \"salesforce_cdc\"\n",
        "\n",
        "    @property\n",
        "    def event_date(self) -> str:\n",
        "        return self.commit_timestamp.date().isoformat()\n",
        "''',\n",
        "    \"sf_to_oracle.ingest\": r'''\n",
        "''',\n",
        "    \"sf_to_oracle.ingest.salesforce_auth\": r'''\n",
        "from dataclasses import dataclass\n",
        "\n",
        "from sf_to_oracle.config.settings import Settings\n",
        "from sf_to_oracle.config.vault import SecretProvider\n",
        "\n",
        "\n",
        "@dataclass(frozen=True)\n",
        "class SalesforceCredentials:\n",
        "    client_id: str\n",
        "    client_secret: str\n",
        "\n",
        "\n",
        "def resolve_salesforce_credentials(settings: Settings, secret_provider: SecretProvider) -> SalesforceCredentials:\n",
        "    if not settings.oci_vault_sf_client_id_secret_ocid or not settings.oci_vault_sf_client_secret_secret_ocid:\n",
        "        raise ValueError(\"Salesforce client ID and client secret OCIDs are required.\")\n",
        "\n",
        "    return SalesforceCredentials(\n",
        "        client_id=secret_provider.get_secret(settings.oci_vault_sf_client_id_secret_ocid),\n",
        "        client_secret=secret_provider.get_secret(settings.oci_vault_sf_client_secret_secret_ocid),\n",
        "    )\n",
        "''',\n",
        "    \"sf_to_oracle.ingest.salesforce_rest\": r'''\n",
        "from datetime import datetime, timezone\n",
        "from typing import Any\n",
        "\n",
        "import requests\n",
        "\n",
        "from sf_to_oracle.ingest.salesforce_auth import SalesforceCredentials\n",
        "from sf_to_oracle.schemas import (\n",
        "    AccountPayload,\n",
        "    CasePayload,\n",
        "    ChangeType,\n",
        "    ContactPayload,\n",
        "    CdcEvent,\n",
        "    EntityName,\n",
        "    OpportunityPayload,\n",
        ")\n",
        "\n",
        "\n",
        "class SalesforceRestClient:\n",
        "    def __init__(self, login_url: str, credentials: SalesforceCredentials, api_version: str = \"60.0\"):\n",
        "        self.login_url = login_url.rstrip(\"/\")\n",
        "        self.credentials = credentials\n",
        "        self.api_version = api_version\n",
        "        self.session_id: str | None = None\n",
        "        self.instance_url: str | None = None\n",
        "\n",
        "    def login(self) -> None:\n",
        "        response = requests.post(\n",
        "            f\"{self.login_url}/services/oauth2/token\",\n",
        "            data={\n",
        "                \"grant_type\": \"client_credentials\",\n",
        "                \"client_id\": self.credentials.client_id,\n",
        "                \"client_secret\": self.credentials.client_secret,\n",
        "            },\n",
        "            timeout=30,\n",
        "        )\n",
        "        if response.status_code >= 400:\n",
        "            raise RuntimeError(response.text)\n",
        "        payload = response.json()\n",
        "        self.session_id = payload[\"access_token\"]\n",
        "        self.instance_url = payload[\"instance_url\"].rstrip(\"/\")\n",
        "\n",
        "    def query(self, soql: str) -> list[dict[str, Any]]:\n",
        "        if not self.session_id or not self.instance_url:\n",
        "            self.login()\n",
        "        assert self.session_id\n",
        "        assert self.instance_url\n",
        "\n",
        "        url = f\"{self.instance_url}/services/data/v{self.api_version}/query\"\n",
        "        records: list[dict[str, Any]] = []\n",
        "        while url:\n",
        "            response = requests.get(\n",
        "                url,\n",
        "                params={\"q\": soql} if not records else None,\n",
        "                headers={\"Authorization\": f\"Bearer {self.session_id}\"},\n",
        "                timeout=30,\n",
        "            )\n",
        "            if response.status_code >= 400:\n",
        "                raise RuntimeError(response.text)\n",
        "            payload = response.json()\n",
        "            records.extend(payload.get(\"records\", []))\n",
        "            next_url = payload.get(\"nextRecordsUrl\")\n",
        "            url = f\"{self.instance_url}{next_url}\" if next_url else \"\"\n",
        "        return records\n",
        "\n",
        "\n",
        "def fetch_salesforce_snapshot_events(client: SalesforceRestClient, limit_per_object: int = 5) -> list[CdcEvent]:\n",
        "    events: list[CdcEvent] = []\n",
        "    events.extend(_account_events(client.query(_limit(\"SELECT Id, Name, Industry, AnnualRevenue, OwnerId, BillingCountry, LastModifiedDate FROM Account ORDER BY LastModifiedDate DESC\", limit_per_object))))\n",
        "    events.extend(_contact_events(client.query(_limit(\"SELECT Id, AccountId, FirstName, LastName, Email, Title, LastModifiedDate FROM Contact WHERE AccountId != null ORDER BY LastModifiedDate DESC\", limit_per_object))))\n",
        "    events.extend(_opportunity_events(client.query(_limit(\"SELECT Id, AccountId, Name, StageName, Amount, CloseDate, LastModifiedDate FROM Opportunity WHERE AccountId != null ORDER BY LastModifiedDate DESC\", limit_per_object))))\n",
        "    events.extend(_case_events(client.query(_limit(\"SELECT Id, AccountId, CaseNumber, Subject, Status, Priority, LastModifiedDate FROM Case WHERE AccountId != null ORDER BY LastModifiedDate DESC\", limit_per_object))))\n",
        "    return events\n",
        "\n",
        "\n",
        "def fetch_salesforce_incremental_events(\n",
        "    client: SalesforceRestClient,\n",
        "    since: datetime,\n",
        "    limit_per_object: int = 200,\n",
        ") -> list[CdcEvent]:\n",
        "    since_soql = _soql_datetime(since)\n",
        "    events: list[CdcEvent] = []\n",
        "    events.extend(_account_events(client.query(_limit(f\"SELECT Id, Name, Industry, AnnualRevenue, OwnerId, BillingCountry, LastModifiedDate FROM Account WHERE LastModifiedDate > {since_soql} ORDER BY LastModifiedDate ASC\", limit_per_object))))\n",
        "    events.extend(_contact_events(client.query(_limit(f\"SELECT Id, AccountId, FirstName, LastName, Email, Title, LastModifiedDate FROM Contact WHERE AccountId != null AND LastModifiedDate > {since_soql} ORDER BY LastModifiedDate ASC\", limit_per_object))))\n",
        "    events.extend(_opportunity_events(client.query(_limit(f\"SELECT Id, AccountId, Name, StageName, Amount, CloseDate, LastModifiedDate FROM Opportunity WHERE AccountId != null AND LastModifiedDate > {since_soql} ORDER BY LastModifiedDate ASC\", limit_per_object))))\n",
        "    events.extend(_case_events(client.query(_limit(f\"SELECT Id, AccountId, CaseNumber, Subject, Status, Priority, LastModifiedDate FROM Case WHERE AccountId != null AND LastModifiedDate > {since_soql} ORDER BY LastModifiedDate ASC\", limit_per_object))))\n",
        "    return sorted(events, key=lambda event: event.commit_timestamp)\n",
        "\n",
        "\n",
        "def _limit(soql: str, limit_per_object: int) -> str:\n",
        "    return f\"{soql} LIMIT {limit_per_object}\"\n",
        "\n",
        "\n",
        "def _soql_datetime(value: datetime) -> str:\n",
        "    return value.astimezone(timezone.utc).strftime(\"%Y-%m-%dT%H:%M:%SZ\")\n",
        "\n",
        "\n",
        "def _commit_timestamp(record: dict[str, Any]) -> datetime:\n",
        "    value = record.get(\"LastModifiedDate\")\n",
        "    if not value:\n",
        "        return datetime.now(timezone.utc)\n",
        "    return datetime.fromisoformat(value.replace(\"Z\", \"+00:00\"))\n",
        "\n",
        "\n",
        "def _account_events(records: list[dict[str, Any]]) -> list[CdcEvent]:\n",
        "    return [\n",
        "        CdcEvent(\n",
        "            entity=EntityName.ACCOUNT,\n",
        "            change_type=ChangeType.UPDATE,\n",
        "            commit_timestamp=_commit_timestamp(record),\n",
        "            record_id=record[\"Id\"],\n",
        "            payload=AccountPayload(\n",
        "                id=record[\"Id\"],\n",
        "                name=record.get(\"Name\") or \"\",\n",
        "                industry=record.get(\"Industry\"),\n",
        "                annual_revenue=record.get(\"AnnualRevenue\"),\n",
        "                owner_id=record.get(\"OwnerId\"),\n",
        "                billing_country=record.get(\"BillingCountry\"),\n",
        "            ),\n",
        "        )\n",
        "        for record in records\n",
        "    ]\n",
        "\n",
        "\n",
        "def _contact_events(records: list[dict[str, Any]]) -> list[CdcEvent]:\n",
        "    return [\n",
        "        CdcEvent(\n",
        "            entity=EntityName.CONTACT,\n",
        "            change_type=ChangeType.UPDATE,\n",
        "            commit_timestamp=_commit_timestamp(record),\n",
        "            record_id=record[\"Id\"],\n",
        "            payload=ContactPayload(\n",
        "                id=record[\"Id\"],\n",
        "                account_id=record[\"AccountId\"],\n",
        "                first_name=record.get(\"FirstName\"),\n",
        "                last_name=record.get(\"LastName\"),\n",
        "                email=record.get(\"Email\"),\n",
        "                title=record.get(\"Title\"),\n",
        "            ),\n",
        "        )\n",
        "        for record in records\n",
        "    ]\n",
        "\n",
        "\n",
        "def _opportunity_events(records: list[dict[str, Any]]) -> list[CdcEvent]:\n",
        "    return [\n",
        "        CdcEvent(\n",
        "            entity=EntityName.OPPORTUNITY,\n",
        "            change_type=ChangeType.UPDATE,\n",
        "            commit_timestamp=_commit_timestamp(record),\n",
        "            record_id=record[\"Id\"],\n",
        "            payload=OpportunityPayload(\n",
        "                id=record[\"Id\"],\n",
        "                account_id=record[\"AccountId\"],\n",
        "                name=record.get(\"Name\") or \"\",\n",
        "                stage_name=record.get(\"StageName\") or \"\",\n",
        "                amount=record.get(\"Amount\"),\n",
        "                close_date=record.get(\"CloseDate\"),\n",
        "            ),\n",
        "        )\n",
        "        for record in records\n",
        "    ]\n",
        "\n",
        "\n",
        "def _case_events(records: list[dict[str, Any]]) -> list[CdcEvent]:\n",
        "    return [\n",
        "        CdcEvent(\n",
        "            entity=EntityName.CASE,\n",
        "            change_type=ChangeType.UPDATE,\n",
        "            commit_timestamp=_commit_timestamp(record),\n",
        "            record_id=record[\"Id\"],\n",
        "            payload=CasePayload(\n",
        "                id=record[\"Id\"],\n",
        "                account_id=record[\"AccountId\"],\n",
        "                case_number=record.get(\"CaseNumber\") or \"\",\n",
        "                subject=record.get(\"Subject\") or \"\",\n",
        "                status=record.get(\"Status\") or \"\",\n",
        "                priority=record.get(\"Priority\"),\n",
        "            ),\n",
        "        )\n",
        "        for record in records\n",
        "    ]\n",
        "''',\n",
        "    \"sf_to_oracle.ingest.salesforce_cdc\": r'''\n",
        "from collections.abc import Iterable\n",
        "\n",
        "from sf_to_oracle.config.settings import Settings\n",
        "from sf_to_oracle.config.vault import SecretProvider\n",
        "from sf_to_oracle.ingest.salesforce_auth import resolve_salesforce_credentials\n",
        "from sf_to_oracle.schemas import CdcEvent\n",
        "\n",
        "\n",
        "class SalesforceCdcSource:\n",
        "    \"\"\"Insertion point for a production Salesforce CDC client.\n",
        "\n",
        "    Use Salesforce Streaming API, Pub/Sub API, or a CometD/Bayeux client here,\n",
        "    normalize each message to `CdcEvent`, and publish it to OCI Streaming.\n",
        "    \"\"\"\n",
        "\n",
        "    def __init__(self, settings: Settings, secret_provider: SecretProvider):\n",
        "        self.settings = settings\n",
        "        self.credentials = resolve_salesforce_credentials(settings, secret_provider)\n",
        "        self.channels = settings.cdc_channels\n",
        "        self.secret_provider = secret_provider\n",
        "\n",
        "    @classmethod\n",
        "    def from_channels(cls, channels: list[str], secret_provider: SecretProvider):\n",
        "        \"\"\"Compatibility helper for tests and examples that only need channel metadata.\"\"\"\n",
        "        settings = Settings(\n",
        "            salesforce_cdc_channels=\",\".join(channels),\n",
        "            oci_vault_sf_username_secret_ocid=\"username\",\n",
        "            oci_vault_sf_password_secret_ocid=\"password\",\n",
        "        )\n",
        "        return cls(settings, secret_provider)\n",
        "\n",
        "    def events(self) -> Iterable[CdcEvent]:\n",
        "        raise NotImplementedError(\"Wire this adapter to your Salesforce CDC client of choice.\")\n",
        "''',\n",
        "    \"sf_to_oracle.stream\": r'''\n",
        "''',\n",
        "    \"sf_to_oracle.stream.kafka\": r'''\n",
        "def _kafka_python_security_protocol(protocol: str) -> str:\n",
        "    return protocol.replace(\"_\", \"_\")\n",
        "\n",
        "\n",
        "def _kafka_python_sasl_mechanism(mechanism: str) -> str:\n",
        "    normalized = mechanism.upper().replace(\"_\", \"-\")\n",
        "    if normalized in {\"SCRAM-SHA-256\", \"SCRAM-SHA-512\", \"PLAIN\"}:\n",
        "        return normalized\n",
        "    raise ValueError(f\"Unsupported Kafka SASL mechanism for kafka-python: {mechanism}\")\n",
        "''',\n",
        "}\n",
        "\n",
        "\n",
        "PACKAGE_NAMES = {\n",
        "    \"sf_to_oracle\",\n",
        "    \"sf_to_oracle.config\",\n",
        "    \"sf_to_oracle.ingest\",\n",
        "    \"sf_to_oracle.schemas\",\n",
        "    \"sf_to_oracle.stream\",\n",
        "}\n",
        "\n",
        "\n",
        "LOAD_ORDER = [\n",
        "    \"sf_to_oracle\",\n",
        "    \"sf_to_oracle.config\",\n",
        "    \"sf_to_oracle.config.kafka_properties\",\n",
        "    \"sf_to_oracle.config.vault\",\n",
        "    \"sf_to_oracle.config.settings\",\n",
        "    \"sf_to_oracle.schemas.cdc\",\n",
        "    \"sf_to_oracle.schemas\",\n",
        "    \"sf_to_oracle.ingest\",\n",
        "    \"sf_to_oracle.ingest.salesforce_auth\",\n",
        "    \"sf_to_oracle.ingest.salesforce_rest\",\n",
        "    \"sf_to_oracle.ingest.salesforce_cdc\",\n",
        "    \"sf_to_oracle.stream\",\n",
        "    \"sf_to_oracle.stream.kafka\",\n",
        "]\n",
        "\n",
        "\n",
        "def _ensure_package(name: str) -> types.ModuleType:\n",
        "    module = sys.modules.get(name)\n",
        "    if module is None:\n",
        "        module = types.ModuleType(name)\n",
        "        module.__package__ = name\n",
        "        module.__path__ = []\n",
        "        sys.modules[name] = module\n",
        "    return module\n",
        "\n",
        "\n",
        "def _register_module(name: str, source: str) -> types.ModuleType:\n",
        "    parent_name, _, child_name = name.rpartition(\".\")\n",
        "    if parent_name:\n",
        "        _ensure_package(parent_name)\n",
        "\n",
        "    module = sys.modules.get(name)\n",
        "    if module is None:\n",
        "        module = types.ModuleType(name)\n",
        "        sys.modules[name] = module\n",
        "\n",
        "    module.__file__ = f\"<embedded:{name}>\"\n",
        "    module.__loader__ = None\n",
        "    module.__package__ = name if name in PACKAGE_NAMES else parent_name\n",
        "    if name in PACKAGE_NAMES:\n",
        "        module.__path__ = []\n",
        "\n",
        "    exec(compile(source, module.__file__, \"exec\"), module.__dict__)\n",
        "\n",
        "    if parent_name:\n",
        "        setattr(sys.modules[parent_name], child_name, module)\n",
        "    return module\n",
        "\n",
        "\n",
        "def install_embedded_sf_to_oracle(verbose: bool = True) -> list[str]:\n",
        "    loaded = []\n",
        "    for name in LOAD_ORDER:\n",
        "        _register_module(name, MODULE_SOURCES[name])\n",
        "        loaded.append(name)\n",
        "    if verbose:\n",
        "        print(f\"Embedded sf_to_oracle loaded: {len(loaded)} modules\")\n",
        "    return loaded\n",
        "\n",
        "\n",
        "install_embedded_sf_to_oracle()\n"
      ]
    },
    {
      "cell_type": "code",
      "execution_count": null,
      "id": "1dba95ca-745b-40d3-825c-a44c9500a88d",
      "metadata": {
        "command_metadata": {
          "end_time": 1781879860304.649,
          "start_time": 1781879858438.4355
        },
        "execution": {
          "iopub.status.busy": "2026-06-19T14:37:39.368Z"
        },
        "result_type": "result",
        "trusted": true,
        "type": "python"
      },
      "outputs": [],
      "source": [
        "# Dependency check: install these in the Workbench environment if this cell fails.\n",
        "required_imports = {\n",
        "    \"pydantic\": \"pydantic\",\n",
        "    \"pydantic_settings\": \"pydantic-settings\",\n",
        "    \"requests\": \"requests\",\n",
        "    \"orjson\": \"orjson\",\n",
        "    \"oci\": \"oci\",\n",
        "    \"kafka\": \"kafka-python\",\n",
        "}\n",
        "missing = []\n",
        "for module_name, package_name in required_imports.items():\n",
        "    try:\n",
        "        __import__(module_name)\n",
        "    except ImportError:\n",
        "        missing.append(package_name)\n",
        "if missing:\n",
        "    raise RuntimeError(f\"Missing notebook dependencies: {', '.join(missing)}\")\n",
        "print(\"Dependencies available\")\n"
      ]
    },
    {
      "cell_type": "code",
      "execution_count": null,
      "id": "ee6e57ae-a6de-4fd8-a2d4-dd4e3ca8b8e6",
      "metadata": {
        "command_metadata": {
          "end_time": 1781879984523.792,
          "start_time": 1781879982415.8665
        },
        "execution": {
          "iopub.status.busy": "2026-06-19T14:39:43.689Z"
        },
        "result_type": "result",
        "trusted": true,
        "type": "python"
      },
      "outputs": [],
      "source": [
        "# Configuration: replace placeholders before publishing.\n",
        "# In Workbench, prefer explicit Kafka credentials instead of relying on a relative config/client.properties file.\n",
        "from sf_to_oracle.config.settings import Settings\n",
        "\n",
        "settings = Settings(\n",
        "    oci_stream_bootstrap_servers=\"<oci streaming with kafka hostname>\",\n",
        "    oci_stream_topic=\"salesforce-cdc\",\n",
        "    oci_stream_group_id=\"sf-oracle-workbench\",\n",
        "    oci_stream_client_properties=None,\n",
        "    oci_stream_username=\"<kafka username>\",\n",
        "    oci_stream_password=\"<kafka password>\",\n",
        "    oci_vault_sf_client_id_secret_ocid=\"<ocid1.vaultsecret.oc1...>\",\n",
        "    oci_vault_sf_client_secret_secret_ocid=\"<ocid1.vaultsecret.oc1...>\",\n",
        "    salesforce_login_url=\"https://<salesforce hostname>.my.salesforce.com\",\n",
        ")\n",
        "\n",
        "placeholders = [\n",
        "    (\"oci_stream_bootstrap_servers\", settings.oci_stream_bootstrap_servers),\n",
        "    (\"oci_stream_username\", settings.oci_stream_username),\n",
        "    (\"oci_stream_password\", settings.oci_stream_password),\n",
        "    (\"oci_vault_sf_client_id_secret_ocid\", settings.oci_vault_sf_client_id_secret_ocid),\n",
        "    (\"oci_vault_sf_client_secret_secret_ocid\", settings.oci_vault_sf_client_secret_secret_ocid),\n",
        "]\n",
        "missing = [name for name, value in placeholders if not value or str(value).startswith(\"<\")]\n",
        "if missing:\n",
        "    raise ValueError(f\"Replace placeholder configuration values: {', '.join(missing)}\")\n",
        "print(f\"Configured topic={settings.oci_stream_topic} bootstrap={settings.oci_stream_bootstrap_servers}\")\n"
      ]
    },
    {
      "cell_type": "code",
      "execution_count": null,
      "id": "e8a386f6-ec53-4454-8a39-ff80f1087703",
      "metadata": {
        "command_metadata": {
          "end_time": 1781880087608.1355,
          "start_time": 1781880085298.3452
        },
        "execution": {
          "iopub.status.busy": "2026-06-19T14:41:26.522Z"
        },
        "result_type": "result",
        "trusted": true,
        "type": "python"
      },
      "outputs": [],
      "source": [
        "# Build clients and validate Kafka topic metadata before publishing.\n",
        "import socket\n",
        "import ssl\n",
        "\n",
        "import orjson\n",
        "from kafka import KafkaProducer\n",
        "from kafka.errors import KafkaTimeoutError\n",
        "\n",
        "from sf_to_oracle.config.vault import OciVaultSecretProvider\n",
        "from sf_to_oracle.ingest.salesforce_auth import resolve_salesforce_credentials\n",
        "from sf_to_oracle.ingest.salesforce_rest import SalesforceRestClient\n",
        "from sf_to_oracle.stream.kafka import _kafka_python_sasl_mechanism, _kafka_python_security_protocol\n",
        "\n",
        "\n",
        "def _parse_bootstrap_server(value: str) -> tuple[str, int]:\n",
        "    first = value.split(\",\", 1)[0].strip()\n",
        "    host, _, port_text = first.rpartition(\":\")\n",
        "    if not host or not port_text:\n",
        "        raise ValueError(f\"Bootstrap server must look like host:port, got {value!r}\")\n",
        "    return host, int(port_text)\n",
        "\n",
        "\n",
        "def assert_tcp_tls_reachable(bootstrap_servers: str, timeout_seconds: int = 15) -> None:\n",
        "    host, port = _parse_bootstrap_server(bootstrap_servers)\n",
        "    with socket.create_connection((host, port), timeout=timeout_seconds) as sock:\n",
        "        context = ssl.create_default_context()\n",
        "        with context.wrap_socket(sock, server_hostname=host):\n",
        "            pass\n",
        "    print(f\"TLS socket reachable: {host}:{port}\")\n",
        "\n",
        "\n",
        "def build_kafka_producer(settings: Settings) -> KafkaProducer:\n",
        "    kafka_properties = settings.kafka_properties\n",
        "    missing = [\n",
        "        name\n",
        "        for name, value in {\n",
        "            \"security.protocol\": kafka_properties.get(\"security.protocol\"),\n",
        "            \"sasl.mechanism\": kafka_properties.get(\"sasl.mechanism\"),\n",
        "            \"sasl.username\": kafka_properties.get(\"sasl.username\"),\n",
        "            \"sasl.password\": kafka_properties.get(\"sasl.password\"),\n",
        "        }.items()\n",
        "        if not value\n",
        "    ]\n",
        "    if missing:\n",
        "        raise ValueError(f\"Missing Kafka configuration values: {', '.join(missing)}\")\n",
        "\n",
        "    return KafkaProducer(\n",
        "        bootstrap_servers=settings.oci_stream_bootstrap_servers,\n",
        "        security_protocol=_kafka_python_security_protocol(kafka_properties[\"security.protocol\"]),\n",
        "        sasl_mechanism=_kafka_python_sasl_mechanism(kafka_properties[\"sasl.mechanism\"]),\n",
        "        sasl_plain_username=kafka_properties[\"sasl.username\"],\n",
        "        sasl_plain_password=kafka_properties[\"sasl.password\"],\n",
        "        key_serializer=lambda value: value.encode(\"utf-8\"),\n",
        "        value_serializer=lambda value: orjson.dumps(value),\n",
        "        api_version_auto_timeout_ms=30000,\n",
        "        request_timeout_ms=30000,\n",
        "        max_block_ms=30000,\n",
        "        retries=3,\n",
        "    )\n",
        "\n",
        "\n",
        "def assert_topic_metadata(producer: KafkaProducer, topic: str) -> None:\n",
        "    try:\n",
        "        partitions = producer.partitions_for(topic)\n",
        "    except KafkaTimeoutError as exc:\n",
        "        raise RuntimeError(\n",
        "            \"Kafka metadata lookup timed out before publishing. Check that the bootstrap host is reachable \"\n",
        "            \"from this Workbench session, the stream/topic name is exact, the stream pool Kafka endpoint \"\n",
        "            \"matches the topic, and the SASL username/password are valid for OCI Streaming.\"\n",
        "        ) from exc\n",
        "    if not partitions:\n",
        "        raise RuntimeError(\n",
        "            f\"Kafka topic metadata was not returned for {topic!r}. Confirm the stream exists in the \"\n",
        "            \"target stream pool and that this principal has permission to use it.\"\n",
        "        )\n",
        "    print(f\"Kafka metadata ready: topic={topic} partitions={sorted(partitions)}\")\n",
        "\n",
        "\n",
        "assert_tcp_tls_reachable(settings.oci_stream_bootstrap_servers)\n",
        "\n",
        "secret_provider = OciVaultSecretProvider(\n",
        "    config_file=\"/Workspace/Shared/Eloi/config\",\n",
        "    profile=\"<PROFILE_NAME>\",\n",
        ")\n",
        "credentials = resolve_salesforce_credentials(settings, secret_provider)\n",
        "salesforce_client = SalesforceRestClient(settings.salesforce_login_url, credentials)\n",
        "\n",
        "producer = build_kafka_producer(settings)\n",
        "assert_topic_metadata(producer, settings.oci_stream_topic)\n",
        "print(\"Clients ready\")\n"
      ]
    },
    {
      "cell_type": "code",
      "execution_count": null,
      "id": "960b21b5-e110-4607-8fdf-6f23953cd4bb",
      "metadata": {
        "command_metadata": {
          "end_time": 1781880120275.6128,
          "start_time": 1781880117470.8823
        },
        "execution": {
          "iopub.status.busy": "2026-06-19T14:41:59.181Z"
        },
        "result_type": "result",
        "trusted": true,
        "type": "python"
      },
      "outputs": [],
      "source": [
        "# Snapshot publish: reads recent Salesforce records and publishes normalized CDC events.\n",
        "from uuid import uuid4\n",
        "\n",
        "from sf_to_oracle.ingest.salesforce_rest import fetch_salesforce_snapshot_events\n",
        "\n",
        "events = fetch_salesforce_snapshot_events(salesforce_client, limit_per_object=2)\n",
        "marker = f\"workbench-snapshot-{uuid4()}\"\n",
        "futures = []\n",
        "for event in events:\n",
        "    event.event_id = f\"{marker}-{event.entity}-{event.record_id}\"\n",
        "    futures.append(\n",
        "        producer.send(\n",
        "            settings.oci_stream_topic,\n",
        "            key=event.record_id,\n",
        "            value=event.model_dump(mode=\"json\"),\n",
        "        )\n",
        "    )\n",
        "for future in futures:\n",
        "    future.get(timeout=30)\n",
        "producer.flush()\n",
        "print(f\"Published {len(events)} Salesforce snapshot events to {settings.oci_stream_topic}\")"
      ]
    },
    {
      "cell_type": "code",
      "execution_count": null,
      "id": "3c0b8c72-74cd-425e-b82f-4d13fddaef34",
      "metadata": {
        "command_metadata": {
          "end_time": 1781880143574.9314,
          "start_time": 1781880140692.5696
        },
        "execution": {
          "iopub.status.busy": "2026-06-19T14:42:22.450Z"
        },
        "result_type": "result",
        "trusted": true,
        "type": "python"
      },
      "outputs": [],
      "source": [
        "# Incremental one-shot publish: adjust the lookback window as needed.\n",
        "from datetime import datetime, timedelta, timezone\n",
        "from uuid import uuid4\n",
        "\n",
        "from sf_to_oracle.ingest.salesforce_rest import fetch_salesforce_incremental_events\n",
        "\n",
        "since = datetime.now(timezone.utc) - timedelta(minutes=10)\n",
        "events = fetch_salesforce_incremental_events(salesforce_client, since=since, limit_per_object=200)\n",
        "futures = []\n",
        "for event in events:\n",
        "    event.event_id = f\"workbench-incremental-{event.entity}-{event.record_id}-{uuid4()}\"\n",
        "    futures.append(\n",
        "        producer.send(\n",
        "            settings.oci_stream_topic,\n",
        "            key=event.record_id,\n",
        "            value=event.model_dump(mode=\"json\"),\n",
        "        )\n",
        "    )\n",
        "for future in futures:\n",
        "    future.get(timeout=30)\n",
        "producer.flush()\n",
        "print(f\"Published {len(events)} Salesforce changes after {since.isoformat()}\")\n"
      ]
    },
    {
      "cell_type": "code",
      "execution_count": null,
      "id": "5c1fce94-102f-4bbb-ac05-e35ad778fd76",
      "metadata": {
        "command_metadata": {
          "end_time": 1781880151500.2732,
          "start_time": 1781880149819.5645
        },
        "execution": {
          "iopub.status.busy": "2026-06-19T14:42:30.293Z"
        },
        "result_type": "result",
        "trusted": true,
        "type": "python"
      },
      "outputs": [],
      "source": [
        "# Cleanup when finished.\n",
        "producer.close()\n",
        "print(\"Producer closed\")\n"
      ]
    }
  ],
  "metadata": {
    "Last_Active_Cell_Index": 7,
    "kernelspec": {
      "display_name": "Python 3",
      "language": "python",
      "name": "python3"
    },
    "language_info": {
      "name": "python",
      "pygments_lexer": "ipython3"
    }
  },
  "nbformat": 4,
  "nbformat_minor": 5
}
