"""CLI entry point. python -m pipeline batch [--as-of DATE] [--crop NAME] [--disease NAME] [--field-id N ...] [--deadline HH:MM] [--skip-existing-advice] [--dry-run] [--workers N] python -m pipeline one --field-id N --disease NAME [--as-of DATE] [--dry-run] `batch` is the daily entry point: it discovers every (field, disease) job implied by config.yaml's `crops:` list and the database, and runs them through a bounded thread pool (see pipeline/batch.py). `one` keeps the original single-field behaviour for ad-hoc debugging. """ from __future__ import annotations import argparse import logging import sys from datetime import date, datetime from pipeline.batch import run_batch from pipeline.config import load_settings from pipeline.errors import NoModelConfiguredError, PipelineError from pipeline.job import run_job from pipeline.observability import setup_tracing, shutdown_tracing from pipeline.output import write_job_output from pipeline.prompts import PromptRegistry from pipeline.resources import Resources from pipeline.stages.disease import find_anmod_id from pipeline.stages.field import resolve_field from pipeline.vocab import load_crop_vocab, load_disease_vocab from pipeline.worklist import Job logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", ) logger = logging.getLogger("pipeline") def _parse_as_of(value: str | None) -> date: if value is None: return date.today() try: return datetime.strptime(value, "%Y-%m-%d").date() except ValueError as exc: raise argparse.ArgumentTypeError( f"Invalid --as-of date '{value}'. Expected YYYY-MM-DD." ) from exc def _run_one(args: argparse.Namespace) -> int: """`python -m pipeline one`: the original single-field pipeline, now built on top of `run_job` so it stays behaviourally identical to `batch`.""" settings = load_settings(field_id_override=args.field_id, disease_name_override=args.disease) setup_tracing(settings) if settings.field_id is None or settings.disease_name is None: raise PipelineError( "`one` needs --field-id and --disease " "(or field_id / disease_name set in config.yaml)." ) crop_vocab = load_crop_vocab(settings.crops_vocab) disease_vocab = load_disease_vocab(settings.diseases_vocab) disease_english = disease_vocab.to_english(settings.disease_name, kind="disease") logger.info( "Starting single-field run as_of=%s field_id=%s disease=%s (%s) dry_run=%s", args.as_of.isoformat(), settings.field_id, settings.disease_name, disease_english, args.dry_run, ) prompts = PromptRegistry(settings.prompts_dir) resources = Resources(settings, prompts) try: conn = resources.sql_connection() field = resolve_field(conn, settings.field_id, crop_vocab) anmod_id = find_anmod_id(conn, settings.field_id, settings.disease_name) resources.load_product_index() job = Job( field_id=settings.field_id, crop_english=field.crop_english, crop_italian=field.crop_italian, model_name=settings.disease_name, disease_english=disease_english, organic=field.organic, station=field.cmplay_station, anmod_id=anmod_id, ) result = run_job(resources, job, args.as_of, dry_run=args.dry_run) finally: resources.close() if result.payload is not None: path = write_job_output(settings.output_dir, args.as_of, result.payload) logger.info("Wrote output -> %s", path) if result.status == "failed": logger.error("%s: %s", result.error_type, result.error_message) return 1 logger.info("Done: status=%s message=%s", result.status, result.message) return 0 def _run_batch(args: argparse.Namespace) -> int: """`python -m pipeline batch`: the daily crop-first multi-field run.""" settings = load_settings() setup_tracing(settings) crop_vocab = load_crop_vocab(settings.crops_vocab) logger.info( "Starting batch as_of=%s crop_filter=%s disease_filter=%s field_filter=%s " "skip_existing_advice=%s dry_run=%s workers=%s", args.as_of.isoformat(), args.crop, args.disease, args.field_id, args.skip_existing_advice, args.dry_run, args.workers or settings.concurrency.workers, ) return run_batch( settings, crop_vocab, args.as_of, crop=args.crop, disease=args.disease, field_ids=args.field_id, deadline_override=args.deadline, skip_existing_advice=args.skip_existing_advice, dry_run=args.dry_run, workers_override=args.workers, ) def main(argv: list[str] | None = None) -> int: parser = argparse.ArgumentParser(description="Daily agronomic advice pipeline.") subparsers = parser.add_subparsers(dest="command") batch_parser = subparsers.add_parser( "batch", help="Run every configured crop/disease pair across all matching fields.", ) batch_parser.add_argument( "--as-of", type=_parse_as_of, default=None, help="Reference date (YYYY-MM-DD). Defaults to today.", ) batch_parser.add_argument( "--crop", default=None, help="Only run this crop (canonical English name).", ) batch_parser.add_argument( "--disease", default=None, help="Only run this disease (canonical English name).", ) batch_parser.add_argument( "--field-id", type=int, action="append", default=None, help="Only run this field ID (repeatable).", ) batch_parser.add_argument( "--deadline", default=None, help="Override schedule.deadline from config.yaml (HH:MM, local time).", ) batch_parser.add_argument( "--skip-existing-advice", action="store_true", help="Skip (field, disease) pairs that already have an advice row for --as-of.", ) batch_parser.add_argument( "--dry-run", action="store_true", help="Build the JSON for every job but skip LLM calls, vector search, and the advice DB write.", ) batch_parser.add_argument( "--workers", type=int, default=None, help="Override concurrency.workers from config.yaml.", ) batch_parser.set_defaults(func=_run_batch) one_parser = subparsers.add_parser( "one", help="Run a single field/disease pair (debugging; the original single-field CLI behaviour).", ) one_parser.add_argument( "--as-of", type=_parse_as_of, default=None, help="Reference date (YYYY-MM-DD). Defaults to today.", ) one_parser.add_argument( "--field-id", type=int, default=None, help="Field ID (defaults to config.yaml's field_id).", ) one_parser.add_argument( "--disease", default=None, help="Disease model name, e.g. PERONOSPORA (defaults to config.yaml's disease_name).", ) one_parser.add_argument( "--dry-run", action="store_true", help="Build the JSON but skip LLM calls, vector search, and the advice DB write.", ) one_parser.set_defaults(func=_run_one) args = parser.parse_args(argv) if not args.command: parser.print_help() return 2 if not isinstance(args.as_of, date): args.as_of = _parse_as_of(None) try: return args.func(args) except NoModelConfiguredError as exc: logger.error("%s", exc) return 2 except PipelineError as exc: logger.error("%s", exc) return 1 except Exception: logger.exception("Unhandled pipeline failure") return 1 finally: # Runs after write_job_output / write_run_report, so a slow or # unreachable Phoenix collector cannot delay the operational # artifacts -- only the process's own exit. shutdown_tracing() if __name__ == "__main__": sys.exit(main())