feat: add workspace RBAC audit foundations

This commit is contained in:
Christian Krakau-Louis
2026-05-23 19:22:56 +02:00
parent 6ddd42bb3e
commit d180f984a6
20 changed files with 1199 additions and 54 deletions
+35 -1
View File
@@ -2,7 +2,7 @@
from typing import List, Optional
from fastapi import APIRouter, Depends, HTTPException, status
from fastapi import APIRouter, Depends, HTTPException, Request, status
from pydantic import BaseModel, Field
from sqlalchemy.orm import Session
@@ -15,6 +15,8 @@ from app.services.api_tokens import (
revoke_api_token,
token_to_dict,
)
from app.services.workspace_audit import record_workspace_audit_log
from app.services.workspaces import assign_default_workspace_to_unscoped_rows
router = APIRouter()
@@ -71,10 +73,12 @@ async def list_api_tokens(
@router.post("", response_model=APITokenCreateResponse, status_code=status.HTTP_201_CREATED)
async def create_public_api_token(
payload: APITokenCreateRequest,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
):
"""Create a scoped API token for read-only automation."""
workspace = assign_default_workspace_to_unscoped_rows(db)
try:
created = create_api_token(db, name=payload.name, scopes=payload.scopes)
except ValueError as exc:
@@ -82,6 +86,21 @@ async def create_public_api_token(
status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
detail=str(exc),
) from exc
record_workspace_audit_log(
db,
workspace=workspace,
action="api_token.created",
entity_type="api_token",
entity_id=created.token.id,
entity_name=created.token.name,
details={
"scopes": sorted(created.token.scopes.split(",")),
"key_prefix": created.token.key_prefix,
},
auth_context=_auth,
request=request,
commit=True,
)
return APITokenCreateResponse(
token=created.secret,
metadata=APITokenResponse(**token_to_dict(created.token)),
@@ -91,13 +110,28 @@ async def create_public_api_token(
@router.delete("/{token_id}", status_code=status.HTTP_200_OK)
async def revoke_public_api_token(
token_id: int,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
):
"""Revoke a scoped API token."""
workspace = assign_default_workspace_to_unscoped_rows(db)
token = db.query(APIToken).filter(APIToken.id == token_id).first()
if not revoke_api_token(db, token_id):
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="API token not found",
)
record_workspace_audit_log(
db,
workspace=workspace,
action="api_token.revoked",
entity_type="api_token",
entity_id=token_id,
entity_name=token.name if token else None,
details={"key_prefix": token.key_prefix if token else None},
auth_context=_auth,
request=request,
commit=True,
)
return {"revoked": True}
+61
View File
@@ -0,0 +1,61 @@
"""Workspace RBAC and audit endpoints."""
from typing import Any, Dict, List, Optional
from fastapi import APIRouter, Depends, Query
from pydantic import BaseModel
from sqlalchemy.orm import Session
from app.core.database import get_db
from app.core.security import require_admin_auth
from app.services.workspace_access import (
PERMISSION_AUDIT_READ,
list_workspace_roles,
require_workspace_permission,
)
from app.services.workspace_audit import list_workspace_audit_logs
from app.services.workspaces import assign_default_workspace_to_unscoped_rows
router = APIRouter()
class WorkspaceRoleResponse(BaseModel):
"""Workspace role and permission definitions."""
roles: List[Dict[str, Any]]
class WorkspaceAuditLogResponse(BaseModel):
"""Workspace audit log list response."""
audit: List[Dict[str, Any]]
@router.get("/roles", response_model=WorkspaceRoleResponse)
async def get_workspace_roles(
_auth: dict = Depends(require_admin_auth),
) -> WorkspaceRoleResponse:
"""Return the supported workspace role definitions."""
return {"roles": list_workspace_roles()}
@router.get("/logs", response_model=WorkspaceAuditLogResponse)
async def get_workspace_audit_logs(
limit: int = Query(50, ge=1, le=200),
action: Optional[str] = None,
entity_type: Optional[str] = None,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> WorkspaceAuditLogResponse:
"""Return recent sanitized audit events for the default workspace."""
require_workspace_permission(_auth, PERMISSION_AUDIT_READ)
workspace = assign_default_workspace_to_unscoped_rows(db)
return {
"audit": list_workspace_audit_logs(
db,
workspace=workspace,
limit=limit,
action=action,
entity_type=entity_type,
)
}
+30 -1
View File
@@ -6,7 +6,7 @@ import logging
from datetime import date, datetime, timezone
from typing import Any, Dict, List, Optional
from fastapi import APIRouter, Depends, HTTPException, Path, Query, status
from fastapi import APIRouter, Depends, HTTPException, Path, Query, Request, status
from fastapi.responses import Response
from pydantic import BaseModel, Field
from sqlalchemy.exc import IntegrityError
@@ -36,6 +36,7 @@ from app.services.report_persistence import (
hydrate_report_store_from_db,
)
from app.services.report_store import ReportStore
from app.services.workspace_audit import record_workspace_audit_log
from app.services.workspaces import (
assign_default_workspace_to_unscoped_rows,
workspace_domain_query,
@@ -1930,8 +1931,10 @@ async def get_domain_selectors(
@router.post("/{domain_id}/selectors", status_code=status.HTTP_201_CREATED)
async def add_domain_selector(
selector_data: SelectorRequest,
request: Request,
domain_id: str = Path(..., title="The domain ID or name"),
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
):
"""Add a DKIM selector to the manual list for a domain.
@@ -1965,15 +1968,29 @@ async def add_domain_selector(
existing.append(selector)
domain_db.dkim_selectors = ",".join(existing)
db.commit()
record_workspace_audit_log(
db,
workspace=workspace,
action="domain.selector_added",
entity_type="domain",
entity_id=domain_db.id,
entity_name=domain_db.name,
details={"selector": selector},
auth_context=_auth,
request=request,
commit=True,
)
return {"selectors": existing}
@router.delete("/{domain_id}/selectors/{selector}", status_code=status.HTTP_200_OK)
async def delete_domain_selector(
request: Request,
domain_id: str = Path(..., title="The domain ID or name"),
selector: str = Path(..., title="The DKIM selector to remove"),
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
):
"""Remove a manually configured DKIM selector from a domain."""
workspace = assign_default_workspace_to_unscoped_rows(db)
@@ -1994,6 +2011,18 @@ async def delete_domain_selector(
existing.remove(selector)
domain_db.dkim_selectors = ",".join(existing)
db.commit()
record_workspace_audit_log(
db,
workspace=workspace,
action="domain.selector_removed",
entity_type="domain",
entity_id=domain_db.id,
entity_name=domain_db.name,
details={"selector": selector},
auth_context=_auth,
request=request,
commit=True,
)
return {"selectors": existing}
+153 -19
View File
@@ -25,6 +25,11 @@ from app.services.gmail_client import GmailClient
from app.services.imap_client import IMAPClient
from app.services.import_history import record_import_attempt
from app.services.microsoft_graph_client import MicrosoftGraphClient
from app.services.workspace_audit import changed_fields, record_workspace_audit_log
from app.services.workspaces import (
assign_default_workspace_to_unscoped_rows,
workspace_mail_source_query,
)
router = APIRouter()
logger = logging.getLogger(__name__)
@@ -331,8 +336,11 @@ def _connection_test_response(
}
def _get_source_or_404(source_id: int, db: Session) -> MailSource:
source = db.query(MailSource).filter(MailSource.id == source_id).first()
def _get_source_or_404(source_id: int, db: Session, workspace=None) -> MailSource:
query = db.query(MailSource).filter(MailSource.id == source_id)
if workspace is not None:
query = query.filter(MailSource.workspace_id == workspace.id)
source = query.first()
if source is None:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
@@ -341,6 +349,30 @@ def _get_source_or_404(source_id: int, db: Session) -> MailSource:
return source
def _audit_mail_source_change(
db: Session,
*,
workspace,
source: MailSource,
action: str,
auth_context: Dict[str, Any],
request: Request,
details: Optional[Dict[str, Any]] = None,
) -> None:
record_workspace_audit_log(
db,
workspace=workspace,
action=action,
entity_type="mail_source",
entity_id=source.id,
entity_name=source.name,
details=details or {"method": source.method},
auth_context=auth_context,
request=request,
commit=True,
)
def _safe_attr(source: MailSource, name: str, default: Any = None) -> Any:
"""Read optional source attributes without letting test doubles invent fields."""
value = getattr(source, name, default)
@@ -572,18 +604,22 @@ async def list_mail_sources(
_auth: dict = Depends(require_admin_auth),
) -> List[MailSourceResponse]:
"""Return all configured mail sources (passwords redacted)."""
sources = db.query(MailSource).order_by(MailSource.id).all()
workspace = assign_default_workspace_to_unscoped_rows(db)
sources = workspace_mail_source_query(db, workspace).order_by(MailSource.id).all()
return [_source_to_response(s) for s in sources]
@router.post("", response_model=MailSourceResponse, status_code=status.HTTP_201_CREATED)
async def create_mail_source(
payload: MailSourceCreate,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> MailSourceResponse:
"""Create a new mail source."""
workspace = assign_default_workspace_to_unscoped_rows(db)
source = MailSource(
workspace_id=workspace.id,
name=payload.name,
method=payload.method.upper(),
server=payload.server,
@@ -605,6 +641,15 @@ async def create_mail_source(
db.add(source)
db.commit()
db.refresh(source)
_audit_mail_source_change(
db,
workspace=workspace,
source=source,
action="mail_source.created",
auth_context=_auth,
request=request,
details={"method": source.method, "enabled": source.enabled},
)
logger.info(
"Created mail source id=%d name=%r method=%r", source.id, source.name, source.method
)
@@ -618,7 +663,8 @@ async def get_mail_source(
_auth: dict = Depends(require_admin_auth),
) -> MailSourceResponse:
"""Return a single mail source by ID (password redacted)."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
return _source_to_response(source)
@@ -630,7 +676,8 @@ async def list_mail_source_imports(
_auth: dict = Depends(require_admin_auth),
) -> List[MailSourceImportResponse]:
"""Return recent sanitized import attempts for one mail source."""
_get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
_get_source_or_404(source_id, db, workspace)
safe_limit = min(max(limit, 1), 100)
rows = (
db.query(MailSourceImport)
@@ -653,7 +700,8 @@ async def fetch_mail_source(
if days < 1 or days > 365:
raise HTTPException(status_code=400, detail="Days parameter must be between 1 and 365")
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
results = _fetch_source(source, db, days)
logger.info(
"Manual fetch for source id=%d: processed=%d reports_found=%d "
@@ -678,11 +726,13 @@ async def fetch_mail_source(
async def update_mail_source(
source_id: int,
payload: MailSourceUpdate,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> MailSourceResponse:
"""Update one or more fields of an existing mail source."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
update_data = payload.model_dump(exclude_unset=True)
if "method" in update_data and update_data["method"]:
@@ -694,6 +744,15 @@ async def update_mail_source(
source.updated_at = datetime.utcnow()
db.commit()
db.refresh(source)
_audit_mail_source_change(
db,
workspace=workspace,
source=source,
action="mail_source.updated",
auth_context=_auth,
request=request,
details={"changed_fields": changed_fields(update_data), "method": source.method},
)
logger.info("Updated mail source id=%d", source.id)
return _source_to_response(source)
@@ -701,28 +760,55 @@ async def update_mail_source(
@router.delete("/{source_id}", status_code=status.HTTP_204_NO_CONTENT)
async def delete_mail_source(
source_id: int,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> None:
"""Delete a mail source permanently."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
source_name = source.name
source_method = source.method
db.delete(source)
db.commit()
record_workspace_audit_log(
db,
workspace=workspace,
action="mail_source.deleted",
entity_type="mail_source",
entity_id=source_id,
entity_name=source_name,
details={"method": source_method},
auth_context=_auth,
request=request,
commit=True,
)
logger.info("Deleted mail source id=%s", _sanitize_for_log(source_id))
@router.post("/{source_id}/toggle", response_model=MailSourceResponse)
async def toggle_mail_source(
source_id: int,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> MailSourceResponse:
"""Toggle the *enabled* flag of a mail source."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
source.enabled = not source.enabled
source.updated_at = datetime.utcnow()
db.commit()
db.refresh(source)
_audit_mail_source_change(
db,
workspace=workspace,
source=source,
action="mail_source.toggled",
auth_context=_auth,
request=request,
details={"enabled": source.enabled},
)
return _source_to_response(source)
@@ -733,7 +819,8 @@ async def test_stored_mail_source( # noqa: C901
_auth: dict = Depends(require_admin_auth),
) -> Dict[str, Any]:
"""Test the connection for an already-stored mail source using its saved credentials."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
if source.method == "GMAIL_API":
if not source.gmail_access_token:
@@ -877,7 +964,8 @@ async def m365_authorize_url(
_auth: dict = Depends(require_admin_auth),
) -> Dict[str, Any]:
"""Return a Microsoft identity platform authorization URL for M365_GRAPH."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
if source.method != "M365_GRAPH":
raise HTTPException(
@@ -908,7 +996,8 @@ async def m365_list_folders(
_auth: dict = Depends(require_admin_auth),
) -> Dict[str, Any]:
"""Return selectable Microsoft 365 mail folders for this source."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
if source.method != "M365_GRAPH":
raise HTTPException(
@@ -1051,11 +1140,13 @@ async def m365_oauth_callback(
async def m365_oauth_callback_post(
source_id: int,
payload: M365CallbackRequest,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> MailSourceResponse:
"""Exchange a Microsoft OAuth2 authorization code for Graph tokens."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
if source.method != "M365_GRAPH":
raise HTTPException(
@@ -1102,6 +1193,15 @@ async def m365_oauth_callback_post(
source.updated_at = datetime.utcnow()
db.commit()
db.refresh(source)
_audit_mail_source_change(
db,
workspace=workspace,
source=source,
action="mail_source.m365_connected",
auth_context=_auth,
request=request,
details={"account": m365_email or "unknown"},
)
logger.info(
"Microsoft 365 OAuth2 tokens saved for source id=%d (account=%s)",
@@ -1119,7 +1219,8 @@ async def m365_fetch_reports(
_auth: dict = Depends(require_admin_auth),
) -> Dict[str, Any]:
"""Manually trigger a Microsoft 365 Graph DMARC report fetch."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
if source.method != "M365_GRAPH":
raise HTTPException(
@@ -1147,11 +1248,13 @@ async def m365_fetch_reports(
@router.delete("/{source_id}/m365/connection", status_code=status.HTTP_204_NO_CONTENT)
async def m365_disconnect(
source_id: int,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> None:
"""Clear the stored Microsoft Graph OAuth2 tokens for this source."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
if source.method != "M365_GRAPH":
raise HTTPException(
@@ -1164,6 +1267,14 @@ async def m365_disconnect(
source.m365_email = None
source.updated_at = datetime.utcnow()
db.commit()
_audit_mail_source_change(
db,
workspace=workspace,
source=source,
action="mail_source.m365_disconnected",
auth_context=_auth,
request=request,
)
logger.info("Microsoft 365 tokens cleared for source id=%d", int(source_id))
@@ -1186,7 +1297,8 @@ async def gmail_authorize_url(
grants access Google redirects back to
``<origin>/mail-sources/<id>/gmail/callback`` with a ``code`` parameter.
"""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
if source.method != "GMAIL_API":
raise HTTPException(
@@ -1311,6 +1423,7 @@ async def gmail_oauth_callback(
async def gmail_oauth_callback_post(
source_id: int,
payload: GmailCallbackRequest,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> MailSourceResponse:
@@ -1321,7 +1434,8 @@ async def gmail_oauth_callback_post(
themselves and post the code here as JSON. Requires the standard
admin authentication.
"""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
if source.method != "GMAIL_API":
raise HTTPException(
@@ -1366,6 +1480,15 @@ async def gmail_oauth_callback_post(
source.updated_at = datetime.utcnow()
db.commit()
db.refresh(source)
_audit_mail_source_change(
db,
workspace=workspace,
source=source,
action="mail_source.gmail_connected",
auth_context=_auth,
request=request,
details={"account": gmail_email or "unknown"},
)
logger.info(
"Gmail OAuth2 tokens saved for source id=%d (account=%s)",
@@ -1387,7 +1510,8 @@ async def gmail_fetch_reports(
Searches Gmail for emails matching the DMARC report heuristic, ingests
any attachments not yet seen, and returns a summary.
"""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
if source.method != "GMAIL_API":
raise HTTPException(
@@ -1458,11 +1582,13 @@ async def gmail_fetch_reports(
@router.delete("/{source_id}/gmail/connection", status_code=status.HTTP_204_NO_CONTENT)
async def gmail_disconnect(
source_id: int,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> None:
"""Revoke / clear the stored Gmail OAuth2 tokens for this source."""
source = _get_source_or_404(source_id, db)
workspace = assign_default_workspace_to_unscoped_rows(db)
source = _get_source_or_404(source_id, db, workspace)
if source.method != "GMAIL_API":
raise HTTPException(
@@ -1475,4 +1601,12 @@ async def gmail_disconnect(
source.gmail_email = None
source.updated_at = datetime.utcnow()
db.commit()
_audit_mail_source_change(
db,
workspace=workspace,
source=source,
action="mail_source.gmail_disconnected",
auth_context=_auth,
request=request,
)
logger.info("Gmail tokens cleared for source id=%d", int(source_id))
+26 -1
View File
@@ -15,7 +15,7 @@ in the ``settings`` database table. Settings are organised into categories:
import logging
from typing import Any, Dict, List, Optional
from fastapi import APIRouter, Depends, HTTPException, status
from fastapi import APIRouter, Depends, HTTPException, Request, status
from pydantic import BaseModel
from sqlalchemy.orm import Session
@@ -36,6 +36,8 @@ from app.services.alert_rules import (
)
from app.services.notifications import send_notification
from app.services.summary_notifications import build_summary, send_summary_notification
from app.services.workspace_audit import record_workspace_audit_log
from app.services.workspaces import assign_default_workspace_to_unscoped_rows
router = APIRouter()
logger = logging.getLogger(__name__)
@@ -342,6 +344,7 @@ def _audit_setting_change(
old_plain: Optional[str],
new_plain: Optional[str],
auth_context: Optional[Dict[str, Any]],
request: Optional[Request] = None,
) -> None:
if not _should_audit_setting(key) or old_plain == new_plain:
return
@@ -352,6 +355,22 @@ def _audit_setting_change(
new_value=_audit_value_for_setting(key, new_plain),
auth_context=auth_context,
)
workspace = assign_default_workspace_to_unscoped_rows(db, commit=False)
record_workspace_audit_log(
db,
workspace=workspace,
action="setting.changed",
entity_type="setting",
entity_id=key,
entity_name=key,
details={
"key": key,
"old_value": _audit_value_for_setting(key, old_plain),
"new_value": _audit_value_for_setting(key, new_plain),
},
auth_context=auth_context,
request=request,
)
def _row_to_dict(row: Setting, redact_secrets: bool = True) -> Dict[str, Any]:
@@ -605,6 +624,7 @@ async def get_setting(
async def update_setting(
key: str,
payload: SettingUpdate,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> SettingResponse:
@@ -629,6 +649,7 @@ async def update_setting(
old_plain=None,
new_plain=new_plain,
auth_context=_auth,
request=request,
)
else:
# For secret keys, only update if not the redacted placeholder
@@ -644,6 +665,7 @@ async def update_setting(
old_plain=old_plain,
new_plain=new_plain,
auth_context=_auth,
request=request,
)
db.commit()
db.refresh(row)
@@ -653,6 +675,7 @@ async def update_setting(
@router.post("/bulk", response_model=List[SettingResponse])
async def bulk_update_settings(
payload: BulkSettingsUpdate,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> List[SettingResponse]:
@@ -681,6 +704,7 @@ async def bulk_update_settings(
old_plain=None,
new_plain=new_plain,
auth_context=_auth,
request=request,
)
else:
# Skip secret placeholder updates
@@ -696,6 +720,7 @@ async def bulk_update_settings(
old_plain=old_plain,
new_plain=new_plain,
auth_context=_auth,
request=request,
)
results.append(_row_to_dict(row))
db.commit()
+57 -1
View File
@@ -2,7 +2,7 @@
from typing import Any, Dict, List, Optional
from fastapi import APIRouter, Depends, HTTPException, Query, status
from fastapi import APIRouter, Depends, HTTPException, Query, Request, status
from pydantic import BaseModel
from sqlalchemy.orm import Session
@@ -18,6 +18,8 @@ from app.services.webhook_events import (
queue_test_webhook,
update_webhook_endpoint,
)
from app.services.workspace_audit import changed_fields, record_workspace_audit_log
from app.services.workspaces import assign_default_workspace_to_unscoped_rows
router = APIRouter()
@@ -121,14 +123,28 @@ async def list_webhook_endpoints(
@router.post("", response_model=WebhookEndpointResponse)
async def create_webhook(
payload: WebhookEndpointCreate,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> Dict[str, Any]:
"""Create an outbound webhook endpoint."""
workspace = assign_default_workspace_to_unscoped_rows(db)
try:
endpoint, raw_secret = create_webhook_endpoint(db, **payload.model_dump())
except ValueError as exc:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(exc)) from exc
record_workspace_audit_log(
db,
workspace=workspace,
action="webhook.created",
entity_type="webhook_endpoint",
entity_id=endpoint.id,
entity_name=endpoint.name,
details={"event_types": payload.event_types, "enabled": endpoint.enabled},
auth_context=_auth,
request=request,
commit=True,
)
body = endpoint_to_dict(endpoint)
body["secret"] = raw_secret
return body
@@ -138,10 +154,12 @@ async def create_webhook(
async def update_webhook(
endpoint_id: int,
payload: WebhookEndpointUpdate,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> Dict[str, Any]:
"""Update an outbound webhook endpoint."""
workspace = assign_default_workspace_to_unscoped_rows(db)
endpoint = db.query(WebhookEndpoint).filter(WebhookEndpoint.id == endpoint_id).first()
if endpoint is None:
raise HTTPException(
@@ -155,6 +173,18 @@ async def update_webhook(
)
except ValueError as exc:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(exc)) from exc
record_workspace_audit_log(
db,
workspace=workspace,
action="webhook.updated",
entity_type="webhook_endpoint",
entity_id=endpoint.id,
entity_name=endpoint.name,
details={"changed_fields": changed_fields(payload.model_dump(exclude_unset=True))},
auth_context=_auth,
request=request,
commit=True,
)
body = endpoint_to_dict(endpoint)
body["secret"] = raw_secret
return body
@@ -163,10 +193,12 @@ async def update_webhook(
@router.delete("/{endpoint_id}", response_model=WebhookEndpointResponse)
async def disable_webhook(
endpoint_id: int,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> Dict[str, Any]:
"""Disable a webhook endpoint without deleting delivery history."""
workspace = assign_default_workspace_to_unscoped_rows(db)
endpoint = db.query(WebhookEndpoint).filter(WebhookEndpoint.id == endpoint_id).first()
if endpoint is None:
raise HTTPException(
@@ -175,6 +207,17 @@ async def disable_webhook(
endpoint.enabled = False
db.commit()
db.refresh(endpoint)
record_workspace_audit_log(
db,
workspace=workspace,
action="webhook.disabled",
entity_type="webhook_endpoint",
entity_id=endpoint.id,
entity_name=endpoint.name,
auth_context=_auth,
request=request,
commit=True,
)
body = endpoint_to_dict(endpoint)
body["secret"] = None
return body
@@ -205,15 +248,28 @@ async def list_webhook_deliveries(
@router.post("/{endpoint_id}/test", response_model=WebhookTestResponse)
async def test_webhook(
endpoint_id: int,
request: Request,
db: Session = Depends(get_db),
_auth: dict = Depends(require_admin_auth),
) -> Dict[str, Any]:
"""Queue and immediately attempt a test delivery for an endpoint."""
workspace = assign_default_workspace_to_unscoped_rows(db)
try:
delivery = queue_test_webhook(db, endpoint_id)
except ValueError as exc:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(exc)) from exc
delivered = deliver_due_webhooks(db, endpoint_id=endpoint_id, limit=1)
record_workspace_audit_log(
db,
workspace=workspace,
action="webhook.tested",
entity_type="webhook_endpoint",
entity_id=endpoint_id,
details={"delivery_id": delivery.id},
auth_context=_auth,
request=request,
commit=True,
)
return {"delivery": delivery_to_dict(delivered[0] if delivered else delivery)}