refactor: introduce TokenSource trait and stop copying tokens as String - #30
Merged
Merged
Conversation
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Token propagation and published API compatibility issues remain unresolved.
Get a fresh assessment by requesting another Copilot review.
Review effort: Balanced
Findings: 6
Open (6)
Retain rotations for late token subscribers · New Preserve ExternalTokenSource API compatibility · New Preserve DatumCloudClient token accessor compatibility · New Subscribe retained clients to shared token rotations · New Retain token updates when no subscribers exist · New Preserve constructor compatibility for external token sources · New
What changed in this PR
Introduces a shared token-source abstraction for Datum clients while retaining secrets in SecretString.
Changes:
- Adds
TokenSourceand in-memoryStaticTokenSource. - Refactors clients to share token sources and observe rotations.
- Simplifies tests by removing unnecessary environment-backed helpers.
| File | Description |
|---|---|
| connect-lib/lib/src/test_util.rs | Adds a reusable static token source fixture. |
| connect-lib/lib/src/project_control_plane.rs | Adopts shared secret-bearing token sources. |
| connect-lib/lib/src/lib.rs | Re-exports the new token-source API. |
| connect-lib/lib/src/heartbeat.rs | Uses the in-memory source in tests. |
| connect-lib/lib/src/datum_cloud/token_source.rs | Defines the trait and static implementation. |
| connect-lib/lib/src/datum_cloud/mod.rs | Stores and exposes shared token sources. |
| connect-lib/lib/src/datum_cloud/external_token_source.rs | Implements the trait for external credentials. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Comment on lines
303
to
+304
| let token = self.token_source.token(); | ||
| self.project_control_plane_client_with_token(project_id, &token) | ||
| self.project_control_plane_client_with_token(project_id, token.expose_secret()) |
DatumCloudClient and ProjectControlPlaneClient were bound to the concrete ExternalTokenSource, so every test that needed a client had to write a shell script to disk, mutate DATUM_CREDENTIALS_HELPER and DATUM_SESSION under a global mutex, and exec the script, even when the helper was nothing to do with what was being tested. Add a TokenSource trait (token, watch, force_refresh) and hold it as Arc<dyn TokenSource> in both clients. ExternalTokenSource implements it unchanged. StaticTokenSource is an in-memory implementation for embedders that already hold a token and for tests; it records force_refresh calls so callers can assert a refresh was requested. DatumCloudClient::with_external_token_source stays as a wrapper so the binary is untouched. The trait returns SecretString and the watch channel carries SecretString. ExternalTokenSource previously unwrapped the secret to a plain String at its boundary and broadcast it that way, and ProjectControlPlaneClient stored a plain String copy. Both now keep the secrecy wrapper until the kube config needs the bytes. Tests in datum_cloud/mod.rs, heartbeat.rs and project_control_plane.rs switch to static_token_source() and drop ENV_LOCK. Only external_token_source.rs keeps the script harness, since the helper exec path is what it tests. This commit was created with the assistance of a LLM.
rawkode
force-pushed
the
claude/token-source-trait
branch
from
September 19, 2026 15:58
8dc401c to
ac80b39
Compare
…ch updates, restore public API compatibility
Three real issues in the TokenSource refactor, found by independently
reading the actual code (the patch/commit this was reported against did
not exist anywhere in this environment or its git history, so nothing
was applied blindly).
1. A ProjectControlPlaneClient built via the normal factory
(DatumCloudClient::project_control_plane_client, used by every
production call site) was constructed through ProjectControlPlaneClient
::new(), which always set token_rx: None. With no token_rx, the
background auth-watch task fell back to watching login_state/
auth_update instead, and auth_update_watch() hands out a receiver
whose sender is dropped before it's returned, so that fallback task
exits almost immediately. A client that is *called fresh* on every use
(which is every current internal call site) is unaffected, since it
reads the token source directly at construction time; a *retained*
client from an embedder is not. new() now always shares
datum.token_source().watch(), which every DatumCloudClient has
unconditionally, so a retained client's kube Client rebuilds itself on
rotation instead of only ever seeing the token it was built with.
Covered by a new #[tokio::test] that rotates a StaticTokenSource after
building a client via the normal factory and asserts the retained
client's access_token() updates on its own.
2. tokio::sync::watch::Sender::send returns early without writing the
value at all when there are zero receivers (confirmed against the
vendored tokio 1.51 source: it checks receiver_count() before calling
send_replace, not after). StaticTokenSource::set and
ExternalTokenSource::swap_token both called send and discarded the
Result, so a rotation that happened while nobody was subscribed was
silently lost — not just unnotified, but never stored — and a later
subscriber would see the stale pre-rotation value. Both now use
send_replace, which writes unconditionally. StaticTokenSource also
drops its separate ArcSwap<SecretString> cache in favour of reading
token() from the watch sender directly, so there is exactly one place
the current value lives. Covered by new regression tests in both
token_source.rs and external_token_source.rs.
3. The refactor silently changed several public signatures that
connect-lib (an embeddable library per its own README) exposed before:
DatumCloudClient::token and ProjectControlPlaneClient::access_token
returned String and are now SecretString; ExternalTokenSource::token/
watch/force_refresh moved from inherent methods to a trait impl, which
changes their behavior under plain dot-call syntax without importing
the trait; ProjectControlPlaneClient::new_with_token_source changed
from taking a concrete ExternalTokenSource to Arc<dyn TokenSource>.
All four now keep their old String-returning / no-import-required
inherent form for compatibility, backed by new _secret-suffixed or
token_secret() methods for callers that want the SecretString. Every
inherent compat method reads the same underlying storage as its trait
counterpart (no parallel cache that could drift). watch() needed an
actual second channel, since a watch::Receiver<T> is fixed to one T at
creation; it is written in the same call as the canonical channel, so
it cannot diverge from it. new_with_token_source(ExternalTokenSource)
is restored as a thin wrapper around a new
new_with_shared_token_source(Arc<dyn TokenSource>), which is what the
normal factory path and embedders holding only a trait object now use.
While adding a regression test for (1), the existing
#[cfg(feature = "integration-tests")] tests in project_control_plane.rs
turned out to have never actually run in an environment without a
pre-installed rustls CryptoProvider: kube::Client::try_from panics (not
Err) without one, which the tests' own `if let Ok(pcp) = pcp { .. }`
skip idiom cannot catch, and separately they were plain #[test] functions
calling a path that spawns a tokio task internally, which panics outside
a runtime. Both are fixed: `rustls` (same version and `ring` feature
`bin/` already depends on) is added as a lib dev-dependency and installed
once per test binary, and the affected tests are converted to
#[tokio::test]. All 7 integration-tests-gated tests now pass for real
instead of silently no-op'ing, confirmed with repeated parallel runs.
Verification: cargo fmt --all --check clean; cargo clippy --workspace
--all-targets exit 0 with and without --features integration-tests;
cargo test --workspace 84 passed, 0 failed (77 default + 7 under
--features integration-tests), repeated three times under
--test-threads=8 with no flakiness.
This commit was created with the assistance of a LLM.
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Concurrent token rebuilds can leave the cached token and kube client inconsistent.
Get a fresh assessment by requesting another Copilot review.
Review effort: Balanced
Findings: 1
Open (10)
Atomically update client and token during rebuilds · New Preserve SecretString through the production client factory · New Synchronize environment access in shared-token tests · New Require successful shared-token client construction · New Unwrap construction before asserting project ID · New Require construction before checking token accessors · New Require construction before checking server URL · New Require construction before checking plugin mode · New Unwrap compatibility constructor result · New Subscribe retained clients to shared token rotations
| // ever seeing the token it was constructed with. Every | ||
| // `DatumCloudClient` has a `TokenSource` unconditionally now, so | ||
| // this is always available. | ||
| let token_rx = Some(datum.token_source().watch()); |
Comment on lines
297
to
+298
| let token = self.token_source.token(); | ||
| self.project_control_plane_client_with_token(project_id, &token) | ||
| self.project_control_plane_client_with_token(project_id, token.expose_secret()) |
| let datum = DatumCloudClient::with_external_token_source( | ||
| let client = Self::build_kube_client(&server_url, initial_token.expose_secret())?; | ||
| let datum = DatumCloudClient::with_token_source( | ||
| crate::ApiEnv::from_env_with_host_override(), |
Comment on lines
+252
to
257
| let result = ProjectControlPlaneClient::new_with_shared_token_source( | ||
| "test-project".to_string(), | ||
| "https://api.datum.net/apis/resourcemanager.miloapis.com/v1alpha1/projects/test-project/control-plane".to_string(), | ||
| token_source, | ||
| ); | ||
| let _ = result; |
Comment on lines
265
to
272
| @@ -218,27 +272,30 @@ mod tests { | |||
| } | |||
Comment on lines
+281
to
289
| let pcp = ProjectControlPlaneClient::new_with_shared_token_source( | ||
| "test-project".to_string(), | ||
| "https://api.datum.net/apis/resourcemanager.miloapis.com/v1alpha1/projects/test-project/control-plane".to_string(), | ||
| token_source, | ||
| ); | ||
| if let Ok(pcp) = pcp { | ||
| assert_eq!(pcp.access_token(), expected_token); | ||
| assert_eq!(pcp.access_token_secret().expose_secret(), expected_token); | ||
| } |
Comment on lines
298
to
305
| @@ -248,11 +305,12 @@ mod tests { | |||
| } | |||
Comment on lines
313
to
320
| @@ -261,4 +319,58 @@ mod tests { | |||
| assert!(pcp.datum.is_plugin_mode()); | |||
| } | |||
Comment on lines
+332
to
+339
| let pcp = ProjectControlPlaneClient::new_with_token_source( | ||
| "test-project".to_string(), | ||
| "https://api.datum.net/apis/resourcemanager.miloapis.com/v1alpha1/projects/test-project/control-plane".to_string(), | ||
| token_source, | ||
| ); | ||
| if let Ok(pcp) = pcp { | ||
| assert_eq!(pcp.access_token(), expected); | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.


Summary
DatumCloudClientandProjectControlPlaneClientwere bound to the concreteExternalTokenSource. Every test that needed a client had to write a shell script to disk, mutateDATUM_CREDENTIALS_HELPERandDATUM_SESSIONunder a global mutex, and exec the script, even when the helper had nothing to do with what was under test. That harness would have been the foundation for the tunnel-service tests planned next, so it is replaced first.Changes
datum_cloud/token_source.rswith theTokenSourcetrait (token,watch,force_refresh) andStaticTokenSource, an in-memory implementation for embedders that already hold a token and for tests. It recordsforce_refreshcalls so a test can assert a refresh was requested.DatumCloudClientandProjectControlPlaneClientholdArc<dyn TokenSource>.with_token_sourceis the new constructor;with_external_token_sourcestays as a wrapper so the binary is untouched.DatumCloudClient::token_source()exposes the shared source for sibling clients.ExternalTokenSourceimplements the trait.swap_tokenandstart_refreshare unchanged.SecretStringand the canonical watch channel carriesSecretString.datum_cloud/mod.rs,heartbeat.rsandproject_control_plane.rsusetest_util::static_token_source()and no longer takeENV_LOCK. Onlyexternal_token_source.rskeeps the script harness, since the helper exec path is what it tests.Follow-up fixes (second commit)
Three real issues, found by reading the actual code rather than trusting an unverifiable external report:
ProjectControlPlaneClient::new()— the path every production caller uses viaDatumCloudClient::project_control_plane_client— never wired up a token watch, so a client that outlives one call site continued using its original token even after the source rotated.new()now always sharesdatum.token_source().watch(). Covered by a#[tokio::test]that rotates aStaticTokenSourceand asserts a retained client picks it up on its own.watch::Sender::sendsilently discarded rotations with zero receivers. Confirmed against the vendored tokio source:sendchecksreceiver_count()before writing the value at all, so a rotation that happened while nobody was subscribed was lost, not just unnotified — a later subscriber saw the stale value.StaticTokenSource::setandExternalTokenSource::swap_tokennow usesend_replace, which writes unconditionally.StaticTokenSourcealso drops its separateArcSwapcache in favour of reading directly from the watch sender, so there is exactly one place the current value lives.DatumCloudClient::token,ProjectControlPlaneClient::access_token, andExternalTokenSource::{token, watch, force_refresh}changed return type or moved from inherent methods to a trait impl (which changes behaviour under plain dot-call syntax without importing the trait). All four keep their oldString-returning, no-import-required inherent form, backed by new_secret-suffixed methods for theSecretStringform.ProjectControlPlaneClient::new_with_token_sourceis restored to take a concreteExternalTokenSource; a newnew_with_shared_token_sourcetakesArc<dyn TokenSource>for the normal factory path and embedders holding only a trait object.While adding the regression test for (1), the pre-existing
#[cfg(feature = "integration-tests")]tests inproject_control_plane.rsturned out to have never actually run in an environment without a pre-installed rustlsCryptoProvider:kube::Client::try_frompanics (notErr) without one, which those tests' ownif let Ok(pcp) = pcp { .. }skip idiom can't catch, and separately they were plain#[test]functions calling a path that spawns a tokio task internally, which panics outside a runtime. Both fixed:rustls(matchingbin/'s existing version andringfeature) added as a lib dev-dependency and installed once per test binary, and the affected tests converted to#[tokio::test]. All 7 now pass for real instead of silently no-op'ing.Behaviour
The binary still constructs
ExternalTokenSource::from_env, callsstart_refresh, and passes it towith_external_token_source— unchanged.Verification
cargo fmt --all --checkcargo clippy --workspace --all-targets--features integration-testscargo test --workspace--features integration-tests)--test-threads=8runsgrep -rln ENV_LOCK lib/srcrepo.rs,env.rsandexternal_token_source.rs, each of which tests env vars directly