mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 13:03:14 +00:00
CDK stack with SES receipt rule, S3 bucket, email processor Lambda (Claude-powered extraction), web UI Lambda with Function URL, and DynamoDB for storage. Includes reprocessing script for missed emails. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
80 lines
2.2 KiB
Python
80 lines
2.2 KiB
Python
"""
|
|
Re-invoke the po-email-processor Lambda for every object still
|
|
sitting under the inbound/ prefix in the email bucket.
|
|
|
|
Usage:
|
|
python scripts/reprocess.py # dry-run (list only)
|
|
python scripts/reprocess.py --execute # actually invoke
|
|
"""
|
|
|
|
import argparse
|
|
import json
|
|
|
|
import boto3
|
|
|
|
sts = boto3.client("sts")
|
|
s3 = boto3.client("s3")
|
|
lambda_client = boto3.client("lambda")
|
|
|
|
FUNCTION_NAME = "po-email-processor"
|
|
|
|
|
|
def get_bucket_name() -> str:
|
|
account_id = sts.get_caller_identity()["Account"]
|
|
return f"po-ingest-emails-{account_id}"
|
|
|
|
|
|
def list_inbound_keys(bucket: str) -> list[str]:
|
|
keys = []
|
|
paginator = s3.get_paginator("list_objects_v2")
|
|
for page in paginator.paginate(Bucket=bucket, Prefix="inbound/"):
|
|
for obj in page.get("Contents", []):
|
|
keys.append(obj["Key"])
|
|
return keys
|
|
|
|
|
|
def build_s3_event(bucket: str, key: str) -> dict:
|
|
return {
|
|
"Records": [
|
|
{
|
|
"s3": {
|
|
"bucket": {"name": bucket},
|
|
"object": {"key": key},
|
|
}
|
|
}
|
|
]
|
|
}
|
|
|
|
|
|
def main():
|
|
parser = argparse.ArgumentParser(description="Reprocess missed PO emails")
|
|
parser.add_argument("--execute", action="store_true", help="Actually invoke the Lambda (default is dry-run)")
|
|
args = parser.parse_args()
|
|
|
|
bucket = get_bucket_name()
|
|
keys = list_inbound_keys(bucket)
|
|
|
|
if not keys:
|
|
print("No objects found under inbound/ — nothing to reprocess.")
|
|
return
|
|
|
|
print(f"Found {len(keys)} email(s) in s3://{bucket}/inbound/\n")
|
|
|
|
for key in keys:
|
|
if args.execute:
|
|
print(f" Invoking 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(f"\nDry run complete. Re-run with --execute to invoke the Lambda.")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|