numinnex commented on code in PR #3776:
URL: https://github.com/apache/iggy/pull/3776#discussion_r3690938920


##########
foreign/python/src/duration.rs:
##########
@@ -0,0 +1,52 @@
+// 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 iggy::prelude::IggyDuration;
+use pyo3::prelude::*;
+use pyo3::types::{PyDelta, PyDeltaAccess};
+use std::time::Duration;
+
+pub fn py_delta_to_iggy_duration(delta: &Py<PyDelta>) -> 
PyResult<IggyDuration> {
+    Python::attach(|py| {
+        let delta = delta.bind(py);
+        // Python normalizes a negative timedelta to negative days plus
+        // non-negative seconds/microseconds, so the sign lives in the sum.
+        let seconds = i64::from(delta.get_days()) * 60 * 60 * 24 + 
i64::from(delta.get_seconds());
+        if seconds < 0 {
+            return Err(PyErr::new::<pyo3::exceptions::PyValueError, _>(
+                "duration must not be negative",
+            ));
+        }
+        let nanos = (delta.get_microseconds() * 1_000) as u32;
+        Ok(IggyDuration::new(Duration::new(seconds as u64, nanos)))
+    })
+}
+
+pub fn iggy_duration_to_py_delta(
+    py: Python<'_>,
+    duration: IggyDuration,
+) -> PyResult<Bound<'_, PyDelta>> {
+    let micros = duration.as_micros();
+    let total_seconds = micros / 1_000_000;
+    let days = i32::try_from(total_seconds / 86_400).map_err(|_| {
+        PyErr::new::<pyo3::exceptions::PyOverflowError, _>(
+            "duration does not fit into a datetime.timedelta",
+        )
+    })?;

Review Comment:
   `IggyDuration::as_micros()` is a truncating cast 
(`core/common/src/utils/duration.rs:94`):
   
   ```rust
   pub fn as_micros(&self) -> u64 {
       self.duration.as_micros() as u64
   }
   ```
   
   `py_delta_to_iggy_duration` accepts `timedelta(days=999_999_999)` — 8.64e13 
seconds, which fits a `Duration` — but that is ~8.64e19 µs against a `u64::MAX` 
of ~1.84e19, so the value wraps here and the getter returns a wrong duration 
rather than raising.
   
   That also leaves the `PyOverflowError` below unreachable: a post-wrap value 
caps at ~213,503 days, well under `i32::MAX`, so the `try_from` can never fail.
   
   Reading the u128 directly fixes both:
   
   ```rust
   let micros = duration.get_duration().as_micros();
   ```
   
   The existing `try_from` guards start doing real work once the input isn't 
pre-wrapped.
   



