aboutsummaryrefslogtreecommitdiff
path: root/weed/mq/messages/message_pipeline_test.go
diff options
context:
space:
mode:
authorchrislu <chris.lu@gmail.com>2022-09-25 17:57:26 -0700
committerchrislu <chris.lu@gmail.com>2022-09-25 17:57:26 -0700
commit24317a97636546db80faf66c1653eab1e1fc2b2a (patch)
treef67fe074353b5c94d0af3dd1198a6d211aeae3b0 /weed/mq/messages/message_pipeline_test.go
parent3da8d7309a2b46861475fb685eb21c7179bf5c44 (diff)
downloadseaweedfs-24317a97636546db80faf66c1653eab1e1fc2b2a.tar.xz
seaweedfs-24317a97636546db80faf66c1653eab1e1fc2b2a.zip
async wait
Diffstat (limited to 'weed/mq/messages/message_pipeline_test.go')
-rw-r--r--weed/mq/messages/message_pipeline_test.go30
1 files changed, 27 insertions, 3 deletions
diff --git a/weed/mq/messages/message_pipeline_test.go b/weed/mq/messages/message_pipeline_test.go
index aaa25e0fd..e7ecfa8f6 100644
--- a/weed/mq/messages/message_pipeline_test.go
+++ b/weed/mq/messages/message_pipeline_test.go
@@ -5,11 +5,11 @@ import (
"time"
)
-func TestAddMessage(t *testing.T) {
+func TestAddMessage1(t *testing.T) {
mp := NewMessagePipeline(0, 3, 10, time.Second, &EmptyMover{})
go func() {
- outChan := mp.OutputChan()
- for mr := range outChan {
+ for i := 0; i < 100; i++ {
+ mr := mp.NextMessageBufferReference()
println(mr.sequence, mr.fileId)
}
}()
@@ -26,4 +26,28 @@ func TestAddMessage(t *testing.T) {
mp.ShutdownStart()
mp.ShutdownWait()
+
+}
+
+func TestAddMessage2(t *testing.T) {
+ mp := NewMessagePipeline(0, 3, 10, time.Second, &EmptyMover{})
+
+ for i := 0; i < 100; i++ {
+ message := &Message{
+ Key: []byte("key"),
+ Content: []byte("data"),
+ Properties: nil,
+ Ts: time.Now(),
+ }
+ mp.AddMessage(message)
+ }
+
+ mp.ShutdownStart()
+ mp.ShutdownWait()
+
+ for i := 0; i < 100; i++ {
+ mr := mp.NextMessageBufferReference()
+ println(mr.sequence, mr.fileId)
+ }
+
}