NAME
EV::Etcd - Async etcd v3 client using native gRPC and EV/libev
SYNOPSIS
use v5.10;
use EV;
use EV::Etcd;
my $client = EV::Etcd->new(
endpoints => ['127.0.0.1:2379'],
);
$client->put('/my/key', 'value', sub {
my ($resp, $err) = @_;
die $err->{message} if $err;
say "Put succeeded, revision: $resp->{header}{revision}";
});
$client->get('/my/key', sub {
my ($resp, $err) = @_;
die $err->{message} if $err;
say "Value: $resp->{kvs}[0]{value}";
});
$client->watch('/my/key', sub {
my ($resp, $err) = @_;
return warn "Watch error: $err->{message}\n" if $err;
for my $event (@{$resp->{events}}) {
say "Event: $event->{type} on $event->{kv}{key}";
}
});
EV::run;
DESCRIPTION
An asynchronous etcd v3 client on the gRPC Core C API and the EV event loop: a gRPC thread waits for completions and wakes the loop, and callbacks run in the Perl thread. It needs etcd 3.4 or later (auth_status needs 3.5).
Every method takes a callback last; methods with options take them as a hash reference just before it. Invalid arguments croak. The callback receives ($response, $error): the response hash and undef, or undef and an error hash (see "ERRORS"). Responses carry a header with cluster_id, member_id, revision and raft_term unless noted otherwise. A kv hash has key, value, create_revision, mod_revision, version and lease.
CONSTRUCTOR
new
my $client = EV::Etcd->new(%options);
- endpoints => [ 'host:port', ... ]
-
Default
['127.0.0.1:2379']; an empty list croaks. An endpoint may carry thehttp://orhttps://scheme etcd prints;https://turns ontls.The client uses one endpoint at a time and moves to the next when it cannot be reached: a unary call fails with UNAVAILABLE, or with DEADLINE_EXCEEDED before a connection was made, a stream has to reconnect because its connection is down, or a keepalive ping goes unanswered. The failing call still reports its error; retrying it reaches the next endpoint. Several failures on one endpoint move the client once.
A member that has lost its leader fails linearizable reads and writes after etcd's request timeout (7 seconds by default), and these failures move the client on too, though a shorter
timeoutends the call on the client first and does not. Errors a working member returns during an election, such asetcdserver: leader changed, do not move it.Streams need a leader as well: a member without one refuses new streams and, after a few election timeouts, ends running watches and keepalives with UNAVAILABLE
etcdserver: no leader. Streams retry such answers a second apart until a leader is back, on the next endpoint when there are several; each restarts themax_retriescount. Once a unary call, or a stream that was running, reportsetcdserver: no leader, pendinglockandelection_campaigncalls sent over the same connection fail with UNAVAILABLE, as they do when a connection that never became ready is abandoned.An endpoint may also be a gRPC target listing several addresses, such as
ipv4:10.0.0.1:2379,10.0.0.2:2379. gRPC moves between them within one connection, so a refused address fails no call, but never away from a member that has lost its leader: list members as separate endpoints to leave one.gRPC connects through a proxy named in
grpc_proxy,https_proxyorhttp_proxy, even to a loopback address; list etcd's hosts inno_proxyto connect directly. - timeout => $seconds
-
RPC timeout in whole seconds, at least 1; default 30.
lockandelection_campaignwait without one, as a timeout could not tell whether they succeeded. - max_retries => $count
-
Reconnection attempts for a stream (
watch,lease_keepalive,election_observe) after a connection failure; default 30, 0 disables reconnecting. Attempts are 0.5 seconds further apart each time, up to 5 seconds, so the default gives up about two minutes after the endpoints start refusing connections, plus gRPC's 20-second connect timeout for one that accepts them but stays silent. Lower it to learn sooner that the endpoints are gone. - keepalive_time => $seconds
- keepalive_timeout => $seconds
-
Seconds between pings while calls or streams are open (default 10, fractions allowed, 0 disables them), and how long a ping may go unanswered (default 10). An unanswered ping closes the connection: its calls fail with UNAVAILABLE and the client moves to the next endpoint. Without pings, an endpoint that goes silent (a hung host, a partition) goes unnoticed until the operating system gives up on the connection. etcd closes connections that ping more often than its
--grpc-keepalive-min-time(5 seconds by default). - health_interval => $seconds
- on_health_change => sub { my ($healthy, $endpoint) = @_; ... }
-
Every
health_intervalseconds (fractions allowed; default 0, off) the client checks its connection state without sending a request, and callson_health_changewhen health changes. Only a failed connection attempt counts as unhealthy; with several endpoints it also moves the client on. Switching endpoints after a failed call is not reported. - auth_token => $token
-
A token from an earlier
authenticate, used without authenticating again. - tls => $bool
-
Connect with TLS; any
tls_*option or anhttps://endpoint turns it on too. Withouttls_ca_file, gRPC uses its default roots, or the PEM bundle named byGRPC_DEFAULT_SSL_ROOTS_FILE_PATH. - tls_ca_file => $path
-
PEM CA certificates to verify the servers with (etcd's
--trusted-ca-file), in place of the default roots. Verification cannot be turned off; for a self-signed server, give its certificate here. - tls_cert_file => $path
- tls_key_file => $path
-
Client certificate and key, given together, for servers started with
--client-cert-auth. - tls_server_name => $name
-
Name to verify the server certificate against and send as SNI, in place of the endpoint's host.
my $client = EV::Etcd->new(
endpoints => ['https://10.0.0.1:2379', 'https://10.0.0.2:2379'],
tls_ca_file => '/etc/etcd/ca.crt',
tls_cert_file => '/etc/etcd/client.crt',
tls_key_file => '/etc/etcd/client.key',
);
ERRORS
{
code => 14, # gRPC status code
status => 'UNAVAILABLE',
message => 'Connection refused',
source => 'range', # the call that failed
retryable => 1,
}
source is the method name, except range for get, lease_ttl for lease_time_to_live, keepalive for lease_keepalive, campaign, proclaim, leader, resign and observe for the election_* calls, and internal for a unary response the client could not read.
retryable is set for UNAVAILABLE, ABORTED, DEADLINE_EXCEEDED and RESOURCE_EXHAUSTED etcdserver: too many requests, but never for a failed lock or election_campaign. Other RESOURCE_EXHAUSTED errors last until someone intervenes, such as etcdserver: mvcc: database space exceeded (see alarm). Retryable means transient, not without effect: a write can be applied after its deadline has passed, so retrying lease_grant, lease_revoke or delete can apply it twice or report what the first attempt did.
Unary calls are not retried. Streams reconnect on their own (see max_retries), to the next endpoint once the current one has failed, and the count restarts once a stream is working again. A stream reports an error when reconnecting is disabled or exhausted, at once for a status a new connection cannot fix (UNAUTHENTICATED, PERMISSION_DENIED, INVALID_ARGUMENT, NOT_FOUND, ALREADY_EXISTS, FAILED_PRECONDITION, OUT_OF_RANGE, UNIMPLEMENTED), and for an error the server sends on it, such as a cancelled or compacted watch or an expired lease, which ends it.
The server can still clean up after a failed lock or election_campaign has reported, deleting the key it made on the lease. Retry with a fresh lease, and revoke the old one (which deletes all its keys) or let it expire, as a stale candidate on it can block the new attempt.
CALLBACKS AND HANDLES
A callback that dies stops neither the loop nor the client: the exception goes to $EV::DIED, a warning by default.
Callbacks live in C structures Perl cannot see, so a closure that captures its client or stream handle keeps it alive until the stream is cancelled or the client destroyed; capture a weakened copy instead. Dropping a handle does not cancel its stream. Destroying a client cancels its calls and streams without calling their callbacks. A client keeps EV::run running while it exists: destroy it or call EV::break to leave the loop. Only EV's default loop runs the callbacks, never a loop from EV::Loop->new.
use Scalar::Util 'weaken';
weaken(my $weak = $client);
my $watch = $client->watch('/jobs', sub {
my ($resp, $err) = @_;
$weak->put('/seen', 1, sub {}) if $weak && !$err;
});
cancel
$handle->cancel($callback);
watch, lease_keepalive and election_observe return a handle (EV::Etcd::Watch, EV::Etcd::Keepalive, EV::Etcd::Observe) whose cancel ends the stream. The callback runs before cancel returns, with an empty hash as the response, and the stream delivers nothing after it. Cancelling again is safe and does the same.
ENCODING
etcd stores keys and values as bytes and the client does no encoding: a string with the UTF-8 flag is stored as its UTF-8 bytes, and responses hold byte strings. Use encode_utf8 and decode_utf8 from Encode at the boundary for character data.
A key or value over 1 MiB croaks before anything is sent. etcd refuses a request over its --max-request-bytes (1.5 MiB by default) with INVALID_ARGUMENT etcdserver: request is too large, and gRPC one 512 KiB larger still with RESOURCE_EXHAUSTED.
KEY-VALUE
put
$client->put($key, $value, [\%opts,] $callback);
- lease => $lease_id
-
Attach the key to a lease.
- prev_kv => $bool
-
Return the previous kv as
prev_kv. - ignore_value => $bool
- ignore_lease => $bool
-
Keep the current value, or lease, and update only the other.
get
$client->get($key, [\%opts,] $callback);
Response keys: kvs (array of kv hashes), count (all keys matched, even beyond limit) and more (true when limit cut the result).
- prefix => $bool
-
Every key with
$keyas prefix; every key at all for an empty one. - range_end => $end
-
Keys from
$keyup to$end, exclusive. - limit => $n
-
At most
$nkeys. - revision => $rev
-
Read at an older revision.
- keys_only => $bool
- count_only => $bool
-
Return keys without values, or only
count. - sort_order => 'ascend' | 'descend'
- sort_target => 'key' | 'version' | 'create' | 'mod' | 'value'
-
Sort the result by
sort_target, insort_order. - serializable => $bool
-
Read from the member's local data: faster, possibly stale.
- min_mod_revision, max_mod_revision, min_create_revision, max_create_revision => $rev
-
Filter by modification or creation revision.
delete
$client->delete($key, [\%opts,] $callback);
Options prefix and range_end as for get (an empty prefix deletes every key), and prev_kv to return the deleted kvs. Response keys: deleted (count) and prev_kvs.
WATCH
watch
my $watch = $client->watch($key, [\%opts,] $callback);
Watch a key or range; returns a handle (see "cancel"). The callback runs for each message, whose events hold hashes with type (PUT or DELETE), kv and, with the prev_kv option, prev_kv. created is true on the first message of each stream, so again after a reconnect.
A watch the server cancels arrives as an error with status CANCELLED and source watch, whatever the cause its message names: a compaction, a permission denied or an expired token. After a compaction, $err->{compact_revision} holds the revision to resume from (0 otherwise).
- prefix => $bool
- range_end => $end
-
As for
get. - start_revision => $rev
-
Start from an older revision instead of the current one.
- prev_kv => $bool
-
Add the previous kv to each event.
- progress_notify => $bool
-
Have the server send empty messages while idle, carrying the current revision.
- watch_id => $id
-
Choose the watch ID instead of letting the server assign one.
- auto_reconnect => $bool
-
Reconnect after a connection failure, resuming from the last revision seen. Default true.
LEASE
lease_grant
$client->lease_grant($ttl, $callback);
Grant a lease for $ttl seconds. Response keys: id and ttl (as granted).
lease_revoke
$client->lease_revoke($lease_id, $callback);
Revoke a lease, deleting every key attached to it.
lease_keepalive
my $keepalive = $client->lease_keepalive($lease_id, [\%opts,] $callback);
Keep a lease refreshed over a stream; returns a handle (see "cancel"). Each refresh calls back with id and ttl. An expired lease ends the stream with a NOT_FOUND error. Option auto_reconnect (default true) reconnects after a connection failure.
lease_time_to_live
$client->lease_time_to_live($lease_id, [\%opts,] $callback);
Response keys: id, ttl (remaining seconds, -1 once expired), granted_ttl and keys, which with the keys option lists the keys attached to the lease.
lease_leases
$client->lease_leases($callback);
Response key leases: an array of hashes with an id.
LOCK
lock
$client->lock($name, $lease_id, $callback);
Acquire the lock $name, held until unlock or until the lease expires or is revoked. The call waits until the lock is free, without the client timeout; destroying the client cancels it. The response key is what unlock takes. On failure, see "ERRORS".
$client->lease_grant(30, sub {
my ($lease, $err) = @_;
die $err->{message} if $err;
$client->lock('my-resource', $lease->{id}, sub {
my ($lock, $err) = @_;
die $err->{message} if $err;
# ... protected work ...
$client->unlock($lock->{key}, sub {});
});
});
unlock
$client->unlock($key, $callback);
AUTHENTICATION
authenticate
$client->authenticate($user, $password, $callback);
On success the client keeps the token (response key token) and sends it with every later call.
Simple tokens expire after --auth-token-ttl seconds unused (300 by default), do not survive a restart, and are timed by each member separately, so one may already have expired on the member the client switches to. JWT tokens expire after the ttl of --auth-token, and calls fail with INVALID_ARGUMENT etcdserver: revision of auth store is old after any user, role or permission change. With an expired or stale token, calls and new streams fail, while running streams may continue for a while; call authenticate again and restart the streams.
auth_enable
$client->auth_enable($callback);
etcd refuses it until a root user with the root role exists.
auth_disable
$client->auth_disable($callback);
Needs root. The client drops its token. Other clients keep theirs, which etcd before 3.4.28 and 3.5.10 rejects with etcdserver: invalid auth token; authenticate on such a client fails with FAILED_PRECONDITION and drops it.
auth_status
$client->auth_status($callback);
Response keys: enabled and auth_revision.
user_add, user_delete, user_change_password, user_get, user_list
$client->user_add($user, $password, $callback);
$client->user_delete($user, $callback);
$client->user_change_password($user, $password, $callback);
$client->user_get($user, $callback); # roles => [...]
$client->user_list($callback); # users => [...]
user_grant_role, user_revoke_role
$client->user_grant_role($user, $role, $callback);
$client->user_revoke_role($user, $role, $callback);
role_add, role_delete, role_get, role_list
$client->role_add($role, $callback);
$client->role_delete($role, $callback);
$client->role_get($role, $callback); # perm => [...]
$client->role_list($callback); # roles => [...]
role_get lists permissions as hashes with perm_type (READ, WRITE or READWRITE), key and range_end.
role_grant_permission, role_revoke_permission
$client->role_grant_permission($role, $perm_type, $key, $range_end, $callback);
$client->role_revoke_permission($role, $key, $range_end, $callback);
$range_end is exclusive; undef means the single key. For a prefix, pass it with its last byte incremented: /app/ gives /app0. "\x00" covers every key from $key on, not only the prefix.
$client->role_grant_permission('app', 'READWRITE', '/app/', '/app0', sub {
my ($resp, $err) = @_;
warn $err->{message} if $err;
});
MAINTENANCE
status
$client->status($callback);
Status of the member the client is connected to. Response keys: version, db_size, db_size_in_use, leader (member ID), raft_index, raft_term, raft_applied_index, is_learner, and errors when the member has any.
compact
$client->compact($revision, [\%opts,] $callback);
Discard all revisions before $revision, irreversibly. With physical => 1 the call returns once the data is removed from the backend rather than once the compaction is committed.
alarm
$client->alarm($action, [\%opts,] $callback);
$action is GET, ACTIVATE or DEACTIVATE. Option alarm is NOSPACE or CORRUPT; the default, NONE, lists every alarm for GET and does nothing otherwise. Option member_id names the member the alarm is recorded for: pass a real one, from GET or member_list. Response key alarms: hashes with member_id, alarm (a number) and alarm_type (its name); etcd 3.4 sends no header.
# After freeing space: clear every alarm
$client->alarm('GET', sub {
my ($resp, $err) = @_;
return warn $err->{message} if $err;
$client->alarm('DEACTIVATE', {
alarm => $_->{alarm_type},
member_id => $_->{member_id},
}, sub { warn $_[1]{message} if $_[1] }) for @{$resp->{alarms}};
});
defragment
$client->defragment($callback);
Defragment the backend of the member the client is connected to. It blocks that member while running. etcd 3.4 and 3.5 send no header, so the response is empty.
hash_kv
$client->hash_kv([$revision,] $callback);
Hash of the store up to $revision (default current), to compare members. Response keys: hash and compact_revision.
move_leader
$client->move_leader($member_id, $callback);
Hand leadership to another voting member. Only the leader accepts it, so the client must be connected to the leader (status then shows leader equal to $resp->{header}{member_id}). etcd sends no header.
ELECTION
election_campaign
$client->election_campaign($name, $lease_id, $value, $callback);
Wait to become leader of $name with $value; leadership lasts as long as the lease. There is no client timeout, and destroying the client cancels the wait. The response leader is a hash (name, key, rev, lease) for election_proclaim and election_resign. On failure, see "ERRORS".
election_leader
$client->election_leader($name, $callback);
Response key kv: the leader's kv. Without a leader, an error.
election_proclaim
$client->election_proclaim($leader, $value, $callback);
Announce a new value as leader.
election_resign
$client->election_resign($leader, $callback);
Give up leadership. Both this and election_proclaim croak unless $leader has a non-empty key and positive rev and lease.
election_observe
my $observe = $client->election_observe($name, [\%opts,] $callback);
Call back with the leader's kv on every change; returns a handle (see "cancel"). Option auto_reconnect (default true) reconnects after a connection failure.
CLUSTER
member_list
$client->member_list([\%opts,] $callback);
Response key members: hashes with id, name, peer_urls, client_urls and is_learner. Option linearizable reads through the leader instead of the member's local view; etcd 3.4 ignores it.
member_add
$client->member_add(\@peer_urls, [\%opts,] $callback);
Option is_learner adds a non-voting member. Response keys: member (the new one) and members.
member_remove, member_update, member_promote
$client->member_remove($member_id, $callback);
$client->member_update($member_id, \@peer_urls, $callback);
$client->member_promote($member_id, $callback); # learner to voter
Response key members.
TRANSACTIONS
txn
$client->txn(
compare => \@compare,
success => \@success,
failure => \@failure,
callback => $callback,
);
$client->txn(\@compare, \@success, \@failure, $callback);
Run success if every comparison holds, otherwise failure, atomically; the response succeeded says which. A comparison names a key and one field, compared with result = (default), !=, < or >:
{ key => $key, value => $expected }
{ key => $key, version => $expected }
{ key => $key, create_revision => $expected }
{ key => $key, mod_revision => $expected, result => '<' }
{ key => $key, lease => $expected }
target (value, version, create, mod or lease) may name the field too, and must agree with it; alone it compares against 0 or an empty value, so { key => $key, target => 'version' }, like { key => $key, version => 0 }, means the key does not exist. Operations, which also take lease (put) and range_end (delete, range):
{ put => { key => $key, value => $value } } # or request_put
{ delete => { key => $key } } # or request_delete_range
{ range => { key => $key } } # or request_range
The response responses holds one hash per operation run, under response_put, response_delete_range (deleted, prev_kvs) or response_range (kvs, count, more).
$client->txn(
compare => [{ key => '/counter', value => '0' }],
success => [{ put => { key => '/counter', value => '1' } }],
callback => sub {
my ($resp, $err) = @_;
say $resp->{succeeded} ? 'Incremented' : 'Already changed';
},
);
CAVEATS
Fork: gRPC's threads do not survive fork(). EV::Etcd starts gRPC with the first client and shuts it down when the last is destroyed, so a child forked while the process holds no client can create its own. The first fork() after shutdown waits up to two seconds for gRPC's threads to finish. Destroying a client whose streams never ran can take seconds to settle; cancel its streams, or run the loop, first. On macOS gRPC stays up until the process exits. A child forked while gRPC is up (a client exists, it is still settling, or on macOS ever since the first client) cannot use etcd: new croaks, as does any call on an inherited client or handle. Inherited clients are inert in the child (they neither fire nor keep its loop running), and destroying one there only frees its Perl side, with a warning. In a server that forks workers, create clients in the workers. A child forked inside an EV::Etcd callback must exec or exit, not return from it.
Signals: on a threaded perl before 5.42, a %SIG handler that runs on a gRPC thread crashes perl. EV::Etcd starts gRPC with all signals blocked, which keeps them on the Perl thread in practice, but gRPC can start a thread later; prefer EV::signal watchers there.
INSTALLATION
Building needs the gRPC Core and protobuf-c C libraries and pkg-config:
apt install libgrpc-dev libgrpc++-dev libprotobuf-c-dev pkg-config
brew install grpc protobuf-c pkg-config
pkg install grpc protobuf-c pkgconf
Most tests need an etcd on 127.0.0.1:2379 and skip without one. They write under their own prefixes and remove the leases, users and roles they create; tests that compact history or add members run only with EV_ETCD_TEST_ETCD=1, for an etcd that exists for testing.
AUTHOR
vividsnow
LICENSE
This library is free software; you can redistribute it and/or modify it under the same terms as Perl itself.