This is an automated email from the ASF dual-hosted git repository.

hubcio pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git


The following commit(s) were added to refs/heads/master by this push:
     new 71c67e432 feat(server-ng): enforce PAT limits and expose Prometheus 
metrics (#3773)
71c67e432 is described below

commit 71c67e432c3256fb82e4453c05374c028d67f3ed
Author: Hubert Gruszecki <[email protected]>
AuthorDate: Mon Aug 3 08:33:31 2026 +0200

    feat(server-ng): enforce PAT limits and expose Prometheus metrics (#3773)
---
 core/metadata/src/stm/user.rs      |  38 +++++
 core/server-ng/Cargo.toml          |   1 -
 core/server-ng/src/bootstrap.rs    |  18 ++-
 core/server-ng/src/dispatch.rs     |  86 +++++++++--
 core/server-ng/src/http.rs         | 160 ++++++++++++--------
 core/server-ng/src/http/metrics.rs | 297 +++++++++++++++++++++++++++++++++++++
 core/server-ng/src/http/state.rs   |   8 +
 core/server-ng/src/http/submit.rs  |  37 +++--
 core/server-ng/src/pat.rs          | 182 ++++++++++++++++++++++-
 core/simulator/src/replica.rs      |   4 +
 10 files changed, 742 insertions(+), 89 deletions(-)

diff --git a/core/metadata/src/stm/user.rs b/core/metadata/src/stm/user.rs
index 66be115d9..18b7a2e56 100644
--- a/core/metadata/src/stm/user.rs
+++ b/core/metadata/src/stm/user.rs
@@ -205,6 +205,18 @@ impl UsersInner {
         tokens.sort_by(|(left, _), (right, _)| left.cmp(right));
         tokens
     }
+
+    /// Count of live personal access tokens for one user.
+    ///
+    /// Single-map read for the ingress-side create cap; the sibling
+    /// `personal_access_tokens_of` allocates and sorts, too heavy for a
+    /// per-request check.
+    #[must_use]
+    pub fn pat_count_of(&self, user_id: UserId) -> usize {
+        self.personal_access_tokens
+            .get(&user_id)
+            .map_or(0, |user_tokens| user_tokens.len())
+    }
 }
 
 impl Users {
@@ -973,6 +985,32 @@ mod tests {
         assert!(users.personal_access_token_index.is_empty());
     }
 
+    #[test]
+    fn pat_count_of_tracks_create_and_delete() {
+        let mut users = UsersInner::new();
+        assert_eq!(users.pat_count_of(9), 0);
+
+        for (name, hash_byte) in [("first", b'a'), ("second", b'b')] {
+            CreatePersonalAccessTokenRequest {
+                user_id: 9,
+                name: WireName::new(name).unwrap(),
+                expiry: 0,
+                token_hash: [hash_byte; PAT_TOKEN_HASH_BYTES],
+            }
+            .apply(&mut users, IggyTimestamp::now());
+        }
+        assert_eq!(users.pat_count_of(9), 2);
+        assert_eq!(users.pat_count_of(1), 0, "count is scoped per user");
+
+        DeletePersonalAccessTokenRequest {
+            user_id: 9,
+            name: WireName::new("first").unwrap(),
+            only_if_expired: false,
+        }
+        .apply(&mut users, IggyTimestamp::now());
+        assert_eq!(users.pat_count_of(9), 1);
+    }
+
     #[test]
     fn delete_keeps_expiry_index_in_sync() {
         let mut users = UsersInner::new();
diff --git a/core/server-ng/Cargo.toml b/core/server-ng/Cargo.toml
index d710066d5..aa6bed2c6 100644
--- a/core/server-ng/Cargo.toml
+++ b/core/server-ng/Cargo.toml
@@ -57,7 +57,6 @@ ignored = [
     "opentelemetry-semantic-conventions",
     "opentelemetry_sdk",
     "papaya",
-    "prometheus-client",
     "rand",
     "ringbuffer",
     "rmp-serde",
diff --git a/core/server-ng/src/bootstrap.rs b/core/server-ng/src/bootstrap.rs
index 799580274..c1e196a1b 100644
--- a/core/server-ng/src/bootstrap.rs
+++ b/core/server-ng/src/bootstrap.rs
@@ -189,6 +189,7 @@ pub fn wire_shell_handlers<B, MJ, S>(
     bus: &B,
     shard_handle: &ShellShardHandle<B, MJ, S>,
     system_config: Arc<NgSystemConfig>,
+    max_tokens_per_user: u32,
 ) -> ShellHandlers
 where
     B: ShellBus,
@@ -204,6 +205,7 @@ where
             shard_handle,
             &sessions,
             system_config,
+            max_tokens_per_user,
         ),
         on_metadata_submit: make_metadata_submit_handler(shard_handle),
         on_list_clients: make_list_clients_handler(&sessions),
@@ -1312,8 +1314,12 @@ async fn shard_main(
             boot_view,
             shard.plane.metadata().client_table.borrow().client_ids(),
         );
-        let on_client_request =
-            make_client_request_handler(&shard, &sessions, 
Arc::clone(&config.system));
+        let on_client_request = make_client_request_handler(
+            &shard,
+            &sessions,
+            Arc::clone(&config.system),
+            config.personal_access_token.max_tokens_per_user,
+        );
         let (accepted_replica, dialed_replica) =
             make_replica_delegation_fns(Rc::clone(&coord), &bus);
         let accepted_client = make_shard_zero_client_accept_fns(coord, &bus, 
on_client_request);
@@ -1741,7 +1747,12 @@ async fn build_shard_for_thread(
         on_list_clients,
         on_partition_read,
         sessions,
-    } = wire_shell_handlers(&bus, &shard_handle, Arc::clone(&config.system));
+    } = wire_shell_handlers(
+        &bus,
+        &shard_handle,
+        Arc::clone(&config.system),
+        config.personal_access_token.max_tokens_per_user,
+    );
     sessions
         .borrow_mut()
         .set_cluster_roster(Rc::new(build_cluster_roster(
@@ -2608,6 +2619,7 @@ async fn start_tcp_runtime(
             http_addr,
             &config.http,
             config.metadata.clients_table_max,
+            config.personal_access_token.max_tokens_per_user,
             &config.cluster,
             Arc::clone(&config.system),
             self_ports,
diff --git a/core/server-ng/src/dispatch.rs b/core/server-ng/src/dispatch.rs
index a44bd8901..a131fdca3 100644
--- a/core/server-ng/src/dispatch.rs
+++ b/core/server-ng/src/dispatch.rs
@@ -116,6 +116,7 @@ pub(crate) fn make_client_request_handler<B, MJ, S>(
     shard: &Rc<ShellShard<B, MJ, S>>,
     sessions: &Rc<RefCell<SessionManager>>,
     system_config: Arc<NgSystemConfig>,
+    max_tokens_per_user: u32,
 ) -> RequestHandler
 where
     B: ShellBus,
@@ -144,6 +145,7 @@ where
             Rc::clone(&shard),
             Rc::clone(&sessions),
             Arc::clone(&system_config),
+            max_tokens_per_user,
             Rc::clone(&queues),
             Rc::clone(&active),
             client_id,
@@ -444,6 +446,7 @@ pub(crate) fn make_deferred_client_request_handler<B, MJ, 
S>(
     shard_handle: &ShellShardHandle<B, MJ, S>,
     sessions: &Rc<RefCell<SessionManager>>,
     system_config: Arc<NgSystemConfig>,
+    max_tokens_per_user: u32,
 ) -> RequestHandler
 where
     B: ShellBus,
@@ -486,7 +489,16 @@ where
                 active.borrow_mut().remove(&client_id);
                 return;
             };
-            drain_client_requests(shard, sessions, system_config, queues, 
active, client_id).await;
+            drain_client_requests(
+                shard,
+                sessions,
+                system_config,
+                max_tokens_per_user,
+                queues,
+                active,
+                client_id,
+            )
+            .await;
         });
     })
 }
