diff --git a/obp-api/src/main/resources/props/sample.props.template b/obp-api/src/main/resources/props/sample.props.template index a9792a942a..303b96f8a5 100644 --- a/obp-api/src/main/resources/props/sample.props.template +++ b/obp-api/src/main/resources/props/sample.props.template @@ -141,6 +141,12 @@ long_endpoint_timeout = 55000 ## Scheduler will be disabled if delay is not set. #transaction_status_scheduler_delay=300 +## Message outbox relay interval in seconds. +## How often the relay sends the queued messages in the message outbox: emails (e.g. "you have been +## granted a Role") and, where Open Corridor is enabled, Interface C messages to the banks. The relay +## always runs. A message that cannot be delivered is retried with backoff, on top of this interval. +#message_outbox.relay_interval_seconds=10 + ## Enable user authentication via the connector #connector.user.authentication=true @@ -262,11 +268,14 @@ write_connector_trace=false ## per-IP limit, IP penalty and the busiest-callers view would then see one address. When a proxy sits ## in front, let it set a header with the client's address (it MUST overwrite any value the client sent, ## e.g. NGINX `proxy_set_header X-Real-IP $remote_addr;`) and trust that header here. -## X-Forwarded-For is also accepted (its leftmost address is used) when the proxy sanitises the chain. +## X-Forwarded-For is also accepted. It carries a chain to which each hop appends the address it received +## the request from (NGINX `proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;`, and API Explorer II, +## Opey and OBP-MCP do the same). OBP-API reads it from the right: it skips the addresses in trust.proxy.peers +## and the first address that is not listed is the client, so every hop must be listed. ## trust.proxy.peers limits whose header is believed: the addresses or CIDR ranges of the proxies (and ## of server-side applications that pass on their users' addresses). A header from any other peer is -## ignored. Unset, the header is believed from anyone, so a caller that reaches OBP-API directly can -## name any address it likes. Deployment Checks (GET /obp/v7.0.0/management/system/diagnostics/deployment, +## ignored. Unset, the header is believed from anyone (and for X-Forwarded-For the leftmost address, which +## the client itself can write, is used), so a caller that reaches OBP-API directly can name any address it likes. Deployment Checks (GET /obp/v7.0.0/management/system/diagnostics/deployment, ## or Observe > Deployment Checks in API Manager) shows whether these are right for the traffic seen. # trust.proxy.enabled=false # trust.proxy.header=X-Real-IP diff --git a/obp-api/src/main/scala/bootstrap/liftweb/Boot.scala b/obp-api/src/main/scala/bootstrap/liftweb/Boot.scala index 012c447af9..869a7a1b9c 100644 --- a/obp-api/src/main/scala/bootstrap/liftweb/Boot.scala +++ b/obp-api/src/main/scala/bootstrap/liftweb/Boot.scala @@ -623,11 +623,10 @@ class Boot extends MdcLoggable { val delay = APIUtil.getPropsAsLongValue("transaction_request_status_scheduler_delay").openOrThrowException("Incorrect value for transaction_request_status_scheduler_delay, please provide number of seconds.") TransactionRequestStatusScheduler.start(delay) } - // Open Corridor: the transactional-outbox relay publishing Interface C messages - // (credit notifications + settlement instructions) to the banks' own vhosts. - if (APIUtil.getPropsAsBoolValue("open_corridor_enabled", false)) { - MessageOutboxRelay.start(APIUtil.getPropsAsLongValue("open_corridor.outbox_relay_interval", 10L)) - } + // The transactional-outbox relay: sends queued emails (e.g. "you have been granted a Role") and, + // when Open Corridor is enabled, publishes Interface C messages (credit notifications + + // settlement instructions) to the banks' own vhosts. + MessageOutboxRelay.start(APIUtil.getPropsAsLongValue("message_outbox.relay_interval_seconds", 10L)) // Chat: emails users an occasional digest of unread messages (computed at // send time from read markers — see ChatEmailDigestScheduler for why this // is not the transactional message outbox). diff --git a/obp-api/src/main/scala/code/api/ResourceDocs1_4_0/SwaggerDefinitionsJSON.scala b/obp-api/src/main/scala/code/api/ResourceDocs1_4_0/SwaggerDefinitionsJSON.scala index 1c588bf76a..65a1237637 100644 --- a/obp-api/src/main/scala/code/api/ResourceDocs1_4_0/SwaggerDefinitionsJSON.scala +++ b/obp-api/src/main/scala/code/api/ResourceDocs1_4_0/SwaggerDefinitionsJSON.scala @@ -3224,6 +3224,7 @@ object SwaggerDefinitionsJSON { duration = 39, source_ip = ExampleValue.ipAddressExample.value, target_ip = ExampleValue.ipAddressExample.value, + forwarded_for = s"${ExampleValue.ipAddressExample.value}, 10.0.0.2, 10.0.0.3", response_body = json.parse("""{"code":401,"message":"OBP-20001: User not logged in. Authentication is required!"}"""), status_code = 401, operation_id = "OBPv4.0.0-getBanks", diff --git a/obp-api/src/main/scala/code/api/dynamic/entity/Http4sDynamicEntity.scala b/obp-api/src/main/scala/code/api/dynamic/entity/Http4sDynamicEntity.scala index ed7228ecf5..deaa70aabe 100644 --- a/obp-api/src/main/scala/code/api/dynamic/entity/Http4sDynamicEntity.scala +++ b/obp-api/src/main/scala/code/api/dynamic/entity/Http4sDynamicEntity.scala @@ -27,7 +27,7 @@ package code.api.dynamic.entity import cats.data.{Kleisli, OptionT} import cats.effect.IO -import code.DynamicData.{DynamicData, DynamicDataProvider, DynamicDataAccessProvider, DynamicDataAccessPermission} +import code.DynamicData.{DynamicData, DynamicDataProvider, DynamicDataAccessProvider, DynamicDataAccessPermission, DynamicDataT} import code.api.Constant.PARAM_LOCALE import code.api.dynamic.entity.helper.{CommunityEntityName, DynamicEntityHelper, DynamicEntityInfo, DynamicEntitySpace, EntityAccessName, EntityName, PublicEntityName} import code.api.dynamic.entity.query.{FieldSpec, InMemoryQueryExecutor, JoinTargetInfo, QueryParamParser, QueryPlan, QueryPlanner} @@ -176,6 +176,91 @@ object Http4sDynamicEntity extends MdcLoggable { (("bank_id" -> DynamicEntitySpace.bankIdOrSystem(bankId)): JObject) merge result else result + private def namesEverySpace(req: Request[IO]): Boolean = + req.attributes.lookup(namesEverySpaceKey).contains(true) + + /** + * This says whether an entity's records are held in OBP's own records table. They are not when a + * method routing sends dynamicEntityProcess for the entity to another connector, which is the same + * test the row-level access guard applies (Http4s600.localBackingOkForRowLevel). Metadata is only + * shown for a locally held record: for one held elsewhere, a leftover local row with the same id + * would describe a different write. + */ + private def isLocallyBacked(entityName: String): Boolean = + !NewStyle.function.getMethodRoutings(Some("dynamicEntityProcess")) + .exists(_.parameters.exists(parameter => parameter.key == "entityName" && parameter.value == entityName)) + + private def recordIdOf(entityName: String, record: JValue): Option[String] = + record \ DynamicEntityHelper.createEntityId(entityName) match { + case JString(recordId) => Some(recordId) + case _ => None + } + + /** The stored rows behind `records`, by record id, for reading their metadata. Empty for an + * entity whose records are not held locally. */ + private def storedRowsByRecordId(bankId: Option[String], entityName: String, records: List[JValue]): Map[String, DynamicDataT] = { + val recordIds = records.flatMap(recordIdOf(entityName, _)) + if (recordIds.isEmpty || !isLocallyBacked(entityName)) Map.empty + else dataVend.getByIds(bankId, entityName, recordIds).flatMap(row => row.dynamicDataId.map(_ -> row)).toMap + } + + private def utcSeconds(date: java.util.Date): String = + java.time.format.DateTimeFormatter.ISO_INSTANT.format(date.toInstant.truncatedTo(java.time.temporal.ChronoUnit.SECONDS)) + + /** + * This is the metadata block of a record in a v7.0.0 response: when it was created and last + * updated, and for each, the user who made the call and the user it was made for. The two ids + * differ only when an agent acted for somebody. A value that was never recorded, because the + * record was written before the columns existed, is null. `showUserIds` is false on the public + * reads, which anyone may call, so they carry the times only. + */ + private def metadataJson(row: DynamicDataT, showUserIds: Boolean): JObject = { + def orNull(value: Option[String]): JValue = value.map(JString(_)).getOrElse(JNull) + def event(at: Option[java.util.Date], userId: Option[String], onBehalfOfUserId: Option[String]): JObject = { + val atField = JField("at", orNull(at.map(utcSeconds))) + if (showUserIds) JObject(atField :: JField("user_id", orNull(userId)) :: JField("on_behalf_of_user_id", orNull(onBehalfOfUserId)) :: Nil) + else JObject(atField :: Nil) + } + JObject( + JField("created", event(row.createdDate, row.createdByUserId, row.createdByOnBehalfOfUserId)) :: + JField("updated", event(row.updatedDate, row.updatedByUserId, row.updatedByOnBehalfOfUserId)) :: Nil) + } + + /** + * This builds the response for one record. At a v7.0.0 URL the record is followed by its + * metadata, when it is held locally; the unversioned URLs return the record alone, as before. + */ + private def singleResponse(req: Request[IO], bankId: Option[String], entityName: String, record: JValue, showUserIds: Boolean = true): JObject = { + val recordField = JField(singleName(entityName), record) + val metadataField = + if (!namesEverySpace(req)) None + else recordIdOf(entityName, record).flatMap(storedRowsByRecordId(bankId, entityName, List(record)).get) + .map(row => JField("metadata", metadataJson(row, showUserIds))) + wrapBankId(req, bankId, JObject(recordField :: metadataField.toList)) + } + + /** + * This builds the response for a list of records. At a v7.0.0 URL each item has the same shape as + * a single record response without bank_id, the record under the entity's name and its metadata + * beside it, so that no record field can collide with the metadata. The unversioned URLs keep + * their items as plain records. + */ + private def listResponse(req: Request[IO], bankId: Option[String], entityName: String, records: JValue, showUserIds: Boolean = true): JObject = + if (!namesEverySpace(req)) wrapBankId(req, bankId, (listName(entityName) -> records)) + else { + val items = records match { + case JArray(values) => values + case _ => Nil + } + val storedRows = storedRowsByRecordId(bankId, entityName, items) + val wrappedItems = items.map { record => + val metadataField = recordIdOf(entityName, record).flatMap(storedRows.get) + .map(row => JField("metadata", metadataJson(row, showUserIds))) + JObject(JField(singleName(entityName), record) :: metadataField.toList) + } + wrapBankId(req, bankId, (listName(entityName) -> JArray(wrappedItems))) + } + private def notFoundMsg(entityName: String, id: String, bankId: Option[String]): String = s"$EntityNotFoundByEntityId Entity: '$entityName', entityId: '$id'" + bankId.map(b => s", bank_id: '$b'").getOrElse("") @@ -420,7 +505,7 @@ object Http4sDynamicEntity extends MdcLoggable { val readableRows = dataVend.getAllCommunity(bankId, entityName).filter(_.dynamicDataId.exists(readable.contains)) val readableJson: JArray = JArray(readableRows.map(r => parse(r.dataJson))) val filtered = filterDynamicObjects(readableJson, queryParams(req)) - wrapBankId(req, bankId, (listName(entityName) -> applyReadRestrictions(filtered, bankId, entityName, Some(u.userId)))) + listResponse(req, bankId, entityName, applyReadRestrictions(filtered, bankId, entityName, Some(u.userId))) } else { val box: Box[JValue] = dataVend.getCommunity(bankId, entityName, id).map(it => parse(it.dataJson)) for { @@ -430,7 +515,7 @@ object Http4sDynamicEntity extends MdcLoggable { } } yield { val singleObject: JValue = unboxResult(box, entityName) - wrapBankId(req, bankId, (singleName(entityName) -> applyReadRestrictions(singleObject, bankId, entityName, Some(u.userId)))) + singleResponse(req, bankId, entityName, applyReadRestrictions(singleObject, bankId, entityName, Some(u.userId))) } } } yield result @@ -451,9 +536,9 @@ object Http4sDynamicEntity extends MdcLoggable { aclVend.allows(bankId, entityName, id, u.userId, DynamicDataAccessPermission.Update) } // Field-level write roles still apply on top of the row ACL. updateJson = preserveRestrictedOnPut(json.asInstanceOf[JObject], existing, writeRestrictedFieldsOf(bankId, entityName)) - box: Box[JValue] = dataVend.updateCommunity(bankId, entityName, updateJson, id).map(it => parse(it.dataJson)) + box: Box[JValue] = dataVend.updateCommunity(bankId, entityName, updateJson, id, Some(u.userId)).map(it => parse(it.dataJson)) singleObject: JValue = unboxResult(box, entityName) - } yield wrapBankId(req, bankId, (singleName(entityName) -> singleObject)) + } yield singleResponse(req, bankId, entityName, singleObject) } private def rowLevelPatch(req: Request[IO], cc: CallContext, bankId: Option[String], entityName: String, id: String): Future[JValue] = { @@ -474,9 +559,9 @@ object Http4sDynamicEntity extends MdcLoggable { existing: Box[JValue] = dataVend.getCommunity(bankId, entityName, id).map(it => parse(it.dataJson)) _ <- Helper.booleanToFuture(notFoundMsg(entityName, id, bankId), 404, cc = callContext) { existing.isDefined } mergedJson = mergePatch(DynamicEntityHelper.definitionOf(bankId, entityName), existing, bodyObj) - box: Box[JValue] = dataVend.updateCommunity(bankId, entityName, mergedJson, id).map(it => parse(it.dataJson)) + box: Box[JValue] = dataVend.updateCommunity(bankId, entityName, mergedJson, id, Some(u.userId)).map(it => parse(it.dataJson)) singleObject: JValue = unboxResult(box, entityName) - } yield wrapBankId(req, bankId, (singleName(entityName) -> singleObject)) + } yield singleResponse(req, bankId, entityName, singleObject) } private def rowLevelDelete(req: Request[IO], cc: CallContext, bankId: Option[String], entityName: String, id: String): Future[JValue] = { @@ -615,10 +700,10 @@ object Http4sDynamicEntity extends MdcLoggable { val legacyFiltered = filterDynamicObjects(resultList, queryParams(req)) applyQueryPlan(legacyFiltered, queryPlan, deIndexedFields(bankId, entityName)) } - wrapBankId(req, bankId, (listName(entityName) -> applyReadRestrictions(filtered, bankId, entityName, userIdOpt))) + listResponse(req, bankId, entityName, applyReadRestrictions(filtered, bankId, entityName, userIdOpt)) } else { val singleObject: JValue = unboxResult(box.asInstanceOf[Box[JValue]], entityName) - wrapBankId(req, bankId, (singleName(entityName) -> applyReadRestrictions(singleObject, bankId, entityName, userIdOpt))) + singleResponse(req, bankId, entityName, applyReadRestrictions(singleObject, bankId, entityName, userIdOpt)) } } } @@ -650,7 +735,7 @@ object Http4sDynamicEntity extends MdcLoggable { userIdOpt.foreach(uid => aclVend.grant(bankId, entityName, rid, uid, canRead = true, canUpdate = true, canDelete = true, canGrant = true, grantedBy = uid)) case _ => } - } yield wrapBankId(req, bankId, (singleName(entityName) -> singleObject)) + } yield singleResponse(req, bankId, entityName, singleObject) } private def genericPut(req: Request[IO], bankId: Option[String], entityName: String, id: String, isPersonalEntity: Boolean): IO[Response[IO]] = @@ -679,7 +764,7 @@ object Http4sDynamicEntity extends MdcLoggable { updateJson = preserveRestrictedOnPut(json.asInstanceOf[JObject], existing.asInstanceOf[Box[JValue]], writeRestrictedFieldsOf(bankId, entityName)) (box: Box[JValue], _) <- NewStyle.function.invokeDynamicConnector(UPDATE, entityName, Some(updateJson), Some(id), bankId, None, userIdOpt, isPersonalEntity, Some(cc)) singleObject: JValue = unboxResult(box, entityName) - } yield wrapBankId(req, bankId, (singleName(entityName) -> singleObject)) + } yield singleResponse(req, bankId, entityName, singleObject) } private def genericPatch(req: Request[IO], bankId: Option[String], entityName: String, id: String, isPersonalEntity: Boolean): IO[Response[IO]] = @@ -713,7 +798,7 @@ object Http4sDynamicEntity extends MdcLoggable { mergedJson = mergePatch(DynamicEntityHelper.definitionOf(bankId, entityName), existing.asInstanceOf[Box[JValue]], bodyObj) (box: Box[JValue], _) <- NewStyle.function.invokeDynamicConnector(UPDATE, entityName, Some(mergedJson), Some(id), bankId, None, userIdOpt, isPersonalEntity, Some(cc)) singleObject: JValue = unboxResult(box, entityName) - } yield wrapBankId(req, bankId, (singleName(entityName) -> singleObject)) + } yield singleResponse(req, bankId, entityName, singleObject) } private def genericDelete(req: Request[IO], bankId: Option[String], entityName: String, id: String, isPersonalEntity: Boolean): IO[Response[IO]] = @@ -765,10 +850,10 @@ object Http4sDynamicEntity extends MdcLoggable { val resultList: JArray = unboxResult(box.asInstanceOf[Box[JArray]], entityName) val legacyFiltered = filterDynamicObjects(resultList, queryParams(req)) val filtered = applyQueryPlan(legacyFiltered, queryPlan, deIndexedFields(bankId, entityName)) - wrapBankId(req, bankId, (listName(entityName) -> applyReadRestrictions(filtered, bankId, entityName, None))) + listResponse(req, bankId, entityName, applyReadRestrictions(filtered, bankId, entityName, None), showUserIds = false) } else { val singleObject: JValue = unboxResult(box.asInstanceOf[Box[JValue]], entityName) - wrapBankId(req, bankId, (singleName(entityName) -> applyReadRestrictions(singleObject, bankId, entityName, None))) + singleResponse(req, bankId, entityName, applyReadRestrictions(singleObject, bankId, entityName, None), showUserIds = false) } } } @@ -797,14 +882,14 @@ object Http4sDynamicEntity extends MdcLoggable { val resultArray = JArray(resultList) val legacyFiltered = filterDynamicObjects(resultArray, queryParams(req)) val filtered = applyQueryPlan(legacyFiltered, queryPlan, deIndexedFields(bankId, entityName)) - wrapBankId(req, bankId, (listName(entityName) -> applyReadRestrictions(filtered, bankId, entityName, Some(u.userId)))) + listResponse(req, bankId, entityName, applyReadRestrictions(filtered, bankId, entityName, Some(u.userId))) } else { val singleResult = DynamicDataProvider.connectorMethodProvider.vend.getCommunity(bankId, entityName, id) val singleObject: JValue = singleResult match { case Full(data) => com.openbankproject.commons.util.JsonAliases.parse(data.dataJson) case _ => throw new RuntimeException(notFoundMsg(entityName, id, bankId)) } - wrapBankId(req, bankId, (singleName(entityName) -> applyReadRestrictions(singleObject, bankId, entityName, Some(u.userId)))) + singleResponse(req, bankId, entityName, applyReadRestrictions(singleObject, bankId, entityName, Some(u.userId))) } } } diff --git a/obp-api/src/main/scala/code/api/dynamic/entity/helper/DynamicEntityHelper.scala b/obp-api/src/main/scala/code/api/dynamic/entity/helper/DynamicEntityHelper.scala index 6a73865884..71c5b41a5e 100644 --- a/obp-api/src/main/scala/code/api/dynamic/entity/helper/DynamicEntityHelper.scala +++ b/obp-api/src/main/scala/code/api/dynamic/entity/helper/DynamicEntityHelper.scala @@ -240,6 +240,22 @@ object DynamicEntityHelper { def v700UrlPrefix(bankId: Option[String]): String = s"/banks/${DynamicEntitySpace.bankIdOrSystem(bankId)}/dynamic-entities" + /** + * This is the example of a record's metadata block in the v7.0.0 docs: an agent created the record + * for a person, and the person then updated it themselves. Without user ids, as the public reads + * return it, it carries the times only. + */ + def exampleRecordMetadata(showUserIds: Boolean): JObject = { + val personUserId = ExampleValue.userIdExample.value + val agentUserId = "a9f2c7e1-3d4b-4a5c-8e6f-7a8b9c0d1e2f" + def event(at: String, userId: String, onBehalfOfUserId: String): JObject = + if (showUserIds) JObject(JField("at", JString(at)) :: JField("user_id", JString(userId)) :: JField("on_behalf_of_user_id", JString(onBehalfOfUserId)) :: Nil) + else JObject(JField("at", JString(at)) :: Nil) + JObject( + JField("created", event("2026-09-30T10:12:00Z", agentUserId, personUserId)) :: + JField("updated", event("2026-09-30T11:40:00Z", personUserId, personUserId)) :: Nil) + } + def createEntityId(entityName: String) = { // (?<=[a-z0-9])(?=[A-Z]) --> mean `Positive Lookbehind (?<=[a-z0-9])` && Positive Lookahead (?=[A-Z]) --> So we can find the space to replace to `_` val regexPattern = "(?<=[a-z0-9])(?=[A-Z])|-" @@ -325,13 +341,22 @@ object DynamicEntityHelper { else bankId.map(b => s"/banks/$b").getOrElse("") val resourceDocUrl = s"$urlPrefix/$entityName" // Response examples. A v7.0.0 response always names its space, `"bank_id": "SYS"` included, where - // the unversioned URLs name a bank only for a bank level entity. - def inSpace(example: JObject): JObject = - if (apiVersion == ApiVersion.v7_0_0) - (("bank_id" -> DynamicEntitySpace.bankIdOrSystem(bankId)): JObject) merge JObject(example.obj.filterNot(_.name == "bank_id")) - else example - val singleExample = inSpace(dynamicEntityInfo.getSingleExample) - val listExample = inSpace(dynamicEntityInfo.getExampleList) + // the unversioned URLs name a bank only for a bank level entity. A v7.0.0 record is also followed + // by its metadata, and each list item has the shape of a single record response without bank_id + // (see Http4sDynamicEntity.singleResponse and listResponse). The public reads show the times only. + def v700Examples(showUserIds: Boolean): (JObject, JObject) = { + val record = JObject(JField(dynamicEntityInfo.idName, JString(ExampleValue.idExample.value)) :: dynamicEntityInfo.getSingleExampleWithoutId.obj) + val recordWithMetadata = List(JField(dynamicEntityInfo.singleName, record), JField("metadata", exampleRecordMetadata(showUserIds))) + val space = JField("bank_id", JString(DynamicEntitySpace.bankIdOrSystem(bankId))) + (JObject(space :: recordWithMetadata), + JObject(space :: JField(listName, JArray(List(JObject(recordWithMetadata)))) :: Nil)) + } + val (singleExample, listExample) = + if (apiVersion == ApiVersion.v7_0_0) v700Examples(showUserIds = true) + else (dynamicEntityInfo.getSingleExample, dynamicEntityInfo.getExampleList) + val (publicSingleExample, publicListExample) = + if (apiVersion == ApiVersion.v7_0_0) v700Examples(showUserIds = false) + else (singleExample, listExample) val myResourceDocUrl = s"$urlPrefix/my/$entityName" @@ -709,7 +734,7 @@ object DynamicEntityHelper { |${dynamicEntityInfo.listQueryDoc(joinsSupported = false)} |""".stripMargin, EmptyBody, - listExample, + publicListExample, List( UnknownError ), @@ -733,7 +758,7 @@ object DynamicEntityHelper { |Authentication is Optional |""".stripMargin, EmptyBody, - singleExample, + publicSingleExample, List( UnknownError ), diff --git a/obp-api/src/main/scala/code/api/util/ApiSession.scala b/obp-api/src/main/scala/code/api/util/ApiSession.scala index b7a559ded3..656801e885 100644 --- a/obp-api/src/main/scala/code/api/util/ApiSession.scala +++ b/obp-api/src/main/scala/code/api/util/ApiSession.scala @@ -84,6 +84,9 @@ case class CallContext( consentMyResources: Option[ConsentMyResources] = None, consumer: Box[Consumer] = Empty, ipAddress: String = "", + // The hops the request passed through: the X-Forwarded-For chain it arrived + // with, followed by the TCP peer. Recorded in API Metrics; see RemoteIpUtil. + forwardedFor: String = "", resourceDocument: Option[ResourceDoc] = None, startTime: Option[Date] = Some(Helpers.now), endTime: Option[Date] = None, @@ -245,7 +248,9 @@ case class CallContext( paginationLimit = this.paginationLimit, consentReferenceId = this.consentReferenceId, certificateTrust = this.certificateTrust, - certificateTrustDetail = this.certificateTrustDetail + certificateTrustDetail = this.certificateTrustDetail, + ipAddress = this.ipAddress, + forwardedFor = this.forwardedFor ) } @@ -344,7 +349,11 @@ case class CallContextLight(gatewayLoginRequestPayload: Option[PayloadOfJwtJSON] paginationLimit : Option[String] = None, consentReferenceId: Option[String] = None, certificateTrust: Option[String] = None, - certificateTrustDetail: Option[String] = None + certificateTrustDetail: Option[String] = None, + // The client address OBP-API decided on (CallContext.ipAddress) + ipAddress: String = "", + // The hops the request passed through (CallContext.forwardedFor) + forwardedFor: String = "" ) trait LoginParam 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 9eeecdb82a..09c03a6db4 100644 --- a/obp-api/src/main/scala/code/api/util/Glossary.scala +++ b/obp-api/src/main/scala/code/api/util/Glossary.scala @@ -3746,6 +3746,28 @@ object Glossary extends MdcLoggable { | |Every v7.0.0 response carries `bank_id`, `SYS` included. The unversioned `/obp/dynamic-entity/[banks/BANK_ID/]...` URLs serve the same records with the same checks, and keep omitting `bank_id` for the system space. | +|**Record metadata from v7.0.0:** +| +|In a v7.0.0 response each record is followed by a `metadata` object, which says when the record was created and last updated, and for each, which User made the call (`user_id`) and which User it was made for (`on_behalf_of_user_id`). The two ids are the same when a User acted for themselves, and differ when an agent acted for somebody through a Consent. The metadata sits beside the record rather than inside it, so it can never collide with a field of the entity: +| +|``` +|{ +| "bank_id": "SYS", +| "soil_sample": { "soil_sample_id": "8f1c2d9e-...", "ph": 6.7 }, +| "metadata": { +| "created": { "at": "2026-09-30T10:12:00Z", "user_id": "a9f2c7e1-...", "on_behalf_of_user_id": "e1a4b6c8-..." }, +| "updated": { "at": "2026-09-30T11:40:00Z", "user_id": "e1a4b6c8-...", "on_behalf_of_user_id": "e1a4b6c8-..." } +| } +|} +|``` +| +|In a list, each item has that same shape without `bank_id`: `{"soil_sample": {...}, "metadata": {...}}`. Times are in UTC, to the second. +| +|* The public reads, which need no login, carry the times only, and never say who wrote a record. +|* A value that was never recorded is null. A record written before this metadata existed has null in `created` for good, and `updated` is filled the next time the record is saved. +|* A record held by another connector (a method routing for `dynamicEntityProcess` names its entity) has no `metadata`, because OBP does not hold that information. +|* The unversioned `/obp/dynamic-entity/...` URLs return the record alone, and their list items stay plain records. +| |Earlier versions keep their separate management URLs for the system space; their Roles are the same ones, granted at `SYS`: | |* POST /management/system-dynamic-entities - Create system level entity @@ -6825,7 +6847,9 @@ object Glossary extends MdcLoggable { |- the User (user id and username) and the Consumer (consumer id, application name and developer email) |- how the caller authenticated (for example DirectLogin, OAuth2, Consent or Anonymous) and, for a call made under a Consent, the Consent reference id |- the correlation id, which is also returned to the caller in the `Correlation-Id` response header and is shared by every Connector call made while serving the call (see [Connector Metrics](/glossary#Connector-Metrics)) - |- the source and target addresses, taken from the `X-Forwarded-For` and `X-Forwarded-Host` request headers + |- `source_ip`: the address of the client, as OBP-API decided it (see [Client IP Address](/glossary#Client-IP-Address)). Records written before this was introduced hold the raw `X-Forwarded-For` header instead + |- `forwarded_for`: the hops the call passed through: the `X-Forwarded-For` list it arrived with, followed by the address of the machine that connected to OBP-API. Entries to the left of the first address OBP-API does not trust may have been written by the caller, so read them as a claim, not a fact + |- `target_ip`: the `X-Forwarded-Host` request header, as sent |- the `api_instance_id` of the OBP-API instance that served the call |- the response body, for selected endpoints only | @@ -6851,6 +6875,61 @@ object Glossary extends MdcLoggable { """) + glossaryItems += GlossaryItem( + title = "Client IP Address", + description = + s""" + |# Client IP Address + | + |The **Client IP Address** is the address of the person or program that really made a call, as opposed to the address of whatever passed the call on to OBP-API. OBP-API uses it for per-address [Rate Limiting](/glossary#Rate-Limiting) (including the documentation limit for callers who are not logged in), for IP penalties, and for the busiest-callers view. + | + |## The problem + | + |OBP-API only sees the machine that opened the connection to it. When a browser talks to API Explorer II, API Explorer II talks to OBP-API; when a User chats with Opey, Opey asks OBP-MCP, and OBP-MCP calls OBP-API. Without help, every one of those calls would appear to come from the same server, so one busy User could use up everyone's limit, and a penalty would lock out everyone at once. + | + |## How the address is passed on: X-Forwarded-For + | + |Each hop adds the address it received the request from to the `X-Forwarded-For` header before passing the request on. The header therefore reads "client, first hop, second hop, ...". For a chat with Opey behind NGINX, the chain that reaches OBP-API looks like this: + | + || Hop | Receives the request from | Appends | + ||---|---|---| + || NGINX in front of API Explorer II | the browser | the browser's address | + || API Explorer II | NGINX | NGINX's address | + || Opey | API Explorer II | API Explorer II's address | + || OBP-MCP | Opey | Opey's address | + || OBP-API | OBP-MCP (the TCP peer) | nothing: it reads the chain | + | + |API Explorer II, Opey and OBP-MCP each append the address of the machine they received the request from, the same way NGINX does. Opey keeps the chain outside the conversation, so it never reaches the language model, and it replaces any address header the model writes into a tool call. On the Berlin Group endpoints, API Explorer II also sends the browser's address as `PSU-IP-Address`. + | + |## How OBP-API stops a false address + | + |Anyone can write anything at the left end of the chain, so OBP-API never simply takes the first entry. It reads the chain **from the right**: it skips every address it trusts, and the first address it does not trust is the client. A caller that OBP-API does not trust cannot name a false address, because the walk stops at that caller's own address. For this to work, every hop (NGINX, API Explorer II, Opey, OBP-MCP) must be on OBP-API's list of trusted addresses, and it must not be possible to reach the services behind NGINX except through the hops in front of them. + | + |OBP-API also accepts an `X-Real-IP` header holding a single address, for a deployment with one proxy that overwrites it. + | + |## Where the address is recorded + | + |Each [API Metrics](/glossary#API-Metrics) record keeps the client address OBP-API decided on (`source_ip`) and the whole list of hops (`forwarded_for`), so a call can be traced back through the apps it passed through. + | + |## What the address does and does not tell you + | + |The client address is the address that opened a connection to the outermost trusted proxy. It cannot be faked by writing a false sending address on the network packets, because opening a connection needs the caller to receive the proxy's reply. It may, however, belong to a home router, a company's or mobile network's shared address, a VPN or Tor exit, rather than to the person's own device. + | + |## On this instance + | + |- Client addresses ${if (APIUtil.getPropsAsBoolValue("trust.proxy.enabled", false)) s"are taken from the `${APIUtil.getPropsValue("trust.proxy.header", "X-Real-IP")}` header" else "are not taken from any header: the machine that opened the connection is treated as the client"}. + |- ${APIUtil.getPropsValue("trust.proxy.peers").toOption.map(_.trim).filter(_.nonEmpty) match { + case Some(peers) => s"The header is believed only from these addresses: $peers." + case None => "The header is believed from any caller, and for `X-Forwarded-For` the leftmost address is used, which the client itself can write." + }} + | + |Deployment Checks (`GET /obp/v7.0.0/management/system/diagnostics/deployment`, Role CanGetConfig) show whether requests arrive with a forwarding header, which machines send it, and whether any were ignored because they came from an address that is not trusted. + | + |See also: [Rate Limiting](/glossary#Rate-Limiting), [API Metrics](/glossary#API-Metrics), [OBP-MCP](/glossary#OBP-MCP). + | +""") + + glossaryItems += GlossaryItem( title = "Connector Metrics", description = @@ -6948,7 +7027,7 @@ object Glossary extends MdcLoggable { | |## Roles | - |Managing Groups needs CanCreateGroupAtOneBank, CanGetGroupsAtOneBank, CanUpdateGroupAtOneBank and CanDeleteGroupAtOneBank at the Group's bank id, or the AllBanks version of each (required for a system level Group). Adding and removing members needs CanAddUserToGroupAtOneBank and CanRemoveUserFromGroupAtOneBank, or their AllBanks versions; syncing a Group's members needs both. + |Managing Groups needs CanCreateGroupAtOneBank, CanGetGroupsAtOneBank, CanUpdateGroupAtOneBank and CanDeleteGroupAtOneBank at the Group's bank id, or the AllBanks version of each (required for a system level Group). Adding and removing members needs CanAddUserToGroupAtOneBank and CanRemoveUserFromGroupAtOneBank, or their AllBanks versions; syncing a Group's members, or a user's Groups, needs both. | |## Endpoints | @@ -6961,6 +7040,8 @@ object Glossary extends MdcLoggable { |- [Remove User from Group](${apiExplorerUrl}/resource-docs/OBPv6.0.0?operationid=OBPv6.0.0-removeUserFromGroup): `DELETE /obp/v6.0.0/users/USER_ID/group-entitlements/GROUP_ID` |- [Get User's Group Memberships](${apiExplorerUrl}/resource-docs/OBPv6.0.0?operationid=OBPv6.0.0-getUserGroupMemberships): `GET /obp/v6.0.0/users/USER_ID/group-entitlements` |- [Sync Group Members](${apiExplorerUrl}/resource-docs/OBPv7.0.0?operationid=OBPv7.0.0-syncGroupMembers): `POST /obp/v7.0.0/management/groups/GROUP_ID/sync-members` + |- [Sync Group Member](${apiExplorerUrl}/resource-docs/OBPv7.0.0?operationid=OBPv7.0.0-syncGroupMember): `POST /obp/v7.0.0/management/groups/GROUP_ID/users/USER_ID/sync`, the same for one member. + |- [Sync User Groups](${apiExplorerUrl}/resource-docs/OBPv7.0.0?operationid=OBPv7.0.0-syncUserGroups): `POST /obp/v7.0.0/management/users/USER_ID/sync-groups`, one user in every Group they are in, and the Entitlements left by Groups since deleted. | |How Roles and Entitlements control access is described ${getGlossaryItemLink("API.Access Control")}. |""".stripMargin) diff --git a/obp-api/src/main/scala/code/api/util/NotificationUtil.scala b/obp-api/src/main/scala/code/api/util/NotificationUtil.scala index 9d6374f1eb..096f3f08eb 100644 --- a/obp-api/src/main/scala/code/api/util/NotificationUtil.scala +++ b/obp-api/src/main/scala/code/api/util/NotificationUtil.scala @@ -29,54 +29,56 @@ package code.api.util import code.api.Constant import code.entitlement.Entitlement +import code.messageoutbox.MessageOutbox import code.users.Users import code.util.Helper.MdcLoggable import com.openbankproject.commons.model.User import net.liftweb.common.Box - - -import scala.collection.immutable.List -import scala.concurrent.Future -import com.openbankproject.commons.ExecutionContext.Implicits.global +import net.liftweb.util.Helpers.tryo object NotificationUtil extends MdcLoggable { - def sendEmailRegardingAssignedRole(userId : String, entitlement: Entitlement): Unit = { - // Fire-and-forget: the user lookup and the SMTP send both block, and the - // grant-entitlement response must not wait on them. - Future { - val user = Users.users.vend.getUserByUserId(userId) - sendEmailRegardingAssignedRole(user, entitlement) - }.failed.foreach(e => - logger.error(s"sendEmailRegardingAssignedRole says: failed for userId=$userId role=${entitlement.roleName}", e) - ) - } + /** + * Queue the "you have been granted a Role" email in the message outbox; the relay sends it. + * + * The row is written in the caller's transaction, so a grant that is rolled back (a request that + * fails or times out) emails nobody, and the SMTP send never runs on a request thread. Sending on + * the shared pool instead let a burst of grants (a Group sync granting dozens of Roles to each + * member) occupy every thread with blocking sends and stall the whole API. + */ + def sendEmailRegardingAssignedRole(userId : String, entitlement: Entitlement): Unit = + sendEmailRegardingAssignedRole(Users.users.vend.getUserByUserId(userId), entitlement) + def sendEmailRegardingAssignedRole(user: Box[User], entitlement: Entitlement): Unit = { - val mailSent = for { + val queued = for { user <- user from <- APIUtil.getPropsValue("mail.api.consumer.registered.sender.address") ?~ "Could not send mail: Missing props param for 'from'" - } yield { - val bodyOfMessage : String = s"""Dear ${user.name}, - | - |You have been granted the entitlement to use ${entitlement.roleName} on ${Constant.HostName} - | - |Cheers - |""".stripMargin - val emailContent = CommonsEmailWrapper.EmailContent( - from = from, - to = List(user.emailAddress), - subject = s"You have been granted the role: ${entitlement.roleName}", - textContent = Some(bodyOfMessage) - ) - // Blocking SMTP send (Transport.send) — only call this off the request - // thread; the userId overload above wraps it in a Future. - CommonsEmailWrapper.sendTextEmail(emailContent) - } - if(mailSent.isEmpty) { + row <- { + val bodyOfMessage : String = s"""Dear ${user.name}, + | + |You have been granted the entitlement to use ${entitlement.roleName} on ${Constant.HostName} + | + |Cheers + |""".stripMargin + tryo(MessageOutbox.enqueueEmail( + subjectId = entitlement.entitlementId, + subjectIdType = MessageOutbox.SUBJECT_TYPE_ENTITLEMENT_ID, + operationName = MessageOutbox.OPERATION_ROLE_GRANTED_EMAIL, + CommonsEmailWrapper.EmailContent( + from = from, + to = List(user.emailAddress), + subject = s"You have been granted the role: ${entitlement.roleName}", + textContent = Some(bodyOfMessage) + ) + )) + } + } yield row + if(queued.isEmpty) { val info = s""" |Sending email is omitted. |User: $user |Props mail.api.consumer.registered.sender.address: ${APIUtil.getPropsValue("mail.api.consumer.registered.sender.address")} + |Reason: $queued |""".stripMargin this.logger.warn(info) } diff --git a/obp-api/src/main/scala/code/api/util/RemoteIpUtil.scala b/obp-api/src/main/scala/code/api/util/RemoteIpUtil.scala index 6738703122..54cb35bf2c 100644 --- a/obp-api/src/main/scala/code/api/util/RemoteIpUtil.scala +++ b/obp-api/src/main/scala/code/api/util/RemoteIpUtil.scala @@ -40,10 +40,14 @@ import code.util.Helper.MdcLoggable * * proxy_set_header X-Real-IP $remote_addr; * - * For `X-Forwarded-For`, the leftmost value is treated as the client. This is only - * trustworthy when the proxy is configured with `set_real_ip_from` + `real_ip_recursive` - * so it sanitises the forwarded chain before forwarding upstream. `X-Real-IP` is the - * simpler choice for single-proxy deployments. + * `X-Forwarded-For` carries a chain: each hop (NGINX, a server-side application such as + * API Explorer II, Opey, OBP-MCP) appends the address it received the request from, so the + * chain reads "client, first hop, second hop, ...". Anyone can write anything at the left + * end, so the chain is read from the right: skip every address in `trust.proxy.peers` and + * the first address that is not trusted is the client. A caller that is not trusted cannot + * name a false address, because it becomes the client itself. This needs every hop listed + * in `trust.proxy.peers`; with the list unset every address counts as trusted and the + * leftmost address is used, which is only safe when the outermost proxy replaces the chain. * * `trust.proxy.peers` closes a gap: without it, the header is believed from whoever sent the * request, so a caller that can reach OBP-API directly (bypassing the proxy) can name any @@ -86,7 +90,7 @@ object RemoteIpUtil extends MdcLoggable { Resolution(peer, peer, ForwardingHeaders.exists(h => getHeader(h).exists(_.trim.nonEmpty)), headerHonoured = false, headerFromUntrustedPeer = false) } else { val headerName = APIUtil.getPropsValue("trust.proxy.header", "X-Real-IP") - val fromHeader = getHeader(headerName).flatMap(raw => extractClientIp(headerName, raw)).map(canonical) + val fromHeader = getHeader(headerName).flatMap(raw => extractClientIp(headerName, raw)) fromHeader match { case None => Resolution(peer, peer, forwardingHeaderPresent = false, headerHonoured = false, headerFromUntrustedPeer = false) case Some(_) if !peerIsTrusted(peer) => Resolution(peer, peer, forwardingHeaderPresent = true, headerHonoured = false, headerFromUntrustedPeer = true) @@ -95,6 +99,17 @@ object RemoteIpUtil extends MdcLoggable { } } + /** The hops a request passed through, for API Metrics: the X-Forwarded-For chain it arrived + * with (all header lines, comma-joined in order) followed by the TCP peer, so the last hop is + * recorded too. This is a record of what arrived, not a decision: entries left of the first + * address that is not trusted may have been written by the client. The client address itself + * comes from [[resolve]]. */ + def forwardedForPath(socketPeer: String, forwardedForHeader: Option[String]): String = { + val peer = canonical(socketPeer) + val incoming = forwardedForHeader.map(_.trim).filter(_.nonEmpty) + (incoming.toList ++ List(peer).filter(_.nonEmpty)).mkString(", ") + } + /** The configured trusted peers, as parsed CIDR ranges (a single address is a /32 or /128). */ def trustedPeers: List[(Array[Byte], Int)] = APIUtil.getPropsValue("trust.proxy.peers").toList @@ -139,15 +154,38 @@ object RemoteIpUtil extends MdcLoggable { else unbracketed } - /** Single-value headers (X-Real-IP) yield the value as-is. - * X-Forwarded-For is comma-separated; the leftmost entry is the original client. */ - private def extractClientIp(headerName: String, raw: String): Option[String] = { - val candidate = - if (headerName.equalsIgnoreCase("X-Forwarded-For")) - raw.split(",").headOption.getOrElse("") - else - raw - val trimmed = candidate.trim - if (trimmed.isEmpty) None else Some(trimmed) + /** The client address a forwarding header names, in canonical form. + * A single-value header (X-Real-IP) yields its value. If the request carried it more than + * once, the values arrive comma-joined and the first is used, as before. + * X-Forwarded-For yields the client found by [[clientFromForwardedFor]]. */ + private def extractClientIp(headerName: String, raw: String): Option[String] = + if (headerName.equalsIgnoreCase("X-Forwarded-For")) clientFromForwardedFor(raw) + else Option(raw.split(",").headOption.getOrElse("").trim).filter(_.nonEmpty).map(canonical) + + /** The client named by an X-Forwarded-For chain, read from the right. + * + * Each hop appends the address it received the request from, so the rightmost entry was + * written by the TCP peer (already checked to be trusted), the next one by the hop before + * it, and so on. Walking leftwards, every address in trust.proxy.peers is a hop that can + * be believed about the entry to its left; the first address that is not in the list is + * the client. Entries further left were written by the client or by hops nobody vouches + * for, and are ignored. + * + * An entry that is not an address (for example "unknown") stops the walk: nothing to its + * left can be believed, so the nearest trusted address to its right is the client, or + * None (meaning the TCP peer) when it is the rightmost entry. + * + * When every entry is trusted, including when trust.proxy.peers is unset, the leftmost + * entry is the client. */ + private[util] def clientFromForwardedFor(raw: String): Option[String] = { + val chainFromTheRight = raw.split(",").map(entry => canonical(entry)).filter(_.nonEmpty).toList.reverse + val firstUntrustedIndex = chainFromTheRight.indexWhere(address => !peerIsTrusted(address)) + if (firstUntrustedIndex < 0) chainFromTheRight.lastOption + else { + val firstUntrusted = chainFromTheRight(firstUntrustedIndex) + if (addressBytes(firstUntrusted).isDefined) Some(firstUntrusted) + else if (firstUntrustedIndex == 0) None + else Some(chainFromTheRight(firstUntrustedIndex - 1)) + } } } diff --git a/obp-api/src/main/scala/code/api/util/WriteMetricUtil.scala b/obp-api/src/main/scala/code/api/util/WriteMetricUtil.scala index 647b308d63..bc1abd9678 100644 --- a/obp-api/src/main/scala/code/api/util/WriteMetricUtil.scala +++ b/obp-api/src/main/scala/code/api/util/WriteMetricUtil.scala @@ -73,6 +73,7 @@ object WriteMetricUtil extends MdcLoggable { responseBodyToWrite: String, sourceIp: String, targetIp: String, + forwardedFor: String, authType: String) private def persistAndPublishMetric(responseBody: Any, cc: CallContextLight): Unit = { @@ -85,8 +86,11 @@ object WriteMetricUtil extends MdcLoggable { implementedByPartialFunction = cc.partialFunctionName, duration = callDuration(cc), responseBodyToWrite = responseBodyForMetric(responseBody, cc), - sourceIp = requestHeaderValue(cc, "x-forwarded-for"), + // The client address OBP-API decided on, and the hops the request passed through + // (see RemoteIpUtil). The raw X-Forwarded-For header is not stored as the source address. + sourceIp = cc.ipAddress, targetIp = requestHeaderValue(cc, "x-forwarded-host"), + forwardedFor = cc.forwardedFor, authType = deriveAuthType(cc) ) @@ -166,6 +170,7 @@ object WriteMetricUtil extends MdcLoggable { responseBodyToWrite, sourceIp, targetIp, + forwardedFor, code.api.Constant.ApiInstanceId, cc.consentReferenceId.orNull, cc.certificateTrust.orNull, diff --git a/obp-api/src/main/scala/code/api/util/http4s/Http4sSupport.scala b/obp-api/src/main/scala/code/api/util/http4s/Http4sSupport.scala index 87e63bc6dc..7a3502712d 100644 --- a/obp-api/src/main/scala/code/api/util/http4s/Http4sSupport.scala +++ b/obp-api/src/main/scala/code/api/util/http4s/Http4sSupport.scala @@ -682,6 +682,10 @@ object Http4sCallContextBuilder { } } yield CallContext( url = request.uri.renderString, + forwardedFor = RemoteIpUtil.forwardedForPath( + request.remoteAddr.map(_.toUriString).getOrElse(""), + requestHeaderValues(request, "X-Forwarded-For") + ), verb = request.method.name, implementedInVersion = apiVersion, correlationId = extractCorrelationId(request), @@ -726,9 +730,14 @@ object Http4sCallContextBuilder { def clientIpResolution(request: Request[IO]): RemoteIpUtil.Resolution = RemoteIpUtil.resolve( request.remoteAddr.map(_.toUriString).getOrElse(""), - name => request.headers.get(CIString(name)).map(_.head.value) + name => requestHeaderValues(request, name) ) + /** A request header's value. A header sent on several lines is joined in order, so an + * X-Forwarded-For line the client wrote cannot hide the lines the proxies added after it. */ + private def requestHeaderValues(request: Request[IO], name: String): Option[String] = + request.headers.get(CIString(name)).map(_.toList.map(_.value).mkString(", ")) + private def extractIpAddress(request: Request[IO]): String = clientIpResolution(request).clientIp /** diff --git a/obp-api/src/main/scala/code/api/v6_0_0/JSONFactory6.0.0.scala b/obp-api/src/main/scala/code/api/v6_0_0/JSONFactory6.0.0.scala index 00b3473c3f..750c3c9c43 100644 --- a/obp-api/src/main/scala/code/api/v6_0_0/JSONFactory6.0.0.scala +++ b/obp-api/src/main/scala/code/api/v6_0_0/JSONFactory6.0.0.scala @@ -468,8 +468,11 @@ case class MetricJsonV600( verb: String, correlation_id: String, duration: Long, + // The client address OBP-API decided on (see the Client IP Address glossary item) source_ip: String, target_ip: String, + // The hops the request passed through: the X-Forwarded-For chain it arrived with, then the TCP peer + forwarded_for: String, response_body: org.json4s.JValue, status_code: Int, operation_id: String, @@ -1768,6 +1771,7 @@ object JSONFactory600 extends CustomJsonFormats with MdcLoggable { duration = metric.getDuration(), source_ip = metric.getSourceIp(), target_ip = metric.getTargetIp(), + forwarded_for = Option(metric.getForwardedFor()).getOrElse(""), response_body = com.openbankproject.commons.util.JsonAliases.parseOpt(metric.getResponseBody()).getOrElse(org.json4s.JString("Not enabled")), status_code = metric.getHttpCode(), operation_id = operationId, diff --git a/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Groups.scala b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Groups.scala index 3c4cca5cf2..2e88e8647e 100644 --- a/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Groups.scala +++ b/obp-api/src/main/scala/code/api/v7_0_0/Http4s700Groups.scala @@ -34,7 +34,7 @@ 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.{APIUtil, ApiRole, CustomJsonFormats} +import code.api.util.{APIUtil, ApiRole, CustomJsonFormats, NewStyle} import code.entitlement.Entitlement import code.group.{GroupMemberships, GroupTrait} import code.users.{Users => UserVend} @@ -42,7 +42,7 @@ import code.util.Helper import com.github.dwickern.macros.NameOf.nameOf import com.openbankproject.commons.ExecutionContext.Implicits.global import com.openbankproject.commons.util.ApiVersion -import net.liftweb.common.Full +import net.liftweb.common.{Empty, Full} import org.http4s._ import org.http4s.dsl.io._ import org.json4s.Formats @@ -64,6 +64,20 @@ case class GroupMembersSyncJsonV700( dry_run: Boolean, members: List[GroupMemberSyncJsonV700] ) +case class UserGroupSyncJsonV700( + group_id: String, + bank_id: Option[String], + group_deleted: Boolean, + entitlements_created: List[String], + entitlements_deleted: List[String], + entitlements_moved: List[GroupMemberRoleMovedJsonV700] +) +case class UserGroupsSyncJsonV700( + user_id: String, + username: String, + dry_run: Boolean, + groups: List[UserGroupSyncJsonV700] +) object Http4s700Groups { implicit val formats: Formats = CustomJsonFormats.formats @@ -83,11 +97,22 @@ object Http4s700Groups { APIUtil.isSuperAdmin(userId) || (holds(addRoles) && holds(removeRoles)) } - /** Bring one member's Entitlements in line with the Group's Roles. Changes nothing when `dryRun`. */ - private def syncMember(group: GroupTrait, userId: String, grantedBy: String, dryRun: Boolean): GroupMemberSyncJsonV700 = { + private def isDryRun(req: Request[IO]): Boolean = + req.uri.query.params.get("dry_run").exists(_.equalsIgnoreCase("true")) + + private def missingRolesMessage = UserHasMissingRoles + addRoles.mkString(" or ") + " and " + removeRoles.mkString(" or ") + + private def usernameOf(userId: String): String = UserVend.users.vend.getUserByUserId(userId).map(_.name).getOrElse("") + + /** + * Bring one member's Entitlements in line with the Group's Roles. Changes nothing when `dryRun`. + * `alsoHeld`: Roles at the Group's bank id an earlier Group of this dry run would have granted. + */ + private def syncMember(group: GroupTrait, userId: String, grantedBy: String, dryRun: Boolean, + alsoHeld: Set[String] = Set.empty): GroupMemberSyncJsonV700 = { val bankId = group.bankId.getOrElse("") val held = Entitlement.entitlement.vend.getEntitlementsByUserId(userId).toList.flatten.filter(_.bankId == bankId) - val heldRoles = held.map(_.roleName).toSet + val heldRoles = held.map(_.roleName).toSet ++ alsoHeld val toCreate = group.listOfRoles.filterNot(heldRoles.contains).distinct val noLongerGranted = held.filter(e => e.groupId.contains(group.groupId) && !group.listOfRoles.contains(e.roleName)) @@ -104,7 +129,7 @@ object Http4s700Groups { } GroupMemberSyncJsonV700( user_id = userId, - username = UserVend.users.vend.getUserByUserId(userId).map(_.name).getOrElse(""), + username = usernameOf(userId), entitlements_created = toCreate, entitlements_deleted = toDelete.map(_._1.roleName).sorted, entitlements_moved = toMove.map { case (e, other) => GroupMemberRoleMovedJsonV700(e.roleName, other.get.groupId) } @@ -116,13 +141,11 @@ object Http4s700Groups { lazy val syncGroupMembers: HttpRoutes[IO] = HttpRoutes.of[IO] { case req @ POST -> `prefixPath` / "management" / "groups" / groupId / "sync-members" => EndpointHelpers.withUser(req) { (user, cc) => - val dryRun = req.uri.query.params.get("dry_run").exists(_.equalsIgnoreCase("true")) + val dryRun = isDryRun(req) for { group <- Future(GroupTrait.group.vend.getGroup(groupId)) .map(APIUtil.unboxFullOrFail(_, Some(cc), s"$UnknownError Group not found", 404)) - _ <- Helper.booleanToFuture( - UserHasMissingRoles + addRoles.mkString(" or ") + " and " + removeRoles.mkString(" or "), - failCode = 403, cc = Some(cc))(mayAddAndRemove(group.bankId, user.userId)) + _ <- Helper.booleanToFuture(missingRolesMessage, failCode = 403, cc = Some(cc))(mayAddAndRemove(group.bankId, user.userId)) _ <- Helper.booleanToFuture(s"$UnknownError Group is not enabled", 400, Some(cc))(group.isEnabled) granted <- Entitlement.entitlement.vend.getEntitlementsByGroupId(groupId) .map(APIUtil.unboxFullOrFail(_, Some(cc), s"$UnknownError Cannot get entitlements", 400)) @@ -178,4 +201,186 @@ object Http4s700Groups { None, http4sPartialFunction = Some(syncGroupMembers) ) + + /** + * The user's Entitlements recorded against a Group that has since been deleted (deleting a Group + * leaves them in place): each is moved to another Group of the user's that grants the Role, or + * deleted. Changes nothing when `dryRun`. + */ + private def syncDeletedGroup(groupId: String, orphans: List[Entitlement], userId: String, dryRun: Boolean): UserGroupSyncJsonV700 = { + val (toMove, toDelete) = orphans + .map(e => (e, GroupMemberships.otherGroupGranting(userId, e.bankId, e.roleName, groupId))) + .partition(_._2.isDefined) + if (!dryRun) { + toMove.foreach { case (e, other) => Entitlement.entitlement.vend.setEntitlementGroupId(e.entitlementId, other.get.groupId) } + toDelete.foreach { case (e, _) => Entitlement.entitlement.vend.deleteEntitlement(Full(e)) } + } + UserGroupSyncJsonV700( + group_id = groupId, + bank_id = orphans.headOption.map(_.bankId).filter(_.nonEmpty), + group_deleted = true, + entitlements_created = Nil, + entitlements_deleted = toDelete.map(_._1.roleName).sorted, + entitlements_moved = toMove.map { case (e, other) => GroupMemberRoleMovedJsonV700(e.roleName, other.get.groupId) } + .sortBy(_.role_name) + ) + } + + // Route: POST /obp/v7.0.0/management/groups/GROUP_ID/users/USER_ID/sync + lazy val syncGroupMember: HttpRoutes[IO] = HttpRoutes.of[IO] { + case req @ POST -> `prefixPath` / "management" / "groups" / groupId / "users" / userId / "sync" => + EndpointHelpers.withUser(req) { (user, cc) => + val dryRun = isDryRun(req) + for { + group <- Future(GroupTrait.group.vend.getGroup(groupId)) + .map(APIUtil.unboxFullOrFail(_, Some(cc), s"$UnknownError Group not found", 404)) + _ <- Helper.booleanToFuture(missingRolesMessage, failCode = 403, cc = Some(cc))(mayAddAndRemove(group.bankId, user.userId)) + _ <- NewStyle.function.findByUserId(userId, Some(cc)) + _ <- Helper.booleanToFuture(s"$UnknownError Group is not enabled", 400, Some(cc))(group.isEnabled) + granted <- Entitlement.entitlement.vend.getEntitlementsByGroupId(groupId) + .map(APIUtil.unboxFullOrFail(_, Some(cc), s"$UnknownError Cannot get entitlements", 400)) + _ <- Helper.booleanToFuture(s"$UnknownError The User is not a member of the Group", 404, Some(cc))( + GroupMemberships.userIdsOfGroup(groupId, granted).contains(userId)) + synced <- Future(syncMember(group, userId, user.userId, dryRun)) + } yield GroupMembersSyncJsonV700(group.groupId, group.bankId, dryRun, List(synced)) + } + } + + resourceDocs += ResourceDoc( + implementedInApiVersion, + nameOf(syncGroupMember), + "POST", + "/management/groups/GROUP_ID/users/USER_ID/sync", + "Sync Group Member", + s"""Bring the Entitlements of one member of a Group in line with the Group's current Roles. + | + |The same as Sync Group Members (POST /management/groups/GROUP_ID/sync-members), for one member only: + | + |- a Role of the Group the member does not hold at the Group's bank id is granted (the member gets + | an email for it), recorded against this Group; + |- an Entitlement this Group granted, for a Role the Group no longer has, is deleted, unless another + | Group the member is in, at the same bank id, still grants that Role: then it is kept and recorded + | against that Group (no email; the member's Roles do not change). + | + |Entitlements granted by hand, or by other Groups, are not touched. The user is not added to or + |removed from the Group. + | + |The user must be a member of the Group (added to it, or holding an Entitlement it granted), or 404. + | + |With `dry_run=true` nothing is changed and the response says what would be. + | + |Requires CanAddUserToGroupAtOneBank or CanAddUserToGroupAtAllBanks, and + |CanRemoveUserFromGroupAtOneBank or CanRemoveUserFromGroupAtAllBanks (the AllBanks Roles for a + |system level Group). + |""".stripMargin, + EmptyBody, + GroupMembersSyncJsonV700( + group_id = "group-id-123", + bank_id = Some("gh.29.uk"), + dry_run = false, + members = List(GroupMemberSyncJsonV700( + user_id = "user-id-123", + username = "felixsmith", + entitlements_created = List("CanGetCustomer"), + entitlements_deleted = List("CanCreateTransaction"), + entitlements_moved = List(GroupMemberRoleMovedJsonV700("CanGetAccount", "group-id-456")) + )) + ), + List($AuthenticatedUserIsRequired, UserHasMissingRoles, UserNotFoundById, UnknownError), + List(apiTagGroup, apiTagUser, apiTagEntitlement), + None, + http4sPartialFunction = Some(syncGroupMember) + ) + + // Route: POST /obp/v7.0.0/management/users/USER_ID/sync-groups + lazy val syncUserGroups: HttpRoutes[IO] = HttpRoutes.of[IO] { + case req @ POST -> `prefixPath` / "management" / "users" / userId / "sync-groups" => + EndpointHelpers.withUser(req) { (user, cc) => + val dryRun = isDryRun(req) + for { + _ <- NewStyle.function.findByUserId(userId, Some(cc)) + held = Entitlement.entitlement.vend.getEntitlementsByUserId(userId).toList.flatten + found = GroupMemberships.groupIdsOfUser(userId).map(id => id -> GroupTrait.group.vend.getGroup(id)) + groups = found.collect { case (_, Full(g)) if g.isEnabled => g } + deleted = found.collect { case (id, Empty) => id -> held.filter(_.groupId.contains(id)) }.filter(_._2.nonEmpty) + bankIds = (groups.map(_.bankId) ++ deleted.flatMap(_._2).map(e => Some(e.bankId).filter(_.nonEmpty))).distinct + _ <- Helper.booleanToFuture(missingRolesMessage, failCode = 403, cc = Some(cc))( + bankIds.forall(mayAddAndRemove(_, user.userId))) + synced <- Future { + val fromDeleted = deleted.map { case (id, orphans) => syncDeletedGroup(id, orphans, userId, dryRun) } + // In a dry run nothing is granted, so a Role an earlier Group would grant is passed on to the later ones. + val (fromGroups, _) = groups.foldLeft((List.empty[UserGroupSyncJsonV700], Map.empty[String, Set[String]])) { + case ((done, planned), g) => + val bankId = g.bankId.getOrElse("") + val m = syncMember(g, userId, user.userId, dryRun, planned.getOrElse(bankId, Set.empty)) + (done :+ UserGroupSyncJsonV700(g.groupId, g.bankId, group_deleted = false, + m.entitlements_created, m.entitlements_deleted, m.entitlements_moved), + planned.updated(bankId, planned.getOrElse(bankId, Set.empty) ++ m.entitlements_created)) + } + fromDeleted ++ fromGroups + } + } yield UserGroupsSyncJsonV700(userId, usernameOf(userId), dryRun, synced) + } + } + + resourceDocs += ResourceDoc( + implementedInApiVersion, + nameOf(syncUserGroups), + "POST", + "/management/users/USER_ID/sync-groups", + "Sync User Groups", + s"""Bring a user's Entitlements in line with every Group they are in, and clear what is left of + |Groups that have been deleted. + | + |For each enabled Group the user is in (added to it, or holding an Entitlement it granted), at any + |bank id, the same as Sync Group Member (POST /management/groups/GROUP_ID/users/USER_ID/sync): + | + |- a Role of the Group the user does not hold at the Group's bank id is granted (the user gets an + | email for it), recorded against that Group; + |- an Entitlement the Group granted, for a Role the Group no longer has, is deleted, unless another + | Group the user is in, at the same bank id, still grants that Role: then it is kept and recorded + | against that Group. + | + |Deleting a Group leaves the Entitlements it granted in place. Each of the user's Entitlements + |recorded against a deleted Group is moved to another Group the user is in that grants the Role at + |that bank id, or else deleted. These are listed with `group_deleted: true`. + | + |Disabled Groups are left as they are. Entitlements granted by hand are not touched. The user is not + |added to or removed from any Group. + | + |With `dry_run=true` nothing is changed and the response says what would be. + | + |Requires CanAddUserToGroupAtOneBank or CanAddUserToGroupAtAllBanks, and + |CanRemoveUserFromGroupAtOneBank or CanRemoveUserFromGroupAtAllBanks, at the bank id of every Group + |involved (the AllBanks Roles for a system level Group). If any is missing, nothing is changed. + |""".stripMargin, + EmptyBody, + UserGroupsSyncJsonV700( + user_id = "user-id-123", + username = "felixsmith", + dry_run = false, + groups = List( + UserGroupSyncJsonV700( + group_id = "group-id-123", + bank_id = Some("gh.29.uk"), + group_deleted = false, + entitlements_created = List("CanGetCustomer"), + entitlements_deleted = List("CanCreateTransaction"), + entitlements_moved = List(GroupMemberRoleMovedJsonV700("CanGetAccount", "group-id-456")) + ), + UserGroupSyncJsonV700( + group_id = "group-id-789", + bank_id = Some("gh.29.uk"), + group_deleted = true, + entitlements_created = Nil, + entitlements_deleted = List("CanGetTransaction"), + entitlements_moved = Nil + ) + ) + ), + List($AuthenticatedUserIsRequired, UserHasMissingRoles, UserNotFoundById, UnknownError), + List(apiTagGroup, apiTagUser, apiTagEntitlement), + None, + http4sPartialFunction = Some(syncUserGroups) + ) } diff --git a/obp-api/src/main/scala/code/dynamicEntity/DynamicDataProvider.scala b/obp-api/src/main/scala/code/dynamicEntity/DynamicDataProvider.scala index 4d795ec23e..5f72c0433b 100644 --- a/obp-api/src/main/scala/code/dynamicEntity/DynamicDataProvider.scala +++ b/obp-api/src/main/scala/code/dynamicEntity/DynamicDataProvider.scala @@ -30,6 +30,7 @@ package code.DynamicData import org.json4s._ import com.openbankproject.commons.model.{Converter, JsonFieldReName} import net.liftweb.common.Box +import java.util.Date import org.json4s.JObject import net.liftweb.util.SimpleInjector @@ -47,6 +48,31 @@ trait DynamicDataT { def bankId: Option[String] def userId: Option[String] def isPersonalEntity: Boolean + /** + * When the record was first saved. None for a record saved before the timestamp columns existed, + * and for an implementation that does not store timestamps. + */ + def createdDate: Option[Date] = None + /** + * When the record was last saved. None for a record not saved since the timestamp columns were + * added, and for an implementation that does not store timestamps. + */ + def updatedDate: Option[Date] = None + /** + * The user who made the call that first saved the record. That is an agent's own user id when an + * agent made the call for somebody else. None for a record saved before the column existed, and for + * an implementation that does not store it. + */ + def createdByUserId: Option[String] = None + /** + * The user the call that first saved the record was made for. It equals createdByUserId when nobody + * was delegating. None in the same cases as createdByUserId. + */ + def createdByOnBehalfOfUserId: Option[String] = None + /** As createdByUserId, for the call that last saved the record. */ + def updatedByUserId: Option[String] = None + /** As createdByOnBehalfOfUserId, for the call that last saved the record. */ + def updatedByOnBehalfOfUserId: Option[String] = None } case class DynamicDataCommons(dynamicEntityName: String, @@ -74,10 +100,19 @@ trait DynamicDataProvider { def getAllDataJsonCommunity(bankId: Option[String], entityName: String): List[JObject] def getCommunity(bankId: Option[String], entityName: String, id: String): Box[DynamicDataT] + /** + * The records of one entity in one space whose ids are in `ids`, whoever owns them. An id with no + * record is left out of the answer rather than reported. This exists so that a response can carry + * each record's metadata after the records themselves were read by some other route, such as the + * connector or the projection tables; it is not an access check. + */ + def getByIds(bankId: Option[String], entityName: String, ids: List[String]): List[DynamicDataT] + // Community mutation methods - operate on a row regardless of owner (used by row-level access, // where the ACL, not ownership, decides who may update/delete). Preserve the row's existing - // userId / isPersonalEntity so its provenance is unchanged. - def updateCommunity(bankId: Option[String], entityName: String, requestBody: JObject, id: String): Box[DynamicDataT] + // userId / isPersonalEntity so its provenance is unchanged. callerUserId is the user making the + // call, recorded as the record's last writer. + def updateCommunity(bankId: Option[String], entityName: String, requestBody: JObject, id: String, callerUserId: Option[String]): Box[DynamicDataT] def deleteCommunity(bankId: Option[String], entityName: String, id: String): Box[Boolean] } diff --git a/obp-api/src/main/scala/code/dynamicEntity/MapppedDynamicDataProvider.scala b/obp-api/src/main/scala/code/dynamicEntity/MapppedDynamicDataProvider.scala index 7b8eef8294..f01743340f 100644 --- a/obp-api/src/main/scala/code/dynamicEntity/MapppedDynamicDataProvider.scala +++ b/obp-api/src/main/scala/code/dynamicEntity/MapppedDynamicDataProvider.scala @@ -41,6 +41,8 @@ import net.liftweb.mapper._ import net.liftweb.util.Helpers.tryo import org.apache.commons.lang3.StringUtils +import java.util.Date + /** * Note on IsPersonalEntity flag: * The IsPersonalEntity flag indicates HOW a record was created (via /my/ endpoint or not), @@ -125,13 +127,13 @@ object MappedDynamicDataProvider extends DynamicDataProvider with CustomJsonForm val idName = getIdName(entityName) val JString(idValue) = (requestBody \ idName).asInstanceOf[JString] val dynamicData: DynamicData = DynamicData.create.DynamicDataId(idValue) - val result = saveOrUpdate(bankId, entityName, requestBody, ownerOf(userId), isPersonalEntity, dynamicData) + val result = saveOrUpdate(bankId, entityName, requestBody, ownerOf(userId), isPersonalEntity, userId, dynamicData) result } override def update(bankId: Option[String], entityName: String, requestBody: JObject, id: String, userId: Option[String], isPersonalEntity: Boolean): Box[DynamicDataT] = { val owner = ownerOf(userId) val dynamicData = get(bankId, entityName, id, owner, isPersonalEntity).openOrThrowException(s"$DynamicDataNotFound dynamicEntityName=$entityName, dynamicDataId=$id").asInstanceOf[DynamicData] - saveOrUpdate(bankId, entityName, requestBody, owner, isPersonalEntity, dynamicData) + saveOrUpdate(bankId, entityName, requestBody, owner, isPersonalEntity, userId, dynamicData) } // Separate method for reference validation - only checks ID and entity name exist @@ -200,14 +202,20 @@ object MappedDynamicDataProvider extends DynamicDataProvider with CustomJsonForm case _ => Failure(notFoundCommunityMessage(entityName, id, bankId)) } - override def updateCommunity(bankId: Option[String], entityName: String, requestBody: JObject, id: String): Box[DynamicDataT] = { + override def updateCommunity(bankId: Option[String], entityName: String, requestBody: JObject, id: String, callerUserId: Option[String]): Box[DynamicDataT] = { val dynamicData = getCommunity(bankId, entityName, id) .openOrThrowException(s"$DynamicDataNotFound dynamicEntityName=$entityName, dynamicDataId=$id") .asInstanceOf[DynamicData] // Preserve the row's existing owner/personal flag — row-level access changes the data, not provenance. - saveOrUpdate(bankId, entityName, requestBody, Option(dynamicData.UserId.get), dynamicData.IsPersonalEntity.get, dynamicData) + saveOrUpdate(bankId, entityName, requestBody, Option(dynamicData.UserId.get), dynamicData.IsPersonalEntity.get, callerUserId, dynamicData) } + override def getByIds(bankId: Option[String], entityName: String, ids: List[String]): List[DynamicDataT] = + // Grouped so that a long list never becomes one very long IN clause. + ids.distinct.grouped(1000).toList.flatMap { someIds => + DynamicData.findAll((whereClauseBankOrSystemAndEntity(bankId, entityName) :+ ByList(DynamicData.DynamicDataId, someIds)): _*) + } + override def deleteCommunity(bankId: Option[String], entityName: String, id: String): Box[Boolean] = { getCommunity(bankId, entityName, id).map { d => val result = d.asInstanceOf[DynamicData].delete_! @@ -222,10 +230,40 @@ object MappedDynamicDataProvider extends DynamicDataProvider with CustomJsonForm DynamicData.find((whereClauseBankOrSystemAndEntity(bankId, dynamicEntityName) :+ byOwnership(isPersonalEntity, userId)): _*).isDefined } - private def saveOrUpdate(bankId: Option[String], entityName: String, requestBody: JObject, userId: Option[String], isPersonalEntity: Boolean, dynamicData: => DynamicData): Box[DynamicData] = { + /** + * This works out the two user ids to record as the writer of a record, for a call made by + * `callerUserId`: the caller itself, which is an agent's own id when an agent made the call, and the + * user the call was made for. They are the same when nobody is delegating. + * + * Both are kept because a record is someone's data, and whether a person wrote it themselves or an + * agent wrote it for them has to be answerable from the record. `ref` names the pair of columns + * being filled, so the attribution log line names them. With no caller, both are null. + */ + private def writersOf(callerUserId: Option[String], ref: code.users.UserReference): (String, String) = + callerUserId.filter(_.nonEmpty) match { + case None => (null, null) + case Some(caller) => + code.users.Users.users.vend.attributionOf(caller, ref) match { + case Full(attribution) => (attribution.userId, attribution.onBehalfOfUserId) + case _ => (caller, caller) + } + } + + /** + * This saves a record, new or existing. `userId` is the owner stored in UserId, already resolved. + * `callerUserId` is the user making the call, not resolved, from which the writer columns are + * filled: the created pair only when the record is new, and the updated pair on every save. + */ + private def saveOrUpdate(bankId: Option[String], entityName: String, requestBody: JObject, userId: Option[String], isPersonalEntity: Boolean, + callerUserId: Option[String], dynamicData: => DynamicData): Box[DynamicData] = { val data: DynamicData = dynamicData tryo { val dataStr = json.compactRender(requestBody) + val isNew = !data.saved_? + val (writerUserId, writerOnBehalfOfUserId) = writersOf(callerUserId, + if (isNew) code.users.UserReference.DynamicData_CreatedByUserId else code.users.UserReference.DynamicData_UpdatedByUserId) + if (isNew) data.CreatedByUserId(writerUserId).CreatedByOnBehalfOfUserId(writerOnBehalfOfUserId) + data.UpdatedByUserId(writerUserId).UpdatedByOnBehalfOfUserId(writerOnBehalfOfUserId) val saved = data.BankId(bankId.getOrElse(DYNAMIC_ENTITY_SYSTEM_LEVEL_BANK_ID)) .DynamicEntityName(entityName) .DataJson(dataStr) @@ -243,7 +281,21 @@ object MappedDynamicDataProvider extends DynamicDataProvider with CustomJsonForm } } -class DynamicData extends DynamicDataT with LongKeyedMapper[DynamicData] with IdPK { +/** + * This class is one record of a Dynamic Entity. + * + * It mixes in CreatedUpdated, which adds the columns createdat and updatedat. createdat is set when + * the record is first saved. updatedat is set again on every save, so each update refreshes it. The + * two columns were added after records already existed, so a record saved before then holds NULL in + * both until it is next updated, and only updatedat is filled at that point: its creation time was + * never recorded and is not guessed. createdDate and updatedDate therefore return an Option. + * + * The four writer columns follow the same rule. CreatedByUserId and CreatedByOnBehalfOfUserId are + * set when the record is first saved; UpdatedByUserId and UpdatedByOnBehalfOfUserId on every save. + * Each pair holds the user who made the call and the user it was made for, which differ only when an + * agent acted for somebody. They are separate from UserId, which is the owner of a personal record. + */ +class DynamicData extends DynamicDataT with LongKeyedMapper[DynamicData] with IdPK with CreatedUpdated { override def getSingleton = DynamicData @@ -271,6 +323,15 @@ class DynamicData extends DynamicDataT with LongKeyedMapper[DynamicData] with Id object IsPersonalEntity extends MappedBoolean(this) + /** The user whose call first saved this record: an agent's own id when an agent made the call. */ + object CreatedByUserId extends MappedString(this, 255) + /** The user that call was made for; the same as CreatedByUserId when nobody was delegating. */ + object CreatedByOnBehalfOfUserId extends MappedString(this, 255) + /** As CreatedByUserId, for the call that last saved this record. */ + object UpdatedByUserId extends MappedString(this, 255) + /** As CreatedByOnBehalfOfUserId, for the call that last saved this record. */ + object UpdatedByOnBehalfOfUserId extends MappedString(this, 255) + override def dynamicDataId: Option[String] = Option(DynamicDataId.get) override def dynamicEntityName: String = DynamicEntityName.get override def dataJson: String = DataJson.get @@ -279,6 +340,12 @@ class DynamicData extends DynamicDataT with LongKeyedMapper[DynamicData] with Id override def bankId: Option[String] = Option(BankId.get).filterNot(_ == DYNAMIC_ENTITY_SYSTEM_LEVEL_BANK_ID) override def userId: Option[String] = Option(UserId.get) override def isPersonalEntity: Boolean = IsPersonalEntity.get + override def createdDate: Option[Date] = Option(createdAt.get) + override def updatedDate: Option[Date] = Option(updatedAt.get) + override def createdByUserId: Option[String] = Option(CreatedByUserId.get).filter(_.nonEmpty) + override def createdByOnBehalfOfUserId: Option[String] = Option(CreatedByOnBehalfOfUserId.get).filter(_.nonEmpty) + override def updatedByUserId: Option[String] = Option(UpdatedByUserId.get).filter(_.nonEmpty) + override def updatedByOnBehalfOfUserId: Option[String] = Option(UpdatedByOnBehalfOfUserId.get).filter(_.nonEmpty) } object DynamicData extends DynamicData with LongKeyedMetaMapper[DynamicData] { diff --git a/obp-api/src/main/scala/code/messageoutbox/MessageOutbox.scala b/obp-api/src/main/scala/code/messageoutbox/MessageOutbox.scala index 0828e106f5..a65a666d23 100644 --- a/obp-api/src/main/scala/code/messageoutbox/MessageOutbox.scala +++ b/obp-api/src/main/scala/code/messageoutbox/MessageOutbox.scala @@ -27,7 +27,25 @@ TESOBE (http://www.tesobe.com/) package code.messageoutbox +import code.api.util.CommonsEmailWrapper.EmailContent import net.liftweb.mapper._ +import org.json4s.native.Serialization +import org.json4s.{DefaultFormats, Formats} + +/** An email as stored in an EMAIL outbox row: everything needed to send it, rendered at enqueue time. */ +case class OutboxEmailPayload( + from: String, + to: List[String], + cc: List[String], + bcc: List[String], + subject: String, + text_content: Option[String], + html_content: Option[String] +) { + def toEmailContent: EmailContent = + EmailContent(from = from, to = to, cc = cc, bcc = bcc, subject = subject, + textContent = text_content, htmlContent = html_content) +} /** * Generic transactional outbox for asynchronous messages OBP-API must deliver. @@ -43,6 +61,11 @@ import net.liftweb.mapper._ * publish behavior to the relay. Types so far: * OPEN_CORRIDOR — Interface C messages to a bank's RabbitMQ vhost * (`target_id` = bank_id, publisher = OpenCorridorPublisher). + * EMAIL — an email, rendered at enqueue time (`payload_json` is an + * OutboxEmailPayload, `target_id` the recipients), sent over + * SMTP by the relay. Queued rather than sent in the request + * so that a rolled back request emails nobody, and so a burst + * of emails never blocks request threads. * * Row lifecycle: * PENDING — not yet delivered; the relay keeps publishing with backoff. @@ -136,12 +159,18 @@ object MessageOutbox extends MessageOutbox with LongKeyedMetaMapper[MessageOutbo val STATUS_STICKY = "STICKY" val TYPE_OPEN_CORRIDOR = "OPEN_CORRIDOR" + val TYPE_EMAIL = "EMAIL" + + val OPERATION_ROLE_GRANTED_EMAIL = "role_granted_email" // subject_id_type holds the OBP id-field name whose value space subject_id // belongs to (exact snake_case field name, e.g. transaction_request_id, // settlement_id, consent_id, customer_id ...). val SUBJECT_TYPE_SETTLEMENT_ID = "settlement_id" val SUBJECT_TYPE_TRANSACTION_REQUEST_ID = "transaction_request_id" + val SUBJECT_TYPE_ENTITLEMENT_ID = "entitlement_id" + + private val emailFormats: Formats = DefaultFormats override def dbTableName = "message_outbox" @@ -166,6 +195,33 @@ object MessageOutbox extends MessageOutbox with LongKeyedMetaMapper[MessageOutbo .Status(STATUS_PENDING) .saveMe() + /** Queue an email; the relay sends it (attachments are not supported). */ + def enqueueEmail(subjectId: String, subjectIdType: String, operationName: String, email: EmailContent): MessageOutbox = { + val payload = OutboxEmailPayload(email.from, email.to, email.cc, email.bcc, email.subject, + email.textContent, email.htmlContent) + enqueue(TYPE_EMAIL, subjectId, subjectIdType, operationName, + email.to.mkString(",").take(255), Serialization.write(payload)(emailFormats)) + } + + def emailPayload(row: MessageOutbox): OutboxEmailPayload = + Serialization.read[OutboxEmailPayload](row.payloadJson)(emailFormats, manifest[OutboxEmailPayload]) + + /** + * Claim a PENDING row for one delivery attempt: counts the attempt, but only if the row is still + * as `row` read it. Returns true for exactly one caller, so with several OBP-API instances each + * running the relay a message is sent by one of them only. It commits at once (outside a request); + * if the claimant dies before sending, the row is still PENDING and is retried after the backoff. + */ + def claimForDelivery(row: MessageOutbox): Boolean = { + import doobie.implicits._ + code.api.util.DoobieUtil.runUpdate( + sql"""UPDATE message_outbox + SET attempts = attempts + 1, updated_at = NOW() + WHERE id = ${row.id.get} + AND status = $STATUS_PENDING + AND attempts = ${row.attempts}""".update.run) == 1 + } + def pending(): List[MessageOutbox] = MessageOutbox.findAll(By(MessageOutbox.Status, STATUS_PENDING)) diff --git a/obp-api/src/main/scala/code/messageoutbox/MessageOutboxRelay.scala b/obp-api/src/main/scala/code/messageoutbox/MessageOutboxRelay.scala index 525c2ed6a7..974c14a632 100644 --- a/obp-api/src/main/scala/code/messageoutbox/MessageOutboxRelay.scala +++ b/obp-api/src/main/scala/code/messageoutbox/MessageOutboxRelay.scala @@ -28,6 +28,7 @@ TESOBE (http://www.tesobe.com/) package code.messageoutbox import code.actorsystem.ObpActorSystem +import code.api.util.CommonsEmailWrapper import code.bankconnectors.opencorridor.OpenCorridorPublisher import code.util.Helper.MdcLoggable import net.liftweb.common.{Box, Failure, Full} @@ -35,6 +36,7 @@ import org.json4s._ import org.json4s.native.Serialization import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicBoolean import scala.concurrent.Await import scala.concurrent.duration._ @@ -66,6 +68,11 @@ import scala.concurrent.duration._ * the bank's CBS refusing the credit itself (unknown account, name * mismatch) — the asynchronous beneficiary refusal, distinct from the * transient CBS-DELIVERY-FAILED. + * + * EMAIL: each row is first claimed (MessageOutbox.claimForDelivery), so with + * several instances only one sends it. Sent → DELIVERED. Not sent → stays + * PENDING and is retried with the same backoff, until maxEmailAttempts, then + * STICKY (an address or mail server an operator must look at). */ object MessageOutboxRelay extends MdcLoggable { @@ -76,6 +83,11 @@ object MessageOutboxRelay extends MdcLoggable { private val maxBackoff = 10.minutes /** Cap on how long one row's publish may block the (serial) relay pass. */ private val perRowTimeout = 60.seconds + /** Attempts after which an EMAIL row that could not be sent goes STICKY. */ + private val maxEmailAttempts = 8 + /** A pass can outlast the interval (many emails, or a slow publish); the scheduler would then + * start another over the same PENDING rows and deliver them twice. */ + private val passRunning = new AtomicBoolean(false) // OPEN_CORRIDOR errors retrying cannot fix. CBS-DELIVERY-FAILED is // deliberately NOT here: a CBS being down is transient, and with credit @@ -98,8 +110,11 @@ object MessageOutboxRelay extends MdcLoggable { interval = scala.concurrent.duration.Duration(intervalSeconds, TimeUnit.SECONDS), runnable = new Runnable { def run(): Unit = - try relayOnePass() - catch { case e: Throwable => logger.error("message outbox relay pass failed", e) } + if (passRunning.compareAndSet(false, true)) { + try relayOnePass() + catch { case e: Throwable => logger.error("message outbox relay pass failed", e) } + finally passRunning.set(false) + } } ) logger.info(s"message outbox relay started (interval ${intervalSeconds}s)") @@ -108,7 +123,11 @@ object MessageOutboxRelay extends MdcLoggable { /** One pass over the PENDING rows that are due (backoff by attempts). */ def relayOnePass(): Unit = { val now = System.currentTimeMillis() - val due = MessageOutbox.pending().filter { row => + // Open Corridor rows are only relayed where Open Corridor is enabled; emails always are. + val openCorridorEnabled = code.api.util.APIUtil.getPropsAsBoolValue("open_corridor_enabled", false) + val due = MessageOutbox.pending() + .filter(row => row.outboxType != MessageOutbox.TYPE_OPEN_CORRIDOR || openCorridorEnabled) + .filter { row => val backoff = (baseBackoff * math.pow(2, math.min(row.attempts, 6)).toLong).min(maxBackoff) row.UpdatedAt.get.getTime + backoff.toMillis <= now || row.attempts == 0 } @@ -118,12 +137,40 @@ object MessageOutboxRelay extends MdcLoggable { def relayRow(row: MessageOutbox): Unit = row.outboxType match { case MessageOutbox.TYPE_OPEN_CORRIDOR => relayOpenCorridorRow(row) + case MessageOutbox.TYPE_EMAIL => relayEmailRow(row) case other => row.Status(MessageOutbox.STATUS_STICKY).Attempts(row.attempts + 1) .LastError(s"no publisher registered for outbox_type '$other'").saveMe() logger.error(s"message outbox row ${row.id.get}: unknown outbox_type '$other' — STICKY") } + private def relayEmailRow(read: MessageOutbox): Unit = + // Another instance's relay may have claimed it first; then it is theirs to send. + if (MessageOutbox.claimForDelivery(read)) { + val row = MessageOutbox.find(net.liftweb.mapper.By(MessageOutbox.id, read.id.get)).openOrThrowException("claimed row") + val sent: Box[String] = + try { + val email = MessageOutbox.emailPayload(row).toEmailContent + if (email.htmlContent.isDefined) CommonsEmailWrapper.sendHtmlEmail(email) + else CommonsEmailWrapper.sendTextEmail(email) + } catch { case e: Throwable => Failure(s"email send failed: ${e.getMessage}") } + sent match { + case Full(messageId) => + row.Status(MessageOutbox.STATUS_DELIVERED).LastError("").LastReplyJson(messageId).saveMe() + logger.info(s"message outbox row ${row.id.get}: ${row.operationName} to ${row.targetId} DELIVERED") + case failure => + val error = failure match { + case Failure(msg, _, _) => msg + case _ => "not sent (see the email log)" + } + // The claim already counted this attempt. + val status = if (row.attempts >= maxEmailAttempts) MessageOutbox.STATUS_STICKY else MessageOutbox.STATUS_PENDING + row.Status(status).LastError(error.take(2000)).saveMe() + logger.warn(s"message outbox row ${row.id.get}: ${row.operationName} to ${row.targetId} " + + s"not sent (attempt ${row.attempts}, now $status): $error") + } + } + private def relayOpenCorridorRow(row: MessageOutbox): Unit = { val replyBox: Box[com.openbankproject.commons.dto.InBoundOpenCorridorReply] = try { diff --git a/obp-api/src/main/scala/code/metrics/APIMetrics.scala b/obp-api/src/main/scala/code/metrics/APIMetrics.scala index 6a50727b24..1177efa07b 100644 --- a/obp-api/src/main/scala/code/metrics/APIMetrics.scala +++ b/obp-api/src/main/scala/code/metrics/APIMetrics.scala @@ -125,6 +125,7 @@ trait APIMetrics { responseBody: String, sourceIp: String, targetIp: String, + forwardedFor: String, apiInstanceId: String, consentReferenceId: String, certificateTrust: String, @@ -148,6 +149,7 @@ trait APIMetrics { responseBody: String, sourceIp: String, targetIp: String, + forwardedFor: String, apiInstanceId: String, consentReferenceId: String, certificateTrust: String, @@ -208,6 +210,7 @@ trait APIMetric { def getResponseBody(): String def getSourceIp(): String def getTargetIp(): String + def getForwardedFor(): String def getApiInstanceId(): String def getConsentReferenceId(): String def getCertificateTrust(): String diff --git a/obp-api/src/main/scala/code/metrics/ElasticsearchMetrics.scala b/obp-api/src/main/scala/code/metrics/ElasticsearchMetrics.scala index 3ceb187362..a16466fd90 100644 --- a/obp-api/src/main/scala/code/metrics/ElasticsearchMetrics.scala +++ b/obp-api/src/main/scala/code/metrics/ElasticsearchMetrics.scala @@ -41,7 +41,7 @@ object ElasticsearchMetrics extends APIMetrics { lazy val es = new elasticsearchMetrics override def saveMetric(userId: String, url: String, date: Date, duration: Long, userName: String, appName: String, developerEmail: String, consumerId: String, implementedByPartialFunction: String, implementedInVersion: String, verb: String, httpCode: Option[Int], correlationId: String, - responseBody: String, sourceIp: String, targetIp: String, apiInstanceId: String, consentReferenceId: String, + responseBody: String, sourceIp: String, targetIp: String, forwardedFor: String, apiInstanceId: String, consentReferenceId: String, certificateTrust: String, certificateTrustDetail: String, authType: String): Unit = { if (APIUtil.getPropsAsBoolValue("allow_elasticsearch", false) && APIUtil.getPropsAsBoolValue("allow_elasticsearch_metrics", false) ) { @@ -53,6 +53,7 @@ object ElasticsearchMetrics extends APIMetrics { responseBody: String, sourceIp: String, targetIp: String, + forwardedFor: String, apiInstanceId: String, consentReferenceId: String, certificateTrust: String, diff --git a/obp-api/src/main/scala/code/metrics/MappedMetrics.scala b/obp-api/src/main/scala/code/metrics/MappedMetrics.scala index 12288848c9..b156e2ed56 100644 --- a/obp-api/src/main/scala/code/metrics/MappedMetrics.scala +++ b/obp-api/src/main/scala/code/metrics/MappedMetrics.scala @@ -140,7 +140,7 @@ object MappedMetrics extends APIMetrics with MdcLoggable{ } override def saveMetric(userId: String, url: String, date: Date, duration: Long, userName: String, appName: String, developerEmail: String, consumerId: String, implementedByPartialFunction: String, implementedInVersion: String, verb: String, httpCode: Option[Int], correlationId: String, - responseBody: String, sourceIp: String, targetIp: String, apiInstanceId: String, consentReferenceId: String, + responseBody: String, sourceIp: String, targetIp: String, forwardedFor: String, apiInstanceId: String, consentReferenceId: String, certificateTrust: String, certificateTrustDetail: String, authType: String): Unit = { // A correlation id is expected on every metric. Rows without one cannot be moved @@ -167,6 +167,7 @@ object MappedMetrics extends APIMetrics with MdcLoggable{ responseBody = responseBody, sourceIp = sourceIp, targetIp = targetIp, + forwardedFor = forwardedFor, apiInstanceId = apiInstanceId, consentReferenceId = consentReferenceId, certificateTrust = certificateTrust, @@ -180,7 +181,7 @@ object MappedMetrics extends APIMetrics with MdcLoggable{ appName: String, developerEmail: String, consumerId: String, implementedByPartialFunction: String, implementedInVersion: String, verb: String, httpCode: Option[Int], correlationId: String, - responseBody: String, sourceIp: String, targetIp: String, + responseBody: String, sourceIp: String, targetIp: String, forwardedFor: String, apiInstanceId: String, consentReferenceId: String, certificateTrust: String, certificateTrustDetail: String, authType: String): Boolean = { @@ -207,6 +208,7 @@ object MappedMetrics extends APIMetrics with MdcLoggable{ .responseBody(responseBody) .sourceIp(sourceIp) .targetIp(targetIp) + .forwardedFor(forwardedFor) .apiInstanceId(apiInstanceId) .consentReferenceId(consentReferenceId) .certificateTrust(certificateTrust) @@ -841,8 +843,17 @@ class MappedMetric extends APIMetric with LongKeyedMapper[MappedMetric] with IdP override def defaultValue = generateUUID() } object responseBody extends MappedText(this) + // The client address OBP-API decided on (CallContext.ipAddress, see RemoteIpUtil). Rows + // written before that change hold the raw X-Forwarded-For header instead. object sourceIp extends MappedString(this, 64) object targetIp extends MappedString(this, 64) + // The hops the request passed through: the X-Forwarded-For chain it arrived with, followed by + // the TCP peer. Client-controlled on the left, so it is cut to this width on write. The archive + // copy (MetricArchive.forwardedFor) MUST keep the same width. + object forwardedFor extends MappedString(this, 512) { + override def dbColumnName = "forwarded_for" + override def defaultValue = null + } object apiInstanceId extends MappedString(this, 255) // Set when the request was authenticated via a consent. Null otherwise. object consentReferenceId extends MappedString(this, 36) { @@ -887,6 +898,7 @@ class MappedMetric extends APIMetric with LongKeyedMapper[MappedMetric] with IdP override def getResponseBody(): String = responseBody.get override def getSourceIp(): String = sourceIp.get override def getTargetIp(): String = targetIp.get + override def getForwardedFor(): String = forwardedFor.get override def getApiInstanceId(): String = apiInstanceId.get override def getConsentReferenceId(): String = consentReferenceId.get override def getCertificateTrust(): String = certificateTrust.get @@ -946,8 +958,17 @@ class MetricArchive extends APIMetric with LongKeyedMapper[MetricArchive] with I override def dbNotNull_? = true } object responseBody extends MappedText(this) + // The client address OBP-API decided on (CallContext.ipAddress, see RemoteIpUtil). Rows + // written before that change hold the raw X-Forwarded-For header instead. object sourceIp extends MappedString(this, 64) object targetIp extends MappedString(this, 64) + // The hops the request passed through: the X-Forwarded-For chain it arrived with, followed by + // the TCP peer. Client-controlled on the left, so it is cut to this width on write. The archive + // copy (MetricArchive.forwardedFor) MUST keep the same width. + object forwardedFor extends MappedString(this, 512) { + override def dbColumnName = "forwarded_for" + override def defaultValue = null + } object apiInstanceId extends MappedString(this, 255) // Set when the request was authenticated via a consent. Null otherwise. object consentReferenceId extends MappedString(this, 36) { @@ -990,6 +1011,7 @@ class MetricArchive extends APIMetric with LongKeyedMapper[MetricArchive] with I override def getResponseBody(): String = responseBody.get override def getSourceIp(): String = sourceIp.get override def getTargetIp(): String = targetIp.get + override def getForwardedFor(): String = forwardedFor.get override def getApiInstanceId(): String = apiInstanceId.get override def getConsentReferenceId(): String = consentReferenceId.get override def getCertificateTrust(): String = certificateTrust.get diff --git a/obp-api/src/main/scala/code/metrics/MetricBatchWriter.scala b/obp-api/src/main/scala/code/metrics/MetricBatchWriter.scala index 0f27ea4621..8396f0709b 100644 --- a/obp-api/src/main/scala/code/metrics/MetricBatchWriter.scala +++ b/obp-api/src/main/scala/code/metrics/MetricBatchWriter.scala @@ -67,6 +67,7 @@ object MetricBatchWriter extends MdcLoggable { responseBody: String, sourceIp: String, targetIp: String, + forwardedFor: String, apiInstanceId: String, consentReferenceId: String, certificateTrust: String, @@ -108,10 +109,46 @@ object MetricBatchWriter extends MdcLoggable { private lazy val telemetry = new code.telemetry.BatchWriterTelemetry("api_metrics") def enqueue(row: MetricRow): Unit = { - queue.add(row) + queue.add(fitToColumns(row)) telemetry.queued() } + /** + * This cuts each text value of a row to the width of its column in the metric table. + * + * The rows of one flush are inserted as a single batch, and the database rejects a value that + * is longer than its column ("value too long for type character varying(N)"). One such value + * would therefore lose every metric in the flush, not only its own. Several values come from + * the caller (X-Forwarded-For, X-Forwarded-Host, the correlation id, the URL), so without this + * any caller could make everyone's metrics disappear with one long header. The widths are read + * from MappedMetric, so they cannot drift from the table definition. + */ + private[metrics] def fitToColumns(row: MetricRow): MetricRow = { + def fit(value: String, column: net.liftweb.mapper.MappedString[MappedMetric]): String = + if (value != null && value.length > column.maxLen) value.substring(0, column.maxLen) else value + val table = MappedMetric + row.copy( + userId = fit(row.userId, table.userId), + url = fit(row.url, table.url), + userName = fit(row.userName, table.userName), + appName = fit(row.appName, table.appName), + developerEmail = fit(row.developerEmail, table.developerEmail), + consumerId = fit(row.consumerId, table.consumerId), + implementedByPartialFunction = fit(row.implementedByPartialFunction, table.implementedByPartialFunction), + implementedInVersion = fit(row.implementedInVersion, table.implementedInVersion), + verb = fit(row.verb, table.verb), + correlationId = fit(row.correlationId, table.correlationId), + sourceIp = fit(row.sourceIp, table.sourceIp), + targetIp = fit(row.targetIp, table.targetIp), + forwardedFor = fit(row.forwardedFor, table.forwardedFor), + apiInstanceId = fit(row.apiInstanceId, table.apiInstanceId), + consentReferenceId = fit(row.consentReferenceId, table.consentReferenceId), + certificateTrust = fit(row.certificateTrust, table.certificateTrust), + certificateTrustDetail = fit(row.certificateTrustDetail, table.certificateTrustDetail), + authType = fit(row.authType, table.authType) + ) + } + /** * Drain the queue and batch-insert all pending metrics via Doobie. */ @@ -141,9 +178,9 @@ object MetricBatchWriter extends MdcLoggable { userid, url, date_c, duration, username, appname, developeremail, consumerid, implementedbypartialfunction, implementedinversion, verb, httpcode, correlationid, - responsebody, sourceip, targetip, apiinstanceid, consent_reference_id, + responsebody, sourceip, targetip, forwarded_for, apiinstanceid, consent_reference_id, certificate_trust, certificate_trust_detail, auth_type - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """ // Use Option[String] so Doobie handles nullable fields via Put[Option[String]] @@ -152,7 +189,7 @@ object MetricBatchWriter extends MdcLoggable { (Option[String], Option[String], Timestamp, Long, Option[String], Option[String], Option[String], Option[String], Option[String], Option[String], Option[String], Int, Option[String], - Option[String], Option[String], Option[String], Option[String], Option[String], + Option[String], Option[String], Option[String], Option[String], Option[String], Option[String], Option[String], Option[String], Option[String]) ](insertSql) @@ -162,7 +199,7 @@ object MetricBatchWriter extends MdcLoggable { r.duration, Option(r.userName), Option(r.appName), Option(r.developerEmail), Option(r.consumerId), Option(r.implementedByPartialFunction), Option(r.implementedInVersion), Option(r.verb), r.httpCode, Option(r.correlationId), - Option(r.responseBody), Option(r.sourceIp), Option(r.targetIp), Option(r.apiInstanceId), + Option(r.responseBody), Option(r.sourceIp), Option(r.targetIp), Option(r.forwardedFor), Option(r.apiInstanceId), Option(r.consentReferenceId), Option(r.certificateTrust), Option(r.certificateTrustDetail), Option(r.authType) diff --git a/obp-api/src/main/scala/code/scheduler/MetricsArchiveScheduler.scala b/obp-api/src/main/scala/code/scheduler/MetricsArchiveScheduler.scala index 38614d31ac..745c203070 100644 --- a/obp-api/src/main/scala/code/scheduler/MetricsArchiveScheduler.scala +++ b/obp-api/src/main/scala/code/scheduler/MetricsArchiveScheduler.scala @@ -251,6 +251,7 @@ object MetricsArchiveScheduler extends MdcLoggable { i.getResponseBody(), i.getSourceIp(), i.getTargetIp(), + i.getForwardedFor(), i.getApiInstanceId(), i.getConsentReferenceId(), i.getCertificateTrust(), diff --git a/obp-api/src/main/scala/code/users/LiftUsers.scala b/obp-api/src/main/scala/code/users/LiftUsers.scala index 8c32e7a698..c1db0f5b61 100644 --- a/obp-api/src/main/scala/code/users/LiftUsers.scala +++ b/obp-api/src/main/scala/code/users/LiftUsers.scala @@ -171,6 +171,9 @@ object LiftUsers extends Users with MdcLoggable{ * caller reference columns written * -------------------------------- ------------------ ----------------------------- * MappedTransactionRequestProvider TransactionRequest mUserId + mOnBehalfOfUserId + * MapperCounterparties Counterparty mCreatedByUserId + mCreatedByOnBehalfOfUserId + * MapppedDynamicDataProvider DynamicData_Created CreatedByUserId + CreatedByOnBehalfOfUserId + * MapppedDynamicDataProvider DynamicData_Updated UpdatedByUserId + UpdatedByOnBehalfOfUserId * * C. One column, two policies. The reference is chosen per process, then passed in. * diff --git a/obp-api/src/main/scala/code/users/UserReference.scala b/obp-api/src/main/scala/code/users/UserReference.scala index 9eb864035f..ea8c9140b2 100644 --- a/obp-api/src/main/scala/code/users/UserReference.scala +++ b/obp-api/src/main/scala/code/users/UserReference.scala @@ -186,6 +186,8 @@ object UserReference { case object UserAuthContextUpdate_UserId extends UserReference(UseOnBehalfOfUserId , "code.context.MappedUserAuthContextUpdate", List("mUserId"), "as UserAuthContext_UserId") case object DynamicEntity_UserId extends UserReference(UseOnBehalfOfUserId , "code.dynamicEntity.DynamicEntity", List("UserId"), "the definition's creator; a definition outlives the Consent that created it") case object DynamicData_UserId extends UserReference(UseOnBehalfOfUserId , "code.DynamicData.DynamicData", List("UserId"), "personal rows, and the one reference where the redirect MUST be symmetric: MapppedDynamicDataProvider resolves on save/update/get/delete alike, because a row keyed by this column on both sides is otherwise written by an agent and then invisible to it") + case object DynamicData_CreatedByUserId extends UserReference(UseOnBehalfOfUserId , "code.DynamicData.DynamicData", List("CreatedByUserId", "CreatedByOnBehalfOfUserId"), "record both: who made the call that created the record, agent included, and the person it was made for; published in the v7.0.0 record metadata") + case object DynamicData_UpdatedByUserId extends UserReference(UseOnBehalfOfUserId , "code.DynamicData.DynamicData", List("UpdatedByUserId", "UpdatedByOnBehalfOfUserId"), "record both: as DynamicData_CreatedByUserId, for the call that last saved the record") case object DynamicDataAccess_UserId extends UserReference(UseOnBehalfOfUserId , "code.DynamicData.DynamicDataAccess", List("UserId"), "row-level ACL. Deliberately still on the consent user today, because the bootstrap grant and the allows check have to agree with each other -- rows strand, nothing leaks. A later Phase 2 row; see the plan") case object DynamicEndpoint_UserId extends UserReference(UseOnBehalfOfUserId , "code.DynamicEndpoint.DynamicEndpoint", List("UserId"), "the dynamic endpoint's creator; outlives the Consent") case object DynamicResourceDoc_CreatedByUserId extends UserReference(UseOnBehalfOfUserId , "code.dynamicResourceDoc.DynamicResourceDoc", List("CreatedByUserId", "UpdatedByUserId"), "and UpdatedByUserId; a dynamic artefact outlives the Consent that created it") @@ -273,6 +275,8 @@ object UserReference { UserAuthContextUpdate_UserId, DynamicEntity_UserId, DynamicData_UserId, + DynamicData_CreatedByUserId, + DynamicData_UpdatedByUserId, DynamicDataAccess_UserId, DynamicEndpoint_UserId, DynamicResourceDoc_CreatedByUserId, diff --git a/obp-api/src/test/scala/code/api/util/AgentDelegationTest.scala b/obp-api/src/test/scala/code/api/util/AgentDelegationTest.scala index 39666e1300..fb2fdbffce 100644 --- a/obp-api/src/test/scala/code/api/util/AgentDelegationTest.scala +++ b/obp-api/src/test/scala/code/api/util/AgentDelegationTest.scala @@ -783,4 +783,56 @@ class AgentDelegationTest extends ServerSetup { storedField(row.mCreatedByOnBehalfOfUserId.get) shouldBe agent2.userId } } + + feature("Dynamic Entity records record both ids (UserReference.DynamicData_CreatedByUserId and DynamicData_UpdatedByUserId)") { + + val bankId = Some("agent-delegation-dynamic-data-bank") + + /** A fresh entity name, so no two scenarios share records. */ + def newEntityName(): String = s"agent_delegation_${generateUUID().take(8).replace("-", "")}" + + def saveAs(callerUserId: String, entityName: String, recordId: String): code.DynamicData.DynamicDataT = + code.DynamicData.DynamicDataProvider.connectorMethodProvider.vend + .save(bankId, entityName, (s"${entityName}_id" -> recordId) ~ ("name" -> "first"), Some(callerUserId), false) + .openOrThrowException("expected the record to be saved") + + def updateAs(callerUserId: String, entityName: String, recordId: String): code.DynamicData.DynamicDataT = + code.DynamicData.DynamicDataProvider.connectorMethodProvider.vend + .update(bankId, entityName, (s"${entityName}_id" -> recordId) ~ ("name" -> "second"), recordId, Some(callerUserId), false) + .openOrThrowException("expected the record to be updated") + + scenario("a record created by an original user names that user in all four columns", AgentDelegationTag) { + val human = createUser() + val record = saveAs(human.userId, newEntityName(), "R1") + record.createdByUserId shouldBe Some(human.userId) + record.createdByOnBehalfOfUserId shouldBe Some(human.userId) + record.updatedByUserId shouldBe Some(human.userId) + record.updatedByOnBehalfOfUserId shouldBe Some(human.userId) + } + + scenario("a record created by a consent user names the agent and its on-behalf-of user", AgentDelegationTag) { + val human = createUser() + val consent = MappedConsent.create.mUserId(human.userId).saveMe() + val agent = createUser(createdByConsentId = Some(consent.consentId)) + val record = saveAs(agent.userId, newEntityName(), "R1") + record.createdByUserId shouldBe Some(agent.userId) + record.createdByOnBehalfOfUserId shouldBe Some(human.userId) + record.updatedByUserId shouldBe Some(agent.userId) + record.updatedByOnBehalfOfUserId shouldBe Some(human.userId) + } + + scenario("an update changes only the updated pair, and the created pair keeps the original writer", AgentDelegationTag) { + val human = createUser() + val consent = MappedConsent.create.mUserId(human.userId).saveMe() + val agent = createUser(createdByConsentId = Some(consent.consentId)) + val entityName = newEntityName() + saveAs(agent.userId, entityName, "R1") + + val updated = updateAs(human.userId, entityName, "R1") + updated.createdByUserId shouldBe Some(agent.userId) + updated.createdByOnBehalfOfUserId shouldBe Some(human.userId) + updated.updatedByUserId shouldBe Some(human.userId) + updated.updatedByOnBehalfOfUserId shouldBe Some(human.userId) + } + } } diff --git a/obp-api/src/test/scala/code/api/util/RemoteIpUtilTest.scala b/obp-api/src/test/scala/code/api/util/RemoteIpUtilTest.scala index b77667c7d8..c0b06bf69e 100644 --- a/obp-api/src/test/scala/code/api/util/RemoteIpUtilTest.scala +++ b/obp-api/src/test/scala/code/api/util/RemoteIpUtilTest.scala @@ -40,7 +40,34 @@ class RemoteIpUtilTest extends ServerSetup { direct.headerHonoured shouldBe false } - scenario("X-Forwarded-For gives its leftmost address, and IPv6 comes back without brackets, in canonical form") { + scenario("X-Forwarded-For is read from the right: the first address not in trust.proxy.peers is the client") { + // NGINX 10.0.0.2, API Explorer II 10.0.0.3, Opey 10.0.0.4 and OBP-MCP 10.0.0.5 are trusted. + setPropsValues("trust.proxy.enabled" -> "true", "trust.proxy.header" -> "X-Forwarded-For", "trust.proxy.peers" -> "10.0.0.0/24") + val chain = "6.6.6.6, 203.0.113.9, 10.0.0.2, 10.0.0.3, 10.0.0.4" + val r = RemoteIpUtil.resolve("10.0.0.5", headers("X-Forwarded-For" -> chain)) + r.clientIp shouldBe "203.0.113.9" + r.headerHonoured shouldBe true + } + + scenario("X-Forwarded-For: a caller that is not trusted becomes the client, whatever it wrote to its left") { + setPropsValues("trust.proxy.enabled" -> "true", "trust.proxy.header" -> "X-Forwarded-For", "trust.proxy.peers" -> "10.0.0.0/24") + // An MCP client at 198.51.100.7 asks OBP-MCP (trusted) to pass on a made-up chain; OBP-MCP appends the caller. + RemoteIpUtil.resolve("10.0.0.5", headers("X-Forwarded-For" -> "6.6.6.6, 10.0.0.2, 198.51.100.7")).clientIp shouldBe "198.51.100.7" + } + + scenario("X-Forwarded-For: an entry that is not an address stops the walk at the nearest trusted hop") { + setPropsValues("trust.proxy.enabled" -> "true", "trust.proxy.header" -> "X-Forwarded-For", "trust.proxy.peers" -> "10.0.0.0/24") + RemoteIpUtil.resolve("10.0.0.5", headers("X-Forwarded-For" -> "203.0.113.9, unknown, 10.0.0.2")).clientIp shouldBe "10.0.0.2" + RemoteIpUtil.resolve("10.0.0.5", headers("X-Forwarded-For" -> "203.0.113.9, unknown")).clientIp shouldBe "10.0.0.5" + } + + scenario("X-Forwarded-For: when every entry is trusted, the leftmost is the client; IPv6 entries are canonicalised") { + setPropsValues("trust.proxy.enabled" -> "true", "trust.proxy.header" -> "X-Forwarded-For", "trust.proxy.peers" -> "10.0.0.0/24, 2001:db8::/32") + RemoteIpUtil.resolve("10.0.0.5", headers("X-Forwarded-For" -> "10.0.0.9, 10.0.0.2")).clientIp shouldBe "10.0.0.9" + RemoteIpUtil.resolve("10.0.0.5", headers("X-Forwarded-For" -> "2001:DB8:0:0:0:0:0:1, [2001:db8::2]")).clientIp shouldBe "2001:db8::1" + } + + scenario("X-Forwarded-For with no peer list gives its leftmost address, and IPv6 comes back without brackets, in canonical form") { setPropsValues("trust.proxy.enabled" -> "true", "trust.proxy.header" -> "X-Forwarded-For", "trust.proxy.peers" -> "") RemoteIpUtil.resolve("10.0.0.5", headers("X-Forwarded-For" -> "203.0.113.9, 10.0.0.1")).clientIp shouldBe "203.0.113.9" setPropsValues("trust.proxy.enabled" -> "false") diff --git a/obp-api/src/test/scala/code/api/v6_0_0/DynamicDataTimestampsTest.scala b/obp-api/src/test/scala/code/api/v6_0_0/DynamicDataTimestampsTest.scala new file mode 100644 index 0000000000..8d248b111e --- /dev/null +++ b/obp-api/src/test/scala/code/api/v6_0_0/DynamicDataTimestampsTest.scala @@ -0,0 +1,128 @@ +/** +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.v6_0_0 + +import code.DynamicData.DynamicDataProvider +import code.api.util.APIUtil +import code.setup.ServerSetup +import net.liftweb.db.DB +import net.liftweb.util.DefaultConnectionIdentifier +import org.json4s.JsonAST.JObject +import org.json4s.JsonDSL._ +import org.scalatest.Tag + +/** + * This suite covers the created and updated timestamps on a Dynamic Entity record. + * + * A record used to carry no timestamps at all, so nobody could tell when it was written or last + * changed. That matters wherever records are someone's own data, for example a farmer's soil samples, + * because the provenance of such a record has to be stated. The records table now has the columns + * createdat and updatedat. A record written before those columns existed holds NULL in both, and + * the provider reports that as None instead of inventing a time. + */ +class DynamicDataTimestampsTest extends ServerSetup { + + object DynamicDataTimestamps extends Tag("DynamicDataTimestamps") + + private def dataProvider = DynamicDataProvider.connectorMethodProvider.vend + + /** The id field name MappedDynamicDataProvider.getIdName derives from the entity name. */ + private def idFieldNameOf(entityName: String): String = + s"${entityName}_Id".replaceAll("(?<=[a-z0-9])(?=[A-Z])|-", "_").toLowerCase + + private def bodyWithId(entityName: String, id: String, name: String): JObject = + (idFieldNameOf(entityName) -> id) ~ ("name" -> name) + + private def uniqueEntityName(prefix: String): String = s"${prefix}${APIUtil.generateUUID().take(8).replace("-", "")}" + + feature("A Dynamic Entity record carries created and updated timestamps") { + + scenario("a new record has both timestamps", DynamicDataTimestamps) { + val entityName = uniqueEntityName("TimestampCreate") + val before = System.currentTimeMillis() + val saved = dataProvider.save(Some("bank_timestamps"), entityName, bodyWithId(entityName, "R1", "first"), None, false) + .openOrThrowException("the record should save") + + saved.createdDate.isDefined should equal(true) + saved.updatedDate.isDefined should equal(true) + // A column may store whole seconds only, so allow a second of rounding either way. + saved.createdDate.get.getTime should be >= (before - 1000) + + And("the timestamps read back from the database") + val readBack = dataProvider.get(Some("bank_timestamps"), entityName, "R1", None, false) + .openOrThrowException("the record should be found") + readBack.createdDate.isDefined should equal(true) + readBack.updatedDate.isDefined should equal(true) + } + + scenario("an update moves the updated timestamp and keeps the created one", DynamicDataTimestamps) { + val entityName = uniqueEntityName("TimestampUpdate") + val saved = dataProvider.save(Some("bank_timestamps"), entityName, bodyWithId(entityName, "R2", "first"), None, false) + .openOrThrowException("the record should save") + val createdOnSave = saved.createdDate.get + + // Longer than one second, so the change shows even where the column stores whole seconds. + Thread.sleep(1100) + dataProvider.update(Some("bank_timestamps"), entityName, bodyWithId(entityName, "R2", "second"), "R2", None, false) + .openOrThrowException("the record should update") + + val readBack = dataProvider.get(Some("bank_timestamps"), entityName, "R2", None, false) + .openOrThrowException("the record should be found") + readBack.createdDate.map(_.getTime / 1000) should equal(Some(createdOnSave.getTime / 1000)) + readBack.updatedDate.get.after(readBack.createdDate.get) should equal(true) + } + + scenario("a record written before the columns existed reports no timestamps", DynamicDataTimestamps) { + val entityName = uniqueEntityName("TimestampLegacy") + dataProvider.save(Some("bank_timestamps"), entityName, bodyWithId(entityName, "R3", "first"), None, false) + .openOrThrowException("the record should save") + + Given("a row whose timestamp columns are NULL, as on an instance upgraded from before they existed") + DB.use(DefaultConnectionIdentifier) { connection => + val statement = connection.prepareStatement( + "UPDATE dynamicdata SET createdat = NULL, updatedat = NULL WHERE dynamicentityname = ?") + try { statement.setString(1, entityName); statement.executeUpdate() } finally statement.close() + } + + Then("both timestamps read as None") + val legacy = dataProvider.get(Some("bank_timestamps"), entityName, "R3", None, false) + .openOrThrowException("the record should be found") + legacy.createdDate should equal(None) + legacy.updatedDate should equal(None) + + When("the record is updated") + dataProvider.update(Some("bank_timestamps"), entityName, bodyWithId(entityName, "R3", "second"), "R3", None, false) + .openOrThrowException("the record should update") + + Then("only the updated timestamp is filled, because the creation time was never recorded") + val updated = dataProvider.get(Some("bank_timestamps"), entityName, "R3", None, false) + .openOrThrowException("the record should be found") + updated.createdDate should equal(None) + updated.updatedDate.isDefined should equal(true) + } + } +} diff --git a/obp-api/src/test/scala/code/api/v7_0_0/DynamicEntityDefinitionTest.scala b/obp-api/src/test/scala/code/api/v7_0_0/DynamicEntityDefinitionTest.scala index 2c6cf0bc48..6bf83ae37a 100644 --- a/obp-api/src/test/scala/code/api/v7_0_0/DynamicEntityDefinitionTest.scala +++ b/obp-api/src/test/scala/code/api/v7_0_0/DynamicEntityDefinitionTest.scala @@ -253,6 +253,120 @@ class DynamicEntityDefinitionTest extends ServerSetupWithTestData { } } + feature("The metadata of a Dynamic Entity record at /obp/v7.0.0/banks/BANK_ID/dynamic-entities/...") { + + def metadataOf(record: JValue, event: String): JValue = record \ "metadata" \ event + + scenario("a record carries when it was created and updated, and by whom", VersionOfApi) { + val entityName = newEntityName() + val dynamicEntityId = createdWithFlags(SYS, entityName) + try { + When("user1 creates a record and reads it back") + val created = makePostRequest((dataAt(SYS) / entityName).POST <@ (user1), write(("name" -> "first"): JValue)) + created.code should equal(201) + val recordId = idOf(created, entityName) + val one = makeGetRequest((dataAt(SYS) / entityName / recordId).GET <@ (user1)) + one.code should equal(200) + + Then("the record is followed by its metadata, naming user1 for both the call and whom it was made for") + (one.body \ entityName \ "name").extract[String] should equal("first") + for (event <- List("created", "updated")) { + (metadataOf(one.body, event) \ "at").extract[String] should fullyMatch regex """\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z""" + (metadataOf(one.body, event) \ "user_id").extract[String] should equal(resourceUser1.userId) + (metadataOf(one.body, event) \ "on_behalf_of_user_id").extract[String] should equal(resourceUser1.userId) + } + And("the metadata is beside the record, never inside it") + (one.body \ entityName \ "metadata") should equal(JNothing) + (created.body \ "metadata" \ "created" \ "user_id").extract[String] should equal(resourceUser1.userId) + + When("the records are listed") + val listed = makeGetRequest((dataAt(SYS) / entityName).GET <@ (user1)) + listed.code should equal(200) + Then("each item is the record under the entity's name, with its metadata beside it") + val item = (listed.body \ s"${entityName}_list").extract[List[JObject]] + .find(o => (o \ entityName \ s"${entityName}_id").extract[String] == recordId) + .getOrElse(fail(s"record $recordId was not listed")) + (metadataOf(item, "created") \ "user_id").extract[String] should equal(resourceUser1.userId) + (metadataOf(item, "updated") \ "on_behalf_of_user_id").extract[String] should equal(resourceUser1.userId) + } finally cascadeDelete(SYS, dynamicEntityId) + } + + scenario("an update moves the updated time and keeps the created time", VersionOfApi) { + val entityName = newEntityName() + val dynamicEntityId = createdWithFlags(SYS, entityName) + try { + val created = makePostRequest((dataAt(SYS) / entityName).POST <@ (user1), write(("name" -> "first"): JValue)) + created.code should equal(201) + val recordId = idOf(created, entityName) + val createdAt = (metadataOf(created.body, "created") \ "at").extract[String] + + val replaced = makePutRequest((dataAt(SYS) / entityName / recordId).PUT <@ (user1), write(("name" -> "second"): JValue)) + replaced.code should equal(200) + (metadataOf(replaced.body, "created") \ "at").extract[String] should equal(createdAt) + (metadataOf(replaced.body, "updated") \ "at").extract[String] should be >= createdAt + } finally cascadeDelete(SYS, dynamicEntityId) + } + + scenario("a public read shows the times only, not who wrote the record", VersionOfApi) { + val entityName = newEntityName() + val dynamicEntityId = createdWithFlags(SYS, entityName, "has_public_access" -> true) + try { + val created = makePostRequest((dataAt(SYS) / entityName).POST <@ (user1), write(("name" -> "open"): JValue)) + created.code should equal(201) + val recordId = idOf(created, entityName) + + val public = makeGetRequest((dataAt(SYS) / "public" / entityName / recordId).GET) + public.code should equal(200) + (metadataOf(public.body, "created") \ "at").extract[String] should not be empty + (metadataOf(public.body, "created") \ "user_id") should equal(JNothing) + (metadataOf(public.body, "updated") \ "on_behalf_of_user_id") should equal(JNothing) + + val publicList = makeGetRequest((dataAt(SYS) / "public" / entityName).GET) + publicList.code should equal(200) + val items = (publicList.body \ s"${entityName}_list").extract[List[JObject]] + items should not be empty + all(items.map(item => metadataOf(item, "created") \ "user_id")) should equal(JNothing) + } finally cascadeDelete(SYS, dynamicEntityId) + } + + scenario("the unversioned URLs return the record alone, with plain list items", VersionOfApi) { + val entityName = newEntityName() + val dynamicEntityId = createdWithFlags(SYS, entityName) + try { + val created = makePostRequest((dataAt(SYS) / entityName).POST <@ (user1), write(("name" -> "first"): JValue)) + created.code should equal(201) + val recordId = idOf(created, entityName) + + val one = makeGetRequest((dynamicEntityData / entityName / recordId).GET <@ (user1)) + one.code should equal(200) + (one.body \ "metadata") should equal(JNothing) + + val listed = makeGetRequest((dynamicEntityData / entityName).GET <@ (user1)) + listed.code should equal(200) + (listed.body \ s"${entityName}_list").extract[List[JObject]].map(o => (o \ s"${entityName}_id").extract[String]) should contain(recordId) + } finally cascadeDelete(SYS, dynamicEntityId) + } + + scenario("the OBPv7.0.0 listing documents the metadata", VersionOfApi) { + val entityName = newEntityName() + val dynamicEntityId = createdWithFlags(SYS, entityName, "has_public_access" -> true) + try { + val functions = List(s"dynamicEntity_getSingle${entityName}_", s"dynamicEntity_getSinglePublic${entityName}_") + val docsResponse = makeGetRequest((v7 / "resource-docs" / "OBPv7.0.0" / "obp") < (d \ "operation_id").extract[String] == id).getOrElse(fail(s"no resource doc $id")) + + val getOne = docWithId(s"OBPv7.0.0-dynamicEntity_getSingle${entityName}_") + (getOne \ "success_response_body" \ "metadata" \ "created" \ "on_behalf_of_user_id") should not equal(JNothing) + val publicGetOne = docWithId(s"OBPv7.0.0-dynamicEntity_getSinglePublic${entityName}_") + (publicGetOne \ "success_response_body" \ "metadata" \ "created" \ "at") should not equal(JNothing) + (publicGetOne \ "success_response_body" \ "metadata" \ "created" \ "user_id") should equal(JNothing) + } finally cascadeDelete(SYS, dynamicEntityId) + } + } + private def dataAt(bankId: String) = v7 / "banks" / bankId / "dynamic-entities" private def flaggedDefinition(entityName: String, flags: (String, Boolean)*): JValue = @@ -287,7 +401,7 @@ class DynamicEntityDefinitionTest extends ServerSetupWithTestData { val listed = makeGetRequest((dataAt(SYS) / entityName).GET <@ (user1)) listed.code should equal(200) (listed.body \ "bank_id").extract[String] should equal(SYS) - (listed.body \ s"${entityName}_list").extract[List[JObject]].map(o => (o \ s"${entityName}_id").extract[String]) should contain(recordId) + (listed.body \ s"${entityName}_list").extract[List[JObject]].map(o => (o \ entityName \ s"${entityName}_id").extract[String]) should contain(recordId) val one = makeGetRequest((dataAt(SYS) / entityName / recordId).GET <@ (user1)) one.code should equal(200) @@ -321,7 +435,7 @@ class DynamicEntityDefinitionTest extends ServerSetupWithTestData { (personal.body \ "bank_id").extract[String] should equal(SYS) val myList = makeGetRequest((dataAt(SYS) / "my" / entityName).GET <@ (user1)) myList.code should equal(200) - (myList.body \ s"${entityName}_list").extract[List[JObject]].map(o => (o \ s"${entityName}_id").extract[String]) should contain(idOf(personal, entityName)) + (myList.body \ s"${entityName}_list").extract[List[JObject]].map(o => (o \ entityName \ s"${entityName}_id").extract[String]) should contain(idOf(personal, entityName)) makeDeleteRequest((dataAt(SYS) / "my" / entityName / idOf(personal, entityName)).DELETE <@ (user1)).code should equal(200) Then("the record can be deleted at SYS") diff --git a/obp-api/src/test/scala/code/api/v7_0_0/GroupMembershipSyncTest.scala b/obp-api/src/test/scala/code/api/v7_0_0/GroupMembershipSyncTest.scala index e40a0ac799..0a00d2f6dd 100644 --- a/obp-api/src/test/scala/code/api/v7_0_0/GroupMembershipSyncTest.scala +++ b/obp-api/src/test/scala/code/api/v7_0_0/GroupMembershipSyncTest.scala @@ -51,6 +51,8 @@ class GroupMembershipSyncTest extends V600ServerSetup { object VersionOfApi extends Tag(ApiVersion.v7_0_0.toString) object SyncGroupMembers extends Tag("syncGroupMembers") + object SyncGroupMember extends Tag("syncGroupMember") + object SyncUserGroups extends Tag("syncUserGroups") object AddEntitlement extends Tag("addEntitlement") private def bankId = testBankId1.value @@ -58,6 +60,8 @@ class GroupMembershipSyncTest extends V600ServerSetup { Entitlement.entitlement.vend.addEntitlement(bank, resourceUser1.userId, role.toString) private def message(r: code.setup.APIResponse) = r.body.extract[ErrorMessage].message private def sync(groupId: String) = v7 / "management" / "groups" / groupId / "sync-members" + private def syncOne(groupId: String, userId: String) = v7 / "management" / "groups" / groupId / "users" / userId / "sync" + private def syncUser(userId: String) = v7 / "management" / "users" / userId / "sync-groups" private def membership(groupId: String) = compact(render("group_id" -> groupId)) private def user2Entitlements = Entitlement.entitlement.vend.getEntitlementsByUserId(resourceUser2.userId).toList.flatten.filter(_.bankId == bankId) @@ -101,6 +105,17 @@ class GroupMembershipSyncTest extends V600ServerSetup { val withAddOnly = makePostRequest(sync(group.groupId).POST <@ (user1), "") withAddOnly.code should equal(403) message(withAddOnly) should startWith(UserHasMissingRoles) + + makePostRequest(syncOne(group.groupId, resourceUser2.userId).POST, "").code should equal(401) + val oneWithAddOnly = makePostRequest(syncOne(group.groupId, resourceUser2.userId).POST <@ (user1), "") + oneWithAddOnly.code should equal(403) + message(oneWithAddOnly) should startWith(UserHasMissingRoles) + + makePostRequest(syncUser(resourceUser2.userId).POST, "").code should equal(401) + Entitlement.entitlement.vend.addEntitlement(bankId, resourceUser2.userId, "CanTestSyncAuthRole", groupId = Some(group.groupId)) + val userWithAddOnly = makePostRequest(syncUser(resourceUser2.userId).POST <@ (user1), "") + userWithAddOnly.code should equal(403) + message(userWithAddOnly) should startWith(UserHasMissingRoles) } scenario("overlapping Groups: membership, sync, and removal", SyncGroupMembers, VersionOfApi) { @@ -160,4 +175,98 @@ class GroupMembershipSyncTest extends V600ServerSetup { roleGroup(r3) should equal(None) } } + + feature("Sync Group Member") { + scenario("one member of a changed Group, and a user who is not a member", SyncGroupMember, VersionOfApi) { + grant(canAddUserToGroupAtAllBanks) + grant(canRemoveUserFromGroupAtAllBanks) + val s1 = "CanTestSyncOneRole1" + val s2 = "CanTestSyncOneRole2" + val g = GroupTrait.group.vend.createGroup(Some(bankId), "sync-one", "", List(s1), isEnabled = true).openOrThrowException("g") + makePostRequest((v6_0_0_Request / "users" / resourceUser2.userId / "group-entitlements").POST <@ (user1), + membership(g.groupId)).code should equal(201) + + When("the Group's Roles change from s1 to s2, and user2 is synced as a dry run") + GroupTrait.group.vend.updateGroup(g.groupId, None, None, Some(List(s2)), None) + val dry = makePostRequest(syncOne(g.groupId, resourceUser2.userId).POST <@ (user1) < (m \ "user_id").extract[String]) should equal(List(resourceUser2.userId)) + (members.head \ "entitlements_created").extract[List[String]] should equal(List(s2)) + (members.head \ "entitlements_deleted").extract[List[String]] should equal(List(s1)) + Then("nothing changed") + roleGroup(s1) should equal(Some(g.groupId)) + roleGroup(s2) should equal(None) + + When("user2 is synced for real") + makePostRequest(syncOne(g.groupId, resourceUser2.userId).POST <@ (user1), "").code should equal(200) + Then("user2 holds s2 from the Group, and not s1") + roleGroup(s2) should equal(Some(g.groupId)) + roleGroup(s1) should equal(None) + + Then("user1, who is not a member, is 404, and so is an unknown user") + makePostRequest(syncOne(g.groupId, resourceUser1.userId).POST <@ (user1), "").code should equal(404) + makePostRequest(syncOne(g.groupId, "no-such-user").POST <@ (user1), "").code should equal(404) + } + } + + feature("Sync User Groups") { + scenario("every Group of a user, including one that was deleted", SyncUserGroups, VersionOfApi) { + grant(canAddUserToGroupAtAllBanks) + grant(canRemoveUserFromGroupAtAllBanks) + val u5 = "CanTestSyncUserRole5" + val u6 = "CanTestSyncUserRole6" + val u7 = "CanTestSyncUserRole7" + val u8 = "CanTestSyncUserRole8" + val u9 = "CanTestSyncUserRole9" + val c = GroupTrait.group.vend.createGroup(Some(bankId), "sync-user-C", "", List(u5, u6), isEnabled = true).openOrThrowException("C") + val d = GroupTrait.group.vend.createGroup(Some(bankId), "sync-user-D", "", List(u6), isEnabled = true).openOrThrowException("D") + val e = GroupTrait.group.vend.createGroup(Some(bankId), "sync-user-E", "", List(u7), isEnabled = true).openOrThrowException("E") + val addUser2 = (v6_0_0_Request / "users" / resourceUser2.userId / "group-entitlements").POST <@ (user1) + List(c, d, e).foreach(g => makePostRequest(addUser2, membership(g.groupId)).code should equal(201)) + roleGroup(u6) should equal(Some(c.groupId)) + + When("E is deleted (its Entitlements stay), C now grants u5 and u8, and D grants u6, u8 and u9") + GroupTrait.group.vend.deleteGroup(e.groupId) + GroupMemberships.removeMembershipsOfGroup(e.groupId) + roleGroup(u7) should equal(Some(e.groupId)) + GroupTrait.group.vend.updateGroup(c.groupId, None, None, Some(List(u5, u8)), None) + GroupTrait.group.vend.updateGroup(d.groupId, None, None, Some(List(u6, u8, u9)), None) + + When("user2's Groups are synced as a dry run") + val dry = makePostRequest(syncUser(resourceUser2.userId).POST <@ (user1) < (g \ "group_id").extract[String] == id).get + def created(id: String) = (of(id) \ "entitlements_created").extract[List[String]] + Then("r8, which both C and D now grant, would be granted once") + (created(c.groupId) ++ created(d.groupId)).count(_ == u8) should equal(1) + created(d.groupId) should contain(u9) + (of(c.groupId) \ "entitlements_moved").children.map(m => (m \ "role_name").extract[String]) should equal(List(u6)) + (of(e.groupId) \ "group_deleted").extract[Boolean] should equal(true) + (of(e.groupId) \ "entitlements_deleted").extract[List[String]] should equal(List(u7)) + Then("nothing changed") + roleGroup(u8) should equal(None) + roleGroup(u6) should equal(Some(c.groupId)) + roleGroup(u7) should equal(Some(e.groupId)) + + When("user2's Groups are synced for real") + makePostRequest(syncUser(resourceUser2.userId).POST <@ (user1), "").code should equal(200) + Then("user2 holds what C and D grant, u6 now recorded against D, and u7 from the deleted E is gone") + roleGroup(u5) should equal(Some(c.groupId)) + List(Some(c.groupId), Some(d.groupId)) should contain(roleGroup(u8)) + roleGroup(u9) should equal(Some(d.groupId)) + roleGroup(u6) should equal(Some(d.groupId)) + roleGroup(u7) should equal(None) + + Then("a second sync changes nothing") + val again = makePostRequest(syncUser(resourceUser2.userId).POST <@ (user1), "") + again.code should equal(200) + (again.body \ "groups").children.foreach { g => + (g \ "entitlements_created").extract[List[String]] shouldBe empty + (g \ "entitlements_deleted").extract[List[String]] shouldBe empty + (g \ "entitlements_moved").children shouldBe empty + } + } + } } diff --git a/obp-api/src/test/scala/code/entitlement/RoleGrantedEmailOutboxTest.scala b/obp-api/src/test/scala/code/entitlement/RoleGrantedEmailOutboxTest.scala new file mode 100644 index 0000000000..54ff2f136d --- /dev/null +++ b/obp-api/src/test/scala/code/entitlement/RoleGrantedEmailOutboxTest.scala @@ -0,0 +1,96 @@ +/** +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.entitlement + +import code.api.util.ApiRole.{canCreateAccount, canCreateCustomer} +import code.messageoutbox.{MessageOutbox, MessageOutboxRelay} +import code.setup.{DefaultUsers, ServerSetup} +import net.liftweb.db.DB +import net.liftweb.mapper.By +import net.liftweb.util.DefaultConnectionIdentifier + +/** + * The "you have been granted a Role" email goes through the message outbox: it is queued in the + * granting transaction and sent by the relay, never on the request thread. + */ +class RoleGrantedEmailOutboxTest extends ServerSetup with DefaultUsers { + + private def emailRowsFor(entitlementId: String): List[MessageOutbox] = + MessageOutbox.findAll(By(MessageOutbox.OutboxType, MessageOutbox.TYPE_EMAIL), By(MessageOutbox.SubjectId, entitlementId)) + + override def beforeEach(): Unit = { + super.beforeEach() + setPropsValues("mail.api.consumer.registered.sender.address" -> "noreply@example.com", "mail.test.mode" -> "true") + } + + feature("Role granted emails are queued in the message outbox") { + + scenario("a grant queues one email, which the relay sends") { + val bankId = testBankId1.value + val entitlement = Entitlement.entitlement.vend.addEntitlement(bankId, resourceUser2.userId, canCreateCustomer.toString) + .openOrThrowException("grant") + + Then("one PENDING email row is queued for the entitlement, addressed to the user") + val queued = emailRowsFor(entitlement.entitlementId) + queued.map(_.status) should equal(List(MessageOutbox.STATUS_PENDING)) + queued.head.operationName should equal(MessageOutbox.OPERATION_ROLE_GRANTED_EMAIL) + MessageOutbox.emailPayload(queued.head).to should equal(List(resourceUser2.emailAddress)) + MessageOutbox.emailPayload(queued.head).subject should include(canCreateCustomer.toString) + + When("the relay runs") + MessageOutboxRelay.relayOnePass() + Then("the email is DELIVERED, after one attempt") + val sent = emailRowsFor(entitlement.entitlementId) + sent.map(_.status) should equal(List(MessageOutbox.STATUS_DELIVERED)) + sent.head.attempts should equal(1) + } + + scenario("a row can be claimed for delivery once") { + val entitlement = Entitlement.entitlement.vend.addEntitlement(testBankId1.value, resourceUser2.userId, canCreateAccount.toString) + .openOrThrowException("grant") + val read = emailRowsFor(entitlement.entitlementId).head + Then("two relays that read the same row: only the first claims it") + MessageOutbox.claimForDelivery(read) should equal(true) + MessageOutbox.claimForDelivery(read) should equal(false) + } + + scenario("a grant that is rolled back queues no email") { + val role = "CanGetCustomersAtOneBank" + intercept[RuntimeException] { + DB.use(DefaultConnectionIdentifier) { _ => + Entitlement.entitlement.vend.addEntitlement(testBankId1.value, resourceUser2.userId, role) + throw new RuntimeException("the request failed after the grant") + } + } + Then("neither the entitlement nor its email exists") + Entitlement.entitlement.vend.getEntitlement(testBankId1.value, resourceUser2.userId, role).isDefined should equal(false) + MessageOutbox.findAll(By(MessageOutbox.OutboxType, MessageOutbox.TYPE_EMAIL)) + .exists(r => MessageOutbox.emailPayload(r).subject.contains(role)) should equal(false) + } + } +} diff --git a/obp-api/src/test/scala/code/metrics/MetricsTest.scala b/obp-api/src/test/scala/code/metrics/MetricsTest.scala index 8abf1b588c..0a59701ee7 100644 --- a/obp-api/src/test/scala/code/metrics/MetricsTest.scala +++ b/obp-api/src/test/scala/code/metrics/MetricsTest.scala @@ -54,6 +54,7 @@ class MetricsTest extends ServerSetup with WipeMetrics { val testVerb = "verb" val testSourceIp = "2001:0db8:3c4d:0015:0000:0000:1a2f:1a2b" val testTargetIp = "2001:0db8:3c4d:0015:0000:0000:1a2f:1a2b" + val testForwardedFor = "203.0.113.9, 10.0.0.2, 10.0.0.3" val testApiInstanceId = "test_instance" val testResponseBody: String = """fbdgbdbg}""".stripMargin @@ -71,6 +72,8 @@ class MetricsTest extends ServerSetup with WipeMetrics { val limit = 100 val limit1 = 101 val limit2 = 102 + val limit3 = 103 + val limit4 = 104 override def beforeEach(): Unit = { super.beforeEach() @@ -92,7 +95,7 @@ class MetricsTest extends ServerSetup with WipeMetrics { scenario("We save a new API metric") { metrics.saveMetric(testUserId,testUrl1, day1, -1L, testUserName, testAppName, testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, - testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testApiInstanceId, null, null, null, null) + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) MetricBatchWriter.flush() val byUrl = metrics.getAllMetrics(List(OBPLimit(limit))).groupBy(_.getUrl()) @@ -107,19 +110,50 @@ class MetricsTest extends ServerSetup with WipeMetrics { metric.getUrl() should equal(testUrl1) } + scenario("The source address and the hops are stored and read back") { + metrics.saveMetric(testUserId, testUrl1, day1, -1L, testUserName, testAppName, + testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, + testVersion, testVerb, None, getCorrelationId(), testResponseBody, "203.0.113.9", testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) + MetricBatchWriter.flush() + + val metric = metrics.getAllMetrics(List(OBPLimit(limit3))).head + metric.getSourceIp() should equal("203.0.113.9") + metric.getForwardedFor() should equal(testForwardedFor) + } + + scenario("A value longer than its column is cut to fit, so it does not lose the other metrics of the flush") { + // A caller can send any X-Forwarded-For or X-Forwarded-Host. Before values were cut to fit, + // one such value failed the whole batch insert and every metric in the flush was lost. + val longChain = List.fill(100)("198.51.100.42").mkString(", ") + val longHost = "h" * 500 + metrics.saveMetric(testUserId, testUrl1, day1, -1L, testUserName, testAppName, + testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp, longHost, longChain, testApiInstanceId, null, null, null, null) + metrics.saveMetric(testUserId, testUrl2, day1, -1L, testUserName, testAppName, + testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp, testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) + MetricBatchWriter.flush() + + val byUrl = metrics.getAllMetrics(List(OBPLimit(limit4))).groupBy(_.getUrl()) + byUrl.keys.size should equal(2) + byUrl(testUrl1).head.getForwardedFor() should equal(longChain.take(MappedMetric.forwardedFor.maxLen)) + byUrl(testUrl1).head.getTargetIp() should equal(longHost.take(MappedMetric.targetIp.maxLen)) + byUrl(testUrl2).head.getForwardedFor() should equal(testForwardedFor) + } + scenario("Group all metrics by url") { metrics.saveMetric(testUserId, testUrl1, day1, -1L, testUserName, testAppName, testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, - testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testApiInstanceId, null, null, null, null) + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) metrics.saveMetric(testUserId, testUrl1, day1, -1L, testUserName, testAppName, testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, - testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testApiInstanceId, null, null, null, null) + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) metrics.saveMetric(testUserId, testUrl1, day2, -1L, testUserName, testAppName, testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, - testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testApiInstanceId, null, null, null, null) + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) metrics.saveMetric(testUserId, testUrl2, day2, -1L, testUserName, testAppName, testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, - testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testApiInstanceId, null, null, null, null) + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) MetricBatchWriter.flush() val byUrl = metrics.getAllMetrics(List(OBPLimit(limit1))).groupBy(_.getUrl()) @@ -140,16 +174,16 @@ class MetricsTest extends ServerSetup with WipeMetrics { scenario("Group all metrics by day") { metrics.saveMetric(testUserId, testUrl1, day1, -1L, testUserName, testAppName, testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, - testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testApiInstanceId, null, null, null, null) + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) metrics.saveMetric(testUserId, testUrl1, day1, -1L, testUserName, testAppName, testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, - testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testApiInstanceId, null, null, null, null) + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) metrics.saveMetric(testUserId, testUrl1, day2, -1L, testUserName, testAppName, testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, - testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testApiInstanceId, null, null, null, null) + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) metrics.saveMetric(testUserId, testUrl2, day2, -1L, testUserName, testAppName, testDeveloperEmail, testConsumerId, testImplementedByPartialFunction, - testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testApiInstanceId, null, null, null, null) + testVersion, testVerb, None, getCorrelationId(), testResponseBody, testSourceIp , testTargetIp, testForwardedFor, testApiInstanceId, null, null, null, null) MetricBatchWriter.flush() val byDay = metrics.getAllMetrics(List(OBPLimit(limit2))).groupBy(APIMetrics.getMetricDay) diff --git a/obp-api/src/test/scala/code/scheduler/MetricsArchiveSchedulerTest.scala b/obp-api/src/test/scala/code/scheduler/MetricsArchiveSchedulerTest.scala index d88d6c661b..dc2bcea9d0 100644 --- a/obp-api/src/test/scala/code/scheduler/MetricsArchiveSchedulerTest.scala +++ b/obp-api/src/test/scala/code/scheduler/MetricsArchiveSchedulerTest.scala @@ -82,6 +82,7 @@ class MetricsArchiveSchedulerTest extends ServerSetup { .responseBody("body") .sourceIp("127.0.0.1") .targetIp("127.0.0.1") + .forwardedFor("203.0.113.9, 127.0.0.1") .apiInstanceId("test") .consentReferenceId("") .saveMe() @@ -120,7 +121,8 @@ class MetricsArchiveSchedulerTest extends ServerSetup { Then("the old row is gone from metric and present in the archive") MappedMetric.find(By(MappedMetric.id, oldRow.id.get)).isDefined should equal(false) - MetricArchive.find(By(MetricArchive.metricId, oldRow.id.get)).isDefined should equal(true) + MetricArchive.find(By(MetricArchive.metricId, oldRow.id.get)).map(_.getForwardedFor()) should equal( + net.liftweb.common.Full("203.0.113.9, 127.0.0.1")) And("the recent row is untouched") MappedMetric.find(By(MappedMetric.id, recentRow.id.get)).isDefined should equal(true) diff --git a/obp-api/src/test/scala/code/setup/LocalMappedConnectorTestSetup.scala b/obp-api/src/test/scala/code/setup/LocalMappedConnectorTestSetup.scala index c8b5e13718..d49ef5a4c0 100644 --- a/obp-api/src/test/scala/code/setup/LocalMappedConnectorTestSetup.scala +++ b/obp-api/src/test/scala/code/setup/LocalMappedConnectorTestSetup.scala @@ -247,6 +247,14 @@ trait LocalMappedConnectorTestSetup extends TestConnectorSetupWithStandardPermis // would clear the whole shared DB and cause cross-shard rate-limit/cache flakiness. try { Redis.deleteKeysByPattern(code.api.Constant.getGlobalCacheNamespacePrefix + "*") + // Deleting the keys also deletes the cache namespace version counters, so each one starts + // again at 1 and a version number seen in the last scenario comes round again. Anything held + // in this JVM under such a version would then look current when it is not: the ResourceDoc + // vocabulary rebuilt for the last scenario's Dynamic Entities, for example, would be reused + // and a new entity's function names dropped from a `functions` filter. So forget those copies + // along with the keys. + code.api.Constant.forgetRecentCacheNamespaceVersions() + code.api.util.ResourceDocVocabulary.refreshDynamic() } catch { case e: Throwable => logger.warn("------------| Redis issue during flushing data |------------") diff --git a/release_notes.md b/release_notes.md index d91140152a..2cfaf2fa5e 100644 --- a/release_notes.md +++ b/release_notes.md @@ -3,6 +3,31 @@ ### Most recent changes at top of file ``` Date Commit Action +30/09/2026 TBD CHANGED in v7.0.0: a Dynamic Entity record at + /obp/v7.0.0/banks/BANK_ID/dynamic-entities/... is followed by a metadata + object: created and updated, each with at (UTC), user_id (the User who made + the call, an agent's own id when an agent made it) and on_behalf_of_user_id + (the User it was made for). List items change shape to + {"": {...}, "metadata": {...}}. Public reads carry the times only. + A record held by another connector has no metadata. The unversioned + /obp/dynamic-entity/ URLs are unchanged. + NEW columns on DynamicData: createdat, updatedat, createdbyuserid, + createdbyonbehalfofuserid, updatedbyuserid, updatedbyonbehalfofuserid. They + are null for records written before this change; updated* is filled on a + record's next save, created* stays null. +29/09/2026 TBD CHANGED: the "you have been granted a Role" email is queued in the message + outbox (outbox_type EMAIL, subject_id = the entitlement_id) in the granting + transaction, and sent by the outbox relay. A grant that is rolled back (a + request that fails or times out) emails nobody, and SMTP sends no longer run + on the shared thread pool, where a burst of grants (e.g. a Group member sync) + could stall the API. The relay now always runs (Open Corridor rows are still + only relayed when open_corridor_enabled=true); an email not sent is retried + with backoff and goes STICKY after 8 attempts, visible and retryable through + GET /management/message-outbox?outbox_type=EMAIL. Each email row is claimed + before it is sent, so several instances never send it twice, and relay + passes no longer overlap. The relay interval prop has a new name: + message_outbox.relay_interval_seconds (default 10), replacing + open_corridor.outbox_relay_interval, which is no longer read. 29/09/2026 TBD NEW in v7.0.0: POST /management/groups/GROUP_ID/sync-members[?dry_run=true] brings the Entitlements of a Group's members in line with its current Roles: grants the Roles they lack, deletes those the Group granted but no longer