Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Aviso Logo Aviso Logo

Introduction

Aviso Server is a real-time notification and historical replay system for data dissemination pipelines. Publishers announce data availability events; subscribers receive those events live as they happen, or replayed from a chosen point in history.

It is designed for environments where timely, reliable notification is critical: scientific computing, operational weather forecasting, large-scale data distribution, and similar domains.


How It Works

At its core, Aviso Server exposes three operations:

graph LR
    P(Publisher) -->|POST /api/v1/notification| A[Aviso Server]
    A -->|store| B[("Backend<br/>JetStream / In-Memory")]
    B -->|fan-out| W1("Subscriber A<br/>watch")
    B -->|fan-out| W2("Subscriber B<br/>watch")
    B -->|history| R("Client<br/>replay")
OperationEndpointDescription
NotifyPOST /api/v1/notificationPublish a notification event to the backend
WatchPOST /api/v1/watchStream live (and optionally historical) events over SSE
ReplayPOST /api/v1/replayStream historical events only, then close

Key Features

Schema-driven validation

Each event type can have a schema that defines which identifier fields are required, their types (date, time, integer, float, enum, polygon), allowed ranges, and topic ordering. Invalid notifications are rejected at the API boundary with a clear error.

Structured topic routing

Aviso builds a deterministic topic string from the notification’s identifier fields. This topic is used to route, store, and filter messages in the backend. Subscribers can use wildcard patterns and constraint objects to filter the stream.

Hybrid filtering

Watch and replay requests are matched using a two-tier strategy: the backend handles coarse routing (broad subject filters), and Aviso applies precise application-level filtering (constraints, spatial checks) on top. This keeps backend subscription counts low while delivering exact results.

Spatial awareness

Polygon and point identifiers are first-class: notifications can carry geographic polygons, and subscribers can filter by polygon intersection or point containment.

Pluggable backends

Aviso abstracts storage behind a NotificationBackend trait. Today two backends ship: JetStream (NATS-backed, durable, production-ready) and In-Memory (single-process, for development and testing).

Server-Sent Events (SSE)

Watch and replay streams use SSE, a simple HTTP streaming protocol supported natively by browsers and all major HTTP clients. The stream includes typed control frames (connection established, replay started/completed, heartbeats, and graceful close reasons).

CloudEvents format

All delivered notifications follow the CloudEvents specification, making them easy to integrate with other event-driven systems.


Use Cases

  • Data availability monitoring: trigger downstream workflows the moment a dataset lands.
  • Operational pipelines: coordinate processing steps across distributed services.
  • Audit and compliance: replay historical events to reconstruct what was published and when.
  • System integration: connect disparate systems through a standardized event interface.

Where to Go Next

If you are new to Aviso, read these pages in order:

  1. Key Concepts: understand the terminology before anything else.
  2. Installation: get the server running.
  3. Getting Started: send your first notification and watch it arrive.
  4. Practical Examples: copy-paste workflows for common scenarios.

If you are configuring for production:

Key Concepts

These are the core terms used throughout Aviso’s documentation and API. Read this page before Getting Started to make the commands easier to follow.


Event Type

An event type is a named category of notification, for example extreme_event, sensor_alert, or data_ready.

Every request to Aviso (notify, watch, or replay) targets exactly one event type. The server uses the event type to:

  • look up the matching schema (validation rules, required fields, topic ordering),
  • route messages to the correct storage stream,
  • apply the correct retention and storage policy.
{ "event_type": "extreme_event", ... }

When a notification_schema is configured, Aviso is strict by default: any event_type that is not in the schema is rejected with 400 UNKNOWN_EVENT_TYPE on /api/v1/notification, /api/v1/watch, and /api/v1/replay. The error body lists the configured event types so clients can self-correct.

When notification_schema is empty or absent (no schema declared at all), Aviso falls back to generic behavior: any event type is accepted, fields are treated as-is, and the topic is built from sorted keys. This mode is intended for local development and quick experiments.

Operators can choose the behavior with notification_schema_strict. In the table below, Schemas refers to notification_schema, and Strict mode refers to notification_schema_strict.

SchemasStrict modeBehavior
DefinedUnsetUnknown event types return HTTP 400.
Empty or absentUnsetGeneric fallback.
AnyfalseGeneric fallback.
AnytrueUnknown event types return HTTP 400.

With strict mode enabled and no schemas defined, all event types are rejected. With strict mode disabled and schemas defined, Aviso logs a startup warning.

Independent of the strict-mode knob, the event_type value that ends up on Prometheus labels and tracing span fields is always bounded: requests whose event_type is not in the configured schema have their observability label collapsed to the literal "generic". In strict mode this collapsing is rarely exercised because unknown event_types are already rejected upstream with 400 UNKNOWN_EVENT_TYPE; in permissive generic-fallback mode it is the main mechanism preventing unbounded label cardinality from user-controlled input.


Identifier

An identifier is a set of key-value pairs that describe what a notification is about. Think of it as structured metadata that uniquely (or approximately) locates a piece of data.

{
  "identifier": {
    "region": "north",
    "run_time": "1200",
    "severity": "4",
    "anomaly": "42.5"
  }
}

Identifiers serve two purposes depending on the operation:

OperationRole of identifier
notifyDeclares the metadata of the notification being published
watchActs as a filter for which notifications to receive
replayActs as a filter for which historical notifications to retrieve

In watch and replay, identifier values can be constraint objects instead of scalars (e.g. {"gte": 5}) for numeric and enum fields. See Streaming Semantics.


Topic

A topic is the internal routing key Aviso builds from an identifier. You rarely construct topics manually; Aviso builds them for you.

Topics are dot-separated strings, for example:

extreme_event.north.1200.4.42%2E5

Each token corresponds to one identifier field, in the order defined by key_order in the schema. The %2E is the dot in 42.5, percent-encoded so it doesn’t conflict with the topic separator. Reserved characters (., *, >, %) in field values are percent-encoded before writing to the backend so they do not break routing. See Topic Encoding.


Schema

A schema configures how Aviso handles a specific event type. Schemas are defined in configuration/config.yaml under notification_schema.

A schema controls:

  • topic.base: the stream/prefix for this event type (e.g. diss, mars)
  • topic.key_order: the order of identifier fields in the topic string
  • identifier.*: validation rules per field (type, required, allowed values, ranges)
  • payload.required: whether a payload is mandatory on notify
  • storage_policy: per-stream retention, size limits, compression (JetStream only)

Example:

notification_schema:
  extreme_event:
    topic:
      base: "extreme_event"
      key_order: ["region", "run_time", "severity", "anomaly", "polygon"]
    identifier:
      region:
        description: "Geographic region label."
        type: EnumHandler
        values: ["north", "south", "east", "west"]
        required: true
      run_time:
        type: TimeHandler
        required: true
      severity:
        description: "Severity level from 1 to 7."
        type: IntHandler
        range: [1, 7]
        required: true
      anomaly:
        type: FloatHandler
        range: [0.0, 100.0]
        required: false
      polygon:
        type: PolygonHandler
        required: false
    payload:
      required: false

Every key listed in key_order must have a corresponding entry in identifier. For notify, every identifier key declared in the schema must be present in the request and every value must pass the handler’s validation (an empty string, an out-of-range integer, an unparseable date, and so on are all rejected). The required flag has no effect on notify; it only changes the behavior of watch and replay. There, a missing key marked required: true returns 400, while a missing key marked required: false is treated as a wildcard. The flag never relaxes value validation; only the presence rules.


Payload

A payload is arbitrary JSON attached to a notification. Aviso treats it as opaque; it stores and replays the value exactly as sent.

Valid payload types: object, array, string, number, boolean, or null.

{ "payload": { "path": "/data/grib2/file.grib2", "size": 1048576 } }

Whether payload is required or optional is controlled per schema by payload.required. See Payload Contract for the full input → storage → replay mapping.


Operations: Notify, Watch, Replay

Aviso exposes three operations, each on its own endpoint:

Notify: POST /api/v1/notification

Publishes a notification to the backend. The identifier must match all required schema fields exactly (no wildcards, no constraints).

Watch: POST /api/v1/watch

Opens a persistent Server-Sent Events (SSE) stream. Receives live notifications as they arrive, optionally starting from a historical point.

  • Omit from_id/from_date for a live-only stream.
  • Provide one of them for historical replay first, then live.

Replay: POST /api/v1/replay

Opens a finite SSE stream of historical notifications only. Requires exactly one of from_id or from_date, and accepts at most one of to_id or to_date to end earlier. Stream closes automatically when history is exhausted or the end point is reached.

For end-to-end working examples of all three operations, including spatial and constraint filtering, see Practical Examples.


CloudEvents

Aviso delivers notifications to watch/replay clients as CloudEvents, a standard envelope format. Each event includes:

  • id: the backend sequence reference (e.g. mars@42), used for targeted delete or resume.
  • type: the Aviso event type string, prefixed with int.ecmwf.aviso. (for example int.ecmwf.aviso.mars).
  • source: the server base URL.
  • data.identifier: the canonicalized identifier.
  • data.payload: the notification payload (or null if omitted).

Backend

The backend is the storage and messaging layer that Aviso delegates to. Two backends are supported:

BackendUse case
jetstreamProduction: durable, replicated, persistent history
in_memoryDevelopment: fast setup, no persistence

The backend is selected via notification_backend.kind in config. See Backends Overview.

Installation

This page covers every way to get Aviso Server running: building from source, using Docker, or deploying to Kubernetes with Helm.


Prerequisites

Rust toolchain

Aviso requires Rust 1.88 or newer. Each published crate declares this minimum supported Rust version in its package metadata. CI also builds with the repository’s newer pinned toolchain.

Install or update Rust via rustup:

curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh

Verify your version:

rustc --version
# rustc 1.88.0 (... ) or newer

System dependencies

Aviso does not require OpenSSL. HTTPS uses rustls with AWS-LC. Building AWS-LC from source requires a native C/C++ toolchain, CMake, and Perl. The runtime uses the platform’s native certificate store, so Linux installations also need a current CA certificate bundle.

On Debian/Ubuntu:

sudo apt-get install -y build-essential cmake perl ca-certificates

On Fedora/RHEL:

sudo dnf install -y gcc gcc-c++ cmake perl ca-certificates

On macOS, install the Command Line Tools and CMake. macOS supplies Perl and the native certificate store:

xcode-select --install
brew install cmake

Build from Source

Clone the repository:

git clone https://github.com/ecmwf/aviso-server.git
cd aviso-server

Development build

Fast to compile, includes debug symbols:

cargo build

Binary location: target/debug/aviso_server

Release build

Optimized for production use:

cargo build --release

Binary location: target/release/aviso_server

Run directly

cargo run                          # development
cargo run --release                # release
./target/release/aviso_server      # pre-built binary

The server loads ./configuration/config.yaml by default. See Configuration for all config loading options.


Docker

The repository includes a multi-stage Dockerfile that produces a minimal distroless image.

Build the image

# Production image (distroless, minimal attack surface)
docker build --target release -t aviso-server:local .

# Debug image (Debian slim, includes bash for troubleshooting)
docker build --target debug -t aviso-server:debug .

Run with Docker

Mount your config file and expose the port:

docker run --rm \
  -p 8000:8000 \
  -v $(pwd)/configuration/config.yaml:/app/configuration/config.yaml:ro \
  aviso-server:local

Or override settings via environment variables (no config mount needed):

docker run --rm \
  -p 8000:8000 \
  -e AVISOSERVER_APPLICATION__HOST=0.0.0.0 \
  -e AVISOSERVER_APPLICATION__PORT=8000 \
  -e AVISOSERVER_NOTIFICATION_BACKEND__KIND=in_memory \
  aviso-server:local

Build targets summary

TargetBase imageSizeUse
releasedistroless/ccminimalProduction
debugdebian:bookworm-slimlargerTroubleshooting

Local JetStream (Docker)

For local development with the JetStream backend, use the provided script to spin up a NATS server with JetStream enabled:

./scripts/init_nats.sh

This script:

  • Generates a private config under ${XDG_STATE_HOME:-$HOME/.local/state}/aviso/
  • Uses a persistent Docker volume named ${CONTAINER_NAME}-data
  • Starts a nats:2.14.6-alpine container on loopback port 4222
  • Waits for the server to be ready and prints a connection summary

Requires: Docker and the nats CLI. Token generation also uses OpenSSL.

Optional environment variables:

NATS_PORT=4222            # NATS client port (default: 4222)
ENABLE_AUTH=true          # Enable token auth (default: false)
MAX_MEMORY=5GB            # JetStream memory limit (default: 5GB)
MAX_STORAGE=10GB          # JetStream file storage limit (default: 10GB)

Example with auth enabled:

ENABLE_AUTH=true ./scripts/init_nats.sh

The script stores the generated token in the private configuration file and never prints it. Supply TOKEN securely to choose your own token; it may contain letters, digits, underscores and hyphens.

Use NATS_IMAGE to override the image. CONFIG_DIR accepts an absolute or relative path; SKIP_DOCKER=1 generates configuration without starting a server. Existing containers are never removed automatically. Choose a new CONTAINER_NAME, or stop and remove the old container explicitly to reuse its data volume. NATS_VOLUME selects an existing volume when needed.

For an isolated test broker, give it separate resources and unused ports:

CONTAINER_NAME=aviso-test DOCKER_NETWORK=aviso-test \
CONFIG_DIR=/tmp/aviso-test NATS_VOLUME=aviso-test-data \
NATS_PORT=14222 NATS_HTTP_PORT=18222 NATS_CLUSTER_PORT=16222 \
./scripts/init_nats.sh
AVISO_RUN_NATS_TESTS=1 NATS_URL=nats://127.0.0.1:14222 cargo test --locked

Run opt-in tests only against a disposable broker, never shared streams. An unreachable broker fails the opted-in tests. The script binds published ports to 127.0.0.1; set NATS_BIND_ADDRESS explicitly to expose another interface. Readiness checks use that address. Wildcard binds (0.0.0.0 or ::) use the corresponding loopback address instead. IPv6 addresses may be bracketed or bare.

After the script completes, configure Aviso to connect:

Without auth (default):

notification_backend:
  kind: jetstream
  jetstream:
    nats_url: "nats://localhost:4222"

With auth, pass the token stored in the private configuration file:

notification_backend:
  kind: jetstream
  jetstream:
    nats_url: "nats://localhost:4222"
    token: "aviso_secure_token_1712345678"

Alternatively, set the token as an environment variable (Aviso reads NATS_TOKEN automatically):

export NATS_TOKEN=aviso_secure_token_1712345678
cargo run

Kubernetes / Helm

For production Kubernetes deployments, use the official Helm chart:

The chart handles:

  • Deployment with configurable replicas
  • ConfigMap-based configuration mounting
  • Service and Ingress setup
  • JetStream connection settings via values

Build Documentation

Aviso docs are built with mdBook.

Install mdBook and the mermaid preprocessor:

cargo install mdbook
cargo install mdbook-mermaid

Serve docs locally with live reload:

mdbook serve docs --open

Build static output to docs/book/:

mdbook build docs

Run Tests

# Unit and integration tests (in-memory backend)
cargo test --workspace

# Include JetStream integration tests (requires running NATS)
AVISO_RUN_NATS_TESTS=1 cargo test --workspace

# Tests must run single-threaded (shared port binding)
cargo test -- --test-threads=1

Getting Started

This guide walks you through running Aviso Server locally and sending your first notification. It assumes you have already completed Installation.

If you haven’t read Key Concepts yet, do that first; it will make the commands below much easier to follow.


1. Choose a Backend

Aviso needs a storage backend before it can accept notifications. For local exploration, in-memory requires zero infrastructure.

GoalBackend
Quick local test, no setupin_memory
Persistent history, realistic behaviorjetstream

The examples on this page use in_memory. To use JetStream locally, see Local JetStream setup below.


2. Configure the Server

Open configuration/config.yaml and make sure it contains at minimum:

application:
  host: "127.0.0.1"
  port: 8000
  base_url: "http://localhost:8000"

notification_backend:
  kind: in_memory
  in_memory:
    max_history_per_topic: 100
    max_topics: 10000

notification_schema:
  my_event:
    topic:
      base: "my_event"
      key_order: ["region", "date"]
    identifier:
      region:
        description: "Geographic region label."
        type: EnumHandler
        values: ["north", "south", "east", "west"]
        required: true
      date:
        type: DateHandler
        required: true
    payload:
      required: false

You can also use environment variables to override any config value without editing the file. See Configuration for the full precedence rules.


3. Start the Server

cargo run

Or with the release binary:

./target/release/aviso_server

You should see structured JSON log output. Once you see a line like:

{
  "level": "INFO",
  "message": "aviso-server listening",
  "address": "127.0.0.1:8000"
}

the server is ready.


4. Check the Health Endpoint

curl -sS http://127.0.0.1:8000/health

Expected response: 200 OK

/health is process liveness only. /ready additionally reflects the notification backend connection (200 when connected, 503 while down or reconnecting) and is the endpoint to use for Kubernetes readiness probes:

curl -sS http://127.0.0.1:8000/ready

5. Open a Watch Stream

Before publishing, open a terminal and start watching for events. This is a live SSE stream; keep it open while you proceed to step 6.

curl -N -X POST "http://127.0.0.1:8000/api/v1/watch" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type": "my_event",
    "identifier": {
      "region": "north",
      "date":   "20250706"
    }
  }'

You will see the SSE connection frame immediately:

event: live-notification
data: {"connection_will_close_in_seconds":3600,"request_id":"a3f1d2c8-9b4e-4f7a-bd56-1c8e2a9d4e3f","timestamp":"2026-03-04T10:00:00Z","topic":"my_event.north.20250706","type":"connection_established"}

The request_id here is unique to this watch request. Save it: you will compare it with the next request’s request_id to confirm each call gets its own.

The stream stays open and will print new events as they arrive.


6. Publish a Notification

In a second terminal, send a notification:

curl -sS -X POST "http://127.0.0.1:8000/api/v1/notification" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type": "my_event",
    "identifier": {
      "region": "north",
      "date":   "20250706"
    },
    "payload": { "note": "data is ready" }
  }'

Expected response:

{
  "status": "success",
  "request_id": "0d4f6758-1ce3-4dda-a0f3-0ccf5fcb50d6",
  "processed_at": "2026-03-04T10:00:00Z"
}

request_id is the per-request UUID for this notify call. Note that it is different from the watch call’s UUID above (a3f1d2c8-...); each HTTP request gets its own. The same value appears in the X-Request-ID HTTP response header and in the corresponding server log lines. Quote it when reporting issues.

Switch back to the watch terminal. The notification should have arrived as a CloudEvent body:

event: live-notification
data: {"data":{"identifier":{...},"payload":{"note":"data is ready"}},"datacontenttype":"application/json","dataschema":"http://localhost:8000/schema/my_event","id":"my_event@1","source":"http://localhost:8000","specversion":"1.0","time":"2026-03-04T10:00:00.123456Z","type":"int.ecmwf.aviso.my_event"}

The source and dataschema fields are derived from the server’s application.base_url. With the default config (no base_url set), they would use http://localhost; the example above matches the snippet in step 2 which sets it to http://localhost:8000. The CloudEvent body also includes the canonicalized identifier and the payload under data.

The CloudEvent id (my_event@1) is the <event_type>@<sequence> reference for replay and admin delete. The replay endpoint takes only the numeric sequence ("1"), not the full string.


7. Replay History

Once you have published a few notifications, you can replay them from a specific point:

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type": "my_event",
    "identifier": {
      "region": "north",
      "date":   "20250706"
    },
    "from_id": "1"
  }'

The stream will emit all matching historical notifications, then close with:

event: connection-closing
data: {"message":"Stream completed","reason":"end_of_stream","request_id":"f7e8d910-2b3c-4d5a-9c8b-1234567890ab","timestamp":"...","topic":"my_event.north.20250706"}

This request_id is yet another distinct UUID, since /replay is a separate HTTP request from /watch and /notification.


Optional: Local JetStream Setup

To test with the JetStream backend (persistent storage, more realistic), see Installation: Local JetStream for the full setup including environment variables and token authentication.


Run the Smoke Test

A Python smoke script covers the full notify → watch → replay cycle. Copy the example config and start the server:

cp configuration/config.yaml.example configuration/config.yaml
cargo run

With auth (default). Start auth-o-tron before running the smoke tests:

python3 -m pip install httpx
./scripts/auth-o-tron-docker.sh
python3 scripts/smoke_test.py

Without auth. Set auth.enabled: false in your config (or remove the auth section), then:

AUTH_ENABLED=false python3 scripts/smoke_test.py

AUTH_ENABLED must match the server’s auth.enabled setting. When false, auth headers are omitted and auth-specific smoke tests are skipped.

Useful overrides:

BASE_URL="http://127.0.0.1:8000" python3 scripts/smoke_test.py
BACKEND="jetstream"               python3 scripts/smoke_test.py
TIMEOUT_SECONDS=12                python3 scripts/smoke_test.py
SMOKE_VERBOSE=1                   python3 scripts/smoke_test.py
python3 scripts/smoke_test.py --verbose

