Files
milvus/pkg/streaming/walimpls/impls/walimplstest/wal.go
T
ababab0d76 feat: [2.6.16] support partial update ops for Array fields (#49811)
pr: [#49328](https://github.com/milvus-io/milvus/pull/49328)
pr: [#49724](https://github.com/milvus-io/milvus/pull/49724)
pr: [#49763](https://github.com/milvus-io/milvus/pull/49763)
pr: [#49698](https://github.com/milvus-io/milvus/pull/49698)
issue: [#49241](https://github.com/milvus-io/milvus/issues/49241)
issue: [#49746](https://github.com/milvus-io/milvus/issues/49746)
issue: [#49634](https://github.com/milvus-io/milvus/issues/49634)

## Summary

Backport the 2.6 partial update op series to `hotfix-2.6.16`:

- support `ARRAY_APPEND` and `ARRAY_REMOVE` partial update ops for Array
fields
- expose `fieldOps` through REST upsert
- preserve existing Array rows when an op payload row is null

## Target branch

Base branch: `hotfix-2.6.16`

## Cherry-picks

- `5849977c408be2abd063d13c318e970bb4515f06` from
[#49328](https://github.com/milvus-io/milvus/pull/49328)
- `017ee8e5d97ec891c442489b1354b60676ca3b15` from
[#49724](https://github.com/milvus-io/milvus/pull/49724)
- `4e10473389843ff2cedc1911083af5a93fa83ec6` from
[#49763](https://github.com/milvus-io/milvus/pull/49763)
- `824c642c71154952bf70eaaf7b57cbcd7d08e67f` backports the applicable
WAL test/recovery stabilization from
`73dc8d4034fd352f5c69fd47266fccc062704feb` /
[#49698](https://github.com/milvus-io/milvus/pull/49698) after omitting
newer rate-limit API changes that do not exist on `hotfix-2.6.16`

The partial update cherry-picks applied cleanly on top of
`milvus/hotfix-2.6.16`.

## Additional revert

- `2753c8defc7a860446d1f5f76b11a50c3cc71550` reverts
  `8ae21f715094abc92e00369e89a7208681eecfed` to restore the 2.6.16
build environment image version after CI reported Conan 2.x in the newer
  image.

## Validation

- `git diff --check milvus/hotfix-2.6.16...HEAD`
- attempted targeted Go test for
`internal/streamingnode/server/wal/adaptor`, blocked locally because
this worktree lacks `rocksdb.pc` and `milvus_core.pc`

PR CI is the validation gate for this backport.

---------

Signed-off-by: Wei Liu <wei.liu@zilliz.com>
Signed-off-by: Zhen Ye <chyezh@outlook.com>
Co-authored-by: Zhen Ye <chyezh@outlook.com>
2026-05-15 08:37:30 +08:00

100 lines
2.7 KiB
Go

//go:build test
// +build test
package walimplstest
import (
"context"
"math/rand"
"github.com/cockroachdb/errors"
"go.uber.org/atomic"
"github.com/milvus-io/milvus/pkg/v2/proto/streamingpb"
"github.com/milvus-io/milvus/pkg/v2/streaming/util/message"
"github.com/milvus-io/milvus/pkg/v2/streaming/util/types"
"github.com/milvus-io/milvus/pkg/v2/streaming/walimpls"
"github.com/milvus-io/milvus/pkg/v2/streaming/walimpls/helper"
"github.com/milvus-io/milvus/pkg/v2/util/typeutil"
)
var (
_ walimpls.WALImpls = &walImpls{}
fenced = typeutil.NewConcurrentSet[string]()
enableFenceError = atomic.NewBool(true)
)
// Reset clears global state of the in-memory WAL test implementation.
func Reset() {
logs = typeutil.NewConcurrentMap[string, *messageLog]()
fenced = typeutil.NewConcurrentSet[string]()
enableFenceError.Store(true)
}
// EnableFenced enables fenced mode for the given channel.
func EnableFenced(channel string) {
fenced.Insert(channel)
}
// DisableFenced disables fenced mode for the given channel.
func DisableFenced(channel string) {
fenced.Remove(channel)
}
type walImpls struct {
helper.WALHelper
datas *messageLog
}
func (w *walImpls) WALName() message.WALName {
return message.WALNameTest
}
func (w *walImpls) Append(ctx context.Context, msg message.MutableMessage) (message.MessageID, error) {
if w.Channel().AccessMode != types.AccessModeRW {
panic("write on a wal that is not in read-write mode")
}
if fenced.Contain(w.Channel().Name) {
return nil, errors.Mark(errors.New("err"), walimpls.ErrFenced)
}
if enableFenceError.Load() && msg.MessageType() != message.MessageTypeTimeTick && rand.Int31n(30) == 0 {
return nil, errors.New("random error")
}
return w.datas.Append(ctx, msg)
}
func (w *walImpls) Read(ctx context.Context, opts walimpls.ReadOption) (walimpls.ScannerImpls, error) {
offset := int64(0)
switch t := opts.DeliverPolicy.GetPolicy().(type) {
case *streamingpb.DeliverPolicy_All:
offset = 0
case *streamingpb.DeliverPolicy_Latest:
offset = w.datas.Len()
case *streamingpb.DeliverPolicy_StartFrom:
id, err := unmarshalTestMessageID(t.StartFrom.GetId())
if err != nil {
return nil, err
}
offset = int64(id)
case *streamingpb.DeliverPolicy_StartAfter:
id, err := unmarshalTestMessageID(t.StartAfter.GetId())
if err != nil {
return nil, err
}
offset = int64(id) + 1
}
return newScannerImpls(
opts, w.datas, int(offset),
), nil
}
func (w *walImpls) Truncate(ctx context.Context, id message.MessageID) error {
if w.Channel().AccessMode != types.AccessModeRW {
panic("truncate on a wal that is not in read-write mode")
}
return nil
}
func (w *walImpls) Close() {
}