all lists on lists.proxmox.com
 help / color / mirror / Atom feed
From: Lukas Wagner <l.wagner@proxmox.com>
To: pdm-devel@lists.proxmox.com
Subject: [PATCH datacenter-manager 18/20] api: subscriptions: use PdmApplication instead of globals
Date: Mon, 17 Aug 2026 14:57:25 +0200	[thread overview]
Message-ID: <20260817125727.454039-19-l.wagner@proxmox.com> (raw)
In-Reply-To: <20260817125727.454039-1-l.wagner@proxmox.com>

Thread State<PdmApplication> through the subscription API handlers and
the daily-update binary, replacing direct calls to
pdm_config::subscriptions/remotes, crate::connection::make_*_client,
and the api_cache free functions with the equivalent methods on the
injected app handle.

Signed-off-by: Lukas Wagner <l.wagner@proxmox.com>
---
 server/src/api/nodes/subscription.rs          |  21 +-
 server/src/api/resources.rs                   |  39 ++-
 server/src/api/subscriptions/mod.rs           | 224 ++++++++++--------
 ...proxmox-datacenter-manager-daily-update.rs |  14 +-
 4 files changed, 176 insertions(+), 122 deletions(-)

diff --git a/server/src/api/nodes/subscription.rs b/server/src/api/nodes/subscription.rs
index 04e4141e..415f77bc 100644
--- a/server/src/api/nodes/subscription.rs
+++ b/server/src/api/nodes/subscription.rs
@@ -2,7 +2,7 @@ use std::collections::{HashMap, HashSet};
 
 use anyhow::{Error, bail};
 
-use proxmox_router::{Permission, Router};
+use proxmox_router::{Permission, Router, State};
 use proxmox_schema::api;
 use proxmox_schema::api_types::NODE_SCHEMA;
 use proxmox_subscription::files::update_apt_auth;
