Files
milvus/internal/querycoordv2/ddl_callbacks_alter_load_info_test.go
e0ae33a599 fix: ack stale alter-load-config on dropped-sentinel channel (#50161)
## Summary

Follow-up to #50002.

After #50002 added the dropped-channel-sentinel guard in
`TargetManager.UpdateCollectionNextTarget`, the new error bubbles up
through `LoadCollectionJob.Execute` and back into the streamingcoord
broadcaster's `alterLoadConfigV2AckCallback`. That callback has no
terminal condition for "channel was dropped" and therefore retries the
broadcast every 5–10 s **forever**, which:

- keeps the broadcast task uncommitted on the secondary,
- blocks rootcoord's collection meta cleanup,
- leaves the half-dropped collection visible in `list_collections` on
the CDC replica even after the primary has fully dropped it.

This PR introduces `merr.ErrChannelDroppedSentinel` so the sentinel
rejection can be distinguished from other `ChannelNotAvailable` causes,
returns it from `TargetManager.UpdateCollectionNextTarget`, and
acks-with-no-op in `alterLoadConfigV2AckCallback` when it appears. The
behavior mirrors the existing `ErrCollectionNotFound` early-return
inside `LoadCollectionJob`, and is safe because
`IsDroppedChannelCheckpoint` is sticky — it is only ever written by
`datacoord.MarkChannelCheckpointDropped` and `UpdateChannelCheckpoint`
refuses to overwrite it (per the design notes of #50002).

### Recovery

The fix is self-healing: the next retry attempt of any stuck
`AlterLoadConfig` broadcast on a CDC secondary will hit the new code
path, return `nil`, and let the broadcast tombstone. Pending rootcoord
meta cleanup then proceeds and the half-dropped collection disappears
from `list_collections`. No manual intervention required on existing
stuck instances after deploy.

issue: #49996

## Test plan

- [x] `pkg/util/merr` `TestWrap`: adds assertions for
`WrapErrChannelDroppedSentinel` and confirms `errors.Is` distinctness
from `ErrChannelNotAvailable` (passes locally)
- [x] `internal/querycoordv2/meta` `TestTargetManager`: 3 existing
sub-tests from #50002 updated to assert the new error type
- [x] `go build` and `go vet` clean on `internal/querycoordv2/...` and
`pkg/util/merr/...`
- [x] `golangci-lint` clean on touched packages (`new-from-rev`)
- [ ] CI

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Signed-off-by: bigsheeper <yihao.dai@zilliz.com>
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-01 16:58:20 +08:00

81 lines
2.9 KiB
Go

// Licensed to the LF AI & Data foundation under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package querycoordv2
import (
"context"
"testing"
"github.com/bytedance/mockey"
"github.com/cockroachdb/errors"
"github.com/stretchr/testify/assert"
"github.com/milvus-io/milvus/internal/querycoordv2/job"
"github.com/milvus-io/milvus/pkg/v3/proto/messagespb"
"github.com/milvus-io/milvus/pkg/v3/streaming/util/message"
"github.com/milvus-io/milvus/pkg/v3/util/merr"
"github.com/milvus-io/milvus/pkg/v3/util/paramtable"
)
func buildAlterLoadConfigBroadcastResult(collectionID int64) message.BroadcastResultAlterLoadConfigMessageV2 {
controlChannel := "_ctrl_channel"
broadcastMsg := message.NewAlterLoadConfigMessageBuilderV2().
WithHeader(&messagespb.AlterLoadConfigMessageHeader{
CollectionId: collectionID,
Replicas: []*messagespb.LoadReplicaConfig{
{ReplicaId: 1, ResourceGroupName: "__default_resource_group"},
},
}).
WithBody(&messagespb.AlterLoadConfigMessageBody{}).
WithBroadcast([]string{controlChannel}).
MustBuildBroadcast()
return message.BroadcastResultAlterLoadConfigMessageV2{
Message: message.MustAsBroadcastAlterLoadConfigMessageV2(broadcastMsg),
Results: map[string]*message.AppendResult{
controlChannel: {},
},
}
}
// TestAlterLoadConfigV2AckCallback verifies that the ack callback swallows the
// dropped-sentinel error (so the broadcaster stops retrying forever) while still
// propagating any other error from the load job.
func TestAlterLoadConfigV2AckCallback(t *testing.T) {
paramtable.Init()
ctx := context.Background()
s := &Server{}
result := buildAlterLoadConfigBroadcastResult(1000)
mockey.PatchConvey("dropped sentinel is acked with no-op", t, func() {
mockey.Mock((*job.LoadCollectionJob).Execute).
Return(merr.WrapErrChannelDroppedSentinel("_ctrl_channel")).Build()
err := s.alterLoadConfigV2AckCallback(ctx, result)
assert.NoError(t, err)
})
mockey.PatchConvey("generic error is propagated", t, func() {
expectedErr := errors.New("broker unavailable")
mockey.Mock((*job.LoadCollectionJob).Execute).Return(expectedErr).Build()
err := s.alterLoadConfigV2AckCallback(ctx, result)
assert.Error(t, err)
assert.True(t, errors.Is(err, expectedErr))
})
}