kpai-analyst / utils.py
Ashad001's picture
fsdfjds
701cf99
Raw
History Blame Contribute Delete
5.14 kB
import os
import uuid
import duckdb
import pandas as pd
from dotenv import load_dotenv
from datetime import datetime, timedelta, timezone
load_dotenv()
est_offset = timedelta(hours=-4)
# Connect to the DuckDB database
conn = duckdb.connect(f"md:?motherduck_token={os.getenv('MOTHERDUCK_TOKEN')}")
# Utility functions for interacting with the database
def create_session(user_id):
session_id = str(uuid.uuid4())
query = """
INSERT INTO SESSIONS (session_id, user_id, session_start)
VALUES (?, ?, ?)
"""
conn.execute(query, (session_id, user_id, datetime.now(timezone(est_offset))))
return session_id
def end_session(session_id):
query = """
UPDATE SESSIONS
SET session_end = ?
WHERE session_id = ?
"""
conn.execute(query, (datetime.now(timezone(est_offset)), session_id))
def get_user_sessions(user_id):
query = """
SELECT * FROM SESSIONS
WHERE user_id = ?
"""
return conn.execute(query, (user_id,)).fetchdf()
def add_response(session_id, user_id, user_input_text, response_text):
response_id = str(uuid.uuid4())
query = """
INSERT INTO RESPONSES (response_id, session_id, user_id, user_input_text, response_text, created_at)
VALUES (?, ?, ?, ?, ?, ?)
"""
conn.execute(query, (response_id, session_id, user_id, user_input_text, response_text, datetime.now(timezone(est_offset))))
return response_id
def add_agent_response(response_id, agent_name, user_input, agent_response_text):
agent_response_id = str(uuid.uuid4())
query = """
INSERT INTO AGENT_RESPONSES (agent_response_id, response_id, agent_name, user_input, agent_response_text, created_at)
VALUES (?, ?, ?, ?, ?, ?)
"""
conn.execute(query, (agent_response_id, response_id, agent_name, user_input, agent_response_text, datetime.now(timezone(est_offset))))
return agent_response_id
def get_response(response_id):
query = """
SELECT * FROM RESPONSES
WHERE response_id = ?
"""
return conn.execute(query, (response_id,)).fetchone()
def get_responses_from_session(session_id):
query = """
SELECT * FROM RESPONSES
WHERE session_id = ?
"""
return conn.execute(query, (session_id,)).fetchdf()
def get_agent_responses(response_id):
query = """
SELECT * FROM AGENT_RESPONSES
WHERE response_id = ?
"""
return conn.execute(query, (response_id,)).fetchdf()
def add_feedback(response_id, user_id, feedback_score, feedback_text):
feedback_id = str(uuid.uuid4())
query = """
INSERT INTO FEEDBACK (feedback_id, response_id, user_id, feedback_score, feedback_text)
VALUES (?, ?, ?, ?, ?)
"""
conn.execute(query, (feedback_id, response_id, user_id, feedback_score, feedback_text))
return feedback_id
def get_feedback(response_id):
query = """
SELECT * FROM FEEDBACK
WHERE response_id = ?
"""
return conn.execute(query, (response_id,)).fetchdf()
def get_user_feedback(user_id):
query = """
SELECT * FROM FEEDBACK
WHERE user_id = ?
"""
return conn.execute(query, (user_id,)).fetchdf()
def get_session_data(session_id):
# Query to get the combined data from responses, agent_responses, and feedback
query = """
SELECT
r.user_input_text AS query,
ar.agent_name AS agent,
ar.agent_response_text AS response,
f.feedback_score AS feedback_score,
f.feedback_text AS feedback_text,
r.created_at AS response_time,
s.session_start AS session_start,
s.session_end AS session_end
FROM
RESPONSES r
JOIN
AGENT_RESPONSES ar ON r.response_id = ar.response_id
LEFT JOIN
FEEDBACK f ON r.response_id = f.response_id
JOIN
SESSIONS s ON r.session_id = s.session_id
WHERE
r.session_id = ?
ORDER BY
r.created_at, ar.agent_name;
"""
# Execute the query and fetch results into a DataFrame
session_df = conn.execute(query, (session_id,)).fetchdf()
if session_df.empty:
return pd.DataFrame()
# Format timestamps
session_df['response_time'] = session_df['response_time'].apply(lambda x: x.strftime('%Y-%m-%d %H:%M:%S'))
session_df['session_start'] = session_df['session_start'].apply(lambda x: x.strftime('%Y-%m-%d %H:%M:%S'))
session_df['session_end'] = session_df['session_end'].apply(lambda x: x.strftime('%Y-%m-%d %H:%M:%S') if x else "Ongoing")
# Combine feedback into a single column if available
session_df['feedback'] = session_df.apply(
lambda row: f"{row['feedback_score']} - {row['feedback_text']}" if row['feedback_score'] and row['feedback_text'] else None,
axis=1
)
# Drop the separate feedback columns
session_df = session_df.drop(columns=['feedback_score', 'feedback_text'])
# Reorder columns for clarity
session_df = session_df[['query', 'agent', 'response', 'feedback', 'response_time', 'session_start']]
return session_df
if __name__ == "__main__":
id = "adec6d43-f4bb-46ab-887c-3a1d492f727f"
df = get_session_data(id)
df.to_csv('session_data.csv', index=False)