270 lines
8.4 KiB
Python
270 lines
8.4 KiB
Python
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"})
|