Async database service using aiosqlite.
| 22 | |
| 23 | |
| 24 | class DatabaseService: |
| 25 | """Async database service using aiosqlite.""" |
| 26 | |
| 27 | def __init__(self, base_dir: Union[str, Path]): |
| 28 | """Initialize the database service. |
| 29 | |
| 30 | Args: |
| 31 | base_dir: Base directory for the database |
| 32 | """ |
| 33 | self.base_dir = Path(base_dir) if isinstance(base_dir, str) else base_dir |
| 34 | self._connection: Optional[aiosqlite.Connection] = None |
| 35 | self._is_connected = False |
| 36 | |
| 37 | async def connect(self) -> None: |
| 38 | """Connect to the database.""" |
| 39 | try: |
| 40 | # Ensure the directory exists |
| 41 | self.base_dir.mkdir(parents=True, exist_ok=True) |
| 42 | |
| 43 | self._connection = await aiosqlite.connect(str(self.base_dir / "database.db")) |
| 44 | self._is_connected = True |
| 45 | except Exception as e: |
| 46 | raise ConnectionError(f"Failed to connect to database: {e}") |
| 47 | |
| 48 | async def disconnect(self) -> None: |
| 49 | """Disconnect from the database.""" |
| 50 | if self._connection: |
| 51 | await self._connection.close() |
| 52 | self._connection = None |
| 53 | self._is_connected = False |
| 54 | |
| 55 | async def execute_query(self, request: QueryRequest) -> ActionResult: |
| 56 | """Execute a SQL query. |
| 57 | |
| 58 | Args: |
| 59 | request: Query request with SQL and parameters |
| 60 | |
| 61 | Returns: |
| 62 | Action result with data and metadata in extra |
| 63 | """ |
| 64 | if not self._is_connected: |
| 65 | return ActionResult( |
| 66 | success=False, |
| 67 | message="Database not connected", |
| 68 | extra={"error": "Database not connected"} |
| 69 | ) |
| 70 | |
| 71 | start_time = time.time() |
| 72 | |
| 73 | try: |
| 74 | # Handle both named parameters (dict) and positional parameters (tuple/list) |
| 75 | params = request.parameters if request.parameters is not None else {} |
| 76 | cursor = await self._connection.execute(request.query, params) |
| 77 | |
| 78 | # Check if it's a SELECT query |
| 79 | if request.query.strip().upper().startswith('SELECT'): |
| 80 | rows = await cursor.fetchall() |
| 81 | columns = [description[0] for description in cursor.description] if cursor.description else [] |
no outgoing calls