feat(python): expose TCP client configuration - #3776
Conversation
IggyClient accepted only a server address, so auto-login and reconnection tuning were unreachable from Python. Without credentials to replay, the SDK's own session recovery never fires and a dropped session surfaces as Unauthenticated on the next call, leaving the application to hand-roll a connect/login/probe loop. TcpConfig mirrors TcpClientConfig field for field and is accepted by the IggyClient constructor alongside the existing address string. AutoLogin carries the credentials without exposing them back to Python, and TcpReconnectionConfig carries the retry policy. Credentials is re-exported from the SDK prelude because AutoLogin::Enabled cannot be constructed without naming it. Closes apache#3742
Round-trip every field through the getters so a default that drifts from the Rust SDK is caught, and assert that neither the password nor a personal access token comes back out of repr. The auto-login tests are the point of the configuration: a privileged call succeeds without a manual login_user() when credentials are configured, and fails without them.
The existing examples all reach for a connection string, which leaves the new config types undiscoverable. This one configures auto-login and reconnection directly and never calls login_user, so the recovery the credentials unlock is visible: restart the server while it runs and the client picks up where it left off.
The README pointed only at the examples directory, so the configuration surface stayed invisible to anyone reading the package page on PyPI.
A negative timedelta normalizes to negative days plus positive seconds, so the old conversion summed to a negative i32 and cast it to u64, turning interval=timedelta(seconds=-1) into u64::MAX seconds: the config constructed fine and the client then slept forever on reconnect. Days arithmetic also overflowed i32 beyond ~68 years, and the reverse conversion stuffed everything into the seconds argument so such values could not read back. Conversion is now fallible, rejects negative input with ValueError at construction, computes in i64, and splits days on the way out. The AutoCommit conversion becomes TryFrom to carry the error. The boolean constructor defaults were literals in the pyo3 signature, so a change to a Rust default would silently not propagate. They are now Option arguments that fall back to TcpClientConfig::default(), the same way the durations already did.
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## master #3776 +/- ##
============================================
- Coverage 76.40% 76.38% -0.02%
Complexity 1046 1046
============================================
Files 1334 1336 +2
Lines 165228 165452 +224
Branches 137583 137660 +77
============================================
+ Hits 126242 126387 +145
+ Misses 35271 35261 -10
- Partials 3715 3804 +89
🚀 New features to boost your workflow:
|
conftest auto-marked every module as integration, so tests explicitly marked unit could not be selected with -m "not integration" even though they need no server. The auto-mark now skips them. New cases pin the duration boundaries (negative rejected, zero legal, beyond the i32 seconds range round-trips) and the README claim that a connection string and TcpConfig reach the same behavior.
The snippet ended with a top-level await; every other sample in the repo wraps in asyncio.run, so paste-and-run failed on the only snippet a PyPI reader sees first.
2b6c6ee to
b3190f4
Compare
|
/ready |
slbotbm
left a comment
There was a problem hiding this comment.
Looks good. Mostly cosmetic changes. One thing though: in the docs, you are declaring the thrown errors as PyValueError and similar types. These are rust types which the python user will not see. Also, there are references to the rust sdk in public docs. Please remove them. Our modelled users are python users, who would not know anything about rust.
| 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", | ||
| ) | ||
| })?; |
There was a problem hiding this comment.
IggyDuration::as_micros() is a truncating cast (core/common/src/utils/duration.rs:94):
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:
let micros = duration.get_duration().as_micros();The existing try_from guards start doing real work once the input isn't pre-wrapped.
There was a problem hiding this comment.
Fixed in 92048cb by reading the std Duration directly, plus a timedelta(days=999_999_999) round-trip test pinning it.
| 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) |
There was a problem hiding this 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.
There was a problem hiding this comment.
Added in dc180e5. The PR description now states the behavior change.
| /// tls_validate_certificate: Whether to validate the server certificate. | ||
| /// Defaults to validating. |
There was a problem hiding this 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."
A separate example for the new config is not needed. The getting-started producer and consumer now build a TcpConfig with auto-login and reconnection instead of a connection string.
The TLS and nodelay options appear as commented-out fields in the snippet instead of prose, and the auto_login and from_connection_string notes are dropped.
Python users see ValueError and RuntimeError rather than the PyO3 exception names, and the wrapped Rust types are an implementation detail.
Docstrings for methods returning Awaitable[None] said they return Ok(()), which does not exist for a Python caller. State the raised exception instead.
IggyDuration::as_micros() truncates the count to u64, so a duration near timedelta.max wrapped to a wrong value instead of surviving the round trip, and the OverflowError guard below could never fire. Read the std Duration directly to keep the u128.
The negative-duration rejection also changed methods that shipped before this branch, such as create_topic's message_expiry, but only the new config classes had coverage.
The tls_validate_certificate docstring was neutral for a flag that accepts any certificate the server presents.
|
@ethanlin01x could you resolve the conflicts? /author |
|
/ready |
hubcio
left a comment
There was a problem hiding this comment.
two things that don't fit on diff lines:
- pre-existing, separate PR: the six sync methods on
IggyConsumer(consumer.rs:58-92) callblocking_lock()while holding the GIL;consume_messagesholds the same mutex for its whole run, so calling e.g.consumer.name()during consumption deadlocks the interpreter (or panics if called from a sync callback body). untouched by this PR, just surfaced while tracing it. - follow-up issue material: the zero-duration hazards below aren't python-specific.
TcpClient::createis the one choke point allTcpClientConfigconstruction sites funnel through (only 2 of 10 go via the builder), and e.g.--tcp-heartbeat-interval noneon the CLI already maps to zero today (IggyDuration::from_strtreats0/none/disabled/unlimitedthe same). a fence there closes every surface at once.
| .map_err(|e| PyErr::new::<pyo3::exceptions::PyValueError, _>(e.to_string()))?; | ||
| inner.auto_login = auto_login.inner.clone(); | ||
| inner.reconnection = reconnection.inner.clone(); | ||
| inner.heartbeat_interval = heartbeat_interval |
There was a problem hiding this comment.
heartbeat_interval=timedelta(0) passes validation (only negatives are rejected) and lands in the sdk heartbeat loop at core/sdk/src/clients/client.rs:269, which does sleep(heartbeat_interval.get_duration()) in an unbounded loop - zero means pinging as fast as round trips complete, for as long as the client lives. nothing downstream reads zero as "disabled". worth rejecting zero here with the same ValueError shape as negatives.
There was a problem hiding this comment.
Fixed in acd01f6 — zero now raises the same ValueError as a negative value.
| inner: RustTcpClientReconnectionConfig { | ||
| enabled: enabled.unwrap_or(defaults.enabled), | ||
| max_retries, | ||
| interval: interval |
There was a problem hiding this comment.
interval=timedelta(0) combined with the default max_retries=None (unlimited) turns the reconnect loop at core/sdk/src/tcp/tcp_client.rs:434-441 into a tight TcpStream::connect spin with one info! line per attempt while the server is down. zero with bounded retries is a legitimate fast-retry policy, so rejecting just the zero + unlimited combination is enough.
There was a problem hiding this comment.
Fixed in acd01f6. Rejects only the zero + unlimited combination.
| with pytest.raises(ValueError, match="negative"): | ||
| construct(negative) | ||
|
|
||
| def test_zero_interval_is_allowed(self): |
There was a problem hiding this comment.
this asserts the zero-interval value from the reconnect-spin problem is legal, so the test will fight the fix later. reestablish_after is the one duration where zero is genuinely meaningful (skips the cooldown - the guard at tcp_client.rs:394 is if elapsed < interval), so retargeting the test there keeps the coverage.
There was a problem hiding this comment.
Retargeted to reestablish_after in acd01f6, as suggested. Split into three: zero reestablish_after, zero interval with bounded retries, and zero interval with reconnection disabled — the three places zero is meaningful.
| builder = builder.init_retries( | ||
| init_retries, | ||
| py_delta_to_iggy_duration(&init_retry_interval), | ||
| py_delta_to_iggy_duration(&init_retry_interval)?, |
There was a problem hiding this comment.
init_retry_interval=timedelta(0) reaches time::interval(interval.get_duration()) at core/sdk/src/clients/consumer.rs:323, and tokio asserts period must be non-zero - a rust panic that surfaces in python as pyo3_async_runtimes.RustPanic, naming neither the argument nor the class. the interval is constructed unconditionally, so it fires every time, even when the stream and topic already exist. since this line now validates the sign, rejecting zero here too is one line.
There was a problem hiding this comment.
Fixed in acd01f6 — rejected on the same line that validates the sign.
| if let Some(polling_retry_interval) = polling_retry_interval { | ||
| builder = | ||
| builder.polling_retry_interval(py_delta_to_iggy_duration(&polling_retry_interval)) | ||
| builder.polling_retry_interval(py_delta_to_iggy_duration(&polling_retry_interval)?) |
There was a problem hiding this comment.
same zero hole: polling_retry_interval=timedelta(0) becomes the retry sleep in core/sdk/src/clients/consumer.rs:690-699 (while !can_poll || ... { trace!(); sleep(...) }) - no syscall in the loop body, so it burns a core, and with auto_join_consumer_group=False the join flag never flips and the spin is permanent. the field is renamed mid-chain to reconnection_retry_interval (consumer_builder.rs:247), in case you grep for it.
| #[pyclass(from_py_object)] | ||
| #[derive(Clone)] | ||
| pub struct TcpConfig { | ||
| auto_login: AutoLogin, |
There was a problem hiding this comment.
auto_login and reconnection are stored twice - as wrapper fields here and inside inner (both assigned in __new__). nothing can make them diverge (no setters), but it's two copies of one truth a reader has to prove safe, and it keeps an extra live copy of the password (SecretString) per config object. both inner fields are pub, so the getters can rebuild the wrappers on demand (AutoLogin { inner: self.inner.auto_login.clone() }), which also makes impl Default for AutoLogin and the Default derive on TcpReconnectionConfig dead.
There was a problem hiding this comment.
Fixed in acd01f6 — the wrapper fields are gone and the getters rebuild from inner. impl Default for AutoLogin and the Default derive on TcpReconnectionConfig went with them.
| tls_validate_certificate: Option<bool>, | ||
| #[gen_stub(override_type(type_repr = "builtins.bool | None"))] nodelay: Option<bool>, | ||
| ) -> PyResult<Self> { | ||
| let defaults = RustTcpClientConfig::default(); |
There was a problem hiding this comment.
TcpClientConfigBuilder is #[derive(Default)] over TcpClientConfig, so the config coming out of build() already carries every default - this defaults is a second identical construction and the unwrap_or(defaults.x) tails re-apply values that are already there. assigning only on Some (or building one exhaustive struct literal sourcing untouched fields from the builder output) drops ~10 lines. two constraints if you do: keep the trimmed server_address from the builder output (build() trims it), and leave max_retries as the raw pass-through it is - None there is a real user value meaning unlimited, not "unset".
There was a problem hiding this comment.
Fixed in acd01f6. The trimmed address and the raw max_retries pass-through are preserved.
| }; | ||
| let tcp_client = TcpClient::create(config) | ||
| .map_err(|e| PyErr::new::<pyo3::exceptions::PyRuntimeError, _>(e.to_string()))?; | ||
| let client = IggyClientBuilder::new() |
There was a problem hiding this comment.
IggyClientBuilder::new().with_client(..).build() plus the map_err guards an error that can't happen - build() only fails when no client was set. RustIggyClient::new(ClientWrapper::Tcp(tcp_client)) is public, infallible and identical (partitioner/encryptor default to None), so this collapses to one line without the dead error path.
|
|
||
| async def main(): | ||
| args: ArgNamespace = parse_args() | ||
| config = build_config(args) |
There was a problem hiding this comment.
build_config() can now raise ValueError (address validation moved to TcpConfig) and it sits outside the try below - --tcp-server-address 127.0.0.1 (no port) passes the argparse url check and produces a raw traceback. same in producer.py, which has no try at all.
There was a problem hiding this comment.
Fixed in 72025b3 — both examples catch the ValueError and print the message instead of a traceback.
| heartbeat_interval=timedelta(seconds=5), | ||
| # tls_enabled=True, | ||
| # tls_domain="localhost", | ||
| # tls_ca_file="core/certs/iggy_ca_cert.pem", |
There was a problem hiding this comment.
this path is relative to the repo root, but a reader of this readme runs from foreign/python - the same file is ../../core/certs/iggy_ca_cert.pem from a sibling dir (that's what the examples readme uses).
The topic API landed on master while this branch was in review: message_expiry and max_topic_size are now IggyExpiry and MaxTopicSize objects rather than a timedelta and an int, and send_messages returns SendMessagesResponse. Took those, and the negative-expiry test follows the new type. Both sides had added their own timedelta conversions, so the two copies master put in consumer.rs and topic.rs give way to the duration module this branch introduced, which is where topic.rs now imports both directions from. The message naming the topic bound is kept at the call site, since the shared conversion has no way to know what the duration is for.
A zero duration reads as "disabled" nowhere in the client: heartbeat_interval pings for as long as the client lives, a reconnection interval with unlimited retries reconnects in a continuous loop, polling_retry_interval spins without a syscall in the loop body, an AutoCommit interval spins and then floods the server with offset stores, and init_retry_interval panics inside the runtime timer without naming the argument that caused it. Each is now rejected where the sign is already validated, except where zero is meaningful: the cooldown before reestablishing, a bounded fast-retry interval, and any interval on a reconnection policy that is switched off. Building the configuration no longer keeps a second copy of the credentials and the reconnection policy beside the one the transport reads, no longer rebuilds the defaults the builder already produced, and no longer routes through a client builder whose only failure mode cannot happen here. Converting a timedelta now goes through the conversion pyo3 ships, keeping the message that does not name Rust types.
The docstrings still described the surface as it behaved before durations were validated: a negative message_expiry or consumer interval now raises ValueError at the call rather than becoming a near-maximum duration, and a zero consumer interval raises it too. A malformed address reaches the user as ValueError through TcpConfig and as RuntimeError through the string form, which the constructor documented as one error. tls_ca_file is silently ignored unless certificate validation is on, and the default reconnection policy retries forever, so an awaited call never returns while the server is down. get_stream and get_topic still described their result as an Option, the one Rust type name left behind when the docstrings were translated for Python readers.
The repr of a configuration printed five of nine fields, dropping every one a TLS handshake is debugged with, so a config that accepts any certificate read exactly like a validating one. Durations printed in a form no constructor accepts, which cost the repr its one job of being pasteable. Both examples built their configuration outside the error handling, where an address without a port reached the user as a traceback rather than as the message the validation produced.
The maximum-interval test named a u64-microsecond boundary that the interval never crosses; what it covers is the day conversion in the getter. The equivalence test claimed both forms of configuration were equivalent while asserting only that both clients authenticate, which is all the client exposes.
The path was written from the repository root, but the snippet around it is run from foreign/python, where the certificate is two levels up. The examples readme already spells it that way.
A max_retries outside the unsigned 32-bit range reached the caller as OverflowError, raised by the argument conversion before any code here ran, so it named neither the argument nor the range. OverflowError is not a ValueError, so a caller guarding construction the way the getting-started examples do never caught it. The count is now taken wide and narrowed here, where the message can say which argument it is and what it accepts.
|
@hubcio On the two that didn't fit on diff lines, the |
|
/ready |
|
@ethanlin01x you can fix them without creation of issue, just mention that it was found in #3776. |
Which issue does this PR address?
Closes #3742
Rationale
The Python binding accepts only a bare server address, so reconnection and auto-login cannot be configured from Python, which makes the SDK's session recovery (#2880) unreachable: a session dropped by a server restart surfaces as
Unauthenticatedon the next call and applications have to hand-roll connect/login/probe retry loops.What changed?
IggyClient(...)only tookhost:port, withAutoLogin::Disabledhardcoded and the reconnection policy untunable. It now also accepts aTcpConfigmirroring the RustTcpClientConfig(auto_login,reconnection,heartbeat_interval, the TLS options,nodelay), keyword-only, with every unset field falling back to the Rust default instead of a value duplicated in the binding. Durations aredatetime.timedeltavalidated at construction, which also fixes the pre-existing conversion that cast a negative timedelta into a huge unsigned duration. This changes shipped behavior: a negative timedelta passed tocreate_topic/update_topic(message_expiry), theconsumer(...)intervals, orAutoCommit.Interval(...)previously became a near-u64::MAXduration and now raisesValueError. Type names follow the maintainer's guidance in the issue (TcpConfig,TcpReconnectionConfig); scope is TCP only, and the bare-address constructor andfrom_connection_stringare unchanged.Local Execution
AI Usage
login_user().