Skip to content

Core API

This is the primary public API for podcast_scraper. Use these functions for programmatic access.

Quick Start

import podcast_scraper

# Create configuration
cfg = podcast_scraper.Config(
    rss="https://example.com/feed.xml",
    output_dir="./transcripts",
    max_episodes=10
)

# Run the pipeline
count, summary = podcast_scraper.run_pipeline(cfg)
print(f"Downloaded {count} transcripts: {summary}")

API Reference

run_pipeline

run_pipeline(cfg: Config) -> Tuple[int, str]

Execute the main podcast scraping pipeline.

This is the primary entry point for programmatic use of podcast_scraper. It orchestrates the complete workflow from RSS feed fetching to transcript generation and optional metadata/summarization.

The pipeline executes the following stages:

  1. Setup output directory (with optional run ID subdirectory)
  2. Fetch and parse RSS feed
  3. Detect speakers (if auto-detection enabled)
  4. Process episodes concurrently:
  5. Download published transcripts
  6. Or queue media for Whisper transcription
  7. Transcribe queued media files sequentially (if Whisper enabled)
  8. Generate metadata documents (if enabled)
  9. Generate episode summaries (if enabled)
  10. Clean up temporary files

Parameters:

Name Type Description Default
cfg Config

Configuration object with all pipeline settings. See Config for available options. Download resilience (HTTP/RSS urllib3 retries, optional episode-level retries, and optional Issue #522 fair-HTTP fields) is controlled via http_retry_*, rss_retry_*, episode_retry_*, host_*, circuit_breaker_*, rss_conditional_get, and rss_cache_dir (defaults and CLI flags are documented in CONFIGURATION.md / CLI.md).

required

Returns:

Type Description
Tuple[int, str]

Tuple[int, str]: A tuple containing:

  • count (int): Number of episodes processed (transcripts saved or planned)
  • summary (str): Human-readable summary message describing the run

Raises:

Type Description
RuntimeError

If output directory cleanup fails when clean_output=True

ValueError

If RSS URL is invalid or feed cannot be parsed

FileNotFoundError

If configuration file references missing files

OSError

If file system operations fail

Example

from podcast_scraper import Config, run_pipeline

cfg = Config( ... rss="https://example.com/feed.xml", ... output_dir="./transcripts", ... max_episodes=10 ... ) count, summary = run_pipeline(cfg) print(f"Downloaded {count} transcripts: {summary}") Downloaded 10 transcripts: Processed 10/50 episodes

Example with Whisper transcription

cfg = Config( ... rss="https://example.com/feed.xml", ... transcribe_missing=True, ... whisper_model="base", ... screenplay=True, ... num_speakers=2 ... ) count, summary = run_pipeline(cfg)

Note

For non-interactive use (daemons, services), consider using the service.run() function instead, which provides structured error handling and return values.

See Also
  • Config: Configuration model with all available options
  • service.run(): Service API with structured error handling
  • load_config_file(): Load configuration from JSON/YAML file
Source code in src/podcast_scraper/workflow/orchestration.py
2936
2937
2938
2939
2940
2941
2942
2943
2944
2945
2946
2947
2948
2949
2950
2951
2952
2953
2954
2955
2956
2957
2958
2959
2960
2961
2962
2963
2964
2965
2966
2967
2968
2969
2970
2971
2972
2973
2974
2975
2976
2977
2978
2979
2980
2981
2982
2983
2984
2985
2986
2987
2988
2989
2990
2991
2992
2993
2994
2995
2996
2997
2998
2999
3000
3001
3002
3003
3004
3005
3006
3007
3008
3009
3010
3011
3012
3013
3014
3015
3016
3017
3018
3019
3020
3021
3022
3023
3024
3025
3026
3027
3028
3029
3030
3031
3032
3033
3034
3035
3036
3037
3038
3039
3040
3041
3042
3043
3044
3045
3046
3047
3048
3049
3050
3051
3052
3053
3054
3055
3056
3057
3058
3059
3060
3061
3062
3063
3064
3065
3066
3067
3068
3069
3070
3071
3072
3073
3074
3075
3076
3077
3078
3079
3080
3081
3082
3083
3084
3085
3086
3087
3088
3089
3090
3091
3092
3093
3094
3095
3096
3097
3098
3099
3100
3101
3102
3103
3104
3105
3106
3107
3108
3109
3110
3111
3112
3113
3114
3115
3116
3117
3118
3119
3120
3121
3122
3123
3124
3125
3126
3127
3128
3129
3130
3131
3132
3133
3134
3135
3136
3137
3138
3139
3140
3141
3142
3143
3144
3145
3146
3147
3148
3149
3150
3151
3152
3153
3154
3155
3156
3157
3158
3159
3160
3161
3162
3163
3164
3165
3166
3167
3168
3169
3170
3171
3172
3173
3174
3175
3176
3177
3178
3179
3180
3181
3182
3183
3184
3185
3186
3187
3188
3189
3190
3191
3192
3193
3194
3195
3196
3197
3198
3199
3200
3201
3202
3203
3204
3205
3206
3207
3208
3209
3210
3211
3212
3213
3214
3215
3216
3217
3218
3219
3220
3221
3222
3223
3224
3225
3226
3227
3228
3229
3230
3231
3232
def run_pipeline(cfg: config.Config) -> Tuple[int, str]:
    """Execute the main podcast scraping pipeline.

    This is the primary entry point for programmatic use of podcast_scraper. It orchestrates
    the complete workflow from RSS feed fetching to transcript generation and optional
    metadata/summarization.

    The pipeline executes the following stages:

    1. Setup output directory (with optional run ID subdirectory)
    2. Fetch and parse RSS feed
    3. Detect speakers (if auto-detection enabled)
    4. Process episodes concurrently:
       - Download published transcripts
       - Or queue media for Whisper transcription
    5. Transcribe queued media files sequentially (if Whisper enabled)
    6. Generate metadata documents (if enabled)
    7. Generate episode summaries (if enabled)
    8. Clean up temporary files

    Args:
        cfg: Configuration object with all pipeline settings. See `Config` for available options.
            Download resilience (HTTP/RSS urllib3 retries, optional episode-level retries,
            and optional Issue #522 fair-HTTP fields) is controlled via ``http_retry_*``,
            ``rss_retry_*``, ``episode_retry_*``, ``host_*``, ``circuit_breaker_*``,
            ``rss_conditional_get``, and ``rss_cache_dir`` (defaults and CLI flags are
            documented in CONFIGURATION.md / CLI.md).

    Returns:
        Tuple[int, str]: A tuple containing:

            - count (int): Number of episodes processed (transcripts saved or planned)
            - summary (str): Human-readable summary message describing the run

    Raises:
        RuntimeError: If output directory cleanup fails when `clean_output=True`
        ValueError: If RSS URL is invalid or feed cannot be parsed
        FileNotFoundError: If configuration file references missing files
        OSError: If file system operations fail

    Example:
        >>> from podcast_scraper import Config, run_pipeline
        >>>
        >>> cfg = Config(
        ...     rss="https://example.com/feed.xml",
        ...     output_dir="./transcripts",
        ...     max_episodes=10
        ... )
        >>> count, summary = run_pipeline(cfg)
        >>> print(f"Downloaded {count} transcripts: {summary}")
        Downloaded 10 transcripts: Processed 10/50 episodes

    Example with Whisper transcription:
        >>> cfg = Config(
        ...     rss="https://example.com/feed.xml",
        ...     transcribe_missing=True,
        ...     whisper_model="base",
        ...     screenplay=True,
        ...     num_speakers=2
        ... )
        >>> count, summary = run_pipeline(cfg)

    Example with metadata and summaries:
        >>> cfg = Config(
        ...     rss="https://example.com/feed.xml",
        ...     generate_metadata=True,
        ...     generate_summaries=True
        ... )
        >>> count, summary = run_pipeline(cfg)

    Note:
        For non-interactive use (daemons, services), consider using the `service.run()`
        function instead, which provides structured error handling and return values.

    See Also:
        - `Config`: Configuration model with all available options
        - `service.run()`: Service API with structured error handling
        - `load_config_file()`: Load configuration from JSON/YAML file
    """
    # #1454: validate the required RSS input at the entry, before any processing, rather than let a
    # URL-less cfg reach `fetch_and_parse_feed` deep in the run and surface via the CLI catch-all as
    # a generic "Unexpected failure". The CLI arg parser already checks this for its own path; this
    # guards every OTHER caller (server jobs, programmatic `run_pipeline(cfg)`) with the same clear,
    # actionable message. rss_urls (the multi-feed field) counts as provided.
    if not getattr(cfg, "rss_url", None) and not getattr(cfg, "rss_urls", None):
        raise ValueError(
            "RSS URL is required: set `rss` (positional / --rss / --rss-file), a `--feeds-spec`, "
            "or `rss_urls` in config before running the pipeline."
        )
    # Install the run-level LLM call fuse for the WHOLE production run, in the main thread, before
    # any stage or worker starts. This is the hard ceiling on total spend: retry_with_metrics ticks
    # it on every attempt, and the fuse is process-global so it is enforced inside the
    # summarization/processing ThreadPoolExecutor workers too (a ContextVar would not reach them).
    # The finer per-episode fuse — which catches a single-episode storm below the run ceiling (the
    # ~3,500-call incident) — is installed per episode in generate_episode_metadata, where that
    # storm (bundled GI evidence → per-pair fallback) actually runs.
    from ..utils import llm_call_fuse

    llm_call_fuse.install_run(getattr(cfg, "llm_max_calls_per_run", 0))

    # GitHub #562: reset gates before setup (setup is outside try/finally below).
    try:
        config.reset_screenplay_issue_562_gates()
    except Exception:  # pragma: no cover - defensive import/cleanup
        logger.debug("reset_screenplay_issue_562_gates (startup) failed", exc_info=True)

    # #1053: resolve this run's correlation id ONCE, up front, so every o11y signal
    # (Loki cost event + logs, Sentry scope, Langfuse trace) for the run stamps the same
    # join key. Process-global because the pipeline is a per-run subprocess.
    from podcast_scraper.utils import correlation

    correlation.set_run_id(correlation.resolve_run_id(cfg.run_id))
    # Mirror the join key onto the Sentry scope so errors correlate too (no-op without Sentry).
    try:
        from podcast_scraper.utils.sentry_init import set_run_tag

        set_run_tag(correlation.get_run_id())
    except Exception:  # pragma: no cover - never block a run on o11y tagging
        logger.debug("sentry run-tag skipped", exc_info=True)

    # Same idea one level up: stamp WHICH profile this run resolved to, and what that
    # profile actually routed each stage to, onto every event the run emits. Resolved here
    # (not from the profile name at read time) because three layers can move routing under a
    # profile — corpus YAML, the feed's own pin, and a per-request override (#1872) — so the
    # name alone no longer implies the routing, and a cost anomaly's first question ("which
    # profile, and was ASR on the DGX or Deepgram?") must be answerable from the event.
    stamp_run_identity(cfg)

    # Step 1: Setup pipeline environment
    effective_output_dir, run_suffix, full_config_string, pipeline_metrics = (
        _setup_pipeline_environment(cfg)
    )

    # GitHub #557: structured incident log (episode/feed scope); default beside run artifacts.
    if not (cfg.incident_log_path or "").strip():
        cfg = cfg.model_copy(
            update={
                "incident_log_path": str(Path(effective_output_dir) / "corpus_incidents.jsonl"),
            }
        )

    monitor_proc: Optional[Any] = None
    py_spy_stop: Optional[Callable[[], None]] = None
    if cfg.monitor:
        from ..monitor.py_spy_listener import start_py_spy_stdin_listener
        from ..monitor.runner import start_monitor_subprocess

        monitor_proc = start_monitor_subprocess(
            pipeline_pid=os.getpid(),
            output_dir=effective_output_dir,
        )
        py_spy_stop = start_py_spy_stdin_listener(
            output_dir=effective_output_dir,
            enabled=True,
        )

    try:
        from ..gi.deps import validate_gil_grounding_dependencies

        validate_gil_grounding_dependencies(cfg)

        # Initialize JSONL emitter if enabled
        jsonl_emitter = _setup_jsonl_emitter(cfg, effective_output_dir, pipeline_metrics)

        # Step 1.5: Preload ML models if configured
        wf_stages.setup.preload_ml_models_if_needed(cfg)

        # Step 1.6: Create all providers once (singleton pattern per run)
        # Providers are created here and passed to stages to avoid redundant initialization
        transcription_provider, speaker_detector, summary_provider = _create_all_providers(cfg)

        # Step 1.7-1.8: Setup logging and device tracking
        _setup_logging_and_devices(
            cfg, transcription_provider, speaker_detector, summary_provider, pipeline_metrics
        )

        # Step 1.5: Create run manifest
        run_manifest = _create_run_manifest(cfg, effective_output_dir)

        # Step 2-4: Fetch and prepare episodes
        maybe_update_pipeline_status(cfg, effective_output_dir, stage="rss_feed_fetch")
        feed, rss_bytes, feed_metadata, episodes = _fetch_and_prepare_episodes(
            cfg, pipeline_metrics
        )

        # Step 5-6.5: Setup pipeline resources
        normalizing_start, host_detection_result, transcription_resources, processing_resources = (
            _setup_pipeline_resources(
                cfg,
                feed,
                episodes,
                effective_output_dir,
                transcription_provider,
                speaker_detector,
                pipeline_metrics,
            )
        )

        # Wrap processing + finalize: JSONL must stay open until _finalize_pipeline
        # calls emit_run_finished (see _finalize_emit_and_save). Closing the emitter in
        # the inner finally was too early and broke run_finished emission.
        interim_checkpoint_manager = _InterimCheckpointManager(
            cfg=cfg,
            output_dir=effective_output_dir,
            pipeline_metrics=pipeline_metrics,
        )
        try:
            try:
                saved = _process_episodes_with_threading(
                    cfg=cfg,
                    episodes=episodes,
                    feed=feed,
                    effective_output_dir=effective_output_dir,
                    run_suffix=run_suffix,
                    feed_metadata=feed_metadata,
                    host_detection_result=host_detection_result,
                    transcription_resources=transcription_resources,
                    processing_resources=processing_resources,
                    pipeline_metrics=pipeline_metrics,
                    summary_provider=summary_provider,
                    transcription_provider=transcription_provider,
                    normalizing_start=normalizing_start,
                    interim_checkpoint_manager=interim_checkpoint_manager,
                )

            finally:
                interim_checkpoint_manager.stop()
                # Step 9.5: Unload models to free memory
                _cleanup_providers(transcription_resources, summary_provider)

            # Step 10-15: Finalize pipeline (metrics, JSONL run_finished, index, …)
            result = _finalize_pipeline(
                cfg=cfg,
                saved=saved,
                transcription_resources=transcription_resources,
                effective_output_dir=effective_output_dir,
                run_suffix=run_suffix,
                pipeline_metrics=pipeline_metrics,
                episodes=episodes,
                jsonl_emitter=jsonl_emitter,
                run_manifest=run_manifest,
                summary_provider=summary_provider,
                transcription_provider=transcription_provider,
            )
        except BaseException:
            if jsonl_emitter is not None:
                try:
                    jsonl_emitter.__exit__(None, None, None)
                except Exception:
                    pass
            raise

        # Step 16: #1058 chunk 3 — corpus-level Topic clustering.
        # Runs AFTER per-episode finalize so every kg.json is on disk
        # before we collect Topic labels across the corpus. Gated on
        # cfg.kg_topic_corpus_clustering (default off; airgapped*
        # overlays flip it on). Non-fatal — a failure here doesn't
        # bring down a successful run.
        if getattr(cfg, "kg_topic_corpus_clustering", False) and not cfg.dry_run:
            try:
                from pathlib import Path as _Path

                from podcast_scraper.kg.topic_clustering import (
                    cluster_and_apply_corpus_topics,
                )

                summary = cluster_and_apply_corpus_topics(_Path(effective_output_dir))
                logger.info(
                    "corpus topic clustering: clusters=%d concept_topics_added=%d "
                    "related_to_edges_added=%d artifacts_mutated=%d",
                    summary.clusters_found,
                    summary.concept_topics_added,
                    summary.related_to_edges_added,
                    summary.artifacts_mutated,
                )
            except Exception as cluster_exc:
                logger.warning(
                    "corpus topic clustering failed (non-fatal): %s",
                    cluster_exc,
                    exc_info=True,
                )

        maybe_update_pipeline_status(cfg, effective_output_dir, stage="done")
        return result
    finally:
        # GitHub #562: allow coercion INFO + screenplay warnings on the next Config / run.
        try:
            config.reset_screenplay_issue_562_gates()
        except Exception:  # pragma: no cover - defensive import/cleanup
            logger.debug("reset_screenplay_issue_562_gates failed", exc_info=True)
        if py_spy_stop is not None:
            py_spy_stop()
        if monitor_proc is not None:
            monitor_proc.join(timeout=30)
            if monitor_proc.is_alive():
                monitor_proc.terminate()
                monitor_proc.join(timeout=5)

load_config_file

load_config_file(path: str) -> Dict[str, Any]

Load configuration from a JSON or YAML file.

This function reads a configuration file and returns a dictionary of configuration values. The file format is auto-detected from the file extension (.json, .yaml, or .yml).

The returned dictionary can be unpacked into the Config constructor to create a configuration object.

Parameters:

Name Type Description Default
path str

Path to configuration file (JSON or YAML). Supports tilde expansion for home directory (e.g., "~/config.yaml").

required

Returns:

Type Description
Dict[str, Any]

Dict[str, Any]: Dictionary containing configuration values from the file. Keys correspond to Config field names (using aliases where applicable).

Raises:

Type Description
ValueError

If any of the following occur:

  • Config path is empty
  • Config file does not exist
  • File format is invalid (not JSON or YAML)
  • JSON parsing fails
  • YAML parsing fails
OSError

If file cannot be read due to permissions or I/O errors

Example

from podcast_scraper import Config, load_config_file, run_pipeline

Load from YAML file

config_dict = load_config_file("config.yaml") cfg = Config(**config_dict) count, summary = run_pipeline(cfg)

Example with JSON

config_dict = load_config_file("config.json") cfg = Config(**config_dict)

Example with direct usage

from podcast_scraper import load_config_file, service

Service API provides load_config_file convenience

result = service.run_from_config_file("config.yaml")

Supported Formats

JSON (.json):

{
  "rss": "https://example.com/feed.xml",
  "output_dir": "./transcripts",
  "max_episodes": 50
}

YAML (.yaml, .yml):

rss: https://example.com/feed.xml
output_dir: ./transcripts
max_episodes: 50
Note
  • Field aliases are supported (e.g., both "rss" and "rss_url" work)
  • See Config documentation for all available configuration options
  • Configuration files should not contain sensitive data (API keys, passwords)
See Also
  • Config: Configuration model and field documentation
  • service.run_from_config_file(): Direct service API from config file
  • Configuration examples: config/examples/config.example.json, config/examples/config.example.yaml
Source code in src/podcast_scraper/config.py
7099
7100
7101
7102
7103
7104
7105
7106
7107
7108
7109
7110
7111
7112
7113
7114
7115
7116
7117
7118
7119
7120
7121
7122
7123
7124
7125
7126
7127
7128
7129
7130
7131
7132
7133
7134
7135
7136
7137
7138
7139
7140
7141
7142
7143
7144
7145
7146
7147
7148
7149
7150
7151
7152
7153
7154
7155
7156
7157
7158
7159
7160
7161
7162
7163
7164
7165
7166
7167
7168
7169
7170
7171
7172
7173
7174
7175
7176
7177
7178
7179
7180
7181
7182
7183
7184
7185
7186
7187
7188
7189
7190
7191
7192
7193
7194
7195
7196
7197
7198
7199
7200
7201
7202
7203
7204
7205
7206
7207
7208
def load_config_file(
    path: str,
) -> Dict[str, Any]:  # noqa: C901 - file parsing handles multiple formats
    """Load configuration from a JSON or YAML file.

    This function reads a configuration file and returns a dictionary of configuration values.
    The file format is auto-detected from the file extension (`.json`, `.yaml`, or `.yml`).

    The returned dictionary can be unpacked into the `Config` constructor to create a
    configuration object.

    Args:
        path: Path to configuration file (JSON or YAML). Supports tilde expansion for
              home directory (e.g., "~/config.yaml").

    Returns:
        Dict[str, Any]: Dictionary containing configuration values from the file.
            Keys correspond to `Config` field names (using aliases where applicable).

    Raises:
        ValueError: If any of the following occur:

            - Config path is empty
            - Config file does not exist
            - File format is invalid (not JSON or YAML)
            - JSON parsing fails
            - YAML parsing fails

        OSError: If file cannot be read due to permissions or I/O errors

    Example:
        >>> from podcast_scraper import Config, load_config_file, run_pipeline
        >>>
        >>> # Load from YAML file
        >>> config_dict = load_config_file("config.yaml")
        >>> cfg = Config(**config_dict)
        >>> count, summary = run_pipeline(cfg)

    Example with JSON:
        >>> config_dict = load_config_file("config.json")
        >>> cfg = Config(**config_dict)

    Example with direct usage:
        >>> from podcast_scraper import load_config_file, service
        >>>
        >>> # Service API provides load_config_file convenience
        >>> result = service.run_from_config_file("config.yaml")

    Supported Formats:
        **JSON** (`.json`):

            {
              "rss": "https://example.com/feed.xml",
              "output_dir": "./transcripts",
              "max_episodes": 50
            }

        **YAML** (`.yaml`, `.yml`):

            rss: https://example.com/feed.xml
            output_dir: ./transcripts
            max_episodes: 50

    Note:
        - Field aliases are supported (e.g., both "rss" and "rss_url" work)
        - See `Config` documentation for all available configuration options
        - Configuration files should not contain sensitive data (API keys, passwords)

    See Also:
        - `Config`: Configuration model and field documentation
        - `service.run_from_config_file()`: Direct service API from config file
        - Configuration examples: `config/examples/config.example.json`,
          `config/examples/config.example.yaml`
    """
    if not path:
        raise ValueError("Config path cannot be empty")

    cfg_path = Path(path).expanduser()
    try:
        resolved = cfg_path.resolve()
    except (OSError, RuntimeError) as exc:
        raise ValueError(f"Invalid config path: {path} ({exc})") from exc

    if not resolved.exists():
        raise ValueError(f"Config file not found: {resolved}")

    suffix = resolved.suffix.lower()
    try:
        text = resolved.read_text(encoding="utf-8")
    except OSError as exc:
        raise ValueError(f"Failed to read config file {resolved}: {exc}") from exc

    if suffix == ".json":
        try:
            data = json.loads(text)
        except json.JSONDecodeError as exc:
            raise ValueError(f"Invalid JSON config file {resolved}: {exc}") from exc
    elif suffix in (".yaml", ".yml"):
        try:
            data = yaml.safe_load(text)
        except yaml.YAMLError as exc:  # type: ignore[attr-defined]
            raise ValueError(f"Invalid YAML config file {resolved}: {exc}") from exc
    else:
        raise ValueError(f"Unsupported config file type: {resolved.suffix}")

    if not isinstance(data, dict):
        raise ValueError("Config file must contain a mapping/object at the top level")

    expanded = _expand_env_vars(data)
    return cast(Dict[str, Any], expanded)

Package Information

Versioning

__version__ module-attribute

__version__ = '2.7.0.dev0'

__api_version__ module-attribute

__api_version__ = '2.7.0'