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 )