""" Reprocess missed procurement-ingest emails by re-invoking an email-processor Lambda with a synthetic S3 event. TARGETED replay is the default and preferred mode: point it at exactly the object(s) you need to replay. # dry-run a single object (default is dry-run — lists, does not invoke) python scripts/reprocess.py --pipeline po --key inbound/2026/msg.eml # actually re-invoke that one object python scripts/reprocess.py --pipeline po --key inbound/2026/msg.eml --execute # replay a prefix, or everything modified since a timestamp python scripts/reprocess.py --pipeline wo --prefix inbound/2026/07/ python scripts/reprocess.py --pipeline po --since 2026-07-16T00:00:00Z --execute FULL-PREFIX sweep is DEMOTED behind an explicit --all flag and carries five load-bearing caveats (see the --all help text and the in-code comment above the sweep branch): python scripts/reprocess.py --pipeline po --all --execute # replays ALL of inbound/ 1. Event-type (async) invocation is CONCURRENT, so sorting does NOT serialize the replay. 2. Switch InvocationType to "RequestResponse" if ordering matters. 3. Metrics get double-counted on replay. 4. Bedrock is re-billed for every ai_fallback email replayed. 5. Out-of-order replay can REGRESS already-merged fields. Dry-run is the default in ALL modes (targeted and --all alike); nothing is invoked unless --execute is passed. """ import argparse import json from datetime import datetime, timezone import boto3 # Per-pipeline resolution: the function to invoke and the raw-email bucket to # list, keyed by --pipeline. Both are derived at runtime from --pipeline + # the caller's account id; there is no PO-hardcoded default. PIPELINES = { "po": {"function": "po-email-processor", "bucket": "po-ingest-emails-{acct}"}, "wo": { "function": "workorder-email-processor", "bucket": "workorder-ingest-emails-{acct}", }, } DEFAULT_PREFIX = "inbound/" def parse_since(value: str) -> datetime: """Parse an ISO-8601 UTC timestamp (accepts a trailing 'Z') into an aware UTC datetime so it can be compared against S3 LastModified.""" normalized = value.replace("Z", "+00:00") dt = datetime.fromisoformat(normalized) if dt.tzinfo is None: dt = dt.replace(tzinfo=timezone.utc) return dt.astimezone(timezone.utc) def list_keys(s3, bucket: str, prefix: str, since: datetime | None) -> list[str]: """Paginate list_objects_v2 under prefix, returning raw object keys. When since is set, only objects with LastModified >= since are kept. The keys returned are the RAW list_objects_v2 obj["Key"] — no URL-encoding or decoding is applied (see build_s3_event). """ keys: list[str] = [] paginator = s3.get_paginator("list_objects_v2") for page in paginator.paginate(Bucket=bucket, Prefix=prefix): for obj in page.get("Contents", []): if since is not None and obj["LastModified"] < since: continue keys.append(obj["Key"]) return keys def build_s3_event(bucket: str, key: str) -> dict: """Build the synthetic S3 notification event the handler consumes. CONTRACT (pinned by tests/test_reprocess_contract.py): the shape is exactly Records[0].s3.bucket.name / Records[0].s3.object.key, and the key is emitted RAW — reprocess applies NO urllib.parse.quote/unquote. A real S3 notification URL-encodes the object key; reprocess builds from the raw list_objects_v2 key (or the raw --key arg), and the handler is what would decode. Replaying a pre-decoded/pre-encoded key would not match what s3.get_object(Bucket, Key) expects. """ return { "Records": [ { "s3": { "bucket": {"name": bucket}, "object": {"key": key}, } } ] } def main(): parser = argparse.ArgumentParser( description=( "Reprocess missed procurement-ingest emails by re-invoking an " "email-processor Lambda with a synthetic S3 event. TARGETED replay " "(--key/--prefix/--since) is the default and preferred mode; " "full-prefix sweep requires the explicit --all flag." ) ) parser.add_argument( "--pipeline", required=True, choices=["po", "wo"], help=( "Which pipeline to replay: 'po' (po-email-processor / " "po-ingest-emails-) or 'wo' (workorder-email-processor / " "workorder-ingest-emails-)." ), ) # Targeting selectors — at least one of {--key, --prefix, --since, --all} # is required (enforced below); --key is exclusive. parser.add_argument( "--key", help="Replay exactly one raw object key (e.g. inbound/2026/msg.eml).", ) parser.add_argument( "--prefix", help=( "Replay every object under this S3 prefix " "(default targeting prefix is inbound/)." ), ) parser.add_argument( "--since", help=( "Only objects with LastModified >= this ISO-8601 UTC timestamp " "(e.g. 2026-07-16T00:00:00Z). Combines with --prefix." ), ) parser.add_argument( "--all", action="store_true", help=( "DEMOTED full-prefix sweep: replay EVERY object under inbound/ " "(not a targeted subset). " "CAVEATS: (1) Event-type (async) invocation is CONCURRENT so " "sorting does NOT serialize the replay; " "(2) switch to RequestResponse if order matters; " "(3) metrics get double-counted on replay; " "(4) Bedrock is re-billed; " "(5) out-of-order replay can REGRESS already-merged fields. " "Prefer --key/--prefix/--since. Still dry-run unless --execute." ), ) parser.add_argument( "--execute", action="store_true", help=( "Actually invoke the Lambda. Default is dry-run (list only). " "Applies to EVERY mode, including --all." ), ) args = parser.parse_args() # Exactly-one target gate: full sweep is ONLY reachable via --all, so an # accidental whole-prefix sweep from running with no selector is impossible. if not any([args.key, args.prefix, args.since, args.all]): parser.error("choose a target: --key, --prefix, --since, or --all (full sweep)") # --key is mutually exclusive with the listing selectors and with --all. if args.key and (args.prefix or args.since or args.all): parser.error("--key cannot be combined with --prefix, --since, or --all") # --all is a distinct MODE (the demoted full-prefix sweep), not a modifier: # reject it alongside the targeting selectors. Without this, the args.all # clobber branch below would silently discard a narrower --prefix/--since # (e.g. `--all --since 2026-07-01` would sweep the ENTIRE corpus, not the # bounded window) and trigger every --all hazard on unintended objects. if args.all and (args.prefix or args.since): parser.error("--all cannot be combined with --prefix or --since") account_id = boto3.client("sts").get_caller_identity()["Account"] function_name = PIPELINES[args.pipeline]["function"] bucket = PIPELINES[args.pipeline]["bucket"].format(acct=account_id) if args.key: # Single raw key — skip listing entirely. The key is passed through # to build_s3_event untransformed (raw-key contract). keys = [args.key] else: s3 = boto3.client("s3") # --all is the DEMOTED full-prefix sweep. Five caveats, all load-bearing: # 1. Event-type (async) invocation is CONCURRENT, so sorting the key # list does NOT serialize the replay. # 2. If ordering matters, switch InvocationType to "RequestResponse" # (serial, waits per call). # 3. Metrics (ParseOutcome EMF, alarms) get double-counted on replay. # 4. Bedrock is re-billed for every ai_fallback email replayed. # 5. Out-of-order replay can REGRESS already-merged fields (a stale # email overwriting a newer merge). # Targeted replay (--key / --prefix / --since) avoids all five at # whole-corpus scale; prefer it. if args.all: prefix = DEFAULT_PREFIX since = None else: # --prefix / --since target. --since without --prefix defaults the # prefix to inbound/ so it filters the raw-email tree, not the bucket. prefix = args.prefix if args.prefix else DEFAULT_PREFIX since = parse_since(args.since) if args.since else None keys = list_keys(s3, bucket, prefix, since) if not keys: print(f"No matching objects in s3://{bucket}/ — nothing to reprocess.") return if args.all: print( f"[--all full sweep] {len(keys)} email(s) under " f"s3://{bucket}/{DEFAULT_PREFIX}\n" ) else: print(f"{len(keys)} email(s) selected in s3://{bucket}/\n") lambda_client = boto3.client("lambda") if args.execute else None for key in keys: if args.execute: print(f" Invoking {function_name} for {key} ... ", end="", flush=True) resp = lambda_client.invoke( FunctionName=function_name, InvocationType="Event", # async — don't wait for each one Payload=json.dumps(build_s3_event(bucket, key)), ) print(f"status {resp['StatusCode']}") else: print(f" [dry-run] {key}") if not args.execute: print("\nDry run complete. Re-run with --execute to invoke the Lambda.") if __name__ == "__main__": main()