From ce0d561adeef6a1dadf29db1a8e680fb005bcbf1 Mon Sep 17 00:00:00 2001 From: Bruno Charest Date: Mon, 21 Sep 2026 06:38:02 -0400 Subject: [PATCH] feat(share,permissions): partage de page par groupes (page_shares) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - app/db.py: ajout colonne shared_with_group_id (FK user_groups, migration backfill v5.x) dans page_shares - app/routers/sharing.py: POST /api/pages/{id}/share accepte group_id (upsert, verif FK groupe), GET /shares expose kind/group_name, synchro bidirectionnelle avec page_permissions (mirror grant/revoke) pour que PermissionManager donne un acces effectif (view/comment/edit) aux membres du groupe; PUT/DELETE gardent le miroir a jour - app/routers/board.py, library.py: received/made incluent les partages via groupes (JOIN group_members) - app/templates/_page_editor_content.html, _page_editor_scripts.html: dialogue Share — invite groups (fetch /api/v2/groups, filtre deja partages), pickInviteGroup, shareInvite(group_id), rendu accessList avec avatar groupe - tests/test_share_groups.py: 10 tests (CRUD groupe, kind, miroir ACL) Chore: inclut evolutions v6.4.0 deja en working copy (webhooks prod, migrations, sync, config/main) pour garder l'arbre coherent. --- app/config.py | 4 + app/db.py | 6 + app/main.py | 13 +- app/migrations.py | 19 ++ app/routers/api_v2.py | 80 ++++++--- app/routers/board.py | 21 ++- app/routers/library.py | 17 +- app/routers/sharing.py | 139 +++++++++++--- app/routers/sync.py | 38 ++-- app/services/webhook_outbound.py | 229 ++++++++++++++++++++++-- app/templates/_page_editor_content.html | 33 +++- app/templates/_page_editor_scripts.html | 49 ++++- tests/test_share_groups.py | 211 ++++++++++++++++++++++ 13 files changed, 774 insertions(+), 85 deletions(-) create mode 100644 tests/test_share_groups.py diff --git a/app/config.py b/app/config.py index f19e60e..5f197ae 100644 --- a/app/config.py +++ b/app/config.py @@ -68,6 +68,10 @@ class Settings(BaseSettings): reminders_enabled: bool = True reminder_scan_interval_seconds: int = 60 + # Webhooks outbound (v6.4.0) — retry of failed deliveries + webhook_retry_enabled: bool = True + webhook_retry_interval_seconds: int = 60 + # Email / SMTP notifications (v4.9.0) — optional. If smtp_host is empty, # email notifications are skipped (only in-app notifications are delivered). smtp_host: str = "" diff --git a/app/db.py b/app/db.py index bc3d475..83136d7 100644 --- a/app/db.py +++ b/app/db.py @@ -442,12 +442,18 @@ def init_db(): id INTEGER PRIMARY KEY AUTOINCREMENT, page_id INTEGER NOT NULL REFERENCES pages(id) ON DELETE CASCADE, shared_with_user_id INTEGER REFERENCES users(id), + shared_with_group_id INTEGER REFERENCES user_groups(id) ON DELETE CASCADE, shared_with_email TEXT DEFAULT '', permission TEXT NOT NULL DEFAULT 'view', created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, created_by INTEGER REFERENCES users(id) ) """) + # v5.x: migration — partage par groupes (colonne manquante sur DB existantes) + try: + conn.execute("ALTER TABLE page_shares ADD COLUMN shared_with_group_id INTEGER REFERENCES user_groups(id) ON DELETE CASCADE") + except sqlite3.OperationalError: + pass # 4) recents table conn.execute(""" CREATE TABLE IF NOT EXISTS recents ( diff --git a/app/main.py b/app/main.py index e031799..d447fb6 100644 --- a/app/main.py +++ b/app/main.py @@ -98,13 +98,22 @@ async def lifespan(_app: FastAPI): from app.services.reminders import reminder_scheduler reminder_task = asyncio.create_task(reminder_scheduler()) + # ── Webhooks outbound (v6.4.0): retry failed deliveries ── + from app.services.webhook_outbound import webhook_retry_scheduler + webhook_task = None + if settings.webhook_retry_enabled: + webhook_task = asyncio.create_task(webhook_retry_scheduler()) + logger.info("FlowDeck v%s started on port %d", dashboard._get_app_version(), settings.app_port) try: yield finally: - for task in (scheduler_task, automation_task, backup_task, projects_task, trash_task, reminder_task): + _tasks = (scheduler_task, automation_task, backup_task, projects_task, trash_task, reminder_task) + if webhook_task is not None: + _tasks = _tasks + (webhook_task,) + for task in _tasks: task.cancel() - for task in (scheduler_task, automation_task, backup_task, projects_task, trash_task, reminder_task): + for task in _tasks: try: await task except asyncio.CancelledError: diff --git a/app/migrations.py b/app/migrations.py index f4b454f..24f9558 100644 --- a/app/migrations.py +++ b/app/migrations.py @@ -904,6 +904,25 @@ def _migration_v630_api_v2(conn: sqlite3.Connection) -> None: conn.execute("CREATE INDEX IF NOT EXISTS idx_idemp_user ON idempotency_keys(user_id, created_at)") +@register(21, "v6.4.0: webhooks prod — event, retry ledger") +def _migration_v640_webhooks_prod(conn: sqlite3.Connection) -> None: + """v6.4.0 — Webhooks v2 production. + + ``webhook_deliveries`` gains ``event`` (which event was delivered) and + ``next_retry_at`` (epoch seconds; picked up by the retry scheduler). + New statuses: ``retrying`` (a later attempt is scheduled) and + ``superseded`` (a retry row replaced this attempt). + """ + _cols = {r[1] for r in conn.execute("PRAGMA table_info(webhook_deliveries)").fetchall()} + if "event" not in _cols: + conn.execute("ALTER TABLE webhook_deliveries ADD COLUMN event TEXT NOT NULL DEFAULT ''") + if "next_retry_at" not in _cols: + conn.execute("ALTER TABLE webhook_deliveries ADD COLUMN next_retry_at REAL") + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_wd_retry ON webhook_deliveries(status, next_retry_at)" + ) + + @register(17, "v6.0.0: sync_version columns") def _migration_sync_version_columns(conn: sqlite3.Connection) -> None: """v6.0.0 — optimistic-concurrency version counters for offline sync. diff --git a/app/routers/api_v2.py b/app/routers/api_v2.py index 3468c9d..41b8fa2 100644 --- a/app/routers/api_v2.py +++ b/app/routers/api_v2.py @@ -2184,15 +2184,32 @@ async def admin_audit_logs_v2(request: Request, limit: int = 50, authorization: rows = conn.execute("SELECT * FROM api_audit_log ORDER BY created_at DESC LIMIT ?", (limit,)).fetchall() return {"logs": [row_to_dict(r) for r in rows]} -# ── Webhooks (CRUD simple) ──────────────────────────────────────────────── +# ── Webhooks (v6.4.0 prod) ──────────────────────────────────────────────── + +@router.get("/webhooks/events") +async def list_webhook_events_v2(request: Request, authorization: str | None = Header(default=None)): + """Catalogue of deliverable events (+ wildcard syntax).""" + user = get_bearer_user(request, authorization) + _v2_rate_check(request, user) + from app.services.webhook_outbound import EVENTS + return {"events": EVENTS, "wildcards": ["*", "page.*", "collection.*"]} + @router.get("/webhooks") async def list_webhooks_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) + limit, offset = parse_pagination(request) with get_conn() as conn: - rows = conn.execute("SELECT * FROM webhook_subscriptions ORDER BY created_at DESC").fetchall() - return {"webhooks": [dict(r) for r in rows]} + total = conn.execute("SELECT COUNT(*) AS n FROM webhook_subscriptions").fetchone()["n"] + rows = conn.execute( + "SELECT * FROM webhook_subscriptions ORDER BY created_at DESC LIMIT ? OFFSET ?", + (limit, offset), + ).fetchall() + return JSONResponse( + content={"webhooks": [dict(r) for r in rows]}, + headers=paginate_headers(total), + ) @router.post("/webhooks") async def create_webhook_v2(request: Request, authorization: str | None = Header(default=None)): @@ -2204,18 +2221,25 @@ async def create_webhook_v2(request: Request, authorization: str | None = Header body = await request.json() except Exception: body = {} + from app.services.webhook_outbound import EVENTS, _event_matches url = (body.get("url") or "").strip() event = (body.get("event") or "page.created").strip() secret = (body.get("secret") or "").strip() if not url or not url.startswith("http"): raise HTTPException(400, "url must start with http") + if not event or (event not in EVENTS and not (event.endswith(".*") or event in ("*", "all"))): + raise HTTPException(400, f"Unknown event '{event}'. See GET /api/v2/webhooks/events") + # make sure the pattern matches at least one known event + if not any(_event_matches(event, e) for e in EVENTS): + raise HTTPException(400, f"Event pattern '{event}' matches no known event") with get_conn() as conn: cur = conn.execute("INSERT INTO webhook_subscriptions (url, event, secret) VALUES (?, ?, ?)", (url, event, secret)) wid = cur.lastrowid conn.commit() row = conn.execute("SELECT * FROM webhook_subscriptions WHERE id=?", (wid,)).fetchone() audit_log(user, "webhook.create", "webhook", wid, url, request) - return {"id": wid, "status": "created", "webhook": dict(row) if row else {}} + return {"id": wid, "status": "created", "webhook": dict(row) if row else {}, + "signature_header": "X-FlowDeck-Signature (HMAC-SHA256, sha256=)" if secret else None} @router.patch("/webhooks/{webhook_id}") async def patch_webhook_v2(webhook_id: int, request: Request, authorization: str | None = Header(default=None)): @@ -2262,24 +2286,40 @@ async def test_webhook_v2(webhook_id: int, request: Request, authorization: str row = conn.execute("SELECT * FROM webhook_subscriptions WHERE id=?", (webhook_id,)).fetchone() if not row: raise HTTPException(404, "Webhook not found") - # fire test ping (best-effort, record delivery) - try: - from app.services.webhook_outbound import fire_event - await fire_event("ping", {"webhook_id": webhook_id, "test": True}) - conn.execute("INSERT INTO webhook_deliveries (webhook_id, status, http_code, payload) VALUES (?, 'delivered', 200, ?)", (webhook_id, json.dumps({"event": "ping"}))) - conn.commit() - except Exception as e: - try: - conn.execute("INSERT INTO webhook_deliveries (webhook_id, status, http_code, error) VALUES (?, 'failed', 0, ?)", (webhook_id, str(e)[:1000])) - conn.commit() - except Exception: - pass - return {"webhook_id": webhook_id, "status": "tested"} + # live delivery via the prod dispatcher (HMAC + retry + journal), + # direct to this subscription only (no wildcard fan-out) + from app.services.webhook_outbound import deliver_to_sub + ok = await deliver_to_sub(webhook_id, row["url"], "ping", + {"webhook_id": webhook_id, "test": True}, + row["secret"] or "") + audit_log(user, "webhook.test", "webhook", webhook_id, f"ok={ok}", request) + return {"webhook_id": webhook_id, "status": "tested", "delivered": ok} @router.get("/webhooks/{webhook_id}/deliveries") -async def list_deliveries_v2(webhook_id: int, request: Request, authorization: str | None = Header(default=None)): +async def list_deliveries_v2(webhook_id: int, request: Request, + status: str | None = None, + authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) + limit, offset = parse_pagination(request) with get_conn() as conn: - rows = conn.execute("SELECT * FROM webhook_deliveries WHERE webhook_id=? ORDER BY created_at DESC LIMIT 50", (webhook_id,)).fetchall() - return {"deliveries": [row_to_dict(r) for r in rows]} + if status: + total = conn.execute( + "SELECT COUNT(*) AS n FROM webhook_deliveries WHERE webhook_id=? AND status=?", + (webhook_id, status)).fetchone()["n"] + rows = conn.execute( + "SELECT * FROM webhook_deliveries WHERE webhook_id=? AND status=? " + "ORDER BY created_at DESC LIMIT ? OFFSET ?", + (webhook_id, status, limit, offset)).fetchall() + else: + total = conn.execute( + "SELECT COUNT(*) AS n FROM webhook_deliveries WHERE webhook_id=?", + (webhook_id,)).fetchone()["n"] + rows = conn.execute( + "SELECT * FROM webhook_deliveries WHERE webhook_id=? " + "ORDER BY created_at DESC LIMIT ? OFFSET ?", + (webhook_id, limit, offset)).fetchall() + return JSONResponse( + content={"deliveries": [row_to_dict(r) for r in rows]}, + headers=paginate_headers(total), + ) diff --git a/app/routers/board.py b/app/routers/board.py index 4d797dc..9a82e57 100644 --- a/app/routers/board.py +++ b/app/routers/board.py @@ -612,12 +612,21 @@ def _load_shared_sidebar_pages(user_id: int) -> tuple[list, list, list, list]: f"ORDER BY updated_at DESC LIMIT 20", scope_params, ).fetchall() - received_rows = conn.execute( - "SELECT DISTINCT p.id, p.title, p.workspace, p.updated_at FROM page_shares s " - "JOIN pages p ON p.id=s.page_id " - "WHERE s.shared_with_user_id=? AND p.deleted_at IS NULL", - (user_id,), - ).fetchall() + try: + received_rows = conn.execute( + "SELECT DISTINCT p.id, p.title, p.workspace, p.updated_at FROM page_shares s " + "JOIN pages p ON p.id=s.page_id " + "LEFT JOIN group_members gm ON gm.group_id = s.shared_with_group_id AND gm.user_id=? " + "WHERE (s.shared_with_user_id=? OR gm.user_id=?) AND p.deleted_at IS NULL", + (user_id, user_id, user_id), + ).fetchall() + except Exception: + received_rows = conn.execute( + "SELECT DISTINCT p.id, p.title, p.workspace, p.updated_at FROM page_shares s " + "JOIN pages p ON p.id=s.page_id " + "WHERE s.shared_with_user_id=? AND p.deleted_at IS NULL", + (user_id,), + ).fetchall() def _entry(r, icon): return { diff --git a/app/routers/library.py b/app/routers/library.py index 91bf6cf..6944777 100644 --- a/app/routers/library.py +++ b/app/routers/library.py @@ -235,15 +235,24 @@ async def library_shared( uid = _get_user_id(request) # Page ids the user shares toward others (nominal page_shares) or receives + # (direct shares + group shares via group_members) with get_conn() as conn: made_rows = conn.execute( "SELECT DISTINCT s.page_id FROM page_shares s WHERE s.created_by=?", (uid,), ).fetchall() - recv_rows = conn.execute( - "SELECT DISTINCT s.page_id FROM page_shares s WHERE s.shared_with_user_id=?", - (uid,), - ).fetchall() + try: + recv_rows = conn.execute( + """SELECT DISTINCT s.page_id FROM page_shares s + LEFT JOIN group_members gm ON gm.group_id = s.shared_with_group_id AND gm.user_id=? + WHERE s.shared_with_user_id=? OR gm.user_id=?""", + (uid, uid, uid), + ).fetchall() + except Exception: + recv_rows = conn.execute( + "SELECT DISTINCT s.page_id FROM page_shares s WHERE s.shared_with_user_id=?", + (uid,), + ).fetchall() made_ids = {r[0] for r in made_rows} recv_ids = {r[0] for r in recv_rows} diff --git a/app/routers/sharing.py b/app/routers/sharing.py index edb5c10..673c3d7 100644 --- a/app/routers/sharing.py +++ b/app/routers/sharing.py @@ -36,18 +36,23 @@ def _slugify(title: str) -> str: @router.post("/pages/{page_id}/share") async def share_page(page_id: int, request: Request): - """Invite a user or email to a page.""" + """Invite a user, an email, or a group to a page.""" user = _require_auth(request) body = await request.json() if request.headers.get("content-type") else {} target_user_id = body.get("user_id") + target_group_id = body.get("group_id") email = body.get("email", "") permission = body.get("permission", "view") if permission not in ("view", "comment", "edit"): raise HTTPException(400, "Invalid permission. Use view, comment, or edit.") - if not target_user_id and not email: - raise HTTPException(400, "Provide user_id or email to share with.") + if not target_user_id and not target_group_id and not email: + raise HTTPException(400, "Provide user_id, group_id or email to share with.") + + # Bridge share permission (view/comment/edit) → granular role + # (viewer/commenter/editor) so page_permissions grants stay in sync. + _SHARE_TO_ROLE = {"view": "viewer", "comment": "commenter", "edit": "editor"} with get_conn() as conn: # Verify page exists @@ -61,18 +66,29 @@ async def share_page(page_id: int, request: Request): if not target: raise HTTPException(404, "Target user not found") + # Verify target group exists if group_id given + if target_group_id: + gtarget = conn.execute("SELECT id FROM user_groups WHERE id=?", (target_group_id,)).fetchone() + if not gtarget: + raise HTTPException(404, "Target group not found") + # Upsert to avoid duplicates: update the existing permission if the same - # target (user or email) is already shared on this page. + # target (user, group or email) is already shared on this page. target_row = None if target_user_id: target_row = conn.execute( "SELECT id FROM page_shares WHERE page_id=? AND shared_with_user_id=?", (page_id, target_user_id), ).fetchone() + elif target_group_id: + target_row = conn.execute( + "SELECT id FROM page_shares WHERE page_id=? AND shared_with_group_id=?", + (page_id, target_group_id), + ).fetchone() elif email: target_row = conn.execute( """SELECT id FROM page_shares - WHERE page_id=? AND shared_with_email=? AND shared_with_user_id IS NULL""", + WHERE page_id=? AND shared_with_email=? AND shared_with_user_id IS NULL AND shared_with_group_id IS NULL""", (page_id, email.strip()), ).fetchone() @@ -85,29 +101,75 @@ async def share_page(page_id: int, request: Request): else: if not email: email = "" - cur = conn.execute( - """INSERT INTO page_shares (page_id, shared_with_user_id, shared_with_email, permission, created_by) - VALUES (?, ?, ?, ?, ?)""", - (page_id, target_user_id, email.strip(), permission, user["id"]), - ) + try: + cur = conn.execute( + """INSERT INTO page_shares (page_id, shared_with_user_id, shared_with_group_id, shared_with_email, permission, created_by) + VALUES (?, ?, ?, ?, ?, ?)""", + (page_id, target_user_id, target_group_id, email.strip(), permission, user["id"]), + ) + except Exception: + # Fallback for DBs where the migration has not run yet + cur = conn.execute( + """INSERT INTO page_shares (page_id, shared_with_user_id, shared_with_email, permission, created_by) + VALUES (?, ?, ?, ?, ?)""", + (page_id, target_user_id, email.strip(), permission, user["id"]), + ) share_id = cur.lastrowid conn.execute("UPDATE pages SET is_shared=1 WHERE id=?", (page_id,)) + # ── Mirror group shares into page_permissions so the ACL used by + # PermissionManager (can_view/edit/comment) grants real access to + # every group member. Best-effort: never break legacy page_shares. + if target_group_id: + try: + _mirror_share_grant(conn, page_id, target_group_id, _SHARE_TO_ROLE[permission], user["id"]) + except Exception: + logger.warning("share→page_permissions mirror failed (page=%s group=%s)", page_id, target_group_id) conn.commit() return { "id": share_id, "page_id": page_id, "shared_with_user_id": target_user_id, + "shared_with_group_id": target_group_id, "shared_with_email": email, "permission": permission, "status": "shared", } +def _mirror_share_grant(conn, page_id: int, group_id: int, role: str, granted_by: int) -> None: + """Upsert a ``page_permissions`` grant mirroring a group ``page_shares`` row. + + Keeps the granular ACL (used by ``PermissionManager``) in sync with what + the share dialog shows, so invited groups get effective view/edit rights. + """ + existing = conn.execute( + "SELECT id FROM page_permissions WHERE page_id=? AND user_id IS NULL AND group_id=?", + (page_id, group_id), + ).fetchone() + if existing: + conn.execute("UPDATE page_permissions SET role=?, granted_by=? WHERE id=?", + (role, granted_by, existing["id"])) + else: + conn.execute( + "INSERT INTO page_permissions (page_id, user_id, group_id, role, granted_by) " + "VALUES (?, NULL, ?, ?, ?)", + (page_id, group_id, role, granted_by), + ) + + +def _mirror_share_revoke(conn, page_id: int, group_id: int) -> None: + """Remove the mirrored grant when a group share is updated away or deleted.""" + conn.execute( + "DELETE FROM page_permissions WHERE page_id=? AND user_id IS NULL AND group_id=?", + (page_id, group_id), + ) + + @router.put("/pages/{page_id}/share/{share_id}", description="Update a share's permission.") async def update_share_permission(page_id: int, share_id: int, request: Request): """Change the permission level of an existing share entry.""" - _require_auth(request) + user = _require_auth(request) body = await request.json() if request.headers.get("content-type") else {} permission = body.get("permission", "") @@ -117,7 +179,7 @@ async def update_share_permission(page_id: int, share_id: int, request: Request) with get_conn() as conn: row = conn.execute( - "SELECT id FROM page_shares WHERE id=? AND page_id=?", + "SELECT id, shared_with_group_id FROM page_shares WHERE id=? AND page_id=?", (share_id, page_id), ).fetchone() if not row: @@ -127,6 +189,18 @@ async def update_share_permission(page_id: int, share_id: int, request: Request) "UPDATE page_shares SET permission=? WHERE id=?", (permission, share_id), ) + # Keep the mirrored ACL grant in sync for group shares. + try: + gid = row["shared_with_group_id"] if "shared_with_group_id" in row.keys() else None + except Exception: + gid = None + if gid: + try: + _mirror_share_grant(conn, page_id, gid, + {"view": "viewer", "comment": "commenter", "edit": "editor"}[permission], + user["id"]) + except Exception: + logger.warning("share→page_permissions mirror failed (share=%s)", share_id) conn.commit() return {"status": "updated", "share_id": share_id, "permission": permission} @@ -139,13 +213,22 @@ async def remove_share(page_id: int, share_id: int, request: Request): with get_conn() as conn: row = conn.execute( - "SELECT id FROM page_shares WHERE id=? AND page_id=?", + "SELECT id, shared_with_group_id FROM page_shares WHERE id=? AND page_id=?", (share_id, page_id), ).fetchone() if not row: raise HTTPException(404, "Share entry not found") conn.execute("DELETE FROM page_shares WHERE id=?", (share_id,)) + try: + gid = row["shared_with_group_id"] if "shared_with_group_id" in row.keys() else None + except Exception: + gid = None + if gid: + try: + _mirror_share_revoke(conn, page_id, gid) + except Exception: + logger.warning("share→page_permissions revoke failed (share=%s)", share_id) # If no more shares, unset is_shared remaining = conn.execute( "SELECT COUNT(*) AS c FROM page_shares WHERE page_id=?", (page_id,) @@ -167,14 +250,25 @@ async def list_shares(page_id: int, request: Request): if not page: raise HTTPException(404, "Page not found") - rows = conn.execute( - """SELECT s.*, u.login, u.full_name, u.avatar_url - FROM page_shares s - LEFT JOIN users u ON s.shared_with_user_id = u.id - WHERE s.page_id=? - ORDER BY s.created_at DESC""", - (page_id,), - ).fetchall() + try: + rows = conn.execute( + """SELECT s.*, u.login, u.full_name, u.avatar_url, g.name AS group_name + FROM page_shares s + LEFT JOIN users u ON s.shared_with_user_id = u.id + LEFT JOIN user_groups g ON s.shared_with_group_id = g.id + WHERE s.page_id=? + ORDER BY s.created_at DESC""", + (page_id,), + ).fetchall() + except Exception: + rows = conn.execute( + """SELECT s.*, u.login, u.full_name, u.avatar_url + FROM page_shares s + LEFT JOIN users u ON s.shared_with_user_id = u.id + WHERE s.page_id=? + ORDER BY s.created_at DESC""", + (page_id,), + ).fetchall() return { "page_id": page_id, @@ -182,6 +276,7 @@ async def list_shares(page_id: int, request: Request): { "id": r["id"], "shared_with_user_id": r["shared_with_user_id"], + "shared_with_group_id": r["shared_with_group_id"] if "shared_with_group_id" in r.keys() else None, "shared_with_email": r["shared_with_email"], "permission": r["permission"], "created_at": r["created_at"], @@ -189,6 +284,8 @@ async def list_shares(page_id: int, request: Request): "user_login": r["login"], "user_full_name": r["full_name"], "user_avatar_url": r["avatar_url"], + "group_name": r["group_name"] if "group_name" in r.keys() else None, + "kind": "group" if (("shared_with_group_id" in r.keys() and r["shared_with_group_id"]) or ("group_name" in r.keys() and r["group_name"])) else "user", } for r in rows ], diff --git a/app/routers/sync.py b/app/routers/sync.py index 703dfad..285edd8 100644 --- a/app/routers/sync.py +++ b/app/routers/sync.py @@ -1,16 +1,18 @@ -"""FlowDeck — /api/v2/sync endpoints (v6.0.0 PWA offline sync). +"""FlowDeck — /api/v2/sync endpoints (v6.0.0 PWA offline sync, Bearer v6.4.0). -Pairs with ``app/services/sync_engine.py``. All routes require an authenticated -session (``flowdeck_session`` cookie). +Auth: ``Authorization: Bearer `` (scopes ``read`` for delta/status, +``write`` for batch). The legacy ``flowdeck_session`` cookie is still accepted +as a fallback so the installed PWA/service worker keeps syncing. """ from __future__ import annotations import logging -from fastapi import APIRouter, HTTPException, Query, Request +from fastapi import APIRouter, Header, HTTPException, Query, Request from fastapi.responses import JSONResponse from app.auth.session import SessionManager +from app.services.api_v2_helpers import get_bearer_user, has_scope from app.services.sync_engine import SyncEngine logger = logging.getLogger(__name__) @@ -20,7 +22,21 @@ router = APIRouter(prefix="/api/v2/sync", tags=["sync"]) _engine = SyncEngine() -def _user(request: Request) -> dict: +def _user(request: Request, authorization: str | None = None, + *, required_scope: str = "read") -> dict: + """Bearer-first auth with session-cookie fallback (offline.js compat).""" + auth = authorization or request.headers.get("authorization") or "" + if auth and auth.lower().startswith("bearer "): + try: + user = get_bearer_user(request, authorization) + except HTTPException: + raise HTTPException(status_code=401, detail="Invalid or expired API token") + if not has_scope(user.get("_token_scopes"), required_scope): + raise HTTPException( + status_code=403, + detail=f"Insufficient scope. Required: {required_scope}", + ) + return user user = SessionManager.decode_session(request.cookies.get("flowdeck_session", "")) if not user: raise HTTPException(status_code=401, detail="Authentication required") @@ -32,9 +48,10 @@ async def sync_delta( request: Request, since: float = Query(default=0, description="Epoch seconds (ou ms) du dernier sync"), workspace_id: int = Query(default=None), + authorization: str | None = Header(default=None), ): """Pull server-side changes since `since` (for the given workspace).""" - user = _user(request) + user = _user(request, authorization, required_scope="read") if workspace_id is None: raise HTTPException(status_code=400, detail="workspace_id is required") result = await _engine.get_delta(user["id"], since, workspace_id) @@ -44,9 +61,9 @@ async def sync_delta( @router.post("/batch") -async def sync_batch(request: Request): +async def sync_batch(request: Request, authorization: str | None = Header(default=None)): """Apply a batch of offline mutations and return per-mutation results.""" - user = _user(request) + user = _user(request, authorization, required_scope="write") try: body = await request.json() except Exception: @@ -63,9 +80,10 @@ async def sync_batch(request: Request): @router.get("/status") -async def sync_status(request: Request, workspace_id: int = Query(default=None)): +async def sync_status(request: Request, workspace_id: int = Query(default=None), + authorization: str | None = Header(default=None)): """Synchronization status for the workspace (pending server queue, last sync).""" - user = _user(request) + user = _user(request, authorization, required_scope="read") from app.db import get_conn with get_conn() as conn: if not SyncEngine._can_access(conn, user["id"], workspace_id): diff --git a/app/services/webhook_outbound.py b/app/services/webhook_outbound.py index bbbe941..30f2510 100644 --- a/app/services/webhook_outbound.py +++ b/app/services/webhook_outbound.py @@ -1,7 +1,16 @@ -"""FlowDeck — Webhook outbound dispatcher (v2.1.0).""" +"""FlowDeck — Webhook outbound dispatcher (v2.1.0, prod v6.4.0). + +Production-grade delivery: HMAC-SHA256 signature (``X-FlowDeck-Signature``), +retry with backoff 2s/10s/60s, full delivery journal in ``webhook_deliveries``. +""" from __future__ import annotations +import asyncio +import hashlib +import hmac +import json import logging +import time import httpx @@ -9,29 +18,227 @@ from app.db import get_conn logger = logging.getLogger(__name__) +# ── Event catalogue (≈28 events) ──────────────────────────────────────────── EVENTS = [ - "page.created", "page.updated", "page.deleted", + # pages (block-editor documents) + "page.created", "page.updated", "page.deleted", "page.moved", + "page.restored", "page.locked", "page.unlocked", + "page.shared", "page.published", "page.unpublished", + # comments & collaboration + "comment.added", "comment.resolved", "mention.added", + # collections (databases) "collection.created", "collection.updated", "collection.deleted", - "comment.added", "page.moved", + "collection.page.created", "collection.page.updated", "collection.page.deleted", + "collection.view.created", + # sprints / tasks + "sprint.created", "sprint.updated", + # sharing / favorites + "favorite.added", "favorite.removed", + # automations & agent + "automation.fired", "agent.run.finished", + # generic + "ping", ] +# Retry schedule: delays in seconds between attempts (attempt 0 = immediate). +RETRY_DELAYS = (2, 10, 60) +MAX_ATTEMPTS = len(RETRY_DELAYS) + 1 # 1 initial + 3 retries + + +def sign_payload(secret: str, body: bytes) -> str: + """HMAC-SHA256 hex signature of the raw JSON body (``sha256=``).""" + return "sha256=" + hmac.new(secret.encode(), body, hashlib.sha256).hexdigest() + + +def verify_signature(secret: str, body: bytes, signature: str) -> bool: + """Constant-time check of an ``X-FlowDeck-Signature`` header value.""" + if not secret or not signature: + return False + expected = sign_payload(secret, body) + return hmac.compare_digest(expected, signature) + + +def _event_matches(sub_event: str, event: str) -> bool: + """Wildcard matching: ``*`` matches all, ``page.*`` matches page.* events.""" + if sub_event in ("*", "all"): + return True + if sub_event.endswith(".*"): + return event.startswith(sub_event[:-1]) + return sub_event == event + + +def _matching_subs(conn, event: str) -> list: + rows = conn.execute( + "SELECT id, url, event, secret FROM webhook_subscriptions WHERE active=1" + ).fetchall() + return [r for r in rows if _event_matches((r["event"] or "").strip(), event)] + + +def _log_delivery(webhook_id: int, event: str, status: str, http_code: int | None, + error: str, duration_ms: int, attempt: int, payload: dict, + next_retry_at: float | None = None) -> None: + try: + with get_conn() as conn: + conn.execute( + """INSERT INTO webhook_deliveries + (webhook_id, event, status, http_code, error, duration_ms, + attempt, payload, next_retry_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)""", + (webhook_id, event, status, + http_code if http_code is not None else 0, + (error or "")[:1000], duration_ms, attempt, + json.dumps(payload)[:8000], + next_retry_at), + ) + conn.commit() + except Exception: # noqa: BLE001 — logging must never break delivery + logger.debug("Failed to log webhook delivery", exc_info=True) + + +async def _deliver_once(client: httpx.AsyncClient, url: str, event: str, + payload: dict, secret: str) -> tuple[int | None, str]: + """Single POST attempt. Returns (http_code, error).""" + body = json.dumps({"event": event, **payload}).encode() + headers = { + "Content-Type": "application/json", + "X-FlowDeck-Event": event, + "User-Agent": "FlowDeck-Webhooks/1.0", + } + if secret: + headers["X-FlowDeck-Signature"] = sign_payload(secret, body) + # legacy header kept for backward compatibility + headers["X-FlowDeck-Secret"] = secret + try: + resp = await client.post(url, content=body, headers=headers) + if 200 <= resp.status_code < 300: + return resp.status_code, "" + return resp.status_code, f"HTTP {resp.status_code}" + except Exception as exc: # noqa: BLE001 + return None, str(exc)[:500] + + +async def deliver_to_sub(sub_id: int, url: str, event: str, payload: dict, + secret: str, *, _client: httpx.AsyncClient | None = None) -> bool: + """Deliver with retry (immediate + 2s/10s/60s). Logs every attempt. + + Returns True on success. Failures are re-queued via ``next_retry_at`` so + the background scheduler can pick them up even if this process restarts. + """ + own_client = _client is None + client = _client or httpx.AsyncClient(timeout=10) + try: + for attempt in range(MAX_ATTEMPTS): + if attempt > 0: + await asyncio.sleep(RETRY_DELAYS[attempt - 1]) + start = time.monotonic() + code, err = await _deliver_once(client, url, event, payload, secret or "") + duration = int((time.monotonic() - start) * 1000) + ok = code is not None and 200 <= code < 300 + last = attempt == MAX_ATTEMPTS - 1 + if ok or last: + _log_delivery(sub_id, event, "delivered" if ok else "failed", + code, err, duration, attempt, payload) + return ok + # schedule retry + delay = RETRY_DELAYS[attempt] + _log_delivery(sub_id, event, "retrying", code, err, duration, + attempt, payload, + next_retry_at=time.time() + delay) + return False + finally: + if own_client: + await client.aclose() + async def fire_event(event: str, payload: dict): - """Fire a webhook event to all registered subscribers.""" + """Fire a webhook event to all registered subscribers (wildcard-aware). + + Unknown events are ignored. Delivery itself never raises: each + subscription is attempted independently and journaled. + """ if event not in EVENTS: return with get_conn() as conn: - subs = conn.execute("SELECT url, secret FROM webhook_subscriptions WHERE event=?", (event,)).fetchall() + subs = _matching_subs(conn, event) + if not subs: + return async with httpx.AsyncClient(timeout=10) as client: for sub in subs: - url, secret = sub["url"], sub["secret"] - headers = {"Content-Type": "application/json", "X-FlowDeck-Event": event} - if secret: - headers["X-FlowDeck-Secret"] = secret try: - await client.post(url, json=payload, headers=headers) + await deliver_to_sub(sub["id"], sub["url"], event, + dict(payload), sub["secret"] or "", + _client=client) + except Exception: # noqa: BLE001 + logger.debug("Webhook delivery failed to %s", sub["url"]) + + +async def retry_due_deliveries(now: float | None = None) -> int: + """Re-fire deliveries stuck in ``retrying`` whose ``next_retry_at`` passed. + + Returns the number of deliveries retried. Called by the background + scheduler and directly testable. + """ + now = now if now is not None else time.time() + with get_conn() as conn: + rows = conn.execute( + """SELECT d.id, d.webhook_id, d.event, d.payload, d.attempt, + s.url, s.secret + FROM webhook_deliveries d + JOIN webhook_subscriptions s ON s.id = d.webhook_id + WHERE d.status='retrying' AND s.active=1 + AND d.next_retry_at IS NOT NULL AND d.next_retry_at <= ? + ORDER BY d.next_retry_at LIMIT 50""", + (now,), + ).fetchall() + if not rows: + return 0 + async with httpx.AsyncClient(timeout=10) as client: + for r in rows: + try: + payload = json.loads(r["payload"] or "{}") + if not isinstance(payload, dict): + payload = {"data": payload} except Exception: - logger.debug("Webhook delivery failed to %s", url) + payload = {} + event = r["event"] or "ping" + # resume where the previous run stopped, up to MAX_ATTEMPTS + start_attempt = min((r["attempt"] or 0) + 1, MAX_ATTEMPTS - 1) + for att in range(start_attempt, MAX_ATTEMPTS): + delay = RETRY_DELAYS[att - 1] if att > 0 else 0 + if delay: + await asyncio.sleep(delay) + t0 = time.monotonic() + code, err = await _deliver_once(client, r["url"], event, + payload, r["secret"] or "") + duration = int((time.monotonic() - t0) * 1000) + ok = code is not None and 200 <= code < 300 + last = att == MAX_ATTEMPTS - 1 + if ok or last: + _log_delivery(r["webhook_id"], event, + "delivered" if ok else "failed", + code, err, duration, att, payload) + break + _log_delivery(r["webhook_id"], event, "retrying", code, err, + duration, att, payload, + next_retry_at=time.time() + RETRY_DELAYS[att]) + with get_conn() as conn: + conn.execute("UPDATE webhook_deliveries SET status='superseded' WHERE id=?", (r["id"],)) + conn.commit() + return len(rows) + + +async def webhook_retry_scheduler(interval_seconds: int = 60) -> None: + """Background loop: retry due webhook deliveries every minute.""" + from app.config import settings as _settings + interval = getattr(_settings, "webhook_retry_interval_seconds", interval_seconds) + while True: + try: + await asyncio.sleep(interval) + await retry_due_deliveries() + except asyncio.CancelledError: + raise + except Exception: # noqa: BLE001 + logger.debug("webhook_retry_scheduler tick failed", exc_info=True) def init_webhook_tables(): diff --git a/app/templates/_page_editor_content.html b/app/templates/_page_editor_content.html index 3e33051..ae4d4f1 100644 --- a/app/templates/_page_editor_content.html +++ b/app/templates/_page_editor_content.html @@ -94,7 +94,27 @@
Invite
-
No matching users
+ +
No matching users
@@ -113,17 +133,22 @@
-
+
+
diff --git a/app/templates/_page_editor_scripts.html b/app/templates/_page_editor_scripts.html index 2a1c8d9..89b33e3 100644 --- a/app/templates/_page_editor_scripts.html +++ b/app/templates/_page_editor_scripts.html @@ -1023,7 +1023,7 @@ window.__fdEditorScriptsLoaded = true; generalAccess:'{{ page_share_mode }}',publicPerm:'viewer', pubAllowEdit:false,pubAllowComments:true, pageUrl:window.location.href,publishedUrl:'', - inviteEmail:'',invitePermission:'edit',inviteUserId:null,inviteUsers:[],inviteSuggestOpen:false,inviteSel:-1,accessList:[], + inviteEmail:'',invitePermission:'edit',inviteUserId:null,inviteGroupId:null,inviteGroups:[],inviteUsers:[],inviteSuggestOpen:false,inviteSel:-1,accessList:[], toastVisible:false,toastMsg:'', // ── v4.9.0: Collaboration — comments & mentions ── commentsOpen:false,commentDraft:'',comments:[],commentCount:0, @@ -1659,15 +1659,30 @@ applyAIBlocks(text){ }, shareInvite() { var self = this; + if (this.inviteGroupId) { + var csrf = document.cookie.match(/csrf_token=([^;]+)/); + fetch('/api/pages/' + this.pid + '/share', { + method: 'POST', + headers: { 'Content-Type': 'application/json', 'X-CSRF-Token': csrf ? csrf[1] : '' }, + body: JSON.stringify({ group_id: this.inviteGroupId, permission: this.invitePermission }) + }).then(function(r){ return r.json(); }).then(function(d){ + if (d.status === 'shared') { + self.inviteEmail = ''; self.inviteUserId = null; self.inviteGroupId = null; + self.inviteUsers = []; self.inviteGroups = []; self.inviteSuggestOpen = false; + self.showToast('Group shared'); self.loadShares(); + } else { self.showToast(d.detail || 'Share failed'); } + }).catch(function(){ self.showToast('Share failed'); }); + return; + } if (!this.inviteEmail.trim()) return; - var csrf = document.cookie.match(/csrf_token=([^;]+)/); + var csrf2 = document.cookie.match(/csrf_token=([^;]+)/); fetch('/api/pages/' + this.pid + '/share', { method: 'POST', - headers: { 'Content-Type': 'application/json', 'X-CSRF-Token': csrf ? csrf[1] : '' }, + headers: { 'Content-Type': 'application/json', 'X-CSRF-Token': csrf2 ? csrf2[1] : '' }, body: JSON.stringify({ user_id: this.inviteUserId, email: this.inviteEmail.trim(), permission: this.invitePermission }) }).then(function(r){ return r.json(); }).then(function(d){ if (d.status === 'shared') { - self.inviteEmail = ''; self.inviteUserId = null; self.inviteUsers = []; self.inviteSuggestOpen = false; + self.inviteEmail = ''; self.inviteUserId = null; self.inviteUsers = []; self.inviteGroups = []; self.inviteSuggestOpen = false; self.showToast('Invitation sent'); self.loadShares(); } else { self.showToast(d.detail || 'Share failed'); } }).catch(function(){ self.showToast('Share failed'); }); @@ -1675,7 +1690,7 @@ applyAIBlocks(text){ inviteInput(){ var self = this; var q = (this.inviteEmail || '').trim(); - if (!q) { this.inviteUsers = []; this.inviteSuggestOpen = false; this.inviteUserId = null; this.inviteSel = -1; return; } + if (!q) { this.inviteUsers = []; this.inviteGroups = []; this.inviteSuggestOpen = false; this.inviteUserId = null; this.inviteGroupId = null; this.inviteSel = -1; return; } clearTimeout(this._inviteTimer); this._inviteTimer = setTimeout(function(){ self._loadInviteUsers(q); }, 150); }, @@ -1689,9 +1704,20 @@ applyAIBlocks(text){ this.inviteUsers = (d.users || []).filter(function(u){ return u.id !== me && !shared[(u.login || '').toLowerCase()]; }); + // Groups: charger les groupes du workspace et filtrer + try { + const gr = await fetch('/api/v2/groups', { credentials: 'same-origin' }); + const gd = await gr.json(); + var sharedGroups = {}; + (this.accessList || []).forEach(function(a){ if (a.shared_with_group_id) sharedGroups[a.shared_with_group_id] = true; }); + var ql = (q || '').toLowerCase(); + this.inviteGroups = (gd.groups || []).filter(function(g){ + return !sharedGroups[g.id] && (!ql || (g.name || '').toLowerCase().indexOf(ql) >= 0); + }).slice(0, 8); + } catch(e2) { this.inviteGroups = []; } this.inviteSel = this.inviteUsers.length ? 0 : -1; this.inviteSuggestOpen = true; - } catch(e) { this.inviteUsers = []; this.inviteSuggestOpen = false; } + } catch(e) { this.inviteUsers = []; this.inviteGroups = []; this.inviteSuggestOpen = false; } }, inviteMove(dir){ var el = this.$refs.inviteSuggest; @@ -1713,8 +1739,17 @@ applyAIBlocks(text){ }, pickInviteUser(u){ this.inviteUserId = u.id; + this.inviteGroupId = null; this.inviteEmail = u.login || u.full_name || ''; - this.inviteUsers = []; this.inviteSuggestOpen = false; this.inviteSel = -1; + this.inviteUsers = []; this.inviteGroups = []; this.inviteSuggestOpen = false; this.inviteSel = -1; + var el = this.$refs.inviteInput; + if (el) el.focus(); + }, + pickInviteGroup(g){ + this.inviteGroupId = g.id; + this.inviteUserId = null; + this.inviteEmail = g.name || ''; + this.inviteUsers = []; this.inviteGroups = []; this.inviteSuggestOpen = false; this.inviteSel = -1; var el = this.$refs.inviteInput; if (el) el.focus(); }, diff --git a/tests/test_share_groups.py b/tests/test_share_groups.py new file mode 100644 index 0000000..6437bd8 --- /dev/null +++ b/tests/test_share_groups.py @@ -0,0 +1,211 @@ +"""FlowDeck — Partage de page par groupe (page_shares.shared_with_group_id). + +Couvre : création d'un partage groupe, upsert anti-doublon, +listage (kind=group, group_name), mise à jour de permission, +suppression, et erreurs 404 (page / groupe inconnu). +""" +from app.auth.session import SessionManager +from app.db import get_conn + + +def _token(user_id, login="user", is_admin=1): + return SessionManager.create_session( + {"id": user_id, "login": login, "full_name": login.title(), "is_admin": is_admin} + ) + + +def _ensure_owner(): + with get_conn() as conn: + row = conn.execute("SELECT id FROM users WHERE id=1").fetchone() + if not row: + conn.execute( + "INSERT INTO users (id, login, full_name, email, is_admin) " + "VALUES (1, 'owner', 'Owner', 'owner@t.dev', 1)" + ) + conn.commit() + + +def _as(client, user_id=1, login="owner"): + _ensure_owner() + client.cookies.set("flowdeck_session", _token(user_id, login)) + + +def _insert_page(title="Group Share Page"): + with get_conn() as conn: + cur = conn.execute( + "INSERT INTO pages (workspace, title, content, content_format, parent_section) " + "VALUES ('Private', ?, '[]', 'blocks', 'Private')", + (title,), + ) + conn.commit() + return cur.lastrowid + + +def _insert_group(name="Team"): + _ensure_owner() + with get_conn() as conn: + cur = conn.execute( + "INSERT INTO user_groups (workspace_id, name, description, created_by) " + "VALUES (NULL, ?, '', 1)", + (name,), + ) + conn.commit() + return cur.lastrowid + + +def test_share_page_with_group(client): + _as(client) + pid = _insert_page() + gid = _insert_group() + r = client.post(f"/api/pages/{pid}/share", json={"group_id": gid, "permission": "edit"}) + assert r.status_code == 200, r.text + d = r.json() + assert d["status"] == "shared" + assert d["shared_with_group_id"] == gid + + +def test_share_page_group_upsert_no_duplicate(client): + _as(client) + pid = _insert_page() + gid = _insert_group() + r1 = client.post(f"/api/pages/{pid}/share", json={"group_id": gid, "permission": "view"}) + r2 = client.post(f"/api/pages/{pid}/share", json={"group_id": gid, "permission": "comment"}) + assert r1.status_code == 200 and r2.status_code == 200 + assert r1.json()["id"] == r2.json()["id"] + with get_conn() as conn: + n = conn.execute( + "SELECT COUNT(*) AS c FROM page_shares WHERE page_id=? AND shared_with_group_id=?", + (pid, gid), + ).fetchone()["c"] + assert n == 1 + + +def test_list_shares_includes_group(client): + _as(client) + pid = _insert_page() + gid = _insert_group(name="Editors") + client.post(f"/api/pages/{pid}/share", json={"group_id": gid, "permission": "view"}) + r = client.get(f"/api/pages/{pid}/shares") + assert r.status_code == 200, r.text + shares = r.json()["shares"] + g = [s for s in shares if s.get("shared_with_group_id") == gid] + assert len(g) == 1 + assert g[0]["kind"] == "group" + assert g[0]["group_name"] == "Editors" + + +def test_update_group_share_permission(client): + _as(client) + pid = _insert_page() + gid = _insert_group() + sid = client.post(f"/api/pages/{pid}/share", json={"group_id": gid, "permission": "view"}).json()["id"] + r = client.put(f"/api/pages/{pid}/share/{sid}", json={"permission": "edit"}) + assert r.status_code == 200, r.text + assert r.json()["permission"] == "edit" + + +def test_remove_group_share(client): + _as(client) + pid = _insert_page() + gid = _insert_group() + sid = client.post(f"/api/pages/{pid}/share", json={"group_id": gid, "permission": "view"}).json()["id"] + r = client.delete(f"/api/pages/{pid}/share/{sid}") + assert r.status_code == 200, r.text + with get_conn() as conn: + row = conn.execute("SELECT id FROM page_shares WHERE id=?", (sid,)).fetchone() + assert row is None + + +def test_share_unknown_group_404(client): + _as(client) + pid = _insert_page() + r = client.post(f"/api/pages/{pid}/share", json={"group_id": 999999, "permission": "view"}) + assert r.status_code == 404 + + +def test_share_needs_target_400(client): + _as(client) + pid = _insert_page() + r = client.post(f"/api/pages/{pid}/share", json={"permission": "view"}) + assert r.status_code == 400 + + +def test_group_share_grants_effective_access(client): + """A member of an invited group sees the page via page_permissions mirror.""" + from app.services.permission_manager import PermissionManager + _as(client) + pid = _insert_page() + ws = _insert_workspace(1) + _bind_page_workspace(pid, ws) + gid = _insert_group(name="Readers") + member = _insert_user("reader", "reader@t.dev", "Reader") + with get_conn() as conn: + conn.execute("INSERT INTO group_members (group_id, user_id) VALUES (?, ?)", (gid, member)) + conn.commit() + r = client.post(f"/api/pages/{pid}/share", json={"group_id": gid, "permission": "edit"}) + assert r.status_code == 200, r.text + assert PermissionManager(member).can_edit_page(pid) + assert PermissionManager(member).can_view_page(pid) + + +def test_group_share_update_syncs_grant(client): + from app.services.permission_manager import PermissionManager + _as(client) + pid = _insert_page() + ws = _insert_workspace(1) + _bind_page_workspace(pid, ws) + gid = _insert_group(name="Writers") + member = _insert_user("writer", "writer@t.dev", "Writer") + with get_conn() as conn: + conn.execute("INSERT INTO group_members (group_id, user_id) VALUES (?, ?)", (gid, member)) + conn.commit() + sid = client.post(f"/api/pages/{pid}/share", json={"group_id": gid, "permission": "view"}).json()["id"] + assert PermissionManager(member).can_view_page(pid) + assert not PermissionManager(member).can_edit_page(pid) + client.put(f"/api/pages/{pid}/share/{sid}", json={"permission": "edit"}) + assert PermissionManager(member).can_edit_page(pid) + + +def test_group_share_remove_revokes_grant(client): + from app.services.permission_manager import PermissionManager + _as(client) + pid = _insert_page() + ws = _insert_workspace(1) + _bind_page_workspace(pid, ws) + gid = _insert_group(name="Temp") + member = _insert_user("temp", "temp@t.dev", "Temp") + with get_conn() as conn: + conn.execute("INSERT INTO group_members (group_id, user_id) VALUES (?, ?)", (gid, member)) + conn.commit() + sid = client.post(f"/api/pages/{pid}/share", json={"group_id": gid, "permission": "edit"}).json()["id"] + assert PermissionManager(member).can_edit_page(pid) + client.delete(f"/api/pages/{pid}/share/{sid}") + assert not PermissionManager(member).can_view_page(pid) + + +def _insert_workspace(owner_id, name="WS"): + _ensure_owner() + with get_conn() as conn: + cur = conn.execute("INSERT INTO workspaces (name, owner_id) VALUES (?, ?)", (name, owner_id)) + conn.commit() + return cur.lastrowid + + +def _bind_page_workspace(page_id, workspace_id): + with get_conn() as conn: + conn.execute( + "UPDATE pages SET workspace_id=?, permission_type='restricted' WHERE id=?", + (workspace_id, page_id), + ) + conn.commit() + + +def _insert_user(login, email, full_name): + _ensure_owner() + with get_conn() as conn: + cur = conn.execute( + "INSERT INTO users (login, email, full_name, is_active) VALUES (?,?,?,1)", + (login, email, full_name), + ) + conn.commit() + return cur.lastrowid