Skip to content

feat(flexlb): support logical multi-engine workers - #1332

Open
Martin7-1 wants to merge 5 commits into
feature/flexlb-mu-devfrom
feature/multi-engine-support
Open

feat(flexlb): support logical multi-engine workers#1332
Martin7-1 wants to merge 5 commits into
feature/flexlb-mu-devfrom
feature/multi-engine-support

Conversation

@Martin7-1

Copy link
Copy Markdown
Collaborator

Background

DashLLM can start multiple independently routable engines behind one frontend when DS_LLM_MULTI_ENGINE_NUM > 1. FlexLB must preserve that engine index from service discovery through worker-status/cache polling, cache matching, scheduling, feedback, rollback, and observability.

Implementation

  • Introduce the shared worker-identity contract:
    • physical frontend: ip:port
    • logical worker: ip:port@engineIndex
    • metrics identity: ip@engineIndex
  • Expand each discovered frontend into multi_engine_num logical workers. Worker-status and cache-status gRPC use worker_status_port + engineIndex; for N=1, an omitted worker_status_port falls back to the engine gRPC port.
  • Keep physical connection targets as plain ip:port; pass WorkerHost to runners so LOCAL_SYNC never parses a logical address as a TCP endpoint.
  • Apply DashServing-aligned physical-group AND health: any missing, duplicate-index, or non-alive sibling removes every sibling from routing. Resource capacity remains per logical worker.
  • Keep LOCAL_SYNC, LOCAL_STANDBY, KVCM cache matching, cache-hit feedback/PV, rollback, and metrics scoped to the logical worker identity. Stale local cache cleanup compares discovered physical frontends against the physical owner recorded for each logical cache entry.
  • Schedule responses preserve engine_index only for N>1. N=1 responses omit it, while all internal cache/subscriber/KVCM identities remain @0.
  • Add endpoint validation: multi_engine_num >= 1, worker_status_port is in [1, 65535], and N>1 requires a configured worker-status port whose complete offset range fits in the TCP range.
  • Add N=1 KVCM compatibility: exact ip:port@0 wins; only KVCM results may fall back to legacy physical ip:port for a single-engine worker. Physical keys remain ignored for N>1 to avoid assigning a multi-engine cache match to the wrong index.
  • Isolate the static load-balancer strategy registry in tests so mock registrations cannot leak across test classes.

Compatibility and rollout

  • New multi-engine deployments must configure multi_engine_num and a base worker_status_port; engines expose one worker-control gRPC endpoint per index.
  • Existing single-engine deployments continue to work with or without worker_status_port; old worker-status responses are represented as index 0.
  • An old KVCM can continue serving single-engine physical host_ip_port values through the N=1 fallback. Multi-engine cache routing requires KVCM/subscriber to return ip:port@engineIndex.
  • The HTTP/gRPC frontend connection address and URI remain unchanged; @index is an internal logical identity, not a network endpoint suffix.

Verification

  • Java 21: ./mvnw verify
  • Focused regression: ./mvnw -pl flexlb-sync -am -Dtest=LoadBalanceStrategyFactoryTest,DefaultRouterTest -Dsurefire.failIfNoSpecifiedTests=false test (20 tests passed)
  • Added coverage for N=2 discovery/offsets, physical AND health, logical cache ownership and rollback, legacy N=1 KVCM fallback, endpoint range boundaries, LOCAL_SYNC failures, and strategy-factory cleanup.
  • N=2 vLLM end-to-end routing was verified for index discovery, worker-status offsets, schedule response propagation, cache identity, and physical-health filtering.

Explicitly deferred

  • Group-level health-transition logs/metrics are not added in this PR.
  • No fixed upper bound is imposed on multi_engine_num; the configured worker-status TCP port range remains the enforced bound.

Rollback

Set multi_engine_num to 1 and restart the deployment. FlexLB resumes one logical worker per frontend and omits engine_index from schedule responses.

@LLLLKKKK LLLLKKKK left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

AI Code Review - PR #1332

Status: LGTM

Summary: P0/0 · P1/0 · P2/16 · P3/7

Reviewed: commit c47bdf062ad3 · 2026-08-25 22:38 UTC+8

lgtm ready to ci

Non-blocking Suggestions

P2

  • KVCM host_ip_port 键语义在 proto 契约与 Java 消费方之间自相矛盾 @ rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/match/kvcm/KvcmCacheMatchProvider.java:31
    • 建议:同步更新 kvcm_meta_service.protohost_ip_port 的注释为逻辑 ip:port@engineIndex(并注明单引擎旧格式在 FlexLB 侧有回退兼容),使 proto 成为唯一契约来源;在 PR description 中把「启用 multi_engine_num > 1 前 KVCM Subscriber 必须先上报逻辑身份」列为显式发布前置条件与灰度顺序。
  • KVCM 键与逻辑 worker 全不匹配时静默按零命中处理,缺可观测手段 @ rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/domain/CacheMatchResult.java:56
    • 建议:补两个计数(或限频 warn):legacy 物理键回退命中次数;「KVCM 返回非空 hostMatches 但无任何键匹配已知 logical worker」的次数(tag 带 multiEngineNum,日志附一个示例键)。并在 04/06 文档把「启用 multi_engine_num > 1 前 KVCM 必须上报 logical key,否则命中恒为 0」列为显式运维前置条件,使灰度阶段即可暴露契约不一致。也可考虑把该兼容判定从 domain record 上移到 KVCM provider 或查询编排层,避免 CacheMatchResult 直接耦合运行态 WorkerStatus
  • logical 到 physical 映射生命周期不完整,孤儿条目单向泄漏且缺失映射走隐式兜底 @ rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/match/localsync/KvCacheManager.java:164
    • 建议:把待清理集合改为 getAllEngineIpPorts()physicalIpPortByLogicalIpPort.keySet() 的并集再做物理活跃性过滤,让孤儿映射同轮回收;把 getOrDefault 兜底改为显式分支,缺失时先 warn 再决定是否判 stale(当前该分支实际不可达,显式化可防后续调用方绕过登记)。也可改由 WorkerStatusProvider 当前活跃 logical 集合驱动清理以免维护第二份状态,并把该 map size 纳入 app.cache.* 便于线上确认回收生效。
  • removeStaleEngineCaches 对空活跃地址列表无保护,服务发现瞬时返回空会清空整份索引 @ rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/match/localsync/KvCacheManager.java:158
    • 建议:对空活跃列表直接跳过清理并打 warn,或要求连续多轮观测到缺失后才执行清理,使缓存元数据清理与 worker 摘除采用一致的宽限期语义。
  • 指标身份契约多处变更(tag key 重命名 + 取值格式 + port 维度丢失)缺少迁移清单 @ rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/telemetry/CacheMetricsReporter.java:333
    • 建议:在 06 文档或发布说明中补一节指标迁移清单:列出被重命名的 tag key(engineIpPortengineIpipengineIp)、取值格式变化涉及的指标族(至少 app.cache.*app.engine.*app.routing.*)与新增的 engineIndex tag,并注明回滚该版本时看板需回退旧格式;同时评估同 IP 多 frontend 部署形态,若存在则为丢失 port 维度的两个 cache 指标补独立 port tag 以避免标签冲突。
  • cache-hit PV 的 worker 字段取值由裸 IP 变为逻辑身份,文档却描述为「继续使用」 @ rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/match/localstandby/LocalStandbyComparisonService.java:141
    • 建议:将该句改为明确的变更声明(旧值为裸 IP,新值为 ip:port@engineIndex,含 N=1),并在 PR 描述的兼容性小节列出该 PV 字段变更;若下游消费方较多,可短期并存一个保留裸 IP 语义的独立字段,给切换留出一个发布周期后再下线。另 CacheHitFeedback.workerIp() 现已无生产调用方,可一并评估是否随之收敛。
  • 路由决策 PV 快照身份迁移不完整,同一记录内 logical 与 physical 混用 @ rtp_llm/flexlb/flexlb-common/src/main/java/org/flexlb/dao/pv/ShortestTtftDecision.java:47
    • 建议:给 WorkerDecision 增加 engineIndex(或直接改用 logical 身份字段),使 workers[] 行可唯一定位逻辑引擎并能与 cacheAffinityDecision 的两个字段字符串对齐;同时在 06 文档/PR 描述中列出 cacheLeaderIpPort / shortestTtftWorkerIpPort 的取值格式变更。若短期不改字段,至少在文档说明多引擎下应改用 cacheLeader / shortestTtftWorker 布尔位而非字符串关联。
  • 已有不可变 WorkerIdentity 却仍以相邻同类型 String 传递三种身份表示 @ rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/match/localsync/KvCacheManager.java:105
    • 建议:让 updateEngineCache / calculateDiff 直接接收 WorkerIdentity(或最小只读身份视图),在方法内部取 getLogicalIpPort() / getPhysicalIpPort() / getIpIndex(),把三种表示的一致性交给类型系统而非形参顺序;CacheHitComparisonResult 同理可只持有 identity 并派生两个字段。若倾向最小改动,可让 KvCacheManager 拿到 DiffResult 后自行上报 diff 指标,calculateDiff 保持仅依赖存储 key 的原签名。
  • multi_engine_num 缺少显式上限,配置笔误会放大同步扇出并使整机永久不可路由 @ rtp_llm/flexlb/flexlb-common/src/main/java/org/flexlb/config/ModelServiceConfiguration.java:64
    • 建议:在 multiEngineNum < 1 判断处一并检查显式上界常量(取值对齐现实部署上限),错误信息同时给出配置值与上限;补一条超上限被拒绝的单测,并在 06 文档字段说明同步该上限,使误配在启动期失败而非退化为运行期静默摘流。
  • 启动日志未输出多引擎派生信息,误配后缺少可对照的排查入口 @ rtp_llm/flexlb/flexlb-common/src/main/java/org/flexlb/config/ModelServiceConfiguration.java:36
    • 建议:在 modelMetaConfig 的 info 日志中补充每个 endpoint 的 addressmulti_engine_num 与派生端口段(或展开后的逻辑 worker 总数);该日志启动期仅执行一次,无热路径噪声风险。可一并考虑在 selectRoutableModelWorkerStatus 过滤掉不完整物理组时输出限频 warn 或计数,使 sibling 不齐导致的整机摘流可观测。
  • 端口段上界「接受」用例依赖字符串 replace,模板变更时会静默退化为永真断言 @ rtp_llm/flexlb/flexlb-common/src/test/java/org/flexlb/config/ModelServiceConfigurationTest.java:85
    • 建议:在 run 回调内追加对生效配置的断言(取出 endpoint 校验 getWorkerStatusPort() 为 65534、getMultiEngineNum() 为 2);或参照同文件 modelConfig(...)"""...""".formatted(...) 写法用带占位符的模板生成 endpoint JSON,不再依赖 replace 命中。
  • CacheMetricsReporterTest 只锁定一个上报方法,其余 engineIp tag 语义变更无覆盖 @ rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/telemetry/CacheMetricsReporterTest.java:26
    • 建议:为上述四个方法各补一条 tag 断言(复用 FlexMetricTags.of(...) + eq(...)),并补 ipIndex == null 时不上报的用例,使「每引擎一条指标」的可观测契约可回归。
  • WorkerStatusTest 未覆盖 engineIndex 传播到 CacheHitFeedback 的路径 @ rtp_llm/flexlb/flexlb-common/src/test/java/org/flexlb/dao/master/WorkerStatusTest.java:34
    • 建议:在产出 feedback 的用例中设置 setEngineIndex(1)(或新增一条),断言 cacheHitFeedbacks().getFirst().logicalWorkerId()127.0.0.1:8080@1ipIndex()127.0.0.1@1,把「engineIndex 端到端进入 feedback」这一契约钉死。
  • KvCacheManagerTest 缺少同物理机多逻辑引擎的失效清理覆盖 @ rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/match/localsync/KvCacheManagerTest.java:51
    • 建议:补一条用例:getAllEngineIpPorts() 返回 10.0.0.1:8080@010.0.0.1:8080@1 并分别登记映射,断言物理机仍在活跃列表时两者都不清理、消失时两者都清理;另补一条「logical key 在 view 中但映射缺失且物理机存活」的用例,明确该兜底分支的预期行为。
  • 缺少「预测集合中不存在该 logical worker」的静默降级负例 @ rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/match/localstandby/LocalStandbyComparisonServiceTest.java:110
    • 建议:补一条用例:预测仅含 10.0.0.1:8080@0 而 feedback 的 engineIndex 为 1,断言 localStandby().hit() 为 0 且 delta() 等于 actualHitTokens,把该降级语义显式固化,后续 key 格式漂移即可被区分出来。
  • CacheMatchResultTest 缺少空 worker 与退化输入的边界覆盖 @ rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/domain/CacheMatchResultTest.java:39
    • 建议:补 hostMatch((WorkerStatus) null) 返回 null 的用例,以及 matchedTokens(0, 4096, 100)matchedTokens(2, 0, 100) 断言为 0,使 legacy 回退与退化输入的适用范围在测试中无歧义。可另加一条 engineIndex>0 且 multiEngineNum==1 的防御性用例(该组合由 normalizeHosts:109 保证不可达,仅用于固化预期)。

