diff --git a/.travis.yml b/.travis.yml index 38b665b5b..bc5c8ba79 100644 --- a/.travis.yml +++ b/.travis.yml @@ -10,6 +10,7 @@ before_install: - wget http://s3.amazonaws.com/influxdb/influxdb_latest_amd64.deb - sudo useradd influxdb - sudo dpkg -i influxdb_latest_amd64.deb + - sudo cp ./scripts/influxdb_conf.toml /opt/influxdb/shared/config.toml - travis_retry sudo service influxdb restart - sudo service influxdb status diff --git a/scripts/influxdb_conf.toml b/scripts/influxdb_conf.toml new file mode 100644 index 000000000..e0c505a7b --- /dev/null +++ b/scripts/influxdb_conf.toml @@ -0,0 +1,187 @@ +# Welcome to the InfluxDB configuration file. + +# If hostname (on the OS) doesn't return a name that can be resolved by the other +# systems in the cluster, you'll have to set the hostname to an IP or something +# that can be resolved here. +# hostname = "" + +bind-address = "0.0.0.0" + +# Once every 24 hours InfluxDB will report anonymous data to m.influxdb.com +# The data includes raft name (random 8 bytes), os, arch and version +# We don't track ip addresses of servers reporting. This is only used +# to track the number of instances running and the versions which +# is very helpful for us. +# Change this option to true to disable reporting. +reporting-disabled = false + +[logging] +# logging level can be one of "debug", "info", "warn" or "error" +level = "info" +file = "/opt/influxdb/shared/log.txt" # stdout to log to standard out + +# Configure the admin server +[admin] +port = 8083 # binding is disabled if the port isn't set +assets = "/opt/influxdb/current/admin" + +# Configure the http api +[api] +port = 8086 # binding is disabled if the port isn't set +# ssl-port = 8084 # Ssl support is enabled if you set a port and cert +# ssl-cert = /path/to/cert.pem + +# connections will timeout after this amount of time. Ensures that clients that misbehave +# and keep alive connections they don't use won't end up connection a million times. +# However, if a request is taking longer than this to complete, could be a problem. +read-timeout = "5s" + +[input_plugins] + + # Configure the graphite api + [input_plugins.graphite] + enabled = false + # port = 2003 + # database = "" # store graphite data in this database + # udp_enabled = true # enable udp interface on the same port as the tcp interface + + # Configure the udp api + #[input_plugins.udp] + #enabled = true + #port = 4444 + # database = "" + + # Configure multiple udp apis each can write to separate db. Just + # repeat the following section to enable multiple udp apis on + # different ports. + [[input_plugins.udp_servers]] # array of tables + enabled = true + port = 5551 + database = "mine" + +# Raft configuration +[raft] +# The raft port should be open between all servers in a cluster. +# However, this port shouldn't be accessible from the internet. + +port = 8090 + +# Where the raft logs are stored. The user running InfluxDB will need read/write access. +dir = "/opt/influxdb/shared/data/raft" + +# election-timeout = "1s" + +[storage] + +dir = "/opt/influxdb/shared/data/db" +# How many requests to potentially buffer in memory. If the buffer gets filled then writes +# will still be logged and once the local storage has caught up (or compacted) the writes +# will be replayed from the WAL +write-buffer-size = 10000 + +# the engine to use for new shards, old shards will continue to use the same engine +default-engine = "rocksdb" + +# The default setting on this is 0, which means unlimited. Set this to something if you want to +# limit the max number of open files. max-open-files is per shard so this * that will be max. +max-open-shards = 0 + +# The default setting is 100. This option tells how many points will be fetched from LevelDb before +# they get flushed into backend. +point-batch-size = 100 + +# The number of points to batch in memory before writing them to leveldb. Lowering this number will +# reduce the memory usage, but will result in slower writes. +write-batch-size = 5000000 + +# The server will check this often for shards that have expired that should be cleared. +retention-sweep-period = "10m" + +[storage.engines.leveldb] + +# Maximum mmap open files, this will affect the virtual memory used by +# the process +max-open-files = 1000 + +# LRU cache size, LRU is used by leveldb to store contents of the +# uncompressed sstables. You can use `m` or `g` prefix for megabytes +# and gigabytes, respectively. +lru-cache-size = "200m" + +[storage.engines.rocksdb] + +# Maximum mmap open files, this will affect the virtual memory used by +# the process +max-open-files = 1000 + +# LRU cache size, LRU is used by rocksdb to store contents of the +# uncompressed sstables. You can use `m` or `g` prefix for megabytes +# and gigabytes, respectively. +lru-cache-size = "200m" + +[storage.engines.hyperleveldb] + +# Maximum mmap open files, this will affect the virtual memory used by +# the process +max-open-files = 1000 + +# LRU cache size, LRU is used by rocksdb to store contents of the +# uncompressed sstables. You can use `m` or `g` prefix for megabytes +# and gigabytes, respectively. +lru-cache-size = "200m" + +[storage.engines.lmdb] + +map-size = "100g" + +[cluster] +# A comma separated list of servers to seed +# this server. this is only relevant when the +# server is joining a new cluster. Otherwise +# the server will use the list of known servers +# prior to shutting down. Any server can be pointed to +# as a seed. It will find the Raft leader automatically. + +# Here's an example. Note that the port on the host is the same as the raft port. +# seed-servers = ["hosta:8090","hostb:8090"] + +# Replication happens over a TCP connection with a Protobuf protocol. +# This port should be reachable between all servers in a cluster. +# However, this port shouldn't be accessible from the internet. + +protobuf_port = 8099 +protobuf_timeout = "2s" # the write timeout on the protobuf conn any duration parseable by time.ParseDuration +protobuf_heartbeat = "200ms" # the heartbeat interval between the servers. must be parseable by time.ParseDuration +protobuf_min_backoff = "1s" # the minimum backoff after a failed heartbeat attempt +protobuf_max_backoff = "10s" # the maxmimum backoff after a failed heartbeat attempt + +# How many write requests to potentially buffer in memory per server. If the buffer gets filled then writes +# will still be logged and once the server has caught up (or come back online) the writes +# will be replayed from the WAL +write-buffer-size = 1000 + +# the maximum number of responses to buffer from remote nodes, if the +# expected number of responses exceed this number then querying will +# happen sequentially and the buffer size will be limited to this +# number +max-response-buffer-size = 100 + +# When queries get distributed out to shards, they go in parallel. This means that results can get buffered +# in memory since results will come in any order, but have to be processed in the correct time order. +# Setting this higher will give better performance, but you'll need more memory. Setting this to 1 will ensure +# that you don't need to buffer in memory, but you won't get the best performance. +concurrent-shard-query-limit = 10 + +[wal] + +dir = "/opt/influxdb/shared/data/wal" +flush-after = 1000 # the number of writes after which wal will be flushed, 0 for flushing on every write +bookmark-after = 1000 # the number of writes after which a bookmark will be created + +# the number of writes after which an index entry is created pointing +# to the offset of the first request, default to 1k +index-after = 1000 + +# the number of requests per one log file, if new requests came in a +# new log file will be created +requests-per-logfile = 10000 diff --git a/src/InfluxDB/Adapter/ConnectableInterface.php b/src/InfluxDB/Adapter/ConnectableInterface.php deleted file mode 100644 index 1b7a79643..000000000 --- a/src/InfluxDB/Adapter/ConnectableInterface.php +++ /dev/null @@ -1,8 +0,0 @@ -options->getHttpSeriesEndpoint(); + try { return $this->httpClient->get($endpoint, $options)->json(); + } catch (\Exception $e) { + var_dump((string)$e->getResponse()->getBody(true)); + die(); + } } } diff --git a/src/InfluxDB/Adapter/UdpAdapter.php b/src/InfluxDB/Adapter/UdpAdapter.php index 2cb026ede..e2dcf8a51 100644 --- a/src/InfluxDB/Adapter/UdpAdapter.php +++ b/src/InfluxDB/Adapter/UdpAdapter.php @@ -3,35 +3,20 @@ namespace InfluxDB\Adapter; use InfluxDB\Options; -class UdpAdapter implements AdapterInterface, ConnectableInterface +class UdpAdapter implements AdapterInterface { private $options; - private $socket; public function __construct(Options $options) { $this->options = $options; } - public function getSocket() - { - return $this->socket; - } - - public function connect() - { - $this->socket = socket_create(AF_INET, SOCK_DGRAM, SOL_UDP); - return $this; - } - - public function disconnect() - { - socket_close($this->getSocket()); - } - public function send($message) { + $socket = socket_create(AF_INET, SOCK_DGRAM, SOL_UDP); $message = json_encode($message); - socket_sendto($this->getSocket(), $message, strlen($message), 0, $this->host, $this->port); + socket_sendto($socket, $message, strlen($message), 0, $this->options->getHost(), $this->options->getPort()); + socket_close($socket); } } diff --git a/src/InfluxDB/Client.php b/src/InfluxDB/Client.php index 501afcefd..2b9017025 100644 --- a/src/InfluxDB/Client.php +++ b/src/InfluxDB/Client.php @@ -20,26 +20,6 @@ class Client return $this->adapter; } - public function connect() - { - $result = false; - if ($this->getAdapter() instanceOf ConnectableInterface) { - $result = $this->getAdapter()->connect(); - } - - return $result; - } - - public function disconnect() - { - $result = false; - if ($this->getAdapter() instanceOf ConnectableInterface) { - $result = $this->getAdapter()->disconnect(); - } - - return $result; - } - public function mark($name, array $values) { $data =[]; diff --git a/tests/InfluxDB/ClientTest.php b/tests/InfluxDB/ClientTest.php index 566150f7d..e4d5b20a1 100644 --- a/tests/InfluxDB/ClientTest.php +++ b/tests/InfluxDB/ClientTest.php @@ -3,11 +3,13 @@ namespace InfluxDB; use InfluxDB\Adapter\GuzzleAdapter as InfluxHttpAdapter; use InfluxDB\Options; +use InfluxDB\Adapter\UdpAdapter; use GuzzleHttp\Client as GuzzleHttpClient; use crodas\InfluxPHP\Client as Crodas; class ClientTest extends \PHPUnit_Framework_TestCase { + private $rawOptions; private $object; private $options; @@ -16,6 +18,7 @@ class ClientTest extends \PHPUnit_Framework_TestCase public function setUp() { $options = include __DIR__ . '/../bootstrap.php'; + $this->rawOptions = $options; $client = new Crodas( $options["tcp"]["host"], @@ -30,6 +33,13 @@ class ClientTest extends \PHPUnit_Framework_TestCase } $client->createDatabase($options["tcp"]["database"]); + try { + $client->deleteDatabase($options["udp"]["database"]); + } catch (\Exception $e) { + // nothing... + } + $client->createDatabase($options["udp"]["database"]); + $this->anotherClient = $client; $tcpOptions = $options["tcp"]; @@ -51,6 +61,9 @@ class ClientTest extends \PHPUnit_Framework_TestCase $this->object = $influx; } + /** + * @group tcp + */ public function testGuzzleHttpApiWorksCorrectly() { $this->object->mark("tcp.test", ["mark" => "element"]); @@ -60,6 +73,9 @@ class ClientTest extends \PHPUnit_Framework_TestCase $this->assertEquals("element", $body[0]["points"][0][2]); } + /** + * @group tcp + */ public function testGuzzleHttpQueryApiWorksCorrectly() { $this->object->mark("tcp.test", ["mark" => "element"]); @@ -71,6 +87,9 @@ class ClientTest extends \PHPUnit_Framework_TestCase $this->assertEquals("element", $body[0]["points"][0][2]); } + /** + * @group tcp + */ public function testGuzzleHttpQueryApiWithMultipleData() { $this->object->mark("tcp.test", ["mark" => "element"]); @@ -83,6 +102,9 @@ class ClientTest extends \PHPUnit_Framework_TestCase $this->assertEquals("tcp.test", $body[0]["name"]); } + /** + * @group tcp + */ public function testGuzzleHttpQueryApiWithTimePrecision() { $this->object->mark("tcp.test", ["mark" => "element"]); @@ -92,4 +114,33 @@ class ClientTest extends \PHPUnit_Framework_TestCase $this->assertCount(1, $body[0]["points"]); $this->assertEquals("tcp.test", $body[0]["name"]); } + + /** + * @group udp + */ + public function testUdpIpWriteData() + { + $rawOptions = $this->rawOptions; + $options = new Options(); + $options->setUsername($rawOptions["udp"]["username"]); + $options->setPassword($rawOptions["udp"]["password"]); + $options->setPort($rawOptions["udp"]["port"]); + + $adapter = new UdpAdapter($options); + $object = new Client(); + $object->setAdapter($adapter); + + $object->mark("udp.test", ["mark" => "element"]); + $object->mark("udp.test", ["mark" => "element1"]); + $object->mark("udp.test", ["mark" => "element2"]); + $object->mark("udp.test", ["mark" => "element3"]); + + // Wait UDP/IP message arrives + usleep(200e3); + + $body = $this->object->query("select * from udp.test"); + + $this->assertCount(4, $body[0]["points"]); + $this->assertEquals("udp.test", $body[0]["name"]); + } } diff --git a/tests/bootstrap.php b/tests/bootstrap.php index cf38911f3..9dd58c2cc 100644 --- a/tests/bootstrap.php +++ b/tests/bootstrap.php @@ -7,5 +7,12 @@ return [ "database" => "mine", "username" => "root", "password" => "root", - ] + ], + "udp" => [ + "host" => "localhost", + "port" => 5551, + "database" => "mine", + "username" => "root", + "password" => "root" + ], ];