blob: 305f5274bbb41f79b22cb290546d5a9282eb52d0 [file]
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
package command
import (
"bytes"
"testing"
"github.com/apache/iggy/foreign/go/contracts"
"github.com/google/uuid"
)
func TestSerialize_TcpFetchMessagesRequest(t *testing.T) {
partitionId := uint32(123)
consumerId, _ := iggcon.NewIdentifier(uint32(42))
streamId, _ := iggcon.NewIdentifier("test_stream_id")
topicId, _ := iggcon.NewIdentifier("test_topic_id")
// Create a sample PollMessages
request := PollMessages{
Consumer: iggcon.NewSingleConsumer(consumerId),
StreamId: streamId,
TopicId: topicId,
PartitionId: &partitionId,
Strategy: iggcon.FirstPollingStrategy(),
Count: 100,
AutoCommit: true,
}
// Serialize the request
serialized, err := request.MarshalBinary()
if err != nil {
t.Error(err)
}
// Expected serialized bytes based on the provided sample request
expected := []byte{
0x01, // Consumer Kind
0x01, // ConsumerId Kind (NumericId)
0x04, // ConsumerId Length (4)
0x2A, 0x00, 0x0, 0x0, // ConsumerId
0x02, // StreamId Kind (StringId)
0x0E, // StreamId Length (14)
0x74, 0x65, 0x73, 0x74, 0x5F, 0x73, 0x74, 0x72, 0x65, 0x61, 0x6D, 0x5F, 0x69, 0x64, // StreamId
0x02, // TopicId Kind (StringId)
0x0D, // TopicId Length (13)
0x74, 0x65, 0x73, 0x74, 0x5F, 0x74, 0x6F, 0x70, 0x69, 0x63, 0x5F, 0x69, 0x64, // TopicId
0x01, // Partition present
0x7B, 0x00, 0x00, 0x00, // PartitionId (123)
0x03, // PollingStrategy Kind
0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // PollingStrategy Value (0)
0x64, 0x00, 0x00, 0x00, // Count (100)
0x01, // AutoCommit
}
// Check if the serialized bytes match the expected bytes
if !areBytesEqual(serialized, expected) {
t.Errorf("Serialized bytes are incorrect. \nExpected:\t%v\nGot:\t\t%v", expected, serialized)
}
}
func areBytesEqual(a, b []byte) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if a[i] != b[i] {
return false
}
}
return true
}
func TestSerialize_SendMessagesRequest(t *testing.T) {
message1 := generateTestMessage("data1")
streamId, _ := iggcon.NewIdentifier("test_stream_id")
topicId, _ := iggcon.NewIdentifier("test_topic_id")
request := SendMessages{
StreamId: streamId,
TopicId: topicId,
Partitioning: iggcon.PartitionId(1),
Messages: []iggcon.IggyMessage{
message1,
},
Compression: iggcon.MESSAGE_COMPRESSION_NONE,
}
// Serialize the request
serialized, err := request.MarshalBinary()
if err != nil {
t.Error(err)
}
// Expected serialized bytes based on the provided sample request
expected := []byte{
0x29, 0x0, 0x0, 0x0, // metadataLength
0x02, // StreamId Kind (StringId)
0x0E, // StreamId Length (14)
0x74, 0x65, 0x73, 0x74, 0x5F, 0x73, 0x74, 0x72, 0x65, 0x61, 0x6D, 0x5F, 0x69, 0x64, // StreamId
0x02, // TopicId Kind (StringId)
0x0D, // TopicId Length (13)
0x74, 0x65, 0x73, 0x74, 0x5F, 0x74, 0x6F, 0x70, 0x69, 0x63, 0x5F, 0x69, 0x64, // TopicId
0x02, // PartitionIdKind
0x04, // Partitioning Length
0x01, 0x00, 0x00, 0x00, // PartitionId (123)
0x01, 0x0, 0x0, 0x0, // MessageCount
0, 0, 0, 0, 120, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, // Index (16*1) bytes
}
expected = append(expected, message1.Header.ToBytes()...)
expected = append(expected, message1.Payload...)
expected = append(expected, message1.UserHeaders...)
// Check if the serialized bytes match the expected bytes
if !bytes.Equal(serialized, expected) {
t.Errorf("Serialized bytes are incorrect. \nExpected:\t%v\nGot:\t\t%v", expected, serialized)
}
}
func createDefaultMessageHeaders() []iggcon.HeaderEntry {
return []iggcon.HeaderEntry{
{Key: iggcon.HeaderKey{Kind: iggcon.String, Value: []byte("HeaderKey1")}, Value: iggcon.HeaderValue{Kind: iggcon.String, Value: []byte("Value 1")}},
{Key: iggcon.HeaderKey{Kind: iggcon.String, Value: []byte("HeaderKey2")}, Value: iggcon.HeaderValue{Kind: iggcon.Uint32, Value: []byte{0x01, 0x02, 0x03, 0x04}}},
}
}
func generateTestMessage(payload string) iggcon.IggyMessage {
msg, _ := iggcon.NewIggyMessage(
[]byte(payload),
iggcon.WithID(uuid.New()),
iggcon.WithUserHeaders(createDefaultMessageHeaders()))
return msg
}