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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions obp-api/src/main/resources/props/sample.props.template
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion obp-api/src/main/scala/code/api/util/APIUtil.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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) =
Expand Down
8 changes: 8 additions & 0 deletions obp-api/src/main/scala/code/api/util/ApiRole.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion obp-api/src/main/scala/code/api/util/DBUtil.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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

/**
Expand Down
1 change: 1 addition & 0 deletions obp-api/src/main/scala/code/api/util/ErrorMessages.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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:"
Expand Down
84 changes: 65 additions & 19 deletions obp-api/src/main/scala/code/api/util/Glossary.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
}
}
}

Expand All @@ -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
Expand All @@ -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) *>
Expand All @@ -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))
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down
9 changes: 6 additions & 3 deletions obp-api/src/main/scala/code/api/v3_0_0/Http4s300.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand All @@ -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
|
Expand Down Expand Up @@ -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)
Expand Down
Loading
Loading