aboutsummaryrefslogtreecommitdiff
path: root/weed/mq/segment
diff options
context:
space:
mode:
authorchrislu <chris.lu@gmail.com>2022-07-18 00:25:14 -0700
committerchrislu <chris.lu@gmail.com>2022-07-18 00:25:14 -0700
commit7f672b37e166a97fc652a40c7fc41d03f1d6e690 (patch)
tree3156f24f800c6bf5513aa9fc874df2fb5bf6ab26 /weed/mq/segment
parent1f2c5ee06ecd54f46169440a56be07c3081407b7 (diff)
downloadseaweedfs-7f672b37e166a97fc652a40c7fc41d03f1d6e690.tar.xz
seaweedfs-7f672b37e166a97fc652a40c7fc41d03f1d6e690.zip
add flatbuffer serde for message
Diffstat (limited to 'weed/mq/segment')
-rw-r--r--weed/mq/segment/message_serde.go48
-rw-r--r--weed/mq/segment/message_serde_test.go47
-rw-r--r--weed/mq/segment/segment_serde.go1
3 files changed, 96 insertions, 0 deletions
diff --git a/weed/mq/segment/message_serde.go b/weed/mq/segment/message_serde.go
new file mode 100644
index 000000000..b69d78cab
--- /dev/null
+++ b/weed/mq/segment/message_serde.go
@@ -0,0 +1,48 @@
+package segment
+
+import (
+ "github.com/chrislusf/seaweedfs/weed/pb/message_fbs"
+ flatbuffers "github.com/google/flatbuffers/go"
+)
+
+func CreateMessage(b *flatbuffers.Builder, producerId int32, producerSeq int64, segmentId int32, segmentSeq int64,
+ eventTsNs int64, recvTsNs int64, properties map[string]string, key []byte, value []byte) {
+ b.Reset()
+
+ var names, values, pairs []flatbuffers.UOffsetT
+ for k, v := range properties {
+ names = append(names, b.CreateString(k))
+ values = append(values, b.CreateString(v))
+ }
+
+ for i, _ := range names {
+ message_fbs.NameValueStart(b)
+ message_fbs.NameValueAddName(b, names[i])
+ message_fbs.NameValueAddValue(b, values[i])
+ pair := message_fbs.NameValueEnd(b)
+ pairs = append(pairs, pair)
+ }
+ message_fbs.MessageStartPropertiesVector(b, len(properties))
+ for i := len(pairs) - 1; i >= 0; i-- {
+ b.PrependUOffsetT(pairs[i])
+ }
+ prop := b.EndVector(len(properties))
+
+ k := b.CreateByteVector(key)
+ v := b.CreateByteVector(value)
+
+ message_fbs.MessageStart(b)
+ message_fbs.MessageAddProducerId(b, producerId)
+ message_fbs.MessageAddProducerSeq(b, producerSeq)
+ message_fbs.MessageAddSegmentId(b, segmentId)
+ message_fbs.MessageAddSegmentSeq(b, segmentSeq)
+ message_fbs.MessageAddEventTsNs(b, eventTsNs)
+ message_fbs.MessageAddRecvTsNs(b, recvTsNs)
+
+ message_fbs.MessageAddProperties(b, prop)
+ message_fbs.MessageAddKey(b, k)
+ message_fbs.MessageAddData(b, v)
+ message := message_fbs.MessageEnd(b)
+
+ b.Finish(message)
+}
diff --git a/weed/mq/segment/message_serde_test.go b/weed/mq/segment/message_serde_test.go
new file mode 100644
index 000000000..7ba0febf0
--- /dev/null
+++ b/weed/mq/segment/message_serde_test.go
@@ -0,0 +1,47 @@
+package segment
+
+import (
+ "github.com/chrislusf/seaweedfs/weed/pb/message_fbs"
+ flatbuffers "github.com/google/flatbuffers/go"
+ "github.com/stretchr/testify/assert"
+ "testing"
+)
+
+func TestMessageSerde(t *testing.T) {
+ b := flatbuffers.NewBuilder(1024)
+
+ prop := make(map[string]string)
+ prop["n1"] = "v1"
+ prop["n2"] = "v2"
+
+ CreateMessage(b, 1, 2, 3, 4, 5, 6, prop,
+ []byte("the primary key"), []byte("body is here"))
+
+ buf := b.FinishedBytes()
+
+ println("serialized size", len(buf))
+
+ m := message_fbs.GetRootAsMessage(buf, 0)
+
+ assert.Equal(t, int32(1), m.ProducerId())
+ assert.Equal(t, int64(2), m.ProducerSeq())
+ assert.Equal(t, int32(3), m.SegmentId())
+ assert.Equal(t, int64(4), m.SegmentSeq())
+ assert.Equal(t, int64(5), m.EventTsNs())
+ assert.Equal(t, int64(6), m.RecvTsNs())
+
+ assert.Equal(t, 2, m.PropertiesLength())
+ nv := &message_fbs.NameValue{}
+ m.Properties(nv, 0)
+ assert.Equal(t, "n1", string(nv.Name()))
+ assert.Equal(t, "v1", string(nv.Value()))
+ m.Properties(nv, 1)
+ assert.Equal(t, "n2", string(nv.Name()))
+ assert.Equal(t, "v2", string(nv.Value()))
+ assert.Equal(t, []byte("the primary key"), m.Key())
+ assert.Equal(t, []byte("body is here"), m.Data())
+
+ m.MutateSegmentSeq(123)
+ assert.Equal(t, int64(123), m.SegmentSeq())
+
+}
diff --git a/weed/mq/segment/segment_serde.go b/weed/mq/segment/segment_serde.go
new file mode 100644
index 000000000..e076271d6
--- /dev/null
+++ b/weed/mq/segment/segment_serde.go
@@ -0,0 +1 @@
+package segment