mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-24 12:41:02 +02:00
Logging strategy updates
This commit is contained in:
parent
22f3f34511
commit
1359eb8a3d
8 changed files with 55 additions and 21 deletions
|
|
@ -4,12 +4,16 @@ Agent manager service completion base class
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import time
|
import time
|
||||||
|
import logging
|
||||||
from prometheus_client import Histogram
|
from prometheus_client import Histogram
|
||||||
|
|
||||||
from .. schema import AgentRequest, AgentResponse, Error
|
from .. schema import AgentRequest, AgentResponse, Error
|
||||||
from .. exceptions import TooManyRequests
|
from .. exceptions import TooManyRequests
|
||||||
from .. base import FlowProcessor, ConsumerSpec, ProducerSpec
|
from .. base import FlowProcessor, ConsumerSpec, ProducerSpec
|
||||||
|
|
||||||
|
# Module logger
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
default_ident = "agent-manager"
|
default_ident = "agent-manager"
|
||||||
|
|
||||||
class AgentService(FlowProcessor):
|
class AgentService(FlowProcessor):
|
||||||
|
|
@ -76,9 +80,9 @@ class AgentService(FlowProcessor):
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
|
||||||
# Apart from rate limits, treat all exceptions as unrecoverable
|
# Apart from rate limits, treat all exceptions as unrecoverable
|
||||||
print(f"on_request Exception: {e}")
|
logger.error(f"Exception in agent service on_request: {e}", exc_info=True)
|
||||||
|
|
||||||
print("Send error response...", flush=True)
|
logger.info("Sending error response...")
|
||||||
|
|
||||||
await flow.producer["response"].send(
|
await flow.producer["response"].send(
|
||||||
AgentResponse(
|
AgentResponse(
|
||||||
|
|
|
||||||
|
|
@ -1,8 +1,13 @@
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
from . request_response_spec import RequestResponse, RequestResponseSpec
|
from . request_response_spec import RequestResponse, RequestResponseSpec
|
||||||
from .. schema import DocumentEmbeddingsRequest, DocumentEmbeddingsResponse
|
from .. schema import DocumentEmbeddingsRequest, DocumentEmbeddingsResponse
|
||||||
from .. knowledge import Uri, Literal
|
from .. knowledge import Uri, Literal
|
||||||
|
|
||||||
|
# Module logger
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
class DocumentEmbeddingsClient(RequestResponse):
|
class DocumentEmbeddingsClient(RequestResponse):
|
||||||
async def query(self, vectors, limit=20, user="trustgraph",
|
async def query(self, vectors, limit=20, user="trustgraph",
|
||||||
collection="default", timeout=30):
|
collection="default", timeout=30):
|
||||||
|
|
@ -17,7 +22,7 @@ class DocumentEmbeddingsClient(RequestResponse):
|
||||||
timeout=timeout
|
timeout=timeout
|
||||||
)
|
)
|
||||||
|
|
||||||
print(resp, flush=True)
|
logger.debug(f"Document embeddings response: {resp}")
|
||||||
|
|
||||||
if resp.error:
|
if resp.error:
|
||||||
raise RuntimeError(resp.error.message)
|
raise RuntimeError(resp.error.message)
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,8 @@ Document embeddings query service. Input is vectors. Output is list of
|
||||||
embeddings.
|
embeddings.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
from .. schema import DocumentEmbeddingsRequest, DocumentEmbeddingsResponse
|
from .. schema import DocumentEmbeddingsRequest, DocumentEmbeddingsResponse
|
||||||
from .. schema import Error, Value
|
from .. schema import Error, Value
|
||||||
|
|
||||||
|
|
@ -11,6 +13,9 @@ from . flow_processor import FlowProcessor
|
||||||
from . consumer_spec import ConsumerSpec
|
from . consumer_spec import ConsumerSpec
|
||||||
from . producer_spec import ProducerSpec
|
from . producer_spec import ProducerSpec
|
||||||
|
|
||||||
|
# Module logger
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
default_ident = "ge-query"
|
default_ident = "ge-query"
|
||||||
|
|
||||||
class DocumentEmbeddingsQueryService(FlowProcessor):
|
class DocumentEmbeddingsQueryService(FlowProcessor):
|
||||||
|
|
@ -47,21 +52,21 @@ class DocumentEmbeddingsQueryService(FlowProcessor):
|
||||||
# Sender-produced ID
|
# Sender-produced ID
|
||||||
id = msg.properties()["id"]
|
id = msg.properties()["id"]
|
||||||
|
|
||||||
print(f"Handling input {id}...", flush=True)
|
logger.debug(f"Handling document embeddings query request {id}...")
|
||||||
|
|
||||||
docs = await self.query_document_embeddings(request)
|
docs = await self.query_document_embeddings(request)
|
||||||
|
|
||||||
print("Send response...", flush=True)
|
logger.debug("Sending document embeddings query response...")
|
||||||
r = DocumentEmbeddingsResponse(documents=docs, error=None)
|
r = DocumentEmbeddingsResponse(documents=docs, error=None)
|
||||||
await flow("response").send(r, properties={"id": id})
|
await flow("response").send(r, properties={"id": id})
|
||||||
|
|
||||||
print("Done.", flush=True)
|
logger.debug("Document embeddings query request completed")
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
|
||||||
print(f"Exception: {e}")
|
logger.error(f"Exception in document embeddings query service: {e}", exc_info=True)
|
||||||
|
|
||||||
print("Send error response...", flush=True)
|
logger.info("Sending error response...")
|
||||||
|
|
||||||
r = DocumentEmbeddingsResponse(
|
r = DocumentEmbeddingsResponse(
|
||||||
error=Error(
|
error=Error(
|
||||||
|
|
|
||||||
|
|
@ -3,10 +3,15 @@
|
||||||
Document embeddings store base class
|
Document embeddings store base class
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
from .. schema import DocumentEmbeddings
|
from .. schema import DocumentEmbeddings
|
||||||
from .. base import FlowProcessor, ConsumerSpec
|
from .. base import FlowProcessor, ConsumerSpec
|
||||||
from .. exceptions import TooManyRequests
|
from .. exceptions import TooManyRequests
|
||||||
|
|
||||||
|
# Module logger
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
default_ident = "document-embeddings-write"
|
default_ident = "document-embeddings-write"
|
||||||
|
|
||||||
class DocumentEmbeddingsStoreService(FlowProcessor):
|
class DocumentEmbeddingsStoreService(FlowProcessor):
|
||||||
|
|
@ -40,7 +45,7 @@ class DocumentEmbeddingsStoreService(FlowProcessor):
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
|
||||||
print(f"Exception: {e}")
|
logger.error(f"Exception in document embeddings store service: {e}", exc_info=True)
|
||||||
raise e
|
raise e
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
|
|
|
||||||
|
|
@ -1,8 +1,13 @@
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
from . request_response_spec import RequestResponse, RequestResponseSpec
|
from . request_response_spec import RequestResponse, RequestResponseSpec
|
||||||
from .. schema import GraphEmbeddingsRequest, GraphEmbeddingsResponse
|
from .. schema import GraphEmbeddingsRequest, GraphEmbeddingsResponse
|
||||||
from .. knowledge import Uri, Literal
|
from .. knowledge import Uri, Literal
|
||||||
|
|
||||||
|
# Module logger
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
def to_value(x):
|
def to_value(x):
|
||||||
if x.is_uri: return Uri(x.value)
|
if x.is_uri: return Uri(x.value)
|
||||||
return Literal(x.value)
|
return Literal(x.value)
|
||||||
|
|
@ -21,7 +26,7 @@ class GraphEmbeddingsClient(RequestResponse):
|
||||||
timeout=timeout
|
timeout=timeout
|
||||||
)
|
)
|
||||||
|
|
||||||
print(resp, flush=True)
|
logger.debug(f"Graph embeddings response: {resp}")
|
||||||
|
|
||||||
if resp.error:
|
if resp.error:
|
||||||
raise RuntimeError(resp.error.message)
|
raise RuntimeError(resp.error.message)
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,8 @@ Graph embeddings query service. Input is vectors. Output is list of
|
||||||
embeddings.
|
embeddings.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
from .. schema import GraphEmbeddingsRequest, GraphEmbeddingsResponse
|
from .. schema import GraphEmbeddingsRequest, GraphEmbeddingsResponse
|
||||||
from .. schema import Error, Value
|
from .. schema import Error, Value
|
||||||
|
|
||||||
|
|
@ -11,6 +13,9 @@ from . flow_processor import FlowProcessor
|
||||||
from . consumer_spec import ConsumerSpec
|
from . consumer_spec import ConsumerSpec
|
||||||
from . producer_spec import ProducerSpec
|
from . producer_spec import ProducerSpec
|
||||||
|
|
||||||
|
# Module logger
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
default_ident = "ge-query"
|
default_ident = "ge-query"
|
||||||
|
|
||||||
class GraphEmbeddingsQueryService(FlowProcessor):
|
class GraphEmbeddingsQueryService(FlowProcessor):
|
||||||
|
|
@ -47,21 +52,21 @@ class GraphEmbeddingsQueryService(FlowProcessor):
|
||||||
# Sender-produced ID
|
# Sender-produced ID
|
||||||
id = msg.properties()["id"]
|
id = msg.properties()["id"]
|
||||||
|
|
||||||
print(f"Handling input {id}...", flush=True)
|
logger.debug(f"Handling graph embeddings query request {id}...")
|
||||||
|
|
||||||
entities = await self.query_graph_embeddings(request)
|
entities = await self.query_graph_embeddings(request)
|
||||||
|
|
||||||
print("Send response...", flush=True)
|
logger.debug("Sending graph embeddings query response...")
|
||||||
r = GraphEmbeddingsResponse(entities=entities, error=None)
|
r = GraphEmbeddingsResponse(entities=entities, error=None)
|
||||||
await flow("response").send(r, properties={"id": id})
|
await flow("response").send(r, properties={"id": id})
|
||||||
|
|
||||||
print("Done.", flush=True)
|
logger.debug("Graph embeddings query request completed")
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
|
||||||
print(f"Exception: {e}")
|
logger.error(f"Exception in graph embeddings query service: {e}", exc_info=True)
|
||||||
|
|
||||||
print("Send error response...", flush=True)
|
logger.info("Sending error response...")
|
||||||
|
|
||||||
r = GraphEmbeddingsResponse(
|
r = GraphEmbeddingsResponse(
|
||||||
error=Error(
|
error=Error(
|
||||||
|
|
|
||||||
|
|
@ -3,10 +3,15 @@
|
||||||
Graph embeddings store base class
|
Graph embeddings store base class
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
from .. schema import GraphEmbeddings
|
from .. schema import GraphEmbeddings
|
||||||
from .. base import FlowProcessor, ConsumerSpec
|
from .. base import FlowProcessor, ConsumerSpec
|
||||||
from .. exceptions import TooManyRequests
|
from .. exceptions import TooManyRequests
|
||||||
|
|
||||||
|
# Module logger
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
default_ident = "graph-embeddings-write"
|
default_ident = "graph-embeddings-write"
|
||||||
|
|
||||||
class GraphEmbeddingsStoreService(FlowProcessor):
|
class GraphEmbeddingsStoreService(FlowProcessor):
|
||||||
|
|
@ -40,7 +45,7 @@ class GraphEmbeddingsStoreService(FlowProcessor):
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
|
||||||
print(f"Exception: {e}")
|
logger.error(f"Exception in graph embeddings store service: {e}", exc_info=True)
|
||||||
raise e
|
raise e
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
|
|
|
||||||
|
|
@ -50,21 +50,21 @@ class TriplesQueryService(FlowProcessor):
|
||||||
# Sender-produced ID
|
# Sender-produced ID
|
||||||
id = msg.properties()["id"]
|
id = msg.properties()["id"]
|
||||||
|
|
||||||
print(f"Handling input {id}...", flush=True)
|
logger.debug(f"Handling triples query request {id}...")
|
||||||
|
|
||||||
triples = await self.query_triples(request)
|
triples = await self.query_triples(request)
|
||||||
|
|
||||||
print("Send response...", flush=True)
|
logger.debug("Sending triples query response...")
|
||||||
r = TriplesQueryResponse(triples=triples, error=None)
|
r = TriplesQueryResponse(triples=triples, error=None)
|
||||||
await flow("response").send(r, properties={"id": id})
|
await flow("response").send(r, properties={"id": id})
|
||||||
|
|
||||||
print("Done.", flush=True)
|
logger.debug("Triples query request completed")
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
|
||||||
print(f"Exception: {e}")
|
logger.error(f"Exception in triples query service: {e}", exc_info=True)
|
||||||
|
|
||||||
print("Send error response...", flush=True)
|
logger.info("Sending error response...")
|
||||||
|
|
||||||
r = TriplesQueryResponse(
|
r = TriplesQueryResponse(
|
||||||
error = Error(
|
error = Error(
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue