Skip to content

feat(config): migrate config to sdk-proto and rebuild the listen state machine - #914

Open
cxhello wants to merge 12 commits into
nacos-group:v3.x-devfrom
cxhello:feat/v3-config
Open

cxhello wants to merge 12 commits into
nacos-group:v3.x-devfrom
cxhello:feat/v3-config

Conversation

@cxhello

@cxhello cxhello commented Sep 1, 2026

Copy link
Copy Markdown
Member

What is this PR for

PR5a of the v3 roadmap (#879, sub-issue #913): migrates the Config module's wire types to nacos-sdk-proto and structurally rebuilds the config listen state machine. Config-side FuzzyWatch (#859) follows separately as PR5b on top of this branch.

Proto migration (same adapter pattern as PR2/PR4)

  • Requests (ConfigQueryRequest, ConfigPublishRequest, ConfigRemoveRequest, ConfigBatchListenRequest) implement ProtoMessage() and encode through PayloadCodec; a wire-parity test pins the protojson key set against the legacy JSON body (guarding the json_name defect class found in PR4).
  • Responses (ConfigQueryResponse, ConfigPublishResponse, ConfigRemoveResponse, ConfigChangeBatchListenResponse) and the ConfigChangeNotifyRequest server push decode through proto_dispatch adapters back into the legacy structs, so downstream code is unchanged. success is derived from resultCode (proto has no such field), matching the PR2 rule.
  • The legacy-JSON fallback path keeps permanent test coverage via a non-migratable stub request type (it previously piggybacked on ConfigQueryRequest, which this PR migrates).
  • Note: legacy ConfigQueryResponse.Tag is a dead bool field (proto's tag is a string); the adapter intentionally does not map it.

Listen state machine rebuild (Java ClientWorker/CacheData parity)

The by-value cacheData copies stored in cache.ConcurrentMap (copy-mutate-set races, single silently-overwritten listener) are replaced by a pointer-based configCacheHolder with per-entry locking:

  • Multiple listeners per config. cacheData.listeners is a list of wraps, each with its own lastCallMd5 watermark. Delivery snapshots state under the entry lock, releases it, then per listener: filter-chain decrypt → recover-wrapped callback → watermark advances only on normal return. A panicking listener is replayed next round and never blocks its siblings; a filter-chain error skips the round without advancing any watermark.
  • Cancel actually reaches the server (cancelling listening is invalid when there is only one configuration #629). CancelListenConfig marks the entry discard and rings the bell; the single executor goroutine sends ConfigBatchListenRequest{listen=false} batches (grouped by taskId, before the listen batches) and only removes an entry after the server confirmed — re-checking under lock that no listener re-registered while the cancel was in flight. This also fixes the corner where cancelling the only listened config sent nothing at all.
  • Listen/Close linearization. ListenConfig re-checks isClosed under the client mutex after committing its listener and rolls back if a concurrent CloseClient won the race; the closed path returns the exported sentinel ErrConfigClientClosed. The bell send selects against the client context so post-close notifications cannot leak goroutines.
  • The executor keeps the existing cadence (bell + 5s poll + 3min full resync) and the taskId sharding (3000 configs per rpc client). Reconnect handling is unchanged: disconnect marks entries out-of-sync, reconnect rings the bell.

Issues

Behavior changes vs v2 (please review deliberately)

  1. Repeated ListenConfig on the same key now appends listeners (Java parity). v2 silently dropped every callback after the first. Code that called ListenConfig repeatedly as an idempotent re-registration (e.g. in reconnect loops) will now accumulate listeners and receive duplicate callbacks — this is the most user-visible change in the PR.
  2. CancelListenConfig now notifies the server and removes all listeners of the key (same key-level granularity as v2's local removal). Watcher-scoped cancellation (handle-based, like FuzzyWatchHandle) is deliberately deferred — proposed for Subscribe/ListenConfig jointly in Repeated subscriptions result in multiple callbacks #655.
  3. ListenConfig on a closed client returns ErrConfigClientClosed instead of silently registering into a dead client.
  4. The push handler acks unknown keys with a success response instead of returning nothing.

Known inherited limitation

A server push that lands between an executor round's response and its sync-flag update can be overwritten and only recovered by the 3-minute full resync (the window is microseconds; the previous code had a strictly larger version of the same window across all task shards). Kept as-is, documented here.

Testing

  • TDD throughout; unit suites -race -count=1 green, including targeted concurrency tests (listen/cancel hammer with invariant checks, cancel-during-drain, revive-during-cancel-in-flight, panic replay, watermark skip).
  • Integration (-tags=integration -race -count=1) green against both Nacos 2.5.2 and 3.2.0, including the new multi-listener / cancel / CAS / server-restart tests.

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
…proto

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
Signed-off-by: cxhello <caixiaohuichn@gmail.com>
…-listener support

Replace the by-value cacheData stored in cache.ConcurrentMap with a
pointer-based configCacheHolder (clients/config_client/config_cache_holder.go),
fixing the long-standing bug where a second ListenConfig call on the same
dataId/group/tenant silently dropped the first listener instead of appending
to it.

- cacheData is now always referenced by pointer and carries its own mutex
  plus a listeners []*listenerWrap slice, so multiple independent listeners
  can be registered on the same key.
- ListenConfig now does get-or-create + an atomic reviveAndAddListener; an
  existing entry is revived (discard=false) rather than replaced, with the
  revive and the listener append happening in a single critical section so a
  concurrent CancelListenConfig can't interleave between them and leave the
  entry discard==true with a non-empty listeners slice. CancelListenConfig
  marks the entry discarded instead of removing it outright, and is a no-op
  on unknown keys.
- configCacheHolder.removeIfDiscarded re-checks discard && no-listeners under
  lock before deleting, so a concurrent revive/addListener can't race a
  removal.
- executeConfigListen/buildListenTask/refreshContentAndCheck and
  config_connection_event_listener.go/config_proxy.go are adapted to read
  from the holder; their single-batch listen=true behavior and per-listener
  notify-on-md5-change semantics are kept equivalent to today's behavior.
  Cancel-batch/listen=false executor semantics and real per-listener
  watermark delivery are left for follow-up tasks.
- Removed the dead ConfigClient.localConfigs field.

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
…oup#629)

executeConfigListen now partitions each round into two independent
batches: discarded entries are sent as a Listen=false batch (grouped
by taskId) before any Listen=true batch, and are only reaped from the
holder via removeIfDiscarded once that batch gets a successful
response; a transport error or non-success response leaves them in
place for the next round to retry. This fixes CancelListenConfig never
actually telling the server to stop pushing a key.

Listen=true batches keep the existing changedConfigs handling
(refreshContentAndCheck) but now mark isSyncWithServer based on the
batch's own caches instead of the full holder snapshot, avoiding
cross-taskId contamination between concurrent batches.

Also make the config-change push handler ack cache misses with a
success response instead of returning nil, matching the documented
"push handler only marks isSyncWithServer=false and rings the bell"
contract.

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
Add cacheData.notifyListeners(chain): snapshots md5/content/encryptedDataKey
and not-yet-caught-up listener wraps under cData.mu, then releases the lock
before running the filter chain and invoking callbacks. Each callback is
recover-wrapped so a panicking listener neither crashes the executor
goroutine nor blocks delivery to other listeners on the same key, and its
watermark only advances to the snapshotted md5 on a normal return -- a panic
leaves it unchanged so the content is replayed on the next round. A filter
chain error skips the whole round without advancing any watermark.

Replace the Task 3 interim notifyListenersIfChanged path (whole-entry,
fire-and-forget goroutines) with this per-listener delivery and rewire
refreshContentAndCheck to call cData.notifyListeners directly.

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
ListenConfig previously registered a listener even after CloseClient
had already torn down the client's context and rpc connection,
leaving the caller with a listener that would never be served again.
Check isClosed under client.mutex at ListenConfig entry (same nacos-group#904
semantics used for naming) and return an error instead. Also add
regression coverage for CancelListenConfig and double CloseClient
staying panic-free after close.

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
… coverage

Add TestIntegrationConfig{MultiListener,CancelStopsPush,CasPublish,
ListenSurvivesRestart} covering: two independent listeners on the same
key both firing on publish; CancelListenConfig actually stopping
further pushes (nacos-group#629); PublishConfig's CasMd5 rejecting stale writes
while leaving content untouched and accepting a correct cas (nacos-group#727);
and a listener surviving a full server restart (nacos-group#694). The restart
test is skipped unless NACOS_CONTAINER_NAME is set and shells out to
plain `docker restart`, matching reconnect_test.go's existing pattern.

Verified locally against both probe-nacos3 (3.2.0, port 8848) and
probe-nacos2 (2.5.2, port 8858) with -race -count=1.

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
ListenConfig checked isClosed under client.mutex but committed the new
listener (getOrCreate + reviveAndAddListener) outside that lock, so a
CloseClient landing in between registered a live listener on an
already-shut-down client while still returning nil. Re-check isClosed
under client.mutex right after the commit and roll back via
cData.markDiscard() when a close won the race.

Export ErrConfigClientClosed (mirroring naming's
ErrFuzzyWatchClientClosed) and return it from both ListenConfig closed
paths so callers can errors.Is instead of matching error text.

Fix asyncNotifyListenConfig leaking a goroutine per call once the
client is closed: it blocked forever sending on listenExecute after
the listen executor loop had already exited via ctx.Done(). Select the
send against ctx.Done() too.

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
TestListenConfigRaceWithCloseClientNeverLeavesLiveListenerOnClosedClient
asserted a <=10% threshold on how often a live listener survived a
nil-error ListenConfig racing CloseClient. That "benign occurrence"
rate is machine-dependent and flaked in CI (52/500 = 10.4% on GitHub
runners), so drop the percentage assertion entirely.

The ghost-listener semantics are already pinned deterministically by
TestListenConfigCommitRollsBackWhenCloseWinsRace. This test now keeps
only the concurrent ListenConfig/CloseClient hammer loop as a -race
exerciser for the Finding 1 locking, with scheduling-independent
invariants: the entry must never be observed with discard==true and
non-empty listeners, and once CloseClient has definitely returned, a
further ListenConfig call must fail with errors.Is(err,
ErrConfigClientClosed) and must not grow the listener count.

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
@codecov-commenter

codecov-commenter commented Sep 1, 2026 •

Copy link
Copy Markdown

⚠️ Please install the 'codecov app svg image' to ensure uploads and comments are reliably processed by Codecov.

Codecov Report

❌ Patch coverage is 95.65217% with 13 lines in your changes missing coverage. Please review.
⚠️ Please upload report for BASE (v3.x-dev@93a9350). Learn more about missing BASE report.

Files with missing lines Patch % Lines
clients/config_client/config_proxy.go 0.00% 7 Missing ⚠️
clients/config_client/config_client.go 94.64% 5 Missing and 1 partial ⚠️
❗ Your organization needs to install the Codecov GitHub app to enable full functionality.
Additional details and impacted files
@@             Coverage Diff             @@
##             v3.x-dev     #914   +/-   ##
===========================================
  Coverage            ?   40.72%           
===========================================
  Files               ?      102           
  Lines               ?     7206           
  Branches            ?        0           
===========================================
  Hits                ?     2935           
  Misses              ?     4095           
  Partials            ?      176           

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@cxhello
cxhello requested a review from Sunrisea September 1, 2026 07:01
@Sunrisea

Copy link
Copy Markdown
Member

这里建议合并前修复一个 P1 回归:不要在配置监听主循环中同步执行用户回调。

notifyListeners 这里直接调用 deliverAndAdvance,调用链是 startInternal → executeConfigListen → refreshContentAndCheck → notifyListeners → listener。同一个 ConfigClient 只有一个监听执行循环,所以影响不只是同一配置下的多个 listener:任意一个慢回调都会阻塞其他配置的查询、通知,以及后续取消订阅请求。

针对当前 head d041737,用 mock RPC 驱动真实 startInternal 循环做了定向测试(-race):

  • 让一轮响应返回 A、B 两个变更,A 回调等待 B 回调通知:B 的查询都不会开始,形成循环等待。
  • A 阻塞期间调用 CancelListenConfig(C),API 返回 nil,但服务端的 listen=false 请求一直无法发出;释放 A 后,B 通知和 C 的取消请求才继续执行。
  • 同期模拟 40 次配置推送,推送均得到成功响应,但 asyncNotifyListenConfig 的待发送 goroutine 持续积压,测试中观察到 41 个。

在 base 93a9350 上,旧实现用 goroutine 执行回调,同样的“A 等待 B 通知”用例可以通过,不需要外部释放 A。因此这是已有 Go 行为的回归,而不只是同步/异步的风格选择。Java 默认回调可以同步执行,但它提供了 Executor 选择,不能据此直接改变 Go 的回调执行语义。

建议保留异步回调,并用每个 listener 的 in-flight 状态避免重复并发投递,成功后再推进 lastCallMd5。请补充覆盖真实监听循环的阻塞回调回归测试,验证一个配置的回调不会阻塞其他配置的处理。

The listen executor is the single goroutine driving queries, notifications
and cancel batches for every config; delivering callbacks inline on it was
a regression from the pre-v3 behavior (which launched every callback with
go), letting one slow callback stall all other configs and piling up one
blocked bell goroutine per pending notification.

Callbacks now run on their own delivery goroutine, gated per listener by an
inFlight flag so a single listener is never invoked concurrently with
itself and a slow listener holds exactly one goroutine. The watermark
advances only on normal return, in the same critical section that clears
the flag; a wrap left trailing the entry's md5 (newer change during the
callback, or a panic) rings the bell and the executor's round-entry sweep
re-notifies promptly. The bell itself is now a 1-buffered coalescing
channel with a non-blocking send, removing the goroutine-per-notification
pile-up entirely.

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
@cxhello

cxhello commented Sep 18, 2026

Copy link
Copy Markdown
Member Author

Verified and fixed — this was a real regression, not a style choice. The base implementation launched every callback with go listener(...), while this branch ran them inline on the single listen executor goroutine, so one slow callback stalled every other config's queries, notifications and cancel batches, and each pending bell parked a goroutine. Both of your test scenarios reproduced exactly as described.

Fix (commit 9dc559b), along the lines you suggested:

  • Callbacks run on their own delivery goroutine, never on the executor. notifyListeners snapshots the entry under its lock, then launches one goroutine per lagging listener.
  • Per-listener inFlight gate. A wrap whose previous delivery has not returned is skipped by subsequent rounds, so a single listener is never invoked concurrently with itself and a slow listener holds exactly one goroutine no matter how many changes arrive meanwhile. lastCallMd5 advances only when the callback returns normally, in the same critical section that clears inFlight.
  • Catch-up path. When a delivery completes and the wrap still trails the entry's current md5 (a newer change landed while the callback ran, or it panicked), the delivery rings the bell and the executor's round-entry sweep re-notifies — so a lagging listener catches up promptly instead of waiting out the poll interval.
  • Coalescing bell. listenExecute is now a 1-buffered channel with a non-blocking send (same shape as the naming holder's Bell), which removes the goroutine-per-notification pile-up you observed entirely — your 40-push scenario now allocates zero goroutines for pending bells.

Regression tests drive the real startInternal loop against a scripted proxy (-race, ×5):

  • TestBlockingCallbackDoesNotBlockOtherConfigs — A's blocked callback; B's change in a later round is still delivered while A is blocked (fails on the previous head, loop stalled).
  • TestCancelProceedsWhileCallbackBlocked — CancelListenConfig gets its listen=false batch sent while A is blocked (fails on the previous head).
  • TestBellStormDoesNotAccumulateGoroutines — 40 notifications with no executor running: goroutine count stays flat (fails on the previous head with +40).
  • TestNoConcurrentDuplicateDeliveryPerListener and TestLaggingWatermarkRedeliveredAfterInFlightCompletes pin the in-flight gate and the catch-up semantics.

Full unit suite and both-version integration (2.5.2 / 3.2.0, -race -count=1) are green.

@Sunrisea

Copy link
Copy Markdown
Member

异步回调这部分修复已复核:相关回归测试本地 -race -count=5 通过,上次指出的跨配置阻塞问题已解决。

还有一个 P2 建议一并修复:持续 panic 的 listener 会触发无间隔重试。

deliverAndAdvance 的这段逻辑在回调 panic 后不推进 lastCallMd5,但仍因 lagging 立即调用 wake();下一轮入口扫描又投递同一内容,形成“panic → wake → 再投递 → panic”的忙循环。容量为 1 的 bell 只能合并待处理通知,无法打断这种逐轮自唤醒。

本地用真实 startInternal 循环做了有界复现(-race):缓存已与服务端同步,没有新推送,listener 前 99 次调用均 panic,第 100 次正常返回以终止探针,约 5ms 就完成了 100 次调用,并反复输出异常堆栈。持续 panic 时这条路径会一直消耗 CPU、刷错误日志。

建议采用最小修改:保留异步投递、inFlight 和每轮补发检查,只收紧主动唤醒条件:

// 在 c.mu 内计算,在锁外调用 wake。
shouldWake := delivered && md5 != c.md5 && !c.discard
  • 成功且执行期间出现了新版本:推进本次水位,立即唤醒以追上最新内容。
  • panic:保留旧水位、清除 inFlight,不因失败主动唤醒,交给正常监听轮次重试。
  • 过滤器失败:保持现有的保留水位、下一轮重试行为。

这与 Java SDK 的失败处理思路一致:CacheData 捕获异常后只清理执行状态,不立即唤醒,后续由 ClientWorker 的监听循环再次检查。没有其他事件时通常等约 5 秒;其他推送仍可能提前触发重试,这不是严格的失败冷却时间,也无需为此新增复杂的重试队列或退避机制。

建议补两个真实监听循环测试:持续 panic 且没有新事件时不会高速自重试;失败回调恢复后,即使服务端配置未再变化,也能在后续轮询中补发成功。现有“成功回调完成后追上新版本”的测试继续保留。

…elf-waking

A failed delivery leaves its watermark lagging by construction, and waking
on that lag redelivers the same content immediately -- which panics again,
a self-sustaining busy loop the coalescing bell cannot damp (measured
12k callback invocations in 600ms). Wake now fires only when a delivery
succeeded and the entry's md5 moved past the delivered snapshot meanwhile;
failed deliveries are retried by the regular poll round's re-notify sweep,
matching the Java client's CacheData/ClientWorker failure handling.

Signed-off-by: cxhello <caixiaohuichn@gmail.com>
@cxhello

cxhello commented Sep 23, 2026

Copy link
Copy Markdown
Member Author

Confirmed and fixed in 48f3294. Reproduced first: the new busy-loop regression test (real startInternal loop, synced entry, persistently panicking listener, one initial bell) recorded 12,093 callback invocations in 600ms on the previous head — the panic→wake→redeliver spin exactly as you described.

The fix is your minimal change, applied verbatim:

shouldWake := delivered && md5 != c.md5 && !c.discard

computed inside c.mu, wake() invoked outside the lock. A successful delivery that fell behind a concurrent update still wakes immediately (the round-1 catch-up semantics are unchanged — TestLaggingWatermarkRedeliveredAfterInFlightCompletes still passes); a panicking delivery keeps its watermark, clears inFlight, and is retried only by the regular poll round's re-notify sweep; the filter-error path is untouched. Comments on deliverAndAdvance/notifyListeners/the executor sweep now state the failed-delivery-rides-the-poll contract with the Java CacheData/ClientWorker parity rationale.

Both requested tests drive the real startInternal loop against a scripted proxy (-race, ×5):

  • TestPersistentlyPanickingListenerDoesNotBusyLoop — synced entry, no new events: a panicking listener gets at most the initial delivery within a 600ms window (12,093 → 1).
  • TestRecoveredListenerIsRedeliveredByPollWithoutServerChange — the listener panics once and then recovers: with no further server change and no further bell, the poll round redelivers successfully and the watermark advances.

Full unit suite and both-version integration (2.5.2 / 3.2.0, -race -count=1) are green.

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.

3 participants