Zhen Ye 4d69898cb2
enhance: support single pchannel level transaction (#35289)
issue: #33285

- support transaction on single wal.
- last confirmed message id can still be used when enable transaction.
- add fence operation for segment allocation interceptor.

---------

Signed-off-by: chyezh <chyezh@outlook.com>
2024-08-19 21:22:56 +08:00

52 lines
1.2 KiB
Go

package message
import (
"time"
"github.com/milvus-io/milvus/pkg/streaming/proto/messagespb"
)
type (
TxnState = messagespb.TxnState
TxnID int64
)
const (
TxnStateBegin TxnState = messagespb.TxnState_TxnBegin
TxnStateInFlight TxnState = messagespb.TxnState_TxnInFlight
TxnStateOnCommit TxnState = messagespb.TxnState_TxnOnCommit
TxnStateCommitted TxnState = messagespb.TxnState_TxnCommitted
TxnStateOnRollback TxnState = messagespb.TxnState_TxnOnRollback
TxnStateRollbacked TxnState = messagespb.TxnState_TxnRollbacked
NonTxnID = TxnID(-1)
)
// NewTxnContextFromProto generates TxnContext from proto message.
func NewTxnContextFromProto(proto *messagespb.TxnContext) *TxnContext {
if proto == nil {
return nil
}
return &TxnContext{
TxnID: TxnID(proto.TxnId),
Keepalive: time.Duration(proto.KeepaliveMilliseconds) * time.Millisecond,
}
}
// TxnContext is the transaction context for message.
type TxnContext struct {
TxnID TxnID
Keepalive time.Duration
}
// IntoProto converts TxnContext to proto message.
func (t *TxnContext) IntoProto() *messagespb.TxnContext {
if t == nil {
return nil
}
return &messagespb.TxnContext{
TxnId: int64(t.TxnID),
KeepaliveMilliseconds: t.Keepalive.Milliseconds(),
}
}