##########
foreign/python/tests/test_client_config.py:
##########
@@ -0,0 +1,303 @@
+# 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.
+
+"""
+Tests for the TCP client configuration surface.
+
+`TcpConfig`, `TcpReconnectionConfig` and `AutoLogin` mirror the Rust SDK
+types, so most of these assert that a value set from Python survives to the
+getters and that unset fields fall back to the Rust defaults. The last class
+proves the point of the configuration: with `auto_login` set, credentials are
+replayed on connect and no manual `login_user()` is needed.
+"""
+
+from collections.abc import Callable
+from datetime import timedelta
+
+import pytest
+
+from apache_iggy import AutoLogin, IggyClient, TcpConfig, TcpReconnectionConfig
+
+from .utils import get_server_config, wait_for_ping, wait_for_server
+
+
[email protected]
+class TestAutoLogin:
+    """Test the credentials carried into the client."""
+
+    def test_disabled_has_no_username(self):
+        """Test that the disabled variant carries no credentials."""
+        auto_login = AutoLogin.disabled()
+
+        assert auto_login.enabled is False
+        assert auto_login.username is None
+
+    def test_username_password_exposes_username_only(self):
+        """Test that the username is readable back but the password is not."""
+        auto_login = AutoLogin.username_password("iggy", "secret")
+
+        assert auto_login.enabled is True
+        assert auto_login.username == "iggy"
+        assert "secret" not in repr(auto_login)
+
+    def test_personal_access_token_hides_the_token(self):
+        """Test that a token login exposes neither a username nor the token."""
+        auto_login = AutoLogin.personal_access_token("secret-token")
+
+        assert auto_login.enabled is True
+        assert auto_login.username is None
+        assert "secret-token" not in repr(auto_login)
+
+
[email protected]
+class TestTcpReconnectionConfig:
+    """Test the reconnection policy."""
+
+    def test_defaults_match_the_rust_sdk(self):
+        """Test that an unconfigured policy reconnects forever, one second 
apart."""
+        reconnection = TcpReconnectionConfig()
+
+        assert reconnection.enabled is True
+        assert reconnection.max_retries is None
+        assert reconnection.interval == timedelta(seconds=1)
+        assert reconnection.reestablish_after == timedelta(seconds=5)
+
+    def test_every_field_round_trips(self):
+        """Test that each configured field is readable back unchanged."""
+        reconnection = TcpReconnectionConfig(
+            enabled=False,
+            max_retries=10,
+            interval=timedelta(milliseconds=250),
+            reestablish_after=timedelta(seconds=30),
+        )
+
+        assert reconnection.enabled is False
+        assert reconnection.max_retries == 10
+        assert reconnection.interval == timedelta(milliseconds=250)
+        assert reconnection.reestablish_after == timedelta(seconds=30)
+
+    def test_arguments_are_keyword_only(self):
+        """Test that the adjacent flags cannot be passed positionally."""
+        with pytest.raises(TypeError):
+            # pyrefly: ignore  # bad-argument-count
+            TcpReconnectionConfig(True)
+
+    @pytest.mark.parametrize(
+        "construct",
+        [
+            lambda duration: TcpReconnectionConfig(interval=duration),
+            lambda duration: TcpReconnectionConfig(reestablish_after=duration),
+        ],
+        ids=["interval", "reestablish_after"],
+    )
+    @pytest.mark.parametrize(
+        "negative",
+        [timedelta(microseconds=-1), timedelta(seconds=-1), 
timedelta(days=-1)],
+    )
+    def test_negative_duration_is_rejected(
+        self,
+        construct: Callable[[timedelta], TcpReconnectionConfig],
+        negative: timedelta,
+    ):
+        """Test that a negative duration fails at construction, not at 
connect."""
+        with pytest.raises(ValueError, match="negative"):
+            construct(negative)

Review Comment:
   The negative-duration rejection also changes methods that already shipped: 
`create_topic` / `update_topic` (`message_expiry`) and 
`consumer(poll_interval=…, polling_retry_interval=…, init_retry_interval=…)`, 
plus `AutoCommit.Interval(…)`. Those previously accepted a negative `timedelta` 
and produced a near-`u64::MAX` duration; they now raise `ValueError`.
   
   That is the right fix, but every negative-duration case here targets the 
three new classes, so the pre-existing surface has no coverage of the new 
behavior. Worth one case on it, e.g. `create_topic(..., 
message_expiry=timedelta(seconds=-1))` raising `ValueError`.
   
   Also worth a line in the PR description: it currently says the conversion is 
fixed, but not that previously-accepted input now raises.
   



