import os import time from datetime import datetime, timezone import boto3 import ulid from botocore.exceptions import ClientError _table = None def _get_table(): global _table if _table is None: _table = boto3.resource("dynamodb").Table(os.environ["TABLE_NAME"]) return _table def _ttl_90_days(): return int(time.time()) + (90 * 24 * 60 * 60) def message_exists(message_id): resp = _get_table().get_item( Key={"pk": f"MSG#{message_id}", "sk": "MSG"}, ProjectionExpression="pk", ) return "Item" in resp def save_message(msg_data): item = { "pk": f"MSG#{msg_data['message_id']}", "sk": "MSG", "thread_id": msg_data["thread_id"], "from_email": msg_data["from_email"], "from_name": msg_data["from_name"], "to_emails": msg_data["to_emails"], "cc_emails": msg_data["cc_emails"], "subject": msg_data["subject"], "snippet": msg_data["snippet"], "internal_date": msg_data["internal_date"], "classification": msg_data["classification"], "classification_reason": msg_data["classification_reason"], "bypassed_wo": msg_data["bypassed_wo"], "notified": msg_data.get("notified", False), "classified_date": msg_data["classified_date"], "slack_ts": msg_data.get("slack_ts", ""), "ttl": _ttl_90_days(), } try: _get_table().put_item( Item=item, ConditionExpression="attribute_not_exists(pk)", ) except ClientError as e: if e.response["Error"]["Code"] == "ConditionalCheckFailedException": return raise def save_thread_state(thread_id, message, adam_email): now_iso = datetime.now(timezone.utc).isoformat() is_adam = message["from_email"] == adam_email if is_adam: _get_table().update_item( Key={"pk": f"THD#{thread_id}", "sk": "THD"}, UpdateExpression=( "SET last_message_from = :from_email, " "last_message_date = :now, " "subject = :subject, " "adam_last_reply_date = :now, " "ttl = :ttl " "REMOVE unanswered_since" ), ExpressionAttributeValues={ ":from_email": message["from_email"], ":now": now_iso, ":subject": message["subject"], ":ttl": _ttl_90_days(), }, ) else: _get_table().update_item( Key={"pk": f"THD#{thread_id}", "sk": "THD"}, UpdateExpression=( "SET last_message_from = :from_email, " "last_message_date = :now, " "subject = :subject, " "unanswered_since = if_not_exists(unanswered_since, :now), " "ttl = :ttl" ), ExpressionAttributeValues={ ":from_email": message["from_email"], ":now": now_iso, ":subject": message["subject"], ":ttl": _ttl_90_days(), }, ) def get_sync_metadata(): resp = _get_table().get_item( Key={"pk": "META#sync", "sk": "META"} ) return resp.get("Item") def save_sync_metadata(history_id): _get_table().put_item( Item={ "pk": "META#sync", "sk": "META", "history_id": str(history_id), "last_sync_at": datetime.now(timezone.utc).isoformat(), "initial_sync_done": True, } ) def get_todays_messages(date_str): items = [] resp = _get_table().query( IndexName="by-date", KeyConditionExpression="classified_date = :d AND sk = :sk", ExpressionAttributeValues={":d": date_str, ":sk": "MSG"}, ) items.extend(resp.get("Items", [])) while "LastEvaluatedKey" in resp: resp = _get_table().query( IndexName="by-date", KeyConditionExpression="classified_date = :d AND sk = :sk", ExpressionAttributeValues={":d": date_str, ":sk": "MSG"}, ExclusiveStartKey=resp["LastEvaluatedKey"], ) items.extend(resp.get("Items", [])) return items def get_conversation(thread_ts): resp = _get_table().get_item( Key={"pk": f"CONV#{thread_ts}", "sk": "CONV"} ) item = resp.get("Item") return item["messages"] if item else [] def save_conversation(thread_ts, messages): _get_table().put_item(Item={ "pk": f"CONV#{thread_ts}", "sk": "CONV", "messages": messages, "updated_at": datetime.now(timezone.utc).isoformat(), "ttl": _ttl_90_days(), }) def save_coordination(channel, coord_thread_ts, conversation_thread_ts, started_by): _get_table().put_item(Item={ "pk": f"COORD#{channel}#{coord_thread_ts}", "sk": "COORD", "conversation_thread_ts": conversation_thread_ts, "channel": channel, "started_by": started_by, "status": "active", "created_at": datetime.now(timezone.utc).isoformat(), "ttl": int(time.time()) + (7 * 24 * 60 * 60), }) def get_coordination(channel, coord_thread_ts): resp = _get_table().get_item( Key={"pk": f"COORD#{channel}#{coord_thread_ts}", "sk": "COORD"} ) item = resp.get("Item") if item and item.get("status") == "active": return item return None def complete_coordination(channel, coord_thread_ts): _get_table().update_item( Key={"pk": f"COORD#{channel}#{coord_thread_ts}", "sk": "COORD"}, UpdateExpression="SET #s = :done", ExpressionAttributeNames={"#s": "status"}, ExpressionAttributeValues={":done": "completed"}, ) def get_unanswered_threads(threshold_hours): cutoff = datetime.now(timezone.utc).timestamp() - (threshold_hours * 3600) items = [] scan_kwargs = { "FilterExpression": "begins_with(pk, :prefix) AND attribute_exists(unanswered_since)", "ExpressionAttributeValues": {":prefix": "THD#"}, } resp = _get_table().scan(**scan_kwargs) items.extend(resp.get("Items", [])) while "LastEvaluatedKey" in resp: scan_kwargs["ExclusiveStartKey"] = resp["LastEvaluatedKey"] resp = _get_table().scan(**scan_kwargs) items.extend(resp.get("Items", [])) return [ item for item in items if datetime.fromisoformat(item["unanswered_since"]).timestamp() < cutoff ] # ── Tasks ────────────────────────────────────────────────────── def create_task(title, due_date=None): task_id = str(ulid.new()) _get_table().put_item(Item={ "pk": f"TASK#{task_id}", "sk": "TASK", "title": title, "due_date": due_date or "", "status": "open", "created_at": datetime.now(timezone.utc).isoformat(), "completed_at": "", "ttl": _ttl_90_days(), }) return task_id def get_tasks(status_filter=None): scan_kwargs = { "FilterExpression": "begins_with(pk, :prefix)", "ExpressionAttributeValues": {":prefix": "TASK#"}, } if status_filter: scan_kwargs["FilterExpression"] += " AND #s = :status" scan_kwargs["ExpressionAttributeNames"] = {"#s": "status"} scan_kwargs["ExpressionAttributeValues"][":status"] = status_filter items = [] resp = _get_table().scan(**scan_kwargs) items.extend(resp.get("Items", [])) while "LastEvaluatedKey" in resp: scan_kwargs["ExclusiveStartKey"] = resp["LastEvaluatedKey"] resp = _get_table().scan(**scan_kwargs) items.extend(resp.get("Items", [])) items.sort(key=lambda x: x["pk"]) return items def update_task_status(task_id, status): update_expr = "SET #s = :status, ttl = :ttl" expr_values = {":status": status, ":ttl": _ttl_90_days()} if status == "completed": update_expr += ", completed_at = :now" expr_values[":now"] = datetime.now(timezone.utc).isoformat() _get_table().update_item( Key={"pk": f"TASK#{task_id}", "sk": "TASK"}, UpdateExpression=update_expr, ExpressionAttributeNames={"#s": "status"}, ExpressionAttributeValues=expr_values, ) def delete_task(task_id): _get_table().delete_item(Key={"pk": f"TASK#{task_id}", "sk": "TASK"})