The smoke script covers:

  • health endpoint
  • replay/watch baseline flows (test_polygon)
  • mars replay with dot-containing identifier values, integer and enum predicates
  • dissemination watch + from_date with dot-containing identifier values
  • read/write auth separation across public, role-restricted, and admin-only streams
  • (optional, off by default) ECPDS plugin allow/deny/notify-bypass (see below)

Optional: end-to-end ECPDS plugin smoke test

If your build has --features ecpds enabled and your config has a working ecpds: block pointing at your real ECPDS servers, the smoke script can verify the plugin end-to-end against that ECPDS. It’s off by default.

Prerequisites:

  • The server must run with auth.enabled: true and auth.mode: direct. The smoke script sends HTTP Basic credentials, which Aviso forwards to auth-o-tron only in direct mode. Trusted-proxy mode would require an upstream proxy to mint a JWT, which is out of scope for the smoke script.

  • Your auth-o-tron config must know two users: an admin user (defaults admin-user / admin-pass) for the NOTIFY-bypass case, and your ECPDS user (ECPDS_ALLOWED_USER / ECPDS_ALLOWED_PASS) for the watch cases.

  • You need a destination value the ECPDS user is entitled to and one they are not (the latter can be a deliberately-fake string).

  • Add a minimal ECPDS test schema to your config, with match_key (e.g. destination) as the only required identifier field. The smoke test sends a minimal request body and does not populate any other required identifier fields. Don’t point it at a richer schema like your production dissemination. Add this dedicated test schema instead:

    notification_schema:
      ecpds_test:
        payload:
          required: false
        topic:
          base: "ecpds_test"
          key_order: ["destination"]
        identifier:
          destination:
            type: StringHandler
            required: true
        auth:
          required: true
          plugins: ["ecpds"]
    

Then:

ECPDS_ENABLED=true \
  ECPDS_EVENT_TYPE=ecpds_test \
  ECPDS_MATCH_KEY=destination \
  ECPDS_ALLOWED_USER="<auth-o-tron-username>" \
  ECPDS_ALLOWED_PASS="<auth-o-tron-password>" \
  ECPDS_ALLOWED_DESTINATION="<destination-the-user-can-read>" \
  ECPDS_DENIED_DESTINATION="NOT-A-REAL-DEST" \
  AUTH_ADMIN_USER=admin-user \
  AUTH_ADMIN_PASS=admin-pass \
  python3 scripts/smoke_test.py

What the three ECPDS smoke cases verify:

CaseWhat it asserts
ecpds: allowed user + entitled destination -> 200POST /api/v1/watch returns HTTP 200 for ECPDS_ALLOWED_USER reading ECPDS_ALLOWED_DESTINATION.
ecpds: allowed user + DENIED destination -> 403Same endpoint returns HTTP 403 for the same user reading ECPDS_DENIED_DESTINATION.
ecpds: notify on ECPDS-protected stream is not gatedPOST /api/v1/notification returns 2xx for AUTH_ADMIN_USER. The plugin is read-only; a 503 here would mean it incorrectly ran on a write.

Troubleshooting:

  • All three skip with [INFO] skipping ECPDS smoke test → check ECPDS_ENABLED=true and that the required env vars are set.
  • The allow case fails with 400 and a “schema validator before the plugin” hint: your ECPDS_EVENT_TYPE schema has additional required identifier fields. Add the minimal test schema above, or simplify the schema you’re pointing at. The schema validator rejecting the request before ECPDS runs is the correct behaviour. The smoke test fails loudly here rather than papering over it.
  • The allow case fails with 503 → the issue is between Aviso and ECPDS rather than at the plugin layer; see the ECPDS Plugin Runbook.
  • The notify-bypass case fails with 401/403 → your AUTH_ADMIN_USER / AUTH_ADMIN_PASS don’t match your auth-o-tron config; that’s an auth setup issue, not an ECPDS issue.

Reporting a Problem

Every aviso response carries an X-Request-ID HTTP header and (for error responses) a request_id field in the JSON body. Streaming responses also include the same UUID in the first SSE event. When something goes wrong, include this id in the bug report so the operator can find the matching server logs in seconds. See API Errors for details.

What’s Next

Defining Notification Schemas

A notification schema describes the shape of an event stream: what identifier fields are accepted, how they are validated, how the storage topic is constructed, whether a payload is required, and who can read or write.

Schemas are defined under the notification_schema key in your configuration file. Each top-level key becomes an event type that clients reference when calling /api/v1/notification, /api/v1/watch, or /api/v1/replay.

notification_schema:
  my_event: # ← event type name
    topic: ...
    identifier: ...
    payload: ...
    auth: ... # optional
    storage_policy: ... # optional, JetStream only

Topic Configuration

A schema should have a topic block that tells Aviso how to build the NATS subject for storage and routing.

topic:
  base: "weather"
  key_order: ["region", "date"]
FieldDescription
baseRoot prefix for the subject. Must be unique across all schemas (case-insensitive).
key_orderOrdered list of identifier field names appended to the base, separated by ..

Bases must match [A-Za-z0-9][A-Za-z0-9_-]*. Startup rejects invalid bases, including event names used as bases without a topic block. See the base contract for details.

Given base: "weather" and key_order: ["region", "date"], a request with region=north and date=20250706 produces the subject:

weather.north.20250706

Values containing reserved characters (., *, >, %) are automatically percent-encoded so they do not interfere with NATS subject routing. See Topic Encoding for details.

When a topic block is configured, startup requires a nonempty key_order. Each entry must name a declared identifier and may appear only once. Every ordinary identifier must be included, even when required: false. Optional fields still need a subject position for watch/replay wildcards and filters. For example, key_order: [region, date] is valid when both fields are declared; key_order: [region, region] is not.

Spatial geometry is the exception. PolygonHandler fields may be omitted because their geometry is stored as metadata. Existing polygon subject positions remain supported. The reserved point_cloud field must never appear in key_order; it uses spatial metadata instead.

There is no request-only authorization exception for ECPDS. Its match_key must appear in a configured topic’s key_order. Checking permission for a request value does not restrict delivered events unless routing also retains that value. Ordinary fields outside key_order are not validation-only fields: their values would be lost from routing and topic-based reconstruction.

These checks apply to configured topic blocks. They do not change the generic fallback used without a topic or schema, including its bare-topic behavior.


Identifier Fields

The identifier map defines the fields that clients can send. Each field specifies a handler type that controls validation and canonicalization.

identifier:
  region:
    type: EnumHandler
    values: ["north", "south", "east", "west"]
    required: true
    description: "Geographic region."
  date:
    type: DateHandler
    required: true

Every field supports these common properties:

PropertyType
typestring
requiredbool
descriptionstring
type Type: string

Handler type (see below). Required.

required Type: bool

Affects watch and replay only: if true, those requests must include this field; if false, missing keys become wildcards. Has no effect on notify, which always requires every declared field. Required.

description Type: string

Human-readable text exposed by GET /api/v1/schema. Optional.

PointCloudHandler is operation-specific. Publishers provide the declared point_cloud field. Watch and replay requests provide a closed polygon instead. That polygon satisfies required: true; subscribers must not send point_cloud.

Handler Types

StringHandler

Accepts any non-empty string. No transformation.

class:
  type: StringHandler
  max_length: 2 # optional: reject strings longer than this
  required: true

DateHandler

Parses dates in multiple formats and canonicalizes to a configured output format.

Accepted inputs: YYYY-MM-DD, YYYYMMDD, YYYY-DDD (day-of-year).

date:
  type: DateHandler
  canonical_format: "%Y%m%d" # output format (default: "%Y%m%d")
  required: false

Use canonical_format to choose the output format:

Output formatExample output
"%Y%m%d"20250706
"%Y-%m-%d"2025-07-06

Invalid dates (e.g. February 30) are rejected.

TimeHandler

Parses times and canonicalizes to four-digit HHMM format.

Accepted inputs: 14:30, 1430, 14, 9:05.

time:
  type: TimeHandler
  required: false

Input 14:30 → stored as 1430. Input 9 → stored as 0900.

EnumHandler

Accepts one value from a predefined list. Matching is case-insensitive; stored in lowercase.

domain:
  type: EnumHandler
  values: ["a", "b", "c"]
  required: false

Input "A" → stored as "a". Input "x" → rejected.

IntHandler

Accepts integer strings. Strips leading zeros for canonical storage.

step:
  type: IntHandler
  range: [0, 100000] # optional: inclusive [min, max] bounds
  required: false

Input "007" → stored as "7". Input "-1" with range: [0, 100] → rejected.

FloatHandler

Accepts floating-point strings. Rejects NaN and Inf.

severity:
  type: FloatHandler
  range: [0.0, 10.0] # optional: inclusive [min, max] bounds
  required: false

Input "3.14" → stored as "3.14". Input "NaN" → rejected.

ExpverHandler

Experiment version handler. Numeric values are zero-padded to four digits; non-numeric values are lowercased.

expver:
  type: ExpverHandler
  default: "0001" # optional: used when the field is empty
  required: false

Input "1" → stored as "0001". Input "test" → stored as "test".

PolygonHandler

Accepts a closed polygon as a JSON array of [latitude,longitude] pairs. The first and last pair must be identical.

Canonical form: [[lat,lon],...,[lat,lon]]. Emitted CloudEvents always use this array form.

polygon:
  type: PolygonHandler
  required: true

Constraints: at least four coordinate pairs including the closing repeat. Coordinates must be finite numbers. Latitude must be in [-90, 90]; longitude must be in [-180, 180]. Aviso checks pair count and closure, not vertex uniqueness, area, or self-intersections. Supply a non-degenerate polygon.

PointCloudHandler

Spatial identifiers use fixed names: PolygonHandler must be named polygon, and PointCloudHandler must be named point_cloud. A schema cannot mix the two handlers or declare multiple geometries. Polygon metadata and existing routed polygon subjects are supported. Neither spatial handler nor the reserved polygon routing position can be an ECPDS match_key; use an ordinary routing identifier such as destination instead.

Accepts a non-empty JSON array of [lat,lon] pairs from providers. Point clouds do not have a string syntax and do not need a closing point. Duplicates are valid. Their order is preserved.

point_cloud:
  type: PointCloudHandler
  required: true
  max_points: 10000
  description: >-
    Publishers provide point_cloud as [[latitude, longitude], ...]. Watch and
    replay requests provide a closed polygon instead. The polygon satisfies
    this required field; subscribers must not send point_cloud.

max_points defaults to 10,000 and cannot exceed 10,000. Canonical point-cloud JSON is also limited to 60 KiB. This is a conservative Aviso interoperability limit informed by NATS-backed header transport and near-limit round-trip tests. It leaves room for Aviso’s other metadata, but is not a protocol-wide NATS header limit. The point-count check runs first.

The handler must use the reserved point_cloud identifier key. Only one is allowed per schema, and the schema must define a topic block. Do not put point_cloud in topic.key_order; startup rejects it. Aviso stores the cloud in spatial metadata instead of the subject.

Point-cloud subscribers use polygon on /watch and /replay. A valid polygon satisfies a required point_cloud field for those operations. A schema with a PointCloudHandler cannot declare polygon, since that key is reserved for the query-time filter.

See Spatial Filtering for usage examples.

Reserved Query-Time Fields

point (built-in)

The point field is a reserved identifier that clients can send on /watch or /replay to filter notifications whose polygon contains the point. Canonical form is [latitude,longitude].

point is not a schema-configurable handler. It is available on schemas that include a PolygonHandler. The /notification endpoint rejects it.

See Spatial Filtering for usage examples.


Alternative Coordinate Format

The HTTP API also accepts comma-separated coordinate strings for polygons and points. Parentheses are optional. For example, the polygon string "(52.5,13.4,52.6,13.5,52.5,13.6,52.4,13.5,52.5,13.4)" and point string "52.55,13.50" represent the arrays in the spatial examples. The same coordinate order and polygon closure rules apply. Strings containing JSON coordinate arrays are also accepted for points and polygons. Polygons and clouds need nested pairs, not flat numeric arrays. GeoJSON objects are not accepted as spatial identifier values. Point clouds have no string format. CloudEvent spatial identifiers are arrays regardless of the input format.

Payload Configuration

Controls whether requests must include a payload field.

payload:
  required: true

The required setting controls whether the payload is mandatory:

Payload requiredBehavior
trueRequests without a payload are rejected (400).
falsePayload is optional; missing payloads are stored as JSON null.

The payload can be any valid JSON value (object, array, string, number, boolean, null). It is stored as-is with no reshaping.

See Payload Contract for full semantics.


Per-Stream Authentication

When global authentication is enabled, individual schemas can require credentials and restrict access by role.

auth:
  required: true
  read_roles:
    localrealm: ["analyst", "consumer"]
  write_roles:
    localrealm: ["producer"]
FieldDefault when omittedEffect
required(none)Must be set explicitly to true or false.
read_rolesAny authenticated user can readMaps realm → role list for watch/replay access.
write_rolesOnly admins can writeMaps realm → role list for notify access.

Use ["*"] as the role list to grant access to all users from a realm.

Admins (users matching global admin_roles) always have both read and write access.

Omitting the entire auth block makes the stream publicly accessible, even when global auth is enabled.

See Authentication for the full access-control matrix and role-matching rules.


Replay Limit

max_historical_notifications Default: inherited from watch_endpoint · Type: positive integer

Optional cap on historical notifications delivered by one replay or replaying watch request. Put it directly under the event schema, not in storage_policy:

notification_schema:
  weather:
    max_historical_notifications: 20000
    # Existing topic, identifier and other schema fields go here.

Omitting this field inherits watch_endpoint.max_historical_notifications (default 10000). An override can raise or lower that value. Zero and unlimited are rejected. Both backends support this setting; it does not change retention. Batch size stays global at watch_endpoint.replay_batch_size (default 100). This operational setting is not exposed by the schema API.

Only notifications that pass request filters and render successfully count. Exactly filling the cap completes normally unless another deliverable notification exists. Truncation closes the request without replay_completed or live delivery. See Historical Replay Limits.

Storage Policy (JetStream Only)

When using the JetStream backend, you can configure per-stream retention limits.

storage_policy:
  retention_time: "7d"
  max_messages: 500000
  max_size: "2Gi"
  allow_duplicates: false
  compression: true
FieldType
retention_timeduration
max_messagesinteger
max_sizesize
allow_duplicatesbool
compressionbool
retention_time Type: duration

Discard messages older than this. Accepts 30m, 1h, 7d, 1w.

max_messages Type: integer

Maximum message count; oldest are discarded when exceeded.

max_size Type: size

Maximum stream size. Accepts 100Mi, 1Gi, etc.

allow_duplicates Type: bool

Allow duplicate message IDs. Default: backend-specific.

compression Type: bool

Enable message-level compression. Default: backend-specific.

All fields are optional. Omitting storage_policy entirely uses backend defaults.

The in-memory backend does not support storage policies.


Complete Example

This example defines a weather alert stream with date/region routing, enum validation, optional payload, and role-restricted access.

notification_schema:
  weather_alert:
    payload:
      required: false

    topic:
      base: "alert"
      key_order: ["region", "severity_level", "date", "issued_by"]

    identifier:
      region:
        description: "Geographic region."
        type: EnumHandler
        values: ["europe", "asia", "africa", "americas", "oceania"]
        required: true
      severity_level:
        description: "Alert severity (1 to 5)."
        type: IntHandler
        range: [1, 5]
        required: true
      date:
        description: "Alert date."
        type: DateHandler
        canonical_format: "%Y%m%d"
        required: true
      issued_by:
        description: "Issuing authority identifier."
        type: StringHandler
        max_length: 64
        required: false

    auth:
      required: true
      read_roles:
        operations: ["*"]
      write_roles:
        operations: ["forecaster", "admin"]

    storage_policy:
      retention_time: "30d"
      max_messages: 100000

With this schema:

  • Publishing a notification with region=europe, severity_level=3, date=2025-07-06, issued_by=forecast produces the subject alert.europe.3.20250706.forecast.
  • Publishers must provide issued_by. Watch/replay clients may omit it to match any issuer because it is declared required: false.
  • Any authenticated user in the operations realm can watch/replay.
  • Only users with the forecaster or admin role can publish.
  • JetStream retains up to 100,000 messages or 30 days, whichever limit is hit first.

Tips

  • Start simple. Define only topic, one or two identifier fields, and payload. Add auth and storage policy later.
  • Use key_order deliberately. Fields in key_order become part of the NATS subject and affect routing granularity. More fields = more specific topics = more efficient filtering, but also more distinct subjects.
  • Choose subscriber requirements. Set required: true when watch/replay clients must supply a field. Publishers always supply every declared field.
  • Keep base short and unique. It is the root of every subject in this stream. Avoid collisions with other schemas.
  • Test with GET /api/v1/schema/{event_type}. This endpoint returns the public view of your schema, showing all identifier fields and their validation rules.

Practical Examples

This section provides copy-paste examples for common workflows.

All examples use the same generic event schema so behavior is easy to compare.

Shared Generic Schema

notification_schema:
  extreme_event:
    topic:
      base: extreme_event
      key_order: [region, run_time, severity, anomaly]
    identifier:
      region:
        description: "Geographic region label."
        type: EnumHandler
        values: ["north", "south", "east", "west"]
        required: true
      run_time:
        type: TimeHandler
        required: true
      severity:
        description: "Severity level from 1 to 7."
        type: IntHandler
        range: [1, 7]
        required: true
      anomaly:
        type: FloatHandler
        range: [0.0, 100.0]
        required: false
      polygon:
        type: PolygonHandler
        required: false
    payload:
      required: false

Shared Assumptions

  • Base URL: http://127.0.0.1:8000
  • Content type: application/json
  • Replay examples use from_id or from_date explicitly.

Identifier Value Style

The examples in this section send scalar identifier values as JSON strings, including numeric ones ("severity":"4", "anomaly":"42.5"). The server canonicalizes scalar values to strings. JSON numbers are accepted and produce the same result. Constraint arguments ({"gte":5}, {"between":[3,7]}) use JSON numbers because they are typed comparison operands.

Notify Identifier Rule

POST /api/v1/notification requires every identifier key declared in the schema, whether it has required: true or required: false. Every value must pass its handler validation.

The flag matters on watch and replay. A missing key marked required: true returns 400. A missing key marked required: false becomes a wildcard. Provided values and constraint objects still pass through handler validation.

The shared schema above declares five keys, so notify examples on the following pages include all five.

Next:

Basic Notify/Watch/Replay

Uses the shared generic schema from Practical Examples.

This page is the quickest way to understand the normal API flow. You first publish (notify), then observe live updates (watch), then read history (replay). If you are onboarding a new environment, start here before trying filters or admin operations. Read the examples in order.

1) Notify

Notify requires every identifier key declared in the schema. The required flag has no effect on notify (every key must be present and every value must pass handler validation); it only affects watch and replay, where keys marked required: false may be omitted and are treated as wildcards. The shared schema declares five keys (region, run_time, severity, anomaly, polygon), so all five appear below.

curl -sS -X POST "http://127.0.0.1:8000/api/v1/notification" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"4",
      "anomaly":"42.5",
      "polygon":[[52.5,13.4],[52.6,13.5],[52.5,13.6],[52.4,13.5],[52.5,13.4]]
    },
    "payload":{"note":"initial forecast"}
  }'

Expected: HTTP 200. Omitting any of the five identifier keys returns 400 with code: INVALID_NOTIFICATION_REQUEST.

2) Watch (Live Only)

curl -N -X POST "http://127.0.0.1:8000/api/v1/watch" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"4",
      "anomaly":"42.5"
    }
  }'

Expected:

  • HTTP 200
  • SSE starts with connection_established
  • only new matching notifications arrive

3) Replay (Historical)

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"4",
      "anomaly":"42.5"
    },
    "from_id":"1"
  }'

Expected:

  • HTTP 200
  • SSE emits replay_started, replay events, replay_completed, then closes

Constraint Filtering

Uses the shared generic schema from Practical Examples.

Constraint filtering lets subscribers express conditions over identifier fields instead of exact values: ranges (severity >= 5), numeric bands, or enum subsets. This page covers seed data, valid constraint requests, and common failure cases.

Seed Notifications

curl -sS -X POST "http://127.0.0.1:8000/api/v1/notification" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north","run_time":"1200","severity":"3","anomaly":"42.5",
      "polygon":[[52.5,13.4],[52.6,13.5],[52.5,13.6],[52.4,13.5],[52.5,13.4]]
    },
    "payload":{"note":"seed-a"}
  }'

curl -sS -X POST "http://127.0.0.1:8000/api/v1/notification" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"south","run_time":"1200","severity":"6","anomaly":"87.2",
      "polygon":[[10.0,10.0],[10.2,10.0],[10.2,10.2],[10.0,10.2],[10.0,10.0]]
    },
    "payload":{"note":"seed-b"}
  }'

Expected:

  • both return HTTP 200

Scalar Value (Implicit eq)

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"south","run_time":"1200","severity":"6","anomaly":"87.2"},
    "from_id":"1"
  }'

Expected:

  • HTTP 200
  • only severity = 6 notifications match

Integer Constraint (gte)

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":{"in":["north","south"]},
      "run_time":"1200",
      "severity":{"gte":5},
      "anomaly":"87.2"
    },
    "from_id":"1"
  }'

Expected:

  • HTTP 200
  • includes severity=6
  • excludes severity=3

