Explorer
/tmp/p2za_parallelize.py
← Zurück ↓ Download
#!/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)')