Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
aabb934
Log cache and Telemetry accept Consumer Scopes (UserOrApplication) so
simonredfern Oct 4, 2026
c7c315d
API Metrics record domain_api_url: the path and query string a caller
simonredfern Oct 4, 2026
5640431
Update DynamicEntityHelper.scala
simonredfern Oct 4, 2026
bf6da8e
spreading the word about Platform Apps
simonredfern Oct 4, 2026
7a4a7e7
Update parity_allowlist.json
simonredfern Oct 4, 2026
fd3245f
fix: let MetricBatchWriter.flush finish one flush before the next starts
simonredfern Oct 4, 2026
8fa3dc8
testfix: enable write_metrics in the DomainApisTest metric scenario
simonredfern Oct 4, 2026
f7a11c1
Aggregate Metrics client / consumer access and related test
simonredfern Oct 5, 2026
26e9f9b
Adding decimal and boolean to attribute types + parity_allowlist for
simonredfern Oct 5, 2026
c78355b
Update parity_allowlist.json
simonredfern Oct 5, 2026
2440780
fix: Open Corridor outbox messages go STICKY after a retry limit
simonredfern Oct 5, 2026
39d15b5
fix: a bare obp_exists[X] / obp_not_exists[X] key (no '=') was dropped
simonredfern Oct 5, 2026
b4ac8c3
feat: asset registry; currency checks read it
simonredfern Oct 5, 2026
9c64de2
fix: keep a bare obp_exists[X] / obp_not_exists[X] key as a join
simonredfern Oct 5, 2026
2ed1614
fix: declare attribution for the asset and files user-id columns
simonredfern Oct 6, 2026
ad3d77f
Merge remote-tracking branch 'upstream/develop' into develop
simonredfern Oct 6, 2026
ac51d3b
refactor: look up request headers only through RequestHeadersUtil,
simonredfern Oct 6, 2026
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
363 changes: 363 additions & 0 deletions ideas/ASSET_REGISTRY.md

Large diffs are not rendered by default.

3 changes: 3 additions & 0 deletions obp-api/src/main/protobuf/metrics_stream.proto
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,9 @@ message MetricEvent {
// The specifics behind certificate_trust: the forwarding proxy's subject for
// "forwarded", the rejection reason for "none". Matches MetricJsonV600.certificate_trust_detail.
string certificate_trust_detail = 22;
// For a call made under a Domain API, the path and query string the caller used, before it was
// rewritten to the OBP URL in `url`; empty for every other call. Matches MetricJsonV600.domain_api_url.
string domain_api_url = 23;
}

// Live tail of API metrics as they are written.
Expand Down
7 changes: 7 additions & 0 deletions obp-api/src/main/resources/props/sample.props.template
Original file line number Diff line number Diff line change
Expand Up @@ -1363,6 +1363,13 @@ featured_apis=elasticSearchWarehouseV300
# Only relevant on an instance that has Open Corridor turned on and settles platform fees.
# open_corridor.platform_bank_id=

# How many delivery attempts an Open Corridor outbox message gets before the relay stops resending
# it and marks it STICKY for operator reconciliation (GET /management/message-outbox, then its
# /retry). Every attempt that leaves the message undelivered counts: a broker that cannot be reached,
# a retryable error reply, and a settlement instruction whose settlement is not yet FINAL. The wait
# between attempts doubles up to 10 minutes, so the default gives up after roughly a day.
# open_corridor.outbox_max_attempts=144


