Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
output_moderation.py132 linesDownload Raw Back to moderation
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