Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 28 additions & 0 deletions docs/telemetry_conventions.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

6 changes: 6 additions & 0 deletions obp-api/src/main/scala/code/api/util/ApiRole.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand Down
1 change: 1 addition & 0 deletions obp-api/src/main/scala/code/api/util/ErrorMessages.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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".
Expand Down
10 changes: 10 additions & 0 deletions obp-api/src/main/scala/code/api/util/http4s/Http4sApp.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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}") }
})
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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" =>
Expand All @@ -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.
//
Expand All @@ -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))
}
}
11 changes: 11 additions & 0 deletions obp-api/src/main/scala/code/api/util/http4s/Http4sSupport.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
Expand All @@ -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))
}
}
Expand Down
3 changes: 3 additions & 0 deletions obp-api/src/main/scala/code/api/v7_0_0/Http4s700.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
120 changes: 120 additions & 0 deletions obp-api/src/main/scala/code/api/v7_0_0/Http4s700TrafficSources.scala
Original file line number Diff line number Diff line change
@@ -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 <http://www.gnu.org/licenses/>.

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)
)
}
Loading
Loading