From 91d28e9d56ffa95cc38c78bbcef884acd04db234 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Wed, 16 Apr 2025 13:47:23 +0100 Subject: [PATCH] Change input_output to flow_processor --- trustgraph-base/trustgraph/base/__init__.py | 2 +- .../trustgraph/base/{input_output.py => flow_processor.py} | 6 +++--- trustgraph-flow/trustgraph/chunking/recursive/chunker.py | 6 +++--- trustgraph-flow/trustgraph/chunking/token/chunker.py | 6 +++--- trustgraph-flow/trustgraph/decoding/pdf/pdf_decoder.py | 6 +++--- 5 files changed, 13 insertions(+), 13 deletions(-) rename trustgraph-base/trustgraph/base/{input_output.py => flow_processor.py} (96%) diff --git a/trustgraph-base/trustgraph/base/__init__.py b/trustgraph-base/trustgraph/base/__init__.py index 79cb5af2..1b04d576 100644 --- a/trustgraph-base/trustgraph/base/__init__.py +++ b/trustgraph-base/trustgraph/base/__init__.py @@ -6,5 +6,5 @@ from . producer import Producer from . publisher import Publisher from . subscriber import Subscriber from . metrics import ProcessorMetrics, ConsumerMetrics, ProducerMetrics -from . input_output import InputOutputProcessor +from . flow_processor import FlowProcessor diff --git a/trustgraph-base/trustgraph/base/input_output.py b/trustgraph-base/trustgraph/base/flow_processor.py similarity index 96% rename from trustgraph-base/trustgraph/base/input_output.py rename to trustgraph-base/trustgraph/base/flow_processor.py index 6c08fa33..8d7cec7a 100644 --- a/trustgraph-base/trustgraph/base/input_output.py +++ b/trustgraph-base/trustgraph/base/flow_processor.py @@ -10,7 +10,7 @@ from .. base import AsyncProcessor, Consumer, Producer from .. base import ProcessorMetrics, ConsumerMetrics, ProducerMetrics -class InputOutputProcessor(AsyncProcessor): +class FlowProcessor(AsyncProcessor): def __init__(self, **params): @@ -23,7 +23,7 @@ class InputOutputProcessor(AsyncProcessor): } ) - super(InputOutputProcessor, self).__init__( + super(FlowProcessor, self).__init__( **params | { "id": self.id, } @@ -138,7 +138,7 @@ class InputOutputProcessor(AsyncProcessor): print("Handled config update") async def start(self): - await super(InputOutputProcessor, self).start() + await super(FlowProcessor, self).start() @staticmethod def add_args(parser, default_subscriber): diff --git a/trustgraph-flow/trustgraph/chunking/recursive/chunker.py b/trustgraph-flow/trustgraph/chunking/recursive/chunker.py index fce1b61f..963298d4 100755 --- a/trustgraph-flow/trustgraph/chunking/recursive/chunker.py +++ b/trustgraph-flow/trustgraph/chunking/recursive/chunker.py @@ -10,13 +10,13 @@ from prometheus_client import Histogram from ... schema import TextDocument, Chunk, Metadata from ... schema import text_ingest_queue, chunk_ingest_queue from ... log_level import LogLevel -from ... base import InputOutputProcessor +from ... base import FlowProcessor module = "chunker" default_subscriber = module -class Processor(InputOutputProcessor): +class Processor(FlowProcessor): def __init__(self, **params): @@ -88,7 +88,7 @@ class Processor(InputOutputProcessor): @staticmethod def add_args(parser): - InputOutputProcessor.add_args(parser, default_subscriber) + FlowProcessor.add_args(parser, default_subscriber) parser.add_argument( '-z', '--chunk-size', diff --git a/trustgraph-flow/trustgraph/chunking/token/chunker.py b/trustgraph-flow/trustgraph/chunking/token/chunker.py index eee16324..d35c003c 100755 --- a/trustgraph-flow/trustgraph/chunking/token/chunker.py +++ b/trustgraph-flow/trustgraph/chunking/token/chunker.py @@ -10,13 +10,13 @@ from prometheus_client import Histogram from ... schema import TextDocument, Chunk, Metadata from ... schema import text_ingest_queue, chunk_ingest_queue from ... log_level import LogLevel -from ... base import InputOutputProcessor +from ... base import FlowProcessor module = "chunker" default_subscriber = module -class Processor(InputOutputProcessor): +class Processor(FlowProcessor): def __init__(self, **params): @@ -87,7 +87,7 @@ class Processor(InputOutputProcessor): @staticmethod def add_args(parser): - InputOutputProcessor.add_args(parser, default_subscriber) + FlowProcessor.add_args(parser, default_subscriber) parser.add_argument( '-z', '--chunk-size', diff --git a/trustgraph-flow/trustgraph/decoding/pdf/pdf_decoder.py b/trustgraph-flow/trustgraph/decoding/pdf/pdf_decoder.py index 08f427aa..8deb0394 100755 --- a/trustgraph-flow/trustgraph/decoding/pdf/pdf_decoder.py +++ b/trustgraph-flow/trustgraph/decoding/pdf/pdf_decoder.py @@ -11,13 +11,13 @@ from langchain_community.document_loaders import PyPDFLoader from ... schema import Document, TextDocument, Metadata from ... schema import document_ingest_queue, text_ingest_queue from ... log_level import LogLevel -from ... base import InputOutputProcessor +from ... base import FlowProcessor module = "pdf-decoder" default_subscriber = module -class Processor(InputOutputProcessor): +class Processor(FlowProcessor): def __init__(self, **params): @@ -78,7 +78,7 @@ class Processor(InputOutputProcessor): @staticmethod def add_args(parser): - InputOutputProcessor.add_args(parser, default_subscriber) + FlowProcessor.add_args(parser, default_subscriber) def run():