diff --git a/obp-api/src/main/resources/props/sample.props.template b/obp-api/src/main/resources/props/sample.props.template
index c960ba902f..ade2a9e657 100644
--- a/obp-api/src/main/resources/props/sample.props.template
+++ b/obp-api/src/main/resources/props/sample.props.template
@@ -1960,6 +1960,10 @@ user_account_is_validated = false
# than 600 so successive runs drift across the wall clock instead of phase-locking
# to every 10th minute and to other periodic schedulers.
# retain_metrics_scheduler_interval_in_seconds = 599
+# The longest date range, in days, that one GET /management/aggregate-metrics call may cover.
+# An aggregate reads every metric row in its range, so the cost grows with the range; a
+# longer range is refused with 400. A call without from_date covers this many days.
+# aggregate_metrics_max_days = 31
# Defines endpoints we want to store responses at Metric table
diff --git a/obp-api/src/main/scala/code/api/util/APIUtil.scala b/obp-api/src/main/scala/code/api/util/APIUtil.scala
index 178f9abe01..14d1e2309e 100644
--- a/obp-api/src/main/scala/code/api/util/APIUtil.scala
+++ b/obp-api/src/main/scala/code/api/util/APIUtil.scala
@@ -5220,7 +5220,7 @@ object APIUtil extends MdcLoggable with CustomJsonFormats{
// The Glossary version belongs in the key: endpoint descriptions embed Glossary text, so a
// Dynamic Glossary Item that overrides a static one must not stay masked by a cached document
// for the rest of the resource-doc / swagger TTL. Reading it is an in-memory lookup that
- // re-checks the database at most once a second.
+ // re-checks the Dynamic Glossary Items every ten minutes and the glossary cache namespace every ten seconds.
s"-contentParam:$contentParam-apiCollectionIdParam:$apiCollectionIdParam-isVersion4OrHigher:$isVersion4OrHigher-glossary:${Glossary.glossaryVersionForCacheKey}".intern()
def getUserLacksRevokePermissionErrorMessage(sourceViewId: ViewId, targetViewId: ViewId) =
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 b503df15bd..b68b73f314 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,14 @@ object ApiRole extends MdcLoggable{
case class CanGetTelemetry(requiresBankId: Boolean = false) extends ApiRole
lazy val canGetTelemetry = CanGetTelemetry()
+ // What has been granted, across all Users and Consumers. CanGetReachableRoles shows only Role names,
+ // never who holds them, so it can be given to a code-review service; CanGetAllScopes shows the Consumers too.
+ case class CanGetAllScopes(requiresBankId: Boolean = false) extends ApiRole
+ lazy val canGetAllScopes = CanGetAllScopes()
+
+ case class CanGetReachableRoles(requiresBankId: Boolean = false) extends ApiRole
+ lazy val canGetReachableRoles = CanGetReachableRoles()
+
// 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
diff --git a/obp-api/src/main/scala/code/api/util/DBUtil.scala b/obp-api/src/main/scala/code/api/util/DBUtil.scala
index a29bfb8db3..e23c435e42 100644
--- a/obp-api/src/main/scala/code/api/util/DBUtil.scala
+++ b/obp-api/src/main/scala/code/api/util/DBUtil.scala
@@ -40,7 +40,7 @@ object DBUtil {
def isSqlServer: Boolean = dbUrl.contains("sqlserver")
/** The endpoint timeout in whole seconds, rounded up, which is the unit JDBC takes. */
- private[util] def queryTimeoutSeconds: Int =
+ private[code] def queryTimeoutSeconds: Int =
math.max(1L, (Constant.longEndpointTimeoutInMillis + 999) / 1000).toInt
/**
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 89b5010035..9944bd3356 100644
--- a/obp-api/src/main/scala/code/api/util/ErrorMessages.scala
+++ b/obp-api/src/main/scala/code/api/util/ErrorMessages.scala
@@ -188,6 +188,7 @@ object ErrorMessages {
val FilterAnonFormatError = s"OBP-10028: anon parameter can only take two values: TRUE or FALSE!"
val FilterDurationFormatError = s"OBP-10029: wrong value for `duration` parameter. Please send a positive integer (=>0)!"
val FilterIsDeletedFormatError = s"OBP-10036: is_deleted parameter can only take two values: TRUE or FALSE!"
+ val AggregateMetricsDateRangeTooLong = "OBP-10069: The date range is too long for aggregate metrics. Set from_date and to_date closer together. "
val InvalidApiVersionString = "OBP-00027: Invalid API Version string. We could not find the version specified."
val IncorrectTriggerName = "OBP-10039: Incorrect Trigger name:"
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 ba4afa12cf..f29998ffd1 100644
--- a/obp-api/src/main/scala/code/api/util/Glossary.scala
+++ b/obp-api/src/main/scala/code/api/util/Glossary.scala
@@ -175,41 +175,87 @@ object Glossary extends MdcLoggable {
// Expansion runs once per Resource Doc per Resource Doc cache TTL, and a cold cache expands
// hundreds of docs in one burst, so they share a lookup map rather than each reading the
- // database. The map is keyed on the Dynamic Glossary Item watermark, and the watermark itself
- // is re-read at most once a second.
- private val GlossaryCacheRecheckMillis = 1000L
- private val cachedItemsByTitle =
- new java.util.concurrent.atomic.AtomicReference[(Long, String, Map[String, GlossaryItem])]((0L, "", Map.empty))
+ // database. The map is keyed on two versions, checked on different schedules:
+ // - the glossary cache namespace counter, a Redis read, at most once every ten seconds. Every
+ // Dynamic Glossary Item write bumps it (see dynamicGlossaryItemsChanged), as does the cache
+ // page in API Manager, so a change reaches every node within those ten seconds.
+ // - the Dynamic Glossary Item watermark, a database read, at most once every ten minutes. This
+ // is the fallback for when the bump could not be made, for example with Redis unreachable.
+ private val GlossaryWatermarkRecheckMillis = 600000L
+ private val GlossaryNamespaceRecheckMillis = 10000L
+
+ /**
+ * This holds the Glossary lookup map together with the two versions it was built from and
+ * when each version was last read. A watermarkCheckedAt of 0 means nothing has been built yet.
+ */
+ private case class GlossaryCacheState(
+ watermarkCheckedAt: Long,
+ watermark: String,
+ namespaceCheckedAt: Long,
+ namespaceVersion: Long,
+ itemsByTitle: Map[String, GlossaryItem]
+ ) {
+ def token: String = s"$watermark-ns$namespaceVersion"
+ }
+
+ private val EmptyGlossaryCacheState = GlossaryCacheState(0L, "", 0L, 0L, Map.empty)
+
+ private val cachedGlossaryState =
+ new java.util.concurrent.atomic.AtomicReference[GlossaryCacheState](EmptyGlossaryCacheState)
/**
* Drops the placeholder lookup cache so a write made on this node is reflected at once, rather
* than on the next watermark re-read. Other nodes still pick the write up via the watermark.
*/
- def invalidateGlossaryItemCache(): Unit = cachedItemsByTitle.set((0L, "", Map.empty))
+ def invalidateGlossaryItemCache(): Unit = cachedGlossaryState.set(EmptyGlossaryCacheState)
+
+ /**
+ * Call this after creating, updating or deleting a Dynamic Glossary Item. It bumps the glossary
+ * cache namespace, which every node re-reads every ten seconds, so all of them serve the change
+ * within that time instead of waiting for the ten-minute watermark check. This node's own cache
+ * is dropped at once.
+ *
+ * Both happen only after the request's transaction has committed. Done earlier, a node could
+ * see the bump, read the Glossary before the write was visible, and keep the old text until
+ * its next watermark check.
+ */
+ def dynamicGlossaryItemsChanged(): Unit =
+ code.api.util.http4s.RequestScopeConnection.afterCommit {
+ code.api.Constant.incrementCacheNamespaceVersion(code.api.Constant.GLOSSARY_NAMESPACE)
+ invalidateGlossaryItemCache()
+ }
private def glossaryState: (String, Map[String, GlossaryItem]) = {
val now = System.currentTimeMillis
- val (checkedAt, version, byTitle) = cachedItemsByTitle.get()
- if (checkedAt != 0L && now - checkedAt < GlossaryCacheRecheckMillis) (version, byTitle)
+ val cached = cachedGlossaryState.get()
+ val isBuilt = cached.watermarkCheckedAt != 0L
+ val watermarkIsFresh = isBuilt && now - cached.watermarkCheckedAt < GlossaryWatermarkRecheckMillis
+ val namespaceIsFresh = isBuilt && now - cached.namespaceCheckedAt < GlossaryNamespaceRecheckMillis
+ if (watermarkIsFresh && namespaceIsFresh) (cached.token, cached.itemsByTitle)
else {
- // 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)
+ // The glossary cache namespace version is part of the token: bumping it reloads the
+ // Glossary here and, because the token is in every resource-docs cache key, rebuilds every
+ // cached document that embeds Glossary text.
+ val (watermarkCheckedAt, watermark) =
+ if (watermarkIsFresh) (cached.watermarkCheckedAt, cached.watermark)
+ else (now, dynamicGlossaryItemsVersion)
+ val (namespaceCheckedAt, namespaceVersion) =
+ if (namespaceIsFresh) (cached.namespaceCheckedAt, cached.namespaceVersion)
+ else (now, code.api.Constant.getCacheNamespaceVersion(code.api.Constant.GLOSSARY_NAMESPACE))
+ val checked = GlossaryCacheState(watermarkCheckedAt, watermark, namespaceCheckedAt, namespaceVersion, cached.itemsByTitle)
+ if (isBuilt && checked.token == cached.token) {
+ cachedGlossaryState.set(checked)
+ (checked.token, checked.itemsByTitle)
} else {
// allGlossaryItems yields one Item per exact title, but titles differing only in case
// survive and collapse together in this case-insensitive map. Reversing makes the first
// of those spellings win, as find() used to.
val items = allGlossaryItems
- val rebuilt = items.reverse.map(item => item.title.toLowerCase -> item).toMap
- cachedItemsByTitle.set((now, currentVersion, rebuilt))
+ val rebuilt = checked.copy(itemsByTitle = items.reverse.map(item => item.title.toLowerCase -> item).toMap)
+ cachedGlossaryState.set(rebuilt)
// Only on a real change, so this reports each edit once rather than on every read.
logStaticOverrides(items.filter(_.shadowsStaticItem))
- (currentVersion, rebuilt)
+ (rebuilt.token, rebuilt.itemsByTitle)
}
}
}
diff --git a/obp-api/src/main/scala/code/api/util/http4s/RequestScopeConnection.scala b/obp-api/src/main/scala/code/api/util/http4s/RequestScopeConnection.scala
index 571d0b02d0..c26288867a 100644
--- a/obp-api/src/main/scala/code/api/util/http4s/RequestScopeConnection.scala
+++ b/obp-api/src/main/scala/code/api/util/http4s/RequestScopeConnection.scala
@@ -87,6 +87,13 @@ import scala.concurrent.Future
* - None → endpoint made zero DB calls; pool unaffected, nothing to commit or close.
* - Some(realConn, _) → commit (or rollback on error/cancel), then close realConn.
*
+ * AFTER-COMMIT ACTIONS: a handler that must tell something else about its write (for example
+ * a cache that other nodes watch) registers the action with afterCommit. Inside a transaction
+ * scope it is held until the transaction has committed, and dropped on rollback, so nothing
+ * hears of a change that is not yet, or never will be, visible in the database. The action
+ * list reaches Future workers the same way the proxy does: an IOLocal on the fiber
+ * (requestAfterCommitLocal) copied into a TTL (currentAfterCommit) by fromFuture.
+ *
* METRIC WRITES: recordMetric runs in IO.blocking (blocking pool, no TTL from compute
* thread). currentProxy.get() returns null there, so RequestAwareConnectionManager
* falls back to the pool — metric writes use a separate connection and commit
@@ -144,6 +151,46 @@ object RequestScopeConnection extends MdcLoggable {
override def childValue(parentValue: Connection): Connection = null
}
+ /**
+ * The after-commit actions of the current transaction scope, held for the request fiber.
+ * None outside a transaction scope (GET/HEAD, background tasks).
+ */
+ val requestAfterCommitLocal: IOLocal[Option[java.util.concurrent.ConcurrentLinkedQueue[() => Unit]]] =
+ IOLocal[Option[java.util.concurrent.ConcurrentLinkedQueue[() => Unit]]](None).unsafeRunSync()(IORuntime.global)
+
+ /**
+ * The same action list as requestAfterCommitLocal, carried to Future workers by TtlRunnable.
+ * childValue returns null for the reason given on currentProxy.
+ */
+ val currentAfterCommit: TransmittableThreadLocal[java.util.concurrent.ConcurrentLinkedQueue[() => Unit]] =
+ new TransmittableThreadLocal[java.util.concurrent.ConcurrentLinkedQueue[() => Unit]]() {
+ override def childValue(parentValue: java.util.concurrent.ConcurrentLinkedQueue[() => Unit]): java.util.concurrent.ConcurrentLinkedQueue[() => Unit] = null
+ }
+
+ /**
+ * Runs `action` once the current request's transaction has committed.
+ *
+ * Use this for anything that announces a write to the outside, such as bumping a cache
+ * namespace that other nodes re-read: done before the commit, another node could act on the
+ * announcement, read the database before the write is visible, and cache the old data.
+ *
+ * Inside a transaction scope (POST, PUT, DELETE, PATCH) the action runs after a successful
+ * commit and is dropped if the transaction rolls back. Outside one (GET, HEAD, background
+ * work) every write has already committed on its own, so the action runs at once.
+ * A failing action is logged and does not affect the response or the other actions.
+ */
+ def afterCommit(action: => Unit): Unit = {
+ val actions = currentAfterCommit.get()
+ if (actions == null) runAfterCommitAction(() => action)
+ else actions.add(() => action)
+ }
+
+ private def runAfterCommitAction(action: () => Unit): Unit =
+ try action()
+ catch {
+ case e: Exception => logger.error(s"afterCommit says: an after-commit action failed: ${e.getMessage}", e)
+ }
+
/**
* Wrap a real Connection in a proxy that no-ops commit, rollback, and close.
* All other methods delegate to the real connection.
@@ -196,11 +243,15 @@ object RequestScopeConnection extends MdcLoggable {
*/
def fromFuture[A](fut: => Future[A]): IO[A] =
ensureProxy.flatMap { proxyOpt =>
- IO.defer {
- proxyOpt.foreach(currentProxy.set) // (1) set TTL on current thread T
- val f = fut // (2) submit Future; TtlRunnable captures proxy from T
- currentProxy.remove() // (3) clear TTL on T — T is clean after this point
- IO.fromFuture(IO.pure(f)) // await the already-submitted future
+ requestAfterCommitLocal.get.flatMap { afterCommitOpt =>
+ IO.defer {
+ proxyOpt.foreach(currentProxy.set) // (1) set TTL on current thread T
+ afterCommitOpt.foreach(currentAfterCommit.set)
+ val f = fut // (2) submit Future; TtlRunnable captures proxy from T
+ currentProxy.remove() // (3) clear TTL on T — T is clean after this point
+ currentAfterCommit.remove()
+ IO.fromFuture(IO.pure(f)) // await the already-submitted future
+ }
}
}
@@ -213,6 +264,8 @@ object RequestScopeConnection extends MdcLoggable {
* and share the winner's proxy, so all fibers use one underlying Connection / one
* transaction. On success: commit then close. On error/cancel: rollback then close.
* If no DB call was made: nothing to commit or close (pool unaffected).
+ * Actions registered with afterCommit run after a successful commit (or, with no DB call,
+ * after the route succeeded) and are dropped on error or cancel.
*
* GET/HEAD must NOT be wrapped (they run on auto-commit vendor connections). Used by
* ResourceDocMiddleware and by services that build their own request scope
@@ -231,9 +284,21 @@ object RequestScopeConnection extends MdcLoggable {
p <- deferred.get.map(_._2)
} yield p
- requestLazyAcquire.set(Some(acquireOnce)).bracket(_ =>
+ val afterCommitActions = new java.util.concurrent.ConcurrentLinkedQueue[() => Unit]()
+
+ // Runs only once the transaction has committed, so whatever an action announces is
+ // already visible to every other reader of the database.
+ val runAfterCommitActions: IO[Unit] = IO.blocking {
+ var action = afterCommitActions.poll()
+ while (action != null) {
+ runAfterCommitAction(action)
+ action = afterCommitActions.poll()
+ }
+ }
+
+ (requestLazyAcquire.set(Some(acquireOnce)) *> requestAfterCommitLocal.set(Some(afterCommitActions))).bracket(_ =>
io.guaranteeCase { outcome =>
- deferred.tryGet.flatMap {
+ val finish: IO[Unit] = deferred.tryGet.flatMap {
case None => IO.unit // no DB calls — pool unaffected
case Some((realConn, _)) =>
requestProxyLocal.set(None) *>
@@ -242,11 +307,16 @@ object RequestScopeConnection extends MdcLoggable {
IO.blocking { realConn.commit() }
case _ =>
IO.blocking { try { realConn.rollback() } catch { case _: Exception => () } }
- }) *>
- IO.blocking { try { realConn.close() } catch { case _: Exception => () } }
+ }).guarantee(
+ IO.blocking { try { realConn.close() } catch { case _: Exception => () } }
+ )
+ }
+ outcome match {
+ case Outcome.Succeeded(_) => finish *> runAfterCommitActions
+ case _ => finish *> IO(afterCommitActions.clear())
}
}
- )(_ => requestLazyAcquire.set(None))
+ )(_ => requestLazyAcquire.set(None) *> requestAfterCommitLocal.set(None))
}
/**
diff --git a/obp-api/src/main/scala/code/api/v1_4_0/JSONFactory1_4_0.scala b/obp-api/src/main/scala/code/api/v1_4_0/JSONFactory1_4_0.scala
index 5a38f26ad0..7d43d06a94 100644
--- a/obp-api/src/main/scala/code/api/v1_4_0/JSONFactory1_4_0.scala
+++ b/obp-api/src/main/scala/code/api/v1_4_0/JSONFactory1_4_0.scala
@@ -606,7 +606,7 @@ object JSONFactory1_4_0 extends MdcLoggable{
// (Superset of upstream's specifiedUrl-only fix in 17faa09ac.)
// The Glossary version belongs in the key too: descriptions embed Glossary text, so a Dynamic
// Glossary Item that overrides a static one must not be masked by an hour-old cache entry.
- // The value is read from an in-memory cache that re-checks the database at most once a second.
+ // The value is read from an in-memory cache that re-checks the Dynamic Glossary Items every ten minutes and the glossary cache namespace every ten seconds.
val cacheKey = LOCALISED_RESOURCE_DOC_PREFIX + s"operationId:${operationId}-locale:$locale- isVersion4OrHigher:$isVersion4OrHigher- includeTechnology:$includeTechnology-requestUrl:${resourceDocUpdatedTags.requestUrl}-specifiedUrl:${resourceDocUpdatedTags.specifiedUrl.getOrElse("")}-glossary:${Glossary.glossaryVersionForCacheKey}".intern()
Caching.memoizeSyncWithImMemory(Some(cacheKey))(CREATE_LOCALISED_RESOURCE_DOC_JSON_TTL.seconds) {
val fieldsDescription =
diff --git a/obp-api/src/main/scala/code/api/v3_0_0/Http4s300.scala b/obp-api/src/main/scala/code/api/v3_0_0/Http4s300.scala
index 0a4c35341b..a931f11878 100644
--- a/obp-api/src/main/scala/code/api/v3_0_0/Http4s300.scala
+++ b/obp-api/src/main/scala/code/api/v3_0_0/Http4s300.scala
@@ -2002,7 +2002,8 @@ object Http4s300 {
for {
_ <- NewStyle.function.hasEntitlement("", user.userId, ApiRole.canReadAggregateMetrics, Some(cc))
httpParams <- NewStyle.function.extractHttpParamsFromUrl(req.uri.renderString)
- (obpQueryParams, _) <- createQueriesByHttpParamsFuture(httpParams, Some(cc))
+ limitedHttpParams <- APIMetrics.limitAggregateMetricsWindow(httpParams, Some(cc))
+ (obpQueryParams, _) <- createQueriesByHttpParamsFuture(limitedHttpParams, Some(cc))
aggregateMetrics <- APIMetrics.apiMetrics.vend.getAllAggregateMetricsFuture(obpQueryParams, false) map {
x => unboxFullOrFail(x, Some(cc), GetAggregateMetricsError)
}
@@ -2026,7 +2027,9 @@ object Http4s300 {
|&verb=GET&anon=false&app_name=MapperPostman
|&exclude_app_names=API-EXPLORER,API-Manager,SOFI,null
|
- |1 from_date (defaults to the day before the current date): eg:from_date=$DateWithMsExampleString
+ |**Date range limit.** One call covers at most ${code.metrics.MetricsProps.aggregateMetricsMaxDays} days of metrics on this instance. A longer range between from_date and to_date is refused with 400 (OBP-10069); to report on a longer period, make one call per period.
+ |
+ |1 from_date (defaults to ${code.metrics.MetricsProps.aggregateMetricsMaxDays} days before to_date): eg:from_date=$DateWithMsExampleString
|
|2 to_date (defaults to the current date) eg:to_date=$DateWithMsExampleString
|
@@ -2061,7 +2064,7 @@ object Http4s300 {
""".stripMargin,
EmptyBody,
aggregateMetricsJSONV300,
- List(AuthenticatedUserIsRequired, UserHasMissingRoles, UnknownError),
+ List(AuthenticatedUserIsRequired, UserHasMissingRoles, AggregateMetricsDateRangeTooLong, UnknownError),
List(apiTagMetric, apiTagAggregateMetrics),
Some(List(canReadAggregateMetrics)),
http4sPartialFunction = Some(getAggregateMetrics)
diff --git a/obp-api/src/main/scala/code/api/v5_1_0/Http4s510.scala b/obp-api/src/main/scala/code/api/v5_1_0/Http4s510.scala
index 0cda869526..0c7b9d6573 100644
--- a/obp-api/src/main/scala/code/api/v5_1_0/Http4s510.scala
+++ b/obp-api/src/main/scala/code/api/v5_1_0/Http4s510.scala
@@ -275,7 +275,8 @@ object Http4s510 {
implicit val cc: code.api.util.CallContext = req.callContext
for {
httpParams <- NewStyle.function.extractHttpParamsFromUrl(req.uri.renderString)
- (obpQueryParams, _) <- createQueriesByHttpParamsFuture(httpParams, Some(cc))
+ limitedHttpParams <- APIMetrics.limitAggregateMetricsWindow(httpParams, Some(cc))
+ (obpQueryParams, _) <- createQueriesByHttpParamsFuture(limitedHttpParams, Some(cc))
aggregateMetrics <- APIMetrics.apiMetrics.vend.getAllAggregateMetricsFuture(obpQueryParams, true)
.map(x => unboxFullOrFail(x, Some(cc), GetAggregateMetricsError))
} yield createAggregateMetricJson(aggregateMetrics).headOption
@@ -296,7 +297,9 @@ object Http4s510 {
|&verb=GET&anon=false&app_name=MapperPostman
|&exclude_app_names=API-EXPLORER,API-Manager,SOFI,null
|
- |1 from_date (defaults to the day before the current date): eg:from_date=$DateWithMsExampleString
+ |**Date range limit.** One call covers at most ${code.metrics.MetricsProps.aggregateMetricsMaxDays} days of metrics on this instance. A longer range between from_date and to_date is refused with 400 (OBP-10069); to report on a longer period, make one call per period.
+ |
+ |1 from_date (defaults to ${code.metrics.MetricsProps.aggregateMetricsMaxDays} days before to_date): eg:from_date=$DateWithMsExampleString
|
|2 to_date (defaults to the current date) eg:to_date=$DateWithMsExampleString
|
@@ -330,7 +333,7 @@ object Http4s510 {
|
""".stripMargin,
EmptyBody, aggregateMetricsJSONV300,
- List(AuthenticatedUserIsRequired, UserHasMissingRoles, UnknownError),
+ List(AuthenticatedUserIsRequired, UserHasMissingRoles, AggregateMetricsDateRangeTooLong, UnknownError),
List(apiTagMetric, apiTagAggregateMetrics),
Some(List(canReadAggregateMetrics)),
http4sPartialFunction = Some(getAggregateMetrics)
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 bc3aee14d3..722f8c3093 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
@@ -724,8 +724,9 @@ object Http4s600 {
throw new Exception(s"$ExcludeParametersNotSupported Parameters found: [${excludes.map(_.name).mkString(", ")}]")
else true
}
- (obpQueryParams, callContext) <- createQueriesByHttpParamsFuture(
+ limitedHttpParams <- APIMetrics.limitAggregateMetricsWindow(
APIMetrics.applyMetricsFromDateDefault(httpParams), cc.callContext)
+ (obpQueryParams, callContext) <- createQueriesByHttpParamsFuture(limitedHttpParams, cc.callContext)
// isNewVersion = true: v6 is include_* style (exclude_* is rejected above). With
// false the include_app_names / include_url_patterns /
// include_implemented_by_partial_functions filters were silently ignored.
@@ -7859,6 +7860,8 @@ object Http4s600 {
|This prevents accidentally querying all metrics since Unix Epoch and ensures reasonable response times.
|For historical/reporting queries, always explicitly specify your desired `from_date`.
|
+ |**Date range limit.** One call covers at most ${code.metrics.MetricsProps.aggregateMetricsMaxDays} days of metrics on this instance. A longer range between from_date and to_date is refused with 400 (OBP-10069); to report on a longer period, make one call per period.
+ |
|**IMPORTANT: Smart Caching & Performance**
|
|This endpoint uses intelligent two-tier caching to optimize performance:
@@ -7943,6 +7946,7 @@ object Http4s600 {
List(
AuthenticatedUserIsRequired,
UserHasMissingRoles,
+ AggregateMetricsDateRangeTooLong,
UnknownError
),
List(apiTagMetric, apiTagAggregateMetrics),
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 6c3262fe3b..135a3c57b3 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
@@ -3710,26 +3710,45 @@ object Http4s700 {
private def normaliseGlossaryTitle(title: String): String =
Option(title).getOrElse("").trim.toLowerCase.replaceAll("[-_/\\s]+", " ")
- // Rendering every Item through Pegdown is too expensive to do per request, and a lazy val
- // would never pick up a Dynamic Glossary Item change, so the merged Glossary is cached against
- // the Dynamic Item watermark. Each entry keeps the database row it came from, so one lookup
- // serves both the Glossary and a single Item without going back to the table.
- private type GlossaryEntry = (Glossary.GlossaryItem, Option[code.glossaryitem.DynamicGlossaryItemTrait])
+ // Rendering every Item through Pegdown is too expensive to do per request: expanding the
+ // placeholders (which can inline whole other Items) and converting the markdown to HTML for
+ // the whole Glossary took several seconds a call in production. A lazy val would never pick
+ // up a Dynamic Glossary Item change, so the merged Glossary, already expanded and rendered,
+ // is cached against the Glossary version token. That token covers the Dynamic Item watermark
+ // and the glossary cache namespace; the first is re-read every ten minutes and the second every ten seconds.
+ // Each entry keeps the database row it came from, so one lookup serves both the Glossary and
+ // a single Item without going back to the table.
+ private case class GlossaryEntry(
+ item: Glossary.GlossaryItem,
+ meta: Option[code.glossaryitem.DynamicGlossaryItemTrait],
+ expandedDescription: JSONFactory700.GlossaryItemDescriptionJsonV700
+ )
private val cachedGlossaryEntries =
new java.util.concurrent.atomic.AtomicReference[Option[(String, List[GlossaryEntry])]](None)
private def glossaryEntries: List[GlossaryEntry] = {
- val version = Glossary.dynamicGlossaryItemsVersion
+ val version = Glossary.glossaryVersionForCacheKey
cachedGlossaryEntries.get() match {
case Some((cachedVersion, entries)) if cachedVersion == version => entries
case _ =>
- val rowsByTitle = DynamicGlossaryItems.dynamicGlossaryItem.vend.getAllDynamicGlossaryItems
- .map(_.map(row => row.title.toLowerCase -> row).toMap)
- .getOrElse(Map.empty[String, code.glossaryitem.DynamicGlossaryItemTrait])
- val entries = APIUtil.getGlossaryItems.map(item => (item, rowsByTitle.get(item.title.toLowerCase)))
- cachedGlossaryEntries.set(Some((version, entries)))
- entries
+ // One request builds the Glossary while the others wait for it, rather than every
+ // request that arrives during a rebuild rendering the whole Glossary again itself.
+ cachedGlossaryEntries.synchronized {
+ cachedGlossaryEntries.get() match {
+ case Some((cachedVersion, entries)) if cachedVersion == version => entries
+ case _ =>
+ val rowsByTitle = DynamicGlossaryItems.dynamicGlossaryItem.vend.getAllDynamicGlossaryItems
+ .map(_.map(row => row.title.toLowerCase -> row).toMap)
+ .getOrElse(Map.empty[String, code.glossaryitem.DynamicGlossaryItemTrait])
+ val entries = APIUtil.getGlossaryItems.map { item =>
+ GlossaryEntry(item, rowsByTitle.get(item.title.toLowerCase),
+ JSONFactory700.createGlossaryItemDescriptionJsonV700(item, expanded = true))
+ }
+ cachedGlossaryEntries.set(Some((version, entries)))
+ entries
+ }
+ }
}
}
@@ -3760,7 +3779,8 @@ object Http4s700 {
// alone found nothing at all. Every word must appear somewhere in the Item, and the
// Items whose titles hold all of them are listed first, so an exact title still leads.
val searchWords: List[String] = search.toList.flatMap(_.split("\\s+")).filter(_.nonEmpty)
- val inSource = glossaryEntries.filter { case (item, _) =>
+ val inSource = glossaryEntries.filter { entry =>
+ val item = entry.item
source match {
case "static" => !item.isDynamic
case "dynamic" => item.isDynamic
@@ -3770,11 +3790,11 @@ object Http4s700 {
val matching =
if (searchWords.isEmpty) inSource
else {
- val (titleMatches, bodyMatches) = inSource.filter { case (item, _) =>
- val titleAndText = s"${item.title}\n${item.textDescription}".toLowerCase
+ val (titleMatches, bodyMatches) = inSource.filter { entry =>
+ val titleAndText = s"${entry.item.title}\n${entry.item.textDescription}".toLowerCase
searchWords.forall(titleAndText.contains)
- }.partition { case (item, _) =>
- val title = item.title.toLowerCase
+ }.partition { entry =>
+ val title = entry.item.title.toLowerCase
searchWords.forall(title.contains)
}
titleMatches ++ bodyMatches
@@ -3785,8 +3805,9 @@ object Http4s700 {
case None => matching
}
JSONFactory700.createGlossaryJsonV700(
- page.map { case (item, meta) =>
- JSONFactory700.createServedGlossaryItemJsonV700(item, meta, expanded = true, includeAuthor = cc.user.isDefined)
+ page.map { entry =>
+ JSONFactory700.createServedGlossaryItemJsonV700(entry.item, entry.meta, entry.expandedDescription,
+ includeAuthor = cc.user.isDefined)
},
totalCount = matching.size
)
@@ -3860,9 +3881,15 @@ object Http4s700 {
// Glossary itself, so this always answers with the text the Glossary is serving.
json <- Future {
val wanted = normaliseGlossaryTitle(titleSegment)
- Box(glossaryEntries.find { case (item, _) => normaliseGlossaryTitle(item.title) == wanted }
- .map { case (item, meta) =>
- JSONFactory700.createServedGlossaryItemJsonV700(item, meta, expanded, includeAuthor = cc.user.isDefined)
+ Box(glossaryEntries.find(entry => normaliseGlossaryTitle(entry.item.title) == wanted)
+ .map { entry =>
+ // The expanded text is already rendered in the cache; the text as authored is
+ // asked for rarely (by an editor) and is one Item, so it is rendered here.
+ val description =
+ if (expanded) entry.expandedDescription
+ else JSONFactory700.createGlossaryItemDescriptionJsonV700(entry.item, expanded = false)
+ JSONFactory700.createServedGlossaryItemJsonV700(entry.item, entry.meta, description,
+ includeAuthor = cc.user.isDefined)
})
}.map(unboxFullOrFail(_, Some(cc), GlossaryItemNotFound, 404))
} yield json
@@ -3934,8 +3961,8 @@ object Http4s700 {
createdByUserId = user.userId
)
}.map(unboxFullOrFail(_, Some(cc), CreateGlossaryItemError, 400))
- // Reflect the write on this node at once; other nodes pick it up via the watermark.
- _ = Glossary.invalidateGlossaryItemCache()
+ // Reflect the write on every node within ten seconds of the commit.
+ _ = Glossary.dynamicGlossaryItemsChanged()
} yield JSONFactory700.createGlossaryItemJsonV700(created)
}
}
@@ -3993,7 +4020,7 @@ object Http4s700 {
DynamicGlossaryItems.dynamicGlossaryItem.vend.updateDynamicGlossaryItem(
titleSegment, body.description, body.overrides_static_item)
}.map(unboxFullOrFail(_, Some(cc), UpdateGlossaryItemError, 400))
- _ = Glossary.invalidateGlossaryItemCache()
+ _ = Glossary.dynamicGlossaryItemsChanged()
} yield JSONFactory700.createGlossaryItemJsonV700(updated)
}
}
@@ -4046,7 +4073,7 @@ object Http4s700 {
.map(unboxFullOrFail(_, Some(cc), GlossaryItemNotFound, 404))
_ <- Future(DynamicGlossaryItems.dynamicGlossaryItem.vend.deleteDynamicGlossaryItem(titleSegment))
.map(unboxFullOrFail(_, Some(cc), DeleteGlossaryItemError, 400))
- _ = Glossary.invalidateGlossaryItemCache()
+ _ = Glossary.dynamicGlossaryItemsChanged()
} yield ()
}
}
@@ -7461,6 +7488,9 @@ object Http4s700 {
// Telemetry, for people; Prometheus reads the separate Telemetry port instead.
resourceDocs ++= Http4s700Telemetry.resourceDocs
+ // What has been granted across all Users and Consumers: every Scope, and the Role names anyone holds.
+ resourceDocs ++= Http4s700Grants.resourceDocs
+
// IP penalties: an operator's temporary per-minute limit on one address.
resourceDocs ++= Http4s700IpPenalties.resourceDocs
diff --git a/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Grants.scala b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Grants.scala
new file mode 100644
index 0000000000..34ae2d8bcd
--- /dev/null
+++ b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Grants.scala
@@ -0,0 +1,131 @@
+/**
+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, UserOrApplication}
+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, Glossary}
+import code.entitlement.Entitlement
+import code.scope.Scope
+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
+
+/**
+ * What has been granted on this instance, across all Users and Consumers: every Scope, and the Role names
+ * anyone holds (as an Entitlement or a Scope). A monitoring or code-review service such as OBP-Sentinel uses
+ * the Role names to learn which endpoints can be reached here, without learning who holds them.
+ *
+ * Declared in its own object to keep Http4s700's initialiser under the JVM's 64KB method limit.
+ */
+object Http4s700Grants {
+
+ 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]()
+
+ // Route: GET /obp/v7.0.0/scopes
+ lazy val getAllScopes: HttpRoutes[IO] = HttpRoutes.of[IO] {
+ case req @ GET -> `prefixPath` / "scopes" =>
+ EndpointHelpers.executeFuture(req) {
+ Scope.scope.vend.getScopesFuture().map(scopes => JSONFactory700.createAllScopesJsonV700(scopes.openOr(Nil)))
+ }
+ }
+
+ resourceDocs += ResourceDoc(
+ implementedInApiVersion,
+ nameOf(getAllScopes),
+ "GET",
+ "/scopes",
+ "Get all Scopes",
+ s"""Get every Scope on this instance, for all Consumers: its `bank_id`, `role_name` and `consumer_id`.
+ |
+ |A Scope is a Role granted to a Consumer (an application) rather than to a User; see
+ |${Glossary.getGlossaryItemLink("API.Endpoint Auth Modes")}. For the Scopes of one Consumer, see Get Scopes for Consumer.
+ |
+ |**Who may call it.** A User with the Role CanGetAllScopes, or an application whose Consumer holds it as a Scope.
+ |""".stripMargin,
+ EmptyBody,
+ JSONFactory700.allScopesJsonV700Example,
+ List($AuthenticatedUserIsRequired, UserHasMissingRoles, UnknownError),
+ List(apiTagScope, apiTagRole),
+ Some(List(canGetAllScopes)),
+ authMode = UserOrApplication,
+ http4sPartialFunction = Some(getAllScopes)
+ )
+
+ // Route: GET /obp/v7.0.0/reachable-roles
+ lazy val getReachableRoles: HttpRoutes[IO] = HttpRoutes.of[IO] {
+ case req @ GET -> `prefixPath` / "reachable-roles" =>
+ EndpointHelpers.executeFuture(req) {
+ for {
+ entitlements <- Entitlement.entitlement.vend.getEntitlementsFuture()
+ scopes <- Scope.scope.vend.getScopesFuture()
+ } yield JSONFactory700.createReachableRolesJsonV700(
+ entitlements.openOr(Nil).map(_.roleName) ++ scopes.openOr(Nil).map(_.roleName))
+ }
+ }
+
+ resourceDocs += ResourceDoc(
+ implementedInApiVersion,
+ nameOf(getReachableRoles),
+ "GET",
+ "/reachable-roles",
+ "Get Reachable Roles",
+ s"""Get the names of the Roles that someone holds on this instance, as an Entitlement (a User) or as a
+ |Scope (a Consumer), each listed once.
+ |
+ |An endpoint that needs a Role can only be reached if someone holds that Role, so together with the
+ |Resource Docs this tells which endpoints can be reached here. Nothing else is returned: no Users, no
+ |Consumers, no bank ids.
+ |
+ |**Who may call it.** A User with the Role CanGetReachableRoles, or an application whose Consumer
+ |holds it as a Scope, such as a code-review service run as a Platform App
+ |(see ${Glossary.getGlossaryItemLink("Platform Apps")}).
+ |""".stripMargin,
+ EmptyBody,
+ JSONFactory700.reachableRolesJsonV700Example,
+ List($AuthenticatedUserIsRequired, UserHasMissingRoles, UnknownError),
+ List(apiTagRole, apiTagEntitlement),
+ Some(List(canGetReachableRoles)),
+ authMode = UserOrApplication,
+ http4sPartialFunction = Some(getReachableRoles)
+ )
+}
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 5bf00346a6..c96123d184 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
@@ -1021,27 +1021,37 @@ object JSONFactory700 extends MdcLoggable with code.api.util.CustomJsonFormats {
case class GlossaryJsonV700(glossary_items: List[GlossaryItemJsonV700], total_count: Int)
/**
- * One Item as the Glossary serves it.
+ * The description of one Item as the Glossary serves it, in markdown and in HTML.
*
* `expanded` controls the Glossary placeholders that Items use to quote each other: true gives
* the text a reader wants, false the text as authored, which is what an editor must PUT back.
- * `includeAuthor` keeps created_by_user_id out of anonymous responses.
+ * This is the expensive part of serving an Item (placeholder expansion and a markdown to HTML
+ * conversion), so the Glossary endpoints cache its result rather than calling it per request.
+ */
+ def createGlossaryItemDescriptionJsonV700(item: Glossary.GlossaryItem, expanded: Boolean): GlossaryItemDescriptionJsonV700 = {
+ val raw = item.description()
+ val markdown = if (expanded) Glossary.expandGlossaryPlaceholders(raw) else raw
+ GlossaryItemDescriptionJsonV700(
+ markdown = markdown,
+ html = PegdownOptions.convertPegdownToHtmlTweaked(markdown)
+ )
+ }
+
+ /**
+ * One Item as the Glossary serves it, with its description already rendered by
+ * createGlossaryItemDescriptionJsonV700. `includeAuthor` keeps created_by_user_id out of
+ * anonymous responses.
*/
def createServedGlossaryItemJsonV700(
item: Glossary.GlossaryItem,
meta: Option[code.glossaryitem.DynamicGlossaryItemTrait],
- expanded: Boolean,
+ description: GlossaryItemDescriptionJsonV700,
includeAuthor: Boolean
- ): GlossaryItemJsonV700 = {
- val raw = item.description()
- val markdown = if (expanded) Glossary.expandGlossaryPlaceholders(raw) else raw
+ ): GlossaryItemJsonV700 =
GlossaryItemJsonV700(
glossary_item_id = meta.map(_.glossaryItemId),
title = item.title,
- description = GlossaryItemDescriptionJsonV700(
- markdown = markdown,
- html = PegdownOptions.convertPegdownToHtmlTweaked(markdown)
- ),
+ description = description,
is_dynamic = item.isDynamic,
overrides_static_item = item.overridesStaticItem,
shadows_static_glossary_item = item.shadowsStaticItem,
@@ -1049,7 +1059,6 @@ object JSONFactory700 extends MdcLoggable with code.api.util.CustomJsonFormats {
created_at = meta.map(_.createdAt),
updated_at = meta.map(_.updatedAt)
)
- }
def createGlossaryJsonV700(items: List[GlossaryItemJsonV700], totalCount: Int): GlossaryJsonV700 =
GlossaryJsonV700(glossary_items = items, total_count = totalCount)
@@ -2257,6 +2266,30 @@ object JSONFactory700 extends MdcLoggable with code.api.util.CustomJsonFormats {
scopes = List(CurrentConsumerScopeJsonV700(role_name = "CanGetDynamicEntityDefinitions", bank_id = "SYS"))
)
+ /** One Scope on this instance: which Consumer holds which Role, and where. */
+ case class ScopeJsonV700(
+ bank_id: String,
+ role_name: String,
+ consumer_id: String
+ )
+
+ case class ScopesJsonV700(scopes: List[ScopeJsonV700])
+
+ def createAllScopesJsonV700(scopes: List[code.scope.Scope]): ScopesJsonV700 =
+ ScopesJsonV700(scopes.map(s => ScopeJsonV700(bank_id = s.bankId, role_name = s.roleName, consumer_id = s.consumerId))
+ .sortBy(s => (s.role_name, s.bank_id, s.consumer_id)))
+
+ lazy val allScopesJsonV700Example = ScopesJsonV700(List(
+ ScopeJsonV700(bank_id = "", role_name = "CanGetTelemetry", consumer_id = ExampleValue.consumerIdExample.value)))
+
+ /** The Role names someone holds, as an Entitlement or a Scope, each once. No Users, Consumers or bank ids. */
+ case class ReachableRolesJsonV700(role_names: List[String])
+
+ def createReachableRolesJsonV700(roleNames: List[String]): ReachableRolesJsonV700 =
+ ReachableRolesJsonV700(roleNames.distinct.sorted)
+
+ lazy val reachableRolesJsonV700Example = ReachableRolesJsonV700(List("CanGetCustomersAtOneBank", "CanGetTelemetry"))
+
case class PasswordPolicyJsonV700(
description: String,
min_length: Int,
diff --git a/obp-api/src/main/scala/code/metrics/APIMetrics.scala b/obp-api/src/main/scala/code/metrics/APIMetrics.scala
index 2b53007666..8ea31eb536 100644
--- a/obp-api/src/main/scala/code/metrics/APIMetrics.scala
+++ b/obp-api/src/main/scala/code/metrics/APIMetrics.scala
@@ -30,6 +30,7 @@ package code.metrics
import java.util.{Calendar, Date}
import code.api.util.{APIUtil, CallContext, OBPQueryParam}
import code.api.util.APIUtil.{HTTPParam, createQueriesByHttpParamsFuture}
+import code.api.util.ErrorMessages.AggregateMetricsDateRangeTooLong
import com.openbankproject.commons.ExecutionContext.Implicits.global
import com.openbankproject.commons.util.ApiVersion
import net.liftweb.common.Box
@@ -77,6 +78,45 @@ object APIMetrics extends SimpleInjector {
}
}
+ /**
+ * This keeps one aggregate-metrics call to at most MetricsProps.aggregateMetricsMaxDays days of
+ * metrics. An aggregate reads every metric row in its range, so a call covering months of
+ * metrics ran past the endpoint timeout and held a database connection the whole time.
+ *
+ * A request without from_date gets one that many days before its end: to_date, or now when
+ * to_date is absent. It is rounded down to the minute so repeated calls share a cache key. A
+ * request whose range is longer than the limit fails with 400. The end of the range counts as
+ * now at the latest, because there are no metrics after now. A date that cannot be parsed is
+ * left alone here, so the usual date format error reports it.
+ */
+ def limitAggregateMetricsWindow(httpParams: List[HTTPParam], callContext: Option[CallContext]): Future[List[HTTPParam]] = {
+ val maxDays = MetricsProps.aggregateMetricsMaxDays
+ val maxMillis = maxDays * 24L * 60 * 60 * 1000
+ val now = System.currentTimeMillis()
+ val hasFromDate = httpParams.exists(p => p.name == "from_date" || p.name == "obp_from_date")
+ val hasToDate = httpParams.exists(p => p.name == "to_date" || p.name == "obp_to_date")
+ val endMillis: Option[Long] =
+ if (hasToDate) APIUtil.getToDate(httpParams).toOption.map(to => math.min(to.value.getTime, now))
+ else Some(now)
+ if (!hasFromDate) {
+ val withFromDate = endMillis.map { end =>
+ val fromMillis = end - maxMillis
+ val fromDate = new Date(fromMillis - fromMillis % 60000L)
+ HTTPParam("from_date", List(APIUtil.DateWithMsFormat.format(fromDate))) :: httpParams
+ }.getOrElse(httpParams)
+ Future.successful(withFromDate)
+ } else {
+ val fromMillis = APIUtil.getFromDate(httpParams).toOption.map(_.value.getTime)
+ val tooLong = (fromMillis, endMillis) match {
+ case (Some(from), Some(end)) => end - from > maxMillis
+ case _ => false
+ }
+ code.util.Helper.booleanToFuture(
+ s"$AggregateMetricsDateRangeTooLong The longest range allowed is $maxDays days.", 400, callContext
+ )(!tooLong).map(_ => httpParams)
+ }
+ }
+
// One shared fetch path for metrics-reading endpoints: builds OBPQueryParams
// from the http params (with the from_date default applied) and runs the query.
// lockedUserIds pins the user filter server-side (for self-service endpoints
diff --git a/obp-api/src/main/scala/code/metrics/DoobieMetricsQueries.scala b/obp-api/src/main/scala/code/metrics/DoobieMetricsQueries.scala
index 4519dcda29..87d4958673 100644
--- a/obp-api/src/main/scala/code/metrics/DoobieMetricsQueries.scala
+++ b/obp-api/src/main/scala/code/metrics/DoobieMetricsQueries.scala
@@ -53,6 +53,17 @@ import scala.concurrent.{ExecutionContext, Future}
*/
object DoobieMetricsQueries {
+ /**
+ * Runs a query and collects its rows, with a JDBC query timeout equal to the endpoint timeout.
+ * When a metrics query outlives its endpoint the caller has already been answered with a 504,
+ * but without a timeout the database carries on and the query keeps its pool connection for as
+ * long as it takes. Metrics queries only read, so cancelling one never interrupts a write.
+ */
+ private def listWithTimeout[A: Read](query: Fragment): ConnectionIO[List[A]] =
+ query.execWith(
+ HPS.setQueryTimeout(DBUtil.queryTimeoutSeconds).flatMap(_ => HPS.executeQuery(HRS.list[A]))
+ )
+
/**
* Get aggregate metrics (count, avg, min, max duration) for the given time range and filters.
*
@@ -112,7 +123,7 @@ object DoobieMetricsQueries {
val conditions = buildFilterConditions(filters, isNewVersion)
val fullQuery = baseQuery ++ conditions
- fullQuery.query[(Long, Option[Double], Option[Double], Option[Double], Long, Long, Long, Long)].to[List].map { rows =>
+ listWithTimeout[(Long, Option[Double], Option[Double], Option[Double], Long, Long, Long, Long)](fullQuery).map { rows =>
rows.map { case (count, avgOpt, minOpt, maxOpt, distinctUsers, distinctConsumers, consentCalls, distinctConsents) =>
AggregateMetrics(
count.toInt,
@@ -197,7 +208,7 @@ object DoobieMetricsQueries {
val fullQuery = baseQuery ++ conditions ++ groupAndOrder ++ limitClause
- fullQuery.query[(Long, String, String)].to[List].map { rows =>
+ listWithTimeout[(Long, String, String)](fullQuery).map { rows =>
rows.map { case (count, partialFunction, version) =>
TopApi(count.toInt, partialFunction, version)
}
@@ -213,6 +224,75 @@ object DoobieMetricsQueries {
* @param filters Filter options
* @return List of TopConsumer sorted by count descending
*/
+ /**
+ * Get the top consumers by API call count, matching a metric row to its consumer by the
+ * consumer's name (metric.appname = consumer.name). This is the grouping v3.1.0's
+ * GET /management/metrics/top-consumers has always used; later versions group by
+ * metric.consumerid instead (see getTopConsumersByConsumerId).
+ *
+ * Only the exclude_* list filters apply here, as they always have for this endpoint.
+ *
+ * @return List of TopConsumer sorted by count descending
+ */
+ def getTopConsumersByAppName(
+ fromDate: Date,
+ toDate: Date,
+ limit: Int,
+ filters: MetricsQueryFilters
+ ): List[TopConsumer] =
+ DoobieUtil.runQuery(buildTopConsumersByAppNameQuery(fromDate, toDate, limit, filters))
+
+ private def buildTopConsumersByAppNameQuery(
+ fromDate: Date,
+ toDate: Date,
+ limit: Int,
+ filters: MetricsQueryFilters
+ ): ConnectionIO[List[TopConsumer]] = {
+ val fromTs = new java.sql.Timestamp(fromDate.getTime)
+ val toTs = new java.sql.Timestamp(toDate.getTime)
+
+ val isSqlServer = DBUtil.isSqlServer
+ val top = if (isSqlServer) fr"TOP($limit)" else fr""
+
+ val baseQuery = fr"SELECT" ++ top ++ fr"""count(*), con.consumerid, m.appname, con.developeremail
+ FROM metric m
+ JOIN consumer con ON m.appname = con.name
+ WHERE m.date_c >= $fromTs
+ AND m.date_c <= $toTs
+ """
+
+ val conditions = List(
+ filters.consumerId.map(v => fr"AND con.consumerid = $v"),
+ filters.userId.map(v => fr"AND m.userid = $v"),
+ filters.implementedByPartialFunction.map(v => fr"AND m.implementedbypartialfunction = $v"),
+ filters.implementedInVersion.map(v => fr"AND m.implementedinversion = $v"),
+ filters.url.map(v => fr"AND m.url = $v"),
+ filters.appName.map(v => fr"AND m.appname = $v"),
+ filters.verb.map(v => fr"AND m.verb = $v"),
+ filters.httpStatusCode.map(v => fr"AND m.httpcode = $v"),
+ filters.anon.map {
+ case true => fr"AND m.userid = 'null'"
+ case false => fr"AND m.userid != 'null'"
+ },
+ filters.excludeAppNames.filter(_.nonEmpty).map(names => buildNotInClause("m.appname", names)),
+ filters.excludeImplementedByPartialFunctions.filter(_.nonEmpty).map(names => buildNotInClause("m.implementedbypartialfunction", names)),
+ filters.excludeUrlPatterns.filter(_.nonEmpty).map(patterns => buildNotLikeClause("m.url", patterns))
+ ).flatten.foldLeft(fr"")(_ ++ _)
+
+ val groupAndOrder = fr"""
+ GROUP BY m.appname, con.developeremail, con.id, con.consumerid
+ ORDER BY count(*) DESC
+ """
+
+ val limitClause = if (isSqlServer) fr"" else fr"LIMIT $limit"
+
+ listWithTimeout[(Long, String, String, String)](baseQuery ++ conditions ++ groupAndOrder ++ limitClause).map { rows =>
+ rows.map { case (count, consumerId, appName, developerEmail) =>
+ TopConsumer(count.toInt, consumerId, appName, developerEmail)
+ }
+ }
+ }
+
/**
* Get top consumers by API call count, grouped by metric.consumerid.
*
@@ -299,7 +379,7 @@ object DoobieMetricsQueries {
val fullQuery = baseQuery ++ conditions ++ groupAndOrder ++ limitClause
- fullQuery.query[(Long, String, String, String)].to[List].map { rows =>
+ listWithTimeout[(Long, String, String, String)](fullQuery).map { rows =>
rows.map { case (count, consumerId, appName, developerEmail) =>
TopConsumer(count.toInt, consumerId, appName, developerEmail)
}
@@ -388,7 +468,7 @@ object DoobieMetricsQueries {
val fullQuery = baseQuery ++ conditions ++ groupAndOrder ++ limitClause
- fullQuery.query[(Long, String, String)].to[List].map { rows =>
+ listWithTimeout[(Long, String, String)](fullQuery).map { rows =>
rows.map { case (count, userId, userName) =>
TopUser(count.toInt, userId, userName)
}
@@ -459,7 +539,7 @@ object DoobieMetricsQueries {
val fullQuery = baseQuery ++ conditions ++ groupAndOrder ++ limitClause
- fullQuery.query[(Long, Long, String, String, String)].to[List].map { rows =>
+ listWithTimeout[(Long, Long, String, String, String)](fullQuery).map { rows =>
rows.map { case (count, _, appName, email, consumerId) =>
TopConsumer(count.toInt, consumerId, appName, email)
}
diff --git a/obp-api/src/main/scala/code/metrics/MappedMetrics.scala b/obp-api/src/main/scala/code/metrics/MappedMetrics.scala
index 4d73b046bf..0b48223d79 100644
--- a/obp-api/src/main/scala/code/metrics/MappedMetrics.scala
+++ b/obp-api/src/main/scala/code/metrics/MappedMetrics.scala
@@ -231,40 +231,26 @@ object MappedMetrics extends APIMetrics with MdcLoggable{
saved
}
- private def trueOrFalse(condition: Boolean): String = if (condition) s"1=1" else s"0=1"
- private def falseOrTrue(condition: Boolean): String = if (condition) s"0=1" else s"1=1"
-
- private def sqlFriendly(value : Option[String]): String = {
- value match {
- case Some(value) => s"'$value'"
- case None => "null"
-
- }
- }
-
- private def sqlFriendlyInt(value : Option[Int]): String = {
- value match {
- case Some(value) => s"$value"
- case None => "null"
- }
- }
-
- /**
- * Formats a Date as an ISO 8601 timestamp string for use in SQL queries.
- * Uses the format yyyy-MM-dd'T'HH:mm:ss.SSS with the 'T' separator, which is
- * universally safe across databases (PostgreSQL, SQL Server, H2, etc.).
- *
- * The 'T' separator is critical for SQL Server compatibility - without it,
- * SQL Server may misinterpret the date based on regional/language settings.
- *
- * @param date The date to format
- * @return ISO 8601 formatted timestamp string (e.g., "2024-01-15T10:30:45.123")
- */
- private def sqlTimestamp(date: Date): String = {
- val sdf = new SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss.SSS")
- sdf.setTimeZone(TimeZone.getTimeZone("UTC"))
- sdf.format(date)
- }
+ /** The metric filters a request asked for, for the queries in DoobieMetricsQueries. */
+ private def metricsQueryFilters(queryParams: List[OBPQueryParam]): MetricsQueryFilters =
+ MetricsQueryFilters(
+ consumerId = queryParams.collectFirst { case OBPConsumerId(value) => value },
+ userId = queryParams.collectFirst { case OBPUserId(value) => value },
+ url = queryParams.collectFirst { case OBPUrl(value) => value },
+ appName = queryParams.collectFirst { case OBPAppName(value) => value },
+ implementedByPartialFunction = queryParams.collectFirst { case OBPImplementedByPartialFunction(value) => value },
+ implementedInVersion = queryParams.collectFirst { case OBPImplementedInVersion(value) => value },
+ verb = queryParams.collectFirst { case OBPVerb(value) => value },
+ anon = queryParams.collectFirst { case OBPAnon(value) => value },
+ correlationId = queryParams.collectFirst { case OBPCorrelationId(value) => value },
+ httpStatusCode = queryParams.collectFirst { case OBPHttpStatusCode(value) => value },
+ excludeAppNames = queryParams.collectFirst { case OBPExcludeAppNames(value) => value },
+ includeAppNames = queryParams.collectFirst { case OBPIncludeAppNames(value) => value },
+ excludeUrlPatterns = queryParams.collectFirst { case OBPExcludeUrlPatterns(value) => value },
+ includeUrlPatterns = queryParams.collectFirst { case OBPIncludeUrlPatterns(value) => value },
+ excludeImplementedByPartialFunctions = queryParams.collectFirst { case OBPExcludeImplementedByPartialFunctions(value) => value },
+ includeImplementedByPartialFunctions = queryParams.collectFirst { case OBPIncludeImplementedByPartialFunctions(value) => value }
+ )
// override def getAllGroupedByUserId(): Map[String, List[APIMetric]] = {
// //TODO: do this all at the db level using an actual group by query
@@ -379,42 +365,6 @@ object MappedMetrics extends APIMetrics with MdcLoggable{
}
}
- private def extendLikeQuery(params: List[String], isLike: Boolean): String = {
- val isLikeQuery = if (isLike) s"" else s"NOT"
-
- if (params.length == 1)
- s"'${params.head}'"
- else
- {
- val sqlList: immutable.Seq[String] = for (i <- 1 to (params.length - 2)) yield
- {
- s" and url ${isLikeQuery} LIKE ('${params(i)}')"
- }
-
- val sqlSingleLine = if (sqlList.length>1)
- sqlList.reduce(_+_)
- else
- s""
-
- s"'${params.head}')"+ sqlSingleLine + s" and url ${isLikeQuery} LIKE ('${params.last}'"
- }
- }
-
-
- /**
- * Example of a Tuple response
- * (List(count, avg, min, max),List(List(7503, 70.3398640543782487, 0, 9039)))
- * First value of the Tuple is a List of field names returned by SQL query.
- * Second value of the Tuple is a List of rows of the result returned by SQL query. Please note it's only one row.
- */
-
- private def extendPrepareStement(startLine: Int, stmt:PreparedStatement, excludeFiledValues : Set[String]) = {
- for(i <- 0 until excludeFiledValues.size) yield {
- stmt.setString(startLine+i, excludeFiledValues.toList(i))
- }
- }
-
-
// Smart caching applied - uses determineMetricsCacheTTL based on query date range
def getAllAggregateMetricsBox(queryParams: List[OBPQueryParam], isNewVersion: Boolean): Box[List[AggregateMetrics]] = {
logger.info(s"getAllAggregateMetricsBox called with ${queryParams.length} query params, isNewVersion=$isNewVersion")
@@ -430,122 +380,15 @@ object MappedMetrics extends APIMetrics with MdcLoggable{
CacheKeyFromArguments.buildCacheKey { Caching.memoizeSyncWithProvider(Some(cacheKey.toString()))(cacheTTL.seconds){
logger.info(s"getAllAggregateMetricsBox - CACHE MISS - Executing database query for aggregate metrics")
val startTime = System.currentTimeMillis()
+ // The query is built by DoobieMetricsQueries with every filter value as a bound parameter.
+ val filters = metricsQueryFilters(queryParams)
val fromDate = queryParams.collect { case OBPFromDate(value) => value }.headOption
val toDate = queryParams.collect { case OBPToDate(value) => value }.headOption
- val consumerId = queryParams.collect { case OBPConsumerId(value) => value }.headOption
- val userId = queryParams.collect { case OBPUserId(value) => value }.headOption
- val url = queryParams.collect { case OBPUrl(value) => value }.headOption
- val appName = queryParams.collect { case OBPAppName(value) => value }.headOption
- val excludeAppNames = queryParams.collect { case OBPExcludeAppNames(value) => value }.headOption
- val includeAppNames = queryParams.collect { case OBPIncludeAppNames(value) => value }.headOption
- val implementedByPartialFunction = queryParams.collect { case OBPImplementedByPartialFunction(value) => value }.headOption
- val implementedInVersion = queryParams.collect { case OBPImplementedInVersion(value) => value }.headOption
- val verb = queryParams.collect { case OBPVerb(value) => value }.headOption
- val anon = queryParams.collect { case OBPAnon(value) => value }.headOption
- val correlationId = queryParams.collect { case OBPCorrelationId(value) => value }.headOption
- val duration = queryParams.collect { case OBPDuration(value) => value }.headOption
- val httpStatusCode = queryParams.collect { case OBPHttpStatusCode(value) => value }.headOption
- val excludeUrlPatterns = queryParams.collect { case OBPExcludeUrlPatterns(value) => value }.headOption
- val includeUrlPatterns = queryParams.collect { case OBPIncludeUrlPatterns(value) => value }.headOption
- val excludeImplementedByPartialFunctions = queryParams.collect { case OBPExcludeImplementedByPartialFunctions(value) => value }.headOption
- val includeImplementedByPartialFunctions = queryParams.collect { case OBPIncludeImplementedByPartialFunctions(value) => value }.headOption
-
- val excludeUrlPatternsList= excludeUrlPatterns.getOrElse(List(""))
- val excludeAppNamesList = excludeAppNames.getOrElse(List("")).map(i => s"'$i'").mkString(",")
- val excludeImplementedByPartialFunctionsList =
- excludeImplementedByPartialFunctions.getOrElse(List("")).map(i => s"'$i'").mkString(",")
-
- val excludeUrlPatternsQueries = extendLikeQuery(excludeUrlPatternsList, false)
-
- val includeUrlPatternsList= includeUrlPatterns.getOrElse(List(""))
- val includeAppNamesList = includeAppNames.getOrElse(List("")).map(i => s"'$i'").mkString(",")
- val includeImplementedByPartialFunctionsList =
- includeImplementedByPartialFunctions.getOrElse(List("")).map(i => s"'$i'").mkString(",")
-
- val includeUrlPatternsQueries = extendLikeQuery(includeUrlPatternsList, true)
- val includeUrlPatternsQueriesSql = s"$includeUrlPatternsQueries"
-
- // The LEFT JOIN attributes consent-borne calls to the granting (on-behalf-of) human:
- // metric.userid records the AUTHENTICATED principal, which under a consent is the
- // consent's own shadow user. COALESCE(consent.muserid, metric.userid) resolves such rows
- // to the granting human at read time, mirroring CallContext.onBehalfOfUserId. (Rows
- // written 2026-08 only, while toLight briefly recorded the human, resolve identically.)
- // The consent side of the join is unique-indexed on consent_reference_id, so the join
- // cannot fan out rows.
- val result = {
- val sqlQuery = if(isNewVersion) // in the version, we use includeXxx instead of excludeXxx, the performance should be better.
- s"""SELECT count(*), avg(duration), min(duration), max(duration),
- count(DISTINCT CASE WHEN COALESCE(c.muserid, m.userid) <> 'null' THEN COALESCE(c.muserid, m.userid) END),
- count(DISTINCT CASE WHEN m.consumerid <> '' AND m.consumerid <> 'null' THEN m.consumerid END),
- count(NULLIF(m.consent_reference_id, '')),
- count(DISTINCT NULLIF(m.consent_reference_id, ''))
- FROM metric m
- LEFT JOIN mappedconsent c ON m.consent_reference_id = c.consent_reference_id
- WHERE date_c >= '${sqlTimestamp(fromDate.get)}'
- AND date_c <= '${sqlTimestamp(toDate.get)}'
- AND (${trueOrFalse(consumerId.isEmpty)} or consumerid = ${sqlFriendly(consumerId)})
- AND (${trueOrFalse(userId.isEmpty)} or userid = ${sqlFriendly(userId)})
- AND (${trueOrFalse(implementedByPartialFunction.isEmpty)} or implementedbypartialfunction = ${sqlFriendly(implementedByPartialFunction)})
- AND (${trueOrFalse(implementedInVersion.isEmpty)} or implementedinversion = ${sqlFriendly(implementedInVersion)})
- AND (${trueOrFalse(url.isEmpty)} or url = ${sqlFriendly(url)})
- AND (${trueOrFalse(appName.isEmpty)} or appname = ${sqlFriendly(appName)})
- AND (${trueOrFalse(verb.isEmpty)} or verb = ${sqlFriendly(verb)})
- AND (${falseOrTrue(anon.isDefined && anon.equals(Some(true)))} or userid = 'null')
- AND (${falseOrTrue(anon.isDefined && anon.equals(Some(false)))} or userid != 'null')
- AND (${trueOrFalse(correlationId.isEmpty)} or correlationId = ${sqlFriendly(correlationId)})
- AND (${trueOrFalse(httpStatusCode.isEmpty)} or httpcode = ${sqlFriendlyInt(httpStatusCode)})
- AND (${trueOrFalse(includeUrlPatterns.isEmpty) } or (url LIKE ($includeUrlPatternsQueriesSql)))
- AND (${trueOrFalse(includeAppNames.isEmpty) } or (appname in ($includeAppNamesList)))
- AND (${trueOrFalse(includeImplementedByPartialFunctions.isEmpty) } or implementedbypartialfunction in ($includeImplementedByPartialFunctionsList))
- """.stripMargin
- else
- s"""SELECT count(*), avg(duration), min(duration), max(duration),
- count(DISTINCT CASE WHEN COALESCE(c.muserid, m.userid) <> 'null' THEN COALESCE(c.muserid, m.userid) END),
- count(DISTINCT CASE WHEN m.consumerid <> '' AND m.consumerid <> 'null' THEN m.consumerid END),
- count(NULLIF(m.consent_reference_id, '')),
- count(DISTINCT NULLIF(m.consent_reference_id, ''))
- FROM metric m
- LEFT JOIN mappedconsent c ON m.consent_reference_id = c.consent_reference_id
- WHERE date_c >= '${sqlTimestamp(fromDate.get)}'
- AND date_c <= '${sqlTimestamp(toDate.get)}'
- AND (${trueOrFalse(consumerId.isEmpty)} or consumerid = ${sqlFriendly(consumerId)})
- AND (${trueOrFalse(userId.isEmpty)} or userid = ${sqlFriendly(userId)})
- AND (${trueOrFalse(implementedByPartialFunction.isEmpty)} or implementedbypartialfunction = ${sqlFriendly(implementedByPartialFunction)})
- AND (${trueOrFalse(implementedInVersion.isEmpty)} or implementedinversion = ${sqlFriendly(implementedInVersion)})
- AND (${trueOrFalse(url.isEmpty)} or url = ${sqlFriendly(url)})
- AND (${trueOrFalse(appName.isEmpty)} or appname = ${sqlFriendly(appName)})
- AND (${trueOrFalse(verb.isEmpty)} or verb = ${sqlFriendly(verb)})
- AND (${falseOrTrue(anon.isDefined && anon.equals(Some(true)))} or userid = 'null')
- AND (${falseOrTrue(anon.isDefined && anon.equals(Some(false)))} or userid != 'null')
- AND (${trueOrFalse(correlationId.isEmpty)} or correlationId = ${sqlFriendly(correlationId)})
- AND (${trueOrFalse(httpStatusCode.isEmpty)} or httpcode = ${sqlFriendlyInt(httpStatusCode)})
- AND (${trueOrFalse(excludeUrlPatterns.isEmpty) } or (url NOT LIKE ($excludeUrlPatternsQueries)))
- AND (${trueOrFalse(excludeAppNames.isEmpty) } or appname not in ($excludeAppNamesList))
- AND (${trueOrFalse(excludeImplementedByPartialFunctions.isEmpty) } or implementedbypartialfunction not in ($excludeImplementedByPartialFunctionsList))
- """.stripMargin
- // Use DBUtil.runQuery which handles SQL Server NVARCHAR properly
- val (_, rows) = DBUtil.runQuery(sqlQuery)
- logger.debug("code.metrics.MappedMetrics.getAllAggregateMetricsBox.sqlQuery --: " + sqlQuery)
- logger.info(s"getAllAggregateMetricsBox - Query executed, returned ${rows.length} rows")
- val sqlResult = rows.map(
- rs => // Map result to case class
- AggregateMetrics(
- tryo(rs(0).toInt).getOrElse(0),
- tryo("%.2f".format(rs(1).toDouble).toDouble).getOrElse(0),
- tryo(rs(2).toDouble).getOrElse(0),
- tryo(rs(3).toDouble).getOrElse(0),
- tryo(rs(4).toInt).getOrElse(0),
- tryo(rs(5).toInt).getOrElse(0),
- tryo(rs(6).toInt).getOrElse(0),
- tryo(rs(7).toInt).getOrElse(0)
- )
- )
- logger.debug("code.metrics.MappedMetrics.getAllAggregateMetricsBox.sqlResult --: " + sqlResult)
- sqlResult
- }
+ val result = tryo(DoobieMetricsQueries.getAggregateMetrics(fromDate.get, toDate.get, filters, isNewVersion))
+ logger.debug("code.metrics.MappedMetrics.getAllAggregateMetricsBox.sqlResult --: " + result)
val elapsedTime = System.currentTimeMillis() - startTime
logger.info(s"getAllAggregateMetricsBox - Query completed in ${elapsedTime}ms")
- tryo(result)
+ result
}}
}
@@ -737,76 +580,12 @@ object MappedMetrics extends APIMetrics with MdcLoggable{
val cacheTTL = determineMetricsCacheTTL(queryParams)
CacheKeyFromArguments.buildCacheKey {Caching.memoizeSyncWithProvider(Some(cacheKey.toString()))(cacheTTL.seconds){
+ // Built by DoobieMetricsQueries with every filter value as a bound parameter; see
+ // getAllAggregateMetricsBox.
val fromDate = queryParams.collect { case OBPFromDate(value) => value }.headOption
val toDate = queryParams.collect { case OBPToDate(value) => value }.headOption
- val consumerId = queryParams.collect { case OBPConsumerId(value) => value }.headOption
- val userId = queryParams.collect { case OBPUserId(value) => value }.headOption
- val url = queryParams.collect { case OBPUrl(value) => value }.headOption
- val appName = queryParams.collect { case OBPAppName(value) => value }.headOption
- val excludeAppNames = queryParams.collect { case OBPExcludeAppNames(value) => value }.headOption
- val implementedByPartialFunction = queryParams.collect { case OBPImplementedByPartialFunction(value) => value }.headOption
- val implementedInVersion = queryParams.collect { case OBPImplementedInVersion(value) => value }.headOption
- val verb = queryParams.collect { case OBPVerb(value) => value }.headOption
- val anon = queryParams.collect { case OBPAnon(value) => value }.headOption
- val correlationId = queryParams.collect { case OBPCorrelationId(value) => value }.headOption
- val duration = queryParams.collect { case OBPDuration(value) => value }.headOption
- val httpStatusCode = queryParams.collect { case OBPHttpStatusCode(value) => value }.headOption
- val excludeUrlPatterns = queryParams.collect { case OBPExcludeUrlPatterns(value) => value }.headOption
- val excludeImplementedByPartialFunctions = queryParams.collect { case OBPExcludeImplementedByPartialFunctions(value) => value }.headOption
- val limit = queryParams.collect { case OBPLimit(value) => value }.headOption.getOrElse("500")
-
- val excludeUrlPatternsList = excludeUrlPatterns.getOrElse(List(""))
- val excludeAppNamesList = excludeAppNames.getOrElse(List("")).map(i => s"'$i'").mkString(",")
- val excludeImplementedByPartialFunctionsList =
- excludeImplementedByPartialFunctions.getOrElse(List("")).map(i => s"'$i'").mkString(",")
-
- val excludeUrlPatternsQueries: String = extendLikeQuery(excludeUrlPatternsList, false)
-
- val (dbUrl, _, _) = DBUtil.getDbConnectionParameters
-
- // MS SQL server has the specific syntax for limiting number of rows
- val msSqlLimit = if (dbUrl.contains("sqlserver")) s"TOP ($limit)" else s""
- // TODO Make it work in case of Oracle database
- val otherDbLimit: String = if (dbUrl.contains("sqlserver")) s"" else s"LIMIT $limit"
- val result: List[TopConsumer] = {
- val sqlQuery =
- s"""SELECT ${msSqlLimit} count(*) as count, consumer.id as consumerprimaryid, metric.appname as appname,
- consumer.developeremail as email, consumer.consumerid as consumerid
- FROM metric, consumer
- WHERE metric.appname = consumer.name
- AND date_c >= '${sqlTimestamp(fromDate.get)}'
- AND date_c <= '${sqlTimestamp(toDate.get)}'
- AND (${trueOrFalse(consumerId.isEmpty)} or consumer.consumerid = ${sqlFriendly(consumerId)})
- AND (${trueOrFalse(userId.isEmpty)} or userid = ${sqlFriendly(userId)})
- AND (${trueOrFalse(implementedByPartialFunction.isEmpty)} or implementedbypartialfunction = ${sqlFriendly(implementedByPartialFunction)})
- AND (${trueOrFalse(implementedInVersion.isEmpty)} or implementedinversion = ${sqlFriendly(implementedInVersion)})
- AND (${trueOrFalse(url.isEmpty)} or url = ${sqlFriendly(url)})
- AND (${trueOrFalse(appName.isEmpty)} or appname = ${sqlFriendly(appName)})
- AND (${trueOrFalse(verb.isEmpty)} or verb = ${sqlFriendly(verb)})
- AND (${falseOrTrue(anon.isDefined && anon.equals(Some(true)))} or userid = null)
- AND (${falseOrTrue(anon.isDefined && anon.equals(Some(false)))} or userid != null)
- AND (${trueOrFalse(httpStatusCode.isEmpty)} or httpcode = ${sqlFriendlyInt(httpStatusCode)})
- AND (${trueOrFalse(excludeUrlPatterns.isEmpty) } or (url NOT LIKE ($excludeUrlPatternsQueries)))
- AND (${trueOrFalse(excludeAppNames.isEmpty) } or appname not in ($excludeAppNamesList))
- AND (${trueOrFalse(excludeImplementedByPartialFunctions.isEmpty) } or implementedbypartialfunction not in ($excludeImplementedByPartialFunctionsList))
- GROUP BY appname, consumer.developeremail, consumer.id, consumer.consumerid
- ORDER BY count DESC
- ${otherDbLimit}
- """.stripMargin
- // Use DBUtil.runQuery which handles SQL Server NVARCHAR properly
- val (_, rows) = DBUtil.runQuery(sqlQuery)
- val sqlResult =
- rows.map { rs => // Map result to case class
- TopConsumer(
- rs(0).toInt,
- rs(4),
- rs(2),
- rs(3)
- )
- }
- sqlResult
- }
- tryo(result)
+ val limit = queryParams.collect { case OBPLimit(value) => value }.headOption.getOrElse(500)
+ tryo(DoobieMetricsQueries.getTopConsumersByAppName(fromDate.get, toDate.get, limit, metricsQueryFilters(queryParams)))
}
}}
diff --git a/obp-api/src/main/scala/code/metrics/MetricsProps.scala b/obp-api/src/main/scala/code/metrics/MetricsProps.scala
index 6047dee126..7576816de8 100644
--- a/obp-api/src/main/scala/code/metrics/MetricsProps.scala
+++ b/obp-api/src/main/scala/code/metrics/MetricsProps.scala
@@ -75,6 +75,12 @@ object MetricsProps {
*/
val RetainMetricsMoveLimitDefault = 10000
+ /**
+ * The longest date range, in days, that one aggregate-metrics call may cover. An aggregate has
+ * to read every metric row in its range, so the cost grows with the range.
+ */
+ val AggregateMetricsMaxDaysDefault = 31
+
def writeMetrics: Boolean =
APIUtil.getPropsAsBoolValue("write_metrics", WriteMetricsDefault)
@@ -92,4 +98,7 @@ object MetricsProps {
def retainMetricsMoveLimit: Int =
APIUtil.getPropsAsIntValue("retain_metrics_move_limit", RetainMetricsMoveLimitDefault)
+
+ def aggregateMetricsMaxDays: Int =
+ APIUtil.getPropsAsIntValue("aggregate_metrics_max_days", AggregateMetricsMaxDaysDefault)
}
diff --git a/obp-api/src/test/scala/code/api/util/http4s/RequestScopeConnectionTest.scala b/obp-api/src/test/scala/code/api/util/http4s/RequestScopeConnectionTest.scala
index 3dd0de3cb7..cf375a0f91 100644
--- a/obp-api/src/test/scala/code/api/util/http4s/RequestScopeConnectionTest.scala
+++ b/obp-api/src/test/scala/code/api/util/http4s/RequestScopeConnectionTest.scala
@@ -43,6 +43,7 @@ import scala.concurrent.{ExecutionContext, Future}
* - RequestScopeConnection.makeProxy — lifecycle methods are no-ops
* - RequestAwareConnectionManager — proxy vs. delegate selection
* - RequestScopeConnection.fromFuture — TTL propagation to Future workers
+ * - RequestScopeConnection.afterCommit — actions held until the transaction succeeds
*
* All tests use JDK dynamic proxy to build trackable mock Connections; no
* mocking framework is needed. The `after` block resets the global TTL so
@@ -60,6 +61,8 @@ class RequestScopeConnectionTest extends FeatureSpec with Matchers with GivenWhe
RequestScopeConnection.currentProxy.set(null)
RequestScopeConnection.requestProxyLocal.set(None).unsafeRunSync()
RequestScopeConnection.requestLazyAcquire.set(None).unsafeRunSync()
+ RequestScopeConnection.currentAfterCommit.set(null)
+ RequestScopeConnection.requestAfterCommitLocal.set(None).unsafeRunSync()
}
// ─── helpers ─────────────────────────────────────────────────────────────────
@@ -329,4 +332,76 @@ class RequestScopeConnectionTest extends FeatureSpec with Matchers with GivenWhe
result shouldBe 42
}
}
+
+ // ─── RequestScopeConnection.afterCommit ──────────────────────────────────────
+
+ // These routes make no database call, so withBusinessDBTransaction borrows no connection and
+ // the scenarios need no database. The ordering they check is the same as when a commit runs.
+ feature("RequestScopeConnection.afterCommit — actions wait for the transaction") {
+
+ scenario("An action registered inside a transaction runs only after the route has succeeded") {
+ Given("A route that registers an action from a Future callback")
+ val ranAfterCommit = new java.util.concurrent.atomic.AtomicBoolean(false)
+ var ranBeforeRouteFinished = true
+ val route: IO[org.http4s.Response[IO]] = for {
+ _ <- RequestScopeConnection.fromFuture {
+ Future(()).map(_ => RequestScopeConnection.afterCommit(ranAfterCommit.set(true)))
+ }
+ _ <- IO { ranBeforeRouteFinished = ranAfterCommit.get() }
+ } yield org.http4s.Response[IO](org.http4s.Status.Ok)
+
+ When("The route runs inside withBusinessDBTransaction")
+ RequestScopeConnection.withBusinessDBTransaction(route).unsafeRunSync()
+
+ Then("The action had not run while the route was still running")
+ ranBeforeRouteFinished shouldBe false
+ And("It ran once the transaction scope finished successfully")
+ ranAfterCommit.get() shouldBe true
+ }
+
+ scenario("An action registered inside a transaction is dropped when the route fails") {
+ Given("A route that registers an action and then fails")
+ val ran = new java.util.concurrent.atomic.AtomicBoolean(false)
+ val route: IO[org.http4s.Response[IO]] = for {
+ _ <- RequestScopeConnection.fromFuture(Future(RequestScopeConnection.afterCommit(ran.set(true))))
+ response <- IO.raiseError[org.http4s.Response[IO]](new RuntimeException("route failed"))
+ } yield response
+
+ When("The route runs inside withBusinessDBTransaction")
+ val outcome = RequestScopeConnection.withBusinessDBTransaction(route).attempt.unsafeRunSync()
+
+ Then("The failure reaches the caller and the action never runs")
+ outcome.isLeft shouldBe true
+ ran.get() shouldBe false
+ }
+
+ scenario("An action registered outside any transaction runs at once") {
+ Given("No transaction scope")
+ val ran = new java.util.concurrent.atomic.AtomicBoolean(false)
+
+ When("An action is registered from a Future")
+ RequestScopeConnection.fromFuture(Future(RequestScopeConnection.afterCommit(ran.set(true)))).unsafeRunSync()
+
+ Then("It has already run")
+ ran.get() shouldBe true
+ }
+
+ scenario("A failing action does not stop the others or fail the request") {
+ Given("A route that registers a failing action followed by a good one")
+ val goodRan = new java.util.concurrent.atomic.AtomicBoolean(false)
+ val route: IO[org.http4s.Response[IO]] = for {
+ _ <- RequestScopeConnection.fromFuture(Future {
+ RequestScopeConnection.afterCommit(throw new RuntimeException("action failed"))
+ RequestScopeConnection.afterCommit(goodRan.set(true))
+ })
+ } yield org.http4s.Response[IO](org.http4s.Status.Ok)
+
+ When("The route runs inside withBusinessDBTransaction")
+ val response = RequestScopeConnection.withBusinessDBTransaction(route).unsafeRunSync()
+
+ Then("The response is the route's and the good action ran")
+ response.status shouldBe org.http4s.Status.Ok
+ goodRan.get() shouldBe true
+ }
+ }
}
diff --git a/obp-api/src/test/scala/code/api/v3_1_0/MetricsTopConsumersTest.scala b/obp-api/src/test/scala/code/api/v3_1_0/MetricsTopConsumersTest.scala
new file mode 100644
index 0000000000..7929e81d12
--- /dev/null
+++ b/obp-api/src/test/scala/code/api/v3_1_0/MetricsTopConsumersTest.scala
@@ -0,0 +1,83 @@
+/**
+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.v3_1_0
+
+import code.api.util.APIUtil
+import code.api.util.APIUtil.OAuth._
+import code.api.util.ApiRole.CanReadMetrics
+import code.entitlement.Entitlement
+import code.metrics.MetricBatchWriter
+import com.openbankproject.commons.util.ApiVersion
+import org.scalatest.Tag
+
+/**
+ * This tests v3.1.0 GET /management/metrics/top-consumers, which matches metric rows to their
+ * consumer by name.
+ */
+class MetricsTopConsumersTest extends V310ServerSetup {
+
+ object VersionOfApi extends Tag(ApiVersion.v3_1_0.toString)
+ object ApiEndpoint1 extends Tag("getMetricsTopConsumers")
+
+ private def metricsDate(millisecondsAgo: Long): String =
+ APIUtil.DateWithMsFormat.format(new java.util.Date(System.currentTimeMillis() - millisecondsAgo))
+
+ private val oneDayInMillis = 24L * 60 * 60 * 1000
+
+ feature(s"test $ApiEndpoint1 version $VersionOfApi - Filters") {
+ scenario("Traffic is counted for its consumer, and filters match their values exactly", ApiEndpoint1, VersionOfApi) {
+ setPropsValues("write_metrics" -> "true")
+ Entitlement.entitlement.vend.addEntitlement("", resourceUser1.userId, CanReadMetrics.toString)
+ // The traffic is on v5.1.0 /banks, whose metric rows record that url, as in the other
+ // metrics tests; only the top-consumers call itself is v3.1.0.
+ (1 to 3).foreach(_ => makeGetRequest((baseRequest / "obp" / "v5.1.0" / "banks").GET <@ (user1)))
+ MetricBatchWriter.flush()
+ val trafficUrl = "/obp/v5.1.0/banks"
+
+ When("We ask for the top consumers of that url by signed-in users")
+ val request = (v3_1_0_Request / "management" / "metrics" / "top-consumers").GET <@ (user1) < List(
+ ("from_date", metricsDate(oneDayInMillis)),
+ ("url", trafficUrl),
+ ("anon", "false"))
+ val response = makeGetRequest(request)
+ Then("user1's consumer is listed with its calls")
+ response.code should equal(200)
+ val topConsumers = response.body.extract[TopConsumersJson].top_consumers
+ topConsumers.map(_.app_name) should contain(testConsumer.name.get)
+ topConsumers.find(_.app_name == testConsumer.name.get).map(_.count).getOrElse(0) should be >= 3
+
+ When("We filter on a url containing a quote, which no metric has")
+ val injected = (v3_1_0_Request / "management" / "metrics" / "top-consumers").GET <@ (user1) < List(
+ ("from_date", metricsDate(oneDayInMillis)),
+ ("url", s"$trafficUrl' OR '1'='1"))
+ val injectedResponse = makeGetRequest(injected)
+ Then("No metric has that url, so no consumer is listed")
+ injectedResponse.code should equal(200)
+ injectedResponse.body.extract[TopConsumersJson].top_consumers shouldBe empty
+ }
+ }
+}
diff --git a/obp-api/src/test/scala/code/api/v6_0_0/AggregateMetricsTest.scala b/obp-api/src/test/scala/code/api/v6_0_0/AggregateMetricsTest.scala
index aa04c2f83e..187c0a0c9e 100644
--- a/obp-api/src/test/scala/code/api/v6_0_0/AggregateMetricsTest.scala
+++ b/obp-api/src/test/scala/code/api/v6_0_0/AggregateMetricsTest.scala
@@ -29,9 +29,10 @@ package code.api.v6_0_0
import code.api.util.APIUtil.OAuth._
import code.api.util.ApiRole.CanReadAggregateMetrics
-import code.api.util.ErrorMessages.{ApplicationNotIdentified, UserHasMissingRoles}
+import code.api.util.APIUtil
+import code.api.util.ErrorMessages.{AggregateMetricsDateRangeTooLong, ApplicationNotIdentified, UserHasMissingRoles}
import code.entitlement.Entitlement
-import code.metrics.MetricBatchWriter
+import code.metrics.{MetricBatchWriter, MetricsProps}
import code.scope.Scope
import com.openbankproject.commons.model.ErrorMessage
import com.openbankproject.commons.util.ApiVersion
@@ -140,4 +141,55 @@ class AggregateMetricsTest extends V600ServerSetup {
} finally Scope.scope.vend.deleteScope(granted)
}
}
+
+ private def metricsDate(millisecondsAgo: Long): String =
+ APIUtil.DateWithMsFormat.format(new java.util.Date(System.currentTimeMillis() - millisecondsAgo))
+
+ private val oneDayInMillis = 24L * 60 * 60 * 1000
+
+ feature(s"test $ApiEndpoint1 version $VersionOfApi - Date range limit") {
+ scenario("A range longer than the limit is refused, and one within it is answered", ApiEndpoint1, VersionOfApi) {
+ Entitlement.entitlement.vend.addEntitlement("", resourceUser1.userId, CanReadAggregateMetrics.toString)
+ val maxDays = MetricsProps.aggregateMetricsMaxDays
+
+ When(s"We ask for ${maxDays + 1} days")
+ val tooLong = (v6_0_0_Request / "management" / "aggregate-metrics").GET <@ (user1) < List(
+ ("from_date", metricsDate((maxDays + 1) * oneDayInMillis)))
+ val tooLongResponse = makeGetRequest(tooLong)
+ Then("We get a 400 naming the limit")
+ tooLongResponse.code should equal(400)
+ tooLongResponse.body.extract[ErrorMessage].message should startWith(AggregateMetricsDateRangeTooLong)
+
+ When(s"We ask for ${maxDays - 1} days")
+ val withinLimit = (v6_0_0_Request / "management" / "aggregate-metrics").GET <@ (user1) < List(
+ ("from_date", metricsDate((maxDays - 1) * oneDayInMillis)))
+ Then("We get an answer")
+ makeGetRequest(withinLimit).code should equal(200)
+
+ When("We ask for an old range that is short enough")
+ val oldButShort = (v6_0_0_Request / "management" / "aggregate-metrics").GET <@ (user1) < List(
+ ("from_date", metricsDate(400 * oneDayInMillis)),
+ ("to_date", metricsDate(390 * oneDayInMillis)))
+ Then("We get an answer, because only the length of the range is limited")
+ makeGetRequest(oldButShort).code should equal(200)
+ }
+ }
+
+ feature(s"test $ApiEndpoint1 version $VersionOfApi - Filter values are matched exactly") {
+ scenario("A filter value containing a quote is matched as text", ApiEndpoint1, VersionOfApi) {
+ setPropsValues("write_metrics" -> "true")
+ Entitlement.entitlement.vend.addEntitlement("", resourceUser1.userId, CanReadAggregateMetrics.toString)
+ makeGetRequest((v5_1_0_Request / "banks").GET <@ (user1))
+ MetricBatchWriter.flush()
+
+ When("We filter on a url containing a quote, which no metric has")
+ val request = (v6_0_0_Request / "management" / "aggregate-metrics").GET <@ (user1) < List(
+ ("from_date", metricsDate(oneDayInMillis)),
+ ("url", "/obp/v5.1.0/banks' OR '1'='1"))
+ val response = makeGetRequest(request)
+ Then("No metric has that url, so the count is 0")
+ response.code should equal(200)
+ response.body.extract[List[AggregateMetricJsonV600]].head.count shouldBe 0
+ }
+ }
}
diff --git a/obp-api/src/test/scala/code/api/v7_0_0/GrantsEndpointTest.scala b/obp-api/src/test/scala/code/api/v7_0_0/GrantsEndpointTest.scala
new file mode 100644
index 0000000000..719c88fc9f
--- /dev/null
+++ b/obp-api/src/test/scala/code/api/v7_0_0/GrantsEndpointTest.scala
@@ -0,0 +1,122 @@
+/**
+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.{CanGetAllScopes, CanGetReachableRoles, CanGetTelemetry}
+import code.api.util.ErrorMessages.{ApplicationNotIdentified, UserHasMissingRoles}
+import code.api.v6_0_0.V600ServerSetup
+import code.api.v7_0_0.JSONFactory700.{ReachableRolesJsonV700, ScopeJsonV700, ScopesJsonV700}
+import code.entitlement.Entitlement
+import code.scope.Scope
+import com.openbankproject.commons.model.ErrorMessage
+import com.openbankproject.commons.util.ApiVersion
+import org.scalatest.Tag
+
+/**
+ * This suite checks GET /obp/v7.0.0/scopes and GET /obp/v7.0.0/reachable-roles: the Roles they need,
+ * and that reachable-roles lists each Role name held as an Entitlement or a Scope once, and nothing else.
+ */
+class GrantsEndpointTest extends V600ServerSetup {
+
+ def v7_0_0_Request = baseRequest / "obp" / "v7.0.0"
+
+ object VersionOfApi extends Tag(ApiVersion.v7_0_0.toString)
+ object GetAllScopes extends Tag("getAllScopes")
+ object GetReachableRoles extends Tag("getReachableRoles")
+
+ private def scopesRequest = v7_0_0_Request / "scopes"
+ private def reachableRequest = v7_0_0_Request / "reachable-roles"
+
+ private def withEntitlement[T](role: String)(body: => T): T = {
+ val entitlement = Entitlement.entitlement.vend.addEntitlement("", resourceUser1.userId, role)
+ try body finally Entitlement.entitlement.vend.deleteEntitlement(entitlement)
+ }
+
+ private def withScope[T](role: String)(body: => T): T = {
+ val scope = Scope.scope.vend.addScope("", testConsumer2.id.get.toString, role)
+ try body finally Scope.scope.vend.deleteScope(scope)
+ }
+
+ feature(s"Get all Scopes - GET /obp/v7.0.0/scopes - $VersionOfApi") {
+
+ scenario("Anonymous access fails with 401", GetAllScopes, VersionOfApi) {
+ val response = makeGetRequest(scopesRequest.GET)
+ response.code should equal(401)
+ response.body.extract[ErrorMessage].message should equal(ApplicationNotIdentified)
+ }
+
+ scenario("A logged-in user without CanGetAllScopes gets 403", GetAllScopes, VersionOfApi) {
+ val response = makeGetRequest(scopesRequest.GET <@ (user1))
+ response.code should equal(403)
+ response.body.extract[ErrorMessage].message should equal(UserHasMissingRoles + CanGetAllScopes)
+ }
+
+ scenario("A user with CanGetAllScopes sees every Consumer's Scopes", GetAllScopes, VersionOfApi) {
+ val response = withScope(CanGetTelemetry.toString) {
+ withEntitlement(CanGetAllScopes.toString) { makeGetRequest(scopesRequest.GET <@ (user1)) }
+ }
+ response.code should equal(200)
+ response.body.extract[ScopesJsonV700].scopes should contain(
+ ScopeJsonV700(bank_id = "", role_name = CanGetTelemetry.toString, consumer_id = testConsumer2.id.get.toString))
+ }
+ }
+
+ feature(s"Get Reachable Roles - GET /obp/v7.0.0/reachable-roles - $VersionOfApi") {
+
+ scenario("Anonymous access fails with 401", GetReachableRoles, VersionOfApi) {
+ val response = makeGetRequest(reachableRequest.GET)
+ response.code should equal(401)
+ response.body.extract[ErrorMessage].message should equal(ApplicationNotIdentified)
+ }
+
+ scenario("A logged-in user without CanGetReachableRoles gets 403", GetReachableRoles, VersionOfApi) {
+ val response = makeGetRequest(reachableRequest.GET <@ (user1))
+ response.code should equal(403)
+ response.body.extract[ErrorMessage].message should equal(UserHasMissingRoles + CanGetReachableRoles)
+ }
+
+ scenario("A Consumer holding it as a Scope gets the Role names held as Entitlements and Scopes, each once", GetReachableRoles, VersionOfApi) {
+ Given("user1 holds CanGetTelemetry as an Entitlement, and testConsumer2 holds it and CanGetReachableRoles as Scopes")
+ val response = withEntitlement(CanGetTelemetry.toString) {
+ withScope(CanGetTelemetry.toString) {
+ withScope(CanGetReachableRoles.toString) { makeGetRequest(reachableRequest.GET <@ (user2)) }
+ }
+ }
+
+ Then("both Role names are listed, once each, and nothing about who holds them")
+ response.code should equal(200)
+ val roleNames = response.body.extract[ReachableRolesJsonV700].role_names
+ roleNames should contain(CanGetTelemetry.toString)
+ roleNames should contain(CanGetReachableRoles.toString)
+ roleNames.count(_ == CanGetTelemetry.toString) should equal(1)
+ roleNames should equal(roleNames.sorted)
+ response.body.values.asInstanceOf[Map[String, Any]].keySet should equal(Set("role_names"))
+ }
+ }
+}
diff --git a/release_notes.md b/release_notes.md
index 136e0e2114..bab02be71b 100644
--- a/release_notes.md
+++ b/release_notes.md
@@ -3,6 +3,19 @@
### Most recent changes at top of file
```
Date Commit Action
+07/10/2026 TBD CHANGED: GET /management/aggregate-metrics (v3.0.0, v5.1.0 and v6.0.0)
+ covers at most 31 days of metrics per call, set by the new prop
+ aggregate_metrics_max_days (default 31). A range between from_date and
+ to_date longer than that returns 400 with the new error OBP-10069. In
+ v3.0.0 and v5.1.0 a call without from_date used to aggregate the whole
+ metric table; it now covers the 31 days before to_date (or before now).
+ v6.0.0 keeps its own default of the last few minutes. Response bodies are
+ unchanged. The query is also given the endpoint timeout, so the database
+ stops it once the caller has already been answered.
+ FIXED: v3.1.0 top-consumers with anon=true or anon=false returned no rows;
+ it now filters as documented. v5.1.0 and v6.0.0 aggregate-metrics with
+ several include_url_patterns now count a url matching any of them, not
+ only one matching all of them.
04/10/2026 TBD CHANGED: GET /obp/v5.1.0/system/log-cache/LEVEL (trace, debug, info, warning,
error, all) and GET /obp/v7.0.0/management/telemetry accept a Consumer that
holds the Role as a Scope, as well as a User who holds it as an Entitlement
diff --git a/scripts/resource_doc_baseline/parity_allowlist.json b/scripts/resource_doc_baseline/parity_allowlist.json
index e81526abd5..3a1e7f30c8 100644
--- a/scripts/resource_doc_baseline/parity_allowlist.json
+++ b/scripts/resource_doc_baseline/parity_allowlist.json
@@ -1007,9 +1007,49 @@
"version": "v6_0_0",
"endpoint": "getAggregateMetrics",
"field": "description",
- "reason": "Documents real new response fields (distinct_user_count etc.) verified present in MappedMetrics.scala, and who may call it: a User with CanReadAggregateMetrics or an application whose Consumer holds it.",
+ "reason": "Documents real new response fields (distinct_user_count etc.) verified present in MappedMetrics.scala, who may call it (a User with CanReadAggregateMetrics or an application whose Consumer holds it), and the per-call date range limit enforced by APIMetrics.limitAggregateMetricsWindow.",
"lift_digest": "21e12b117560d354f2ee58bb845e9a5d39d9f47721da8195868cf7a6ed9a216d",
- "http4s_digest": "a1631529945acf190378edeff834bbdef9040aee86cb0282008bb7ee6ac539a8"
+ "http4s_digest": "b340b1bf1f8fddc121972bb407c327a196d85429f72f6cf80864595f7a76fcce"
+ },
+ {
+ "version": "v6_0_0",
+ "endpoint": "getAggregateMetrics",
+ "field": "errorResponseBodies",
+ "reason": "AggregateMetricsDateRangeTooLong (400) is returned by APIMetrics.limitAggregateMetricsWindow when from_date and to_date are further apart than aggregate_metrics_max_days.",
+ "lift_digest": "f06cd9a2114625e9b6a37b8e57fa91e7bdb6586b9899fcd9265713aed3e873a4",
+ "http4s_digest": "8f16b7299b9d91cda6127a8a1b18fd8ef7bda1ef93177f22d159fcad017a4428"
+ },
+ {
+ "version": "v3_0_0",
+ "endpoint": "getAggregateMetrics",
+ "field": "description",
+ "reason": "Documents the per-call date range limit enforced by APIMetrics.limitAggregateMetricsWindow, and corrects the from_date default: it was 1 January 1970, not the day before, and is now aggregate_metrics_max_days before to_date.",
+ "lift_digest": "2407b22d7271118a65e34d3400d7af72580aeb8a6cb9521103a991fa464cf3f1",
+ "http4s_digest": "e19a45086f2e479d8cf6744f950678ab224c0460ecad0de8847645450ff4c6ef"
+ },
+ {
+ "version": "v3_0_0",
+ "endpoint": "getAggregateMetrics",
+ "field": "errorResponseBodies",
+ "reason": "AggregateMetricsDateRangeTooLong (400) is returned by APIMetrics.limitAggregateMetricsWindow when from_date and to_date are further apart than aggregate_metrics_max_days.",
+ "lift_digest": "f06cd9a2114625e9b6a37b8e57fa91e7bdb6586b9899fcd9265713aed3e873a4",
+ "http4s_digest": "8f16b7299b9d91cda6127a8a1b18fd8ef7bda1ef93177f22d159fcad017a4428"
+ },
+ {
+ "version": "v5_1_0",
+ "endpoint": "getAggregateMetrics",
+ "field": "description",
+ "reason": "Documents the per-call date range limit enforced by APIMetrics.limitAggregateMetricsWindow, and corrects the from_date default: it was 1 January 1970, not the day before, and is now aggregate_metrics_max_days before to_date.",
+ "lift_digest": "dc8534c5f01cfb8aa8a495f490d136e257788f51f06322ed87d08157831bb874",
+ "http4s_digest": "85dfb572d4097570ef0ab17ed201f6204dfff63c6c589c9c1730a3982ba96744"
+ },
+ {
+ "version": "v5_1_0",
+ "endpoint": "getAggregateMetrics",
+ "field": "errorResponseBodies",
+ "reason": "AggregateMetricsDateRangeTooLong (400) is returned by APIMetrics.limitAggregateMetricsWindow when from_date and to_date are further apart than aggregate_metrics_max_days.",
+ "lift_digest": "f06cd9a2114625e9b6a37b8e57fa91e7bdb6586b9899fcd9265713aed3e873a4",
+ "http4s_digest": "8f16b7299b9d91cda6127a8a1b18fd8ef7bda1ef93177f22d159fcad017a4428"
},
{
"version": "v6_0_0",