P3

  • CacheMatchResult 两个 hostMatch 重载语义分歧,易被后续调用方误用为严格版本 @ rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/domain/CacheMatchResult.java:45
    • 建议:给两个方法取语义化的不同名字(如 exactHostMatch(String)resolveHostMatch(WorkerStatus)),或把 hostMatch(String) 收窄为包级可见,并在 javadoc 顶部显式标注「不含 KVCM legacy 回退,路由路径请使用 WorkerStatus 重载」。
  • WorkerCacheUpdateResult 的 javadoc 交叉引用指向已不存在的 WorkerStatus 字段 @ rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/domain/WorkerCacheUpdateResult.java:23
    • 建议:改为引用 org.flexlb.dao.master.WorkerIdentity#logicalIpPort(或 WorkerStatus#getLogicalIpPort()),与本 PR 确立的「身份统一由 WorkerIdentity 提供」约定保持一致。
  • LOCAL_STANDBY 测试夹具仍用物理 key,与生产 key 契约脱节 @ rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/match/localstandby/LocalStandbyCacheMatchProviderTest.java:48
    • 建议:将该用例的 mock key 与断言统一改为逻辑格式(如 10.0.0.1:8080@0),物理 key 用例仅保留在 KVCM 兼容场景;并可把 LocalSyncCacheMatchProviderTest 各用例重复的 mock/provider 构造抽到 @BeforeEach,减少后续身份格式变更的漏改面。
  • 新增 engine-index 用例的期望值与 feedback 字段重合,判别力弱于命名声明 @ rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/match/localstandby/LocalStandbyComparisonServiceTest.java:133
    • 建议:把 feedback 的 predictedHitTokens / deltaHitTokens 改成与预测推导值不同的数值(例如 4096/4904),使断言只能由「按索引查预测」这一条路径满足。
  • WorkerIdentityTest 混入 WorkerHost 断言,且未覆盖多引擎新字段 @ rtp_llm/flexlb/flexlb-common/src/test/java/org/flexlb/dao/master/WorkerIdentityTest.java:32
    • 建议:将该用例移入新建的 WorkerHostTest 并补 getWorkerStatusPort() / getMultiEngineNum() 断言,WorkerIdentityTest 仅保留 WorkerIdentity 自身用例;如 equals/hashCode 用于集合去重或参数匹配,另补一条等值用例。
  • 端口上界常量与既有字面量重复,端口范围校验错误语义不统一 @ rtp_llm/flexlb/flexlb-common/src/main/java/org/flexlb/config/ModelServiceConfiguration.java:23
    • 建议:将 TCP 端口上界提升为共享常量(本 PR 已在 CommonConstants 增补 LOGICAL_WORKER_ENGINE_INDEX_SEPARATOR,可顺带放入,或提供 isValidTcpPort 小工具)供两处共用;并考虑在 validateOptimizer 中对端口补一致的范围校验,让配置面统一走启动期 fail-fast。
  • address 为空时端点级校验的错误信息退化为 null @ rtp_llm/flexlb/flexlb-common/src/main/java/org/flexlb/config/ModelServiceConfiguration.java:65
    • 建议:在 validateEngineEndpointConfiguration 开头先校验 endpoint.getAddress() 非空,或在拼接消息时用 group/role 等可辨识信息兜底,保证任何端点级错误信息都能指向唯一 endpoint。

Checklist Findings (15 fail / 26 total)

