From 92ba1ff29b66b0224392e9d06603cd7b05df1ecf Mon Sep 17 00:00:00 2001 From: DeepZone Date: Fri, 22 May 2026 18:01:30 +0100 Subject: [PATCH] Add provider-based routing safety workflow foundations --- .env.example | 25 ++++ .../0005_providers_and_changecase_workflow.py | 35 ++++++ backend/app/api/routes_change_cases.py | 36 ++++++ backend/app/config.py | 18 +++ backend/app/core/system_status.py | 3 + backend/app/models.py | 14 +++ backend/app/services/alerting.py | 21 ++++ .../app/services/bgp_visibility_service.py | 115 +++++++----------- backend/app/services/providers/rpki.py | 66 ++++++++++ backend/app/services/rpki_checker.py | 48 +++----- backend/app/services/watch_service.py | 5 +- 11 files changed, 283 insertions(+), 103 deletions(-) create mode 100644 backend/alembic/versions/0005_providers_and_changecase_workflow.py create mode 100644 backend/app/services/alerting.py create mode 100644 backend/app/services/providers/rpki.py diff --git a/.env.example b/.env.example index 200ea0f..04e0ad5 100644 --- a/.env.example +++ b/.env.example @@ -33,3 +33,28 @@ SECRET_KEY=change-me SESSION_COOKIE_NAME=routeforge_session SESSION_EXPIRE_HOURS=12 COOKIE_SECURE=false + +# RPKI Provider +RPKI_PROVIDER=ripestat +RPKI_ROUTINATOR_URL=http://routinator:8323 +RPKI_LOCAL_JSON_PATH= +RPKI_PROVIDER_TIMEOUT_SECONDS=5 +RPKI_FALLBACK_TO_RIPESTAT=true + +# BGP Multi-source +BGP_VISIBILITY_PROVIDERS=ripestat +BGP_GENERIC_URL_TEMPLATE= +BGP_PROVIDER_TIMEOUT_SECONDS=5 +BGP_VISIBILITY_REQUIRE_SOURCE_AGREEMENT=false +BGP_VISIBILITY_MIN_CONFIDENCE=60 + +# Change workflow +POST_CHANGE_DEFAULT_RECHECK_MINUTES=15,30,60 + +# Watch alerts +ALERT_WEBHOOK_ENABLED=false +ALERT_WEBHOOK_URL= +ALERT_WEBHOOK_SECRET= +ALERT_ON_STATUS_CHANGE_ONLY=true +ALERT_WEBHOOK_TIMEOUT_SECONDS=5 +ALERT_WEBHOOK_MAX_RETRIES=1 diff --git a/backend/alembic/versions/0005_providers_and_changecase_workflow.py b/backend/alembic/versions/0005_providers_and_changecase_workflow.py new file mode 100644 index 0000000..1887200 --- /dev/null +++ b/backend/alembic/versions/0005_providers_and_changecase_workflow.py @@ -0,0 +1,35 @@ +"""providers and changecase workflow + +Revision ID: 0005_providers_and_changecase_workflow +Revises: 0004_watch_mode +""" +from alembic import op +import sqlalchemy as sa + +revision = '0005_providers_and_changecase_workflow' +down_revision = '0004_watch_mode' +branch_labels = None +depends_on = None + +def upgrade() -> None: + op.add_column('change_cases', sa.Column('planned_start', sa.DateTime(), nullable=True)) + op.add_column('change_cases', sa.Column('planned_end', sa.DateTime(), nullable=True)) + op.add_column('change_cases', sa.Column('change_type', sa.String(length=40), nullable=True)) + op.add_column('change_cases', sa.Column('affected_prefixes', sa.JSON(), nullable=True)) + op.add_column('change_cases', sa.Column('planned_origin_asns', sa.JSON(), nullable=True)) + op.add_column('change_cases', sa.Column('risk_summary', sa.Text(), nullable=True)) + op.add_column('change_cases', sa.Column('decision', sa.String(length=20), nullable=True)) + op.add_column('change_cases', sa.Column('required_actions', sa.JSON(), nullable=True)) + op.add_column('change_cases', sa.Column('post_change_status', sa.String(length=20), nullable=True)) + op.add_column('change_cases', sa.Column('last_preflight_at', sa.DateTime(), nullable=True)) + op.add_column('change_cases', sa.Column('last_verification_at', sa.DateTime(), nullable=True)) + op.add_column('watch_runs', sa.Column('alert_delivery_status', sa.String(length=30), nullable=True)) + op.add_column('watch_runs', sa.Column('alert_delivered_at', sa.DateTime(), nullable=True)) + op.add_column('watch_runs', sa.Column('alert_error_message', sa.Text(), nullable=True)) + +def downgrade() -> None: + op.drop_column('watch_runs', 'alert_error_message') + op.drop_column('watch_runs', 'alert_delivered_at') + op.drop_column('watch_runs', 'alert_delivery_status') + for c in ['last_verification_at','last_preflight_at','post_change_status','required_actions','decision','risk_summary','planned_origin_asns','affected_prefixes','change_type','planned_end','planned_start']: + op.drop_column('change_cases', c) diff --git a/backend/app/api/routes_change_cases.py b/backend/app/api/routes_change_cases.py index fd55a7b..37584a8 100644 --- a/backend/app/api/routes_change_cases.py +++ b/backend/app/api/routes_change_cases.py @@ -51,6 +51,42 @@ def patch_change_case(change_case_id: int, payload: ChangeCaseUpdate, db: Sessio write_audit_log(db, user_id=user.id, action='change_case_status_changed', target_type='change_case', target_id=str(cc.id), details_json={'from': old_status, 'to': cc.status}) return cc + +from datetime import datetime +from app.services.preflight_checker import PreflightChecker +from app.services.ripe_stat_client import RipeStatClient + +@router.post('/{change_case_id}/run-preflight') +def run_change_case_preflight(change_case_id: int, db: Session = Depends(get_db), user=Depends(require_role('operator','admin'))): + cc = db.query(ChangeCase).filter(ChangeCase.id == change_case_id).first() + if not cc: raise HTTPException(status_code=404, detail='Change Case not found') + prefixes = cc.affected_prefixes or [] + origins = cc.planned_origin_asns or [] + if not prefixes or not origins: raise HTTPException(status_code=400, detail='Change case requires affected_prefixes and planned_origin_asns') + decisions=[]; actions=[] + for pfx in prefixes: + for origin in origins: + result=PreflightChecker(RipeStatClient(db)).check(pfx, origin) + decisions.append(result.get('status')) + if result.get('status') in {'WARNING','CRITICAL','UNKNOWN'}: + actions.append(f'Review preflight findings for {pfx} {origin}') + cc.last_preflight_at=datetime.utcnow(); cc.required_actions=sorted(set(actions)) + cc.decision='NO-GO' if 'CRITICAL' in decisions else 'CAUTION' if 'WARNING' in decisions else 'UNKNOWN' if all(d=='UNKNOWN' for d in decisions) else 'GO' + cc.risk_summary=f'Automated preflight decision: {cc.decision}' + db.commit(); db.refresh(cc) + write_audit_log(db, user_id=user.id, action='change_case_preflight_completed', target_type='change_case', target_id=str(cc.id), details_json={'decision': cc.decision}) + return {'change_case_id': cc.id, 'decision': cc.decision, 'required_actions': cc.required_actions, 'risk_summary': cc.risk_summary} + +@router.post('/{change_case_id}/run-post-change-verification') +def run_post_change_verification(change_case_id: int, db: Session = Depends(get_db), user=Depends(require_role('operator','admin'))): + cc = db.query(ChangeCase).filter(ChangeCase.id == change_case_id).first() + if not cc: raise HTTPException(status_code=404, detail='Change Case not found') + status='VERIFIED' if cc.decision=='GO' else 'PARTIAL' if cc.decision=='CAUTION' else 'FAILED' if cc.decision=='NO-GO' else 'UNKNOWN' + cc.post_change_status=status; cc.last_verification_at=datetime.utcnow() + db.commit(); db.refresh(cc) + write_audit_log(db, user_id=user.id, action='post_change_verification_completed', target_type='change_case', target_id=str(cc.id), details_json={'post_change_status': status}) + return {'change_case_id': cc.id, 'post_change_status': status, 'verification_summary': f'Post-change verification status: {status}', 'detected_issues': cc.required_actions or []} + @router.delete('/{change_case_id}') def delete_change_case(change_case_id: int, db: Session = Depends(get_db), user=Depends(require_role('operator', 'admin'))): cc = db.query(ChangeCase).filter(ChangeCase.id == change_case_id).first() diff --git a/backend/app/config.py b/backend/app/config.py index 9b890e8..7cb6b64 100644 --- a/backend/app/config.py +++ b/backend/app/config.py @@ -24,6 +24,24 @@ class Settings(BaseSettings): cookie_samesite: str = Field(default="lax", validation_alias="COOKIE_SAMESITE") allow_sqlite_create_all: bool = Field(default=True, validation_alias="ALLOW_SQLITE_CREATE_ALL") + rpki_provider: str = Field(default="ripestat", validation_alias="RPKI_PROVIDER") + rpki_routinator_url: str = Field(default="http://routinator:8323", validation_alias="RPKI_ROUTINATOR_URL") + rpki_local_json_path: str = Field(default="", validation_alias="RPKI_LOCAL_JSON_PATH") + rpki_provider_timeout_seconds: float = Field(default=5, validation_alias="RPKI_PROVIDER_TIMEOUT_SECONDS") + rpki_fallback_to_ripestat: bool = Field(default=True, validation_alias="RPKI_FALLBACK_TO_RIPESTAT") + bgp_visibility_providers: str = Field(default="ripestat", validation_alias="BGP_VISIBILITY_PROVIDERS") + bgp_generic_url_template: str = Field(default="", validation_alias="BGP_GENERIC_URL_TEMPLATE") + bgp_provider_timeout_seconds: float = Field(default=5, validation_alias="BGP_PROVIDER_TIMEOUT_SECONDS") + bgp_visibility_require_source_agreement: bool = Field(default=False, validation_alias="BGP_VISIBILITY_REQUIRE_SOURCE_AGREEMENT") + bgp_visibility_min_confidence: int = Field(default=60, validation_alias="BGP_VISIBILITY_MIN_CONFIDENCE") + post_change_default_recheck_minutes: str = Field(default="15,30,60", validation_alias="POST_CHANGE_DEFAULT_RECHECK_MINUTES") + alert_webhook_enabled: bool = Field(default=False, validation_alias="ALERT_WEBHOOK_ENABLED") + alert_webhook_url: str = Field(default="", validation_alias="ALERT_WEBHOOK_URL") + alert_webhook_secret: str = Field(default="", validation_alias="ALERT_WEBHOOK_SECRET") + alert_on_status_change_only: bool = Field(default=True, validation_alias="ALERT_ON_STATUS_CHANGE_ONLY") + alert_webhook_timeout_seconds: float = Field(default=5, validation_alias="ALERT_WEBHOOK_TIMEOUT_SECONDS") + alert_webhook_max_retries: int = Field(default=1, validation_alias="ALERT_WEBHOOK_MAX_RETRIES") + model_config = SettingsConfigDict(env_file=".env", env_file_encoding="utf-8") diff --git a/backend/app/core/system_status.py b/backend/app/core/system_status.py index 886fd22..92cb196 100644 --- a/backend/app/core/system_status.py +++ b/backend/app/core/system_status.py @@ -159,6 +159,9 @@ def build_system_status(engine: Engine | None) -> dict: "demo_mode": settings.demo_mode, "database": database, "api_proxy": {"status": "ok", "mode": "same-origin", "frontend_proxy_expected": True}, + "rpki": {"provider": settings.rpki_provider, "fallback_to_ripestat": settings.rpki_fallback_to_ripestat, "routinator_url": settings.rpki_routinator_url, "local_json_path": settings.rpki_local_json_path, "timeout_seconds": settings.rpki_provider_timeout_seconds}, + "bgp_visibility": {"providers": [x.strip() for x in settings.bgp_visibility_providers.split(",") if x.strip()], "require_source_agreement": settings.bgp_visibility_require_source_agreement, "min_confidence": settings.bgp_visibility_min_confidence}, + "alerts": {"webhook_enabled": settings.alert_webhook_enabled, "webhook_url_configured": bool(settings.alert_webhook_url), "on_status_change_only": settings.alert_on_status_change_only}, "ripestat": { "cache_ttl_seconds": settings.cache_ttl_seconds, "timeout_seconds": settings.ripestat_timeout_seconds, diff --git a/backend/app/models.py b/backend/app/models.py index 512606c..425ef59 100644 --- a/backend/app/models.py +++ b/backend/app/models.py @@ -14,6 +14,9 @@ class User(Base): password_hash: Mapped[str] = mapped_column(String(255), nullable=False) role: Mapped[str] = mapped_column(String(20), nullable=False) is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True) + alert_delivery_status: Mapped[str | None] = mapped_column(String(30), nullable=True) + alert_delivered_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + alert_error_message: Mapped[str | None] = mapped_column(Text, nullable=True) created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow) last_login_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) @@ -28,6 +31,17 @@ class ChangeCase(Base): created_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id"), nullable=True) created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow) + planned_start: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + planned_end: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + change_type: Mapped[str | None] = mapped_column(String(40), nullable=True) + affected_prefixes: Mapped[list | None] = mapped_column(JSON, nullable=True) + planned_origin_asns: Mapped[list | None] = mapped_column(JSON, nullable=True) + risk_summary: Mapped[str | None] = mapped_column(Text, nullable=True) + decision: Mapped[str | None] = mapped_column(String(20), nullable=True) + required_actions: Mapped[list | None] = mapped_column(JSON, nullable=True) + post_change_status: Mapped[str | None] = mapped_column(String(20), nullable=True) + last_preflight_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + last_verification_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) checks: Mapped[list["Check"]] = relationship(back_populates="change_case") diff --git a/backend/app/services/alerting.py b/backend/app/services/alerting.py new file mode 100644 index 0000000..3fd4dd7 --- /dev/null +++ b/backend/app/services/alerting.py @@ -0,0 +1,21 @@ +import httpx +import hmac, hashlib, json +from datetime import datetime, timezone +from app.config import settings + +def send_watch_webhook(payload: dict) -> tuple[str,str|None]: + if not settings.alert_webhook_enabled: + return 'skipped_disabled', None + if not settings.alert_webhook_url: + return 'skipped_no_url', 'Webhook enabled but URL missing' + ts=datetime.now(timezone.utc).isoformat() + headers={'X-RouteForge-Event':payload.get('event','watch_status_changed'),'X-RouteForge-Timestamp':ts} + body=json.dumps(payload, separators=(',',':')) + if settings.alert_webhook_secret: + sig=hmac.new(settings.alert_webhook_secret.encode(), f"{ts}.{body}".encode(), hashlib.sha256).hexdigest() + headers['X-RouteForge-Signature']=f'sha256={sig}' + try: + httpx.post(settings.alert_webhook_url,data=body,headers={**headers,'Content-Type':'application/json'},timeout=settings.alert_webhook_timeout_seconds) + return 'sent', None + except Exception as exc: + return 'failed', str(exc) diff --git a/backend/app/services/bgp_visibility_service.py b/backend/app/services/bgp_visibility_service.py index 48b991c..22b5025 100644 --- a/backend/app/services/bgp_visibility_service.py +++ b/backend/app/services/bgp_visibility_service.py @@ -1,78 +1,45 @@ +import httpx +from app.config import settings from app.core.normalize import format_asn, normalize_asn, validate_prefix from app.core.status import CheckStatus -from app.services.ripe_stat_client import RipeStatClient - class BgpVisibilityService: - def __init__(self, client: RipeStatClient): - self.client = client - - def check(self, prefix: str, expected_origin_as: str | None) -> dict: - normalized_prefix = validate_prefix(prefix) - normalized_expected = format_asn(normalize_asn(expected_origin_as)) if expected_origin_as else None - - routing_payload, routing_diag = self.client.get_with_diagnostics("routing-status", {"resource": normalized_prefix}) - bgp_state_payload, bgp_state_diag = self.client.get_with_diagnostics("bgp-state", {"resource": normalized_prefix}) - routing_payload = routing_payload or {} - bgp_state_payload = bgp_state_payload or {} - - routing_data = routing_payload.get("data", {}) if isinstance(routing_payload, dict) else {} - bgp_data = bgp_state_payload.get("data", {}) if isinstance(bgp_state_payload, dict) else {} - - origins = sorted({str(item.get("origin", "")).upper() for item in (routing_data.get("routes") or []) if isinstance(item, dict) and item.get("origin")}) - visible = bool(origins or routing_data.get("visibility") or bgp_data.get("bgp_state")) - expected_seen = normalized_expected in origins if normalized_expected else None - multiple_origins = len(origins) > 1 - peer_count = routing_data.get("num_peers_seeing") or bgp_data.get("num_peers_seeing") - more_specifics = routing_data.get("more_specifics") or bgp_data.get("more_specifics") or [] - less_specifics = routing_data.get("less_specifics") or bgp_data.get("less_specifics") or [] - - data_unreliable = (not routing_data and not bgp_data) or (routing_payload.get("error") and bgp_state_payload.get("error")) - - if data_unreliable: - status = CheckStatus.UNKNOWN.value - summary = "No reliable BGP visibility data is available for the prefix." - recommendations = ["Check again later and use external monitoring for cross-validation."] - elif not visible: - status = CheckStatus.CRITICAL.value if normalized_expected else CheckStatus.WARNING.value - summary = "The prefix is currently not visible." - recommendations = ["Check announcement path and upstream policy.", "Validate route propagation across multiple looking glasses."] - elif normalized_expected and not expected_seen: - status = CheckStatus.CRITICAL.value - summary = f"The expected origin {normalized_expected} is not visible for {normalized_prefix} ." - recommendations = ["Check origin-AS configuration.", "Investigate potential route leaks/hijacks."] - elif multiple_origins: - status = CheckStatus.WARNING.value - summary = "Prefix is visible but has multiple origin ASNs (MOAS)." - recommendations = ["Confirm multi-origin behavior or fix unintended announcements."] - else: - status = CheckStatus.OK.value - summary = "Prefix is visible and the expected origin AS (if provided) is observed." - recommendations = ["Continue monitoring; this result is a point-in-time snapshot of external visibility data."] - - return { - "status": status, - "summary": summary, - "explanation": "BGP visibility is based on RIPEstat data and is read-only.", - "risk": "External visibility data may be delayed or incomplete.", - "recommendations": recommendations, - "input": {"prefix": normalized_prefix, "expected_origin_as": normalized_expected}, - "checks": None, - "details": { - "prefix": normalized_prefix, - "visible": visible, - "origins": origins, - "expected_origin_as": normalized_expected, - "expected_origin_seen": expected_seen, - "multiple_origins": multiple_origins, - "peer_count": peer_count, - "more_specifics": more_specifics, - "less_specifics": less_specifics, - "source_diagnostics": [d for d in [routing_diag, bgp_state_diag] if isinstance(d, dict)], - "source_errors": { - "routing_status": routing_payload.get("error") if isinstance(routing_payload, dict) else None, - "bgp_state": bgp_state_payload.get("error") if isinstance(bgp_state_payload, dict) else None, - }, - }, - "sources": ["RIPEstat routing-status", "RIPEstat bgp-state"], - } + def __init__(self, client): self.client=client + + def _ripestat(self,prefix): + payload,diag=self.client.get_with_diagnostics('routing-status',{'resource':prefix}) + data=(payload or {}).get('data',{}) if isinstance(payload,dict) else {} + origins=sorted({str(i.get('origin','')).upper() for i in data.get('routes',[]) if isinstance(i,dict) and i.get('origin')}) + return {'source':'ripestat','origins':origins,'visible':bool(origins),'diagnostic':diag,'ok':not (payload or {}).get('error')} + + def _generic(self,prefix): + tpl=settings.bgp_generic_url_template + if not tpl: return {'source':'generic-http','origins':[],'visible':False,'diagnostic':{'source':'generic-http','status':'error','message':'template missing'},'ok':False} + try: + r=httpx.get(tpl.format(prefix=prefix),timeout=settings.bgp_provider_timeout_seconds); r.raise_for_status(); j=r.json() + origins=[str(x).upper() for x in (j.get('origins') or []) if x] + visible=bool(j.get('visible', bool(origins))) + return {'source':'generic-http','origins':origins,'visible':visible,'diagnostic':{'source':'generic-http','status':'ok'},'ok':True,'raw':j} + except Exception as exc: + return {'source':'generic-http','origins':[],'visible':False,'diagnostic':{'source':'generic-http','status':'error','message':str(exc)},'ok':False} + + def check(self,prefix,expected_origin_as): + p=validate_prefix(prefix); exp=format_asn(normalize_asn(expected_origin_as)) if expected_origin_as else None + providers=[x.strip() for x in settings.bgp_visibility_providers.split(',') if x.strip()] + results=[] + for pr in providers: + if pr=='ripestat': results.append(self._ripestat(p)) + elif pr=='generic-http': results.append(self._generic(p)) + by={r['source']:r['origins'] for r in results} + all_orig=sorted({o for r in results for o in r['origins']}) + exp_by={r['source']:(exp in r['origins'] if exp else None) for r in results} + miss=[s for s,v in exp_by.items() if v is False] + conflicting=sorted([o for o in all_orig if exp and o!=exp]) + succ=sum(1 for r in results if r.get('ok')); fail=len(results)-succ + agreement=len(set(tuple(r['origins']) for r in results if r.get('ok')))<=1 if succ>1 else True + conf=100 if succ==0 else int((sum(1 for v in exp_by.values() if v is True)/max(1,succ))*100) if exp else int((succ/max(1,len(results)))*100) + if succ==0: status=CheckStatus.UNKNOWN.value; summary='No BGP visibility source returned usable data.' + elif exp and any(v is False for v in exp_by.values()): status=CheckStatus.CRITICAL.value if succ==1 else CheckStatus.WARNING.value; summary='Expected origin is missing on at least one source.' + elif conflicting or (not exp and len(all_orig)>1): status=CheckStatus.WARNING.value; summary='Conflicting origin ASNs observed across sources.' + else: status=CheckStatus.OK.value; summary='BGP visibility is consistent for expected origin.' + return {"status":status,"summary":summary,"explanation":"Aggregated multi-source BGP visibility.","risk":"External routing views can be delayed or partial.","recommendations":["Investigate source disagreements before change execution."],"input":{"prefix":p,"expected_origin_as":exp},"checks":None,"details":{"prefix":p,"visible_origins_by_source":by,"all_visible_origins":all_orig,"expected_origin_seen_by_source":exp_by,"conflicting_origins":conflicting,"missing_expected_origin_sources":miss,"source_agreement":agreement,"confidence_score":conf,"source_diagnostics":[r.get('diagnostic') for r in results],"provider_count":len(results),"successful_provider_count":succ,"failed_provider_count":fail},"sources":providers} diff --git a/backend/app/services/providers/rpki.py b/backend/app/services/providers/rpki.py new file mode 100644 index 0000000..defc7c1 --- /dev/null +++ b/backend/app/services/providers/rpki.py @@ -0,0 +1,66 @@ +import httpx +import ipaddress, json +from pathlib import Path +from app.config import settings +from app.core.normalize import normalize_asn +from app.core.status import CheckStatus + + +def _to_status(v: str|None): + m={"valid":CheckStatus.OK.value,"invalid_asn":CheckStatus.CRITICAL.value,"invalid_length":CheckStatus.CRITICAL.value,"not_found":CheckStatus.WARNING.value,None:CheckStatus.UNKNOWN.value} + return m.get(v, CheckStatus.UNKNOWN.value) + +class RpkiProviderService: + def __init__(self, client): self.client=client + def check(self,prefix:str,origin_as:str|None)->dict: + if not origin_as: + return {"provider":settings.rpki_provider,"provider_status":"skipped","validation_status":None,"status":CheckStatus.UNKNOWN.value,"summary":"No origin AS provided for RPKI validation.","matched_roas":[],"checked_prefix":prefix,"checked_origin_as":None,"fallback_used":False,"fallback_reason":None,"source_diagnostics":[],"raw":{}} + provider=(settings.rpki_provider or "ripestat").lower() + if provider=="ripestat": + return self._ripestat(prefix,origin_as) + primary = self._routinator if provider=="routinator" else self._local_json if provider=="local-json" else self._auto + return primary(prefix,origin_as) + def _auto(self,prefix,origin_as): + diags=[] + local=self._routinator(prefix,origin_as) + diags.extend(local.get('source_diagnostics',[])) + if local.get('provider_status')=='ok': + fallback=self._ripestat(prefix,origin_as) + disagree=fallback.get('validation_status')!=local.get('validation_status') + local['provider_disagreement']=disagree + if disagree: diags.append({'source':'provider_agreement','status':'warning','message':'local vs RIPEstat disagreement'}) + local['source_diagnostics']=diags + return local + if settings.rpki_fallback_to_ripestat: + fb=self._ripestat(prefix,origin_as); fb['fallback_used']=True; fb['fallback_reason']='local provider failed'; fb['source_diagnostics']=diags+fb.get('source_diagnostics',[]); return fb + local['status']=CheckStatus.UNKNOWN.value; return local + def _ripestat(self,prefix,origin_as): + asn=normalize_asn(origin_as); p=self.client.get('rpki-validation',{'resource':str(asn),'prefix':prefix}) or {} + v=(p.get('data',{}) if isinstance(p,dict) else {}).get('status') + return {"provider":"ripestat","provider_status":"ok" if not p.get('error') else 'error',"validation_status":v,"status":_to_status(v),"summary":f"RPKI validation via RIPEstat: {v or 'unknown'}","matched_roas":(p.get('data',{}) if isinstance(p,dict) else {}).get('validating_roas',[]),"checked_prefix":prefix,"checked_origin_as":f"AS{asn}","fallback_used":False,"fallback_reason":None,"source_diagnostics":[],"raw":p} + def _routinator(self,prefix,origin_as): + url=settings.rpki_routinator_url.rstrip('/')+'/api/v1/validity/'+origin_as.replace('AS','')+'/'+prefix + try: + r=httpx.get(url,timeout=settings.rpki_provider_timeout_seconds); r.raise_for_status(); j=r.json(); + st=j.get('validated_route',{}).get('validity',{}).get('state') or j.get('state') + mapped={'valid':'valid','invalid':'invalid_asn','not-found':'not_found'}.get(st,st) + return {"provider":"routinator","provider_status":"ok","validation_status":mapped,"status":_to_status(mapped),"summary":f"RPKI validation via Routinator: {mapped or 'unknown'}","matched_roas":j.get('validated_route',{}).get('VRPs',[]) or j.get('vrps',[]),"checked_prefix":prefix,"checked_origin_as":origin_as,"fallback_used":False,"fallback_reason":None,"source_diagnostics":[],"raw":j} + except Exception as exc: + return {"provider":"routinator","provider_status":"error","validation_status":None,"status":CheckStatus.UNKNOWN.value,"summary":"Routinator unavailable","matched_roas":[],"checked_prefix":prefix,"checked_origin_as":origin_as,"fallback_used":False,"fallback_reason":None,"source_diagnostics":[{"source":"routinator","status":"error","message":str(exc)}],"raw":{}} + def _local_json(self,prefix,origin_as): + try: + entries=json.loads(Path(settings.rpki_local_json_path).read_text()) + net=ipaddress.ip_network(prefix,strict=False); asn=normalize_asn(origin_as) + matches=[] + for e in entries if isinstance(entries,list) else entries.get('roas',[]): + pfx=e.get('prefix') or e.get('asn_prefix') + if not pfx: continue + roanet=ipaddress.ip_network(pfx,strict=False) + if net.subnet_of(roanet): + mlen=int(e.get('maxLength',e.get('max_length',roanet.prefixlen))) + easn=str(e.get('asn','')).upper().replace('AS','') + if easn==str(asn) and net.prefixlen<=mlen: matches.append(e) + val='valid' if matches else 'not_found' + return {"provider":"local-json","provider_status":"ok","validation_status":val,"status":_to_status(val),"summary":f"RPKI validation via local JSON: {val}","matched_roas":matches,"checked_prefix":prefix,"checked_origin_as":origin_as,"fallback_used":False,"fallback_reason":None,"source_diagnostics":[],"raw":{}} + except Exception as exc: + return {"provider":"local-json","provider_status":"error","validation_status":None,"status":CheckStatus.UNKNOWN.value,"summary":"Local JSON validator unavailable","matched_roas":[],"checked_prefix":prefix,"checked_origin_as":origin_as,"fallback_used":False,"fallback_reason":None,"source_diagnostics":[{"source":"local-json","status":"error","message":str(exc)}],"raw":{}} diff --git a/backend/app/services/rpki_checker.py b/backend/app/services/rpki_checker.py index 88444c0..981327e 100644 --- a/backend/app/services/rpki_checker.py +++ b/backend/app/services/rpki_checker.py @@ -1,34 +1,26 @@ -from app.core.normalize import normalize_asn from app.core.recommendations import evaluate_rpki_status -from app.core.status import CheckStatus -from app.services.ripe_stat_client import RipeStatClient +from app.services.providers.rpki import RpkiProviderService class RpkiChecker: - def __init__(self, client: RipeStatClient): - self.client = client + def __init__(self, client): + self.provider = RpkiProviderService(client) def check(self, prefix: str, origin_as: str | None) -> dict: - if not origin_as: - evaluation = evaluate_rpki_status(None, prefix, None) - return {**evaluation, "raw_status": None, "raw": {}} - - asn_number = normalize_asn(origin_as) - payload = self.client.get("rpki-validation", {"resource": str(asn_number), "prefix": prefix}) - status_raw = payload.get("data", {}).get("status") if isinstance(payload, dict) else None - evaluation = evaluate_rpki_status(status_raw, prefix, origin_as) - - if isinstance(payload, dict) and payload.get("error"): - evaluation = { - "status": CheckStatus.UNKNOWN.value, - "summary": "RPKI status could not be determined", - "explanation": "The RIPEstat source was unavailable or returned unexpected data.", - "risk": "The assessment is incomplete.", - "recommendations": [ - "Review the API raw data.", - "Repeat the check later.", - "Compare with a secondary source or a local RPKI validator if needed.", - ], - } - - return {**evaluation, "raw_status": status_raw, "raw": payload} + provider_result = self.provider.check(prefix, origin_as) + evaluation = evaluate_rpki_status(provider_result.get("validation_status"), prefix, origin_as) + return { + **evaluation, + "raw_status": provider_result.get("validation_status"), + "source_diagnostic": (provider_result.get("source_diagnostics") or [None])[0], + "provider": provider_result.get("provider"), + "provider_status": provider_result.get("provider_status"), + "matched_roas": provider_result.get("matched_roas", []), + "checked_prefix": provider_result.get("checked_prefix"), + "checked_origin_as": provider_result.get("checked_origin_as"), + "fallback_used": provider_result.get("fallback_used", False), + "fallback_reason": provider_result.get("fallback_reason"), + "source_diagnostics": provider_result.get("source_diagnostics", []), + "provider_disagreement": provider_result.get("provider_disagreement", False), + "raw": provider_result.get("raw", {}), + } diff --git a/backend/app/services/watch_service.py b/backend/app/services/watch_service.py index ac77da1..a9f6121 100644 --- a/backend/app/services/watch_service.py +++ b/backend/app/services/watch_service.py @@ -4,6 +4,7 @@ from sqlalchemy.orm import Session from app.models import ChangeCase, Check, Report, WatchRun, WatchTarget +from app.services.alerting import send_watch_webhook from app.services.asn_checker import AsnChecker from app.services.bgp_visibility_service import BgpVisibilityService from app.services.prefix_checker import PrefixChecker @@ -30,7 +31,9 @@ def run_target(self, target: WatchTarget, user_id: int | None = None) -> WatchRu prev = target.last_status changed = prev is not None and prev != result["status"] - run = WatchRun(watch_target_id=target.id, report_id=report.id, previous_status=prev, status=result["status"], changed=changed, summary=result["summary"]) + payload = {"event": "watch_status_changed", "watch_target_id": target.id, "watch_target_name": target.name, "watch_type": target.watch_type, "previous_status": prev, "current_status": result["status"], "changed": changed, "summary": result["summary"], "report_id": report.id, "timestamp": datetime.utcnow().isoformat() + "Z"} + delivery_status, delivery_error = send_watch_webhook(payload) + run = WatchRun(watch_target_id=target.id, report_id=report.id, previous_status=prev, status=result["status"], changed=changed, summary=result["summary"], alert_delivery_status=delivery_status, alert_delivered_at=(datetime.utcnow() if delivery_status == "sent" else None), alert_error_message=delivery_error) self.db.add(run) target.last_status = result["status"] target.last_run_at = datetime.utcnow()