ANSP-role client for consuming Digital NOTAM events via SWIM infrastructure. Connects to one or more SWIM Subscription Managers, creates subscriptions based on configuration, receives DNOTAM events via AMQP 1.0 over TLS, validates against AIXM 5.1.1 XSD, persists to MongoDB, and distributes to Kafka topics by business intent.
- Configuration-driven subscriptions, declares desired subscriptions as JSON, reconciles them automatically against the provider
- Multi-provider support, connects to multiple SWIM providers simultaneously, each with independent mTLS and AMQP configuration
- AMQP 1.0 consumption, receives DNOTAM events with mTLS, backpressure control, and batch staging via Kafka inbox
- XSD validation, validates incoming AIXM 5.1.1 messages against schema before processing
- MongoDB persistence, stores events and subscriptions with idempotency cache and dead letter queue
- Kafka distribution, routes events to 6 topic categories by business intent (closures, restrictions, surface conditions, airspace, hazards, unknown)
- Subscription renewal, automatic renewal before expiration with configurable threshold
- Heartbeat monitoring, detects provider heartbeat loss per subscription
- Outbox pattern, reliable event delivery with recovery and cleanup schedulers
- GraphQL API, query events by scenario, airport, time range
- Observability, OpenTelemetry tracing, Prometheus metrics, structured logging
-
Java 21
-
Maven 3.9+
-
Podman Desktop, includes the Podman engine and a graphical interface for managing containers and compose stacks. Any OCI-compatible runtime with Compose support also works.
-
mkcert, local certificate authority
macOS
brew install mkcert
Fedora / RHEL
sudo dnf install nss-tools curl -Lo mkcert https://dl.filippo.io/mkcert/latest?for=linux/amd64 chmod +x mkcert sudo mv mkcert /usr/local/bin/Debian / Ubuntu
sudo apt install libnss3-tools curl -Lo mkcert https://dl.filippo.io/mkcert/latest?for=linux/amd64 chmod +x mkcert sudo mv mkcert /usr/local/bin/Linux (Homebrew)
brew install mkcert
Windows (Chocolatey)
choco install mkcert choco install openssl
Windows (Scoop)
scoop bucket add extras scoop install mkcert scoop install openssl
On Windows,
opensslis also bundled with Git for Windows and is usually already on the PATH.
The consumer connects to the validator's Artemis broker via mTLS on port 5671 by default. Generate the certificates before starting the infrastructure.
macOS / Linux
./certs/generate.shWindows (PowerShell)
.\certs\generate.ps1This is a one-time step per machine. It uses mkcert to create a local CA (installed into your system trust store), then generates:
certs/broker.p12: Artemis broker keystore, mounted into the validator containercerts/ca-truststore.p12: CA truststore for Artemis, used to verify the consumer client certcerts/keystore.jks: consumer client keystore, used by the Quarkus dev profilecerts/truststore.jks: consumer CA truststore, used by the Quarkus dev profile
The broker certificate covers localhost, 127.0.0.1, dnotam-consumer-validator-artemis, artemis.127.0.0.1.nip.io, and dnotam-consumer-validator-artemis.127.0.0.1.nip.io. No /etc/hosts entries needed.
The compose.yml at the root of this project brings up everything the consumer needs: a fake SWIM provider (Subscription Manager API + AMQP broker + AIXM event generator), MongoDB, and Kafka.
The Artemis broker is built from src/local-dev/artemis/. Use --build the first time, or after any change to those files:
podman compose up --build -dOn subsequent runs, when nothing in src/local-dev/artemis/ has changed, the cached image is used:
podman compose up -dTip: Quarkus Dev Services can also provision databases and brokers automatically during
./mvnw quarkus:dev, as an alternative to thecompose.yml.
Services started:
| Service | Port | Description |
|---|---|---|
dnotam-consumer-validator |
8084 | Mock Subscription Manager REST API + event generator |
dnotam-consumer-validator-artemis |
5671 (AMQPS/mTLS), 5672 (plain), 8161 | AMQP broker (fake provider side) |
dnotam-consumer-validator-mariadb |
3307 | Validator persistence |
dnotam-consumer-mongodb |
27017 | Consumer event store |
dnotam-consumer-mongo-express |
9081 | MongoDB web UI |
kafka |
9092 | Kafka broker (KRaft) |
dnotam-consumer-akhq |
9080 | Kafka web UI (AKHQ) |
The dnotam-consumer-validator simulates a real SWIM provider: it exposes the Subscription Manager REST API, manages an Artemis broker, and periodically publishes real AIXM DNOTAM events. The consumer connects to it exactly as it would connect to a production EUROCONTROL provider.
The compose uses the pre-built image quay.io/masales/swim-dnotam-consumer-validator:latest, pulled fresh on every podman compose up.
If you want to run the validator from source instead of the pre-built image, clone swim-dnotam-consumer-validator, start its infrastructure separately, and run:
./mvnw quarkus:devThen update src/main/resources/application-dev.properties to point swim.providers at the locally running validator instead of the compose container.
./mvnw quarkus:devAdd -Ddebug=false to skip the remote debug port (5005) if you do not need a debugger attached:
./mvnw quarkus:dev -Ddebug=falseTo make the REST API accessible from other machines on the same network (useful on headless servers or VMs), bind to all interfaces:
./mvnw quarkus:dev -Dquarkus.http.host=0.0.0.0 -Ddebug=falseThe API will be reachable at http://<your-ip>:8080.
The dev profile connects to the infrastructure above. The consumer will subscribe to the validator, receive AIXM events, persist them to MongoDB, and publish to Kafka.
| URL | Description |
|---|---|
| http://localhost:8080 | REST API |
| http://localhost:8080/swagger-ui | Swagger UI |
| http://localhost:8080/q/graphql-ui | GraphQL playground |
| http://localhost:8080/q/health | Health |
| http://localhost:8161 | Artemis console (admin / admin) |
| http://localhost:9081 | MongoDB UI |
| http://localhost:9080 | Kafka UI (AKHQ) |
Wait ~30 seconds for the consumer to reconcile subscriptions with the validator. Then:
# Check overall health
curl -s http://localhost:8080/q/health | jq .status
# List subscriptions, at least one should be ACTIVE
curl -s http://localhost:8080/api/v1/subscriptions | jq '.[] | {id, status, topic}'
# List subscriptions, confirmed ACTIVE means events are flowing
curl -s "http://localhost:8080/api/v1/subscriptions" | jq '.[] | {id, status, topic}'Or use GraphQL:
curl -s -X POST http://localhost:8080/graphql \
-H "Content-Type: application/json" \
-d '{"query":"{ dnotamEvents(limit: 5) { eventId notamNumber validTimeStart } }"}' | jq .The consumer is working correctly when:
GET /q/health/readyreturns{"status":"UP"}- At least one subscription shows
"status":"ACTIVE" - Events appear in MongoDB (http://localhost:9081 →
swim-dnotamdb →eventscollection) - Events appear in Kafka topics (http://localhost:9080, AKHQ)
The default dev profile connects to the validator's Artemis on port 5671 using mTLS. This matches production behaviour and is the recommended approach.
If you need to temporarily bypass mTLS (for example, to isolate a connectivity issue), edit application-dev.properties to switch to the plain AMQP port:
# switch provider amqpBroker port from 5671 to 5672 and set sslEnabled to false
swim.providers=[{"providerId":"validator","subscriptionManager":{"url":"http://localhost:8084","tls":null,...},"amqpBroker":{"host":"localhost","port":5672,"sslEnabled":false,"username":"admin","password":"admin","tls":null}}]
swim.amqp.host=localhost
swim.amqp.port=5672
swim.amqp.ssl.enabled=falseThe plain AMQP port 5672 on the validator container is always available alongside 5671. Revert to mTLS once the issue is identified.
Events are routed to topics based on business intent:
| Topic | Scenarios | Use case |
|---|---|---|
dnotam-events-closure-topic |
RWY.CLS, AD.CLS, TWY.CLS, APN.CLS, STAND.CLS |
Flight cancellation/diversion |
dnotam-events-restriction-topic |
RWY.LIM, AD.LIM, RCP.CHG, STAND.LIM, STAND.STS |
Weight/performance calculation |
dnotam-events-surface-condition-topic |
SFC.CON |
Braking action, landing distance |
dnotam-events-airspace-topic |
SAA.ACT, SAA.NEW |
Flight routing |
dnotam-events-hazards-navaids-topic |
OBS.NEW, NAV.UNS, WLD.HZD |
Approach procedures |
dnotam-events-others-topic |
UNKNOWN |
Monitoring, schema drift detection |
| Method | Endpoint | Description |
|---|---|---|
GET |
/api/v1/subscriptions |
List all subscriptions |
GET |
/api/v1/subscriptions/active |
List active subscriptions |
POST |
/api/v1/subscriptions |
Create manual subscription |
PUT |
/api/v1/subscriptions/{id} |
Update status (ACTIVE/PAUSED) |
DELETE |
/api/v1/subscriptions/{id} |
Delete subscription |
GET |
/api/v1/topics |
List topics from ConfigMap |
GET |
/api/v1/subscriptions/{id}/events |
List events (paginated) |
GET |
/api/v1/subscriptions/{id}/events/count |
Count events |
GET |
/api/v1/subscriptions/{id}/events/range |
Events by date range |
GET |
/api/v1/events/{messageId} |
Get event by AMQP message ID |
GET |
/api/v1/dlq |
List dead letter queue (paginated) |
GET |
/api/v1/dlq/count |
Count DLQ messages |
GET |
/api/v1/stats |
Consumer statistics |
Swagger UI available at /swagger-ui.
A ready-to-import Postman collection is available in src/test/postman/:
| File | Description |
|---|---|
SWIM-DNOTAM-Consumer.postman_collection.json |
Full REST API coverage: subscriptions, events, operational |
SWIM-Local.postman_environment.json |
Local environment: consumer_url=http://localhost:8080, validator_url=http://localhost:8084 |
- Open Postman → Import → select both files.
- Select SWIM Local as the active environment.
- Run requests in order following the demo flow described in the collection description.
- Subscriptions / Create — creates a subscription and stores
sub_idautomatically. - Subscriptions / List all — confirms the subscription exists.
- Subscriptions / Pause — pauses; verify with List active (should not appear).
- Subscriptions / Resume — resumes; verify with List active (appears again).
- Events / List by subscription — once DNOTAM events arrive, browse them here.
- Operational / Consumer statistics — overall health and counters.
- Operational / DLQ count — must be
0on a healthy consumer. - Subscriptions / Delete — clean up.
All requests can also be run with curl. Examples:
# Create a subscription
curl -s -X POST http://localhost:8080/api/v1/subscriptions \
-H "Content-Type: application/json" \
-d '{"topic":"DNOTAM/v1","eventScenario":["RWY.CLS"],"airportHeliport":["EADD"],"description":"Test"}' | jq .
# List all subscriptions
curl -s http://localhost:8080/api/v1/subscriptions | jq .
# Consumer statistics
curl -s http://localhost:8080/api/v1/stats | jq .
# DLQ count (should be 0)
curl -s http://localhost:8080/api/v1/dlq/count | jq .Available at /graphql with interactive playground at /q/graphql-ui.
query {
dnotamEvents(scenario: "RWY.CLS", airport: "LPPT", limit: 50) {
eventId
notamNumber
validTimeStart
validTimeEnd
}
}The full provider configuration is a JSON array set via SWIM_PROVIDERS. Each entry configures one external SWIM provider.
[
{
"providerId": "my-provider",
"subscriptionManager": {
"url": "https://sm.provider.example",
"tls": {
"trustStorePath": "/certs/truststore.jks",
"trustStorePassword": "changeit",
"keyStorePath": "/certs/keystore.jks",
"keyStorePassword": "changeit"
}
},
"amqpBroker": {
"host": "amqp.provider.example",
"port": 5671,
"sslEnabled": true,
"tls": {
"trustStorePath": "/certs/truststore.jks",
"trustStorePassword": "changeit",
"keyStorePath": "/certs/keystore.jks",
"keyStorePassword": "changeit"
}
}
}
][
{
"topic": "DNOTAM/v1",
"provider": "my-provider",
"eventScenario": ["RWY.CLS", "AD.CLS"],
"airportHeliport": ["LPPT", "EHAM"],
"description": "Runway and aerodrome closures"
}
]| Variable | Default | Description |
|---|---|---|
SWIM_PROVIDERS |
[{...localhost...}] |
JSON array of SWIM provider configurations (SM URL + AMQP broker) |
DNOTAM_SUBSCRIPTIONS |
: | JSON array of desired subscriptions |
DNOTAM_DELETE_AND_RECREATE |
true |
Recreate subscriptions on startup |
MONGODB_URI |
mongodb://localhost:27017 |
MongoDB connection string |
MONGODB_DATABASE |
swim-dnotam |
MongoDB database name |
KAFKA_BOOTSTRAP_SERVERS |
kafka-kafka-bootstrap:9092 |
Kafka bootstrap servers |
SWIM_VALIDATION_ENABLED |
true |
Enable XSD validation against AIXM 5.1.1 |
SWIM_VALIDATION_FAIL_ON_NULLBODY |
false |
Reject messages with null body |
SWIM_RENEWAL_ENABLED |
true |
Enable automatic subscription renewal |
SWIM_RENEWAL_CHECK_INTERVAL |
5m |
How often to check for expiring subscriptions |
SWIM_RENEWAL_THRESHOLD |
1h |
Renew when less than this time remains |
RECONCILIATION_RETRY_INTERVAL |
10s |
Retry interval for subscription reconciliation |
RECONCILIATION_RETRY_INITIAL_DELAY |
30s |
Initial delay before first reconciliation attempt |
SWIM_SCHEDULER_INITIAL_DELAY |
30s |
Startup delay before schedulers begin |
SWIM_OUTBOX_RECOVERY_INTERVAL |
30s |
Outbox recovery sweep interval |
SWIM_OUTBOX_MAX_RETRIES |
5 |
Max retry attempts for outbox delivery |
SWIM_IDEMPOTENCY_CACHE_SIZE |
100000 |
Max entries in idempotency cache |
SWIM_IDEMPOTENCY_CACHE_TTL |
24H |
Cache TTL for processed message IDs |
OTEL_ENABLED |
true |
Enable OpenTelemetry tracing |
OTEL_ENDPOINT |
http://localhost:4317 |
OTLP collector endpoint |
PROMETHEUS_ENABLED |
true |
Enable Prometheus metrics at /q/metrics |
VERTX_WORKER_POOL_SIZE |
100 |
Vert.x worker thread pool size |
THREAD_POOL_MAX_THREADS |
250 |
Max application thread pool size |
Pre-built multi-arch images (linux/amd64 + linux/arm64):
quay.io/masales/swim-dnotam-consumer:latest
Run with Podman (or any OCI runtime):
podman run -p 8080:8080 \
-e MONGODB_URI=mongodb://host:27017 \
-e KAFKA_BOOTSTRAP_SERVERS=host:9092 \
-e SWIM_PROVIDERS='[{...}]' \
-e DNOTAM_SUBSCRIPTIONS='[{...}]' \
quay.io/masales/swim-dnotam-consumer:latest./mvnw clean package -DskipTestsmake jvm # JVM multi-arch image, build + push (fastest)
make native-amd64 # Native amd64, build + push (run on amd64 machine)
make native-arm64 # Native arm64, build + push (run on arm64 machine)
make manifest # Create multi-arch manifest from registry images
make push # Push manifest to registryOverride registry or tag: make jvm REGISTRY=quay.io/myorg TAG=v1.2.3
Run make deps to see which sibling repos to install first.
Run fast, no containers, no I/O:
mvn testRun with real infrastructure via Testcontainers (MongoDB, Kafka/Redpanda, Artemis, WireMock). These test the full consumer pipeline in a JVM process.
mvn verify -DskipITs=falseRun one project at a time. Integration tests bind to host ports; running multiple projects in parallel causes port conflicts.
Compile the application to a GraalVM native executable and run the test suite against the binary. This catches issues that JVM-mode tests cannot detect:
| Issue | Why JVM misses it | Caught by native IT |
|---|---|---|
| CDI raw-type injection | ArC wires types lazily in JVM | Wired at build time in native; startup fails if broken |
| MongoDB POJO codec | Reflection available by default in JVM | Classes without @RegisterForReflection excluded in native |
| Quarkus cache registration | Caches registered on first use in JVM | Must be declared at build time in native via @CacheResult |
mvn verify -PnativeNative compilation takes 5-15 minutes depending on the machine. Run only when validating a new image before deployment or after changes to CDI beans, MongoDB projections, or cache configuration.
The native IT tests live in src/test/java/.../native_it/DnotamConsumerNativeIT.java and are skipped by default (skipITs=true). The native Maven profile sets skipITs=false and enables quarkus.native.enabled=true.
| Endpoint | Description |
|---|---|
/q/health/live |
Liveness probe |
/q/health/ready |
Readiness probe |
/q/health |
Combined status |
Helm chart in src/main/helm/ with CRC and production values.
For operator-based deployment (single CR), see swim-operator.
| Project | Why you need it |
|---|---|
| swim-digital-notam-consumer-validator | Fake SWIM provider for local testing: included in compose.yml |
| swim-developer-extensions | Kafka outbox routers that forward events from this consumer to downstream systems |
| swim-aixm-model | AIXM 5.1.1 JAXB bindings used internally |
| swim-developer-framework | Core framework this service is built on |
| swim-developer-tools | Infrastructure tooling: cert generation, full-stack compose, pipelines |
Licensed under the Apache License 2.0.