Browse Source

feat(装配): 接入真实输入源和本地持久化输出

生产 Host 启动时按配置建立旧库 SOCKS 隧道,启动失败会关闭隧道,正常 shutdown 同样释放进程和数据库资源。

根据配置装配本地向量 decode、旧版知识模型、小红书与知乎检索适配器;保留现有 HTTP/File Adapter 作为显式回退路径。

内部 test 环境将新任务的 list、detail 和 log 读取指向本地 runtime 数据库,并关闭两张未创建的历史表查询;生产环境继续保留只读库和历史兼容行为。

继续强制 Runner 提供可重开验证的持久 TraceStore,同时将 Task Ledger、Framework Artifact、Raw Artifact、Task Contract 和 Owner Lease 绑定到统一 agent_data_root。

更新生产装配测试,确认持久化 Task Ledger 使用预期本地目录。
SamLee 1 ngày trước cách đây
mục cha
commit
f24a87cd83

+ 105 - 20
script_build_host/src/script_build_host/production.py

@@ -8,6 +8,7 @@ from typing import Any
 
 import httpx
 from agent import FileSystemArtifactStore, FileSystemTaskStore, FileSystemTraceStore
+from sqlalchemy.engine import make_url
 
 from script_build_host.adapters import (
     DatabaseFirstPromptSource,
@@ -16,6 +17,10 @@ from script_build_host.adapters import (
     FileKnowledgeRetrievalAdapter,
     FilePersonaSource,
     HttpDecodeRetrievalAdapter,
+    LegacyExternalRetrievalAdapter,
+    LegacyLlmKnowledgeRetrievalAdapter,
+    LegacyOpenRouterClient,
+    LegacyVectorDecodeRetrievalAdapter,
     PatternRetrievalAdapter,
     SafeHttpClient,
     SafeImageAdapter,
@@ -41,6 +46,7 @@ from script_build_host.infrastructure.manifests import SettingsRuntimeManifestPr
 from script_build_host.infrastructure.outbound import OutboundPolicy
 from script_build_host.infrastructure.ownership import FencedCommandGate, OwnerLease
 from script_build_host.infrastructure.raw_artifacts import FileRawArtifactStore
+from script_build_host.infrastructure.socks_tunnel import ManagedSocks4aTunnel
 from script_build_host.infrastructure.trace_verifier import FileSystemTraceStoreVerifier
 from script_build_host.repositories import (
     LegacySqlAlchemyInputReader,
@@ -74,17 +80,25 @@ class ProductionHost:
     runtime_manifest: SettingsRuntimeManifestProvider
     trace_store: Any
     trace_store_verifier: DurableTraceStoreVerifier
+    read_database_tunnel: ManagedSocks4aTunnel | None = None
     require_trace_attachments: bool = False
     _closed: bool = False
 
     async def validate_startup(self) -> None:
         """Resolve configured endpoints and sources before serving traffic."""
 
-        await self.runtime_manifest.load()
-        await self.trace_store_verifier.verify(
-            self.trace_store,
-            require_attachments=self.require_trace_attachments,
-        )
+        if self.read_database_tunnel is not None:
+            await self.read_database_tunnel.start()
+        try:
+            await self.runtime_manifest.load()
+            await self.trace_store_verifier.verify(
+                self.trace_store,
+                require_attachments=self.require_trace_attachments,
+            )
+        except Exception:
+            if self.read_database_tunnel is not None:
+                await self.read_database_tunnel.stop()
+            raise
 
     async def close(self) -> None:
         if self._closed:
@@ -92,6 +106,8 @@ class ProductionHost:
         self._closed = True
         await self.http_client.aclose()
         await self.database.dispose()
+        if self.read_database_tunnel is not None:
+            await self.read_database_tunnel.stop()
 
 
 def compose_production_host(
@@ -113,16 +129,40 @@ def compose_production_host(
             )
 
     decode_path: Path | None = None
-    if not settings.decode_endpoint:
+    vector_decode = (
+        not settings.decode_endpoint
+        and settings.decode_raw_root is not None
+        and (settings.decode_index_root / "manifest.json").is_file()
+        and (settings.decode_index_root / "meta.jsonl").is_file()
+        and (settings.decode_index_root / "embeddings.npy").is_file()
+    )
+    if not settings.decode_endpoint and not vector_decode:
         if settings.environment.lower() != "test":
-            raise RuntimeError("SCRIPT_BUILD_DECODE_ENDPOINT is required in production")
+            raise RuntimeError(
+                "SCRIPT_BUILD_DECODE_ENDPOINT or a complete local vector index is required"
+            )
         decode_path = _one_json_source(settings.decode_index_root, "decode")
     knowledge_path = _one_json_source(settings.knowledge_root, "knowledge")
 
+    read_database_tunnel: ManagedSocks4aTunnel | None = None
+    proxy = settings.read_socks_proxy()
+    if proxy is not None:
+        direct_read_url = make_url(settings.direct_read_dsn())
+        if not direct_read_url.host:
+            raise RuntimeError("the direct read database URL has no host")
+        read_database_tunnel = ManagedSocks4aTunnel(
+            proxy_host=proxy[0],
+            proxy_port=proxy[1],
+            target_host=direct_read_url.host,
+            target_port=direct_read_url.port or 3306,
+            listen_host=settings.read_database_tunnel_host,
+            listen_port=settings.read_database_tunnel_port,
+        )
     database = create_database_sessions(settings)
     outbound_policy = OutboundPolicy(
-        frozenset(settings.outbound_allowed_hosts),
-        frozenset(settings.outbound_allowed_ports),
+        allowed_hosts=frozenset(settings.outbound_allowed_hosts),
+        allowed_ports=frozenset(settings.outbound_allowed_ports),
+        allowed_http_hosts=frozenset(settings.outbound_allowed_http_hosts),
     )
     http_client = httpx.AsyncClient(
         transport=runtime.http_transport,
@@ -134,6 +174,16 @@ def compose_production_host(
         timeout_seconds=settings.request_timeout_seconds,
         max_response_bytes=settings.max_response_bytes,
     )
+    openrouter: LegacyOpenRouterClient | None = None
+    if settings.openrouter_api_key is not None and settings.embedding_endpoint:
+        openrouter = LegacyOpenRouterClient(
+            safe_http,
+            api_key=settings.openrouter_api_key.get_secret_value(),
+            embedding_endpoint=settings.embedding_endpoint,
+            chat_endpoint=settings.openrouter_chat_endpoint,
+            embedding_model=settings.embedding_model,
+            embedding_dimension=settings.embedding_dimension,
+        )
 
     retrieval_adapters: dict[str, Any] = {}
     if settings.pattern_endpoint:
@@ -148,16 +198,44 @@ def compose_production_host(
             settings.decode_endpoint,
             safe_http,
         )
+    elif vector_decode and settings.decode_raw_root is not None:
+        retrieval_adapters["decode"] = LegacyVectorDecodeRetrievalAdapter(
+            index_root=settings.decode_index_root,
+            raw_root=settings.decode_raw_root,
+            openrouter=openrouter,
+            min_score=settings.decode_min_score,
+        )
     elif decode_path is not None:
         retrieval_adapters["decode"] = FileDecodeRetrievalAdapter(decode_path)
     else:
         raise AssertionError("decode configuration was validated before resource creation")
-    if settings.external_endpoint:
+    if (
+        settings.xhs_search_endpoint
+        and settings.xhs_detail_endpoint
+        and settings.zhihu_search_endpoint
+    ):
+        retrieval_adapters["external"] = LegacyExternalRetrievalAdapter(
+            xhs_search_endpoint=settings.xhs_search_endpoint,
+            xhs_detail_endpoint=settings.xhs_detail_endpoint,
+            zhihu_search_endpoint=settings.zhihu_search_endpoint,
+            http=safe_http,
+            openrouter=openrouter,
+            image_model=settings.external_image_model,
+        )
+    elif settings.external_endpoint:
         retrieval_adapters["external"] = ExternalRetrievalAdapter(
             settings.external_endpoint,
             safe_http,
         )
-    retrieval_adapters["knowledge"] = FileKnowledgeRetrievalAdapter(knowledge_path)
+    retrieval_adapters["knowledge"] = (
+        LegacyLlmKnowledgeRetrievalAdapter(
+            knowledge_path,
+            openrouter=openrouter,
+            model=settings.knowledge_model,
+        )
+        if openrouter is not None
+        else FileKnowledgeRetrievalAdapter(knowledge_path)
+    )
 
     snapshots = SqlAlchemyInputSnapshotRepository(database.write)
     bindings = SqlAlchemyMissionBindingRepository(database.write)
@@ -187,14 +265,17 @@ def compose_production_host(
         knowledge_path=knowledge_path,
     )
     data_root = settings.agent_data_root.resolve()
+    fresh_local_outputs = settings.environment.lower() == "test"
+    output_reads = database.write if fresh_local_outputs else database.read
     raw_artifacts = FileRawArtifactStore(data_root)
     candidate_workspaces = SqlAlchemyCandidateWorkspaceRepository(database.write, artifacts)
     legacy_projection = LegacyDetailProjectionService(
-        database.read,
+        output_reads,
         runtime.build_authorizer,
         topic_reader=legacy_input,
         bindings=bindings,
         snapshots=snapshots,
+        include_legacy_decision_traces=not fresh_local_outputs,
     )
     final_uow = SqlAlchemyFinalPublicationUnitOfWork(
         database.final,
@@ -257,19 +338,23 @@ def compose_production_host(
             final_publications=final_publications,
             legacy_projection=legacy_projection,
             legacy_api=LegacyScriptBuildApiService(
-                database.read, database.write, runtime.build_authorizer
+                output_reads,
+                database.write,
+                runtime.build_authorizer,
+                include_legacy_logs=not fresh_local_outputs,
             ),
             http_command_journal=HttpCommandJournal(database.write),
         )
     )
     host = ProductionHost(
-        composition,
-        database,
-        http_client,
-        runtime_manifest,
-        trace_store,
-        trace_store_verifier,
-        settings.enable_image_understanding,
+        composition=composition,
+        database=database,
+        http_client=http_client,
+        runtime_manifest=runtime_manifest,
+        trace_store=trace_store,
+        trace_store_verifier=trace_store_verifier,
+        read_database_tunnel=read_database_tunnel,
+        require_trace_attachments=settings.enable_image_understanding,
     )
     composition.app.state.production_host = host
     composition.app.router.add_event_handler("startup", host.validate_startup)

+ 1 - 0
script_build_host/tests/test_production_composition.py

@@ -51,6 +51,7 @@ def _settings(tmp_path: Path, **overrides: object) -> ScriptBuildSettings:
     values: dict[str, object] = {
         "environment": "test",
         "read_database_url": f"sqlite+aiosqlite:///{tmp_path / 'read.db'}",
+        "read_database_socks_proxy": None,
         "write_database_url": f"sqlite+aiosqlite:///{tmp_path / 'write.db'}",
         "agent_data_root": tmp_path / "agent-data",
         "persona_root": tmp_path / "persona",