Float Constraint (between)

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"3",
      "anomaly":{"between":[40.0,50.0]}
    },
    "from_id":"1"
  }'

Expected:

  • HTTP 200
  • includes anomaly=42.5

Float eq Is Exact (No Tolerance)

Float eq and in are exact comparisons. This keeps behavior deterministic across replay/live and avoids hidden tolerance windows.

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"3",
      "anomaly":{"eq":42.5}
    },
    "from_id":"1"
  }'

Expected:

  • HTTP 200
  • only notifications with exactly anomaly=42.5 match
  • NaN/inf values are rejected by float validation/constraints

Enum Constraint (in)

curl -N -X POST "http://127.0.0.1:8000/api/v1/watch" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":{"in":["south","west"]},
      "run_time":"1200",
      "severity":"6",
      "anomaly":"87.2"
    }
  }'

Expected:

  • HTTP 200
  • live notifications pass only for regions in ["south","west"]

Invalid: Two Operators in One Constraint Object

curl -sS -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":{"gte":4,"lt":7},
      "anomaly":"42.5"
    },
    "from_id":"1"
  }'

Expected:

  • HTTP 400
  • message says constraint object must contain exactly one operator

Invalid: Constraint Object on /notification

All five identifier keys are present so the request fails specifically on the constraint object, not on a missing key.

curl -sS -X POST "http://127.0.0.1:8000/api/v1/notification" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":{"gte":4},
      "anomaly":"42.5",
      "polygon":[[52.5,13.4],[52.6,13.5],[52.5,13.6],[52.4,13.5],[52.5,13.4]]
    },
    "payload":{"note":"should-fail"}
  }'

Expected:

  • HTTP 400

Spatial Filtering

Uses the shared generic schema from Practical Examples.

Spatial filtering has two modes:

  • identifier.polygon: keep notifications whose polygon intersects the request polygon.
  • identifier.point: keep notifications whose polygon contains the request point.

This matters when many notifications share similar non-spatial identifiers and you need geographic precision.

Seed Notifications

These two notifications have different polygon shapes and anomaly values.

curl -sS -X POST "http://127.0.0.1:8000/api/v1/notification" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"4",
      "anomaly":"42.5",
      "polygon":[[52.5,13.4],[52.6,13.5],[52.5,13.6],[52.4,13.5],[52.5,13.4]]
    },
    "payload":{"note":"poly-a"}
  }'

curl -sS -X POST "http://127.0.0.1:8000/api/v1/notification" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"4",
      "anomaly":"87.2",
      "polygon":[[10.0,10.0],[10.2,10.0],[10.2,10.2],[10.0,10.2],[10.0,10.0]]
    },
    "payload":{"note":"poly-b"}
  }'

Expected:

  • both requests return HTTP 200

Replay with Polygon Intersection Filter

This request polygon intersects poly-a but not poly-b.

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"4",
      "polygon":[
        [52.52,13.45],[52.62,13.55],[52.52,13.65],
        [52.42,13.55],[52.52,13.45]
      ]
    },
    "from_id":"1"
  }'

Expected:

  • HTTP 200
  • replay includes poly-a
  • replay excludes poly-b

Replay with Point Containment Filter

The point below is inside poly-a and outside poly-b.

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"4",
      "point":[52.55,13.50]
    },
    "from_id":"1"
  }'

Expected:

  • HTTP 200
  • replay includes poly-a
  • replay excludes poly-b

Replay Without Spatial Filter

No polygon and no point means no spatial narrowing.

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"4"
    },
    "from_id":"1"
  }'

Expected:

  • HTTP 200
  • replay may include both poly-a and poly-b because only non-spatial fields are applied

Invalid: polygon and point Together

curl -sS -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{
      "region":"north",
      "run_time":"1200",
      "severity":"4",
      "polygon":[[52.5,13.4],[52.6,13.5],[52.5,13.6],[52.4,13.5],[52.5,13.4]],
      "point":[52.55,13.50]
    },
    "from_id":"1"
  }'

Expected:

  • HTTP 400
  • validation error says both spatial filters cannot be used together

Point-Cloud Filtering

Point-cloud streams let a provider publish many locations in one notification. A subscriber supplies a polygon and receives the notification when any cloud point lies inside the polygon or on its boundary.

Coordinates always use latitude first, then longitude. The canonical JSON forms are:

  • point: [lat,lon]
  • polygon: [[lat,lon],...], at least four pairs, first pair repeated last
  • point cloud: [[lat,lon],...], with at least one point

Point clouds accept JSON arrays only. Emitted CloudEvents always use arrays for all spatial identifier values. See the alternative format for point and polygon strings.

Schema

notification_schema:
  observations:
    topic:
      base: observations
      key_order: [date]
    identifier:
      date:
        type: DateHandler
        canonical_format: "%Y%m%d"
        required: true
      point_cloud:
        type: PointCloudHandler
        required: true
        max_points: 10000
        description: >-
          Publishers provide point_cloud as [[latitude, longitude], ...].
          Watch and replay requests provide a closed polygon instead. The
          polygon satisfies this required field; subscribers must not send
          point_cloud.
    payload:
      required: true

The handler must use the reserved point_cloud key. A schema can contain only one PointCloudHandler. The configured max_points defaults to 10,000 and cannot exceed 10,000.

Do not put point_cloud in topic.key_order. Startup rejects that configuration. Cloud coordinates stay in the spatial_point_cloud backend metadata header. The subject contains only normal routing fields. Aviso also stores a spatial_bbox header for coarse rejection before exact matching.

Notify

The provider sends the cloud on /notification:

curl -sS -X POST "http://127.0.0.1:8000/api/v1/notification" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"observations",
    "identifier":{
      "date":"20260826",
      "point_cloud":[
        [52.52,13.40],
        [48.14,11.58],
        [53.55,9.99]
      ]
    },
    "payload":{"source":"stations"}
  }'

Clouds do not need a closing repeat. Duplicate points are valid, and Aviso preserves their order. Every latitude and longitude must be finite. Latitude must be in [-90, 90]; longitude must be in [-180, 180].

The canonical point-cloud JSON is limited to 60 KiB. This conservative Aviso interoperability limit is informed by NATS-backed header transport and near-limit round-trip tests. It leaves room for Aviso’s other metadata, but is not a protocol-wide NATS header limit. The point-count limit is checked first.

Watch Or Replay

Subscribers use the reserved polygon query field. They do not send point_cloud:

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"observations",
    "identifier":{
      "date":"20260826",
      "polygon":[
        [52.40,13.20],
        [52.40,13.70],
        [52.70,13.70],
        [52.70,13.20],
        [52.40,13.20]
      ]
    },
    "from_id":"0"
  }'

The polygon satisfies a required point_cloud field for watch and replay. Aviso first rejects non-overlapping bounding boxes. It then tests points in their stored order and stops at the first match. Points on an edge or vertex match. A cloud whose points are all outside does not match.

Requests are rejected when they send point_cloud to watch or replay, declare polygon in a point-cloud schema, or combine incompatible spatial filters.

CloudEvent Output

The event reconstructs the provider identifier as JSON:

{
  "data": {
    "identifier": {
      "date": "20260826",
      "point_cloud": [
        [52.52, 13.4],
        [48.14, 11.58],
        [53.55, 9.99]
      ]
    }
  }
}

The request polygon is a filter. It is not substituted into the emitted identifier.

Replay Start and End Points

Uses the shared generic schema from Practical Examples.

Replay start parameters control where historical delivery begins, and the optional end parameters control where it stops. Choose from_id and to_id when you track sequence progress; choose from_date and to_date when you track wall-clock time. These examples cover valid forms and the common invalid combinations that return 400. Use this page to validate client retry and resume logic.

Replay from Sequence (from_id)

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"},
    "from_id":"10"
  }'

Expected:

  • HTTP 200
  • replay starts from sequence 10 (inclusive)

Replay from Time (from_date) RFC3339

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"},
    "from_date":"2026-03-01T12:00:00Z"
  }'

Expected:

  • HTTP 200
  • replay starts from that UTC timestamp (inclusive)

Replay from Time (from_date) Unix Seconds

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"},
    "from_date":"1740509903"
  }'

Expected:

  • HTTP 200

Replay from Time (from_date) Unix Milliseconds

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"},
    "from_date":"1740509903710"
  }'

Expected:

  • HTTP 200

Replay a Sequence Range (to_id)

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"},
    "from_id":"10",
    "to_id":"20"
  }'

Expected:

  • HTTP 200
  • replay delivers the matching notifications with sequences 10 to 20 (both inclusive), then closes

Replay a Time Window (to_date)

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"},
    "from_date":"2026-03-01T00:00:00Z",
    "to_date":"2026-03-02T00:00:00Z"
  }'

Expected:

  • HTTP 200
  • replay delivers the matching notifications stored from the start of 1 March to the start of 2 March 2026 UTC (both inclusive), then closes

to_date accepts the same formats as from_date, and either end parameter can follow either start parameter.

Invalid Replay Start Combinations

Missing Both from_id and from_date

curl -sS -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"}
  }'

Expected:

  • HTTP 400

Both from_id and from_date Provided

curl -sS -X POST "http://127.0.0.1:8000/api/v1/watch" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"},
    "from_id":"5",
    "from_date":"2026-03-01T12:00:00Z"
  }'

Expected:

  • HTTP 400

Invalid Replay End Combinations

Both to_id and to_date Provided

curl -sS -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"},
    "from_id":"5",
    "to_id":"20",
    "to_date":"2026-03-01T12:00:00Z"
  }'

Expected:

  • HTTP 400

to_id Lower Than from_id

curl -sS -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"},
    "from_id":"20",
    "to_id":"10"
  }'

Expected:

  • HTTP 400

End Parameter on a Watch

curl -sS -X POST "http://127.0.0.1:8000/api/v1/watch" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type":"extreme_event",
    "identifier":{"region":"north","run_time":"1200","severity":"4","anomaly":"42.5"},
    "from_id":"5",
    "to_id":"20"
  }'

Expected:

  • HTTP 400

Admin Operations (Practical)

See full reference in Admin Operations.

These operations are for cleanup and recovery, not normal data flow. Use delete when one bad record must be removed; use wipe when resetting a stream or environment. Because these endpoints are destructive, validate IDs and stream names carefully before execution.

Delete One Notification by ID

curl -X DELETE "http://127.0.0.1:8000/api/v1/admin/notification/extreme_event@42"

Expected:

  • 200 if it exists
  • 404 if stream/sequence does not exist
  • 400 for invalid format

Wipe One Stream

curl -X DELETE "http://127.0.0.1:8000/api/v1/admin/wipe/stream" \
  -H "Content-Type: application/json" \
  -d '{"stream_name":"extreme_event"}'

The name is case-insensitive and accepts the event type or the backend stream name (extreme_event and EXTREME_EVENT target the same stream).

Expected:

  • 200 if the stream exists
  • 404 if no stream matches; the message lists the configured event types
  • stream definition remains, messages are removed

Wipe All Streams

curl -X DELETE "http://127.0.0.1:8000/api/v1/admin/wipe/all"

Expected:

  • 200
  • all stream data is removed

Backends Overview

Aviso abstracts all storage and messaging behind a NotificationBackend trait. Two implementations ship out of the box.


Which Backend Should I Use?

flowchart TD
    A["What do you need?"] --> B{"Persistent history<br/>across restarts?"}
    B -->|Yes| C[jetstream]
    B -->|No| D{"Multiple replicas<br/>or pods?"}
    D -->|Yes| C
    D -->|No| E{"Production<br/>workload?"}
    E -->|Yes| C
    E -->|No| F[in_memory]

    C:::jet
    F:::mem

    classDef jet fill:#1a6b3a,color:#fff,stroke:#0d4a27
    classDef mem fill:#1a4d6b,color:#fff,stroke:#0d3347
RequirementRecommended backend
Persistent history across restartsjetstream
Replay endpoint supportjetstream (or in_memory for local/node-local use)
Live watch streaming supportjetstream (or in_memory for local/node-local use)
Multi-replica deploymentjetstream
Quick local experimentation with minimal setupin_memory

Capability Comparison

CapabilityJetStreamIn-Memory
Durable storageYesNo (data lost on restart)
Replay supportYesYes (node-local only)
Live watch supportYesYes (node-local fan-out)
Multi-replica / HAYes (clustered NATS)No
Per-schema storage policyYesNo (rejected at startup)
Cross-instance consistencyYesNo

Backend Details

InMemory Backend

Intended use

in_memory backend is best for:

  • local development,
  • schema/validation testing,
  • lightweight experimentation where persistence is not required.

Behavior

  • Data is process-memory only and is lost on restart.
  • Topic/message limits are enforced with eviction.
  • No shared state across replicas or pods.
  • Supports live watch subscriptions (live-only delivery).
  • Supports replay batch retrieval for from_id and from_date, with optional to_id and to_date end points.
  • Uses in-process fanout only, so subscriptions/replay are node-local.

Historical delivery uses the same request-wide watch_endpoint.max_historical_notifications cap as JetStream, after request filtering and successful rendering. A schema’s max_historical_notifications can override the global cap. This is separate from retention and batch size. See Historical Replay Limits for truncation controls and watch behavior.

Watch creates its broadcast receiver and captures the last allocated sequence under the same lock used to store and publish notifications. Replay reads only up to that inclusive bound; live delivery starts above it. Replay-only captures the bound under that lock without creating a receiver. Eviction or deletion can still remove history during replay, and a lagging live receiver can lose queued notifications. The bound fixes the sequence range, not the stored contents.

Configuration

notification_backend.kind: in_memory

Available knobs:

  • max_history_per_topic (default 1)
  • max_topics (default 10000)
  • enable_metrics (default false)

Per-schema storage_policy fields are currently not supported on in_memory and are rejected at startup.

Production suitability

Not recommended for production because:

  • no durability,
  • no HA replication,
  • no cross-instance consistency,
  • replay/watch history is limited to local in-memory retention.

JetStream Backend

The jetstream backend is the production-oriented storage implementation. It connects to a NATS server with JetStream enabled and uses it for durable message storage, replay, and live streaming.


Intended Use

Use jetstream when you need:

  • durable storage that survives server restarts
  • replay across multiple server instances
  • live streaming with cluster-wide fan-out
  • configurable retention, size limits, and compression

Local Test Setup

Start a NATS + JetStream instance via Docker:

./scripts/init_nats.sh

Then configure Aviso:

notification_backend:
  kind: jetstream
  jetstream:
    nats_url: "nats://localhost:4222"

For full setup options including authentication and storage limits, see Installation: Local JetStream.


Core Behavior

  • Connects to the configured NATS server on startup (with retry).
  • Creates JetStream streams on demand, one per topic base (e.g. MARS, DISS, POLYGON).
  • Publishes notifications directly to JetStream subjects using the encoded wire format.
  • Uses pull consumers for replay batching (from_id, from_date).
  • Uses push consumers for live watch subscriptions.
  • Reconciles existing streams against current config when they are first accessed.

Pull batches use watch_endpoint.replay_batch_size independently of the request-wide max_historical_notifications delivery cap. The shared SSE layer applies that cap after request filtering and successful rendering, not inside each backend batch. A schema’s max_historical_notifications can override the global cap independently of its storage policy. Watch captures its history bound from the initial DeliverNew consumer create response’s delivered.stream_sequence, before any pulls or consumer-info refresh. The consumer and bound therefore share one creation point, including when the tail was deleted or the stream is empty at a nonzero sequence. A retry uses the final successful consumer’s bound. Replay-only reads the stream’s last sequence during setup without creating a live consumer. Neither operation freezes retention or deletion.

See Historical Replay Limits for truncation controls and watch behavior.


Configuration Reference

All fields live under notification_backend.jetstream.

Connection & startup

SettingDefault
nats_urlnats://localhost:4222
tokenNone
timeout_seconds30
retry_attempts3
nats_url Default: nats://localhost:4222

NATS server URL.

token Default: None

Token auth; falls back to NATS_TOKEN environment variable.

timeout_seconds Default: 30

Per-attempt connection timeout (> 0).

retry_attempts Default: 3

Startup connection attempts before backend init fails (> 0).

Runtime reconnect

SettingDefault
max_reconnect_attemptsunlimited
reconnect_delay_ms2000
max_reconnect_attempts Default: unlimited

Unset and 0 both mean unlimited reconnect retries; set a positive value only if you explicitly want the client to give up (the backend then stays disconnected until a process restart).

Subscription creation uses a bounded retry loop: unset means five attempts, 0 means one attempt, and a positive value sets the attempt limit.

reconnect_delay_ms Default: 2000

Delay between reconnect attempts and startup connect retries (> 0).

Publish resilience

publish_retry_attempts Default: 5

Retries for transient channel closed publish failures (> 0).

publish_retry_base_delay_ms Default: 150

Base backoff in ms for publish retries; grows exponentially per attempt (> 0).

Stream defaults

These apply to every stream created by Aviso unless overridden by a per-schema storage_policy.

max_messages Default: None

Stream message cap (maps to max_messages).

max_bytes Default: None

Stream size cap in bytes (maps to max_bytes).

retention_time Default: None

Default max age: duration literal (s, m, h, d, w; e.g. 30d).

storage_type Default: file

file or memory, parsed as typed enum at config load.

Omitting this setting requests file. Storage cannot change on an existing stream. If its storage differs, operations that ensure the stream fail with the stream name and current and requested types, before any mutable settings change. Aviso does not delete or recreate the stream; stored messages are left intact. Use the stream’s current type or arrange a separate migration.

replicas Default: None

Stream replica count.

retention_policy Default: limits

limits or interest. workqueue is rejected at startup.

discard_policy Default: old

old or new, parsed as typed enum.

Fail-fast validation: storage_type, retention_policy, and discard_policy are parsed as typed enums during configuration loading. Invalid values fail startup immediately, before any streams are created.

workqueue retention is not supported because it cannot support Aviso’s independent watch/replay consumers. Startup rejects it for backend defaults, including streams with schema storage policies. Schema storage policies inherit the backend retention policy; they cannot override it. Use limits for history bounded by configured limits. interest remains accepted with its existing NATS interest-based retention semantics; it does not guarantee retained history when there are no interested consumers. This validation does not migrate or delete existing streams.

Full example

notification_backend:
  kind: jetstream
  jetstream:
    nats_url: "nats://localhost:4222"
    timeout_seconds: 30
    retry_attempts: 3
    reconnect_delay_ms: 2000
    publish_retry_attempts: 5
    publish_retry_base_delay_ms: 150
    storage_type: file
    retention_policy: limits
    discard_policy: old

Stream Management

Stream creation

On first access (e.g. first publish for a given event type), Aviso creates a JetStream stream with the following settings applied:

  • storage_type, retention_policy, discard_policy
  • max_messages, max_bytes, retention_time → max_age
  • replicas

The stream subject binding is set to <base>.> (e.g. mars.>) to capture all topics under that base.

Reconciliation of existing streams

When a stream already exists and is accessed by Aviso, it is reconciled: the current stream config is compared against the desired config and mutable fields are updated if drift is detected:

  • limits (retention, size, message count)
  • compression
  • duplicate window
  • replicas
  • subject binding

Aviso leaves settings alone when they already have the intended effect, even if NATS reports a default differently. Equivalent defaults do not trigger an update.

If JetStream rejects an update (e.g. the field is not editable in the current server/stream state), the operation fails. Aviso does not report success while using stale retention settings. If another replica creates the stream during creation, Aviso reloads and reconciles that stream before proceeding.

Precedence

Backend-level defaults are applied first, then per-schema storage_policy overrides for that stream:

Values under notification_schema.<event_type>.storage_policy override the matching notification_backend.jetstream defaults.

Policy lookup matches the topic base without regard to ASCII case, just as startup validation does. NATS subjects themselves remain case-sensitive. Retention must be positive and fit signed 64-bit nanoseconds. The largest whole-second literal is 9223372036s; larger values fail validation.

Shortening retention deletes messages older than the new window, including messages stored before the change. They are no longer available for replay. This does not require deleting or recreating the stream. Keep all Aviso replicas on the same configuration so they do not repeatedly change each other’s policy. When shortening retention below the existing duplicate-detection window, Aviso also shortens that window to satisfy NATS’s limit. A smaller existing window is preserved.

Applying config changes to existing streams

Aviso uses the configuration loaded at startup. After editing config.yaml, restart Aviso or roll out the updated configuration to all replicas. Existing streams are then reconciled when accessed, for example by a publish or a new watch/replay request. Aviso does not sweep all streams at startup or in the background.

Compression applies to future file-storage writes at the block level. Changing the setting does not automatically recompress existing history. Aviso provides no automatic history migration.

Deleting a stream loses its stored messages. Recreating it starts an empty stream; it does not rewrite or restore history. The wipe_stream admin endpoint also removes messages, but preserves the stream configuration. Neither is a compression migration.


Verifying Effective Stream Policy

Use the nats CLI to inspect the stream config after a publish or reconcile:

# Replace POLYGON with your stream name (MARS, DISS, etc.)
nats --server nats://localhost:4222 stream info POLYGON

Fields to check:

CLI fieldConfig field
Max Ageretention_time
Max Messagesmax_messages
Max Bytesmax_bytes / per-schema max_size
Max Messages Per Subjectallow_duplicates: 1 = disabled, -1 = enabled
CompressionNone or S2

