feat: wire neo4j projection into worker

This commit is contained in:
张胜宇
2026-09-11 20:19:34 +08:00
parent 6c605ec3c1
commit 5e848f5690
2 changed files with 24 additions and 1 deletions
@@ -11,7 +11,7 @@ from app.core.conversation_privacy import sanitize_customer_service_message
class Neo4jQueryDriver(Protocol):
async def execute_query(self, query: str, **parameters: Any) -> Any: ...
async def execute_query(self, *args: Any, **kwargs: Any) -> Any: ...
@dataclass(frozen=True)
+23
View File
@@ -1,9 +1,12 @@
import argparse
import asyncio
import logging
from typing import Any
from app.core.config import get_settings
from app.infrastructure.db import engine
from app.infrastructure.neo4j_profile_projection import Neo4jProfileProjection
from app.worker.memory_sync_outbox_worker import MemorySyncOutboxWorker
from app.worker.offsite_mail_worker import OffsiteMailWorker
from app.worker.runtime import WorkerRuntime
@@ -14,10 +17,28 @@ async def serve(*, once: bool = False) -> None:
settings = get_settings()
runtime = WorkerRuntime(settings=settings)
offsite_worker = OffsiteMailWorker(settings)
neo4j_driver: Any | None = None
memory_sync_worker: MemorySyncOutboxWorker | None = None
if settings.neo4j_password:
from neo4j import AsyncGraphDatabase
neo4j_driver = AsyncGraphDatabase.driver(
settings.neo4j_uri,
auth=(settings.neo4j_username, settings.neo4j_password),
)
memory_sync_worker = MemorySyncOutboxWorker({
"neo4j": Neo4jProfileProjection(neo4j_driver).upsert,
})
else:
logger.warning("Neo4j password not configured; profile projection remains pending")
try:
while True:
try:
worked = await runtime.run_once()
if memory_sync_worker is not None:
worked = await memory_sync_worker.run_once(
target_store="neo4j"
) or worked
worked = await offsite_worker.run_once() or worked
except Exception:
# 常驻 Worker 不能因为"某一轮"的异常就整体退出:数据库抖动、
@@ -35,6 +56,8 @@ async def serve(*, once: bool = False) -> None:
await asyncio.sleep(settings.worker_poll_seconds)
finally:
await offsite_worker.close()
if neo4j_driver is not None:
await neo4j_driver.close()
await engine.dispose()