|
| 1 | +%% This Source Code Form is subject to the terms of the Mozilla Public |
| 2 | +%% License, v. 2.0. If a copy of the MPL was not distributed with this |
| 3 | +%% file, You can obtain one at https://mozilla.org/MPL/2.0/. |
| 4 | +%% |
| 5 | +%% Copyright (c) 2007-2024 Broadcom. All Rights Reserved. The term “Broadcom” refers to Broadcom Inc. and/or its subsidiaries. All rights reserved. |
| 6 | +%% |
| 7 | + |
| 8 | +%% There is no open source Erlang RabbitMQ Stream client. |
| 9 | +%% Therefore, we have to build the Stream protocol commands manually. |
| 10 | + |
| 11 | +-module(stream_test_utils). |
| 12 | + |
| 13 | +-compile([export_all, nowarn_export_all]). |
| 14 | + |
| 15 | +-include_lib("amqp10_common/include/amqp10_framing.hrl"). |
| 16 | + |
| 17 | +-define(RESPONSE_CODE_OK, 1). |
| 18 | + |
| 19 | +connect(Config, Node) -> |
| 20 | + StreamPort = rabbit_ct_broker_helpers:get_node_config(Config, Node, tcp_port_stream), |
| 21 | + {ok, Sock} = gen_tcp:connect("localhost", StreamPort, [{active, false}, {mode, binary}]), |
| 22 | + |
| 23 | + C0 = rabbit_stream_core:init(0), |
| 24 | + PeerPropertiesFrame = rabbit_stream_core:frame({request, 1, {peer_properties, #{}}}), |
| 25 | + ok = gen_tcp:send(Sock, PeerPropertiesFrame), |
| 26 | + {{response, 1, {peer_properties, _, _}}, C1} = receive_stream_commands(Sock, C0), |
| 27 | + |
| 28 | + ok = gen_tcp:send(Sock, rabbit_stream_core:frame({request, 1, sasl_handshake})), |
| 29 | + {{response, _, {sasl_handshake, _, _}}, C2} = receive_stream_commands(Sock, C1), |
| 30 | + Username = <<"guest">>, |
| 31 | + Password = <<"guest">>, |
| 32 | + Null = 0, |
| 33 | + PlainSasl = <<Null:8, Username/binary, Null:8, Password/binary>>, |
| 34 | + ok = gen_tcp:send(Sock, rabbit_stream_core:frame({request, 2, {sasl_authenticate, <<"PLAIN">>, PlainSasl}})), |
| 35 | + {{response, 2, {sasl_authenticate, _}}, C3} = receive_stream_commands(Sock, C2), |
| 36 | + {{tune, DefaultFrameMax, _}, C4} = receive_stream_commands(Sock, C3), |
| 37 | + |
| 38 | + ok = gen_tcp:send(Sock, rabbit_stream_core:frame({response, 0, {tune, DefaultFrameMax, 0}})), |
| 39 | + ok = gen_tcp:send(Sock, rabbit_stream_core:frame({request, 3, {open, <<"/">>}})), |
| 40 | + {{response, 3, {open, _, _ConnectionProperties}}, C5} = receive_stream_commands(Sock, C4), |
| 41 | + {ok, Sock, C5}. |
| 42 | + |
| 43 | +create_stream(Sock, C0, Stream) -> |
| 44 | + CreateStreamFrame = rabbit_stream_core:frame({request, 1, {create_stream, Stream, #{}}}), |
| 45 | + ok = gen_tcp:send(Sock, CreateStreamFrame), |
| 46 | + {{response, 1, {create_stream, ?RESPONSE_CODE_OK}}, C1} = receive_stream_commands(Sock, C0), |
| 47 | + {ok, C1}. |
| 48 | + |
| 49 | +declare_publisher(Sock, C0, Stream, PublisherId) -> |
| 50 | + DeclarePublisherFrame = rabbit_stream_core:frame({request, 1, {declare_publisher, PublisherId, <<>>, Stream}}), |
| 51 | + ok = gen_tcp:send(Sock, DeclarePublisherFrame), |
| 52 | + {{response, 1, {declare_publisher, ?RESPONSE_CODE_OK}}, C1} = receive_stream_commands(Sock, C0), |
| 53 | + {ok, C1}. |
| 54 | + |
| 55 | +subscribe(Sock, C0, Stream, SubscriptionId, InitialCredit) -> |
| 56 | + SubscribeFrame = rabbit_stream_core:frame({request, 1, {subscribe, SubscriptionId, Stream, _OffsetSpec = first, InitialCredit, _Props = #{}}}), |
| 57 | + ok = gen_tcp:send(Sock, SubscribeFrame), |
| 58 | + {{response, 1, {subscribe, ?RESPONSE_CODE_OK}}, C1} = receive_stream_commands(Sock, C0), |
| 59 | + {ok, C1}. |
| 60 | + |
| 61 | +publish(Sock, C0, PublisherId, Sequence0, Payloads) -> |
| 62 | + SeqIds = lists:seq(Sequence0, Sequence0 + length(Payloads) - 1), |
| 63 | + Messages = [simple_entry(Seq, P) |
| 64 | + || {Seq, P} <- lists:zip(SeqIds, Payloads)], |
| 65 | + {ok, SeqIds, C1} = publish_entries(Sock, C0, PublisherId, length(Messages), Messages), |
| 66 | + {ok, C1}. |
| 67 | + |
| 68 | +publish_entries(Sock, C0, PublisherId, MsgCount, Messages) -> |
| 69 | + PublishFrame1 = rabbit_stream_core:frame({publish, PublisherId, MsgCount, Messages}), |
| 70 | + ok = gen_tcp:send(Sock, PublishFrame1), |
| 71 | + {{publish_confirm, PublisherId, SeqIds}, C1} = receive_stream_commands(Sock, C0), |
| 72 | + {ok, SeqIds, C1}. |
| 73 | + |
| 74 | +%% Streams contain AMQP 1.0 encoded messages. |
| 75 | +%% In this case, the AMQP 1.0 encoded message contains a single data section. |
| 76 | +simple_entry(Sequence, Body) |
| 77 | + when is_binary(Body) -> |
| 78 | + DataSect = iolist_to_binary(amqp10_framing:encode_bin(#'v1_0.data'{content = Body})), |
| 79 | + DataSectSize = byte_size(DataSect), |
| 80 | + <<Sequence:64, 0:1, DataSectSize:31, DataSect:DataSectSize/binary>>. |
| 81 | + |
| 82 | +%% Streams contain AMQP 1.0 encoded messages. |
| 83 | +%% In this case, the AMQP 1.0 encoded message consists of an application-properties section and a data section. |
| 84 | +simple_entry(Sequence, Body, AppProps) |
| 85 | + when is_binary(Body) -> |
| 86 | + AppPropsSect = iolist_to_binary(amqp10_framing:encode_bin(AppProps)), |
| 87 | + DataSect = iolist_to_binary(amqp10_framing:encode_bin(#'v1_0.data'{content = Body})), |
| 88 | + Sects = <<AppPropsSect/binary, DataSect/binary>>, |
| 89 | + SectSize = byte_size(Sects), |
| 90 | + <<Sequence:64, 0:1, SectSize:31, Sects:SectSize/binary>>. |
| 91 | + |
| 92 | +%% Here, each AMQP 1.0 encoded message consists of an application-properties section and a data section. |
| 93 | +%% All data sections are delivered uncompressed in 1 batch. |
| 94 | +sub_batch_entry_uncompressed(Sequence, Bodies) -> |
| 95 | + Batch = lists:foldl(fun(Body, Acc) -> |
| 96 | + AppProps = #'v1_0.application_properties'{ |
| 97 | + content = [{{utf8, <<"my key">>}, {utf8, <<"my value">>}}]}, |
| 98 | + Sect0 = iolist_to_binary(amqp10_framing:encode_bin(AppProps)), |
| 99 | + Sect1 = iolist_to_binary(amqp10_framing:encode_bin(#'v1_0.data'{content = Body})), |
| 100 | + Sect = <<Sect0/binary, Sect1/binary>>, |
| 101 | + <<Acc/binary, 0:1, (byte_size(Sect)):31, Sect/binary>> |
| 102 | + end, <<>>, Bodies), |
| 103 | + Size = byte_size(Batch), |
| 104 | + <<Sequence:64, 1:1, 0:3, 0:4, (length(Bodies)):16, Size:32, Size:32, Batch:Size/binary>>. |
| 105 | + |
| 106 | +%% Here, each AMQP 1.0 encoded message contains a single data section. |
| 107 | +%% All data sections are delivered in 1 gzip compressed batch. |
| 108 | +sub_batch_entry_compressed(Sequence, Bodies) -> |
| 109 | + Uncompressed = lists:foldl(fun(Body, Acc) -> |
| 110 | + Bin = iolist_to_binary(amqp10_framing:encode_bin(#'v1_0.data'{content = Body})), |
| 111 | + <<Acc/binary, Bin/binary>> |
| 112 | + end, <<>>, Bodies), |
| 113 | + Compressed = zlib:gzip(Uncompressed), |
| 114 | + CompressedLen = byte_size(Compressed), |
| 115 | + <<Sequence:64, 1:1, 1:3, 0:4, (length(Bodies)):16, (byte_size(Uncompressed)):32, |
| 116 | + CompressedLen:32, Compressed:CompressedLen/binary>>. |
| 117 | + |
| 118 | +receive_stream_commands(Sock, C0) -> |
| 119 | + case rabbit_stream_core:next_command(C0) of |
| 120 | + empty -> |
| 121 | + case gen_tcp:recv(Sock, 0, 5000) of |
| 122 | + {ok, Data} -> |
| 123 | + C1 = rabbit_stream_core:incoming_data(Data, C0), |
| 124 | + case rabbit_stream_core:next_command(C1) of |
| 125 | + empty -> |
| 126 | + {ok, Data2} = gen_tcp:recv(Sock, 0, 5000), |
| 127 | + rabbit_stream_core:next_command( |
| 128 | + rabbit_stream_core:incoming_data(Data2, C1)); |
| 129 | + Res -> |
| 130 | + Res |
| 131 | + end; |
| 132 | + {error, Err} -> |
| 133 | + ct:fail("error receiving stream data ~w", [Err]) |
| 134 | + end; |
| 135 | + Res -> |
| 136 | + Res |
| 137 | + end. |
0 commit comments