Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
263 changes: 263 additions & 0 deletions data-analytics/lambda_functions/kpi_api_analytics.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,263 @@
import json
import os
import boto3
import psycopg2
from psycopg2.extras import RealDictCursor


SCHEMA_NAME = "virginia_dev_saayam_rdbms"


SLA = {
"target_days": 10,
"target_hours": 240,
"warning_days": 8.33,
"warning_hours": 200
}

VALID_TIME_RANGES = {"7D", "30D", "1Y", "All", "Custom"}


def get_default_response():
return {
"request_status_distribution": [],
"total_requests": 0,
"average_resolution_time_by_category": [],
"sla": SLA
}


def build_response(status_code, body):
return {
"statusCode": status_code,
"headers": {
"Content-Type": "application/json",
"Access-Control-Allow-Origin": "*"
},
"body": json.dumps(body, default=str)
}


def get_db_connection():
"""
Connects to the Virginia RDS PostgreSQL DB.

For local testing, set LOCAL_DB_TEST=1 and provide DB_HOST, DB_NAME,
DB_USER, DB_PASSWORD, DB_PORT as environment variables instead of
pulling credentials from SSM (SSM is only reachable from within AWS).
"""
if os.environ.get("LOCAL_DB_TEST") == "1":
return psycopg2.connect(
host=os.environ["DB_HOST"],
database=os.environ["DB_NAME"],
user=os.environ["DB_USER"],
password=os.environ["DB_PASSWORD"],
port=os.environ.get("DB_PORT", "5432"),
sslmode="require"
)

ssm = boto3.client("ssm", region_name="us-east-1")

response = ssm.get_parameter(
Name="/dev/saayam/db/Virginia/Analytics/user",
WithDecryption=True
)

creds = json.loads(response["Parameter"]["Value"])
db_name = creds["DATABASE NAME"]
return psycopg2.connect(
host=creds["HOST"],
database=db_name,
user=creds["USERNAME"],
password=creds["PASSWORD"],
port=creds["PORT"],
sslmode="require"
)


def build_date_filter(time_range, start_date=None, end_date=None, date_column="r.submission_date"):
"""
Builds a parameterized SQL date condition based on time_range.

Returns a tuple: (sql_condition, params)
- sql_condition is a string starting with "AND ..." (or "" for All),
meant to be appended after a "WHERE 1=1" clause.
- params is a list of query parameters to pass alongside the condition.

Raises ValueError if time_range == "Custom" and start_date/end_date
are not both provided.
"""
if not time_range:
time_range = "All"

time_range = time_range.strip()

if time_range not in VALID_TIME_RANGES:
# Unrecognized value: fall back to no filter rather than failing the request.
time_range = "All"

if time_range == "7D":
return f"AND {date_column} >= CURRENT_DATE - INTERVAL '7 days'", []

if time_range == "30D":
return f"AND {date_column} >= CURRENT_DATE - INTERVAL '30 days'", []

if time_range == "1Y":
return f"AND {date_column} >= CURRENT_DATE - INTERVAL '1 year'", []

if time_range == "Custom":
if not start_date or not end_date:
raise ValueError("Custom time_range requires both start_date and end_date")
return f"AND {date_column} BETWEEN %s AND %s", [start_date, end_date]

# "All"
return "", []


def fetch_request_status_distribution(cursor, time_range, start_date=None, end_date=None):
date_condition, params = build_date_filter(time_range, start_date, end_date)

query = f"""
SELECT
rs.req_status AS status,
COUNT(r.req_id) AS count
FROM {SCHEMA_NAME}.request r
JOIN {SCHEMA_NAME}.request_status rs
ON r.req_status_id = rs.req_status_id
WHERE 1=1
{date_condition}
GROUP BY rs.req_status
ORDER BY rs.req_status;
"""

cursor.execute(query, params)
rows = cursor.fetchall()

return [
{
"status": row["status"],
"count": int(row["count"])
}
for row in rows
]


def fetch_total_requests(cursor, time_range, start_date=None, end_date=None):
date_condition, params = build_date_filter(time_range, start_date, end_date)

query = f"""
SELECT COUNT(req_id) AS total_requests
FROM {SCHEMA_NAME}.request r
WHERE 1=1
{date_condition};
"""

cursor.execute(query, params)
row = cursor.fetchone()

return int(row["total_requests"]) if row and row["total_requests"] is not None else 0


def fetch_average_resolution_time_by_category(cursor, time_range, start_date=None, end_date=None):
date_condition, params = build_date_filter(time_range, start_date, end_date)

query = f"""
SELECT
hc.cat_name AS category,
ROUND(
AVG(
EXTRACT(EPOCH FROM (r.serviced_date - r.submission_date)) / 3600
)::numeric,
2
) AS avg_hours
FROM {SCHEMA_NAME}.request r
JOIN {SCHEMA_NAME}.help_categories hc
ON r.req_cat_id = hc.cat_id
JOIN {SCHEMA_NAME}.request_status rs
ON r.req_status_id = rs.req_status_id
WHERE r.submission_date IS NOT NULL
AND r.serviced_date IS NOT NULL
AND r.serviced_date >= r.submission_date
AND UPPER(rs.req_status) IN ('COMPLETED', 'RESOLVED')
{date_condition}
GROUP BY hc.cat_name
ORDER BY avg_hours DESC;
"""

cursor.execute(query, params)
rows = cursor.fetchall()

return [
{
"category": row["category"],
"avg_hours": float(row["avg_hours"]) if row["avg_hours"] is not None else 0
}
for row in rows
]


def lambda_handler(event, context):
conn = None
cursor = None
response_body = get_default_response()

event = event or {}
time_range = event.get("time_range", "All")
start_date = event.get("start_date")
end_date = event.get("end_date")

try:
conn = get_db_connection()
cursor = conn.cursor(cursor_factory=RealDictCursor)

try:
response_body["request_status_distribution"] = fetch_request_status_distribution(
cursor, time_range, start_date, end_date
)
except Exception as error:
print(f"Status distribution query failed: {error}")
response_body["request_status_distribution"] = []

try:
response_body["total_requests"] = fetch_total_requests(
cursor, time_range, start_date, end_date
)
except Exception as error:
print(f"Total request count query failed: {error}")
response_body["total_requests"] = 0

try:
response_body["average_resolution_time_by_category"] = fetch_average_resolution_time_by_category(
cursor, time_range, start_date, end_date
)
except Exception as error:
print(f"Average resolution time query failed: {error}")
response_body["average_resolution_time_by_category"] = []

return build_response(200, response_body)

except Exception as error:
print(f"DB connection failed: {error}")
return build_response(500, response_body)

finally:
if cursor:
cursor.close()
if conn:
conn.close()


if __name__ == "__main__":
test_payloads = [
{},
{"time_range": "7D"},
{"time_range": "30D"},
{"time_range": "1Y"},
{"time_range": "All"},
{"time_range": "Custom", "start_date": "2026-05-01", "end_date": "2026-05-31"},
]

for payload in test_payloads:
print(f"\n--- Testing payload: {payload} ---")
result = lambda_handler(payload, None)
print(result)