MCPcopy Create free account
hub / github.com/MemTensor/MemOS / ContextThreadPoolExecutor

Class ContextThreadPoolExecutor

src/memos/context/context.py:244–313  ·  view source on GitHub ↗

ThreadPoolExecutor that automatically propagates the main thread's trace_id to worker threads.

Source from the content-addressed store, hash-verified

242
243
244class 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(

Callers 15

_restart_poolMethod · 0.90
__init__Method · 0.90
batch_handlerMethod · 0.90
batch_handlerMethod · 0.90
batch_handlerMethod · 0.90
searchMethod · 0.90
addMethod · 0.90
get_sub_answersMethod · 0.90
_search_with_engineMethod · 0.90
add_graph_edgesMethod · 0.90
soft_deleteMethod · 0.90

Calls

no outgoing calls