@@ -620,10 +632,12 @@ where
 // An unbound transport sending a replicated frame therefore gets the
 // empty-reply fail-fast below and must log in.
 
+#[allow(clippy::too_many_arguments)]
 fn enqueue_client_request<B, MJ, S>(
     shard: Rc<ShellShard<B, MJ, S>>,
     sessions: Rc<RefCell<SessionManager>>,
     system_config: Arc<NgSystemConfig>,
+    max_tokens_per_user: u32,
     queues: ClientRequestQueues,
     active: ActiveClientRequests,
     client_id: u128,
@@ -645,7 +659,16 @@ fn enqueue_client_request<B, MJ, S>(
 
     let bus = shard.bus.clone();
     bus.spawn(async move {
-        drain_client_requests(shard, sessions, system_config, queues, active, 
client_id).await;
+        drain_client_requests(
+            shard,
+            sessions,
+            system_config,
+            max_tokens_per_user,
+            queues,
+            active,
+            client_id,
+        )
+        .await;
     });
 }
 
@@ -654,6 +677,7 @@ async fn drain_client_requests<B, MJ, S>(
     shard: Rc<ShellShard<B, MJ, S>>,
     sessions: Rc<RefCell<SessionManager>>,
     system_config: Arc<NgSystemConfig>,
+    max_tokens_per_user: u32,
     queues: ClientRequestQueues,
     active: ActiveClientRequests,
     client_id: u128,
@@ -667,7 +691,15 @@ async fn drain_client_requests<B, MJ, S>(
         let Some(message) = pop_next_client_request(&queues, &active, 
client_id) else {
             return;
         };
-        handle_client_request(&shard, &sessions, &system_config, client_id, 
message).await;
+        handle_client_request(
+            &shard,
+            &sessions,
+            &system_config,
+            max_tokens_per_user,
+            client_id,
+            message,
+        )
+        .await;
     }
 }
 
@@ -696,6 +728,7 @@ async fn handle_client_request<B, MJ, S>(
     shard: &Rc<ShellShard<B, MJ, S>>,
     sessions: &Rc<RefCell<SessionManager>>,
     system_config: &Arc<NgSystemConfig>,
+    max_tokens_per_user: u32,
     transport_client_id: u128,
     message: Message<iggy_binary_protocol::GenericHeader>,
 ) where
@@ -856,19 +889,48 @@ async fn handle_client_request<B, MJ, S>(
             new_header.session = bound_session;
         }
     });
