diff --git a/trustgraph-flow/trustgraph/extract/kg/agent/extract.py b/trustgraph-flow/trustgraph/extract/kg/agent/extract.py index af2679e2..6f49bc27 100644 --- a/trustgraph-flow/trustgraph/extract/kg/agent/extract.py +++ b/trustgraph-flow/trustgraph/extract/kg/agent/extract.py @@ -96,17 +96,17 @@ class Processor(FlowProcessor): def to_uri(self, text): return TRUSTGRAPH_ENTITIES + urllib.parse.quote(text) - def emit_triples(self, pub, metadata, triples): + async def emit_triples(self, pub, metadata, triples): tpls = Triples() tpls.metadata = metadata tpls.triples = triples - pub.publish(tpls) + await pub.send(tpls) - def emit_entity_contexts(self, pub, metadata, entity_contexts): + async def emit_entity_contexts(self, pub, metadata, entity_contexts): ecs = EntityContexts() ecs.metadata = metadata - ecs.entity_contexts = entity_contexts - pub.publish(ecs) + ecs.entities = entity_contexts + await pub.send(ecs) def parse_json(self, text): json_match = re.search(r'```(?:json)?(.*?)```', text, re.DOTALL) @@ -162,14 +162,11 @@ class Processor(FlowProcessor): question = prompt ) - if agent_response.error: - raise RuntimeError(f"Agent error: {agent_response.error}") - print("response:", agent_response, flush=True) # Parse JSON response try: - extraction_data = self.parse_json(agent_response.answer) + extraction_data = self.parse_json(agent_response) except json.JSONDecodeError as e: raise ValueError(f"Invalid JSON response from agent: {e}") @@ -180,12 +177,12 @@ class Processor(FlowProcessor): # Emit outputs if triples: - self.emit_triples(flow("triples"), msg.metadata, triples) + await self.emit_triples(flow("triples"), v.metadata, triples) if entity_contexts: - self.emit_entity_contexts( + await self.emit_entity_contexts( flow("entity-contexts"), - msg.metadata, + v.metadata, entity_contexts ) @@ -200,91 +197,98 @@ class Processor(FlowProcessor): # Process definitions for defn in data.get("definitions", []): + entity_uri = self.to_uri(defn["entity"]) # Add entity label triples.append(Triple( - subject=entity_uri, - predicate=RDF_LABEL, - object=defn["entity"] + subject=Value(value=entity_uri, is_uri=True), + predicate=Value(value=RDF_LABEL, is_uri=True), + object=Value(value=defn["entity"], is_uri=False), )) # Add definition triples.append(Triple( - subject=entity_uri, - predicate=DEFINITION, - object=defn["definition"] + subject=Value(value=entity_uri, is_uri=True), + predicate=Value(value=DEFINITION, is_uri=True), + object=Value(value=defn["definition"], is_uri=False), )) # Add subject-of relationship to document if metadata.id: triples.append(Triple( - subject=entity_uri, - predicate=SUBJECT_OF, - object=metadata.id + subject=Value(value=entity_uri, is_uri=True), + predicate=Value(value=SUBJECT_OF, is_uri=True), + object=metadata.id, )) # Create entity context for embeddings entity_contexts.append(EntityContext( - entity=entity_uri, - context=defn["context"] + entity=Value(value=entity_uri, is_uri=True), + context=defn["definition"] )) # Process relationships for rel in data.get("relationships", []): + subject_uri = self.to_uri(rel["subject"]) predicate_uri = self.to_uri(rel["predicate"]) + + subject_value = Value(value=subject_uri, is_uri=True) + predicate_value = Value(value=predicate_uri, is_uri=True) + if data.get("object-entity", False): + object_value = Value(value=predicate_uri, is_uri=True) + else: + object_value = Value(value=predicate_uri, is_uri=False) # Add subject and predicate labels triples.append(Triple( - subject=subject_uri, - predicate=RDF_LABEL, - object=rel["subject"] + subject=subject_value, + predicate=Value(value=RDF_LABEL, is_uri=True), + object=Value(value=rel["subject"], is_uri=False), )) triples.append(Triple( - subject=predicate_uri, - predicate=RDF_LABEL, - object=rel["predicate"] + subject=predicate_value, + predicate=Value(value=RDF_LABEL, is_uri=True), + object=Value(value=rel["predicate"], is_uri=False), )) # Handle object (entity vs literal) if rel.get("object-entity", True): object_uri = self.to_uri(rel["object"]) triples.append(Triple( - subject=object_uri, - predicate=RDF_LABEL, - object=rel["object"] + subject=object_value, + predicate=Value(value=RDF_LABEL, is_uri=True), + object=rel["object"], )) - else: - object_uri = rel["object"] # Keep as literal # Add the main relationship triple triples.append(Triple( - subject=subject_uri, - predicate=predicate_uri, - object=object_uri + subject=subject_value, + predicate=predicate_value, + object=object_value )) # Add subject-of relationships to document if metadata.id: triples.append(Triple( - subject=subject_uri, - predicate=SUBJECT_OF, - object=metadata.id + subject=subject_value, + predicate=Value(value=SUBJECT_OF, is_uri=True), + object=metadata.id, )) triples.append(Triple( - subject=predicate_uri, - predicate=SUBJECT_OF, - object=metadata.id + subject=predicate_value, + predicate=Value(SUBJECT_OF, is_uri=True), + object=metadata.id, )) if rel.get("object-entity", True): triples.append(Triple( - subject=object_uri, - predicate=SUBJECT_OF, - object=metadata.id + subject=object_value, + predicate=Value(SUBJECT_OF, is_uri=True), + object=metadata.id, )) return triples, entity_contexts