General Principles Checklist

  • [6.1] Architecture — 兼容性:外部 HTTP/RPC API、持久数据、配置、环境迁移安全 → issue 路由决策 PV 快照身份迁移不完整,同一记录内 logical 与 physical 混用
    CacheAffinityDecision 的 javadoc(:34-35)已声明两个 worker 字段使用逻辑 ip:port@engineIndexCacheAffinityFirstStrategy:108-109/139-140 也确实改用 getLogicalIpPort()。但同一 record 内 WorkerDecision(:47-50)仍是 ip + port 两个字段、无 engineIndex(ShortestTTFTStrategy:604-607worker.getIp() / getPort())。该记录经 PvLogData.shortestTtftDecisions 落 PV 日志:N>1 时 workers[] 会出现多行 ip/port 完全相同、无法区分引擎的候选;且原先按 ip+":"+port 关联 cacheLeaderIpPort 的离线分析会因新增 @index 后缀失配。
  • [6.1] Architecture — 可观测性:日志/指标/超时可操作、非噪声 → issue CacheMetricsReporterTest 只锁定一个上报方法,其余 engineIp tag 语义变更无覆盖
    该文件是本 PR 为 CacheMetricsReporter 新增的唯一测试,全文只有一条 reportsCacheUpdateLatencyWithIndexedEngineIp(:26-33),断言 reportUpdateEngineBlockCacheRT 的 tag(key 由 engineIpPort 改为 engineIp、值为 10.0.0.8@1)。同一提交把 reportEngineLocalMetrics:127reportCacheHitMetrics:158reportKvcmSelectedMatch:171reportCacheDiffMetrics:353engineIp 取值同样改为 ipIndex,这四个方法以及 ipIndex == null 的早退分支(:120、:348)均无任何断言。多引擎下若某条链路仍传入无索引 IP,指标会把兄弟引擎聚合成一条曲线而测试不失败。
  • [6.1] Architecture — 回滚路径:风险行为存在运维回滚手段 → issue 指标身份契约多处变更(tag key 重命名 + 取值格式 + port 维度丢失)缺少迁移清单
    本 PR 同时做了三类指标契约变更:一是 tag key 重命名,CACHE_UPDATE_ENGINE_BLOCK_CACHE_RTengineIpPort 改为 engineIp(:333),ExpirationCleanerTASK_REMOVEDip 改为 engineIp;二是约 20 个指标的 engineIp 取值由 ip / ip:port 改为 ip@engineIndexreportEngineLocalMetrics:127reportCacheHitMetrics:158reportKvcmSelectedMatch:171reportCacheDiffMetrics:353EngineHealthReporterResourceMonitorReporter 等),N=1 也带 @0;三是 select detail 新增 engineIndex tag。按旧 key/取值精确匹配或按 : 拆端口的看板与告警会静默断流。`06-configuratio
  • [6.1] Architecture — 状态不变量:创建/更新/失败/重试/回滚路径有效 → issue multi_engine_num 缺少显式上限,配置笔误会放大同步扇出并使整机永久不可路由
    校验只拒绝 multiEngineNum < 1(:64),唯一隐式天花板是端口段检查(:84),故 "worker_status_port":1,"multi_engine_num":60000 可通过。该值随后线性放大三处资源:RoutingServiceDiscovery:101discoveredHosts.size() * multiEngineNum、每个逻辑 worker 一条 GrpcWorkerStatusRunner/GrpcCacheStatusCheckRunner 连接、EngineWorkerStatus:91new boolean[expected]。更严重的是 EngineWorkerStatus.isHealthyPhysicalGroup:82-96 要求 siblings 数恰好等于 N,20 误写成 200 时该物理机组永远不健康、被整体摘除路由,而启动期不报错。
  • [6.1] Architecture — 错误语义:fail-fast/retry/fallback/silent 行为显式 → issue address 为空时端点级校验的错误信息退化为 null
    validateServiceRoute 的循环先调用 validateEngineEndpointConfiguration(endpoint)(:54)再调用 serviceDiscovery.validate(endpoint)(:55),而 address 非空校验位于 RoutingServiceDiscovery:70-71。因此当某 endpoint 同时缺失 addressmulti_engine_num 非法时,抛出的是 ... multi_engine_num must be greater than zero: null(:65-67),运维在多 endpoint 配置里无法据此定位是哪个 endpoint 出错;worker_status_port 的两条错误信息(:71-74、:85-88)同样以 address 结尾,存在相同问题。
  • [6.1] Quality — PR description 说明动机与设计 → issue cache-hit PV 的 worker 字段取值由裸 IP 变为逻辑身份,文档却描述为「继续使用」
    result() 构造 CacheHitComparisonResultworker 实参由 feedback.workerIp()(裸 IP,如 10.0.0.1)改为 feedback.logicalWorkerId()10.0.0.1:8080@0,N=1 也带 @0)。该 JSON 由 GrpcWorkerStatusRunner 直接 pvLogger.info(json) 落盘,是命中率离线分析的 worker 维度,所有存量部署该字段取值都会变化。而 06-configuration-and-observability.md:169-170 写的是「路由、KVCM、cache key 和 cache-hit PV 的 worker 字段继续使用完整 ip:port@engineIndex」,把一次取值格式变更描述成保持不变,读者无法据此判断需要改造下游解析,按 IP 精确匹配或聚合该字段的报表/离线任务会在发布后静默失配。
  • [6.1] Software Engineering — DRY:重复非平凡逻辑被抽取或显式复用 → issue 端口上界常量与既有字面量重复,端口范围校验错误语义不统一
    新增私有常量 MAX_TCP_PORT = 65_535(:23),而 flexlb-sync/.../optimizer/OptimizerAddressResolver.java:105-111 对同一领域概念仍用裸字面量 port <= 0 || port > 65535 做等价判定,并以 log.warn 静默跳过。同一约束在两个模块各自编码,后续调整(如收窄到非特权端口)需多点同步;且同为 MODEL_SERVICE_CONFIG 字段,worker_status_port 越界会启动失败,optimizer 端口越界只在运行期告警丢弃,错误语义不一致。
  • [6.1] Software Engineering — ISP:调用方不依赖无关大接口 → issue 已有不可变 WorkerIdentity 却仍以相邻同类型 String 传递三种身份表示
    updateEngineCache(engineIPort, physicalIpPort, ipIndex, role, newCacheBlocks)(:105-110)把已聚合的身份拆回 4 个相邻 String;EngineLocalView.calculateDiff(engineIPort, ipIndex, newCacheBlocks, role)(:57-58)同样有两个相邻同类型 String,其中 ipIndex 不参与任何存储或差分、仅透传给 reportCacheDiffMetrics:85-86,而调用方本已持有该值并直接上报 reportEngineLocalMetricsCacheHitComparisonResult 亦把 workeripIndex 作为相邻 record 组件。任意两处调换实参都能通过编译,后果是指标 tag 维度错位、或 physicalIpPortByLogicalIpPort 记下错误映射进而误删活跃引擎缓存,只能靠单测断言字面量发现。
  • [6.1] Software Engineering — KISS/YAGNI:无投机性抽象 → issue 已有不可变 WorkerIdentity 却仍以相邻同类型 String 传递三种身份表示
    updateEngineCache(engineIPort, physicalIpPort, ipIndex, role, newCacheBlocks)(:105-110)把已聚合的身份拆回 4 个相邻 String;EngineLocalView.calculateDiff(engineIPort, ipIndex, newCacheBlocks, role)(:57-58)同样有两个相邻同类型 String,其中 ipIndex 不参与任何存储或差分、仅透传给 reportCacheDiffMetrics:85-86,而调用方本已持有该值并直接上报 reportEngineLocalMetricsCacheHitComparisonResult 亦把 workeripIndex 作为相邻 record 组件。任意两处调换实参都能通过编译,后果是指标 tag 维度错位、或 physicalIpPortByLogicalIpPort 记下错误映射进而误删活跃引擎缓存,只能靠单测断言字面量发现。
  • [6.1] Software Engineering — SRP:模块/类职责单一 → issue WorkerIdentityTest 混入 WorkerHost 断言,且未覆盖多引擎新字段
    workerHostExposesItsStoredIdentityRepresentations(:32-43)实际验证的是 WorkerHost 的 10 参构造与身份 getter 委托,与文件名声明的被测类不一致;后续维护者按类名检索 WorkerHost 的测试会误判其无覆盖。该用例只断言身份三元组与 engineIndex,未断言构造入参中的 workerStatusPort=18003multiEngineNum=2——这两个值正是多引擎端口派生与 sibling 完整性判定(isHealthyPhysicalGroup:87-96)的输入,位置参数错位不会被发现;WorkerIdentity 由 Lombok 生成的 equals/hashCode(WorkerIdentity.java:22,被 CacheHitFeedback 相等性间接依赖)亦无覆盖。
  • [6.1] Tests — 分布式/跨平台变更有对应覆盖 → issue KvCacheManagerTest 缺少同物理机多逻辑引擎的失效清理覆盖
    新增的 keepsLogicalCacheWhenItsPhysicalEngineRemainsDiscoverable(:51)与 removesLogicalCacheWhenItsPhysicalEngineDisappears(:68)都只注册一个逻辑引擎 10.0.0.1:8080@0。而多引擎正是本次改造目标:removeStaleEngineCaches:158-166 入参是物理地址、内部靠 physicalIpPortByLogicalIpPort 反查,同一物理机下 @0/@1 两条映射必须同时保留或同时清理。当前用例无法发现「只清理了一个兄弟索引」或「映射未按 logical key 建立导致误删存活缓存」;getOrDefault 兜底的「映射缺失且物理机存活」方向亦未覆盖。
  • [6.1] Tests — 新逻辑有聚焦单测 + 相关集成/smoke 测试 → issue 新增 engine-index 用例的期望值与 feedback 字段重合,判别力弱于命名声明
    feedback 传入的 predictedHitTokens=8192deltaHitTokens=808(:133-135),与由 @1 预测(local(2) × blockSize 4096,inputTokens 12000,actual 9000)推导出的 localStandby hit/delta 完全相同。因此 assertEquals(8192, result.localStandby().hit())(:141)虽能排除「错误回退到 index 0」(@0 会得到 4096/4904),却无法区分「按 engineIndex 命中预测」与「直接复制 routing 预测字段」两种实现。
  • [6.1] Tests — 边界 case 覆盖(空、单元素、最大值) → issue WorkerIdentityTest 混入 WorkerHost 断言,且未覆盖多引擎新字段
    workerHostExposesItsStoredIdentityRepresentations(:32-43)实际验证的是 WorkerHost 的 10 参构造与身份 getter 委托,与文件名声明的被测类不一致;后续维护者按类名检索 WorkerHost 的测试会误判其无覆盖。该用例只断言身份三元组与 engineIndex,未断言构造入参中的 workerStatusPort=18003multiEngineNum=2——这两个值正是多引擎端口派生与 sibling 完整性判定(isHealthyPhysicalGroup:87-96)的输入,位置参数错位不会被发现;WorkerIdentity 由 Lombok 生成的 equals/hashCode(WorkerIdentity.java:22,被 CacheHitFeedback 相等性间接依赖)亦无覆盖。

RTP-LLM Checklist

  • [I] 代码质量 — 删除或重命名内部 file、registry entry、model name、metric enum、op binding、plugin symbol 时,必须全仓搜索消费者,并提供替代实现、迁移说明或 smoke 覆盖;只有暴露到 HTTP/RPC/config/persisted format 时才按外部兼容性处理 → issue WorkerCacheUpdateResult 的 javadoc 交叉引用指向已不存在的 WorkerStatus 字段
    新增注释在 :23 引用 org.flexlb.dao.master.WorkerStatus#engineIndex,但本 PR 已把 engineIndexWorkerStatus 移入 WorkerIdentity:全仓检索显示该字段仅声明于 WorkerIdentity.java:31WorkerStatus 只剩委托给 workerIdentitygetEngineIndex():112 / setEngineIndex():116。该引用为悬空链接(当前 pom 未启用 maven-javadoc-plugin,故不影响构建,仅误导读者并在开启 javadoc 校验时告警)。
  • [I] 代码质量 — 同一功能用统一工具函数 → issue 端口上界常量与既有字面量重复,端口范围校验错误语义不统一
    新增私有常量 MAX_TCP_PORT = 65_535(:23),而 flexlb-sync/.../optimizer/OptimizerAddressResolver.java:105-111 对同一领域概念仍用裸字面量 port <= 0 || port > 65535 做等价判定,并以 log.warn 静默跳过。同一约束在两个模块各自编码,后续调整(如收窄到非特权端口)需多点同步;且同为 MODEL_SERVICE_CONFIG 字段,worker_status_port 越界会启动失败,optimizer 端口越界只在运行期告警丢弃,错误语义不一致。

Strengths

  • WorkerIdentity 作为不可变值对象一次性预计算 physical / logical / ipIndex 三种表示(WorkerIdentity.java:39-50),WorkerStatussetIp/setPort/setEngineIndex 整份原子替换 volatile 引用(WorkerStatus.java:35,94-120),从根上消除「三种派生表示来自不同快照」的竞态,优于在各调用点分散拼接字符串。
  • KVCM 兼容策略收敛克制、边界清晰:hostMatch(WorkerStatus) 仅在「logical 未命中 + source == KVCM + multiEngineNum == 1」三条件同时成立时回退 physical key(CacheMatchResult.java:56-67),N>1 与非 KVCM source 一律 exact match,避免把整台 frontend 的命中错误归属到某个 index。
  • 身份切换无半迁移:LOCAL_SYNC 两级索引、LOCAL_STANDBY 候选集/rollback/routed-request 记账、cache-hit 对比查表、worker 状态表、路由回滚全部切到 logical key,写入侧与读取侧同源;全仓已无 getEngineIpPort 残留消费者。
  • 多引擎机制自洽:normalizeHostsworkerStatusBasePort + engineIndex 展开逻辑 worker,两个 gRPC runner 各自连自己的 workerStatusPort,每个逻辑引擎的 cache/worker 状态天然独立,不会出现 N 个逻辑 worker 抓同一份数据。
  • N=1 语义保守:内部保留 @0 身份、仅省略 wire engine_indexServerStatus.java:36-47),rollback 与 cache 归属在单引擎下同样精确,没有留下「索引可空」分支;selectRoutableModelWorkerStatus 在 N=1 时退化为原有 alive 过滤,存量部署行为不变。
  • 配置校验严谨:(long) workerStatusPort + multiEngineNum - 1 > MAX_TCP_PORTModelServiceConfiguration.java:84)显式上转 long 规避 int 溢出绕过;N=1 提前 return 不强制基准端口,与 normalizeHosts:106-108 的 grpcPort 回退闭合,存量单引擎配置零改动可启动;四类错误信息各自独立并附带 endpoint address。
  • 兄弟引擎串号这一最核心风险有正面覆盖:isolatesMatchesBetweenSiblingEngineIndexesappliesCacheMatchRollbackOnlyToMatchingSiblingIndexmatchesLocalStandbyPredictionByEngineIndex 构造同 IP 同端口、仅 engineIndex 不同的两个 worker,实现若退回物理 key 会立即失败;ModelServiceConfigurationTest 成对覆盖了端口段上界两侧(65534+2-1 通过、65535+2-1 拒绝)。
  • 身份语义在 javadoc 中逐参数标注(CacheMatchProviderEngineLocalView.java:51-55KvCacheManager.java:97-103GlobalCacheIndexKvcmCacheMatchProvider),并同步更新 5 篇架构文档的三种身份表示对照表,后续维护者不需要靠猜字符串格式。

}

/**
* Returns KVCM matches without rewriting their keys. KVCM {@code host_ip_port} values follow

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] KVCM host_ip_port 键语义在 proto 契约与 Java 消费方之间自相矛盾

本 PR 把 KVCM 返回键的语义改为逻辑身份:KvcmCacheMatchProvider.java:31-32 的 javadoc 写 "KVCM host_ip_port values follow the logical ip:port@engineIndex identity reported by the Subscriber",04-worker-sync-and-cache.md:103 同步声明。但键的生产方契约定义 flexlb-grpc/src/main/proto/kvcm_meta_service.proto:80 未改,仍写 host_ip_port "Must match WorkerStatus#getIpPort()",而 WorkerStatus.java:588-589 明确返回物理 ip:portKvcmGrpcClient 原样透传该键作 map key,因此 N>1 时 KVCM 侧若按 proto 注释实现即全量 miss。跨服务契约在同一仓库内给出两个互斥定义。

建议: 同步更新 kvcm_meta_service.protohost_ip_port 的注释为逻辑 ip:port@engineIndex(并注明单引擎旧格式在 FlexLB 侧有回退兼容),使 proto 成为唯一契约来源;在 PR description 中把「启用 multi_engine_num > 1 前 KVCM Subscriber 必须先上报逻辑身份」列为显式发布前置条件与灰度顺序。

* worker has an unambiguous logical {@code @0} identity, so it can use that legacy entry
* when no exact logical match exists. Multi-engine workers require an exact logical match.
*/
public HostCacheMatch hostMatch(WorkerStatus workerStatus) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] KVCM 键与逻辑 worker 全不匹配时静默按零命中处理,缺可观测手段

物理回退仅在 source == KVCM && multiEngineNum == 1 生效(:61-66)。若 KVCM 仍上报物理键而运维先把 multi_engine_num 调到 2,则 logical 全 miss 且回退条件不成立,ShortestTTFTStrategy:234/306/889WeightedCacheLoadBalancer:144 拿到 null,整组 worker 的 cache 亲和归零、TTFT 退化为全量 prefill。该路径既无日志也无指标:MetricConstant 仅有 app.cache.match.standby.fallback.qps(语义是 KVCM 不可用)与 app.cache.kvcm.query.retry.qps,既无法区分「确实无命中」与「键格式不匹配」,也没有 legacy 回退命中计数以判断迁移是否完成、兼容分支何时可删。

建议: 补两个计数(或限频 warn):legacy 物理键回退命中次数;「KVCM 返回非空 hostMatches 但无任何键匹配已知 logical worker」的次数(tag 带 multiEngineNum,日志附一个示例键)。并在 04/06 文档把「启用 multi_engine_num > 1 前 KVCM 必须上报 logical key,否则命中恒为 0」列为显式运维前置条件,使灰度阶段即可暴露契约不一致。也可考虑把该兼容判定从 domain record 上移到 KVCM provider 或查询编排层,避免 CacheMatchResult 直接耦合运行态 WorkerStatus

Set<String> activePhysicalIpPorts = new HashSet<>(activeEngineIpPorts);
Set<String> staleEngineIpPorts = new HashSet<>(engineLocalView.getAllEngineIpPorts());
staleEngineIpPorts.removeAll(new HashSet<>(activeEngineIpPorts));
staleEngineIpPorts.removeIf(engineIpPort ->

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] logical 到 physical 映射生命周期不完整,孤儿条目单向泄漏且缺失映射走隐式兜底

physicalIpPortByLogicalIpPort 仅在 updateEngineCache:115-117 写入,在 stale 清理命中(:171)与 clear() 时删除;待清理集合只来自 engineLocalView.getAllEngineIpPorts():163。而 EngineLocalView.removeCacheBlock:138-140 在某引擎 block 集合清空后会 engineViews.remove(engineIPort),该 logical key 此后不再出现于该集合,映射条目永久残留;pod 轮转/扩缩容下条目数按历史 ip:port@index 去重集合单向增长且无指标暴露(功能不受影响,下轮 update 会重新登记)。同时 getOrDefault(engineIpPort, engineIpPort):166 把「映射缺失」隐式当物理地址处理,落入该分支即判 stale 并清空该引擎全部缓存元数据,无日志、无显式分支。

建议: 把待清理集合改为 getAllEngineIpPorts()physicalIpPortByLogicalIpPort.keySet() 的并集再做物理活跃性过滤,让孤儿映射同轮回收;把 getOrDefault 兜底改为显式分支,缺失时先 warn 再决定是否判 stale(当前该分支实际不可达,显式化可防后续调用方绕过登记)。也可改由 WorkerStatusProvider 当前活跃 logical 集合驱动清理以免维护第二份状态,并把该 map size 纳入 app.cache.* 便于线上确认回收生效。

*
* @param activeEngineIpPorts active physical engine addresses in {@code ip:port} format
*/
public void removeStaleEngineCaches(Collection<String> activeEngineIpPorts) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] removeStaleEngineCaches 对空活跃地址列表无保护,服务发现瞬时返回空会清空整份索引

EngineAddressResolver.updateEndpointHosts:96-97 在某 endpoint 的 hostList 为 null/空时 domainHostsMap.remove(endpoint),随后以聚合结果通知 listener;LocalCacheEngineAddressListener.onAddressUpdate:34-38 只拦 null、不拦空列表。因此聚合为空时全部 logical key 都会被判 stale 并清空两级索引,命中率骤降且需一整轮同步才能重建。对比 worker 状态摘除路径存在宽限期(EngineSyncRunner 按实际同步间隔的倍数延迟摘除),本清理路径没有等价保护。该行为在本次改动前已存在,但本 PR 正在重写该方法的 staleness 判定语义,适合一并收敛。

建议: 对空活跃列表直接跳过清理并打 warn,或要求连续多轮观测到缺失后才执行清理,使缓存元数据清理与 worker 摘除采用一致的宽限期语义。

FlexMetricTags tags = FlexMetricTags.of("engineIpPort", engineIpPort, "role", role, "success", success);
public void reportUpdateEngineBlockCacheRT(
String ipIndex, String role, long startTime, String success) {
FlexMetricTags tags = FlexMetricTags.of("engineIp", ipIndex, "role", role, "success", success);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] 指标身份契约多处变更(tag key 重命名 + 取值格式 + port 维度丢失)缺少迁移清单

本 PR 同时做了三类指标契约变更:一是 tag key 重命名,CACHE_UPDATE_ENGINE_BLOCK_CACHE_RTengineIpPort 改为 engineIp(:333),ExpirationCleanerTASK_REMOVEDip 改为 engineIp;二是约 20 个指标的 engineIp 取值由 ip / ip:port 改为 ip@engineIndexreportEngineLocalMetrics:127reportCacheHitMetrics:158reportKvcmSelectedMatch:171reportCacheDiffMetrics:353EngineHealthReporterResourceMonitorReporter 等),N=1 也带 @0;三是 select detail 新增 engineIndex tag。按旧 key/取值精确匹配或按 : 拆端口的看板与告警会静默断流。`06-configura...

