diff --git a/tests/test-load-pdf b/tests/test-load-pdf index c57ebcc1..838a57ce 100755 --- a/tests/test-load-pdf +++ b/tests/test-load-pdf @@ -9,7 +9,7 @@ from trustgraph.schema import Document, Metadata client = pulsar.Client("pulsar://localhost:6650", listener_name="localhost") prod = client.create_producer( - topic="persistent://tg/flow/document-load:0002", + topic="persistent://tg/flow/document-load:0000", schema=JsonSchema(Document), chunking_enabled=True, ) diff --git a/tests/test-load-text b/tests/test-load-text index 754458aa..83006c6d 100755 --- a/tests/test-load-text +++ b/tests/test-load-text @@ -14,7 +14,7 @@ prod = client.create_producer( chunking_enabled=True, ) -path = "docs/README.cats" +path = "../trustgraph/docs/README.cats" with open(path, "r") as f: # blob = base64.b64encode(f.read()).decode("utf-8") diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/manager.py b/trustgraph-flow/trustgraph/gateway/dispatch/manager.py index 9f83b3ca..66095e67 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/manager.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/manager.py @@ -25,7 +25,7 @@ request_response_dispatchers = { } receive_dispatchers = { - "embeddings": TriplesStream, + "triples-store": TriplesStream, } class TestDispatcher: @@ -127,6 +127,9 @@ class DispatcherManager: flow = params.get("flow") kind = params.get("kind") + # FIXME: What?!?! + kind += "-store" + if flow not in self.flows: raise RuntimeError("Invalid flow") @@ -150,7 +153,7 @@ class DispatcherManager: ws = ws, running = running, # FIXME! - queue = qconfig["response"], + queue = qconfig, consumer = f"api-gateway-{flow}-{kind}-request", subscriber = f"api-gateway-{flow}-{kind}-request", ) diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/triples_stream.py b/trustgraph-flow/trustgraph/gateway/dispatch/triples_stream.py index 66d61930..92e89cbe 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/triples_stream.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/triples_stream.py @@ -3,7 +3,7 @@ import asyncio import queue import uuid -from ... schema import Triples, EmbeddingsResponse +from ... schema import Triples from ... base import Subscriber from . serialize import serialize_triples @@ -37,8 +37,7 @@ class TriplesStream: subs = Subscriber( client = self.pulsar_client, topic = self.queue, consumer_name = self.consumer, subscription = self.subscriber, -# schema = Triples - schema = EmbeddingsResponse + schema = Triples ) await subs.start() @@ -50,12 +49,7 @@ class TriplesStream: try: resp = await asyncio.wait_for(q.get(), timeout=0.5) -# await self.ws.send_json(serialize_triples(resp)) - - print("GOT MESSAGE!!!") - - - await self.ws.send_json(str(resp)) + await self.ws.send_json(serialize_triples(resp)) except TimeoutError: continue