Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 11 additions & 10 deletions pkg/client/client_engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ import (
// DATA_LAYOUT_TYPE_SHARDED constructs a single shardGroupBackend (EC).
// Without forwarding this field, callers would default to the proto3 zero
// value (REPLICATED), and EC volumes would be silently miscreated as RAID1.
func (c *SPDKClient) EngineCreate(name, volumeName, frontend string, specSize uint64, replicaAddressMap map[string]string, portCount int32, salvageRequested bool, snapshotMaxCount int32, dataLayoutType spdkrpc.DataLayoutType) (*api.Engine, error) {
func (c *SPDKClient) EngineCreate(name, volumeName, frontend string, specSize uint64, replicaAddressMap map[string]string, portCount int32, salvageRequested bool, snapshotMaxCount int32, dataLayoutType spdkrpc.DataLayoutType, transportType spdkrpc.TransportType) (*api.Engine, error) {
if name == "" {
return nil, fmt.Errorf("failed to start engine: missing required parameter name")
}
Expand All @@ -36,15 +36,16 @@ func (c *SPDKClient) EngineCreate(name, volumeName, frontend string, specSize ui
defer cancel()

resp, err := client.EngineCreate(ctx, &spdkrpc.EngineCreateRequest{
Name: name,
VolumeName: volumeName,
Frontend: frontend,
SpecSize: specSize,
ReplicaAddressMap: replicaAddressMap,
PortCount: portCount,
SalvageRequested: salvageRequested,
SnapshotMaxCount: snapshotMaxCount,
DataLayoutType: dataLayoutType,
Name: name,
VolumeName: volumeName,
Frontend: frontend,
SpecSize: specSize,
ReplicaAddressMap: replicaAddressMap,
PortCount: portCount,
SalvageRequested: salvageRequested,
SnapshotMaxCount: snapshotMaxCount,
DataLayoutType: dataLayoutType,
TransportType: transportType,
})
if err != nil {
return nil, errors.Wrap(err, "failed to start engine")
Expand Down
15 changes: 8 additions & 7 deletions pkg/client/client_replica.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ import (
)

// ReplicaCreate creates and starts a replica in the specified lvstore.
func (c *SPDKClient) ReplicaCreate(name, lvsName, lvsUUID string, specSize uint64, portCount int32, backingImageName string) (*api.Replica, error) {
func (c *SPDKClient) ReplicaCreate(name, lvsName, lvsUUID string, specSize uint64, portCount int32, backingImageName string, transportType spdkrpc.TransportType) (*api.Replica, error) {
if name == "" || lvsName == "" || lvsUUID == "" {
return nil, fmt.Errorf("failed to start SPDK replica: missing required parameters")
}
Expand All @@ -25,12 +25,13 @@ func (c *SPDKClient) ReplicaCreate(name, lvsName, lvsUUID string, specSize uint6
defer cancel()

resp, err := client.ReplicaCreate(ctx, &spdkrpc.ReplicaCreateRequest{
Name: name,
LvsName: lvsName,
LvsUuid: lvsUUID,
SpecSize: specSize,
PortCount: portCount,
BackingImageName: backingImageName,
Name: name,
LvsName: lvsName,
LvsUuid: lvsUUID,
SpecSize: specSize,
PortCount: portCount,
BackingImageName: backingImageName,
TransportType: transportType,
})
if err != nil {
return nil, errors.Wrap(err, "failed to start SPDK replica")
Expand Down
3 changes: 2 additions & 1 deletion pkg/spdk/backend_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package spdk
import (
"fmt"

spdktypes "github.com/longhorn/go-spdk-helper/pkg/spdk/types"
"github.com/longhorn/types/pkg/generated/spdkrpc"
grpccodes "google.golang.org/grpc/codes"
grpcstatus "google.golang.org/grpc/status"
Expand Down Expand Up @@ -84,7 +85,7 @@ func (s *TestSuite) TestIsShardedEngine(c *C) {

for name, tc := range testCases {
fmt.Println("Testing isShardedEngine:", name)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e.backends = tc.backends
c.Assert(isShardedEngine(e), Equals, tc.expected, Commentf("case %q", name))
}
Expand Down
28 changes: 15 additions & 13 deletions pkg/spdk/backup_restore_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@ import (
"strings"
"time"

spdktypes "github.com/longhorn/go-spdk-helper/pkg/spdk/types"

"github.com/longhorn/longhorn-spdk-engine/pkg/api"
lhtypes "github.com/longhorn/longhorn-spdk-engine/pkg/types"

Expand All @@ -27,7 +29,7 @@ func newTestReplicaBackend(name, address string, mode lhtypes.Mode) *replicaBack
func (s *TestSuite) TestEnsureReplicaModeForInfoUpdateRWQualifies(c *C) {
fmt.Println("Testing ensureReplicaModeForInfoUpdate: RW mode qualifies for info update")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
rs := newTestReplicaBackend("replica-1", "10.0.0.1:1234", lhtypes.ModeRW)

ok := e.ensureReplicaModeForInfoUpdate("replica-1", rs)
Expand All @@ -39,7 +41,7 @@ func (s *TestSuite) TestEnsureReplicaModeForInfoUpdateRWQualifies(c *C) {
func (s *TestSuite) TestEnsureReplicaModeForInfoUpdateWOQualifies(c *C) {
fmt.Println("Testing ensureReplicaModeForInfoUpdate: WO mode qualifies for info update")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
rs := newTestReplicaBackend("replica-1", "10.0.0.1:1234", lhtypes.ModeWO)

ok := e.ensureReplicaModeForInfoUpdate("replica-1", rs)
Expand All @@ -51,7 +53,7 @@ func (s *TestSuite) TestEnsureReplicaModeForInfoUpdateWOQualifies(c *C) {
func (s *TestSuite) TestEnsureReplicaModeForInfoUpdateERRDoesNotQualify(c *C) {
fmt.Println("Testing ensureReplicaModeForInfoUpdate: ERR mode does not qualify")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
rs := newTestReplicaBackend("replica-1", "10.0.0.1:1234", lhtypes.ModeERR)

ok := e.ensureReplicaModeForInfoUpdate("replica-1", rs)
Expand All @@ -63,7 +65,7 @@ func (s *TestSuite) TestEnsureReplicaModeForInfoUpdateERRDoesNotQualify(c *C) {
func (s *TestSuite) TestEnsureReplicaModeForInfoUpdateUnexpectedModeDowngradesToERR(c *C) {
fmt.Println("Testing ensureReplicaModeForInfoUpdate: unexpected mode is downgraded to ERR")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
rs := newTestReplicaBackend("replica-1", "10.0.0.1:1234", lhtypes.Mode("UNKNOWN"))

ok := e.ensureReplicaModeForInfoUpdate("replica-1", rs)
Expand All @@ -78,7 +80,7 @@ func (s *TestSuite) TestEnsureReplicaModeForInfoUpdateUnexpectedModeDowngradesTo
func (s *TestSuite) TestCheckAndUpdateInfoFromReplicasNoLockEmptyMap(c *C) {
fmt.Println("Testing checkAndUpdateInfoFromReplicasNoLock: empty backends does not panic")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e.backends = map[string]Backend{}

// Should not panic with empty map
Expand All @@ -88,7 +90,7 @@ func (s *TestSuite) TestCheckAndUpdateInfoFromReplicasNoLockEmptyMap(c *C) {
func (s *TestSuite) TestCheckAndUpdateInfoFromReplicasNoLockAllERRSkipped(c *C) {
fmt.Println("Testing checkAndUpdateInfoFromReplicasNoLock: all-ERR replicas are skipped without network calls")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e.backends = map[string]Backend{
"replica-1": newTestReplicaBackend("replica-1", "10.0.0.1:1234", lhtypes.ModeERR),
"replica-2": newTestReplicaBackend("replica-2", "10.0.0.2:1234", lhtypes.ModeERR),
Expand Down Expand Up @@ -129,7 +131,7 @@ func (s *TestSuite) TestEngineFrontendTeardownRestoreInitiatorMarksStopped(c *C)
func (s *TestSuite) TestEngineReplicaAddRejectedDuringRestore(c *C) {
fmt.Println("Testing Engine.ReplicaAdd is rejected while restore is in progress")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendEmpty, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendEmpty, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e.State = lhtypes.InstanceStateRunning
e.IsRestoring = true

Expand All @@ -142,7 +144,7 @@ func (s *TestSuite) TestEngineReplicaAddRejectedDuringRestore(c *C) {
func (s *TestSuite) TestRecordBackupRestoreStartErrorExposedInRestoreStatus(c *C) {
fmt.Println("Testing restore start errors are exposed through Engine.RestoreStatus")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendEmpty, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendEmpty, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e.backends = map[string]Backend{
"replica-1": newTestReplicaBackend("replica-1", "10.0.0.1:1234", lhtypes.ModeRW),
"replica-2": newTestReplicaBackend("replica-2", "10.0.0.2:1234", lhtypes.ModeRW),
Expand Down Expand Up @@ -171,7 +173,7 @@ func (s *TestSuite) TestRecordBackupRestoreStartErrorExposedInRestoreStatus(c *C
func (s *TestSuite) TestRecordBackupRestoreStartErrorPreservesLastRestored(c *C) {
fmt.Println("Testing restore start errors preserve last restored backup")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendEmpty, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendEmpty, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e.restore = NewEngineRestore(nil, "s3://backupbucket@us-east-1/backupstore?backup=backup-old&volume=vol-a", "backup-old", e, nil)
e.restore.FinishRestore()

Expand All @@ -187,7 +189,7 @@ func (s *TestSuite) TestRecordBackupRestoreStartErrorPreservesLastRestored(c *C)
func (s *TestSuite) TestCheckAndUpdateInfoFromReplicasNoLockAppliesBackendView(c *C) {
fmt.Println("Testing checkAndUpdateInfoFromReplicasNoLock applies SnapshotMap/Head/ActualSize from Backend.Get()")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
u := newFakeBackend("r1", "10.0.0.1:1234")
u.SetMode(lhtypes.ModeRW)
headLvol := &api.Lvol{Name: "vol-head", Parent: "snap-1"}
Expand All @@ -212,7 +214,7 @@ func (s *TestSuite) TestCheckAndUpdateInfoFromReplicasNoLockAppliesBackendView(c
func (s *TestSuite) TestCheckAndUpdateInfoFromReplicasNoLockMarksERROnGetError(c *C) {
fmt.Println("Testing checkAndUpdateInfoFromReplicasNoLock marks backend ERR when Backend.Get() returns an error")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
u := newFakeBackend("r1", "10.0.0.1:1234")
u.SetMode(lhtypes.ModeRW)
u.ViewErr = errors.New("backend unavailable")
Expand All @@ -226,7 +228,7 @@ func (s *TestSuite) TestCheckAndUpdateInfoFromReplicasNoLockMarksERROnGetError(c
func (s *TestSuite) TestResolveReplicaAncestorRoutesBackingImageThroughBackend(c *C) {
fmt.Println("Testing resolveReplicaAncestor calls BackingImageGet via the Backend interface")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
u := newFakeBackend("r1", "10.0.0.1:1234")
u.SetMode(lhtypes.ModeRW)
biSnap := &api.Lvol{Name: "bi-snap"}
Expand All @@ -252,7 +254,7 @@ func (s *TestSuite) TestResolveReplicaAncestorRoutesBackingImageThroughBackend(c
func (s *TestSuite) TestResolveReplicaAncestorMarksERROnBackingImageError(c *C) {
fmt.Println("Testing resolveReplicaAncestor marks backend ERR when BackingImageGet fails")

e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
e := NewEngine("engine-a", "vol-a", lhtypes.FrontendSPDKTCPBlockdev, spdktypes.NvmeTransportTypeTCP, 10, make(chan interface{}, 1), defaultTestSnapshotMaxCount, nil)
u := newFakeBackend("r1", "10.0.0.1:1234")
u.SetMode(lhtypes.ModeRW)
u.BackingImageGetErr = errors.New("backing image not found")
Expand Down
33 changes: 25 additions & 8 deletions pkg/spdk/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,14 @@ type Engine struct {
ActualSize uint64
Frontend string

// TransportType is the NVMe-oF transport this engine uses to connect to its
// replicas' exposed head bdevs. It must match the transport each replica
// used to expose (Replica.TransportType). Set from EngineCreateRequest.
// data_engine_transport; unset (empty) is treated as TCP for backward
// compatibility. Only the internal engine<->replica data fabric honors this;
// the host-facing NVMe-TCP frontend and all EC/transient paths stay TCP.
TransportType spdktypes.NvmeTransportType

ctrlrLossTimeout int
fastIOFailTimeoutSec int
// backends maps each peer's name to its Backend. Each one provides a base
Expand Down Expand Up @@ -180,7 +188,7 @@ type Engine struct {
newServiceClient ServiceClientFactory
}

func NewEngine(engineName, volumeName, frontend string, specSize uint64, engineUpdateCh chan interface{}, snapshotMaxCount int32, newServiceClient ServiceClientFactory) *Engine {
func NewEngine(engineName, volumeName, frontend string, transportType spdktypes.NvmeTransportType, specSize uint64, engineUpdateCh chan interface{}, snapshotMaxCount int32, newServiceClient ServiceClientFactory) *Engine {
log := logrus.StandardLogger().WithFields(logrus.Fields{
"engineName": engineName,
"volumeName": volumeName,
Expand All @@ -207,6 +215,10 @@ func NewEngine(engineName, volumeName, frontend string, specSize uint64, engineU
Frontend: frontend,
SpecSize: specSize,

// Empty (unset) transportType is treated as TCP by nvmeTransportFromProto,
// preserving historical behavior for callers that do not select a transport.
TransportType: transportType,

// TODO: support user-defined values
ctrlrLossTimeout: replicaCtrlrLossTimeoutSec,
fastIOFailTimeoutSec: replicaFastIOFailTimeoutSec,
Expand Down Expand Up @@ -364,8 +376,10 @@ func (e *Engine) createNVMeTCPTarget(spdkClient *spdkclient.Client, superiorPort

e.log.Infof("Starting to expose RAID bdev for engine target %v on %v:%v with initial ANA state %v, cntlid %v, nsUUID %v",
e.Name, e.NvmeTcpTarget.IP, e.NvmeTcpTarget.Port, initialANAState, cntlid, nsUUID)
// TODO: the host-facing frontend is always exposed over NVMe-TCP for now.
// Adapt it to e.TransportType once the initiator/host side supports RDMA.
if err := spdkClient.StartExposeBdevWithANAState(e.NvmeTcpTarget.Nqn, e.Name, e.NvmeTcpTarget.Nguid, nsUUID,
e.NvmeTcpTarget.IP, strconv.Itoa(int(e.NvmeTcpTarget.Port)), spdkANAState, cntlid, cntlid); err != nil {
e.NvmeTcpTarget.IP, strconv.Itoa(int(e.NvmeTcpTarget.Port)), spdktypes.NvmeTransportTypeTCP, spdkANAState, cntlid, cntlid); err != nil {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

NIT: Add a TODO comment for the Frontend adaptation

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

renamed the struct field to TransportType on both Engine and Replica, and the client-method params to transportType

// No need to release ports here. The engine will be marked as ERR by
// Create's deferred error handler, and Delete will release the ports
// when the user cleans up this engine.
Expand All @@ -387,7 +401,7 @@ func (e *Engine) connectReplicas(spdkClient *spdkclient.Client, replicaAddressMa
for replicaName, replicaAddr := range replicaAddressMap {
e.backends[replicaName] = backendFactory(replicaName, replicaAddr)

bdevName, err := connectNVMfBdev(spdkClient, replicaName, replicaAddr, e.ctrlrLossTimeout, e.fastIOFailTimeoutSec, maxRetries, retryInterval)
bdevName, err := connectNVMfBdev(spdkClient, replicaName, replicaAddr, e.TransportType, e.ctrlrLossTimeout, e.fastIOFailTimeoutSec, maxRetries, retryInterval)
if err != nil {
e.log.WithError(err).Warnf("Failed to get bdev from replica %s with address %s during engine creation, will mark the mode to ERR and continue", replicaName, replicaAddr)
e.backends[replicaName].SetMode(types.ModeERR)
Expand Down Expand Up @@ -1041,7 +1055,7 @@ func (e *Engine) replicaAddStart(spdkClient *spdkclient.Client,
}

// Add rebuilding replica head bdev to the base bdev list of the RAID bdev
dstHeadLvolBdevName, err := connectNVMfBdev(spdkClient, dstReplicaName, dstHeadLvolAddress, e.ctrlrLossTimeout, e.fastIOFailTimeoutSec, maxRetries, retryInterval)
dstHeadLvolBdevName, err := connectNVMfBdev(spdkClient, dstReplicaName, dstHeadLvolAddress, e.TransportType, e.ctrlrLossTimeout, e.fastIOFailTimeoutSec, maxRetries, retryInterval)
if err != nil {
return nil, startUpdateRequired, nil, err
}
Expand Down Expand Up @@ -1878,7 +1892,7 @@ func (e *Engine) replicaSnapshotOperation(spdkClient *spdkclient.Client, replica
if err := replicaStatus.SnapshotRevert(snapshotName); err != nil {
return err
}
bdevName, err := connectNVMfBdev(spdkClient, replicaName, replicaStatus.Address(), e.ctrlrLossTimeout, e.fastIOFailTimeoutSec, maxRetries, retryInterval)
bdevName, err := connectNVMfBdev(spdkClient, replicaName, replicaStatus.Address(), e.TransportType, e.ctrlrLossTimeout, e.fastIOFailTimeoutSec, maxRetries, retryInterval)
if err != nil {
return err
}
Expand Down Expand Up @@ -2919,7 +2933,7 @@ func (e *Engine) Expand(spdkClient *spdkclient.Client, size uint64) (err error)
e.log.Infof("Starting to expose RAID bdev for engine target %v on %v:%v with ANA state %v, cntlid %v, nsUUID %v",
e.Name, e.NvmeTcpTarget.IP, e.NvmeTcpTarget.Port, currentANAState, cntlid, nsUUID)
if err := spdkClient.StartExposeBdevWithANAState(e.NvmeTcpTarget.Nqn, e.Name, e.NvmeTcpTarget.Nguid, nsUUID,
e.NvmeTcpTarget.IP, strconv.Itoa(int(e.NvmeTcpTarget.Port)),
e.NvmeTcpTarget.IP, strconv.Itoa(int(e.NvmeTcpTarget.Port)), spdktypes.NvmeTransportTypeTCP,
spdkANAState, cntlid, cntlid); err != nil {
return errors.Wrapf(err, "failed to start exposing RAID bdev for engine target %v", e.Name)
}
Expand Down Expand Up @@ -3161,7 +3175,7 @@ func (e *Engine) expandSingleReplica(spdkClient *spdkclient.Client, replicaName
return err
}

_, err = connectNVMfBdev(spdkClient, replicaName, replicaStatus.Address(), e.ctrlrLossTimeout, e.fastIOFailTimeoutSec, maxRetries, retryInterval)
_, err = connectNVMfBdev(spdkClient, replicaName, replicaStatus.Address(), e.TransportType, e.ctrlrLossTimeout, e.fastIOFailTimeoutSec, maxRetries, retryInterval)
return err
}

Expand Down Expand Up @@ -3829,7 +3843,10 @@ func validateAndGetSingleNvmeInfo(replicaName string, bdev *spdktypes.BdevInfo)
}

func validateNvmeTransport(replicaName, bdevName string, nvmeInfo spdktypes.NvmeNamespaceInfo) error {
if !strings.EqualFold(string(nvmeInfo.Trid.Trtype), string(spdktypes.NvmeTransportTypeTCP)) {
// The internal engine<->replica fabric may run over TCP or RDMA (RoCEv2).
// Both are valid; only reject transports we never use (e.g. FC, PCIe).
if !strings.EqualFold(string(nvmeInfo.Trid.Trtype), string(spdktypes.NvmeTransportTypeTCP)) &&
!strings.EqualFold(string(nvmeInfo.Trid.Trtype), string(spdktypes.NvmeTransportTypeRDMA)) {
return fmt.Errorf(
"found invalid transport type %s in a remote NVMe base bdev %s during replica %s mode validation",
nvmeInfo.Trid.Trtype, bdevName, replicaName,
Expand Down
Loading