Registers an abonnement on the `objecten` kanaal pointing at the existing
webhook sink, writes a RegisterRecord exactly as the ACL does on approval, and
waits for the delivery. That is the whole publish chain in one assertion:
Objecten → its celery worker → NRC → nrc-beat → the callback.
Fails today at the first hop, which is the point:
NRC POST /api/v1/abonnement → 400: {"name":"naam","code":"kanaal_naam",
"reason":"Kanaal met deze naam bestaat niet."}
Also splits S-19b (#150) into #152/#153 in BACKLOG.md — the notification wiring
and the projection re-source are independently deployable (CLAUDE.md §13).
120 lines
4.8 KiB
Python
120 lines
4.8 KiB
Python
#!/usr/bin/env python3
|
|
"""S-19b-1 (#152): driver for the Objecten → NRC notification check.
|
|
|
|
Registers an abonnement on the `objecten` kanaal pointing at the webhook sink, then writes a
|
|
RegisterRecord object exactly as the ACL's ObjectenGateway does (S-19a). The caller
|
|
(run-objecten-notifications-check.sh) watches the sink for the delivery — this only sets it up,
|
|
and prints `REFERENCE <value>` for the caller to grep on.
|
|
|
|
Delivery exercises the whole chain: Objecten → its celery worker → NRC → nrc-beat → the callback.
|
|
Anything missing (broker, worker, kanaal, notifications config) shows up as a non-delivery.
|
|
|
|
Stdlib only so it runs in a bare python:3-slim container on the compose network.
|
|
"""
|
|
import base64
|
|
import hashlib
|
|
import hmac
|
|
import json
|
|
import os
|
|
import sys
|
|
import time
|
|
import urllib.error
|
|
import urllib.request
|
|
|
|
OBJECTEN = os.environ["OBJECTEN"] # http://objecten:8000
|
|
OBJECTEN_TOKEN = os.environ["OBJECTEN_TOKEN"]
|
|
OBJECTTYPEN = os.environ["OBJECTTYPEN"] # http://objecttypen:8000
|
|
OBJECTTYPEN_TOKEN = os.environ["OBJECTTYPEN_TOKEN"]
|
|
NRC_BASE = os.environ["NRC_BASE"] # http://<nrc-ip>:8000
|
|
SINK_CALLBACK = os.environ["SINK_CALLBACK"] # http://<sink-ip>:9000/
|
|
SINK_AUTH = os.environ["SINK_AUTH"]
|
|
CLIENT_ID = os.environ.get("NRC_CLIENT_ID", "big-reference-seed")
|
|
SECRET = os.environ.get("NRC_SECRET", "insecure-dev-secret-change-me")
|
|
KANAAL = "objecten"
|
|
|
|
|
|
def mint():
|
|
"""The HS256 JWT NRC expects (same shape as infra/local/register-abonnement.py)."""
|
|
def seg(d):
|
|
return base64.urlsafe_b64encode(json.dumps(d).encode()).rstrip(b"=")
|
|
|
|
payload = seg({
|
|
"iss": CLIENT_ID, "iat": int(time.time()), "client_id": CLIENT_ID,
|
|
"user_id": CLIENT_ID, "user_representation": CLIENT_ID,
|
|
})
|
|
signing_input = seg({"typ": "JWT", "alg": "HS256"}) + b"." + payload
|
|
signature = base64.urlsafe_b64encode(
|
|
hmac.new(SECRET.encode(), signing_input, hashlib.sha256).digest()).rstrip(b"=")
|
|
return (signing_input + b"." + signature).decode()
|
|
|
|
|
|
def nrc(method, url, body=None):
|
|
"""Call NRC. `url` may be a path or an absolute URL (the list returns absolute ones)."""
|
|
data = json.dumps(body).encode() if body is not None else None
|
|
req = urllib.request.Request(
|
|
url if url.startswith("http") else f"{NRC_BASE}{url}", data=data, method=method,
|
|
headers={"Authorization": f"Bearer {mint()}", "Content-Type": "application/json"})
|
|
try:
|
|
with urllib.request.urlopen(req, timeout=15) as r:
|
|
return json.load(r) if r.length != 0 else {}
|
|
except urllib.error.HTTPError as e:
|
|
# The body carries the reason (e.g. an unregistered kanaal); the status alone does not.
|
|
raise SystemExit(f"FAIL — NRC {method} {url} → {e.code}: {e.read().decode(errors='replace')[:400]}")
|
|
|
|
|
|
def token_api(base, token, method, path, body=None, crs=False):
|
|
data = json.dumps(body).encode() if body is not None else None
|
|
headers = {"Authorization": f"Token {token}"}
|
|
if body is not None:
|
|
headers["Content-Type"] = "application/json"
|
|
if crs:
|
|
headers["Accept-Crs"] = "EPSG:4326"
|
|
if body is not None:
|
|
headers["Content-Crs"] = "EPSG:4326"
|
|
req = urllib.request.Request(f"{base}{path}", data=data, method=method, headers=headers)
|
|
with urllib.request.urlopen(req, timeout=15) as r:
|
|
return json.load(r) if r.length != 0 else {}
|
|
|
|
|
|
def subscribe():
|
|
"""Register an abonnement on the objecten kanaal, replacing a stale one for the same callback."""
|
|
# NRC returns a bare list here, not a paginated envelope.
|
|
for existing in nrc("GET", "/api/v1/abonnement") or []:
|
|
if existing.get("callbackUrl") == SINK_CALLBACK:
|
|
nrc("DELETE", existing["url"])
|
|
nrc("POST", "/api/v1/abonnement", {
|
|
"callbackUrl": SINK_CALLBACK,
|
|
"auth": SINK_AUTH,
|
|
"kanalen": [{"naam": KANAAL, "filters": {}}],
|
|
})
|
|
print(f">> abonnement on '{KANAAL}' -> {SINK_CALLBACK}")
|
|
|
|
|
|
def objecttype_url():
|
|
results = token_api(OBJECTTYPEN, OBJECTTYPEN_TOKEN, "GET", "/api/v2/objecttypes").get("results", [])
|
|
match = next((o for o in results if o.get("name") == "RegisterRecord"), None)
|
|
if not match:
|
|
print("FAIL — no RegisterRecord objecttype in Objecttypen", file=sys.stderr)
|
|
raise SystemExit(1)
|
|
return match["url"]
|
|
|
|
|
|
def main():
|
|
subscribe()
|
|
reference = f"NOTIF-{int(time.time())}"
|
|
created = token_api(OBJECTEN, OBJECTEN_TOKEN, "POST", "/api/v2/objects", {
|
|
"type": objecttype_url(),
|
|
"record": {
|
|
"typeVersion": 1,
|
|
"data": {"id": f"zaak-{reference}", "status": "INGESCHREVEN", "reference": reference},
|
|
"startAt": time.strftime("%Y-%m-%d"),
|
|
},
|
|
}, crs=True)
|
|
print(f">> wrote RegisterRecord {created['url']}")
|
|
print(f"REFERENCE {reference}")
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|