@@ -18,6 +18,7 @@ use pdm_api_types::subscription::{
 use crate::api::resources::{
     fetch_complete_subscription_info_for_remote, get_subscription_info_for_remote,
 };
+use crate::context::PdmApplication;
 
 const PRODUCT_URL: &str = "https://pdm.proxmox.com/faq.html";
 const APT_AUTH_FN: &str = "/etc/apt/auth.conf.d/pdm.conf";
@@ -31,13 +32,14 @@ fn apt_auth_file_opts() -> CreateOptions {
     CreateOptions::new().perm(mode).owner(nix::unistd::ROOT)
 }
 
-async fn get_all_subscription_infos()
--> Result<HashMap<String, (RemoteType, HashMap<String, Option<NodeSubscriptionInfo>>)>, Error> {
+async fn get_all_subscription_infos(
+    app: &PdmApplication,
+) -> Result<HashMap<String, (RemoteType, HashMap<String, Option<NodeSubscriptionInfo>>)>, Error> {
     let (remotes_config, _digest) = pdm_config::remotes::config()?;
 
     let mut subscription_info = HashMap::new();
     for (remote_name, remote) in remotes_config.iter() {
-        match get_subscription_info_for_remote(remote, 24 * 60 * 60).await {
+        match get_subscription_info_for_remote(app, remote, 24 * 60 * 60).await {
             Ok(info) => {
                 subscription_info.insert(remote_name.to_string(), (remote.ty, info));
             }
@@ -118,8 +120,8 @@ fn check_counts(stats: &SubscriptionStatistics) -> Result<(), Error> {
     }
 )]
 /// Return subscription status
-pub async fn get_subscription() -> Result<PdmSubscriptionInfo, Error> {
-    let infos = get_all_subscription_infos().await?;
+pub async fn get_subscription(app: State<PdmApplication>) -> Result<PdmSubscriptionInfo, Error> {
+    let infos = get_all_subscription_infos(&app).await?;
 
     let statistics = count_subscriptions(&infos);
 
@@ -156,8 +158,8 @@ pub async fn get_subscription() -> Result<PdmSubscriptionInfo, Error> {
     },
 )]
 /// Update subscription information
-pub async fn check_subscription() -> Result<(), Error> {
-    let infos = get_all_subscription_infos().await?;
+pub async fn check_subscription(app: State<PdmApplication>) -> Result<(), Error> {
+    let infos = get_all_subscription_infos(&app).await?;
     let stats = count_subscriptions(&infos);
 
     if let Err(err) = check_counts(&stats) {
@@ -189,7 +191,8 @@ pub async fn check_subscription() -> Result<(), Error> {
                         );
                         continue 'outer;
                     };
-                    let node_info = fetch_complete_subscription_info_for_remote(remote).await?;
+                    let node_info =
+                        fetch_complete_subscription_info_for_remote(&app, remote).await?;
                     let Some(info) = node_info.iter().find_map(|(node, val)| {
                         if let Some(info) = val.as_ref() {
                             if info.status == SubscriptionStatus::Active
diff --git a/server/src/api/resources.rs b/server/src/api/resources.rs
index 09d2b88d..168f6b68 100644
--- a/server/src/api/resources.rs
+++ b/server/src/api/resources.rs
@@ -22,7 +22,7 @@ use pdm_api_types::{Authid, CachedLocationInfo, PRIV_RESOURCE_AUDIT, VIEW_ID_SCH
 use pdm_search::{Search, SearchTerm};
 use proxmox_access_control::CachedUserInfo;
 use proxmox_router::{
-    Permission, Router, RpcEnvironment, SubdirMap, http_bail, list_subdirs_api_method,
+    Permission, Router, RpcEnvironment, State, SubdirMap, http_bail, list_subdirs_api_method,
 };
 use proxmox_rrd_api_types::RrdTimeframe;
 use proxmox_schema::{api, parse_boolean};
@@ -31,6 +31,8 @@ use proxmox_subscription::SubscriptionStatus;
 use pve_api_types::{ClusterResource, ClusterResourceNetworkType, ClusterResourceType};
 use serde::{Deserialize, Serialize};
 
+use crate::api_cache::ApiCache;
+use crate::context::PdmApplication;
 use crate::metric_collection::top_entities;
 use crate::{api_cache, connection, views};
 
@@ -686,6 +688,7 @@ pub async fn get_subscription_status(
     verbose: bool,
     view: Option<String>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<Vec<RemoteSubscriptions>, Error> {
     let (remotes_config, _) = pdm_config::remotes::config()?;
 
@@ -711,10 +714,11 @@ pub async fn get_subscription_status(
         }
 
         let view = view.clone();
+        let app_clone = app.clone();
 
         let future = async move {
             let (node_status, error) =
-                match get_subscription_info_for_remote(&remote, max_age).await {
+                match get_subscription_info_for_remote(&app_clone, &remote, max_age).await {
                     Ok(mut node_status) => {
                         node_status.retain(|node, _| {
                             if let Some(view) = &view {
@@ -846,17 +850,21 @@ struct CachedSubscriptionState {
 /// If recent enough cached data is available, it is returned
 /// instead of calling out to the remote.
 pub async fn get_subscription_info_for_remote(
+    app: &PdmApplication,
     remote: &Remote,
     max_age: u64,
 ) -> Result<HashMap<String, Option<NodeSubscriptionInfo>>, Error> {
-    if let Some(cached_subscription) = get_cached_subscription_info(&remote.id, max_age).await? {
+    if let Some(cached_subscription) =
+        get_cached_subscription_info(app.api_cache(), &remote.id, max_age).await?
+    {
         Ok(cached_subscription.node_info)
     } else {
-        let node_info = fetch_remote_subscription_info(remote).await?;
+        let node_info = fetch_remote_subscription_info(app, remote).await?;
         let now = proxmox_time::epoch_i64();
 
         if let Some(existing_state) =
-            update_cached_subscription_info(&remote.id, node_info.clone(), now).await?
+            update_cached_subscription_info(app.api_cache(), &remote.id, node_info.clone(), now)
+                .await?
         {
             // Somebody else updated the cache while we performed the API request,
             // return the more recent data instead of the data we just fetched.
@@ -871,21 +879,24 @@ pub async fn get_subscription_info_for_remote(
 /// The cache will be updated, but never read from. This guarantees that the `serverid` is set, as
 /// it cannot be stored in the cache.
 pub async fn fetch_complete_subscription_info_for_remote(
+    app: &PdmApplication,
     remote: &Remote,
 ) -> Result<HashMap<String, Option<NodeSubscriptionInfo>>, Error> {
-    let node_info = fetch_remote_subscription_info(remote).await?;
+    let node_info = fetch_remote_subscription_info(app, remote).await?;
     let now = proxmox_time::epoch_i64();
-    let _ = update_cached_subscription_info(&remote.id, node_info.clone(), now).await?;
+    let _ = update_cached_subscription_info(app.api_cache(), &remote.id, node_info.clone(), now)
+        .await?;
     Ok(node_info)
 }
 
 const SUBSCRIPTION_STATE_CACHE_KEY: &str = "subscription-state";
 
 async fn get_cached_subscription_info(
+    api_cache: &ApiCache,
     remote: &str,
     max_age: u64,
 ) -> Result<Option<CachedSubscriptionState>, Error> {
-    let cache = api_cache::read_remote(remote).await?;
+    let cache = api_cache.read_remote(remote).await?;
     let subscription_state = cache
         .get_with_max_age(SUBSCRIPTION_STATE_CACHE_KEY, max_age as i64)
         .await
@@ -897,8 +908,8 @@ async fn get_cached_subscription_info(
 }
 
 /// Drop the cached subscription state for a remote, forcing the next read to refetch.
-pub async fn invalidate_subscription_info_for_remote(remote_id: &str) {
-    let cache = match api_cache::write_remote(remote_id).await {
+pub async fn invalidate_subscription_info_for_remote(api_cache: &ApiCache, remote_id: &str) {
+    let cache = match api_cache.write_remote(remote_id).await {
         Ok(cache) => cache,
         Err(err) => {
             log::error!("could not open API cache for {remote_id}: {err}");
@@ -916,11 +927,12 @@ pub async fn invalidate_subscription_info_for_remote(remote_id: &str) {
 /// stored state as `Ok(Some(state))`. If the data that was passed in replaced the cache
 /// entry, `Ok(None)` is returned.
 async fn update_cached_subscription_info(
+    api_cache: &ApiCache,
     remote: &str,
     node_info: HashMap<String, Option<NodeSubscriptionInfo>>,
     now: i64,
 ) -> Result<Option<CachedSubscriptionState>, Error> {
-    let cache = api_cache::write_remote(remote).await?;
+    let cache = api_cache.write_remote(remote).await?;
 
     Ok(cache
         .set_if_newer_with_timestamp(
@@ -967,12 +979,13 @@ fn map_node_subscription_list_to_state(
 
 /// Fetch remote resources and map to pdm-native data types.
 async fn fetch_remote_subscription_info(
+    app: &PdmApplication,
     remote: &Remote,
 ) -> Result<HashMap<String, Option<NodeSubscriptionInfo>>, Error> {
     let mut list = HashMap::new();
     match remote.ty {
         RemoteType::Pve => {
-            let client = connection::make_pve_client(remote)?;
+            let client = app.client_factory().make_pve_client(remote)?;
 
             let nodes = client.list_nodes().await?;
             let mut futures = Vec::with_capacity(nodes.len());
@@ -1013,7 +1026,7 @@ async fn fetch_remote_subscription_info(
             }
         }
         RemoteType::Pbs => {
-            let client = connection::make_pbs_client(remote)?;
+            let client = app.client_factory().make_pbs_client(remote)?;
 
             let info = client.get_subscription().await.ok().map(|info| {
                 let level = SubscriptionLevel::from_key(info.key.as_deref());
diff --git a/server/src/api/subscriptions/mod.rs b/server/src/api/subscriptions/mod.rs
index 81582945..922b43f7 100644
--- a/server/src/api/subscriptions/mod.rs
+++ b/server/src/api/subscriptions/mod.rs
@@ -13,7 +13,8 @@ use proxmox_access_control::CachedUserInfo;
 use proxmox_config_digest::ConfigDigest;
 use proxmox_log::{info, warn};
 use proxmox_router::{
-    Permission, Router, RpcEnvironment, SubdirMap, http_bail, http_err, list_subdirs_api_method,
+    Permission, Router, RpcEnvironment, State, SubdirMap, http_bail, http_err,
+    list_subdirs_api_method,
 };
 use proxmox_schema::api;
 use proxmox_section_config::typed::SectionConfigData;
@@ -33,6 +34,7 @@ use crate::api::remotes::RemoteIterator;
 use crate::api::resources::{
     get_subscription_info_for_remote, invalidate_subscription_info_for_remote,
 };
+use crate::context::PdmApplication;
 
 pub const ROUTER: Router = Router::new()
     .get(&list_subdirs_api_method!(SUBDIRS))
@@ -131,14 +133,18 @@ fn key_not_found(key: &str) -> Error {
 /// additionally gated on per-remote `PRIV_RESOURCE_AUDIT` so that an operator who can audit the
 /// pool but not a specific remote does not learn which keys are pinned to it (and through that,
 /// the existence and rough size of that remote's deployment).
-fn list_keys(rpcenv: &mut dyn RpcEnvironment) -> Result<Vec<SubscriptionKeyEntry>, Error> {
+fn list_keys(
+    rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
+) -> Result<Vec<SubscriptionKeyEntry>, Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
         .context("no authid available")?
         .parse()?;
     let user_info = CachedUserInfo::new()?;
 
-    let (config, digest) = pdm_config::subscriptions::config()?;
+    let (config, digest) = app.subscription_key_config().read()?;
+
     rpcenv["digest"] = digest.to_hex().into();
     Ok(config
         .into_iter()
@@ -193,6 +199,7 @@ async fn add_keys(
     keys: Vec<String>,
     digest: Option<ConfigDigest>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<AddKeysResult, Error> {
     if keys.is_empty() {
         http_bail!(BAD_REQUEST, "no keys provided");
@@ -229,8 +236,8 @@ async fn add_keys(
 
     let added = entries.len() as u32;
     let new_digest = tokio::task::spawn_blocking(move || -> Result<ConfigDigest, Error> {
-        let _lock = pdm_config::subscriptions::lock_config()?;
-        let (mut config, config_digest) = pdm_config::subscriptions::config()?;
+        let _lock = app.subscription_key_config().lock()?;
+        let (mut config, config_digest) = app.subscription_key_config().read()?;
         config_digest.detect_modification(digest.as_ref())?;
 
         // `insert` returns the previous entry when one existed; treat that as the duplicate
@@ -244,7 +251,7 @@ async fn add_keys(
             }
         }
 
-        pdm_config::subscriptions::save_config(&config)
+        app.subscription_key_config().write(&config)
     })
     .await??;
     rpcenv["digest"] = new_digest.to_hex().into();
@@ -270,14 +277,18 @@ async fn add_keys(
 /// Bound entries are hidden from operators who cannot audit the bound remote (mirrors the
 /// `list_keys` filter); the response is the same 404 either way so a probe cannot distinguish
 /// "key exists but you cannot see it" from "key not in pool".
-fn get_key(key: String, rpcenv: &mut dyn RpcEnvironment) -> Result<SubscriptionKeyEntry, Error> {
+fn get_key(
+    key: String,
+    rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
+) -> Result<SubscriptionKeyEntry, Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
         .context("no authid available")?
         .parse()?;
     let user_info = CachedUserInfo::new()?;
 
-    let (config, digest) = pdm_config::subscriptions::config()?;
+    let (config, digest) = app.subscription_key_config().read()?;
     rpcenv["digest"] = digest.to_hex().into();
     let mut entry = config
         .get(&key)
@@ -322,6 +333,7 @@ async fn delete_key(
     key: String,
     digest: Option<ConfigDigest>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<(), Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
@@ -333,7 +345,7 @@ async fn delete_key(
     // operator with only PRIV_SYS_MODIFY should not be able to probe live subscription state on
     // a remote they cannot audit. Read the entry once without the lock for this gate; the
     // authoritative read happens under the spawn_blocking section below.
-    let (pre_config, pre_digest) = pdm_config::subscriptions::config()?;
+    let (pre_config, pre_digest) = app.subscription_key_config().read()?;
     let Some(pre_entry) = pre_config.get(&key) else {
         return Err(key_not_found(&key));
     };
@@ -352,7 +364,7 @@ async fn delete_key(
     let pre_binding = pre_entry.remote.as_deref().zip(pre_entry.node.as_deref());
     // Owned bool so the orphan guard inside spawn_blocking does not borrow `pre_config`.
     let pre_had_binding = pre_binding.is_some();
-    let synced_block = check_synced_assignment_for_unassign(&key, pre_binding).await?;
+    let synced_block = check_synced_assignment_for_unassign(&app, &key, pre_binding).await?;
     drop(pre_config);
 
     // The lock + sync IO runs on a blocking thread so the async runtime is free for other work
@@ -361,8 +373,8 @@ async fn delete_key(
     // cross the boundary; reconstructing it is cheap (it just reads the shared ACL cache).
     let new_digest = tokio::task::spawn_blocking(move || -> Result<ConfigDigest, Error> {
         let user_info = CachedUserInfo::new()?;
-        let _lock = pdm_config::subscriptions::lock_config()?;
-        let (mut config, config_digest) = pdm_config::subscriptions::config()?;
+        let _lock = app.subscription_key_config().lock()?;
+        let (mut config, config_digest) = app.subscription_key_config().read()?;
         config_digest.detect_modification(digest.as_ref())?;
 
         let Some(entry) = config.get(&key) else {
@@ -404,14 +416,14 @@ async fn delete_key(
         // Save the authoritative pool config first: an interrupted remove must not leave a `key`
         // entry whose signed blob is gone. A stale shadow blob with no main entry is benign, as
         // readers do not consult it.
-        let new_digest = pdm_config::subscriptions::save_config(&config)?;
+        let new_digest = app.subscription_key_config().write(&config)?;
         // Best-effort shadow cleanup. The shadow only caches signed info, so a corrupt or
         // otherwise unparseable shadow must not block removing the key from the pool: drop the
         // cached blob when the shadow loads, and leave the orphan entry behind otherwise.
-        match pdm_config::subscriptions::shadow_config() {
+        match app.subscription_key_config().read_shadow() {
             Ok(mut shadow) => {
                 shadow.remove(&key);
-                if let Err(err) = pdm_config::subscriptions::save_shadow(&shadow) {
+                if let Err(err) = app.subscription_key_config().write_shadow(&shadow) {
                     warn!("key '{key}' removed from pool, but updating its shadow failed: {err}");
                 }
             }
@@ -457,6 +469,7 @@ async fn set_assignment(
     node: String,
     digest: Option<ConfigDigest>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<(), Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
@@ -475,7 +488,7 @@ async fn set_assignment(
     // orphans whatever live subscription the old remote still ran. Same shape and same guard
     // as delete_key / clear_assignment; only fires when the binding actually moves (re-set to
     // the same target leaves the OLD binding intact and carries no orphan risk).
-    let (pre_config, pre_digest) = pdm_config::subscriptions::config()?;
+    let (pre_config, pre_digest) = app.subscription_key_config().read()?;
     let pre_entry = pre_config.get(&key);
     let pre_binding = pre_entry.and_then(|e| e.remote.as_deref().zip(e.node.as_deref()));
     let rebind_moves_binding = match pre_binding {
@@ -498,7 +511,7 @@ async fn set_assignment(
     }
     let pre_had_binding = pre_binding.is_some();
     let synced_block = if rebind_moves_binding {
-        check_synced_assignment_for_unassign(&key, pre_binding).await?
+        check_synced_assignment_for_unassign(&app, &key, pre_binding).await?
     } else {
         None
     };
@@ -509,8 +522,8 @@ async fn set_assignment(
     // under the lock.
     let new_digest = tokio::task::spawn_blocking(move || -> Result<ConfigDigest, Error> {
         let user_info = CachedUserInfo::new()?;
-        let _lock = pdm_config::subscriptions::lock_config()?;
-        let (mut config, config_digest) = pdm_config::subscriptions::config()?;
+        let _lock = app.subscription_key_config().lock()?;
+        let (mut config, config_digest) = app.subscription_key_config().read()?;
         config_digest.detect_modification(digest.as_ref())?;
 
         let Some(stored_entry) = config.get(&key).cloned() else {
@@ -559,7 +572,7 @@ async fn set_assignment(
             );
         }
 
-        let (remotes_config, _) = pdm_config::remotes::config()?;
+        let (remotes_config, _) = app.remote_config().read()?;
         let remote_entry = remotes_config
             .get(&remote)
             .ok_or_else(|| http_err!(NOT_FOUND, "remote '{remote}' not found"))?;
@@ -600,7 +613,7 @@ async fn set_assignment(
             entry.pending_clear = false;
         }
 
-        pdm_config::subscriptions::save_config(&config)
+        app.subscription_key_config().write(&config)
     })
     .await??;
     rpcenv["digest"] = new_digest.to_hex().into();
@@ -631,6 +644,7 @@ async fn clear_assignment(
     key: String,
     digest: Option<ConfigDigest>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<(), Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
@@ -641,7 +655,7 @@ async fn clear_assignment(
     // Authorise against the entry's bound remote BEFORE hitting the network. An operator with
     // only PRIV_SYS_MODIFY should not be able to probe live subscription state on a remote
     // they cannot audit. The authoritative re-check happens after the lock below.
-    let (pre_config, pre_digest) = pdm_config::subscriptions::config()?;
+    let (pre_config, pre_digest) = app.subscription_key_config().read()?;
     let pre_entry = pre_config.get(&key);
     if let Some(pre_entry) = pre_entry {
         if let Some(assigned_remote) = pre_entry.remote.as_deref() {
@@ -662,7 +676,7 @@ async fn clear_assignment(
     let pre_binding = pre_entry.and_then(|e| e.remote.as_deref().zip(e.node.as_deref()));
     // Owned bool so the orphan guard inside spawn_blocking does not borrow `pre_config`.
     let pre_had_binding = pre_binding.is_some();
-    let synced_block = check_synced_assignment_for_unassign(&key, pre_binding).await?;
+    let synced_block = check_synced_assignment_for_unassign(&app, &key, pre_binding).await?;
     drop(pre_config);
 
     // The lock + sync IO runs on a blocking thread so the async runtime is free for other work
@@ -671,8 +685,8 @@ async fn clear_assignment(
     // boundary; reconstructing it is cheap (it just reads the shared ACL cache).
     let new_digest = tokio::task::spawn_blocking(move || -> Result<ConfigDigest, Error> {
         let user_info = CachedUserInfo::new()?;
-        let _lock = pdm_config::subscriptions::lock_config()?;
-        let (mut config, config_digest) = pdm_config::subscriptions::config()?;
+        let _lock = app.subscription_key_config().lock()?;
+        let (mut config, config_digest) = app.subscription_key_config().read()?;
         config_digest.detect_modification(digest.as_ref())?;
 
         let Some(stored_entry) = config.get(&key).cloned() else {
@@ -720,7 +734,7 @@ async fn clear_assignment(
         // not re-trigger a stale teardown.
         entry.pending_clear = false;
 
-        pdm_config::subscriptions::save_config(&config)
+        app.subscription_key_config().write(&config)
     })
     .await??;
     rpcenv["digest"] = new_digest.to_hex().into();
@@ -739,17 +753,19 @@ async fn clear_assignment(
 /// parallel rebind between pre-read and here cannot redirect us at a remote the caller has no
 /// AUDIT on.
 async fn check_synced_assignment_for_unassign(
+    app: &PdmApplication,
     key: &str,
     binding: Option<(&str, &str)>,
 ) -> Result<Option<(String, String)>, Error> {
     let Some((prev_remote, prev_node)) = binding else {
         return Ok(None);
     };
-    let (remotes_config, _) = pdm_config::remotes::config()?;
+    let (remotes_config, _) = app.remote_config().read()?;
     let Some(remote_entry) = remotes_config.get(prev_remote) else {
         return Ok(None);
     };
-    let live = match get_subscription_info_for_remote(remote_entry, FRESH_NODE_STATUS_MAX_AGE).await
+    let live = match get_subscription_info_for_remote(app, remote_entry, FRESH_NODE_STATUS_MAX_AGE)
+        .await
     {
         Ok(v) => v,
         Err(_) => return Ok(None),
@@ -767,13 +783,18 @@ async fn check_synced_assignment_for_unassign(
 
 /// Push a single key to its assigned remote node. Operates on a borrowed `Remote` so the
 /// caller can fetch the remotes-config once and reuse it.
-async fn push_key_to_remote(remote: &Remote, key: &str, node_name: &str) -> Result<(), Error> {
+async fn push_key_to_remote(
+    app: &PdmApplication,
+    remote: &Remote,
+    key: &str,
+    node_name: &str,
+) -> Result<(), Error> {
     let product_type =
         ProductType::from_key(key).ok_or_else(|| format_err!("unrecognised key format: {key}"))?;
 
     match product_type {
         ProductType::Pve => {
-            let client = crate::connection::make_pve_client(remote)?;
+            let client = app.client_factory().make_pve_client(remote)?;
             client
                 .set_subscription(
                     node_name,
@@ -784,7 +805,7 @@ async fn push_key_to_remote(remote: &Remote, key: &str, node_name: &str) -> Resu
                 .await?;
         }
         ProductType::Pbs => {
-            let client = crate::connection::make_pbs_client(remote)?;
+            let client = app.client_factory().make_pbs_client(remote)?;
             client
                 .set_subscription(proxmox_subscription::SetSubscription {
                     key: key.to_string(),
@@ -806,17 +827,18 @@ async fn push_key_to_remote(remote: &Remote, key: &str, node_name: &str) -> Resu
 
 /// Tear down a node's subscription via the remote's `/nodes/{node}/subscription` endpoint.
 async fn delete_subscription_on_remote(
+    app: &PdmApplication,
     remote: &Remote,
     product_type: ProductType,
     node_name: &str,
 ) -> Result<(), Error> {
     match product_type {
         ProductType::Pve => {
-            let client = crate::connection::make_pve_client(remote)?;
+            let client = app.client_factory().make_pve_client(remote)?;
             client.delete_subscription(node_name).await?;
         }
         ProductType::Pbs => {
-            let client = crate::connection::make_pbs_client(remote)?;
+            let client = app.client_factory().make_pbs_client(remote)?;
             client.delete_subscription().await?;
         }
         ProductType::Pmg | ProductType::Pom => {
@@ -831,13 +853,14 @@ async fn delete_subscription_on_remote(
 /// Trigger a fresh shop-side subscription check on `remote`/`node` and return once the remote
 /// has stored the result. Equivalent to the per-product "Check" button, just driven through PDM.
 async fn check_subscription_on_remote(
+    app: &PdmApplication,
     remote: &Remote,
     product_type: ProductType,
     node_name: &str,
 ) -> Result<(), Error> {
     match product_type {
         ProductType::Pve => {
-            let client = crate::connection::make_pve_client(remote)?;
+            let client = app.client_factory().make_pve_client(remote)?;
             client
                 .update_subscription(
                     node_name,
@@ -846,7 +869,7 @@ async fn check_subscription_on_remote(
                 .await?;
         }
         ProductType::Pbs => {
-            let client = crate::connection::make_pbs_client(remote)?;
+            let client = app.client_factory().make_pbs_client(remote)?;
             client
                 .check_subscription(proxmox_subscription::UpdateSubscription { force: Some(true) })
                 .await?;
@@ -896,6 +919,7 @@ async fn queue_clear(
     node: String,
     digest: Option<ConfigDigest>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<(), Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
@@ -916,8 +940,8 @@ async fn queue_clear(
     // The lock + sync IO runs on a blocking thread so the async runtime stays free for other
     // work even when /etc/proxmox-datacenter-manager/subscriptions is on slow storage.
     let new_digest = tokio::task::spawn_blocking(move || -> Result<ConfigDigest, Error> {
-        let _lock = pdm_config::subscriptions::lock_config()?;
-        let (mut config, config_digest) = pdm_config::subscriptions::config()?;
+        let _lock = app.subscription_key_config().lock()?;
+        let (mut config, config_digest) = app.subscription_key_config().read()?;
         config_digest.detect_modification(digest.as_ref())?;
 
         let bound_id = config
@@ -943,7 +967,7 @@ async fn queue_clear(
         }
         entry.pending_clear = true;
 
-        pdm_config::subscriptions::save_config(&config)
+        app.subscription_key_config().write(&config)
     })
     .await??;
     rpcenv["digest"] = new_digest.to_hex().into();
@@ -975,6 +999,7 @@ async fn revert_pending_clear(
     node: String,
     digest: Option<ConfigDigest>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<(), Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
@@ -989,8 +1014,8 @@ async fn revert_pending_clear(
     )?;
 
     let new_digest = tokio::task::spawn_blocking(move || -> Result<ConfigDigest, Error> {
-        let _lock = pdm_config::subscriptions::lock_config()?;
-        let (mut config, config_digest) = pdm_config::subscriptions::config()?;
+        let _lock = app.subscription_key_config().lock()?;
+        let (mut config, config_digest) = app.subscription_key_config().read()?;
         config_digest.detect_modification(digest.as_ref())?;
 
         let bound_id = config
@@ -1012,7 +1037,7 @@ async fn revert_pending_clear(
         }
         entry.pending_clear = false;
 
-        pdm_config::subscriptions::save_config(&config)
+        app.subscription_key_config().write(&config)
     })
     .await??;
     rpcenv["digest"] = new_digest.to_hex().into();
@@ -1046,6 +1071,7 @@ async fn check_subscription(
     remote: String,
     node: String,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<(), Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
@@ -1059,7 +1085,7 @@ async fn check_subscription(
         false,
     )?;
 
-    let (remotes_config, _) = pdm_config::remotes::config()?;
+    let (remotes_config, _) = app.remote_config().read()?;
     let remote_entry = remotes_config
         .get(&remote)
         .ok_or_else(|| http_err!(NOT_FOUND, "remote '{remote}' not found"))?;
@@ -1069,10 +1095,10 @@ async fn check_subscription(
         pdm_api_types::remotes::RemoteType::Pbs => ProductType::Pbs,
     };
 
-    check_subscription_on_remote(remote_entry, product_type, &node)
+    check_subscription_on_remote(&app, remote_entry, product_type, &node)
         .await
         .map_err(|err| http_err!(BAD_REQUEST, "check failed on {remote}/{node}: {err}"))?;
-    invalidate_subscription_info_for_remote(&remote).await;
+    invalidate_subscription_info_for_remote(app.api_cache(), &remote).await;
     Ok(())
 }
 
@@ -1115,6 +1141,7 @@ async fn adopt_key(
     node: String,
     digest: Option<ConfigDigest>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<(), Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
@@ -1129,15 +1156,15 @@ async fn adopt_key(
     )?;
 
     // Pre-fetch digest to catch a parallel set_assignment during the live read below.
-    let (_pre_config, pre_digest) = pdm_config::subscriptions::config()?;
+    let (_pre_config, pre_digest) = app.subscription_key_config().read()?;
 
     // Fetch live state before grabbing the config lock so the network call does not pin the
     // lock for the duration of a remote query.
-    let (remotes_config, _) = pdm_config::remotes::config()?;
+    let (remotes_config, _) = app.remote_config().read()?;
     let remote_entry = remotes_config
         .get(&remote)
         .ok_or_else(|| http_err!(NOT_FOUND, "remote '{remote}' not found"))?;
-    let live = get_subscription_info_for_remote(remote_entry, FRESH_NODE_STATUS_MAX_AGE)
+    let live = get_subscription_info_for_remote(&app, remote_entry, FRESH_NODE_STATUS_MAX_AGE)
         .await
         .map_err(|err| {
             http_err!(
@@ -1159,8 +1186,8 @@ async fn adopt_key(
     // The lock + sync IO runs on a blocking thread so the async runtime stays free for other
     // work even when /etc/proxmox-datacenter-manager/subscriptions is on slow storage.
     let new_digest = tokio::task::spawn_blocking(move || -> Result<ConfigDigest, Error> {
-        let _lock = pdm_config::subscriptions::lock_config()?;
-        let (mut config, config_digest) = pdm_config::subscriptions::config()?;
+        let _lock = app.subscription_key_config().lock()?;
+        let (mut config, config_digest) = app.subscription_key_config().read()?;
         config_digest.detect_modification(digest.as_ref())?;
         if config_digest != pre_digest {
             http_bail!(
@@ -1225,7 +1252,7 @@ async fn adopt_key(
             config.insert(live_current_key, entry);
         }
 
-        pdm_config::subscriptions::save_config(&config)
+        app.subscription_key_config().write(&config)
     })
     .await??;
     rpcenv["digest"] = new_digest.to_hex().into();
@@ -1269,6 +1296,7 @@ async fn adopt_key(
 async fn adopt_all(
     digest: Option<ConfigDigest>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<Vec<AdoptedEntry>, Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
@@ -1278,7 +1306,7 @@ async fn adopt_all(
     // Use a fresh node-status snapshot: a cached entry from minutes ago could miss a live
     // subscription that was just installed on a remote, or vice-versa, claim a subscription
     // that has since been removed. Adopting bogus or already-cleared keys would be a footgun.
-    let node_statuses = collect_node_status(FRESH_NODE_STATUS_MAX_AGE, rpcenv).await?;
+    let node_statuses = collect_node_status(&app, FRESH_NODE_STATUS_MAX_AGE, rpcenv).await?;
 
     // Lock + sync IO under spawn_blocking. The closure re-resolves the candidate set under the
     // lock: a parallel admin's Assign / Adopt between the network read above and the lock
@@ -1287,8 +1315,8 @@ async fn adopt_all(
     let (adopted, new_digest_opt) = tokio::task::spawn_blocking(
         move || -> Result<(Vec<AdoptedEntry>, Option<ConfigDigest>), Error> {
             let user_info = CachedUserInfo::new()?;
-            let _lock = pdm_config::subscriptions::lock_config()?;
-            let (mut config, config_digest) = pdm_config::subscriptions::config()?;
+            let _lock = app.subscription_key_config().lock()?;
+            let (mut config, config_digest) = app.subscription_key_config().read()?;
             config_digest.detect_modification(digest.as_ref())?;
 
             let mut adopted: Vec<AdoptedEntry> = Vec::new();
@@ -1374,7 +1402,7 @@ async fn adopt_all(
             let new_digest = if adopted.is_empty() {
                 None
             } else {
-                Some(pdm_config::subscriptions::save_config(&config)?)
+                Some(app.subscription_key_config().write(&config)?)
             };
             Ok((adopted, new_digest))
         },
@@ -1415,14 +1443,16 @@ async fn adopt_all(
 async fn node_status(
     max_age: Option<u64>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<Vec<RemoteNodeStatus>, Error> {
-    collect_node_status(max_age.unwrap_or(PANEL_NODE_STATUS_MAX_AGE), rpcenv).await
+    collect_node_status(&app, max_age.unwrap_or(PANEL_NODE_STATUS_MAX_AGE), rpcenv).await
 }
 
 /// Shared helper: fan out subscription queries to all remotes the caller has audit privilege on,
 /// in parallel, reusing the per-remote API cache via `get_subscription_info_for_remote`.
 /// Joins the results with the key-pool assignment table.
 async fn collect_node_status(
+    app: &PdmApplication,
     max_age: u64,
     rpcenv: &mut dyn RpcEnvironment,
 ) -> Result<Vec<RemoteNodeStatus>, Error> {
@@ -1432,18 +1462,17 @@ async fn collect_node_status(
         .parse()?;
     let user_info = CachedUserInfo::new()?;
 
-    let visible_remotes: Vec<(String, Remote)> =
-        RemoteIterator::new(pdm_config::remotes::instance())?
-            .any_privs(&user_info, &auth_id, PRIV_RESOURCE_AUDIT)
-            .into_iter()
-            .collect();
+    let visible_remotes: Vec<(String, Remote)> = RemoteIterator::new(app.remote_config())?
+        .any_privs(&user_info, &auth_id, PRIV_RESOURCE_AUDIT)
+        .into_iter()
+        .collect();
 
-    let (keys_config, _) = pdm_config::subscriptions::config()?;
+    let (keys_config, _) = app.subscription_key_config().read()?;
 
     // `get_subscription_info_for_remote` re-uses the per-remote API cache so this
     // fan-out is safe to run concurrently.
     let fetch = visible_remotes.iter().map(|(name, remote)| async move {
-        let res = get_subscription_info_for_remote(remote, max_age).await;
+        let res = get_subscription_info_for_remote(app, remote, max_age).await;
         (name.clone(), remote.ty, res)
     });
     let results = join_all(fetch).await;
@@ -1528,9 +1557,12 @@ async fn collect_node_status(
 /// The response carries nested `AutoAssignProposal` data; clients must submit follow-up
 /// `bulk_assign` calls with an `application/json` body, the form-urlencoded path cannot encode
 /// the nested structure.
-async fn auto_assign(rpcenv: &mut dyn RpcEnvironment) -> Result<AutoAssignProposal, Error> {
-    let node_statuses = collect_node_status(FRESH_NODE_STATUS_MAX_AGE, rpcenv).await?;
-    let (config, keys_digest) = pdm_config::subscriptions::config()?;
+async fn auto_assign(
+    rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
+) -> Result<AutoAssignProposal, Error> {
+    let node_statuses = collect_node_status(&app, FRESH_NODE_STATUS_MAX_AGE, rpcenv).await?;
+    let (config, keys_digest) = app.subscription_key_config().read()?;
     let assignments = compute_proposals(&config, &node_statuses);
     Ok(AutoAssignProposal {
         assignments,
@@ -1566,13 +1598,14 @@ async fn auto_assign(rpcenv: &mut dyn RpcEnvironment) -> Result<AutoAssignPropos
 async fn bulk_assign(
     proposal: AutoAssignProposal,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<Vec<ProposedAssignment>, Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
         .context("no authid available")?
         .parse()?;
 
-    let node_statuses = collect_node_status(FRESH_NODE_STATUS_MAX_AGE, rpcenv).await?;
+    let node_statuses = collect_node_status(&app, FRESH_NODE_STATUS_MAX_AGE, rpcenv).await?;
     let live_digest = hash_node_status(&node_statuses);
     if live_digest != proposal.node_status_digest {
         http_bail!(
@@ -1587,10 +1620,10 @@ async fn bulk_assign(
     let (applied, new_digest_opt) = tokio::task::spawn_blocking(
         move || -> Result<(Vec<ProposedAssignment>, Option<ConfigDigest>), Error> {
             let user_info = CachedUserInfo::new()?;
-            let _lock = pdm_config::subscriptions::lock_config()?;
-            let (mut config, config_digest) = pdm_config::subscriptions::config()?;
+            let _lock = app.subscription_key_config().lock()?;
+            let (mut config, config_digest) = app.subscription_key_config().read()?;
             config_digest.detect_modification(Some(&proposal.keys_digest))?;
-            let (remotes_config, _) = pdm_config::remotes::config()?;
+            let (remotes_config, _) = app.remote_config().read()?;
 
             let mut applied = Vec::with_capacity(proposal.assignments.len());
             for p in &proposal.assignments {
@@ -1652,7 +1685,7 @@ async fn bulk_assign(
             let new_digest = if applied.is_empty() {
                 None
             } else {
-                Some(pdm_config::subscriptions::save_config(&config)?)
+                Some(app.subscription_key_config().write(&config)?)
             };
             Ok((applied, new_digest))
         },
@@ -1788,6 +1821,7 @@ fn compute_proposals(
 async fn apply_pending(
     digest: Option<ConfigDigest>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<Option<String>, Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
@@ -1795,11 +1829,11 @@ async fn apply_pending(
         .parse()?;
     let user_info = CachedUserInfo::new()?;
 
-    let (_, config_digest) = pdm_config::subscriptions::config()?;
+    let (_, config_digest) = app.subscription_key_config().read()?;
     config_digest.detect_modification(digest.as_ref())?;
 
-    let node_statuses = collect_node_status(FRESH_NODE_STATUS_MAX_AGE, rpcenv).await?;
-    let pending = compute_pending(&user_info, &auth_id, &node_statuses)?;
+    let node_statuses = collect_node_status(&app, FRESH_NODE_STATUS_MAX_AGE, rpcenv).await?;
+    let pending = compute_pending(&app, &user_info, &auth_id, &node_statuses)?;
 
     if pending.is_empty() {
         return Ok(None);
@@ -1811,7 +1845,7 @@ async fn apply_pending(
         None,
         auth_id.to_string(),
         true,
-        move |_worker| async move { run_apply_pending(worker_auth).await },
+        move |_worker| async move { run_apply_pending(&app, worker_auth).await },
     )?;
 
     Ok(Some(upid))
@@ -1822,12 +1856,12 @@ async fn apply_pending(
 /// The worker re-reads remotes and the pool config so a reassign or removal between the API call
 /// returning a UPID and the worker firing is honoured (pushing the old key to a node after the
 /// operator retracted the assignment was a real footgun).
-async fn run_apply_pending(auth_id: Authid) -> Result<(), Error> {
+async fn run_apply_pending(app: &PdmApplication, auth_id: Authid) -> Result<(), Error> {
     let user_info = CachedUserInfo::new()?;
-    let (remotes_config, _) = pdm_config::remotes::config()?;
+    let (remotes_config, _) = app.remote_config().read()?;
 
-    let node_statuses = collect_status_uncached(&remotes_config).await;
-    let pending = compute_pending(&user_info, &auth_id, &node_statuses)?;
+    let node_statuses = collect_status_uncached(app, &remotes_config).await;
+    let pending = compute_pending(app, &user_info, &auth_id, &node_statuses)?;
 
     if pending.is_empty() {
         info!("apply-pending: nothing to do (state changed since the API call)");
@@ -1845,7 +1879,7 @@ async fn run_apply_pending(auth_id: Authid) -> Result<(), Error> {
         // branch) makes the at-start snapshot stale, and a parallel admin's Discard Pending
         // between worker start and this iteration must cancel a planned op rather than have us
         // execute it against a flag the operator just retracted.
-        let (config, _) = pdm_config::subscriptions::config()?;
+        let (config, _) = app.subscription_key_config().read()?;
         if !pool_assignment_still_valid(&config, &entry) {
             info!(
                 "skipping {}/{}: pool entry changed before worker ran",
@@ -1885,7 +1919,7 @@ async fn run_apply_pending(auth_id: Authid) -> Result<(), Error> {
                     continue;
                 };
                 info!("pushing {redacted} to {}/{}...", entry.remote, entry.node);
-                if let Err(err) = push_key_to_remote(remote, &entry.key, &entry.node).await {
+                if let Err(err) = push_key_to_remote(app, remote, &entry.key, &entry.node).await {
                     warn!(
                         "push of {redacted} to {}/{} failed: {err}",
                         entry.remote, entry.node
@@ -1914,7 +1948,7 @@ async fn run_apply_pending(auth_id: Authid) -> Result<(), Error> {
                         entry.remote, entry.node
                     );
                     if let Err(err) =
-                        delete_subscription_on_remote(remote, product_type, &entry.node).await
+                        delete_subscription_on_remote(app, remote, product_type, &entry.node).await
                     {
                         warn!(
                             "clear of {redacted} on {}/{} failed: {err}",
@@ -1937,9 +1971,10 @@ async fn run_apply_pending(auth_id: Authid) -> Result<(), Error> {
                 let entry_key = entry.key.clone();
                 let entry_remote = entry.remote.clone();
                 let entry_node = entry.node.clone();
+                let app = app.clone();
                 let pool_update = tokio::task::spawn_blocking(move || -> Result<(), Error> {
-                    let _lock = pdm_config::subscriptions::lock_config()?;
-                    let (mut updated, _) = pdm_config::subscriptions::config()?;
+                    let _lock = app.subscription_key_config().lock()?;
+                    let (mut updated, _) = app.subscription_key_config().read()?;
                     if let Some(stored) = updated.get_mut(&entry_key) {
                         if stored.remote.as_deref() == Some(entry_remote.as_str())
                             && stored.node.as_deref() == Some(entry_node.as_str())
@@ -1950,7 +1985,7 @@ async fn run_apply_pending(auth_id: Authid) -> Result<(), Error> {
                         }
                     }
                     // Worker context: no `rpcenv` to set, post-save digest is unused here.
-                    let _ = pdm_config::subscriptions::save_config(&updated)?;
+                    let _ = app.subscription_key_config().write(&updated)?;
                     Ok(())
                 })
                 .await
@@ -1970,7 +2005,7 @@ async fn run_apply_pending(auth_id: Authid) -> Result<(), Error> {
             }
         }
         info!("  success");
-        invalidate_subscription_info_for_remote(&entry.remote).await;
+        invalidate_subscription_info_for_remote(app.api_cache(), &entry.remote).await;
         ok += 1;
     }
 
@@ -2031,6 +2066,7 @@ async fn run_apply_pending(auth_id: Authid) -> Result<(), Error> {
 async fn clear_pending(
     digest: Option<ConfigDigest>,
     rpcenv: &mut dyn RpcEnvironment,
+    app: State<PdmApplication>,
 ) -> Result<ClearPendingResult, Error> {
     let auth_id: Authid = rpcenv
         .get_auth_id()
@@ -2038,11 +2074,11 @@ async fn clear_pending(
         .parse()?;
     let user_info = CachedUserInfo::new()?;
 
-    let (_, pre_digest) = pdm_config::subscriptions::config()?;
+    let (_, pre_digest) = app.subscription_key_config().read()?;
     pre_digest.detect_modification(digest.as_ref())?;
 
-    let node_statuses = collect_node_status(FRESH_NODE_STATUS_MAX_AGE, rpcenv).await?;
-    let pending = compute_pending(&user_info, &auth_id, &node_statuses)?;
+    let node_statuses = collect_node_status(&app, FRESH_NODE_STATUS_MAX_AGE, rpcenv).await?;
+    let pending = compute_pending(&app, &user_info, &auth_id, &node_statuses)?;
 
     if pending.is_empty() {
         return Ok(ClearPendingResult { cleared: 0 });
@@ -2052,8 +2088,8 @@ async fn clear_pending(
     // operations.
     let (cleared, new_digest_opt) =
         tokio::task::spawn_blocking(move || -> Result<(u32, Option<ConfigDigest>), Error> {
-            let _lock = pdm_config::subscriptions::lock_config()?;
-            let (mut config, locked_digest) = pdm_config::subscriptions::config()?;
+            let _lock = app.subscription_key_config().lock()?;
+            let (mut config, locked_digest) = app.subscription_key_config().read()?;
             locked_digest.detect_modification(digest.as_ref())?;
 
             let mut cleared: u32 = 0;
@@ -2087,7 +2123,7 @@ async fn clear_pending(
             }
 
             let new_digest = if cleared > 0 {
-                Some(pdm_config::subscriptions::save_config(&config)?)
+                Some(app.subscription_key_config().write(&config)?)
             } else {
                 None
             };
@@ -2120,11 +2156,12 @@ enum PendingOp {
 }
 
 fn compute_pending(
+    app: &PdmApplication,
     user_info: &CachedUserInfo,
     auth_id: &Authid,
     node_statuses: &[RemoteNodeStatus],
 ) -> Result<Vec<PendingEntry>, Error> {
-    let (config, _) = pdm_config::subscriptions::config()?;
+    let (config, _) = app.subscription_key_config().read()?;
 
     Ok(config
         .iter()
@@ -2182,10 +2219,11 @@ fn pool_assignment_still_valid(
 /// Like [`collect_node_status`] but bypasses the auth filter, for the apply-pending worker
 /// which gates each entry through its own per-remote priv check based on the persisted pool plan.
 async fn collect_status_uncached(
+    app: &PdmApplication,
     remotes_config: &SectionConfigData<Remote>,
 ) -> Vec<RemoteNodeStatus> {
     let fetch = remotes_config.iter().map(|(name, remote)| async move {
-        let res = get_subscription_info_for_remote(remote, FRESH_NODE_STATUS_MAX_AGE).await;
+        let res = get_subscription_info_for_remote(app, remote, FRESH_NODE_STATUS_MAX_AGE).await;
         (name.to_string(), remote.ty, res)
     });
     let results = join_all(fetch).await;
diff --git a/server/src/bin/proxmox-datacenter-manager-daily-update.rs b/server/src/bin/proxmox-datacenter-manager-daily-update.rs
index 314b3399..b5989109 100644
--- a/server/src/bin/proxmox-datacenter-manager-daily-update.rs
+++ b/server/src/bin/proxmox-datacenter-manager-daily-update.rs
@@ -2,11 +2,11 @@ use anyhow::Error;
 use serde_json::json;
 
 //use proxmox_notify::context::pbs::PBS_CONTEXT;
-use proxmox_router::{ApiHandler, RpcEnvironment, cli::*};
+use proxmox_router::{ApiHandler, RpcEnvironment, State, cli::*};
 use proxmox_subscription::SubscriptionStatus;
 use proxmox_sys::fs::CreateOptions;
 
-use server::api;
+use server::{api, context::PdmApplication};
 
 async fn wait_for_local_worker(upid_str: &str) -> Result<(), Error> {
     let upid: pbs_api_types::UPID = upid_str.parse()?;
@@ -22,11 +22,11 @@ async fn wait_for_local_worker(upid_str: &str) -> Result<(), Error> {
 }
 
 /// Daily update
-async fn do_update(rpcenv: &mut dyn RpcEnvironment) -> Result<(), Error> {
-    if let Err(err) = &api::nodes::subscription::check_subscription().await {
+async fn do_update(rpcenv: &mut dyn RpcEnvironment, app: PdmApplication) -> Result<(), Error> {
+    if let Err(err) = &api::nodes::subscription::check_subscription(State(app.clone())).await {
         log::error!("Error checking subscription - {err}");
     }
-    match api::nodes::subscription::get_subscription().await {
+    match api::nodes::subscription::get_subscription(State(app)).await {
         Ok(info) if info.info.status == SubscriptionStatus::Active => {}
         Ok(info) => {
             log::warn!(
@@ -101,9 +101,9 @@ async fn run(rpcenv: &mut dyn RpcEnvironment) -> Result<(), Error> {
     proxmox_product_config::init(pdm_config::api_user()?, pdm_config::priv_user()?);
     proxmox_acme_api::init(pdm_buildcfg::configdir!("/acme"), false)?;
 
-    server::context::init()?;
+    let app = server::context::init()?;
 
-    do_update(rpcenv).await
+    do_update(rpcenv, app).await
 }
 
 fn main() {
-- 
2.47.3





  parent reply	other threads:[~2026-08-17 12:58 UTC|newest]

Thread overview: 28+ messages / expand[flat|nested]  mbox.gz  Atom feed  top
2026-08-17 12:57 [PATCH datacenter-manager/proxmox 00/20] inject application context via API macro for easier integration testing Lukas Wagner
2026-08-17 12:57 ` [PATCH proxmox 01/20] router: introduce shared state Lukas Wagner
2026-08-17 13:26   ` Lukas Wagner
2026-08-20 11:23   ` Lukas Wagner
2026-08-21 13:59   ` Robert Obkircher
2026-08-17 12:57 ` [PATCH proxmox 02/20] rest-server: allow to inject " Lukas Wagner
2026-08-17 12:57 ` [PATCH proxmox 03/20] api-macro: support shared state extraction type Lukas Wagner
2026-08-21 13:59   ` Robert Obkircher
2026-08-17 12:57 ` [PATCH datacenter-manager 04/20] context: promote context to a dir-style module Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 05/20] pdm-config: remotes: rename trait methods to read/write/lock Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 06/20] pdm-config: subscriptions: " Lukas Wagner
2026-08-21 14:00   ` Robert Obkircher
2026-08-17 12:57 ` [PATCH datacenter-manager 07/20] remote iterator: pass remote config reader explicitly Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 08/20] context: introduce a ContextFactory to build application context Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 09/20] context: establish PdmApplication object Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 10/20] context: register PdmApplication in router Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 11/20] parallel fetcher: pass arguments to closure in a single type Lukas Wagner
2026-08-21 14:00   ` Robert Obkircher
2026-08-17 12:57 ` [PATCH datacenter-manager 12/20] parallel fetcher: support a custom client factory Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 13/20] api: sdn: use PdmApplication handle for accessing remotes Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 14/20] tests: add helpers for building API-handler-level integration tests Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 15/20] tests: add example tests for SDN API routes Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 16/20] api-cache: add wrapper type Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 17/20] context: provide api-cache on the app object Lukas Wagner
2026-08-17 12:57 ` Lukas Wagner [this message]
2026-08-17 12:57 ` [PATCH datacenter-manager 19/20] pdm-config: subscriptions: drop unused accessor functions Lukas Wagner
2026-08-17 12:57 ` [PATCH datacenter-manager 20/20] tests: add example tests for remote subscription management Lukas Wagner
2026-08-20 14:54 ` superseded: [PATCH datacenter-manager/proxmox 00/20] inject application context via API macro for easier integration testing Lukas Wagner

Reply instructions:

You may reply publicly to this message via plain-text email
using any one of the following methods:

* Save the following mbox file, import it into your mail client,
  and reply-to-all from there: mbox

  Avoid top-posting and favor interleaved quoting:
  https://en.wikipedia.org/wiki/Posting_style#Interleaved_style

* Reply using the --to, --cc, and --in-reply-to
  switches of git-send-email(1):

  git send-email \
    --in-reply-to=20260817125727.454039-19-l.wagner@proxmox.com \
    --to=l.wagner@proxmox.com \
    --cc=pdm-devel@lists.proxmox.com \
    /path/to/YOUR_REPLY

  https://kernel.org/pub/software/scm/git/docs/git-send-email.html

* If your mail client supports setting the In-Reply-To header
  via mailto: links, try the mailto: link
Be sure your reply has a Subject: header at the top and a blank line before the message body.
This is an external index of several public inboxes,
see mirroring instructions on how to clone and mirror
all data and code used by this external index.
Service provided by Proxmox Server Solutions GmbH | Privacy | Legal