diff --git a/.gitignore b/.gitignore index 5ca81560..486ef6cc 100644 --- a/.gitignore +++ b/.gitignore @@ -36,3 +36,6 @@ docker-compose.local.yml # SSO *.pem *.crt + +# Local Claude Code skills +.claude/commands/ diff --git a/server/mergin/app.py b/server/mergin/app.py index 77d5a5ac..146767a4 100644 --- a/server/mergin/app.py +++ b/server/mergin/app.py @@ -154,6 +154,7 @@ def create_app(public_keys: List[str] = None) -> Flask: """Factory function to create Flask app instance""" from itsdangerous import BadTimeSignature, BadSignature + from .audit import register as register_audit from .auth import auth_required, decode_token, register as register_auth from .auth.models import User from .sync.app import register as register_sync @@ -180,6 +181,9 @@ def create_app(public_keys: List[str] = None) -> Flask: csrf.init_app(app.app) login_manager.init_app(app.app) + # register audit module + register_audit(app.app) + # register auth blueprint register_auth(app.app) diff --git a/server/mergin/audit/__init__.py b/server/mergin/audit/__init__.py new file mode 100644 index 00000000..145b1f84 --- /dev/null +++ b/server/mergin/audit/__init__.py @@ -0,0 +1,5 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +from .app import emit, register diff --git a/server/mergin/audit/app.py b/server/mergin/audit/app.py new file mode 100644 index 00000000..b5c5dfbb --- /dev/null +++ b/server/mergin/audit/app.py @@ -0,0 +1,78 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +import datetime +import enum +import json +import logging + +from flask import Flask, current_app, has_app_context + +from .events import AuditEvent, EventType +from .sinks import NullSink + +logger = logging.getLogger(__name__) + + +def register(app: Flask) -> None: + """Wire the audit module into a Flask app. + + Stores the sink in app.extensions["audit"] so emit() has one consistent lookup path. + """ + app.extensions["audit"] = {"sink": NullSink()} + + +def _json_safe(metadata: dict) -> dict: + """Sanitize metadata to a JSON-safe structure using the app's own JSON + provider. + """ + + def default(o): + try: + return current_app.json.default(o) + except TypeError: + if isinstance(o, enum.Enum): + return o.value + return str(o) + + return json.loads(json.dumps(metadata, default=default)) + + +def emit( + event_type: EventType, + actor_id=None, + actor_email=None, + actor_ua=None, + actor_device=None, + actor_ip=None, + target_user_id=None, + target_project_id=None, + target_workspace_id=None, + **metadata, +) -> None: + """Emit one audit event to the configured sink. + + Set at least one of target_user_id, target_project_id, target_workspace_id to identify the target. + Extra keyword arguments become the metadata dict — sanitized to a JSON-safe + form so callers never need to serialize values (e.g. datetimes) themselves. + """ + if not has_app_context() or "audit" not in current_app.extensions: + return + event = AuditEvent( + event_type=event_type, + actor_id=actor_id, + actor_email=actor_email, + actor_ua=actor_ua, + actor_device=actor_device, + actor_ip=actor_ip, + happened_at=datetime.datetime.utcnow(), + target_user_id=target_user_id, + target_project_id=target_project_id, + target_workspace_id=target_workspace_id, + metadata=_json_safe(metadata), + ) + try: + current_app.extensions["audit"]["sink"].write(event) + except Exception: + logger.warning("Failed to emit audit event %s", event_type, exc_info=True) diff --git a/server/mergin/audit/events.py b/server/mergin/audit/events.py new file mode 100644 index 00000000..0e276497 --- /dev/null +++ b/server/mergin/audit/events.py @@ -0,0 +1,29 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +import datetime +import uuid +from dataclasses import dataclass, field +from typing import Any + +# Noun.verb dot-notation string, e.g. "user.login.succeeded". +# Each module defines its own str enum; the sink stores the raw string. +EventType = str + + +@dataclass(frozen=True) +class AuditEvent: + event_type: EventType + actor_ip: str | None + happened_at: datetime.datetime + actor_id: int | None + actor_email: str | None + actor_ua: str | None + actor_device: str | None # X-Device-Id header; set by mobile/QGIS clients + target_user_id: int | None # set when the target is a user + target_project_id: uuid.UUID | None # set when the target is a project + target_workspace_id: ( + int | None + ) # workspace the event belongs to; set for project and workspace events + metadata: dict[str, Any] = field(default_factory=dict) diff --git a/server/mergin/audit/listeners.py b/server/mergin/audit/listeners.py new file mode 100644 index 00000000..11634394 --- /dev/null +++ b/server/mergin/audit/listeners.py @@ -0,0 +1,89 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +""" +Utilities for writing SQLAlchemy-based audit listeners in any module. +""" + +import logging +from contextlib import contextmanager + +from sqlalchemy import inspect as sa_inspect +from sqlalchemy.orm import ColumnProperty +from flask import has_request_context, request, current_app +from flask_login import current_user + +from ..utils import get_ip, get_user_agent, get_device_id +from .app import emit + +logger = logging.getLogger(__name__) + + +@contextmanager +def audit_session_flags(session, **flags): + """Context manager that sets db.session.info flags for the duration of a block + and removes them in a finally clause so they never leak on exceptions.""" + session.info.update(flags) + try: + yield + finally: + for key in flags: + session.info.pop(key, None) + + +def request_context(): + """Return the three request-derived actor kwargs: user_agent, device_id, ip. + + Use **request_context() in explicit emit() calls so adding a new request + field only requires changing this one function. + """ + if not has_request_context(): + return dict(actor_ua=None, actor_device=None, actor_ip=None) + return dict( + actor_ua=get_user_agent(request), + actor_device=get_device_id(request), + actor_ip=get_ip(request), + ) + + +def actor_context(): + """Return full actor kwargs for emit() drawn from the current request context. + + Used by SQLAlchemy listeners where current_user is the actor. + """ + actor_id = None + actor_email = None + if has_request_context() and hasattr( + current_app._get_current_object(), "login_manager" + ): + try: + if current_user.is_authenticated: + actor_id = current_user.id + actor_email = current_user.email + except Exception: + pass + return dict(actor_id=actor_id, actor_email=actor_email, **request_context()) + + +def field_changes(target, skip=frozenset()): + """Return flat old_/new_ context for all changed non-skipped column fields. + + Only column attributes are included — relationships are skipped because their + history entries are ORM instances, not JSON-serializable values. + """ + mapper = sa_inspect(type(target)) + ctx = {} + for attr in sa_inspect(target).attrs: + if attr.key in skip: + continue + if not isinstance(mapper.attrs[attr.key], ColumnProperty): + continue + hist = attr.history + if hist.has_changes(): + old = hist.deleted[0] if hist.deleted else None + new = hist.added[0] if hist.added else None + if old != new: + ctx[f"old_{attr.key}"] = old + ctx[f"new_{attr.key}"] = new + return ctx diff --git a/server/mergin/audit/sinks.py b/server/mergin/audit/sinks.py new file mode 100644 index 00000000..14b85002 --- /dev/null +++ b/server/mergin/audit/sinks.py @@ -0,0 +1,21 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +from abc import ABC, abstractmethod + +from .events import AuditEvent + + +class AbstractSink(ABC): + """Interface all audit sinks must implement.""" + + @abstractmethod + def write(self, event: AuditEvent) -> None: ... + + +class NullSink(AbstractSink): + """Default sink — discards all events.""" + + def write(self, event: AuditEvent) -> None: + pass diff --git a/server/mergin/auth/app.py b/server/mergin/auth/app.py index 09bab97e..638ec2ff 100644 --- a/server/mergin/auth/app.py +++ b/server/mergin/auth/app.py @@ -13,7 +13,11 @@ from .commands import add_commands from .config import Configuration +from .listeners import register_listeners from .models import User, _check_dummy_password +from ..audit import emit +from ..audit.listeners import actor_context, request_context +from .events import AuthEventType # signal for other versions to listen to user_account_closed = signal("user_account_closed") @@ -39,6 +43,7 @@ def register(app): app.blueprints["/"].name = "auth" app.blueprints["auth"] = app.blueprints.pop("/") add_commands(app) + register_listeners() _permissions = {} @@ -125,6 +130,13 @@ def authenticate(login, password): db.session.commit() if duration is not None: send_account_locked_email(current_app, user, duration) + emit( + AuthEventType.USER_LOCKED, + **request_context(), + target_user_id=user.id, + target_email=user.email, + locked_until=user.locked_until.isoformat(), + ) return None diff --git a/server/mergin/auth/commands.py b/server/mergin/auth/commands.py index af6ed6af..d2ae9844 100644 --- a/server/mergin/auth/commands.py +++ b/server/mergin/auth/commands.py @@ -37,6 +37,7 @@ def create(username, password, is_admin, email): # pylint: disable=W0612 user = User(username=username, passwd=password, is_admin=is_admin, email=email) user.active = True + db.session.info["audit_user_creation_source"] = "cli" db.session.add(user) db.session.commit() click.secho("User created", fg="green") diff --git a/server/mergin/auth/controller.py b/server/mergin/auth/controller.py index 66c9b258..9994a65a 100644 --- a/server/mergin/auth/controller.py +++ b/server/mergin/auth/controller.py @@ -38,6 +38,9 @@ ApiLoginForm, ) from ..app import db +from ..audit import emit +from ..audit.listeners import actor_context, request_context +from .events import AuthEventType from ..sync.models import Project from ..sync.utils import files_size @@ -159,8 +162,26 @@ def login_public(): # noqa: E501 data = user_profile(user) data["session"] = {"token": token, "expire": expire} LoginHistory.add_record(user.id, request) + emit( + AuthEventType.USER_LOGIN_SUCCEEDED, + actor_id=user.id, + actor_email=user.email, + **request_context(), + target_user_id=user.id, + login_method="password", + token_expires_at=expire.isoformat(), + ) return data else: + audit_user = user or User.get_by_login(form.login.data) + emit( + AuthEventType.USER_LOGIN_FAILED, + **request_context(), + target_user_id=audit_user.id if audit_user else None, + login=form.login.data, + reason="account_inactive" if user else "invalid_credentials", + login_method="password", + ) abort(401, "Invalid username or password") abort(400, _extract_first_error(form.errors)) @@ -172,6 +193,15 @@ def close_user_account(): shared projects as well clean references to created projects. """ current_user.inactivate() + emit( + AuthEventType.USER_MARKED_FOR_DELETION, + **actor_context(), + target_user_id=current_user.id, + target_email=current_user.email, + scheduled_for_deletion_at=( + current_user.removal_at.isoformat() if current_user.removal_at else None + ), + ) # emit signal to be caught elsewhere user_account_closed.send(current_user) return NoContent, 204 @@ -228,8 +258,25 @@ def login(): # pylint: disable=W0613,W0612 login_user(user) if not os.path.isfile(current_app.config["MAINTENANCE_FILE"]): LoginHistory.add_record(user.id, request) + emit( + AuthEventType.USER_LOGIN_SUCCEEDED, + actor_id=user.id, + actor_email=user.email, + **request_context(), + target_user_id=user.id, + login_method="password", + ) return "", 200 else: + audit_user = user or User.get_by_login(form.login.data) + emit( + AuthEventType.USER_LOGIN_FAILED, + **request_context(), + target_user_id=audit_user.id if audit_user else None, + login=form.login.data, + reason="account_inactive" if user else "invalid_credentials", + login_method="password", + ) abort(401, "Invalid username or password") return jsonify(form.errors), 401 @@ -244,11 +291,36 @@ def admin_login(): # pylint: disable=W0613,W0612 if user: if user.active and user.is_admin: login_user(user) + emit( + AuthEventType.USER_LOGIN_SUCCEEDED, + actor_id=user.id, + actor_email=user.email, + **request_context(), + target_user_id=user.id, + login_method="password", + ) LoginHistory.add_record(user.id, request) return "", 200 else: + emit( + AuthEventType.USER_LOGIN_FAILED, + **request_context(), + target_user_id=user.id, + login=form.login.data, + reason="insufficient_permissions", + login_method="password", + ) abort(403, "You do not have permissions") else: + audit_user = User.get_by_login(form.login.data) + emit( + AuthEventType.USER_LOGIN_FAILED, + **request_context(), + target_user_id=audit_user.id if audit_user else None, + login=form.login.data, + reason="invalid_credentials", + login_method="password", + ) abort(401, "Invalid username or password") @@ -269,6 +341,12 @@ def change_password(): # pylint: disable=W0613,W0612 current_user.assign_password(form.password.data) db.session.add(current_user) db.session.commit() + emit( + AuthEventType.USER_PASSWORD_CHANGED, + **actor_context(), + target_user_id=current_user.id, + target_email=current_user.email, + ) return "", 200 return jsonify(form.errors), 400 @@ -295,6 +373,12 @@ def password_reset(): # pylint: disable=W0613,W0612 user = User.query.filter( func.lower(User.email) == func.lower(form.email.data.strip()) ).one_or_none() + emit( + AuthEventType.USER_PASSWORD_RESET_REQUESTED, + **request_context(), + target_user_id=user.id if user else None, + target_email=form.email.data.strip(), + ) if user and user.active and user.can_edit_profile: send_confirmation_email( current_app, @@ -309,12 +393,42 @@ def password_reset(): # pylint: disable=W0613,W0612 def confirm_new_password(token): # pylint: disable=W0613,W0612 email = confirm_token(token, salt=current_app.config["SECURITY_PASSWORD_SALT"]) if not email: + emit( + AuthEventType.USER_PASSWORD_RESET_FAILED, + **request_context(), + target_user_id=None, + target_email=None, + reason="invalid_token", + ) abort(400, "Invalid token") - user = User.query.filter_by(email=email).first_or_404() + user = User.query.filter_by(email=email).first() + if not user: + emit( + AuthEventType.USER_PASSWORD_RESET_FAILED, + **request_context(), + target_user_id=None, + target_email=email, + reason="user_not_found", + ) + abort(404) if not user.active: + emit( + AuthEventType.USER_PASSWORD_RESET_FAILED, + **request_context(), + target_user_id=user.id, + target_email=user.email, + reason="account_inactive", + ) abort(400, "Account is not active") if not user.can_edit_profile: + emit( + AuthEventType.USER_PASSWORD_RESET_FAILED, + **request_context(), + target_user_id=user.id, + target_email=user.email, + reason="profile_edit_disabled", + ) abort(403, CANNOT_EDIT_PROFILE_MSG) form = UserPasswordForm.from_json(request.json) @@ -323,6 +437,12 @@ def confirm_new_password(token): # pylint: disable=W0613,W0612 user.reset_lockout() db.session.add(user) db.session.commit() + emit( + AuthEventType.USER_PASSWORD_RESET_COMPLETED, + **request_context(), + target_user_id=user.id, + target_email=user.email, + ) return "", 200 return jsonify(form.errors), 400 @@ -366,6 +486,12 @@ def unlock_account(token: str): # pylint: disable=W0613,W0612 user.reset_lockout() db.session.commit() + emit( + AuthEventType.USER_UNLOCKED, + **request_context(), + target_user_id=user.id, + target_email=user.email, + ) return "", 200 @@ -412,6 +538,7 @@ def register_user(): # pylint: disable=W0613,W0612 form = UserRegistrationForm() form.username.data = User.generate_username(form.email.data) if form.is_submitted() and form.validate(): + db.session.info["audit_user_creation_source"] = "admin" user = User.create(form.username.data, form.email.data, form.password.data) user_created.send(user, source="admin") token = generate_confirmation_token( @@ -457,13 +584,27 @@ def update_user(username): # pylint: disable=W0613,W0612 abort(400, "Unable to assign super admin role") user = User.query.filter_by(username=username).first_or_404("User not found") + old_active = user.active form.update_obj(user) - - # remove inactive since flag for ban or re-activation user.inactive_since = None - db.session.add(user) db.session.commit() + + if old_active and not user.active: + emit( + AuthEventType.USER_DEACTIVATED, + **actor_context(), + target_user_id=user.id, + target_email=user.email, + ) + elif not old_active and user.active: + emit( + AuthEventType.USER_RESTORED, + **actor_context(), + target_user_id=user.id, + target_email=user.email, + ) + return jsonify(UserSchema().dump(user)) @@ -471,8 +612,22 @@ def update_user(username): # pylint: disable=W0613,W0612 def delete_user(username): # pylint: disable=W0613,W0612 user = User.query.filter_by(username=username).first_or_404("User not found") user.inactivate() + emit( + AuthEventType.USER_MARKED_FOR_DELETION, + **actor_context(), + target_user_id=user.id, + target_email=user.email, + scheduled_for_deletion_at=( + user.removal_at.isoformat() if user.removal_at else None + ), + ) user_account_closed.send(user) - # force 'delete' user + emit( + AuthEventType.USER_DELETED, + **actor_context(), + target_user_id=user.id, + target_email=user.email, + ) user.anonymize() return "", 204 @@ -565,6 +720,7 @@ def create_user(): if not form.validate(): return jsonify(form.errors), 400 + db.session.info["audit_user_creation_source"] = "api" user = User.create( form.username.data, form.email.data, diff --git a/server/mergin/auth/events.py b/server/mergin/auth/events.py new file mode 100644 index 00000000..86f23821 --- /dev/null +++ b/server/mergin/auth/events.py @@ -0,0 +1,33 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +from enum import Enum + + +class AuthEventType(str, Enum): + # authentication events + USER_LOGIN_SUCCEEDED = "user.login.succeeded" + USER_LOGIN_FAILED = "user.login.failed" + USER_PASSWORD_CHANGED = "user.password.changed" + # password reset flow (three separate events; all unauthenticated) + USER_PASSWORD_RESET_REQUESTED = "user.password.reset_requested" + USER_PASSWORD_RESET_COMPLETED = "user.password.reset_completed" + USER_PASSWORD_RESET_FAILED = "user.password.reset_failed" + # general CRUD (SQLAlchemy listeners) + USER_CREATED = "user.created" + USER_UPDATED = "user.updated" + # lifecycle events (explicit emit only; active/inactive_since excluded from user.updated) + USER_MARKED_FOR_DELETION = ( + "user.marked_for_deletion" # user or admin triggers deletion flow + ) + USER_DEACTIVATED = "user.deactivated" # admin sets active=False without deletion + USER_RESTORED = ( + "user.restored" # admin re-activates after deactivation or marked_for_deletion + ) + USER_DELETED = "user.deleted" # personal data permanently erased + # lockout events (explicit emit) + USER_LOCKED = "user.locked" # account locked after too many failed logins + USER_UNLOCKED = "user.unlocked" # self-service token-based unlock + # admin panel access (SQLAlchemy listener + explicit on CLI create) + USER_ADMIN_PANEL_ACCESS_CHANGED = "user.admin_panel_access.changed" diff --git a/server/mergin/auth/listeners.py b/server/mergin/auth/listeners.py new file mode 100644 index 00000000..88de3765 --- /dev/null +++ b/server/mergin/auth/listeners.py @@ -0,0 +1,98 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +from typing import Any + +from sqlalchemy import event, inspect as sa_inspect +from sqlalchemy.orm import object_session + +from ..audit import emit +from ..audit.listeners import actor_context, field_changes +from .events import AuthEventType +from .models import User + +# Fields excluded from user.updated audit events: +# - sensitive values that must never appear in logs (passwd) +# - high-frequency operational fields (last_signed_in, registration_date) +# - lifecycle state fields covered by dedicated events (active, inactive_since) +# - is_admin covered by the dedicated user.admin_panel_access.changed event +# frozenset prevents accidental mutation of module-level state. +_EXCLUDED_FROM_USER_UPDATED = frozenset( + { + "passwd", + "last_signed_in", + "registration_date", + "active", + "inactive_since", + "is_admin", + } +) + + +def _on_user_created(_mapper: Any, _connection: Any, target: User) -> None: + """Emit USER_CREATED after a User row is inserted. + + Actor attribution: + - Authenticated request (admin creates user): actor comes from actor_context(). + - Unauthenticated request (self-registration): no session yet, so the new user + themselves is used as the actor rather than leaving it null. + - No request context (Celery/system import): actor stays null to signal a system action. + """ + ctx = actor_context() + session = object_session(target) + source = session.info.get("audit_user_creation_source") if session else None + # For self-registration there is no prior actor — attribute the creation to the new user. + if not ctx.get("actor_email") and source == "self_registration": + ctx["actor_email"] = target.email + ctx["actor_id"] = target.id + emit( + AuthEventType.USER_CREATED, + **ctx, + target_user_id=target.id, + target_email=target.email, + target_username=target.username, + source=source, + ) + if target.is_admin: + emit( + AuthEventType.USER_ADMIN_PANEL_ACCESS_CHANGED, + **ctx, + target_user_id=target.id, + old_is_admin=None, + new_is_admin=True, + ) + + +def _on_user_updated(_mapper, _connection, target): + if object_session(target).info.get("audit_skip_user_update"): + return + changes = field_changes(target, _EXCLUDED_FROM_USER_UPDATED) + if changes: + emit( + AuthEventType.USER_UPDATED, + **actor_context(), + target_user_id=target.id, + target_email=target.email, + **changes, + ) + # is_admin is excluded from field_changes; handle it with a dedicated event. + is_admin_hist = sa_inspect(target).attrs.is_admin.history + if is_admin_hist.has_changes(): + old = is_admin_hist.deleted[0] if is_admin_hist.deleted else None + new = is_admin_hist.added[0] if is_admin_hist.added else None + if old != new: + emit( + AuthEventType.USER_ADMIN_PANEL_ACCESS_CHANGED, + **actor_context(), + target_user_id=target.id, + old_is_admin=old, + new_is_admin=new, + ) + + +def register_listeners(): + if event.contains(User, "after_insert", _on_user_created): + return + event.listen(User, "after_insert", _on_user_created) + event.listen(User, "after_update", _on_user_updated) diff --git a/server/mergin/auth/models.py b/server/mergin/auth/models.py index f80ccbe5..7ae80513 100644 --- a/server/mergin/auth/models.py +++ b/server/mergin/auth/models.py @@ -11,8 +11,10 @@ from sqlalchemy import or_, func, text from ..app import db +from ..audit.listeners import audit_session_flags from ..sync.models import ProjectUser -from ..sync.utils import get_user_agent, get_ip, get_device_id, is_reserved_word +from ..sync.utils import is_reserved_word +from ..utils import get_ip, get_user_agent, get_device_id MAX_USERNAME_LENGTH = 50 @@ -233,8 +235,10 @@ def inactivate(self) -> None: """ from ..sync.models import AccessRequest, RequestStatus - # remove explicit permissions - ProjectUser.query.filter(ProjectUser.user_id == self.id).delete() + db.session.info["project_member_delete_reason"] = "user_deleted" + # Remove explicit project permissions; PROJECT_MEMBER_DELETED events fire automatically via the after_delete listener. + for m in ProjectUser.query.filter(ProjectUser.user_id == self.id).all(): + db.session.delete(m) # decline all access requests for req in ( @@ -248,18 +252,21 @@ def inactivate(self) -> None: self.active = False self.inactive_since = datetime.datetime.utcnow() db.session.commit() + db.session.info.pop("project_member_delete_reason", None) def anonymize(self): """Anonymize user object in database - remove personal information""" ts = round(datetime.datetime.utcnow().timestamp() * 1000) del_str = f"deleted_{ts}" - self.active = False - self.username = del_str - self.email = None - self.passwd = None - self.first_name = None - self.last_name = None - db.session.commit() + # Suppress user.updated — these changes are covered by the USER_DELETED event. + with audit_session_flags(db.session, audit_skip_user_update=True): + self.active = False + self.username = del_str + self.email = None + self.passwd = None + self.first_name = None + self.last_name = None + db.session.commit() @classmethod def get_by_login(cls, login: str) -> Optional[User]: diff --git a/server/mergin/auth/tasks.py b/server/mergin/auth/tasks.py index 3c408d35..5e8d8737 100644 --- a/server/mergin/auth/tasks.py +++ b/server/mergin/auth/tasks.py @@ -7,6 +7,9 @@ from ..celery import celery from ..app import db +from ..audit import emit +from .app import user_account_closed +from .events import AuthEventType from .models import User from .config import Configuration @@ -22,4 +25,12 @@ def anonymize_removed_users(): User.username.op("~")("^(?!deleted_\d{13})"), ).all() for user in users: + emit( + AuthEventType.USER_DELETED, + target_user_id=user.id, + target_email=user.email, + ) + # Remove project/workspace memberships + user.inactivate() + user_account_closed.send(user) user.anonymize() diff --git a/server/mergin/sync/app.py b/server/mergin/sync/app.py index e97f7cc3..1b6a923d 100644 --- a/server/mergin/sync/app.py +++ b/server/mergin/sync/app.py @@ -7,6 +7,7 @@ from .commands import add_commands from .config import Configuration from .db_events import register_events +from .listeners import register_listeners def register(app: Flask): @@ -40,3 +41,4 @@ def register(app: Flask): add_commands(app) register_events() + register_listeners() diff --git a/server/mergin/sync/db_events.py b/server/mergin/sync/db_events.py index 48a1756d..03e0ff0b 100644 --- a/server/mergin/sync/db_events.py +++ b/server/mergin/sync/db_events.py @@ -6,9 +6,12 @@ from flask import current_app, abort from sqlalchemy import event -from .models import ProjectVersion +from .events import SyncEventType +from .models import ProjectUser, ProjectVersion from .tasks import optimize_storage from ..app import db +from ..audit import emit +from ..audit.listeners import actor_context def check(session): @@ -22,11 +25,38 @@ def optimize_gpgk_storage(mapper, connection, project_version): optimize_storage.delay(project_version.project_id) +def on_project_member_deleted(mapper, connection, project_user: ProjectUser): + """Emit PROJECT_MEMBER_DELETED whenever a ProjectUser row is deleted via the ORM.""" + if db.session.info.get("suppress_project_member_deleted"): + return + + from .models import Project + + project = db.session.get(Project, project_user.project_id) + workspace = project.workspace if project else None + ws_name = workspace.name if workspace else None + reason = db.session.info.get("project_member_delete_reason", "removed") + emit( + SyncEventType.PROJECT_MEMBER_DELETED, + **actor_context(), + target_project_id=project_user.project_id, + target_workspace_id=project.workspace_id if project else None, + target_user_id=project_user.user_id, + target_email=project_user.user.email if project_user.user else None, + workspace_name=ws_name, + project_name=project.name if project else None, + role=project_user.role, + reason=reason, + ) + + def register_events(): event.listen(db.session, "before_commit", check) event.listen(ProjectVersion, "after_insert", optimize_gpgk_storage) + event.listen(ProjectUser, "after_delete", on_project_member_deleted) def remove_events(): event.remove(db.session, "before_commit", check) - event.listen(ProjectVersion, "after_insert", optimize_gpgk_storage) + event.remove(ProjectVersion, "after_insert", optimize_gpgk_storage) + event.remove(ProjectUser, "after_delete", on_project_member_deleted) diff --git a/server/mergin/sync/events.py b/server/mergin/sync/events.py new file mode 100644 index 00000000..37450284 --- /dev/null +++ b/server/mergin/sync/events.py @@ -0,0 +1,32 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +from enum import Enum + + +class SyncEventType(str, Enum): + # automatic CRUD events (SQLAlchemy listeners) + PROJECT_CREATED = "project.created" # also emitted on clone + PROJECT_UPDATED = "project.updated" + # transfer request events + PROJECT_TRANSFER_REQUEST_INITIATED = "project.transfer_request.initiated" + PROJECT_TRANSFER_REQUEST_RECEIVED = "project.transfer_request.received" + PROJECT_TRANSFER_REQUEST_COMPLETED = "project.transfer_request.completed" + PROJECT_TRANSFER_REQUEST_ACCEPTED = "project.transfer_request.accepted" + PROJECT_TRANSFER_REQUEST_CANCELED = "project.transfer_request.canceled" + PROJECT_TRANSFER_REQUEST_REJECTED = "project.transfer_request.rejected" + # lifecycle events (explicit emit) + PROJECT_MARKED_FOR_DELETION = "project.marked_for_deletion" + PROJECT_RESTORED = "project.restored" + PROJECT_DELETED = "project.deleted" + # membership events (explicit emit) + PROJECT_MEMBER_ADDED = "project.member.added" + PROJECT_MEMBER_UPDATED = "project.member.updated" + PROJECT_MEMBER_DELETED = "project.member.deleted" + # access request events (explicit emit) + PROJECT_ACCESS_REQUEST_INITIATED = "project.access_request.initiated" + PROJECT_ACCESS_REQUEST_ACCEPTED = "project.access_request.accepted" + PROJECT_ACCESS_REQUEST_CANCELED = "project.access_request.canceled" + # data events (explicit emit) + PROJECT_VERSION_CREATED = "project.version.created" diff --git a/server/mergin/sync/listeners.py b/server/mergin/sync/listeners.py new file mode 100644 index 00000000..6be11011 --- /dev/null +++ b/server/mergin/sync/listeners.py @@ -0,0 +1,65 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +from sqlalchemy import event +from sqlalchemy.orm import object_session + +from ..audit import emit +from ..audit.listeners import actor_context, field_changes +from .events import SyncEventType +from .models import Project + +# Fields excluded from project.updated audit events — either auto-computed on every +# version push (disk_usage, latest_version, tags), operational metadata (updated, +# storage_params, locked_until), or covered by dedicated events with their own emit +# (removed_at/removed_by are suppressed via audit_skip_project_update instead). +# frozenset prevents accidental mutation of module-level state. +_EXCLUDED_FROM_PROJECT_UPDATED = frozenset( + { + "disk_usage", + "latest_version", + "updated", + "storage_params", + "tags", + } +) + + +def _on_project_created(_mapper, _connection, target): + if object_session(target).info.get("audit_skip_project_create"): + return + emit( + SyncEventType.PROJECT_CREATED, + **actor_context(), + target_project_id=target.id, + target_workspace_id=target.workspace_id, + project_name=target.name, + workspace_name=target.workspace.name, + is_public=target.public, + creator=target.creator_id, + ) + + +def _on_project_updated(_mapper, _connection, target): + if object_session(target).info.get("audit_skip_project_update"): + return + changes = field_changes(target, _EXCLUDED_FROM_PROJECT_UPDATED) + if not changes: + return + emit( + SyncEventType.PROJECT_UPDATED, + **actor_context(), + target_project_id=target.id, + target_workspace_id=target.workspace_id, + project_name=target.name, + workspace_name=target.workspace.name, + **changes, + ) + + +def register_listeners(): + if event.contains(Project, "after_insert", _on_project_created): + return + event.listen(Project, "after_insert", _on_project_created) + event.listen(Project, "after_update", _on_project_updated) diff --git a/server/mergin/sync/models.py b/server/mergin/sync/models.py index 3817acd6..733a0272 100644 --- a/server/mergin/sync/models.py +++ b/server/mergin/sync/models.py @@ -44,6 +44,9 @@ from .interfaces import WorkspaceRole from .storages.disk import copy_file, move_to_tmp from ..app import db +from ..audit import emit +from ..audit.listeners import actor_context +from .events import SyncEventType from .storages import DiskStorage from .utils import ( LOG_BASE, @@ -316,6 +319,9 @@ def delete(self, removed_by: int = None): # do nothing if the project has been already deleted if not self.storage_params: return + # Set before any mutation below + db.session.info["audit_skip_project_update"] = True + project_name = self.name self.name = f"{self.name}_{str(self.id)}" # make sure remove_at is not null as it is used as filter for APIs if not self.removed_at: @@ -343,6 +349,8 @@ def delete(self, removed_by: int = None): db.session.execute( delta_table.delete().where(delta_table.c.project_id == self.id) ) + + db.session.info["project_member_delete_reason"] = "project_deleted" self.project_users.clear() access_requests = ( AccessRequest.query.filter_by(project_id=self.id) @@ -351,7 +359,19 @@ def delete(self, removed_by: int = None): ) for req in access_requests: req.resolve(status=RequestStatus.DECLINED, resolved_by=self.removed_by) + project_delete_reason = db.session.info.get("project_delete_reason") + emit( + SyncEventType.PROJECT_DELETED, + **actor_context(), + target_project_id=self.id, + target_workspace_id=self.workspace_id, + workspace_name=self.workspace.name, + project_name=project_name, + **({"reason": project_delete_reason} if project_delete_reason else {}), + ) db.session.commit() + db.session.info.pop("project_member_delete_reason", None) + db.session.info.pop("audit_skip_project_update", None) project_deleted.send(self) def _member(self, user_id: int) -> Optional[ProjectUser]: diff --git a/server/mergin/sync/private_api_controller.py b/server/mergin/sync/private_api_controller.py index fbe7b5cf..1a1bb679 100644 --- a/server/mergin/sync/private_api_controller.py +++ b/server/mergin/sync/private_api_controller.py @@ -11,8 +11,12 @@ from sqlalchemy import text from ..app import db +from ..audit import emit +from ..audit.listeners import actor_context from ..auth import auth_required +from ..auth.models import User from .forms import AccessPermissionForm +from .events import SyncEventType from .models import ( Project, AccessRequest, @@ -62,6 +66,15 @@ def create_project_access_request(namespace, project_name): # noqa: E501 access_request = AccessRequest(project, current_user.id) db.session.add(access_request) db.session.commit() + emit( + SyncEventType.PROJECT_ACCESS_REQUEST_INITIATED, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + workspace_name=project.workspace.name, + project_name=project.name, + access_request_id=access_request.id, + ) # notify project owners owners = current_app.project_handler.get_email_receivers(project) for owner in owners: @@ -99,8 +112,19 @@ def decline_project_access_request(request_id): # noqa: E501 project_role == ProjectRole.OWNER or current_user.id == access_request.requested_by ): + requester = User.query.get(access_request.requested_by) access_request.resolve(RequestStatus.DECLINED, current_user.id) db.session.commit() + emit( + SyncEventType.PROJECT_ACCESS_REQUEST_CANCELED, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + target_email=requester.email if requester else None, + workspace_name=project.workspace.name, + project_name=project.name, + access_request_id=access_request.id, + ) return "", 200 abort(403, "You don't have permissions to remove project access request") @@ -124,7 +148,30 @@ def accept_project_access_request(request_id): project = access_request.project project_role = ProjectPermissions.get_user_project_role(project, current_user) if project_role == ProjectRole.OWNER: + requester = User.query.get(access_request.requested_by) access_request.accept(permission) + emit( + SyncEventType.PROJECT_ACCESS_REQUEST_ACCEPTED, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + target_email=requester.email if requester else None, + workspace_name=project.workspace.name, + project_name=project.name, + access_request_id=access_request.id, + target_user_id=requester.id if requester else None, + role=permission, + ) + emit( + SyncEventType.PROJECT_MEMBER_ADDED, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + target_email=requester.email if requester else None, + workspace_name=project.workspace.name, + project_name=project.name, + role=permission, + ) return "", 200 abort(403, "You don't have permissions to accept project access request") @@ -228,6 +275,14 @@ def restore_project(id): # noqa: E501 project.removed_at = None project.removed_by = None db.session.commit() + emit( + SyncEventType.PROJECT_RESTORED, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + workspace_name=project.workspace.name, + project_name=project.name, + ) return "", 201 @@ -288,7 +343,9 @@ def unsubscribe_project(id): # pylint: disable=W0612 project.unset_role(current_user.id) db.session.add(project) + db.session.info["project_member_delete_reason"] = "left" db.session.commit() + db.session.info.pop("project_member_delete_reason", None) return NoContent, 200 diff --git a/server/mergin/sync/public_api_controller.py b/server/mergin/sync/public_api_controller.py index 34a2d28f..e9b1cac2 100644 --- a/server/mergin/sync/public_api_controller.py +++ b/server/mergin/sync/public_api_controller.py @@ -35,8 +35,11 @@ from mergin.sync.forms import project_name_validation from .interfaces import WorkspaceRole from ..app import db +from ..audit import emit +from ..audit.listeners import actor_context, audit_session_flags from ..auth import auth_required from ..auth.models import User +from .events import SyncEventType from .models import ( FileSyncErrorType, FileDiff, @@ -78,16 +81,13 @@ ) from .utils import ( generate_checksum, - get_ip, - get_user_agent, generate_location, is_valid_uuid, - get_device_id, is_versioned_file, prepare_download_response, - get_device_id, wkb2wkt, ) +from ..utils import get_ip, get_user_agent, get_device_id from .errors import StorageLimitHit, ProjectLocked from ..utils import format_time_delta @@ -224,6 +224,9 @@ def add_project(namespace): # noqa: E501 template_name = request.json.get("template", None) if template_name: + # Set flag before the template query — p is already in the session via the + # workspace backref, so any query triggers autoflush and fires after_insert. + db.session.info["audit_skip_project_create"] = True template = ( Project.query.filter(Project.creator.has(username="TEMPLATES")) .filter(Project.name == template_name) @@ -243,7 +246,6 @@ def add_project(namespace): # noqa: E501 change=PushChangeType.CREATE, ) ) - else: template = None version_name = 0 @@ -267,6 +269,28 @@ def add_project(namespace): # noqa: E501 db.session.add(p) db.session.add(version) db.session.commit() + if template_name: + db.session.info.pop("audit_skip_project_create", None) + emit( + SyncEventType.PROJECT_CREATED, + **actor_context(), + target_project_id=p.id, + target_workspace_id=p.workspace_id, + project_name=p.name, + workspace_name=p.workspace.name, + is_public=p.public, + creator=p.creator_id, + created_from_template=template_name, + ) + emit( + SyncEventType.PROJECT_VERSION_CREATED, + **actor_context(), + target_project_id=p.id, + target_workspace_id=p.workspace_id, + workspace_name=p.workspace.name, + project_name=p.name, + version=ProjectVersion.to_v_name(version_name), + ) project_version_created.send(version) return NoContent, 200 @@ -285,7 +309,19 @@ def delete_project(namespace, project_name): # noqa: E501 :rtype: None """ project = require_project(namespace, project_name, ProjectPermissions.Delete) - project.schedule_deletion(removed_by=current_user.id) + with audit_session_flags(db.session, audit_skip_project_update=True): + project.schedule_deletion(removed_by=current_user.id) + emit( + SyncEventType.PROJECT_MARKED_FOR_DELETION, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + workspace_name=project.workspace.name, + project_name=project.name, + scheduled_for_deletion_at=( + project.removed_at.isoformat() if project.removed_at else None + ), + ) return NoContent, 200 @@ -984,6 +1020,15 @@ def project_push(namespace, project_name): f"A project version {ProjectVersion.to_v_name(next_version)} for project: {project.id} created. " f"Transaction id: {upload.transaction_id}. No upload." ) + emit( + SyncEventType.PROJECT_VERSION_CREATED, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + workspace_name=project.workspace.name, + project_name=project.name, + version=ProjectVersion.to_v_name(next_version), + ) project_version_created.send(pv) push_finished.send(pv) return jsonify(ProjectSchema().dump(project)), 200 @@ -1140,6 +1185,15 @@ def push_finish(transaction_id): logging.info( f"Push finished for project: {project.id}, project version: {v_next_version}, transaction id: {transaction_id}." ) + emit( + SyncEventType.PROJECT_VERSION_CREATED, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + workspace_name=project.workspace.name, + project_name=project.name, + version=v_next_version, + ) project_version_created.send(pv) push_finished.send(pv) except (psycopg2.Error, OSError, IntegrityError) as err: @@ -1257,6 +1311,7 @@ def clone_project(namespace, project_name): # noqa: E501 ) p.updated = datetime.utcnow() db.session.add(p) + db.session.info["audit_skip_project_create"] = True files_to_exclude = current_app.config.get("EXCLUDED_CLONE_FILENAMES", []) try: @@ -1297,6 +1352,28 @@ def clone_project(namespace, project_name): # noqa: E501 ) db.session.add(project_version) db.session.commit() + db.session.info.pop("audit_skip_project_create", None) + emit( + SyncEventType.PROJECT_CREATED, + **actor_context(), + target_project_id=p.id, + target_workspace_id=p.workspace_id, + project_name=p.name, + workspace_name=ws.name, + is_public=p.public, + creator=p.creator_id, + cloned_from=str(cloned_project.id), + ) + if version >= 1: + emit( + SyncEventType.PROJECT_VERSION_CREATED, + **actor_context(), + target_project_id=p.id, + target_workspace_id=p.workspace_id, + workspace_name=ws.name, + project_name=p.name, + version=ProjectVersion.to_v_name(version), + ) project_version_created.send(project_version) return NoContent, 200 diff --git a/server/mergin/sync/public_api_v2_controller.py b/server/mergin/sync/public_api_v2_controller.py index e7806865..d27fdbce 100644 --- a/server/mergin/sync/public_api_v2_controller.py +++ b/server/mergin/sync/public_api_v2_controller.py @@ -33,6 +33,9 @@ UploadError, ) from .files import ChangesSchema, DeltaChangeRespSchema, ProjectFileSchema +from .events import SyncEventType +from ..audit import emit +from ..audit.listeners import actor_context, audit_session_flags from .forms import project_name_validation from .models import ( FileDiff, @@ -58,12 +61,10 @@ from .schemas_v2 import ProjectSchema as ProjectSchemaV2 from .storages.disk import move_to_tmp, save_to_file from .utils import ( - get_device_id, - get_ip, - get_user_agent, get_chunk_location, prepare_download_response, ) +from ..utils import get_ip, get_user_agent, get_device_id from .tasks import remove_transaction_chunks from .workspace import WorkspaceRole from ..utils import parse_order_params, get_schema_fields_map @@ -76,8 +77,19 @@ def schedule_delete_project(id): rest. """ project = require_project_by_uuid(id, ProjectPermissions.Delete) - project.schedule_deletion(removed_by=current_user.id) - + with audit_session_flags(db.session, audit_skip_project_update=True): + project.schedule_deletion(removed_by=current_user.id) + emit( + SyncEventType.PROJECT_MARKED_FOR_DELETION, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + workspace_name=project.workspace.name, + project_name=project.name, + scheduled_for_deletion_at=( + project.removed_at.isoformat() if project.removed_at else None + ), + ) return NoContent, 204 @@ -150,6 +162,16 @@ def add_project_collaborator(id): project.set_role(user.id, ProjectRole(request.json["role"])) db.session.commit() + emit( + SyncEventType.PROJECT_MEMBER_ADDED, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + target_email=user.email, + workspace_name=project.workspace.name, + project_name=project.name, + role=request.json["role"], + ) data = ProjectMemberSchema().dump(project.get_member(user.id)) return data, 201 @@ -159,11 +181,23 @@ def update_project_collaborator(id, user_id): """Update project collaborator""" project = require_project_by_uuid(id, ProjectPermissions.Update) user = User.query.filter_by(id=user_id, active=True).first_or_404() - if not project.get_role(user_id): + old_role = project.get_role(user_id) + if not old_role: abort(404) project.set_role(user.id, ProjectRole(request.json["role"])) db.session.commit() + emit( + SyncEventType.PROJECT_MEMBER_UPDATED, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + target_email=user.email, + workspace_name=project.workspace.name, + project_name=project.name, + old_role=old_role.value, + new_role=request.json["role"], + ) data = ProjectMemberSchema().dump(project.get_member(user.id)) return data, 200 @@ -172,11 +206,14 @@ def update_project_collaborator(id, user_id): def remove_project_collaborator(id, user_id): """Remove project collaborator""" project = require_project_by_uuid(id, ProjectPermissions.Update) - if not project.get_role(user_id): + removed_role = project.get_role(user_id) + if not removed_role: abort(404) project.unset_role(user_id) + db.session.info["project_member_delete_reason"] = "removed" db.session.commit() + db.session.info.pop("project_member_delete_reason", None) return NoContent, 204 @@ -340,6 +377,15 @@ def create_project_version(id): os.renames(temp_files_dir, version_dir) db.session.commit() + emit( + SyncEventType.PROJECT_VERSION_CREATED, + **actor_context(), + target_project_id=project.id, + target_workspace_id=project.workspace_id, + workspace_name=project.workspace.name, + project_name=project.name, + version=v_next_version, + ) # remove used chunks only after commit — chunks belong to the now-committed version if to_be_added_files or to_be_updated_files: diff --git a/server/mergin/sync/tasks.py b/server/mergin/sync/tasks.py index 480222e6..9ebeb4c8 100644 --- a/server/mergin/sync/tasks.py +++ b/server/mergin/sync/tasks.py @@ -11,12 +11,14 @@ from zipfile import ZIP_DEFLATED, ZipFile from flask import current_app +from .events import SyncEventType from .models import Project, ProjectVersion, FileHistory from .storages.disk import move_to_tmp from .config import Configuration from .utils import get_chunk_location, remove_outdated_files from ..celery import celery from ..app import db +from ..audit import emit @celery.task diff --git a/server/mergin/sync/utils.py b/server/mergin/sync/utils.py index 6dd7abe1..e1c06678 100644 --- a/server/mergin/sync/utils.py +++ b/server/mergin/sync/utils.py @@ -103,32 +103,6 @@ def get_blacklisted_files(blacklist): return [p for p in blacklist if not p.endswith("/")] -def get_user_agent(request): - """Return user agent from request headers - - In case of browser client a parsed version from werkzeug utils is returned else raw value of header. - """ - if request.user_agent.browser and request.user_agent.platform: - client = request.user_agent.browser.capitalize() - version = request.user_agent.version - system = request.user_agent.platform.capitalize() - return f"{client}/{version} ({system})" - else: - return request.user_agent.string - - -def get_ip(request): - """Returns request's IP address based on X_FORWARDED_FOR header - from proxy webserver (which should always be the case) - """ - forwarded_ips = request.environ.get( - "HTTP_X_FORWARDED_FOR", request.environ.get("REMOTE_ADDR", "untrackable") - ) - # seems like we get list of IP addresses from AWS infra (beginning with external IP address of client, followed by some internal IP) - ip = forwarded_ips.split(",")[0] - return ip - - def generate_location(): """Return random location where project is saved on disk @@ -257,11 +231,6 @@ def split_project_path(project_path): return workspace_name, project_name -def get_device_id(request: Request) -> Optional[str]: - """Get device uuid from http header X-Device-Id""" - return request.headers.get("X-Device-Id") - - def files_size(): """Get total size of all files""" from mergin.app import db diff --git a/server/mergin/tests/fixtures.py b/server/mergin/tests/fixtures.py index 5d719878..c1a5ab8f 100644 --- a/server/mergin/tests/fixtures.py +++ b/server/mergin/tests/fixtures.py @@ -17,7 +17,7 @@ from ..stats.app import register from ..stats.models import MerginInfo from . import test_project, test_workspace_id, test_project_dir, TMP_DIR -from .utils import login_as_admin, initialize, cleanup, file_info +from .utils import login_as_admin, initialize, cleanup, file_info, ListSink from ..sync.files import files_changes_from_upload thisdir = os.path.dirname(os.path.realpath(__file__)) @@ -98,6 +98,16 @@ def client(app): return client +@pytest.fixture(scope="function") +def audit_capture(app): + """Replace the app's audit sink with an in-memory ListSink for the duration of the test.""" + sink = ListSink() + old = app.extensions["audit"]["sink"] + app.extensions["audit"]["sink"] = sink + yield sink + app.extensions["audit"]["sink"] = old + + @pytest.fixture(scope="function") def diff_project(app): """Modify testing project to contain some history with diffs. Geodiff lib is used to handle changes. diff --git a/server/mergin/tests/test_audit_app.py b/server/mergin/tests/test_audit_app.py new file mode 100644 index 00000000..a986935b --- /dev/null +++ b/server/mergin/tests/test_audit_app.py @@ -0,0 +1,28 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +import json +from unittest.mock import patch + +from ..app import db +from ..auth.events import AuthEventType +from ..auth.models import LoginHistory +from .utils import add_user + + +def test_emit_sanitizes_metadata_from_real_lockout_event(app, client, audit_capture): + """Test event gets serialized properly""" + user = add_user("lockme", "pass123") + for _ in range(4): + db.session.add(LoginHistory(user.id, "test-ua", "127.0.0.1", successful=False)) + db.session.commit() + + with patch.dict(app.config, {"LOCKOUT_POLICY": "5:300,10:3600"}): + client.post("/app/auth/login", json={"login": "lockme", "password": "wrong"}) + + e = audit_capture.one(AuthEventType.USER_UPDATED) + json.dumps(e.metadata) + + assert e.metadata["old_locked_until"] is None + assert isinstance(e.metadata["new_locked_until"], str) diff --git a/server/mergin/tests/test_audit_events.py b/server/mergin/tests/test_audit_events.py new file mode 100644 index 00000000..112e9953 --- /dev/null +++ b/server/mergin/tests/test_audit_events.py @@ -0,0 +1,600 @@ +# Copyright (C) Lutra Consulting Limited +# +# SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-MerginMaps-Commercial + +"""Contract tests: one test per defined EventType to verify the event is emitted +with the required fields. Each test exercises the minimum code path needed to +trigger the event — it is not a functional test of that path.""" + +from unittest.mock import patch + +from ..app import db +from ..auth.app import generate_confirmation_token, generate_unlock_token +from ..auth.events import AuthEventType +from ..auth.models import User +from ..sync.events import SyncEventType +from ..sync.models import AccessRequest, Project, ProjectRole +from . import DEFAULT_USER, test_project, test_workspace_id +from .utils import add_user, create_project, create_workspace, login + + +# --------------------------------------------------------------------------- +# Auth events +# --------------------------------------------------------------------------- + + +def test_user_login_succeeded(client, audit_capture): + login(client, DEFAULT_USER[0], DEFAULT_USER[1]) + + e = audit_capture.one(AuthEventType.USER_LOGIN_SUCCEEDED) + assert e.actor_email == f"{DEFAULT_USER[0]}@mergin.com" + assert ( + e.target_user_id == e.actor_id + ) # actor and target are the same person on login + + +def test_user_login_failed_invalid_credentials(client, audit_capture): + client.post( + "/app/auth/login", json={"login": "mergin", "password": "wrongpassword"} + ) + + e = audit_capture.one(AuthEventType.USER_LOGIN_FAILED) + assert e.metadata["reason"] == "invalid_credentials" + assert e.metadata["login"] == "mergin" + assert e.actor_id is None + + +def test_user_login_failed_account_inactive(client, audit_capture): + user = add_user("inactive_user", "pass123") + user.active = False + db.session.commit() + + client.post( + "/app/auth/login", json={"login": "inactive_user", "password": "pass123"} + ) + + assert ( + audit_capture.one(AuthEventType.USER_LOGIN_FAILED).metadata["reason"] + == "account_inactive" + ) + + +def test_user_password_changed(client, audit_capture): + user = add_user("pwduser", "oldpass123") + login(client, "pwduser", "oldpass123") + + client.post( + "/app/auth/change-password", + json={ + "old_password": "oldpass123", + "password": "New#pass456", + "confirm": "New#pass456", + }, + ) + + e = audit_capture.one(AuthEventType.USER_PASSWORD_CHANGED) + assert e.target_user_id == user.id + assert e.actor_id == user.id + + +def test_user_password_reset(app, client, audit_capture): + user = User.query.filter_by(username=DEFAULT_USER[0]).first() + token = generate_confirmation_token( + app, user.email, app.config["SECURITY_PASSWORD_SALT"] + ) + + client.post( + f"/app/auth/reset-password/{token}", + json={"password": "NewPass#123", "confirm": "NewPass#123"}, + ) + + e = audit_capture.one(AuthEventType.USER_PASSWORD_RESET_COMPLETED) + assert e.target_user_id == user.id + assert e.metadata["target_email"] == user.email + + +def test_user_created(audit_capture): + user = add_user("newuser", "pass123") + + e = audit_capture.one(AuthEventType.USER_CREATED) + assert e.target_user_id == user.id + assert e.metadata["target_email"] == "newuser@mergin.com" + + +def test_user_updated(audit_capture): + user = add_user("editme", "pass123") + db.session.refresh(user) + user.email = "updated@mergin.com" + user.passwd = ( + "newpassword" # excluded from user.updated — must never appear in audit + ) + db.session.commit() + + e = audit_capture.one(AuthEventType.USER_UPDATED) + assert e.metadata["new_email"] == "updated@mergin.com" + assert e.metadata["old_email"] == "editme@mergin.com" + assert "new_passwd" not in e.metadata + assert "old_passwd" not in e.metadata + + +def test_listener_null_actor_outside_request(audit_capture): + """Listeners fired from a Celery-like context (no active request) emit null actor + fields — the event is recorded as a system action with no user attributed.""" + add_user(username="systemcreated", password="pass123") + + e = audit_capture.one(AuthEventType.USER_CREATED) + assert e.actor_id is None + assert e.actor_email is None + assert e.actor_ip is None + + +def test_user_marked_for_deletion_by_user(client, audit_capture): + user = add_user("selfdelete", "pass123") + login(client, "selfdelete", "pass123") + + client.delete("/v1/user") + + e = audit_capture.one(AuthEventType.USER_MARKED_FOR_DELETION) + assert e.target_user_id == user.id + assert e.actor_id == user.id + + +def test_user_marked_for_deletion_by_admin(client, audit_capture): + user = add_user("admindelete", "pass123") + + client.delete(f"/app/admin/user/{user.username}") + + assert ( + audit_capture.one(AuthEventType.USER_MARKED_FOR_DELETION).target_user_id + == user.id + ) + + +def test_user_deactivated(client, audit_capture): + user = add_user("todeactivate", "pass123") + + client.patch(f"/app/admin/user/{user.username}", json={"active": False}) + + assert audit_capture.one(AuthEventType.USER_DEACTIVATED).target_user_id == user.id + + +def test_user_restored(client, audit_capture): + user = add_user("torestore", "pass123") + user.active = False + db.session.commit() + + client.patch(f"/app/admin/user/{user.username}", json={"active": True}) + + assert audit_capture.one(AuthEventType.USER_RESTORED).target_user_id == user.id + + +def test_user_deleted(client, audit_capture): + user = add_user("todelete", "pass123") + + client.delete(f"/app/admin/user/{user.username}") + + e = audit_capture.one(AuthEventType.USER_DELETED) + assert e.target_user_id == user.id + assert e.metadata["target_email"] == "todelete@mergin.com" + + +def test_user_locked(app, client, audit_capture): + from mergin.auth.models import LoginHistory + + user = add_user("lockme", "pass123") + # Insert 4 failed LoginHistory records so the next failed attempt hits threshold 5. + for _ in range(4): + db.session.add(LoginHistory(user.id, "test-ua", "127.0.0.1", successful=False)) + db.session.commit() + + with patch.dict(app.config, {"LOCKOUT_POLICY": "5:300,10:3600"}): + client.post("/app/auth/login", json={"login": "lockme", "password": "wrong"}) + + e = audit_capture.one(AuthEventType.USER_LOCKED) + assert e.target_user_id == user.id + assert "locked_until" in e.metadata + + +def test_user_unlocked(app, client, audit_capture): + import datetime + + user = add_user("unlockme", "pass123") + user.failed_login_attempts = 5 + user.locked_until = datetime.datetime.utcnow() + datetime.timedelta(seconds=300) + db.session.commit() + + token = generate_unlock_token(app, user) + client.post(f"/app/auth/unlock-account/{token}") + + e = audit_capture.one(AuthEventType.USER_UNLOCKED) + assert e.target_user_id == user.id + + +# --------------------------------------------------------------------------- +# Sync / project events +# --------------------------------------------------------------------------- + + +def test_project_created(audit_capture): + user = add_user("projowner", "pass123") + ws = create_workspace() + project = create_project("myproject", ws, user) + + e = audit_capture.one(SyncEventType.PROJECT_CREATED) + assert e.target_project_id == project.id + assert e.target_workspace_id == test_workspace_id + assert e.metadata["project_name"] == "myproject" + assert e.metadata["workspace_name"] == "mergin" + + +def test_project_created_from_template(client, audit_capture): + # Re-assign test_project's creator to the reserved TEMPLATES user so it + # becomes a template project (this is how the app identifies templates). + template_user = add_user("TEMPLATES", "pass123") + template = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + template.creator = template_user + db.session.commit() + + client.post( + f"/v1/project/{template.workspace.name}", + json={"name": "from_template", "template": test_project}, + ) + + # filter to the new project only (template re-assignment fires project.updated) + new_project = Project.query.filter_by(name="from_template").first() + events = [ + e + for e in audit_capture.of_type(SyncEventType.PROJECT_CREATED) + if e.target_project_id == new_project.id + ] + assert len(events) == 1 + assert events[0].metadata.get("created_from_template") == test_project + + +def test_project_created_from_clone(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + ws = project.workspace + + client.post( + f"/v1/project/clone/{ws.name}/{test_project}", + json={"namespace": ws.name, "project": "cloned_project"}, + ) + + e = audit_capture.one(SyncEventType.PROJECT_CREATED) + assert e.metadata.get("cloned_from") == str(project.id) + + +def test_project_updated(audit_capture): + user = add_user("projupdater", "pass123") + ws = create_workspace() + project = create_project("updateme", ws, user) + db.session.refresh(project) + + project.public = True + db.session.commit() + + e = audit_capture.one(SyncEventType.PROJECT_UPDATED) + assert e.metadata["new_public"] is True + assert e.metadata["old_public"] is False + + +def test_project_marked_for_deletion(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + + client.post(f"/v2/projects/{project.id}/scheduleDelete") + + e = audit_capture.one(SyncEventType.PROJECT_MARKED_FOR_DELETION) + assert e.target_project_id == project.id + assert e.target_workspace_id == project.workspace_id + + +def test_project_marked_for_deletion_via_v1_api(client, audit_capture): + """Deleting a project via the legacy v1 API (used by QGIS plugin) emits PROJECT_MARKED_FOR_DELETION.""" + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + workspace_name = project.workspace.name + + client.delete(f"/v1/project/{workspace_name}/{project.name}") + + e = audit_capture.one(SyncEventType.PROJECT_MARKED_FOR_DELETION) + assert e.target_project_id == project.id + assert e.target_workspace_id == project.workspace_id + + +def test_project_restored(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + client.post(f"/v2/projects/{project.id}/scheduleDelete") + audit_capture.events.clear() + + client.post(f"/app/project/removed-project/restore/{project.id}") + + assert ( + audit_capture.one(SyncEventType.PROJECT_RESTORED).target_project_id + == project.id + ) + + +def test_project_deleted(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + + client.delete(f"/v2/projects/{project.id}") + + assert ( + audit_capture.one(SyncEventType.PROJECT_DELETED).target_project_id == project.id + ) + # delete() renames/clears the project internally; must not also emit project.updated + assert len(audit_capture.of_type(SyncEventType.PROJECT_UPDATED)) == 0 + + +def test_project_member_added(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + user = add_user("newmember", "pass123") + + client.post( + f"/v2/projects/{project.id}/collaborators", + json={"user": user.email, "role": ProjectRole.READER.value}, + ) + + e = audit_capture.one(SyncEventType.PROJECT_MEMBER_ADDED) + assert e.target_project_id == project.id + assert e.metadata["target_email"] == user.email + assert e.metadata["role"] == ProjectRole.READER.value + + +def test_project_member_updated(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + user = add_user("updatemember", "pass123") + project.set_role(user.id, ProjectRole.READER) + db.session.commit() + + client.patch( + f"/v2/projects/{project.id}/collaborators/{user.id}", + json={"role": ProjectRole.EDITOR.value}, + ) + + e = audit_capture.one(SyncEventType.PROJECT_MEMBER_UPDATED) + assert e.metadata["old_role"] == ProjectRole.READER.value + assert e.metadata["new_role"] == ProjectRole.EDITOR.value + + +def test_project_member_deleted(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + user = add_user("removemember", "pass123") + project.set_role(user.id, ProjectRole.READER) + db.session.commit() + + client.delete(f"/v2/projects/{project.id}/collaborators/{user.id}") + + e = audit_capture.one(SyncEventType.PROJECT_MEMBER_DELETED) + assert e.metadata["target_email"] == user.email + assert e.metadata["reason"] == "removed" + + +def test_project_member_left(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + user = add_user("leavemember", "pass123") + project.set_role(user.id, ProjectRole.READER) + db.session.commit() + login(client, user.username, "pass123") + + client.post(f"/app/project/unsubscribe/{project.id}") + + e = audit_capture.one(SyncEventType.PROJECT_MEMBER_DELETED) + assert e.metadata["target_email"] == user.email + assert e.metadata["reason"] == "left" + + +def test_project_access_request_created(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + user = add_user("requester", "pass123") + login(client, "requester", "pass123") + + client.post(f"/app/project/access-request/{project.workspace.name}/{project.name}") + + e = audit_capture.one(SyncEventType.PROJECT_ACCESS_REQUEST_INITIATED) + assert e.target_project_id == project.id + assert e.actor_id == user.id + + +def test_project_access_request_accepted(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + requester = add_user("acceptrequester", "pass123") + access_request = AccessRequest(project, requester.id) + db.session.add(access_request) + db.session.commit() + + client.post( + f"/app/project/access-request/accept/{access_request.id}", + json={"permissions": "read"}, + ) + + assert ( + audit_capture.one(SyncEventType.PROJECT_ACCESS_REQUEST_ACCEPTED).metadata[ + "target_email" + ] + == requester.email + ) + # accepting also fires project.member.added + assert ( + audit_capture.one(SyncEventType.PROJECT_MEMBER_ADDED).metadata["target_email"] + == requester.email + ) + + +def test_project_access_request_rejected(client, audit_capture): + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + requester = add_user("rejectrequester", "pass123") + access_request = AccessRequest(project, requester.id) + db.session.add(access_request) + db.session.commit() + + client.delete(f"/app/project/access-request/{access_request.id}") + + assert ( + audit_capture.one(SyncEventType.PROJECT_ACCESS_REQUEST_CANCELED).metadata[ + "target_email" + ] + == requester.email + ) + + +def test_project_version_created(client, audit_capture): + from .utils import file_info + from . import test_project_dir + + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + # A remove-only push needs no chunk uploads so it's self-contained. + data = { + "version": "v1", + "changes": { + "added": [], + "updated": [], + "removed": [file_info(test_project_dir, "test3.txt")], + }, + } + resp = client.post(f"/v2/projects/{project.id}/versions", json=data) + assert resp.status_code == 201 + + e = audit_capture.one(SyncEventType.PROJECT_VERSION_CREATED) + assert e.target_project_id == project.id + assert e.metadata["version"] == "v2" + + +def test_project_version_created_v1_no_upload(client, audit_capture): + """V1 push with only removals takes the no-upload fast path in project_push.""" + from .utils import file_info + from . import test_project_dir + + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + data = { + "version": "v1", + "changes": { + "added": [], + "updated": [], + "removed": [file_info(test_project_dir, "test3.txt")], + }, + } + resp = client.post( + f"/v1/project/push/{project.workspace.name}/{project.name}", json=data + ) + assert resp.status_code == 200 + + e = audit_capture.one(SyncEventType.PROJECT_VERSION_CREATED) + assert e.target_project_id == project.id + assert e.metadata["version"] == "v2" + + +def test_project_version_created_v1_push_finish(client, audit_capture): + """V1 push with file uploads goes through push_finish.""" + import os + from .utils import file_info + from . import test_project_dir + + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + filename = "test.qgs" + filepath = os.path.join(test_project_dir, filename) + # "updated" because the fixture already uploaded all test_project_dir files at v1. + data = { + "version": "v1", + "changes": { + "added": [], + "updated": [file_info(test_project_dir, filename)], + "removed": [], + }, + } + resp = client.post( + f"/v1/project/push/{project.workspace.name}/{project.name}", json=data + ) + assert resp.status_code == 200 + upload_id = resp.json["transaction"] + for chunk_id in data["changes"]["updated"][0]["chunks"]: + with open(filepath, "rb") as f: + client.post( + f"/v1/project/push/chunk/{upload_id}/{chunk_id}", + data=f.read(1024), + headers={"Content-Type": "application/octet-stream"}, + ) + resp = client.post(f"/v1/project/push/finish/{upload_id}") + assert resp.status_code == 200 + + e = audit_capture.one(SyncEventType.PROJECT_VERSION_CREATED) + assert e.target_project_id == project.id + assert e.metadata["version"] == "v2" + + +def test_project_version_created_from_template(client, audit_capture): + """Creating a project from a template emits PROJECT_VERSION_CREATED for the v1.""" + template_user = add_user("TEMPLATES", "pass123") + template = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + template.creator = template_user + db.session.commit() + audit_capture.events.clear() + + client.post( + f"/v1/project/{template.workspace.name}", + json={"name": "from_template_audit", "template": test_project}, + ) + + new_project = Project.query.filter_by(name="from_template_audit").first() + events = [ + e + for e in audit_capture.of_type(SyncEventType.PROJECT_VERSION_CREATED) + if e.target_project_id == new_project.id + ] + assert len(events) == 1 + assert events[0].metadata["version"] == "v1" + + +def test_project_version_created_from_clone(client, audit_capture): + """Cloning a non-empty project emits PROJECT_VERSION_CREATED for the v1.""" + project = Project.query.filter_by( + workspace_id=test_workspace_id, name=test_project + ).first() + ws = project.workspace + + client.post( + f"/v1/project/clone/{ws.name}/{test_project}", + json={"namespace": ws.name, "project": "cloned_audit"}, + ) + + cloned = Project.query.filter_by(name="cloned_audit").first() + events = [ + e + for e in audit_capture.of_type(SyncEventType.PROJECT_VERSION_CREATED) + if e.target_project_id == cloned.id + ] + assert len(events) == 1 + assert events[0].metadata["version"] == "v1" diff --git a/server/mergin/tests/test_project_controller.py b/server/mergin/tests/test_project_controller.py index 1a0c76aa..3e5e12ae 100644 --- a/server/mergin/tests/test_project_controller.py +++ b/server/mergin/tests/test_project_controller.py @@ -2072,21 +2072,12 @@ def test_get_projects_by_uuids(client): user = User.query.filter_by(username="mergin").first() test_workspace = create_workspace() p1 = create_project("foo", test_workspace, user) - user2 = add_user("user2", "ilovemergin") - test_workspace_2 = create_workspace() - test_workspace_2._id = ( - 2 # FIXME: This should be refactored due to only one workspace in CE - ) - p2 = create_project("foo", test_workspace_2, user2) - uuids = ",".join([str(p1.id), str(p2.id), "1234"]) + uuids = ",".join([str(p1.id), "1234"]) resp = client.get(f"/v1/project/by_uuids?uuids={uuids}") assert resp.status_code == 200 - assert str(p1.id) in resp.json # user has access to - assert ( - str(p2.id) not in resp.json - ) # belongs to user2, and user does not have access - assert "1234" not in resp.json # invalid id + assert str(p1.id) in resp.json + assert "1234" not in resp.json # invalid id is excluded uuids = ",".join([str(uuid.uuid4()) for _ in range(0, 11)]) resp = client.get(f"/v1/project/by_uuids?uuids={uuids}") diff --git a/server/mergin/tests/utils.py b/server/mergin/tests/utils.py index 57f67e80..52b3e970 100644 --- a/server/mergin/tests/utils.py +++ b/server/mergin/tests/utils.py @@ -4,6 +4,7 @@ import json import shutil +from typing import List import pysqlite3 import uuid import math @@ -405,3 +406,27 @@ def logout(client): """Test helper to log out the client""" resp = client.get(url_for("/.mergin_auth_controller_logout")) assert resp.status_code == 200 + + +class ListSink: + """In-memory audit sink for use in automated tests. + + Install via the audit_capture fixture; do not use in production code. + """ + + def __init__(self): + self.events: List = [] + + def write(self, event) -> None: + json.dumps(event.metadata) + self.events.append(event) + + def of_type(self, event_type) -> List: + """Return all captured events matching event_type.""" + return [e for e in self.events if e.event_type == event_type] + + def one(self, event_type): + """Assert exactly one event of event_type was captured and return it.""" + events = self.of_type(event_type) + assert len(events) == 1, f"Expected 1 {event_type} event, got {len(events)}" + return events[0] diff --git a/server/mergin/utils.py b/server/mergin/utils.py index aa878ffe..550bd1ce 100644 --- a/server/mergin/utils.py +++ b/server/mergin/utils.py @@ -116,6 +116,33 @@ def parse_order_params( return order_by_params +def get_user_agent(request) -> str: + """Return user agent from request headers. + + For browser clients returns a parsed summary; otherwise the raw header value. + """ + if request.user_agent.browser and request.user_agent.platform: + client = request.user_agent.browser.capitalize() + version = request.user_agent.version + system = request.user_agent.platform.capitalize() + return f"{client}/{version} ({system})" + return request.user_agent.string + + +def get_ip(request) -> str: + """Return the client IP address, respecting X-Forwarded-For from a proxy.""" + forwarded_ips = request.environ.get( + "HTTP_X_FORWARDED_FOR", request.environ.get("REMOTE_ADDR", "untrackable") + ) + # AWS infra may send a comma-separated list; the first entry is the real client IP + return forwarded_ips.split(",")[0] + + +def get_device_id(request) -> Optional[str]: + """Return the device UUID from the X-Device-Id header, or None if absent.""" + return request.headers.get("X-Device-Id") + + def format_time_delta(delta: timedelta) -> str: """Format timedelta difference approximately in days or hours""" days = round(delta.total_seconds() / (24 * 3600))