EV-Kafka
view release on metacpan or search on metacpan
eg/tls_sasl.pl view on Meta::CPAN
#!/usr/bin/env perl
# TLS + SASL/SCRAM-SHA-256 example.
#
# Set KAFKA_BROKER to a TLS listener (e.g. host:9093). Provide the CA cert
# in KAFKA_TLS_CA, and credentials via KAFKA_USER and KAFKA_PASS. The
# client warns at construction if SASL is configured without TLS.
#
# Quick local test with Redpanda:
# docker run -e RP_BOOTSTRAP_USER=admin:secret123 \
# redpandadata/redpanda:latest start --kafka-addr=...
use strict;
use warnings;
use EV;
use EV::Kafka;
lib/EV/Kafka.pm view on Meta::CPAN
$cfg->{_txn_gen} = 0; # bumped on every txn state transition
$cfg->{_conn_state} = {}; # "$conn" => {auto_reconnect, was_ready}
# Consumer state
$cfg->{assignments} = []; # [{topic, partition, offset}]
$cfg->{fetch_active} = 0;
$cfg->{group} = undef;
my $self = bless { cfg => $cfg, loop => $loop }, "${class}::Client";
# Warn on credentials over plaintext.
if ($cfg->{sasl} && !$cfg->{tls}) {
my $mech = $cfg->{sasl}{mechanism} // '';
if ($mech eq 'PLAIN' || $mech =~ /^SCRAM-/) {
warn "EV::Kafka: SASL $mech configured without TLS â "
. "credentials will be sent over plaintext\n";
}
}
return $self;
}
package EV::Kafka::Client;
use EV;
use Carp 'croak';
use Scalar::Util 'weaken';
t/24_conn_liveness.t view on Meta::CPAN
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] }
t/24_conn_liveness.t view on Meta::CPAN
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;
( run in 1.022 second using v1.01-cache-2.11-cpan-007c89162af )