已合并
feat: 支持链路 mTLS 证书认证及身份持久化与 AgentServer 证书注入 #523
feat: 支持链路 mTLS 证书认证及身份持久化与 AgentServer 证书注入 #523
已合并
Wal1et创建于 11 天前
46 个文件变更+6010-346
@@ -39,3 +39,11 @@ AGENT_RUNTIME_DEFAULT_NAMESPACE=default
39# AGENT_RUNTIME_EVAL_LLM_MAX_TOKENS=1638439# AGENT_RUNTIME_EVAL_LLM_MAX_TOKENS=16384
40# AGENT_RUNTIME_EVAL_LLM_DISABLE_THINKING=false40# AGENT_RUNTIME_EVAL_LLM_DISABLE_THINKING=false
41# AGENT_RUNTIME_EVAL_POD_BUDGET=041# AGENT_RUNTIME_EVAL_POD_BUDGET=0
42+ 
43+# ---- 内部 HTTP/SSE 链路 mTLS(默认关闭,保持现有 HTTP)----
44+# 未设置或 off 均不启用。仅显式 enforce 启用 HTTPS/mTLS;observe 只做预检。
45+# 通过配套 deploy.sh 部署时,在 .env.custom 设置同名开关,由工具注入各组件。
46+JIUWENSWARM_LINK_MTLS_MODE=off
47+# 证书、绑定 ID/epoch、AgentServer Secret、headless DNS 及挂载路径由部署层管理。
48+# 不需要手工填写 CA/CERT/KEY 路径,也不要通过修改 ID/epoch 实现换证。
49+# 使用说明:docs/zh/内部链路mTLS接入.md
@@ -40,6 +40,7 @@ def main(argv: Sequence[str] | None = None) -> None:
40 from openjiuwen_runtime.service.config import ServiceConfig40 from openjiuwen_runtime.service.config import ServiceConfig
41 41 
42 from .config import AgentRuntimeConfig42 from .config import AgentRuntimeConfig
43+ from .link_mtls import LinkMTLSConfig
43 from .logsetup import configure_logging44 from .logsetup import configure_logging
44 from .main import create_app45 from .main import create_app
45 46 
@@ -49,7 +50,8 @@ def main(argv: Sequence[str] | None = None) -> None:
49 50 
50 settings = ServiceConfig.from_env()51 settings = ServiceConfig.from_env()
51 arc = AgentRuntimeConfig.from_env()52 arc = AgentRuntimeConfig.from_env()
52- application = create_app(settings, arc)53+ link_mtls = LinkMTLSConfig.from_env()
54+ application = create_app(settings, arc, link_mtls_config=link_mtls)
53 55 
54 import uvicorn56 import uvicorn
55 57 
@@ -58,6 +60,7 @@ def main(argv: Sequence[str] | None = None) -> None:
58 host=settings.host,60 host=settings.host,
59 port=settings.port,61 port=settings.port,
60 log_level=os.getenv("AGENT_RUNTIME_LOG_LEVEL", "info").strip().lower(),62 log_level=os.getenv("AGENT_RUNTIME_LOG_LEVEL", "info").strip().lower(),
63+ **link_mtls.uvicorn_ssl_kwargs(),
61 )64 )
62 65 
63 66 
@@ -28,8 +28,6 @@ from openjiuwen_runtime.service.config import ServiceConfig
28 28 
29from . import errors as app_errors29from . import errors as app_errors
30from .config import RM_KEY_PREFIX, SERVICE_PREFIX, SM_KEY_PREFIX, AgentRuntimeConfig30from .config import RM_KEY_PREFIX, SERVICE_PREFIX, SM_KEY_PREFIX, AgentRuntimeConfig
31-from .visualization_api import register_visualization_api
32-from .metrics import MetricsRegistry, request_metrics_middleware
33from .evaluation.collector import (31from .evaluation.collector import (
34 FLUSH_INTERVAL_SEC,32 FLUSH_INTERVAL_SEC,
35 FLUSH_TIMEOUT_SEC,33 FLUSH_TIMEOUT_SEC,
@@ -39,6 +37,12 @@ from .evaluation.collector import (
39from .evaluation.evaluator import Evaluator37from .evaluation.evaluator import Evaluator
40from .evaluation.llm import LLMClient38from .evaluation.llm import LLMClient
41from .evaluation.state import EvaluationState39from .evaluation.state import EvaluationState
40+from .link_binding_state import (
41+ LINK_BINDING_STATE_TABLE_DEF,
42+ sync_local_link_binding_state,
43+)
44+from .link_mtls import LinkMTLSConfig, LinkMTLSError, LinkMTLSMode
45+from .metrics import MetricsRegistry, request_metrics_middleware
42from .resource_manager.facade import ResourceManagerFacade46from .resource_manager.facade import ResourceManagerFacade
43from .resource_manager.k8s import FakeK8sPodClient, RealK8sPodClient47from .resource_manager.k8s import FakeK8sPodClient, RealK8sPodClient
44from .resource_manager.orchestrator import ResourceOrchestrator48from .resource_manager.orchestrator import ResourceOrchestrator
@@ -55,6 +59,7 @@ from .session_manager.handlers import register_handlers
55from .session_manager.orchestrator import SessionOrchestrator59from .session_manager.orchestrator import SessionOrchestrator
56from .session_manager.state import SessionState60from .session_manager.state import SessionState
57from .session_manager.sweeper import SessionSweeper61from .session_manager.sweeper import SessionSweeper
62+from .visualization_api import register_visualization_api
58 63 
59logger = logging.getLogger("agent_runtime")64logger = logging.getLogger("agent_runtime")
60 65 
@@ -83,7 +88,9 @@ TICK_TIMEOUTS = {
83 88 
84 89 
85def build_resources(90def build_resources(
86- settings: ServiceConfig, arc: AgentRuntimeConfig91+ settings: ServiceConfig,
92+ arc: AgentRuntimeConfig,
93+ link_mtls: LinkMTLSConfig | None = None,
87) -> tuple[Any, Any, Any]:94) -> tuple[Any, Any, Any]:
88 """构造共享物理资源(redis client / db handler / k8s client)。95 """构造共享物理资源(redis client / db handler / k8s client)。
89 96 
@@ -98,11 +105,16 @@ def build_resources(
98 105 
99 redis_client = FakeRedis()106 redis_client = FakeRedis()
100 # 文件型 SQLite(:memory: 在连接池下会丢表;local 模式仅供开发调试)107 # 文件型 SQLite(:memory: 在连接池下会丢表;local 模式仅供开发调试)
101- db = SQLiteHandler(os.getenv("AGENT_RUNTIME_SQLITE_PATH", "./agent_runtime_local.db"))108+ db = SQLiteHandler(
109+ os.getenv("AGENT_RUNTIME_SQLITE_PATH", "./agent_runtime_local.db")
110+ )
102 k8s = FakeK8sPodClient(default_namespace=arc.default_namespace)111 k8s = FakeK8sPodClient(default_namespace=arc.default_namespace)
103 return redis_client, db, k8s112 return redis_client, db, k8s
104 113 
105- from openjiuwen_runtime.service.bootstrap import build_db_handler, build_redis_client114+ from openjiuwen_runtime.service.bootstrap import (
115+ build_db_handler,
116+ build_redis_client,
117+ )
106 118 
107 redis_client = build_redis_client(settings)119 redis_client = build_redis_client(settings)
108 db = build_db_handler(settings)120 db = build_db_handler(settings)
@@ -112,7 +124,9 @@ def build_resources(
112 "(OPENJIUWEN_SERVICE_DB_TYPE=mysql|postgresql)"124 "(OPENJIUWEN_SERVICE_DB_TYPE=mysql|postgresql)"
113 )125 )
114 k8s = RealK8sPodClient(126 k8s = RealK8sPodClient(
115- kubeconfig=arc.kubeconfig, default_namespace=arc.default_namespace127+ kubeconfig=arc.kubeconfig,
128+ default_namespace=arc.default_namespace,
129+ link_mtls_config=link_mtls,
116 )130 )
117 return redis_client, db, k8s131 return redis_client, db, k8s
118 132 
@@ -128,6 +142,7 @@ class OrchestratorSystemContext(SystemContext):
128 k8s: Any,142 k8s: Any,
129 settings: ServiceConfig,143 settings: ServiceConfig,
130 arc: AgentRuntimeConfig,144 arc: AgentRuntimeConfig,
145+ link_mtls: LinkMTLSConfig | None = None,
131 instance_id: str | None = None,146 instance_id: str | None = None,
132 owns_resources: bool = True,147 owns_resources: bool = True,
133 ) -> None:148 ) -> None:
@@ -136,9 +151,12 @@ class OrchestratorSystemContext(SystemContext):
136 db=db,151 db=db,
137 settings=settings,152 settings=settings,
138 key_prefix=SM_KEY_PREFIX,153 key_prefix=SM_KEY_PREFIX,
139- table_definitions=[SERVICE_CONFIG_TEMPLATE_TABLE_DEF,154+ table_definitions=[
140- SERVICE_CONFIG_CONTAINER_TABLE_DEF,155+ SERVICE_CONFIG_TEMPLATE_TABLE_DEF,
141- ROUTING_SCOPE_TABLE_DEF],156+ SERVICE_CONFIG_CONTAINER_TABLE_DEF,
157+ ROUTING_SCOPE_TABLE_DEF,
158+ LINK_BINDING_STATE_TABLE_DEF,
159+ ],
142 instance_id=instance_id,160 instance_id=instance_id,
143 _owns_db=owns_resources,161 _owns_db=owns_resources,
144 _owns_redis=owns_resources,162 _owns_redis=owns_resources,
@@ -158,6 +176,7 @@ class OrchestratorSystemContext(SystemContext):
158 _owns_redis=False,176 _owns_redis=False,
159 )177 )
160 self.arc = arc178 self.arc = arc
179+ self.link_mtls = link_mtls or LinkMTLSConfig.from_env()
161 self.k8s = k8s180 self.k8s = k8s
162 self._jobs: list[Any] = []181 self._jobs: list[Any] = []
163 self._bind_modules()182 self._bind_modules()
@@ -175,11 +194,14 @@ class OrchestratorSystemContext(SystemContext):
175 194 
176 self.sm_facade = SessionManagerFacade(sm_state)195 self.sm_facade = SessionManagerFacade(sm_state)
177 196 
178- rm_orchestrator = ResourceOrchestrator(rm_state, self.k8s, telemetry=telemetry)197+ rm_orchestrator = ResourceOrchestrator(
198+ rm_state, self.k8s, telemetry=telemetry, link_mtls_config=self.link_mtls
199+ )
179 self.rm_facade = ResourceManagerFacade(rm_orchestrator)200 self.rm_facade = ResourceManagerFacade(rm_orchestrator)
180 201 
181 self.sm_config_store = ConfigStore(202 self.sm_config_store = ConfigStore(
182- self.db, sm_state,203+ self.db,
204+ sm_state,
183 push_pool_config=self.rm_facade.update_pool_config,205 push_pool_config=self.rm_facade.update_pool_config,
184 known_rm_scopes=self.rm_facade.known_scope_ids,206 known_rm_scopes=self.rm_facade.known_scope_ids,
185 bump_generation=self.rm_facade.bump_generation,207 bump_generation=self.rm_facade.bump_generation,
@@ -193,17 +215,25 @@ class OrchestratorSystemContext(SystemContext):
193 )215 )
194 self.sm_sweeper = SessionSweeper(sm_state, self.rm_facade)216 self.sm_sweeper = SessionSweeper(sm_state, self.rm_facade)
195 self.rm_sweeper = ResourceSweeper(217 self.rm_sweeper = ResourceSweeper(
196- rm_state, self.k8s, self.sm_facade, orchestrator=rm_orchestrator,218+ rm_state,
219+ self.k8s,
220+ self.sm_facade,
221+ orchestrator=rm_orchestrator,
197 event_sink=eval_state.bump_event,222 event_sink=eval_state.bump_event,
198 )223 )
199 self.eval_state = eval_state224 self.eval_state = eval_state
200 self.eval_telemetry = telemetry225 self.eval_telemetry = telemetry
201 self.eval_collector = EvaluationCollector(226 self.eval_collector = EvaluationCollector(
202- eval_state=eval_state, sm_state=sm_state, rm_state=rm_state,227+ eval_state=eval_state,
228+ sm_state=sm_state,
229+ rm_state=rm_state,
203 )230 )
204 self.evaluator = Evaluator(231 self.evaluator = Evaluator(
205- collector=self.eval_collector, llm=LLMClient.from_arc(self.arc),232+ collector=self.eval_collector,
206- state=eval_state, arc=self.arc, instance_id=self.instance_id,233+ llm=LLMClient.from_arc(self.arc),
234+ state=eval_state,
235+ arc=self.arc,
236+ instance_id=self.instance_id,
207 )237 )
208 self._telemetry_task: asyncio.Task | None = None238 self._telemetry_task: asyncio.Task | None = None
209 239 
@@ -212,43 +242,51 @@ class OrchestratorSystemContext(SystemContext):
212 arc = self.arc242 arc = self.arc
213 jobs = [243 jobs = [
214 self.create_single_leader_job(244 self.create_single_leader_job(
215- name="sm_sweep", on_tick=self.sm_sweeper.sweep_once,245+ name="sm_sweep",
216- interval_sec=arc.sweep_interval, lock_key="agent_runtime:job:sm_sweep",246+ on_tick=self.sm_sweeper.sweep_once,
247+ interval_sec=arc.sweep_interval,
248+ lock_key="agent_runtime:job:sm_sweep",
217 tick_timeout_sec=TICK_TIMEOUTS["sm_sweep"],249 tick_timeout_sec=TICK_TIMEOUTS["sm_sweep"],
218 ),250 ),
219 self.rm_sysctx.create_single_leader_job(251 self.rm_sysctx.create_single_leader_job(
220- name="rm_autoscale", on_tick=self.rm_sweeper.autoscale_once,252+ name="rm_autoscale",
253+ on_tick=self.rm_sweeper.autoscale_once,
221 interval_sec=arc.autoscale_interval,254 interval_sec=arc.autoscale_interval,
222 lock_key="agent_runtime:job:rm_autoscale",255 lock_key="agent_runtime:job:rm_autoscale",
223 tick_timeout_sec=TICK_TIMEOUTS["rm_autoscale"],256 tick_timeout_sec=TICK_TIMEOUTS["rm_autoscale"],
224 ),257 ),
225 self.rm_sysctx.create_single_leader_job(258 self.rm_sysctx.create_single_leader_job(
226- name="rm_reclaim", on_tick=self.rm_sweeper.reclaim_once,259+ name="rm_reclaim",
260+ on_tick=self.rm_sweeper.reclaim_once,
227 interval_sec=arc.reclaim_interval,261 interval_sec=arc.reclaim_interval,
228 lock_key="agent_runtime:job:rm_reclaim",262 lock_key="agent_runtime:job:rm_reclaim",
229 tick_timeout_sec=TICK_TIMEOUTS["rm_reclaim"],263 tick_timeout_sec=TICK_TIMEOUTS["rm_reclaim"],
230 ),264 ),
231 self.rm_sysctx.create_single_leader_job(265 self.rm_sysctx.create_single_leader_job(
232- name="rm_watch", on_tick=self.rm_sweeper.watch_once,266+ name="rm_watch",
267+ on_tick=self.rm_sweeper.watch_once,
233 interval_sec=arc.watch_interval,268 interval_sec=arc.watch_interval,
234 lock_key="agent_runtime:job:rm_watch",269 lock_key="agent_runtime:job:rm_watch",
235 tick_timeout_sec=TICK_TIMEOUTS["rm_watch"],270 tick_timeout_sec=TICK_TIMEOUTS["rm_watch"],
236 ),271 ),
237 self.rm_sysctx.create_single_leader_job(272 self.rm_sysctx.create_single_leader_job(
238- name="rm_reconcile", on_tick=self.rm_sweeper.reconcile_once,273+ name="rm_reconcile",
274+ on_tick=self.rm_sweeper.reconcile_once,
239 interval_sec=arc.reconcile_interval,275 interval_sec=arc.reconcile_interval,
240 lock_key="agent_runtime:job:rm_reconcile",276 lock_key="agent_runtime:job:rm_reconcile",
241 tick_timeout_sec=TICK_TIMEOUTS["rm_reconcile"],277 tick_timeout_sec=TICK_TIMEOUTS["rm_reconcile"],
242 ),278 ),
243 # 系统自评估(sys_eval 全局单副本产报告,任意副本可读):279 # 系统自评估(sys_eval 全局单副本产报告,任意副本可读):
244 self.create_single_leader_job(280 self.create_single_leader_job(
245- name="sys_sample", on_tick=self.eval_collector.sample_once,281+ name="sys_sample",
282+ on_tick=self.eval_collector.sample_once,
246 interval_sec=max(arc.eval_sample_interval, 5),283 interval_sec=max(arc.eval_sample_interval, 5),
247 lock_key="agent_runtime:job:sys_sample",284 lock_key="agent_runtime:job:sys_sample",
248 tick_timeout_sec=TICK_TIMEOUTS["sys_sample"],285 tick_timeout_sec=TICK_TIMEOUTS["sys_sample"],
249 ),286 ),
250 self.create_single_leader_job(287 self.create_single_leader_job(
251- name="sys_eval", on_tick=self.evaluator.evaluate_once,288+ name="sys_eval",
289+ on_tick=self.evaluator.evaluate_once,
252 interval_sec=max(arc.eval_interval, 30),290 interval_sec=max(arc.eval_interval, 30),
253 lock_key="agent_runtime:job:sys_eval",291 lock_key="agent_runtime:job:sys_eval",
254 tick_timeout_sec=TICK_TIMEOUTS["sys_eval"],292 tick_timeout_sec=TICK_TIMEOUTS["sys_eval"],
@@ -260,24 +298,29 @@ class OrchestratorSystemContext(SystemContext):
260 298 
261 async def start(self) -> None:299 async def start(self) -> None:
262 await super().start()300 await super().start()
301+ await sync_local_link_binding_state(self.db, self.link_mtls)
263 await self.rm_sysctx.start()302 await self.rm_sysctx.start()
264 # 启动即重建路由快照(消除首次 route 的冷启动窗口;失败降级到首次 route 重建)303 # 启动即重建路由快照(消除首次 route 的冷启动窗口;失败降级到首次 route 重建)
265 try:304 try:
266 await self.sm_config_store.ensure_snapshot()305 await self.sm_config_store.ensure_snapshot()
267 except Exception: # noqa: BLE001306 except Exception: # noqa: BLE001
268- self.logger.exception("routing snapshot rebuild failed at startup "307+ self.logger.exception(
269- "(defer to first route)")308+ "routing snapshot rebuild failed at startup (defer to first route)"
309+ )
270 try:310 try:
271 await self.k8s.start()311 await self.k8s.start()
272 except Exception: # noqa: BLE001 - k8s 不可用只影响扩缩容,不阻断启动312 except Exception: # noqa: BLE001 - k8s 不可用只影响扩缩容,不阻断启动
273- self.logger.exception("kubernetes client start failed (scale in/out degraded)")313+ self.logger.exception(
314+ "kubernetes client start failed (scale in/out degraded)"
315+ )
274 self._jobs = self._build_jobs()316 self._jobs = self._build_jobs()
275 interval_by_name = self._job_intervals()317 interval_by_name = self._job_intervals()
276 for job in self._jobs:318 for job in self._jobs:
277 await job.start()319 await job.start()
278 self.logger.info(320 self.logger.info(
279 "background job registered: name=%s interval_sec=%s tick_timeout_sec=%s",321 "background job registered: name=%s interval_sec=%s tick_timeout_sec=%s",
280- job.name, interval_by_name.get(job.name, "?"),322+ job.name,
323+ interval_by_name.get(job.name, "?"),
281 TICK_TIMEOUTS.get(job.name),324 TICK_TIMEOUTS.get(job.name),
282 )325 )
283 # 计数缓冲 flusher:每副本独立(选主 job 会漏非 leader 副本的缓冲),326 # 计数缓冲 flusher:每副本独立(选主 job 会漏非 leader 副本的缓冲),
@@ -290,21 +333,38 @@ class OrchestratorSystemContext(SystemContext):
290 "reclaim=%ss watch=%ss reconcile=%ss "333 "reclaim=%ss watch=%ss reconcile=%ss "
291 "default_session_ttl=%ss eval_sample=%ss eval=%ss eval_llm=%s "334 "default_session_ttl=%ss eval_sample=%ss eval=%ss eval_llm=%s "
292 "kubeconfig=%s",335 "kubeconfig=%s",
293- self.arc.mode, self.arc.default_namespace,336+ self.arc.mode,
294- self.arc.sweep_interval, self.arc.autoscale_interval,337+ self.arc.default_namespace,
295- self.arc.reclaim_interval, self.arc.watch_interval,338+ self.arc.sweep_interval,
339+ self.arc.autoscale_interval,
340+ self.arc.reclaim_interval,
341+ self.arc.watch_interval,
296 self.arc.reconcile_interval,342 self.arc.reconcile_interval,
297 self.arc.default_session_ttl,343 self.arc.default_session_ttl,
298- self.arc.eval_sample_interval, self.arc.eval_interval,344+ self.arc.eval_sample_interval,
299- "enabled" if (self.arc.eval_llm_base_url and self.arc.eval_llm_model)345+ self.arc.eval_interval,
346+ "enabled"
347+ if (self.arc.eval_llm_base_url and self.arc.eval_llm_model)
300 else "disabled",348 else "disabled",
301 "set" if self.arc.kubeconfig else "in-cluster",349 "set" if self.arc.kubeconfig else "in-cluster",
302 )350 )
303 self.logger.info(351 self.logger.info(
304 "agent-runtime started: instance=%s mode=%s port=%s",352 "agent-runtime started: instance=%s mode=%s port=%s",
305- self.instance_id, self.arc.mode,353+ self.instance_id,
354+ self.arc.mode,
306 getattr(self.settings, "port", "?"),355 getattr(self.settings, "port", "?"),
307 )356 )
357+ self.logger.info(
358+ "link mTLS: mode=%s mtls_deployment_id=%s mtls_binding_id=%s epoch=%s",
359+ self.link_mtls.mode.value,
360+ self.link_mtls.identity.mtls_deployment_id
361+ if self.link_mtls.identity
362+ else "-",
363+ self.link_mtls.identity.mtls_binding_id if self.link_mtls.identity else "-",
364+ self.link_mtls.identity.mtls_binding_epoch
365+ if self.link_mtls.identity
366+ else "-",
367+ )
308 368 
309 def _job_intervals(self) -> dict[str, int]:369 def _job_intervals(self) -> dict[str, int]:
310 """诊断用:job 名 → 调度间隔(arc)。"""370 """诊断用:job 名 → 调度间隔(arc)。"""
@@ -350,7 +410,11 @@ class OrchestratorSystemContext(SystemContext):
350 lock_key = f"agent_runtime:job:{job.name}"410 lock_key = f"agent_runtime:job:{job.name}"
351 token = await self.redis.get(lock_key)411 token = await self.redis.get(lock_key)
352 if token:412 if token:
353- token = token.decode() if isinstance(token, (bytes, bytearray)) else str(token)413+ token = (
414+ token.decode()
415+ if isinstance(token, (bytes, bytearray))
416+ else str(token)
417+ )
354 leader = token.removeprefix(f"{job.name}:").rsplit(":", 1)[0]418 leader = token.removeprefix(f"{job.name}:").rsplit(":", 1)[0]
355 entry["leader"] = {419 entry["leader"] = {
356 "instance_id": leader,420 "instance_id": leader,
@@ -410,6 +474,7 @@ def create_app(
410 resources: tuple[Any, Any, Any] | None = None,474 resources: tuple[Any, Any, Any] | None = None,
411 instance_id: str | None = None,475 instance_id: str | None = None,
412 own_resources: bool = True,476 own_resources: bool = True,
477+ link_mtls_config: LinkMTLSConfig | None = None,
413) -> App:478) -> App:
414 """构造唯一 App(/api/session)并注册 5 个对外 handler。479 """构造唯一 App(/api/session)并注册 5 个对外 handler。
415 480 
@@ -418,19 +483,44 @@ def create_app(
418 生产路径不传,行为与原先完全一致。483 生产路径不传,行为与原先完全一致。
419 """484 """
420 485 
486+ link_mtls = link_mtls_config or LinkMTLSConfig.from_env()
487+ 
421 def ctx_factory() -> OrchestratorSystemContext:488 def ctx_factory() -> OrchestratorSystemContext:
422 if resources is not None:489 if resources is not None:
423 redis_client, db, k8s = resources490 redis_client, db, k8s = resources
424 else:491 else:
425- redis_client, db, k8s = build_resources(settings, arc)492+ redis_client, db, k8s = build_resources(settings, arc, link_mtls)
426 return OrchestratorSystemContext(493 return OrchestratorSystemContext(
427- redis_client=redis_client, db=db, k8s=k8s,494+ redis_client=redis_client,
428- settings=settings, arc=arc,495+ db=db,
496+ k8s=k8s,
497+ settings=settings,
498+ arc=arc,
499+ link_mtls=link_mtls,
429 instance_id=instance_id,500 instance_id=instance_id,
430 owns_resources=own_resources,501 owns_resources=own_resources,
431 )502 )
432 503 
433- app = App(ctx_factory, prefix=SERVICE_PREFIX, enable_ws=False, title="agent-runtime")504+ app = App(
505+ ctx_factory, prefix=SERVICE_PREFIX, enable_ws=False, title="agent-runtime"
506+ )
507+ 
508+ @app.asgi.middleware("http")
509+ async def _link_binding_guard(request: Request, call_next): # noqa: ANN001, ANN202
510+ try:
511+ link_mtls.authorize_request(request)
512+ except LinkMTLSError as exc:
513+ logger.warning("link binding rejected: %s", exc)
514+ return JSONResponse(
515+ status_code=403,
516+ content={
517+ "ok": False,
518+ "error_code": "LINK_BINDING_MISMATCH",
519+ "error_message": str(exc),
520+ },
521+ )
522+ return await call_next(request)
523+ 
434 # 请求汇总日志 + 指标(touch/cleanup 从此每请求一行;registry 供 /visualization/stats)524 # 请求汇总日志 + 指标(touch/cleanup 从此每请求一行;registry 供 /visualization/stats)
435 registry = MetricsRegistry()525 registry = MetricsRegistry()
436 app.use(request_metrics_middleware(registry))526 app.use(request_metrics_middleware(registry))
@@ -438,6 +528,8 @@ def create_app(
438 _register_healthz(app)528 _register_healthz(app)
439 register_visualization_api(app, registry=registry)529 register_visualization_api(app, registry=registry)
440 register_handlers(app)530 register_handlers(app)
531+ if link_mtls.mode is LinkMTLSMode.OBSERVE:
532+ logger.info("link mTLS observe preflight complete; data path remains HTTP")
441 return app533 return app
442 534 
443 535 
@@ -479,7 +571,9 @@ def _collect_runtime_identity() -> dict[str, Any]:
479 except OSError:571 except OSError:
480 ns = ""572 ns = ""
481 if not ns:573 if not ns:
482- ns = (os.getenv("AGENT_RUNTIME_DEFAULT_NAMESPACE") or "default").strip() or "default"574+ ns = (
575+ os.getenv("AGENT_RUNTIME_DEFAULT_NAMESPACE") or "default"
576+ ).strip() or "default"
483 pod = (os.getenv("HOSTNAME") or "").strip()577 pod = (os.getenv("HOSTNAME") or "").strip()
484 return {578 return {
485 "namespace": ns,579 "namespace": ns,
Mapplications/agent_runtime/src/agent_runtime/resource_manager/k8s.py+152-1文件内容审核中,请稍后刷新重试
Mapplications/agent_runtime/src/agent_runtime/resource_manager/sweeper.py+7-1文件内容审核中,请稍后刷新重试
Mapplications/agent_runtime/src/agent_runtime/session_manager/handlers.py+10-1文件内容审核中,请稍后刷新重试
Mapplications/agent_runtime/tests/resource_manager/test_k8s_pod_body.py+563-201文件内容审核中,请稍后刷新重试
@@ -77,6 +77,12 @@ CLAWMANAGER_CONFIG_ENC_REQUIRED=false
77CLAWMANAGER_CONFIG_SIGN_ENABLED=false77CLAWMANAGER_CONFIG_SIGN_ENABLED=false
78CLAWMANAGER_CONFIG_SIGN_ALG=Ed2551978CLAWMANAGER_CONFIG_SIGN_ALG=Ed25519
79 79 
80+# 可选 Manager 参考实现到 Gateway / Runtime 的内部 mTLS。
81+# 未设置或 off 均不启用;仅显式 enforce 启用。配套部署工具管理 manager 角色材料。
82+# 受管地址按部署配置解析为 HTTPS;外部地址须显式 https:// 并通过 CA/SAN/指纹校验。
83+# 当前从部署身份读取绑定 ID/epoch,不需另填证书路径或手改绑定数据库。
84+JIUWENSWARM_LINK_MTLS_MODE=off
85+ 
80# 本地拉起时写入 Gateway 的数据库环境(可选)86# 本地拉起时写入 Gateway 的数据库环境(可选)
81# JIUWENCLAW_GATEWAY_DB_TYPE=mysql87# JIUWENCLAW_GATEWAY_DB_TYPE=mysql
82# JIUWENCLAW_GATEWAY_DB_HOST=localhost88# JIUWENCLAW_GATEWAY_DB_HOST=localhost
@@ -273,12 +273,12 @@ async def purge_runtime_instance_data(
273 from manager_server.core.instance_resource.runtime_config_sync import (273 from manager_server.core.instance_resource.runtime_config_sync import (
274 sync_runtime_config,274 sync_runtime_config,
275 )275 )
276- from manager_server.infrastructure.config import settings
277 276 
278 jid = str(jiuwenclaw_id or "").strip()277 jid = str(jiuwenclaw_id or "").strip()
279 if not jid:278 if not jid:
280 return {"purged": False}279 return {"purged": False}
281- if not settings.agent_runtime_endpoint.strip():280+ from manager_server.core.instance_resource.runtime_config_sync import resolve_runtime_endpoint
281+ if not await resolve_runtime_endpoint(handler, jid):
282 logger.info(282 logger.info(
283 "[InstanceDataLifecycle] runtime purge skipped jiuwenclaw_id=%s "283 "[InstanceDataLifecycle] runtime purge skipped jiuwenclaw_id=%s "
284 "(AGENT_RUNTIME_ENDPOINT empty)",284 "(AGENT_RUNTIME_ENDPOINT empty)",
@@ -7,33 +7,34 @@ from openjiuwen_runtime.foundation.db.sqlalchemy_handler import SQLAlchemyHandle
7from sqlalchemy import inspect, text7from sqlalchemy import inspect, text
8from sqlalchemy.exc import DBAPIError8from sqlalchemy.exc import DBAPIError
9 9 
10+from manager_server.models.application_config_models import (
11+ _MEMORY_CONFIG_TABLE_DEF,
12+ _TASK_MEMORY_CONFIG_TABLE_DEF,
13+ LOG_MASKING_RULE_TABLE_DEF,
14+ LOGGING_CONFIG_TABLE_DEF,
15+)
16+from manager_server.models.instance_access_models import INSTANCE_ACCESS_TABLE_DEFINITIONS
10from manager_server.models.instance_models import INSTANCE_INFO_TABLE_DEF17from manager_server.models.instance_models import INSTANCE_INFO_TABLE_DEF
18+from manager_server.models.instance_resource_models import INSTANCE_RESOURCE_TABLE_DEFINITIONS
19+from manager_server.models.jid_template_ref_models import (
20+ JID_TEMPLATE_REF_TABLE_DEF,
21+)
11from manager_server.models.key_models import (22from manager_server.models.key_models import (
12 INSTANCE_ENC_PUBKEY_TABLE_DEF,23 INSTANCE_ENC_PUBKEY_TABLE_DEF,
13 MANAGER_IDENTITY_TABLE_DEF,24 MANAGER_IDENTITY_TABLE_DEF,
14)25)
15-from manager_server.models.application_config_models import (26+from manager_server.models.link_binding_models import INSTANCE_LINK_BINDING_TABLE_DEF
16- LOG_MASKING_RULE_TABLE_DEF,
17- LOGGING_CONFIG_TABLE_DEF,
18- _TASK_MEMORY_CONFIG_TABLE_DEF,
19- _MEMORY_CONFIG_TABLE_DEF,
20-)
21-from manager_server.models.jid_template_ref_models import (
22- JID_TEMPLATE_REF_TABLE_DEF,
23-)
24-from manager_server.models.instance_access_models import INSTANCE_ACCESS_TABLE_DEFINITIONS
25-from manager_server.models.instance_resource_models import INSTANCE_RESOURCE_TABLE_DEFINITIONS
26from manager_server.models.template_models import (27from manager_server.models.template_models import (
27 A2A_ACCESS_POLICY_TEMPLATE_TABLE_DEF,28 A2A_ACCESS_POLICY_TEMPLATE_TABLE_DEF,
28 A2A_DISCOVERY_SETTINGS_TABLE_DEF,29 A2A_DISCOVERY_SETTINGS_TABLE_DEF,
29- A2A_OUTBOUND_TEMPLATE_TABLE_DEF,
30 A2A_OUTBOUND_DISCOVERY_TABLE_DEF,30 A2A_OUTBOUND_DISCOVERY_TABLE_DEF,
31+ A2A_OUTBOUND_TEMPLATE_TABLE_DEF,
31 AGENT_TEMPLATE_TABLE_DEF,32 AGENT_TEMPLATE_TABLE_DEF,
32 EMBEDDING_TEMPLATE_TABLE_DEF,33 EMBEDDING_TEMPLATE_TABLE_DEF,
33 EXTENSION_CONFIG_TEMPLATE_TABLE_DEF,34 EXTENSION_CONFIG_TEMPLATE_TABLE_DEF,
35+ MCP_TEMPLATE_TABLE_DEF,
34 MODEL_TEMPLATE_TABLE_DEF,36 MODEL_TEMPLATE_TABLE_DEF,
35 PERMISSIONS_TEMPLATE_TABLE_DEF,37 PERMISSIONS_TEMPLATE_TABLE_DEF,
36- MCP_TEMPLATE_TABLE_DEF,
37 SERVICE_CONFIG_CONTAINER_TABLE_DEF,38 SERVICE_CONFIG_CONTAINER_TABLE_DEF,
38 SERVICE_CONFIG_TEMPLATE_TABLE_DEF,39 SERVICE_CONFIG_TEMPLATE_TABLE_DEF,
39 SKILL_PREBUILT_TEMPLATE_TABLE_DEF,40 SKILL_PREBUILT_TEMPLATE_TABLE_DEF,
@@ -41,6 +42,7 @@ from manager_server.models.template_models import (
41 42 
42ALL_TABLE_DEFINITIONS = (43ALL_TABLE_DEFINITIONS = (
43 INSTANCE_INFO_TABLE_DEF,44 INSTANCE_INFO_TABLE_DEF,
45+ INSTANCE_LINK_BINDING_TABLE_DEF,
44 MANAGER_IDENTITY_TABLE_DEF,46 MANAGER_IDENTITY_TABLE_DEF,
45 INSTANCE_ENC_PUBKEY_TABLE_DEF,47 INSTANCE_ENC_PUBKEY_TABLE_DEF,
46 _TASK_MEMORY_CONFIG_TABLE_DEF,48 _TASK_MEMORY_CONFIG_TABLE_DEF,
@@ -5,16 +5,21 @@ Gateway / Runtime 存活由 Manager 周期探活 ``*_config_host`` 健康检查
5 5 
6from __future__ import annotations6from __future__ import annotations
7 7 
8-from typing import Annotated, Any8+from typing import Annotated
9 9 
10from fastapi import APIRouter, Depends, HTTPException, Query10from fastapi import APIRouter, Depends, HTTPException, Query
11from openjiuwen_runtime.foundation.db.handler import DBHandler11from openjiuwen_runtime.foundation.db.handler import DBHandler
12 12 
13from manager_server.core.instance import InstanceService13from manager_server.core.instance import InstanceService
14-from manager_server.infrastructure.db import get_db_handler14+from manager_server.core.instance.link_binding_service import (
15+ InstanceLinkBindingService,
16+ LinkBindingConflict,
17+)
15from manager_server.core.instance.pod_status_cache import (18from manager_server.core.instance.pod_status_cache import (
16 get_pod_status_snapshot,19 get_pod_status_snapshot,
17)20)
21+from manager_server.infrastructure.db import get_db_handler
22+from manager_server.routers.deps import AdminUser
18from manager_server.schedulers.heartbeat_scanner import scan_instance_health_once23from manager_server.schedulers.heartbeat_scanner import scan_instance_health_once
19from manager_server.schemas.common_schemas import ResponseModel24from manager_server.schemas.common_schemas import ResponseModel
20from manager_server.schemas.instance_schemas import (25from manager_server.schemas.instance_schemas import (
@@ -22,6 +27,7 @@ from manager_server.schemas.instance_schemas import (
22 InstanceListQuery,27 InstanceListQuery,
23 InstanceUpdateBody,28 InstanceUpdateBody,
24)29)
30+from manager_server.schemas.link_binding_schemas import LinkBindingCreateBody
25 31 
26instance_router = APIRouter()32instance_router = APIRouter()
27 33 
@@ -30,6 +36,10 @@ def _svc(handler: DBHandler) -> InstanceService:
30 return InstanceService(handler)36 return InstanceService(handler)
31 37 
32 38 
39+def _link_svc(handler: DBHandler) -> InstanceLinkBindingService:
40+ return InstanceLinkBindingService(handler)
41+ 
42+ 
33def _request_volume_value(bv: dict, key: str, legacy_key: str | None = None) -> int:43def _request_volume_value(bv: dict, key: str, legacy_key: str | None = None) -> int:
34 value = bv.get(key)44 value = bv.get(key)
35 if value is None and legacy_key is not None:45 if value is None and legacy_key is not None:
@@ -41,9 +51,7 @@ def _normalize_request_volume(bv: dict) -> dict:
41 return {51 return {
42 "gateway_queued": _request_volume_value(bv, "gateway_queued"),52 "gateway_queued": _request_volume_value(bv, "gateway_queued"),
43 "gateway_running": _request_volume_value(bv, "gateway_running"),53 "gateway_running": _request_volume_value(bv, "gateway_running"),
44- "service_manager_queued": _request_volume_value(54+ "service_manager_queued": _request_volume_value(bv, "service_manager_queued", "sm_queued"),
45- bv, "service_manager_queued", "sm_queued"
46- ),
47 "service_manager_routing": _request_volume_value(55 "service_manager_routing": _request_volume_value(
48 bv, "service_manager_routing", "sm_routing"56 bv, "service_manager_routing", "sm_routing"
49 ),57 ),
@@ -58,16 +66,12 @@ def _normalize_request_volume(bv: dict) -> dict:
58 66 
59 67 
60def _build_request_volume_summary(bv: dict) -> dict:68def _build_request_volume_summary(bv: dict) -> dict:
61- queued_requests = _request_volume_value(69+ queued_requests = _request_volume_value(bv, "gateway_queued") + _request_volume_value(
62- bv, "gateway_queued"
63- ) + _request_volume_value(
64 bv, "service_manager_queued", "sm_queued"70 bv, "service_manager_queued", "sm_queued"
65 )71 )
66 running_requests = _request_volume_value(72 running_requests = _request_volume_value(
67 bv, "service_manager_running", "sm_running"73 bv, "service_manager_running", "sm_running"
68- ) or _request_volume_value(74+ ) or _request_volume_value(bv, "gateway_running")
69- bv, "gateway_running"
70- )
71 return {75 return {
72 "queued_requests": queued_requests,76 "queued_requests": queued_requests,
73 "running_requests": running_requests,77 "running_requests": running_requests,
@@ -126,9 +130,7 @@ async def update_instance(
126 130 
127 131 
128@instance_router.get("/{jiuwenclaw_id}", response_model=ResponseModel)132@instance_router.get("/{jiuwenclaw_id}", response_model=ResponseModel)
129-async def get_instance(133+async def get_instance(jiuwenclaw_id: str, handler: Annotated[DBHandler, Depends(get_db_handler)]):
130- jiuwenclaw_id: str, handler: Annotated[DBHandler, Depends(get_db_handler)]
131-):
132 svc = _svc(handler)134 svc = _svc(handler)
133 row = await svc.get(jiuwenclaw_id)135 row = await svc.get(jiuwenclaw_id)
134 if row is None:136 if row is None:
@@ -150,6 +152,63 @@ async def delete_instance(
150 return ResponseModel(code=200, message="success", data={"deleted": True})152 return ResponseModel(code=200, message="success", data={"deleted": True})
151 153 
152 154 
155+@instance_router.put("/{jiuwenclaw_id}/link-binding", response_model=ResponseModel)
156+async def bind_instance_link(
157+ jiuwenclaw_id: str,
158+ body: LinkBindingCreateBody,
159+ handler: Annotated[DBHandler, Depends(get_db_handler)],
160+ admin: AdminUser,
161+):
162+ """建立当前有效的 Gateway ↔ Runtime 绑定;完全相同请求幂等。"""
163+ _ = admin
164+ try:
165+ binding = await _link_svc(handler).bind(jiuwenclaw_id, body)
166+ except LookupError as exc:
167+ raise HTTPException(status_code=404, detail=str(exc)) from exc
168+ except LinkBindingConflict as exc:
169+ raise HTTPException(status_code=409, detail=str(exc)) from exc
170+ return ResponseModel(code=200, message="success", data=binding.model_dump())
171+ 
172+ 
173+@instance_router.get("/{jiuwenclaw_id}/link-binding", response_model=ResponseModel)
174+async def get_instance_link_binding(
175+ jiuwenclaw_id: str,
176+ handler: Annotated[DBHandler, Depends(get_db_handler)],
177+ admin: AdminUser,
178+):
179+ _ = admin
180+ binding = await _link_svc(handler).get(jiuwenclaw_id)
181+ if binding is None:
182+ raise HTTPException(status_code=404, detail="link binding not found")
183+ return ResponseModel(code=200, message="success", data=binding.model_dump())
184+ 
185+ 
186+@instance_router.delete("/{jiuwenclaw_id}/link-binding", response_model=ResponseModel)
187+async def unbind_instance_link(
188+ jiuwenclaw_id: str,
189+ handler: Annotated[DBHandler, Depends(get_db_handler)],
190+ admin: AdminUser,
191+ updated_by: str = Query("system", min_length=1, max_length=64),
192+):
193+ _ = admin
194+ try:
195+ binding, rotation_required = await _link_svc(handler).unbind(
196+ jiuwenclaw_id, updated_by=updated_by
197+ )
198+ except LookupError as exc:
199+ raise HTTPException(status_code=404, detail=str(exc)) from exc
200+ except LinkBindingConflict as exc:
201+ raise HTTPException(status_code=409, detail=str(exc)) from exc
202+ return ResponseModel(
203+ code=200,
204+ message="success",
205+ data={
206+ **binding.model_dump(),
207+ "rotation_required": rotation_required,
208+ },
209+ )
210+ 
211+ 
153@instance_router.get("/{jiuwenclaw_id}/pods", response_model=ResponseModel)212@instance_router.get("/{jiuwenclaw_id}/pods", response_model=ResponseModel)
154async def get_instance_pods(213async def get_instance_pods(
155 jiuwenclaw_id: str,214 jiuwenclaw_id: str,
@@ -75,6 +75,7 @@ async def test_purge_runtime_pushes_empty_projection():
75@pytest.mark.asyncio75@pytest.mark.asyncio
76async def test_purge_runtime_skips_when_endpoint_empty():76async def test_purge_runtime_skips_when_endpoint_empty():
77 handler = AsyncMock()77 handler = AsyncMock()
78+ handler.get = AsyncMock(return_value=None) # Neither instance-specific nor global endpoint.
78 with (79 with (
79 patch(80 patch(
80 "manager_server.infrastructure.config.settings.agent_runtime_endpoint",81 "manager_server.infrastructure.config.settings.agent_runtime_endpoint",