|
1 | | -from __future__ import annotations |
2 | | - |
3 | | -import asyncio |
4 | | -import logging |
5 | | -from typing import Any |
6 | | - |
7 | | -from fastapi import FastAPI |
8 | | - |
9 | | -from src.core.interfaces.session_service_interface import ISessionService |
10 | | - |
11 | | -logger = logging.getLogger(__name__) |
12 | | - |
13 | | - |
14 | | -class AppLifecycle: |
15 | | - """Handles application lifecycle events. |
16 | | -
|
17 | | - This class manages startup and shutdown tasks for the application. |
18 | | - """ |
19 | | - |
20 | | - def __init__(self, app: FastAPI, config: dict[str, Any]) -> None: |
21 | | - """Initialize the lifecycle manager. |
22 | | -
|
23 | | - Args: |
24 | | - app: The FastAPI application |
25 | | - config: The application configuration |
26 | | - """ |
27 | | - self.app = app |
28 | | - self.config = config |
29 | | - self._background_tasks: list[asyncio.Task] = [] |
30 | | - |
31 | | - async def startup(self) -> None: |
32 | | - """Perform startup tasks. |
33 | | -
|
34 | | - This method is called during application startup. |
35 | | - """ |
36 | | - logger.info("Starting application lifecycle...") |
37 | | - |
38 | | - # Start background tasks |
39 | | - self._start_background_tasks() |
40 | | - |
41 | | - async def shutdown(self) -> None: |
42 | | - """Perform shutdown tasks. |
43 | | -
|
44 | | - This method is called during application shutdown. |
45 | | - """ |
46 | | - logger.info("Shutting down application lifecycle...") |
47 | | - |
48 | | - # Stop background tasks |
49 | | - await self._stop_background_tasks() |
50 | | - |
51 | | - # Close any remaining connections |
52 | | - await self._close_connections() |
53 | | - |
54 | | - def _start_background_tasks(self) -> None: |
55 | | - """Start background tasks.""" |
56 | | - # Start session cleanup task |
57 | | - if self.config.get("session_cleanup_enabled", False): |
58 | | - interval = self.config.get( |
59 | | - "session_cleanup_interval", 3600 |
60 | | - ) # 1 hour default |
61 | | - max_age = self.config.get("session_max_age", 86400) # 1 day default |
62 | | - |
63 | | - task = asyncio.create_task( |
64 | | - self._session_cleanup_task(interval, max_age), name="session_cleanup" |
65 | | - ) |
66 | | - self._background_tasks.append(task) |
67 | | - logger.info( |
68 | | - f"Started session cleanup task (interval: {interval}s, max age: {max_age}s)" |
69 | | - ) |
70 | | - |
71 | | - async def _stop_background_tasks(self) -> None: |
72 | | - """Stop background tasks.""" |
73 | | - for task in self._background_tasks: |
74 | | - if not task.done(): |
75 | | - task.cancel() |
76 | | - try: |
77 | | - await task |
78 | | - except asyncio.CancelledError: |
79 | | - logger.info(f"Cancelled background task: {task.get_name()}") |
80 | | - |
81 | | - async def _close_connections(self) -> None: |
82 | | - """Close any remaining connections.""" |
83 | | - # Any connection cleanup code would go here |
84 | | - |
85 | | - async def _session_cleanup_task(self, interval: int, max_age: int) -> None: |
86 | | - """Background task for cleaning up expired sessions. |
87 | | -
|
88 | | - Args: |
89 | | - interval: The interval in seconds between cleanup runs |
90 | | - max_age: The maximum age in seconds for sessions |
91 | | - """ |
92 | | - try: |
93 | | - while True: |
94 | | - await asyncio.sleep(interval) |
95 | | - |
96 | | - try: |
97 | | - # Get service provider |
98 | | - provider = self.app.state.service_provider |
99 | | - if not provider: |
100 | | - logger.warning( |
101 | | - "Service provider not available for session cleanup" |
102 | | - ) |
103 | | - continue |
104 | | - |
105 | | - # Get session service |
106 | | - session_service = provider.get_service(ISessionService) |
107 | | - if not session_service: |
108 | | - logger.warning("Session service not available for cleanup") |
109 | | - continue |
110 | | - |
111 | | - # Perform cleanup |
112 | | - deleted_count = 0 |
113 | | - if hasattr(session_service, "cleanup_expired_sessions"): |
114 | | - deleted_count = await session_service.cleanup_expired_sessions( |
115 | | - max_age |
116 | | - ) |
117 | | - |
118 | | - if deleted_count > 0 and logger.isEnabledFor(logging.INFO): |
119 | | - logger.info(f"Cleaned up {deleted_count} expired sessions") |
120 | | - |
121 | | - except Exception as e: |
122 | | - logger.error(f"Error during session cleanup: {e!s}") |
123 | | - |
124 | | - except asyncio.CancelledError: |
125 | | - logger.debug("Session cleanup task cancelled") |
126 | | - raise |
| 1 | +from __future__ import annotations |
| 2 | + |
| 3 | +import asyncio |
| 4 | +import logging |
| 5 | +from contextlib import suppress |
| 6 | +from typing import Any |
| 7 | + |
| 8 | +from fastapi import FastAPI |
| 9 | + |
| 10 | +from src.core.interfaces.session_service_interface import ISessionService |
| 11 | +from src.core.interfaces.wire_capture_interface import IWireCapture |
| 12 | + |
| 13 | +logger = logging.getLogger(__name__) |
| 14 | + |
| 15 | + |
| 16 | +class AppLifecycle: |
| 17 | + """Handles application lifecycle events. |
| 18 | +
|
| 19 | + This class manages startup and shutdown tasks for the application. |
| 20 | + """ |
| 21 | + |
| 22 | + def __init__(self, app: FastAPI, config: dict[str, Any]) -> None: |
| 23 | + """Initialize the lifecycle manager. |
| 24 | +
|
| 25 | + Args: |
| 26 | + app: The FastAPI application |
| 27 | + config: The application configuration |
| 28 | + """ |
| 29 | + self.app = app |
| 30 | + self.config = config |
| 31 | + self._background_tasks: list[asyncio.Task] = [] |
| 32 | + |
| 33 | + async def startup(self) -> None: |
| 34 | + """Perform startup tasks. |
| 35 | +
|
| 36 | + This method is called during application startup. |
| 37 | + """ |
| 38 | + logger.info("Starting application lifecycle...") |
| 39 | + |
| 40 | + # Start background tasks |
| 41 | + self._start_background_tasks() |
| 42 | + |
| 43 | + async def shutdown(self) -> None: |
| 44 | + """Perform shutdown tasks. |
| 45 | +
|
| 46 | + This method is called during application shutdown. |
| 47 | + """ |
| 48 | + logger.info("Shutting down application lifecycle...") |
| 49 | + |
| 50 | + # Stop background tasks |
| 51 | + await self._stop_background_tasks() |
| 52 | + |
| 53 | + # Close any remaining connections |
| 54 | + await self._close_connections() |
| 55 | + |
| 56 | + def _start_background_tasks(self) -> None: |
| 57 | + """Start background tasks.""" |
| 58 | + # Start session cleanup task |
| 59 | + if self.config.get("session_cleanup_enabled", False): |
| 60 | + interval = self.config.get( |
| 61 | + "session_cleanup_interval", 3600 |
| 62 | + ) # 1 hour default |
| 63 | + max_age = self.config.get("session_max_age", 86400) # 1 day default |
| 64 | + |
| 65 | + task = asyncio.create_task( |
| 66 | + self._session_cleanup_task(interval, max_age), name="session_cleanup" |
| 67 | + ) |
| 68 | + self._background_tasks.append(task) |
| 69 | + logger.info( |
| 70 | + f"Started session cleanup task (interval: {interval}s, max age: {max_age}s)" |
| 71 | + ) |
| 72 | + |
| 73 | + async def _stop_background_tasks(self) -> None: |
| 74 | + """Stop background tasks.""" |
| 75 | + for task in self._background_tasks: |
| 76 | + if not task.done(): |
| 77 | + task.cancel() |
| 78 | + try: |
| 79 | + await task |
| 80 | + except asyncio.CancelledError: |
| 81 | + logger.info(f"Cancelled background task: {task.get_name()}") |
| 82 | + |
| 83 | + async def _close_connections(self) -> None: |
| 84 | + """Close any remaining connections.""" |
| 85 | + # Get service provider |
| 86 | + provider = getattr(self.app.state, "service_provider", None) |
| 87 | + if not provider: |
| 88 | + return |
| 89 | + |
| 90 | + # Get wire capture service and shut it down |
| 91 | + wire_capture_service = provider.get_service(IWireCapture) |
| 92 | + if wire_capture_service and hasattr(wire_capture_service, "shutdown"): |
| 93 | + await wire_capture_service.shutdown() |
| 94 | + |
| 95 | + async def _session_cleanup_task(self, interval: int, max_age: int) -> None: |
| 96 | + """Background task for cleaning up expired sessions. |
| 97 | +
|
| 98 | + Args: |
| 99 | + interval: The interval in seconds between cleanup runs |
| 100 | + max_age: The maximum age in seconds for sessions |
| 101 | + """ |
| 102 | + try: |
| 103 | + while True: |
| 104 | + await asyncio.sleep(interval) |
| 105 | + |
| 106 | + try: |
| 107 | + # Get service provider |
| 108 | + provider = self.app.state.service_provider |
| 109 | + if not provider: |
| 110 | + logger.warning( |
| 111 | + "Service provider not available for session cleanup" |
| 112 | + ) |
| 113 | + continue |
| 114 | + |
| 115 | + # Get session service |
| 116 | + session_service = provider.get_service(ISessionService) |
| 117 | + if not session_service: |
| 118 | + logger.warning("Session service not available for cleanup") |
| 119 | + continue |
| 120 | + |
| 121 | + # Perform cleanup |
| 122 | + deleted_count = 0 |
| 123 | + with suppress(AttributeError): |
| 124 | + deleted_count = await session_service.cleanup_expired(max_age) |
| 125 | + |
| 126 | + if deleted_count > 0 and logger.isEnabledFor(logging.INFO): |
| 127 | + logger.info(f"Cleaned up {deleted_count} expired sessions") |
| 128 | + |
| 129 | + except Exception as e: |
| 130 | + logger.error(f"Error during session cleanup: {e!s}") |
| 131 | + |
| 132 | + except asyncio.CancelledError: |
| 133 | + logger.debug("Session cleanup task cancelled") |
| 134 | + raise |
0 commit comments