AnyEvent-InfluxDB
view release on metacpan or search on metacpan
lib/AnyEvent/InfluxDB.pm view on Meta::CPAN
#ABSTRACT: An asynchronous library for InfluxDB time-series database
use strict;
use warnings;
package AnyEvent::InfluxDB;
our $AUTHORITY = 'cpan:AJGB';
$AnyEvent::InfluxDB::VERSION = '1.0.2.0';
use AnyEvent;
use AnyEvent::HTTP;
use URI;
use URI::QueryParam;
use JSON qw(decode_json);
use List::MoreUtils qw(zip);
use URI::Encode::XS qw( uri_encode );
use Moo;
has [qw( ssl_options username password jwt on_request )] => (
is => 'ro',
predicate => 1,
);
has 'server' => (
is => 'rw',
default => 'http://localhost:8086',
);
has '_is_ssl' => (
is => 'lazy',
);
has '_tls_ctx' => (
is => 'lazy',
);
has '_server_uri' => (
is => 'lazy',
);
sub _build__tls_ctx {
my ($self) = @_;
# no ca/hostname checks
return 'low' unless $self->has_ssl_options;
# create ctx
require AnyEvent::TLS;
return AnyEvent::TLS->new( %{ $self->ssl_options } );
}
sub _build__is_ssl {
my ($self) = @_;
return $self->server =~ /^https/;
}
sub _build__server_uri {
my ($self) = @_;
my $url = URI->new( $self->server, 'http' );
if ( $self->has_username && $self->has_password ) {
$url->query_param( 'u' => $self->username );
$url->query_param( 'p' => $self->password );
}
return $url;
}
sub _make_url {
my ($self, $path, $params) = @_;
my $url = $self->_server_uri->clone;
$url->path($path);
while ( my ($k, $v) = each %$params ) {
$url->query_param( $k => $v );
}
return $url;
}
sub _http_request {
my $cb = pop;
my ($self, $method, $url, $post_data) = @_;
if ($self->has_on_request) {
$self->on_request->($method, $url, $post_data);
}
my %args = (
headers => {
referer => undef,
'user-agent' => "AnyEvent-InfluxDB/0.13",
}
);
if ($self->has_jwt) {
$args{headers}->{Authorization} = 'Bearer '. $self->jwt;
}
if ( $method eq 'POST' ) {
if ( defined $post_data ) {
$args{'body'} = $post_data;
} else {
if ( my $q = $url->query_param_delete('q') ) {
$args{headers}{'content-type'} = 'application/x-www-form-urlencoded';
$args{body} = 'q='. uri_encode($q);
}
}
}
if ( $self->_is_ssl ) {
$args{tls_ctx} = $self->_tls_ctx;
}
my $guard;
$guard = http_request
$method => $url->as_string,
%args,
sub {
$cb->(@_);
undef $guard;
};
};
sub ping {
my ($self, %args) = @_;
my $url = $self->_make_url('/ping', {
(
exists $args{wait_for_leader} ?
( wait_for_leader => $args{wait_for_leader} )
:
()
)
});
$self->_http_request( GET => $url,
sub {
my ($body, $headers) = @_;
if ( $headers->{Status} eq '204' ) {
$args{on_success}->( $headers->{'x-influxdb-version'} );
} else {
$args{on_error}->( $headers->{Reason} || $body );
}
}
);
}
sub _to_line {
my $data = shift;
my $t = $data->{tags} || {};
my $f = $data->{fields} || {};
return $data->{measurement}
.(
$t ?
','.
join(',',
map {
join('=', $_, $t->{$_})
} sort { $a cmp $b } keys %$t
)
:
''
)
. ' '
.(
join(',',
lib/AnyEvent/InfluxDB.pm view on Meta::CPAN
$db->write(
database => 'mydb',
data => [
map {
+{
measurement => $measurement,
tags => {
label => $_->label,
},
fields => {
value => $_->value,
uom => '"'. $_->uom .'"',
},
}
} @perfdata
],
on_success => sub { print "$line written\n"; },
on_error => sub { print "$line error: @_\n"; },
);
$hdl->on_drain(
sub {
$hdl->fh->close;
undef $hdl;
}
);
},
);
};
EV::run;
=head1 DESCRIPTION
Asynchronous client library for InfluxDB time-series database L<https://influxdb.com>.
This version is meant to be used with InfluxDB v1.0.0 or newer.
=head1 METHODS
=head2 new
my $db = AnyEvent::InfluxDB->new(
server => 'http://localhost:8086',
# authenticate using Basic credentials
username => 'admin',
password => 'password',
# or use JWT token
jwt => 'JWT_TOKEN_BLOB'
);
Returns object representing given InfluDB C<server> connected using optionally
provided username C<username> and password C<password>.
Default value of C<server> is C<http://localhost:8086>.
If the server protocol is C<https> then by default no validation of remote
host certificate is performed. This can be changed by setting C<ssl_options>
parameter with any options accepted by L<AnyEvent::TLS>.
my $db = AnyEvent::InfluxDB->new(
...
ssl_options => {
verify => 1,
verify_peername => 'https',
ca_file => '/path/to/cacert.pem',
}
);
As an debugging aid the C<on_request> code reference may also be provided. It will
be executed before each request with the method name, url and POST data if set.
my $db = AnyEvent::InfluxDB->new(
...
on_request => sub {
my ($method, $url, $post_data) = @_;
print "$method $url\n";
print "$post_data\n" if $post_data;
}
);
=for Pod::Coverage has_jwt jwt has_on_request has_password has_ssl_options has_username on_request password server ssl_options username
=head2 ping
$cv = AE::cv;
$db->ping(
wait_for_leader => 2,
on_success => $cv,
on_error => sub {
$cv->croak("Failed to ping cluster leader: @_");
}
);
my $version = $cv->recv;
Checks the leader of the cluster to ensure that the leader is available and ready.
The optional parameter C<wait_for_leader> specifies the number of seconds to wait
before returning a response.
The required C<on_success> code reference is executed if request was successful
with the value of C<X-Influxdb-Version> response header as argument,
otherwise executes the required C<on_error> code reference with the value of
C<Reason> response header as argument.
=head2 Managing Data
=head3 write
$cv = AE::cv;
$db->write(
database => 'mydb',
precision => 's',
rp => 'last_day',
consistency => 'quorum',
data => [
# line protocol formatted
'cpu_load,host=server02,region=eu-east sensor="top",value=0.64 1456097956',
# or as a hash
{
measurement => 'cpu_load',
tags => {
host => 'server02',
region => 'eu-east',
},
fields => {
value => '0.64',
sensor => q{"top"},
},
time => time()
}
],
on_success => $cv,
on_error => sub {
$cv->croak("Failed to write data: @_");
}
);
$cv->recv;
( run in 0.874 second using v1.01-cache-2.11-cpan-acf6aa7dc9e )