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 the "http://" or "https://" scheme etcd prints; "https://" turns on "tls". 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 "timeout" ends the call on the client first and does not. Errors a working member returns during an election, such as "etcdserver: 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 the "max_retries" count. Once a unary call, or a stream that was running, reports "etcdserver: no leader", pending "lock" and "election_campaign" calls 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_proxy" or "http_proxy", even to a loopback address; list etcd's hosts in "no_proxy" to connect directly. timeout => $seconds RPC timeout in whole seconds, at least 1; default 30. "lock" and "election_campaign" wait 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_interval" seconds (fractions allowed; default 0, off) the client checks its connection state without sending a request, and calls "on_health_change" when 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 an "https://" endpoint turns it on too. Without "tls_ca_file", gRPC uses its default roots, or the PEM bundle named by "GRPC_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 $key as prefix; every key at all for an empty one. range_end => $end Keys from $key up to $end, exclusive. limit => $n At most $n keys. 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", in "sort_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.