Files

184 lines
5.6 KiB
Go
Raw Permalink Normal View History

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
}