diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/core_export.py b/trustgraph-flow/trustgraph/gateway/dispatch/core_export.py index 941ce5d8..61b0bcbc 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/core_export.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/core_export.py @@ -2,8 +2,12 @@ import asyncio import uuid import msgpack +import logging from . knowledge import KnowledgeRequestor +# Module logger +logger = logging.getLogger(__name__) + class CoreExport: def __init__(self, pulsar_client): @@ -84,7 +88,7 @@ class CoreExport: except Exception as e: - print("Exception:", e) + logger.error(f"Core export exception: {e}", exc_info=True) finally: diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/core_import.py b/trustgraph-flow/trustgraph/gateway/dispatch/core_import.py index b819d286..b32fb7f7 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/core_import.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/core_import.py @@ -3,8 +3,12 @@ import asyncio import json import uuid import msgpack +import logging from . knowledge import KnowledgeRequestor +# Module logger +logger = logging.getLogger(__name__) + class CoreImport: def __init__(self, pulsar_client): @@ -80,14 +84,14 @@ class CoreImport: await kr.process(msg) except Exception as e: - print("Exception:", e) + logger.error(f"Core import exception: {e}", exc_info=True) await error(str(e)) finally: await kr.stop() - print("All done.") + logger.info("Core import completed") response = await ok() await response.write_eof() diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/document_load.py b/trustgraph-flow/trustgraph/gateway/dispatch/document_load.py index 101e9b41..7e38877c 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/document_load.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/document_load.py @@ -1,11 +1,15 @@ import base64 +import logging from ... schema import Document, Metadata from ... messaging import TranslatorRegistry from . sender import ServiceSender +# Module logger +logger = logging.getLogger(__name__) + class DocumentLoad(ServiceSender): def __init__(self, pulsar_client, queue): @@ -18,6 +22,6 @@ class DocumentLoad(ServiceSender): self.translator = TranslatorRegistry.get_request_translator("document") def to_request(self, body): - print("Document received") + logger.info("Document received") return self.translator.to_pulsar(body) diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/manager.py b/trustgraph-flow/trustgraph/gateway/dispatch/manager.py index b32a6253..9ec7b0ab 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/manager.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/manager.py @@ -2,6 +2,10 @@ import asyncio from aiohttp import web import uuid +import logging + +# Module logger +logger = logging.getLogger(__name__) from . config import ConfigRequestor from . flow import FlowRequestor @@ -92,12 +96,12 @@ class DispatcherManager: self.dispatchers = {} async def start_flow(self, id, flow): - print("Start flow", id) + logger.info(f"Starting flow {id}") self.flows[id] = flow return async def stop_flow(self, id, flow): - print("Stop flow", id) + logger.info(f"Stopping flow {id}") del self.flows[id] return diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/text_load.py b/trustgraph-flow/trustgraph/gateway/dispatch/text_load.py index 8f30c8de..36922c89 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/text_load.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/text_load.py @@ -1,11 +1,15 @@ import base64 +import logging from ... schema import TextDocument, Metadata from ... messaging import TranslatorRegistry from . sender import ServiceSender +# Module logger +logger = logging.getLogger(__name__) + class TextLoad(ServiceSender): def __init__(self, pulsar_client, queue): @@ -18,6 +22,6 @@ class TextLoad(ServiceSender): self.translator = TranslatorRegistry.get_request_translator("text-document") def to_request(self, body): - print("Text document received") + logger.info("Text document received") return self.translator.to_pulsar(body) diff --git a/trustgraph-flow/trustgraph/gateway/endpoint/socket.py b/trustgraph-flow/trustgraph/gateway/endpoint/socket.py index 1bfec637..c912a460 100644 --- a/trustgraph-flow/trustgraph/gateway/endpoint/socket.py +++ b/trustgraph-flow/trustgraph/gateway/endpoint/socket.py @@ -74,24 +74,24 @@ class SocketEndpoint: self.listener(ws, dispatcher, running) ) - print("Created taskgroup, waiting...") + logger.debug("Created task group, waiting for completion...") # Wait for threads to complete - print("Task group closed") + logger.debug("Task group closed") # Finally? await dispatcher.destroy() except ExceptionGroup as e: - print("Exception group:", flush=True) + logger.error("Exception group occurred:", exc_info=True) for se in e.exceptions: - print(" Type:", type(se), flush=True) - print(f" Exception: {se}", flush=True) + logger.error(f" Exception type: {type(se)}") + logger.error(f" Exception: {se}") except Exception as e: - print("Socket exception:", e, flush=True) + logger.error(f"Socket exception: {e}", exc_info=True) await ws.close()