Replay Behavior

  • Sequence replay (from_id): starts from that sequence number, inclusive.
  • Time replay (from_date): uses JetStream start-time delivery policy.
  • The API enforces mutual exclusivity: from_id and from_date cannot both be present.
  • End point (to_id or to_date): lowers the replay’s end sequence. For to_date, a pull consumer that starts one nanosecond after it reads one message; the replay ends at the sequence before that message.

Smoke Test (JetStream Mode)

python3 -m pip install httpx

BACKEND=jetstream \
NATS_URL=nats://localhost:4222 \
JETSTREAM_POLICY_STREAM_NAME=POLYGON \
EXPECT_MAX_MESSAGES=500000 \
EXPECT_MAX_BYTES=2147483648 \
EXPECT_MAX_MESSAGES_PER_SUBJECT=1 \
EXPECT_COMPRESSION=None \
python3 scripts/smoke_test.py

Operational Caveats

Run the opt-in reconnect outage test from the repository root:

bash scripts/test_jetstream_reconnect.sh

It requires Docker and uses its own NATS 2.14.6 container with a random loopback port. It stops and restarts only that container, checks that bounded reconnect gives up while unset and 0 recover, then removes the container. It does not use NATS_URL or shared NATS storage.

  • Startup connectivity is controlled by timeout_seconds + retry_attempts.
  • Runtime reconnect is controlled by max_reconnect_attempts and reconnect_delay_ms.
  • Publish retry is a narrow resilience path for transient channel closed failures; non-transient failures fail fast.
  • retry_attempts applies only to startup; post-startup reconnect uses the reconnect settings.
  • Reconnect retries are unlimited unless max_reconnect_attempts is set to a positive value. A bounded value means the client gives up permanently once exhausted and the backend stays disconnected until a process restart, while the HTTP surface keeps serving.
  • GET /ready reflects the connection state: 200 while the NATS connection is live, 503 while it is down or reconnecting. Point Kubernetes readiness probes at /ready so traffic routes away during a backend outage and resumes on reconnect; keep liveness on /health (process-only) so pods are not killed during an outage the client recovers from by itself.
  • max_reconnect_attempts also bounds subscription-creation retries, where unset means 5 attempts and 0 means one attempt (a subscribe call has a caller waiting on it, so it never retries forever).

Deployment Modes

Local experimentation

Recommended backend:

  • in_memory for quick local request/validation testing.

Characteristics:

  • No persistence: data is lost on process restart.
  • Single-process state only.
  • Not suitable for horizontal scaling or replica failover.
  • Supports replay and watch in-process, limited by local memory retention.
  • For local JetStream testing, see Installation: Local JetStream.
  • For a quick end-to-end behavior check, see Getting Started: Run the Smoke Test.

Production-like / persistent mode

Recommended backend:

  • jetstream

Characteristics:

  • Durable message storage.
  • Retention and size limits.
  • Replica support (requires clustered NATS setup).
  • Supports replay and live streaming workflows.

For Kubernetes, use the Aviso Helm chart.

Selection guideline

  • Need persistence/replay/streaming robustness: use jetstream.
  • Need fastest setup for local functional checks only: use in_memory (node-local replay/watch).

Deploying with Helm

The server provides the application and its configuration format. The chart defines Kubernetes resources and supplies default values. Keep your own deployment values separately; no additional configuration repository is required.

Prepare the chart

You need a Kubernetes cluster, kubectl, and Helm with OCI support. The chart and its dependencies, as well as the server image, must be accessible from your environment. The default image and one chart dependency are hosted at eccr.ecmwf.int; check registry access before installing. Helm registry authentication does not supply image-pull credentials to Kubernetes nodes.

To install from a chart checkout:

git clone https://github.com/ecmwf/aviso-chart.git
helm repo add nats https://nats-io.github.io/k8s/helm/charts/
helm dependency build ./aviso-chart

Select the chart revision you intend to deploy before building dependencies. helm dependency build uses the chart’s lockfile. Dependencies must be downloadable even when their components are disabled in your values.

Start with a development deployment

This example runs one Aviso replica and one NATS server with JetStream. Your cluster needs a default StorageClass that can provision a persistent volume. Create values-demo.yaml:

fullnameOverride: aviso
replicaCount: 1

nats:
  enabled: true
  config:
    jetstream:
      enabled: true
      fileStore:
        enabled: true
        pvc:
          enabled: true
          size: 10Gi

extraEnv:
  - name: AVISOSERVER_CONFIG_FILE
    value: /etc/aviso_server/config.yaml

config:
  notification_backend:
    kind: jetstream
    jetstream:
      nats_url: nats://aviso-nats:4222
      storage_type: file
      replicas: 1

The NATS URL assumes the Helm release name aviso, used in the commands below. Change the URL if you choose another release name.

NATS stores notifications on the persistent volume, so they can survive pod restarts while that volume is retained. This uses chart defaults for the remaining settings, with no authentication or ingress. It is a development setup, not a highly available production deployment.

Validate and inspect the rendered resources before installing:

helm lint ./aviso-chart -f values-demo.yaml
helm template aviso ./aviso-chart --namespace aviso -f values-demo.yaml

helm upgrade --install aviso ./aviso-chart \
  --namespace aviso \
  --create-namespace \
  -f values-demo.yaml \
  --wait --timeout 5m

The fullnameOverride setting names the Aviso Service and Deployment aviso. For a temporary local check, forward the Service port to your machine:

kubectl -n aviso port-forward service/aviso 8000:8000

Keep this command running while testing. It exposes Aviso on your machine’s loopback interface and stops forwarding when you press Ctrl+C. It does not provide a shared endpoint for other users.

In another terminal, check liveness and backend-connection readiness:

curl --fail http://127.0.0.1:8000/health
curl --fail http://127.0.0.1:8000/ready

Applications inside the cluster can use the Kubernetes Service directly. For access from outside the cluster, configure the chart’s ingress settings with your hostname and TLS configuration, or use your platform’s supported exposure method. An Ingress requires an ingress controller; creating the resource alone does not make the service reachable. Configure authentication before exposing Aviso, and ensure proxy buffering and timeouts support long-lived SSE connections.

Understand the configuration layers

Helm merges chart defaults with your -f files in command-line order. Later files take precedence. Maps merge, but lists such as extraEnv are replaced as a whole. You can use one values file or layer several; their names and organization are yours to choose.

Top-level chart values configure Kubernetes resources. The config: map contains server settings and becomes a ConfigMap mounted at /etc/aviso_server/config.yaml.

The server then applies its own configuration loading rules. Without an explicit file selector, configuration included in the image can be merged with the mounted file. Omitting a setting from an overlay does not remove it from an earlier source.

The AVISOSERVER_CONFIG_FILE setting in the example selects the mounted file as the only file source. Field-level AVISOSERVER_* environment variables still override its settings. This does not remove Helm’s chart defaults.

Use Helm values for settings that also affect rendered resources, such as the application port. Changing the server port only through an environment variable does not update the chart’s container ports or probes.

Supply credentials through Secrets

Do not put credentials in config:: it is stored in a ConfigMap, not a Secret. Provision Secrets separately and reference their keys through extraEnv. For example, a deployment using token-authenticated NATS could use this values fragment:

extraEnv:
  - name: AVISOSERVER_CONFIG_FILE
    value: /etc/aviso_server/config.yaml
  - name: AVISOSERVER_NOTIFICATION_BACKEND__JETSTREAM__TOKEN
    valueFrom:
      secretKeyRef:
        name: aviso-nats-credentials
        key: token

Create that Secret in the release namespace before installation. This fragment supplies credentials only; it does not select or deploy JetStream. Keep any other required extraEnv entries in the same list.

Prepare for production

The example configures both sides: nats.enabled deploys NATS, while config.notification_backend tells Aviso how to connect to it. Neither setting automatically configures the other. To use an existing NATS deployment, disable the dependency and supply its URL and credentials.

Choose retention and capacity limits for your notification schemas through storage_policy. Schema settings override backend defaults per field. The demo intentionally leaves retention unspecified rather than choosing one age limit for every stream. Other limits still apply, including keeping only the latest notification per subject by default. See the JetStream backend guide for storage policy details.

Aviso pod replicas and JetStream stream replicas are separate settings. A single NATS server with persistent storage is not a highly available cluster. Before exposing the service, also configure authentication and ingress/TLS for your environment.

Apply configuration changes

Use helm upgrade with your chosen values files. The chart includes a checksum of the rendered ConfigMap in the pod template, so configuration changes through Helm trigger a rollout. The server reads configuration at startup; this is not live reload.

Editing the ConfigMap directly does not update that checksum. The mounted file also uses subPath, so existing containers do not receive projected ConfigMap updates. Keep configuration changes in your Helm values.

Updating a referenced Secret’s contents does not automatically restart pods. After rotating credentials, restart the Deployment so its environment variables are loaded again:

kubectl -n aviso rollout restart deployment/aviso
kubectl -n aviso rollout status deployment/aviso

Configuration

Key Behaviors to Know First

These rules affect how everything behaves at runtime:

  • Environment variables always win. They override any YAML value regardless of which file it came from.
  • Replay and watch behavior is controlled by request parameters, not static config switches.
  • Invalid policy values fail startup immediately. storage_type, retention_policy, and discard_policy are parsed as typed enums; bad values are caught before any streams are created.
  • Per-schema storage_policy is validated at startup against the selected backend’s capabilities. Unsupported fields (for example retention_time on in_memory) cause a startup failure with a clear error.
  • JetStream stream changes require a restart or rollout. Aviso uses the configuration loaded at startup. After loading the updated config, it reconciles existing streams when accessed, not through an all-stream sweep. Compression affects future file-storage writes, not an automatic rewrite of history. Deleting and recreating a stream loses its stored messages.
  • /api/v1/schema responses are client-focused. Internal storage_policy settings are not exposed.

Loading Precedence

Configuration is loaded in this order (later sources override earlier ones):

  1. ./configuration/config.yaml
  2. /etc/aviso_server/config.yaml
  3. $HOME/.aviso_server/config.yaml
  4. Environment variables (highest precedence)

If AVISOSERVER_CONFIG_FILE is set, only that single file is loaded (steps 1 through 3 are skipped). Environment variables still override values from the file.

Environment variable format

Prefix: AVISOSERVER_ Nested separator: __

AVISOSERVER_APPLICATION__HOST=0.0.0.0
AVISOSERVER_APPLICATION__PORT=8000
AVISOSERVER_NOTIFICATION_BACKEND__KIND=jetstream
AVISOSERVER_NOTIFICATION_BACKEND__JETSTREAM__NATS_URL=nats://localhost:4222

Config File Structure

The top-level sections are:

SectionPurpose
applicationServer host, port, static files path
loggingLog level and format
authAuthentication mode, secrets, admin roles
notification_backendBackend selection and backend-specific settings
notification_schemaPer-event-type validation, topic ordering, storage policy
metricsOptional Prometheus metrics server
watch_endpointSSE heartbeat, connection limits, replay batch settings

notification_backend.kind selects the storage implementation:

  • jetstream: production backend (NATS JetStream).
  • in_memory: development backend (process-local, no persistence).

Backend Details


For full field-level documentation of every config option, see Configuration Reference.

Configuration Reference

This page documents runtime-relevant configuration fields and defaults.

Topic Wire Format

  • Topic wire subjects always use . as separator.
  • Per-schema topic.separator is no longer used.
  • Token values are percent-encoded for reserved chars (., *, >, %) before writing to backend subjects.

See Topic Encoding for rules and examples.

application

SettingDefault
hostnone
portnone
base_urlhttp://localhost
static_files_path/app/static
homepagepublic client and server links
host Default: none · Type: string

Bind address.

port Default: none · Type: u16

Bind port.

base_url Default: http://localhost · Type: string

Used in generated CloudEvent source links.

static_files_path Default: /app/static · Type: string

Static asset root for homepage assets. The directory holds index.html, the homepage, and the files it links: logo.png, which is also the PNG favicon, and favicon.svg. Both icons are the documentation’s own. A custom directory needs all three files.

homepage Default: public project links · Type: object

The homepage puts the Aviso client card before the Aviso server card. Each card links to its documentation and GitHub repository. A separate HTTP API reference section opens this server’s Swagger UI.

All four fields are optional strings. Omitting the object or any field keeps the corresponding default:

FieldDefault
client_documentation_urlhttps://sites.ecmwf.int/docs/aviso-client/main/
client_repository_urlhttps://github.com/ecmwf/aviso-client
server_documentation_urlhttps://sites.ecmwf.int/docs/aviso-server/main/
server_repository_urlhttps://github.com/ecmwf/aviso-server

Set these under application.homepage in YAML. Environment overrides take priority over file values, using these exact keys:

AVISOSERVER_APPLICATION__HOMEPAGE__CLIENT_DOCUMENTATION_URL
AVISOSERVER_APPLICATION__HOMEPAGE__CLIENT_REPOSITORY_URL
AVISOSERVER_APPLICATION__HOMEPAGE__SERVER_DOCUMENTATION_URL
AVISOSERVER_APPLICATION__HOMEPAGE__SERVER_REPOSITORY_URL

Startup rejects URLs that cannot be parsed as absolute HTTP or HTTPS URLs with a host, or that contain a username or password. The error names the configuration field without printing its value. Links are HTML-escaped when rendered, including quotes and query-string ampersands.

logging

SettingDefault
levelinfo
formatimplementation default
level Default: info · Type: string

One of trace, debug, info, warn, error. Unknown values fall back to info instead of failing startup. Used as the application-wide level when RUST_LOG is unset.

format Default: implementation default · Type: string

Kept for compatibility; output is OTel-aligned JSON.

Runtime override via RUST_LOG

If the RUST_LOG environment variable is set, it takes priority over logging.level and gives the operator full EnvFilter directive syntax for runtime triage without a code change. Examples:

RUST_LOG=info,aviso_server=debug
RUST_LOG=warn,aviso_server::auth=trace
RUST_LOG=info,aviso_server::sse=debug,actix_web=warn

A malformed RUST_LOG value is reported on stderr at startup and the server falls back to logging.level. The most common parse failures are an empty target before = (for example RUST_LOG==warn) and a non-level value after = (for example RUST_LOG=info,aviso_server=verbose).

A missing comma like RUST_LOG=info aviso_server=debug does not trigger the fallback. EnvFilter parses the whole string as a single target name with a space, and the directive ends up matching nothing instead of failing loudly. If a RUST_LOG value looks correct but no logs appear, double-check the commas first.

RUST_LOG="" (empty string) is treated as if RUST_LOG were unset and falls back to logging.level. Without this guard EnvFilter::try_new("") silently succeeds with a filter that matches nothing and silences the entire process. This is a real failure mode under deployment systems that export unset variables as empty strings, such as the Kubernetes downward API or docker-compose’s ${VAR:-}.

When RUST_LOG is unset, the default filter combines logging.level with a small set of mute directives so that framework internals do not flood operational logs:

DirectiveEffect
actix_web=warnCaps Actix-web request lifecycle logs at warn (worker started, accepting, etc.).
actix_server=warnCaps Actix-server lifecycle logs at warn.
async_nats=infoCaps the NATS client at info; trace/debug per-message chatter stays off.

These mute directives are pinned by unit tests, only apply when RUST_LOG is unset, and only apply when the directive’s level is more restrictive than logging.level. With logging.level=warn or logging.level=error the directives are skipped entirely so they never raise the per-target ceiling above what the operator chose; with logging.level=info the two actix_*=warn directives narrow framework chatter while async_nats=info is skipped (it would be neutral); with logging.level=debug or logging.level=trace all three directives apply. Setting RUST_LOG opts out of all of them and gives the operator full directive control.

Push-based export via logging.otlp

Logs are always written to stdout as OTel-aligned JSON. With an otlp block the server additionally pushes every log record to an OpenTelemetry collector over OTLP, for clusters where log collection is push-based instead of scraping container output.

SettingDefault
enabledfalse
endpointnone
protocol"grpc"
enabled Default: false · Type: bool

Turns OTLP log export on. Startup fails when enabled without an endpoint.

endpoint Default: none · Type: string

Collector endpoint. A missing scheme defaults to http://. For protocol: http the OTLP path /v1/logs is appended when absent.

protocol Default: "grpc" · Type: "grpc"|"http"

Transport. Collectors conventionally listen on 4317 for gRPC and 4318 for HTTP.

logging:
  level: info
  format: json
  otlp:
    enabled: true
    endpoint: "http://otel-collector.observability.svc:4317"
    protocol: grpc

Operational behavior:

  • Export runs on a background batch thread with a bounded queue. A slow or unreachable collector never blocks request handling; overflow drops records from the export path only, and stdout remains complete.
  • Export errors are reported by the SDK’s internal diagnostics on stdout, so a broken collector connection is visible in kubectl logs.
  • Export health is measurable on the Prometheus endpoint: aviso_otlp_export_failures_total counts failed batch exports and aviso_otlp_suppressed_log_records_total counts records withheld by the redaction guard. Alert on a sustained non-zero rate of either.
  • Exported records carry the same request context as stdout records: request_id, username, auth_realm, event_type and topic, taken from the request when the log event does not set them itself. Each key appears once, with the event’s own value when it has one.
  • The global filter (logging.level / RUST_LOG) applies to both sinks, so the collector receives the same event stream as stdout. The export transport’s own targets (opentelemetry*, tonic, hyper, h2, tower, reqwest) are excluded from the export path to prevent feedback loops; they still appear on stdout.
  • Exported records carry the same resource identity as stdout records (service.name, service.version, and k8s.namespace.name / k8s.pod.name when the corresponding environment variables are set).
  • Redaction on the export path is stricter than stdout: record bodies get the same pattern redaction, but a record carrying a sensitive attribute key (password, secret, token, authorization, api_key) or a URL value with embedded credentials is withheld from export entirely. The stdout copy of the same record keeps field-level [REDACTED] markers, so no information is lost to operators.
  • On shutdown the server flushes buffered records before exiting.

The endpoint can also be injected without a config file change via environment overrides, for example AVISOSERVER_LOGGING__OTLP__ENDPOINT=http://collector:4317.

auth

Authentication is optional. When disabled (default), all API endpoints are publicly accessible only if schemas do not define stream auth rules. Startup fails if global auth is disabled while a schema sets auth.required=true or non-empty auth.read_roles/auth.write_roles.

When enabled:

  • Admin endpoints always require a valid JWT and an admin role.
  • Stream endpoints (notify, watch, replay) enforce authentication only when the target schema has auth.required: true.
  • Schema endpoints (/api/v1/schema) are always public.
  • In trusted_proxy mode, Aviso validates Authorization: Bearer <jwt> locally with jwt_secret.
SettingDefault
enabledfalse
mode"direct"
auth_o_tron_url""
jwt_secret""
admin_roles{}
timeout_ms5000
enabled Default: false · Type: bool

Set to true to enable authentication.

mode Default: "direct" · Type: "direct"|"trusted_proxy"

direct: forward credentials to auth-o-tron. trusted_proxy: validate forwarded JWT locally.

auth_o_tron_url Default: "" · Type: string

auth-o-tron base URL. Required when enabled=true and mode=direct.

jwt_secret Default: "" · Type: string

Shared HMAC secret for JWT validation. Required when enabled=true. Not exposed via /api/v1/schema endpoints and redacted when auth settings are serialized or logged.

admin_roles Default: {} · Type: map<string, string[]>

