EV-Kafka
view release on metacpan or search on metacpan
src/EV__Kafka.xs view on Meta::CPAN
done:
hv_store(result, "topics", 6, newRV_inc((SV*)topics_av), 0);
return sv_2mortal(newRV_noinc((SV*)result));
}
/* ================================================================
* RecordBatch decoder (for Fetch responses)
* ================================================================ */
/* Decode records from a RecordBatch, push them as hashrefs onto records_av.
* Returns number of records decoded, or -1 on error. */
static int kf_decode_record_batch(pTHX_ const char *data, size_t len,
AV *records_av, int64_t *out_base_offset)
{
const char *p = data;
const char *end = data + len;
int n;
if (end - p < 12) return -1;
int64_t base_offset = kf_read_i64(p); p += 8;
if (out_base_offset) *out_base_offset = base_offset;
t/16_codec_roundtrip.t view on Meta::CPAN
plan tests => 20;
# Single record, no compression.
{
my $bytes = EV::Kafka::_test_encode_batch(
[{ key => 'k1', value => 'v1' }]
);
ok length($bytes) > 0, 'single record encodes to non-empty bytes';
my $decoded = EV::Kafka::_test_decode_batch($bytes);
ok $decoded, 'single record decodes';
is scalar @$decoded, 1, 'one record back';
is $decoded->[0]{key}, 'k1', 'key round-trips';
is $decoded->[0]{value}, 'v1', 'value round-trips';
}
# Many records.
{
my @recs = map { { key => "k$_", value => "v$_" } } 1..50;
my $bytes = EV::Kafka::_test_encode_batch(\@recs);
my $decoded = EV::Kafka::_test_decode_batch($bytes);
is scalar @$decoded, 50, '50-record batch round-trips';
is $decoded->[0]{value}, 'v1', 'first record value preserved';
is $decoded->[49]{value}, 'v50', 'last record value preserved';
}
# Headers preserved.
{
my $bytes = EV::Kafka::_test_encode_batch(
[{ key => 'k', value => 'v', headers => { 'h1' => 'a', 'h2' => 'b' } }]
);
my $decoded = EV::Kafka::_test_decode_batch($bytes);
is_deeply $decoded->[0]{headers}, { h1 => 'a', h2 => 'b' },
'headers round-trip';
}
# Null key / null value.
{
my $bytes = EV::Kafka::_test_encode_batch(
[{ key => undef, value => undef }]
);
my $decoded = EV::Kafka::_test_decode_batch($bytes);
ok !defined $decoded->[0]{key}, 'null key preserved';
ok !defined $decoded->[0]{value}, 'null value preserved';
}
# Idempotent producer fields are encoded but stripped on decode (decoder
# doesn't surface producer_id/epoch â that's broker-only state).
{
my $bytes = EV::Kafka::_test_encode_batch(
[{ key => 'k', value => 'v' }],
{ producer_id => 42, producer_epoch => 3, base_sequence => 100 },
);
my $decoded = EV::Kafka::_test_decode_batch($bytes);
is scalar @$decoded, 1, 'idempotent batch decodes';
is $decoded->[0]{value}, 'v', 'payload survives idempotent encode';
}
# Compression: gzip.
SKIP: {
skip "gzip support not built", 2 unless eval {
my $b = EV::Kafka::_test_encode_batch(
[{ key => 'k', value => 'gzip-payload' }],
{ compression => 1 },
);
defined $b;
};
my $bytes = EV::Kafka::_test_encode_batch(
[ map { { key => "k$_", value => "gzip-$_" } } 1..20 ],
{ compression => 1 },
);
ok length($bytes) > 0, 'gzip-compressed batch encodes';
my $decoded = EV::Kafka::_test_decode_batch($bytes);
is scalar @$decoded, 20, 'gzip batch round-trips';
}
# Compression: lz4.
SKIP: {
skip "lz4 support not built", 2 unless eval {
my $b = EV::Kafka::_test_encode_batch(
[{ key => 'k', value => 'lz4-payload' }],
{ compression => 3 },
);
defined $b;
};
my $bytes = EV::Kafka::_test_encode_batch(
[ map { { key => "k$_", value => "lz4-$_" } } 1..20 ],
{ compression => 3 },
);
ok length($bytes) > 0, 'lz4-compressed batch encodes';
my $decoded = EV::Kafka::_test_decode_batch($bytes);
is scalar @$decoded, 20, 'lz4 batch round-trips';
}
# Highly compressible payload â exercises decompression buffer growth.
SKIP: {
skip "lz4 not built", 1 unless eval {
defined EV::Kafka::_test_encode_batch([{key=>'k',value=>'x'}], {compression=>3});
};
my $payload = 'a' x 100_000; # 100k of 'a' compresses ~100x
my $bytes = EV::Kafka::_test_encode_batch(
[{ key => 'big', value => $payload }],
{ compression => 3 },
);
my $decoded = EV::Kafka::_test_decode_batch($bytes);
is length($decoded->[0]{value}), length($payload),
'highly-compressible lz4 payload survives the doubling decode loop';
}
# Compression: zstd (code 4).
SKIP: {
skip "zstd not built", 1 unless eval {
my $b = EV::Kafka::_test_encode_batch([{key=>'k',value=>'z'}], {compression=>4});
my $r = EV::Kafka::_test_decode_batch($b);
$r && $r->[0]{value} eq 'z';
};
my $bytes = EV::Kafka::_test_encode_batch(
[ map { { key => "k$_", value => "zstd-$_" } } 1..20 ],
{ compression => 4 },
);
my $decoded = EV::Kafka::_test_decode_batch($bytes);
is scalar @$decoded, 20, 'zstd batch round-trips';
}
# Compression: snappy (code 2).
SKIP: {
skip "snappy not built", 1 unless eval {
my $b = EV::Kafka::_test_encode_batch([{key=>'k',value=>'s'}], {compression=>2});
my $r = EV::Kafka::_test_decode_batch($b);
$r && $r->[0]{value} eq 's';
};
my $bytes = EV::Kafka::_test_encode_batch(
[ map { { key => "k$_", value => "snappy-$_" } } 1..20 ],
{ compression => 2 },
);
my $decoded = EV::Kafka::_test_decode_batch($bytes);
is scalar @$decoded, 20, 'snappy batch round-trips';
}
t/18_response_parsers.t view on Meta::CPAN
my $body =
i32(1) # topics array len
. kstr('mytopic')
. i32(1) # partitions array len
. i32(0) # partition
. i16(0) # error_code
. i64(42); # base_offset
my $r = EV::Kafka::_test_parse_response('produce', 0, $body);
ok ref $r eq 'HASH', 'produce v0 response parses to hashref';
is $r->{topics}[0]{topic}, 'mytopic',
'produce v0 topic name decoded';
is $r->{topics}[0]{partitions}[0]{base_offset}, 42,
'produce v0 base_offset decoded';
}
# --- Metadata response v0: brokers + topics ---
{
my $body =
i32(1) # brokers array len
. i32(0) # node_id
. kstr('host1')
. i32(9092) # port
. i32(1) # topics array len
t/18_response_parsers.t view on Meta::CPAN
. kstr('mytopic')
. i32(1) # partitions array len
. i16(0) # partition error_code
. i32(0) # partition id
. i32(0) # leader id
. i32(0) # replicas array len
. i32(0); # isr array len
my $r = EV::Kafka::_test_parse_response('metadata', 0, $body);
ok ref $r eq 'HASH', 'metadata v0 parses';
is scalar @{$r->{brokers}}, 1, 'one broker';
is $r->{brokers}[0]{host}, 'host1', 'broker host decoded';
is $r->{brokers}[0]{port}, 9092, 'broker port decoded';
is $r->{topics}[0]{name}, 'mytopic', 'topic name decoded';
is $r->{topics}[0]{partitions}[0]{leader}, 0, 'partition leader decoded';
}
# --- Truncated metadata response: should not crash, returns partial data ---
{
my $body = i32(1); # claims one broker but no broker bytes follow
my $r = EV::Kafka::_test_parse_response('metadata', 0, $body);
ok ref $r eq 'HASH', 'truncated metadata returns parsed-so-far hash';
is scalar @{$r->{brokers} // []}, 0,
'truncated metadata yields empty brokers list (parser bails on bounds check)';
}
# --- Heartbeat response v0: just throttle(?) + error_code ---
{
my $body = i16(0); # error_code = 0, no throttle in v0
my $r = EV::Kafka::_test_parse_response('heartbeat', 0, $body);
ok ref $r eq 'HASH', 'heartbeat v0 parses';
is $r->{error_code}, 0, 'heartbeat error_code decoded';
}
# --- LeaveGroup response v0: error_code only ---
{
my $body = i16(0);
my $r = EV::Kafka::_test_parse_response('leave_group', 0, $body);
ok ref $r eq 'HASH', 'leave_group v0 parses';
}
# --- EndTxn response v0: throttle_ms(i32) + error_code ---
t/18_response_parsers.t view on Meta::CPAN
is $r->{host}, 'coord-host', 'find_coordinator host';
}
# --- SyncGroup response v0: error_code + assignment(BYTES) ---
{
my $assignment = "\x00\x01" . # version
i32(0); # zero topics
my $body = i16(0) . i32(length $assignment) . $assignment;
my $r = EV::Kafka::_test_parse_response('sync_group', 0, $body);
ok ref $r eq 'HASH', 'sync_group v0 parses';
is $r->{error_code}, 0, 'sync_group error_code decoded';
}
# --- OffsetFetch response v0: topics array with partitions ---
{
my $body =
i32(1) # topics
. kstr('mytopic')
. i32(1) # partitions
. i32(0) # partition
. pack('q>', 100) # offset
t/18_response_parsers.t view on Meta::CPAN
}
# --- InitProducerId response v1: throttle prefix is added at v1+ ---
{
my $body =
i32(0) # throttle_time_ms (v1+)
. i16(0) # error_code
. pack('q>', 12345) # producer_id
. i16(0); # producer_epoch
my $r = EV::Kafka::_test_parse_response('init_producer_id', 1, $body);
is $r->{producer_id}, 12345, 'init_producer_id pid decoded';
is $r->{producer_epoch}, 0, 'init_producer_id epoch decoded';
}
t/21_malformed_responses.t view on Meta::CPAN
. i16(0) # error_code
. i64(0) # high_watermark
. i64(0) # last_stable_offset (v4+)
. i32(0) # aborted_transactions count (v4+)
. i32(length $batch) # record_set BYTES size
. $batch;
my $r = EV::Kafka::_test_parse_response('fetch', 4, $body);
ok ref $r eq 'HASH',
'fetch v4 with batch_length 0x7FFFFFF9 returns without crashing';
is_deeply $r->{topics}[0]{partitions}[0]{records}, [],
'oversized batch_length rejected, no records decoded';
}
# --- Metadata v9 with a hostile compact-string length for broker host ---
# uvarint 80 80 80 80 10 decodes to raw = 2^32; len = raw - 1 used to wrap
# to -1 in the int32_t cast, escaping as a negative strlen
# (panic: sv_setpvn_fresh called with negative strlen -1).
{
my $body =
i32(0) # throttle_time_ms
. "\x02" # brokers: compact array, count+1 = 2 => 1 broker
t/25_lz4_frame.t view on Meta::CPAN
my $batch = EV::Kafka::_test_encode_batch($records,
{ compression => 3, timestamp => 1000 });
is unpack('H8', substr($batch, 61, 4)), '04224d18',
'L7: LZ4-compressed batch starts with the frame magic 04 22 4D 18';
}
# --- L7: frame roundtrip through the real decode path ----------------------
{
my $batch = EV::Kafka::_test_encode_batch($records,
{ compression => 3, timestamp => 1000 });
my $decoded = EV::Kafka::_test_decode_batch($batch);
is_deeply $decoded, [{ offset => 0, timestamp => 1000,
key => 'hello', value => 'world' }],
'L7: frame format roundtrips through kf_decode_record_batch';
}
# --- L7: a reference frame (golden bytes) decodes --------------------------
# Produced once by the frame encoder; pinned here so the test does not
# depend on the build's own encoder agreeing with itself.
{
my $golden = pack 'H*',
'000000000000000000000051000000000255d2dad3000300000000000000'
. '00000003e800000000000003e8ffffffffffffffffffffffffffff0000'
. '000104224d184040c011000080200000000a68656c6c6f0a776f726c64'
. '0000000000';
my $decoded = EV::Kafka::_test_decode_batch($golden);
is_deeply $decoded, [{ offset => 0, timestamp => 1000,
key => 'hello', value => 'world' }],
'L7: reference LZ4 frame decodes to the original record';
}
# --- L8: failed decompression must not parse payload as records ------------
{
# Valid uncompressed batch, but attributes claim LZ4: the plaintext
# payload is not a frame, decompression fails â must return undef
# (pre-fix it fell through and delivered the plaintext as records).
my $batch = EV::Kafka::_test_encode_batch($records,
xt/fuzz_structured.t view on Meta::CPAN
my $r = EV::Kafka::_test_parse_response('produce', 7, $bytes);
ok((ref $r eq 'HASH')
&& $r->{topics}[0]{partitions}[0]{base_offset} == 42,
'template produce v7 parses to expected structure');
}
{
my ($bytes) = tpl_fetch_v4();
my $r = EV::Kafka::_test_parse_response('fetch', 4, $bytes);
my $recs = (ref $r eq 'HASH') ? $r->{topics}[0]{partitions}[0]{records} : undef;
ok((ref $recs eq 'ARRAY') && @$recs == 1 && $recs->[0]{key} eq 'k',
'template fetch v4 parses, record batch decoded');
}
{
my ($bytes) = tpl_fetch_v7();
my $r = EV::Kafka::_test_parse_response('fetch', 7, $bytes);
my $recs = (ref $r eq 'HASH') ? $r->{topics}[0]{partitions}[0]{records} : undef;
ok((ref $recs eq 'ARRAY') && @$recs == 1 && $recs->[0]{value} eq 'v',
'template fetch v7 parses, record batch decoded');
}
{
my ($bytes) = tpl_list_offsets_v1();
my $r = EV::Kafka::_test_parse_response('list_offsets', 1, $bytes);
ok((ref $r eq 'HASH')
&& $r->{topics}[0]{partitions}[0]{offset} == 7,
'template list_offsets v1 parses to expected structure');
}
{
my ($bytes) = tpl_batch();
( run in 2.663 seconds using v1.01-cache-2.11-cpan-2c0d6866c4f )