diff --git a/trustgraph-base/trustgraph/base/flow_processor.py b/trustgraph-base/trustgraph/base/flow_processor.py index fdeb5950..6d3ba64f 100644 --- a/trustgraph-base/trustgraph/base/flow_processor.py +++ b/trustgraph-base/trustgraph/base/flow_processor.py @@ -4,6 +4,7 @@ # configuration service which can't manage itself. import json +import logging from pulsar.schema import JsonSchema @@ -14,6 +15,9 @@ from .. log_level import LogLevel from . async_processor import AsyncProcessor from . flow import Flow +# Module logger +logger = logging.getLogger(__name__) + # Parent class for configurable processors, configured with flows by # the config service class FlowProcessor(AsyncProcessor): @@ -34,7 +38,7 @@ class FlowProcessor(AsyncProcessor): # Array of specifications: ConsumerSpec, ProducerSpec, SettingSpec self.specifications = [] - print("Service initialised.") + logger.info("Service initialised.") # Register a configuration variable def register_specification(self, spec): @@ -44,19 +48,19 @@ class FlowProcessor(AsyncProcessor): async def start_flow(self, flow, defn): self.flows[flow] = Flow(self.id, flow, self, defn) await self.flows[flow].start() - print("Started flow: ", flow) + logger.info(f"Started flow: {flow}") # Stop processing for a new flow async def stop_flow(self, flow): if flow in self.flows: await self.flows[flow].stop() del self.flows[flow] - print("Stopped flow: ", flow, flush=True) + logger.info(f"Stopped flow: {flow}") # Event handler - called for a configuration change async def on_configure_flows(self, config, version): - print("Got config version", version, flush=True) + logger.info(f"Got config version {version}") # Skip over invalid data if "flows-active" not in config: return @@ -69,7 +73,7 @@ class FlowProcessor(AsyncProcessor): else: - print("No configuration settings for me.", flush=True) + logger.debug("No configuration settings for me.") flow_config = {} # Get list of flows which should be running and are currently @@ -88,7 +92,7 @@ class FlowProcessor(AsyncProcessor): if flow not in wanted_flows: await self.stop_flow(flow) - print("Handled config update") + logger.info("Handled config update") # Start threads, just call parent async def start(self): diff --git a/trustgraph-flow/trustgraph/storage/doc_embeddings/pinecone/write.py b/trustgraph-flow/trustgraph/storage/doc_embeddings/pinecone/write.py index 0d8bac83..1851a243 100644 --- a/trustgraph-flow/trustgraph/storage/doc_embeddings/pinecone/write.py +++ b/trustgraph-flow/trustgraph/storage/doc_embeddings/pinecone/write.py @@ -9,9 +9,13 @@ from pinecone.grpc import PineconeGRPC, GRPCClientConfig import time import uuid import os +import logging from .... base import DocumentEmbeddingsStoreService +# Module logger +logger = logging.getLogger(__name__) + default_ident = "de-write" default_api_key = os.getenv("PINECONE_API_KEY", "not-specified") default_cloud = "aws" @@ -104,10 +108,10 @@ class Processor(DocumentEmbeddingsStoreService): self.create_index(index_name, dim) except Exception as e: - print("Pinecone index creation failed") + logger.error("Pinecone index creation failed") raise e - print(f"Index {index_name} created", flush=True) + logger.info(f"Index {index_name} created") self.last_index_name = index_name diff --git a/trustgraph-flow/trustgraph/storage/doc_embeddings/qdrant/write.py b/trustgraph-flow/trustgraph/storage/doc_embeddings/qdrant/write.py index d65a75eb..6005df1f 100644 --- a/trustgraph-flow/trustgraph/storage/doc_embeddings/qdrant/write.py +++ b/trustgraph-flow/trustgraph/storage/doc_embeddings/qdrant/write.py @@ -7,9 +7,13 @@ from qdrant_client import QdrantClient from qdrant_client.models import PointStruct from qdrant_client.models import Distance, VectorParams import uuid +import logging from .... base import DocumentEmbeddingsStoreService +# Module logger +logger = logging.getLogger(__name__) + default_ident = "de-write" default_store_uri = 'http://localhost:6333' @@ -60,7 +64,7 @@ class Processor(DocumentEmbeddingsStoreService): ), ) except Exception as e: - print("Qdrant collection creation failed") + logger.error("Qdrant collection creation failed") raise e self.last_collection = collection diff --git a/trustgraph-flow/trustgraph/storage/graph_embeddings/pinecone/write.py b/trustgraph-flow/trustgraph/storage/graph_embeddings/pinecone/write.py index e575d12a..f73cfd22 100755 --- a/trustgraph-flow/trustgraph/storage/graph_embeddings/pinecone/write.py +++ b/trustgraph-flow/trustgraph/storage/graph_embeddings/pinecone/write.py @@ -9,9 +9,13 @@ from pinecone.grpc import PineconeGRPC, GRPCClientConfig import time import uuid import os +import logging from .... base import GraphEmbeddingsStoreService +# Module logger +logger = logging.getLogger(__name__) + default_ident = "ge-write" default_api_key = os.getenv("PINECONE_API_KEY", "not-specified") default_cloud = "aws" @@ -103,10 +107,10 @@ class Processor(GraphEmbeddingsStoreService): self.create_index(index_name, dim) except Exception as e: - print("Pinecone index creation failed") + logger.error("Pinecone index creation failed") raise e - print(f"Index {index_name} created", flush=True) + logger.info(f"Index {index_name} created") self.last_index_name = index_name diff --git a/trustgraph-flow/trustgraph/storage/graph_embeddings/qdrant/write.py b/trustgraph-flow/trustgraph/storage/graph_embeddings/qdrant/write.py index ecefee4f..903702c7 100755 --- a/trustgraph-flow/trustgraph/storage/graph_embeddings/qdrant/write.py +++ b/trustgraph-flow/trustgraph/storage/graph_embeddings/qdrant/write.py @@ -7,9 +7,13 @@ from qdrant_client import QdrantClient from qdrant_client.models import PointStruct from qdrant_client.models import Distance, VectorParams import uuid +import logging from .... base import GraphEmbeddingsStoreService +# Module logger +logger = logging.getLogger(__name__) + default_ident = "ge-write" default_store_uri = 'http://localhost:6333' @@ -50,7 +54,7 @@ class Processor(GraphEmbeddingsStoreService): ), ) except Exception as e: - print("Qdrant collection creation failed") + logger.error("Qdrant collection creation failed") raise e self.last_collection = cname diff --git a/trustgraph-flow/trustgraph/storage/rows/cassandra/write.py b/trustgraph-flow/trustgraph/storage/rows/cassandra/write.py index a84aefde..e8948668 100755 --- a/trustgraph-flow/trustgraph/storage/rows/cassandra/write.py +++ b/trustgraph-flow/trustgraph/storage/rows/cassandra/write.py @@ -8,6 +8,7 @@ import base64 import os import argparse import time +import logging from cassandra.cluster import Cluster from cassandra.auth import PlainTextAuthProvider from ssl import SSLContext, PROTOCOL_TLSv1_2 @@ -17,6 +18,9 @@ from .... schema import rows_store_queue from .... log_level import LogLevel from .... base import Consumer +# Module logger +logger = logging.getLogger(__name__) + module = "rows-write" ssl_context = SSLContext(PROTOCOL_TLSv1_2) @@ -111,7 +115,7 @@ class Processor(Consumer): except Exception as e: - print("Exception:", str(e), flush=True) + logger.error(f"Exception: {str(e)}", exc_info=True) # If there's an error make sure to do table creation etc. self.tables.remove(name) diff --git a/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py b/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py index f8396692..ac790bcc 100755 --- a/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py +++ b/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py @@ -8,10 +8,14 @@ import base64 import os import argparse import time +import logging from .... direct.cassandra import TrustGraph from .... base import TriplesStoreService +# Module logger +logger = logging.getLogger(__name__) + default_ident = "triples-write" default_graph_host='localhost' @@ -61,7 +65,7 @@ class Processor(TriplesStoreService): table=message.metadata.collection, ) except Exception as e: - print("Exception", e, flush=True) + logger.error(f"Exception: {e}", exc_info=True) time.sleep(1) raise e diff --git a/trustgraph-flow/trustgraph/storage/triples/falkordb/write.py b/trustgraph-flow/trustgraph/storage/triples/falkordb/write.py index defb7d69..b71c247b 100755 --- a/trustgraph-flow/trustgraph/storage/triples/falkordb/write.py +++ b/trustgraph-flow/trustgraph/storage/triples/falkordb/write.py @@ -8,11 +8,15 @@ import base64 import os import argparse import time +import logging from falkordb import FalkorDB from .... base import TriplesStoreService +# Module logger +logger = logging.getLogger(__name__) + default_ident = "triples-write" default_graph_url = 'falkor://falkordb:6379' @@ -38,7 +42,7 @@ class Processor(TriplesStoreService): def create_node(self, uri): - print("Create node", uri) + logger.debug(f"Create node {uri}") res = self.io.query( "MERGE (n:Node {uri: $uri})", @@ -47,14 +51,14 @@ class Processor(TriplesStoreService): }, ) - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=res.nodes_created, time=res.run_time_ms )) def create_literal(self, value): - print("Create literal", value) + logger.debug(f"Create literal {value}") res = self.io.query( "MERGE (n:Literal {value: $value})", @@ -63,14 +67,14 @@ class Processor(TriplesStoreService): }, ) - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=res.nodes_created, time=res.run_time_ms )) def relate_node(self, src, uri, dest): - print("Create node rel", src, uri, dest) + logger.debug(f"Create node rel {src} {uri} {dest}") res = self.io.query( "MATCH (src:Node {uri: $src}) " @@ -83,14 +87,14 @@ class Processor(TriplesStoreService): }, ) - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=res.nodes_created, time=res.run_time_ms )) def relate_literal(self, src, uri, dest): - print("Create literal rel", src, uri, dest) + logger.debug(f"Create literal rel {src} {uri} {dest}") res = self.io.query( "MATCH (src:Node {uri: $src}) " @@ -103,7 +107,7 @@ class Processor(TriplesStoreService): }, ) - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=res.nodes_created, time=res.run_time_ms )) diff --git a/trustgraph-flow/trustgraph/storage/triples/memgraph/write.py b/trustgraph-flow/trustgraph/storage/triples/memgraph/write.py index 9079923e..fa0260ac 100755 --- a/trustgraph-flow/trustgraph/storage/triples/memgraph/write.py +++ b/trustgraph-flow/trustgraph/storage/triples/memgraph/write.py @@ -8,11 +8,15 @@ import base64 import os import argparse import time +import logging from neo4j import GraphDatabase from .... base import TriplesStoreService +# Module logger +logger = logging.getLogger(__name__) + default_ident = "triples-write" default_graph_host = 'bolt://memgraph:7687' @@ -55,49 +59,49 @@ class Processor(TriplesStoreService): # and this process will restart several times until Pulsar arrives, # so should be safe - print("Create indexes...", flush=True) + logger.info("Create indexes...") try: session.run( "CREATE INDEX ON :Node", ) except Exception as e: - print(e, flush=True) + logger.warning(f"Index create failure: {e}") # Maybe index already exists - print("Index create failure ignored", flush=True) + logger.warning("Index create failure ignored") try: session.run( "CREATE INDEX ON :Node(uri)" ) except Exception as e: - print(e, flush=True) + logger.warning(f"Index create failure: {e}") # Maybe index already exists - print("Index create failure ignored", flush=True) + logger.warning("Index create failure ignored") try: session.run( "CREATE INDEX ON :Literal", ) except Exception as e: - print(e, flush=True) + logger.warning(f"Index create failure: {e}") # Maybe index already exists - print("Index create failure ignored", flush=True) + logger.warning("Index create failure ignored") try: session.run( "CREATE INDEX ON :Literal(value)" ) except Exception as e: - print(e, flush=True) + logger.warning(f"Index create failure: {e}") # Maybe index already exists - print("Index create failure ignored", flush=True) + logger.warning("Index create failure ignored") - print("Index creation done", flush=True) + logger.info("Index creation done") def create_node(self, uri): - print("Create node", uri) + logger.debug(f"Create node {uri}") summary = self.io.execute_query( "MERGE (n:Node {uri: $uri})", @@ -105,14 +109,14 @@ class Processor(TriplesStoreService): database_=self.db, ).summary - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=summary.counters.nodes_created, time=summary.result_available_after )) def create_literal(self, value): - print("Create literal", value) + logger.debug(f"Create literal {value}") summary = self.io.execute_query( "MERGE (n:Literal {value: $value})", @@ -120,14 +124,14 @@ class Processor(TriplesStoreService): database_=self.db, ).summary - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=summary.counters.nodes_created, time=summary.result_available_after )) def relate_node(self, src, uri, dest): - print("Create node rel", src, uri, dest) + logger.debug(f"Create node rel {src} {uri} {dest}") summary = self.io.execute_query( "MATCH (src:Node {uri: $src}) " @@ -137,14 +141,14 @@ class Processor(TriplesStoreService): database_=self.db, ).summary - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=summary.counters.nodes_created, time=summary.result_available_after )) def relate_literal(self, src, uri, dest): - print("Create literal rel", src, uri, dest) + logger.debug(f"Create literal rel {src} {uri} {dest}") summary = self.io.execute_query( "MATCH (src:Node {uri: $src}) " @@ -154,7 +158,7 @@ class Processor(TriplesStoreService): database_=self.db, ).summary - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=summary.counters.nodes_created, time=summary.result_available_after )) diff --git a/trustgraph-flow/trustgraph/storage/triples/neo4j/write.py b/trustgraph-flow/trustgraph/storage/triples/neo4j/write.py index 5293ee1e..e1913c14 100755 --- a/trustgraph-flow/trustgraph/storage/triples/neo4j/write.py +++ b/trustgraph-flow/trustgraph/storage/triples/neo4j/write.py @@ -8,10 +8,14 @@ import base64 import os import argparse import time +import logging from neo4j import GraphDatabase from .... base import TriplesStoreService +# Module logger +logger = logging.getLogger(__name__) + default_ident = "triples-write" default_graph_host = 'bolt://neo4j:7687' @@ -55,40 +59,40 @@ class Processor(TriplesStoreService): # and this process will restart several times until Pulsar arrives, # so should be safe - print("Create indexes...", flush=True) + logger.info("Create indexes...") try: session.run( "CREATE INDEX Node_uri FOR (n:Node) ON (n.uri)", ) except Exception as e: - print(e, flush=True) + logger.warning(f"Index create failure: {e}") # Maybe index already exists - print("Index create failure ignored", flush=True) + logger.warning("Index create failure ignored") try: session.run( "CREATE INDEX Literal_value FOR (n:Literal) ON (n.value)", ) except Exception as e: - print(e, flush=True) + logger.warning(f"Index create failure: {e}") # Maybe index already exists - print("Index create failure ignored", flush=True) + logger.warning("Index create failure ignored") try: session.run( "CREATE INDEX Rel_uri FOR ()-[r:Rel]-() ON (r.uri)", ) except Exception as e: - print(e, flush=True) + logger.warning(f"Index create failure: {e}") # Maybe index already exists - print("Index create failure ignored", flush=True) + logger.warning("Index create failure ignored") - print("Index creation done", flush=True) + logger.info("Index creation done") def create_node(self, uri): - print("Create node", uri) + logger.debug(f"Create node {uri}") summary = self.io.execute_query( "MERGE (n:Node {uri: $uri})", @@ -96,14 +100,14 @@ class Processor(TriplesStoreService): database_=self.db, ).summary - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=summary.counters.nodes_created, time=summary.result_available_after )) def create_literal(self, value): - print("Create literal", value) + logger.debug(f"Create literal {value}") summary = self.io.execute_query( "MERGE (n:Literal {value: $value})", @@ -111,14 +115,14 @@ class Processor(TriplesStoreService): database_=self.db, ).summary - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=summary.counters.nodes_created, time=summary.result_available_after )) def relate_node(self, src, uri, dest): - print("Create node rel", src, uri, dest) + logger.debug(f"Create node rel {src} {uri} {dest}") summary = self.io.execute_query( "MATCH (src:Node {uri: $src}) " @@ -128,14 +132,14 @@ class Processor(TriplesStoreService): database_=self.db, ).summary - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=summary.counters.nodes_created, time=summary.result_available_after )) def relate_literal(self, src, uri, dest): - print("Create literal rel", src, uri, dest) + logger.debug(f"Create literal rel {src} {uri} {dest}") summary = self.io.execute_query( "MATCH (src:Node {uri: $src}) " @@ -145,7 +149,7 @@ class Processor(TriplesStoreService): database_=self.db, ).summary - print("Created {nodes_created} nodes in {time} ms.".format( + logger.debug("Created {nodes_created} nodes in {time} ms.".format( nodes_created=summary.counters.nodes_created, time=summary.result_available_after ))