slbotbm commented on code in PR #4018:
URL: https://github.com/apache/iggy/pull/4018#discussion_r3935551603
##########
foreign/python/src/client.rs:
##########
@@ -157,6 +158,25 @@ impl IggyClient {
})
}
+ /// Get the statistics and details of the server and its running process.
+ ///
+ /// Returns:
+ /// An awaitable that resolves to `Stats`.
+ ///
+ /// Raises:
+ /// RuntimeError: If the request fails.
Review Comment:
The Rust `SystemClient::get_stats` contract requires authentication and the
`read_servers` permission. Document both. Test before connect/login, after
disconnect, without permission, and with `read_servers`/`manage_servers`.
##########
foreign/python/src/stats.rs:
##########
@@ -0,0 +1,328 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use crate::duration::iggy_duration_to_py_delta;
+use iggy::prelude::{
+ CacheMetrics as RustCacheMetrics, CacheMetricsKey as RustCacheMetricsKey,
Stats as RustStats,
+};
+use pyo3::prelude::*;
+use pyo3::types::{PyDelta, PyDict};
+use pyo3_stub_gen::derive::{gen_stub_pyclass, gen_stub_pymethods};
+
+/// Key identifying the partition a `CacheMetrics` entry belongs to.
+///
+/// Hashable and comparable, so it can key the `Stats.cache_metrics` dict.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+#[gen_stub_pyclass]
+#[pyclass(eq, frozen, hash, skip_from_py_object)]
+pub struct CacheMetricsKey {
+ /// The unique identifier (numeric) of the stream.
+ #[pyo3(get)]
+ pub stream_id: u32,
+ /// The unique identifier (numeric) of the topic within the stream.
+ #[pyo3(get)]
+ pub topic_id: u32,
+ /// The unique identifier (numeric) of the partition within the topic.
+ #[pyo3(get)]
+ pub partition_id: u32,
+}
+
+impl From<&RustCacheMetricsKey> for CacheMetricsKey {
+ fn from(key: &RustCacheMetricsKey) -> Self {
+ Self {
+ stream_id: key.stream_id,
+ topic_id: key.topic_id,
+ partition_id: key.partition_id,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetricsKey {
+ #[new]
+ fn new(stream_id: u32, topic_id: u32, partition_id: u32) -> Self {
+ Self {
+ stream_id,
+ topic_id,
+ partition_id,
+ }
+ }
+
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetricsKey(stream_id={}, topic_id={}, partition_id={})",
+ self.stream_id, self.topic_id, self.partition_id
+ )
+ }
+}
+
+/// Cache metrics for a specific partition.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct CacheMetrics {
+ /// Number of cache hits.
+ #[pyo3(get)]
+ pub hits: u64,
+ /// Number of cache misses.
+ #[pyo3(get)]
+ pub misses: u64,
+ /// Hit ratio (hits / (hits + misses)).
+ #[pyo3(get)]
+ pub hit_ratio: f32,
+}
+
+impl From<&RustCacheMetrics> for CacheMetrics {
+ fn from(metrics: &RustCacheMetrics) -> Self {
+ Self {
+ hits: metrics.hits,
+ misses: metrics.misses,
+ hit_ratio: metrics.hit_ratio,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetrics {
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetrics(hits={}, misses={}, hit_ratio={})",
+ self.hits, self.misses, self.hit_ratio
+ )
+ }
+}
+
+/// The statistics and details of the server and its running process.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct Stats {
+ pub(crate) inner: RustStats,
+ /// Converted once here so that every `cache_metrics` access returns the
+ /// same dict instead of re-collecting the whole map.
+ cache_metrics: Py<PyDict>,
+}
+
+impl From<RustStats> for Stats {
+ fn from(stats: RustStats) -> Self {
+ let cache_metrics = Python::attach(|py| {
+ let dict = PyDict::new(py);
+ for (key, metrics) in &stats.cache_metrics {
+ dict.set_item(CacheMetricsKey::from(key),
CacheMetrics::from(metrics))
+ .expect("insert cache metrics entry");
+ }
+ dict.unbind()
+ });
+ Self {
+ inner: stats,
+ cache_metrics,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl Stats {
+ /// The unique identifier of the server process.
+ #[getter]
+ pub fn process_id(&self) -> u32 {
+ self.inner.process_id
+ }
+
+ /// The CPU usage of the server process, in percent.
+ #[getter]
+ pub fn cpu_usage(&self) -> f32 {
+ self.inner.cpu_usage
+ }
+
+ /// The total CPU usage of the system, in percent.
+ #[getter]
+ pub fn total_cpu_usage(&self) -> f32 {
+ self.inner.total_cpu_usage
+ }
+
+ /// The memory usage of the server process, in bytes.
+ #[getter]
+ pub fn memory_usage(&self) -> u64 {
+ self.inner.memory_usage.as_bytes_u64()
+ }
+
+ /// The total memory of the system, in bytes.
+ #[getter]
+ pub fn total_memory(&self) -> u64 {
+ self.inner.total_memory.as_bytes_u64()
+ }
+
+ /// The available memory of the system, in bytes.
+ #[getter]
Review Comment:
Document that CPU is a per-shard delta and the first sample can be zero.
Also document that runtime and start time have whole-second precision.
##########
foreign/python/tests/test_stats.py:
##########
Review Comment:
Add a single-node topology test with multiple streams/topics/partitions,
consumer groups, and a second client. Verify increases in `streams_count`,
`topics_count`, `partitions_count`, `segments_count`, `consumer_groups_count`,
and `clients_count`; after cleanup use `<=` deltas.
##########
foreign/python/src/stats.rs:
##########
@@ -0,0 +1,328 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use crate::duration::iggy_duration_to_py_delta;
+use iggy::prelude::{
+ CacheMetrics as RustCacheMetrics, CacheMetricsKey as RustCacheMetricsKey,
Stats as RustStats,
+};
+use pyo3::prelude::*;
+use pyo3::types::{PyDelta, PyDict};
+use pyo3_stub_gen::derive::{gen_stub_pyclass, gen_stub_pymethods};
+
+/// Key identifying the partition a `CacheMetrics` entry belongs to.
+///
+/// Hashable and comparable, so it can key the `Stats.cache_metrics` dict.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+#[gen_stub_pyclass]
+#[pyclass(eq, frozen, hash, skip_from_py_object)]
+pub struct CacheMetricsKey {
+ /// The unique identifier (numeric) of the stream.
+ #[pyo3(get)]
+ pub stream_id: u32,
+ /// The unique identifier (numeric) of the topic within the stream.
+ #[pyo3(get)]
+ pub topic_id: u32,
+ /// The unique identifier (numeric) of the partition within the topic.
+ #[pyo3(get)]
+ pub partition_id: u32,
+}
+
+impl From<&RustCacheMetricsKey> for CacheMetricsKey {
+ fn from(key: &RustCacheMetricsKey) -> Self {
+ Self {
+ stream_id: key.stream_id,
+ topic_id: key.topic_id,
+ partition_id: key.partition_id,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetricsKey {
+ #[new]
+ fn new(stream_id: u32, topic_id: u32, partition_id: u32) -> Self {
+ Self {
+ stream_id,
+ topic_id,
+ partition_id,
+ }
+ }
+
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetricsKey(stream_id={}, topic_id={}, partition_id={})",
+ self.stream_id, self.topic_id, self.partition_id
+ )
+ }
+}
+
+/// Cache metrics for a specific partition.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct CacheMetrics {
+ /// Number of cache hits.
+ #[pyo3(get)]
+ pub hits: u64,
+ /// Number of cache misses.
+ #[pyo3(get)]
+ pub misses: u64,
+ /// Hit ratio (hits / (hits + misses)).
+ #[pyo3(get)]
+ pub hit_ratio: f32,
+}
+
+impl From<&RustCacheMetrics> for CacheMetrics {
+ fn from(metrics: &RustCacheMetrics) -> Self {
+ Self {
+ hits: metrics.hits,
+ misses: metrics.misses,
+ hit_ratio: metrics.hit_ratio,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetrics {
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetrics(hits={}, misses={}, hit_ratio={})",
+ self.hits, self.misses, self.hit_ratio
+ )
+ }
+}
+
+/// The statistics and details of the server and its running process.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct Stats {
+ pub(crate) inner: RustStats,
+ /// Converted once here so that every `cache_metrics` access returns the
+ /// same dict instead of re-collecting the whole map.
+ cache_metrics: Py<PyDict>,
+}
+
+impl From<RustStats> for Stats {
+ fn from(stats: RustStats) -> Self {
+ let cache_metrics = Python::attach(|py| {
+ let dict = PyDict::new(py);
+ for (key, metrics) in &stats.cache_metrics {
+ dict.set_item(CacheMetricsKey::from(key),
CacheMetrics::from(metrics))
+ .expect("insert cache metrics entry");
+ }
+ dict.unbind()
+ });
+ Self {
+ inner: stats,
+ cache_metrics,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl Stats {
+ /// The unique identifier of the server process.
+ #[getter]
+ pub fn process_id(&self) -> u32 {
+ self.inner.process_id
+ }
+
+ /// The CPU usage of the server process, in percent.
+ #[getter]
+ pub fn cpu_usage(&self) -> f32 {
+ self.inner.cpu_usage
+ }
+
+ /// The total CPU usage of the system, in percent.
+ #[getter]
+ pub fn total_cpu_usage(&self) -> f32 {
+ self.inner.total_cpu_usage
+ }
+
+ /// The memory usage of the server process, in bytes.
+ #[getter]
+ pub fn memory_usage(&self) -> u64 {
+ self.inner.memory_usage.as_bytes_u64()
+ }
+
+ /// The total memory of the system, in bytes.
+ #[getter]
+ pub fn total_memory(&self) -> u64 {
+ self.inner.total_memory.as_bytes_u64()
+ }
+
+ /// The available memory of the system, in bytes.
+ #[getter]
+ pub fn available_memory(&self) -> u64 {
+ self.inner.available_memory.as_bytes_u64()
+ }
+
+ /// The run time of the server process.
+ #[getter]
+ #[gen_stub(override_return_type(type_repr = "datetime.timedelta",
imports=("datetime")))]
+ pub fn run_time<'a>(&self, py: Python<'a>) -> PyResult<Bound<'a, PyDelta>>
{
+ iggy_duration_to_py_delta(py, self.inner.run_time)
+ }
+
+ /// The start time of the server process, in microseconds since the Unix
epoch.
+ #[getter]
+ pub fn start_time(&self) -> u64 {
+ self.inner.start_time.as_micros()
+ }
+
+ /// The total number of bytes read.
+ #[getter]
+ pub fn read_bytes(&self) -> u64 {
+ self.inner.read_bytes.as_bytes_u64()
+ }
+
+ /// The total number of bytes written.
+ #[getter]
+ pub fn written_bytes(&self) -> u64 {
+ self.inner.written_bytes.as_bytes_u64()
+ }
+
+ /// The total size of the messages, in bytes.
+ #[getter]
+ pub fn messages_size_bytes(&self) -> u64 {
+ self.inner.messages_size_bytes.as_bytes_u64()
+ }
+
+ /// The total number of streams.
+ #[getter]
+ pub fn streams_count(&self) -> u32 {
+ self.inner.streams_count
+ }
+
+ /// The total number of topics.
+ #[getter]
+ pub fn topics_count(&self) -> u32 {
+ self.inner.topics_count
+ }
+
+ /// The total number of partitions.
+ #[getter]
+ pub fn partitions_count(&self) -> u32 {
+ self.inner.partitions_count
+ }
+
+ /// The total number of segments.
+ #[getter]
+ pub fn segments_count(&self) -> u32 {
+ self.inner.segments_count
+ }
+
+ /// The total number of messages.
+ #[getter]
+ pub fn messages_count(&self) -> u64 {
+ self.inner.messages_count
+ }
+
+ /// The total number of connected clients.
+ #[getter]
+ pub fn clients_count(&self) -> u32 {
+ self.inner.clients_count
+ }
+
+ /// The total number of consumer groups.
+ #[getter]
+ pub fn consumer_groups_count(&self) -> u32 {
+ self.inner.consumer_groups_count
+ }
+
+ /// The name of the host the server runs on.
+ #[getter]
+ pub fn hostname(&self) -> String {
+ self.inner.hostname.clone()
+ }
+
+ /// The name of the operating system.
+ #[getter]
+ pub fn os_name(&self) -> String {
+ self.inner.os_name.clone()
+ }
+
+ /// The version of the operating system.
+ #[getter]
+ pub fn os_version(&self) -> String {
+ self.inner.os_version.clone()
+ }
+
+ /// The version of the kernel.
+ #[getter]
+ pub fn kernel_version(&self) -> String {
+ self.inner.kernel_version.clone()
+ }
+
+ /// The version of the Iggy server.
+ #[getter]
+ pub fn iggy_server_version(&self) -> String {
+ self.inner.iggy_server_version.clone()
+ }
+
+ /// The numeric semantic version of the Iggy server, or `None` when
unknown.
+ /// E.g. 1.2.3 -> 1002003 (major * 1000000 + minor * 1000 + patch).
+ #[getter]
+ #[gen_stub(override_return_type(type_repr = "builtins.int | None"))]
+ pub fn iggy_server_semver(&self) -> Option<u32> {
+ self.inner.iggy_server_semver
+ }
+
+ /// Cache metrics per partition.
+ ///
+ /// Built once when the stats snapshot is created; every access returns the
+ /// same dict.
+ #[getter]
+ #[gen_stub(override_return_type(type_repr =
"builtins.dict[CacheMetricsKey, CacheMetrics]"))]
+ pub fn cache_metrics(&self, py: Python<'_>) -> Py<PyDict> {
+ self.cache_metrics.clone_ref(py)
+ }
Review Comment:
The server currently returns an empty map. Document that and assert
`stats.cache_metrics == {}`; keep the key/hash test separate.
##########
foreign/python/src/stats.rs:
##########
@@ -0,0 +1,328 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use crate::duration::iggy_duration_to_py_delta;
+use iggy::prelude::{
+ CacheMetrics as RustCacheMetrics, CacheMetricsKey as RustCacheMetricsKey,
Stats as RustStats,
+};
+use pyo3::prelude::*;
+use pyo3::types::{PyDelta, PyDict};
+use pyo3_stub_gen::derive::{gen_stub_pyclass, gen_stub_pymethods};
+
+/// Key identifying the partition a `CacheMetrics` entry belongs to.
+///
+/// Hashable and comparable, so it can key the `Stats.cache_metrics` dict.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+#[gen_stub_pyclass]
+#[pyclass(eq, frozen, hash, skip_from_py_object)]
+pub struct CacheMetricsKey {
+ /// The unique identifier (numeric) of the stream.
+ #[pyo3(get)]
+ pub stream_id: u32,
+ /// The unique identifier (numeric) of the topic within the stream.
+ #[pyo3(get)]
+ pub topic_id: u32,
+ /// The unique identifier (numeric) of the partition within the topic.
+ #[pyo3(get)]
+ pub partition_id: u32,
+}
+
+impl From<&RustCacheMetricsKey> for CacheMetricsKey {
+ fn from(key: &RustCacheMetricsKey) -> Self {
+ Self {
+ stream_id: key.stream_id,
+ topic_id: key.topic_id,
+ partition_id: key.partition_id,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetricsKey {
+ #[new]
+ fn new(stream_id: u32, topic_id: u32, partition_id: u32) -> Self {
+ Self {
+ stream_id,
+ topic_id,
+ partition_id,
+ }
+ }
+
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetricsKey(stream_id={}, topic_id={}, partition_id={})",
+ self.stream_id, self.topic_id, self.partition_id
+ )
+ }
+}
+
+/// Cache metrics for a specific partition.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct CacheMetrics {
+ /// Number of cache hits.
+ #[pyo3(get)]
+ pub hits: u64,
+ /// Number of cache misses.
+ #[pyo3(get)]
+ pub misses: u64,
+ /// Hit ratio (hits / (hits + misses)).
+ #[pyo3(get)]
+ pub hit_ratio: f32,
+}
+
+impl From<&RustCacheMetrics> for CacheMetrics {
+ fn from(metrics: &RustCacheMetrics) -> Self {
+ Self {
+ hits: metrics.hits,
+ misses: metrics.misses,
+ hit_ratio: metrics.hit_ratio,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetrics {
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetrics(hits={}, misses={}, hit_ratio={})",
+ self.hits, self.misses, self.hit_ratio
+ )
+ }
+}
+
+/// The statistics and details of the server and its running process.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct Stats {
+ pub(crate) inner: RustStats,
+ /// Converted once here so that every `cache_metrics` access returns the
+ /// same dict instead of re-collecting the whole map.
+ cache_metrics: Py<PyDict>,
+}
+
+impl From<RustStats> for Stats {
+ fn from(stats: RustStats) -> Self {
+ let cache_metrics = Python::attach(|py| {
+ let dict = PyDict::new(py);
+ for (key, metrics) in &stats.cache_metrics {
+ dict.set_item(CacheMetricsKey::from(key),
CacheMetrics::from(metrics))
+ .expect("insert cache metrics entry");
+ }
+ dict.unbind()
+ });
+ Self {
+ inner: stats,
+ cache_metrics,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl Stats {
+ /// The unique identifier of the server process.
+ #[getter]
+ pub fn process_id(&self) -> u32 {
+ self.inner.process_id
+ }
+
+ /// The CPU usage of the server process, in percent.
+ #[getter]
+ pub fn cpu_usage(&self) -> f32 {
+ self.inner.cpu_usage
+ }
+
+ /// The total CPU usage of the system, in percent.
+ #[getter]
+ pub fn total_cpu_usage(&self) -> f32 {
+ self.inner.total_cpu_usage
+ }
+
+ /// The memory usage of the server process, in bytes.
+ #[getter]
+ pub fn memory_usage(&self) -> u64 {
+ self.inner.memory_usage.as_bytes_u64()
+ }
+
+ /// The total memory of the system, in bytes.
+ #[getter]
+ pub fn total_memory(&self) -> u64 {
+ self.inner.total_memory.as_bytes_u64()
+ }
+
+ /// The available memory of the system, in bytes.
+ #[getter]
Review Comment:
Also, document that these are not an atomic snapshot.
##########
foreign/python/src/stats.rs:
##########
@@ -0,0 +1,328 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use crate::duration::iggy_duration_to_py_delta;
+use iggy::prelude::{
+ CacheMetrics as RustCacheMetrics, CacheMetricsKey as RustCacheMetricsKey,
Stats as RustStats,
+};
+use pyo3::prelude::*;
+use pyo3::types::{PyDelta, PyDict};
+use pyo3_stub_gen::derive::{gen_stub_pyclass, gen_stub_pymethods};
+
+/// Key identifying the partition a `CacheMetrics` entry belongs to.
+///
+/// Hashable and comparable, so it can key the `Stats.cache_metrics` dict.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+#[gen_stub_pyclass]
+#[pyclass(eq, frozen, hash, skip_from_py_object)]
+pub struct CacheMetricsKey {
+ /// The unique identifier (numeric) of the stream.
+ #[pyo3(get)]
+ pub stream_id: u32,
+ /// The unique identifier (numeric) of the topic within the stream.
+ #[pyo3(get)]
+ pub topic_id: u32,
+ /// The unique identifier (numeric) of the partition within the topic.
+ #[pyo3(get)]
+ pub partition_id: u32,
+}
+
+impl From<&RustCacheMetricsKey> for CacheMetricsKey {
+ fn from(key: &RustCacheMetricsKey) -> Self {
+ Self {
+ stream_id: key.stream_id,
+ topic_id: key.topic_id,
+ partition_id: key.partition_id,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetricsKey {
+ #[new]
+ fn new(stream_id: u32, topic_id: u32, partition_id: u32) -> Self {
+ Self {
+ stream_id,
+ topic_id,
+ partition_id,
+ }
+ }
+
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetricsKey(stream_id={}, topic_id={}, partition_id={})",
+ self.stream_id, self.topic_id, self.partition_id
+ )
+ }
+}
+
+/// Cache metrics for a specific partition.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct CacheMetrics {
+ /// Number of cache hits.
+ #[pyo3(get)]
+ pub hits: u64,
+ /// Number of cache misses.
+ #[pyo3(get)]
+ pub misses: u64,
+ /// Hit ratio (hits / (hits + misses)).
+ #[pyo3(get)]
+ pub hit_ratio: f32,
+}
+
+impl From<&RustCacheMetrics> for CacheMetrics {
+ fn from(metrics: &RustCacheMetrics) -> Self {
+ Self {
+ hits: metrics.hits,
+ misses: metrics.misses,
+ hit_ratio: metrics.hit_ratio,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetrics {
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetrics(hits={}, misses={}, hit_ratio={})",
+ self.hits, self.misses, self.hit_ratio
+ )
+ }
+}
+
+/// The statistics and details of the server and its running process.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct Stats {
+ pub(crate) inner: RustStats,
+ /// Converted once here so that every `cache_metrics` access returns the
+ /// same dict instead of re-collecting the whole map.
+ cache_metrics: Py<PyDict>,
+}
+
+impl From<RustStats> for Stats {
+ fn from(stats: RustStats) -> Self {
+ let cache_metrics = Python::attach(|py| {
+ let dict = PyDict::new(py);
+ for (key, metrics) in &stats.cache_metrics {
+ dict.set_item(CacheMetricsKey::from(key),
CacheMetrics::from(metrics))
+ .expect("insert cache metrics entry");
+ }
+ dict.unbind()
+ });
+ Self {
+ inner: stats,
+ cache_metrics,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl Stats {
+ /// The unique identifier of the server process.
+ #[getter]
+ pub fn process_id(&self) -> u32 {
+ self.inner.process_id
+ }
+
+ /// The CPU usage of the server process, in percent.
+ #[getter]
+ pub fn cpu_usage(&self) -> f32 {
+ self.inner.cpu_usage
+ }
+
+ /// The total CPU usage of the system, in percent.
+ #[getter]
+ pub fn total_cpu_usage(&self) -> f32 {
+ self.inner.total_cpu_usage
+ }
+
+ /// The memory usage of the server process, in bytes.
+ #[getter]
+ pub fn memory_usage(&self) -> u64 {
+ self.inner.memory_usage.as_bytes_u64()
+ }
+
+ /// The total memory of the system, in bytes.
+ #[getter]
+ pub fn total_memory(&self) -> u64 {
+ self.inner.total_memory.as_bytes_u64()
+ }
+
+ /// The available memory of the system, in bytes.
+ #[getter]
+ pub fn available_memory(&self) -> u64 {
+ self.inner.available_memory.as_bytes_u64()
+ }
+
+ /// The run time of the server process.
+ #[getter]
+ #[gen_stub(override_return_type(type_repr = "datetime.timedelta",
imports=("datetime")))]
+ pub fn run_time<'a>(&self, py: Python<'a>) -> PyResult<Bound<'a, PyDelta>>
{
+ iggy_duration_to_py_delta(py, self.inner.run_time)
+ }
+
+ /// The start time of the server process, in microseconds since the Unix
epoch.
+ #[getter]
+ pub fn start_time(&self) -> u64 {
+ self.inner.start_time.as_micros()
+ }
+
+ /// The total number of bytes read.
+ #[getter]
+ pub fn read_bytes(&self) -> u64 {
+ self.inner.read_bytes.as_bytes_u64()
+ }
+
+ /// The total number of bytes written.
+ #[getter]
+ pub fn written_bytes(&self) -> u64 {
+ self.inner.written_bytes.as_bytes_u64()
+ }
+
+ /// The total size of the messages, in bytes.
+ #[getter]
+ pub fn messages_size_bytes(&self) -> u64 {
+ self.inner.messages_size_bytes.as_bytes_u64()
+ }
+
+ /// The total number of streams.
+ #[getter]
+ pub fn streams_count(&self) -> u32 {
+ self.inner.streams_count
+ }
+
+ /// The total number of topics.
+ #[getter]
+ pub fn topics_count(&self) -> u32 {
+ self.inner.topics_count
+ }
+
+ /// The total number of partitions.
+ #[getter]
+ pub fn partitions_count(&self) -> u32 {
+ self.inner.partitions_count
+ }
+
+ /// The total number of segments.
+ #[getter]
+ pub fn segments_count(&self) -> u32 {
+ self.inner.segments_count
+ }
+
+ /// The total number of messages.
+ #[getter]
+ pub fn messages_count(&self) -> u64 {
+ self.inner.messages_count
+ }
+
+ /// The total number of connected clients.
+ #[getter]
+ pub fn clients_count(&self) -> u32 {
+ self.inner.clients_count
+ }
+
+ /// The total number of consumer groups.
+ #[getter]
+ pub fn consumer_groups_count(&self) -> u32 {
+ self.inner.consumer_groups_count
+ }
+
+ /// The name of the host the server runs on.
+ #[getter]
+ pub fn hostname(&self) -> String {
+ self.inner.hostname.clone()
+ }
+
+ /// The name of the operating system.
+ #[getter]
+ pub fn os_name(&self) -> String {
+ self.inner.os_name.clone()
+ }
+
+ /// The version of the operating system.
+ #[getter]
+ pub fn os_version(&self) -> String {
+ self.inner.os_version.clone()
+ }
+
+ /// The version of the kernel.
+ #[getter]
+ pub fn kernel_version(&self) -> String {
+ self.inner.kernel_version.clone()
+ }
+
+ /// The version of the Iggy server.
+ #[getter]
+ pub fn iggy_server_version(&self) -> String {
+ self.inner.iggy_server_version.clone()
+ }
+
+ /// The numeric semantic version of the Iggy server, or `None` when
unknown.
+ /// E.g. 1.2.3 -> 1002003 (major * 1000000 + minor * 1000 + patch).
+ #[getter]
+ #[gen_stub(override_return_type(type_repr = "builtins.int | None"))]
+ pub fn iggy_server_semver(&self) -> Option<u32> {
+ self.inner.iggy_server_semver
+ }
+
+ /// Cache metrics per partition.
+ ///
+ /// Built once when the stats snapshot is created; every access returns the
+ /// same dict.
+ #[getter]
+ #[gen_stub(override_return_type(type_repr =
"builtins.dict[CacheMetricsKey, CacheMetrics]"))]
+ pub fn cache_metrics(&self, py: Python<'_>) -> Py<PyDict> {
+ self.cache_metrics.clone_ref(py)
+ }
+
+ /// The number of threads in the server process.
+ #[getter]
+ pub fn threads_count(&self) -> u32 {
+ self.inner.threads_count
+ }
+
+ /// The available (free) disk space for the data directory, in bytes.
+ #[getter]
+ pub fn free_disk_space(&self) -> u64 {
+ self.inner.free_disk_space.as_bytes_u64()
+ }
+
+ /// The total disk space for the data directory, in bytes.
+ #[getter]
+ pub fn total_disk_space(&self) -> u64 {
+ self.inner.total_disk_space.as_bytes_u64()
+ }
Review Comment:
Document that disk values can be zero when the data path or probe is
unavailable. Also assert `free_disk_space <= total_disk_space` in tests.
##########
foreign/python/tests/test_stats.py:
##########
@@ -0,0 +1,103 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+import datetime
+
+import pytest
+
+from apache_iggy import CacheMetrics, CacheMetricsKey, IggyClient
+from apache_iggy import SendMessage as Message
+
+
+class TestStats:
+ """Test server statistics retrieval."""
+
+ @pytest.mark.asyncio
+ async def test_get_stats(self, iggy_client: IggyClient, unique_name):
+ """Sending messages moves the server counts reported by get_stats."""
+ stats_before = await iggy_client.get_stats()
+
+ stream_name = unique_name()
+ topic_name = unique_name()
+ await iggy_client.create_stream(stream_name)
+ await iggy_client.create_topic(
+ stream=stream_name, name=topic_name, partitions_count=1
+ )
+ await iggy_client.send_messages(
+ stream=stream_name,
+ topic=topic_name,
+ partitioning=0,
+ messages=[Message(f"stats message {i}") for i in range(3)],
+ )
+
+ stats = await iggy_client.get_stats()
+
+ # `>=` rather than exact equality: the counters are server-global, so
+ # concurrently running tests (e.g. under pytest-xdist) may bump them
too.
+ assert stats.streams_count >= stats_before.streams_count + 1
+ assert stats.topics_count >= stats_before.topics_count + 1
+ assert stats.partitions_count >= stats_before.partitions_count + 1
+ assert stats.messages_count >= stats_before.messages_count + 3
+ assert stats.clients_count >= 1
+
+ assert stats.iggy_server_version
+ assert stats.hostname
+ assert stats.process_id > 0
+ assert stats.start_time > 0
+ assert stats.total_memory > 0
+ assert stats.total_disk_space > 0
+
+ assert isinstance(stats.run_time, datetime.timedelta)
+ assert stats.run_time >= stats_before.run_time
+
+ assert f"streams_count={stats.streams_count}" in repr(stats)
+ assert stats.hostname in repr(stats)
+
+ @pytest.mark.asyncio
+ async def test_get_stats_cache_metrics_dict(self, iggy_client: IggyClient):
+ """cache_metrics is a dict keyed by CacheMetricsKey, and repeated
+ accesses return the same dict rather than rebuilding it."""
+ stats = await iggy_client.get_stats()
+
+ assert isinstance(stats.cache_metrics, dict)
+ # The getter must not re-collect the map on every access.
+ assert stats.cache_metrics is stats.cache_metrics
+ # The server does not populate cache metrics yet (`GetStats` replies
+ # with an empty map), so entries are only checked when present.
+ for key, metrics in stats.cache_metrics.items():
+ assert isinstance(key, CacheMetricsKey)
+ assert isinstance(metrics, CacheMetrics)
Review Comment:
Assert the current empty server response here. Keep key equality and hashing
in the separate no-server test.
##########
foreign/python/tests/test_stats.py:
##########
@@ -0,0 +1,103 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+import datetime
+
+import pytest
+
+from apache_iggy import CacheMetrics, CacheMetricsKey, IggyClient
+from apache_iggy import SendMessage as Message
+
+
+class TestStats:
+ """Test server statistics retrieval."""
+
+ @pytest.mark.asyncio
+ async def test_get_stats(self, iggy_client: IggyClient, unique_name):
+ """Sending messages moves the server counts reported by get_stats."""
+ stats_before = await iggy_client.get_stats()
+
+ stream_name = unique_name()
+ topic_name = unique_name()
+ await iggy_client.create_stream(stream_name)
+ await iggy_client.create_topic(
+ stream=stream_name, name=topic_name, partitions_count=1
+ )
+ await iggy_client.send_messages(
+ stream=stream_name,
+ topic=topic_name,
+ partitioning=0,
+ messages=[Message(f"stats message {i}") for i in range(3)],
+ )
+
+ stats = await iggy_client.get_stats()
+
+ # `>=` rather than exact equality: the counters are server-global, so
+ # concurrently running tests (e.g. under pytest-xdist) may bump them
too.
+ assert stats.streams_count >= stats_before.streams_count + 1
+ assert stats.topics_count >= stats_before.topics_count + 1
+ assert stats.partitions_count >= stats_before.partitions_count + 1
+ assert stats.messages_count >= stats_before.messages_count + 3
+ assert stats.clients_count >= 1
+
+ assert stats.iggy_server_version
+ assert stats.hostname
+ assert stats.process_id > 0
+ assert stats.start_time > 0
+ assert stats.total_memory > 0
+ assert stats.total_disk_space > 0
+
+ assert isinstance(stats.run_time, datetime.timedelta)
+ assert stats.run_time >= stats_before.run_time
Review Comment:
Also assert that `messages_size_bytes` increases after sending messages.
##########
foreign/python/tests/test_stats.py:
##########
@@ -0,0 +1,103 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+import datetime
+
+import pytest
+
+from apache_iggy import CacheMetrics, CacheMetricsKey, IggyClient
+from apache_iggy import SendMessage as Message
+
+
+class TestStats:
+ """Test server statistics retrieval."""
+
+ @pytest.mark.asyncio
+ async def test_get_stats(self, iggy_client: IggyClient, unique_name):
+ """Sending messages moves the server counts reported by get_stats."""
+ stats_before = await iggy_client.get_stats()
+
+ stream_name = unique_name()
+ topic_name = unique_name()
+ await iggy_client.create_stream(stream_name)
+ await iggy_client.create_topic(
+ stream=stream_name, name=topic_name, partitions_count=1
+ )
+ await iggy_client.send_messages(
+ stream=stream_name,
+ topic=topic_name,
+ partitioning=0,
+ messages=[Message(f"stats message {i}") for i in range(3)],
+ )
+
+ stats = await iggy_client.get_stats()
+
+ # `>=` rather than exact equality: the counters are server-global, so
+ # concurrently running tests (e.g. under pytest-xdist) may bump them
too.
+ assert stats.streams_count >= stats_before.streams_count + 1
+ assert stats.topics_count >= stats_before.topics_count + 1
+ assert stats.partitions_count >= stats_before.partitions_count + 1
+ assert stats.messages_count >= stats_before.messages_count + 3
+ assert stats.clients_count >= 1
+
+ assert stats.iggy_server_version
+ assert stats.hostname
+ assert stats.process_id > 0
+ assert stats.start_time > 0
+ assert stats.total_memory > 0
+ assert stats.total_disk_space > 0
+
+ assert isinstance(stats.run_time, datetime.timedelta)
+ assert stats.run_time >= stats_before.run_time
Review Comment:
In addition, ddd checks for `threads_count > 0`, `available_memory <=
total_memory`, `free_disk_space <= total_disk_space`, and non-empty `os_name`,
`os_version`, and `kernel_version`. A second call should keep static fields
stable and runtime nondecreasing.
##########
foreign/python/src/stats.rs:
##########
@@ -0,0 +1,328 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use crate::duration::iggy_duration_to_py_delta;
+use iggy::prelude::{
+ CacheMetrics as RustCacheMetrics, CacheMetricsKey as RustCacheMetricsKey,
Stats as RustStats,
+};
+use pyo3::prelude::*;
+use pyo3::types::{PyDelta, PyDict};
+use pyo3_stub_gen::derive::{gen_stub_pyclass, gen_stub_pymethods};
+
+/// Key identifying the partition a `CacheMetrics` entry belongs to.
+///
+/// Hashable and comparable, so it can key the `Stats.cache_metrics` dict.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+#[gen_stub_pyclass]
+#[pyclass(eq, frozen, hash, skip_from_py_object)]
+pub struct CacheMetricsKey {
+ /// The unique identifier (numeric) of the stream.
+ #[pyo3(get)]
+ pub stream_id: u32,
+ /// The unique identifier (numeric) of the topic within the stream.
+ #[pyo3(get)]
+ pub topic_id: u32,
+ /// The unique identifier (numeric) of the partition within the topic.
+ #[pyo3(get)]
+ pub partition_id: u32,
+}
+
+impl From<&RustCacheMetricsKey> for CacheMetricsKey {
+ fn from(key: &RustCacheMetricsKey) -> Self {
+ Self {
+ stream_id: key.stream_id,
+ topic_id: key.topic_id,
+ partition_id: key.partition_id,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetricsKey {
+ #[new]
+ fn new(stream_id: u32, topic_id: u32, partition_id: u32) -> Self {
+ Self {
+ stream_id,
+ topic_id,
+ partition_id,
+ }
+ }
+
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetricsKey(stream_id={}, topic_id={}, partition_id={})",
+ self.stream_id, self.topic_id, self.partition_id
+ )
+ }
+}
+
+/// Cache metrics for a specific partition.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct CacheMetrics {
+ /// Number of cache hits.
+ #[pyo3(get)]
+ pub hits: u64,
+ /// Number of cache misses.
+ #[pyo3(get)]
+ pub misses: u64,
+ /// Hit ratio (hits / (hits + misses)).
+ #[pyo3(get)]
+ pub hit_ratio: f32,
+}
+
+impl From<&RustCacheMetrics> for CacheMetrics {
+ fn from(metrics: &RustCacheMetrics) -> Self {
+ Self {
+ hits: metrics.hits,
+ misses: metrics.misses,
+ hit_ratio: metrics.hit_ratio,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetrics {
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetrics(hits={}, misses={}, hit_ratio={})",
+ self.hits, self.misses, self.hit_ratio
+ )
+ }
+}
+
+/// The statistics and details of the server and its running process.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct Stats {
+ pub(crate) inner: RustStats,
+ /// Converted once here so that every `cache_metrics` access returns the
+ /// same dict instead of re-collecting the whole map.
+ cache_metrics: Py<PyDict>,
+}
+
+impl From<RustStats> for Stats {
+ fn from(stats: RustStats) -> Self {
+ let cache_metrics = Python::attach(|py| {
+ let dict = PyDict::new(py);
+ for (key, metrics) in &stats.cache_metrics {
+ dict.set_item(CacheMetricsKey::from(key),
CacheMetrics::from(metrics))
+ .expect("insert cache metrics entry");
+ }
+ dict.unbind()
+ });
+ Self {
+ inner: stats,
+ cache_metrics,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl Stats {
+ /// The unique identifier of the server process.
+ #[getter]
+ pub fn process_id(&self) -> u32 {
+ self.inner.process_id
+ }
+
+ /// The CPU usage of the server process, in percent.
+ #[getter]
+ pub fn cpu_usage(&self) -> f32 {
+ self.inner.cpu_usage
+ }
+
+ /// The total CPU usage of the system, in percent.
+ #[getter]
+ pub fn total_cpu_usage(&self) -> f32 {
+ self.inner.total_cpu_usage
+ }
+
+ /// The memory usage of the server process, in bytes.
+ #[getter]
+ pub fn memory_usage(&self) -> u64 {
+ self.inner.memory_usage.as_bytes_u64()
+ }
+
+ /// The total memory of the system, in bytes.
+ #[getter]
+ pub fn total_memory(&self) -> u64 {
+ self.inner.total_memory.as_bytes_u64()
+ }
+
+ /// The available memory of the system, in bytes.
+ #[getter]
Review Comment:
Also, mirror the Rust docs: CPU is cpuset-scoped and memory is cgroup-scoped.
##########
foreign/python/src/stats.rs:
##########
@@ -0,0 +1,328 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use crate::duration::iggy_duration_to_py_delta;
+use iggy::prelude::{
+ CacheMetrics as RustCacheMetrics, CacheMetricsKey as RustCacheMetricsKey,
Stats as RustStats,
+};
+use pyo3::prelude::*;
+use pyo3::types::{PyDelta, PyDict};
+use pyo3_stub_gen::derive::{gen_stub_pyclass, gen_stub_pymethods};
+
+/// Key identifying the partition a `CacheMetrics` entry belongs to.
+///
+/// Hashable and comparable, so it can key the `Stats.cache_metrics` dict.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+#[gen_stub_pyclass]
+#[pyclass(eq, frozen, hash, skip_from_py_object)]
+pub struct CacheMetricsKey {
+ /// The unique identifier (numeric) of the stream.
+ #[pyo3(get)]
+ pub stream_id: u32,
+ /// The unique identifier (numeric) of the topic within the stream.
+ #[pyo3(get)]
+ pub topic_id: u32,
+ /// The unique identifier (numeric) of the partition within the topic.
+ #[pyo3(get)]
+ pub partition_id: u32,
+}
+
+impl From<&RustCacheMetricsKey> for CacheMetricsKey {
+ fn from(key: &RustCacheMetricsKey) -> Self {
+ Self {
+ stream_id: key.stream_id,
+ topic_id: key.topic_id,
+ partition_id: key.partition_id,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetricsKey {
+ #[new]
+ fn new(stream_id: u32, topic_id: u32, partition_id: u32) -> Self {
+ Self {
+ stream_id,
+ topic_id,
+ partition_id,
+ }
+ }
+
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetricsKey(stream_id={}, topic_id={}, partition_id={})",
+ self.stream_id, self.topic_id, self.partition_id
+ )
+ }
+}
+
+/// Cache metrics for a specific partition.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct CacheMetrics {
+ /// Number of cache hits.
+ #[pyo3(get)]
+ pub hits: u64,
+ /// Number of cache misses.
+ #[pyo3(get)]
+ pub misses: u64,
+ /// Hit ratio (hits / (hits + misses)).
+ #[pyo3(get)]
+ pub hit_ratio: f32,
+}
+
+impl From<&RustCacheMetrics> for CacheMetrics {
+ fn from(metrics: &RustCacheMetrics) -> Self {
+ Self {
+ hits: metrics.hits,
+ misses: metrics.misses,
+ hit_ratio: metrics.hit_ratio,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl CacheMetrics {
+ fn __repr__(&self) -> String {
+ format!(
+ "CacheMetrics(hits={}, misses={}, hit_ratio={})",
+ self.hits, self.misses, self.hit_ratio
+ )
+ }
+}
+
+/// The statistics and details of the server and its running process.
+#[gen_stub_pyclass]
+#[pyclass]
+pub struct Stats {
+ pub(crate) inner: RustStats,
+ /// Converted once here so that every `cache_metrics` access returns the
+ /// same dict instead of re-collecting the whole map.
+ cache_metrics: Py<PyDict>,
+}
+
+impl From<RustStats> for Stats {
+ fn from(stats: RustStats) -> Self {
+ let cache_metrics = Python::attach(|py| {
+ let dict = PyDict::new(py);
+ for (key, metrics) in &stats.cache_metrics {
+ dict.set_item(CacheMetricsKey::from(key),
CacheMetrics::from(metrics))
+ .expect("insert cache metrics entry");
+ }
+ dict.unbind()
+ });
+ Self {
+ inner: stats,
+ cache_metrics,
+ }
+ }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl Stats {
+ /// The unique identifier of the server process.
+ #[getter]
+ pub fn process_id(&self) -> u32 {
+ self.inner.process_id
+ }
+
+ /// The CPU usage of the server process, in percent.
+ #[getter]
+ pub fn cpu_usage(&self) -> f32 {
+ self.inner.cpu_usage
+ }
+
+ /// The total CPU usage of the system, in percent.
+ #[getter]
+ pub fn total_cpu_usage(&self) -> f32 {
+ self.inner.total_cpu_usage
+ }
+
+ /// The memory usage of the server process, in bytes.
+ #[getter]
+ pub fn memory_usage(&self) -> u64 {
+ self.inner.memory_usage.as_bytes_u64()
+ }
+
+ /// The total memory of the system, in bytes.
+ #[getter]
+ pub fn total_memory(&self) -> u64 {
+ self.inner.total_memory.as_bytes_u64()
+ }
+
+ /// The available memory of the system, in bytes.
+ #[getter]
+ pub fn available_memory(&self) -> u64 {
+ self.inner.available_memory.as_bytes_u64()
+ }
+
+ /// The run time of the server process.
+ #[getter]
+ #[gen_stub(override_return_type(type_repr = "datetime.timedelta",
imports=("datetime")))]
+ pub fn run_time<'a>(&self, py: Python<'a>) -> PyResult<Bound<'a, PyDelta>>
{
+ iggy_duration_to_py_delta(py, self.inner.run_time)
+ }
+
+ /// The start time of the server process, in microseconds since the Unix
epoch.
+ #[getter]
+ pub fn start_time(&self) -> u64 {
+ self.inner.start_time.as_micros()
+ }
+
+ /// The total number of bytes read.
+ #[getter]
+ pub fn read_bytes(&self) -> u64 {
+ self.inner.read_bytes.as_bytes_u64()
+ }
+
+ /// The total number of bytes written.
+ #[getter]
+ pub fn written_bytes(&self) -> u64 {
+ self.inner.written_bytes.as_bytes_u64()
+ }
+
+ /// The total size of the messages, in bytes.
+ #[getter]
+ pub fn messages_size_bytes(&self) -> u64 {
+ self.inner.messages_size_bytes.as_bytes_u64()
+ }
+
+ /// The total number of streams.
+ #[getter]
+ pub fn streams_count(&self) -> u32 {
+ self.inner.streams_count
+ }
+
+ /// The total number of topics.
+ #[getter]
+ pub fn topics_count(&self) -> u32 {
+ self.inner.topics_count
+ }
+
+ /// The total number of partitions.
+ #[getter]
+ pub fn partitions_count(&self) -> u32 {
+ self.inner.partitions_count
+ }
+
+ /// The total number of segments.
+ #[getter]
+ pub fn segments_count(&self) -> u32 {
+ self.inner.segments_count
+ }
+
+ /// The total number of messages.
+ #[getter]
+ pub fn messages_count(&self) -> u64 {
+ self.inner.messages_count
+ }
+
+ /// The total number of connected clients.
+ #[getter]
+ pub fn clients_count(&self) -> u32 {
+ self.inner.clients_count
+ }
+
+ /// The total number of consumer groups.
+ #[getter]
+ pub fn consumer_groups_count(&self) -> u32 {
+ self.inner.consumer_groups_count
+ }
+
+ /// The name of the host the server runs on.
+ #[getter]
+ pub fn hostname(&self) -> String {
+ self.inner.hostname.clone()
+ }
+
+ /// The name of the operating system.
+ #[getter]
+ pub fn os_name(&self) -> String {
+ self.inner.os_name.clone()
+ }
+
+ /// The version of the operating system.
+ #[getter]
+ pub fn os_version(&self) -> String {
+ self.inner.os_version.clone()
+ }
+
+ /// The version of the kernel.
+ #[getter]
+ pub fn kernel_version(&self) -> String {
+ self.inner.kernel_version.clone()
+ }
+
+ /// The version of the Iggy server.
+ #[getter]
+ pub fn iggy_server_version(&self) -> String {
+ self.inner.iggy_server_version.clone()
+ }
+
+ /// The numeric semantic version of the Iggy server, or `None` when
unknown.
+ /// E.g. 1.2.3 -> 1002003 (major * 1000000 + minor * 1000 + patch).
+ #[getter]
+ #[gen_stub(override_return_type(type_repr = "builtins.int | None"))]
+ pub fn iggy_server_semver(&self) -> Option<u32> {
+ self.inner.iggy_server_semver
+ }
Review Comment:
Add a binding unit test for the `None` branch when a `RustStats` conversion
can be constructed directly.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]