##########
foreign/python/src/config.rs:
##########
@@ -0,0 +1,379 @@
+// 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 iggy::prelude::{
+    AutoLogin as RustAutoLogin, Credentials as RustCredentials,
+    TcpClientConfig as RustTcpClientConfig, TcpClientConfigBuilder,
+    TcpClientReconnectionConfig as RustTcpClientReconnectionConfig,
+};
+use pyo3::prelude::*;
+use pyo3::types::PyDelta;
+use pyo3_stub_gen::derive::{gen_stub_pyclass, gen_stub_pymethods};
+use pyo3_stub_gen::impl_stub_type;
+use secrecy::SecretString;
+use std::sync::Arc;
+
+use crate::duration::{iggy_duration_to_py_delta, py_delta_to_iggy_duration};
+
+/// The credentials replayed by the client every time it (re)connects.
+///
+/// `IggyClient` only recovers a lost session when it has credentials to 
replay,
+/// so a long-running consumer should pass one of the enabled variants.
+#[gen_stub_pyclass]
+#[pyclass(from_py_object)]
+#[derive(Clone)]
+pub struct AutoLogin {
+    pub(crate) inner: RustAutoLogin,
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl AutoLogin {
+    /// No automatic login. `login_user()` must be called by hand after every 
connect.
+    #[staticmethod]
+    fn disabled() -> Self {
+        Self {
+            inner: RustAutoLogin::Disabled,
+        }
+    }
+
+    /// Log in with the given username and password on every connect.
+    #[staticmethod]
+    fn username_password(username: String, password: String) -> Self {
+        Self {
+            inner: RustAutoLogin::Enabled(RustCredentials::UsernamePassword(
+                username,
+                SecretString::from(password),
+            )),
+        }
+    }
+
+    /// Log in with the given personal access token on every connect.
+    #[staticmethod]
+    fn personal_access_token(token: String) -> Self {
+        Self {
+            inner: RustAutoLogin::Enabled(RustCredentials::PersonalAccessToken(
+                SecretString::from(token),
+            )),
+        }
+    }
+
+    /// Whether automatic login is enabled.
+    #[getter]
+    fn enabled(&self) -> bool {
+        matches!(self.inner, RustAutoLogin::Enabled(_))
+    }
+
+    /// The username to log in with, or `None` for the disabled and token 
variants.
+    #[gen_stub(override_return_type(type_repr = "builtins.str | None"))]
+    #[getter]
+    fn username(&self) -> Option<String> {
+        match &self.inner {
+            RustAutoLogin::Enabled(RustCredentials::UsernamePassword(username, 
_)) => {
+                Some(username.clone())
+            }
+            _ => None,
+        }
+    }
+
+    fn __repr__(&self) -> String {
+        match &self.inner {
+            RustAutoLogin::Disabled => "AutoLogin.disabled()".to_owned(),
+            RustAutoLogin::Enabled(RustCredentials::UsernamePassword(username, 
_)) => {
+                format!("AutoLogin.username_password({username:?}, ...)")
+            }
+            RustAutoLogin::Enabled(RustCredentials::PersonalAccessToken(_)) => 
{
+                "AutoLogin.personal_access_token(...)".to_owned()
+            }
+        }
+    }
+}
+
+impl Default for AutoLogin {
+    fn default() -> Self {
+        Self::disabled()
+    }
+}
+
+/// How the TCP client reconnects after the connection to the server is lost.
+#[gen_stub_pyclass]
+#[pyclass(from_py_object)]
+#[derive(Clone, Default)]
+pub struct TcpReconnectionConfig {
+    pub(crate) inner: RustTcpClientReconnectionConfig,
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl TcpReconnectionConfig {
+    /// Constructs a reconnection policy, defaulting every unset field to the
+    /// value the Rust SDK uses.
+    ///
+    /// Args:
+    ///     enabled: Whether to reconnect at all. Defaults to enabled.
+    ///     max_retries: Attempts before giving up, or `None` for unlimited.
+    ///     interval: Delay between attempts. Defaults to 1 second.
+    ///     reestablish_after: Cooldown before reconnecting after a previously
+    ///         successful connection. Defaults to 5 seconds.
+    ///
+    /// Raises:
+    ///     PyValueError: If a duration is negative.
+    #[new]
+    #[pyo3(signature = (*, enabled=None, max_retries=None, interval=None, 
reestablish_after=None))]
+    fn new(
+        #[gen_stub(override_type(type_repr = "builtins.bool | None"))] 
enabled: Option<bool>,
+        #[gen_stub(override_type(type_repr = "builtins.int | None"))] 
max_retries: Option<u32>,
+        #[gen_stub(override_type(type_repr = "datetime.timedelta | None", 
imports=("datetime")))]
+        interval: Option<Py<PyDelta>>,
+        #[gen_stub(override_type(type_repr = "datetime.timedelta | None", 
imports=("datetime")))]
+        reestablish_after: Option<Py<PyDelta>>,
+    ) -> PyResult<Self> {
+        let defaults = RustTcpClientReconnectionConfig::default();
+        Ok(Self {
+            inner: RustTcpClientReconnectionConfig {
+                enabled: enabled.unwrap_or(defaults.enabled),
+                max_retries,
+                interval: interval
+                    .as_ref()
+                    .map(py_delta_to_iggy_duration)
+                    .transpose()?
+                    .unwrap_or(defaults.interval),
+                reestablish_after: reestablish_after
+                    .as_ref()
+                    .map(py_delta_to_iggy_duration)
+                    .transpose()?
+                    .unwrap_or(defaults.reestablish_after),
+            },
+        })
+    }
+
+    #[getter]
+    fn enabled(&self) -> bool {
+        self.inner.enabled
+    }
+
+    #[gen_stub(override_return_type(type_repr = "builtins.int | None"))]
+    #[getter]
+    fn max_retries(&self) -> Option<u32> {
+        self.inner.max_retries
+    }
+
+    #[gen_stub(override_return_type(type_repr = "datetime.timedelta", 
imports=("datetime")))]
+    #[getter]
+    fn interval<'a>(&self, py: Python<'a>) -> PyResult<Bound<'a, PyDelta>> {
+        iggy_duration_to_py_delta(py, self.inner.interval)
+    }
+
+    #[gen_stub(override_return_type(type_repr = "datetime.timedelta", 
imports=("datetime")))]
+    #[getter]
+    fn reestablish_after<'a>(&self, py: Python<'a>) -> PyResult<Bound<'a, 
PyDelta>> {
+        iggy_duration_to_py_delta(py, self.inner.reestablish_after)
+    }
+
+    fn __repr__(&self) -> String {
+        let max_retries = match self.inner.max_retries {
+            Some(max_retries) => max_retries.to_string(),
+            None => "None".to_owned(),
+        };
+        format!(
+            "TcpReconnectionConfig(enabled={}, max_retries={max_retries}, 
interval={}, reestablish_after={})",
+            if self.inner.enabled { "True" } else { "False" },
+            self.inner.interval.as_human_time_string(),
+            self.inner.reestablish_after.as_human_time_string(),
+        )
+    }
+}
+
+/// Configuration for the TCP transport, accepted by `IggyClient(...)`.
+///
+/// Mirrors `TcpClientConfig` in the Rust SDK. Every field is keyword-only and
+/// falls back to the same default the Rust SDK uses.
+#[gen_stub_pyclass]
+#[pyclass(from_py_object)]
+#[derive(Clone)]
+pub struct TcpConfig {
+    auto_login: AutoLogin,
+    reconnection: TcpReconnectionConfig,
+    inner: Arc<RustTcpClientConfig>,
+}
+
+impl TcpConfig {
+    /// The configuration in the shape `TcpClient::create` expects.
+    pub(crate) fn client_config(&self) -> Arc<RustTcpClientConfig> {
+        self.inner.clone()
+    }
+}
+
+#[gen_stub_pymethods]
+#[pymethods]
+impl TcpConfig {
+    /// Constructs a TCP configuration, defaulting every unset field to the 
value
+    /// the Rust SDK uses.
+    ///
+    /// Args:
+    ///     server_address: `host:port` of the Iggy server. Defaults to 
`127.0.0.1:8090`.
+    ///     auto_login: Credentials replayed on every connect. Defaults to 
`AutoLogin.disabled()`.
+    ///     reconnection: Reconnection policy. Defaults to 
`TcpReconnectionConfig()`.
+    ///     heartbeat_interval: Interval of heartbeats sent by the client. 
Defaults to 5 seconds.
+    ///     tls_enabled: Whether to connect over TLS. Defaults to disabled.
+    ///     tls_domain: Domain to validate the certificate against. Empty 
means it is
+    ///         taken from `server_address`.
+    ///     tls_ca_file: Path to the CA file for TLS.
+    ///     tls_validate_certificate: Whether to validate the server 
certificate.
+    ///         Defaults to validating.

Review Comment:
   Mirroring `TcpClientConfigBuilder::with_tls_validate_certificate` is in 
scope, and the connection-string path hardcoding `true` 
(`core/common/src/types/configuration/tcp_config/tcp_client_config.rs:72-73`) 
guards a different case — strings arriving from config files or user input, not 
a config built in code. So no objection to exposing it.
   
   The docstring is neutral for a flag that accepts any certificate the server 
presents, though. One sentence on the cost would help, e.g. "Disabling this 
accepts any certificate the server presents, including self-signed and 
mismatched ones; intended for local development only."
   



-- 
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]

Reply via email to