server: update process tools to save all events into a jsonl file

This commit is contained in:
2023-08-30 19:20:46 +02:00
committed by Mathieu Virbel
parent df078f7bd6
commit bdf7fe6ebc
3 changed files with 59 additions and 30 deletions

View File

@@ -1,10 +1,10 @@
from .base import Processor, ThreadedProcessor, Pipeline # noqa: F401
from .types import AudioFile, Transcript, Word, TitleSummary, FinalSummary # noqa: F401
from .audio_file_writer import AudioFileWriterProcessor # noqa: F401
from .audio_chunker import AudioChunkerProcessor # noqa: F401
from .audio_file_writer import AudioFileWriterProcessor # noqa: F401
from .audio_merge import AudioMergeProcessor # noqa: F401
from .audio_transcript import AudioTranscriptProcessor # noqa: F401
from .audio_transcript_auto import AudioTranscriptAutoProcessor # noqa: F401
from .base import Pipeline, PipelineEvent, Processor, ThreadedProcessor # noqa: F401
from .transcript_final_summary import TranscriptFinalSummaryProcessor # noqa: F401
from .transcript_liner import TranscriptLinerProcessor # noqa: F401
from .transcript_topic_detector import TranscriptTopicDetectorProcessor # noqa: F401
from .transcript_final_summary import TranscriptFinalSummaryProcessor # noqa: F401
from .types import AudioFile, FinalSummary, TitleSummary, Transcript, Word # noqa: F401

View File

@@ -3,15 +3,23 @@ from concurrent.futures import ThreadPoolExecutor
from typing import Any
from uuid import uuid4
from pydantic import BaseModel
from reflector.logger import logger
class PipelineEvent(BaseModel):
processor: str
uid: str
data: Any
class Processor:
INPUT_TYPE: type = None
OUTPUT_TYPE: type = None
WARMUP_EVENT: str = "WARMUP_EVENT"
def __init__(self, callback=None, custom_logger=None):
self.name = self.__class__.__name__
self._processors = []
self._callbacks = []
if callback:
@@ -67,6 +75,10 @@ class Processor:
return default
async def emit(self, data):
if self.pipeline:
await self.pipeline.emit(
PipelineEvent(processor=self.name, uid=self.uid, data=data)
)
for callback in self._callbacks:
await callback(data)
for processor in self._processors: