184 lines
5.6 KiB
Go
184 lines
5.6 KiB
Go
|
|
package notify
|
||
|
|
|
||
|
|
import (
|
||
|
|
"context"
|
||
|
|
"errors"
|
||
|
|
"io"
|
||
|
|
"testing"
|
||
|
|
"time"
|
||
|
|
)
|
||
|
|
|
||
|
|
func TestBulkRuntimeAdoptFailureFinalizesCandidate(t *testing.T) {
|
||
|
|
tests := []struct {
|
||
|
|
name string
|
||
|
|
adopt func(*bulkRuntime, string, *bulkHandle) error
|
||
|
|
}{
|
||
|
|
{
|
||
|
|
name: "inbound",
|
||
|
|
adopt: func(runtime *bulkRuntime, scope string, bulk *bulkHandle) error {
|
||
|
|
return runtime.adoptInbound(scope, bulk)
|
||
|
|
},
|
||
|
|
},
|
||
|
|
{
|
||
|
|
name: "reserved",
|
||
|
|
adopt: func(runtime *bulkRuntime, scope string, bulk *bulkHandle) error {
|
||
|
|
return runtime.adoptReserved(scope, bulk)
|
||
|
|
},
|
||
|
|
},
|
||
|
|
}
|
||
|
|
|
||
|
|
for _, test := range tests {
|
||
|
|
t.Run(test.name, func(t *testing.T) {
|
||
|
|
runtime := newBulkRuntime("cblk")
|
||
|
|
scope := clientFileScope()
|
||
|
|
existing := newBulkHandle(context.Background(), runtime, scope, BulkOpenRequest{
|
||
|
|
BulkID: "duplicate",
|
||
|
|
DataID: 1,
|
||
|
|
}, 0, nil, nil, 0, nil, nil, nil, nil, nil)
|
||
|
|
if err := runtime.registerInbound(scope, existing); err != nil {
|
||
|
|
t.Fatalf("register existing bulk: %v", err)
|
||
|
|
}
|
||
|
|
defer existing.markReset(io.ErrClosedPipe)
|
||
|
|
|
||
|
|
candidate := newBulkHandle(context.Background(), runtime, scope, BulkOpenRequest{
|
||
|
|
BulkID: "duplicate",
|
||
|
|
DataID: 3,
|
||
|
|
ChunkSize: 4,
|
||
|
|
WindowBytes: 4,
|
||
|
|
MaxInFlight: 1,
|
||
|
|
}, 0, nil, nil, 0, nil, nil, nil,
|
||
|
|
func(context.Context, *bulkHandle, uint64, []byte, bool) (int, error) {
|
||
|
|
return 0, nil
|
||
|
|
},
|
||
|
|
func(*bulkHandle, int64, int) error { return nil },
|
||
|
|
)
|
||
|
|
if err := test.adopt(runtime, scope, candidate); !errors.Is(err, errBulkAlreadyExists) {
|
||
|
|
t.Fatalf("adopt error = %v, want %v", err, errBulkAlreadyExists)
|
||
|
|
}
|
||
|
|
if err := candidate.resetErrSnapshot(); !errors.Is(err, errBulkAlreadyExists) {
|
||
|
|
t.Fatalf("candidate reset error = %v, want %v", err, errBulkAlreadyExists)
|
||
|
|
}
|
||
|
|
waitBulkWorkerStopped(t, "write", candidate.writeWorkerDone)
|
||
|
|
waitBulkWorkerStopped(t, "release", candidate.releaseWorkerDone)
|
||
|
|
})
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
func TestSharedBulkFinalizeDoesNotReleaseDedicatedLane(t *testing.T) {
|
||
|
|
client := NewClient().(*ClientCommon)
|
||
|
|
laneID := client.reserveBulkDedicatedLane()
|
||
|
|
defer client.releaseBulkDedicatedLane(laneID)
|
||
|
|
|
||
|
|
bulk := newBulkHandle(context.Background(), nil, clientFileScope(), BulkOpenRequest{
|
||
|
|
BulkID: "shared",
|
||
|
|
DataID: 1,
|
||
|
|
}, 0, nil, nil, 0, nil, nil, nil, nil, nil)
|
||
|
|
bulk.setClientSnapshotOwner(client)
|
||
|
|
bulk.finalize()
|
||
|
|
|
||
|
|
if got := clientDedicatedLaneActiveBulks(client, laneID); got != 1 {
|
||
|
|
t.Fatalf("shared finalize changed lane %d active bulks to %d, want 1", laneID, got)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
func TestDedicatedLaneLeaseRetainsExactLaneAndReleasesOnce(t *testing.T) {
|
||
|
|
client := NewClient().(*ClientCommon)
|
||
|
|
laneID := client.retainBulkDedicatedLane(1)
|
||
|
|
bulk := newBulkHandle(context.Background(), nil, clientFileScope(), BulkOpenRequest{
|
||
|
|
BulkID: "dedicated",
|
||
|
|
DataID: 1,
|
||
|
|
Dedicated: true,
|
||
|
|
DedicatedLaneID: laneID,
|
||
|
|
}, 0, nil, nil, 0, nil, nil, nil, nil, nil)
|
||
|
|
bulk.setClientSnapshotOwner(client)
|
||
|
|
bulk.markDedicatedLaneReserved()
|
||
|
|
|
||
|
|
otherLaneID := client.reserveBulkDedicatedLane()
|
||
|
|
if otherLaneID == laneID {
|
||
|
|
t.Fatalf("new lane reservation overwrote retained lane %d", laneID)
|
||
|
|
}
|
||
|
|
defer client.releaseBulkDedicatedLane(otherLaneID)
|
||
|
|
|
||
|
|
bulk.finalize()
|
||
|
|
bulk.finalize()
|
||
|
|
if got := clientDedicatedLaneActiveBulks(client, laneID); got != 0 {
|
||
|
|
t.Fatalf("dedicated lane %d active bulks after duplicate finalize = %d, want 0", laneID, got)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
func TestInboundDedicatedRetainFailureFinalizesCandidate(t *testing.T) {
|
||
|
|
client := NewClient().(*ClientCommon)
|
||
|
|
runtime := newBulkRuntime("retain-failure")
|
||
|
|
route := clientSessionRoute{epoch: 1}
|
||
|
|
bulk := newBulkHandle(context.Background(), runtime, clientFileScope(), BulkOpenRequest{
|
||
|
|
BulkID: "retain-failure",
|
||
|
|
DataID: 1,
|
||
|
|
Dedicated: true,
|
||
|
|
DedicatedLaneID: 9,
|
||
|
|
ChunkSize: 4,
|
||
|
|
WindowBytes: 4,
|
||
|
|
MaxInFlight: 1,
|
||
|
|
}, 1, nil, nil, 0, nil, nil, nil,
|
||
|
|
func(context.Context, *bulkHandle, uint64, []byte, bool) (int, error) {
|
||
|
|
return 0, nil
|
||
|
|
},
|
||
|
|
func(*bulkHandle, int64, int) error { return nil },
|
||
|
|
)
|
||
|
|
bulk.setClientSnapshotOwner(client)
|
||
|
|
bulk.setClientSessionRoute(route)
|
||
|
|
|
||
|
|
err := client.retainBulkDedicatedLaneAtRoute(bulk.dedicatedLaneIDSnapshot(), route)
|
||
|
|
if err == nil {
|
||
|
|
t.Fatal("retain should fail for an unavailable session route")
|
||
|
|
}
|
||
|
|
bulk.markReset(err)
|
||
|
|
|
||
|
|
if got := bulk.resetErrSnapshot(); got == nil {
|
||
|
|
t.Fatal("retain failure did not set candidate reset error")
|
||
|
|
}
|
||
|
|
if _, ok := runtime.lookup(clientFileScope(), bulk.ID()); ok {
|
||
|
|
t.Fatal("unadopted candidate was registered in bulk runtime")
|
||
|
|
}
|
||
|
|
waitBulkWorkerStopped(t, "write", bulk.writeWorkerDone)
|
||
|
|
waitBulkWorkerStopped(t, "release", bulk.releaseWorkerDone)
|
||
|
|
}
|
||
|
|
|
||
|
|
func TestFinalizedBulkRejectsDedicatedSenderInstall(t *testing.T) {
|
||
|
|
bulk := newBulkHandle(context.Background(), nil, clientFileScope(), BulkOpenRequest{
|
||
|
|
BulkID: "finalized-sender",
|
||
|
|
DataID: 1,
|
||
|
|
Dedicated: true,
|
||
|
|
}, 0, nil, nil, 0, nil, nil, nil, nil, nil)
|
||
|
|
bulk.finalize()
|
||
|
|
|
||
|
|
sender := &bulkDedicatedSender{}
|
||
|
|
if got := bulk.installDedicatedSender(sender); got != nil {
|
||
|
|
t.Fatalf("install sender after finalize = %p, want nil", got)
|
||
|
|
}
|
||
|
|
if got := bulk.dedicatedSenderSnapshot(); got != nil {
|
||
|
|
t.Fatalf("finalized bulk retained sender %p", got)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
func waitBulkWorkerStopped(t *testing.T, name string, done <-chan struct{}) {
|
||
|
|
t.Helper()
|
||
|
|
if done == nil {
|
||
|
|
t.Fatalf("%s worker was not started", name)
|
||
|
|
}
|
||
|
|
select {
|
||
|
|
case <-done:
|
||
|
|
case <-time.After(time.Second):
|
||
|
|
t.Fatalf("%s worker did not stop", name)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
func clientDedicatedLaneActiveBulks(client *ClientCommon, laneID uint32) int {
|
||
|
|
client.bulkDedicatedSidecarMu.Lock()
|
||
|
|
defer client.bulkDedicatedSidecarMu.Unlock()
|
||
|
|
lane := client.bulkDedicatedLanes[normalizeBulkDedicatedLaneID(laneID)]
|
||
|
|
if lane == nil {
|
||
|
|
return 0
|
||
|
|
}
|
||
|
|
return lane.activeBulks
|
||
|
|
}
|