建议: 在 06 文档或发布说明中补一节指标迁移清单:列出被重命名的 tag key(engineIpPortengineIpipengineIp)、取值格式变化涉及的指标族(至少 app.cache.*app.engine.*app.routing.*)与新增的 engineIndex tag,并注明回滚该版本时看板需回退旧格式;同时评估同 IP 多 frontend 部署形态,若存在则为丢失 port 维度的两个 cache 指标补独立 port tag 以避免标签冲突。

Checklist: [6.1] 回滚路径:风险行为存在运维回滚手段

@@ -60,8 +62,10 @@ void waitsForStandbyHashBeforeMatching() throws Exception {
}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

📍 实际位置 rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/match/localstandby/LocalStandbyCacheMatchProviderTest.java:48(不在 diff 展示范围内,就近挂载)

[P3] LOCAL_STANDBY 测试夹具仍用物理 key,与生产 key 契约脱节

waitsForStandbyHashBeforeMatchingcacheManager.findMatchingEngines 返回 Map.of("10.0.0.1:8080", 1)(:48),并按 result.hostMatch("10.0.0.1:8080")(:57)断言。但 LocalStandbyCacheManager.findMatchingEngines 的 javadoc 与实现已明确 key 为 ip:port@engineIndex,且 hostMatch(WorkerStatus) 对非 KVCM 来源不再回退物理 key,因此该夹具与生产数据形态不符。同一 PR 内 LocalStandbyComparisonServiceTest 与同文件 updatesRequestDerivedCacheMetadataAsynchronously(:66)已统一使用逻辑 key,此处遗漏容易被后续用例照抄。

