chainshift-dashboard / core /supabase_action_items.py
GitHub Action
Sync from GitHub
ef78361
Raw
History Blame Contribute Delete
6.5 kB
"""Supabase action items & reports queries.
Module-level functions extracted from supabase_client.py.
All functions use get_supabase_client() from the parent module.
"""
import streamlit as st
def get_action_items(
campaign_id: int,
status: str | None = None,
category: str | None = None,
page: int = 1,
page_size: int = 50,
order_by: str = "created_at",
desc: bool = True,
) -> tuple[list[dict], int]:
"""Get action items for a campaign with filters and sorting.
Args:
order_by: Column to sort by (created_at, priority, category, status)
desc: True for descending, False for ascending
Returns (items, total_count).
"""
from core.supabase_client import get_supabase_client
try:
client = get_supabase_client()
query = (
client.table("action_items")
.select(
"id, campaign_id, trigger_rule_id, category, priority, label, "
"evidence, llm_recommendation, status, assignee_email, "
"created_at, completed_at",
count="planned",
)
.eq("campaign_id", campaign_id)
)
if status:
query = query.eq("status", status)
if category:
query = query.eq("category", category)
offset = (page - 1) * page_size
query = (
query.order(order_by, desc=desc)
.range(offset, offset + page_size - 1)
)
result = query.execute()
return result.data or [], result.count or 0
except Exception:
return [], 0
@st.cache_data(ttl=60)
def get_action_item_stats(campaign_id: int) -> dict:
"""Get action item stats via RPC (COUNT FILTER pattern)."""
_empty = {"pending": 0, "in_progress": 0, "completed": 0, "archived": 0, "total": 0}
try:
from core.supabase_client import get_supabase_client
client = get_supabase_client()
result = client.rpc(
"get_action_item_stats", {"p_campaign_id": campaign_id}
).execute()
data = result.data
if isinstance(data, list) and data:
data = data[0]
if isinstance(data, str):
import json
data = json.loads(data)
return data if isinstance(data, dict) else _empty
except Exception:
return _empty
def update_action_item_status(item_id: str, new_status: str) -> bool:
"""Update action item status. Returns True on success."""
try:
from core.supabase_client import get_supabase_client
client = get_supabase_client()
result = (
client.table("action_items")
.update({"status": new_status})
.eq("id", item_id)
.execute()
)
return bool(result.data)
except Exception:
return False
def delete_action_item(item_id: str) -> bool:
"""Delete action item. Returns True on success."""
try:
from core.supabase_client import get_supabase_client
client = get_supabase_client()
result = (
client.table("action_items")
.delete()
.eq("id", item_id)
.execute()
)
return bool(result.data)
except Exception:
return False
def save_action_items_batch(
campaign_id: int,
user_id: str,
items: list[dict],
report_id: str | None = None,
) -> dict:
"""Save action items with dedup (check existing active items).
Args:
campaign_id: Campaign ID
user_id: User UUID
items: List of dicts with trigger_rule_id, category, priority, label, evidence
report_id: Optional report UUID
Returns:
{"created": N, "skipped": N}
"""
try:
from core.supabase_client import get_supabase_client
client = get_supabase_client()
# Batch dedup check: single query instead of N queries
trigger_ids = [item["trigger_rule_id"] for item in items]
existing_result = (
client.table("action_items")
.select("trigger_rule_id")
.eq("campaign_id", campaign_id)
.in_("trigger_rule_id", trigger_ids)
.in_("status", ["pending", "in_progress"])
.execute()
)
existing_set = {r["trigger_rule_id"] for r in (existing_result.data or [])}
rows_to_insert = []
skipped = 0
for item in items:
if item["trigger_rule_id"] in existing_set:
skipped += 1
continue
rows_to_insert.append({
"campaign_id": campaign_id,
"user_id": user_id,
"report_id": report_id,
"trigger_rule_id": item["trigger_rule_id"],
"category": item["category"],
"priority": item["priority"],
"label": item["label"],
"evidence": item.get("evidence"),
"llm_recommendation": item.get("llm_recommendation"),
"status": "pending",
})
created = 0
if rows_to_insert:
client.table("action_items").insert(rows_to_insert).execute()
created = len(rows_to_insert)
return {"created": created, "skipped": skipped}
except Exception as e:
return {"created": 0, "skipped": 0, "error": str(e)}
def get_campaign_date_range(campaign_id: int) -> tuple[str, str] | None:
"""Get first and last data dates for a campaign via RPC (single query).
Returns:
Tuple of (first_date, last_date) as strings, or None if no data
"""
from core.supabase_client import get_supabase_client
client = get_supabase_client()
result = client.rpc("get_campaign_date_range_agg", {"p_campaign_id": campaign_id}).execute()
data = result.data
# PostgREST may wrap json return as [dict] or dict
if isinstance(data, list) and data:
data = data[0]
if isinstance(data, dict) and data.get("first_date") and data.get("last_date"):
return (data["first_date"], data["last_date"])
return None
def get_report_history_count(campaign_id: int) -> int:
"""Get total count of generated HTML reports for a campaign."""
from core.supabase_client import get_supabase_client
client = get_supabase_client()
result = (
client.table("html_reports")
.select("id", count="planned")
.eq("campaign_id", campaign_id)
.execute()
)
return result.count or 0