-    let (request, raw_pat_token) =
-        match maybe_rewrite_pat_request(sessions, transport_client_id, 
request) {
-            Ok(rewritten) => rewritten,
-            Err(error) => {
+    let (request, raw_pat_token) = match maybe_rewrite_pat_request(
+        sessions,
+        transport_client_id,
+        max_tokens_per_user,
+        |user_id| {
+            shard
+                .plane
+                .metadata()
+                .mux_stm
+                .users()
+                .read(|users| users.pat_count_of(user_id))
+        },
+        request,
+    ) {
+        Ok(rewritten) => rewritten,
+        Err(error) => {
+            // Pre-consensus rejection (token cap reached, malformed body, or a
+            // lost session binding): deny fast with the typed code. A silent
+            // drop would wedge every later request on the connection until the
+            // socket read timeout.
+            warn!(
+                transport_client_id,
+                error = %error,
+                operation = ?header.operation,
+                "denying personal-access-token request"
+            );
+            let commit = current_metadata_commit(shard);
+            let reply = build_deny_reply(&header, transport_client_id, 0, 
commit, error.as_code());
+            if let Err(send_error) = shard
+                .bus
+                .send_to_client(transport_client_id, 
reply.into_generic().into_frozen())
+                .await
+            {
                 warn!(
                     transport_client_id,
-                    error = %error,
-                    operation = ?header.operation,
-                    "dropping request with invalid PAT replication context"
+                    error = %send_error,
+                    "failed to send personal-access-token deny reply"
                 );
-                return;
             }
-        };
+            return;
+        }
+    };
     // Hash raw passwords and, for ChangePassword, verify the current password
     // on the primary before replication; see `crate::users`. Replicas store 
the
     // hash directly. A wrong current password is not denied here: it rides
diff --git a/core/server-ng/src/http.rs b/core/server-ng/src/http.rs
index 44f6f33d2..23946535c 100644
--- a/core/server-ng/src/http.rs
+++ b/core/server-ng/src/http.rs
@@ -29,6 +29,7 @@ mod forward;
 mod handlers;
 mod jwks;
 mod jwt;
+mod metrics;
 mod reads;
 mod reply;
 mod session;
@@ -91,11 +92,13 @@ use crate::server_error::ServerNgError;
 /// Returns [`ServerNgError`] if the JWT manager cannot be built from
 /// `http_config.jwt`, the `[http.cors]` config is invalid, the `[http.tls]`
 /// credentials cannot be loaded, or the listener cannot bind to `addr`.
+#[allow(clippy::too_many_arguments)]
 pub async fn start(
     shard: &Rc<ServerNgShard>,
     addr: SocketAddr,
     http_config: &HttpConfig,
     clients_table_max: usize,
+    max_tokens_per_user: u32,
     cluster: &ClusterConfig,
     system_config: Arc<NgSystemConfig>,
     self_ports: TransportPorts,
@@ -129,6 +132,9 @@ pub async fn start(
         .enabled
         .then(|| configure_cors(&http_config.cors))
         .transpose()?;
+    // Same early-fail rule for the scrape path: axum panics on a route
+    // without a leading '/', so reject it as a config error instead.
+    let metrics_endpoint = metrics::validated_endpoint(&http_config.metrics)?;
     let (listener, bound_addr) = client_listener::tcp::bind(addr).await?;
 
     let state: HttpState = SendWrapper::new(Rc::new(HttpInner {
@@ -154,10 +160,18 @@ pub async fn start(
             metadata_view: 
Arc::new(AtomicU64::new(crate::cluster_meta::METADATA_VIEW_UNKNOWN)),
         },
         max_http_sessions: 
crate::http::session::max_http_sessions(clients_table_max),
+        max_tokens_per_user,
         in_flight_writes: Cell::new(0),
         forward,
+        metrics: metrics::HttpMetrics::init(),
     }));
-    let router = router(state, max_request_size, cors, http_config.web_ui);
+    let router = router(
+        state,
+        max_request_size,
+        cors,
+        metrics_endpoint,
+        http_config.web_ui,
+    );
 
     if http_config.tls.enabled {
         let server_config = 
tls::load_http_tls_server_config(&http_config.tls)?;
@@ -217,66 +231,27 @@ const PING_PATH: &str = "/ping";
 /// 405 if it reached the router; the outermost `CorsLayer` answers it first
 /// instead, and stamps the CORS response headers over every reply, including
 /// the inner layer's `iggy-view`.
+///
+/// `metrics_endpoint`, present only when `[http.metrics]` is enabled, mounts
+/// the public scrape route among the local routes (a scrape must describe the
+/// serving node, never a forwarded primary) and switches on the
+/// request-counting layer.
 fn router(
     state: HttpState,
     max_request_size: usize,
     cors: Option<CorsLayer>,
+    metrics_endpoint: Option<String>,
     web_ui: bool,
 ) -> Router {
     // Cloned for the response layer so `Iggy-View` reads the live view per
     // response; the original `state` is moved into `with_state` below.
     let view_source = state.clone();
-    let forwardable = Router::new()
-        .route("/users", get(get_users).post(create_user))
-        .route(
-            "/users/{user_id}",
-            get(get_user).put(update_user).delete(delete_user),
-        )
-        .route("/users/{user_id}/password", put(change_password))
-        .route("/users/{user_id}/permissions", put(update_permissions))
-        // Static `logout` outranks the `{user_id}` capture in axum's matcher, 
so
-        // `DELETE /users/logout` never misroutes to `delete_user`.
-        .route("/users/logout", delete(logout_user))
-        .route("/streams", get(get_streams).post(create_stream))
-        .route(
-            "/streams/{stream_id}",
-            get(get_stream).put(update_stream).delete(delete_stream),
-        )
-        .route("/streams/{stream_id}/purge", delete(purge_stream))
-        .route(
-            "/streams/{stream_id}/topics",
-            get(get_topics).post(create_topic),
-        )
-        .route(
-            "/streams/{stream_id}/topics/{topic_id}",
-            get(get_topic).put(update_topic).delete(delete_topic),
-        )
-        .route(
-            "/streams/{stream_id}/topics/{topic_id}/purge",
-            delete(purge_topic),
-        )
-        .route(
-            "/streams/{stream_id}/topics/{topic_id}/partitions",
-            post(create_partitions).delete(delete_partitions),
-        )
-        .route(
-            "/streams/{stream_id}/topics/{topic_id}/partitions/{partition_id}",
-            delete(delete_segments),
-        )
-        .route(
-            "/streams/{stream_id}/topics/{topic_id}/consumer-groups",
-            get(get_cgs).post(create_cg),
-        )
-        .route(
-            
"/streams/{stream_id}/topics/{topic_id}/consumer-groups/{group_id}",
-            get(get_cg).delete(delete_cg),
-        )
-        .route("/personal-access-tokens", get(get_pats).post(create_pat))
-        .route("/personal-access-tokens/{name}", delete(delete_pat))
-        .route_layer(from_fn_with_state(
-            state.clone(),
-            forward::forward_to_primary,
-        ));
+    // Counter handle taken before `state` moves below; the counting layer
+    // shares it with the registry the scrape handler encodes.
+    let http_requests = metrics_endpoint
+        .is_some()
+        .then(|| state.metrics.request_counter());
+    let forwardable = forwardable_routes(state.clone());
     // The partition-plane routes (produce, consumer-offset writes) stay local:
     // each partition is its own consensus group whose primary can diverge from
     // the metadata primary, so forwarding them to the metadata primary would
@@ -309,6 +284,10 @@ fn router(
         .route("/cluster/metadata", get(get_cluster_metadata))
         .route("/clients", get(get_clients))
         .route("/clients/{client_id}", get(get_client));
+    let local = match &metrics_endpoint {
+        Some(endpoint) => local.route(endpoint, get(metrics::get_metrics)),
+        None => local,
+    };
     let router = Router::new()
         .merge(forwardable)
         .merge(local)
@@ -316,12 +295,14 @@ fn router(
         .layer(DefaultBodyLimit::max(max_request_size))
         .layer(from_fn(move |request: Request, next: Next| {
             let view_source = view_source.clone();
-            // `/ping` is the sole success route reached without proving a
-            // credential, so it must not leak the cluster-internal view number
-            // (the anon-leak gate). Every other route authenticates before its
-            // handler, so a success/redirect there is an authed flow that may
-            // carry the header; the login routes prove credentials on success.
-            let suppress_view = request.uri().path() == PING_PATH;
+            // `/ping` and the metrics scrape are the success routes reached
+            // without proving a credential, so they must not leak the
+            // cluster-internal view number (the anon-leak gate). Every other
+            // route authenticates before its handler, so a success/redirect
+            // there is an authed flow that may carry the header; the login
+            // routes prove credentials on success.
+            let suppress_view = request.uri().path() == PING_PATH
+                || metrics_endpoint.as_deref() == Some(request.uri().path());
             async move {
                 let response = next.run(request).await;
                 if suppress_view {
@@ -336,10 +317,73 @@ fn router(
         Some(cors) => router.layer(cors),
         None => router,
     };
+    // Outermost, mirroring the legacy server's wiring: every request is
+    // counted, including CORS preflights answered by the layer beneath.
+    let router = match http_requests {
+        Some(http_requests) => router.layer(from_fn(move |request: Request, 
next: Next| {
+            http_requests.inc();
+            async move { next.run(request).await }
+        })),
+        None => router,
+    };
 
     merge_web_ui(router, web_ui)
 }
 
+/// The control-plane route table: every write here commits through the
+/// metadata consensus group, so the shared `forward_to_primary` route layer
+/// (holding `state`) relays them on a follower.
+fn forwardable_routes(state: HttpState) -> Router<HttpState> {
+    Router::new()
+        .route("/users", get(get_users).post(create_user))
+        .route(
+            "/users/{user_id}",
+            get(get_user).put(update_user).delete(delete_user),
+        )
+        .route("/users/{user_id}/password", put(change_password))
+        .route("/users/{user_id}/permissions", put(update_permissions))
+        // Static `logout` outranks the `{user_id}` capture in axum's matcher, 
so
+        // `DELETE /users/logout` never misroutes to `delete_user`.
+        .route("/users/logout", delete(logout_user))
+        .route("/streams", get(get_streams).post(create_stream))
+        .route(
+            "/streams/{stream_id}",
+            get(get_stream).put(update_stream).delete(delete_stream),
+        )
+        .route("/streams/{stream_id}/purge", delete(purge_stream))
+        .route(
+            "/streams/{stream_id}/topics",
+            get(get_topics).post(create_topic),
+        )
+        .route(
+            "/streams/{stream_id}/topics/{topic_id}",
+            get(get_topic).put(update_topic).delete(delete_topic),
+        )
+        .route(
+            "/streams/{stream_id}/topics/{topic_id}/purge",
+            delete(purge_topic),
+        )
+        .route(
+            "/streams/{stream_id}/topics/{topic_id}/partitions",
+            post(create_partitions).delete(delete_partitions),
+        )
+        .route(
+            "/streams/{stream_id}/topics/{topic_id}/partitions/{partition_id}",
+            delete(delete_segments),
+        )
+        .route(
+            "/streams/{stream_id}/topics/{topic_id}/consumer-groups",
+            get(get_cgs).post(create_cg),
+        )
+        .route(
+            
"/streams/{stream_id}/topics/{topic_id}/consumer-groups/{group_id}",
+            get(get_cg).delete(delete_cg),
+        )
+        .route("/personal-access-tokens", get(get_pats).post(create_pat))
+        .route("/personal-access-tokens/{name}", delete(delete_pat))
+        .route_layer(from_fn_with_state(state, forward::forward_to_primary))
+}
+
 /// Stamp `Connection: close` on a 413, which rejects its request body unread.
 ///
 /// hyper decides to close only once it notices the dropped body, and that is
diff --git a/core/server-ng/src/http/metrics.rs 
b/core/server-ng/src/http/metrics.rs
new file mode 100644
index 000000000..bb18a8a9b
--- /dev/null
+++ b/core/server-ng/src/http/metrics.rs
@@ -0,0 +1,297 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! The `[http.metrics]` scrape surface: the legacy-parity metric registry
+//! (entity gauges plus the request counter), its public scrape handler, and
+//! the config gate deciding whether the route is mounted.
+
+use axum::extract::State;
+use configs::http::HttpMetricsConfig;
+use consensus::MetadataHandle;
+use iggy_common::IggyError;
+use metadata::impls::metadata::StreamsFrontend;
+use prometheus_client::encoding::text::encode;
+use prometheus_client::metrics::counter::Counter;
+use prometheus_client::metrics::gauge::Gauge;
+use prometheus_client::registry::Registry;
+use send_wrapper::SendWrapper;
+use tracing::error;
+
+use crate::http::state::HttpState;
+
+/// The legacy server's metric set, registered under the same names and help
+/// texts so existing dashboards and alerts keep working unchanged.
+///
+/// Unlike the legacy server, the entity gauges are not counted at mutation
+/// sites: [`get_metrics`] samples the live state on every scrape, so a gauge
+/// can never drift from the state it describes.
+pub(in crate::http) struct HttpMetrics {
+    registry: Registry,
+    http_requests: Counter,
+    streams: Gauge,
+    topics: Gauge,
+    partitions: Gauge,
+    segments: Gauge,
+    messages: Gauge,
+    users: Gauge,
+    clients: Gauge,
+}
+
+impl HttpMetrics {
+    pub(in crate::http) fn init() -> Self {
+        let mut registry = Registry::default();
+        let http_requests = Counter::default();
+        let streams = Gauge::default();
+        let topics = Gauge::default();
+        let partitions = Gauge::default();
+        let segments = Gauge::default();
+        let messages = Gauge::default();
+        let users = Gauge::default();
+        let clients = Gauge::default();
+        registry.register(
+            "http_requests",
+            "total count of http_requests",
+            http_requests.clone(),
+        );
+        registry.register("streams", "total count of streams", 
streams.clone());
+        registry.register("topics", "total count of topics", topics.clone());
+        registry.register(
+            "partitions",
+            "total count of partitions",
+            partitions.clone(),
+        );
+        registry.register("segments", "total count of segments", 
segments.clone());
+        registry.register("messages", "total count of messages", 
messages.clone());
+        registry.register("users", "total count of users", users.clone());
+        registry.register("clients", "total count of clients", 
clients.clone());
+        Self {
+            registry,
+            http_requests,
+            streams,
+            topics,
+            partitions,
+            segments,
+            messages,
+            users,
+            clients,
+        }
+    }
+
+    /// Handle for the router's request-counting layer. The counter is
+    /// `Arc`-backed, so bumping the clone bumps the registered metric.
+    pub(in crate::http) fn request_counter(&self) -> Counter {
+        self.http_requests.clone()
+    }
+
+    fn formatted_output(&self) -> String {
+        let mut buffer = String::new();
+        if let Err(error) = encode(&mut buffer, &self.registry) {
+            error!(%error, "failed to encode metrics");
+        }
+        buffer
+    }
+}
+
+/// Resolve the configured scrape path: `None` when `[http.metrics]` is
+/// disabled, so the route is never mounted and the endpoint answers 404.
+///
+/// axum's `Router::route` panics on a path without a leading `/`, so an
+/// enabled endpoint missing one is rejected as a configuration error before
+/// the router is assembled.
+///
+/// # Errors
+///
+/// Returns [`IggyError::InvalidConfiguration`] when metrics are enabled and
+/// the endpoint does not start with `/`.
+pub(in crate::http) fn validated_endpoint(
+    config: &HttpMetricsConfig,
+) -> Result<Option<String>, IggyError> {
+    if !config.enabled {
+        return Ok(None);
+    }
+    if !config.endpoint.starts_with('/') {
+        error!(
+            endpoint = %config.endpoint,
+            "invalid http.metrics.endpoint: the path must start with '/'"
+        );
+        return Err(IggyError::InvalidConfiguration);
+    }
+    Ok(Some(config.endpoint.clone()))
+}
+
+/// `GET <http.metrics.endpoint>`: the metric set in prometheus text
+/// exposition. Public - reached without proving a credential, exactly like the
+/// legacy endpoint.
+///
+/// The entity gauges sample the same reads `/stats` serves: the metadata STM
+/// stream and user maps plus the stats-registry rollups, whose partition-plane
+/// increments are relaxed, so scraped values are approximate while writes are
+/// in flight. The clients count scatter-gathers the per-shard session managers
+/// exactly like `GET /clients` and turns partial when a shard misses the reply
+/// deadline.
+pub(in crate::http) async fn get_metrics(State(state): State<HttpState>) -> 
String {
+    let (streams_count, topics_count, partitions_count, segments_count, 
messages_count) = state
+        .shard
+        .plane
+        .metadata()
+        .mux_stm
+        .streams()
+        .read(|streams| {
+            let mut topics_count = 0u64;
+            let mut partitions_count = 0u64;
+            let mut segments_count = 0u64;
+            let mut messages_count = 0u64;
+            for (_, stream) in &streams.items {
+                topics_count = topics_count.saturating_add(stream.topics.len() 
as u64);
+                segments_count = segments_count
+                    
.saturating_add(u64::from(stream.stats.segments_count_inconsistent()));
+                messages_count =
+                    
messages_count.saturating_add(stream.stats.messages_count_inconsistent());
+                for (_, topic) in &stream.topics {
+                    partitions_count =
+                        partitions_count.saturating_add(topic.partitions.len() 
as u64);
+                }
+            }
+            (
+                streams.items.len() as u64,
+                topics_count,
+                partitions_count,
+                segments_count,
+                messages_count,
+            )
+        });
+    let users_count = state
+        .shard
+        .plane
+        .metadata()
+        .mux_stm
+        .users()
+        .read(|users| users.items.len() as u64);
+    let clients_count = 
SendWrapper::new(state.shard.list_all_clients()).await.len() as u64;
+
+    let metrics = &state.metrics;
+    metrics.streams.set(gauge_value(streams_count));
+    metrics.topics.set(gauge_value(topics_count));
+    metrics.partitions.set(gauge_value(partitions_count));
+    metrics.segments.set(gauge_value(segments_count));
+    metrics.messages.set(gauge_value(messages_count));
+    metrics.users.set(gauge_value(users_count));
+    metrics.clients.set(gauge_value(clients_count));
+    metrics.formatted_output()
+}
+
+/// Clamp a count into the gauge's `i64` domain; only `messages` can pass
+/// `i64::MAX` even in theory, the rest are bounded far below it.
+fn gauge_value(count: u64) -> i64 {
+    i64::try_from(count).unwrap_or(i64::MAX)
+}
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+
+    const PARITY_METRIC_NAMES: [&str; 8] = [
+        "http_requests",
+        "streams",
+        "topics",
+        "partitions",
+        "segments",
+        "messages",
+        "users",
+        "clients",
+    ];
+
+    fn metrics_config(enabled: bool, endpoint: &str) -> HttpMetricsConfig {
+        HttpMetricsConfig {
+            enabled,
+            endpoint: endpoint.to_owned(),
+        }
+    }
+
+    #[test]
+    fn formatted_output_exposes_every_parity_metric() {
+        let metrics = HttpMetrics::init();
+        let output = metrics.formatted_output();
+        for name in PARITY_METRIC_NAMES {
+            assert!(
+                output.contains(&format!("# TYPE {name} ")),
+                "metric {name} missing from exposition:\n{output}"
+            );
+        }
+        assert!(
+            output.ends_with("# EOF\n"),
+            "missing exposition trailer:\n{output}"
+        );
+    }
+
+    #[test]
+    fn scraped_values_land_in_the_exposition() {
+        let metrics = HttpMetrics::init();
+        metrics.streams.set(1);
+        metrics.topics.set(2);
+        metrics.partitions.set(3);
+        metrics.segments.set(4);
+        metrics.messages.set(5);
+        metrics.users.set(6);
+        metrics.clients.set(7);
+        metrics.request_counter().inc();
+        let output = metrics.formatted_output();
+        for line in [
+            "streams 1",
+            "topics 2",
+            "partitions 3",
+            "segments 4",
+            "messages 5",
+            "users 6",
+            "clients 7",
+            "http_requests_total 1",
+        ] {
+            assert!(
+                output.contains(&format!("\n{line}\n")),
+                "expected `{line}` in exposition:\n{output}"
+            );
+        }
+    }
+
+    #[test]
+    fn gauge_value_clamps_past_i64_range() {
+        assert_eq!(gauge_value(42), 42);
+        assert_eq!(gauge_value(u64::MAX), i64::MAX);
+    }
+
+    #[test]
+    fn validated_endpoint_disabled_yields_none() {
+        assert!(matches!(
+            validated_endpoint(&metrics_config(false, "/metrics")),
+            Ok(None)
+        ));
+    }
+
+    #[test]
+    fn validated_endpoint_returns_enabled_path() {
+        let endpoint = validated_endpoint(&metrics_config(true, 
"/metrics")).unwrap();
+        assert_eq!(endpoint.as_deref(), Some("/metrics"));
+    }
+
+    #[test]
+    fn validated_endpoint_rejects_missing_leading_slash() {
+        assert!(matches!(
+            validated_endpoint(&metrics_config(true, "metrics")),
+            Err(IggyError::InvalidConfiguration)
+        ));
+    }
+}
diff --git a/core/server-ng/src/http/state.rs b/core/server-ng/src/http/state.rs
index 0c5c16dc9..3971f9aeb 100644
--- a/core/server-ng/src/http/state.rs
+++ b/core/server-ng/src/http/state.rs
@@ -42,6 +42,7 @@ use crate::dispatch::submit_register_on_owner;
 use crate::http::error::{AuthError, ReadError, primary_redirect_location};
 use crate::http::forward::ForwardState;
 use crate::http::jwt::JwtManager;
+use crate::http::metrics::HttpMetrics;
 use crate::http::session::{
     BarrierEntry, FIRST_REQUEST_ID, FRESH_ENTRY_WATERMARK, HttpSession, 
RegistrationBarrier,
     forget_if_same, live_entry, sweep_expired,
@@ -87,6 +88,10 @@ pub(in crate::http) struct HttpInner {
     /// clients out of the shared VSR client table. Read by `resolve_session`
     /// when admitting a fresh session.
     pub(in crate::http) max_http_sessions: usize,
+    /// Configured `[personal_access_token] max_tokens_per_user`, enforced
+    /// pre-consensus by the PAT rewrite inside the write submit (the cap is
+    /// config-derived, so it must never branch inside the replicated apply).
+    pub(in crate::http) max_tokens_per_user: u32,
     /// Awaited partition writes currently in flight across all sessions, gated
     /// by [`MAX_IN_FLIGHT_WRITES_GLOBAL`]. Only [`InFlightWriteGuard`] touches
     /// it, so every admission is paired with exactly one release.
@@ -94,6 +99,9 @@ pub(in crate::http) struct HttpInner {
     /// Follower-to-primary forwarding context: outbound client, scheme, body
     /// bound, and its own in-flight budget (see `http::forward`).
     pub(in crate::http) forward: ForwardState,
+    /// Legacy-parity metric registry served by the scrape route; the router's
+    /// counting layer holds a clone of its request counter.
+    pub(in crate::http) metrics: HttpMetrics,
 }
 
 impl HttpInner {
diff --git a/core/server-ng/src/http/submit.rs 
b/core/server-ng/src/http/submit.rs
index 82a715bac..e50b3373e 100644
--- a/core/server-ng/src/http/submit.rs
+++ b/core/server-ng/src/http/submit.rs
@@ -23,10 +23,12 @@ use std::rc::Rc;
 use std::time::{Duration, Instant};
 
 use bytes::Bytes;
+use consensus::MetadataHandle;
 use futures::channel::oneshot;
 use iggy_binary_protocol::consensus::Command2;
 use iggy_binary_protocol::{GenericHeader, Operation, RequestHeader};
 use iggy_common::IggyError;
+use metadata::impls::metadata::StreamsFrontend;
 use server_common::Message;
 use tracing::warn;
 
@@ -106,10 +108,12 @@ pub(in crate::http) async fn submit_committed(
     let shard = Rc::clone(&state.shard);
     let task_session = Rc::clone(session);
     let body = body.to_vec();
+    let max_tokens_per_user = state.max_tokens_per_user;
     // Detached so a client disconnect cannot abandon the gate mid-submit;
     // the write runs to completion regardless of handler liveness.
     compio::runtime::spawn(async move {
-        let result = submit_gated(&shard, &task_session, operation, 
&body).await;
+        let result =
+            submit_gated(&shard, &task_session, operation, 
max_tokens_per_user, &body).await;
         // A failed send means the handler died mid-await; the submit itself
         // already completed, which is the invariant that matters.
         let _ = result_slot.send(result);
@@ -157,7 +161,8 @@ pub(in crate::http) async fn submit_committed(
 /// re-registers cleanly.
 ///
 /// Both shard-0 request rewrites run here before consensus, mirroring the TCP
-/// dispatch path: the PAT rewrite mints a raw token and replicates only its 
hash
+/// dispatch path: the PAT rewrite enforces the per-user token cap, then mints
+/// a raw token and replicates only its hash
 /// (`CreatePersonalAccessToken`), and the user-password rewrite hashes the new
 /// password and, for `ChangePassword`, strips the current one and verifies it
 /// against the stored hash (`CreateUser` / `ChangePassword`). Both are no-ops
@@ -165,15 +170,17 @@ pub(in crate::http) async fn submit_committed(
 /// write. A third resolution handles `DeleteSegments`, which is not itself a
 /// consensus op: it is rewritten to the metadata `TruncatePartition` that 
commits
 /// the trim (see [`resolve_delete_segments_truncate`]), so the truncate rides
-/// this session's gate id. A rejection (a malformed body or an unresolved
-/// delete-segments namespace) still spends the id like any other exit. A
-/// failed current-password check is not a rejection here: the op still commits
-/// with an emptied new-password sentinel that the replicated apply grades to
-/// `InvalidCredentials` (see `verify_and_rewrite_change_password`).
+/// this session's gate id. A rejection (a malformed body, a caller at the
+/// personal-access-token cap, or an unresolved delete-segments namespace) 
still
+/// spends the id like any other exit. A failed current-password check is not a
+/// rejection here: the op still commits with an emptied new-password sentinel
+/// that the replicated apply grades to `InvalidCredentials` (see
+/// `verify_and_rewrite_change_password`).
 async fn submit_gated(
     shard: &Rc<ServerNgShard>,
     session: &HttpSession,
     operation: Operation,
+    max_tokens_per_user: u32,
     body: &[u8],
 ) -> Result<(RequestHeader, Message<GenericHeader>, Option<String>), 
WriteError> {
     let mut next_request_id = session.gate.lock().await;
@@ -192,8 +199,20 @@ async fn submit_gated(
         request_id,
         body,
     );
-    let (message, raw_token) =
-        rewrite_pat_request_for_user(session.user_id, 
message).map_err(WriteError::Rejected)?;
+    let (message, raw_token) = rewrite_pat_request_for_user(
+        session.user_id,
+        max_tokens_per_user,
+        |user_id| {
+            shard
+                .plane
+                .metadata()
+                .mux_stm
+                .users()
+                .read(|users| users.pat_count_of(user_id))
+        },
+        message,
+    )
+    .map_err(WriteError::Rejected)?;
     let message =
         maybe_rewrite_user_password_request(shard, 
message).map_err(WriteError::Rejected)?;
     // `DeleteSegments` is not itself a consensus op: resolve it to the 
metadata
diff --git a/core/server-ng/src/pat.rs b/core/server-ng/src/pat.rs
index 98781c296..dc22f2000 100644
--- a/core/server-ng/src/pat.rs
+++ b/core/server-ng/src/pat.rs
@@ -40,6 +40,8 @@ use std::rc::Rc;
 pub(crate) fn maybe_rewrite_pat_request(
     sessions: &Rc<RefCell<SessionManager>>,
     transport_client_id: u128,
+    max_tokens_per_user: u32,
+    pat_count_of: impl FnOnce(u32) -> usize,
     request: Message<RequestHeader>,
 ) -> Result<(Message<RequestHeader>, Option<String>), IggyError> {
     let user_id = match request.header().operation {
@@ -49,20 +51,24 @@ pub(crate) fn maybe_rewrite_pat_request(
             .ok_or(IggyError::Unauthenticated)?,
         _ => return Ok((request, None)),
     };
-    rewrite_pat_request_for_user(user_id, request)
+    rewrite_pat_request_for_user(user_id, max_tokens_per_user, pat_count_of, 
request)
 }
 
 /// PAT rewrite for a caller that already holds the authenticated `user_id`.
 ///
 /// The HTTP listener authenticates against its own session table rather than
 /// the transport `SessionManager`, so it resolves the id itself and calls this
-/// directly. `CreatePersonalAccessToken` mints the raw token + hash and hands
-/// the raw token back; `DeletePersonalAccessToken` is rewritten to the
-/// replicated form scoped to `user_id` (`only_if_expired` false); every other
-/// operation passes through unchanged with `None` (the HTTP write core routes
-/// all of its ops here, so the non-PAT arm is a no-op, not `unreachable`).
+/// directly. `CreatePersonalAccessToken` enforces `max_tokens_per_user`
+/// against `pat_count_of` (the caller's committed-STM read, invoked lazily
+/// and only for creates), then mints the raw token + hash and hands the raw
+/// token back; `DeletePersonalAccessToken` is rewritten to the replicated
+/// form scoped to `user_id` (`only_if_expired` false); every other operation
+/// passes through unchanged with `None` (the HTTP write core routes all of
+/// its ops here, so the non-PAT arm is a no-op, not `unreachable`).
 pub(crate) fn rewrite_pat_request_for_user(
     user_id: u32,
+    max_tokens_per_user: u32,
+    pat_count_of: impl FnOnce(u32) -> usize,
     request: Message<RequestHeader>,
 ) -> Result<(Message<RequestHeader>, Option<String>), IggyError> {
     let body = request_body(&request);
@@ -71,6 +77,20 @@ pub(crate) fn rewrite_pat_request_for_user(
         Operation::CreatePersonalAccessToken => {
             let wire = WireCreatePersonalAccessTokenRequest::decode_from(body)
                 .map_err(|_| IggyError::InvalidCommand)?;
+            // The cap is config-derived, so it must be enforced here at
+            // ingress: a config-dependent branch inside the replicated apply
+            // would diverge replicas configured differently. Denying before
+            // consensus also runs it ahead of the in-apply duplicate-name
+            // grading, so an at-limit duplicate name reports the limit,
+            // matching the legacy server's ordering. Reads the local
+            // committed count only: concurrent in-flight creates can briefly
+            // overshoot, the same soft cap legacy has.
+            if pat_count_of(user_id) >= max_tokens_per_user as usize {
+                return Err(IggyError::PersonalAccessTokensLimitReached(
+                    user_id,
+                    max_tokens_per_user,
+                ));
+            }
             // Primary mints the raw token + hash here and ships the hash
             // through consensus. Replicas decode the hash directly. Doing
             // this inside `CreatePersonalAccessTokenRequest::apply` would
@@ -118,3 +138,153 @@ fn mint_pat_raw_and_hash() -> (String, [u8; 64]) {
     out.copy_from_slice(bytes);
     (raw, out)
 }
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+    use iggy_binary_protocol::WireName;
+    use iggy_common::IggyTimestamp;
+    use metadata::stm::StateHandler;
+    use metadata::stm::user::{PAT_TOKEN_HASH_BYTES, UsersInner};
+
+    const USER_ID: u32 = 7;
+    const MAX_TOKENS: u32 = 3;
+    const AT_LIMIT_TOKENS: [(&str, u8); 3] = [("one", b'a'), ("two", b'b'), 
("three", b'c')];
+
+    fn create_pat_request(name: &str) -> Message<RequestHeader> {
+        let header_len = std::mem::size_of::<RequestHeader>();
+        let mut template = Message::<RequestHeader>::new(header_len);
+        let header = bytemuck::checked::try_from_bytes_mut::<RequestHeader>(
+            &mut template.as_mut_slice()[..header_len],
+        )
+        .expect("zeroed bytes are a valid request header");
+        header.operation = Operation::CreatePersonalAccessToken;
+        let body = WireCreatePersonalAccessTokenRequest {
+            name: WireName::new(name).expect("test token name fits the wire 
bound"),
+            expiry: 0,
+        }
+        .to_bytes();
+        rewrite_request_body(&template, &body).expect("test body fits a 
request message")
+    }
+
+    fn users_with_tokens(tokens: &[(&str, u8)]) -> UsersInner {
+        let mut users = UsersInner::new();
+        for (name, hash_byte) in tokens {
+            ReplicatedCreatePersonalAccessTokenRequest {
+                user_id: USER_ID,
+                name: WireName::new(*name).expect("test token name fits the 
wire bound"),
+                expiry: 0,
+                token_hash: [*hash_byte; PAT_TOKEN_HASH_BYTES],
+            }
+            .apply(&mut users, IggyTimestamp::now());
+        }
+        users
+    }
+
+    #[test]
+    fn given_count_below_limit_when_rewrite_should_mint_and_stamp_user() {
+        let (rewritten, raw_token) = rewrite_pat_request_for_user(
+            USER_ID,
+            MAX_TOKENS,
+            |_| (MAX_TOKENS - 1) as usize,
+            create_pat_request("fresh"),
+        )
+        .expect("below-limit create passes the gate");
+
+        assert!(
+            raw_token.is_some(),
+            "create must hand back the one-time raw token"
+        );
+        let replicated =
+            
ReplicatedCreatePersonalAccessTokenRequest::decode_from(request_body(&rewritten))
+                .expect("replicated body round-trips");
+        assert_eq!(replicated.user_id, USER_ID);
+    }
+
+    #[test]
+    fn given_count_at_limit_when_rewrite_should_deny_with_limit_error() {
+        let error = rewrite_pat_request_for_user(
+            USER_ID,
+            MAX_TOKENS,
+            |_| MAX_TOKENS as usize,
+            create_pat_request("one-too-many"),
+        )
+        .expect_err("at-limit create must be denied");
+
+        assert!(
+            matches!(
+                error,
+                IggyError::PersonalAccessTokensLimitReached(USER_ID, 
MAX_TOKENS)
+            ),
+            "expected PersonalAccessTokensLimitReached({USER_ID}, 
{MAX_TOKENS}), got {error:?}"
+        );
+    }
+
+    #[test]
+    fn 
given_at_limit_duplicate_name_when_rewrite_should_report_limit_not_duplicate() {
+        let users = users_with_tokens(&AT_LIMIT_TOKENS);
+
+        // "one" already exists, but the ingress gate must fire before the
+        // in-apply duplicate grading ever could (legacy ordering).
+        let error = rewrite_pat_request_for_user(
+            USER_ID,
+            MAX_TOKENS,
+            |user_id| users.pat_count_of(user_id),
+            create_pat_request("one"),
+        )
+        .expect_err("at-limit create must be denied even for a duplicate 
name");
+
+        assert!(
+            matches!(error, IggyError::PersonalAccessTokensLimitReached(..)),
+            "expected the limit error before duplicate grading, got {error:?}"
+        );
+    }
+
+    #[test]
+    fn given_deletion_freed_a_slot_when_rewrite_should_mint_again() {
+        let mut users = users_with_tokens(&AT_LIMIT_TOKENS);
+        ReplicatedDeletePersonalAccessTokenRequest {
+            user_id: USER_ID,
+            name: WireName::new("one").expect("test token name fits the wire 
bound"),
+            only_if_expired: false,
+        }
+        .apply(&mut users, IggyTimestamp::now());
+
+        let result = rewrite_pat_request_for_user(
+            USER_ID,
+            MAX_TOKENS,
+            |user_id| users.pat_count_of(user_id),
+            create_pat_request("replacement"),
+        );
+
+        assert!(
+            result.is_ok(),
+            "a freed slot must admit the next create, got {:?}",
+            result.err()
+        );
+    }
+
+    #[test]
+    fn given_delete_op_when_at_limit_should_pass_the_gate() {
+        let header_len = std::mem::size_of::<RequestHeader>();
+        let mut template = Message::<RequestHeader>::new(header_len);
+        let header = bytemuck::checked::try_from_bytes_mut::<RequestHeader>(
+            &mut template.as_mut_slice()[..header_len],
+        )
+        .expect("zeroed bytes are a valid request header");
+        header.operation = Operation::DeletePersonalAccessToken;
+        let body = WireDeletePersonalAccessTokenRequest {
+            name: WireName::new("one").expect("test token name fits the wire 
bound"),
+        }
+        .to_bytes();
+        let request =
+            rewrite_request_body(&template, &body).expect("test body fits a 
request message");
+
+        // The cap gates creates only; a delete is how an at-limit user frees a
+        // slot, so gating it too would lock the account at the cap forever.
+        let (_, raw_token) =
+            rewrite_pat_request_for_user(USER_ID, MAX_TOKENS, |_| MAX_TOKENS 
as usize, request)
+                .expect("delete must pass the create cap untouched");
+        assert!(raw_token.is_none(), "delete mints no token");
+    }
+}
diff --git a/core/simulator/src/replica.rs b/core/simulator/src/replica.rs
index 53e24a647..ace319b92 100644
--- a/core/simulator/src/replica.rs
+++ b/core/simulator/src/replica.rs
@@ -17,6 +17,7 @@
 
 use crate::bus::{SharedSimOutbox, SimOutbox};
 use crate::deps::{MemStorage, SimJournal, SimMuxStateMachine, SimSnapshot};
+use configs::server::PersonalAccessTokenConfig;
 use configs::server_ng::NgSystemConfig;
 use consensus::{ConsensusClock, LocalPipeline, Sequencer, VsrConsensus, 
VsrState};
 use iggy_common::IggyByteSize;
@@ -309,6 +310,9 @@ pub fn new_shard(
             &SharedSimOutbox(Rc::clone(bus)),
             &shard_handle,
             Arc::new(NgSystemConfig::default()),
+            // Default-config PAT cap, like the system config above, so sim
+            // ingress admits exactly what a default-configured server does.
+            PersonalAccessTokenConfig::default().max_tokens_per_user,
         )
     } else {
         ShellHandlers::noop()

Reply via email to