建议: 将该用例的 mock key 与断言统一改为逻辑格式(如 10.0.0.1:8080@0),物理 key 用例仅保留在 KVCM 兼容场景;并可把 LocalSyncCacheMatchProviderTest 各用例重复的 mock/provider 构造抽到 @BeforeEach,减少后续身份格式变更的漏改面。


CacheHitFeedback feedback = new CacheHitFeedback(
"cache_hit_comparison", "request-index-1", "LOCAL_STANDBY", "PREFILL", "default",
"10.0.0.1", 8080, 1, "running", 12000, 4096, 8192,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P3] 新增 engine-index 用例的期望值与 feedback 字段重合,判别力弱于命名声明

feedback 传入的 predictedHitTokens=8192deltaHitTokens=808(:133-135),与由 @1 预测(local(2) × blockSize 4096,inputTokens 12000,actual 9000)推导出的 localStandby hit/delta 完全相同。因此 assertEquals(8192, result.localStandby().hit())(:141)虽能排除「错误回退到 index 0」(@0 会得到 4096/4904),却无法区分「按 engineIndex 命中预测」与「直接复制 routing 预测字段」两种实现。

建议: 把 feedback 的 predictedHitTokens / deltaHitTokens 改成与预测推导值不同的数值(例如 4096/4904),使断言只能由「按索引查预测」这一条路径满足。

