diff --git a/.claude/workflows/phase-7-ops-recovery.js b/.claude/workflows/phase-7-ops-recovery.js new file mode 100644 index 0000000..31e5a57 --- /dev/null +++ b/.claude/workflows/phase-7-ops-recovery.js @@ -0,0 +1,661 @@ +export const meta = { + name: 'phase-7-ops-recovery', + description: 'Phase 7 of the procurement-ingest refactor (docs/refactor-evaluation.md): ops/recovery + dependency hygiene. Generalize scripts/reprocess.py (--pipeline po|wo, --key/--prefix/--since; TARGETED replay primary, full-prefix DEMOTED behind --all with documented caveats) + a synthetic-event-shape contract test pinning the raw-key/no-URL-decode S3 event. New docs/runbook-dlq-recovery.md for the no-redrive async-destination DLQ (receive->key->re-invoke->verify->purge; 14d DLQ / 90d S3 windows; sender-auth + ai_fallback_rejected drops never reach the DLQ), linked from README alarms. Dependency hygiene: drop vendored boto3 from both email-processor requirements (handbook empty-with-comment form; runtime copy suffices per lambda-template.md), simplify bundling to cp-only where the vendored dep is gone, exact-pin moto==, pin-then-add dependabot entries for /tests + po/web_ui + po/site_extractor, reduce wo/web_ui dead manifest to empty-with-comment. Optional local-only rm of the untracked 44 MB package/ dir. Gates on Phase 2 merged (../lambdas asset root). Ops tooling + dep hygiene, no auth/untrusted-input/IAM change. Committed locally, never pushed.', + phases: [ + { title: 'Setup', detail: 'verify Phase 2 on base (Code.from_asset("../lambdas")), branch feature/phase-7-ops-recovery', model: 'haiku' }, + { title: 'Recon', detail: '4 mappers: reprocess.py + event shape, dependency/manifest/dependabot state, post-Phase-2 bundling + package/ dir, DLQ/alarm/S3-lifecycle names (read-only AWS)' }, + { title: 'Spec', detail: 'serial fable spec: reprocess CLI + caveats, runbook contents, requirements/moto/dependabot edits, cp-only bundling form, package/ cleanup, contract test, verify rules' }, + { title: 'Implement', detail: 'opus: scripts/ + reprocess contract test; sonnet: requirements + dependabot + cdk bundling; opus: runbook + README — disjoint files', model: 'opus' }, + { title: 'Verify', detail: 'mechanical gates + 3 fable lenses (bundling-integrity, reprocess-safety, dependency)', model: 'sonnet' }, + { title: 'Fix', detail: 'opus fixer, full re-verify, max 3 rounds', model: 'opus' }, + { title: 'Package', detail: 'single commit via -F (no push); optional local rm of package/ reported, not committed', model: 'sonnet' }, + ], +} + +// ---------------------------------------------------------------- constants + +const REPO = '/Users/adammoussa/Documents/repositories/seahaven/procurement-ingest' +const BRANCH = 'feature/phase-7-ops-recovery' +let _args = args +if (typeof _args === 'string') { + try { _args = JSON.parse(_args) } catch (e) { _args = null } +} +const BASE = (_args && _args.base) || 'main' + +const CONSTRAINTS = ` +PINNED BEHAVIORAL CONSTRAINTS (docs/refactor-evaluation.md Phase 7 — violating any is a build failure): +1. GENERALIZE scripts/reprocess.py: add --pipeline po|wo (selects the + function name + bucket per pipeline — po-email-processor / + po-ingest-emails-, workorder-email-processor / its bucket), plus + --key (single object), --prefix, and --since (time filter on + LastModified). TARGETED replay (--key / --prefix / --since) is the + PRIMARY, default mode. FULL-PREFIX replay is DEMOTED behind an explicit + --all flag. The --all caveats MUST be documented BOTH in --help text and + in code comments: Event-type (async) invocation is CONCURRENT so sorting + does NOT serialize; switch to RequestResponse if order matters; metrics + get double-counted on replay; Bedrock is re-billed; out-of-order replay + can REGRESS already-merged fields. Keep the existing dry-run-by-default / + --execute safety (targeted replay stays dry-run unless --execute). +2. SYNTHETIC-EVENT-SHAPE CONTRACT TEST (§4.7): pin the exact S3 event shape + reprocess emits — Records[0].s3.bucket.name + Records[0].s3.object.key — + and assert the key is the RAW object key with NO URL-decoding (a real S3 + notification URL-encodes the key; reprocess builds from the raw + list_objects_v2 key, and the handler is what decodes — replaying a + pre-decoded key would double-decode). Test lives under scripts' tests or + tests/ and runs in the standard pytest collection. +3. NEW docs/runbook-dlq-recovery.md: the async on-failure destination DLQ + has NO console redrive-to-source. Documented procedure: receive-message + -> extract the S3 key from the event body -> targeted re-invoke (via the + generalized reprocess.py --key) -> verify the write -> purge the message. + State the recovery WINDOWS explicitly: 14-day DLQ breadcrumb retention, + 90-day raw-email S3 retention (the inbound/ lifecycle rule OVERRIDES the + table RETAIN policy — S3 is the real replay floor). Document that + sender-auth rejections AND ai_fallback_rejected drops INTENTIONALLY never + reach the DLQ (they are fail-closed SKIPS, not errors — no retry, no DLQ + message). Link the runbook from the README alarms section. +4. DEPENDENCY HYGIENE — DROP vendored boto3 from BOTH email-processor + requirements.txt: replace 'boto3>=1.43.47' with the handbook + empty-with-comment form (a comment explaining the Lambda runtime provides + boto3 per lambda-template.md — only-boto3 functions correctly use the + runtime copy; no pinned third-party dep remains). Same empty-with-comment + reduction for wo/web_ui/requirements.txt (a dead manifest today). +5. SIMPLIFY BUNDLING TO CP-ONLY where the vendored dep is gone: once the + email-processor requirements are empty-with-comment, 'pip install -r ... ' + installs nothing, so the pip step is removable and the command becomes + cp-only. This is ONLY safe BECAUSE there are no third-party binary deps + left. Do NOT regress: the Phase 2 exclude lists + (['**/__pycache__/**','**/tests/**','**/package/**']) STAY; the widened + '../lambdas' asset root STAYS; if a Phase 3 'cp shared/*.py' line is + present on the base it STAYS. The '--platform manylinux2014_aarch64 + --only-binary=:all:' pin is only meaningful while a pip install runs — its + accidental removal caused the PR #34 outage, so removing the WHOLE pip + step (deliberately, because nothing is installed) is the only acceptable + way it disappears; you may NOT keep a pip install while dropping the pin. + If recon finds ANY third-party dep still required by an email processor, + cp-only is a blocker — keep the pinned pip step. +6. EXACT-PIN moto: change tests/requirements.txt 'moto>=5.0.0' to + 'moto==' (resolve the actually-installed/current 5.x + version; floor pins make the dependabot entry a no-op). +7. DEPENDABOT: add pip entries for '/tests', '/lambdas/po/web_ui', and + '/lambdas/po/site_extractor' — but PIN THOSE MANIFESTS FIRST. po/web_ui + and po/site_extractor have NO requirements.txt today, so a dependabot + entry pointed at them is a no-op until a PINNED manifest exists; create + each manifest (exact-pinned boto3, or empty-with-comment ONLY if the + dependabot entry would then be pointless — a dependabot entry needs at + least one pinned dep to act on). Match the existing dependabot.yml entry + shape (weekly, minor-and-patch group). tests/ gets its entry once moto is + exact-pinned (constraint 6). +8. DELETE the untracked 44 MB lambdas/po/email_processor/package/ dir as an + OPTIONAL local-cleanup step: it is UNTRACKED / local-only, so this is a + filesystem 'rm -rf', NOT a 'git rm' (there is nothing tracked to commit). + With Phase 2's '**/package/**' exclude already deployed, the deletion is + asset-hash-NEUTRAL. The Package agent PERFORMS and REPORTS it; it is NOT + part of the commit. +9. UNTOUCHABLE: no handler / lambda runtime code change under lambdas/*/ + email_processor/*.py (this is dep + ops + bundling only). No alarm, IAM, + table, env, runtime, memory, timeout, or logical-ID delta on either stack + — the ONLY cdk delta permitted is the bundling command string on the two + email processors (+ any resulting asset-hash / CDK metadata change). The + dropped-vendored-boto3 functions MUST still 'import boto3' at runtime from + the Lambda runtime copy. +10. AWS access is READ-ONLY (get-function, get-function-event-invoke-config, + SQS get-queue-attributes, S3 get-bucket-lifecycle-configuration, + describe-alarms). NEVER cdk deploy, never invoke, never mutate, never + purge a real queue. +` + +const PREAMBLE = ` +You are one of several agents building refactor Phase 7 in the git repo at +${REPO} on branch ${BRANCH} (already checked out — do NOT switch branches, +do NOT create branches, do NOT commit, NEVER push, do NOT run cdk deploy or +touch AWS resources beyond read-only calls). +Authoritative spec: docs/refactor-evaluation.md, section "Phase 7 — Ops/recovery + dependency hygiene". +Work ONLY in the files you are told you own; other agents are concurrently +editing other files in this same working tree. +${CONSTRAINTS} +Your final message is consumed by an orchestrator script, not a human — +return only the structured data requested. +` + +// ------------------------------------------------------------------ schemas + +const RECON = { + type: 'object', + required: ['summary', 'facts'], + properties: { + summary: { type: 'string' }, + facts: { type: 'array', items: { type: 'string' } }, + blockers: { type: 'array', items: { type: 'string' } }, + }, +} + +const SPEC = { + type: 'object', + required: ['reprocessCli', 'contractTest', 'runbook', 'dependencyEdits', 'bundlingEdits', 'packageCleanup', 'verifyRules', 'notes'], + properties: { + reprocessCli: { type: 'string', description: 'the complete generalized scripts/reprocess.py design: argparse surface (--pipeline po|wo, --key/--prefix/--since, --all, --execute), per-pipeline function-name/bucket resolution, which modes are dry-run-by-default, and the EXACT --help + code-comment caveat text for --all (concurrency/no-serialize, RequestResponse-if-order-matters, double-count, Bedrock re-bill, out-of-order regression)' }, + contractTest: { type: 'string', description: 'the synthetic-event-shape contract test: file path, what it asserts about Records[0].s3.bucket.name / object.key, and the exact raw-key / no-URL-decoding assertion (encoded vs raw key example)' }, + runbook: { type: 'string', description: 'the docs/runbook-dlq-recovery.md outline: the no-redrive procedure steps (receive->key->targeted re-invoke via reprocess --key->verify->purge), the 14d-DLQ / 90d-S3 windows with the lifecycle-overrides-RETAIN note, the intentional sender-auth + ai_fallback_rejected non-DLQ statement, and the exact README alarms-section link edit' }, + dependencyEdits: { type: 'string', description: 'per requirements.txt: exact new contents (both email-processor empty-with-comment, wo/web_ui empty-with-comment, po/web_ui + po/site_extractor new pinned manifests, tests/ moto==); the resolved moto version; the exact dependabot.yml additions matching the existing entry shape' }, + bundlingEdits: { type: 'string', description: 'both email-processor bundling command strings BEFORE and AFTER, verbatim: the cp-only result, proof that the pip step is safely removable (no third-party dep remains), and confirmation the exclude lists + widened root + any Phase-3 shared cp line are preserved' }, + packageCleanup: { type: 'string', description: 'the package/ dir cleanup: the exact rm command, why it is a filesystem rm not git rm (untracked), and why it is asset-hash-neutral given the Phase 2 exclude — this is a REPORTED local step, not a commit' }, + verifyRules: { type: 'string', description: 'how Verify judges success: which pytest tests must pass (incl. the new contract test), ruff scope now including scripts/, both cdk synth expectations, the dependabot yaml-parse check, requirements resolve/import check, and the runtime import-boto3 proof for the cp-only assets' }, + notes: { type: 'string' }, + }, +} + +const IMPL = { + type: 'object', + required: ['filesChanged', 'summary', 'checksRun'], + properties: { + filesChanged: { type: 'array', items: { type: 'string' } }, + summary: { type: 'string' }, + checksRun: { type: 'string' }, + blockers: { type: 'array', items: { type: 'string' } }, + }, +} + +const CHECKS = { + type: 'object', + required: ['passed', 'details'], + properties: { + passed: { type: 'boolean' }, + details: { type: 'string' }, + scopeViolations: { type: 'array', items: { type: 'string' } }, + }, +} + +const FINDINGS = { + type: 'object', + required: ['findings'], + properties: { + findings: { + type: 'array', + items: { + type: 'object', + required: ['title', 'severity', 'confirmed', 'evidence', 'fix'], + properties: { + title: { type: 'string' }, + severity: { enum: ['critical', 'high', 'medium', 'low'] }, + confirmed: { type: 'boolean' }, + evidence: { type: 'string' }, + fix: { type: 'string' }, + }, + }, + }, + }, +} + +// ------------------------------------------------------------------- setup + +phase('Setup') +const setup = await agent(` +In ${REPO}: +1. SEQUENCING GATE — Phase 2 must be on ${BASE}. Phase 7 "can run parallel + to 4-6" (largely independent), BUT the cp-only bundling simplification and + the package/ cleanup touch the same cdk bundling that Phases 2/3 + established, so this phase must NOT race Phase 2. git fetch origin, then + pick the base ref: origin/${BASE} if that remote ref exists, otherwise the + local branch ${BASE} (a stacked local-only base is expected and fine). + Verify on the base ref that BOTH cdk/po_stack.py and cdk/wo_stack.py + contain Code.from_asset("../lambdas") for the email processors + (git show :cdk/po_stack.py | grep -n '\\.\\./lambdas', same for + wo_stack.py). If either is missing, STOP with a blocker naming Phase 2 as + unmet and do nothing else. + NOTE (report, do not block): Phase 7 is otherwise independent of Phases + 4-6. If Phase 4 (cdk/common.py dedup) lands on ${BASE} FIRST, the + stack-touching bundling edits here must be rebased onto the moved bundling + code — flag that in your facts so the implement agents know whether the + bundling lives inline or in a common helper. +2. Verify clean working tree (untracked .coverage / .claude/ / the local + 44 MB lambdas/po/email_processor/package/ dir are fine; any OTHER dirt = + blocker, never stash or discard). +3. git checkout ${BASE}; then git pull --ff-only ONLY if the branch has an + upstream (a local-only base skips the pull — not a blocker); then + git checkout -b ${BRANCH} +4. gh pr list --state open --json number,title,headRefName (overlap check — + especially any in-flight Phase 2/3/4 stack PR). +Return facts: HEAD sha, Phase 2 gate evidence, whether Phase 4 has landed +(bundling inline vs helper), open PRs, blockers. +`, { label: 'setup:branch', model: 'haiku', schema: RECON }) + +if (!setup || (setup.blockers && setup.blockers.length)) { + return { aborted: 'setup blockers', blockers: setup ? setup.blockers : ['setup agent died'], facts: setup ? setup.facts : [] } +} +log(`Branch ${BRANCH} ready off ${BASE}. ${setup.summary}`) + +// ------------------------------------------------------------------- recon + +phase('Recon') +const recon = await parallel([ + () => agent(`${PREAMBLE} +Read-only recon of scripts/reprocess.py and the event shape it must pin: +1. Quote scripts/reprocess.py verbatim with file:line — the argparse surface + today (--execute only), FUNCTION_NAME, get_bucket_name (po-only, + po-ingest-emails-), list_inbound_keys (paginator, Prefix='inbound/'), + build_s3_event (the exact Records/s3/bucket.name/object.key shape), the + Event (async) InvocationType. +2. Find the WO analogues so --pipeline wo can resolve: the workorder + email-processor function name and its email bucket naming (grep the cdk + stacks / function_name / bucket_name). State them exactly. +3. Determine how the email-processor handler reads the key from an S3 event: + does it URL-decode Records[].s3.object.key (urllib.parse.unquote_plus)? + Quote the handler line. This decides the contract test's no-double-decode + assertion (reprocess must emit the RAW key so the handler decode path is + exercised exactly once — a pre-decoded key double-decodes '+'/'%20'). +4. Is there ANY existing test that touches scripts/reprocess.py? (grep tests/ + and scripts/). Where should the new synthetic-event contract test live so + it is collected by the standard pytest run (repo-root pytest collects + tests/ and both email_processor tests/ roots)? Note whether scripts/ is on + sys.path for import. +15-25 precise facts.`, + { label: 'recon:reprocess', model: 'sonnet', phase: 'Recon', schema: RECON }), + + () => agent(`${PREAMBLE} +Read-only recon of the dependency / manifest / dependabot state: +1. Quote every requirements.txt under lambdas/ and tests/ verbatim with path: + lambdas/po/email_processor (boto3>=1.43.47), lambdas/wo/email_processor + (boto3>=1.43.47), lambdas/wo/web_ui (boto3>=1.43.47 — dead manifest), + tests/ (moto>=5.0.0), cdk/. Confirm po/web_ui and po/site_extractor have + NO requirements.txt (they must be CREATED + pinned before a dependabot + entry is not a no-op). +2. Resolve the moto version to exact-pin: what 5.x is actually installed / + currently resolves (pip show moto or the lockfile / venv)? State the exact + version string to pin. +3. Quote .github/dependabot.yml verbatim: the existing pip entries (/cdk, + /lambdas/po/email_processor, /lambdas/wo/email_processor, /lambdas/wo/web_ui) + and the github-actions entry — the exact shape (weekly interval, + minor-and-patch group) the three NEW entries must match. +4. Confirm via lambda-template.md (handbook) the exact empty-with-comment + form for only-boto3 functions (the Lambda runtime provides boto3) — quote + the canonical comment wording so the edits match the handbook. +5. Does anything under lambdas/*/email_processor actually import a THIRD-PARTY + package other than boto3 (grep imports; anthropic/requests/etc.)? If yes, + cp-only is unsafe — that is a blocker for constraint 5. +15-25 facts.`, + { label: 'recon:dependencies', model: 'sonnet', phase: 'Recon', schema: RECON }), + + () => agent(`${PREAMBLE} +Read-only recon of the post-Phase-2 CDK bundling + the package/ dir: +1. Quote BOTH email-processor Code.from_asset blocks verbatim with current + file:line (cdk/po_stack.py ~241, cdk/wo_stack.py ~240): asset root, the + FULL exclude list, the entire bundling command list — the pip install + '--platform manylinux2014_aarch64 --only-binary=:all: -r + /email_processor/requirements.txt -t /asset-output' step, the + '&&', the 'cp /email_processor/*.py /asset-output/' step, and + whether a Phase-3 'cp shared/*.py /asset-output/' line is already present + on this branch's base. Note the explanatory NOTE comment in po_stack. +2. State the exact cp-only rewrite for each (drop the pip step entirely, + keep every cp line + the widened root + the exclude list). Confirm the + pip step is the ONLY thing removed. +3. The three plain from_asset calls (po/web_ui, po/site_extractor, + wo/web_ui) — quote them; they are NOT bundling-simplified here (no pip + step), just confirm they are untouched by this phase. +4. lambdas/po/email_processor/package/ — du -sh it, confirm it is UNTRACKED + (git status / git ls-files must not list it), and confirm the Phase 2 + '**/package/**' exclude already keeps it out of the asset hash (so its + removal is asset-hash-neutral). +5. If Phase 4 has landed (bundling moved into a cdk/common helper) note the + new location. +12-20 facts.`, + { label: 'recon:cdk-bundling', model: 'haiku', phase: 'Recon', schema: RECON }), + + () => agent(`${PREAMBLE} +Read-only AWS + docs recon for the runbook (region us-east-1, READ-ONLY): +1. Both email-processor DLQs: exact SQS queue names/ARNs and the + MessageRetentionPeriod (confirm the ~14-day breadcrumb window). Find them + via the stacks (dead_letter_queue / make_processor_dlq) and + aws sqs get-queue-attributes. Also the async event-invoke config / + on-failure destination if present (get-function-event-invoke-config) — + confirm it is an async on-failure destination DLQ with NO console + redrive-to-source. +2. Both raw-email S3 buckets: exact names and the inbound/ prefix lifecycle + rule (get-bucket-lifecycle-configuration) — confirm the ~90-day expiry + that overrides the table RETAIN policy (S3 is the replay floor). +3. The relevant CloudWatch alarm names the runbook references: the + sender-auth-rejected alarm(s) and the ai_fallback_rejected alarm(s) per + pipeline (describe-alarms / grep the stacks) — so the runbook states that + those two drop classes are fail-closed SKIPS that never reach the DLQ. +4. Quote the README alarms section (its heading + surrounding lines) so the + Implement agent can add the runbook link in the right place, matching + style. +Return names/ARNs/windows as facts; flag any name you could not resolve so +the runbook can be validated against real infra later.`, + { label: 'recon:dlq-alarms', model: 'sonnet', phase: 'Recon', schema: RECON }), +]) + +const reconOk = recon.filter(Boolean) +const pack = reconOk.map(r => `## ${r.summary}\n${r.facts.join('\n')}`).join('\n\n') +const reconBlockers = reconOk.flatMap(r => r.blockers || []) + .filter(b => b && !/^\s*(none|n\/a)\b/i.test(b)) +log(`Recon complete: ${reconOk.length}/4 mappers, ${reconBlockers.length} blockers`) +if (reconBlockers.length) { + return { aborted: 'recon blockers (likely a surviving third-party dep blocking cp-only, or an unresolvable pipeline name)', blockers: reconBlockers, reconPack: pack } +} + +// -------------------------------------------------------------------- spec + +phase('Spec') +const spec = await agent(`${PREAMBLE} +You are the SPEC agent — the single authority that pins every contested +decision BEFORE parallel implementation (parallel leaves cannot see each +other's choices). Using the recon pack below plus your own reads of the +actual files, produce the binding implementation spec: +- reprocessCli: the complete generalized scripts/reprocess.py — argparse + surface (--pipeline po|wo required; --key / --prefix / --since targeted + modes; --all for full-prefix; --execute preserving dry-run-by-default), + the per-pipeline function-name + bucket resolution table, and the EXACT + --help + in-code caveat text for --all (all five caveats from constraint 1 + verbatim). TARGETED is the default path; --all is the ONLY way to sweep the + whole prefix. +- contractTest: the synthetic-event-shape test — file path (collected by the + standard pytest run), the assertions on Records[0].s3.bucket.name and + object.key, and the raw-key / no-URL-decoding assertion with a concrete + encoded-vs-raw example (e.g. a key with a space or '+'). +- runbook: the full docs/runbook-dlq-recovery.md outline with the real + DLQ/queue/alarm/bucket names from recon, the no-redrive procedure steps + (receive -> extract S3 key from event body -> targeted re-invoke via + reprocess.py --key -> verify -> purge), the 14d-DLQ / 90d-S3 windows + + lifecycle-overrides-RETAIN note, the intentional sender-auth + + ai_fallback_rejected non-DLQ statement, and the exact README alarms-section + link edit. +- dependencyEdits: exact new contents of every manifest — both + email-processor requirements.txt (handbook empty-with-comment), wo/web_ui + (empty-with-comment), po/web_ui + po/site_extractor (NEW; decide pinned dep + vs whether the dependabot entry is worthwhile — a dependabot entry needs at + least one pinned dep to act on, so pin boto3 exactly if you want the entry + live), tests/ (moto== from recon) — plus the exact + dependabot.yml additions (three entries) matching the existing shape. +- bundlingEdits: both email-processor bundling command strings BEFORE and + AFTER, verbatim — the cp-only result, an explicit proof line that no + third-party dep remains (so removing the whole pip step is safe, NOT a + regression of the manylinux pin), and confirmation the exclude lists + + widened '../lambdas' root + any Phase-3 'cp shared/*.py' line survive + untouched. If Phase 4 moved bundling into a helper, target that location. +- packageCleanup: the exact rm -rf command for + lambdas/po/email_processor/package/, the note that it is a filesystem rm + (UNTRACKED — never git rm) performed and REPORTED by the Package agent, not + committed, and why it is asset-hash-neutral (Phase 2 exclude). +- verifyRules: exactly how Verify judges success — the pytest set incl. the + new contract test, ruff now INCLUDING scripts/, both cdk synth, + dependabot.yml yaml-parse, requirements resolve/import, and the runtime + 'import boto3' proof for a cp-only staged asset. +Recon pack:\n${pack}`, + { label: 'spec:pin-ops-recovery', phase: 'Spec', schema: SPEC }) + +if (!spec) return { aborted: 'spec agent died — rerun workflow', reconBlockers } +const specBlock = `BINDING SPEC (from the spec agent — implement EXACTLY this):\n${JSON.stringify(spec, null, 2)}` +log('Spec pinned: reprocess CLI + caveats, contract test, runbook, dependency edits, cp-only bundling, package cleanup') + +// --------------------------------------------------------------- implement + +phase('Implement') +const impl = await parallel([ + () => agent(`${PREAMBLE} +YOU OWN: scripts/reprocess.py and the new synthetic-event contract test file +ONLY (per spec.contractTest's path — do not touch any other test). Do not +touch cdk/, requirements, dependabot, docs, or README. +Task: rewrite scripts/reprocess.py per spec.reprocessCli — add --pipeline +po|wo, --key/--prefix/--since (TARGETED, the default), --all (DEMOTED +full-prefix, with all five caveats in BOTH --help and code comments), +preserving dry-run-by-default / --execute. Keep build_s3_event emitting the +RAW key (no URL-decoding). Then add the contract test per spec.contractTest +pinning the event shape + the raw-key / no-double-decode assertion. +Run before returning: ruff check scripts && ruff format scripts --check, +python3 scripts/reprocess.py --help (must render the caveats), and +pytest -q --no-cov -k (tests land in parallel — if +collection of other roots is red from another agent's in-flight edit, run +your test file directly; poll up to ~10 min before reporting a blocker). +${specBlock}`, + { label: 'impl:reprocess', model: 'opus', phase: 'Implement', schema: IMPL }), + + () => agent(`${PREAMBLE} +YOU OWN: every requirements.txt (lambdas/po/email_processor, +lambdas/wo/email_processor, lambdas/wo/web_ui, tests/, and the two NEW +manifests lambdas/po/web_ui/requirements.txt + +lambdas/po/site_extractor/requirements.txt), .github/dependabot.yml, and the +bundling regions of cdk/po_stack.py + cdk/wo_stack.py ONLY (only the +email-processor from_asset bundling command — alarms, IAM, tables, env, +memory, timeout, the three plain from_asset calls are all off-limits). +Task: apply spec.dependencyEdits (drop vendored boto3 -> empty-with-comment +in both email-processor + wo/web_ui; create the two pinned po manifests; +exact-pin moto==) + the three dependabot entries + spec.bundlingEdits +(cp-only: remove ONLY the pip step from both email-processor commands, +preserving the exclude lists, widened root, and any Phase-3 shared cp line). +Do NOT delete package/ (that is the Package agent's reported step). +Run before returning: ruff check cdk, python3 -c 'import yaml,pathlib; +yaml.safe_load(pathlib.Path(".github/dependabot.yml").read_text())' (parses), +cd cdk && npx cdk synth po-ingest -q -o /tmp/phase7-synth && +npx cdk synth workorder-ingest -q -o /tmp/phase7-synth (artifact-id +selectors). Docker bundling runs — confirm each staged email-processor asset +still contains handler.py + its siblings (and shared/*.py if Phase 3 landed) +and that 'import boto3' resolves at runtime from the Lambda runtime copy (the +staged asset need not vendor it). Note asset hashes. If the other agents' +edits have not landed the synth is unaffected (bundling is yours) — but poll +~10 min if Docker is slow before a blocker. +${specBlock}`, + { label: 'impl:deps-bundling', model: 'sonnet', phase: 'Implement', schema: IMPL }), + + () => agent(`${PREAMBLE} +YOU OWN: docs/runbook-dlq-recovery.md (NEW) and README.md ONLY. +Task A: write docs/runbook-dlq-recovery.md per spec.runbook — the no-redrive +async-destination DLQ procedure (receive-message -> extract the S3 key from +the event body -> targeted re-invoke via 'python scripts/reprocess.py +--pipeline --key --execute' -> verify the write -> purge the +message), the real queue/bucket/alarm names from the spec, the recovery +windows (14-day DLQ breadcrumb, 90-day raw-email S3 with the +lifecycle-overrides-RETAIN note), and the explicit statement that sender-auth +rejections AND ai_fallback_rejected drops are fail-closed SKIPS that never +reach the DLQ (no retry, no DLQ message). +Task B: README — link the new runbook from the alarms section per +spec.runbook's exact edit, and note reprocess.py is now pipeline-general +(--pipeline / --key / --prefix / --since targeted; --all demoted). Match +existing README style. +Run before returning: confirm the runbook links resolve and the README edit +is in the alarms section. (No pytest owned by you; if you want a doc-lint, +ruff does not lint .md.) +${specBlock}`, + { label: 'impl:runbook-readme', model: 'opus', phase: 'Implement', schema: IMPL }), +]) + +const implOk = impl.filter(Boolean) +const implBlockers = implOk.flatMap(r => r.blockers || []) +log(`Implement complete: ${implOk.length}/3 agents, blockers: ${implBlockers.length}`) + +// ---------------------------------------------------- verify + fix loop + +const EXPECTED_SCOPE = [ + 'scripts/reprocess.py', + 'tests/', + 'docs/runbook-dlq-recovery.md', + 'lambdas/po/email_processor/requirements.txt', + 'lambdas/wo/email_processor/requirements.txt', + 'lambdas/wo/web_ui/requirements.txt', + 'lambdas/po/web_ui/requirements.txt', + 'lambdas/po/site_extractor/requirements.txt', + '.github/dependabot.yml', + 'cdk/po_stack.py', + 'cdk/wo_stack.py', + 'README.md', +] + +const mechanicalPrompt = `${PREAMBLE} +Independent re-verification — trust nothing self-reported. Run ALL gates, +quoting failures verbatim: +1. pytest -q --no-cov (repo root, all roots green) — INCLUDING the new + synthetic-event contract test (name it and quote its pass line). +2. ruff check . && ruff format --check . — scripts/ is now IN scope (it was + previously omitted); confirm scripts/reprocess.py lints and formats clean. +3. cd cdk && npx cdk synth po-ingest -q && npx cdk synth workorder-ingest -q + (Docker bundling runs; both must succeed). +4. cd cdk && npx cdk diff po-ingest ; npx cdk diff workorder-ingest — + constraint 9: the ONLY resource delta allowed is the two email + processors' Code/S3Key (asset) + CDK metadata. Any alarm/IAM/env/runtime/ + memory/timeout/handler-prop/logical-ID delta = FAIL. Paste the diff + summaries. +5. DEPENDABOT: python3 -c 'import yaml,pathlib; d=yaml.safe_load(pathlib.Path(".github/dependabot.yml").read_text()); print([u["directory"] for u in d["updates"]])' + — the three new dirs (/tests, /lambdas/po/web_ui, /lambdas/po/site_extractor) + present, shape matches existing entries. +6. MANIFESTS: every email-processor + wo/web_ui requirements.txt is the + empty-with-comment form (no pinned third-party dep); tests/requirements.txt + is moto==; po/web_ui + po/site_extractor manifests EXIST and are + PINNED (a dependabot entry pointed at a floor/absent manifest is a no-op — + FAIL if unpinned). pip install --dry-run (or pip-compile check) resolves + each without error. +7. RUNTIME IMPORT: no lambda handler .py under lambdas/*/email_processor + changed (git diff ${BASE}...HEAD -- must be empty for those); 'import boto3' + still resolves — prove the Lambda runtime copy suffices (the cp-only staged + asset need not vendor boto3). +8. RUNBOOK: docs/runbook-dlq-recovery.md exists, states the 14d/90d windows, + the no-redrive procedure, and the sender-auth + ai_fallback_rejected + non-DLQ note; README alarms section links it. +9. git status --porcelain scope check: every modified/added path under + ${EXPECTED_SCOPE.join(', ')} (untracked .coverage/.claude/package/ + tolerated; package/ must NOT be staged). +passed=true only if all green. YOU MAY NOT edit files.` + +const lenses = [ + { key: 'bundling-integrity', prompt: `${PREAMBLE} +ADVERSARIAL REVIEW — bundling-integrity lens. The bundling was simplified to +cp-only and vendored boto3 was dropped; try to prove a runtime ImportError +was introduced (the PR #34 / PR #105 failure class). (1) cd cdk && synth both +stacks; inside each staged email-processor asset dir list the files and run +python3 -c "import handler" with that dir alone on sys.path (env-var stubs as +needed) — does every runtime-imported module still ship (handler siblings + +shared/*.py if Phase 3 landed)? (2) Does 'import boto3' succeed against the +Lambda runtime copy — i.e. is boto3 genuinely NOT needed in the zip, or does +some code path need a newer boto3 than the runtime ships? (3) Are the Phase +2/3 excludes (['**/__pycache__/**','**/tests/**','**/package/**']) STILL +present and is the widened '../lambdas' root intact? (4) Was the manylinux +pin removed ONLY as part of removing the whole pip step (because nothing is +installed), NOT while keeping a pip install? A cp-only command that still +tries to pip-install without the pin, or a surviving third-party dep, is a +critical finding. (5) DETERMINISM: synth po-ingest twice into fresh -o dirs — +identical asset hashes. confirmed=true only with a concrete ImportError +sketch or file:line proof.` }, + { key: 'reprocess-safety', prompt: `${PREAMBLE} +ADVERSARIAL REVIEW — reprocess-safety lens. Attack the generalized +reprocess.py. (1) Is the DEFAULT behavior TARGETED (—key/—prefix/—since), and +is it IMPOSSIBLE to sweep the whole inbound/ prefix without the explicit +--all flag? Trace the arg parsing — can a missing/empty --prefix silently +fall through to a full-prefix replay? That would be a critical finding +(accidental mass re-invoke = mass Bedrock re-bill + merged-field regression). +(2) Are ALL FIVE --all caveats present in BOTH --help and code comments +(concurrency/no-serialize, RequestResponse-if-order-matters, double-count, +Bedrock re-bill, out-of-order regression)? (3) Does build_s3_event emit the +RAW key with NO URL-decoding, and does the contract test actually FAIL if +someone adds an unquote_plus (mutate it on a scratch copy and confirm red)? +A key with a space/'+' must round-trip raw so the handler decodes exactly +once. (4) --pipeline wo resolves the correct function name + bucket (not the +PO defaults)? confirmed=true only with file:line or reproduced-failure +evidence.` }, + { key: 'dependency', prompt: `${PREAMBLE} +ADVERSARIAL REVIEW — dependency lens. (1) moto: is tests/requirements.txt +EXACT-pinned (moto==), not a floor (>=)? A floor makes the tests/ +dependabot entry a no-op. (2) The three NEW dependabot dirs: do +/lambdas/po/web_ui and /lambdas/po/site_extractor now have PINNED manifests +that actually exist? A dependabot entry pointed at an absent or floor-pinned +manifest does nothing — that is a finding. Cross-check dependabot.yml +directory strings against the real file paths (a typo'd directory silently +no-ops). (3) The email-processor + wo/web_ui requirements are the handbook +empty-with-comment form — confirm no pinned third-party dep is left that +would then require the pip step back (which would contradict the cp-only +bundling). (4) Does the empty-with-comment wording match lambda-template.md +(runtime provides boto3)? confirmed=true only with file:line evidence.` }, +] + +let round = 0 +let checks = null +let confirmed = [] +while (round < 3) { + phase('Verify') + const results = await parallel([ + () => agent(mechanicalPrompt, { label: `verify:mechanical-r${round}`, model: 'sonnet', phase: 'Verify', schema: CHECKS }), + ...lenses.map(l => () => + agent(l.prompt, { label: `verify:${l.key}-r${round}`, phase: 'Verify', schema: FINDINGS })), + ]) + checks = results[0] + confirmed = results.slice(1).filter(Boolean) + .flatMap(r => r.findings || []) + .filter(f => f.confirmed && f.severity !== 'low') + const green = checks && checks.passed + log(`Verify round ${round}: mechanical ${checks && checks.passed ? 'GREEN' : 'RED'}, confirmed findings: ${confirmed.length}`) + if (green && confirmed.length === 0) break + + round += 1 + if (round >= 3) break + phase('Fix') + await agent(`${PREAMBLE} +You are the fix agent — you may edit files under: ${EXPECTED_SCOPE.join(', ')}. +Fix EVERY item below minimally; the binding spec and 10 pinned constraints +still hold (a finding that conflicts with a constraint is reported, not +"fixed" — the constraint wins, esp. constraint 5's cp-only-only-if-no-dep +rule, constraint 8's do-NOT-git-rm-package, and constraint 9's no-handler / +no-alarm change). Re-run the specific failing gate/test per fix. +MECHANICAL:\n${checks ? checks.details : '(agent died — rerun all gates)'} +CONFIRMED FINDINGS:\n${JSON.stringify(confirmed, null, 2)} +${specBlock}`, + { label: `fix:round-${round}`, model: 'opus', phase: 'Fix', schema: IMPL }) +} + +const verifyClean = checks && checks.passed && confirmed.length === 0 +if (!verifyClean) { + return { + status: 'NEEDS ATTENTION — verify not clean after 3 rounds; branch left uncommitted', + branch: BRANCH, + mechanical: checks, + unresolvedFindings: confirmed, + implBlockers, + reconBlockers, + spec, + } +} + +// ----------------------------------------------------------------- package + +phase('Package') +const commit = await agent(`${PREAMBLE.replace('do NOT commit, ', '')} +YOU are the commit agent: +1. Read ~/Documents/repositories/seahaven/engineering-handbook/commit-messages.md + and follow it exactly. +2. OPTIONAL LOCAL CLEANUP (constraint 8): rm -rf + lambdas/po/email_processor/package/ — it is UNTRACKED, so this is a + filesystem rm, NOT git rm; there is NOTHING to stage from it and it is + asset-hash-neutral (Phase 2 exclude). PERFORM it and REPORT the reclaimed + ~44 MB in your summary. It is NOT part of the commit. +3. git add only paths under: ${EXPECTED_SCOPE.join(', ')} and + .claude/workflows/phase-7-ops-recovery.js. NOT .coverage, NOT package/ + (it is gone / untracked either way). Verify the staged set with + git status — the two NEW manifests and the new runbook MUST be staged. +4. ONE commit; write the message to /tmp/phase7-commit-msg.txt and use + git commit -F /tmp/phase7-commit-msg.txt (backticks in -m get eaten by + zsh). Suggested subject: + "feat: ops/recovery tooling + dependency hygiene — generalized reprocess, DLQ runbook, cp-only bundling (refactor phase 7)" + Body: the reprocess generalization (targeted-primary, --all demoted with + caveats) + contract test, the DLQ runbook + windows, the dropped vendored + boto3 / cp-only bundling, the moto exact-pin + three dependabot entries + with pinned manifests, and the reported (not committed) package/ cleanup. + NO AI attribution / Co-Authored-By lines. +5. Do NOT push. Return commit sha + shortstat + the package/ cleanup result + in summary.`, + { label: 'package:commit', model: 'sonnet', phase: 'Package', schema: IMPL }) + +return { + status: 'BUILT — committed locally, NOT pushed', + branch: BRANCH, + base: BASE, + commit: commit ? commit.summary : 'commit agent died — commit manually', + spec: { reprocessCli: spec.reprocessCli, runbook: spec.runbook, dependencyEdits: spec.dependencyEdits, bundlingEdits: spec.bundlingEdits, packageCleanup: spec.packageCleanup }, + implementation: implOk.map(r => r.summary), + filesChanged: implOk.flatMap(r => r.filesChanged), + verifyRounds: round + 1, + blockers: implBlockers.concat(reconBlockers), + outstandingGates: [ + '/sh-security-review NOT required (ops tooling + dependency hygiene only — no untrusted-input parsing, no auth code, no IAM/policy change). The pre-push deterministic scanners still run as the unattended backstop.', + 'cross-family cross_review.py NOT required (no IAM/policy change; no Lambda handler signature / event-shape / return-contract change — reprocess.py emits the SAME S3 event shape it always did, pinned by the new contract test).', + 'push + PR + gh pr checks green', + 'deploy-then-merge for the bundling/dep changes: deploy from branch, smoke green, one real email per pipeline (the dropped-vendored-boto3 email processors MUST still import boto3 at runtime from the Lambda runtime copy), expect ONE benign asset-hash redeploy from the package/ cleanup + cp-only change, THEN merge.', + 'reprocess.py + docs/runbook-dlq-recovery.md are ops/docs (no deploy risk), but the runbook should be VALIDATED against the real DLQ / alarm / bucket names before it is relied on in an incident.', + 'Confluence "AWS Architecture Map" — Phase 7 adds the DLQ-recovery runbook to the inventory; note the new runbook when the resource inventory is next refreshed.', + ], +} diff --git a/.github/dependabot.yml b/.github/dependabot.yml index 7c6531e..92ce50a 100644 --- a/.github/dependabot.yml +++ b/.github/dependabot.yml @@ -36,6 +36,33 @@ updates: update-types: - "minor" - "patch" + - package-ecosystem: "pip" + directory: "/tests" + schedule: + interval: "weekly" + groups: + minor-and-patch: + update-types: + - "minor" + - "patch" + - package-ecosystem: "pip" + directory: "/lambdas/po/web_ui" + schedule: + interval: "weekly" + groups: + minor-and-patch: + update-types: + - "minor" + - "patch" + - package-ecosystem: "pip" + directory: "/lambdas/po/site_extractor" + schedule: + interval: "weekly" + groups: + minor-and-patch: + update-types: + - "minor" + - "patch" - package-ecosystem: "github-actions" directory: "/" schedule: diff --git a/README.md b/README.md index 8fb8937..3d25999 100644 --- a/README.md +++ b/README.md @@ -138,7 +138,7 @@ The `-sender-auth-rejected` alarm closes the silent-drop gap in INFRA-107: a The `-duration` and `-throttles` alarms for `po-email-processor` and `workorder-email-processor` supersede the orphaned, CLI-created `Lambda-Duration-*` / `Lambda-Throttles-*` alarms (deleted post-deploy). -**DLQ alarms** (`AWS/SQS`): `po-email-processor-dlq-messages` and `workorder-email-processor-dlq-messages` fire when any message is visible on an email-processor DLQ (`ApproximateNumberOfMessagesVisible` Maximum, 5 min, `> 0`, eval 1) — a message there means an email was dropped after Lambda exhausted its async retries. +**DLQ alarms** (`AWS/SQS`): `po-email-processor-dlq-messages` and `workorder-email-processor-dlq-messages` fire when any message is visible on an email-processor DLQ (`ApproximateNumberOfMessagesVisible` Maximum, 5 min, `> 0`, eval 1) — a message there means an email was dropped after Lambda exhausted its async retries. Recovery from a DLQ message (no console redrive) is documented in the [DLQ recovery runbook](docs/runbook-dlq-recovery.md). **Parse-outcome metric + fallback-rate alarm (workorder-ingest):** the WO processor writes one CloudWatch **EMF** line per email to namespace `Seahaven/WorkorderIngest`, metric `ParseOutcome` (Unit Count, value 1), dimensioned by `ParseMethod` (`template` | `ai_fallback` | `ai_fallback_rejected`) and `TemplateId` (`update_plaintext` | `assign_html` | `unknown`). `ai_fallback_rejected` counts AI-fallback output that failed the fail-closed `validate_ai_fallback()` gate (schema/enum/date contract on raw Bedrock output — prompt-injection defence) and was dropped without a DynamoDB write. Non-dimension EMF properties `ReasonCode` and `work_order_id` are queryable in Logs Insights but not promoted to metrics (kept low-cardinality). EMF is used instead of `PutMetricData` so there is no extra sync call / latency / IAM grant on the async hot path (the role already has `logs:PutLogEvents`). The alarm `workorder-email-processor-template-fallback-rate` fires when the AI-fallback share of parses — rejected fallback parses included, so a drift outage whose AI output also fails the gate cannot lower the observed rate while dropping mail — exceeds **15%** sustained (a `MathExpression` with `FILL(...,0)` and a ≥10-sample volume floor over 15-minute periods, eval 3 / datapoints 2) — catching Hexagon template-drift coverage collapse while the volume floor + `FILL` prevent low-volume false pages / `INSUFFICIENT_DATA`. ALARM-only `SnsAction` to `site-alerts`, no OK action, `NOT_BREACHING`. The 15-minute period is a deliberate deviation from the 5-minute house style to accumulate a stable denominator at the low ~760/day volume. A second alarm, `workorder-email-processor-ai-fallback-rejected`, pages on the rejected series itself (≥1 rejection per 5-min period, 2 of the last 6 periods — the sender-auth-rejected sparse-arrival idiom) because a gate rejection drops mail without error/retry/DLQ and would otherwise be silent. @@ -305,10 +305,11 @@ All three roots are discovered by `pytest.ini` (`testpaths`). ## Scripts -**Reprocess PO emails** (re-run parser against all emails still in S3): +**Reprocess emails** (re-invoke an email-processor with a synthetic S3 event). `reprocess.py` is now **pipeline-general**: `--pipeline po|wo` selects the function + raw-email bucket. **Targeted replay** (`--key` one object, `--prefix`, or `--since` a `LastModified` timestamp) is the default, preferred mode; the full-prefix sweep is demoted behind an explicit `--all`. Every mode is dry-run unless `--execute`. See the [DLQ recovery runbook](docs/runbook-dlq-recovery.md) for the targeted single-key re-invoke flow. ```bash -python scripts/reprocess.py # dry-run -python scripts/reprocess.py --execute # invoke po-email-processor for each +python scripts/reprocess.py --pipeline po --key inbound/2026/msg.eml # targeted, dry-run +python scripts/reprocess.py --pipeline po --key inbound/2026/msg.eml --execute # targeted, re-invoke one +python scripts/reprocess.py --pipeline wo --all --execute # demoted full-prefix sweep (see --all caveats) ``` **Backfill verified sites** (one-time scan of historical POs): diff --git a/cdk/po_stack.py b/cdk/po_stack.py index 9946a94..4cd0614 100644 --- a/cdk/po_stack.py +++ b/cdk/po_stack.py @@ -246,8 +246,6 @@ class PoIngestStack(Stack): command=[ "bash", "-c", - "pip install --platform manylinux2014_aarch64 --only-binary=:all: " - "-r po/email_processor/requirements.txt -t /asset-output && " # NOTE: non-recursive glob (not `cp -r`) so tests/ and # the stale package/ dir are never shipped -- only # top-level .py siblings of handler.py. This replaces @@ -261,6 +259,9 @@ class PoIngestStack(Stack): # handler.py's first-party imports and asserts this # command ships all of them, so a future revert back # to an allowlist that omits a sibling fails CI. + # pip step removed in Phase 7: requirements.txt is now + # empty (boto3 comes from the Lambda runtime), so nothing + # is installed and the manylinux pin has nothing to pin. "cp po/email_processor/*.py /asset-output/", ], ), @@ -601,7 +602,14 @@ class PoIngestStack(Stack): architecture=lambda_.Architecture.ARM_64, handler="handler.handler", code=lambda_.Code.from_asset( - "../lambdas/po/web_ui", exclude=["**/__pycache__/**"] + # requirements.txt is excluded from the bundle: it exists only + # as a Dependabot anchor (git-based scan sees it), never pip- + # installed (this is a plain non-bundled asset) and never needed + # at runtime (boto3 comes from the Lambda runtime). Excluding it + # keeps the deployed asset hash neutral vs base while the manifest + # still lands in git for Dependabot. + "../lambdas/po/web_ui", + exclude=["**/__pycache__/**", "requirements.txt"], ), timeout=Duration.seconds(60), memory_size=256, @@ -682,7 +690,14 @@ class PoIngestStack(Stack): architecture=lambda_.Architecture.ARM_64, handler="handler.handler", code=lambda_.Code.from_asset( - "../lambdas/po/site_extractor", exclude=["**/__pycache__/**"] + # requirements.txt is excluded from the bundle: it exists only + # as a Dependabot anchor (git-based scan sees it), never pip- + # installed (this is a plain non-bundled asset) and never needed + # at runtime (boto3 comes from the Lambda runtime). Excluding it + # keeps the deployed asset hash neutral vs base while the manifest + # still lands in git for Dependabot. + "../lambdas/po/site_extractor", + exclude=["**/__pycache__/**", "requirements.txt"], ), timeout=Duration.seconds(60), memory_size=256, diff --git a/cdk/wo_stack.py b/cdk/wo_stack.py index f302b00..a99be30 100644 --- a/cdk/wo_stack.py +++ b/cdk/wo_stack.py @@ -245,8 +245,9 @@ class WorkorderIngestStack(Stack): command=[ "bash", "-c", - "pip install --platform manylinux2014_aarch64 --only-binary=:all: " - "-r wo/email_processor/requirements.txt -t /asset-output && " + # pip step removed in Phase 7: requirements.txt is now empty + # (boto3 comes from the Lambda runtime), so nothing is installed + # and the manylinux pin has nothing to pin. cp-only is safe. "cp wo/email_processor/*.py /asset-output/", ], ), diff --git a/docs/runbook-dlq-recovery.md b/docs/runbook-dlq-recovery.md new file mode 100644 index 0000000..d945f03 --- /dev/null +++ b/docs/runbook-dlq-recovery.md @@ -0,0 +1,117 @@ +# DLQ Recovery Runbook — Email-Processor Dead-Letter Queues + +Operational procedure for draining an email-processor dead-letter queue (DLQ) +after a failed async parse. There is **no console redrive-to-source** for these +queues; recovery is a manual, targeted re-invoke via `scripts/reprocess.py`. + +## Scope / mechanism + +Both email processors set `dead_letter_queue=` on the `lambda_.Function` +construct. This is the legacy per-function **Lambda `DeadLetterConfig`** +(asynchronous-invocation DLQ), **not** an EventInvokeConfig on-failure +Destination — `aws lambda get-function-event-invoke-config` returns +`ResourceNotFoundException` for both functions (no destination config exists). A +failed async invocation lands on the DLQ only **after** Lambda exhausts its +automatic retries. + +There is **no console redrive-to-source**: the SQS console's "redrive to source" +applies only to SQS-to-SQS DLQ relationships, not to a Lambda `DeadLetterConfig` +target. Recovery is manual, via targeted re-invoke. + +## Resource inventory (acct 328440206208, us-east-1) + +| Pipeline | Function | DLQ queue name | Raw-email bucket | DLQ alarm | +|---|---|---|---|---| +| PO | `po-email-processor` | `po-ingest-EmailProcessorDlqA753DED5-az8LUZE3ubtz` | `po-ingest-emails-328440206208` | `po-email-processor-dlq-messages` | +| WO | `workorder-email-processor` | `WorkorderIngestStack-EmailProcessorDlqA753DED5-Q8H555LrqSU1` | `workorder-ingest-emails-328440206208` | `workorder-email-processor-dlq-messages` | + +DLQ URLs are `https://sqs.us-east-1.amazonaws.com/328440206208/`. +Both DLQs: 14-day retention, SSE, TLS-enforced, `VisibilityTimeout` 30s. + +## Recovery procedure (no redrive — receive → extract key → targeted re-invoke → verify → purge) + +1. **Trigger.** The `-dlq-messages` alarm fires + (`ApproximateNumberOfMessagesVisible` Maximum, 5 min, `> 0`, eval 1). + +2. **RECEIVE** the message (do not purge yet): + + ```bash + aws sqs receive-message \ + --queue-url \ + --max-number-of-messages 1 \ + --visibility-timeout 120 \ + --wait-time-seconds 5 + ``` + + Capture the `ReceiptHandle` from the response. + +3. **EXTRACT the S3 key.** The `DeadLetterConfig` message `Body` is the original + async invocation payload — the S3 event JSON. Read + `Records[0].s3.bucket.name` and `Records[0].s3.object.key` from the `Body`. + This is the **raw** object key that reprocess/S3 emitted (no URL-decoding + applied). + +4. **RE-INVOKE (targeted; dry-run first).** Confirm the key with a dry-run, then + execute: + + ```bash + # PO queue: + python scripts/reprocess.py --pipeline po --key '' # dry-run + python scripts/reprocess.py --pipeline po --key '' --execute # re-invoke + + # WO queue: + python scripts/reprocess.py --pipeline wo --key '' --execute + ``` + + This re-invokes the **same** function with the same raw-key synthetic S3 + event (a single object — **not** `--all`). + +5. **VERIFY the write.** Confirm the downstream effect landed before proceeding: + the DynamoDB item exists / was updated (`purchase-orders` for PO, + `WorkOrders` for WO), and the function's log group shows a clean parse (no new + error, no new DLQ message). Do not proceed until verified. + +6. **PURGE the one message.** Delete only the processed message by its + `ReceiptHandle`: + + ```bash + aws sqs delete-message --queue-url --receipt-handle '' + ``` + + Do **not** `purge-queue` — that would drop unexamined breadcrumbs. AWS access + is otherwise read-only; `delete-message` on a DLQ you are actively draining is + the one write this runbook performs. + +## Recovery windows + +- **DLQ breadcrumb retention: 14 days** (`MessageRetentionPeriod=1209600s`, + confirmed live on both queues). After 14 days the breadcrumb is gone. +- **Raw-email S3 retention: 90 days.** Both buckets have + `RemovalPolicy.RETAIN`, **but** a single Enabled lifecycle rule + (`Expiration.Days=90`, `Filter.Prefix=''`) expires objects bucket-wide. **The + lifecycle rule OVERRIDES the RETAIN policy — S3 is the real replay floor:** a + raw email is gone at ~90 days regardless of the table RETAIN policy. + +Because 14 days (DLQ) < 90 days (S3), any object referenced by a live DLQ +breadcrumb is always still in S3, so DLQ replay within its 14-day window is never +blocked by S3 expiry. The 90-day floor binds only for replays reconstructed from +other sources (e.g. logs) after the breadcrumb has expired. + +## Intentional non-DLQ drops (do NOT hunt for these in the DLQ) + +Two drop classes **never** produce a DLQ message because they are fail-closed +**skips, not errors** — the handler returns normally (no raise, no retry, no DLQ +message): + +- **Sender-auth rejections.** Logged as a structured `sender_auth_rejected` + warning and skipped; covered by the `-sender-auth-rejected` alarm + (log-metric filter), **not** the DLQ. +- **`ai_fallback_rejected` drops.** AI-fallback output that failed the + fail-closed `validate_ai_fallback()` gate, emitted as a + `ParseMethod=ai_fallback_rejected` EMF datapoint and dropped without a + DynamoDB write; covered by the `-ai-fallback-rejected` alarm, **not** the + DLQ. + +If mail is missing but the DLQ is empty, check those two alarms / log filters — +the email was intentionally rejected. Re-invoking it via reprocess will just be +rejected again; fix the sender-auth config or the upstream email, not the DLQ. diff --git a/lambdas/po/email_processor/requirements.txt b/lambdas/po/email_processor/requirements.txt index dd0c8a3..22d2b5f 100644 --- a/lambdas/po/email_processor/requirements.txt +++ b/lambdas/po/email_processor/requirements.txt @@ -1 +1,2 @@ -boto3>=1.43.47 +# Per-function dependencies. Leave empty if the function uses only boto3 and stdlib. +# boto3 is provided by the Lambda Python runtime (lambda-template.md) — do not vendor it. diff --git a/lambdas/po/site_extractor/requirements.txt b/lambdas/po/site_extractor/requirements.txt new file mode 100644 index 0000000..2a345ad --- /dev/null +++ b/lambdas/po/site_extractor/requirements.txt @@ -0,0 +1 @@ +boto3==1.43.31 diff --git a/lambdas/po/web_ui/requirements.txt b/lambdas/po/web_ui/requirements.txt new file mode 100644 index 0000000..2a345ad --- /dev/null +++ b/lambdas/po/web_ui/requirements.txt @@ -0,0 +1 @@ +boto3==1.43.31 diff --git a/lambdas/wo/email_processor/requirements.txt b/lambdas/wo/email_processor/requirements.txt index dd0c8a3..22d2b5f 100644 --- a/lambdas/wo/email_processor/requirements.txt +++ b/lambdas/wo/email_processor/requirements.txt @@ -1 +1,2 @@ -boto3>=1.43.47 +# Per-function dependencies. Leave empty if the function uses only boto3 and stdlib. +# boto3 is provided by the Lambda Python runtime (lambda-template.md) — do not vendor it. diff --git a/scripts/reprocess.py b/scripts/reprocess.py index 38920ce..1037cd6 100644 --- a/scripts/reprocess.py +++ b/scripts/reprocess.py @@ -1,39 +1,93 @@ """ -Re-invoke the po-email-processor Lambda for every object still -sitting under the inbound/ prefix in the email bucket. +Reprocess missed procurement-ingest emails by re-invoking an email-processor +Lambda with a synthetic S3 event. -Usage: - python scripts/reprocess.py # dry-run (list only) - python scripts/reprocess.py --execute # actually invoke +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 -sts = boto3.client("sts") -s3 = boto3.client("s3") -lambda_client = boto3.client("lambda") +# 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}", + }, +} -FUNCTION_NAME = "po-email-processor" +DEFAULT_PREFIX = "inbound/" -def get_bucket_name() -> str: - account_id = sts.get_caller_identity()["Account"] - return f"po-ingest-emails-{account_id}" +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_inbound_keys(bucket: str) -> list[str]: - keys = [] +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="inbound/"): + 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": [ { @@ -47,28 +101,134 @@ def build_s3_event(bucket: str, key: str) -> dict: def main(): - parser = argparse.ArgumentParser(description="Reprocess missed PO emails") + 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)", + help=( + "Actually invoke the Lambda. Default is dry-run (list only). " + "Applies to EVERY mode, including --all." + ), ) args = parser.parse_args() - bucket = get_bucket_name() - keys = list_inbound_keys(bucket) + # 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("No objects found under inbound/ — nothing to reprocess.") + print(f"No matching objects in s3://{bucket}/ — nothing to reprocess.") return - print(f"Found {len(keys)} email(s) in s3://{bucket}/inbound/\n") + 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 for {key} ... ", end="", flush=True) + print(f" Invoking {function_name} for {key} ... ", end="", flush=True) resp = lambda_client.invoke( - FunctionName=FUNCTION_NAME, + FunctionName=function_name, InvocationType="Event", # async — don't wait for each one Payload=json.dumps(build_s3_event(bucket, key)), ) diff --git a/tests/requirements.txt b/tests/requirements.txt index 021a495..0c521bd 100644 --- a/tests/requirements.txt +++ b/tests/requirements.txt @@ -1 +1 @@ -moto>=5.0.0 +moto==5.2.2 diff --git a/tests/test_reprocess_contract.py b/tests/test_reprocess_contract.py new file mode 100644 index 0000000..a899966 --- /dev/null +++ b/tests/test_reprocess_contract.py @@ -0,0 +1,124 @@ +"""Synthetic-event-shape contract test for scripts/reprocess.py (refactor §4.7). + +Pins the exact S3 event shape reprocess emits and asserts the object key is +the RAW list_objects_v2 key — reprocess applies NO URL-encoding or decoding. + +Why this matters: a real S3 event notification URL-encodes the object key, and +the handler is where any decode would live. reprocess builds its synthetic +event from the RAW list_objects_v2 obj["Key"] (or the raw --key arg), so it +must emit that key byte-for-byte. Replaying a pre-decoded (or pre-encoded) key +would not match what s3.get_object(Bucket, Key) expects downstream. + +Recon note (honored narrowly): neither the PO nor the WO handler imports +urllib.parse or calls unquote/unquote_plus today — the raw event key goes +straight into s3.get_object. So the "handler decodes -> replaying a decoded +key double-decodes" failure mode is NOT reproducible against today's handlers +(they never decode once). This test therefore guards ONLY reprocess's own +contract: emit the raw list_objects_v2 key with no transformation. That keeps +replay single-decode-correct if a future handler ever adds unquote_plus. +""" + +import importlib.util +import os +import pathlib +import sys + +import pytest + +# reprocess.py builds its boto3 clients inside main()/functions, so importing +# it is side-effect-free (no AWS calls, no clients at module scope). Set dummy +# region/creds defensively anyway so import stays offline-safe even if that +# ever changes. build_s3_event itself makes no AWS calls. +os.environ.setdefault("AWS_DEFAULT_REGION", "us-east-1") +os.environ.setdefault("AWS_ACCESS_KEY_ID", "testing") +os.environ.setdefault("AWS_SECRET_ACCESS_KEY", "testing") +os.environ.setdefault("AWS_SESSION_TOKEN", "testing") + +# scripts/ is not a package (no __init__.py) and is not on sys.path, and +# tests/conftest.py only wires up per-handler dirs — so load reprocess.py by +# file path with importlib rather than `import scripts.reprocess`. +_p = pathlib.Path(__file__).resolve().parents[1] / "scripts" / "reprocess.py" +_spec = importlib.util.spec_from_file_location("reprocess", _p) +reprocess = importlib.util.module_from_spec(_spec) +_spec.loader.exec_module(reprocess) + + +def test_build_s3_event_shape_and_raw_key(): + bucket = "po-ingest-emails-328440206208" + # Deliberately contains a SPACE, a literal '+', and a literal '%41'. + # A real S3 event notification would deliver this key encoded as + # "inbound/2026/AB+12%2B34+%2541.eml" + # (space->'+', '+'->'%2B', '%'->'%25'); the handler is where any decode + # would live. reprocess builds from the RAW list_objects_v2 key, so it + # must emit "inbound/2026/AB 12+34 %41.eml" unchanged. + raw_key = "inbound/2026/AB 12+34 %41.eml" + + event = reprocess.build_s3_event(bucket, raw_key) + + # 1. Exact top-level shape: a single Records entry. + assert list(event.keys()) == ["Records"] + assert len(event["Records"]) == 1 + + # 2. Bucket name is nested at Records[0].s3.bucket.name. + assert event["Records"][0]["s3"]["bucket"]["name"] == bucket + + # 3. Object key is nested at Records[0].s3.object.key, byte-for-byte. + assert event["Records"][0]["s3"]["object"]["key"] == raw_key + + # 4. RAW / no-URL-encoding, no-URL-decoding on the emitted key. + emitted = event["Records"][0]["s3"]["object"]["key"] + assert emitted == raw_key # reprocess applies NO transformation + # space NOT percent-encoded (a real S3 notification would send %20): + assert " " in emitted and "%20" not in emitted + # literal '+' preserved, not turned into a space: + assert "+" in emitted + # '%41' NOT decoded to 'A': + assert "%41" in emitted and "A .eml" not in emitted + + +# --- Selector mutual-exclusion: --all is a MODE, not a modifier --------------- +# reprocess.main() reaches its first AWS call (sts get_caller_identity) only +# AFTER the argparse mutual-exclusion gate, so an invalid selector combination +# raises SystemExit(2) with no AWS access needed. This pins that `--all` cannot +# be silently combined with a narrower --prefix/--since (which the args.all +# clobber branch would otherwise discard -> unintended whole-corpus sweep). +def _expect_argparse_reject(monkeypatch, capsys, argv): + monkeypatch.setattr(sys, "argv", ["reprocess.py", *argv]) + with pytest.raises(SystemExit) as exc: + reprocess.main() + assert exc.value.code == 2 # argparse error exit code + return capsys.readouterr().err + + +def test_all_rejects_prefix(monkeypatch, capsys): + err = _expect_argparse_reject( + monkeypatch, + capsys, + ["--pipeline", "po", "--all", "--prefix", "inbound/2026/07/"], + ) + assert "--all cannot be combined with --prefix or --since" in err + + +def test_all_rejects_since(monkeypatch, capsys): + err = _expect_argparse_reject( + monkeypatch, + capsys, + ["--pipeline", "po", "--all", "--since", "2026-07-01T00:00:00Z"], + ) + assert "--all cannot be combined with --prefix or --since" in err + + +def test_all_rejects_prefix_and_since_together(monkeypatch, capsys): + _expect_argparse_reject( + monkeypatch, + capsys, + [ + "--pipeline", + "po", + "--all", + "--prefix", + "inbound/2026/", + "--since", + "2026-07-01T00:00:00Z", + ], + )