Realm-scoped roles for admin endpoints (/api/v1/admin/*). Must contain at least one realm with non-empty roles when enabled=true.

timeout_ms Default: 5000 · Type: u64

Timeout for auth-o-tron requests (milliseconds). Must be > 0.

Per-stream auth (notification_schema.<event_type>.auth)

SettingDefault
required(none)
read_roles(none)
write_roles(none)
plugins(none)
required Default: (none) · Type: bool

Must be explicitly set whenever an auth block is present. When true, the stream requires authentication.

read_roles Default: (none) · Type: map<string, string[]>

Realm-scoped roles for read access (watch/replay). When omitted, any authenticated user can read. Use ["*"] as the role list to grant realm-wide access.

write_roles Default: (none) · Type: map<string, string[]>

Realm-scoped roles for write access (notify). When omitted, only users matching global admin_roles can write. Use ["*"] as the role list to grant realm-wide access.

plugins Default: (none) · Type: string[]

Optional list of authorization plugins to run after role-based checks. Currently supported: "ecpds" (requires --features ecpds build). On a build without the required feature, startup fails with a clear error pointing at the offending stream. (Silent skip would widen access.) Empty plugins: [] is rejected; omit the field instead. Plugins only run when auth.required is true.

See Authentication for detailed setup, client usage, and error responses.

ecpds

Optional ECPDS destination authorization, available when built with --features ecpds. Add "ecpds" to a stream’s auth.plugins list to check destination access on watch and replay requests.

username Default: none · Type: nonempty string

Service account username used for HTTP Basic Auth to ECPDS.

password Default: none · Type: nonempty string

Service account password used for HTTP Basic Auth to ECPDS. It is redacted in configuration debug output. The schema discovery API does not expose the top-level ecpds settings.

servers Default: none · Type: list of URL strings

Use HTTPS to protect credentials and destination lookups. HTTP is accepted only for local testing with 127.0.0.1, [::1], or localhost; other HTTP addresses fail startup validation.

servers is a list of base URL strings, without query strings or fragments. Path prefixes such as https://proxy.example/ecpds-api/ are supported. Aviso appends /ecpds/v1/destination/list?id=<username> to each base URL.

match_key Default: none · Type: string

Set match_key to an ordinary identifier such as destination, declared in the schema with required: true. The name must not contain whitespace, /, or NUL. When the schema defines a topic, include this field in topic.key_order so delivery is filtered by the authorized destination.

Spatial identifiers cannot be match keys: PolygonHandler, PointCloudHandler, and the field name polygon are not allowed. Spatial matching does not enforce access to an exact destination value.

target_field Default: "name" · Type: string

Selects a JSON field from each ECPDS destination record. Records missing that field are skipped. To investigate missing destinations, set RUST_LOG=info,aviso_ecpds=debug and look for auth.ecpds.fetch.skipped_record events.

cache_ttl_seconds Default: 300 · Unit: seconds · Minimum: 1

How long to cache a user’s destination list before fetching it again. Use a whole number of seconds.

max_entries Default: 10000 · Unit: users · Minimum: 1

Maximum number of users in the destination cache. Use a whole number. The cache uses TinyLFU eviction when it needs to make room.

request_timeout_seconds Default: 30 · Unit: seconds · Minimum: 1

Maximum time for the whole ECPDS request, from DNS lookup through reading the response body. Use a whole number of seconds.

connect_timeout_seconds Default: 5 · Unit: seconds · Minimum: 1

Maximum time to establish the connection, including TCP and TLS. This counts toward the total request timeout; it is not extra time. Use a whole number of seconds.

partial_outage_policy Default: "strict" · Values: "strict", "any_success"

Controls what happens when an ECPDS server is unavailable:

  • strict: every configured server must respond successfully. If any fails, the destination lookup fails with HTTP 503.
  • any_success: combine destinations from the servers that respond successfully. The lookup fails if none succeeds.

In both modes, Aviso combines the returned destination lists. With any_success, destinations known only to an unavailable server may be missing. See Partial outage policy for the trade-off.

See ECPDS Destination Authorization for setup and runtime behavior, and the ECPDS runbook for troubleshooting.

metrics

Optional Prometheus metrics endpoint. When enabled, a separate HTTP server serves /metrics on an internal port for scraping by Prometheus/ServiceMonitor. This keeps metrics isolated from the public API.

SettingDefault
enabledfalse
host"127.0.0.1"
portnone
enabled Default: false · Type: bool

Enable the metrics endpoint.

host Default: "127.0.0.1" · Type: string

Bind address for the metrics server. Defaults to loopback to avoid public exposure.

port Default: none · Type: u16

Required when enabled=true. Must differ from application.port.

Exposed metrics:

aviso_build_info Type: gauge · Labels: version

Constant 1 with the server version as a label; join on it in dashboards to annotate deploys.

aviso_http_requests_total Type: counter · Labels: route, method, status_code

HTTP requests on the main server by matched route pattern (e.g. /api/v1/schema/{event_type}). Reserved label values: unrouted requests (404 scans) collapse into route="unmatched", requests failing with a service-level error (no route information available) record route="error", and non-standard HTTP methods collapse into method="other". The label is named route (not endpoint) to avoid colliding with the Prometheus Operator target label endpoint.

aviso_http_request_duration_seconds Type: histogram · Labels: route, method

Request duration until response headers are ready. For the SSE routes (/api/v1/watch, /api/v1/replay) this is stream setup latency, not connection lifetime; see aviso_sse_connection_duration_seconds.

aviso_http_requests_in_flight Type: gauge · Labels: method

HTTP requests currently being processed, by method. Labelled by method only because the matched route pattern is not known until routing completes (after the request is already in flight). Distinguishes “slow because busy” from “slow because a downstream/backend stalled”.

aviso_backend_operations_total Type: counter · Labels: backend, operation, outcome

Notification-backend operations at the trait boundary. operation ∈ {publish, get_batch, wipe_stream, wipe_all, delete_message}; outcome ∈ {ok, error}. subscribe_to_topic is excluded (its work happens lazily as the stream is polled).

aviso_backend_operation_duration_seconds Type: histogram · Labels: backend, operation, outcome

Caller-observed backend operation latency (same labels as aviso_backend_operations_total). This is the metric to watch when notification throughput plateaus while pods are underused: it isolates backend (NATS/JetStream) latency from app CPU.

aviso_notifications_total Type: counter · Labels: event_type, status

Total notification requests. status ∈ {success, error, rejected}; requests failing before schema validation record event_type="unknown".

aviso_sse_connections_active Type: gauge · Labels: route, event_type

Currently active SSE connections. route ∈ {/api/v1/watch, /api/v1/replay}.

aviso_sse_connections_total Type: counter · Labels: route, event_type

Total SSE connections opened.

aviso_sse_unique_users_active Type: gauge · Labels: route

Distinct users with active SSE connections.

aviso_sse_events_sent_total Type: counter · Labels: route, event_type

Notification events delivered to SSE clients. Heartbeats, control events, and close frames are not counted.

aviso_sse_stream_errors_total Type: counter · Labels: route, event_type

Error events emitted into SSE streams after the response started (typed stream errors and notification rendering failures); these are invisible to aviso_http_requests_total because the stream already returned 200.

aviso_sse_connection_duration_seconds Type: histogram · Labels: route

SSE connection lifetime, observed when the connection closes (buckets 1s-24h). Long-lived open connections appear in aviso_sse_connections_active, not here, until they close.

aviso_auth_requests_total Type: counter · Labels: mode, outcome

Authentication attempts. mode ∈ {direct, trusted_proxy}; outcome ∈ {success, unauthorized, forbidden, service_unavailable}.

The SSE and HTTP request metrics share a route label whose values are real route patterns (e.g. /api/v1/watch), so a single dashboard route variable spans both. Like the ECPDS counters below, the bounded label combinations of aviso_auth_requests_total, aviso_notifications_total (including one series per configured stream), and aviso_backend_operations_total / aviso_backend_operation_duration_seconds (per active backend) are pre-initialised at zero on startup so rate(...) > 0 alert rules evaluate against existing series.

A binary built with --features ecpds registers the following five metrics. The unlabelled counters and the gauge appear as Prometheus series at process startup. The two labelled counters (access_decisions_total, fetch_total) are pre-initialised at startup with every documented outcome value, so each outcome label appears as a series at zero before any ECPDS traffic; this lets alert rules of the form rate(metric{outcome="error"}[5m]) > 0 start evaluating on a known-zero baseline rather than on a missing series.

aviso_ecpds_cache_hits_total Type: counter · Labels: (none)

ECPDS destination cache hits (requests served from cache without an upstream call).

aviso_ecpds_cache_misses_total Type: counter · Labels: (none)

ECPDS destination cache misses (requests not served from cache). Includes coalesced waiters that did not trigger an upstream call themselves; aviso_ecpds_fetch_total is the right metric for “actual upstream calls”.

aviso_ecpds_cache_size Type: gauge · Labels: (none)

Number of usernames in the ECPDS destination cache, sampled from moka after eviction passes. Expired entries are pruned by moka asynchronously, so this gauge can briefly include not-yet-pruned expired entries until the next pending-tasks run.

aviso_ecpds_access_decisions_total Type: counter · Labels: outcome

Access decisions. outcome ∈ {allow, deny_destination, deny_match_key_missing, unavailable, admin_bypass, error}.

aviso_ecpds_fetch_total Type: counter · Labels: outcome

Upstream fetch outcomes (recorded once per access check whose request actually ran the upstream call; coalesced waiters do not contribute). outcome ∈ {success, http_401, http_403, http_4xx, http_5xx, invalid_response, unreachable}.

Process-level metrics (CPU, memory, open FDs) are automatically collected on Linux.

notification_backend

FieldTypeDefaultNotes
kindstringnonejetstream or in_memory.
in_memoryobjectoptionalUsed when kind = in_memory.
jetstreamobjectoptionalUsed when kind = jetstream.

notification_backend.in_memory

FieldTypeDefaultNotes
max_history_per_topicusize1Retained messages per topic in memory.
max_topicsusize10000Max tracked topics before LRU-style eviction.
enable_metricsboolfalseEnables extra internal metrics logs.

See InMemory Backend for operational caveats.

notification_backend.jetstream

nats_url Default: nats://localhost:4222 · Type: string

NATS connection URL.

token Default: None · Type: string?

Token auth; NATS_TOKEN env fallback.

timeout_seconds Default: 30 · Type: u64?

NATS connection timeout for each startup connect attempt (> 0).

retry_attempts Default: 3 · Type: u32?

Startup connect attempts before backend init fails (> 0).

max_messages Default: None · Type: i64?

Stream message cap.

max_bytes Default: None · Type: i64?

Stream size cap in bytes.

retention_time Default: None · Type: string?

Default stream max age (s, m, h, d, w; for example 30d).

storage_type Default: file · Type: string?

file or memory (parsed as typed enum at config load).

Omitting this setting requests file. Existing streams must use the requested type: a mismatch fails stream setup before any mutable settings change. The error names the stream and its current and requested types. Aviso does not delete or recreate streams, so messages remain intact. Use the current type or arrange a separate migration.

replicas Default: None · Type: usize?

Stream replicas.

retention_policy Default: limits · Type: string?

limits/interest. workqueue fails startup because independent watch/replay consumers are not supported.

discard_policy Default: old · Type: string?

old/new (parsed as typed enum at config load).

max_reconnect_attempts Default: unlimited · Type: u32?

Mapped to NATS max_reconnects; unset and 0 both mean unlimited. A positive value makes the client give up permanently once exhausted.

Subscription creation uses a bounded retry loop: unset means five attempts, 0 means one attempt, and a positive value sets the attempt limit. Startup connection attempts are controlled separately by retry_attempts.

reconnect_delay_ms Default: 2000 · Type: u64?

Reconnect delay and startup connect retry backoff (> 0).

publish_retry_attempts Default: 5 · Type: u32?

Retry attempts for transient publish channel closed failures (> 0).

publish_retry_base_delay_ms Default: 150 · Type: u64?

Base backoff in milliseconds for publish retries (> 0).

See JetStream Backend for detailed behavior.

notification_schema_strict

Controls how the server treats event_type values that are not declared in notification_schema.

SettingDefault
notification_schema_strictderived
notification_schema_strict Default: derived · Type: bool?

When unset, the effective value is true if notification_schema is non-empty, false otherwise. Set to true to force strict rejection even with no schema (deny-all “drain” mode). Set to false to preserve the legacy permissive generic fallback even with a declared schema; a startup warning is emitted in that case.

In strict mode, POST /api/v1/notification, POST /api/v1/watch, and POST /api/v1/replay reject any event_type not present in notification_schema with 400 UNKNOWN_EVENT_TYPE. The error body is:

{
  "code": "UNKNOWN_EVENT_TYPE",
  "error": "unknown_event_type",
  "message": "unknown event type 'X'",
  "configured_event_types": ["dissemination", "mars", "test_polygon"],
  "request_id": "<uuid>"
}

configured_event_types is sorted for stable diffing in client tooling.

The same flag also bounds Prometheus / tracing label cardinality. Whenever effective strict mode is off (either notification_schema_strict is explicitly false, or it is unset with an empty/absent notification_schema so the startup default resolves to non-strict), a request whose event_type is not in the schema reaches the generic-fallback path and has its recorded event_type label collapsed to the literal "generic" instead of being persisted as user-controlled input.

notification_schema.<event_type>.payload

Schema-level payload contract for notify requests.

FieldTypeExampleNotes
requiredbooltrueWhen true, /notification rejects requests without payload.

Behavior details and edge cases are documented in Payload Contract.

notification_schema.<event_type>.max_historical_notifications

Optional positive integer overriding watch_endpoint.max_historical_notifications for this event type. Omit it to inherit the global cap (default 10000). Zero and unlimited are rejected. This field sits outside storage_policy and works with both backends. See Replay Limit for an example and Historical Replay Limits for the wire behavior.

notification_schema.<event_type>.storage_policy

Optional per-schema storage settings validated at startup against selected backend capabilities.

SettingExample
retention_time7d, 12h, 30m
max_messages100000
max_size512Mi, 2G
allow_duplicatestrue
compressiontrue
retention_time Example: 7d, 12h, 30m · Type: string

Duration literal (s, m, h, d, w).

max_messages Example: 100000 · Type: integer

Must be > 0.

max_size Example: 512Mi, 2G · Type: string

Size literal (K, Ki, M, Mi, G, Gi, T, Ti).

allow_duplicates Example: true · Type: bool

Backend support is capability-gated.

compression Example: true · Type: bool

Backend support is capability-gated.

Field behavior:

  • retention_time overrides backend-level retention for the schema stream.
  • max_messages overrides backend-level message cap for the schema stream.
  • max_size overrides backend-level byte cap for the schema stream.
  • allow_duplicates = false maps to one message per subject (latest kept); true removes this cap.
  • compression = true enables stream compression when backend supports it.

Startup behavior:

  • Schema storage policies inherit the backend retention_policy; there is no per-schema override. workqueue fails startup because it does not support independent watch/replay consumers. Existing streams are not migrated or deleted by this validation.
  • Invalid retention_time/max_size format fails startup.
  • Unsupported fields for selected backend fail startup.
  • Validation happens before backend initialization.
  • With in_memory, all storage_policy fields are currently unsupported (startup fails if provided).

Runtime application behavior:

  • Aviso uses the configuration loaded at startup. Config edits require a restart or rollout to all replicas; editing the file alone does not change policy.
  • storage_policy is applied on stream create and reconciled for existing JetStream streams when those streams are accessed by Aviso. There is no all-stream sweep at startup or in the background.
  • Aviso-managed stream subject binding is also reconciled to the expected <base>.> pattern.
  • Mutable fields (retention/limits/compression/duplicates/replicas) are updated when drift is detected.
  • Compression applies to future file-storage writes at the block level. Changing it does not automatically recompress existing history.
  • Deleting and recreating a stream loses its stored messages; it does not rewrite history. Aviso provides no automatic history migration.

Example:

notification_backend:
  kind: jetstream
  jetstream:
    nats_url: "nats://localhost:4222"
    publish_retry_attempts: 5
    publish_retry_base_delay_ms: 150

notification_schema:
  dissemination:
    topic:
      base: "diss"
      key_order: ["destination", "target", "class", "expver", "domain", "date", "time", "stream", "step"]
    storage_policy:
      retention_time: "7d"
      max_messages: 2000000
      max_size: "10Gi"
      allow_duplicates: true
      compression: true

watch_endpoint

sse_heartbeat_interval_sec Default: 30 · Type: u64

SSE heartbeat period.

connection_max_duration_sec Default: 3600 · Type: u64

Maximum live watch duration.

replay_batch_size Default: 100 · Type: usize

Historical backend fetch batch size, independent of the request-wide delivery cap. Filtering can leave a batch with no notifications to deliver; replay keeps advancing through history.

max_historical_notifications Default: 10000 · Type: usize

Maximum historical notifications delivered per replay or replaying watch request, across all batches and after identifier constraints, spatial filters and successful CloudEvent rendering. Both backends enforce the same cap. Live-only watches are unaffected. A schema can override this default with notification_schema.<event_type>.max_historical_notifications, outside storage_policy. Omitting the schema field inherits the global value.

The value must be a positive integer; zero and unlimited are rejected. The server emits notification_replay_limit_reached only when it finds a renderable notification beyond the cap. Exactly filling the quota is not truncation. A truncated request closes without replay_completed or a transition to live delivery. See Historical Replay Limits.

replay_batch_delay_ms Default: 100 · Type: u64

Delay between historical replay batches.

concurrent_notification_processing Default: 15 · Type: usize

Live stream CloudEvent conversion concurrency.

Custom config file path

Set AVISOSERVER_CONFIG_FILE to use a specific config file instead of the default search cascade:

AVISOSERVER_CONFIG_FILE=/path/to/config.yaml cargo run

When set, only this file is loaded as a file source (startup fails if it does not exist). The default locations (./configuration/config.yaml, /etc/aviso_server/config.yaml, $HOME/.aviso_server/config.yaml) are skipped. AVISOSERVER_* field-level overrides still apply on top.

Environment override examples

AVISOSERVER_APPLICATION__HOST=0.0.0.0
AVISOSERVER_APPLICATION__PORT=8000
AVISOSERVER_NOTIFICATION_BACKEND__KIND=jetstream
AVISOSERVER_NOTIFICATION_BACKEND__JETSTREAM__NATS_URL=nats://localhost:4222
AVISOSERVER_NOTIFICATION_BACKEND__JETSTREAM__TOKEN=secret
AVISOSERVER_WATCH_ENDPOINT__REPLAY_BATCH_SIZE=200
AVISOSERVER_AUTH__ENABLED=true
AVISOSERVER_AUTH__JWT_SECRET=secret
AVISOSERVER_METRICS__ENABLED=true
AVISOSERVER_METRICS__PORT=9090

Authentication

Authentication is optional. When enabled, Aviso supports two modes:

  • Direct. Aviso forwards Bearer or Basic credentials to auth-o-tron, which returns a signed JWT.
  • Trusted proxy. An upstream reverse proxy authenticates the user and forwards a signed JWT; Aviso validates it locally.

How It Works

  1. Client sends credentials to Aviso.
  2. Middleware resolves user identity:
    • direct: forwards the Authorization header to auth-o-tron GET /authenticate and receives a JWT back.
    • trusted_proxy: validates the forwarded Authorization: Bearer <jwt> locally using jwt_secret.
  3. Username, realm, and roles are extracted from JWT claims and attached to the request.
  4. Route handlers enforce per-stream auth rules on notify, watch, and replay.
  5. Admin endpoints (/api/v1/admin/*) always require a valid JWT with an admin role.

Schema endpoints (GET /api/v1/schema, GET /api/v1/schema/{event_type}) are always publicly accessible, even when auth is enabled.

Quick Start (Direct Mode)

1. Start auth-o-tron

# Foreground (Ctrl+C to stop):
./scripts/auth-o-tron-docker.sh start

# Background:
./scripts/auth-o-tron-docker.sh start --detach

By default this runs auth-o-tron 0.3.7 with scripts/example_auth_config.yaml, bound to 127.0.0.1:8080. Use AUTH_O_TRON_PORT and AUTH_O_TRON_CONTAINER_NAME for an isolated instance. Set AUTH_O_TRON_BIND_ADDRESS explicitly to expose another interface. The launcher requires Python 3 to validate the bind IP and port before replacing an existing container. IPv6 addresses can be bare (::1) or bracketed ([::1]). Scoped IPv6 addresses such as fe80::1%eth0 are not supported by Docker and are rejected before container replacement. To check the bundled users against a running local instance without printing tokens:

AVISO_TEST_AUTH_O_TRON_URL=http://127.0.0.1:8080 \
cargo test --locked --test auth_o_tron_live

To use your own config:

AUTH_O_TRON_CONFIG_FILE=/path/to/auth-config.yaml ./scripts/auth-o-tron-docker.sh start

The bundled example config defines three local test users in realm localrealm:

UserPasswordRole
admin-useradmin-passadmin
reader-userreader-passreader
producer-userproducer-passproducer

2. Enable auth in config

auth:
  enabled: true
  mode: direct
  auth_o_tron_url: "http://localhost:8080"
  jwt_secret: "your-shared-secret"   # must match auth-o-tron jwt.secret
  admin_roles:
    localrealm: ["admin"]
  timeout_ms: 5000

Roles are realm-scoped: admin_roles maps each realm name to its authorized role list. A user must belong to a listed realm and hold one of that realm’s roles.

3. Run aviso-server

Auth is now enforced for:

  • Admin endpoints (/api/v1/admin/*) always require auth and an admin role.
  • Stream endpoints (/api/v1/notification, /api/v1/watch, /api/v1/replay) require auth only when the target schema sets auth.required: true.

For full field-level documentation, see the auth section in Configuration Reference.

Trusted Proxy Mode

Use trusted_proxy when Aviso sits behind a reverse proxy or API gateway that handles authentication. The proxy authenticates the user (via OIDC, SAML, etc.) and forwards a signed JWT to Aviso.

Aviso validates the forwarded Authorization: Bearer <jwt> locally using jwt_secret. Username and roles are read directly from JWT claims; no outbound call to auth-o-tron is made.

auth:
  enabled: true
  mode: trusted_proxy
  jwt_secret: "shared-signing-secret"
  admin_roles:
    ecmwf: ["admin"]

auth_o_tron_url is not required in this mode.

Per-Stream Authentication

Streams support separate read and write access controls. Read access governs /watch and /replay; write access governs /notification.

Configure authentication per stream in your notification schema:

notification_schema:
  # Public: no auth section means anonymous access
  public_events:
    payload:
      required: true
    topic:
      base: "public"

  # Authenticated: any valid user can read, only admins can write
  internal_events:
    payload:
      required: true
    topic:
      base: "internal"
    auth:
      required: true

  # Separate read/write roles
  sensor_data:
    payload:
      required: true
    topic:
      base: "sensor"
    auth:
      required: true
      read_roles:
        internal: ["analyst", "consumer"]
        external: ["partner"]
      write_roles:
        internal: ["producer"]

  # Realm-wide read access using wildcard, restricted write
  shared_events:
    payload:
      required: true
    topic:
      base: "shared"
    auth:
      required: true
      read_roles:
        internal: ["*"]
        external: ["analyst"]
      write_roles:
        internal: ["producer", "operator"]

Read vs. write access defaults

auth.requiredread_roleswrite_rolesRead (watch/replay)Write (notify)
false or omitted(any)(any)AnyoneAnyone
trueomittedomittedAny authenticated userAdmins only
truesetomittedMust match read_rolesAdmins only
trueomittedsetAny authenticated userMust match write_roles or be admin
truesetsetMust match read_rolesMust match write_roles or be admin

Admins (users matching global admin_roles) always have both read and write access.

Role matching rules

Both read_roles and write_roles map realm names to role lists. A user’s realm claim from the JWT must match a key in the map, and the user must hold at least one of that realm’s listed roles.

  • Wildcard "*". Use ["*"] as the role list to grant access to all users from a realm, regardless of their specific roles.
  • Omitted role list. When read_roles is omitted, any authenticated user can read. When write_roles is omitted, only admins can write.

When a per-stream auth block is present, auth.required must be explicitly set to either true or false.

ECPDS Destination Authorization

When built with --features ecpds, Aviso supports an optional authorization plugin that checks whether a user has access to a specific ECPDS destination before allowing watch or replay requests. The plugin is read-only: it never runs on notify.

Enabling the plugin

  1. Build Aviso with the ecpds feature:

    cargo build --release --features ecpds
    

    On a build without this feature, any YAML containing plugins: ["ecpds"] is rejected at startup with an error pointing at the offending stream. This is deliberate: silently skipping the plugin would widen access.

  2. Add a top-level ecpds section to your config with ECPDS service credentials:

    ecpds:
      username: "ecpds-service-account"
      password: "service-password"
      servers:
        - "https://ecpds-primary.ecmwf.int"
        - "https://ecpds-secondary.ecmwf.int"
      match_key: "destination"
      target_field: "name"            # default: "name"
      cache_ttl_seconds: 300          # default: 300 (5 min)
      max_entries: 10000              # default: 10000
      request_timeout_seconds: 30     # default: 30
      connect_timeout_seconds: 5      # default: 5
      partial_outage_policy: strict   # default: strict; alternative: any_success
    
  3. Enable the plugin on a stream by adding plugins: ["ecpds"] to its auth block. Minimal canonical shape:

    notification_schema:
      dissemination:
        payload:
          required: true
        topic:
          base: "diss"
          key_order: ["destination", "target", "class", "expver", "domain", "date", "time", "stream", "step"]
        identifier:
          destination:
            type: StringHandler
            required: true                # MUST be required
          # ... other fields ...
        auth:
          required: true                  # MUST be true
          plugins: ["ecpds"]
    

read_roles is optional. If you want a realm-wide gate before ECPDS even runs (e.g. block users from realms you don’t trust to query ECPDS in the first place), add it; if not, omit it and the plugin runs for every authenticated user.

The plugin requires (and startup validation enforces):

  • match_key (default "destination") is present in the schema’s identifier and marked required: true there. Use an ordinary destination identifier, not a geometry: PolygonHandler, PointCloudHandler, and the field name polygon are not allowed as match keys. Spatial filters find matching areas; they do not enforce access to an exact destination value. When the schema defines a topic, include the match key in topic.key_order so notifications are filtered by the authorized destination.
  • auth.required is true. The plugin runs after standard stream auth, so plugins on a stream where auth.required is false would never execute.

How it works at runtime

  1. Standard role-based stream auth runs first. If it fails (missing token, wrong realm/role), the request fails before the ECPDS plugin sees it.
  2. The plugin extracts the match_key value (e.g. destination) from the request’s canonicalised identifier.
  3. It looks up the user’s destination list in an in-process cache. If absent, it queries the configured ECPDS servers in parallel, then merges per the partial_outage_policy.
  4. If the requested destination is in the user’s list, the request proceeds. Otherwise, 403 Forbidden.
  5. Users matching the global auth.admin_roles bypass step 2-4 entirely.

Partial-outage policy

When more than one ECPDS server is configured, the user’s effective destination list is always the union of every per-server response. ECMWF ECPDS deployments are typically federated (e.g. diss-monitor and aux-monitor cover different destination namespaces), so a user’s full entitlement is the combination of what each server reports. The partial_outage_policy field only governs how tolerant the merge is when one of those servers fails or times out.

ValueBehaviourOperational implication
strict (default)Every configured server must reply successfully within the per-request timeout. The destination list is the union of their responses. Any one server failing fails the whole lookup with 503.A single ECPDS server going down takes the plugin to 503. The trade is: 503 (try again later) is preferred over 403 (definitely no access) when we can’t be sure we saw the user’s complete entitlement set.
any_successTake the union of whichever servers responded successfully within the per-request timeout. Failed servers are silently dropped from the merge. Only fails if no server responded usefully.Keeps serving during a partial outage. The cost: if a user’s only entitlement to a destination lived on an unreachable server, that user will see 403 until the server is back, even though their access is genuinely valid.

Error responses

CodeHTTP Status
FORBIDDEN403
SERVICE_UNAVAILABLE503
INTERNAL_ERROR500
FORBIDDEN HTTP Status: 403

User does not have access to the requested destination, or the required identifier field is missing. Tracing event: auth.ecpds.check.denied (with reason ∈ {DestinationNotInList, MatchKeyMissing}).

SERVICE_UNAVAILABLE HTTP Status: 503

Upstream / network problem: the lookup failed under the active partial_outage_policy. Tracing event: auth.ecpds.check.unavailable. The cause is on the aviso_ecpds_fetch_total{outcome=…} metric (e.g. unreachable, http_401, http_4xx, http_5xx, invalid_response). Investigate ECPDS, the network, and the service-account credentials.

INTERNAL_ERROR HTTP Status: 500

Aviso-side server error: missing AuthSettings in app_data, no checker registered, or an unexpected plugin error. Tracing event: auth.ecpds.check.error (with error_kind for the misconfiguration cases). Investigate Aviso, not ECPDS.

Caching

Destination lists are cached per user for cache_ttl_seconds (default 300 seconds). The cache holds at most max_entries users (default 10 000) and uses moka’s TinyLFU eviction policy when full (TinyLFU mixes recency with admission frequency, so a one-shot scan does not flush the working set). Successful results are cached. Errors are not cached. A short ECPDS outage will not get extended by stale 503s sitting in the cache.

The cache is single-flight: when many requests for the same user arrive at the same time and the user is not yet cached, only one upstream call goes to ECPDS. The rest wait for that one call’s result. This protects ECPDS when many SSE clients reconnect at once.

The cache lives in process memory. Restarting Aviso clears it. Multiple replicas have independent caches.

What is not checked

notify (write) is never gated by ECPDS. The plugin applies only to reads (watch, replay).

No retries by design

Aviso does not retry failed ECPDS calls. A 503 is the signal to investigate ECPDS itself, not to bump timeouts. See the ECPDS runbook for triage steps.

For the full ecpds field reference, see the ecpds section in Configuration Reference. For metrics and tracing event names, see the ECPDS runbook.

Admin Endpoints

Admin endpoints always require authentication and one of the configured admin_roles, regardless of per-stream settings:

  • DELETE /api/v1/admin/notification/{id}
  • DELETE /api/v1/admin/wipe/stream
  • DELETE /api/v1/admin/wipe/all

See Admin Operations for request/response details.

Audit attribution in logs

Every notify, watch, replay, and admin request records who performed it. The request span carries two fields that appear on every log event the request produces, across all layers (the API-level line, the SSE stream events, and the backend events they trigger):

  • username: the authenticated username, or anonymous when authentication is disabled or the request carried no credentials.
  • auth_realm: the realm the identity authenticated through (for example localrealm or ecmwf), or none when the token carries no realm claim.

Both fields are always present, so log queries can rely on them unconditionally: attributes.username: producer-pgen finds every notification that identity published, and the request_id shared by the same events joins the API-level line with the backend lines it triggered.

Disabling Authentication

auth:
  enabled: false

Or omit the auth section entirely. When auth is disabled, all endpoints are publicly accessible.

Startup fails if global auth is disabled while any schema defines auth.required: true or non-empty auth.read_roles/auth.write_roles. Remove stream-level auth blocks before disabling global auth.

Client Usage

Bearer token (both modes)

# Watch an authenticated stream:
curl -N -H "Authorization: Bearer <jwt-token>" \
  -X POST http://localhost:8000/api/v1/watch \
  -H "Content-Type: application/json" \
  -d '{"event_type": "private_events", "identifier": {}}'

Basic credentials (direct mode only)

# Notify with Basic auth:
curl -X POST http://localhost:8000/api/v1/notification \
  -u "admin-user:admin-pass" \
  -H "Content-Type: application/json" \
  -d '{"event_type": "ops_events", "identifier": {"event_type": "deploy"}, "payload": "ok"}'

In direct mode, Aviso forwards Basic credentials to auth-o-tron, which authenticates the user and returns a JWT. The response JWT is validated and used for authorization.

Error Responses

Auth errors use a subset of the standard API error shape with three fields (code, error, message; no details):

{
  "code": "UNAUTHORIZED",
  "error": "unauthorized",
  "message": "Authorization header is required"
}
CodeHTTP StatusWhen
UNAUTHORIZED401Missing Authorization header, invalid token format, expired or bad signature.
FORBIDDEN403Valid credentials but user lacks the required role for the stream or admin endpoint.
SERVICE_UNAVAILABLE503auth-o-tron is unreachable or returned an unexpected error (direct mode only).

A 401 response includes a WWW-Authenticate header indicating the supported scheme (Bearer in trusted-proxy mode; Bearer, Basic in direct mode).

ECPDS Plugin Runbook

Use this page to investigate ECPDS authorization problems on watch or replay requests. For setup, see ECPDS Destination Authorization.

At a glance

  • The plugin is read-only (watch, replay). The notify endpoint is never gated by ECPDS.
  • The plugin fails closed: it will never accidentally allow a request. The status code distinguishes where the problem is. 503 Service Unavailable means the ECPDS check could not reach a verdict (an upstream / partial-outage problem); investigate ECPDS and the network. 500 Internal Server Error means the plugin itself hit a server-side bug or a misconfiguration on Aviso’s side (missing AuthSettings, no checker registered, an unexpected plugin error); investigate Aviso. The full mapping is in the response codes table below.
  • The plugin does not retry. A 503 is the signal to investigate ECPDS; a 500 is the signal to investigate Aviso.
  • The cache lives in process memory. Restarting Aviso clears it. Replicas have independent caches.
  • The default partial_outage_policy is strict: every configured ECPDS server must respond successfully or the call fails with 503. A single ECPDS server going away takes the whole plugin down. This is intentional. The destination list itself is the union of every server’s response under both policies; the choice is purely about how tolerant we are of per-server failures.

Response codes the plugin emits

HTTPWhere to look
200Allowed
403Authorization
503ECPDS or the network
500Aviso
200 Allowed Destination access approved.

The destination is in the user’s ECPDS allow-list.

Tracing event: auth.ecpds.check.allowed

403 Authorization denied Check the requested destination and match key.

The destination is not in the user’s allow-list (reason=DestinationNotInList), or the request omitted the configured match_key field (reason=MatchKeyMissing).

Tracing event: auth.ecpds.check.denied

503 ECPDS or network failure The destination lookup could not reach a verdict.

The combined ECPDS responses could not satisfy the active partial_outage_policy. Check ECPDS availability, network connectivity and service-account credentials. The event’s fetch_outcome field helps narrow down the cause.

Tracing event: auth.ecpds.check.unavailable

500 Aviso error Investigate Aviso, not ECPDS.

A server-side bug or local misconfiguration prevented the check. Possible causes include missing AuthSettings or EcpdsChecker in app_data, or an unexpected plugin error.

Tracing event: auth.ecpds.check.error

Symptom and first checks

Start with one affected request and find its plugin event in the logs. An HTTP 403 or 503 alone does not show that ECPDS caused it. Use the same time window and Aviso replica when comparing logs with metrics. The username values show who was affected, not what caused the problem.

Watch or replay returns 503

event_name=auth.ecpds.check.unavailable confirms that the plugin could not get a usable destination list under the configured partial_outage_policy.

  1. Find event_name=auth.ecpds.fetch.failed for that user and time. Read server and error to identify the failing ECPDS server and its error.
  2. Check that server from the Aviso host, using the configured service account and the affected username as the lookup ID. For connection or timeout errors, check the URL, DNS and connectivity. For HTTP 401 or 403, verify the service account’s credentials and access. For other HTTP errors, inspect the status: check the URL for 404, throttling for 429, and upstream logs for 5xx. For an invalid response, inspect the returned body and its success value rather than assuming the API changed.
  3. Check partial_outage_policy: strict needs every server to succeed; any_success needs at least one. Use the per-server failure logs to see which servers need attention.

Metrics: aviso_ecpds_access_decisions_total{outcome="unavailable"} counts requests that the plugin rejected with 503. aviso_ecpds_fetch_total, grouped by outcome, counts fetch attempts across the configured servers, not individual server calls. Labels include unreachable, http_401, http_403, http_4xx, http_5xx and invalid_response. The request log uses fetch_outcome with values such as Unreachable or Unauthorized; the fetch failure log uses error, not outcome or fetch_outcome.

Failed lookups are not cached, but concurrent requests for the same user can share one fetch. The request and fetch counters need not rise together. Under any_success, a failure label on the fetch metric can also accompany a usable list, so it does not by itself mean a request returned 503.

Watch or replay returns 403

event_name=auth.ecpds.check.denied with reason=DestinationNotInList means the requested destination was absent from the list Aviso used for that user. That list may have come from cache, not a new ECPDS call. If the reason is MatchKeyMissing, use the next section instead.

  1. Check the event’s username and event_type against the intended user and schema. Compare the request’s destination with the configured match_key and any schema rules that change its value before the check.
  2. Query the configured ECPDS servers using Aviso’s service account and that username as the lookup ID. Check that the destination record has active: true and a string value in target_field. Debug events auth.ecpds.fetch.skipped_inactive and auth.ecpds.fetch.skipped_record identify servers whose records were excluded. With any_success, also check auth.ecpds.fetch.failed: a failed server’s destinations are absent.
  3. Read cache_outcome on the denial. hit means Aviso reused the user’s cached list. If ECPDS access was recently changed, retry after cache_ttl_seconds expires. Check the same replica, since each has its own cache.

Metric: aviso_ecpds_access_decisions_total{outcome="deny_destination"} counts denied requests, including repeated checks against a cached list. Those cache hits do not increase aviso_ecpds_fetch_total. This does not tell you how many users are affected; use the denial logs for that.

Logs show MatchKeyMissing

event_name=auth.ecpds.check.denied with reason=MatchKeyMissing means the configured match_key was absent from the processed request identifiers. The plugin returns 403 before consulting the cache or ECPDS.

  1. Use event_type to identify the schema. Compare ecpds.match_key with its identifier field name and the actual request body.
  2. Check the deployed version and the configuration loaded at startup. Startup validation requires the match key to exist in the schema with required: true; normal request validation rejects an omitted required field before the plugin runs. This event alone does not explain how the key went missing.
  3. If those settings match, keep the request ID and a redacted request example for an Aviso bug report. Investigate how request processing passed identifiers without the key to the checker, rather than changing ECPDS permissions.

Metric: aviso_ecpds_access_decisions_total{outcome="deny_match_key_missing"} counts these denials. The denial event has cache_outcome="none"; this path does not increase the cache or fetch counters.

Requests succeed, but no ECPDS checks appear

An absent auth.ecpds.check.allowed event does not prove the plugin is off. Admin requests bypass the destination check, and log filters can hide events.

  1. Confirm that you are looking at new watch or replay requests for the intended event_type, not an already-open stream or notify traffic. Check the logs and metrics for the replica handling those requests.
  2. Check aviso_ecpds_access_decisions_total{outcome="admin_bypass"}. An increase explains why there are no allow events for admin requests. The matching auth.ecpds.admin.bypass event is debug-level; the normal auth.ecpds.check.allowed event is info-level. Check the logging filters.
  3. For a non-admin request, verify that the deployed schema’s auth block contains plugins: ["ecpds"] and required: true. Startup rejects an ECPDS plugin reference if the binary lacks the ecpds feature or the schema has auth.required: false; those are not silent bypass settings.

Metrics: compare changes in aviso_ecpds_access_decisions_total by outcome, not just allow. If all aviso_ecpds_* series are missing, verify that metrics are enabled and the scrape reaches the right replica before checking whether the binary was built with --features ecpds.

Starting a watch or replay is slow

event_name=auth.ecpds.cache.miss means a request could not use a cached list. It may have fetched from ECPDS or waited for another request’s fetch. A miss alone does not prove ECPDS caused the delay.

  1. Check aviso_http_request_duration_seconds for route="/api/v1/watch" or route="/api/v1/replay". This measures time until response headers are ready, not how long the stream stays open.
  2. For a slow request, inspect cache_outcome on its auth.ecpds.check.* result event. hit means no upstream fetch was needed; miss_fetched means this request fetched; miss_coalesced means it shared another request’s fetch. Check auth.ecpds.fetch.failed for timeout or connection errors. Cache hit/miss and fetch success events require debug logging.
  3. Compare the rates of aviso_ecpds_cache_misses_total and aviso_ecpds_cache_hits_total on that replica. Check recent restarts or traffic moving between replicas before tuning the cache. Compare cache_ttl_seconds with how often users reconnect, and aviso_ecpds_cache_size with max_entries. The size gauge is an approximate count of cached usernames, not proof of eviction.

Use aviso_ecpds_fetch_total to check whether more misses also mean more fetches; shared fetches count once. If slow requests are cache hits, continue with the request’s other logs rather than assuming the cache needs resizing.

Tracing event reference

Every event uses the codebase’s standard structured shape (service_name, service_version, event_name, plus event-specific fields). The list below covers each event with a one-line meaning. Field-value details follow.

auth.ecpds.check.started Level: debug

The plugin started checking access for a request.

auth.ecpds.check.allowed Level: info

The plugin allowed the request.

auth.ecpds.check.denied Level: warn

The plugin denied the request. See reason field.

auth.ecpds.check.unavailable Level: warn

The plugin failed to reach a verdict. See fetch_outcome field.

auth.ecpds.check.error Level: error

An unexpected error in the plugin. See error_kind or error field.

auth.ecpds.admin.bypass Level: debug

An admin user skipped the ECPDS check. Demoted from info because admin bypass is configured behaviour, not an event SREs alert on; the aviso_ecpds_access_decisions_total{outcome="admin_bypass"} Prometheus counter still records every occurrence unconditionally.

auth.ecpds.cache.hit Level: debug

The destination list came from cache.

auth.ecpds.cache.miss Level: debug

The destination list was not in cache; a fetch was triggered.

auth.ecpds.fetch.succeeded Level: debug

A fetch to one ECPDS server succeeded.

auth.ecpds.fetch.failed Level: warn

A fetch to one ECPDS server failed. See error field.

auth.ecpds.fetch.skipped_inactive Level: debug

One or more ECPDS records returned by a single server had active != true (false, missing, or not a boolean) and got dropped from the user’s allow-list. Carries server_index, server, username, skipped, total. Demoted from info because every ECPDS fetch routinely returns inactive records and the skip behaviour is the documented contract; flip to debug only when investigating a denied user whose expected destination appears in this skip count.

auth.ecpds.fetch.skipped_record Level: debug

One or more ECPDS records returned by a single server were active but missing the configured target_field and got dropped. Carries server_index, server, username, target_field, skipped, total so on-call can pinpoint which ECPDS server is producing the malformed records. Demoted from info on the same grounds as skipped_inactive.

Common fields

Most events carry event_type (the schema name) and username (the JWT subject). Per-server events (auth.ecpds.fetch.succeeded, .failed, and .skipped_record) also carry server_index (zero-based) and server (the parsed URL).

Field value reference

Some events carry a typed enum field. The values you will see in logs are listed below. They are spelled exactly as shown.

  • reason (on auth.ecpds.check.denied):
    • DestinationNotInList: the user is not entitled to the requested destination.
    • MatchKeyMissing: the request body did not include the configured match-key field.
  • fetch_outcome (on auth.ecpds.check.unavailable):
    • Unauthorized, Forbidden: an ECPDS server returned 401 or 403.
    • ClientError: an ECPDS server returned a 4xx other than 401 or 403 (commonly 404 for a misconfigured base URL or 429 for throttling).
    • ServerError: an ECPDS server returned 5xx.
    • InvalidResponse: an ECPDS server returned a body the parser could not read.
    • Unreachable: network or timeout failure.
  • cache_outcome (on every auth.ecpds.check.* event: .allowed, .denied, .unavailable, .error):
    • hit: served from cache.
    • miss_coalesced: the cache was empty for this key but a concurrent caller’s fetch was in flight; this request waited on it.
    • miss_fetched: this request ran the upstream fetch itself. The merged per-server result of that fetch is recorded as the outcome label on the aviso_ecpds_fetch_total metric, and on the auth.ecpds.check.unavailable event also as fetch_outcome (see above). It is intentionally NOT inlined into cache_outcome so log filters keyed on cache_outcome:miss_fetched stay stable as new FetchOutcome variants are added.
    • none: cache lookup was deliberately skipped because the request fell at the MatchKeyMissing deny path before any cache call ran. Only appears on auth.ecpds.check.denied events alongside reason=MatchKeyMissing.

How to confirm “config error vs. upstream outage”

  1. Is the ECPDS plugin even compiled in? Check /metrics for aviso_ecpds_* series. The unlabelled counters and gauge plus the pre-initialised label values on aviso_ecpds_access_decisions_total and aviso_ecpds_fetch_total register at process startup whenever the binary is built with --features ecpds, regardless of whether an ecpds: config block exists. If the series are absent, the binary does not have the feature, OR the metrics endpoint itself is disabled (metrics.enabled: false in your config). If the series exist but aviso_ecpds_access_decisions_total{outcome="allow"} plus outcome="deny_*" are all flat at zero under load, the plugin is compiled in but no stream actually opts in via plugins: ["ecpds"].

  2. Are the configured server URLs reachable from this Aviso host? Run this from the same host as Aviso:

    curl -i -u "<service-username>:<service-password>" \
         "https://<your-ecpds-host>/ecpds/v1/destination/list?id=<some-test-username>"
    
    • 200 with a JSON destinationList: ECPDS is up and credentials are valid. Problem is on the Aviso side.
    • 401 or 403: service-account credentials are wrong (rotated, revoked, typoed).
    • 5xx or hang: ECPDS itself is broken.
    • DNS error or connection refused: network-level issue.
  3. Is one specific user being denied while others succeed? Run the curl above with that user’s id and compare with the destination they tried to read.

When an ECPDS server is unavailable

With the default partial_outage_policy: strict, every configured ECPDS server must respond successfully when Aviso fetches a user’s destination list. If one fails, requests needing a fresh list receive HTTP 503. Requests that can use a cached list can still be checked against that list.

With partial_outage_policy: any_success, Aviso can use the lists from the servers that respond successfully. A destination known only to a failed server may be missing, so a user could be denied access they would normally have. If no server succeeds, the lookup still fails with HTTP 503.

Both policies combine the destination lists returned by ECPDS. The choice is whether a lookup requires all servers or at least one. See Partial outage policy for more details.

Streaming Semantics

This page defines the exact behavior of the watch and replay endpoints, including start points, spatial filtering, identifier constraints, and SSE lifecycle events.


Request ID Correlation

Every HTTP response carries an X-Request-ID header with a per-request UUID. The same UUID is also embedded in the JSON data: payload of certain SSE events so that a client which only sees the body (not headers) can still quote it back when reporting a problem.

A note on terminology before the tables: SSE has two related but distinct labels for an event. The event: line is the SSE-level type that an EventSource.addEventListener(name, handler) call dispatches on. The data.type field (when present) is an aviso-level discriminator inside the JSON body, used for the cases where we reuse a single event: name for several control events. A wire example:

event: live-notification
data: {"type":"connection_established","topic":"...","timestamp":"...","connection_will_close_in_seconds":3600,"request_id":"<uuid>"}

The event: line above is live-notification (so an EventSource client listens for live-notification); the data.type value is connection_established (so the client distinguishes it from a normal notification, whose type is the CloudEvent type int.ecmwf.aviso.<event_type>).

In the tables below, Event is the SSE event: line and Type is the data.type field.

In-stream events that include request_id:

EventTypeSent
live-notificationconnection_establishedFirst event of a live-only watch
replay-controlreplay_startedFirst event of a stream that begins with replay
errornone; the error field identifies itOn a backend or CloudEvent-creation failure mid-stream
connection-closingnone; the reason field identifies itLast event, on a graceful close

In-stream events that intentionally do not include request_id:

EventTypeSent
live-notificationCloudEvent typeEvery live notification
replayCloudEvent typeEvery replayed notification
heartbeatnoneEvery few seconds
replay-controlreplay_completed, notification_replay_limit_reachedAt replay phase boundaries

The response header and the first event already carry the UUID, so repeating it would add bytes without adding information.

The first event of any stream is guaranteed to carry the request_id (a live-notification event with data.type = "connection_established" for live-only watches, or a replay-control event with data.type = "replay_started" for any stream that begins with replay).

Historical Replay Limits

watch_endpoint.max_historical_notifications limits historical notifications delivered by one request, not one backend batch. Identifier matching, identifier constraints, spatial filtering and successful CloudEvent rendering happen before quota accounting. Render failures emit errors without consuming the quota. A schema’s max_historical_notifications overrides the global cap; omitting it inherits the global default of 10000. replay_batch_size controls fetching independently of this quota. Excluded notifications do not consume it, even when entire batches are excluded.

After filling the quota, replay looks for one more renderable matching notification. If history is exhausted first, replay completes normally. Otherwise the server emits a replay-control event with these data fields:

{"type":"notification_replay_limit_reached","topic":"example.*","max_allowed":2,"message":"Historical replay limited to 2 messages. Additional historical messages may be available but were not retrieved.","timestamp":"2026-01-01T00:00:00Z"}

The limit control is followed by connection-closing with reason end_of_stream. There is no replay_completed event and a watch does not transition to live delivery. Treat this as a history gap, not successful catch-up. max_allowed is always present and reports the effective request cap, including a schema override. A new request gets a fresh quota.

Caps must be positive integers. Zero and unlimited are not accepted. If fetching a historical batch fails, including during lookahead, the server emits event: error and closes without replay_completed, a limit control or live delivery. This is failed catch-up, not exhausted history. Individual CloudEvent rendering failures remain nonfatal and do not consume the quota.

Live-only watches are not limited.

Replay/Live Boundary

A watch captures an inclusive history sequence bound H atomically with live subscription creation. Replay reads only sequences <= H, then emits replay_completed before delivering live sequences > H. Publications during catch-up belong to live delivery, not to later historical batches. No last-delivered watermark is used to discard queued live notifications.

A replay-only request captures H during setup without creating a live subscription. Its sequence range stays fixed across every page, including quota lookahead. Messages above H cannot consume the replay quota or prove truncation. A deleted or filtered message at H does not prevent completion.

An end point (to_id or to_date) lowers H for that request. See End Point for Historical Events.

This snapshots a sequence range, not immutable storage. Deletion, overwrite and retention can remove history before it is read. In-memory live queues can lag; JetStream retention can remove queued live messages before delivery. The boundary prevents replay/live overlap, but does not guarantee lossless delivery.

Sequence and date start cursors constrain history only. A start beyond H (or a future date) can produce empty replay, after which a watch still delivers new notifications from its subscription creation point. Setup failures return an HTTP error, not a successful replay completion.

Reconnecting after disconnect

If a stream drops (network blip, client restart, connection-closing with reason max_duration_reached, etc.), the recommended reconnect protocol is:

  1. Read the top-level CloudEvent id from each notification’s SSE data: body. It has the form <event_type>@<sequence>, such as extreme_event@123. The numeric suffix is the sequence; normal notifications do not have a separate sequence field. Control events, including connection_established, are not notification checkpoints.
  2. Save progress only after successfully processing the notification. If processing concurrently, do not advance the checkpoint past unfinished notifications.
  3. Issue a fresh POST /api/v1/watch (or /api/v1/replay) with the same event type and filters. Set from_id to the saved sequence plus 1, encoded as a decimal JSON string, and omit from_date. The client computes this increment; the server treats from_id as inclusive.

For example, after processing and saving a checkpoint for extreme_event@123, send a new POST /api/v1/watch request with this JSON body, keeping the same identifier filters as the original request:

{
  "event_type": "extreme_event",
  "identifier": {},
  "from_id": "124"
}

Here, identifier: {} represents a request without identifier filters. Replace it with your original filters; schemas with required filters do not accept an empty identifier object.

This avoids requesting the checkpoint notification again, but does not guarantee duplicate-free processing. A crash after applying an effect but before saving progress can cause that notification to be processed again. Use idempotent handlers or deduplicate notifications within the same backend history.

Recovery depends on the requested history still being available. Retention or deletion can remove notifications, and in-memory history is node-local and disappears on restart. A saved cursor is not portable across unrelated or reset backend histories. Reconnecting cannot recover history that is no longer stored and does not guarantee lossless delivery. Each reconnect is still subject to the historical replay limits above.

If you need time-based catch-up rather than sequence-based, use from_date instead of from_id (see Start Point for Historical Events below for accepted formats).

SSE Stream Lifecycle

Every streaming response (watch or replay) goes through a typed lifecycle:

stateDiagram-v2
    [*] --> Connected : SSE connection established
    Connected --> Replaying : from_id or from_date provided
    Connected --> Live : no replay parameters (watch only)
    Replaying --> Live : replay_completed (watch only)
    Replaying --> Closed : end_of_stream (replay only)
    Live --> Closed : max_duration_reached
    Live --> Closed : server_shutdown
    Closed --> [*]

Close reasons emitted in the final connection-closing SSE event:

ReasonTrigger
end_of_streamReplay finished, was truncated, or failed during batch retrieval
max_duration_reachedconnection_max_duration_sec elapsed on a watch stream
server_shutdownServer is shutting down gracefully

POST /api/v1/watch

  • If both from_id and from_date are omitted:
    • stream is live-only (new notifications from now onward).
  • If exactly one replay parameter is present:
    • historical replay starts first, then transitions to live stream.
  • If both are present:
    • request is rejected with 400.
  • to_id and to_date are not accepted; they are rejected with 400.
flowchart TD
    A["Watch Request"] --> B{"from_id or<br/>from_date?"}
    B -->|neither| C["Live-only stream"]
    B -->|exactly one| D["Historical replay<br/>then live stream"]
    B -->|both| E["400 Bad Request"]

    style E fill:#8b1a1a,color:#fff
    style C fill:#1a6b3a,color:#fff
    style D fill:#1a4d6b,color:#fff

POST /api/v1/replay

  • Requires exactly one replay start parameter:
    • from_id (sequence-based), or
    • from_date (time-based).
  • If both are missing or both are present:
    • request is rejected with 400.
  • Accepts at most one replay end parameter:
    • to_id (sequence-based), or
    • to_date (time-based).
  • Stream closes with end_of_stream when history is exhausted or the end point is reached.

Start Point for Historical Events

from_id is an unsigned 64-bit sequence number encoded as a JSON string. Historical delivery starts at that number (inclusive), subject to the request filters and available history. Send a decimal string, not a JSON number or the full CloudEvent id. The value "0" starts from the beginning of available matching history.

from_date accepts any of these formats:

FormatExample
RFC3339 with timezone2025-01-15T10:00:00Z
RFC3339 with offset2025-01-15T10:00:00+02:00
Space-separated with timezone2025-01-15 10:00:00+00:00
Naive datetime (interpreted as UTC)2025-01-15T10:00:00
Unix seconds (≤ 11 digits)1740509903
Unix milliseconds (≥ 12 digits)1740509903710

All inputs are normalized to UTC internally.


End Point for Historical Events

A replay can stop before the last stored notification. to_id and to_date mirror from_id and from_date, and are accepted only by POST /api/v1/replay.

to_id is an unsigned 64-bit sequence number encoded as a JSON string. Delivery ends at that number (inclusive). Unlike from_id, the value "0" has no special meaning: sequences start at 1, so it produces an empty replay. to_date accepts the same formats as from_date and ends delivery with the last notification stored at or before that time (inclusive). Like from_date, it refers to the time the server stored the notification, which is the CloudEvent time.

The replay ends at whichever comes first: the end point, or the last notification stored when the replay starts. An end point in the future therefore ends at the last stored notification, and an end point before every stored notification produces an empty replay that completes normally.

The end point bounds history before filtering, like H. The replay limit counts only notifications inside it, so a window with no more notifications than the limit is never truncated.

Requests are rejected with 400 when:

  • both to_id and to_date are present;
  • to_id is lower than from_id, or to_date is earlier than from_date.

A sequence end may follow a date start, and a date end may follow a sequence start. These are not compared, and may produce an empty replay.

The replay_started event reports the end point. end_sequence is the last sequence the replay can deliver, which may be lower than the requested to_id, and to_date repeats the requested time (truncated to whole seconds). Both are null without an end point.


Spatial Filter Model

Spatial filtering applies on top of the identifier field filters. Think of it in two layers:

  • Non-spatial identifier fields (time, date, class, etc.) narrow candidates by topic routing.
  • Spatial fields (polygon or point) further narrow that candidate set geographically.
flowchart LR
    A["Candidate<br/>notifications"] -->|"non-spatial<br/>identifier filter"| B["Topic-matched<br/>subset"]
    B -->|"spatial filter<br/>if provided"| C["Final<br/>results"]

Rules

identifier.polygonidentifier.pointResult
providedomittedpolygon-intersects-polygon filter
omittedprovidedpoint-inside-notification-polygon filter
omittedomittedno spatial filter
providedprovided400 Bad Request
  • identifier.polygon: keep notifications whose stored polygon intersects the request polygon.
  • identifier.point: keep notifications whose stored polygon contains the request point.
  • Both together: invalid (the request is rejected).

Point-cloud schemas use a different provider and subscriber pair. Providers send identifier.point_cloud on /notification. Subscribers send identifier.polygon on /watch or /replay. A notification matches when any cloud point is inside the request polygon or on an edge or vertex. Aviso checks the stored bounding box first, then stops at the first matching point.

The request polygon satisfies a required point_cloud field for watch and replay. Subscribers cannot send point_cloud. See Point-Cloud Filtering for the schema and request contract.


Identifier Constraints (watch / replay)

For schema-backed event types, identifier fields in watch/replay requests accept constraint objects instead of (or in addition to) scalar values. A scalar value is treated as an implicit eq constraint.

For schema-backed event types, unknown identifier names return HTTP 400 rather than being ignored. The supported spatial filters (point for polygon schemas and polygon for point-cloud schemas) remain valid.

Constraint objects are rejected on /notification: publish concrete identifier values in the shapes accepted by their handlers, including arrays for spatial identifiers, rather than watch/replay predicates.

Supported operators by field type

HandlerOperators
IntHandlereq, in, gt, gte, lt, lte, between
FloatHandlereq, in, gt, gte, lt, lte, between
EnumHandlereq, in

Notes

  • between expects exactly two values [min, max] and is inclusive on both ends.
  • Float constraints reject NaN and inf; only finite values are valid.
  • Float eq and in use exact numeric equality; no tolerance window is applied.
  • A constraint object must contain exactly one operator; combining operators in a single object is rejected.

Examples

{ "severity": { "gte": 5 } }
{ "severity": { "between": [3, 7] } }
{ "region":   { "in": ["north", "south"] } }
{ "anomaly":  { "lt": 50.0 } }

Backend Behavior

BackendHistorical replayLive watch
in_memoryNode-local only; clears on restartNode-local fan-out
jetstreamDurable; survives restartsCluster-wide fan-out

SSE Timestamp Format

All control, heartbeat, and close event timestamps use canonical UTC second precision:

YYYY-MM-DDTHH:MM:SSZ

Example: 2026-02-25T18:58:23Z


Replay Payload Shape

  • Replay and watch CloudEvent output always includes data.payload.
  • If a notify request omitted payload (optional schema), replay returns data.payload = null.
  • Payload values are not reshaped; scalar strings remain strings, objects remain objects.

See Payload Contract for the full input → storage → output mapping.


For end-to-end examples, see:

Payload Contract

This page defines the canonical payload behavior for Aviso notifications.

Scope

  • Applies to POST /api/v1/notification input.
  • Applies to stored backend payload representation.
  • Applies to replay/watch CloudEvent output payload field.

Schema Configuration

Per event schema, payload configuration is:

payload:
  required: true # or false

There is no payload.type list.

Canonical Rules

  1. Payload values are JSON values: object, array, string, number, boolean, null.
  2. If payload.required = true and request omits payload, request is rejected (400).
  3. If payload.required = false and request omits payload, Aviso stores canonical JSON null.
  4. Aviso does not wrap or reshape payload values (for example no auto-wrapping into {"data": ...}).

Input to Storage to Replay Mapping

Notify request payloadStored payloadReplay/Watch CloudEvent data.payload
omitted (optional schema)nullnull
"forecast complete""forecast complete""forecast complete"
424242
truetruetrue
["a","b"]["a","b"]["a","b"]
{"note":"ok"}{"note":"ok"}{"note":"ok"}

Failure Cases

  • Missing required payload:
    • HTTP 400
    • validation error (INVALID_NOTIFICATION_REQUEST)
  • Malformed JSON request body:
    • HTTP 400
    • parse error (INVALID_JSON)

Consumer Guidance

  • Treat data.payload as dynamic JSON.
  • If your client requires object-only payloads, normalize on the client side.
    • Example strategy: for non-object payloads, convert to {"data": <payload>} in your consumer.

Topic Encoding

Aviso uses a single backend-agnostic wire format for topics across all backends.


Why This Exists

NATS subject tokenization uses . as the separator. Wildcards (*, >) also operate on dot-delimited tokens. If topic field values contain any of these reserved characters, they would silently break routing and filtering.

Example of the problem:

Logical value:  1.45
Naive subject:  mars.od.1.45      (looks like 4 tokens, not 3)

To prevent this, Aviso percent-encodes each token value before assembling the wire subject.


Encoding Rules

Topic Bases

Configured topic.base values must match [A-Za-z0-9][A-Za-z0-9_-]*. The first character is an ASCII letter or digit. Remaining characters may also be underscores or hyphens. Empty bases, dots, percent signs, wildcards, spaces and Unicode are rejected at startup. Bases must be unique across schemas, ignoring ASCII case. Without a topic block, the event name is the base and must meet the same rules. Invalid generic request event names receive a validation 4xx response before storage or streaming starts.

JetStream uses the ASCII-uppercase base as its stream name and <base>.> as its subject binding, preserving the base’s case in subjects. For example, Weather_v2 binds Weather_v2.> to WEATHER_V2. Publish, replay, watch and admin operations target the same uppercase name. Admin cleanup also accepts legacy backend stream names, such as _WEATHER or WEATHER%2EV1, without applying the logical-base restriction. Existing streams are not renamed. This contract does not add support for bare subjects without an identifier token.

Identifier Values

The base restriction does not apply to identifier values. Decimal values such as 1.45 and strings such as a.b or a%2Eb retain the encoding below.

Only four characters are reserved and must be encoded:

CharacterEncoded formReason
.%2ENATS token separator
*%2ANATS single-token wildcard
>%3ENATS multi-token wildcard
%%25Escape character itself (keeps decoding unambiguous)

All other characters pass through unchanged.


Encode → Wire → Decode Flow

flowchart LR
    A["Logical value<br/>e.g. 1.45"] -->|encode| B["Wire token<br/>e.g. 1%2E45"]
    B -->|assemble| C["Wire subject<br/>e.g. extreme_event.north.1%2E45"]
    C -->|"split on '.'"| D["Wire tokens"]
    D -->|decode each| E["Logical tokens<br/>e.g. 1.45"]

    style A fill:#2a4a6b,color:#fff
    style C fill:#1a6b3a,color:#fff
    style E fill:#2a4a6b,color:#fff

The decoder is strict: malformed %HH sequences (e.g. %GG) are rejected, not passed through.


Examples

Encoding

Logical valueWire token
1.451%2E45
1*341%2A34
1>01%3E0
1%251%2525

Decoding (single pass)

Wire tokenLogical value
1%2E451.45
1%2A341*34
1%25251%25
1%251%

Impact on Wildcard Matching

Watch and replay requests use a two-step filter:

  1. Backend coarse filter: operates on wire subjects (NATS wildcard matching).
  2. App-level wildcard match: operates on decoded logical tokens.

Both steps are safe with reserved characters because the app layer always decodes before matching. Subscribers never need to think about encoding in their filter values; Aviso handles it transparently.


Invariants

  • Wire subject separator is always .
  • One shared codec is used for all backends (JetStream and In-Memory)
  • Encoding is applied per-token, not per-subject
  • Decoding is a strict single-pass operation

API Errors

Errors produced by aviso’s own handlers use a consistent JSON error object on 4xx and 5xx responses. A small number of 4xx responses are framework-level fallbacks rather than aviso-handled errors and use a different shape; see Framework-Level Fallbacks below.

Response Shape

{
  "code": "INVALID_REPLAY_REQUEST",
  "details": "Replay endpoint requires either from_id or from_date parameter...",
  "error": "Invalid Replay Request",
  "message": "Replay endpoint requires either from_id or from_date parameter...",
  "request_id": "0d4f6758-1ce3-4dda-a0f3-0ccf5fcb50d6"
}

Wire field order is alphabetical because serde_json::Map uses a BTreeMap. The examples below match what curl actually emits.

Fields:

  • code: stable machine-readable error code.
  • error: human-readable error category.
  • message: top-level failure message (safe for client display).
  • details: deepest/root detail available.
  • request_id: per-request UUID. The same value is also returned in the X-Request-ID HTTP response header and in every server-side log line for this request. Quoting it is the easiest way to ask the operator to look up the corresponding traces.

Notes:

  • error_chain is logged server-side for diagnostics, but is not returned in API responses.
  • message is always present.
  • details is present on 4xx and 5xx errors emitted from the notification/watch/replay request path (parse, validation, processing, and SSE stream initialization). It is intentionally omitted from authentication errors (401/403/503) and from the streaming-auth helpers (forbidden, ECPDS service unavailable, internal misconfiguration), where the upstream service or authorization plugin does not provide a stable detailed message.
  • request_id is present on every error response body produced by aviso’s own handlers (notification, watch, replay, schema, admin, auth and ECPDS authorization helpers). The X-Request-ID HTTP response header carries the same UUID on every response (success, aviso error, or framework-level fallback); see Framework-Level Fallbacks for the fallback cases where a JSON body field is not available.
  • SSE stream initialization failures additionally include topic (the decoded logical topic) so the operator can scope log queries faster.
  • The admin endpoints (/api/v1/admin/*) and POST /api/v1/notification return typed responses where request_id is a struct field rather than a free-form JSON key, but the field name and value are the same.

How to report a problem

Capture either of these and pass them to the operator:

  1. The X-Request-ID HTTP response header. Visible to curl -i, browser devtools, every reverse proxy, and most log aggregators. Present on every response, success or failure.
  2. The request_id field in any error response body. Same UUID as the header.

For streaming responses (/api/v1/watch, /api/v1/replay), the same UUID also appears in the JSON data: payload of the very first event and in any error or connection-closing event the stream emits before terminating. The exact wire shape depends on the stream variant:

  • For a live-only watch, the first event has SSE event: live-notification and a JSON body with "type": "connection_established".
  • For a stream that begins with replay, the first event has SSE event: replay-control and a JSON body with "type": "replay_started".

In both cases the request_id field is in the JSON body alongside type. This means a user who only sees the open-ended SSE body (no header parsing) can still recover the id without running the request again. See Streaming Semantics for the full event-by-event payload table, including the SSE event: versus data.type distinction (relevant for clients using EventSource.addEventListener).

Error Telemetry Events

These event_name values are emitted in structured logs:

Event NameLevelTrigger
api.request.parse.failedwarnJSON parse/shape/unknown-field failure before domain validation.
api.request.validation.failedwarnDomain/request validation failure (400).
api.request.processing.failederrorServer-side processing/storage failure (500).
stream.sse.initialization.failederrorReplay/watch SSE initialization failure (500).

Every event carries request_id. The formatter additionally propagates event_type and topic from the surrounding request span when the handler has recorded them on the span before emitting the error log. Whether they appear depends on which step failed:

  • api.request.parse.failed never carries event_type or topic. The request body is rejected before either is known.
  • api.request.validation.failed sometimes carries event_type. Validation steps that run after the handler has parsed the schema (notify-side process_notification_request failures) include it. Steps that run before, namely the watch/replay request validator and the notify-side endpoint-mismatch check, do not.
  • api.request.processing.failed carries event_type. Storage-write failures additionally carry topic.
  • stream.sse.initialization.failed carries both event_type and topic.

In all cases, filter on request_id first; treat event_type and topic as auxiliary filters where present.

Error Code Reference

CodeHTTP StatusMeaning
INVALID_JSON400Request body is not valid JSON.
UNKNOWN_FIELD400Request contains fields outside API contract.
INVALID_REQUEST_SHAPE400JSON structure cannot be deserialized into request model.
INVALID_NOTIFICATION_REQUEST400Notification request failed business validation.
INVALID_WATCH_REQUEST400Watch request failed validation (replay/spatial/schema rules).
INVALID_REPLAY_REQUEST400Replay request failed validation (start or end cursor/spatial/schema rules).
UNAUTHORIZED401Missing or invalid credentials (no token, bad format, expired, bad signature).
FORBIDDEN403Valid credentials but user lacks the required role.
NOTIFICATION_PROCESSING_FAILED500Notification processing pipeline failed before storage.
NOTIFICATION_STORAGE_FAILED500Backend write operation failed.
SSE_STREAM_INITIALIZATION_FAILED500Replay/watch SSE stream could not be created.
INTERNAL_ERROR500Fallback internal error code (reserved).
SERVICE_UNAVAILABLE503Auth service (auth-o-tron) unreachable or returned an unexpected error.

Examples

Invalid replay request:

{
  "code": "INVALID_REPLAY_REQUEST",
  "details": "Cannot specify both from_id and from_date...",
  "error": "Invalid Replay Request",
  "message": "Cannot specify both from_id and from_date...",
  "request_id": "0d4f6758-1ce3-4dda-a0f3-0ccf5fcb50d6"
}

Auth error (missing credentials on a protected stream). Auth errors use four fields (code, error, message, request_id); details is not included:

{
  "code": "UNAUTHORIZED",
  "error": "unauthorized",
  "message": "Authentication is required for this stream",
  "request_id": "0d4f6758-1ce3-4dda-a0f3-0ccf5fcb50d6"
}

SSE initialization failure:

{
  "code": "SSE_STREAM_INITIALIZATION_FAILED",
  "details": "nats connect failed: timeout",
  "error": "SSE stream creation failed",
  "message": "Failed to create stream consumer",
  "request_id": "0d4f6758-1ce3-4dda-a0f3-0ccf5fcb50d6",
  "topic": "test_polygon.*.1200"
}

Framework-Level Fallbacks

A few 4xx responses are produced by the Actix HTTP framework itself before any aviso handler runs, so the body shape is whatever the framework defaults to (typically text/plain or empty) and not the JSON object documented above:

StatusTrigger
404 Not FoundRequest path does not match any registered aviso route.
405 Method Not AllowedPath matches an aviso route but the HTTP method does not.
400 Bad Request (rare)Request fails framework-level checks (malformed Content-Length, etc.) before reaching aviso’s body parsers.

For these cases:

  • code, error, message, details, and request_id JSON fields are not in the body.
  • The X-Request-ID HTTP response header is still set on 404 and 405 (the middleware stack runs before Actix’s default route-mismatch responses), so an operator can correlate the request with server logs by header alone. Errors raised even earlier in the HTTP stack (e.g., malformed Content-Length, TLS handshake failure) bypass aviso’s middleware entirely and produce no X-Request-ID; in those cases the request never reaches aviso, no log line is generated, and no correlation is possible.
  • The aviso code reference table above does not apply.

If you need a stable JSON shape on these paths, hit a known-good route (GET /health is the simplest); a 4xx from there indicates an actual aviso-handled error and follows the documented contract.

Admin Operations

Admin endpoints are destructive. Restrict access in production.

When authentication is enabled, all admin endpoints require a valid credential and one of the configured admin_roles. Add -H "Authorization: Bearer <token>" or -u user:pass (direct mode) to the curl examples below.

Delete One Notification

DELETE /api/v1/admin/notification/{notification_id}

notification_id format:

  • <stream>@<sequence> (canonical)
  • <event_type>@<sequence> (alias; resolved through configured schema topic.base)

How notification_id maps to your schema

Delete IDs use the stream key plus backend sequence number.

If your schema contains:

notification_schema:
  mars:
    topic:
      base: "mars"
  dissemination:
    topic:
      base: "diss"
  test_polygon:
    topic:
      base: "polygon"

Then valid delete IDs include:

  • mars@42 (event type alias and stream key are same)
  • dissemination@42 (alias form, resolved to stream key diss)
  • diss@42 (canonical stream key form)
  • test_polygon@306 (alias form, resolved to stream key polygon)
  • polygon@306 (canonical stream key form)

Example: replay ID then delete

Replay returns CloudEvent IDs like mars@1:

curl -N -X POST "http://127.0.0.1:8000/api/v1/replay" \
  -H "Content-Type: application/json" \
  -d '{
    "event_type": "mars",
    "identifier": {
      "class": "od",
      "expver": "0001",
      "domain": "g",
      "date": "20250706",
      "time": "1200",
      "stream": "enfo",
      "step": "1"
    },
    "from_id": "1"
  }'

If one replayed event has "id":"mars@1", delete it with:

curl -X DELETE "http://127.0.0.1:8000/api/v1/admin/notification/mars@1"

Response behavior

  • 200: notification deleted.
  • 404: stream/sequence pair not found.
  • 400: invalid ID format (<name>@<positive-integer> required).

The response body has the shape:

{
  "success": true,
  "message": "Notification deleted",
  "notification_id": "mars@42",
  "request_id": "<uuid>"
}

Failure responses (404, 400) keep the same fields with success: false and a descriptive message.

Invalid examples:

  • mars (missing @sequence)
  • mars@0 (sequence must be > 0)
  • mars@abc (sequence must be an integer)

Wipe Endpoints

  • DELETE /api/v1/admin/wipe/stream
  • DELETE /api/v1/admin/wipe/all

These endpoints remove many messages at once and should be used with extreme caution.

Wipe One Stream

DELETE /api/v1/admin/wipe/stream

Request body:

{
  "stream_name": "mars"
}

stream_name accepts the event type as configured in the notification schema or the backend stream name, case-insensitively: mars and MARS wipe the same stream.

Example:

curl -X DELETE "http://127.0.0.1:8000/api/v1/admin/wipe/stream" \
  -H "Content-Type: application/json" \
  -d '{"stream_name":"mars"}'

What it does:

  • Removes all stored messages for the selected stream.
  • Keeps the stream definition/configuration in place.
  • New notifications can still be written to that stream immediately after wipe.

When no stream matches the name, the response is 404 and the message lists the configured event types, so a typo is distinguishable from a backend failure (500).

When to use:

  • You want to reset one event family (mars, diss, polygon) without affecting others.

Wipe All Streams

DELETE /api/v1/admin/wipe/all

Example:

curl -X DELETE "http://127.0.0.1:8000/api/v1/admin/wipe/all"

What it does:

  • Removes all stored messages from all streams managed by the backend.
  • Leaves service configuration intact, but data history is gone.

When to use:

  • Local/dev reset before a fresh test run.
  • Operational emergency cleanup where full history removal is intended.

Which Admin Operation Should I Use?

  • Delete one notification (/admin/notification/{id}):
    • Use when you know the exact sequence to remove.
  • Wipe one stream (/admin/wipe/stream):
    • Use when one stream is polluted and others must remain untouched.
  • Wipe all (/admin/wipe/all):
    • Use only when complete history reset is intended.

Wipe Response Shape

Both wipe endpoints return the same field set: success, message, request_id. The message value differs:

DELETE /api/v1/admin/wipe/stream:

{
  "success": true,
  "message": "Successfully wiped stream: MARS",
  "request_id": "<uuid>"
}

DELETE /api/v1/admin/wipe/all:

{
  "success": true,
  "message": "Successfully wiped all data",
  "request_id": "<uuid>"
}

Failure responses keep the same fields with success: false and a descriptive message.

Architecture

Aviso Server is built around three operations (Notify, Watch, and Replay) that share a common validation and schema layer but diverge at the backend interaction.


System Overview

graph TB
    subgraph Clients
        P(Publisher)
        W(Watcher)
        R(Replayer)
    end

    subgraph "Aviso Server"
        direction TB
        AM["Auth Middleware<br/>(optional)"]
        RT["Routes<br/>HTTP handlers"]
        VP["Validation &<br/>Processing"]
        NC["Notification<br/>Core"]
        BE["Backend<br/>Abstraction"]
    end

    AOT["auth-o-tron<br/>(external)"]

    subgraph Backend
        JS[("JetStream<br/>NATS")]
        IM[("In-Memory<br/>Process")]
    end

    P -->|POST /api/v1/notification| AM
    W -->|POST /api/v1/watch| AM
    R -->|POST /api/v1/replay| AM

    AM -.->|verify credentials| AOT
    AM --> RT
    RT --> VP
    VP --> NC
    NC --> BE
    BE --> JS
    BE --> IM

    JS -.->|SSE stream| W
    JS -.->|SSE stream| R
    IM -.->|SSE stream| W
    IM -.->|SSE stream| R

Notify Request Flow

When a publisher sends POST /api/v1/notification:

sequenceDiagram
    participant C as Publisher
    participant A as Auth Middleware
    participant R as Route Handler
    participant V as Validator
    participant P as Processor
    participant T as Topic Builder
    participant B as Backend

    C->>A: POST /api/v1/notification (JSON)
    alt stream requires auth
        A->>A: resolve user (JWT or auth-o-tron)
        A-->>C: 401/403 if unauthorized
    end
    A->>R: forward request (+ user identity)
    R->>V: parse & shape-check JSON
    V-->>R: 400 if malformed
    R->>P: process_notification_request()
    P->>P: look up event schema
    P->>P: validate each identifier field
    P->>P: canonicalize values (dates, enums)
    P->>T: build_topic_with_schema()
    T-->>P: topic string (e.g. mars.od.0001.g.20250706.1200)
    P->>B: put_message_with_headers()
    B-->>C: 200 { status, request_id, processed_at }

Key steps:

  1. Parse: raw JSON bytes are deserialized; unknown fields are rejected (UNKNOWN_FIELD).
  2. Validate: each identifier field is checked against its ValidationRules (type, range, enum values).
  3. Canonicalize: values are normalized (for example dates to YYYYMMDD, enums to lowercase).
  4. Build topic: fields are ordered per key_order, reserved chars are percent-encoded.
  5. Store: the message is written to the backend with the encoded topic as the subject.

Watch Request Flow

POST /api/v1/watch opens a persistent SSE stream. It optionally starts with a historical replay phase before transitioning to live delivery.

sequenceDiagram
    participant C as Subscriber
    participant A as Auth Middleware
    participant R as Route Handler
    participant P as Stream Processor
    participant F as Hybrid Filter
    participant B as Backend

    C->>A: POST /api/v1/watch (JSON)
    alt stream requires auth
        A->>A: resolve user (JWT or auth-o-tron)
        A-->>C: 401/403 if unauthorized
    end
    A->>R: forward request (+ user identity)
    R->>P: process_request (ValidationConfig::for_watch)
    P->>P: allow optional fields & constraint objects
    P->>P: analyze_watch_pattern() → coarse + precise patterns

    alt has from_id or from_date
        P->>B: fetch historical batch
        B-->>P: NotificationMessage[]
        P->>F: apply wildcard + constraint + spatial filter
        F-->>C: SSE: replay_started → events → replay_completed
    end

    P->>B: subscribe(coarse_pattern)
    loop live stream
        B-->>P: live NotificationMessage
        P->>F: apply precise filter
        F-->>C: SSE: notification event
    end

    C-->>R: disconnect / timeout
    R-->>C: SSE: connection-closing

Replay Request Flow

POST /api/v1/replay is like watch but historical-only; the stream closes when history ends, or at the optional end point (to_id or to_date).

sequenceDiagram
    participant C as Client
    participant A as Auth Middleware
    participant R as Route Handler
    participant P as Stream Processor
    participant B as Backend

    C->>A: POST /api/v1/replay (JSON + from_id or from_date, optional to_id or to_date)
    alt stream requires auth
        A->>A: resolve user (JWT or auth-o-tron)
        A-->>C: 401/403 if unauthorized
    end
    A->>R: forward request (+ user identity)
    R->>P: process_request (ValidationConfig::for_replay)
    P->>B: resolve end sequence (history end, lowered by EndAt)
    P->>B: batch fetch from StartAt::Sequence or StartAt::Date
    loop batches
        B-->>P: NotificationMessage[]
        P->>P: filter + convert to CloudEvent
        P-->>C: SSE: notification events
    end
    P-->>C: SSE: replay_completed → connection-closing (end_of_stream)

SSE Streaming Pipeline

The streaming layer (src/sse/) is built around typed values rather than raw strings, which keeps the lifecycle explicit and the endpoint logic thin.

Cursor types describe how a start point is represented internally:

  • StartAt::LiveOnly: no history, subscribe immediately.
  • StartAt::Sequence(u64): start from a specific backend sequence number.
  • StartAt::Date(DateTime<Utc>): start from a UTC timestamp.

End types describe where a replay stops:

  • EndAt::Latest: at the last notification stored when the replay starts.
  • EndAt::Sequence(u64): at a backend sequence number, inclusive.
  • EndAt::Date(DateTime<Utc>): at the last notification stored at or before a UTC timestamp.

Frame types are what the stream produces before rendering to SSE wire format:

  • Control frames: connection_established, replay_started, replay_completed, notification_replay_limit_reached.
  • Notification frames: a decoded CloudEvent ready for delivery.
  • Heartbeat frames: periodic keep-alive.
  • Error frames: non-fatal stream errors.
  • Close frame: carries one of end_of_stream, max_duration_reached, server_shutdown.

Lifecycle (shutdown token, max duration, natural end) is applied once in a shared wrapper, so individual endpoint handlers don’t need to reimplement it.


Component Map

ComponentPathRole
Routessrc/routes/Thin HTTP handlers: parse request, delegate, return response
Authsrc/auth/Middleware, JWT validation, role matching, auth-o-tron client
Handlerssrc/handlers/Shared parsing, validation, and processing logic
Notification coresrc/notification/Schema registry, topic builder/codec/parser, wildcard matcher, spatial
Backend abstractionsrc/notification_backend/NotificationBackend trait + JetStream and InMemory implementations
SSE layersrc/sse/Stream composition, typed frames, heartbeats, lifecycle
CloudEventssrc/cloudevents/Converts stored messages into CloudEvent envelope
Configurationsrc/configuration/Config loading, schema validation, global snapshot
Error modelsrc/error.rsStable HTTP error codes and structured responses

Hybrid Filtering

Watch subscriptions use a two-tier strategy to balance backend load with filter precision:

graph LR
    A[Watch Request] --> B[analyze_watch_pattern]
    B --> C["Coarse pattern<br/>e.g. mars.*.*.*"]
    B --> D["Precise pattern<br/>full decoded topic"]
    C -->|backend subscription| E[(NATS JetStream)]
    E -->|candidate messages| F[App-level filter]
    D --> F
    F -->|matched| G[SSE client]
    F -->|rejected| H[dropped]
  • The coarse pattern is sent to the backend as the subscription subject filter. It uses NATS wildcards and covers a superset of the desired messages.
  • The precise pattern is applied in-process on decoded topics + constraint objects + spatial checks. Only messages that pass both layers reach the client.

This avoids creating one backend subscription per unique topic while still delivering exact results.


JetStream Backend Internals

ModulePathResponsibility
Confignotification_backend/jetstream/config.rsDecode and validate JetStream settings
Connectionnotification_backend/jetstream/connection.rsNATS connect with retry
Streamsnotification_backend/jetstream/streams.rsCreate and reconcile streams
Publishernotification_backend/jetstream/publisher.rsPublish with retry on transient failures
Subscribernotification_backend/jetstream/subscriber.rsConsumer-based live subscriptions
Replaynotification_backend/jetstream/replay.rsPull consumer batch retrieval
Adminnotification_backend/jetstream/admin.rsWipe and delete operations

Backend Development Guide

This guide explains how to add a new notification backend in Aviso.

Required Contract

A backend must implement NotificationBackend in src/notification_backend/mod.rs.

Core requirements:

  • Implement publish/replay/subscribe/admin methods.
  • Implement capabilities() and return a stable BackendCapabilities map.
  • Keep startup/shutdown behavior explicit and logged.

subscribe_to_topic returns Subscription { stream, history_end }. Capture the receiver and inclusive history bound at the same logical creation point as publishing. The returned live stream must start strictly above that bound. Never combine a subscription with a later, independently sampled stream tail. history_end(topic) captures the replay-only bound without a live subscription.

first_sequence_after(topic, at) returns a sequence that bounds a replay ending at at: a message of topic has a lower sequence exactly when it was stored at or before at. Return None when nothing was stored after at. A replay with to_date ends one sequence before the returned value. The first message stored after at on the topic’s backend subject is a valid answer, even when it belongs to another topic. Use the storage time that becomes the CloudEvent time.

get_messages_batch must enforce BatchParams.end_sequence before filtering, rendering or pagination. Preserve that bound when advancing the start cursor. Completion must not depend on finding a message exactly at the bound: it may be deleted, overwritten or excluded. The sequence range is finite, but storage contents can still change during replay.

Storage Policy Compatibility

Per-schema storage policy is validated at startup before backend initialization.

Validation entry point:

  • configuration::validate_schema_storage_policy_support(...)

Capability source:

  • notification_backend::capabilities_for_backend_kind(...)

If a schema requests unsupported fields, startup fails fast with a clear error. Do not silently ignore unsupported storage policy fields.

Capability Checklist

When adding backend <new_backend>:

  1. Add backend kind support in capabilities_for_backend_kind.
  2. Add NotificationBackend::capabilities() implementation.
  3. Ensure capability values match real backend behavior.
  4. Add tests for:
    • capability map values
    • accepted storage-policy fields
    • rejected storage-policy fields
  5. Add backend docs page and update summary links.

Minimal Capability Example

BackendCapabilities {
    retention_time: true,
    max_messages: true,
    max_size: false,
    allow_duplicates: false,
    compression: false,
}

Meaning:

  • retention_time and max_messages can be used in schema storage policy.
  • max_size, allow_duplicates, and compression must be rejected at startup.

Testing Expectations

  • Unit tests should verify capability flags are stable.
  • Validation tests should verify fail-fast messages for unsupported fields.
  • Integration tests should use test-local config/schema fixtures, not developer-local YAML files.

Run bash scripts/test_replay_boundary.sh for the full test suite with live JetStream coverage on an isolated NATS 2.14.6 container. Add --features ecpds for that build. The script allocates a loopback port and removes its own server and disposable storage on exit; it does not use an existing NATS instance.