diff --git a/docs/telemetry_conventions.md b/docs/telemetry_conventions.md
index 1445bf026c..34637f6f55 100644
--- a/docs/telemetry_conventions.md
+++ b/docs/telemetry_conventions.md
@@ -312,3 +312,31 @@ none:
The full picture, and the decision to leave these routes as they are for now, is in
`docs/resource_doc_and_endpoint_consistency_status.md`.
+## 14. Where traffic is coming from (not Telemetry, by design)
+
+Telemetry never names a caller (section 6), so it cannot say who is sending the traffic. That question
+is answered by `code.telemetry.TrafficSources`, kept apart on purpose:
+
+- **Three tables per minute:** the busiest Consumers (200 slots, authenticated requests), the busiest
+ client IP addresses (200 slots, every request), and the busiest pairs of caller and endpoint (500
+ slots; the caller is the Consumer when one was authenticated, otherwise the address). 15 minutes are
+ kept; reads merge the last 1, 5 or 15.
+- **Bounded memory, estimated counts.** Each table is a `HeavyHitters` summary of fixed size (the Space-Saving algorithm, credited below): every key
+ with more than 1/capacity of a minute's requests is guaranteed to be in it, and each count is within
+ its stated `error` of the truth.
+- **Recorded once per request, at the outermost layer** (`Http4sApp.httpApp`), so refused and
+ unmatched requests count too. Inner layers fill in a per-request note (endpoint, Consumer, refusal).
+ Paths no endpoint serves are grouped as one endpoint, `unmatched`.
+- **Never exported to Prometheus, never written anywhere.** Read only through
+ `GET /obp/v7.0.0/management/traffic/top-callers` (Role `CanGetTrafficSources`) and the "Where traffic
+ is coming from" panel of the API Manager Telemetry page, which links an address to the IP penalty form.
+ IP addresses are personal data: they stay in memory for at most 15 minutes.
+
+Credit for the ideas used:
+- Space-Saving: Ahmed Metwally, Divyakant Agrawal and Amr El Abbadi, "Efficient Computation of Frequent
+ and Top-k Elements in Data Streams", ICDT 2005, LNCS 3363, pages 398-412.
+- Its predecessor, the first frequent-items algorithm in bounded space: Jayadev Misra and David Gries,
+ "Finding Repeated Elements", Science of Computer Programming 2(2), 1982, pages 143-152.
+- Merging per-minute summaries into a window: Pankaj K. Agarwal, Graham Cormode, Zengfeng Huang,
+ Jeff M. Phillips, Zhewei Wei and Ke Yi, "Mergeable Summaries", PODS 2012, pages 23-34.
+
diff --git a/obp-api/src/main/scala/code/api/util/ApiRole.scala b/obp-api/src/main/scala/code/api/util/ApiRole.scala
index 9b87e38b36..c590525fe5 100644
--- a/obp-api/src/main/scala/code/api/util/ApiRole.scala
+++ b/obp-api/src/main/scala/code/api/util/ApiRole.scala
@@ -580,6 +580,12 @@ object ApiRole extends MdcLoggable{
case class CanDeleteIpPenalty(requiresBankId: Boolean = false) extends ApiRole
lazy val canDeleteIpPenalty = CanDeleteIpPenalty()
+ // Shows which Consumers and client IP addresses are sending the most traffic to the instance
+ // (TrafficSources). About the instance, so held at the empty bank id. It names Consumers and IP
+ // addresses, which is why it is a Role of its own and not part of CanGetTelemetry.
+ case class CanGetTrafficSources(requiresBankId: Boolean = false) extends ApiRole
+ lazy val canGetTrafficSources = CanGetTrafficSources()
+
case class CanGetSignalStats(requiresBankId: Boolean = false) extends ApiRole
lazy val canGetSignalStats = CanGetSignalStats()
diff --git a/obp-api/src/main/scala/code/api/util/ErrorMessages.scala b/obp-api/src/main/scala/code/api/util/ErrorMessages.scala
index 596a6270cf..acc523e232 100644
--- a/obp-api/src/main/scala/code/api/util/ErrorMessages.scala
+++ b/obp-api/src/main/scala/code/api/util/ErrorMessages.scala
@@ -143,6 +143,7 @@ object ErrorMessages {
val InvalidIpAddress = "OBP-10063: Invalid IP address. Give an IPv4 or IPv6 address, not a host name."
val IpPenaltyAlreadyExists = "OBP-10064: This IP address already has a penalty. Remove it first to change it."
val IpPenaltyNotFound = "OBP-10065: This IP address has no penalty."
+ val InvalidTrafficWindow = "OBP-10067: Invalid window. Use window=1, 5 or 15 (minutes)."
val InvalidIpPenalty = "OBP-10066: Invalid IP penalty. per_minute_limit must be 0 or more, duration_minutes between 1 and 10080 (one week), and reason between 1 and 255 characters."
// Not an error: the text of the X-Rate-Limit-Warning header a self-service endpoint returns in
// shadow mode. SCOPE and LIMIT are replaced at runtime, e.g. "signup" and "5 per hour".
diff --git a/obp-api/src/main/scala/code/api/util/http4s/Http4sApp.scala b/obp-api/src/main/scala/code/api/util/http4s/Http4sApp.scala
index cbfbbe91d0..3a89c068a8 100644
--- a/obp-api/src/main/scala/code/api/util/http4s/Http4sApp.scala
+++ b/obp-api/src/main/scala/code/api/util/http4s/Http4sApp.scala
@@ -203,7 +203,11 @@ object Http4sApp extends MdcLoggable {
// trusted forwarder naming it. Unconditional on purpose — the deployment where the answer is
// least obvious, a proxy forwarding the header over a hop with no client certificate, enables
// no TLS middleware at all, so neither step can live in Http4sServer's mtls.enabled branch.
+ // The traffic note travels with the request; inner layers fill it in (see TrafficSources).
+ val note = new code.telemetry.TrafficSources.Note
val req = CallerCertificate.resolveCaller(Psd2CertIngress.canonicalize(rawReq))
+ .withAttribute(Http4sRequestAttributes.trafficNoteKey, note)
+ val startNanos = System.nanoTime()
// Self-service rate limiting (sign-up, password reset, consent requests, consumer
// registration, lookups, signal channel creation) runs here, before routing, keyed by the
// client IP. In shadow mode it only adds X-Rate-Limit-* headers; in enforce mode a trip
@@ -220,6 +224,12 @@ object Http4sApp extends MdcLoggable {
.withEntity(body.getBytes("UTF-8"))
.withHeaders(Headers(Header.Raw(CIString("Content-Type"), "application/json; charset=utf-8"))))))
}
+ .flatTap(resp => IO {
+ // Where the traffic is coming from: every request, served, refused or unmatched, once.
+ try code.telemetry.TrafficSources.record(note, Http4sCallContextBuilder.clientIp(req), resp.status.code,
+ (System.nanoTime() - startNanos) / 1000000L)
+ catch { case e: Throwable => logger.debug(s"Http4sApp says: could not record traffic: ${e.getMessage}") }
+ })
}
}
}
diff --git a/obp-api/src/main/scala/code/api/util/http4s/Http4sResourceDocs.scala b/obp-api/src/main/scala/code/api/util/http4s/Http4sResourceDocs.scala
index a97eb4ee72..da7f1c94a1 100644
--- a/obp-api/src/main/scala/code/api/util/http4s/Http4sResourceDocs.scala
+++ b/obp-api/src/main/scala/code/api/util/http4s/Http4sResourceDocs.scala
@@ -712,9 +712,16 @@ object Http4sResourceDocs extends MdcLoggable {
* answer outside ResourceDocMiddleware, which is where every other endpoint is timed; their docs
* are declared in ResourceDocs1_4_0, at v1.4.0, whatever version prefix the request used.
*/
- private def timed(handlerName: String)(response: IO[Response[IO]]): IO[Response[IO]] =
- code.telemetry.Telemetry.timeEndpoint(
- APIUtil.buildOperationId(ApiVersion.v1_4_0, handlerName), ApiVersion.v1_4_0.apiShortVersion)(response)
+ private def timed(req: Request[IO], handlerName: String)(response: IO[Response[IO]]): IO[Response[IO]] =
+ timedAs(req, APIUtil.buildOperationId(ApiVersion.v1_4_0, handlerName), ApiVersion.v1_4_0.apiShortVersion)(response)
+
+ private def timedAs(req: Request[IO], operationId: String, apiVersion: String)(response: IO[Response[IO]]): IO[Response[IO]] = {
+ Http4sRequestAttributes.trafficNote(req).foreach { note =>
+ note.operationId = Some(operationId)
+ note.apiVersion = Some(apiVersion)
+ }
+ code.telemetry.Telemetry.timeEndpoint(operationId, apiVersion)(response)
+ }
val routes: HttpRoutes[IO] = HttpRoutes.of[IO] {
case req @ GET -> Root / "obp" / prefix / "resource-docs" / requestedApiVersionString / "obp" =>
@@ -725,10 +732,10 @@ object Http4sResourceDocs extends MdcLoggable {
case "v4.0.0" | "v5.0.0" | "v5.1.0" | "v6.0.0" => true
case _ => false
}
- timed("getResourceDocsObp")(handleGetResourceDocsObp(req, prefix, requestedApiVersionString, isVersion4OrHigher = isV4OrHigher))
+ timed(req, "getResourceDocsObp")(handleGetResourceDocsObp(req, prefix, requestedApiVersionString, isVersion4OrHigher = isV4OrHigher))
case req @ GET -> Root / "obp" / prefix / "resource-docs" / requestedApiVersionString / "swagger" =>
- timed("getResourceDocsSwagger")(handleGetResourceDocsSwagger(req, prefix, requestedApiVersionString))
+ timed(req, "getResourceDocsSwagger")(handleGetResourceDocsSwagger(req, prefix, requestedApiVersionString))
// OpenAPI 3.1 JSON and YAML — served for every URL prefix.
//
@@ -742,18 +749,17 @@ object Http4sResourceDocs extends MdcLoggable {
// for non-v6 prefixes; `isVersion4OrHigher` is hardcoded `true` inside the
// handlers because the OpenAPI converter always consumes the v4-shape input.
case req @ GET -> Root / "obp" / prefix / "resource-docs" / requestedApiVersionString / "openapi" =>
- timed("getResourceDocsOpenAPI31")(handleGetResourceDocsOpenAPI31(req, prefix, requestedApiVersionString))
+ timed(req, "getResourceDocsOpenAPI31")(handleGetResourceDocsOpenAPI31(req, prefix, requestedApiVersionString))
// No ResourceDoc describes the YAML form, so it has no operation id and is not timed.
case req @ GET -> Root / "obp" / prefix / "resource-docs" / requestedApiVersionString / "openapi.yaml" =>
handleGetResourceDocsOpenAPI31Yaml(req, prefix, requestedApiVersionString)
case req @ GET -> Root / "obp" / prefix / "banks" / bankIdStr / "resource-docs" / requestedApiVersionString / "obp" =>
- timed("getBankLevelDynamicResourceDocsObp")(handleGetBankLevelDynamicResourceDocsObp(req, prefix, bankIdStr, requestedApiVersionString))
+ timed(req, "getBankLevelDynamicResourceDocsObp")(handleGetBankLevelDynamicResourceDocsObp(req, prefix, bankIdStr, requestedApiVersionString))
case req @ GET -> Root / "obp" / _ / "message-docs" / connector / "swagger2.0" =>
- code.telemetry.Telemetry.timeEndpoint(
- APIUtil.buildOperationId(ApiVersion.v3_1_0, "getMessageDocsSwagger"), ApiVersion.v3_1_0.apiShortVersion
- )(handleGetMessageDocsSwagger(req, connector))
+ timedAs(req, APIUtil.buildOperationId(ApiVersion.v3_1_0, "getMessageDocsSwagger"), ApiVersion.v3_1_0.apiShortVersion)(
+ handleGetMessageDocsSwagger(req, connector))
}
}
diff --git a/obp-api/src/main/scala/code/api/util/http4s/Http4sSupport.scala b/obp-api/src/main/scala/code/api/util/http4s/Http4sSupport.scala
index 2d9aebbaa7..f6a74258b5 100644
--- a/obp-api/src/main/scala/code/api/util/http4s/Http4sSupport.scala
+++ b/obp-api/src/main/scala/code/api/util/http4s/Http4sSupport.scala
@@ -82,6 +82,17 @@ object Http4sRequestAttributes {
val callContextKey: Key[CallContext] =
Key.newKey[IO, CallContext].unsafeRunSync()(cats.effect.unsafe.IORuntime.global)
+ /**
+ * Vault key for the traffic note of a request: what inner layers learn about it (its endpoint, its
+ * Consumer, a refusal) for TrafficSources, which Http4sApp records once the response is ready.
+ * Installed by Http4sApp on every request; bridge hops keep it, because `withUri` keeps attributes.
+ */
+ val trafficNoteKey: Key[code.telemetry.TrafficSources.Note] =
+ Key.newKey[IO, code.telemetry.TrafficSources.Note].unsafeRunSync()(cats.effect.unsafe.IORuntime.global)
+
+ /** The request's traffic note, when Http4sApp installed one (tests that call routes directly have none). */
+ def trafficNote(req: Request[IO]): Option[code.telemetry.TrafficSources.Note] = req.attributes.lookup(trafficNoteKey)
+
/**
* Vault key for caching the (already-read) request body across bridge cascade hops.
*
diff --git a/obp-api/src/main/scala/code/api/util/http4s/ResourceDocMiddleware.scala b/obp-api/src/main/scala/code/api/util/http4s/ResourceDocMiddleware.scala
index b710cded52..791b014731 100644
--- a/obp-api/src/main/scala/code/api/util/http4s/ResourceDocMiddleware.scala
+++ b/obp-api/src/main/scala/code/api/util/http4s/ResourceDocMiddleware.scala
@@ -177,6 +177,10 @@ object ResourceDocMiddleware extends MdcLoggable {
// api_enabled_versions. Fall through so the Lift bridge can serve or 404.
OptionT.none[IO, Response[IO]]
case Some(resourceDoc) =>
+ Http4sRequestAttributes.trafficNote(req).foreach { note =>
+ note.operationId = Some(resourceDoc.operationId)
+ note.apiVersion = Some(resourceDoc.implementedInApiVersion.apiShortVersion)
+ }
val ccWithDoc = ResourceDocMatcher.attachToCallContext(cc, resourceDoc)
val pathParams = ResourceDocMatcher.extractPathParams(req.uri.path, resourceDoc)
// Validate first (read-only, outside any transaction), then run business logic.
@@ -188,6 +192,14 @@ object ResourceDocMiddleware extends MdcLoggable {
case Left(errorResponse) =>
IO.pure(Option(errorResponse))
case Right(enrichedReq) =>
+ // The caller is known now: record its Consumer for TrafficSources.
+ for {
+ note <- Http4sRequestAttributes.trafficNote(req)
+ consumer <- enrichedReq.attributes.lookup(Http4sRequestAttributes.callContextKey).flatMap(_.consumer.toOption)
+ } {
+ note.consumerId = Some(consumer.consumerId.get)
+ note.consumerName = Some(consumer.name.get)
+ }
val routeIO =
routes.run(enrichedReq)
.map(ensureJsonContentType)
diff --git a/obp-api/src/main/scala/code/api/util/http4s/SelfServiceRateLimitMiddleware.scala b/obp-api/src/main/scala/code/api/util/http4s/SelfServiceRateLimitMiddleware.scala
index 7299e64f00..77694cccb0 100644
--- a/obp-api/src/main/scala/code/api/util/http4s/SelfServiceRateLimitMiddleware.scala
+++ b/obp-api/src/main/scala/code/api/util/http4s/SelfServiceRateLimitMiddleware.scala
@@ -139,7 +139,9 @@ object SelfServiceRateLimitMiddleware extends MdcLoggable {
if (ipPenaltyManagementPath.findFirstIn(req.uri.path.renderString).isDefined) IO.pure(None)
else IO.blocking(IpPenalties.check(Http4sCallContextBuilder.clientIp(req)))
penaltyRefusal.flatMap {
- case Some(refusal) => IO.pure(penaltyResponse(refusal))
+ case Some(refusal) =>
+ Http4sRequestAttributes.trafficNote(req).foreach(_.refusedBy = Some("ip_penalty"))
+ IO.pure(penaltyResponse(refusal))
case None => applyScopes(req)(run)
}
}
@@ -162,7 +164,9 @@ object SelfServiceRateLimitMiddleware extends MdcLoggable {
case None => run(req)
case Some(scope) =>
IO.blocking(SelfServiceRateLimiter.check(scope, Http4sCallContextBuilder.clientIp(req), "ip")).flatMap {
- case Blocked(s, _, exceeded) => IO.pure(blockedResponse(s, exceeded))
+ case Blocked(s, _, exceeded) =>
+ Http4sRequestAttributes.trafficNote(req).foreach(_.refusedBy = Some(s))
+ IO.pure(blockedResponse(s, exceeded))
case outcome => run(req).map(resp => decorate(resp, outcome))
}
}
diff --git a/obp-api/src/main/scala/code/api/v7_0_0/Http4s700.scala b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700.scala
index da808df931..b255a08561 100644
--- a/obp-api/src/main/scala/code/api/v7_0_0/Http4s700.scala
+++ b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700.scala
@@ -7249,6 +7249,9 @@ object Http4s700 {
// IP penalties: an operator's temporary per-minute limit on one address.
resourceDocs ++= Http4s700IpPenalties.resourceDocs
+ // Where traffic is coming from: the busiest Consumers, addresses, and callers and endpoints.
+ resourceDocs ++= Http4s700TrafficSources.resourceDocs
+
val allRoutes: HttpRoutes[IO] = {
val sorted = resourceDocs
.sortBy(rd => -rd.requestUrl.split("/").count(_.nonEmpty))
diff --git a/obp-api/src/main/scala/code/api/v7_0_0/Http4s700TrafficSources.scala b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700TrafficSources.scala
new file mode 100644
index 0000000000..c939d4169d
--- /dev/null
+++ b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700TrafficSources.scala
@@ -0,0 +1,120 @@
+/**
+Open Bank Project - API
+Copyright (C) 2011-2026, TESOBE GmbH.
+
+This program is free software: you can redistribute it and/or modify
+it under the terms of the GNU Affero General Public License as published by
+the Free Software Foundation, either version 3 of the License, or
+(at your option) any later version.
+
+This program is distributed in the hope that it will be useful,
+but WITHOUT ANY WARRANTY; without even the implied warranty of
+MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+GNU Affero General Public License for more details.
+
+You should have received a copy of the GNU Affero General Public License
+along with this program. If not, see .
+
+Email: contact@tesobe.com
+TESOBE GmbH.
+Osloer Strasse 16/17
+Berlin 13359, Germany
+
+This product includes software developed at
+TESOBE (http://www.tesobe.com/)
+
+ */
+package code.api.v7_0_0
+
+import cats.effect.IO
+import code.api.Constant.ApiPathZero
+import code.api.util.APIUtil.{EmptyBody, ResourceDoc}
+import code.api.util.ApiRole._
+import code.api.util.ApiTag._
+import code.api.util.ErrorMessages._
+import code.api.util.http4s.Http4sRequestAttributes.EndpointHelpers
+import code.api.util.CustomJsonFormats
+import code.telemetry.TrafficSources
+import code.util.Helper
+import com.github.dwickern.macros.NameOf.nameOf
+import com.openbankproject.commons.ExecutionContext.Implicits.global
+import com.openbankproject.commons.util.ApiVersion
+import org.http4s._
+import org.http4s.dsl.io._
+import org.json4s.Formats
+
+import scala.collection.mutable.ArrayBuffer
+import scala.concurrent.Future
+
+/**
+ * This object holds the v7.0.0 endpoint that shows where the traffic on this instance is coming
+ * from: the busiest Consumers, client IP addresses, and callers and endpoints (see
+ * [[code.telemetry.TrafficSources]]).
+ *
+ * It is declared in its own object to keep Http4s700's initialiser under the JVM's 64KB method limit.
+ */
+object Http4s700TrafficSources {
+
+ implicit val formats: Formats = CustomJsonFormats.formats
+
+ private val implementedInApiVersion = ApiVersion.v7_0_0
+ private val prefixPath = Root / ApiPathZero.toString / implementedInApiVersion.toString
+
+ val resourceDocs = ArrayBuffer[ResourceDoc]()
+
+ private val AllowedWindows = Set(1, 5, 15)
+
+ // Route: GET /obp/v7.0.0/management/traffic/top-callers
+ lazy val getTrafficSources: HttpRoutes[IO] = HttpRoutes.of[IO] {
+ case req @ GET -> `prefixPath` / "management" / "traffic" / "top-callers" =>
+ EndpointHelpers.withUser(req) { (_, cc) =>
+ val window = req.uri.query.params.get("window").map(_.trim).getOrElse("5")
+ for {
+ _ <- Helper.booleanToFuture(s"$InvalidTrafficWindow Current value is $window", failCode = 400, cc = Some(cc)) {
+ window.forall(_.isDigit) && window.nonEmpty && AllowedWindows.contains(window.toInt)
+ }
+ } yield JSONFactory700Operations.createTrafficSourcesJson(window.toInt)
+ }
+ }
+
+ resourceDocs += ResourceDoc(
+ implementedInApiVersion,
+ nameOf(getTrafficSources),
+ "GET",
+ "/management/traffic/top-callers",
+ "Get Top Callers",
+ s"""Where the traffic on this OBP-API instance is coming from: the busiest Consumers, the busiest
+ |client IP addresses, and the busiest pairs of caller and endpoint, over the last `window`
+ |minutes (1, 5 or 15; default 5).
+ |
+ |- `consumers`: requests that authenticated a Consumer, by Consumer.
+ |- `addresses`: every request, authenticated or not, by client IP address, with the Consumers seen from it.
+ |- `callers_and_endpoints`: by caller (the Consumer when one was authenticated, otherwise the IP address)
+ | and endpoint. The endpoint is an operation id, `unmatched` for paths no endpoint serves (one value for
+ | all of them), `refused:LIMITER` for requests refused before routing, or `other`.
+ |
+ |Each list shows at most ${JSONFactory700Operations.TrafficRowsShown} rows, busiest first. Counts are
+ |estimates: `requests` is within `error` of the true count. `status_*`, `refused` and `unmatched` are exact
+ |from when the row started being tracked in each minute.
+ |
+ |How it is counted. Each minute has three tables of fixed size (${TrafficSources.ConsumerSlots} Consumers,
+ |${TrafficSources.AddressSlots} addresses, ${TrafficSources.CallerEndpointSlots} pairs), kept with the
+ |Space-Saving algorithm of Ahmed Metwally, Divyakant Agrawal and Amr El Abbadi ("Efficient Computation of
+ |Frequent and Top-k Elements in Data Streams", ICDT 2005), which refines the frequent-items algorithm of
+ |Jayadev Misra and David Gries (1982). Any caller with more than 1/${TrafficSources.ConsumerSlots} of a
+ |minute's requests is guaranteed to appear. Minutes are merged into the window as mergeable summaries
+ |(Agarwal, Cormode, Huang, Phillips, Wei and Yi, PODS 2012). ${TrafficSources.MinutesKept} minutes are kept.
+ |
+ |Each instance counts only its own traffic, in memory. Behind a load balancer the response describes the
+ |instance that answered (`api_instance_id`). Nothing is written to the database or to Prometheus, and IP
+ |addresses are forgotten after ${TrafficSources.MinutesKept} minutes. IP addresses are personal data: this
+ |Role should be granted only to people who handle incidents.
+ |""".stripMargin,
+ EmptyBody,
+ JSONFactory700Operations.trafficSourcesJsonV700Example,
+ List($AuthenticatedUserIsRequired, UserHasMissingRoles, InvalidTrafficWindow, UnknownError),
+ List(apiTagRateLimits, apiTagSystem),
+ Some(List(canGetTrafficSources)),
+ http4sPartialFunction = Some(getTrafficSources)
+ )
+}
diff --git a/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory700Operations.scala b/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory700Operations.scala
index 5e92e9084a..25e9dac5eb 100644
--- a/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory700Operations.scala
+++ b/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory700Operations.scala
@@ -30,7 +30,7 @@ import java.util.Date
import code.api.Constant
import code.api.util.{APIUtil, ExampleValue, IpPenalties}
-import code.telemetry.Telemetry
+import code.telemetry.{Telemetry, TrafficSources}
/*
* The JSON of the v7.0.0 operations endpoints (Telemetry and IP penalties).
@@ -80,6 +80,66 @@ case class IpPenaltyJsonV700(
case class IpPenaltiesJsonV700(ip_penalties: List[IpPenaltyJsonV700])
+// ===== Where traffic is coming from =====
+
+/** A Consumer among the busiest. `requests` is within `error` of the true count. */
+case class TrafficConsumerJsonV700(
+ consumer_id: String,
+ application_name: String,
+ requests: Long,
+ error: Long,
+ status_2xx: Long,
+ status_4xx: Long,
+ status_5xx: Long,
+ refused: Long,
+ unmatched: Long,
+ endpoints: List[String],
+ last_ip_address: String,
+ first_seen: Date,
+ last_seen: Date
+)
+
+/** A client IP address among the busiest, whatever credentials its requests carried. */
+case class TrafficAddressJsonV700(
+ ip_address: String,
+ requests: Long,
+ error: Long,
+ status_2xx: Long,
+ status_4xx: Long,
+ status_5xx: Long,
+ refused: Long,
+ unmatched: Long,
+ endpoints: List[String],
+ consumer_ids: List[String],
+ first_seen: Date,
+ last_seen: Date
+)
+
+/** A pair of caller (a Consumer, or an IP address for other requests) and endpoint among the busiest. */
+case class TrafficCallerEndpointJsonV700(
+ caller_kind: String,
+ caller: String,
+ endpoint: String,
+ api_version: Option[String],
+ requests: Long,
+ error: Long,
+ status_2xx: Long,
+ status_4xx: Long,
+ status_5xx: Long,
+ refused: Long,
+ mean_duration_ms: Long,
+ max_duration_ms: Long,
+ last_seen: Date
+)
+
+case class TrafficSourcesJsonV700(
+ api_instance_id: String,
+ window_minutes: Int,
+ consumers: List[TrafficConsumerJsonV700],
+ addresses: List[TrafficAddressJsonV700],
+ callers_and_endpoints: List[TrafficCallerEndpointJsonV700]
+)
+
/** This object builds the JSON above, and holds the examples the ResourceDocs show. */
object JSONFactory700Operations {
@@ -137,4 +197,72 @@ object JSONFactory700Operations {
created_at = APIUtil.DateWithMsExampleObject, expires_at = APIUtil.DateWithMsExampleObject)
lazy val ipPenaltiesJsonV700Example = IpPenaltiesJsonV700(List(ipPenaltyJsonV700Example))
+
+ // ===== Where traffic is coming from =====
+
+ val TrafficRowsShown = 50
+
+ private def endpointsOf(sets: List[scala.collection.mutable.LinkedHashSet[String]]): List[String] =
+ sets.flatten.distinct.take(TrafficSources.EndpointsKeptPerCaller)
+
+ /** The busiest Consumers, addresses, and callers and endpoints of this instance, over the last `windowMinutes`. */
+ def createTrafficSourcesJson(windowMinutes: Int): TrafficSourcesJsonV700 = {
+ import TrafficSources._
+ val consumers = TrafficSources.consumers(windowMinutes).take(TrafficRowsShown).map { m =>
+ val latest = m.details.maxBy(_.lastSeen)
+ TrafficConsumerJsonV700(
+ consumer_id = m.key,
+ application_name = m.details.map(_.name).find(_.nonEmpty).getOrElse(""),
+ requests = m.requests, error = m.error,
+ status_2xx = m.details.map(_.status2xx).sum, status_4xx = m.details.map(_.status4xx).sum,
+ status_5xx = m.details.map(_.status5xx).sum, refused = m.details.map(_.refused).sum,
+ unmatched = m.details.map(_.unmatched).sum,
+ endpoints = endpointsOf(m.details.map(_.endpoints)),
+ last_ip_address = latest.lastIp,
+ first_seen = new Date(m.details.map(_.firstSeen).min), last_seen = new Date(latest.lastSeen))
+ }
+ val addresses = TrafficSources.addresses(windowMinutes).take(TrafficRowsShown).map { m =>
+ TrafficAddressJsonV700(
+ ip_address = m.key,
+ requests = m.requests, error = m.error,
+ status_2xx = m.details.map(_.status2xx).sum, status_4xx = m.details.map(_.status4xx).sum,
+ status_5xx = m.details.map(_.status5xx).sum, refused = m.details.map(_.refused).sum,
+ unmatched = m.details.map(_.unmatched).sum,
+ endpoints = endpointsOf(m.details.map(_.endpoints)),
+ consumer_ids = m.details.flatMap(_.consumers).distinct.take(ConsumersKeptPerAddress),
+ first_seen = new Date(m.details.map(_.firstSeen).min), last_seen = new Date(m.details.map(_.lastSeen).max))
+ }
+ val callerEndpoints = TrafficSources.callerEndpoints(windowMinutes).take(TrafficRowsShown).map { m =>
+ val (caller, endpoint) = m.key
+ val counted = m.details.map(d => d.status2xx + d.status4xx + d.status5xx).sum
+ TrafficCallerEndpointJsonV700(
+ caller_kind = caller.kind, caller = caller.value, endpoint = endpoint,
+ api_version = m.details.map(_.apiVersion).find(_.nonEmpty),
+ requests = m.requests, error = m.error,
+ status_2xx = m.details.map(_.status2xx).sum, status_4xx = m.details.map(_.status4xx).sum,
+ status_5xx = m.details.map(_.status5xx).sum, refused = m.details.map(_.refused).sum,
+ mean_duration_ms = if (counted > 0) m.details.map(_.totalDurationMillis).sum / counted else 0L,
+ max_duration_ms = m.details.map(_.maxDurationMillis).max,
+ last_seen = new Date(m.details.map(_.lastSeen).max))
+ }
+ TrafficSourcesJsonV700(Constant.ApiInstanceId, windowMinutes, consumers, addresses, callerEndpoints)
+ }
+
+ lazy val trafficSourcesJsonV700Example = TrafficSourcesJsonV700(
+ api_instance_id = "obp_4f6b3c2a-9d1e-4b7a-8c5f-2e1d0a9b8c7d",
+ window_minutes = 5,
+ consumers = List(TrafficConsumerJsonV700(
+ consumer_id = ExampleValue.consumerIdExample.value, application_name = "Mobile Banking App",
+ requests = 12400, error = 0, status_2xx = 12310, status_4xx = 85, status_5xx = 5, refused = 0, unmatched = 0,
+ endpoints = List("OBPv7.0.0-getBanks", "OBPv6.0.0-getCoreAccountById"), last_ip_address = "198.51.100.23",
+ first_seen = APIUtil.DateWithMsExampleObject, last_seen = APIUtil.DateWithMsExampleObject)),
+ addresses = List(TrafficAddressJsonV700(
+ ip_address = "203.0.113.42", requests = 52400, error = 300, status_2xx = 1200, status_4xx = 51100, status_5xx = 100,
+ refused = 0, unmatched = 38000, endpoints = List("unmatched", "OBPv1.4.0-getResourceDocsObp"), consumer_ids = Nil,
+ first_seen = APIUtil.DateWithMsExampleObject, last_seen = APIUtil.DateWithMsExampleObject)),
+ callers_and_endpoints = List(TrafficCallerEndpointJsonV700(
+ caller_kind = "ip", caller = "203.0.113.42", endpoint = "unmatched", api_version = None,
+ requests = 38000, error = 250, status_2xx = 0, status_4xx = 38000, status_5xx = 0, refused = 0,
+ mean_duration_ms = 3, max_duration_ms = 41, last_seen = APIUtil.DateWithMsExampleObject))
+ )
}
diff --git a/obp-api/src/main/scala/code/telemetry/HeavyHitters.scala b/obp-api/src/main/scala/code/telemetry/HeavyHitters.scala
new file mode 100644
index 0000000000..b34cb80a8d
--- /dev/null
+++ b/obp-api/src/main/scala/code/telemetry/HeavyHitters.scala
@@ -0,0 +1,111 @@
+/**
+Open Bank Project - API
+Copyright (C) 2011-2026, TESOBE GmbH.
+
+This program is free software: you can redistribute it and/or modify
+it under the terms of the GNU Affero General Public License as published by
+the Free Software Foundation, either version 3 of the License, or
+(at your option) any later version.
+
+This program is distributed in the hope that it will be useful,
+but WITHOUT ANY WARRANTY; without even the implied warranty of
+MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+GNU Affero General Public License for more details.
+
+You should have received a copy of the GNU Affero General Public License
+along with this program. If not, see .
+
+Email: contact@tesobe.com
+TESOBE GmbH.
+Osloer Strasse 16/17
+Berlin 13359, Germany
+
+This product includes software developed at
+TESOBE (http://www.tesobe.com/)
+
+ */
+package code.telemetry
+
+import scala.collection.mutable
+
+/**
+ * This class keeps the heavy hitters of a stream (its most frequent keys) in a table of fixed size, using the Space-Saving
+ * algorithm (Metwally, Agrawal and El Abbadi, "Efficient Computation of Frequent and Top-k Elements
+ * in Data Streams", 2005).
+ *
+ * The problem it solves: counting every key of a stream (every caller, every caller and endpoint)
+ * takes memory without limit, and a caller can inflate it on purpose by varying its key. This table
+ * never grows past `capacity`. With N offers counted, every key offered more than N / capacity times
+ * is guaranteed to be in it, and each entry's `count` is at most its `error` above the true count.
+ *
+ * How an offer is counted:
+ * - a key already in the table: its count goes up by one;
+ * - a new key and a free slot: it is added with count 1 and error 0;
+ * - a new key and a full table: the entry with the smallest count (call it min) is replaced by the
+ * new key, with count min + 1 and error min. The new entry's details start from this offer.
+ *
+ * The paper finds the smallest entry with a linked structure of count buckets. At the sizes used
+ * here (hundreds of slots), a scan only when a full table meets a new key costs microseconds, and is
+ * simpler to get right.
+ *
+ * `D` is a mutable details record kept per entry (status counts, durations and so on). Details are
+ * exact from the moment the key entered the table; only `count` covers the key's whole history
+ * within its bound.
+ *
+ * All methods are synchronised: one instance is shared by every request thread of its minute.
+ *
+ * References (credit to the authors of the ideas used here):
+ * - Space-Saving, the algorithm this class implements: Ahmed Metwally, Divyakant Agrawal and
+ * Amr El Abbadi, "Efficient Computation of Frequent and Top-k Elements in Data Streams",
+ * Proceedings of the 10th International Conference on Database Theory (ICDT 2005), Lecture Notes
+ * in Computer Science 3363, Springer, 2005, pages 398-412.
+ * - The earlier algorithm it refines, the first for finding frequent items in bounded space:
+ * Jayadev Misra and David Gries, "Finding Repeated Elements", Science of Computer Programming 2(2),
+ * 1982, pages 143-152.
+ * - Combining per-minute tables into a longer window (see TrafficSources) relies on these summaries
+ * being mergeable: Pankaj K. Agarwal, Graham Cormode, Zengfeng Huang, Jeff M. Phillips, Zhewei Wei
+ * and Ke Yi, "Mergeable Summaries", Proceedings of the 31st ACM Symposium on Principles of
+ * Database Systems (PODS 2012), pages 23-34.
+ */
+final class HeavyHitters[K, D](val capacity: Int, newDetails: () => D) {
+
+ require(capacity > 0, "HeavyHitters capacity must be positive")
+
+ final class Entry private[HeavyHitters] (val key: K, var count: Long, var error: Long, val details: D)
+
+ private val entries = mutable.HashMap.empty[K, Entry]
+ private var offered: Long = 0L
+
+ /** Counts one offer of `key` and applies `update` to its details. */
+ def offer(key: K)(update: D => Unit): Unit = synchronized {
+ offered += 1
+ val entry = entries.get(key) match {
+ case Some(existing) =>
+ existing.count += 1
+ existing
+ case None if entries.size < capacity =>
+ val added = new Entry(key, 1L, 0L, newDetails())
+ entries.put(key, added)
+ added
+ case None =>
+ val smallest = entries.valuesIterator.minBy(_.count)
+ entries.remove(smallest.key)
+ val replacement = new Entry(key, smallest.count + 1, smallest.count, newDetails())
+ entries.put(key, replacement)
+ replacement
+ }
+ update(entry.details)
+ }
+
+ /** True when every slot is taken, so a key missing from the table may still have up to [[minCount]] offers. */
+ def isFull: Boolean = synchronized(entries.size >= capacity)
+
+ /** The smallest count in the table (0 when empty). */
+ def minCount: Long = synchronized(if (entries.isEmpty) 0L else entries.valuesIterator.map(_.count).min)
+
+ /** How many offers the table has counted. */
+ def totalOffered: Long = synchronized(offered)
+
+ /** The entries, as a copy taken under the lock (the details are shared; read them, do not change them). */
+ def snapshot: List[Entry] = synchronized(entries.values.toList)
+}
diff --git a/obp-api/src/main/scala/code/telemetry/TrafficSources.scala b/obp-api/src/main/scala/code/telemetry/TrafficSources.scala
new file mode 100644
index 0000000000..c06839b657
--- /dev/null
+++ b/obp-api/src/main/scala/code/telemetry/TrafficSources.scala
@@ -0,0 +1,230 @@
+/**
+Open Bank Project - API
+Copyright (C) 2011-2026, TESOBE GmbH.
+
+This program is free software: you can redistribute it and/or modify
+it under the terms of the GNU Affero General Public License as published by
+the Free Software Foundation, either version 3 of the License, or
+(at your option) any later version.
+
+This program is distributed in the hope that it will be useful,
+but WITHOUT ANY WARRANTY; without even the implied warranty of
+MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+GNU Affero General Public License for more details.
+
+You should have received a copy of the GNU Affero General Public License
+along with this program. If not, see .
+
+Email: contact@tesobe.com
+TESOBE GmbH.
+Osloer Strasse 16/17
+Berlin 13359, Germany
+
+This product includes software developed at
+TESOBE (http://www.tesobe.com/)
+
+ */
+package code.telemetry
+
+import java.util.concurrent.atomic.AtomicReference
+
+import scala.collection.mutable
+
+/**
+ * This object answers "where is the traffic on this instance coming from, right now?": the busiest
+ * Consumers, the busiest client IP addresses, and the busiest pairs of caller and endpoint, over the
+ * last 1, 5 or 15 minutes.
+ *
+ * It is not Telemetry in the Prometheus sense and is never exported there: it names Consumers and
+ * IP addresses, which must not become series (docs/telemetry_conventions.md, section 6). It is read
+ * only through a Role-gated endpoint, lives only in this instance's memory, and forgets everything
+ * after 15 minutes. It needs no database, so it keeps working when API Metrics are off or the
+ * database is the thing under strain.
+ *
+ * Memory is bounded whatever the traffic: each minute has three [[HeavyHitters]] tables (200
+ * Consumers, 200 addresses, 500 caller and endpoint pairs), and 15 minutes are kept. Counts are
+ * estimates with a stated bound; see HeavyHitters for the algorithm (Space-Saving) and its authors.
+ *
+ * Every request is recorded once, at the outermost layer (Http4sApp), which sees refused and
+ * unmatched requests as well as served ones. Inner layers fill in a [[Note]] carried as a request
+ * attribute: ResourceDocMiddleware and the resource-docs routes set the endpoint and, after
+ * authentication, the Consumer; the IP limiters set a refusal.
+ */
+object TrafficSources {
+
+ val ConsumerSlots = 200
+ val AddressSlots = 200
+ val CallerEndpointSlots = 500
+ val MinutesKept = 15
+ val EndpointsKeptPerCaller = 32
+ val ConsumersKeptPerAddress = 8
+
+ // ===== What inner layers tell the recording point =====
+
+ /** Filled in by inner layers while a request is served; read once when it is recorded. */
+ final class Note {
+ @volatile var operationId: Option[String] = None
+ @volatile var apiVersion: Option[String] = None
+ @volatile var consumerId: Option[String] = None
+ @volatile var consumerName: Option[String] = None
+ @volatile var refusedBy: Option[String] = None
+ }
+
+ // ===== Keys =====
+
+ sealed trait Caller { def kind: String; def value: String }
+ final case class ConsumerCaller(consumerId: String) extends Caller { val kind = "consumer"; def value: String = consumerId }
+ final case class AddressCaller(ipAddress: String) extends Caller { val kind = "ip"; def value: String = ipAddress }
+
+ /** Requests to paths no endpoint serves are grouped under one endpoint: scanners send endless distinct paths. */
+ val UnmatchedEndpoint = "unmatched"
+ /** Requests that no documented endpoint answered and that were not a 404, e.g. status pages and CORS preflight. */
+ val OtherEndpoint = "other"
+
+ // ===== Details kept per entry =====
+
+ class StatusCounts {
+ var status2xx = 0L; var status4xx = 0L; var status5xx = 0L; var refused = 0L; var unmatched = 0L
+ def add(status: Int, isRefused: Boolean, isUnmatched: Boolean): Unit = {
+ if (status >= 200 && status < 300) status2xx += 1
+ else if (status >= 400 && status < 500) status4xx += 1
+ else if (status >= 500) status5xx += 1
+ if (isRefused) refused += 1
+ if (isUnmatched) unmatched += 1
+ }
+ }
+
+ final class ConsumerDetails extends StatusCounts {
+ var name = ""
+ val endpoints = mutable.LinkedHashSet.empty[String]
+ var lastIp = ""
+ var firstSeen = 0L; var lastSeen = 0L
+ }
+
+ final class AddressDetails extends StatusCounts {
+ val endpoints = mutable.LinkedHashSet.empty[String]
+ val consumers = mutable.LinkedHashSet.empty[String]
+ var firstSeen = 0L; var lastSeen = 0L
+ }
+
+ final class CallerEndpointDetails extends StatusCounts {
+ var apiVersion = ""
+ var totalDurationMillis = 0L; var maxDurationMillis = 0L
+ var lastSeen = 0L
+ }
+
+ // ===== One minute =====
+
+ final class Minute(val startMillis: Long) {
+ val consumers = new HeavyHitters[String, ConsumerDetails](ConsumerSlots, () => new ConsumerDetails)
+ val addresses = new HeavyHitters[String, AddressDetails](AddressSlots, () => new AddressDetails)
+ val callerEndpoints = new HeavyHitters[(Caller, String), CallerEndpointDetails](CallerEndpointSlots, () => new CallerEndpointDetails)
+ }
+
+ private def minuteStart(millis: Long): Long = millis - millis % 60000L
+
+ private val minutes = new AtomicReference[List[Minute]](Nil) // newest first
+
+ private def minuteFor(now: Long): Minute = {
+ val start = minuteStart(now)
+ minutes.get() match {
+ case current :: _ if current.startMillis == start => current
+ case _ => synchronized {
+ minutes.get() match {
+ case current :: _ if current.startMillis == start => current
+ case kept =>
+ val fresh = new Minute(start)
+ minutes.set((fresh :: kept).filter(_.startMillis > start - MinutesKept * 60000L))
+ fresh
+ }
+ }
+ }
+ }
+
+ private def addCapped(set: mutable.LinkedHashSet[String], value: String, limit: Int): Unit =
+ if (set.contains(value) || set.size < limit) set += value
+
+ /**
+ * Records one request. `note` is what inner layers filled in; `status` and `durationMillis` are the
+ * response's. A request counts under its Consumer when one was authenticated, and always under its
+ * client address.
+ */
+ def record(note: Note, ipAddress: String, status: Int, durationMillis: Long, now: Long = System.currentTimeMillis()): Unit = {
+ val minute = minuteFor(now)
+ val isRefused = note.refusedBy.isDefined || status == 429
+ val endpoint = note.refusedBy.map(limiter => s"refused:$limiter")
+ .orElse(note.operationId)
+ .getOrElse(if (status == 404) UnmatchedEndpoint else OtherEndpoint)
+ val isUnmatched = endpoint == UnmatchedEndpoint
+ val address = Option(ipAddress).filter(_.nonEmpty).getOrElse("unknown")
+
+ note.consumerId.foreach { consumerId =>
+ minute.consumers.offer(consumerId) { d =>
+ d.add(status, isRefused, isUnmatched)
+ note.consumerName.foreach(d.name = _)
+ addCapped(d.endpoints, endpoint, EndpointsKeptPerCaller)
+ d.lastIp = address
+ if (d.firstSeen == 0L) d.firstSeen = now
+ d.lastSeen = now
+ }
+ }
+ minute.addresses.offer(address) { d =>
+ d.add(status, isRefused, isUnmatched)
+ addCapped(d.endpoints, endpoint, EndpointsKeptPerCaller)
+ note.consumerId.foreach(addCapped(d.consumers, _, ConsumersKeptPerAddress))
+ if (d.firstSeen == 0L) d.firstSeen = now
+ d.lastSeen = now
+ }
+ val caller: Caller = note.consumerId.map(ConsumerCaller(_)).getOrElse(AddressCaller(address))
+ minute.callerEndpoints.offer((caller, endpoint)) { d =>
+ d.add(status, isRefused, isUnmatched)
+ note.apiVersion.foreach(d.apiVersion = _)
+ d.totalDurationMillis += durationMillis
+ d.maxDurationMillis = math.max(d.maxDurationMillis, durationMillis)
+ d.lastSeen = now
+ }
+ }
+
+ // ===== Reading: merging the minutes of a window =====
+
+ /**
+ * One key's totals over a window. The true number of requests is within `error` of `requests`:
+ * `error` adds up each minute's overcount for the key and, for each full minute table the key was
+ * missing from, that minute's smallest count (the key may have been there below it).
+ */
+ final case class Merged[K, D](key: K, requests: Long, error: Long, details: List[D])
+
+ private def merge[K, D](tables: List[HeavyHitters[K, D]]): List[Merged[K, D]] = {
+ val snapshots = tables.map(t => (t, t.snapshot.map(e => e.key -> e).toMap, t.isFull, t.minCount))
+ val keys = snapshots.flatMap(_._2.keys).distinct
+ keys.map { key =>
+ var requests = 0L; var error = 0L; val details = List.newBuilder[D]
+ snapshots.foreach { case (_, byKey, full, min) =>
+ byKey.get(key) match {
+ case Some(entry) => requests += entry.count; error += entry.error; details += entry.details
+ case None if full => error += min
+ case None => ()
+ }
+ }
+ Merged(key, requests, error, details.result())
+ }.sortBy(m => -m.requests)
+ }
+
+ /** The window's minutes, newest first, and the time its oldest minute started. */
+ private def window(minutesBack: Int, now: Long): List[Minute] = {
+ val from = minuteStart(now) - (minutesBack - 1) * 60000L
+ minutes.get().filter(_.startMillis >= from)
+ }
+
+ def consumers(minutesBack: Int, now: Long = System.currentTimeMillis()): List[Merged[String, ConsumerDetails]] =
+ merge(window(minutesBack, now).map(_.consumers))
+
+ def addresses(minutesBack: Int, now: Long = System.currentTimeMillis()): List[Merged[String, AddressDetails]] =
+ merge(window(minutesBack, now).map(_.addresses))
+
+ def callerEndpoints(minutesBack: Int, now: Long = System.currentTimeMillis()): List[Merged[(Caller, String), CallerEndpointDetails]] =
+ merge(window(minutesBack, now).map(_.callerEndpoints))
+
+ /** Forget everything (tests). */
+ def clear(): Unit = minutes.set(Nil)
+}
diff --git a/obp-api/src/test/scala/code/api/v7_0_0/TrafficSourcesEndpointTest.scala b/obp-api/src/test/scala/code/api/v7_0_0/TrafficSourcesEndpointTest.scala
new file mode 100644
index 0000000000..12bdca1f3c
--- /dev/null
+++ b/obp-api/src/test/scala/code/api/v7_0_0/TrafficSourcesEndpointTest.scala
@@ -0,0 +1,71 @@
+package code.api.v7_0_0
+
+import code.api.util.APIUtil.OAuth._
+import code.api.util.ApiRole.CanGetTrafficSources
+import code.api.util.ErrorMessages.{AuthenticatedUserIsRequired, UserHasMissingRoles}
+import code.api.v6_0_0.V600ServerSetup
+import code.entitlement.Entitlement
+import com.openbankproject.commons.model.ErrorMessage
+import com.openbankproject.commons.util.ApiVersion
+import org.scalatest.Tag
+
+/**
+ * This suite checks GET /obp/v7.0.0/management/traffic/top-callers: that it needs its Role, and that
+ * real requests appear in the right tables (an authenticated one under its Consumer, every one under
+ * its address, and an unknown path as `unmatched`).
+ */
+class TrafficSourcesEndpointTest extends V600ServerSetup {
+
+ def v7_0_0_Request = baseRequest / "obp" / "v7.0.0"
+
+ object VersionOfApi extends Tag(ApiVersion.v7_0_0.toString)
+ object ApiEndpoint extends Tag("getTrafficSources")
+
+ private def topCallers = v7_0_0_Request / "management" / "traffic" / "top-callers"
+
+ private def withRole[T](body: => T): T = {
+ val entitlement = Entitlement.entitlement.vend.addEntitlement("", resourceUser1.userId, CanGetTrafficSources.toString)
+ try body finally Entitlement.entitlement.vend.deleteEntitlement(entitlement)
+ }
+
+ feature(s"Get Top Callers - GET /obp/v7.0.0/management/traffic/top-callers - $VersionOfApi") {
+
+ scenario("anonymous access is 401, and a user without the Role gets 403", ApiEndpoint, VersionOfApi) {
+ val anonymous = makeGetRequest(topCallers.GET)
+ anonymous.code should equal(401)
+ anonymous.body.extract[ErrorMessage].message should equal(AuthenticatedUserIsRequired)
+ val noRole = makeGetRequest(topCallers.GET <@ (user1))
+ noRole.code should equal(403)
+ noRole.body.extract[ErrorMessage].message should equal(UserHasMissingRoles + CanGetTrafficSources)
+ }
+
+ scenario("requests appear under their Consumer, their address, and their endpoint", ApiEndpoint, VersionOfApi) {
+ makeGetRequest((v7_0_0_Request / "banks").GET <@ (user1)).code should equal(200)
+ makeGetRequest((v7_0_0_Request / "banks").GET).code should equal(200)
+ makeGetRequest((v7_0_0_Request / "no-such-endpoint-traffic-probe").GET).code should equal(404)
+
+ val response = withRole(makeGetRequest(topCallers.GET <@ (user1) < List("window" -> "1")))
+ response.code should equal(200)
+ val traffic = response.body.extract[TrafficSourcesJsonV700]
+ traffic.window_minutes should equal(1)
+
+ Then("the authenticated request is under the test Consumer")
+ val consumer = traffic.consumers.find(_.consumer_id == testConsumer.consumerId.get)
+ consumer should not be empty
+ consumer.get.endpoints.exists(_.endsWith("-getBanks")) shouldBe true
+
+ And("every request is under this client's address, the unknown path as unmatched")
+ traffic.addresses should not be empty
+ traffic.callers_and_endpoints.exists(p => p.caller_kind == "ip" && p.endpoint == "unmatched") shouldBe true
+ traffic.callers_and_endpoints.exists(p => p.caller_kind == "consumer" && p.caller == testConsumer.consumerId.get &&
+ p.endpoint.endsWith("-getBanks")) shouldBe true
+ }
+
+ scenario("a window other than 1, 5 or 15 is refused with 400", ApiEndpoint, VersionOfApi) {
+ withRole {
+ makeGetRequest(topCallers.GET <@ (user1) < List("window" -> "7")).code should equal(400)
+ makeGetRequest(topCallers.GET <@ (user1) < List("window" -> "abc")).code should equal(400)
+ }
+ }
+ }
+}
diff --git a/obp-api/src/test/scala/code/telemetry/HeavyHittersTest.scala b/obp-api/src/test/scala/code/telemetry/HeavyHittersTest.scala
new file mode 100644
index 0000000000..1baa31fe7b
--- /dev/null
+++ b/obp-api/src/test/scala/code/telemetry/HeavyHittersTest.scala
@@ -0,0 +1,51 @@
+package code.telemetry
+
+import org.scalatest.{FlatSpec, Matchers}
+
+import scala.collection.mutable
+
+/**
+ * This suite checks the guarantees of HeavyHitters, the Space-Saving algorithm (Metwally, Agrawal and El Abbadi, ICDT 2005): below
+ * capacity it counts exactly; a key with more than N / capacity offers is always kept; and every
+ * entry's true count lies between `count - error` and `count`.
+ */
+class HeavyHittersTest extends FlatSpec with Matchers {
+
+ final class Hits { var n = 0 }
+
+ "HeavyHitters" should "count exactly while it has free slots" in {
+ val table = new HeavyHitters[String, Hits](10, () => new Hits)
+ List("a", "b", "a", "c", "a").foreach(k => table.offer(k)(_.n += 1))
+ table.snapshot.map(e => e.key -> (e.count, e.error)).toMap shouldBe Map("a" -> (3L, 0L), "b" -> (1L, 0L), "c" -> (1L, 0L))
+ table.isFull shouldBe false
+ }
+
+ it should "keep a heavy key through a flood of distinct keys, within its error bound" in {
+ val capacity = 10
+ val table = new HeavyHitters[String, Hits](capacity, () => new Hits)
+ val truth = mutable.Map.empty[String, Long].withDefaultValue(0L)
+ def offer(k: String): Unit = { truth(k) += 1; table.offer(k)(_.n += 1) }
+ // 300 offers of one key among 1,000 distinct keys: 300 > 1,300 / 10, so it must be kept.
+ (1 to 1000).foreach { i => offer(s"noise-$i"); if (i % 10 == 0) (1 to 3).foreach(_ => offer("heavy")) }
+
+ val entries = table.snapshot
+ entries.size shouldBe capacity
+ entries.map(_.key) should contain("heavy")
+ entries.foreach { e =>
+ withClue(s"${e.key}: count ${e.count}, error ${e.error}, true ${truth(e.key)}: ") {
+ truth(e.key) should be <= e.count
+ truth(e.key) should be >= (e.count - e.error)
+ }
+ }
+ table.totalOffered shouldBe 1300L
+ }
+
+ it should "start a replacement entry's details afresh" in {
+ val table = new HeavyHitters[String, Hits](1, () => new Hits)
+ table.offer("a")(_.n += 1)
+ table.offer("a")(_.n += 1)
+ table.offer("b")(_.n += 1) // replaces "a": count 3, error 2, details from this offer only
+ val only = table.snapshot.head
+ (only.key, only.count, only.error, only.details.n) shouldBe (("b", 3L, 2L, 1))
+ }
+}
diff --git a/obp-api/src/test/scala/code/telemetry/TrafficSourcesTest.scala b/obp-api/src/test/scala/code/telemetry/TrafficSourcesTest.scala
new file mode 100644
index 0000000000..a48cad20a9
--- /dev/null
+++ b/obp-api/src/test/scala/code/telemetry/TrafficSourcesTest.scala
@@ -0,0 +1,65 @@
+package code.telemetry
+
+import code.telemetry.TrafficSources.{AddressCaller, ConsumerCaller, Note}
+import org.scalatest.{BeforeAndAfterEach, FlatSpec, Matchers}
+
+/**
+ * This suite checks what TrafficSources records for a request and how it merges minutes into a
+ * window. Times are passed in, so the minutes are chosen by the test.
+ */
+class TrafficSourcesTest extends FlatSpec with Matchers with BeforeAndAfterEach {
+
+ override def beforeEach(): Unit = TrafficSources.clear()
+
+ private val minute = 60000L
+ private val t0 = 1790000000000L - 1790000000000L % minute // the start of a minute
+
+ private def note(operationId: Option[String] = None, consumer: Option[String] = None, refusedBy: Option[String] = None): Note = {
+ val n = new Note
+ n.operationId = operationId
+ n.apiVersion = operationId.map(_ => "v7.0.0")
+ n.consumerId = consumer
+ n.consumerName = consumer.map(c => s"app $c")
+ n.refusedBy = refusedBy
+ n
+ }
+
+ "TrafficSources" should "count an authenticated request under its Consumer and its address, and pair the Consumer with the endpoint" in {
+ TrafficSources.record(note(Some("OBPv7.0.0-getBanks"), Some("consumer-1")), "198.51.100.1", 200, 12, t0)
+
+ val consumers = TrafficSources.consumers(1, t0)
+ consumers.map(c => (c.key, c.requests)) shouldBe List(("consumer-1", 1L))
+ consumers.head.details.head.name shouldBe "app consumer-1"
+ consumers.head.details.head.lastIp shouldBe "198.51.100.1"
+
+ val addresses = TrafficSources.addresses(1, t0)
+ addresses.map(_.key) shouldBe List("198.51.100.1")
+ addresses.head.details.head.consumers.toList shouldBe List("consumer-1")
+
+ TrafficSources.callerEndpoints(1, t0).map(_.key) shouldBe List((ConsumerCaller("consumer-1"), "OBPv7.0.0-getBanks"))
+ }
+
+ it should "count an anonymous request under its address only, and group unknown paths as unmatched" in {
+ TrafficSources.record(note(), "203.0.113.9", 404, 2, t0)
+ TrafficSources.consumers(1, t0) shouldBe empty
+ TrafficSources.addresses(1, t0).head.details.head.unmatched shouldBe 1L
+ TrafficSources.callerEndpoints(1, t0).map(_.key) shouldBe List((AddressCaller("203.0.113.9"), TrafficSources.UnmatchedEndpoint))
+ }
+
+ it should "record a refusal under the limiter that refused" in {
+ TrafficSources.record(note(refusedBy = Some("ip_penalty")), "203.0.113.9", 429, 1, t0)
+ val pair = TrafficSources.callerEndpoints(1, t0).head
+ pair.key._2 shouldBe "refused:ip_penalty"
+ pair.details.head.refused shouldBe 1L
+ }
+
+ it should "merge minutes into the window asked for, and leave out older minutes" in {
+ (0 until 3).foreach { m =>
+ TrafficSources.record(note(Some("OBPv7.0.0-getBanks")), "203.0.113.9", 200, 5, t0 + m * minute)
+ }
+ val now = t0 + 2 * minute
+ TrafficSources.addresses(1, now).head.requests shouldBe 1L
+ TrafficSources.addresses(5, now).head.requests shouldBe 3L
+ TrafficSources.addresses(5, now).head.error shouldBe 0L
+ }
+}
diff --git a/release_notes.md b/release_notes.md
index b1c99d43b2..e7004d9ef8 100644
--- a/release_notes.md
+++ b/release_notes.md
@@ -3,6 +3,13 @@
### Most recent changes at top of file
```
Date Commit Action
+28/09/2026 TBD NEW in v7.0.0: GET /management/traffic/top-callers?window=1|5|15, where the
+ traffic on the answering instance is coming from: the busiest Consumers,
+ client IP addresses (every request, authenticated or not), and callers and
+ endpoints, over the last 1, 5 or 15 minutes. Estimated counts from
+ fixed-size tables (the Space-Saving algorithm of Metwally, Agrawal and
+ El Abbadi, 2005), in memory only, 15 minutes kept. New Role
+ CanGetTrafficSources (empty bank id). New error code OBP-10067.
28/09/2026 TBD NEW cache namespaces message_docs and glossary, shown and invalidated like
the others (system/cache in API Manager, POST
/management/cache/namespaces/invalidate).