Files
wehub-resource-sync 498b235461
Build and test / Build and test AMD64 Ubuntu 22.04 (push) Failing after 0s
Publish Builder / amazonlinux2023 (push) Failing after 1s
Build and test / UT for Go (push) Has been skipped
Publish KRTE Images / KRTE (push) Failing after 1s
Build and test / Integration Test (push) Has been skipped
Build and test / Upload Code Coverage (push) Has been skipped
Publish Builder / rockylinux9 (push) Failing after 1s
Publish Builder / ubuntu22.04 (push) Failing after 0s
Publish Builder / ubuntu24.04 (push) Failing after 0s
Publish Gpu Builder / publish-gpu-builder (push) Failing after 1s
Publish Test Images / PyTest (push) Failing after 0s
Build and test / UT for Cpp (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 12:31:17 +08:00

145 lines
5.9 KiB
Go

package adaptor
import (
"fmt"
"github.com/apache/pulsar-client-go/pulsar"
rawKafka "github.com/confluentinc/confluent-kafka-go/kafka"
rawWP "github.com/zilliztech/woodpecker/woodpecker/log"
"github.com/milvus-io/milvus-proto/go-api/v3/commonpb"
"github.com/milvus-io/milvus/pkg/v3/mq/common"
rawRocksmq "github.com/milvus-io/milvus/pkg/v3/mq/mqimpl/rocksmq/client"
"github.com/milvus-io/milvus/pkg/v3/mq/mqimpl/rocksmq/server"
mqkafka "github.com/milvus-io/milvus/pkg/v3/mq/msgstream/mqwrapper/kafka"
mqpulsar "github.com/milvus-io/milvus/pkg/v3/mq/msgstream/mqwrapper/pulsar"
mqwoodpecker "github.com/milvus-io/milvus/pkg/v3/mq/msgstream/mqwrapper/wp"
"github.com/milvus-io/milvus/pkg/v3/streaming/util/message"
msgkafka "github.com/milvus-io/milvus/pkg/v3/streaming/walimpls/impls/kafka"
msgpulsar "github.com/milvus-io/milvus/pkg/v3/streaming/walimpls/impls/pulsar"
"github.com/milvus-io/milvus/pkg/v3/streaming/walimpls/impls/rmq"
msgwoodpecker "github.com/milvus-io/milvus/pkg/v3/streaming/walimpls/impls/wp"
"github.com/milvus-io/milvus/pkg/v3/util/merr"
)
// MustGetMQWrapperIDFromMessage converts message.MessageID to common.MessageID
// TODO: should be removed in future after common.MessageID is removed
// Deprecated
func MustGetMQWrapperIDFromMessage(messageID message.MessageID) common.MessageID {
if id, ok := messageID.(interface{ PulsarID() pulsar.MessageID }); ok {
return mqpulsar.NewPulsarID(id.PulsarID())
} else if id, ok := messageID.(interface{ RmqID() int64 }); ok {
return &server.RmqID{MessageID: id.RmqID()}
} else if id, ok := messageID.(interface{ KafkaID() rawKafka.Offset }); ok {
return mqkafka.NewKafkaID(int64(id.KafkaID()))
} else if id, ok := messageID.(interface{ WoodpeckerID() *rawWP.LogMessageId }); ok {
return mqwoodpecker.NewWoodpeckerID(id.WoodpeckerID())
}
panic("unsupported now")
}
func MustGetMQWrapperIDAndWALNameFromMessage(messageID message.MessageID) (common.MessageID, commonpb.WALName) {
if id, ok := messageID.(interface{ PulsarID() pulsar.MessageID }); ok {
return mqpulsar.NewPulsarID(id.PulsarID()), commonpb.WALName_Pulsar
} else if id, ok := messageID.(interface{ RmqID() int64 }); ok {
return &server.RmqID{MessageID: id.RmqID()}, commonpb.WALName_RocksMQ
} else if id, ok := messageID.(interface{ KafkaID() rawKafka.Offset }); ok {
return mqkafka.NewKafkaID(int64(id.KafkaID())), commonpb.WALName_Kafka
} else if id, ok := messageID.(interface{ WoodpeckerID() *rawWP.LogMessageId }); ok {
return mqwoodpecker.NewWoodpeckerID(id.WoodpeckerID()), commonpb.WALName_WoodPecker
}
panic("unsupported now")
}
// MustGetMessageIDFromMQWrapperID converts common.MessageID to message.MessageID
// TODO: should be removed in future after common.MessageID is removed
func MustGetMessageIDFromMQWrapperID(commonMessageID common.MessageID) message.MessageID {
if id, ok := commonMessageID.(interface{ PulsarID() pulsar.MessageID }); ok {
return msgpulsar.NewPulsarID(id.PulsarID())
} else if id, ok := commonMessageID.(*server.RmqID); ok {
return rmq.NewRmqID(id.MessageID)
} else if id, ok := commonMessageID.(*mqkafka.KafkaID); ok {
return msgkafka.NewKafkaID(rawKafka.Offset(id.MessageID))
} else if id, ok := commonMessageID.(interface{ WoodpeckerID() *rawWP.LogMessageId }); ok {
return msgwoodpecker.NewWpID(id.WoodpeckerID())
}
return nil
}
// DeserializeToMQWrapperID deserializes messageID bytes to common.MessageID
// TODO: should be removed in future after common.MessageID is removed
func DeserializeToMQWrapperID(msgID []byte, walName string) (common.MessageID, error) {
switch walName {
case "pulsar", commonpb.WALName_Pulsar.String():
pulsarID, err := mqpulsar.DeserializePulsarMsgID(msgID)
if err != nil {
return nil, err
}
return mqpulsar.NewPulsarID(pulsarID), nil
case "rocksmq", commonpb.WALName_RocksMQ.String():
rID := server.DeserializeRmqID(msgID)
return &server.RmqID{MessageID: rID}, nil
case "kafka", commonpb.WALName_Kafka.String():
kID := mqkafka.DeserializeKafkaID(msgID)
return mqkafka.NewKafkaID(kID), nil
case "woodpecker", commonpb.WALName_WoodPecker.String():
wID, err := mqwoodpecker.DeserializeWoodpeckerMsgID(msgID)
if err != nil {
return nil, err
}
return mqwoodpecker.NewWoodpeckerID(wID), nil
default:
return nil, merr.WrapErrParameterInvalidMsg("unsupported mq type %s", walName)
}
}
func MustGetMessageIDFromMQWrapperIDBytesWithWALName(walName message.WALName, msgIDBytes []byte) message.MessageID {
wName := walName
if wName == message.WALNameUnknown {
walName = message.MustGetDefaultWALName()
}
var commonMsgID common.MessageID
switch walName {
case message.WALNameRocksmq:
id := server.DeserializeRmqID(msgIDBytes)
commonMsgID = &server.RmqID{MessageID: id}
case message.WALNamePulsar:
msgID, err := mqpulsar.DeserializePulsarMsgID(msgIDBytes)
if err != nil {
panic(err)
}
commonMsgID = mqpulsar.NewPulsarID(msgID)
case message.WALNameKafka:
id := mqkafka.DeserializeKafkaID(msgIDBytes)
commonMsgID = mqkafka.NewKafkaID(id)
case message.WALNameWoodpecker:
msgID, err := mqwoodpecker.DeserializeWoodpeckerMsgID(msgIDBytes)
if err != nil {
panic(err)
}
commonMsgID = mqwoodpecker.NewWoodpeckerID(msgID)
default:
panic("unsupported now")
}
return MustGetMessageIDFromMQWrapperID(commonMsgID)
}
func MustGetEarliestMessageIDFromMQType(walName commonpb.WALName) (common.MessageID, commonpb.WALName) {
switch walName {
case commonpb.WALName_Pulsar:
pulsarID := pulsar.EarliestMessageID()
return mqpulsar.NewPulsarID(pulsarID), commonpb.WALName_Pulsar
case commonpb.WALName_RocksMQ:
rID := rawRocksmq.EarliestMessageID()
return &server.RmqID{MessageID: rID}, commonpb.WALName_RocksMQ
case commonpb.WALName_Kafka:
kID := int64(rawKafka.OffsetBeginning)
return mqkafka.NewKafkaID(kID), commonpb.WALName_Kafka
case commonpb.WALName_WoodPecker:
wID := rawWP.EarliestLogMessageID()
return mqwoodpecker.NewWoodpeckerID(&wID), commonpb.WALName_WoodPecker
default:
panic(fmt.Sprintf("unsupported mq type %s", walName))
}
}