// nolint: err113 package queue import ( "bytes" "context" "errors" "fmt" "testing" "time" "github.com/number571/go-peer/pkg/crypto/asymmetric" "github.com/number571/go-peer/pkg/crypto/scheme/layer1" "github.com/number571/go-peer/pkg/crypto/scheme/layer2" "github.com/number571/go-peer/pkg/crypto/scheme/layer2/hybrid" "github.com/number571/go-peer/pkg/payload" testutils "github.com/number571/go-peer/test/utils" ) const ( tcQueueCap = 16 tcMsgSize = (8 << 10) tcMsgBody = "hello, world!" ) func TestError(t *testing.T) { t.Parallel() str := "value" err := &SQueueError{str} if err.Error() != errPrefix+str { t.Fatal("incorrect err.Error()") } } func TestSettings(t *testing.T) { t.Parallel() for i := 0; i < 4; i++ { testSettings(t, i) } } func testSettings(t *testing.T, n int) { defer func() { if r := recover(); r == nil { t.Fatal("nothing panics") } }() switch n { case 0: _ = NewSettings(&SSettings{ FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: layer1.NewSettings(&layer1.SSettings{}), }), FQueuePeriod: 500 * time.Millisecond, FConsumersCap: 1, }) case 1: _ = NewSettings(&SSettings{ FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: layer1.NewSettings(&layer1.SSettings{}), }), FQueuePoolCap: [2]uint64{tcQueueCap, tcQueueCap}, FConsumersCap: 1, }) case 2: _ = NewSettings(&SSettings{ FQueuePoolCap: [2]uint64{tcQueueCap, tcQueueCap}, FQueuePeriod: 500 * time.Millisecond, FConsumersCap: 1, }) case 3: _ = NewSettings(&SSettings{ FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: layer1.NewSettings(&layer1.SSettings{}), }), FQueuePoolCap: [2]uint64{tcQueueCap, tcQueueCap}, FQueuePeriod: 500 * time.Millisecond, }) } } func TestRunStopQueue(t *testing.T) { t.Parallel() privKey := asymmetric.NewPrivKey() scheme, _ := hybrid.NewScheme( privKey, tcMsgSize, ) queue := NewQBProblemProcessor( NewSettings(&SSettings{ FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: layer1.NewSettings(&layer1.SSettings{}), }), FQueuePoolCap: [2]uint64{tcQueueCap, tcQueueCap}, FQueuePeriod: 100 * time.Millisecond, FConsumersCap: 1, }), scheme, ) ctx1, cancel1 := context.WithCancel(context.Background()) defer cancel1() go func() { if err := queue.Run(ctx1); err != nil && !errors.Is(err, context.Canceled) { t.Error(err) } }() err := testutils.TryN(50, 10*time.Millisecond, func() error { sett := queue.GetSettings() sQueue := queue.(*sQBProblemProcessor) if uint64(len(sQueue.fRandPool.fQueue)) == sett.GetQueuePoolCap()[0] { return nil } return errors.New("len(void queue) != max capacity") //nolint:err113 }) if err != nil { t.Fatal(err) } ctx2, cancel2 := context.WithCancel(context.Background()) defer cancel2() go func() { if err := queue.Run(ctx2); err == nil { t.Error("success run already running queue") } }() pubKey := privKey.GetPubKey() pldBytes := payload.NewPayload64(0, []byte(tcMsgBody)).ToBytes() for i := 0; i < tcQueueCap; i++ { if err := queue.EnqueueMessage(pubKey, pldBytes); err != nil { t.Fatal(err) } } // after full queue for i := 0; i < 2*tcQueueCap; i++ { if err := queue.EnqueueMessage(pubKey, pldBytes); err != nil { return } } t.Fatal("success enqueue message with max capacity") } func TestQueue(t *testing.T) { t.Parallel() privKey := asymmetric.NewPrivKey() queue := NewQBProblemProcessor( NewSettings(&SSettings{ FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: 10, }), }), FQueuePoolCap: [2]uint64{tcQueueCap, tcQueueCap}, FQueuePeriod: 100 * time.Millisecond, FConsumersCap: 1, }), func() layer2.IScheme { scheme, _ := hybrid.NewScheme( privKey, tcMsgSize, ) return scheme }(), ) sett := queue.GetSettings() if sett.GetQueuePoolCap() != [2]uint64{tcQueueCap, tcQueueCap} { t.Fatal("sett.GetMainCapacity() != tcQueueCap") } if err := testQueue(queue, privKey); err != nil { t.Fatal(err) } } func testQueue(queue IQBProblemProcessor, privKey asymmetric.IPrivKey) error { ctx, cancel := context.WithCancel(context.Background()) defer func() { cancel() time.Sleep(200 * time.Millisecond) }() go func() { if err := queue.Run(ctx); err != nil && !errors.Is(err, context.Canceled) { return } }() pubKey := privKey.GetPubKey() pldBytes := payload.NewPayload64(0, []byte(tcMsgBody)).ToBytes() if err := queue.EnqueueMessage(pubKey, pldBytes); err != nil { return err } // wait minimum one generated message time.Sleep(300 * time.Millisecond) // auto fill queue enabled only if QB=true msgs := make([]layer1.IMessage, 0, 3) for i := 0; i < 3; i++ { msgs = append(msgs, queue.DequeueMessage(ctx)) } for i := 0; i < len(msgs)-1; i++ { for j := i + 1; j < len(msgs); j++ { if bytes.Equal(msgs[i].GetHash(), msgs[j].GetHash()) { return fmt.Errorf("hash of messages equals (%d and %d)", i, i) //nolint:err113 } } } notClosed := make(chan bool) go func() { // test close with parallel dequeue msg := queue.DequeueMessage(ctx) notClosed <- (msg != nil) }() cancel() if <-notClosed { return errors.New("success dequeue with close") //nolint:err113 } return nil }