Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
ops_trace_manager.py750 linesDownload Raw Back to ops
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