AI_Agro_Support/pipeline/__main__.py
Arsham Mirehvandi e588d97e0f Add observability support with OpenTelemetry tracing integration
- Updated .env.example and config.yaml to include observability settings.
- Added a new Phoenix service in docker-compose.yml for self-hosted tracing.
- Enhanced README.md with instructions on enabling and using observability features.
- Implemented tracing in the pipeline, including job spans and LLM stage spans.
- Introduced ObservabilitySettings class in config.py for better configuration management.
- Updated job.py and resources.py to support tracing without affecting fault isolation.
- Minor adjustments to other files for compatibility with the new observability features.
2026-08-25 22:06:54 +02:00

230 lines
7.9 KiB
Python

"""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())