Skip to content

Commit 9701ef6

Browse files
authored
✨ qontract-api: events framework (#5420)
* ✨ qontract-api: events framework * replace DIY event framework with faststream * feat(01-foundation-configuration-02): extend SlackSettings with notification fields - Add notification_channel field (str | None, default None) - Add notification_workspace field (str | None, default None) - Add secret_path field (str, default 'app-sre/slack/bot-token') - Add 5 comprehensive tests for new fields - Maintain backwards compatibility with existing deployments * test(01-01): add failing tests for chat_post_message - Add ChatPostMessageResponse model to models.py with ts, channel, thread_ts fields - Add 7 test functions covering success, thread reply, auto-join, error handling, truncation, and hooks - All tests fail as expected (RED state) because chat_post_message method does not exist yet * feat(01-01): implement chat_post_message - Add chat_post_message method to SlackApi with @invoke_with_hooks decorator - Auto-join logic handles not_in_channel error with conversations_join + retry - channel_not_found logged at ERROR level and re-raised - Other errors logged at WARNING level and re-raised - Text truncation at 10,000 characters with "... [truncated]" suffix - Export ChatPostMessageResponse from __init__.py - All tests pass (7 new + 23 existing) * refactor(01-01): clean up chat_post_message - Extract magic value 10000 to MAX_MESSAGE_LENGTH constant - Use logger.exception instead of logger.error when re-raising - Add type checking for response.data to satisfy mypy - Fix exception chaining with 'from None' - Update test to check logger.exception instead of logger.error - All tests pass, ruff check/format pass, mypy passes * feat(02-01): create shared slack module with workspace client and factory function - Relocate SlackWorkspaceClient to qontract_api/slack/ - Add chat_post_message method to SlackWorkspaceClient - Create create_slack_workspace_client factory function following PagerDuty pattern - Factory resolves Secret via SecretManager * feat(02-01): migrate imports and add tests for shared slack module - Convert old modules to re-export shims for backward compatibility - Update service.py to use create_slack_workspace_client function - Update tasks.py to pass cache directly to service - Migrate test imports to new location - Add new tests for factory function and chat_post_message - Use lazy imports to avoid circular dependencies - All 43 tests passing * feat(03-01): add POST /api/v1/external/slack/chat endpoint - Create frozen Pydantic request/response models with Secret, icon, and username fields - Implement POST /chat router with ValueError (404) and SlackApiError (502) handling - Add channel name → ID resolution with # prefix stripping in SlackWorkspaceClient - Rename SlackApi.chat_post_message channel → channel_id for type clarity - Add icon_emoji, icon_url, username params through the full call chain - Replace msg_kwargs dict with explicit typed parameters for mypy safety - Improve general_exception_handler to preserve HTTPException status/detail - Register slack router in api_v1 - Add 5 endpoint tests + 2 workspace client tests (channel resolution, hash prefix) * fixup! feat(02-01): create shared slack module with workspace client and factory function * feat(04-01): add SubscriberSettings to config and clean up old SlackSettings fields - Add SubscriberSettings model with required slack_channel, slack_workspace, slack_token_path - Add optional qontract_api_url (defaults to http://localhost:8000) and qontract_api_token_path - Add Settings.subscriber field (defaults to None for backward compatibility) - Remove deprecated notification_channel, notification_workspace, secret_path from SlackSettings - Add comprehensive tests for SubscriberSettings and Settings.subscriber * feat(04-01): create event formatter registry with generic formatter, metrics, and tests - Add EventFormatter protocol for type-safe formatter registration - Implement GenericEventFormatter with emoji mapping (error, fail, create, update, delete) - Format events as human-readable Slack messages with emoji, type, source, and JSON data dump - Create formatter registry with format_event and register_formatter functions - Add Prometheus metrics: events_received, events_posted, events_failed counters and event_processing_duration histogram - Add comprehensive tests for formatter registry and generic formatter with complex data - Fix ClassVar annotation for EMOJI_MAP to satisfy ruff linter * feat(04-02): create qontract-api client and subscriber handler - Add qontract_api_token field to SubscriberSettings for direct token auth - Create _client.py with _get_client() and post_to_slack() functions - Rewrite _subscriptions.py event_handler with per-event error isolation - Increment Prometheus metrics for received, posted, failed events - Record event processing duration in histogram * test(04-02): add comprehensive subscriber handler tests - Create conftest.py with sample_event and error_event fixtures - Add test_subscriptions.py with 8 test cases covering: * Event processing flow (format + post) * Per-event error isolation (SUB-02) * Prometheus metrics (received, posted, failed, duration) * Exception handling in both format and post stages - All 16 subscriber tests pass * subscriber: slack username and emoji * delete unused event_sink integration * subscriber: prometheus * subscriber: logging * subscriber: openshift template * fixup! delete unused event_sink integration * fixup! subscriber: logging * fix tests
1 parent 2e11083 commit 9701ef6

72 files changed

Lines changed: 3450 additions & 252 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.
Lines changed: 199 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,199 @@
1+
# ADR-018: Event-Driven Communication Pattern
2+
3+
**Status:** Accepted
4+
**Date:** 2026-02-10
5+
**Authors:** cassing
6+
7+
## Context
8+
9+
qontract-api performs actions on external systems (e.g., updating Slack usergroup members) that other systems need to be notified about. Currently there is no mechanism for cross-system event notification between qontract-api and reconcile integrations.
10+
11+
- qontract-api celery tasks execute reconciliation actions but results are only stored as task results
12+
- There is no way for reconcile integrations to react to actions performed by qontract-api
13+
- Future use cases require audit logging, notification chains, and event-driven workflows
14+
- The solution must support multiple consumers and be extensible to new event types
15+
16+
## Decision
17+
18+
Use [faststream](https://faststream.ag2.ai/) with [CloudEvents](https://cloudevents.io/) for event-driven communication between qontract-api (producer) and subscriber processes (consumers). The shared event model and synchronous publishing wrapper live in `qontract_utils`. The subscriber runs as a separate ASGI process via faststream's `AsgiFastStream`.
19+
20+
### Key Points
21+
22+
- **CloudEvents standard** -- the `Event` model extends `cloudevents.pydantic.v2.event.CloudEvent`, providing an industry-standard, self-describing event envelope (`specversion`, `type`, `source`, `id`, `time`, `datacontenttype`, `data`)
23+
- **faststream as messaging framework** -- handles Redis broker connections, message serialization/deserialization, subscriber registration via decorators, and AsyncAPI specification generation
24+
- **Synchronous publisher** (`qontract_utils.events.RedisBroker`) -- wraps faststream's async `RedisBroker` in a sync context manager so that celery workers (sync) can publish events without async overhead
25+
- **Subscriber as ASGI app** (`qontract_api.subscriber`) -- runs as a separate uvicorn process with health check endpoints and auto-generated AsyncAPI documentation
26+
- **`EventManager`** (`qontract_api.event_manager`) -- encapsulates publishing with fire-and-forget semantics: failures are logged but never propagate to the producer task
27+
- **Feature-flagged** -- event publishing is controlled via `QAPI_EVENTS__ENABLED` (default: `true`) and the channel name via `QAPI_EVENTS__CHANNEL` (default: `"main"`)
28+
29+
## Alternatives Considered
30+
31+
### Alternative 1: DIY Redis Streams with Protocol Abstraction
32+
33+
Custom `EventPublisher`/`EventConsumer` protocols with factory functions, Redis Streams backend (XADD/XREADGROUP/XACK), and manual consumer group management.
34+
35+
**Pros:**
36+
37+
- Full control over implementation details
38+
- No external framework dependency
39+
40+
**Cons:**
41+
42+
- Significant boilerplate: custom protocols, factories, consumer group management, serialization
43+
- No standard event format -- custom schema requires documentation and versioning
44+
- No auto-generated API documentation
45+
- Consumer lifecycle (pending messages, acknowledgment, group creation) must be hand-coded
46+
47+
### Alternative 2: AWS SNS + SQS
48+
49+
SNS for publishing (fan-out), SQS for consuming (durable queues).
50+
51+
**Pros:**
52+
53+
- Native fan-out: multiple SQS queues subscribe to one topic
54+
- Managed infrastructure with high availability
55+
56+
**Cons:**
57+
58+
- App-interface (external-resources) doesn't support SNS with SQS subscriptions
59+
- SNS wraps messages in an envelope that must be unwrapped
60+
- AWS credentials must be managed separately
61+
62+
### Alternative 3: faststream with CloudEvents (Selected)
63+
64+
faststream handles the messaging infrastructure (broker connections, serialization, subscriber lifecycle). CloudEvents provides the standardized event envelope.
65+
66+
**Pros:**
67+
68+
- Minimal boilerplate: decorator-based subscriber registration, automatic serialization
69+
- Industry-standard event format (CloudEvents 1.0) -- no custom schema to maintain
70+
- Auto-generated AsyncAPI documentation at `/docs/asyncapi`
71+
- Built-in health check endpoints for Kubernetes readiness/liveness
72+
- Sync publisher wrapper is thin (~40 lines) -- all heavy lifting delegated to faststream
73+
- Uses existing Redis infrastructure -- no additional setup needed
74+
75+
**Cons:**
76+
77+
- External framework dependency (faststream)
78+
- Redis is a single point of failure (unlike managed SNS/SQS)
79+
80+
## Consequences
81+
82+
### Positive
83+
84+
- No additional infrastructure required -- uses existing Redis
85+
- Decoupled communication between qontract-api and subscribers
86+
- CloudEvents standard enables interoperability with external systems and tooling
87+
- AsyncAPI documentation auto-generated for service discovery
88+
- Subscriber runs as independent process -- can be scaled separately
89+
- Minimal code to maintain: the sync publisher wrapper and event model are ~50 lines total
90+
91+
### Negative
92+
93+
- Redis as single point of failure
94+
- **Mitigation:** Redis is already critical infrastructure for caching; adding events doesn't change the risk profile
95+
- External dependency on faststream
96+
- **Mitigation:** The sync publisher wrapper isolates the dependency; only the subscriber directly uses faststream decorators
97+
98+
## Implementation Guidelines
99+
100+
### Event Model
101+
102+
```python
103+
from qontract_utils.events import Event
104+
105+
event = Event(
106+
source="qontract-api",
107+
type="qontract-api.slack-usergroups.update_users",
108+
data={"workspace": "coreos", "usergroup": "team-a", "users": ["alice"]},
109+
datacontenttype="application/json",
110+
)
111+
```
112+
113+
### Publishing Events (Producer -- Sync)
114+
115+
Use `EventManager` in qontract-api celery tasks:
116+
117+
```python
118+
from qontract_api.event_manager import get_event_manager
119+
from qontract_utils.events import Event
120+
121+
event_manager = get_event_manager()
122+
if event_manager:
123+
event_manager.publish_event(
124+
Event(
125+
source=__name__,
126+
type=f"qontract-api.slack-usergroups.{action.action_type}",
127+
data=action.model_dump(mode="json"),
128+
datacontenttype="application/json",
129+
)
130+
)
131+
```
132+
133+
For direct publishing (e.g., scripts or tests):
134+
135+
```python
136+
from qontract_utils.events import Event, RedisBroker
137+
138+
with RedisBroker("redis://localhost:6379") as broker:
139+
broker.publish(
140+
Event(
141+
source="my-script",
142+
type="test.ping",
143+
data={"hello": "world"},
144+
datacontenttype="application/json",
145+
),
146+
channel="main",
147+
)
148+
```
149+
150+
### Subscribing to Events (Consumer -- Async)
151+
152+
Add handlers in `qontract_api/qontract_api/subscriber/_subscriptions.py`:
153+
154+
```python
155+
from qontract_utils.events import Event
156+
157+
from ._base import broker
158+
159+
160+
@broker.subscriber("main")
161+
async def base_handler(event: Event) -> None:
162+
print(event)
163+
```
164+
165+
The subscriber runs as a separate process:
166+
167+
```bash
168+
QAPI_START_MODE=subscriber uvicorn qontract_api.subscriber:app
169+
```
170+
171+
### Configuration (qontract-api)
172+
173+
Environment variables (prefix `QAPI_`):
174+
175+
| Variable | Default | Description |
176+
|---|---|---|
177+
| `QAPI_EVENTS__ENABLED` | `true` | Enable/disable event publishing |
178+
| `QAPI_EVENTS__CHANNEL` | `"main"` | Redis channel name for events |
179+
| `QAPI_START_MODE` | `"api"` | Process mode: `api`, `worker`, or `subscriber` |
180+
181+
The Redis connection is reused from the existing cache backend (`cache_broker_url`).
182+
183+
### Checklist
184+
185+
- [ ] Event publishing is behind a feature flag (`QAPI_EVENTS__ENABLED`)
186+
- [ ] Publishing failures do not break the producer (`EventManager` catches all exceptions)
187+
- [ ] Events use CloudEvents format with `source`, `type`, `data`, and `datacontenttype`
188+
- [ ] Event types use dot-separated naming (`qontract-api.integration.action`)
189+
- [ ] New subscribers are registered in `qontract_api/subscriber/_subscriptions.py`
190+
191+
## References
192+
193+
- Related ADRs: ADR-011 (Dependency Injection), ADR-012 (Typed Models), ADR-014 (Three-Layer Architecture)
194+
- Event model and sync broker: `qontract_utils/qontract_utils/events/`
195+
- Subscriber ASGI app: `qontract_api/qontract_api/subscriber/`
196+
- Event manager: `qontract_api/qontract_api/event_manager/`
197+
- [faststream documentation](https://faststream.ag2.ai/)
198+
- [CloudEvents specification](https://cloudevents.io/)
199+
- [AsyncAPI specification](https://www.asyncapi.com/)

‎docs/adr/README.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ An Architecture Decision Record (ADR) documents an important architecture decisi
2727
| [ADR-015](ADR-015-cache-update-strategy.md) | Cache Update Instead of Invalidation | Accepted |
2828
| [ADR-016](ADR-016-two-tier-cache.md) | Two-Tier Cache Architecture (Memory + Redis) | Accepted |
2929
| [ADR-017](ADR-017-factory-pattern.md) | Factory Pattern | Accepted |
30+
| [ADR-018](ADR-018-event-driven-communication.md) | Event-Driven Communication Pattern | Accepted |
3031

3132
## ADR Categories
3233

@@ -54,6 +55,7 @@ An Architecture Decision Record (ADR) documents an important architecture decisi
5455
- [ADR-015](ADR-015-cache-update-strategy.md) - Cache update strategy (update vs invalidation)
5556
- [ADR-016](ADR-016-two-tier-cache.md) - Two-tier cache (memory + Redis) for performance
5657
- [ADR-017](ADR-017-factory-pattern.md) - Factory pattern for extensibility
58+
- [ADR-018](ADR-018-event-driven-communication.md) - Event-driven communication (SNS/SQS)
5759

5860
### Code Conventions
5961

‎openshift/qontract-api.yaml‎

Lines changed: 156 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -160,6 +160,135 @@ objects:
160160
app.kubernetes.io/component: api
161161
app.kubernetes.io/name: qontract-api
162162

163+
164+
# --------- SUBSCRIBER DEPLOYMENT --------------
165+
- apiVersion: policy/v1
166+
kind: PodDisruptionBudget
167+
metadata:
168+
name: qontract-api-subscriber
169+
spec:
170+
minAvailable: 1
171+
selector:
172+
matchLabels:
173+
app.kubernetes.io/component: subscriber
174+
175+
- apiVersion: v1
176+
kind: ServiceAccount
177+
imagePullSecrets: "${{IMAGE_PULL_SECRETS}}"
178+
metadata:
179+
name: qontract-api-subscriber
180+
labels:
181+
app.kubernetes.io/component: subscriber
182+
app.kubernetes.io/name: qontract-api
183+
184+
- apiVersion: apps/v1
185+
kind: Deployment
186+
metadata:
187+
annotations:
188+
ignore-check.kube-linter.io/unset-cpu-requirements: "no cpu limits"
189+
labels:
190+
app.kubernetes.io/component: subscriber
191+
app.kubernetes.io/name: qontract-api
192+
name: qontract-api-subscriber
193+
spec:
194+
replicas: ${{QAPI_SUBSCRIBER_REPLICAS}}
195+
selector:
196+
matchLabels:
197+
app.kubernetes.io/component: subscriber
198+
app.kubernetes.io/name: qontract-api
199+
template:
200+
metadata:
201+
labels:
202+
app.kubernetes.io/component: subscriber
203+
app.kubernetes.io/name: qontract-api
204+
spec:
205+
restartPolicy: Always
206+
serviceAccountName: qontract-api-subscriber
207+
affinity:
208+
podAntiAffinity:
209+
preferredDuringSchedulingIgnoredDuringExecution:
210+
- podAffinityTerm:
211+
labelSelector:
212+
matchExpressions:
213+
- key: app.kubernetes.io/component
214+
operator: In
215+
values:
216+
- subscriber
217+
topologyKey: topology.kubernetes.io/zone
218+
weight: 100
219+
- podAffinityTerm:
220+
labelSelector:
221+
matchExpressions:
222+
- key: app.kubernetes.io/component
223+
operator: In
224+
values:
225+
- subscriber
226+
topologyKey: kubernetes.io/hostname
227+
weight: 1
228+
containers:
229+
- env:
230+
- name: QAPI_START_MODE
231+
value: subscriber
232+
- name: QAPI_SUBSCRIBER_METRICS_PORT
233+
value: "${QAPI_SUBSCRIBER_METRICS_PORT}"
234+
- name: QAPI_SENTRY_DSN
235+
valueFrom:
236+
secretKeyRef:
237+
name: ${QAPI_SENTRY_SECRET_NAME}
238+
key: dsn
239+
envFrom:
240+
- secretRef:
241+
name: qontract-api-secret
242+
optional: true
243+
- configMapRef:
244+
name: qontract-api-config
245+
optional: true
246+
image: "${IMAGE}:${IMAGE_TAG}"
247+
name: subscriber
248+
ports:
249+
- containerPort: ${{QAPI_SUBSCRIBER_METRICS_PORT}}
250+
readinessProbe:
251+
httpGet:
252+
path: /health/ready
253+
port: ${{QAPI_APP_PORT}}
254+
periodSeconds: 60
255+
timeoutSeconds: 5
256+
livenessProbe:
257+
httpGet:
258+
path: /health/live
259+
port: ${{QAPI_APP_PORT}}
260+
periodSeconds: 15
261+
timeoutSeconds: 5
262+
startupProbe:
263+
tcpSocket:
264+
port: ${{QAPI_APP_PORT}}
265+
initialDelaySeconds: 5
266+
periodSeconds: 5
267+
timeoutSeconds: 5
268+
resources:
269+
requests:
270+
cpu: ${{QAPI_SUBSCRIBER_CPU_REQUESTS}}
271+
memory: ${{QAPI_SUBSCRIBER_MEMORY_REQUESTS}}
272+
limits:
273+
memory: ${{QAPI_SUBSCRIBER_MEMORY_LIMITS}}
274+
275+
# ---------- SUBSCRIBER SERVICE (PROMETHEUS) -----------
276+
- apiVersion: v1
277+
kind: Service
278+
metadata:
279+
labels:
280+
app.kubernetes.io/component: subscriber
281+
app.kubernetes.io/name: qontract-api
282+
name: qontract-api-subscriber
283+
spec:
284+
ports:
285+
- name: "http"
286+
port: ${{QAPI_SUBSCRIBER_METRICS_PORT}}
287+
targetPort: ${{QAPI_SUBSCRIBER_METRICS_PORT}}
288+
selector:
289+
app.kubernetes.io/component: subscriber
290+
app.kubernetes.io/name: qontract-api
291+
163292
# --------- WORKER DEPLOYMENT --------------
164293
- apiVersion: policy/v1
165294
kind: PodDisruptionBudget
@@ -365,6 +494,33 @@ parameters:
365494
value: "100m"
366495
required: true
367496

497+
# Subscriber config
498+
- name: QAPI_SUBSCRIBER_METRICS_PORT
499+
description: Port to expose the web app on
500+
value: "8080"
501+
required: true
502+
503+
## Subscriber Pod limits
504+
- name: QAPI_SUBSCRIBER_REPLICAS
505+
description: Subscriber replicas
506+
value: "3"
507+
required: true
508+
509+
- name: QAPI_SUBSCRIBER_MEMORY_REQUESTS
510+
description: Subscriber memory requests
511+
value: "200Mi"
512+
required: true
513+
514+
- name: QAPI_SUBSCRIBER_MEMORY_LIMITS
515+
description: Subscriber memory limits
516+
value: "200Mi"
517+
required: true
518+
519+
- name: QAPI_SUBSCRIBER_CPU_REQUESTS
520+
description: Subscriber cpu requests
521+
value: "100m"
522+
required: true
523+
368524
# Worker config
369525
- name: QAPI_WORKER_METRICS_PORT
370526
description: Port to expose the web app on

‎qontract_api/Dockerfile‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,6 @@ ENV \
5050
# bypass dynamic versioning (git not available)
5151
UV_DYNAMIC_VERSIONING_BYPASS="0.1.0"
5252

53-
COPY --from=base --chown=1001:root $APP_ROOT $APP_ROOT
5453
COPY --chown=1001:root qontract_api qontract_api
5554
COPY --chown=1001:root qontract_api_client qontract_api_client
5655
COPY --chown=1001:root qontract_utils qontract_utils

0 commit comments

Comments
 (0)