From e71c90ca0784cb813a39ffd82eca0d44342d82a4 Mon Sep 17 00:00:00 2001 From: simonredfern Date: Mon, 28 Sep 2026 12:14:45 +0200 Subject: [PATCH 1/3] IP penalty for shortlived bad acting IP addresses --- docs/telemetry_conventions.md | 2 + .../main/scala/code/api/util/ApiRole.scala | 11 ++ .../scala/code/api/util/ErrorMessages.scala | 6 + .../main/scala/code/api/util/Glossary.scala | 2 + .../scala/code/api/util/IpPenalties.scala | 180 ++++++++++++++++++ .../SelfServiceRateLimitMiddleware.scala | 33 +++- .../scala/code/api/v7_0_0/Http4s700.scala | 3 + .../api/v7_0_0/Http4s700IpPenalties.scala | 175 +++++++++++++++++ .../code/api/v7_0_0/JSONFactory7.0.0.scala | 31 +++ .../code/telemetry/TelemetryBindings.scala | 5 + .../code/api/v7_0_0/IpPenaltiesTest.scala | 150 +++++++++++++++ release_notes.md | 10 + 12 files changed, 606 insertions(+), 2 deletions(-) create mode 100644 obp-api/src/main/scala/code/api/util/IpPenalties.scala create mode 100644 obp-api/src/main/scala/code/api/v7_0_0/Http4s700IpPenalties.scala create mode 100644 obp-api/src/test/scala/code/api/v7_0_0/IpPenaltiesTest.scala diff --git a/docs/telemetry_conventions.md b/docs/telemetry_conventions.md index 777dd7033c..1445bf026c 100644 --- a/docs/telemetry_conventions.md +++ b/docs/telemetry_conventions.md @@ -288,6 +288,8 @@ thread and class figures on its own port. Remove it once the Micrometer JVM bind | `obp.api.batch_writer.queue.depth` | gauge | `writer` | rows queued minus rows written or lost (the queue's own `size()` walks it) | | `obp.api.batch_writer.flushes` | timer | `writer`, `result` | each flush that had rows | | `obp.api.redis_cache.gets`, `obp.api.redis_cache.sets` | counter | `cache` (`static_resource_docs`, `dynamic_resource_docs`, `all_resource_docs`, `static_swagger`, `financial_products`, `api_products`), `result` (`hit`, `miss`, `error` for gets; `success`, `error` for sets) | `Caching.tryGet` / `trySet`; `error` means Redis was unreachable | +| `obp.api.ip_penalties.active` | gauge | | how many addresses have an IP penalty (never which ones) | +| `obp.api.ip_penalties.refused` | counter | | requests refused with 429 OBP-10062 | | `obp.api.self_service_rate_limit.checks` | counter | `scope`, `outcome` (`allowed`, `warned`, `blocked`, `skipped`) | `SelfServiceRateLimiter.check`; in shadow mode, `warned` shows how often a scope would refuse real traffic | | `obp.api.connector.calls` | timer, fixed buckets | `connector`, `connector_method`, `result` | the Connector proxy (`code/bankconnectors/package.scala`) | | `obp.api.redis.commands` | timer | `command`, `result` | `Redis.use` | 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 27236f080d..9b87e38b36 100644 --- a/obp-api/src/main/scala/code/api/util/ApiRole.scala +++ b/obp-api/src/main/scala/code/api/util/ApiRole.scala @@ -569,6 +569,17 @@ object ApiRole extends MdcLoggable{ case class CanGetTelemetry(requiresBankId: Boolean = false) extends ApiRole lazy val canGetTelemetry = CanGetTelemetry() + // IP penalties restrict an address on the whole instance, which belongs to no bank, so these + // Roles are held at the empty bank id. + case class CanCreateIpPenalty(requiresBankId: Boolean = false) extends ApiRole + lazy val canCreateIpPenalty = CanCreateIpPenalty() + + case class CanGetIpPenalties(requiresBankId: Boolean = false) extends ApiRole + lazy val canGetIpPenalties = CanGetIpPenalties() + + case class CanDeleteIpPenalty(requiresBankId: Boolean = false) extends ApiRole + lazy val canDeleteIpPenalty = CanDeleteIpPenalty() + 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 f4531dfe74..596a6270cf 100644 --- a/obp-api/src/main/scala/code/api/util/ErrorMessages.scala +++ b/obp-api/src/main/scala/code/api/util/ErrorMessages.scala @@ -138,6 +138,12 @@ object ErrorMessages { // OBP-10061 the authentication limiter (AuthRateLimiter, inside the credential check, keyed by IP and account) val TooManyRequestsSelfService = "OBP-10060: Too Many Requests for a self-service endpoint." val TooManyRequestsAuth = "OBP-10061: Too Many Requests for authentication. Too many login attempts from this address or for this account." + // OBP-10062 an IP penalty (IpPenalties, before everything else, an operator's temporary limit on one address) + val TooManyRequestsIpPenalty = "OBP-10062: Too Many Requests. This address is under a temporary rate limit set by an operator." + 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 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". // See SelfServiceRateLimiter.warningMessage. diff --git a/obp-api/src/main/scala/code/api/util/Glossary.scala b/obp-api/src/main/scala/code/api/util/Glossary.scala index a140ab3bc0..59998a4074 100644 --- a/obp-api/src/main/scala/code/api/util/Glossary.scala +++ b/obp-api/src/main/scala/code/api/util/Glossary.scala @@ -763,6 +763,8 @@ object Glossary extends MdcLoggable { |2. **Authentication limiter** (`auth.rate_limit.*`) runs inside the credential check of Direct Login, DAuth, Gateway Login and SIWE, before the password or token is verified, keyed by IP address and by account. It defends against brute force, credential stuffing and lockout attacks. Trip code: `OBP-10061`. |3. **Consumer quota** (the limits described above) runs after authentication, keyed by Consumer, or by IP address with a single hourly ceiling for anonymous calls. It is the commercial and fair-use quota. Trip code: `OBP-10018`. | + |Before all three, an operator can put a single IP address under a temporary **IP penalty**: a per-minute limit on every endpoint, for a set time, for example during a scan or denial-of-service attempt (`POST /obp/v7.0.0/management/ip-penalties`, Role CanCreateIpPenalty). A per-minute limit of 0 refuses every request. Penalties are always enforced, shared by every instance, and disappear when they expire; the penalty endpoints themselves are never refused, so a mistake can be undone. Trip code: `OBP-10062`. + | |A login attempt is counted by the authentication limiter only; it is not a self-service scope, so no attempt is counted twice. Every limiter counts in Redis and fails open: a Redis outage never blocks a call. | |### Self-service rate limiting (per IP address, before any credential) diff --git a/obp-api/src/main/scala/code/api/util/IpPenalties.scala b/obp-api/src/main/scala/code/api/util/IpPenalties.scala new file mode 100644 index 0000000000..13ea8539ae --- /dev/null +++ b/obp-api/src/main/scala/code/api/util/IpPenalties.scala @@ -0,0 +1,180 @@ +/** +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.util + +import java.util.concurrent.atomic.AtomicReference + +import code.api.JedisMethod +import code.api.cache.Redis +import code.util.Helper.MdcLoggable +import com.google.common.net.InetAddresses +import com.openbankproject.commons.util.JsonAliases.{compactRender, parse} +import org.json4s.{Extraction, Formats} + +import scala.util.Try + +/** + * This object keeps the list of IP addresses an operator has put under a temporary rate limit, a + * "penalty", and checks requests against it. + * + * It exists for incidents like the NMB scan of 2026-09-23: when one address is hammering the + * instance, an operator can slow that address right down for a while, on every endpoint, without a + * restart or a props change. A penalty is a per-minute limit for one address (0 refuses every + * request) with an expiry time; when it expires it disappears by itself, so nothing stays + * restricted by mistake. Penalties are always enforced: an operator set them on purpose. + * + * Each penalty is one Redis key with a TTL, so every instance sees the same list. Each instance + * reads it into a local copy at most every [[RefreshMillis]], so checking a request costs a map + * lookup, not a Redis call, even under attack. Reading fails open: with Redis unreachable there are + * no penalties, and requests are not refused because of the list. + * + * The address is the one `Http4sCallContextBuilder.clientIp` resolves, the same one the other + * per-IP limits use. Traffic that reaches OBP-API through a proxy which does not pass on the client + * address shows the proxy's address, and penalising that would restrict every user behind it. + */ +object IpPenalties extends MdcLoggable { + + implicit private val formats: Formats = CustomJsonFormats.formats + + /** One penalty. Times are epoch milliseconds. `perMinuteLimit` 0 refuses every request. */ + final case class Penalty( + ipAddress: String, + perMinuteLimit: Long, + reason: String, + createdByUserId: String, + createdAtMillis: Long, + expiresAtMillis: Long + ) + + val RefreshMillis: Long = 5000L + val MaxDurationMinutes: Long = 7L * 24 * 60 + val MaxReasonLength: Int = 255 + + private def keyPrefix: String = s"${code.api.Constant.getGlobalCacheNamespacePrefix}ip_penalty_" + private def keyFor(ipAddress: String): String = keyPrefix + ipAddress + private def counterKeyFor(ipAddress: String): String = s"${keyPrefix}count_$ipAddress" + + /** An IPv4 or IPv6 literal, never a host name (so no lookup happens). Returned in canonical form. */ + def canonicalAddress(value: String): Option[String] = + Option(value).map(_.trim).filter(v => v.nonEmpty && InetAddresses.isInetAddress(v)) + .map(v => InetAddresses.toAddrString(InetAddresses.forString(v))) + + // ===== The list ===== + + private final case class Snapshot(loadedAtMillis: Long, penalties: Map[String, Penalty]) + private val snapshot = new AtomicReference[Snapshot](Snapshot(0L, Map.empty)) + + private def readAll(): Map[String, Penalty] = + Redis.scanKeys(s"$keyPrefix*") + .filterNot(_.startsWith(s"${keyPrefix}count_")) + .flatMap(key => Try(Redis.use(JedisMethod.GET, key, None, None)).toOption.flatten) + .flatMap(json => Try(parse(json).extract[Penalty]).toOption) + .filter(_.expiresAtMillis > System.currentTimeMillis()) + .map(p => p.ipAddress -> p).toMap + + /** The penalties in force, from the local copy (refreshed at most every [[RefreshMillis]]). */ + def active(): Map[String, Penalty] = { + val now = System.currentTimeMillis() + val current = snapshot.get() + if (now - current.loadedAtMillis < RefreshMillis) current.penalties.filter(_._2.expiresAtMillis > now) + else { + val penalties = try readAll() catch { + case e: Throwable => + logger.debug(s"IpPenalties.active says: could not read penalties, treating the list as empty: ${e.getMessage}") + Map.empty[String, Penalty] + } + snapshot.set(Snapshot(now, penalties)) + penalties + } + } + + /** Forget the local copy, so the next check reads Redis again (after a change on this instance). */ + def refresh(): Unit = snapshot.set(Snapshot(0L, Map.empty)) + + def find(ipAddress: String): Option[Penalty] = canonicalAddress(ipAddress).flatMap(active().get) + + /** Every penalty in Redis now, soonest to expire first. For the management endpoints, which must not show a stale copy. */ + def listAll(): List[Penalty] = + Try(readAll()).getOrElse(Map.empty[String, Penalty]).values.toList.sortBy(_.expiresAtMillis) + + /** Whether the address has a penalty in Redis now. */ + def exists(ipAddress: String): Boolean = + canonicalAddress(ipAddress).exists(address => Try(readAll().contains(address)).getOrElse(false)) + + /** Adds a penalty. Fails when the address already has one: remove it first to change it. */ + def add(ipAddress: String, perMinuteLimit: Long, durationMinutes: Long, reason: String, createdByUserId: String): Either[String, Penalty] = + canonicalAddress(ipAddress) match { + case None => Left(ErrorMessages.InvalidIpAddress) + case Some(address) if Try(readAll().contains(address)).getOrElse(false) => Left(ErrorMessages.IpPenaltyAlreadyExists) + case Some(address) => + val now = System.currentTimeMillis() + val penalty = Penalty(address, perMinuteLimit, reason, createdByUserId, now, now + durationMinutes * 60000L) + Redis.use(JedisMethod.SET, keyFor(address), Some((durationMinutes * 60).toInt), Some(compactRender(Extraction.decompose(penalty)))) + refresh() + logger.warn(s"IpPenalties.add says: $address limited to $perMinuteLimit per minute for $durationMinutes minutes by user $createdByUserId: $reason") + Right(penalty) + } + + /** Removes a penalty. False when the address had none. */ + def remove(ipAddress: String): Boolean = + canonicalAddress(ipAddress) match { + case Some(address) if Try(readAll().contains(address)).getOrElse(false) => + Redis.use(JedisMethod.DELETE, keyFor(address), None, None) + Redis.use(JedisMethod.DELETE, counterKeyFor(address), None, None) + refresh() + logger.warn(s"IpPenalties.remove says: penalty on $address removed") + true + case _ => false + } + + // ===== Checking a request ===== + + /** A request from a penalised address that went over its limit, with seconds until the minute resets. */ + final case class Refusal(penalty: Penalty, retryAfterSeconds: Long) + + /** + * Counts a request from `ipAddress` against its penalty, if it has one. Returns a refusal when the + * address is over its per-minute limit; otherwise None. Addresses without a penalty cost a map + * lookup and nothing else. + */ + def check(ipAddress: String): Option[Refusal] = + find(ipAddress).flatMap { penalty => + val (ttl, current) = + if (penalty.perMinuteLimit == 0L) (60L, 1L) // refuse every request; no need to count + else RateLimitingUtil.incrementCounter(counterKeyFor(penalty.ipAddress), RateLimitingPeriod.PER_MINUTE) + // current == -1: Redis unavailable for the counter. Fail open, like the other limiters. + if (current >= 0 && current > penalty.perMinuteLimit) { + code.telemetry.Telemetry.counter("obp.api.ip_penalties.refused").increment() + Some(Refusal(penalty, math.max(1L, ttl))) + } else None + } + + /** The 429 body text for a refused request. */ + def refusedMessage(refusal: Refusal): String = + s"${ErrorMessages.TooManyRequestsIpPenalty} This address is limited to ${refusal.penalty.perMinuteLimit} requests per minute until " + + s"${java.time.Instant.ofEpochMilli(refusal.penalty.expiresAtMillis)}." +} 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 be15d3eaf3..7299e64f00 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 @@ -28,7 +28,7 @@ TESOBE (http://www.tesobe.com/) package code.api.util.http4s import cats.effect.IO -import code.api.util.SelfServiceRateLimiter +import code.api.util.{IpPenalties, SelfServiceRateLimiter} import code.api.util.SelfServiceRateLimiter.{Blocked, Outcome, Skipped, Warned, Window} import code.util.Helper.MdcLoggable import org.http4s.{Header, Headers, Method, Request, Response, Status} @@ -128,7 +128,36 @@ object SelfServiceRateLimitMiddleware extends MdcLoggable { } /** Wrap the application. */ - def apply(req: Request[IO])(run: Request[IO] => IO[Response[IO]]): IO[Response[IO]] = + /** + * The penalty management endpoints are never refused because of a penalty, so an operator who + * penalised the wrong address (their own, or a proxy's) can always undo it. + */ + private val ipPenaltyManagementPath = "^/obp/[^/]+/management/ip-penalties(/.*)?$".r + + def apply(req: Request[IO])(run: Request[IO] => IO[Response[IO]]): IO[Response[IO]] = { + val penaltyRefusal: IO[Option[IpPenalties.Refusal]] = + 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 None => applyScopes(req)(run) + } + } + + private def penaltyResponse(refusal: IpPenalties.Refusal): Response[IO] = { + val escaped = IpPenalties.refusedMessage(refusal).replace("\\", "\\\\").replace("\"", "\\\"") + Response[IO](status = Status.TooManyRequests) + .withEntity(s"""{"code":429,"message":"$escaped"}""".getBytes("UTF-8")) + .withHeaders(Headers( + Header.Raw(CIString("Content-Type"), "application/json; charset=utf-8"), + Header.Raw(CIString("Retry-After"), refusal.retryAfterSeconds.toString), + Header.Raw(CIString(LimitHeader), refusal.penalty.perMinuteLimit.toString), + Header.Raw(CIString(RemainingHeader), "0"), + Header.Raw(CIString(ResetHeader), refusal.retryAfterSeconds.toString) + )) + } + + private def applyScopes(req: Request[IO])(run: Request[IO] => IO[Response[IO]]): IO[Response[IO]] = scopeFor(req) match { case None => run(req) case Some(scope) => 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 a24cf0a720..da808df931 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 @@ -7246,6 +7246,9 @@ object Http4s700 { // Telemetry, for people; Prometheus reads the separate Telemetry port instead. resourceDocs ++= Http4s700Telemetry.resourceDocs + // IP penalties: an operator's temporary per-minute limit on one address. + resourceDocs ++= Http4s700IpPenalties.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/Http4s700IpPenalties.scala b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700IpPenalties.scala new file mode 100644 index 0000000000..4ad67bfd74 --- /dev/null +++ b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700IpPenalties.scala @@ -0,0 +1,175 @@ +/** +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, IpPenalties} +import code.api.v7_0_0.JSONFactory700.{IpPenaltiesJsonV700, IpPenaltyJsonV700, PostIpPenaltyJsonV700} +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 endpoints that manage IP penalties: temporary per-minute limits an + * operator puts on one IP address during an incident (see [[code.api.util.IpPenalties]]). + * + * It is declared in its own object to keep Http4s700's initialiser under the JVM's 64KB method limit. + */ +object Http4s700IpPenalties { + + 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 penaltyDescription = + s"""An IP penalty limits one IP address to a number of requests per minute, on every endpoint of + |this instance, until it expires. A `per_minute_limit` of 0 refuses every request. A refused + |request gets 429 `${TooManyRequestsIpPenalty.takeWhile(_ != ':')}`. Penalties are always + |enforced, are shared by every instance (they are kept in Redis), and disappear by themselves + |when they expire. The endpoints under `/management/ip-penalties` are never refused because of a + |penalty, so a mistake can always be undone. + | + |The address is the client address OBP-API resolves for the request. Behind a proxy that does not + |pass on the client's address, that is the proxy's address, and a penalty on it restricts + |everyone behind the proxy.""".stripMargin + + // Route: POST /obp/v7.0.0/management/ip-penalties + lazy val createIpPenalty: HttpRoutes[IO] = HttpRoutes.of[IO] { + case req @ POST -> `prefixPath` / "management" / "ip-penalties" => + EndpointHelpers.withUserAndBodyCreated[PostIpPenaltyJsonV700, IpPenaltyJsonV700](req) { (user, body, cc) => + for { + _ <- Helper.booleanToFuture(InvalidIpPenalty, failCode = 400, cc = Some(cc)) { + body.per_minute_limit >= 0 && + body.duration_minutes >= 1 && body.duration_minutes <= IpPenalties.MaxDurationMinutes && + body.reason != null && body.reason.trim.nonEmpty && body.reason.length <= IpPenalties.MaxReasonLength + } + _ <- Helper.booleanToFuture(s"$InvalidIpAddress Current value is ${body.ip_address}", failCode = 400, cc = Some(cc)) { + IpPenalties.canonicalAddress(body.ip_address).isDefined + } + _ <- Helper.booleanToFuture(IpPenaltyAlreadyExists, failCode = 409, cc = Some(cc)) { + !IpPenalties.exists(body.ip_address) + } + added <- Future(IpPenalties.add(body.ip_address, body.per_minute_limit, body.duration_minutes, body.reason.trim, user.userId)) + _ <- Helper.booleanToFuture(added.left.getOrElse(UnknownError), failCode = 409, cc = Some(cc))(added.isRight) + } yield JSONFactory700.createIpPenaltyJson(added.toOption.get) + } + } + + resourceDocs += ResourceDoc( + implementedInApiVersion, + nameOf(createIpPenalty), + "POST", + "/management/ip-penalties", + "Create IP Penalty", + s"""Put one IP address under a temporary per-minute limit, for example during a denial-of-service + |or scanning incident. + | + |$penaltyDescription + | + |`duration_minutes` is from 1 to ${IpPenalties.MaxDurationMinutes} (one week). An address can have one + |penalty at a time: to change it, delete it and create it again (409 when one exists). + |""".stripMargin, + JSONFactory700.postIpPenaltyJsonV700Example, + JSONFactory700.ipPenaltyJsonV700Example, + List($AuthenticatedUserIsRequired, UserHasMissingRoles, InvalidJsonFormat, InvalidIpPenalty, InvalidIpAddress, + IpPenaltyAlreadyExists, UnknownError), + List(apiTagRateLimits, apiTagSystem), + Some(List(canCreateIpPenalty)), + http4sPartialFunction = Some(createIpPenalty) + ) + + // Route: GET /obp/v7.0.0/management/ip-penalties + lazy val getIpPenalties: HttpRoutes[IO] = HttpRoutes.of[IO] { + case req @ GET -> `prefixPath` / "management" / "ip-penalties" => + EndpointHelpers.withUser(req) { (_, _) => + Future(IpPenaltiesJsonV700(IpPenalties.listAll().map(JSONFactory700.createIpPenaltyJson))) + } + } + + resourceDocs += ResourceDoc( + implementedInApiVersion, + nameOf(getIpPenalties), + "GET", + "/management/ip-penalties", + "Get IP Penalties", + s"""Get the IP penalties in force, soonest to expire first. + | + |$penaltyDescription + |""".stripMargin, + EmptyBody, + JSONFactory700.ipPenaltiesJsonV700Example, + List($AuthenticatedUserIsRequired, UserHasMissingRoles, UnknownError), + List(apiTagRateLimits, apiTagSystem), + Some(List(canGetIpPenalties)), + http4sPartialFunction = Some(getIpPenalties) + ) + + // Route: DELETE /obp/v7.0.0/management/ip-penalties/IP_ADDRESS + lazy val deleteIpPenalty: HttpRoutes[IO] = HttpRoutes.of[IO] { + case req @ DELETE -> `prefixPath` / "management" / "ip-penalties" / ipAddress => + EndpointHelpers.withUserDelete(req) { (_, cc) => + for { + removed <- Future(IpPenalties.remove(ipAddress)) + _ <- Helper.booleanToFuture(s"$IpPenaltyNotFound Current value is $ipAddress", failCode = 404, cc = Some(cc))(removed) + } yield () + } + } + + resourceDocs += ResourceDoc( + implementedInApiVersion, + nameOf(deleteIpPenalty), + "DELETE", + "/management/ip-penalties/IP_ADDRESS", + "Delete IP Penalty", + s"""Remove the penalty on one IP address before it expires. + | + |$penaltyDescription + |""".stripMargin, + EmptyBody, + EmptyBody, + List($AuthenticatedUserIsRequired, UserHasMissingRoles, IpPenaltyNotFound, UnknownError), + List(apiTagRateLimits, apiTagSystem), + Some(List(canDeleteIpPenalty)), + http4sPartialFunction = Some(deleteIpPenalty) + ) +} diff --git a/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory7.0.0.scala b/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory7.0.0.scala index 2505309962..47aedf1f87 100644 --- a/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory7.0.0.scala +++ b/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory7.0.0.scala @@ -2772,6 +2772,37 @@ object JSONFactory700 extends MdcLoggable with code.api.util.CustomJsonFormats { ) lazy val apiProductSubscriptionsJsonV700Example = ApiProductSubscriptionsJsonV700(List(apiProductSubscriptionJsonV700Example)) + // ===== IP penalties ===== + + /** Request body of POST /management/ip-penalties. `per_minute_limit` 0 refuses every request. */ + case class PostIpPenaltyJsonV700(ip_address: String, per_minute_limit: Long, duration_minutes: Long, reason: String) + + case class IpPenaltyJsonV700( + ip_address: String, + per_minute_limit: Long, + reason: String, + created_by_user_id: String, + created_at: Date, + expires_at: Date + ) + case class IpPenaltiesJsonV700(ip_penalties: List[IpPenaltyJsonV700]) + + def createIpPenaltyJson(penalty: code.api.util.IpPenalties.Penalty): IpPenaltyJsonV700 = + IpPenaltyJsonV700(penalty.ipAddress, penalty.perMinuteLimit, penalty.reason, penalty.createdByUserId, + new Date(penalty.createdAtMillis), new Date(penalty.expiresAtMillis)) + + lazy val postIpPenaltyJsonV700Example = PostIpPenaltyJsonV700( + ip_address = "203.0.113.42", per_minute_limit = 10, duration_minutes = 60, + reason = "Vulnerability scan: about 900 requests a minute to resource-docs with changing filters") + + lazy val ipPenaltyJsonV700Example = IpPenaltyJsonV700( + ip_address = "203.0.113.42", per_minute_limit = 10, + reason = "Vulnerability scan: about 900 requests a minute to resource-docs with changing filters", + created_by_user_id = ExampleValue.userIdExample.value, + created_at = APIUtil.DateWithMsExampleObject, expires_at = APIUtil.DateWithMsExampleObject) + + lazy val ipPenaltiesJsonV700Example = IpPenaltiesJsonV700(List(ipPenaltyJsonV700Example)) + // ===== Telemetry ===== /** One meter: its Micrometer name, type, unit, tags and current values (count, total_time, max, value ...). */ diff --git a/obp-api/src/main/scala/code/telemetry/TelemetryBindings.scala b/obp-api/src/main/scala/code/telemetry/TelemetryBindings.scala index 49f1394ccb..e0c8987884 100644 --- a/obp-api/src/main/scala/code/telemetry/TelemetryBindings.scala +++ b/obp-api/src/main/scala/code/telemetry/TelemetryBindings.scala @@ -47,8 +47,13 @@ object TelemetryBindings { bindLogging() bindMessageDocs() bindRedisLogger() + bindIpPenalties() } + /** How many addresses are under an operator's temporary limit. Never which ones: an address is not a tag. */ + private def bindIpPenalties(): Unit = + Telemetry.gauge("obp.api.ip_penalties.active")(code.api.util.IpPenalties.active().size.toDouble) + /** The log dispatch pool (Helper.MdcLoggable) and log masking. */ private def bindLogging(): Unit = { Telemetry.gauge("obp.api.log.dispatch.queue.depth")(Helper.mdcLogQueueDepth.toDouble) diff --git a/obp-api/src/test/scala/code/api/v7_0_0/IpPenaltiesTest.scala b/obp-api/src/test/scala/code/api/v7_0_0/IpPenaltiesTest.scala new file mode 100644 index 0000000000..d24943df14 --- /dev/null +++ b/obp-api/src/test/scala/code/api/v7_0_0/IpPenaltiesTest.scala @@ -0,0 +1,150 @@ +/** +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 code.api.util.APIUtil.OAuth._ +import code.api.util.ApiRole.{CanCreateIpPenalty, CanDeleteIpPenalty, CanGetIpPenalties} +import code.api.util.ErrorMessages.{AuthenticatedUserIsRequired, UserHasMissingRoles} +import code.api.util.IpPenalties +import code.api.v6_0_0.V600ServerSetup +import code.api.v7_0_0.JSONFactory700.{IpPenaltiesJsonV700, IpPenaltyJsonV700, PostIpPenaltyJsonV700} +import code.entitlement.Entitlement +import com.openbankproject.commons.model.ErrorMessage +import com.openbankproject.commons.util.ApiVersion +import org.json4s.native.Serialization.write +import org.scalatest.Tag + +/** + * This suite checks the IP penalty endpoints (create, list, delete, each behind its own Role) and the + * penalty itself: a penalised address is refused with 429 OBP-10062 once over its per-minute limit, + * on any endpoint, while the penalty endpoints stay reachable so a mistake can be undone. + */ +class IpPenaltiesTest 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("ipPenalties") + + private def penalties = v7_0_0_Request / "management" / "ip-penalties" + + private def withRoles[T](body: => T): T = { + val granted = List(CanCreateIpPenalty, CanGetIpPenalties, CanDeleteIpPenalty) + .map(role => Entitlement.entitlement.vend.addEntitlement("", resourceUser1.userId, role.toString)) + try body finally granted.foreach(Entitlement.entitlement.vend.deleteEntitlement) + } + + private def post(json: PostIpPenaltyJsonV700) = makePostRequest(penalties.POST <@ (user1), write(json)) + + // Addresses from the documentation ranges (RFC 5737), never a real client. + private val documentationAddress = "203.0.113.77" + + feature(s"IP penalties - /obp/v7.0.0/management/ip-penalties - $VersionOfApi") { + + scenario("anonymous access is 401 and a user without the Roles gets 403", ApiEndpoint, VersionOfApi) { + makeGetRequest(penalties.GET).code should equal(401) + makeGetRequest(penalties.GET).body.extract[ErrorMessage].message should equal(AuthenticatedUserIsRequired) + + val get = makeGetRequest(penalties.GET <@ (user1)) + get.code should equal(403) + get.body.extract[ErrorMessage].message should equal(UserHasMissingRoles + CanGetIpPenalties) + + post(PostIpPenaltyJsonV700(documentationAddress, 10, 60, "test")).code should equal(403) + makeDeleteRequest((penalties / documentationAddress).DELETE <@ (user1)).code should equal(403) + } + + scenario("create, list and delete a penalty", ApiEndpoint, VersionOfApi) { + withRoles { + try { + val created = post(PostIpPenaltyJsonV700(documentationAddress, 10, 60, "scanning resource-docs")) + created.code should equal(201) + val penalty = created.body.extract[IpPenaltyJsonV700] + penalty.ip_address should equal(documentationAddress) + penalty.per_minute_limit should equal(10L) + penalty.created_by_user_id should equal(resourceUser1.userId) + (penalty.expires_at.getTime - penalty.created_at.getTime) should equal(60L * 60 * 1000) + + Then("a second penalty on the same address is a conflict") + post(PostIpPenaltyJsonV700(documentationAddress, 5, 30, "again")).code should equal(409) + + Then("the list shows it") + val listed = makeGetRequest(penalties.GET <@ (user1)) + listed.code should equal(200) + listed.body.extract[IpPenaltiesJsonV700].ip_penalties.map(_.ip_address) should contain(documentationAddress) + + Then("delete removes it, and deleting again is 404") + makeDeleteRequest((penalties / documentationAddress).DELETE <@ (user1)).code should equal(204) + makeDeleteRequest((penalties / documentationAddress).DELETE <@ (user1)).code should equal(404) + } finally IpPenalties.remove(documentationAddress) + } + } + + scenario("invalid input is refused with 400", ApiEndpoint, VersionOfApi) { + withRoles { + post(PostIpPenaltyJsonV700("example.com", 10, 60, "a host name is not an address")).code should equal(400) + post(PostIpPenaltyJsonV700(documentationAddress, -1, 60, "negative limit")).code should equal(400) + post(PostIpPenaltyJsonV700(documentationAddress, 10, 0, "no duration")).code should equal(400) + post(PostIpPenaltyJsonV700(documentationAddress, 10, IpPenalties.MaxDurationMinutes + 1, "too long")).code should equal(400) + post(PostIpPenaltyJsonV700(documentationAddress, 10, 60, "")).code should equal(400) + IpPenalties.exists(documentationAddress) shouldBe false + } + } + + scenario("a penalised address is refused on any endpoint once over its limit, but can still manage penalties", ApiEndpoint, VersionOfApi) { + // The test server may see this client as either loopback address. + val loopbackAddresses = List("127.0.0.1", "::1") + try { + loopbackAddresses.foreach(address => IpPenalties.add(address, 1, 5, "IpPenaltiesTest", resourceUser1.userId)) + IpPenalties.refresh() + + val first = makeGetRequest((v7_0_0_Request / "banks").GET) + val second = makeGetRequest((v7_0_0_Request / "banks").GET) + first.code should equal(200) + second.code should equal(429) + second.body.extract[ErrorMessage].message should startWith("OBP-10062") + + Then("the penalty endpoints still answer, so the penalty can be removed") + withRoles { + makeGetRequest(penalties.GET <@ (user1)).code should equal(200) + } + } finally { + loopbackAddresses.foreach(IpPenalties.remove) + } + + Then("once removed, the address is served again") + makeGetRequest((v7_0_0_Request / "banks").GET).code should equal(200) + } + + scenario("addresses are compared in canonical form, and host names are not addresses", ApiEndpoint, VersionOfApi) { + IpPenalties.canonicalAddress("2001:DB8:0:0:0:0:0:1") shouldBe Some("2001:db8::1") + IpPenalties.canonicalAddress(" 203.0.113.5 ") shouldBe Some("203.0.113.5") + IpPenalties.canonicalAddress("example.com") shouldBe None + IpPenalties.canonicalAddress("") shouldBe None + } + } +} diff --git a/release_notes.md b/release_notes.md index de3ff124e7..3f14a8bd82 100644 --- a/release_notes.md +++ b/release_notes.md @@ -3,6 +3,16 @@ ### Most recent changes at top of file ``` Date Commit Action +28/09/2026 TBD NEW: IP penalties, an operator's temporary per-minute limit on one IP + address, on every endpoint, checked before all other rate limiters. + per_minute_limit 0 refuses every request. Always enforced, kept in Redis + (shared by every instance), removed automatically on expiry (at most one + week). A refused request gets 429 OBP-10062. + NEW in v7.0.0: POST, GET /management/ip-penalties and + DELETE /management/ip-penalties/IP_ADDRESS, with the new Roles + CanCreateIpPenalty, CanGetIpPenalties, CanDeleteIpPenalty (empty bank id). + The penalty endpoints are never refused, so a mistake can be undone. + NEW error codes: OBP-10062 to OBP-10066. 28/09/2026 TBD NEW rate-limit scope "documentation" in the self-service (per client IP) limiter, covering every public documentation read under any version prefix: resource-docs, message-docs, api/glossary, api/tags, api/versions, From 026c6e946c98680d1c824fae384286cd332a41e6 Mon Sep 17 00:00:00 2001 From: simonredfern Date: Mon, 28 Sep 2026 12:29:18 +0200 Subject: [PATCH 2/3] fix run_trivy --- .github/workflows/run_trivy.yml | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/.github/workflows/run_trivy.yml b/.github/workflows/run_trivy.yml index edd69cbd9b..766ed28d02 100644 --- a/.github/workflows/run_trivy.yml +++ b/.github/workflows/run_trivy.yml @@ -20,7 +20,9 @@ env: jobs: build: runs-on: ubuntu-latest - if: ${{ github.event_name == 'workflow_dispatch' || github.event.workflow_run.conclusion == 'success' }} + # Only scan when build_container.yml actually pushed an image: its docker job + # is gated on the same variable, so without it there is nothing to pull. + if: ${{ vars.ENABLE_CONTAINER_BUILDING == 'true' && (github.event_name == 'workflow_dispatch' || github.event.workflow_run.conclusion == 'success') }} steps: - uses: actions/checkout@v4 @@ -40,8 +42,9 @@ jobs: path: .trivy key: ${{ runner.os }}-trivy-db-${{ steps.trivy-db.outputs.sha }} - name: Run Trivy vulnerability scanner - uses: aquasecurity/trivy-action@ed142fd0673e97e23eac54620cfb913e5ce36c25 + uses: aquasecurity/trivy-action@ed142fd0673e97e23eac54620cfb913e5ce36c25 # v0.36.0 with: + version: "v0.74.0" image-ref: "docker.io/${{ env.DOCKER_HUB_ORGANIZATION }}/${{ env.DOCKER_HUB_REPOSITORY }}:${{ github.event_name == 'workflow_dispatch' && (inputs.image_sha || github.sha) || github.event.workflow_run.head_sha }}" format: "template" template: "@/contrib/sarif.tpl" From abc05741a79b98ddb4c3a2c944e0ac33a2240554 Mon Sep 17 00:00:00 2001 From: simonredfern Date: Mon, 28 Sep 2026 13:35:51 +0200 Subject: [PATCH 3/3] Cache namespaces message_docs and glossary: message docs and JSON Schemas (Redis and in-process) and the in-memory Glossary follow a namespace version to follow the approach of other documentation caches. Move Telemetry and IP penalty JSON types to JSONFactory700Operations json4s readability guard test. --- .../main/scala/code/api/cache/Caching.scala | 8 +- .../scala/code/api/constant/constant.scala | 39 ++++- .../main/scala/code/api/util/Glossary.scala | 9 +- .../code/api/util/JsonSchemaGenerator.scala | 17 ++- .../code/api/util/ResourceDocFilters.scala | 8 + .../api/v2_2_0/MessageDocsJsonCache.scala | 27 ++-- .../scala/code/api/v6_0_0/Http4s600.scala | 5 +- .../code/api/v6_0_0/JSONFactory6.0.0.scala | 15 +- .../api/v7_0_0/Http4s700IpPenalties.scala | 11 +- .../code/api/v7_0_0/Http4s700Telemetry.scala | 4 +- .../code/api/v7_0_0/JSONFactory7.0.0.scala | 91 ------------ .../api/v7_0_0/JSONFactory700Operations.scala | 140 ++++++++++++++++++ .../code/api/JsonFactorySignatureTest.scala | 55 +++++++ .../DocumentationCacheNamespacesTest.scala | 77 ++++++++++ .../code/api/v7_0_0/IpPenaltiesTest.scala | 1 - .../api/v7_0_0/TelemetryEndpointTest.scala | 1 - release_notes.md | 13 +- 17 files changed, 397 insertions(+), 124 deletions(-) create mode 100644 obp-api/src/main/scala/code/api/v7_0_0/JSONFactory700Operations.scala create mode 100644 obp-api/src/test/scala/code/api/JsonFactorySignatureTest.scala create mode 100644 obp-api/src/test/scala/code/api/cache/DocumentationCacheNamespacesTest.scala diff --git a/obp-api/src/main/scala/code/api/cache/Caching.scala b/obp-api/src/main/scala/code/api/cache/Caching.scala index 39ec880b42..24b18792da 100644 --- a/obp-api/src/main/scala/code/api/cache/Caching.scala +++ b/obp-api/src/main/scala/code/api/cache/Caching.scala @@ -143,13 +143,19 @@ object Caching extends MdcLoggable { def setAllResourceDocCache(key: String, value: String): Unit = trySet("all_resource_docs", ALL_RESOURCE_DOC_CACHE_KEY_PREFIX, key, GET_DYNAMIC_RESOURCE_DOCS_TTL, value) - // Also holds the connector JSON Schemas served by v6.0.0 message-docs/CONNECTOR/json-schema. def getStaticSwaggerDocCache(key: String): Option[String] = tryGet("static_swagger", STATIC_SWAGGER_DOC_CACHE_KEY_PREFIX, key, GET_STATIC_RESOURCE_DOCS_TTL) def setStaticSwaggerDocCache(key: String, value: String): Unit = trySet("static_swagger", STATIC_SWAGGER_DOC_CACHE_KEY_PREFIX, key, GET_STATIC_RESOURCE_DOCS_TTL, value) + // The rendered message docs and the JSON Schema of each connector (namespace message_docs). + def getMessageDocsCache(key: String): Option[String] = + tryGet("message_docs", MESSAGE_DOCS_CACHE_KEY_PREFIX, key, GET_STATIC_RESOURCE_DOCS_TTL) + + def setMessageDocsCache(key: String, value: String): Unit = + trySet("message_docs", MESSAGE_DOCS_CACHE_KEY_PREFIX, key, GET_STATIC_RESOURCE_DOCS_TTL, value) + // Fail-safe wrappers around Redis.use. If Redis is unreachable (dev without a // running Redis, transient failure, etc.) we treat it as a miss and recompute instead of failing // the whole request. diff --git a/obp-api/src/main/scala/code/api/constant/constant.scala b/obp-api/src/main/scala/code/api/constant/constant.scala index fb8dc9ca94..30199ea5f9 100644 --- a/obp-api/src/main/scala/code/api/constant/constant.scala +++ b/obp-api/src/main/scala/code/api/constant/constant.scala @@ -164,6 +164,32 @@ object Constant extends MdcLoggable { } } + /** How long [[recentCacheNamespaceVersion]] trusts its copy of a namespace version before re-reading Redis. */ + final val RecentNamespaceVersionMillis = 1000L + + private val recentNamespaceVersions = new java.util.concurrent.ConcurrentHashMap[String, (Long, Long)]() + + /** + * The version of a cache namespace, re-read from Redis at most once every + * [[RecentNamespaceVersionMillis]]. For caches held in each instance's memory (message docs, + * JSON Schemas, the Glossary), which put the version in their own keys so that bumping the + * namespace reaches them too, without a Redis read on every request. After a bump, the instance + * that made it sees the new version at once and the others within a second. + */ + def recentCacheNamespaceVersion(namespaceId: String): Long = { + val now = System.currentTimeMillis() + Option(recentNamespaceVersions.get(namespaceId)) match { + case Some((version, readAt)) if now - readAt < RecentNamespaceVersionMillis => version + case _ => + val version = getCacheNamespaceVersion(namespaceId) + recentNamespaceVersions.put(namespaceId, (version, now)) + version + } + } + + /** Forget the local copies of namespace versions, so the next read goes to Redis (tests). */ + def forgetRecentCacheNamespaceVersions(): Unit = recentNamespaceVersions.clear() + /** * Increment the version counter for a cache namespace. * This effectively invalidates all cached keys in that namespace by making them unreachable. @@ -182,6 +208,8 @@ object Constant extends MdcLoggable { val newVersion = Redis.use(JedisMethod.INCR, versionKey, None, None) .map(_.toLong) logger.info(s"Cache namespace version incremented: ${namespaceId} -> ${newVersion.getOrElse("unknown")}") + // This instance sees the new version at once; the others within RecentNamespaceVersionMillis. + newVersion.foreach(v => recentNamespaceVersions.put(namespaceId, (v, System.currentTimeMillis()))) newVersion } catch { case e: Throwable => @@ -340,6 +368,12 @@ object Constant extends MdcLoggable { final val CONNECTOR_INBOUND_NAMESPACE = "connector_inbound" final val FINANCIAL_PRODUCTS_NAMESPACE = "financial_products" final val API_PRODUCTS_NAMESPACE = "api_products" + // The rendered message docs of each connector (GET /message-docs/CONNECTOR) and each connector's + // JSON Schema (GET /message-docs/CONNECTOR/json-schema), in Redis and in each instance's memory. + final val MESSAGE_DOCS_NAMESPACE = "message_docs" + // The Glossary as each instance holds it in memory. Bumping it reloads the Glossary at once and + // rebuilds every cached resource-docs document, because their keys carry the Glossary version. + final val GLOSSARY_NAMESPACE = "glossary" // List of all versioned cache namespaces final val ALL_CACHE_NAMESPACES = List( @@ -357,7 +391,9 @@ object Constant extends MdcLoggable { CONNECTOR_OUTBOUND_NAMESPACE, CONNECTOR_INBOUND_NAMESPACE, FINANCIAL_PRODUCTS_NAMESPACE, - API_PRODUCTS_NAMESPACE + API_PRODUCTS_NAMESPACE, + MESSAGE_DOCS_NAMESPACE, + GLOSSARY_NAMESPACE ) // Cache key prefixes with global namespace and versioning for easy invalidation @@ -368,6 +404,7 @@ object Constant extends MdcLoggable { def STATIC_RESOURCE_DOC_CACHE_KEY_PREFIX: String = getVersionedCachePrefix(RD_STATIC_NAMESPACE) def ALL_RESOURCE_DOC_CACHE_KEY_PREFIX: String = getVersionedCachePrefix(RD_ALL_NAMESPACE) def STATIC_SWAGGER_DOC_CACHE_KEY_PREFIX: String = getVersionedCachePrefix(SWAGGER_STATIC_NAMESPACE) + def MESSAGE_DOCS_CACHE_KEY_PREFIX: String = getVersionedCachePrefix(MESSAGE_DOCS_NAMESPACE) final val CREATE_LOCALISED_RESOURCE_DOC_JSON_TTL: Int = APIUtil.getPropsValue(s"createLocalisedResourceDocJson.cache.ttl.seconds", "3600").toInt final val GET_DYNAMIC_RESOURCE_DOCS_TTL: Int = APIUtil.getPropsValue(s"dynamicResourceDocsObp.cache.ttl.seconds", "3600").toInt final val GET_STATIC_RESOURCE_DOCS_TTL: Int = APIUtil.getPropsValue(s"staticResourceDocsObp.cache.ttl.seconds", "3600").toInt diff --git a/obp-api/src/main/scala/code/api/util/Glossary.scala b/obp-api/src/main/scala/code/api/util/Glossary.scala index 59998a4074..ff49ca3911 100644 --- a/obp-api/src/main/scala/code/api/util/Glossary.scala +++ b/obp-api/src/main/scala/code/api/util/Glossary.scala @@ -168,7 +168,11 @@ object Glossary extends MdcLoggable { val (checkedAt, version, byTitle) = cachedItemsByTitle.get() if (checkedAt != 0L && now - checkedAt < GlossaryCacheRecheckMillis) (version, byTitle) else { - val currentVersion = dynamicGlossaryItemsVersion + // The glossary cache namespace version is part of the token: bumping it (for example from + // the cache page in API Manager) reloads the Glossary here and, because the token is in + // every resource-docs cache key, rebuilds every cached document that embeds Glossary text. + val currentVersion = + s"$dynamicGlossaryItemsVersion-ns${code.api.Constant.recentCacheNamespaceVersion(code.api.Constant.GLOSSARY_NAMESPACE)}" if (checkedAt != 0L && currentVersion == version) { cachedItemsByTitle.set((now, version, byTitle)) (version, byTitle) @@ -188,6 +192,9 @@ object Glossary extends MdcLoggable { private def glossaryItemsByTitle: Map[String, GlossaryItem] = glossaryState._2 + /** Glossary items this instance holds in memory now (static and dynamic), for the cache page. */ + def loadedItemCount: Int = glossaryItemsByTitle.size + /** * A token for Resource Doc cache keys. It changes whenever a Dynamic Glossary Item is added, * changed or removed, so a cached endpoint description that embeds Glossary text is rebuilt diff --git a/obp-api/src/main/scala/code/api/util/JsonSchemaGenerator.scala b/obp-api/src/main/scala/code/api/util/JsonSchemaGenerator.scala index 222033f13f..dc25f8ab9c 100644 --- a/obp-api/src/main/scala/code/api/util/JsonSchemaGenerator.scala +++ b/obp-api/src/main/scala/code/api/util/JsonSchemaGenerator.scala @@ -59,12 +59,13 @@ object JsonSchemaGenerator { * recompute if Redis is unreachable or slow -- this in-memory layer doesn't depend on * Redis at all, so it stays a working safety net even when Redis is the one struggling. * - * The key is the connector name only. It must not be derived from `messageDocs`: turning the - * whole list (with every example message) into a key string costs megabytes per call. - * Callers always pass the named connector's own message docs. + * The key is the connector name (plus the `message_docs` cache namespace version, see + * `localKey`). It must not be derived from `messageDocs`: turning the whole list (with every + * example message) into a key string costs megabytes per call. Callers always pass the named + * connector's own message docs. */ def messageDocsToJsonSchema(messageDocs: List[MessageDoc], connectorName: String): JObject = - try schemaCache.get(connectorName, new Callable[JObject] { + try schemaCache.get(localKey(connectorName), new Callable[JObject] { def call(): JObject = { generatorCallsCounter.incrementAndGet() messageDocsToJsonSchemaUncached(messageDocs, connectorName) @@ -76,11 +77,19 @@ object JsonSchemaGenerator { case e: com.google.common.util.concurrent.UncheckedExecutionException if e.getCause != null => throw e.getCause } + // The key also carries the `message_docs` cache namespace version, so bumping that namespace + // rebuilds the schema on every instance (the version is re-read at most once a second). + private def localKey(connectorName: String): String = + s"${code.api.Constant.recentCacheNamespaceVersion(code.api.Constant.MESSAGE_DOCS_NAMESPACE)}|$connectorName" + private val schemaCache: Cache[String, JObject] = code.telemetry.Telemetry.monitorCache( CacheBuilder.newBuilder().maximumSize(64L).recordStats().build[String, JObject](), "json_schema") private val generatorCallsCounter = new AtomicLong(0) + /** Schemas held in this instance's memory. */ + def cacheSize: Long = schemaCache.size() + /** Cache hits and misses, for tests and monitoring. */ def cacheStats: CacheStats = schemaCache.stats() diff --git a/obp-api/src/main/scala/code/api/util/ResourceDocFilters.scala b/obp-api/src/main/scala/code/api/util/ResourceDocFilters.scala index 7f2d943ead..2f5d5ed807 100644 --- a/obp-api/src/main/scala/code/api/util/ResourceDocFilters.scala +++ b/obp-api/src/main/scala/code/api/util/ResourceDocFilters.scala @@ -136,6 +136,14 @@ object ResourceDocVocabulary { Vocabulary(staticVocabulary.tags ++ dynamic.tags, staticVocabulary.functions ++ dynamic.functions) } + /** A line for the cache page: the list's size and when its dynamic part was last rebuilt. */ + def describe(): String = { + val vocabulary = current() + val rebuiltAt = dynamicSnapshot.get().map(s => java.time.Instant.ofEpochMilli(s.builtAtMillis).toString).getOrElse("not yet") + s"Tag and function list for resource-docs filters: ${vocabulary.tags.size} tags, ${vocabulary.functions.size} functions; " + + s"its dynamic part follows this namespace and was last rebuilt at $rebuiltAt" + } + /** Forget the dynamic part, so the next call rebuilds it (tests, and after bulk dynamic changes). */ def refreshDynamic(): Unit = dynamicSnapshot.set(None) } diff --git a/obp-api/src/main/scala/code/api/v2_2_0/MessageDocsJsonCache.scala b/obp-api/src/main/scala/code/api/v2_2_0/MessageDocsJsonCache.scala index f150b916f1..7c873920e6 100644 --- a/obp-api/src/main/scala/code/api/v2_2_0/MessageDocsJsonCache.scala +++ b/obp-api/src/main/scala/code/api/v2_2_0/MessageDocsJsonCache.scala @@ -21,10 +21,9 @@ import org.json4s.JValue * 1. In-process: a bounded Guava cache holding the immutable JValue. It keeps working when * Redis is down, so an unreachable Redis can never send every request back into * reflection. - * 2. Shared: the same Redis-backed store the resource-doc and swagger endpoints use - * (`Caching.getStaticSwaggerDocCache`, same key prefix and TTL, same fail-safe behaviour: - * an unreachable Redis is a miss, never an error). It lets replicas and restarts reuse - * one instance's work. + * 2. Shared: Redis, in the `message_docs` cache namespace (`Caching.getMessageDocsCache`, + * same TTL as the resource docs, same fail-safe behaviour: an unreachable Redis is a miss, + * never an error). It lets replicas and restarts reuse one instance's work. * * Contract: * - Key: the connector name, and only after it has been resolved to a real connector. @@ -35,9 +34,12 @@ import org.json4s.JValue * - A failure is never cached; the next request retries. * - Redis is written only after a successful generation. An unparsable Redis value is * treated as a miss and regenerated. - * - Invalidation: `invalidateAll()` clears the in-process level. The Redis level expires by - * `staticResourceDocsObp.cache.ttl.seconds`. Nothing in production mutates a connector's - * message docs after start-up, so nothing calls `invalidateAll()` there. + * - Invalidation: bumping the `message_docs` cache namespace (for example from the cache page + * in API Manager) reaches both levels on every instance: the Redis keys carry the namespace + * version, and so do the in-process keys, which read it at most once a second + * (`Constant.recentCacheNamespaceVersion`). The Redis level also expires by + * `staticResourceDocsObp.cache.ttl.seconds`. `invalidateAll()` clears this instance's + * in-process level only (tests). */ object MessageDocsJsonCache extends Loggable { private val MaxEntries = 64L @@ -49,8 +51,8 @@ object MessageDocsJsonCache extends Loggable { } object RedisStore extends SharedStore { - def get(key: String): Option[String] = Caching.getStaticSwaggerDocCache(key) - def set(key: String, value: String): Unit = Caching.setStaticSwaggerDocCache(key, value) + def get(key: String): Option[String] = Caching.getMessageDocsCache(key) + def set(key: String, value: String): Unit = Caching.setMessageDocsCache(key, value) } private def sharedKey(connectorName: String) = s"message-docs-v2.2.0-$connectorName" @@ -64,8 +66,13 @@ object MessageDocsJsonCache extends Loggable { private val sharedHitsCounter = new AtomicLong(0) private val sharedSetsCounter = new AtomicLong(0) + // The in-process key carries the namespace version, so a bump makes old entries unreachable here + // too; they then age out of the bounded cache. + private def localKey(connectorName: String): String = + s"${code.api.Constant.recentCacheNamespaceVersion(code.api.Constant.MESSAGE_DOCS_NAMESPACE)}|$connectorName" + def getOrCompute(connectorName: String, store: SharedStore = RedisStore)(generate: => JValue): JValue = - try cache.get(connectorName, new Callable[JValue] { + try cache.get(localKey(connectorName), new Callable[JValue] { def call(): JValue = { val key = sharedKey(connectorName) sharedGetsCounter.incrementAndGet() diff --git a/obp-api/src/main/scala/code/api/v6_0_0/Http4s600.scala b/obp-api/src/main/scala/code/api/v6_0_0/Http4s600.scala index 24c61d1efe..e04faeedbc 100644 --- a/obp-api/src/main/scala/code/api/v6_0_0/Http4s600.scala +++ b/obp-api/src/main/scala/code/api/v6_0_0/Http4s600.scala @@ -1238,6 +1238,7 @@ object Http4s600 { (Constant.STATIC_RESOURCE_DOC_CACHE_KEY_PREFIX, "Static resource documentation", Constant.GET_STATIC_RESOURCE_DOCS_TTL.toString, "Resource Documentation"), (Constant.ALL_RESOURCE_DOC_CACHE_KEY_PREFIX, "All resource documentation", Constant.GET_STATIC_RESOURCE_DOCS_TTL.toString, "Resource Documentation"), (Constant.STATIC_SWAGGER_DOC_CACHE_KEY_PREFIX, "Swagger documentation", Constant.GET_STATIC_RESOURCE_DOCS_TTL.toString, "Resource Documentation"), + (Constant.MESSAGE_DOCS_CACHE_KEY_PREFIX, "Message docs and connector JSON Schemas", Constant.GET_STATIC_RESOURCE_DOCS_TTL.toString, "Resource Documentation"), (Constant.CONNECTOR_PREFIX, "Connector method names and metadata", "3600", "Connector"), (Constant.METRICS_STABLE_PREFIX, "Stable metrics (historical)", "86400", "Metrics"), (Constant.METRICS_RECENT_PREFIX, "Recent metrics", "7", "Metrics"), @@ -2543,7 +2544,7 @@ object Http4s600 { case req @ GET -> `prefixPath` / "message-docs" / connector / "json-schema" => EndpointHelpers.executeAndRespond(req) { implicit cc => val cacheKey = s"message-docs-json-schema-$connector" - val cacheValueFromRedis = code.api.cache.Caching.getStaticSwaggerDocCache(cacheKey) + val cacheValueFromRedis = code.api.cache.Caching.getMessageDocsCache(cacheKey) for { jsonSchema <- if (cacheValueFromRedis.isDefined) { NewStyle.function.tryons(s"$UnknownError Cannot parse cached JSON Schema.", 400, Some(cc)) { @@ -2559,7 +2560,7 @@ object Http4s600 { val schema = code.api.util.JsonSchemaGenerator.messageDocsToJsonSchema( connectorObject.messageDocs.toList, connector) val schemaString = com.openbankproject.commons.util.JsonAliases.compactRender(schema) - code.api.cache.Caching.setStaticSwaggerDocCache(cacheKey, schemaString) + code.api.cache.Caching.setMessageDocsCache(cacheKey, schemaString) schema } } diff --git a/obp-api/src/main/scala/code/api/v6_0_0/JSONFactory6.0.0.scala b/obp-api/src/main/scala/code/api/v6_0_0/JSONFactory6.0.0.scala index 78a43091a7..00b3473c3f 100644 --- a/obp-api/src/main/scala/code/api/v6_0_0/JSONFactory6.0.0.scala +++ b/obp-api/src/main/scala/code/api/v6_0_0/JSONFactory6.0.0.scala @@ -2424,7 +2424,7 @@ object JSONFactory600 extends CustomJsonFormats with MdcLoggable { Constant.CALL_COUNTER_NAMESPACE -> ("Rate limit call counters", "Rate Limiting"), Constant.RL_ACTIVE_NAMESPACE -> ("Active rate limit states", "Rate Limiting"), Constant.RD_LOCALISED_NAMESPACE -> ("Localized resource docs", "API Documentation"), - Constant.RD_DYNAMIC_NAMESPACE -> ("Dynamic resource docs", "API Documentation"), + Constant.RD_DYNAMIC_NAMESPACE -> (s"Dynamic resource docs. ${code.api.util.ResourceDocVocabulary.describe()}", "API Documentation"), Constant.RD_STATIC_NAMESPACE -> ("Static resource docs", "API Documentation"), Constant.RD_ALL_NAMESPACE -> ("All resource docs", "API Documentation"), Constant.SWAGGER_STATIC_NAMESPACE -> ("Static Swagger docs", "API Documentation"), @@ -2433,7 +2433,16 @@ object JSONFactory600 extends CustomJsonFormats with MdcLoggable { Constant.METRICS_RECENT_NAMESPACE -> ("Recent metrics data", "Metrics"), Constant.ABAC_RULE_NAMESPACE -> ("ABAC rule cache", "Authorization"), Constant.FINANCIAL_PRODUCTS_NAMESPACE -> ("Financial product list (bank-scoped and all-banks)", "Products"), - Constant.API_PRODUCTS_NAMESPACE -> ("Api product list (all banks)", "Products") + Constant.API_PRODUCTS_NAMESPACE -> ("Api product list (all banks)", "Products"), + Constant.MESSAGE_DOCS_NAMESPACE -> ("Message docs and connector JSON Schemas (Redis, and each instance's memory)", "API Documentation"), + Constant.GLOSSARY_NAMESPACE -> ("Glossary items held in each instance's memory; bumping also rebuilds cached resource docs", "API Documentation") + ) + + // Caches held in this instance's memory outside the shared in-memory store, which follow their + // namespace's version through their own keys. Counted here so the page shows their real size. + val inProcessEntries: Map[String, () => Long] = Map( + Constant.MESSAGE_DOCS_NAMESPACE -> (() => code.api.v2_2_0.MessageDocsJsonCache.size + code.api.util.JsonSchemaGenerator.cacheSize), + Constant.GLOSSARY_NAMESPACE -> (() => code.api.util.Glossary.loadedItemCount.toLong) ) var redisAvailable = true @@ -2481,7 +2490,7 @@ object JSONFactory600 extends CustomJsonFormats with MdcLoggable { } try { - memoryKeyCount = InMemory.countKeys(pattern) + memoryKeyCount = InMemory.countKeys(pattern) + inProcessEntries.get(namespaceId).map(count => count().toInt).getOrElse(0) totalKeys += memoryKeyCount if (memoryKeyCount > 0 && redisKeyCount == 0) { diff --git a/obp-api/src/main/scala/code/api/v7_0_0/Http4s700IpPenalties.scala b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700IpPenalties.scala index 4ad67bfd74..bbf2440820 100644 --- a/obp-api/src/main/scala/code/api/v7_0_0/Http4s700IpPenalties.scala +++ b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700IpPenalties.scala @@ -34,7 +34,6 @@ import code.api.util.ApiTag._ import code.api.util.ErrorMessages._ import code.api.util.http4s.Http4sRequestAttributes.EndpointHelpers import code.api.util.{CustomJsonFormats, IpPenalties} -import code.api.v7_0_0.JSONFactory700.{IpPenaltiesJsonV700, IpPenaltyJsonV700, PostIpPenaltyJsonV700} import code.util.Helper import com.github.dwickern.macros.NameOf.nameOf import com.openbankproject.commons.ExecutionContext.Implicits.global @@ -91,7 +90,7 @@ object Http4s700IpPenalties { } added <- Future(IpPenalties.add(body.ip_address, body.per_minute_limit, body.duration_minutes, body.reason.trim, user.userId)) _ <- Helper.booleanToFuture(added.left.getOrElse(UnknownError), failCode = 409, cc = Some(cc))(added.isRight) - } yield JSONFactory700.createIpPenaltyJson(added.toOption.get) + } yield JSONFactory700Operations.createIpPenaltyJson(added.toOption.get) } } @@ -109,8 +108,8 @@ object Http4s700IpPenalties { |`duration_minutes` is from 1 to ${IpPenalties.MaxDurationMinutes} (one week). An address can have one |penalty at a time: to change it, delete it and create it again (409 when one exists). |""".stripMargin, - JSONFactory700.postIpPenaltyJsonV700Example, - JSONFactory700.ipPenaltyJsonV700Example, + JSONFactory700Operations.postIpPenaltyJsonV700Example, + JSONFactory700Operations.ipPenaltyJsonV700Example, List($AuthenticatedUserIsRequired, UserHasMissingRoles, InvalidJsonFormat, InvalidIpPenalty, InvalidIpAddress, IpPenaltyAlreadyExists, UnknownError), List(apiTagRateLimits, apiTagSystem), @@ -122,7 +121,7 @@ object Http4s700IpPenalties { lazy val getIpPenalties: HttpRoutes[IO] = HttpRoutes.of[IO] { case req @ GET -> `prefixPath` / "management" / "ip-penalties" => EndpointHelpers.withUser(req) { (_, _) => - Future(IpPenaltiesJsonV700(IpPenalties.listAll().map(JSONFactory700.createIpPenaltyJson))) + Future(IpPenaltiesJsonV700(IpPenalties.listAll().map(JSONFactory700Operations.createIpPenaltyJson))) } } @@ -137,7 +136,7 @@ object Http4s700IpPenalties { |$penaltyDescription |""".stripMargin, EmptyBody, - JSONFactory700.ipPenaltiesJsonV700Example, + JSONFactory700Operations.ipPenaltiesJsonV700Example, List($AuthenticatedUserIsRequired, UserHasMissingRoles, UnknownError), List(apiTagRateLimits, apiTagSystem), Some(List(canGetIpPenalties)), diff --git a/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Telemetry.scala b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Telemetry.scala index ef831b0142..c297e2757f 100644 --- a/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Telemetry.scala +++ b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Telemetry.scala @@ -66,7 +66,7 @@ object Http4s700Telemetry { case req @ GET -> `prefixPath` / "management" / "telemetry" => EndpointHelpers.withUser(req) { (_, _) => val namePrefix = req.uri.query.params.get("name_prefix").filter(_.nonEmpty) - Future(JSONFactory700.createTelemetryJson(namePrefix)) + Future(JSONFactory700Operations.createTelemetryJson(namePrefix)) } } @@ -112,7 +112,7 @@ object Http4s700Telemetry { |for example `?name_prefix=obp.api.endpoint`. |""".stripMargin, EmptyBody, - JSONFactory700.telemetryJsonV700Example, + JSONFactory700Operations.telemetryJsonV700Example, List($AuthenticatedUserIsRequired, UserHasMissingRoles, UnknownError), List(apiTagApi, apiTagSystem), Some(List(canGetTelemetry)), diff --git a/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory7.0.0.scala b/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory7.0.0.scala index 47aedf1f87..ec48051a30 100644 --- a/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory7.0.0.scala +++ b/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory7.0.0.scala @@ -2771,95 +2771,4 @@ object JSONFactory700 extends MdcLoggable with code.api.util.CustomJsonFormats { attributes = Some(List(apiProductSubscriptionAttributeResponseJsonV700Example)) ) lazy val apiProductSubscriptionsJsonV700Example = ApiProductSubscriptionsJsonV700(List(apiProductSubscriptionJsonV700Example)) - - // ===== IP penalties ===== - - /** Request body of POST /management/ip-penalties. `per_minute_limit` 0 refuses every request. */ - case class PostIpPenaltyJsonV700(ip_address: String, per_minute_limit: Long, duration_minutes: Long, reason: String) - - case class IpPenaltyJsonV700( - ip_address: String, - per_minute_limit: Long, - reason: String, - created_by_user_id: String, - created_at: Date, - expires_at: Date - ) - case class IpPenaltiesJsonV700(ip_penalties: List[IpPenaltyJsonV700]) - - def createIpPenaltyJson(penalty: code.api.util.IpPenalties.Penalty): IpPenaltyJsonV700 = - IpPenaltyJsonV700(penalty.ipAddress, penalty.perMinuteLimit, penalty.reason, penalty.createdByUserId, - new Date(penalty.createdAtMillis), new Date(penalty.expiresAtMillis)) - - lazy val postIpPenaltyJsonV700Example = PostIpPenaltyJsonV700( - ip_address = "203.0.113.42", per_minute_limit = 10, duration_minutes = 60, - reason = "Vulnerability scan: about 900 requests a minute to resource-docs with changing filters") - - lazy val ipPenaltyJsonV700Example = IpPenaltyJsonV700( - ip_address = "203.0.113.42", per_minute_limit = 10, - reason = "Vulnerability scan: about 900 requests a minute to resource-docs with changing filters", - created_by_user_id = ExampleValue.userIdExample.value, - created_at = APIUtil.DateWithMsExampleObject, expires_at = APIUtil.DateWithMsExampleObject) - - lazy val ipPenaltiesJsonV700Example = IpPenaltiesJsonV700(List(ipPenaltyJsonV700Example)) - - // ===== Telemetry ===== - - /** One meter: its Micrometer name, type, unit, tags and current values (count, total_time, max, value ...). */ - case class TelemetryMeterJsonV700( - name: String, - `type`: String, - base_unit: Option[String], - tags: Map[String, String], - measurements: Map[String, Double] - ) - - /** The separate port Prometheus scrapes, as configured on this instance. */ - case class TelemetryPortJsonV700(enabled: Boolean, port: Int, path: String) - - case class TelemetryJsonV700( - api_instance_id: String, - git_commit: String, - port: TelemetryPortJsonV700, - meters: List[TelemetryMeterJsonV700] - ) - - /** This instance's Telemetry, from the same registry the separate port serves, optionally limited to names starting with `namePrefix`. */ - def createTelemetryJson(namePrefix: Option[String]): TelemetryJsonV700 = { - import scala.jdk.CollectionConverters._ - val settings = code.telemetry.Telemetry.portSettings - val meters = code.telemetry.Telemetry.registry.getMeters.asScala.toList - .filter(meter => namePrefix.forall(prefix => meter.getId.getName.startsWith(prefix))) - .map { meter => - val id = meter.getId - TelemetryMeterJsonV700( - name = id.getName, - `type` = id.getType.name.toLowerCase, - base_unit = Option(id.getBaseUnit), - tags = id.getTags.asScala.map(tag => tag.getKey -> tag.getValue).toMap, - // A gauge whose source has gone reads NaN, which is not valid JSON. - measurements = meter.measure().asScala - .filter(measurement => java.lang.Double.isFinite(measurement.getValue)) - .map(measurement => measurement.getStatistic.name.toLowerCase -> measurement.getValue).toMap) - } - .sortBy(meter => (meter.name, meter.tags.toList.sorted.mkString(","))) - TelemetryJsonV700( - api_instance_id = Constant.ApiInstanceId, - git_commit = APIUtil.gitCommit, - port = TelemetryPortJsonV700(settings.enabled, settings.port, code.telemetry.Telemetry.ScrapePath), - meters = meters) - } - - lazy val telemetryJsonV700Example = TelemetryJsonV700( - api_instance_id = "obp_4f6b3c2a-9d1e-4b7a-8c5f-2e1d0a9b8c7d", - git_commit = "3286937795b4d0c2e1f6a8b9c0d1e2f3a4b5c6d7", - port = TelemetryPortJsonV700(enabled = true, port = code.telemetry.Telemetry.DefaultPort, path = code.telemetry.Telemetry.ScrapePath), - meters = List( - TelemetryMeterJsonV700("cache.gets", "counter", None, Map("cache" -> "json_schema", "result" -> "hit"), Map("count" -> 118.0)), - TelemetryMeterJsonV700("jvm.threads.live", "gauge", Some("threads"), Map.empty, Map("value" -> 64.0)), - TelemetryMeterJsonV700("obp.api.endpoint.requests", "timer", Some("seconds"), - Map("operation" -> "OBPv7.0.0-getBanks", "api_version" -> "v7.0.0", "status" -> "2xx"), - Map("count" -> 42.0, "total_time" -> 1.26, "max" -> 0.081)) - ) - ) } 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 new file mode 100644 index 0000000000..5e92e9084a --- /dev/null +++ b/obp-api/src/main/scala/code/api/v7_0_0/JSONFactory700Operations.scala @@ -0,0 +1,140 @@ +/** +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 java.util.Date + +import code.api.Constant +import code.api.util.{APIUtil, ExampleValue, IpPenalties} +import code.telemetry.Telemetry + +/* + * The JSON of the v7.0.0 operations endpoints (Telemetry and IP penalties). + * + * These case classes are declared at package level rather than inside `object JSONFactory700`, + * which already holds over a hundred nested case classes. json4s reads a case class's field names + * from the compiled Scala signature; for a nested class that is the enclosing object's signature, + * so the object's signature grows with every class added to it. Keeping new, self-contained groups + * of JSON types in their own file stops that growth. JsonFactorySignatureTest checks that json4s + * can read every case class nested in a JSONFactory object. + */ + +// ===== Telemetry ===== + +/** One meter: its Micrometer name, type, unit, tags and current values (count, total_time, max, value ...). */ +case class TelemetryMeterJsonV700( + name: String, + `type`: String, + base_unit: Option[String], + tags: Map[String, String], + measurements: Map[String, Double] +) + +/** The separate port Prometheus scrapes, as configured on this instance. */ +case class TelemetryPortJsonV700(enabled: Boolean, port: Int, path: String) + +case class TelemetryJsonV700( + api_instance_id: String, + git_commit: String, + port: TelemetryPortJsonV700, + meters: List[TelemetryMeterJsonV700] +) + +// ===== IP penalties ===== + +/** Request body of POST /management/ip-penalties. `per_minute_limit` 0 refuses every request. */ +case class PostIpPenaltyJsonV700(ip_address: String, per_minute_limit: Long, duration_minutes: Long, reason: String) + +case class IpPenaltyJsonV700( + ip_address: String, + per_minute_limit: Long, + reason: String, + created_by_user_id: String, + created_at: Date, + expires_at: Date +) + +case class IpPenaltiesJsonV700(ip_penalties: List[IpPenaltyJsonV700]) + +/** This object builds the JSON above, and holds the examples the ResourceDocs show. */ +object JSONFactory700Operations { + + /** This instance's Telemetry, from the same registry the separate port serves, optionally limited to names starting with `namePrefix`. */ + def createTelemetryJson(namePrefix: Option[String]): TelemetryJsonV700 = { + import scala.jdk.CollectionConverters._ + val settings = Telemetry.portSettings + val meters = Telemetry.registry.getMeters.asScala.toList + .filter(meter => namePrefix.forall(prefix => meter.getId.getName.startsWith(prefix))) + .map { meter => + val id = meter.getId + TelemetryMeterJsonV700( + name = id.getName, + `type` = id.getType.name.toLowerCase, + base_unit = Option(id.getBaseUnit), + tags = id.getTags.asScala.map(tag => tag.getKey -> tag.getValue).toMap, + // A gauge whose source has gone reads NaN, which is not valid JSON. + measurements = meter.measure().asScala + .filter(measurement => java.lang.Double.isFinite(measurement.getValue)) + .map(measurement => measurement.getStatistic.name.toLowerCase -> measurement.getValue).toMap) + } + .sortBy(meter => (meter.name, meter.tags.toList.sorted.mkString(","))) + TelemetryJsonV700( + api_instance_id = Constant.ApiInstanceId, + git_commit = APIUtil.gitCommit, + port = TelemetryPortJsonV700(settings.enabled, settings.port, Telemetry.ScrapePath), + meters = meters) + } + + lazy val telemetryJsonV700Example = TelemetryJsonV700( + api_instance_id = "obp_4f6b3c2a-9d1e-4b7a-8c5f-2e1d0a9b8c7d", + git_commit = "3286937795b4d0c2e1f6a8b9c0d1e2f3a4b5c6d7", + port = TelemetryPortJsonV700(enabled = true, port = Telemetry.DefaultPort, path = Telemetry.ScrapePath), + meters = List( + TelemetryMeterJsonV700("cache.gets", "counter", None, Map("cache" -> "json_schema", "result" -> "hit"), Map("count" -> 118.0)), + TelemetryMeterJsonV700("jvm.threads.live", "gauge", Some("threads"), Map.empty, Map("value" -> 64.0)), + TelemetryMeterJsonV700("obp.api.endpoint.requests", "timer", Some("seconds"), + Map("operation" -> "OBPv7.0.0-getBanks", "api_version" -> "v7.0.0", "status" -> "2xx"), + Map("count" -> 42.0, "total_time" -> 1.26, "max" -> 0.081)) + ) + ) + + def createIpPenaltyJson(penalty: IpPenalties.Penalty): IpPenaltyJsonV700 = + IpPenaltyJsonV700(penalty.ipAddress, penalty.perMinuteLimit, penalty.reason, penalty.createdByUserId, + new Date(penalty.createdAtMillis), new Date(penalty.expiresAtMillis)) + + lazy val postIpPenaltyJsonV700Example = PostIpPenaltyJsonV700( + ip_address = "203.0.113.42", per_minute_limit = 10, duration_minutes = 60, + reason = "Vulnerability scan: about 900 requests a minute to resource-docs with changing filters") + + lazy val ipPenaltyJsonV700Example = IpPenaltyJsonV700( + ip_address = "203.0.113.42", per_minute_limit = 10, + reason = "Vulnerability scan: about 900 requests a minute to resource-docs with changing filters", + created_by_user_id = ExampleValue.userIdExample.value, + created_at = APIUtil.DateWithMsExampleObject, expires_at = APIUtil.DateWithMsExampleObject) + + lazy val ipPenaltiesJsonV700Example = IpPenaltiesJsonV700(List(ipPenaltyJsonV700Example)) +} diff --git a/obp-api/src/test/scala/code/api/JsonFactorySignatureTest.scala b/obp-api/src/test/scala/code/api/JsonFactorySignatureTest.scala new file mode 100644 index 0000000000..0c29cae5c1 --- /dev/null +++ b/obp-api/src/test/scala/code/api/JsonFactorySignatureTest.scala @@ -0,0 +1,55 @@ +package code.api + +import java.io.File + +import org.scalatest.{FlatSpec, Matchers} + +import scala.io.Source + +/** + * This suite fails when json4s cannot read the field names of a case class declared inside a + * JSONFactory object. + * + * json4s finds a case class's field names in the compiled Scala signature (for a nested class, the + * enclosing object's). When it cannot, every request or response using that class fails with + * "Can't find ScalaSig", a 400 or 500 on any endpoint whose body or response uses it. This suite + * asks json4s to describe every nested case class, which is the lookup that would fail. + */ +class JsonFactorySignatureTest extends FlatSpec with Matchers { + + private val sourceRoot: File = + List(new File("src/main/scala"), new File("obp-api/src/main/scala")).find(_.isDirectory) + .getOrElse(throw new IllegalStateException("JsonFactorySignatureTest cannot find obp-api/src/main/scala")) + + private def walk(dir: File): List[File] = + Option(dir.listFiles).toList.flatten.flatMap(f => if (f.isDirectory) walk(f) else List(f)) + + /** Fully qualified names of every `object JSONFactory...` in the source. */ + private lazy val factoryObjects: List[String] = + walk(sourceRoot).filter(f => f.getName.startsWith("JSONFactory") && f.getName.endsWith(".scala")).flatMap { file => + val source = Source.fromFile(file, "UTF-8") + val text = try source.mkString finally source.close() + val pkg = """(?m)^package\s+([\w.]+)""".r.findFirstMatchIn(text).map(_.group(1)) + """(?m)^object\s+(JSONFactory\w*)""".r.findAllMatchIn(text).map(_.group(1)).toList.flatMap(name => pkg.map(p => s"$p.$name")) + } + + "The JSONFactory objects" should "be found" in { + factoryObjects should contain("code.api.v7_0_0.JSONFactory700") + factoryObjects.size should be > 5 + } + + they should "have every nested case class readable by json4s" in { + val unreadable = factoryObjects.flatMap { name => + Class.forName(name + "$").getDeclaredClasses.toList + .filter(c => classOf[Product].isAssignableFrom(c) && !c.getName.endsWith("$")) + .flatMap { caseClass => + try { org.json4s.reflect.Reflector.describe(org.json4s.reflect.Reflector.scalaTypeOf(caseClass)); None } + catch { case e: Throwable => Some(s"${caseClass.getName}: ${e.getMessage}") } + } + } + withClue("json4s cannot read these case classes. Declare new case classes at package level, not inside the " + + "JSONFactory object (see JSONFactory700Operations.scala).\n" + unreadable.take(20).mkString("\n") + "\n") { + unreadable shouldBe empty + } + } +} diff --git a/obp-api/src/test/scala/code/api/cache/DocumentationCacheNamespacesTest.scala b/obp-api/src/test/scala/code/api/cache/DocumentationCacheNamespacesTest.scala new file mode 100644 index 0000000000..a3b5bdd2b8 --- /dev/null +++ b/obp-api/src/test/scala/code/api/cache/DocumentationCacheNamespacesTest.scala @@ -0,0 +1,77 @@ +package code.api.cache + +import java.util.UUID + +import code.api.Constant +import code.api.util.{Glossary, JsonSchemaGenerator} +import code.api.v2_2_0.MessageDocsJsonCache +import code.setup.ServerSetup +import org.json4s.JsonAST.{JObject, JString} + +/** + * This suite checks that the documentation caches held in each instance's memory follow their + * cache namespace: bumping `message_docs` rebuilds the message docs and the connector JSON Schema, + * and bumping `glossary` changes the Glossary token that every resource-docs cache key carries. + * + * Bumping needs Redis (the version counter lives there); without it the scenarios cancel. + */ +class DocumentationCacheNamespacesTest extends ServerSetup { + + private def bump(namespaceId: String): Unit = { + val bumped = Constant.incrementCacheNamespaceVersion(namespaceId) + if (bumped.isEmpty) cancel(s"Redis is not reachable, so the $namespaceId namespace cannot be bumped") + } + + private def freshConnector(): String = s"namespace-probe-${UUID.randomUUID().toString.take(8)}" + + feature("The documentation cache namespaces") { + + scenario("message_docs and glossary are listed with the other namespaces") { + Constant.ALL_CACHE_NAMESPACES should contain allOf (Constant.MESSAGE_DOCS_NAMESPACE, Constant.GLOSSARY_NAMESPACE) + } + + scenario("bumping message_docs rebuilds the message docs response") { + val connector = freshConnector() + var builds = 0 + // An object, like real message docs: the cache stores the rendered text and parses it back. + def build() = { builds += 1; JObject("build" -> JString(s"$builds")) } + MessageDocsJsonCache.getOrCompute(connector)(build()) + MessageDocsJsonCache.getOrCompute(connector)(build()) + builds shouldBe 1 + + bump(Constant.MESSAGE_DOCS_NAMESPACE) + MessageDocsJsonCache.getOrCompute(connector)(build()) + builds shouldBe 2 + } + + scenario("bumping message_docs rebuilds the connector JSON Schema") { + val connector = freshConnector() + val before = JsonSchemaGenerator.generatorCalls + JsonSchemaGenerator.messageDocsToJsonSchema(Nil, connector) + JsonSchemaGenerator.messageDocsToJsonSchema(Nil, connector) + JsonSchemaGenerator.generatorCalls - before shouldBe 1 + + bump(Constant.MESSAGE_DOCS_NAMESPACE) + JsonSchemaGenerator.messageDocsToJsonSchema(Nil, connector) + JsonSchemaGenerator.generatorCalls - before shouldBe 2 + } + + scenario("bumping glossary changes the Glossary token in resource-docs cache keys") { + Glossary.invalidateGlossaryItemCache() + val before = Glossary.glossaryVersionForCacheKey + bump(Constant.GLOSSARY_NAMESPACE) + Glossary.invalidateGlossaryItemCache() + Glossary.glossaryVersionForCacheKey should not equal before + } + + scenario("another instance sees a bump within a second, without a Redis read per request") { + val first = Constant.recentCacheNamespaceVersion(Constant.MESSAGE_DOCS_NAMESPACE) + // Simulate a bump made on another instance: the counter changes in Redis, not locally. + val versionKey = s"${Constant.getGlobalCacheNamespacePrefix}cache_version_${Constant.MESSAGE_DOCS_NAMESPACE}" + if (Redis.use(code.api.JedisMethod.INCR, versionKey, None, None).isEmpty) cancel("Redis is not reachable") + Constant.recentCacheNamespaceVersion(Constant.MESSAGE_DOCS_NAMESPACE) shouldBe first // still the local copy + Thread.sleep(Constant.RecentNamespaceVersionMillis + 100) + Constant.recentCacheNamespaceVersion(Constant.MESSAGE_DOCS_NAMESPACE) shouldBe first + 1 + } + } +} diff --git a/obp-api/src/test/scala/code/api/v7_0_0/IpPenaltiesTest.scala b/obp-api/src/test/scala/code/api/v7_0_0/IpPenaltiesTest.scala index d24943df14..d904c96771 100644 --- a/obp-api/src/test/scala/code/api/v7_0_0/IpPenaltiesTest.scala +++ b/obp-api/src/test/scala/code/api/v7_0_0/IpPenaltiesTest.scala @@ -32,7 +32,6 @@ import code.api.util.ApiRole.{CanCreateIpPenalty, CanDeleteIpPenalty, CanGetIpPe import code.api.util.ErrorMessages.{AuthenticatedUserIsRequired, UserHasMissingRoles} import code.api.util.IpPenalties import code.api.v6_0_0.V600ServerSetup -import code.api.v7_0_0.JSONFactory700.{IpPenaltiesJsonV700, IpPenaltyJsonV700, PostIpPenaltyJsonV700} import code.entitlement.Entitlement import com.openbankproject.commons.model.ErrorMessage import com.openbankproject.commons.util.ApiVersion diff --git a/obp-api/src/test/scala/code/api/v7_0_0/TelemetryEndpointTest.scala b/obp-api/src/test/scala/code/api/v7_0_0/TelemetryEndpointTest.scala index 4c73ec9865..c59f9cbe88 100644 --- a/obp-api/src/test/scala/code/api/v7_0_0/TelemetryEndpointTest.scala +++ b/obp-api/src/test/scala/code/api/v7_0_0/TelemetryEndpointTest.scala @@ -32,7 +32,6 @@ import code.api.util.APIUtil.OAuth._ import code.api.util.ApiRole.CanGetTelemetry import code.api.util.ErrorMessages.{AuthenticatedUserIsRequired, UserHasMissingRoles} import code.api.v6_0_0.V600ServerSetup -import code.api.v7_0_0.JSONFactory700.TelemetryJsonV700 import code.entitlement.Entitlement import com.openbankproject.commons.model.ErrorMessage import com.openbankproject.commons.util.ApiVersion diff --git a/release_notes.md b/release_notes.md index 3f14a8bd82..b1c99d43b2 100644 --- a/release_notes.md +++ b/release_notes.md @@ -3,7 +3,18 @@ ### Most recent changes at top of file ``` Date Commit Action -28/09/2026 TBD NEW: IP penalties, an operator's temporary per-minute limit on one IP +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). + message_docs holds the rendered message docs and the connector JSON + Schemas, in Redis (moved out of swagger_static) and in each instance's + memory; bumping it rebuilds both on every instance within a second. + glossary: bumping it reloads the Glossary each instance holds in memory + and rebuilds every cached resource-docs document, whose keys carry the + Glossary version. + The rd_dynamic row now also reports the tag and function list used to + check resource-docs filters (size, and when it was last rebuilt). +28/09/2026 e71c90ca0 NEW: IP penalties, an operator's temporary per-minute limit on one IP address, on every endpoint, checked before all other rate limiters. per_minute_limit 0 refuses every request. Always enforced, kept in Redis (shared by every instance), removed automatically on expiry (at most one