#!/usr/bin/env python3
"""
P2Z-A Optimierung 1: Parallelisiere die drei LLM-Operationen pro Edge
(dedupe / dates / contradictions) wieder — sie sind unabhängig voneinander
(dedupe liefert resolved_edge, dates/contradictions arbeiten auf dem
ORIGINALEN extracted_edge und deren Ergebnisse werden nach dem Gather nur
kombiniert). Das sequential-Patch hatte den originalen asyncio.gather
durch serielle Aufrufe ersetzt — das ist der Hauptgrund für 40-50+ LLM-Calls
mit langer Wall-Clock-Zeit.
Wir patchen NUR den install_sequential_graphiti_path in graphiti_worker.py:
statt der seriellen Kette wieder asyncio.gather (Original-Graphiti-Verhalten),
behalten aber die Progress-Instrumentation bei.
Zusätzlich: entity_resolution parallelisieren (resolve_extracted_node pro
Node mit asyncio.gather statt for-Schleife).
"""
import subprocess
r = subprocess.run(['docker', 'exec', 'graphiti-service', 'cat', '/app/graphiti_worker.py'],
capture_output=True, text=True)
src = r.stdout
# --- Fix 1: edge operations parallel (restore gather, keep instrumentation) ---
old_edge = ''' new_edge_sequential = """resolved_edge = await dedupe_extracted_edge(llm_client, extracted_edge, related_edges)
valid_at, invalid_at = await extract_edge_dates(
llm_client, extracted_edge, current_episode, previous_episodes
)
invalidation_candidates = await get_edge_contradictions(
llm_client, extracted_edge, existing_edges
)"""'''
new_edge = ''' new_edge_sequential = """resolved_edge, (valid_at, invalid_at), invalidation_candidates = await asyncio.gather(
dedupe_extracted_edge(llm_client, extracted_edge, related_edges),
extract_edge_dates(llm_client, extracted_edge, current_episode, previous_episodes),
get_edge_contradictions(llm_client, extracted_edge, existing_edges),
)"""'''
assert old_edge in src, 'edge block not found'
src = src.replace(old_edge, new_edge)
# --- Fix 2: node resolution parallel ---
old_nodes = ''' async def sequential_resolve_extracted_nodes(llm_client, extracted_nodes, existing_nodes_lists):
uuid_map = {}
resolved_nodes = []
results = []
for extracted_node, existing_nodes in zip(extracted_nodes, existing_nodes_lists):
results.append(await original_resolve_extracted_node(llm_client, extracted_node, existing_nodes))
for result in results:
uuid_map.update(result[1])
resolved_nodes.append(result[0])
return resolved_nodes, uuid_map'''
new_nodes = ''' async def sequential_resolve_extracted_nodes(llm_client, extracted_nodes, existing_nodes_lists):
import asyncio as _aio
results = await _aio.gather(*[
original_resolve_extracted_node(llm_client, extracted_node, existing_nodes)
for extracted_node, existing_nodes in zip(extracted_nodes, existing_nodes_lists)
])
uuid_map = {}
resolved_nodes = []
for result in results:
uuid_map.update(result[1])
resolved_nodes.append(result[0])
return resolved_nodes, uuid_map'''
assert old_nodes in src, 'nodes block not found'
src = src.replace(old_nodes, new_nodes)
# write patched file into container and restart it
import subprocess as sp
proc = sp.run(['docker', 'exec', 'graphiti-service', 'cp', '/app/graphiti_worker.py', '/app/graphiti_worker.py.pre-p2za'], capture_output=True)
w = sp.run(['docker', 'exec', '-i', 'graphiti-service', 'sh', '-c',
'cat > /app/graphiti_worker.py'], input=src.encode(), capture_output=True)
assert w.returncode == 0, w.stderr
print('graphiti_worker.py patched (edge ops + node resolution parallelisiert)')