Skip to content

Python API

Everything the CLI does is available as plain Python functions. Use the API from notebooks, Airflow or Dagster tasks, Snowflake Python procedures, or your own scripts.

Sessions and access tokens

from lht.user.auth import create_session
from lht.user.salesforce_auth import get_salesforce_access_info

session = create_session(connection_name="my_snowflake")        # snowflake.snowpark.Session
access_info = get_salesforce_access_info("my_salesforce")       # {"access_token": ..., "instance_url": ...}

Pass no name to use the primary connection of each type.

Headless credentials

In CI or an orchestrator, read secrets from your secret manager and pass them in directly. Nothing touches ~/.solomo.

import os
from lht.user.auth import create_session
from lht.user.salesforce_auth import get_salesforce_access_info_from_credentials

session = create_session({
    "account": os.environ["SNOWFLAKE_ACCOUNT"],
    "user": os.environ["SNOWFLAKE_USER"],
    "role": "LHT_ROLE",
    "warehouse": "COMPUTE_WH",
    "database": "SALESFORCE",
    "schema": "RAW",
    "private_key_file": "/run/secrets/snowflake_key.p8",
    "private_key_passphrase": os.environ.get("SNOWFLAKE_KEY_PASSPHRASE"),
})

# Client Credentials flow
access_info = get_salesforce_access_info_from_credentials({
    "auth_flow": "client_credentials",
    "client_id": os.environ["SF_CLIENT_ID"],
    "client_key": os.environ["SF_CLIENT_SECRET"],
    "my_domain": "acme",                      # or "acme--dev.sandbox"
})

# ...or JWT Bearer flow
access_info = get_salesforce_access_info_from_credentials({
    "auth_flow": "jwt_bearer",
    "client_id": os.environ["SF_CLIENT_ID"],
    "username": "integration@acme.com",
    "private_key_pem": os.environ["SF_PRIVATE_KEY"],
    "sandbox": False,
})

Sync Salesforce → Snowflake

from lht.salesforce.intelligent_sync import sync_sobject_intelligent

result = sync_sobject_intelligent(
    session=session,
    access_info=access_info,
    sobject="Account",
    schema="RAW",
    table="ACCOUNT",
    match_field="ID",            # merge key (default)
    where_clause=None,           # e.g. "IsPersonAccount = false"
    force_full_sync=False,       # True rebuilds the table
    use_stage=False,             # True + stage_name loads through a Snowflake stage
    stage_name=None,
    delete_job=True,             # delete the Bulk API job when done
)

result is a dict:

{
    "sobject": "Account",
    "target_table": "RAW.ACCOUNT",
    "sync_method": "bulk_api_incremental",   # or bulk_api_full, bulk_api_stage_*
    "estimated_records": 1500,
    "actual_records": 1487,
    "sync_duration_seconds": 45.2,
    "last_modified_date": Timestamp("2026-01-15 10:30:00"),
    "sync_timestamp": Timestamp("2026-01-16 14:20:00"),
    "success": True,
    "error": None,
}

To sync several objects in a loop:

for sobject in ["Account", "Contact", "Opportunity", "Case"]:
    r = sync_sobject_intelligent(session, access_info, sobject, schema="RAW", table=sobject.upper())
    print(f"{sobject}: {r['actual_records']} rows via {r['sync_method']}")

Reverse ETL: Snowflake → Salesforce

from lht.salesforce import retl

# Column names must be Salesforce field API names.
retl.upsert(session, access_info, sobject="Account",
            query="SELECT External_Id__c, Rating, Industry FROM ANALYTICS.ACCOUNT_SCORES",
            field="External_Id__c", batch_size=25000, clear_nulls=False)

retl.update(session, access_info, sobject="Contact",
            query="SELECT Id, Title FROM STAGE.CONTACT_FIXES", clear_nulls=True)

retl.insert(session, access_info, sobject="Task",
            query="SELECT WhoId, Subject, ActivityDate FROM STAGE.NEW_TASKS")

retl.delete(session, access_info, sobject="Lead",
            query="SELECT Id FROM STAGE.LEADS_TO_DELETE", field="Id")

clear_nulls=True sends NULL as #N/A, which clears the field in Salesforce. Without it, a NULL is sent as an empty cell, which Salesforce treats as "leave unchanged".

Merge records

from lht.salesforce.merge import merge

summary = merge(session, access_info, sobject="Account",
                query="SELECT MasterId, LoserId FROM DEDUPE.ACCOUNT_PAIRS",
                dry_run=True)

See Merging records for the validation rules.

Bulk API 2.0 jobs

from lht.salesforce import jobs

jobs.list_bulk_api_jobs(access_info)
jobs.get_bulk_api_job(access_info, "750xx000000abcDAAQ")
jobs.get_ingest_job_results(access_info, "750xx000000abcDAAQ")   # successful / failed / unprocessed CSVs
jobs.delete_bulk_api_job(access_info, "750xx000000abcDAAQ")