File size: 8,992 Bytes
54eb2ce | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 | from sqlalchemy import select
from sqlalchemy.exc import SQLAlchemyError, DBAPIError
from fastapi import HTTPException
from datetime import date
from ..database.db import session_pool
from ..errors import DatabaseConnectionError
from ..model.chat_session import Session
from ..schema.session import SessionCreate, SessionUpdate, SessionResponse
from ..core.logger import SingletonLogger
logger = SingletonLogger().get_logger()
async def create_session(user_id: int, payload: SessionCreate) -> SessionResponse:
"""Create a new session for a user."""
try:
async with session_pool() as session:
new_session = Session(
user_id=user_id,
paper_id=payload.paper_id,
title=payload.title,
started_at=payload.started_at,
ended_at=payload.ended_at,
device_type=payload.device_type,
)
session.add(new_session)
await session.commit()
await session.refresh(new_session)
return SessionResponse(
id=new_session.id,
user_id=new_session.user_id,
title=new_session.title,
started_at=new_session.started_at,
ended_at=new_session.ended_at,
device_type=new_session.device_type,
paper_id=new_session.paper_id,
)
except DBAPIError as e:
logger.exception(
f"Database connection error creating session for user_id={user_id}: {str(e)}"
)
raise DatabaseConnectionError(str(e))
except SQLAlchemyError as e:
logger.error(
f"Database error creating session for user_id={user_id}: {str(e)}"
)
raise HTTPException(status_code=500, detail="Failed to create session")
except Exception as e:
logger.error(
f"Unexpected error creating session for user_id={user_id}: {str(e)}"
)
raise HTTPException(status_code=500, detail="Internal server error")
async def get_session(session_id: int, user_id: int) -> SessionResponse:
"""Retrieve a session by ID, ensuring it belongs to the user."""
try:
async with session_pool() as session:
result = await session.execute(
select(Session).where(
Session.id == session_id, Session.user_id == user_id
)
)
sess = result.scalar_one_or_none()
if not sess:
raise HTTPException(status_code=404, detail="Session not found")
return SessionResponse(
id=sess.id,
user_id=sess.user_id,
paper_id=sess.paper_id,
title=sess.title,
started_at=sess.started_at,
ended_at=sess.ended_at,
device_type=sess.device_type,
)
except HTTPException:
raise
except DBAPIError as e:
logger.exception(
f"Database connection error retrieving session id={session_id}: {str(e)}"
)
raise DatabaseConnectionError(str(e))
except SQLAlchemyError as e:
logger.error(f"Database error retrieving session id={session_id}: {str(e)}")
raise HTTPException(status_code=500, detail="Failed to retrieve session")
except Exception as e:
logger.error(f"Unexpected error retrieving session id={session_id}: {str(e)}")
raise HTTPException(status_code=500, detail="Internal server error")
async def get_sessions_by_user(user_id: int) -> list[SessionResponse]:
"""Retrieve all sessions for a user."""
try:
async with session_pool() as session:
result = await session.execute(
select(Session)
.where(Session.user_id == user_id)
.order_by(Session.started_at.desc())
)
sessions = result.scalars().all()
return [
SessionResponse(
id=sess.id,
user_id=sess.user_id,
paper_id=sess.paper_id,
title=sess.title,
started_at=sess.started_at,
ended_at=sess.ended_at,
device_type=sess.device_type,
)
for sess in sessions
]
except DBAPIError as e:
logger.exception(
f"Database connection error retrieving sessions for user_id={user_id}: {str(e)}"
)
raise DatabaseConnectionError(str(e))
except SQLAlchemyError as e:
logger.error(
f"Database error retrieving sessions for user_id={user_id}: {str(e)}"
)
raise HTTPException(status_code=500, detail="Failed to retrieve sessions")
except Exception as e:
logger.error(
f"Unexpected error retrieving sessions for user_id={user_id}: {str(e)}"
)
raise HTTPException(status_code=500, detail="Internal server error")
async def get_session_by_user_id(user_id: int):
"""Fetch the most recent session for a given user ID."""
try:
async with session_pool() as session:
result = await session.execute(
select(Session)
.where(Session.user_id == user_id)
.order_by(Session.updated_at.desc())
)
sess = result.scalars().first()
return sess
except Exception as e:
logger.error(f"Error fetching session for user_id {user_id}: {e}")
return None
async def update_session(
session_id: int, user_id: int, payload: SessionUpdate
) -> SessionResponse:
"""Update an existing session with new data."""
try:
async with session_pool() as session:
result = await session.execute(
select(Session).where(
Session.id == session_id, Session.user_id == user_id
)
)
sess = result.scalar_one_or_none()
if not sess:
raise HTTPException(status_code=404, detail="Session not found")
# Update fields only if provided
if payload.title is not None:
sess.title = payload.title
if payload.started_at is not None:
sess.started_at = payload.started_at
if payload.ended_at is not None:
sess.ended_at = payload.ended_at
if payload.device_type is not None:
sess.device_type = payload.device_type
await session.commit()
await session.refresh(sess)
return SessionResponse(
id=sess.id,
user_id=sess.user_id,
paper_id=sess.paper_id,
title=sess.title,
started_at=sess.started_at,
ended_at=sess.ended_at,
device_type=sess.device_type,
)
except HTTPException:
raise
except DBAPIError as e:
logger.exception(
f"Database connection error updating session id={session_id}: {str(e)}"
)
raise DatabaseConnectionError(str(e))
except SQLAlchemyError as e:
logger.error(f"Database error updating session id={session_id}: {str(e)}")
raise HTTPException(status_code=500, detail="Failed to update session")
except Exception as e:
logger.error(f"Unexpected error updating session id={session_id}: {str(e)}")
raise HTTPException(status_code=500, detail="Internal server error")
async def delete_session(session_id: int, user_id: int) -> None:
"""Delete a session by its ID."""
try:
async with session_pool() as session:
result = await session.execute(
select(Session).where(
Session.id == session_id, Session.user_id == user_id
)
)
sess = result.scalar_one_or_none()
if not sess:
raise HTTPException(status_code=404, detail="Session not found")
await session.delete(sess)
await session.commit()
except HTTPException:
raise
except DBAPIError as e:
logger.exception(
f"Database connection error deleting session id={session_id}: {str(e)}"
)
raise DatabaseConnectionError(str(e))
except SQLAlchemyError as e:
logger.error(f"Database error deleting session id={session_id}: {str(e)}")
raise HTTPException(status_code=500, detail="Failed to delete session")
except Exception as e:
logger.error(f"Unexpected error deleting session id={session_id}: {str(e)}")
raise HTTPException(status_code=500, detail="Internal server error")
|