From 9b9a95ec173d21568bf4e251aaecbdc1943dd55f Mon Sep 17 00:00:00 2001 From: Yassine Rhouma Date: Fri, 24 Jul 2026 11:33:43 +0200 Subject: [PATCH 1/6] refactor(dva-api): slim to pure HTTP orchestrator --- .../hu/bme/mit/ftsrg/dva/api/Application.kt | 18 --- .../mit/ftsrg/dva/api/resource/Evaluation.kt | 11 -- .../mit/ftsrg/dva/api/resource/Templates.kt | 18 --- .../hu/bme/mit/ftsrg/dva/api/resource/VLAs.kt | 17 --- .../ftsrg/dva/api/route/evaluationRoutes.kt | 73 ---------- .../mit/ftsrg/dva/api/route/templateRoutes.kt | 93 ------------- .../bme/mit/ftsrg/dva/api/route/vlaRoutes.kt | 128 ------------------ 7 files changed, 358 deletions(-) delete mode 100644 dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/Evaluation.kt delete mode 100644 dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/Templates.kt delete mode 100644 dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/VLAs.kt delete mode 100644 dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/evaluationRoutes.kt delete mode 100644 dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/templateRoutes.kt delete mode 100644 dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/vlaRoutes.kt diff --git a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/Application.kt b/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/Application.kt index 0426d21..4a31918 100644 --- a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/Application.kt +++ b/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/Application.kt @@ -1,15 +1,10 @@ package hu.bme.mit.ftsrg.dva.api -import com.rabbitmq.client.Connection -import com.rabbitmq.client.ConnectionFactory import hu.bme.mit.ftsrg.dva.api.db.* import hu.bme.mit.ftsrg.dva.api.err.addHandlers -import hu.bme.mit.ftsrg.dva.api.rabbit.connectWithRetry import hu.bme.mit.ftsrg.dva.api.route.* import hu.bme.mit.ftsrg.dva.log.ReqestLogRepo import hu.bme.mit.ftsrg.dva.log.VerifRequestLogRepo -import hu.bme.mit.ftsrg.dva.vla.TemplateRepo -import hu.bme.mit.ftsrg.dva.vla.VLARepo import io.ktor.client.* import io.ktor.client.engine.cio.CIO import io.ktor.http.* @@ -61,15 +56,7 @@ fun Application.installPlugins() { } fun Application.configureKoin() { - val rabbitHost = environment.config.property("rabbitmq.host").getString() - val appModule = module { - single { - ConnectionFactory().run { - host = rabbitHost - connectWithRetry(logger = log) - } - } single { HttpClient(CIO) { install(ClientContentNegotiation) { @@ -80,10 +67,8 @@ fun Application.configureKoin() { } } } - single { PgTemplateRepo() } single { PgRequestLogRepo() } single { PgVerifRequestLogRepo() } - single { PgVLARepo() } } serverInstall(Koin) { modules(appModule) } @@ -91,9 +76,6 @@ fun Application.configureKoin() { fun Application.addRoutes() { docRoutes(openapiPath = environment.config.property("swagger.openapiFile").getString()) - templateRoutes() aovRoutes() - vlaRoutes() - evaluationRoutes() infoRoutes() } diff --git a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/Evaluation.kt b/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/Evaluation.kt deleted file mode 100644 index 18e2d1e..0000000 --- a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/Evaluation.kt +++ /dev/null @@ -1,11 +0,0 @@ -package hu.bme.mit.ftsrg.dva.api.resource - -import io.ktor.resources.* - -@Suppress("unused") -@Resource("/evaluate") -class Evaluation { - - @Resource(path = "from-template") - class FromTemplate(val parent: Evaluation = Evaluation()) -} \ No newline at end of file diff --git a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/Templates.kt b/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/Templates.kt deleted file mode 100644 index 5058902..0000000 --- a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/Templates.kt +++ /dev/null @@ -1,18 +0,0 @@ -package hu.bme.mit.ftsrg.dva.api.resource - -import io.ktor.resources.* -import kotlin.uuid.ExperimentalUuidApi -import kotlin.uuid.Uuid - -@Suppress("unused") -@Resource("/template") -class Templates { - - @OptIn(ExperimentalUuidApi::class) - @Resource("{id}") - class Id(val parent: Templates = Templates(), val id: Uuid) { - - @Resource("render") - class Render(val parent: Id) - } -} \ No newline at end of file diff --git a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/VLAs.kt b/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/VLAs.kt deleted file mode 100644 index 925a95d..0000000 --- a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/resource/VLAs.kt +++ /dev/null @@ -1,17 +0,0 @@ -package hu.bme.mit.ftsrg.dva.api.resource - -import io.ktor.resources.* -import kotlin.uuid.ExperimentalUuidApi -import kotlin.uuid.Uuid - -@Suppress("unused") -@Resource("/vla") -class VLAs { - - @OptIn(ExperimentalUuidApi::class) - @Resource("{id}") - class Id(val parent: VLAs = VLAs(), val id: Uuid) - - @Resource("from-templates") - class FromTemplates(val parent: VLAs = VLAs()) -} \ No newline at end of file diff --git a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/evaluationRoutes.kt b/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/evaluationRoutes.kt deleted file mode 100644 index 6f27e63..0000000 --- a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/evaluationRoutes.kt +++ /dev/null @@ -1,73 +0,0 @@ -package hu.bme.mit.ftsrg.dva.api.route - -import hu.bme.mit.ftsrg.dva.api.resource.Evaluation -import hu.bme.mit.ftsrg.dva.api.service.render -import hu.bme.mit.ftsrg.dva.dto.ErrDTO -import hu.bme.mit.ftsrg.dva.evaluation.Evaluate -import hu.bme.mit.ftsrg.dva.evaluation.EvaluateFromTemplate -import hu.bme.mit.ftsrg.dva.vla.Template -import hu.bme.mit.ftsrg.dva.vla.TemplateRepo -import hu.bme.mit.ftsrg.odcs.DataQuality -import io.ktor.client.HttpClient -import io.ktor.client.request.* -import io.ktor.client.statement.bodyAsText -import io.ktor.http.ContentType -import io.ktor.http.HttpStatusCode.Companion.InternalServerError -import io.ktor.http.HttpStatusCode.Companion.NotFound -import io.ktor.http.contentType -import io.ktor.server.application.Application -import io.ktor.server.request.receive -import io.ktor.server.resources.post -import io.ktor.server.response.respond -import io.ktor.server.response.respondText -import io.ktor.server.routing.routing -import org.koin.ktor.ext.inject -import kotlin.uuid.ExperimentalUuidApi - -@OptIn(ExperimentalUuidApi::class) -fun Application.evaluationRoutes() { - val httpClient by inject() - val templateRepo by inject() - - val processingURL = environment.config.property("processing.url").getString() - - routing { - post { - val request = call.receive() - val response = httpClient.post("${processingURL}/evaluate") { - setBody(request) - contentType(ContentType.Application.Json) - } - - call.respondText( - text = response.bodyAsText(), - contentType = ContentType.Application.Json, - status = response.status - ) - } - - post { - val evaluationReq = call.receive() - val template: Template = templateRepo.byID(evaluationReq.templateID) ?: run { - call.respond(NotFound) - return@post - } - val quality: DataQuality = template.render(evaluationReq.templateModel) ?: run { - call.respond(InternalServerError, ErrDTO(type = "UNKNOWN", title = "Failed to render template")) - return@post - } - - val request = Evaluate(requirement = quality, data = evaluationReq.data) - val response = httpClient.post("${processingURL}/evaluate") { - setBody(request) - contentType(ContentType.Application.Json) - } - - call.respondText( - text = response.bodyAsText(), - contentType = ContentType.Application.Json, - status = response.status - ) - } - } -} diff --git a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/templateRoutes.kt b/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/templateRoutes.kt deleted file mode 100644 index c29205f..0000000 --- a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/templateRoutes.kt +++ /dev/null @@ -1,93 +0,0 @@ -package hu.bme.mit.ftsrg.dva.api.route - -import hu.bme.mit.ftsrg.dva.api.resource.Templates -import hu.bme.mit.ftsrg.dva.api.service.render -import hu.bme.mit.ftsrg.dva.dto.ErrDTO -import hu.bme.mit.ftsrg.dva.dto.IDDTO -import hu.bme.mit.ftsrg.dva.vla.TemplateNew -import hu.bme.mit.ftsrg.dva.vla.TemplatePatch -import hu.bme.mit.ftsrg.dva.vla.TemplateRepo -import io.ktor.http.HttpStatusCode.Companion.BadRequest -import io.ktor.http.HttpStatusCode.Companion.Created -import io.ktor.http.HttpStatusCode.Companion.InternalServerError -import io.ktor.http.HttpStatusCode.Companion.NoContent -import io.ktor.http.HttpStatusCode.Companion.NotFound -import io.ktor.server.application.* -import io.ktor.server.request.* -import io.ktor.server.resources.* -import io.ktor.server.resources.patch -import io.ktor.server.resources.post -import io.ktor.server.response.* -import io.ktor.server.routing.* -import kotlinx.serialization.json.JsonObject -import org.koin.ktor.ext.inject -import kotlin.uuid.ExperimentalUuidApi - -@OptIn(ExperimentalUuidApi::class) -fun Application.templateRoutes() { - val repo by inject() - - routing { - get { - call.respond(repo.all()) - } - - get { req -> - val template = repo.byID(req.id) ?: run { - call.respond(NotFound) - return@get - } - call.respond(template) - } - - post { - val templateReq = call.receive() - - val template = repo.add(templateReq) ?: run { - call.respond(InternalServerError, ErrDTO(type = "UNKNOWN", title = "Failed to create template")) - return@post - } - call.respond(Created, IDDTO(template.id.toString())) - } - - patch { req -> - val patch = call.receive() - if (req.id != patch.id) { - call.respond( - status = BadRequest, - ErrDTO(type = "BAD_REQUEST", title = "ID path parameter does not match ID in body") - ) - } - val updatedTemplate = repo.update(patch) ?: run { - call.respond(NotFound) - return@patch - } - call.respond(updatedTemplate) - } - - delete { req -> - if (repo.remove(req.id)) { - call.respond(NoContent) - } else { - call.respond(NotFound) - } - } - - delete { - repo.removeAll() - call.respond(NoContent) - } - - post { req -> - val template = repo.byID(req.parent.id) ?: run { - call.respond(NotFound) - return@post - } - val model = call.receive() - val renderedQuality = template.render(model) ?: run { - call.respond(BadRequest, ErrDTO(type = "BAD_REQUEST", title = "Failed to render template")) - } - call.respond(renderedQuality) - } - } -} \ No newline at end of file diff --git a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/vlaRoutes.kt b/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/vlaRoutes.kt deleted file mode 100644 index 15c762f..0000000 --- a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/vlaRoutes.kt +++ /dev/null @@ -1,128 +0,0 @@ -package hu.bme.mit.ftsrg.dva.api.route - -import hu.bme.mit.ftsrg.dva.api.resource.VLAs -import hu.bme.mit.ftsrg.dva.api.service.render -import hu.bme.mit.ftsrg.dva.dto.ErrDTO -import hu.bme.mit.ftsrg.dva.dto.IDDTO -import hu.bme.mit.ftsrg.dva.vla.TemplateRepo -import hu.bme.mit.ftsrg.dva.vla.VLANew -import hu.bme.mit.ftsrg.dva.vla.VLANewFromTemplates -import hu.bme.mit.ftsrg.dva.vla.VLARepo -import hu.bme.mit.ftsrg.odcs.DataQuality -import io.ktor.http.HttpStatusCode.Companion.Created -import io.ktor.http.HttpStatusCode.Companion.InternalServerError -import io.ktor.http.HttpStatusCode.Companion.NotFound -import io.ktor.server.application.Application -import io.ktor.server.request.* -import io.ktor.server.resources.* -import io.ktor.server.resources.post -import io.ktor.server.response.* -import io.ktor.server.routing.routing -import kotlinx.serialization.json.* -import org.koin.ktor.ext.inject -import kotlin.uuid.ExperimentalUuidApi - -@OptIn(ExperimentalUuidApi::class) -fun Application.vlaRoutes() { - val templateRepo by inject() - val vlaRepo by inject() - - routing { - get { - call.respond(vlaRepo.all()) - } - - get { req -> - val vla = vlaRepo.byID(req.id) ?: run { - call.respond(NotFound) - return@get - } - call.respond(vla) - } - - post { - val vlaReq = call.receive() - val vla = buildJsonObject { - put("apiVersion", "v3.0.2") - put("kind", "DataContract") - put("version", "0.1.0") - put("status", "active") - - vlaReq.apply { - description?.let { put("description", it) } - servers?.let { put("servers", it) } - schema?.let { put("schema", it) } - quality?.let { put("quality", Json.encodeToJsonElement(it)) } - price?.let { put("price", it) } - team?.let { put("team", it) } - roles?.let { put("roles", it) } - slaProperties?.let { put("slaProperties", it) } - support?.let { put("support", it) } - tags?.let { put("tags", it) } - } - } - val id = vlaRepo.add(vla) ?: run { - call.respond(InternalServerError, ErrDTO(type = "UNKNOWN", title = "Failed to create VLA")) - return@post - } - - call.respond(Created, IDDTO(id.toString())) - } - - post { - val vlaReq = call.receive() - val originalVLA = buildJsonObject { - put("apiVersion", "v3.0.2") - put("kind", "DataContract") - put("version", "0.1.0") - put("status", "active") - - vlaReq.apply { - description?.let { put("description", it) } - servers?.let { put("servers", it) } - schema?.let { put("schema", it) } - quality?.let { put("quality", Json.encodeToJsonElement(it)) } - price?.let { put("price", it) } - team?.let { put("team", it) } - roles?.let { put("roles", it) } - slaProperties?.let { put("slaProperties", it) } - support?.let { put("support", it) } - tags?.let { put("tags", it) } - } - } - - // Extend VLA with rendered templates - val qualityRequirements = mutableListOf() - vlaReq.qualityTemplates?.forEach { templateInstantiation -> - val template = templateRepo.byID(templateInstantiation.id) ?: run { - // TODO: respond with more detail? - call.respond(NotFound) - return@post - } - val quality = template.render(templateInstantiation.model) ?: run { - call.respond(InternalServerError, ErrDTO(type = "UNKNOWN", title = "Failed to render template")) - return@post - } - qualityRequirements.add(quality) - } - val extendedVLA = buildJsonObject { - originalVLA.forEach { k, v -> put(k, v) } - putJsonArray("quality") { - originalVLA["quality"]?.jsonArray?.forEach { add(it) } - qualityRequirements.forEach { add(Json.encodeToJsonElement(it)) } - } - } - - val id = vlaRepo.add(extendedVLA) ?: run { - call.respond(InternalServerError, ErrDTO(type = "UNKNOWN", title = "Failed to create VLA")) - return@post - } - - call.respond(Created, IDDTO(id.toString())) - } - - delete { - vlaRepo.removeAll() - } - } -} \ No newline at end of file From 2ce055285743d87f2bbd699e603978ffc874e921 Mon Sep 17 00:00:00 2001 From: Yassine Rhouma Date: Fri, 24 Jul 2026 11:38:18 +0200 Subject: [PATCH 2/6] refactor(dva-processing): remove RabbitMQ consumer and database writes --- dva-processing/src/dva_processing/config.py | 17 +--- dva-processing/src/dva_processing/main.py | 70 ++------------ .../src/dva_processing/processing.py | 94 +------------------ .../src/dva_processing/rmq_consumer.py | 92 ------------------ 4 files changed, 12 insertions(+), 261 deletions(-) delete mode 100644 dva-processing/src/dva_processing/rmq_consumer.py diff --git a/dva-processing/src/dva_processing/config.py b/dva-processing/src/dva_processing/config.py index cacf09e..d495199 100644 --- a/dva-processing/src/dva_processing/config.py +++ b/dva-processing/src/dva_processing/config.py @@ -2,26 +2,11 @@ from pydantic import BaseModel -# RabbitMQ queue name for AoV requests -QUEUE_NAME = "ATTESTATION_REQUESTS" - -# RabbitMQ server hostname -RABBITMQ_HOST = env.get("DVA_RABBITMQ_HOST", default="localhost") - -# Postgres connection data -PG_URL = env.get("DVA_POSTGRES_URL", default="postgresql://localhost:5432") -PG_USER = env.get("DVA_POSTGRES_USER", default="postgres") -PG_PASS = env.get("DVA_POSTGRES_PASSWORD", default="postgres") - -# Log level (must be supported by structlog) LOG_LEVEL = env.get("DVA_LOG_LEVEL", default="warn") -# ACA-Py Controller URL -ACA_PY_CONTROLLER_URL = env.get("DVA_ACA_PY_CONTROLLER_URL", default="localhost:8050") - class Configuration(BaseModel): log_level: str = LOG_LEVEL -cfg = Configuration() +cfg = Configuration() \ No newline at end of file diff --git a/dva-processing/src/dva_processing/main.py b/dva-processing/src/dva_processing/main.py index bb0a556..0713568 100644 --- a/dva-processing/src/dva_processing/main.py +++ b/dva-processing/src/dva_processing/main.py @@ -1,35 +1,11 @@ -import threading - -import plac import uvicorn import dva_processing.config import dva_processing.http import dva_processing.log -import dva_processing.rmq_consumer -@plac.flg( - "verbose", - help="Be more verbose (INFO loglevel)", - abbrev="v", -) -@plac.flg( - "debug", - help="Enable DEBUG log verbosity", - abbrev="d", -) -@plac.flg( - "no_http", - help="Do not start the HTTP server", - abbrev="H", -) -@plac.flg( - "no_rmq", - help="Do not start the RabbitMQ consumer", - abbrev="R", -) -def main(verbose=False, debug=False, no_http=False, no_rmq=False): +def main(verbose=False, debug=False): if debug: dva_processing.config.cfg.log_level = "debug" elif verbose: @@ -37,44 +13,18 @@ def main(verbose=False, debug=False, no_http=False, no_rmq=False): dva_processing.log.setup_logging() logger = dva_processing.log.get_logger() - if no_http and no_rmq: - logger.error("Both HTTP and RMQ are disabled; nothing to run. Exiting.") - return - - threads = [] - - if not no_http: - http_thread = threading.Thread( - target=uvicorn.run, - args=(dva_processing.http.app,), - kwargs=dict(host="0.0.0.0", port=5000, log_level="info"), - daemon=True, - ) - threads.append(http_thread) - else: - logger.info("HTTP server disabled in CLI options") - - if not no_rmq: - rmq_thread = threading.Thread( - target=dva_processing.rmq_consumer.run, - daemon=True, - ) - threads.append(rmq_thread) - else: - logger.info("RabbitMQ message consumer disabled in CLI options") - - for t in threads: - t.start() - try: - for t in threads: - t.join() - except KeyboardInterrupt: - logger.info("Exiting due to keyboard interrupt") + logger.info("Starting DVA PROCESSING (HTTP only)") + uvicorn.run( + dva_processing.http.app, + host="0.0.0.0", + port=5000, + log_level="info", + ) def cli(): - plac.call(main) + main() if __name__ == "__main__": - cli() + cli() \ No newline at end of file diff --git a/dva-processing/src/dva_processing/processing.py b/dva-processing/src/dva_processing/processing.py index 1cf57ca..e6339d1 100644 --- a/dva-processing/src/dva_processing/processing.py +++ b/dva-processing/src/dva_processing/processing.py @@ -1,15 +1,6 @@ -import json -from typing import Any - -import psycopg as pg - -from .config import PG_PASS, PG_URL, PG_USER from .eval import eval_requirement from .log import get_logger from .model import ( - AoVGenerationRequest, - AoVGenerationRequestPayload, - AoVRequest, EvaluationRequest, EvaluationResult, Requirement, @@ -30,87 +21,4 @@ def handle_eval_request(request: EvaluationRequest) -> EvaluationResult: success=False, details=None, error=str(e), - ) - - -def handle_aov_request(request: AoVRequest) -> AoVGenerationRequest: - logger.debug("Handling an AoV request", request=request) - contract: dict[str, Any] = request.contract - - if "vla" not in contract or "schema" not in contract["vla"]: - logger.warning("No VLA in contract or no requirements in VLA; ignoring") - return None - - # Evaluate all requirements - results: list[EvaluationResult] = [] - any_evaluations = False - if len(contract["vla"]["schema"]) == 0: - logger.warning("No schema items found in VLA") - for schema_item in contract["vla"]["schema"]: - if "quality" not in schema_item: - continue - - requirement_dict: dict - for requirement_dict in schema_item["quality"]: - any_evaluations = True - try: - requirement = Requirement(**requirement_dict) - result: EvaluationResult = eval_requirement(request.data, requirement) - except Exception as e: - logger.warning( - "An error was thrown the evaluation of a requirement; tolerating", - error=e, - ) - result = EvaluationResult( - engine=None, timestamp=now(), success=False, error=str(e) - ) - finally: - results.append(result) - - if not any_evaluations: - logger.warning("Nothing was evaluated from this VLA") - - all_success: bool = all(x.success for x in results) - - # Log evaluation results to psql database - try: - with pg.connect(f"{PG_URL}?user={PG_USER}&password={PG_PASS}") as conn: - conn.execute( - """ - UPDATE request_logs - SET - evaluation_passing = %s, - evaluation_date = %s, - evaluation_results = %s - WHERE request_id = %s - """, - ( - all_success, - now().isoformat(), - json.dumps([r.model_dump_json() for r in results]), - request.id, - ), - ) - logger.info( - f"Successfully updated PostgreSQL entry for request {request.id}", - overall_result=all_success, - request_id=request.id, - ) - except Exception as e: - logger.error( - f"Failed to update request log entry for request {request.id}", error=e - ) - - # Return AoV generation request for ACA-Py - return AoVGenerationRequest( - request_id=request.id, - exchange_id=request.exchangeID, - contract_id=contract["id"], - subject=contract["dataProvider"], - issuer_id=request.attesterID, - payload=AoVGenerationRequestPayload( - success=all_success, - results=results, - ), - target="self", - ) + ) \ No newline at end of file diff --git a/dva-processing/src/dva_processing/rmq_consumer.py b/dva-processing/src/dva_processing/rmq_consumer.py deleted file mode 100644 index af67cba..0000000 --- a/dva-processing/src/dva_processing/rmq_consumer.py +++ /dev/null @@ -1,92 +0,0 @@ -from collections.abc import Generator -from time import sleep - -import requests -import structlog -from pika import BasicProperties, BlockingConnection, ConnectionParameters -from pika.adapters.blocking_connection import BlockingChannel -from pika.spec import Basic - -from .config import ACA_PY_CONTROLLER_URL, QUEUE_NAME, RABBITMQ_HOST -from .log import get_logger -from .model import AoVGenerationRequest, AoVRequest -from .processing import handle_aov_request - -logger: structlog.BoundLogger = get_logger() - - -def run() -> None: - conn: BlockingConnection = connect_with_retry( - ConnectionParameters(host=RABBITMQ_HOST, heartbeat=60) - ) - chan: BlockingChannel = conn.channel() - chan.queue_declare(queue=QUEUE_NAME) - - def callback( - chan: BlockingChannel, - method: Basic.Deliver, - props: BasicProperties, - body: bytes, - ) -> None: - aov_request = AoVRequest.model_validate_json( - bytearray(body).decode(encoding="utf-8") - ) - logger.info("Received AoV request via RMQ", request=aov_request) - - aov_gen_request: AoVGenerationRequest = handle_aov_request(aov_request) - if aov_gen_request is None or not aov_gen_request.payload.success: - logger.warning("Not all checks were successful; not sending to ACA-Py") - return - - logger.info( - "Sending AoV VC generation request to ACA-Py", request=aov_gen_request - ) - requests.post( - f"{ACA_PY_CONTROLLER_URL}/generate_aov", - json=aov_gen_request.model_dump(mode="json"), - ) - - chan.basic_consume(queue=QUEUE_NAME, auto_ack=True, on_message_callback=callback) - - logger.info("Waiting for RMQ messages") - chan.start_consuming() - - -def connect_with_retry( - params: ConnectionParameters, - max_retries: int = -1, - initial_delay_s: float = 1, - max_delay_s: float = 60, - backoff_multiplier: float = 2.0, -) -> BlockingConnection: - delay: float = initial_delay_s - attempt: int - for attempt in range_or_infinity(max_retries): - try: - logger.debug( - f"Connecting to RabbitMQ at {params.host}:{params.port}; attempt {attempt + 1}" - ) - return BlockingConnection(params) - except Exception as e: - logger.error(f"Connection attempt failed: {type(e).__name__}", error=e) - - if max_retries != -1 and attempt >= max_retries: - logger.error( - f"Unable to connect to RabbitMQ at {params.host} after {max_retries} attempts" - ) - raise e - - logger.debug(f"Retrying RMQ connection in {delay}s") - sleep(delay) - - delay = min(delay * backoff_multiplier, max_delay_s) - - -def range_or_infinity(n: int) -> Generator[int, None, None]: - if n == -1: - i = 0 - while True: - yield i - i = i + 1 - else: - yield from range(n) From 23a9eafea64318b350435e7d58bc680c1e81a541 Mon Sep 17 00:00:00 2001 From: Yassine Rhouma Date: Fri, 24 Jul 2026 11:38:44 +0200 Subject: [PATCH 3/6] refactor(dva-api): remove orphan rabbitmq block from application.yaml --- dva-api/api/src/main/resources/application.yaml | 3 --- 1 file changed, 3 deletions(-) diff --git a/dva-api/api/src/main/resources/application.yaml b/dva-api/api/src/main/resources/application.yaml index f2af4bc..4859bc7 100644 --- a/dva-api/api/src/main/resources/application.yaml +++ b/dva-api/api/src/main/resources/application.yaml @@ -13,9 +13,6 @@ postgres: user: "$DVA_POSTGRES_USER:postgres" password: "$DVA_POSTGRES_PASSWORD:postgres" -rabbitmq: - host: "$DVA_RABBITMQ_HOST:localhost" - processing: url: "$DVA_PROCESSING_URL:http://localhost:5000" From 626ec1ed706009d718b866e6e1538aa8859744a2 Mon Sep 17 00:00:00 2001 From: Yassine Rhouma Date: Fri, 24 Jul 2026 11:41:58 +0200 Subject: [PATCH 4/6] refactor(test-env): clean docker-compose topology to one container per role --- test-env/common-services.yml | 7 ------- test-env/compose.yml | 22 ---------------------- 2 files changed, 29 deletions(-) diff --git a/test-env/common-services.yml b/test-env/common-services.yml index 7861819..8dc3d30 100755 --- a/test-env/common-services.yml +++ b/test-env/common-services.yml @@ -9,13 +9,6 @@ services: context: ../ dockerfile: ./dva-processing/Dockerfile init: true - rabbit: - image: rabbitmq:4-management-alpine - healthcheck: - test: rabbitmq-diagnostics -q check_port_connectivity - interval: 5s - timeout: 5s - retries: 3 postgres: image: postgres:17-alpine environment: diff --git a/test-env/compose.yml b/test-env/compose.yml index cd7463e..2e435ae 100755 --- a/test-env/compose.yml +++ b/test-env/compose.yml @@ -6,13 +6,11 @@ services: service: dva-api environment: DVA_PROCESSING_URL: http://dva-processing-provider:5000 - DVA_RABBITMQ_HOST: rabbit-provider DVA_POSTGRES_URL: postgresql://postgres-provider:5432/dva DVA_ACA_PY_CONTROLLER_URL: http://dva-aca-py-controller-provider:8050 ports: [9091:9090] depends_on: dva-processing-provider: {condition: service_healthy} - rabbit-provider: {condition: service_healthy} postgres-provider: {condition: service_healthy} aca-py-provider: {condition: service_healthy} dva-aca-py-controller-provider: {condition: service_healthy} @@ -22,13 +20,11 @@ services: service: dva-api environment: DVA_PROCESSING_URL: http://dva-processing-consumer:5000 - DVA_RABBITMQ_HOST: rabbit-consumer DVA_POSTGRES_URL: postgresql://postgres-consumer:5432/dva DVA_ACA_PY_CONTROLLER_URL: http://dva-aca-py-controller-consumer:8050 ports: [9092:9090] depends_on: dva-processing-consumer: {condition: service_healthy} - rabbit-consumer: {condition: service_healthy} postgres-consumer: {condition: service_healthy} aca-py-consumer: {condition: service_healthy} dva-aca-py-controller-consumer: {condition: service_healthy} @@ -38,12 +34,10 @@ services: service: dva-api environment: DVA_PROCESSING_URL: http://dva-processing-vla-manager:5000 - DVA_RABBITMQ_HOST: rabbit-vla-manager DVA_POSTGRES_URL: postgresql://postgres-vla-manager:5432/dva ports: [9099:9090] depends_on: dva-processing-vla-manager: {condition: service_healthy} - rabbit-vla-manager: {condition: service_healthy} postgres-vla-manager: {condition: service_healthy} dva-processing-provider: extends: @@ -51,7 +45,6 @@ services: service: dva-processing environment: DVA_LOG_LEVEL: debug - DVA_RABBITMQ_HOST: rabbit-provider DVA_POSTGRES_URL: postgresql://postgres-provider:5432/dva DVA_ACA_PY_CONTROLLER_URL: http://dva-aca-py-controller-provider:8050 healthcheck: @@ -60,7 +53,6 @@ services: timeout: 5s retries: 3 depends_on: - rabbit-consumer: {condition: service_healthy} postgres-provider: {condition: service_healthy} dva-processing-consumer: extends: @@ -68,7 +60,6 @@ services: service: dva-processing environment: DVA_LOG_LEVEL: debug - DVA_RABBITMQ_HOST: rabbit-consumer DVA_POSTGRES_URL: postgresql://postgres-consumer:5432/dva DVA_ACA_PY_CONTROLLER_URL: http://dva-aca-py-controller-consumer:8050 healthcheck: @@ -77,7 +68,6 @@ services: timeout: 5s retries: 3 depends_on: - rabbit-consumer: {condition: service_healthy} postgres-consumer: {condition: service_healthy} dva-processing-vla-manager: extends: @@ -91,18 +81,6 @@ services: interval: 5s timeout: 5s retries: 3 - rabbit-provider: - extends: - file: common-services.yml - service: rabbit - rabbit-consumer: - extends: - file: common-services.yml - service: rabbit - rabbit-vla-manager: - extends: - file: common-services.yml - service: rabbit postgres-provider: extends: file: common-services.yml From 1e4d9e3d2eacc14dc2d9fa32c40edd454081f40a Mon Sep 17 00:00:00 2001 From: Yassine Rhouma Date: Fri, 24 Jul 2026 11:48:49 +0200 Subject: [PATCH 5/6] fix(dva-api): relay upstream status codes instead of masking 400 as 502 --- .../bme/mit/ftsrg/dva/api/route/aovRoutes.kt | 26 ++++++++++++------- 1 file changed, 17 insertions(+), 9 deletions(-) diff --git a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/aovRoutes.kt b/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/aovRoutes.kt index 1c87ac9..f129528 100644 --- a/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/aovRoutes.kt +++ b/dva-api/api/src/main/kotlin/hu/bme/mit/ftsrg/dva/api/route/aovRoutes.kt @@ -3,6 +3,7 @@ package hu.bme.mit.ftsrg.dva.api.route import com.rabbitmq.client.Connection import com.rabbitmq.client.MessageProperties import hu.bme.mit.ftsrg.dva.api.resource.Attestations +import hu.bme.mit.ftsrg.dva.dto.ErrDTO import hu.bme.mit.ftsrg.dva.dto.IDDTO import hu.bme.mit.ftsrg.dva.dto.aov.ACAPyPresentationRequestDTO import hu.bme.mit.ftsrg.dva.dto.aov.ACAPyPresentationResponseDTO @@ -103,18 +104,25 @@ fun Application.aovRoutes() { ) ) } - val acaPyResp: ACAPyPresentationResponseDTO = resp.body() + if (!resp.status.isSuccess()) { + call.respond( + resp.status, + ErrDTO(type = "ACAPY_${resp.status.value}", title = resp.bodyAsText()), + ) + } else { + val acaPyResp: ACAPyPresentationResponseDTO = resp.body() - if (verifLogEntity != null) { - verifsRepo.update( - VerifRequestLogPatch( - id = verifLogEntity.id, - presentationRequestData = acaPyResp.aov, + if (verifLogEntity != null) { + verifsRepo.update( + VerifRequestLogPatch( + id = verifLogEntity.id, + presentationRequestData = acaPyResp.aov, + ) ) - ) - } + } - call.respond(status = resp.status, message = acaPyResp) + call.respond(status = resp.status, message = acaPyResp) + } } } } From cc72d593c4190540ca651f4efa5b0e6c9417a5c6 Mon Sep 17 00:00:00 2001 From: Yassine Rhouma Date: Mon, 27 Jul 2026 12:57:47 +0200 Subject: [PATCH 6/6] test(dva-api): drop obsolete TemplateRoutes/VLARoutes tests relocated to vla-manager-api --- .../ftsrg/dva/api/route/TemplateRoutesTest.kt | 210 ------------------ .../mit/ftsrg/dva/api/route/VLARoutesTest.kt | 159 ------------- 2 files changed, 369 deletions(-) delete mode 100644 dva-api/api/src/test/kotlin/hu/bme/mit/ftsrg/dva/api/route/TemplateRoutesTest.kt delete mode 100644 dva-api/api/src/test/kotlin/hu/bme/mit/ftsrg/dva/api/route/VLARoutesTest.kt diff --git a/dva-api/api/src/test/kotlin/hu/bme/mit/ftsrg/dva/api/route/TemplateRoutesTest.kt b/dva-api/api/src/test/kotlin/hu/bme/mit/ftsrg/dva/api/route/TemplateRoutesTest.kt deleted file mode 100644 index 5ea2206..0000000 --- a/dva-api/api/src/test/kotlin/hu/bme/mit/ftsrg/dva/api/route/TemplateRoutesTest.kt +++ /dev/null @@ -1,210 +0,0 @@ -package hu.bme.mit.ftsrg.dva.api.route - -import hu.bme.mit.ftsrg.dva.api.testutil.createTestClient -import hu.bme.mit.ftsrg.dva.api.testutil.setupTestApplication -import hu.bme.mit.ftsrg.dva.dto.IDDTO -import hu.bme.mit.ftsrg.dva.vla.* -import hu.bme.mit.ftsrg.odcs.DataQuality -import io.ktor.client.call.* -import io.ktor.client.request.* -import io.ktor.http.* -import io.ktor.http.ContentType.* -import io.ktor.http.HttpStatusCode.Companion.Created -import io.ktor.http.HttpStatusCode.Companion.NoContent -import io.ktor.http.HttpStatusCode.Companion.NotFound -import io.ktor.http.HttpStatusCode.Companion.OK -import io.ktor.server.application.install -import io.ktor.server.testing.* -import kotlinx.serialization.json.buildJsonObject -import kotlinx.serialization.json.put -import kotlinx.serialization.json.putJsonObject -import org.junit.jupiter.api.Assertions.assertEquals -import org.junit.jupiter.api.Assertions.assertNotNull -import org.junit.jupiter.api.Test -import org.koin.dsl.module -import org.koin.ktor.plugin.Koin -import kotlin.uuid.ExperimentalUuidApi -import kotlin.uuid.Uuid - -@OptIn(ExperimentalUuidApi::class) -class TemplateRoutesTest { - - @Test - fun `should return list of templates`() = testApplication { - setupApplication() - - val client = createTestClient() - client.get("/template").apply { - assertEquals(OK, status) - // Fake template repository seeds itself with 1 template - assertEquals(1, body>().size) - } - } - - @Test - fun `should return template by ID`() = testApplication { - setupApplication() - - val client = createTestClient() - client.get("/template/${Uuid.NIL}").apply { - assertEquals(OK, status) - assertNotNull(body