Underground-Digital/Workflow-Engine
0
1import json2import logging3import os4import queue5import threading6import time7from datetime import timedelta8from typing import Any, Optional, Union9from uuid import UUID10 11from flask import current_app12 13from core.helper.encrypter import decrypt_token, encrypt_token, obfuscated_token14from core.ops.entities.config_entity import (15 LangfuseConfig,16 LangSmithConfig,17 TracingProviderEnum,18)19from core.ops.entities.trace_entity import (20 DatasetRetrievalTraceInfo,21 GenerateNameTraceInfo,22 MessageTraceInfo,23 ModerationTraceInfo,24 SuggestedQuestionTraceInfo,25 ToolTraceInfo,26 TraceTaskName,27 WorkflowTraceInfo,28)29from core.ops.langfuse_trace.langfuse_trace import LangFuseDataTrace30from core.ops.langsmith_trace.langsmith_trace import LangSmithDataTrace31from core.ops.utils import get_message_data32from extensions.ext_database import db33from models.model import App, AppModelConfig, Conversation, Message, MessageAgentThought, MessageFile, TraceAppConfig34from models.workflow import WorkflowAppLog, WorkflowRun35from tasks.ops_trace_task import process_trace_tasks36 37provider_config_map = {38 TracingProviderEnum.LANGFUSE.value: {39 "config_class": LangfuseConfig,40 "secret_keys": ["public_key", "secret_key"],41 "other_keys": ["host", "project_key"],42 "trace_instance": LangFuseDataTrace,43 },44 TracingProviderEnum.LANGSMITH.value: {45 "config_class": LangSmithConfig,46 "secret_keys": ["api_key"],47 "other_keys": ["project", "endpoint"],48 "trace_instance": LangSmithDataTrace,49 },50}51 52 53class OpsTraceManager:54 @classmethod55 def encrypt_tracing_config(56 cls, tenant_id: str, tracing_provider: str, tracing_config: dict, current_trace_config=None57 ):58 """59 Encrypt tracing config.60 :param tenant_id: tenant id61 :param tracing_provider: tracing provider62 :param tracing_config: tracing config dictionary to be encrypted63 :param current_trace_config: current tracing configuration for keeping existing values64 :return: encrypted tracing configuration65 """66 # Get the configuration class and the keys that require encryption67 config_class, secret_keys, other_keys = (68 provider_config_map[tracing_provider]["config_class"],69 provider_config_map[tracing_provider]["secret_keys"],70 provider_config_map[tracing_provider]["other_keys"],71 )72 73 new_config = {}74 # Encrypt necessary keys75 for key in secret_keys:76 if key in tracing_config:77 if "*" in tracing_config[key]:78 # If the key contains '*', retain the original value from the current config79 new_config[key] = current_trace_config.get(key, tracing_config[key])80 else:81 # Otherwise, encrypt the key82 new_config[key] = encrypt_token(tenant_id, tracing_config[key])83 84 for key in other_keys:85 new_config[key] = tracing_config.get(key, "")86 87 # Create a new instance of the config class with the new configuration88 encrypted_config = config_class(**new_config)89 return encrypted_config.model_dump()90 91 @classmethod92 def decrypt_tracing_config(cls, tenant_id: str, tracing_provider: str, tracing_config: dict):93 """94 Decrypt tracing config95 :param tenant_id: tenant id96 :param tracing_provider: tracing provider97 :param tracing_config: tracing config98 :return:99 """100 config_class, secret_keys, other_keys = (101 provider_config_map[tracing_provider]["config_class"],102 provider_config_map[tracing_provider]["secret_keys"],103 provider_config_map[tracing_provider]["other_keys"],104 )105 new_config = {}106 for key in secret_keys:107 if key in tracing_config:108 new_config[key] = decrypt_token(tenant_id, tracing_config[key])109 110 for key in other_keys:111 new_config[key] = tracing_config.get(key, "")112 113 return config_class(**new_config).model_dump()114 115 @classmethod116 def obfuscated_decrypt_token(cls, tracing_provider: str, decrypt_tracing_config: dict):117 """118 Decrypt tracing config119 :param tracing_provider: tracing provider120 :param decrypt_tracing_config: tracing config121 :return:122 """123 config_class, secret_keys, other_keys = (124 provider_config_map[tracing_provider]["config_class"],125 provider_config_map[tracing_provider]["secret_keys"],126 provider_config_map[tracing_provider]["other_keys"],127 )128 new_config = {}129 for key in secret_keys:130 if key in decrypt_tracing_config:131 new_config[key] = obfuscated_token(decrypt_tracing_config[key])132 133 for key in other_keys:134 new_config[key] = decrypt_tracing_config.get(key, "")135 return config_class(**new_config).model_dump()136 137 @classmethod138 def get_decrypted_tracing_config(cls, app_id: str, tracing_provider: str):139 """140 Get decrypted tracing config141 :param app_id: app id142 :param tracing_provider: tracing provider143 :return:144 """145 trace_config_data: TraceAppConfig = (146 db.session.query(TraceAppConfig)147 .filter(TraceAppConfig.app_id == app_id, TraceAppConfig.tracing_provider == tracing_provider)148 .first()149 )150 151 if not trace_config_data:152 return None153 154 # decrypt_token155 tenant_id = db.session.query(App).filter(App.id == app_id).first().tenant_id156 decrypt_tracing_config = cls.decrypt_tracing_config(157 tenant_id, tracing_provider, trace_config_data.tracing_config158 )159 160 return decrypt_tracing_config161 162 @classmethod163 def get_ops_trace_instance(164 cls,165 app_id: Optional[Union[UUID, str]] = None,166 ):167 """168 Get ops trace through model config169 :param app_id: app_id170 :return:171 """172 if isinstance(app_id, UUID):173 app_id = str(app_id)174 175 if app_id is None:176 return None177 178 app: App = db.session.query(App).filter(App.id == app_id).first()179 180 if app is None:181 return None182 183 app_ops_trace_config = json.loads(app.tracing) if app.tracing else None184 185 if app_ops_trace_config is None:186 return None187 188 tracing_provider = app_ops_trace_config.get("tracing_provider")189 190 if tracing_provider is None or tracing_provider not in provider_config_map:191 return None192 193 # decrypt_token194 decrypt_trace_config = cls.get_decrypted_tracing_config(app_id, tracing_provider)195 if app_ops_trace_config.get("enabled"):196 trace_instance, config_class = (197 provider_config_map[tracing_provider]["trace_instance"],198 provider_config_map[tracing_provider]["config_class"],199 )200 tracing_instance = trace_instance(config_class(**decrypt_trace_config))201 return tracing_instance202 203 return None204 205 @classmethod206 def get_app_config_through_message_id(cls, message_id: str):207 app_model_config = None208 message_data = db.session.query(Message).filter(Message.id == message_id).first()209 conversation_id = message_data.conversation_id210 conversation_data = db.session.query(Conversation).filter(Conversation.id == conversation_id).first()211 212 if conversation_data.app_model_config_id:213 app_model_config = (214 db.session.query(AppModelConfig)215 .filter(AppModelConfig.id == conversation_data.app_model_config_id)216 .first()217 )218 elif conversation_data.app_model_config_id is None and conversation_data.override_model_configs:219 app_model_config = conversation_data.override_model_configs220 221 return app_model_config222 223 @classmethod224 def update_app_tracing_config(cls, app_id: str, enabled: bool, tracing_provider: str):225 """226 Update app tracing config227 :param app_id: app id228 :param enabled: enabled229 :param tracing_provider: tracing provider230 :return:231 """232 # auth check233 if tracing_provider not in provider_config_map and tracing_provider is not None:234 raise ValueError(f"Invalid tracing provider: {tracing_provider}")235 236 app_config: App = db.session.query(App).filter(App.id == app_id).first()237 app_config.tracing = json.dumps(238 {239 "enabled": enabled,240 "tracing_provider": tracing_provider,241 }242 )243 db.session.commit()244 245 @classmethod246 def get_app_tracing_config(cls, app_id: str):247 """248 Get app tracing config249 :param app_id: app id250 :return:251 """252 app: App = db.session.query(App).filter(App.id == app_id).first()253 if not app.tracing:254 return {"enabled": False, "tracing_provider": None}255 app_trace_config = json.loads(app.tracing)256 return app_trace_config257 258 @staticmethod259 def check_trace_config_is_effective(tracing_config: dict, tracing_provider: str):260 """261 Check trace config is effective262 :param tracing_config: tracing config263 :param tracing_provider: tracing provider264 :return:265 """266 config_type, trace_instance = (267 provider_config_map[tracing_provider]["config_class"],268 provider_config_map[tracing_provider]["trace_instance"],269 )270 tracing_config = config_type(**tracing_config)271 return trace_instance(tracing_config).api_check()272 273 @staticmethod274 def get_trace_config_project_key(tracing_config: dict, tracing_provider: str):275 """276 get trace config is project key277 :param tracing_config: tracing config278 :param tracing_provider: tracing provider279 :return:280 """281 config_type, trace_instance = (282 provider_config_map[tracing_provider]["config_class"],283 provider_config_map[tracing_provider]["trace_instance"],284 )285 tracing_config = config_type(**tracing_config)286 return trace_instance(tracing_config).get_project_key()287 288 @staticmethod289 def get_trace_config_project_url(tracing_config: dict, tracing_provider: str):290 """291 get trace config is project key292 :param tracing_config: tracing config293 :param tracing_provider: tracing provider294 :return:295 """296 config_type, trace_instance = (297 provider_config_map[tracing_provider]["config_class"],298 provider_config_map[tracing_provider]["trace_instance"],299 )300 tracing_config = config_type(**tracing_config)301 return trace_instance(tracing_config).get_project_url()302 303 304class TraceTask:305 def __init__(306 self,307 trace_type: Any,308 message_id: Optional[str] = None,309 workflow_run: Optional[WorkflowRun] = None,310 conversation_id: Optional[str] = None,311 user_id: Optional[str] = None,312 timer: Optional[Any] = None,313 **kwargs,314 ):315 self.trace_type = trace_type316 self.message_id = message_id317 self.workflow_run = workflow_run318 self.conversation_id = conversation_id319 self.user_id = user_id320 self.timer = timer321 self.kwargs = kwargs322 self.file_base_url = os.getenv("FILES_URL", "http://127.0.0.1:5001")323 324 self.app_id = None325 326 def execute(self):327 return self.preprocess()328 329 def preprocess(self):330 preprocess_map = {331 TraceTaskName.CONVERSATION_TRACE: lambda: self.conversation_trace(**self.kwargs),332 TraceTaskName.WORKFLOW_TRACE: lambda: self.workflow_trace(333 self.workflow_run, self.conversation_id, self.user_id334 ),335 TraceTaskName.MESSAGE_TRACE: lambda: self.message_trace(self.message_id),336 TraceTaskName.MODERATION_TRACE: lambda: self.moderation_trace(self.message_id, self.timer, **self.kwargs),337 TraceTaskName.SUGGESTED_QUESTION_TRACE: lambda: self.suggested_question_trace(338 self.message_id, self.timer, **self.kwargs339 ),340 TraceTaskName.DATASET_RETRIEVAL_TRACE: lambda: self.dataset_retrieval_trace(341 self.message_id, self.timer, **self.kwargs342 ),343 TraceTaskName.TOOL_TRACE: lambda: self.tool_trace(self.message_id, self.timer, **self.kwargs),344 TraceTaskName.GENERATE_NAME_TRACE: lambda: self.generate_name_trace(345 self.conversation_id, self.timer, **self.kwargs346 ),347 }348 349 return preprocess_map.get(self.trace_type, lambda: None)()350 351 # process methods for different trace types352 def conversation_trace(self, **kwargs):353 return kwargs354 355 def workflow_trace(self, workflow_run: WorkflowRun, conversation_id, user_id):356 workflow_id = workflow_run.workflow_id357 tenant_id = workflow_run.tenant_id358 workflow_run_id = workflow_run.id359 workflow_run_elapsed_time = workflow_run.elapsed_time360 workflow_run_status = workflow_run.status361 workflow_run_inputs = workflow_run.inputs_dict362 workflow_run_outputs = workflow_run.outputs_dict363 workflow_run_version = workflow_run.version364 error = workflow_run.error or ""365 366 total_tokens = workflow_run.total_tokens367 368 file_list = workflow_run_inputs.get("sys.file") or []369 query = workflow_run_inputs.get("query") or workflow_run_inputs.get("sys.query") or ""370 371 # get workflow_app_log_id372 workflow_app_log_data = (373 db.session.query(WorkflowAppLog)374 .filter_by(tenant_id=tenant_id, app_id=workflow_run.app_id, workflow_run_id=workflow_run.id)375 .first()376 )377 workflow_app_log_id = str(workflow_app_log_data.id) if workflow_app_log_data else None378 # get message_id379 message_data = (380 db.session.query(Message.id)381 .filter_by(conversation_id=conversation_id, workflow_run_id=workflow_run_id)382 .first()383 )384 message_id = str(message_data.id) if message_data else None385 386 metadata = {387 "workflow_id": workflow_id,388 "conversation_id": conversation_id,389 "workflow_run_id": workflow_run_id,390 "tenant_id": tenant_id,391 "elapsed_time": workflow_run_elapsed_time,392 "status": workflow_run_status,393 "version": workflow_run_version,394 "total_tokens": total_tokens,395 "file_list": file_list,396 "triggered_form": workflow_run.triggered_from,397 "user_id": user_id,398 }399 400 workflow_trace_info = WorkflowTraceInfo(401 workflow_data=workflow_run.to_dict(),402 conversation_id=conversation_id,403 workflow_id=workflow_id,404 tenant_id=tenant_id,405 workflow_run_id=workflow_run_id,406 workflow_run_elapsed_time=workflow_run_elapsed_time,407 workflow_run_status=workflow_run_status,408 workflow_run_inputs=workflow_run_inputs,409 workflow_run_outputs=workflow_run_outputs,410 workflow_run_version=workflow_run_version,411 error=error,412 total_tokens=total_tokens,413 file_list=file_list,414 query=query,415 metadata=metadata,416 workflow_app_log_id=workflow_app_log_id,417 message_id=message_id,418 start_time=workflow_run.created_at,419 end_time=workflow_run.finished_at,420 )421 422 return workflow_trace_info423 424 def message_trace(self, message_id):425 message_data = get_message_data(message_id)426 if not message_data:427 return {}428 conversation_mode = db.session.query(Conversation.mode).filter_by(id=message_data.conversation_id).first()429 conversation_mode = conversation_mode[0]430 created_at = message_data.created_at431 inputs = message_data.message432 433 # get message file data434 message_file_data = db.session.query(MessageFile).filter_by(message_id=message_id).first()435 file_list = []436 if message_file_data and message_file_data.url is not None:437 file_url = f"{self.file_base_url}/{message_file_data.url}" if message_file_data else ""438 file_list.append(file_url)439 440 metadata = {441 "conversation_id": message_data.conversation_id,442 "ls_provider": message_data.model_provider,443 "ls_model_name": message_data.model_id,444 "status": message_data.status,445 "from_end_user_id": message_data.from_account_id,446 "from_account_id": message_data.from_account_id,447 "agent_based": message_data.agent_based,448 "workflow_run_id": message_data.workflow_run_id,449 "from_source": message_data.from_source,450 "message_id": message_id,451 }452 453 message_tokens = message_data.message_tokens454 455 message_trace_info = MessageTraceInfo(456 message_id=message_id,457 message_data=message_data.to_dict(),458 conversation_model=conversation_mode,459 message_tokens=message_tokens,460 answer_tokens=message_data.answer_tokens,461 total_tokens=message_tokens + message_data.answer_tokens,462 error=message_data.error or "",463 inputs=inputs,464 outputs=message_data.answer,465 file_list=file_list,466 start_time=created_at,467 end_time=created_at + timedelta(seconds=message_data.provider_response_latency),468 metadata=metadata,469 message_file_data=message_file_data,470 conversation_mode=conversation_mode,471 )472 473 return message_trace_info474 475 def moderation_trace(self, message_id, timer, **kwargs):476 moderation_result = kwargs.get("moderation_result")477 inputs = kwargs.get("inputs")478 message_data = get_message_data(message_id)479 if not message_data:480 return {}481 metadata = {482 "message_id": message_id,483 "action": moderation_result.action,484 "preset_response": moderation_result.preset_response,485 "query": moderation_result.query,486 }487 488 # get workflow_app_log_id489 workflow_app_log_id = None490 if message_data.workflow_run_id:491 workflow_app_log_data = (492 db.session.query(WorkflowAppLog).filter_by(workflow_run_id=message_data.workflow_run_id).first()493 )494 workflow_app_log_id = str(workflow_app_log_data.id) if workflow_app_log_data else None495 496 moderation_trace_info = ModerationTraceInfo(497 message_id=workflow_app_log_id or message_id,498 inputs=inputs,499 message_data=message_data.to_dict(),500 flagged=moderation_result.flagged,501 action=moderation_result.action,502 preset_response=moderation_result.preset_response,503 query=moderation_result.query,504 start_time=timer.get("start"),505 end_time=timer.get("end"),506 metadata=metadata,507 )508 509 return moderation_trace_info510 511 def suggested_question_trace(self, message_id, timer, **kwargs):512 suggested_question = kwargs.get("suggested_question")513 message_data = get_message_data(message_id)514 if not message_data:515 return {}516 metadata = {517 "message_id": message_id,518 "ls_provider": message_data.model_provider,519 "ls_model_name": message_data.model_id,520 "status": message_data.status,521 "from_end_user_id": message_data.from_account_id,522 "from_account_id": message_data.from_account_id,523 "agent_based": message_data.agent_based,524 "workflow_run_id": message_data.workflow_run_id,525 "from_source": message_data.from_source,526 }527 528 # get workflow_app_log_id529 workflow_app_log_id = None530 if message_data.workflow_run_id:531 workflow_app_log_data = (532 db.session.query(WorkflowAppLog).filter_by(workflow_run_id=message_data.workflow_run_id).first()533 )534 workflow_app_log_id = str(workflow_app_log_data.id) if workflow_app_log_data else None535 536 suggested_question_trace_info = SuggestedQuestionTraceInfo(537 message_id=workflow_app_log_id or message_id,538 message_data=message_data.to_dict(),539 inputs=message_data.message,540 outputs=message_data.answer,541 start_time=timer.get("start"),542 end_time=timer.get("end"),543 metadata=metadata,544 total_tokens=message_data.message_tokens + message_data.answer_tokens,545 status=message_data.status,546 error=message_data.error,547 from_account_id=message_data.from_account_id,548 agent_based=message_data.agent_based,549 from_source=message_data.from_source,550 model_provider=message_data.model_provider,551 model_id=message_data.model_id,552 suggested_question=suggested_question,553 level=message_data.status,554 status_message=message_data.error,555 )556 557 return suggested_question_trace_info558 559 def dataset_retrieval_trace(self, message_id, timer, **kwargs):560 documents = kwargs.get("documents")561 message_data = get_message_data(message_id)562 if not message_data:563 return {}564 565 metadata = {566 "message_id": message_id,567 "ls_provider": message_data.model_provider,568 "ls_model_name": message_data.model_id,569 "status": message_data.status,570 "from_end_user_id": message_data.from_account_id,571 "from_account_id": message_data.from_account_id,572 "agent_based": message_data.agent_based,573 "workflow_run_id": message_data.workflow_run_id,574 "from_source": message_data.from_source,575 }576 577 dataset_retrieval_trace_info = DatasetRetrievalTraceInfo(578 message_id=message_id,579 inputs=message_data.query or message_data.inputs,580 documents=[doc.model_dump() for doc in documents],581 start_time=timer.get("start"),582 end_time=timer.get("end"),583 metadata=metadata,584 message_data=message_data.to_dict(),585 )586 587 return dataset_retrieval_trace_info588 589 def tool_trace(self, message_id, timer, **kwargs):590 tool_name = kwargs.get("tool_name")591 tool_inputs = kwargs.get("tool_inputs")592 tool_outputs = kwargs.get("tool_outputs")593 message_data = get_message_data(message_id)594 if not message_data:595 return {}596 tool_config = {}597 time_cost = 0598 error = None599 tool_parameters = {}600 created_time = message_data.created_at601 end_time = message_data.updated_at602 agent_thoughts: list[MessageAgentThought] = message_data.agent_thoughts603 for agent_thought in agent_thoughts:604 if tool_name in agent_thought.tools:605 created_time = agent_thought.created_at606 tool_meta_data = agent_thought.tool_meta.get(tool_name, {})607 tool_config = tool_meta_data.get("tool_config", {})608 time_cost = tool_meta_data.get("time_cost", 0)609 end_time = created_time + timedelta(seconds=time_cost)610 error = tool_meta_data.get("error", "")611 tool_parameters = tool_meta_data.get("tool_parameters", {})612 metadata = {613 "message_id": message_id,614 "tool_name": tool_name,615 "tool_inputs": tool_inputs,616 "tool_outputs": tool_outputs,617 "tool_config": tool_config,618 "time_cost": time_cost,619 "error": error,620 "tool_parameters": tool_parameters,621 }622 623 file_url = ""624 message_file_data = db.session.query(MessageFile).filter_by(message_id=message_id).first()625 if message_file_data:626 message_file_id = message_file_data.id if message_file_data else None627 type = message_file_data.type628 created_by_role = message_file_data.created_by_role629 created_user_id = message_file_data.created_by630 file_url = f"{self.file_base_url}/{message_file_data.url}"631 632 metadata.update(633 {634 "message_file_id": message_file_id,635 "created_by_role": created_by_role,636 "created_user_id": created_user_id,637 "type": type,638 }639 )640 641 tool_trace_info = ToolTraceInfo(642 message_id=message_id,643 message_data=message_data.to_dict(),644 tool_name=tool_name,645 start_time=timer.get("start") if timer else created_time,646 end_time=timer.get("end") if timer else end_time,647 tool_inputs=tool_inputs,648 tool_outputs=tool_outputs,649 metadata=metadata,650 message_file_data=message_file_data,651 error=error,652 inputs=message_data.message,653 outputs=message_data.answer,654 tool_config=tool_config,655 time_cost=time_cost,656 tool_parameters=tool_parameters,657 file_url=file_url,658 )659 660 return tool_trace_info661 662 def generate_name_trace(self, conversation_id, timer, **kwargs):663 generate_conversation_name = kwargs.get("generate_conversation_name")664 inputs = kwargs.get("inputs")665 tenant_id = kwargs.get("tenant_id")666 start_time = timer.get("start")667 end_time = timer.get("end")668 669 metadata = {670 "conversation_id": conversation_id,671 "tenant_id": tenant_id,672 }673 674 generate_name_trace_info = GenerateNameTraceInfo(675 conversation_id=conversation_id,676 inputs=inputs,677 outputs=generate_conversation_name,678 start_time=start_time,679 end_time=end_time,680 metadata=metadata,681 tenant_id=tenant_id,682 )683 684 return generate_name_trace_info685 686 687trace_manager_timer = None688trace_manager_queue = queue.Queue()689trace_manager_interval = int(os.getenv("TRACE_QUEUE_MANAGER_INTERVAL", 5))690trace_manager_batch_size = int(os.getenv("TRACE_QUEUE_MANAGER_BATCH_SIZE", 100))691 692 693class TraceQueueManager:694 def __init__(self, app_id=None, user_id=None):695 global trace_manager_timer696 697 self.app_id = app_id698 self.user_id = user_id699 self.trace_instance = OpsTraceManager.get_ops_trace_instance(app_id)700 self.flask_app = current_app._get_current_object()701 if trace_manager_timer is None:702 self.start_timer()703 704 def add_trace_task(self, trace_task: TraceTask):705 global trace_manager_timer, trace_manager_queue706 try:707 if self.trace_instance:708 trace_task.app_id = self.app_id709 trace_manager_queue.put(trace_task)710 except Exception as e:711 logging.error(f"Error adding trace task: {e}")712 finally:713 self.start_timer()714 715 def collect_tasks(self):716 global trace_manager_queue717 tasks = []718 while len(tasks) < trace_manager_batch_size and not trace_manager_queue.empty():719 task = trace_manager_queue.get_nowait()720 tasks.append(task)721 trace_manager_queue.task_done()722 return tasks723 724 def run(self):725 try:726 tasks = self.collect_tasks()727 if tasks:728 self.send_to_celery(tasks)729 except Exception as e:730 logging.error(f"Error processing trace tasks: {e}")731 732 def start_timer(self):733 global trace_manager_timer734 if trace_manager_timer is None or not trace_manager_timer.is_alive():735 trace_manager_timer = threading.Timer(trace_manager_interval, self.run)736 trace_manager_timer.name = f"trace_manager_timer_{time.strftime('%Y-%m-%d %H:%M:%S', time.localtime())}"737 trace_manager_timer.daemon = False738 trace_manager_timer.start()739 740 def send_to_celery(self, tasks: list[TraceTask]):741 with self.flask_app.app_context():742 for task in tasks:743 trace_info = task.execute()744 task_data = {745 "app_id": task.app_id,746 "trace_info_type": type(trace_info).__name__,747 "trace_info": trace_info.model_dump() if trace_info else {},748 }749 process_trace_tasks.delay(task_data)750 