procurement-ingest/scripts/reprocess.py
Adam Moussa 75fe91c198
Ops/recovery tooling + dependency hygiene (refactor phase 7) (#110)
* feat: ops/recovery tooling + dependency hygiene (refactor phase 7)

Generalize scripts/reprocess.py from a PO-only full-sweep script into a
pipeline-general recovery tool. Targeted replay (--key/--prefix/--since)
is now the default, and the full inbound/ sweep is demoted behind an
explicit --all that documents its five hazards (async concurrency does
not serialize, use RequestResponse if order matters, metric double-count,
Bedrock re-bill, out-of-order field regression). --pipeline po|wo resolves
the correct function + bucket; dry-run-by-default / --execute is preserved.
A new tests/test_reprocess_contract.py pins the synthetic S3 event shape
and asserts the raw list_objects_v2 key is emitted untransformed (the
handler is the single decode point; a pre-decoded key would corrupt keys
containing spaces or '+').

Add docs/runbook-dlq-recovery.md: the async on-failure DLQ has no console
redrive-to-source, so it documents the receive -> extract key -> targeted
reprocess --key -> verify -> purge procedure, the real recovery windows
(14-day DLQ breadcrumb, 90-day raw-email S3 that overrides the table
RETAIN policy and is the true replay floor), and that sender-auth and
ai_fallback_rejected drops are fail-closed skips that never reach the DLQ.
Linked from the README alarms and scripts sections.

Drop the vendored boto3 floor pin from both email-processor requirements
(the Lambda runtime provides boto3; lambda-template.md empty-with-comment
form). With nothing left to install, the email-processor bundling becomes
cp-only -- the whole pip step is removed, which is the only acceptable way
the manylinux2014_aarch64 pin disappears (removing the pin while keeping a
pip install caused the PR #34 x86-wheel outage). Exact-pin moto==5.2.2 and
add pinned po/web_ui + po/site_extractor manifests (excluded from their
bundles, so hash-neutral) so their new Dependabot entries have something
to act on; add Dependabot entries for /tests, /lambdas/po/web_ui, and
/lambdas/po/site_extractor.

cdk diff is confined to exactly the two email processors' asset hashes on
both stacks. The wo/web_ui dead-manifest reduction was deliberately left
out: that manifest already ships inside the plain (non-bundled) WebUI
asset on main, so reducing or excluding it would redeploy workorder-web-ui
for no functional change -- deferred to keep the blast radius to the two
intended targets.

The untracked 44 MB lambdas/po/email_processor/package/ dir was removed
from the filesystem (asset-hash-neutral given Phase 2's package/ exclude);
it is untracked, so there is nothing to commit for it.

* Reject --all combined with --prefix/--since in reprocess.py

--all is a distinct mode (the demoted full-prefix sweep), but the args.all
branch unconditionally set prefix=inbound/ and since=None, so passing it
alongside a narrower selector silently discarded that selector. `--all
--since 2026-07-01` swept the entire corpus instead of the bounded window,
triggering every documented --all hazard (Bedrock re-bill, metric double-
count, merged-field regression) on objects the operator never targeted --
contradicting the tool's safety goal. Add the missing mutual-exclusion
guard alongside the existing --key one, and pin --all+--prefix,
--all+--since, and all three together as argparse rejections.
2026-07-20 12:53:34 -04:00

244 lines
9.7 KiB
Python

"""
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-<acct>) or 'wo' (workorder-email-processor / "
"workorder-ingest-emails-<acct>)."
),
)
# 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()