mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 08:23:14 +00:00
* 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.
244 lines
9.7 KiB
Python
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()
|