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) <= 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) < "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) <. + +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",