Checklist: [6.1] 新逻辑有聚焦单测 + 相关集成/smoke 测试

}

@Test
void workerHostExposesItsStoredIdentityRepresentations() {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P3] WorkerIdentityTest 混入 WorkerHost 断言,且未覆盖多引擎新字段

workerHostExposesItsStoredIdentityRepresentations(:32-43)实际验证的是 WorkerHost 的 10 参构造与身份 getter 委托,与文件名声明的被测类不一致;后续维护者按类名检索 WorkerHost 的测试会误判其无覆盖。该用例只断言身份三元组与 engineIndex,未断言构造入参中的 workerStatusPort=18003multiEngineNum=2——这两个值正是多引擎端口派生与 sibling 完整性判定(isHealthyPhysicalGroup:87-96)的输入,位置参数错位不会被发现;WorkerIdentity 由 Lombok 生成的 equals/hashCode(WorkerIdentity.java:22,被 CacheHitFeedback 相等性间接依赖)亦无覆盖。

建议: 将该用例移入新建的 WorkerHostTest 并补 getWorkerStatusPort() / getMultiEngineNum() 断言,WorkerIdentityTest 仅保留 WorkerIdentity 自身用例;如 equals/hashCode 用于集合去重或参数匹配,另补一条等值用例。

Checklist: [6.1] SRP:模块/类职责单一;[6.1] 边界 case 覆盖(空、单元素、最大值)

