-
Notifications
You must be signed in to change notification settings - Fork 9.9k
Expand file tree
/
Copy pathdeps.py
More file actions
268 lines (180 loc) · 8.57 KB
/
Copy pathdeps.py
File metadata and controls
268 lines (180 loc) · 8.57 KB
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
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
from __future__ import annotations
from contextlib import asynccontextmanager
from http import HTTPStatus
from typing import TYPE_CHECKING
from fastapi import HTTPException
from lfx.log.logger import logger
from langflow.services.schema import ServiceType
if TYPE_CHECKING:
from collections.abc import AsyncGenerator
from lfx.services.settings.service import SettingsService
from sqlmodel.ext.asyncio.session import AsyncSession
from langflow.services.cache.service import AsyncBaseCacheService, CacheService
from langflow.services.chat.service import ChatService
from langflow.services.database.service import DatabaseService
from langflow.services.job_queue.service import JobQueueService
from langflow.services.session.service import SessionService
from langflow.services.state.service import StateService
from langflow.services.storage.service import StorageService
from langflow.services.store.service import StoreService
from langflow.services.task.service import TaskService
from langflow.services.telemetry.service import TelemetryService
from langflow.services.tracing.service import TracingService
from langflow.services.variable.service import VariableService
def get_service(service_type: ServiceType, default=None):
"""Retrieves the service instance for the given service type.
Args:
service_type (ServiceType): The type of service to retrieve.
default (ServiceFactory, optional): The default ServiceFactory to use if the service is not found.
Defaults to None.
Returns:
Any: The service instance.
"""
from lfx.services.manager import get_service_manager
service_manager = get_service_manager()
if not service_manager.are_factories_registered():
# ! This is a workaround to ensure that the service manager is initialized
# ! Not optimal, but it works for now
from langflow.services.manager import ServiceManager
service_manager.register_factories(ServiceManager.get_factories())
return service_manager.get(service_type, default)
def get_telemetry_service() -> TelemetryService:
"""Retrieves the TelemetryService instance from the service manager.
Returns:
TelemetryService: The TelemetryService instance.
"""
from langflow.services.telemetry.factory import TelemetryServiceFactory
return get_service(ServiceType.TELEMETRY_SERVICE, TelemetryServiceFactory())
def get_tracing_service() -> TracingService:
"""Retrieves the TracingService instance from the service manager.
Returns:
TracingService: The TracingService instance.
"""
from langflow.services.tracing.factory import TracingServiceFactory
return get_service(ServiceType.TRACING_SERVICE, TracingServiceFactory())
def get_state_service() -> StateService:
"""Retrieves the StateService instance from the service manager.
Returns:
The StateService instance.
"""
from langflow.services.state.factory import StateServiceFactory
return get_service(ServiceType.STATE_SERVICE, StateServiceFactory())
def get_storage_service() -> StorageService:
"""Retrieves the storage service instance.
Returns:
The storage service instance.
"""
from langflow.services.storage.factory import StorageServiceFactory
return get_service(ServiceType.STORAGE_SERVICE, default=StorageServiceFactory())
def get_variable_service() -> VariableService:
"""Retrieves the VariableService instance from the service manager.
Returns:
The VariableService instance.
"""
from langflow.services.variable.factory import VariableServiceFactory
return get_service(ServiceType.VARIABLE_SERVICE, VariableServiceFactory())
def is_settings_service_initialized() -> bool:
"""Check if the SettingsService is already initialized without triggering initialization.
Returns:
bool: True if the SettingsService is already initialized, False otherwise.
"""
from lfx.services.manager import get_service_manager
return ServiceType.SETTINGS_SERVICE in get_service_manager().services
def get_settings_service() -> SettingsService:
"""Retrieves the SettingsService instance.
If the service is not yet initialized, it will be initialized before returning.
Returns:
The SettingsService instance.
Raises:
ValueError: If the service cannot be retrieved or initialized.
"""
from lfx.services.settings.factory import SettingsServiceFactory
return get_service(ServiceType.SETTINGS_SERVICE, SettingsServiceFactory())
def get_db_service() -> DatabaseService:
"""Retrieves the DatabaseService instance from the service manager.
Returns:
The DatabaseService instance.
"""
from langflow.services.database.factory import DatabaseServiceFactory
return get_service(ServiceType.DATABASE_SERVICE, DatabaseServiceFactory())
async def get_session() -> AsyncGenerator[AsyncSession, None]:
"""Retrieves an async session from the database service.
Yields:
AsyncSession: An async session object.
"""
async with session_scope() as session:
yield session
@asynccontextmanager
async def session_scope() -> AsyncGenerator[AsyncSession, None]:
"""Context manager for managing an async session scope.
This context manager is used to manage an async session scope for database operations.
It ensures that the session is properly committed if no exceptions occur,
and rolled back if an exception is raised.
Yields:
AsyncSession: The async session object.
Raises:
Exception: If an error occurs during the session scope.
"""
db_service = get_db_service()
async with db_service.with_session() as session:
try:
yield session
await session.commit()
except Exception as e:
await session.rollback()
# Log at appropriate level based on error type
if isinstance(e, HTTPException):
if HTTPStatus.BAD_REQUEST.value <= e.status_code < HTTPStatus.INTERNAL_SERVER_ERROR.value:
# Client errors (4xx) - log at info level
await logger.ainfo(f"Client error during session scope: {e.status_code}: {e.detail}")
else:
# Server errors (5xx) or other - log at error level
await logger.aexception("An error occurred during the session scope.", exception=e)
else:
# Non-HTTP exceptions - log at error level
await logger.aexception("An error occurred during the session scope.", exception=e)
raise
def get_cache_service() -> CacheService | AsyncBaseCacheService:
"""Retrieves the cache service from the service manager.
Returns:
The cache service instance.
"""
from langflow.services.cache.factory import CacheServiceFactory
return get_service(ServiceType.CACHE_SERVICE, CacheServiceFactory())
def get_shared_component_cache_service() -> CacheService:
"""Retrieves the cache service from the service manager.
Returns:
The cache service instance.
"""
from langflow.services.shared_component_cache.factory import SharedComponentCacheServiceFactory
return get_service(ServiceType.SHARED_COMPONENT_CACHE_SERVICE, SharedComponentCacheServiceFactory())
def get_session_service() -> SessionService:
"""Retrieves the session service from the service manager.
Returns:
The session service instance.
"""
from langflow.services.session.factory import SessionServiceFactory
return get_service(ServiceType.SESSION_SERVICE, SessionServiceFactory())
def get_task_service() -> TaskService:
"""Retrieves the TaskService instance from the service manager.
Returns:
The TaskService instance.
"""
from langflow.services.task.factory import TaskServiceFactory
return get_service(ServiceType.TASK_SERVICE, TaskServiceFactory())
def get_chat_service() -> ChatService:
"""Get the chat service instance.
Returns:
ChatService: The chat service instance.
"""
return get_service(ServiceType.CHAT_SERVICE)
def get_store_service() -> StoreService:
"""Retrieves the StoreService instance from the service manager.
Returns:
StoreService: The StoreService instance.
"""
return get_service(ServiceType.STORE_SERVICE)
def get_queue_service() -> JobQueueService:
"""Retrieves the QueueService instance from the service manager."""
from langflow.services.job_queue.factory import JobQueueServiceFactory
return get_service(ServiceType.JOB_QUEUE_SERVICE, JobQueueServiceFactory())