Quellcode durchsuchen

test: implement coedition

clovis vor 6 Tagen
Ursprung
Commit
2b1b8c7f86

+ 2 - 0
.env.example

@@ -6,6 +6,8 @@ REFRESH_TOKEN_EXPIRE_MINUTES=40320
 BACKEND_CORS_ORIGINS=["http://localhost:3000","http://localhost:8001"]
 ALLOWED_HOSTS=["localhost", "127.0.0.1"]
 
+REDIS_URL=redis://localhost:6379
+
 DEFAULT_DATABASE_HOSTNAME=localhost
 DEFAULT_DATABASE_USER=rDGJeEDqAz
 DEFAULT_DATABASE_PASSWORD=XsPQhCoEfOQZueDjsILetLDUvbvSxAMnrVtgVZpmdcSssUgbvs

+ 0 - 31
Dockerfile

@@ -1,31 +0,0 @@
-FROM python:3.11.1-slim-bullseye
-
-ENV PYTHONUNBUFFERED 1
-WORKDIR /build
-
-# Create venv, add it to path and install requirements
-RUN python -m venv /venv
-ENV PATH="/venv/bin:$PATH"
-
-COPY requirements.txt .
-RUN pip install -r requirements.txt
-
-# Install uvicorn server
-RUN pip install uvicorn[standard]
-
-# Copy the rest of app
-COPY app app
-COPY alembic alembic
-COPY alembic.ini .
-COPY pyproject.toml .
-COPY init.sh .
-
-# Create new user to run app process as unprivilaged user
-RUN addgroup --gid 1001 --system uvicorn && \
-    adduser --gid 1001 --shell /bin/false --disabled-password --uid 1001 uvicorn
-
-# Run init.sh script then start uvicorn
-RUN chown -R uvicorn:uvicorn /build
-CMD bash init.sh && \
-    runuser -u uvicorn -- /venv/bin/uvicorn app.main:app --app-dir /build --host 0.0.0.0 --port 8000 --workers 2 --loop uvloop
-EXPOSE 8000

+ 19 - 25
Readme.md

@@ -1,32 +1,27 @@
 # bdlg-2023-server
 
-This repository hosts an API that helps volunteer managers plan festival
-activities on a schedule, associate volunteers to time slots, and send SMS
-reminders during the event. Volunteer/slot data can also be bootstrapped
-from a Google Sheet import.
+This repository hosts an API that helps volunteer managers plan festival activities on a schedule, associate volunteers
+to time slots, and send SMS reminders during the event. Volunteer/slot data can also be bootstrapped from a Google Sheet
+import.
 
 ## Core concepts
 
-- **Organization** — the top-level tenant. Every user belongs to zero or
-  more organizations (multi-org membership is supported). Every project
-  belongs to exactly one organization, which defines who can access it.
+- **Organization** — the top-level tenant. Every user belongs to zero or more organizations (multi-org membership is
+  supported). Every project belongs to exactly one organization, which defines who can access it.
 - **Global roles** — `user` (default) or `super_admin`. Only `super_admin`
   can create/edit organizations and manage organization membership.
 - **Organization roles** (per user, per organization):
-    - `org_admin` — full control over the organization's projects and can
-      change member roles within their own organization.
-    - `respo_benevole` — manages volunteer allocation across the whole
-      project: slots, templates, volunteers, groups, and SMS.
-    - `respo_commission` — read access to the whole project plan; write
-      access limited to templates/slots belonging to the commission(s) they
-      are a member of. Cannot manage volunteers, groups, or SMS directly.
+    - `org_admin` — full control over the organization's projects and can change member roles within their own
+      organization.
+    - `respo_benevole` — manages volunteer allocation across the whole project: slots, templates, volunteers, groups,
+      and SMS.
+    - `respo_commission` — read access to the whole project plan; write access limited to templates/slots belonging to
+      the commission (s) they are a member of. Cannot manage volunteers, groups, or SMS directly.
 - **Commission** (Pôle) — represents a team/area within a project (e.g.
-  "Bar", "Accueil"). A commission's contact info is derived from its
-  members' own profile (`name` + `phone_number`) rather than stored as
-  free text — see `SlotTemplate.responsible_override` for the manual
-  exception case.
-- **Volunteer group** — a saved, searchable set of volunteers within a
-  project, usable for bulk slot assignment and group SMS.
+  "Bar", "Accueil"). A commission's contact info is derived from its members' own profile (`name` + `phone_number`)
+  rather than stored as free text — see `SlotTemplate.responsible_override` for the manual exception case.
+- **Volunteer group** — a saved, searchable set of volunteers within a project, usable for bulk slot assignment and
+  group SMS.
 
 ## Getting started
 
@@ -60,13 +55,12 @@ Run specific tests
 ## Update requirements
 
 ```
-> poetry lock
-> poetry export -f requirements.txt --output requirements.txt --without-hashes
-> poetry export -f requirements.txt --output requirements-dev.txt --without-hashes --with dev
+poetry lock
+poetry export -f requirements.txt --output requirements.txt --without-hashes
+poetry export -f requirements.txt --output requirements-dev.txt --without-hashes --with dev
 ```
 
-Regenerate after any backend route or schema change — see the frontend
-repo's README for details.
+Regenerate after any backend route or schema change — see the frontend repo's README for details.
 
 ## Credit
 

+ 43 - 0
app/api/SSEBroadcasterRoute.py

