Repository navigation
Expand file tree
/
Copy pathQueryApi.php
More file actions
120 lines (96 loc) · 2.8 KB
/
Copy pathQueryApi.php
File metadata and controls
120 lines (96 loc) · 2.8 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
<?php
namespace InfluxDB2;
use InfluxDB2\Model\Dialect;
use InfluxDB2\Model\Query;
use Psr\Http\Message\ResponseInterface;
class QueryApi extends DefaultApi
{
private $DEFAULT_DIALECT;
/**
* QueryApi constructor.
* @param array $options
*/
public function __construct(array $options)
{
parent::__construct($options);
$this->DEFAULT_DIALECT = new Dialect([
'header' => true,
'delimiter' => ',',
'comment_prefix' => '#',
'annotations' => ['datatype', 'group', 'default']
]);
}
/**
*
* @param string $query query the flux query to execute. The data could be represent by string, Query
* @param null $org specifies the source organization
* @param null $dialect csv dialect
* @return string
*/
public function queryRaw(string $query, $org = null, $dialect = null): ?string
{
$result = $this->postQuery($query, $org, $dialect ?: $this->DEFAULT_DIALECT);
if ($result == null)
{
return null;
}
return $result->getBody()->getContents();
}
/**
* @param $query
* @param $org
* @param $dialect
* @return FluxTable[]
*/
public function query($query, $org = null, $dialect = null)
{
$response = $this->postQuery($query, $org, $dialect ?: $this->DEFAULT_DIALECT);
if ($response == null) {
return null;
}
$parser = new FluxCsvParser($response->getBody());
$parser->parse();
return $parser->tables;
}
/**
* @param $query
* @param $org
* @param $dialect
* @return FluxCsvParser generator
*/
public function queryStream($query, $org = null, $dialect = null): ?FluxCsvParser
{
$response = $this->postQuery($query, $org, $dialect ?: $this->DEFAULT_DIALECT);
if ($response == null)
{
return null;
}
return new FluxCsvParser($response->getBody(), true);
}
private function postQuery($query, $org, $dialect): ?ResponseInterface
{
$orgParam = $org ?: $this->options["org"];
$this->check("org", $orgParam);
$payload = $this->generatePayload($query, $dialect);
$queryParams = ["org" => $orgParam];
if ($payload == null)
{
return null;
}
return $this->post($payload->__toString(), "/api/v2/query", $queryParams);
}
private function generatePayload($query, $dialect)
{
if ((!isset($query) || trim($query) === '')) {
return null;
}
if ($query instanceof Query) {
return $query;
}
return new Query([
'query' => $query,
'dialect' => $dialect,
'type' => null
]);
}
}