public class ModelServiceConfiguration {

private static final String MODEL_SERVICE_CONFIG = "MODEL_SERVICE_CONFIG";
private static final int MAX_TCP_PORT = 65_535;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P3] 端口上界常量与既有字面量重复,端口范围校验错误语义不统一

新增私有常量 MAX_TCP_PORT = 65_535(:23),而 flexlb-sync/.../optimizer/OptimizerAddressResolver.java:105-111 对同一领域概念仍用裸字面量 port <= 0 || port > 65535 做等价判定,并以 log.warn 静默跳过。同一约束在两个模块各自编码,后续调整(如收窄到非特权端口)需多点同步;且同为 MODEL_SERVICE_CONFIG 字段,worker_status_port 越界会启动失败,optimizer 端口越界只在运行期告警丢弃,错误语义不一致。

建议: 将 TCP 端口上界提升为共享常量(本 PR 已在 CommonConstants 增补 LOGICAL_WORKER_ENGINE_INDEX_SEPARATOR,可顺带放入,或提供 isValidTcpPort 小工具)供两处共用;并考虑在 validateOptimizer 中对端口补一致的范围校验,让配置面统一走启动期 fail-fast。

Checklist: [6.1] DRY:重复非平凡逻辑被抽取或显式复用;[I] 同一功能用统一工具函数

private void validateEngineEndpointConfiguration(Endpoint endpoint) {
int multiEngineNum = endpoint.getMultiEngineNum();
if (multiEngineNum < 1) {
throw new IllegalArgumentException(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P3] address 为空时端点级校验的错误信息退化为 null

validateServiceRoute 的循环先调用 validateEngineEndpointConfiguration(endpoint)(:54)再调用 serviceDiscovery.validate(endpoint)(:55),而 address 非空校验位于 RoutingServiceDiscovery:70-71。因此当某 endpoint 同时缺失 addressmulti_engine_num 非法时,抛出的是 ... multi_engine_num must be greater than zero: null(:65-67),运维在多 endpoint 配置里无法据此定位是哪个 endpoint 出错;worker_status_port 的两条错误信息(:71-74、:85-88)同样以 address 结尾,存在相同问题。

建议:validateEngineEndpointConfiguration 开头先校验 endpoint.getAddress() 非空,或在拼接消息时用 group/role 等可辨识信息兜底,保证任何端点级错误信息都能指向唯一 endpoint。

Checklist: [6.1] 错误语义:fail-fast/retry/fallback/silent 行为显式

@Martin7-1
Martin7-1 removed the request for review from jianglan89 August 31, 2026 05:55
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants