Underground-Digital/Workflow-Engine
0
1import logging2import threading3import time4from typing import Any, Optional5 6from flask import Flask, current_app7from pydantic import BaseModel, ConfigDict8 9from configs import dify_config10from core.app.apps.base_app_queue_manager import AppQueueManager, PublishFrom11from core.app.entities.queue_entities import QueueMessageReplaceEvent12from core.moderation.base import ModerationAction, ModerationOutputsResult13from core.moderation.factory import ModerationFactory14 15logger = logging.getLogger(__name__)16 17 18class ModerationRule(BaseModel):19 type: str20 config: dict[str, Any]21 22 23class OutputModeration(BaseModel):24 tenant_id: str25 app_id: str26 27 rule: ModerationRule28 queue_manager: AppQueueManager29 30 thread: Optional[threading.Thread] = None31 thread_running: bool = True32 buffer: str = ""33 is_final_chunk: bool = False34 final_output: Optional[str] = None35 model_config = ConfigDict(arbitrary_types_allowed=True)36 37 def should_direct_output(self) -> bool:38 return self.final_output is not None39 40 def get_final_output(self) -> str:41 return self.final_output or ""42 43 def append_new_token(self, token: str) -> None:44 self.buffer += token45 46 if not self.thread:47 self.thread = self.start_thread()48 49 def moderation_completion(self, completion: str, public_event: bool = False) -> str:50 self.buffer = completion51 self.is_final_chunk = True52 53 result = self.moderation(tenant_id=self.tenant_id, app_id=self.app_id, moderation_buffer=completion)54 55 if not result or not result.flagged:56 return completion57 58 if result.action == ModerationAction.DIRECT_OUTPUT:59 final_output = result.preset_response60 else:61 final_output = result.text62 63 if public_event:64 self.queue_manager.publish(QueueMessageReplaceEvent(text=final_output), PublishFrom.TASK_PIPELINE)65 66 return final_output67 68 def start_thread(self) -> threading.Thread:69 buffer_size = dify_config.MODERATION_BUFFER_SIZE70 thread = threading.Thread(71 target=self.worker,72 kwargs={73 "flask_app": current_app._get_current_object(),74 "buffer_size": buffer_size if buffer_size > 0 else dify_config.MODERATION_BUFFER_SIZE,75 },76 )77 78 thread.start()79 80 return thread81 82 def stop_thread(self):83 if self.thread and self.thread.is_alive():84 self.thread_running = False85 86 def worker(self, flask_app: Flask, buffer_size: int):87 with flask_app.app_context():88 current_length = 089 while self.thread_running:90 moderation_buffer = self.buffer91 buffer_length = len(moderation_buffer)92 if not self.is_final_chunk:93 chunk_length = buffer_length - current_length94 if 0 <= chunk_length < buffer_size:95 time.sleep(1)96 continue97 98 current_length = buffer_length99 100 result = self.moderation(101 tenant_id=self.tenant_id, app_id=self.app_id, moderation_buffer=moderation_buffer102 )103 104 if not result or not result.flagged:105 continue106 107 if result.action == ModerationAction.DIRECT_OUTPUT:108 final_output = result.preset_response109 self.final_output = final_output110 else:111 final_output = result.text + self.buffer[len(moderation_buffer) :]112 113 # trigger replace event114 if self.thread_running:115 self.queue_manager.publish(QueueMessageReplaceEvent(text=final_output), PublishFrom.TASK_PIPELINE)116 117 if result.action == ModerationAction.DIRECT_OUTPUT:118 break119 120 def moderation(self, tenant_id: str, app_id: str, moderation_buffer: str) -> Optional[ModerationOutputsResult]:121 try:122 moderation_factory = ModerationFactory(123 name=self.rule.type, app_id=app_id, tenant_id=tenant_id, config=self.rule.config124 )125 126 result: ModerationOutputsResult = moderation_factory.moderation_for_outputs(moderation_buffer)127 return result128 except Exception as e:129 logger.error("Moderation Output error: %s", e)130 131 return None132 