EV-Kafka
view release on metacpan or search on metacpan
t/24_conn_liveness.t view on Meta::CPAN
use strict;
use warnings;
use Test::More;
use IO::Socket::INET;
use Errno qw(EAGAIN EWOULDBLOCK);
use EV;
use EV::Kafka;
# Regression tests for connection liveness:
# L1 handshake phases (ApiVersions/TLS/SASL) are covered by the connect
# timer â a peer that accepts but stays silent fails the conn
# L2 auto_reconnect keeps retrying after a SYNCHRONOUS connect failure
# L5 SASL config errors ("unsupported mechanism", "username required")
# disconnect instead of wedging the state machine
# L9 undef/empty SASL credentials are stored as NULL, not "", without a
# warning; tied ones are fetched
# L6 acks=0 produce callbacks fire (empty success = handed to socket)
# L10 reconnect uses capped exponential backoff with jitter, first retry
# still exactly reconnect_delay_ms
#
# Mock-broker pattern as in t/20-t/23 â no live broker needed.
plan tests => 20;
sub i16 { pack 'n', $_[0] }
sub i32 { pack 'N', $_[0] }
my $apis_body =
i16(0)
. i32(2)
. i16(0) . i16(0) . i16(7) # API_PRODUCE, v0..v7
. i16(18) . i16(0) . i16(0); # API_API_VERSIONS, v0..v0
# SaslHandshake v1 success: no error, mechanisms = ["PLAIN"]
my $sasl_handshake_body = i16(0) . i32(1) . i16(5) . 'PLAIN';
sub mock_broker {
my (%opt) = @_;
my $respond = $opt{respond} || {};
my $server = IO::Socket::INET->new(
LocalAddr => '127.0.0.1', LocalPort => 0, Listen => 1,
Proto => 'tcp', ReuseAddr => 1,
) or BAIL_OUT "cannot bind localhost listener: $!";
$server->blocking(0);
my $ctx = { server => $server };
$ctx->{accept_w} = EV::io fileno($server), EV::READ, sub {
my $fh = $server->accept or return;
$fh->blocking(0);
$ctx->{client_fh} = $fh;
my $incoming = '';
$ctx->{read_w} = EV::io fileno($fh), EV::READ, sub {
my $buf;
my $n = sysread $fh, $buf, 4096;
if (!defined $n) {
return if $!{EAGAIN} || $!{EWOULDBLOCK};
undef $ctx->{read_w};
return;
}
if ($n == 0) { undef $ctx->{read_w}; return; }
$incoming .= $buf;
while (length($incoming) >= 4) {
my $size = unpack 'N', substr($incoming, 0, 4);
last if length($incoming) < 4 + $size;
my $api = unpack 'n', substr($incoming, 4, 2);
my $corr = unpack 'N', substr($incoming, 8, 4);
substr($incoming, 0, 4 + $size) = '';
my $body = $api == 18 ? $apis_body : $respond->{$api};
next unless defined $body;
my $payload = i32($corr) . $body;
syswrite $fh, i32(length $payload) . $payload;
}
};
};
t/24_conn_liveness.t view on Meta::CPAN
Proto => 'tcp', ReuseAddr => 1,
) or BAIL_OUT "cannot bind: $!";
my $port = $sock->sockport;
close $sock;
return $port;
}
# --- L1: a peer that accepts but stays silent fails the handshake ---------
{
my $server = IO::Socket::INET->new(
LocalAddr => '127.0.0.1', LocalPort => 0, Listen => 1,
Proto => 'tcp', ReuseAddr => 1,
) or BAIL_OUT "cannot bind: $!";
$server->blocking(0);
my ($fh, $accept_w);
$accept_w = EV::io fileno($server), EV::READ, sub {
$fh = $server->accept; # accept and stay silent forever
};
my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', undef);
my ($err, $disconnected);
$conn->on_error(sub { $err = $_[0]; });
$conn->on_disconnect(sub { $disconnected++; EV::break });
$conn->connect('127.0.0.1', $server->sockport, 0.3);
my $safety = EV::timer 4, 0, sub { diag "L1 timed out"; EV::break };
EV::run;
like $err, qr/handshake timeout/,
'L1: silent peer fails with handshake timeout';
is $disconnected, 1, 'L1: conn disconnects after handshake timeout';
close $server;
}
# --- L5/L9: SASL config errors disconnect instead of wedging --------------
for my $case (
# [mechanism, username, password, expected error]
['GSSAPI', 'u', 'p', 'unsupported SASL mechanism'],
['SCRAM-SHA-256', undef, undef, 'SCRAM: username required'],
['SCRAM-SHA-256', '', '', 'SCRAM: username required'], # L9: '' == undef
) {
my ($mech, $user, $pass, $expect) = @$case;
my ($port, $broker) = mock_broker(respond => { 17 => $sasl_handshake_body });
my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', undef);
my ($err, $disconnected);
$conn->on_error(sub { $err = $_[0]; });
$conn->on_disconnect(sub { $disconnected++; EV::break });
my @warn;
{ local $SIG{__WARN__} = sub { push @warn, @_ }; $conn->sasl($mech, $user, $pass) }
is "@warn", '', "L9: $mech sasl() does not warn";
$conn->connect('127.0.0.1', $port, 5.0);
my $safety = EV::timer 4, 0, sub { diag "L5 case $mech timed out"; EV::break };
EV::run;
like $err, qr/^\Q$expect\E/, "L5: $mech config error surfaces";
is $disconnected, 1, "L5: $mech config error disconnects the conn";
$conn->DESTROY;
}
# --- L9: tied SASL credentials are fetched, not taken for undef ------------
{
package TiedCred;
sub TIESCALAR { my ($c, $v) = @_; bless \$v, $c }
sub FETCH { ${ $_[0] } }
}
{
my ($port, $broker) = mock_broker(respond => { 17 => $sasl_handshake_body });
my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', undef);
my $err;
$conn->on_error(sub { $err = $_[0]; EV::break });
$conn->on_disconnect(sub { EV::break });
tie my $user, 'TiedCred', 'u';
tie my $pass, 'TiedCred', 'p';
$conn->sasl('SCRAM-SHA-256', $user, $pass);
$conn->connect('127.0.0.1', $port, 5.0);
my $safety = EV::timer 1, 0, sub { EV::break };
EV::run;
unlike $err // '', qr/username required/, 'L9: tied SASL credentials reach SCRAM';
$conn->DESTROY;
}
# --- L9: sasl(undef) turns SASL off without a warning -----------------------
{
my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', undef);
my @warn;
{ local $SIG{__WARN__} = sub { push @warn, @_ }; $conn->sasl(undef) }
is "@warn", '', 'L9: sasl(undef) does not warn';
$conn->DESTROY;
}
# --- L6: acks=0 produce callback fires with empty success -----------------
{
my ($conn, $broker) = ready_conn();
my ($fired, $res, $error);
$conn->produce('t', 0, 'k', 'v', { acks => 0 }, sub {
$fired++;
($res, $error) = @_;
});
is $fired, 1,
'L6: acks=0 callback fires immediately (handed to the socket)';
ok((ref $res eq 'HASH' && !keys %$res), 'L6: empty success result');
ok !defined($error), 'L6: no error';
$conn->DESTROY;
}
# --- L2: auto_reconnect survives a synchronous connect failure ------------
{
my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', undef);
my @errors;
$conn->auto_reconnect(1, 100);
$conn->on_error(sub { push @errors, $_[0]; EV::break if @errors >= 2 });
$conn->connect('x' x 300, 9092, 1.0); # getaddrinfo fails synchronously
my $safety = EV::timer 3, 0, sub { EV::break };
EV::run;
is(scalar(@errors) >= 2 ? 2 : scalar(@errors), 2,
'L2: auto_reconnect keeps retrying after synchronous resolve failure');
like $errors[0], qr/^resolve:/, 'L2: error mentions resolve';
$conn->DESTROY;
}
# --- L10: reconnect backoff, first retry unchanged -------------------------
{
my $port = free_port(); # nothing listening: connection refused
my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', undef);
my @times;
$conn->auto_reconnect(1, 200);
$conn->on_error(sub { push @times, EV::now; EV::break if @times >= 4 });
$conn->connect('127.0.0.1', $port, 1.0);
my $safety = EV::timer 4, 0, sub { EV::break };
EV::run;
# failures: e0 ~ 0 (refused), e1 ~ 0.2s (first retry, exact base),
# e2 ~ 0.6s, e3 ~ 1.4s with 2x/4x backoff and +/-25% jitter
cmp_ok scalar(@times), '>=', 4, 'L10: four failures observed';
( run in 0.980 second using v1.01-cache-2.11-cpan-007c89162af )