# -- Scopes -----------------------------------------------------
# Scopes can be used to limit the APIs a Consumer can call.
Expand Down
11 changes: 11 additions & 0 deletions obp-api/src/main/scala/bootstrap/liftweb/Boot.scala
Original file line number Diff line number Diff line change
Expand Up @@ -327,6 +327,10 @@ class Boot extends MdcLoggable {
// Toggle off via routing_schemes.seed_defaults_at_boot=false in environments that don't want defaults.
code.routingscheme.RoutingSchemeSeed.runIfEnabled()

// Idempotent seed of the asset registry (ISO currencies, precious metals, accounting units,
// XBT, ADA, ETH) at the decimal places OBP uses today. Currency validation and decimal places read it.
code.asset.AssetSeed.run()

// Report which static Glossary Items the database is currently displacing. A developer editing
// Glossary.scala has no other way to find out that their text is being overridden.
code.api.util.Glossary.logStaticOverrides()
Expand Down Expand Up @@ -1053,6 +1057,11 @@ object ToSchemify extends MdcLoggable {
DynamicData,
DynamicDataAccess,
code.api.dynamic.entity.projection.DynamicEntityIndex,
// Files: written in full because Boot imports java.io.File.
code.files.File,
code.files.FileContent,
code.files.FileAttachment,
code.files.FileAccess,
DynamicEndpoint,
AccountIdMapping,
DirectDebit,
Expand Down Expand Up @@ -1141,6 +1150,8 @@ object ToSchemify extends MdcLoggable {
Organisation,
RoutingScheme,
BankSupportedRoutingScheme,
code.asset.Asset,
code.asset.AssetStatusHistory,
code.glossaryitem.DynamicGlossaryItem,
code.platformapp.PlatformApp,
code.platformapp.PlatformAppRequiredScope,
Expand Down
3 changes: 1 addition & 2 deletions obp-api/src/main/scala/code/api/DirectLoginRoutes.scala
Original file line number Diff line number Diff line change
Expand Up @@ -60,8 +60,7 @@ object DirectLoginRoutes {
* instead of Lift's thread-local `S.request`.
*/
private def parseDirectLoginParams(cc: CallContext): Map[String, String] = {
def find(name: String): Option[String] = cc.requestHeaders
.find(_.name.equalsIgnoreCase(name))
def find(name: String): Option[String] = code.api.util.RequestHeadersUtil.find(cc.requestHeaders, name)
.flatMap(_.values.headOption)
val directLoginHeader = find("DirectLogin")
val authHeader = find("Authorization")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3232,7 +3232,8 @@ object SwaggerDefinitionsJSON {
consent_reference_id = Some(ExampleValue.consentReferenceIdExample.value),
auth_type = Some("Consent"),
certificate_trust = Some("forwarded"),
certificate_trust_detail = Some("cn=nginx-prod-1,ou=edge,o=tesobe gmbh,c=de")
certificate_trust_detail = Some("cn=nginx-prod-1,ou=edge,o=tesobe gmbh,c=de"),
domain_api_url = None
)
lazy val metricsJsonV600 = MetricsJsonV600(
metrics = List(metricJsonV600)
Expand Down
2 changes: 1 addition & 1 deletion obp-api/src/main/scala/code/api/dauth.scala
Original file line number Diff line number Diff line change
Expand Up @@ -134,7 +134,7 @@ object DAuth extends MdcLoggable {

// Check if the request (access token or request token) is valid and return a tuple
def getDAuthToken(requestHeaders: List[HTTPParam]) : Option[List[String]] = {
requestHeaders.find(_.name.equalsIgnoreCase(APIUtil.DAuthHeaderKey)).map(_.values)
code.api.util.RequestHeadersUtil.find(requestHeaders, APIUtil.DAuthHeaderKey).map(_.values)
}

def getOrCreateResourceUser(jwtPayload: String, callContext: Option[CallContext]) : Box[(User, Option[CallContext])] = {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,8 +54,11 @@ import org.json4s.JsonAST.{JObject, JValue}
*/
object DomainApiPaths {

/** What the front door records on a request it rewrote: which Domain API, and the path that was called. */
case class DomainApiCall(domainApiId: String, basePath: String, calledPath: String)
/**
* What the front door records on a request it rewrote: which Domain API, and the URL that was called (path
* and query string, the shape of CallContext.url), which API Metrics record as `domain_api_url`.
*/
case class DomainApiCall(domainApiId: String, basePath: String, calledUrl: String)

val domainApiCallKey: org.typelevel.vault.Key[DomainApiCall] =
org.typelevel.vault.Key.newKey[IO, DomainApiCall].unsafeRunSync()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ object Http4sDomainApi extends MdcLoggable {
case Nil => OptionT.none[IO, Response[IO]]
case _ =>
val marked = req.withAttribute(domainApiCallKey,
DomainApiCall(route.domainApiId, route.basePath, req.uri.path.renderString))
DomainApiCall(route.domainApiId, route.basePath, req.uri.renderString))
Http4sDynamicEndpoint.wrappedRoutesDynamicEndpoint.run(withPath(marked, DomainApiPaths.dynamicResourceDocPath(route.bankId, rest)))
.orElse(Http4sDynamicEntity.wrappedRoutesDynamicEntityV700.run(withPath(marked, DomainApiPaths.dynamicEntityPath(route.bankId, rest))))
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -303,8 +303,60 @@ object DynamicEntityHelper {

def operationToResourceDoc: Map[(DynamicEntityOperation, String), ResourceDoc] = docsIn(implementedInApiVersion)

/**
* This cache keeps the ResourceDocs that Dynamic Entities generate, so that a request does not
* rebuild them.
*
* The problem it solves: every Dynamic Entity call finds its own ResourceDoc by looking its operation
* up in [[operationToResourceDoc]] (Http4sDynamicEntity calls it for each request). Without this
* cache that lookup built the docs from scratch: for every entity in every space, one doc per
* operation (get all, get one, create, update, patch, delete, and the my, public and community
* variants), each with its description and example bodies, on every request, only to pick one.
*
* What it keeps: one entry per API version, because the same entities are documented twice. v4.0.0
* documents the unversioned /obp/dynamic-entity/... URLs, and v7.0.0 documents the
* /obp/v7.0.0/banks/BANK_ID/dynamic-entities/... URLs (see [[v700Doc]]). Each entry holds the built
* docs together with the definitions map they were built from.
*
* How it knows when to rebuild: [[docsIn]] asks for the current [[definitionsMap]] and reuses the
* kept docs only if that is the very same object they were built from (`eq`: the same instance, not
* equal contents). So the docs have no expiry of their own; they follow the definitions map cache.
* definitionsMap returns the same object until its time-to-live runs out
* (dynamicEntity.definitions_map.cache.ttl.seconds) or until [[forgetDefinitions]] runs because a
* definition was created, updated or deleted on this node. Either way a new map object appears, the
* identity check fails, and the docs are rebuilt once. In test mode the time-to-live is 0, so every
* call builds a new map and the docs are never reused: tests always see current docs, and never
* exercise this cache.
*
* The docs are shared, as static ones are. ResourceDoc has mutable fields, such as specifiedUrl.
* When every caller got freshly built docs, writing to one affected nobody else; now every request
* gets the same objects, so a write is seen by every other request. The only writer is the
* resource-docs listing (ResourceDocsAPIMethods), which sets specifiedUrl afresh each time. Two
* listings running at once may write the same doc concurrently, which is harmless because a given doc
* always gets the same value: a v7.0.0 doc its v7.0.0 URL, a v4.0.0 doc its dynamic-entity URL. Code
* that sets another field per request, or sets specifiedUrl to a value that depends on the request,
* would leak between requests and must copy the doc first.
*
* Concurrency: two requests that miss at the same moment both build and both store; the last store
* wins, and both are correct. A request still holding an older map may store docs built from it after
* a newer entry; the next caller then fails the identity check and rebuilds, so the stale docs are not
* served for long.
*/
private val docsCache = new java.util.concurrent.ConcurrentHashMap[ScannedApiVersion, (Map[(String, String), DynamicEntityInfo], Map[(DynamicEntityOperation, String), ResourceDoc])]()

/** Every entity's docs in one API version: v4.0.0 for the unversioned URLs, v7.0.0 for the v7.0.0 ones. */
private def docsIn(apiVersion: ScannedApiVersion): Map[(DynamicEntityOperation, String), ResourceDoc] = {
val definitions = definitionsMap
Option(docsCache.get(apiVersion)) match {
case Some((builtFrom, docs)) if builtFrom eq definitions => docs
case _ =>
val docs = buildDocsIn(apiVersion, definitions)
docsCache.put(apiVersion, (definitions, docs))
docs
}
}

private def buildDocsIn(apiVersion: ScannedApiVersion, definitions: Map[(String, String), DynamicEntityInfo]): Map[(DynamicEntityOperation, String), ResourceDoc] = {
val addPrefix = APIUtil.getPropsAsBoolValue("dynamic_entities_have_prefix", true)

// record exists tag names, to avoid duplicated dynamic tag name.
Expand Down Expand Up @@ -345,7 +397,7 @@ object DynamicEntityHelper {
ApiTag(tagName)
}
val fun: DynamicEntityInfo => mutable.Map[(DynamicEntityOperation, String), ResourceDoc] = createDocs(apiTag, apiVersion)
val docs: Iterable[((DynamicEntityOperation, String), ResourceDoc)] = definitionsMap.values.flatMap(fun)
val docs: Iterable[((DynamicEntityOperation, String), ResourceDoc)] = definitions.values.flatMap(fun)
docs.toMap
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -145,9 +145,12 @@ object QueryParamParser {

private def parseJoins(params: Map[String, List[String]]): Either[QueryError, List[RawJoin]] = {
// NotExistsKey is tried first; `obp_not_exists[...]` never matches `obp_exists[...]` so order is safe either way.
// A bare key (`?obp_exists[CHILD]`, no `=`) arrives with no values; treat it as one empty value
// (no predicate), the same as `?obp_exists[CHILD]=`, rather than silently dropping the join.
def orEmpty(values: List[String]): List[String] = if (values.isEmpty) List("") else values
val perKey: List[Either[QueryError, List[RawJoin]]] = params.toList.collect {
case (NotExistsKey(child), values) => traverse(values)(parseOneJoin(Quantifier.NotExists, child, _))
case (ExistsKey(child), values) => traverse(values)(parseOneJoin(Quantifier.Exists, child, _))
case (NotExistsKey(child), values) => traverse(orEmpty(values))(parseOneJoin(Quantifier.NotExists, child, _))
case (ExistsKey(child), values) => traverse(orEmpty(values))(parseOneJoin(Quantifier.Exists, child, _))
}
sequence(perKey).map(_.flatten)
}
Expand Down
5 changes: 2 additions & 3 deletions obp-api/src/main/scala/code/api/siwe.scala
Original file line number Diff line number Diff line change
Expand Up @@ -127,12 +127,11 @@ object SIWE extends MdcLoggable {
val SiweHeaderKey = "SIWE"

def hasSiweHeader(requestHeaders: List[HTTPParam]): Boolean =
requestHeaders.exists(_.name.equalsIgnoreCase(SiweHeaderKey))
code.api.util.RequestHeadersUtil.exists(requestHeaders, SiweHeaderKey)

/** Parse `SIWE: token=<key>` → Some(key). Mirrors DirectLogin's `token=` parsing. */
def getSiweToken(requestHeaders: List[HTTPParam]): Option[String] = {
val raw = requestHeaders
.find(_.name.equalsIgnoreCase(SiweHeaderKey))
val raw = code.api.util.RequestHeadersUtil.find(requestHeaders, SiweHeaderKey)
.flatMap(_.values.headOption)
.getOrElse("")
raw.split(",").map(_.trim).flatMap { entry =>
Expand Down
Loading
Loading