From 9ec7e9a2b9dad3fe663d5fd9f46302d63f1a2b7b Mon Sep 17 00:00:00 2001 From: simonredfern Date: Mon, 28 Sep 2026 15:21:31 +0200 Subject: [PATCH] Top callers: HeavyHitters (Space-Saving, Metwally, Agrawal, El Abbadi 2005) tables per minute of busiest Consumers, client IP addresses and caller+endpoint, counting every request including refused and unknown paths; v7.0.0 GET /management/traffic/top-callers?window=1|5|15 with CanGetTrafficSources. --- docs/telemetry_conventions.md | 28 +++ .../main/scala/code/api/util/ApiRole.scala | 6 + .../scala/code/api/util/ErrorMessages.scala | 1 + .../code/api/util/http4s/Http4sApp.scala | 10 + .../api/util/http4s/Http4sResourceDocs.scala | 26 +- .../code/api/util/http4s/Http4sSupport.scala | 11 + .../util/http4s/ResourceDocMiddleware.scala | 12 + .../SelfServiceRateLimitMiddleware.scala | 8 +- .../scala/code/api/v7_0_0/Http4s700.scala | 3 + .../api/v7_0_0/Http4s700TrafficSources.scala | 120 +++++++++ .../api/v7_0_0/JSONFactory700Operations.scala | 130 +++++++++- .../scala/code/telemetry/HeavyHitters.scala | 111 +++++++++ .../scala/code/telemetry/TrafficSources.scala | 230 ++++++++++++++++++ .../v7_0_0/TrafficSourcesEndpointTest.scala | 71 ++++++ .../code/telemetry/HeavyHittersTest.scala | 51 ++++ .../code/telemetry/TrafficSourcesTest.scala | 65 +++++ release_notes.md | 7 + 17 files changed, 877 insertions(+), 13 deletions(-) create mode 100644 obp-api/src/main/scala/code/api/v7_0_0/Http4s700TrafficSources.scala create mode 100644 obp-api/src/main/scala/code/telemetry/HeavyHitters.scala create mode 100644 obp-api/src/main/scala/code/telemetry/TrafficSources.scala create mode 100644 obp-api/src/test/scala/code/api/v7_0_0/TrafficSourcesEndpointTest.scala create mode 100644 obp-api/src/test/scala/code/telemetry/HeavyHittersTest.scala create mode 100644 obp-api/src/test/scala/code/telemetry/TrafficSourcesTest.scala 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) < "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) < "7")).code should equal(400) + makeGetRequest(topCallers.GET <@ (user1) < "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).