Skip to content

Commit bb2eb4e

Browse files
committed
AmqpBankBroker, MessageOutbox instead of OC
1 parent a071a33 commit bb2eb4e

15 files changed

Lines changed: 724 additions & 378 deletions

‎obp-api/src/main/scala/bootstrap/liftweb/Boot.scala‎

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -135,7 +135,8 @@ import code.transactionRequestAttribute.TransactionRequestAttribute
135135
import code.transactionStatusScheduler.TransactionRequestStatusScheduler
136136
import code.transaction_types.MappedTransactionType
137137
import code.transactionattribute.MappedTransactionAttribute
138-
import code.bankconnectors.opencorridor.{OpenCorridorBankBroker, OpenCorridorOutbox, OpenCorridorOutboxRelay}
138+
import code.amqpbroker.AmqpBankBroker
139+
import code.messageoutbox.{MessageOutbox, MessageOutboxRelay}
139140
import code.transactionrequests.{MappedTransactionRequest, MappedTransactionRequestTypeCharge, TransactionRequestReasons}
140141
import code.usercustomerlinks.MappedUserCustomerLink
141142
import code.customerlinks.CustomerLink
@@ -564,7 +565,7 @@ class Boot extends MdcLoggable {
564565
// Open Corridor: the transactional-outbox relay publishing Interface C messages
565566
// (credit notifications + settlement instructions) to the banks' own vhosts.
566567
if (APIUtil.getPropsAsBoolValue("open_corridor_enabled", false)) {
567-
OpenCorridorOutboxRelay.start(APIUtil.getPropsAsLongValue("open_corridor.outbox_relay_interval", 10L))
568+
MessageOutboxRelay.start(APIUtil.getPropsAsLongValue("open_corridor.outbox_relay_interval", 10L))
568569
}
569570
APIUtil.getPropsAsLongValue("database_messages_scheduler_interval") match {
570571
case Full(i) => DatabaseDriverScheduler.start(i)
@@ -994,8 +995,8 @@ object ToSchemify extends MdcLoggable {
994995
MappedCounterpartyWhereTag,
995996
MappedTransactionRequest,
996997
TransactionRequestAttribute,
997-
OpenCorridorBankBroker,
998-
OpenCorridorOutbox,
998+
AmqpBankBroker,
999+
MessageOutbox,
9991000
MappedMetric,
10001001
MetricArchive,
10011002
MetricsArchiveRun,
Lines changed: 103 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,103 @@
1+
package code.amqpbroker
2+
3+
import net.liftweb.common.Box
4+
import net.liftweb.mapper._
5+
6+
/**
7+
* Per-bank AMQP broker coordinates — where OBP-API publishes messages destined
8+
* for a bank's own infrastructure.
9+
*
10+
* Named by transport, not by consumer: the fields (host/port/vhost/credentials)
11+
* are AMQP 0-9-1 concepts, and any feature that needs to push AMQP messages to
12+
* a specific bank resolves its coordinates here. The first consumer is Open
13+
* Corridor Interface C: each onboarded bank's Bank Node consumes on its OWN
14+
* vhost (e.g. `/bank.ke.01.kcs`) with its own credentials — permission
15+
* isolation is enforced at the broker level, and publishing is keyed by
16+
* bank_id through this registry (populated at onboarding).
17+
*
18+
* Transport coordinates only: the bank's on-chain settlement address is NOT
19+
* stored here — it is the CARDANO account routing on the bank's
20+
* OBP-INCOMING-SETTLEMENT-ACCOUNT.
21+
*/
22+
class AmqpBankBroker extends LongKeyedMapper[AmqpBankBroker] with IdPK {
23+
def getSingleton = AmqpBankBroker
24+
25+
object BankId extends MappedString(this, 255) {
26+
override def dbColumnName = "bank_id"
27+
}
28+
object Host extends MappedString(this, 255) {
29+
override def dbColumnName = "host"
30+
}
31+
object Port extends MappedInt(this) {
32+
override def dbColumnName = "port"
33+
override def defaultValue = 5672
34+
}
35+
object VirtualHost extends MappedString(this, 255) {
36+
override def dbColumnName = "virtual_host"
37+
}
38+
object Username extends MappedString(this, 255) {
39+
override def dbColumnName = "username"
40+
}
41+
/** Write-only: accepted on registration, never echoed by any endpoint. */
42+
object Password extends MappedString(this, 255) {
43+
override def dbColumnName = "password"
44+
}
45+
object UseSsl extends MappedBoolean(this) {
46+
override def dbColumnName = "use_ssl"
47+
override def defaultValue = false
48+
}
49+
object CreatedAt extends MappedDateTime(this) {
50+
override def dbColumnName = "created_at"
51+
override def defaultValue = new java.util.Date()
52+
}
53+
object UpdatedAt extends MappedDateTime(this) {
54+
override def dbColumnName = "updated_at"
55+
override def defaultValue = new java.util.Date()
56+
}
57+
58+
def bankId: String = BankId.get
59+
def host: String = Host.get
60+
def port: Int = Port.get
61+
def virtualHost: String = VirtualHost.get
62+
def username: String = Username.get
63+
def password: String = Password.get
64+
def useSsl: Boolean = UseSsl.get
65+
66+
override def save: Boolean = {
67+
UpdatedAt(new java.util.Date())
68+
super.save
69+
}
70+
}
71+
72+
object AmqpBankBroker extends AmqpBankBroker with LongKeyedMetaMapper[AmqpBankBroker] {
73+
override def dbTableName = "amqp_bank_broker"
74+
75+
override def dbIndexes: List[BaseIndex[AmqpBankBroker]] = UniqueIndex(BankId) :: super.dbIndexes
76+
77+
def findByBankId(bankId: String): Box[AmqpBankBroker] =
78+
AmqpBankBroker.find(By(AmqpBankBroker.BankId, bankId))
79+
80+
/** Upsert the broker coordinates for a bank (one row per bank, enforced by the unique index). */
81+
def upsert(
82+
bankId: String,
83+
host: String,
84+
port: Int,
85+
virtualHost: String,
86+
username: String,
87+
password: String,
88+
useSsl: Boolean
89+
): AmqpBankBroker = {
90+
val row = findByBankId(bankId).getOrElse(AmqpBankBroker.create.BankId(bankId))
91+
row
92+
.Host(host)
93+
.Port(port)
94+
.VirtualHost(virtualHost)
95+
.Username(username)
96+
.Password(password)
97+
.UseSsl(useSsl)
98+
.saveMe()
99+
}
100+
101+
def deleteByBankId(bankId: String): Boolean =
102+
AmqpBankBroker.bulkDelete_!!(By(AmqpBankBroker.BankId, bankId))
103+
}

‎obp-api/src/main/scala/code/api/util/ApiRole.scala‎

Lines changed: 15 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -212,16 +212,27 @@ object ApiRole extends MdcLoggable{
212212
case class CanAttachOpenCorridorPromise(requiresBankId: Boolean = true) extends ApiRole
213213
lazy val canAttachOpenCorridorPromise = CanAttachOpenCorridorPromise()
214214

215-
// Open Corridor: operator role for registering each onboarded bank's RabbitMQ broker
216-
// coordinates (host/port/vhost/credentials) in the per-bank publish registry.
217-
case class CanConfigureOpenCorridorBroker(requiresBankId: Boolean = false) extends ApiRole
218-
lazy val canConfigureOpenCorridorBroker = CanConfigureOpenCorridorBroker()
215+
// Operator role for registering each onboarded bank's AMQP broker coordinates
216+
// (host/port/vhost/credentials) in the per-bank publish registry. Transport
217+
// registry, not corridor-specific; Open Corridor Interface C is the first consumer.
218+
case class CanConfigureAmqpBankBroker(requiresBankId: Boolean = false) extends ApiRole
219+
lazy val canConfigureAmqpBankBroker = CanConfigureAmqpBankBroker()
219220

220221
// Open Corridor: operator role for the settle-pair trigger — nets a bank pair's
221222
// PENDING promises, posts the net Transaction and enqueues the Interface C messages.
222223
case class CanSettleOpenCorridor(requiresBankId: Boolean = true) extends ApiRole
223224
lazy val canSettleOpenCorridor = CanSettleOpenCorridor()
224225

226+
// Operator role for reading the generic message outbox (delivery states,
227+
// sticky failures) across all outbox types.
228+
case class CanGetMessageOutbox(requiresBankId: Boolean = false) extends ApiRole
229+
lazy val canGetMessageOutbox = CanGetMessageOutbox()
230+
231+
// Operator role for re-queuing a STICKY message-outbox row after
232+
// reconciliation — flips it back to PENDING for the relay to redeliver.
233+
case class CanRetryMessageOutbox(requiresBankId: Boolean = false) extends ApiRole
234+
lazy val canRetryMessageOutbox = CanRetryMessageOutbox()
235+
225236
case class CanAddSocialMediaHandle(requiresBankId: Boolean = true) extends ApiRole
226237
lazy val canAddSocialMediaHandle = CanAddSocialMediaHandle()
227238

‎obp-api/src/main/scala/code/api/util/ErrorMessages.scala‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -889,11 +889,13 @@ object ErrorMessages {
889889
val OpenCorridorPromiseTypeMismatch = "OBP-40051: The Transaction Request is not of type OPEN_CORRIDOR_PROMISE."
890890
val OpenCorridorPromiseNotPending = "OBP-40052: The Open Corridor promise Transaction Request is not in PENDING status."
891891
val OpenCorridorPromiseEvidenceConflict = "OBP-40053: Open Corridor promise evidence is already attached to this Transaction Request with different values. Evidence cannot be overwritten."
892-
val OpenCorridorBankBrokerNotConfigured = "OBP-40054: No Open Corridor broker is configured for this bank. Register the bank's RabbitMQ coordinates first."
892+
val AmqpBankBrokerNotConfigured = "OBP-40054: No AMQP broker is configured for this bank. Register the bank's AMQP coordinates first."
893893
val OpenCorridorPublishFailed = "OBP-40055: Could not publish the Open Corridor message to the bank's broker or no reply arrived in time."
894894
val OpenCorridorSettlementAddressMissing = "OBP-40056: The creditor bank has no settlement address registered in its Open Corridor broker registration, so the settlement instruction cannot be addressed."
895895
val OpenCorridorDisabled = "OBP-40057: Open Corridor is not enabled on this API instance. Set open_corridor_enabled=true in the props."
896896
val OpenCorridorSettlementNotFound = "OBP-40058: No Open Corridor settlement with this SETTLEMENT_ID exists for this bank."
897+
val MessageOutboxRowNotFound = "OBP-40059: No message outbox row with this OUTBOX_ID exists."
898+
val MessageOutboxRowNotSticky = "OBP-40060: The message outbox row is not STICKY. Only STICKY rows can be re-queued; PENDING rows retry automatically."
897899
// Exceptions (OBP-50XXX)
898900
val UnknownError = "OBP-50000: Unknown Error."
899901
val FutureTimeoutException = "OBP-50001: Future Timeout Exception."

0 commit comments

Comments
 (0)