Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -165,7 +165,12 @@ impl Client {
(url, p)
};

params.push(("precision", opts.precision.as_str().to_string()));
let precision_str = if opts.use_v2_api {
opts.precision.as_v2_str()
} else {
opts.precision.as_str()
};
params.push(("precision", precision_str.to_string()));

// Compress once; each attempt re-sends the same (Arc-backed) Bytes.
let (final_body, compressed) = maybe_gzip(body, opts.gzip_threshold).await?;
Expand Down
25 changes: 24 additions & 1 deletion src/precision.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ pub enum Precision {
}

impl Precision {
/// Returns the API query-parameter string for this precision.
/// Returns the v3 API query-parameter string for this precision.
pub fn as_str(self) -> &'static str {
match self {
Precision::Nanosecond => "nanosecond",
Expand All @@ -28,6 +28,16 @@ impl Precision {
}
}

/// Returns the v2 API query-parameter string for this precision.
pub fn as_v2_str(self) -> &'static str {
match self {
Precision::Nanosecond => "ns",
Precision::Microsecond => "us",
Precision::Millisecond => "ms",
Precision::Second => "s",
}
}

/// Number of nanoseconds in one unit of this precision.
pub(crate) fn nanos_per_unit(self) -> i64 {
match self {
Expand Down Expand Up @@ -80,4 +90,17 @@ mod tests {
assert_eq!(p.scale_timestamp(ns), scaled);
}
}

#[test]
fn v2_str_roundtrip() {
for (p, expected) in [
(Precision::Nanosecond, "ns"),
(Precision::Microsecond, "us"),
(Precision::Millisecond, "ms"),
(Precision::Second, "s"),
] {
assert_eq!(p.as_v2_str(), expected);
assert_eq!(expected.parse::<Precision>().unwrap(), p);
}
}
}
63 changes: 62 additions & 1 deletion tests/client.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
use std::time::{SystemTime, UNIX_EPOCH};

use influxdb3_client::{Client, ClientConfig, Error, Point, Row, Value, WriteOptions};
use influxdb3_client::{Client, ClientConfig, Error, Point, Precision, Row, Value, WriteOptions};

const MEASUREMENT: &str = "rust_e2e";
const LOCATION: &str = "sun-valley-1";
Expand Down Expand Up @@ -183,6 +183,67 @@ async fn partial_write_error() -> Result<(), Box<dyn std::error::Error>> {
Ok(())
}

#[tokio::test]
async fn write_and_query_data_with_v2_precision() -> Result<(), Box<dyn std::error::Error>> {
let Some(config) = testing_config() else {
eprintln!("skipping e2e test: TESTING_INFLUXDB_* env vars are not set");
return Ok(());
};

let v2_write_options = WriteOptions {
precision: Precision::Second,
default_tags: config.write_options.default_tags.clone(),
gzip_threshold: config.write_options.gzip_threshold,
no_sync: config.write_options.no_sync,
accept_partial: config.write_options.accept_partial,
use_v2_api: true,
tag_order: config.write_options.tag_order,
batch_size: config.write_options.batch_size,
max_inflight: config.write_options.max_inflight,
};

let v2_config = ClientConfig::builder()
.host(config.host.clone())
.database(config.database.clone())
.token(config.token.unwrap())
.write_options(v2_write_options)
.build()
.expect("v2 config should be built correctly");

let client = Client::new(v2_config).await?;
let test_id = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos() as i64;

let point = Point::new(MEASUREMENT)
.tag("location", LOCATION)
.field("temp", 15.5_f64)
.field("index", 80_i64)
.field("uindex", 800_u64)
.field("valid", true)
.field("testId", test_id)
.field("text", "a1")
.timestamp_nanos(test_id);

client.write(vec![point]).await?;

let row = query_written_point(&client, test_id)
.await?
.unwrap_or_else(|| panic!("expected to query back point with test_id={test_id}"));

assert_eq!(row["location"].as_str(), Some(LOCATION));
assert_eq!(row["temp"].as_f64(), Some(15.5));
assert_eq!(row["index"].as_i64(), Some(80));
assert_eq!(row["uindex"], Value::U64(800));
assert_eq!(row["valid"].as_bool(), Some(true));
assert_eq!(row["testId"].as_i64(), Some(test_id));
assert_eq!(row["text"].as_str(), Some("a1"));
assert_eq!(
row["time"],
Value::Timestamp((test_id / 1_000_000_000) * 1_000_000_000)
);

Ok(())
}

async fn query_written_point(
client: &Client,
test_id: i64,
Expand Down
74 changes: 72 additions & 2 deletions tests/write_tests.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
use influxdb3_client::error::LineError;
/// Write-path integration tests against a mockito HTTP server.
use influxdb3_client::{Client, ClientConfig, Error, Point, Precision};
use influxdb3_client::{Client, ClientConfig, Error, Point, Precision, WriteOptions};
use mockito::{Matcher, Server};

async fn make_client(server: &Server) -> Client {
Expand Down Expand Up @@ -52,7 +52,7 @@ async fn v2_write_uses_bucket_query_parameter() {
.mock("POST", "/api/v2/write")
.match_query(Matcher::AllOf(vec![
Matcher::UrlEncoded("bucket".into(), "testdb".into()),
Matcher::UrlEncoded("precision".into(), "nanosecond".into()),
Matcher::UrlEncoded("precision".into(), "ns".into()),
]))
.match_header("Authorization", "Bearer test-token")
.match_header("Content-Type", Matcher::Regex("text/plain.*".into()))
Expand Down Expand Up @@ -586,3 +586,73 @@ async fn test_write_error_classification() {
_m.assert_async().await;
}
}

#[tokio::test]
async fn test_v2_precision() {
struct V2PrecisionCase {
url_encode: &'static str,
precision: Precision,
}

let precision_cases = vec![
V2PrecisionCase {
url_encode: "ns",
precision: Precision::Nanosecond,
},
V2PrecisionCase {
url_encode: "us",
precision: Precision::Microsecond,
},
V2PrecisionCase {
url_encode: "ms",
precision: Precision::Millisecond,
},
V2PrecisionCase {
url_encode: "s",
precision: Precision::Second,
},
];

for case in precision_cases {
let mut server = Server::new_async().await;

let _m = server
.mock("POST", "/api/v2/write")
.match_query(Matcher::AllOf(vec![Matcher::UrlEncoded(
"precision".into(),
case.url_encode.into(),
)]))
.with_status(204)
.expect_at_least(1)
.create_async()
.await;

let write_options = WriteOptions {
precision: case.precision,
default_tags: Default::default(),
gzip_threshold: None,
no_sync: false,
accept_partial: false,
use_v2_api: true,
tag_order: vec![],
batch_size: 0,
max_inflight: 0,
};

let client = Client::new(
ClientConfig::builder()
.host(server.url())
.database("testdb")
.token("test-token")
.write_options(write_options)
.build()
.unwrap(),
)
.await
.unwrap();

let _ = client.write("cpu usage=1.0").await;

_m.assert_async().await;
}
}
Loading