{
  "files": [
    {
      "path": "src/pipeline.py",
      "content": "\"\"\"ETL Pipeline with extract, transform, load stages.\"\"\"\nimport asyncio\nimport logging\nfrom dataclasses import dataclass\nfrom typing import Any\nfrom pathlib import Path\n\nlogger = logging.getLogger(__name__)\n\n@dataclass\nclass PipelineConfig:\n    source_path: str\n    target_path: str\n    batch_size: int = 1000\n    dry_run: bool = False\n\nasync def extract(config: PipelineConfig) -> list[dict[str, Any]]:\n    \"\"\"Extract data from source.\"\"\"\n    import json\n    path = Path(config.source_path)\n    if not path.exists():\n        raise FileNotFoundError(f\"Source not found: {path}\")\n    with open(path) as f:\n        data = json.load(f)\n    logger.info(\"Extracted %d records from %s\", len(data), path)\n    return data if isinstance(data, list) else [data]\n\ndef transform(records: list[dict[str, Any]]) -> list[dict[str, Any]]:\n    \"\"\"Transform records: clean, validate, enrich.\"\"\"\n    result = []\n    for record in records:\n        cleaned = {k.strip().lower(): v for k, v in record.items() if v is not None}\n        if \"id\" not in cleaned:\n            continue\n        cleaned[\"_processed\"] = True\n        result.append(cleaned)\n    logger.info(\"Transformed %d -> %d records\", len(records), len(result))\n    return result\n\nasync def load(records: list[dict[str, Any]], config: PipelineConfig) -> int:\n    \"\"\"Load transformed records to target.\"\"\"\n    import json\n    if config.dry_run:\n        logger.info(\"DRY RUN: would load %d records to %s\", len(records), config.target_path)\n        return len(records)\n    path = Path(config.target_path)\n    path.parent.mkdir(parents=True, exist_ok=True)\n    with open(path, \"w\") as f:\n        json.dump(records, f, indent=2)\n    logger.info(\"Loaded %d records to %s\", len(records), path)\n    return len(records)\n\nasync def run_pipeline(config: PipelineConfig) -> dict[str, Any]:\n    \"\"\"Run the full ETL pipeline.\"\"\"\n    raw = await extract(config)\n    transformed = transform(raw)\n    count = await load(transformed, config)\n    return {\"extracted\": len(raw), \"transformed\": len(transformed), \"loaded\": count}\n",
      "type": "module"
    },
    {
      "path": "src/main.py",
      "content": "\"\"\"ETL Pipeline CLI entrypoint.\"\"\"\nimport asyncio\nimport argparse\nimport logging\nfrom .pipeline import PipelineConfig, run_pipeline\n\ndef main() -> None:\n    parser = argparse.ArgumentParser(description=\"ETL Pipeline\")\n    parser.add_argument(\"source\", help=\"Source file path\")\n    parser.add_argument(\"target\", help=\"Target file path\")\n    parser.add_argument(\"--batch-size\", type=int, default=1000)\n    parser.add_argument(\"--dry-run\", action=\"store_true\")\n    parser.add_argument(\"--verbose\", \"-v\", action=\"store_true\")\n    args = parser.parse_args()\n\n    logging.basicConfig(level=logging.DEBUG if args.verbose else logging.INFO)\n    config = PipelineConfig(\n        source_path=args.source,\n        target_path=args.target,\n        batch_size=args.batch_size,\n        dry_run=args.dry_run,\n    )\n    result = asyncio.run(run_pipeline(config))\n    print(f\"Pipeline complete: {result}\")\n\nif __name__ == \"__main__\":\n    main()\n",
      "type": "entrypoint"
    },
    {
      "path": "pyproject.toml",
      "content": "[project]\nname = \"etl-pipeline\"\nversion = \"0.1.0\"\nrequires-python = \">=3.11\"\ndependencies = []\n\n[project.scripts]\netl = \"src.main:main\"\n",
      "type": "config"
    }
  ]
}
