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)