view release on metacpan or search on metacpan
lib/EV/Kafka.pm view on Meta::CPAN
. "credentials will be sent over plaintext\n";
}
}
return $self;
}
package EV::Kafka::Client;
use EV;
use Carp 'croak';
use Scalar::Util 'weaken';
sub _any_conn {
my ($self) = @_;
my $cfg = $self->{cfg};
return undef if $cfg->{closed};
my $conn = $cfg->{bootstrap_conn};
for my $c (values %{$cfg->{conns}}) {
if ($c->connected) { $conn = $c; last }
}
return ($conn && $conn->connected) ? $conn : undef;
lib/EV/Kafka.pm view on Meta::CPAN
my $info = $cfg->{broker_map}{$node_id};
return undef unless $info;
my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', $self->{loop});
$self->_configure_conn($conn);
$conn->auto_reconnect(1, 1000);
$cfg->{_conn_state}{"$conn"}{auto_reconnect} = 1;
$cfg->{conns}{$node_id} = $conn;
weaken(my $weak = $self);
$conn->on_connect(sub {
my $s = $weak or return;
my $cfg = $s->{cfg};
return if $cfg->{closed};
$cfg->{_conn_state}{"$conn"}{was_ready} = 1;
$s->_drain_pending_for($node_id);
});
$conn->connect($info->{host}, $info->{port}, 10.0);
return $conn;
lib/EV/Kafka.pm view on Meta::CPAN
$self->_wire_conn_handlers($conn);
}
# Error/disconnect reporting for client-managed conns. An
# auto_reconnect conn reports its loss ONCE via on_disconnect; reconnect
# attempts stay silent so on_error (default: die) doesn't storm.
sub _wire_conn_handlers {
my ($self, $conn) = @_;
my $cfg = $self->{cfg};
$cfg->{_conn_state}{"$conn"} //= { auto_reconnect => 0, was_ready => 0 };
weaken(my $weak_cfg = $cfg);
$conn->on_error(sub {
my $cfg = $weak_cfg or return;
return if $cfg->{closed};
my $st = $cfg->{_conn_state}{"$conn"} // {};
return if $st->{auto_reconnect} && !$conn->connected;
$cfg->{on_error}->($_[0]) if $cfg->{on_error};
});
$conn->on_disconnect(sub {
my $cfg = $weak_cfg or return;
return if $cfg->{closed};
lib/EV/Kafka.pm view on Meta::CPAN
}
my ($host, $port) = @{$bs[$idx]};
my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', $self->{loop});
$self->_configure_conn($conn);
# Keeps the attempt conn alive until on_connect/on_error fires,
# without a closure cycle.
$cfg->{_bootstrap_attempt} = $conn;
weaken(my $weak = $self);
# Replace the forwarding on_error from _configure_conn with the
# try-next-broker handler for the duration of the bootstrap.
$conn->on_error(sub {
my $s = $weak or return;
my $cfg = $s->{cfg};
return if $cfg->{closed};
$cfg->{_bootstrap_attempt} = undef;
$s->_bootstrap_try($idx + 1);
});
lib/EV/Kafka.pm view on Meta::CPAN
sub _refresh_metadata {
my ($self) = @_;
my $cfg = $self->{cfg};
return if $cfg->{closed} || $cfg->{meta_pending};
$cfg->{meta_pending} = 1;
my $conn = $self->_any_conn;
unless ($conn) { $cfg->{meta_pending} = 0; return }
weaken(my $weak = $self);
$conn->metadata(undef, sub {
my ($meta, $err) = @_;
my $s = $weak or return;
my $cfg = $s->{cfg};
$cfg->{meta_pending} = 0;
return if $cfg->{closed};
if ($err) {
$cfg->{on_error}->("metadata: $err") if $cfg->{on_error};
return;
}
lib/EV/Kafka.pm view on Meta::CPAN
$s->_arm_metadata_timer;
});
}
sub _arm_metadata_timer {
my ($self) = @_;
my $cfg = $self->{cfg};
return if $cfg->{_meta_timer};
my $interval = $cfg->{metadata_refresh} || 0;
return if $interval <= 0;
weaken(my $weak = $self);
$cfg->{_meta_timer} = EV::timer $interval, $interval, sub {
return unless $weak;
$weak->_refresh_metadata unless $weak->{cfg}{meta_pending};
};
}
sub _disarm_metadata_timer {
my ($self) = @_;
undef $self->{cfg}{_meta_timer};
}
lib/EV/Kafka.pm view on Meta::CPAN
$cfg->{pending_ops} = \@keep;
$cfg->{on_error}->($msg) if !$reported && $cfg->{on_error};
return;
}
$cfg->{meta_pending} = 1;
my $conn = $self->_any_conn;
unless ($conn) { $cfg->{meta_pending} = 0; return }
weaken(my $weak = $self);
$conn->metadata([$topic], sub {
my ($meta, $err) = @_;
my $s = $weak or return;
my $cfg = $s->{cfg};
$cfg->{meta_pending} = 0;
return if $cfg->{closed};
if ($err) {
$cfg->{on_error}->("metadata: $err") if $cfg->{on_error};
return;
}
lib/EV/Kafka.pm view on Meta::CPAN
$cfg->{_pid_cbs} = [];
push @{$cfg->{_pid_cbs}}, $cb if $cb;
$self->_pid_attempt(1);
}
sub _pid_attempt {
my ($self, $attempt) = @_;
my $cfg = $self->{cfg};
return if $cfg->{closed};
weaken(my $weak = $self);
my $retry = sub {
my ($err) = @_;
my $s = $weak or return;
if ($attempt < 3) {
my $t; $t = EV::timer 0.5, 0, sub {
undef $t;
$weak->_pid_attempt($attempt + 1) if $weak;
};
} else {
$s->_pid_finish($err);
lib/EV/Kafka.pm view on Meta::CPAN
if (ref $a eq 'CODE') { $cb = $a }
elsif (ref $a eq 'HASH') { %opts = %$a }
}
for my $k (keys %opts) {
croak "EV::Kafka: produce: unknown option '$k'"
unless $k eq 'partition' || $k eq 'headers';
}
my $cfg = $self->{cfg};
croak "EV::Kafka: client is closed" if $cfg->{closed};
weaken(my $weak = $self);
# ensure we have metadata
unless ($cfg->{meta}) {
push @{$cfg->{pending_ops}}, {
topic => $topic,
cb => $cb,
run => sub { $weak->produce($topic, $key, $value, @rest) if $weak },
};
$self->_refresh_metadata unless $cfg->{meta_pending};
return;
lib/EV/Kafka.pm view on Meta::CPAN
my $batch_bytes = 0;
for my $b (@$batch) {
$batch_bytes += length($b->{rec}{value} // '') + length($b->{rec}{key} // '') + 20;
}
if ($batch_bytes >= $cfg->{batch_size}) {
$self->_flush_batch($topic, $partition, $conn);
} elsif (!$cfg->{_linger_active}) {
# start linger timer
$cfg->{_linger_active} = 1;
weaken(my $weak = $self);
$cfg->{_linger_timer} = EV::timer $cfg->{linger_ms} / 1000.0, 0, sub {
$cfg->{_linger_active} = 0;
$weak->_flush_all_batches if $weak;
};
}
}
sub _flush_batch {
my ($self, $topic, $partition, $conn) = @_;
my $cfg = $self->{cfg};
lib/EV/Kafka.pm view on Meta::CPAN
$self->_add_txn_partition($topic, $partition) if $cfg->{_txn_active};
# retry count persists on the batch across re-queues
$cfg->{_batch_retries}{$bkey} //= 3;
# Transaction generation at send time: a late response must never
# re-queue a batch from an aborted/committed transaction.
my $txn_gen = $cfg->{_txn_gen};
weaken(my $weak_self = $self);
$conn->produce_batch($topic, $partition, \@records, \%popts, sub {
my ($result, $err) = @_;
delete $cfg->{_inflight}{$bkey} if $idempotent;
my $retriable = 0;
my $fatal_seq = 0;
my $dup_seq = 0;
if (!$err && $result && ref $result->{topics} eq 'ARRAY') {
for my $t (@{$result->{topics}}) {
for my $p (@{$t->{partitions} // []}) {
lib/EV/Kafka.pm view on Meta::CPAN
my ($topic, $partition) = split /:/, $bkey, 2;
my $leader_id = $self->_get_leader($topic, $partition);
unless (defined $leader_id) { $skipped++; next }
my $conn = $self->_get_or_create_conn($leader_id);
unless ($conn && $conn->connected) { $skipped++; next }
$self->_flush_batch($topic, $partition, $conn);
}
# re-arm timer if batches were skipped (connection not yet ready)
if ($skipped && keys %{$cfg->{batches}}) {
$cfg->{_linger_active} = 1;
weaken(my $weak = $self);
$cfg->{_linger_timer} = EV::timer 0.1, 0, sub {
$cfg->{_linger_active} = 0;
$weak->_flush_all_batches if $weak;
};
}
}
sub produce_many {
my ($self, $messages, $cb) = @_;
croak "EV::Kafka: client is closed" if $self->{cfg}{closed};
lib/EV/Kafka.pm view on Meta::CPAN
my ($self, $cb) = @_;
my $cfg = $self->{cfg};
croak "EV::Kafka: client is closed" if $cfg->{closed};
# Empty assignment is a successful empty poll, not a dropped callback.
if (!@{$cfg->{assignments}}) {
$cb->() if $cb;
return;
}
unless ($cfg->{meta}) {
weaken(my $weak = $self);
push @{$cfg->{pending_ops}}, {
cb => $cb,
run => sub { $weak->poll($cb) if $weak },
};
$self->_refresh_metadata unless $cfg->{meta_pending};
return;
}
# Group assignments by leader for multi-partition fetch
my %by_leader; # leader_id => { topic => [{partition, offset, assign_ref}] }
lib/EV/Kafka.pm view on Meta::CPAN
next unless defined $leader_id;
push @{$by_leader{$leader_id}{$a->{topic}}}, {
partition => $a->{partition},
offset => $a->{offset},
_assign => $a,
};
}
my $dispatched = 0;
my $first_err;
weaken(my $weak = $self);
for my $leader_id (keys %by_leader) {
my $conn = $self->_get_or_create_conn($leader_id);
next unless $conn && $conn->connected;
$dispatched++;
# build fetch_multi argument: {topic => [{partition, offset}]}
my %fetch_arg;
my %assign_map; # "topic:partition" => [assignment ref, requested offset]
for my $topic (keys %{$by_leader{$leader_id}}) {
for my $p (@{$by_leader{$leader_id}{$topic}}) {
lib/EV/Kafka.pm view on Meta::CPAN
auto_commit => $opts{auto_commit} // 1,
auto_offset_reset => $opts{auto_offset_reset} // 'earliest',
group_instance_id => $opts{group_instance_id},
coordinator => undef,
heartbeat_timer => undef,
state => 'init',
};
# Step 1: ensure we have metadata
unless ($cfg->{meta}) {
weaken(my $weak = $self);
push @{$cfg->{pending_ops}}, {
run => sub { $weak->_group_start if $weak },
};
$self->_refresh_metadata unless $cfg->{meta_pending};
return;
}
$self->_group_start;
}
lib/EV/Kafka.pm view on Meta::CPAN
my ($self) = @_;
my $cfg = $self->{cfg};
return if $cfg->{closed};
my $g = $cfg->{group} or return;
my $conn = $self->_any_conn;
return $self->_group_retry('no broker connection for group start')
unless $conn;
$g->{state} = 'finding';
weaken(my $weak = $self);
$conn->find_coordinator($g->{group_id}, sub {
my ($res, $err) = @_;
my $s = $weak or return;
my $cfg = $s->{cfg};
return if $cfg->{closed};
my $g = $cfg->{group} or return;
if ($err || $res->{error_code}) {
my $msg = $err || "FindCoordinator error: $res->{error_code}";
$cfg->{on_error}->($msg) if $cfg->{on_error};
# retry after delay
lib/EV/Kafka.pm view on Meta::CPAN
my ($self, $msg) = @_;
my $cfg = $self->{cfg};
return if $cfg->{closed};
my $g = $cfg->{group} or return;
if (++$g->{_start_retries} > 5) {
$g->{_start_retries} = 0;
$g->{state} = 'stopped';
$cfg->{on_error}->($msg) if $cfg->{on_error};
return;
}
weaken(my $weak = $self);
my $t; $t = EV::timer 1, 0, sub { undef $t; $weak->_group_start if $weak };
}
sub _group_join {
my ($self) = @_;
my $cfg = $self->{cfg};
return if $cfg->{closed};
my $g = $cfg->{group} or return;
my $coord = $g->{coordinator} or return;
$g->{state} = 'joining';
weaken(my $weak = $self);
$coord->join_group(
$g->{group_id}, $g->{member_id},
$g->{topics}, sub {
my ($res, $err) = @_;
my $s = $weak or return;
my $cfg = $s->{cfg};
return if $cfg->{closed};
my $g = $cfg->{group} or return;
if ($err) {
$s->_group_retry("JoinGroup: $err");
lib/EV/Kafka.pm view on Meta::CPAN
}
sub _group_sync {
my ($self, $assignments) = @_;
my $cfg = $self->{cfg};
return if $cfg->{closed};
my $g = $cfg->{group} or return;
my $coord = $g->{coordinator} or return;
$g->{state} = 'syncing';
weaken(my $weak = $self);
my $sync_cb = sub {
my ($res, $err) = @_;
my $s = $weak or return;
my $cfg = $s->{cfg};
return if $cfg->{closed};
my $g = $cfg->{group} or return;
if ($err) {
$s->_group_retry("SyncGroup: $err");
return;
}
lib/EV/Kafka.pm view on Meta::CPAN
# Build topics array for offset_fetch
my %by_topic;
for my $a (@$assignments) {
push @{$by_topic{$a->{topic}}}, $a->{partition};
}
my @topics;
for my $t (sort keys %by_topic) {
push @topics, { topic => $t, partitions => $by_topic{$t} };
}
weaken(my $weak = $self);
$coord->offset_fetch($g->{group_id}, \@topics, sub {
my ($res, $err) = @_;
my $s = $weak or return;
my $cfg = $s->{cfg};
return if $cfg->{closed};
if (!$err && $res && ref $res->{topics} eq 'ARRAY') {
for my $t (@{$res->{topics}}) {
for my $p (@{$t->{partitions} // []}) {
next if $p->{error_code};
next if $p->{offset} < 0; # no committed offset
lib/EV/Kafka.pm view on Meta::CPAN
}
sub _start_heartbeat {
my ($self) = @_;
my $cfg = $self->{cfg};
return if $cfg->{closed};
my $g = $cfg->{group} or return;
# Weak $self only; group state is re-derived per tick, so a stale
# group after re-subscribe is never acted on.
weaken(my $weak = $self);
$g->{heartbeat_timer} = EV::timer $g->{heartbeat_interval}, $g->{heartbeat_interval}, sub {
my $s = $weak or return;
my $cfg = $s->{cfg};
return if $cfg->{closed};
my $g = $cfg->{group} or return;
return unless $g->{state} eq 'stable';
my $coord = $g->{coordinator};
return unless $coord && $coord->connected;
$coord->heartbeat($g->{group_id}, $g->{generation}, $g->{member_id}, sub {
lib/EV/Kafka.pm view on Meta::CPAN
sub _start_fetch_loop {
my ($self) = @_;
my $cfg = $self->{cfg};
return if $cfg->{closed} || $cfg->{fetch_active};
$cfg->{fetch_active} = 1;
# Generation-tagged in-flight flag: a completion from a previous
# loop instance must not clear THIS loop's flag.
my $gen = ++$cfg->{_fetch_gen};
$cfg->{_fetch_in_flight} = 0;
weaken(my $weak = $self);
$cfg->{fetch_timer} = EV::timer 0, 0.1, sub {
return unless $weak && $cfg->{fetch_active};
# skip ticks while a prior poll round is still in flight
return if $cfg->{_fetch_in_flight};
$cfg->{_fetch_in_flight} = $gen;
$weak->poll(sub {
$cfg->{_fetch_in_flight} = 0
if ($cfg->{_fetch_in_flight} // 0) == $gen;
});
};
src/ppport.h view on Meta::CPAN
sv_resetpvn|5.017005||Viu
SvRMAGICAL|5.003007||Viu
SvRMAGICAL_off|5.003007||Viu
SvRMAGICAL_on|5.003007||Viu
SvROK|5.003007|5.003007|
SvROK_off|5.003007|5.003007|
SvROK_on|5.003007|5.003007|
SvRV|5.003007|5.003007|
SvRV_const|5.010001||Viu
SvRV_set|5.009003|5.003007|p
sv_rvunweaken|5.027004|5.027004|
sv_rvweaken|5.006000|5.006000|
SvRVx|5.003007||Viu
SvRX|5.009005|5.003007|p
SvRXOK|5.009005|5.003007|p
SV_SAVED_COPY|5.009005||Viu
SvSCREAM|5.003007||Viu
SvSCREAM_off|5.003007||Viu
SvSCREAM_on|5.003007||Viu
sv_setbool|5.035004|5.035004|
sv_setbool_mg|5.035004|5.035004|
sv_setgid|5.019001||Viu
t/22_lifecycle.t view on Meta::CPAN
use strict;
use warnings;
use Test::More;
use IO::Socket::INET;
use Errno qw(EAGAIN EWOULDBLOCK);
use Scalar::Util qw(weaken);
use B;
use EV;
use EV::Kafka;
# Regression tests for connection object lifecycle and memory ownership:
# M1 response mortals are reclaimed per event, not when EV::run returns
# M2 DESTROY fires every pending callback (not just the first)
# M5 a custom EV::Loop is refcounted by the conn
# M6 explicit DESTROY inerts the blessed ref; stale use croaks
# M7 conn_emit_error warns (not croaks) when no on_error is installed
t/22_lifecycle.t view on Meta::CPAN
push @{ $ctx->{timers} },
EV::timer(0.2, 0, sub { $send->($corr, $meta_body) });
}
}
});
my $conn = EV::Kafka::Conn::_new('EV::Kafka::Conn', undef);
my $saved;
$conn->on_error(sub { diag "M1 error: $_[0]"; EV::break });
$conn->on_connect(sub {
$conn->metadata(undef, sub { $saved = $_[0]; weaken($saved); });
$conn->metadata(undef, sub {
ok !defined($saved),
'M1: response mortal is reclaimed when its event ends';
EV::break;
});
});
$conn->connect('127.0.0.1', $port, 5.0);
my $t = timeout_w();
EV::run;
}
t/22_lifecycle.t view on Meta::CPAN
$conn->on_connect(sub { $ready = 1; EV::break });
$conn->connect('127.0.0.1', $port, 5.0);
my $t = timeout_w();
EV::run;
ok $ready, 'M2: handshake complete';
my (@fired, @refs);
for (1..3) {
my $cb = sub { push @fired, $_[1] };
push @refs, $cb;
weaken($refs[-1]);
$conn->metadata(undef, $cb);
}
is $conn->pending, 3, 'M2: three requests in flight';
$conn->DESTROY;
is scalar(@fired), 3, 'M2: DESTROY fires every pending callback';
is_deeply \@fired, [('destroyed') x 3], 'M2: all report "destroyed"';
is scalar(grep { defined } @refs), 0, 'M2: callback CVs are released';
}
t/26_client_lifetime.t view on Meta::CPAN
use strict;
use warnings;
use Test::More;
use FindBin;
use lib "$FindBin::Bin/lib";
use EVKafkaTest;
use Scalar::Util 'weaken';
use EV;
use EV::Kafka;
# Regression tests for the client-lifetime contract (CLIENT LIFETIME in
# the EV::Kafka POD):
# L1 connect -> produce -> flush works against the mock broker
# L2 a connected client is collectible when dropped (no cycles)
# L3 pre-metadata produce no longer pins the client (pending_ops
# cycle); dropping fails the queued callback with an error
# L4 dropping mid-connect fails the connect AND queued produce
t/26_client_lifetime.t view on Meta::CPAN
# --- L1/L2: full lifecycle, then collect --------------------------------
{
my ($port, $broker) = mock_broker(topics => { t1 => 1 });
my (@errs, @ccb, @pcb);
my $flushed = 0;
my $k = EV::Kafka->new(
brokers => "127.0.0.1:$port",
on_error => sub { push @errs, $_[0] },
);
weaken(my $weak = $k);
$k->connect(sub {
my ($meta, $err) = @_;
push @ccb, [$meta, $err];
$k->produce('t1', 'k', 'v', sub { push @pcb, [@_] });
$k->flush(sub { $flushed = 1; EV::break });
});
my $t = timeout_w();
EV::run;
ok $ccb[0] && $ccb[0][0]{brokers} && !$ccb[0][1],
t/26_client_lifetime.t view on Meta::CPAN
ok !defined($weak), 'L2: dropped client is collected (no cycles)';
}
# --- L3: pending_ops cycle, no broker at all ----------------------------
{
my @pcb;
my $k = EV::Kafka->new(
brokers => '127.0.0.1:1', # nothing listening; never connect
on_error => sub { },
);
weaken(my $weak = $k);
$k->produce('t1', 'k', 'v', sub { push @pcb, [@_] });
undef $k;
ok !defined($weak), 'L3: pre-metadata produce does not pin the client';
is scalar(@pcb), 1, 'L3: queued produce callback fired at teardown';
like $pcb[0][1] // '', qr/client closed/,
'L3: queued produce callback got the teardown error';
}
# --- L4: dropped mid-connect ---------------------------------------------
{
t/26_client_lifetime.t view on Meta::CPAN
});
return undef;
} },
);
my (@ccb, @pcb, @warns);
local $SIG{__WARN__} = sub { push @warns, $_[0] };
my $k = EV::Kafka->new(
brokers => "127.0.0.1:$port",
on_error => sub { },
);
weaken(my $weak = $k);
$k->connect(sub { push @ccb, [@_] });
$k->produce('t1', 'k', 'v', sub { push @pcb, [@_] });
# run until the metadata request is in flight, then drop the client
my $wait; $wait = EV::timer 0, 0.005, sub {
return unless grep { $_->{api} == 3 } @{$broker->{requests}};
undef $wait;
EV::break;
};
my $t = timeout_w();
t/26_client_lifetime.t view on Meta::CPAN
) or BAIL_OUT "cannot bind: $!";
my $dead_port = $dead->sockport;
close $dead;
my ($port, $broker) = mock_broker(topics => { t1 => 1 });
my @ccb;
my $k = EV::Kafka->new(
brokers => "127.0.0.1:$dead_port,127.0.0.1:$port",
on_error => sub { },
);
weaken(my $weak = $k);
$k->connect(sub { push @ccb, [@_]; EV::break });
my $t = timeout_w();
EV::run;
ok $ccb[0] && $ccb[0][0]{brokers}, 'L5: failover connect succeeded';
is scalar(@ccb), 1, 'L5: connect callback fired exactly once';
undef $k;
ok !defined($weak), 'L5: client collectible after failover connect';
}
# --- L6: armed metadata timer does not pin the client --------------------
{
my ($port, $broker) = mock_broker(topics => { t1 => 1 });
my $k = EV::Kafka->new(
brokers => "127.0.0.1:$port",
metadata_refresh => 0.1,
on_error => sub { },
);
weaken(my $weak = $k);
$k->connect(sub { });
my $ticks = 0;
my $wait; $wait = EV::timer 0, 0.02, sub {
my $n = grep { $_->{api} == 3 } @{$broker->{requests}};
if ($n >= 2) { $ticks = $n; undef $wait; EV::break }
};
my $t = timeout_w();
EV::run;
cmp_ok $ticks, '>=', 2, 'L6: periodic metadata refresh ticked';
undef $k;