@@ -0,0 +1,43 @@
+import asyncio
+import json
+
+from fastapi import Request, Response
+from fastapi.routing import APIRoute
+
+from app.core.events import broadcast
+
+
+class SSEBroadcasterRoute(APIRoute):
+    def get_route_handler(self):
+        original_route_handler = super().get_route_handler()
+
+        async def custom_route_handler(request: Request) -> Response:
+            # 1. Execute the route and let FastAPI serialize the response
+            response = await original_route_handler(request)
+
+            openapi_extra = self.openapi_extra or {}
+            event_name = openapi_extra.get("sse_event")
+
+            # 2. Check if an event is declared and the request was successful
+            if event_name and response.status_code in (200, 201):
+                project_id = request.path_params.get("project_id")
+                client_id = request.headers.get("X-Client-ID", "unknown")
+
+                if request.method == "DELETE":
+                    # For DELETE, the response body is often empty. We broadcast the path parameters.
+                    item_json = json.dumps(request.path_params)
+                else:
+                    # For POST/PUT, response.body is already serialized as a JSON string by FastAPI
+                    item_json = response.body.decode("utf-8")
+
+                # Assemble the raw payload string
+                payload_str = (
+                    f'{{"client_id": "{client_id}", "event": "{event_name}", "item": {item_json}}}'
+                )
+
+                channel = f"project_{project_id}"
+                asyncio.create_task(broadcast.publish(channel, payload_str))
+
+            return response
+
+        return custom_route_handler

+ 2 - 0
app/api/api.py

