ThreadPoolExecutor that automatically propagates the main thread's trace_id to worker threads.
| 242 | |
| 243 | |
| 244 | class ContextThreadPoolExecutor(ThreadPoolExecutor): |
| 245 | """ |
| 246 | ThreadPoolExecutor that automatically propagates the main thread's trace_id to worker threads. |
| 247 | """ |
| 248 | |
| 249 | def submit(self, fn: Callable[..., T], *args: Any, **kwargs: Any) -> Any: |
| 250 | """ |
| 251 | Submit a callable to be executed with the given arguments. |
| 252 | Automatically propagates the current thread's context to the worker thread. |
| 253 | """ |
| 254 | main_trace_id = get_current_trace_id() |
| 255 | main_api_path = get_current_api_path() |
| 256 | main_env = get_current_env() |
| 257 | main_user_type = get_current_user_type() |
| 258 | main_user_name = get_current_user_name() |
| 259 | main_context = get_current_context() |
| 260 | |
| 261 | @functools.wraps(fn) |
| 262 | def wrapper(*args: Any, **kwargs: Any) -> Any: |
| 263 | if main_context: |
| 264 | # Create and set new context in worker thread |
| 265 | child_context = RequestContext( |
| 266 | trace_id=main_trace_id, |
| 267 | api_path=main_api_path, |
| 268 | env=main_env, |
| 269 | user_type=main_user_type, |
| 270 | user_name=main_user_name, |
| 271 | ) |
| 272 | child_context._data = main_context._data.copy() |
| 273 | set_request_context(child_context) |
| 274 | |
| 275 | return fn(*args, **kwargs) |
| 276 | |
| 277 | return super().submit(wrapper, *args, **kwargs) |
| 278 | |
| 279 | def map( |
| 280 | self, |
| 281 | fn: Callable[..., T], |
| 282 | *iterables: Any, |
| 283 | timeout: float | None = None, |
| 284 | chunksize: int = 1, |
| 285 | ) -> Any: |
| 286 | """ |
| 287 | Returns an iterator equivalent to map(fn, iter). |
| 288 | Automatically propagates the current thread's context to worker threads. |
| 289 | """ |
| 290 | main_trace_id = get_current_trace_id() |
| 291 | main_api_path = get_current_api_path() |
| 292 | main_env = get_current_env() |
| 293 | main_user_type = get_current_user_type() |
| 294 | main_user_name = get_current_user_name() |
| 295 | main_context = get_current_context() |
| 296 | |
| 297 | @functools.wraps(fn) |
| 298 | def wrapper(*args: Any, **kwargs: Any) -> Any: |
| 299 | if main_context: |
| 300 | # Create and set new context in worker thread |
| 301 | child_context = RequestContext( |
no outgoing calls