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 2.323 seconds using v1.01-cache-2.11-cpan-007c89162af )