@@ -5,6 +5,7 @@ from app.api.endpoints import (
     commissions,
     organizations,
     project,
+    project_stream,
     slots,
     sms,
     sms_sender,
@@ -20,6 +21,7 @@ api_router.include_router(auth.router, prefix="/auth", tags=["auth"])
 api_router.include_router(users.router, prefix="/users", tags=["users"])
 api_router.include_router(organizations.router, prefix="/organizations", tags=["organization"])
 api_router.include_router(project.router, tags=["project"])
+api_router.include_router(project_stream.router, tags=["project"])
 api_router.include_router(commissions.router, tags=["commissions"])
 api_router.include_router(slots.router, tags=["slot"])
 api_router.include_router(tags.router, tags=["tag"])

+ 25 - 5
app/api/endpoints/commissions.py

@@ -48,7 +48,11 @@ async def list_project_commissions(
     return results.scalars().all()
 
 
-@router.post("/commission", response_model=CommissionResponse)
+@router.post(
+    "/commission",
+    response_model=CommissionResponse,
+    openapi_extra={"sse_event": "commission_created"},
+)
 async def create_commission(
     project_id: UUID,
     new_commission: CommissionCreateRequest,
@@ -75,7 +79,11 @@ async def get_commission(
     return _get_commission_or_404(session, project_id, commission_id)
 
 
-@router.post("/commission/{commission_id}", response_model=CommissionResponse)
+@router.post(
+    "/commission/{commission_id}",
+    response_model=CommissionResponse,
+    openapi_extra={"sse_event": "commission_updated"},
+)
 async def update_commission(
     project_id: UUID,
     commission_id: UUID,
@@ -91,7 +99,10 @@ async def update_commission(
     return commission
 
 
-@router.delete("/commission/{commission_id}")
+@router.delete(
+    "/commission/{commission_id}",
+    openapi_extra={"sse_event": "commission_deleted"},
+)
 async def delete_commission(
     project_id: UUID,
     commission_id: UUID,
@@ -104,7 +115,11 @@ async def delete_commission(
     session.commit()
 
 
-@router.post("/commission/{commission_id}/members", response_model=CommissionResponse)
+@router.post(
+    "/commission/{commission_id}/members",
+    response_model=CommissionResponse,
+    openapi_extra={"sse_event": "commission_updated"},
+)
 async def add_members_to_commission(
     project_id: UUID,
     commission_id: UUID,
@@ -129,7 +144,11 @@ async def add_members_to_commission(
     return _get_commission_or_404(session, project_id, commission_id)
 
 
-@router.delete("/commission/{commission_id}/member/{user_id}", response_model=CommissionResponse)
+@router.delete(
+    "/commission/{commission_id}/member/{user_id}",
+    response_model=CommissionResponse,
+    openapi_extra={"sse_event": "commission_updated"},
+)
 async def remove_member_from_commission(
     project_id: UUID,
     commission_id: UUID,
@@ -152,6 +171,7 @@ async def remove_member_from_commission(
 @router.post(
     "/commission/{commission_id}/invite-member",
     response_model=CommissionResponse,
+    openapi_extra={"sse_event": "commission_updated"},
 )
 async def invite_commission_member(
     project_id: UUID,

+ 12 - 3
app/api/endpoints/project.py

@@ -22,6 +22,7 @@ from app.models import (
     Volunteer,
     VolunteerGroup,
 )
+from app.schemas.objects import PlanningConstraints
 from app.schemas.requests import (
     ProjectConstraintUpdateRequest,
     ProjectCreateRequest,
@@ -128,7 +129,11 @@ async def get_project(
     return project
 
 
-@router.post("/project/{project_id}", response_model=ProjectListResponse)
+@router.post(
+    "/project/{project_id}",
+    response_model=ProjectListResponse,
+    openapi_extra={"sse_event": "project_updated"},
+)
 async def update_project(
     project_id: UUID,
     edit_project: ProjectUpdateRequest,
@@ -143,7 +148,11 @@ async def update_project(
     return p
 
 
-@router.post("/project/{project_id}/constraints", response_model=ProjectResponse)
+@router.post(
+    "/project/{project_id}/constraints",
+    openapi_extra={"sse_event": "project_constraints_updated"},
+    response_model=PlanningConstraints,
+)
 async def update_project_constraints(
     project_id: UUID,
     payload: ProjectConstraintUpdateRequest,
@@ -156,7 +165,7 @@ async def update_project_constraints(
     project.constraints = payload.model_dump()
     session.commit()
     session.refresh(project)
-    return project
+    return project.constraints
 
 
 @router.post("/project/{project_id}/import-gsheet", response_model=ProjectResponse)

+ 49 - 0
app/api/endpoints/project_stream.py

@@ -0,0 +1,49 @@
+import json
+
+from fastapi import APIRouter, Depends, Request
+from fastapi.responses import StreamingResponse
+
+from app.api import deps
+from app.core.events import broadcast
+from app.models import OrgRole, User
+
+router = APIRouter()
+
+
+@router.get("/project/{project_id}/stream", summary="Real-Time SSE Stream")
+async def project_stream(
+    project_id: str,
+    client_id: str,
+    request: Request,
+    current_user: User = Depends(
+        deps.require_org_role(OrgRole.ORG_ADMIN, OrgRole.RESPO_BENEVOLE, OrgRole.RESPO_COMMISSION)
+    ),
+):
+    """
+    Connect to this endpoint to receive real-time updates.
+
+    **Required Query Parameters:**
+    - `client_id`: A unique UUID generated by the frontend to prevent echoing.
+    """
+
+    async def event_generator():
+        channel = f"project_{project_id}"
+
+        # Subscribe to the Redis channel specific to this project
+        async with broadcast.subscribe(channel=channel) as subscriber:
+            async for event in subscriber:
+                # If the user closes the tab or navigates away, cleanly close the connection
+                if await request.is_disconnected():
+                    break
+                # Parse the raw string payload from Redis
+                message = json.loads(event.message)
+
+                # Echo prevention: Do not send the event back to the client that caused it
+                if message.get("client_id") == client_id:
+                    continue
+
+                # Yield the properly formatted Server-Sent Event (SSE)
+                yield f"event: {message['event']}\ndata: {json.dumps(message['item'])}\n\n"
+
+    # Return as text/event-stream so the browser knows it's an ongoing connection
+    return StreamingResponse(event_generator(), media_type="text/event-stream")

+ 9 - 4
app/api/endpoints/slots.py

@@ -5,6 +5,7 @@ from sqlalchemy import delete, select
 from sqlalchemy.orm import Session
 
 from app.api import deps
+from app.api.SSEBroadcasterRoute import SSEBroadcasterRoute
 from app.api.utils import assert_project_exists_or_404, update_object_from_payload, verify_id_list
 from app.models import (
     OrgRole,
@@ -20,7 +21,9 @@ from app.schemas.requests import (
 )
 from app.schemas.responses import SlotResponse
 
-router = APIRouter(prefix="/project/{project_id}", tags=["project"])
+router = APIRouter(
+    route_class=SSEBroadcasterRoute, prefix="/project/{project_id}", tags=["project"]
+)
 
 READ_ROLES = (OrgRole.ORG_ADMIN, OrgRole.RESPO_BENEVOLE, OrgRole.RESPO_COMMISSION)
 WRITE_ROLES = (OrgRole.ORG_ADMIN, OrgRole.RESPO_BENEVOLE, OrgRole.RESPO_COMMISSION)
@@ -51,7 +54,7 @@ async def list_project_slots(
     return results.scalars().all()
 
 
-@router.post("/slot", response_model=SlotResponse)
+@router.post("/slot", response_model=SlotResponse, openapi_extra={"sse_event": "slot_created"})
 async def create_slot(
     project_id: UUID,
     new_slot: SlotCreateRequest,
@@ -85,7 +88,9 @@ async def create_slot(
     return slot
 
 
-@router.post("/slot/{slot_id}", response_model=SlotResponse)
+@router.post(
+    "/slot/{slot_id}", response_model=SlotResponse, openapi_extra={"sse_event": "slot_updated"}
+)
 async def update_slot(
     project_id: UUID,
     slot_id: UUID,
@@ -132,7 +137,7 @@ async def update_slot(
     return slot
 
 
-@router.delete("/slot/{slot_id}")
+@router.delete("/slot/{slot_id}", openapi_extra={"sse_event": "slot_deleted"})
 async def delete_slot(
     project_id: UUID,
     slot_id: UUID,

+ 5 - 3
app/api/endpoints/tags.py

@@ -34,7 +34,7 @@ async def list_project_tags(
     return p.tags
 
 
-@router.post("/tag", response_model=TagResponse)
+@router.post("/tag", response_model=TagResponse, openapi_extra={"sse_event": "tag_created"})
 async def create_tag(
     project_id: UUID,
     payload: TagCreateRequest,
@@ -67,7 +67,9 @@ async def create_tag(
     return tag
 
 
-@router.post("/tag/{tag_id}", response_model=TagResponse)
+@router.post(
+    "/tag/{tag_id}", response_model=TagResponse, openapi_extra={"sse_event": "tag_updated"}
+)
 async def update_tag(
     project_id: UUID,
     tag_id: UUID,
@@ -126,7 +128,7 @@ async def list_tagged_slot(
     return list(slot_set)
 
 
-@router.delete("/tag/{tag_id}")
+@router.delete("/tag/{tag_id}", openapi_extra={"sse_event": "tag_deleted"})
 async def delete_tag(
     project_id: UUID,
     tag_id: UUID,

+ 9 - 3
app/api/endpoints/templates.py

@@ -37,7 +37,9 @@ async def list_project_templates(
     return p.templates
 
 
-@router.post("/template", response_model=TemplateResponse)
+@router.post(
+    "/template", response_model=TemplateResponse, openapi_extra={"sse_event": "template_created"}
+)
 async def create_template(
     project_id: UUID,
     payload: TemplateCreateRequest,
@@ -70,7 +72,11 @@ async def create_template(
     return template
 
 
-@router.post("/template/{template_id}", response_model=TemplateResponse)
+@router.post(
+    "/template/{template_id}",
+    response_model=TemplateResponse,
+    openapi_extra={"sse_event": "template_updated"},
+)
 async def update_template(
     project_id: UUID,
     template_id: UUID,
@@ -113,7 +119,7 @@ async def update_template(
     return template
 
 
-@router.delete("/template/{template_id}")
+@router.delete("/template/{template_id}", openapi_extra={"sse_event": "template_deleted"})
 async def delete_template(
     project_id: UUID,
     template_id: UUID,

+ 30 - 35
app/api/endpoints/volunteer_groups.py

@@ -9,14 +9,14 @@ from app.api.utils import (  # adjust import path if different
     assert_project_exists_or_404,
     verify_id_list,
 )
-from app.models import OrgRole, Slot, Sms, User, Volunteer, VolunteerGroup
+from app.models import OrgRole, Sms, User, Volunteer, VolunteerGroup
 from app.schemas.requests import (
     GroupMembershipRequest,
     GroupSmsRequest,
     VolunteerGroupCreateRequest,
     VolunteerGroupUpdateRequest,
 )
-from app.schemas.responses import SMSResponse, VolunteerGroupResponse, VolunteerResponse
+from app.schemas.responses import SMSResponse, VolunteerGroupResponse
 
 router = APIRouter(prefix="/project/{project_id}", tags=["volunteer_groups"])
 
@@ -42,7 +42,11 @@ async def list_project_groups(
     return results.scalars().all()
 
 
-@router.post("/group", response_model=VolunteerGroupResponse)
+@router.post(
+    "/group",
+    response_model=VolunteerGroupResponse,
+    openapi_extra={"sse_event": "volunteer_group_created"},
+)
 async def create_group(
     project_id: UUID,
     new_group: VolunteerGroupCreateRequest,
@@ -58,7 +62,10 @@ async def create_group(
     return group
 
 
-@router.get("/group/{group_id}", response_model=VolunteerGroupResponse)
+@router.get(
+    "/group/{group_id}",
+    response_model=VolunteerGroupResponse,
+)
 async def get_group(
     project_id: UUID,
     group_id: UUID,
@@ -69,7 +76,11 @@ async def get_group(
     return _get_group_or_404(session, project_id, group_id)
 
 
-@router.post("/group/{group_id}", response_model=VolunteerGroupResponse)
+@router.post(
+    "/group/{group_id}",
+    response_model=VolunteerGroupResponse,
+    openapi_extra={"sse_event": "volunteer_group_updated"},
+)
 async def update_group(
     project_id: UUID,
     group_id: UUID,
@@ -86,7 +97,10 @@ async def update_group(
     return group
 
 
-@router.delete("/group/{group_id}")
+@router.delete(
+    "/group/{group_id}",
+    openapi_extra={"sse_event": "volunteer_group_deleted"},
+)
 async def delete_group(
     project_id: UUID,
     group_id: UUID,
@@ -99,7 +113,11 @@ async def delete_group(
     session.commit()
 
 
-@router.post("/group/{group_id}/volunteers", response_model=VolunteerGroupResponse)
+@router.post(
+    "/group/{group_id}/volunteers",
+    response_model=VolunteerGroupResponse,
+    openapi_extra={"sse_event": "volunteer_group_updated"},
+)
 async def add_volunteers_to_group(
     project_id: UUID,
     group_id: UUID,
@@ -126,7 +144,11 @@ async def add_volunteers_to_group(
     return group
 
 
-@router.delete("/group/{group_id}/volunteer/{volunteer_id}", response_model=VolunteerGroupResponse)
+@router.delete(
+    "/group/{group_id}/volunteer/{volunteer_id}",
+    response_model=VolunteerGroupResponse,
+    openapi_extra={"sse_event": "volunteer_group_updated"},
+)
 async def remove_volunteer_from_group(
     project_id: UUID,
     group_id: UUID,
@@ -144,33 +166,6 @@ async def remove_volunteer_from_group(
     return group
 
 
-@router.post("/group/{group_id}/add-to-slot/{slot_id}", response_model=list[VolunteerResponse])
-async def add_group_to_slot(
-    project_id: UUID,
-    group_id: UUID,
-    slot_id: UUID,
-    current_user: User = Depends(deps.require_org_role(*MANAGE_GROUPS)),
-    session: Session = Depends(deps.get_session),
-):
-    """Bulk-assigns every member of the group to the slot in one call.
-    This is a one-time copy, not a live link: adding someone to the group
-    later does NOT retroactively add them to slots already assigned."""
-    group = _get_group_or_404(session, project_id, group_id)
-
-    slot = session.get(Slot, slot_id)
-    if slot is None or slot.project_id != str(project_id):
-        raise HTTPException(status_code=404, detail="Slot not found in this project")
-
-    existing_ids = {v.id for v in slot.volunteers}
-    for volunteer in group.volunteers:
-        if volunteer.id not in existing_ids:
-            slot.volunteers.append(volunteer)
-
-    session.commit()
-    session.refresh(slot)
-    return slot.volunteers
-
-
 @router.post("/group/{group_id}/send-sms", response_model=list[SMSResponse])
 async def send_sms_to_group(
     project_id: UUID,

+ 12 - 3
app/api/endpoints/volunteers.py

@@ -29,7 +29,9 @@ async def list_project_volunteers(
     return results.scalars().all()
 
 
-@router.post("/volunteer", response_model=VolunteerResponse)
+@router.post(
+    "/volunteer", response_model=VolunteerResponse, openapi_extra={"sse_event": "volunteer_created"}
+)
 async def create_volunteer(
     project_id: UUID,
     new_volunteer: VolunteerCreateRequest,
@@ -60,7 +62,11 @@ async def create_volunteer(
     return volunteer
 
 
-@router.post("/volunteer/{volunteer_id}", response_model=VolunteerResponse)
+@router.post(
+    "/volunteer/{volunteer_id}",
+    response_model=VolunteerResponse,
+    openapi_extra={"sse_event": "volunteer_updated"},
+)
 async def update_volunteer(
     project_id: UUID,
     volunteer_id: UUID,
@@ -96,7 +102,10 @@ async def update_volunteer(
     return volunteer
 
 
-@router.delete("/volunteer/{volunteer_id}")
+@router.delete(
+    "/volunteer/{volunteer_id}",
+    openapi_extra={"sse_event": "volunteer_deleted"},
+)
 async def delete_volunteer(
     project_id: UUID,
     volunteer_id: UUID,

+ 3 - 0
app/core/config.py

@@ -64,6 +64,9 @@ class Settings(BaseSettings):
     EMAIL_FROM_ADDRESS: str = "no-reply@example.com"
     EMAIL_FROM_NAME: str = "BDLG Planner"
 
+    # REDIS
+    REDIS_URL: str = "redis://localhost:6379"
+
     # POSTGRESQL DEFAULT DATABASE
     DEFAULT_DATABASE_HOSTNAME: str
     DEFAULT_DATABASE_USER: str

+ 6 - 0
app/core/events.py

@@ -0,0 +1,6 @@
+from broadcaster import Broadcast
+
+from app.core.config import settings
+
+# Instantiate the broadcaster using the URL from your environment variables
+broadcast = Broadcast(settings.REDIS_URL)

+ 99 - 0
app/main.py

@@ -1,11 +1,24 @@
 """Main FastAPI app instance declaration."""
 
+import typing
+from contextlib import asynccontextmanager
+
 from fastapi import FastAPI
 from fastapi.middleware.cors import CORSMiddleware
 from fastapi.middleware.trustedhost import TrustedHostMiddleware
+from fastapi.openapi.utils import get_openapi
 
 from app.api.api import api_router
 from app.core import config
+from app.core.events import broadcast
+
+
+@asynccontextmanager
+async def lifespan(app: FastAPI):
+    await broadcast.connect()
+    yield
+    await broadcast.disconnect()
+
 
 app = FastAPI(
     title=config.settings.PROJECT_NAME,
@@ -13,6 +26,7 @@ app = FastAPI(
     description=config.settings.DESCRIPTION,
     openapi_url="/openapi.json",
     docs_url="/",
+    lifespan=lifespan,
 )
 app.include_router(api_router)
 
@@ -28,3 +42,88 @@ app.add_middleware(
 # Guards against HTTP Host Header attacks
 if config.settings.ENVIRONMENT != "PYTEST":
     app.add_middleware(TrustedHostMiddleware, allowed_hosts=config.settings.ALLOWED_HOSTS)
+
+
+def custom_openapi():
+    if app.openapi_schema:
+        return app.openapi_schema
+
+    openapi_schema = get_openapi(
+        title=config.settings.PROJECT_NAME,
+        version=config.settings.VERSION,
+        description=config.settings.DESCRIPTION,
+        routes=app.routes,
+    )
+
+    one_of_schemas = []
+
+    for route in app.routes:
+        if hasattr(route, "openapi_extra") and route.openapi_extra:
+            event_name = route.openapi_extra.get("sse_event")
+            if not event_name:
+                continue
+
+            method = list(route.methods)[0] if route.methods else "GET"
+            payload_schema = {}
+
+            # Determine payload schema for DELETE (Path params)
+            model = getattr(route, "response_model", None)
+            if method == "DELETE" and model is None:
+                props = {}
+                for param in getattr(route.dependant, "path_params", []):
+                    props[param.name] = {"type": "string"}
+                payload_schema = {"type": "object", "properties": props}
+
+            # Determine payload schema for POST/PUT (Response Model)
+            else:
+                origin = typing.get_origin(model)
+                if origin is list:
+                    item_model = typing.get_args(model)[0]
+                    if hasattr(item_model, "__name__"):
+                        model_name = item_model.__name__
+                        payload_schema = {
+                            "type": "array",
+                            "items": {"$ref": f"#/components/schemas/{model_name}"},
+                        }
+                elif hasattr(model, "__name__"):
+                    model_name = model.__name__
+                    payload_schema = {"$ref": f"#/components/schemas/{model_name}"}
+
+            # Build the discriminated union object with traceability in the description
+            event_schema = {
+                "type": "object",
+                "description": f"Triggered by `{method}` `{route.path}`",
+                "properties": {
+                    "event": {"type": "string", "enum": [event_name]},
+                    "data": payload_schema,
+                },
+                "required": ["event", "data"],
+            }
+            one_of_schemas.append(event_schema)
+
+    if one_of_schemas:
+        if "components" not in openapi_schema:
+            openapi_schema["components"] = {"schemas": {}}
+        elif "schemas" not in openapi_schema["components"]:
+            openapi_schema["components"]["schemas"] = {}
+
+        openapi_schema["components"]["schemas"]["SseEventPayload"] = {
+            "title": "SseEventPayload",
+            "description": "Discriminated union of all possible SSE events and their payloads.",
+            "oneOf": one_of_schemas,
+        }
+
+        for path, path_item in openapi_schema["paths"].items():
+            if path.endswith("/stream") and "get" in path_item:
+                path_item["get"]["responses"]["200"]["content"] = {
+                    "text/event-stream": {
+                        "schema": {"$ref": "#/components/schemas/SseEventPayload"}
+                    }
+                }
+
+    app.openapi_schema = openapi_schema
+    return app.openapi_schema
+
+
+# 5. Override the default OpenAPI method
+app.openapi = custom_openapi

+ 86 - 0
app/tests/test_projects_stream.py

@@ -0,0 +1,86 @@
+import json
+from unittest.mock import AsyncMock, patch
+
+import pytest
+from httpx import AsyncClient
+
+# Adjust this import to match where you instantiated your broadcast object!
+from app.main import app
+from app.models import OrgRole, Project
+from app.tests.conftest import default_project_id
+
+pytestmark = pytest.mark.asyncio
+
+
+class TestProjectStream:
+    @pytest.mark.parametrize(
+        "role, expected_status",
+        [
+            (OrgRole.ORG_ADMIN, 200),
+            (OrgRole.RESPO_BENEVOLE, 200),
+            (OrgRole.RESPO_COMMISSION, 200),
+            (None, 403),
+        ],
+    )
+    async def test_role_access(
+        self, client: AsyncClient, default_project: Project, make_org_user, role, expected_status
+    ):
+        """Test that only authorized roles can connect to the stream."""
+        _, headers = make_org_user(role=role)
+        url = app.url_path_for("project_stream", project_id=default_project_id)
+
+        # We mock 'broadcast.subscribe' so it immediately returns an empty async iterator
+        # This prevents the endpoint from hanging forever in the test
+        with patch("app.core.events.broadcast.subscribe") as mock_subscribe:
+            # Create an async generator that yields nothing and closes
+            async def mock_generator():
+                return
+                yield
+
+            mock_subscriber = AsyncMock()
+            mock_subscriber.__aenter__.return_value = mock_generator()
+            mock_subscribe.return_value = mock_subscriber
+
+            async with client.stream(
+                "GET", f"{url}?client_id=test-auth", headers=headers
+            ) as response:
+                assert response.status_code == expected_status
+
+    async def test_receives_sse_events(
+        self, client: AsyncClient, default_project: Project, make_org_user
+    ):
+        """Test that events from the broadcaster are correctly formatted as SSE."""
+        _, headers = make_org_user(role=OrgRole.ORG_ADMIN)
+        url = app.url_path_for("project_stream", project_id=default_project_id)
+
+        # Create a fake Redis event
+        class FakeEvent:
+            def __init__(self):
+                self.message = json.dumps(
+                    {
+                        "client_id": "someone-else",
+                        "event": "slot_updated",
+                        "item": {"id": "123", "title": "Installation"},
+                    }
+                )
+
+        with patch("app.core.events.broadcast.subscribe") as mock_subscribe:
+            # Make the mocked subscribe yield our FakeEvent, then stop
+            async def mock_generator():
+                yield FakeEvent()
+
+            mock_subscriber = AsyncMock()
+            mock_subscriber.__aenter__.return_value = mock_generator()
+            mock_subscribe.return_value = mock_subscriber
+
+            async with client.stream(
+                "GET", f"{url}?client_id=listener-client", headers=headers
+            ) as response:
+                assert response.status_code == 200
+
+                # Read the response
+                lines = [line async for line in response.aiter_lines() if line.strip()]
+
+                # Verify the formatting is standard SSE format
+                assert lines[0] == "event: slot_updated"
+                assert "Installation" in lines[1]

+ 1 - 78
app/tests/test_volunteer_groups.py

@@ -1,5 +1,4 @@
 import uuid
-from datetime import datetime, timedelta
 
 import pytest
 from httpx import AsyncClient
@@ -8,7 +7,7 @@ from sqlalchemy.orm import Session
 
 from app.core.session import session as session_maker
 from app.main import app
-from app.models import Organization, OrgRole, Project, Slot, Volunteer, VolunteerGroup
+from app.models import Organization, OrgRole, Project, Volunteer, VolunteerGroup
 from app.tests.conftest import default_project_id, default_slot_id
 from app.tests.shared_access import SharedProjectAccessTests
 
@@ -64,7 +63,6 @@ VOLUNTEER_GROUP_ROUTES = [
     ("DELETE", "delete_group", route_kwargs_2, None),
     ("POST", "add_volunteers_to_group", route_kwargs_2, {"volunteer_ids": []}),
     ("DELETE", "remove_volunteer_from_group", {**route_kwargs_2, "volunteer_id": "VOL"}, None),
-    ("POST", "add_group_to_slot", {**route_kwargs_2, "slot_id": "SLOT"}, None),
     ("POST", "send_sms_to_group", route_kwargs_2, {"content": "coucou"}),
 ]
 
@@ -362,81 +360,6 @@ class TestGroupMembership:
         assert response.json()["volunteers_id"] == [v2.id]
 
 
-class TestAddGroupToSlot:
-    async def test_bulk_assigns_all_group_members(
-        self,
-        client: AsyncClient,
-        default_project: Project,
-        default_group: VolunteerGroup,
-        two_volunteers,
-        make_org_user,
-        session: Session,
-    ):
-        v1, v2 = two_volunteers
-        group = session.get(VolunteerGroup, default_group.id)
-        group.volunteers.append(session.get(Volunteer, v1.id))
-        group.volunteers.append(session.get(Volunteer, v2.id))
-        slot = Slot(
-            project_id=default_project.id,
-            title="Garde du Graal",
-            starting_time=datetime.now() + timedelta(hours=1),
-            ending_time=datetime.now() + timedelta(hours=2),
-        )
-        session.add(slot)
-        session.commit()
-
-        _, headers = make_org_user(role=OrgRole.ORG_ADMIN)
-        response = await client.post(
-            app.url_path_for(
-                "add_group_to_slot",
-                project_id=default_project.id,
-                group_id=default_group.id,
-                slot_id=slot.id,
-            ),
-            headers=headers,
-        )
-        assert response.status_code == 200
-        ids = [v["id"] for v in response.json()]
-        assert sorted(ids) == sorted([v1.id, v2.id])
-
-    async def test_slot_from_other_project_not_found(
-        self,
-        client: AsyncClient,
-        default_project: Project,
-        default_group: VolunteerGroup,
-        make_org_user,
-        session: Session,
-    ):
-        other_org = Organization(id=str(uuid.uuid4()), name="Other Org")
-        session.add(other_org)
-        session.commit()
-        other_project = Project(
-            name="Other Project 3", is_public=False, organization_id=other_org.id
-        )
-        session.add(other_project)
-        session.commit()
-        stray_slot = Slot(
-            project_id=other_project.id,
-            title="Stray slot",
-            starting_time=datetime.now(),
-            ending_time=datetime.now() + timedelta(hours=1),
-        )
-        session.add(stray_slot)
-        session.commit()
-
-        _, headers = make_org_user(role=OrgRole.ORG_ADMIN)
-        response = await client.post(
-            app.url_path_for(
-                "add_group_to_slot",
-                project_id=default_project.id,
-                group_id=default_group.id,
-                slot_id=stray_slot.id,
-            ),
-            headers=headers,
-        )
-        assert response.status_code == 404
-
-
 class TestSendSmsToGroup:
     async def test_sends_to_each_member_with_automatic_sms(
         self,

+ 0 - 36
docker-compose.dev.yml

@@ -1,36 +0,0 @@
-version: "3.7"
-
-# Database + Webserver (under http, for testing setup on localhost:80)
-#
-# docker-compose -f docker-compose.dev.yml up -d
-#
-
-services:
-  postgres:
-    restart: unless-stopped
-    image: postgres:latest
-    volumes:
-      - postgres_data:/var/lib/postgresql/data
-    env_file:
-      - .env
-    environment:
-      - POSTGRES_DB=${DEFAULT_DATABASE_DB}
-      - POSTGRES_USER=${DEFAULT_DATABASE_USER}
-      - POSTGRES_PASSWORD=${DEFAULT_DATABASE_PASSWORD}
-  web:
-    depends_on:
-      - postgres
-    restart: "unless-stopped"
-    build:
-      context: ./
-      dockerfile: Dockerfile
-    env_file:
-      - .env
-    environment:
-      - DEFAULT_DATABASE_HOSTNAME=postgres
-      - DEFAULT_DATABASE_PORT=5432
-    ports:
-      - 80:8000
-
-volumes:
-  postgres_data:

+ 0 - 40
docker-compose.yml

@@ -1,40 +0,0 @@
-version: "3.7"
-
-# For local development, only database is running
-#
-# docker-compose up -d
-# uvicorn app.main:app --reload
-#
-
-services:
-  default_database:
-    restart: unless-stopped
-    image: postgres:latest
-    volumes:
-      - default_database_data:/var/lib/postgresql/data
-    environment:
-      - POSTGRES_DB=${DEFAULT_DATABASE_DB}
-      - POSTGRES_USER=${DEFAULT_DATABASE_USER}
-      - POSTGRES_PASSWORD=${DEFAULT_DATABASE_PASSWORD}
-    env_file:
-      - .env
-    ports:
-      - "${DEFAULT_DATABASE_PORT}:5432"
-
-  test_database:
-    restart: unless-stopped
-    image: postgres:latest
-    volumes:
-      - test_database_data:/var/lib/postgresql/data
-    environment:
-      - POSTGRES_DB=${TEST_DATABASE_DB}
-      - POSTGRES_USER=${TEST_DATABASE_USER}
-      - POSTGRES_PASSWORD=${TEST_DATABASE_PASSWORD}
-    env_file:
-      - .env
-    ports:
-      - "${TEST_DATABASE_PORT}:5432"
-
-volumes:
-  test_database_data:
-  default_database_data:

+ 59 - 1
poetry.lock

@@ -68,6 +68,19 @@ typing_extensions = {version = ">=4.5", markers = "python_version < \"3.13\""}
 [package.extras]
 trio = ["trio (>=0.26.1)"]
 
+[[package]]
+name = "async-timeout"
+version = "5.0.1"
+description = "Timeout context manager for asyncio programs"
+optional = false
+python-versions = ">=3.8"
+groups = ["main"]
+markers = "python_full_version < \"3.11.3\""
+files = [
+    {file = "async_timeout-5.0.1-py3-none-any.whl", hash = "sha256:39e3809566ff85354557ec2398b55e096c8364bacac9405a7a1fa429e77fe76c"},
+    {file = "async_timeout-5.0.1.tar.gz", hash = "sha256:d9321a7a3d5a6a5e187e824d2fa0793ce379a202935782d555d6e9d2735677d3"},
+]
+
 [[package]]
 name = "bcrypt"
 version = "4.0.1"
@@ -103,6 +116,28 @@ files = [
 tests = ["pytest (>=3.2.1,!=3.3.0)"]
 typecheck = ["mypy"]
 
+[[package]]
+name = "broadcaster"
+version = "0.3.1"
+description = "Simple broadcast channels."
+optional = false
+python-versions = ">=3.8"
+groups = ["main"]
+files = [
+    {file = "broadcaster-0.3.1-py3-none-any.whl", hash = "sha256:433023ab6b6b4a8da9cbba95910eff52b1e767141419659be287cfd49f2a3ecb"},
+    {file = "broadcaster-0.3.1.tar.gz", hash = "sha256:35d1e174a98346d184cc6558979548340e877ae841c65c7e354153a84fe675d6"},
+]
+
+[package.dependencies]
+anyio = ">=3.4.0,<5"
+redis = {version = "*", optional = true, markers = "extra == \"redis\""}
+
+[package.extras]
+kafka = ["aiokafka"]
+postgres = ["asyncpg"]
+redis = ["redis"]
+test = ["pytest", "pytest-asyncio"]
+
 [[package]]
 name = "certifi"
 version = "2025.8.3"
@@ -1568,6 +1603,29 @@ files = [
     {file = "pyyaml-6.0.2.tar.gz", hash = "sha256:d584d9ec91ad65861cc08d42e834324ef890a082e591037abe114850ff7bbc3e"},
 ]
 
+[[package]]
+name = "redis"
+version = "8.0.1"
+description = "Python client for Redis database and key-value store"
+optional = false
+python-versions = ">=3.10"
+groups = ["main"]
+files = [
+    {file = "redis-8.0.1-py3-none-any.whl", hash = "sha256:47daa35a058c23468d6437f17a8c76882cb316b838ef763036af99b96cedd743"},
+    {file = "redis-8.0.1.tar.gz", hash = "sha256:afc5a7a2f5a084f5b1880dec548dd45be17db7e43c82a30d84f952aefb05cfb0"},
+]
+
+[package.dependencies]
+async-timeout = {version = ">=4.0.3", markers = "python_full_version < \"3.11.3\""}
+
+[package.extras]
+circuit-breaker = ["pybreaker (>=1.4.0)"]
+hiredis = ["hiredis (>=3.2.0)"]
+jwt = ["pyjwt (>=2.13.0)"]
+ocsp = ["cryptography (>=36.0.1)", "pyopenssl (>=20.0.1)", "requests (>=2.31.0)"]
+otel = ["opentelemetry-api (>=1.39.1)", "opentelemetry-exporter-otlp-proto-http (>=1.39.1)", "opentelemetry-sdk (>=1.39.1)"]
+xxhash = ["xxhash (>=3.6.0,<3.7.0)"]
+
 [[package]]
 name = "requests"
 version = "2.32.5"
@@ -2238,4 +2296,4 @@ files = [
 [metadata]
 lock-version = "2.1"
 python-versions = "^3.11"
-content-hash = "4283413c272bc2cdda2f7d7c8ee1dff14f83335daed0fa7afda5c289d9fa34b4"
+content-hash = "cd4255985724c76c2c60d622c0f550423b604d4b8973bf34737e8990989a9494"

+ 1 - 0
pyproject.toml

@@ -22,6 +22,7 @@ psycopg2 = "^2.9.9"
 ruff = "^0.16.0"
 bcrypt = "4.0.1"
 aiosmtplib = "^5.1.2"
+broadcaster = {extras = ["redis"], version = "^0.3.1"}
 
 [tool.poetry.group.dev.dependencies]
 coverage = "^7.1.0"

+ 4 - 0
requirements-dev.txt

@@ -1,7 +1,10 @@
+aiosmtplib==5.1.2 ; python_version >= "3.11" and python_version < "4.0"
 alembic==1.16.5 ; python_version >= "3.11" and python_version < "4.0"
 annotated-types==0.7.0 ; python_version >= "3.11" and python_version < "4.0"
 anyio==4.10.0 ; python_version >= "3.11" and python_version < "4.0"
+async-timeout==5.0.1 ; python_version >= "3.11" and python_full_version < "3.11.3"
 bcrypt==4.0.1 ; python_version >= "3.11" and python_version < "4.0"
+broadcaster==0.3.1 ; python_version >= "3.11" and python_version < "4.0"
 certifi==2025.8.3 ; python_version >= "3.11" and python_version < "4.0"
 cffi==2.0.0 ; python_version >= "3.11" and python_version < "4.0" and platform_python_implementation != "PyPy"
 cfgv==3.4.0 ; python_version >= "3.11" and python_version < "4.0"
@@ -49,6 +52,7 @@ python-dotenv==1.1.1 ; python_version >= "3.11" and python_version < "4.0"
 python-multipart==0.0.9 ; python_version >= "3.11" and python_version < "4.0"
 pytz==2025.2 ; python_version >= "3.11" and python_version < "4.0"
 pyyaml==6.0.2 ; python_version >= "3.11" and python_version < "4.0"
+redis==8.0.1 ; python_version >= "3.11" and python_version < "4.0"
 requests==2.32.5 ; python_version >= "3.11" and python_version < "4.0"
 ruff==0.16.0 ; python_version >= "3.11" and python_version < "4.0"
 six==1.17.0 ; python_version >= "3.11" and python_version < "4.0"

+ 4 - 0
requirements.txt

@@ -1,7 +1,10 @@
+aiosmtplib==5.1.2 ; python_version >= "3.11" and python_version < "4.0"
 alembic==1.16.5 ; python_version >= "3.11" and python_version < "4.0"
 annotated-types==0.7.0 ; python_version >= "3.11" and python_version < "4.0"
 anyio==4.10.0 ; python_version >= "3.11" and python_version < "4.0"
+async-timeout==5.0.1 ; python_version >= "3.11" and python_full_version < "3.11.3"
 bcrypt==4.0.1 ; python_version >= "3.11" and python_version < "4.0"
+broadcaster==0.3.1 ; python_version >= "3.11" and python_version < "4.0"
 certifi==2025.8.3 ; python_version >= "3.11" and python_version < "4.0"
 cffi==2.0.0 ; python_version >= "3.11" and python_version < "4.0" and platform_python_implementation != "PyPy"
 charset-normalizer==3.4.3 ; python_version >= "3.11" and python_version < "4.0"
@@ -30,6 +33,7 @@ python-dateutil==2.9.0.post0 ; python_version >= "3.11" and python_version < "4.
 python-dotenv==1.1.1 ; python_version >= "3.11" and python_version < "4.0"
 python-multipart==0.0.9 ; python_version >= "3.11" and python_version < "4.0"
 pytz==2025.2 ; python_version >= "3.11" and python_version < "4.0"
+redis==8.0.1 ; python_version >= "3.11" and python_version < "4.0"
 requests==2.32.5 ; python_version >= "3.11" and python_version < "4.0"
 ruff==0.16.0 ; python_version >= "3.11" and python_version < "4.0"
 six==1.17.0 ; python_version >= "3.11" and python_version < "4.0"