From 59c904c9b0bccfe77e8e5dc109f4362b4e37e5a4 Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 01:39:40 +0100 Subject: [PATCH] feat: automate event lifecycle transitions on start/end time - Add events.end_date (nullable timestamptz) via migration. - Add ActivateDueEvents / ArchiveEndedEvents queries: draft -> scheduled when event_date has passed, scheduled -> archived (+ archived_at) when end_date has passed. - Add a polling worker (app/worker/event_lifecycle) that runs both transitions every EVENT_LIFECYCLE_POLL_INTERVAL_SECONDS (default 60s), wired into make run-workers. - Thread end_date through EventCreate/EventResponse/UserEventResponse and CreateEvent/GetUserEvents. Events with no end_date set are never auto-archived and stay in scheduled until archived manually via the existing endpoint. --- app/core/config.py | 1 + app/schema/request/web/event.py | 1 + app/schema/response/web/event.py | 2 + app/service/event.py | 1 + app/worker/event_lifecycle/__init__.py | 0 app/worker/event_lifecycle/main.py | 36 +++++++++++ db/generated/event_participant.py | 15 +++-- db/generated/events.py | 60 +++++++++++++++---- db/generated/models.py | 1 + db/queries/event_participant.sql | 9 +-- db/queries/events.sql | 24 ++++++-- makefile | 1 + .../sql/down/add_end_date_to_events.sql | 1 + migrations/sql/up/add_end_date_to_events.sql | 1 + .../afb0a93aa21e_add_end_date_to_events.py | 27 +++++++++ 15 files changed, 155 insertions(+), 25 deletions(-) create mode 100644 app/worker/event_lifecycle/__init__.py create mode 100644 app/worker/event_lifecycle/main.py create mode 100644 migrations/sql/down/add_end_date_to_events.sql create mode 100644 migrations/sql/up/add_end_date_to_events.sql create mode 100644 migrations/versions/afb0a93aa21e_add_end_date_to_events.py diff --git a/app/core/config.py b/app/core/config.py index 35ec0daf..c9270f03 100644 --- a/app/core/config.py +++ b/app/core/config.py @@ -34,6 +34,7 @@ class Settings(BaseSettings): POSTGRES_PORT: int = 5432 PHOTO_APPROVAL_TIMEOUT_DAYS: int = 7 + EVENT_LIFECYCLE_POLL_INTERVAL_SECONDS: int = 60 # Mobile auth/session defaults MOBILE_SESSION_LIMIT: int = 3 diff --git a/app/schema/request/web/event.py b/app/schema/request/web/event.py index d03b46ff..c7c05e3d 100644 --- a/app/schema/request/web/event.py +++ b/app/schema/request/web/event.py @@ -5,6 +5,7 @@ class EventCreate(BaseModel): name: str event_date: datetime + end_date: Optional[datetime] = None status: Optional[str] = "draft" class JoinEventRequest(BaseModel): diff --git a/app/schema/response/web/event.py b/app/schema/response/web/event.py index 4334fc44..01a64266 100644 --- a/app/schema/response/web/event.py +++ b/app/schema/response/web/event.py @@ -10,6 +10,7 @@ class EventResponse(BaseModel): name: str event_code: str event_date: datetime + end_date: Optional[datetime] = None status: str created_by: uuid.UUID created_at: datetime @@ -30,6 +31,7 @@ class UserEventResponse(BaseModel): id: uuid.UUID name: str event_date: datetime + end_date: Optional[datetime] = None status: str joined_at: datetime diff --git a/app/service/event.py b/app/service/event.py index e697ddd8..fb6165e9 100644 --- a/app/service/event.py +++ b/app/service/event.py @@ -33,6 +33,7 @@ async def create_event(self, req: EventCreate, creator_id: uuid.UUID) -> EventRe name=req.name, event_code=code_created, event_date=req.event_date, + end_date=req.end_date, status=req.status or "draft", created_by=creator_id ) diff --git a/app/worker/event_lifecycle/__init__.py b/app/worker/event_lifecycle/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/app/worker/event_lifecycle/main.py b/app/worker/event_lifecycle/main.py new file mode 100644 index 00000000..7ab11b8d --- /dev/null +++ b/app/worker/event_lifecycle/main.py @@ -0,0 +1,36 @@ +import asyncio + +from app.core.config import settings +from app.core.logger import logger +from app.infra.database import engine +from db.generated import events as event_queries + + +async def run_lifecycle_pass() -> None: + async with engine.begin() as conn: + querier = event_queries.AsyncQuerier(conn) + + activated = [event_id async for event_id in querier.activate_due_events()] + if activated: + logger.info("event_lifecycle: activated %d event(s): %s", len(activated), activated) + + archived = [event_id async for event_id in querier.archive_ended_events()] + if archived: + logger.info("event_lifecycle: archived %d event(s): %s", len(archived), archived) + + +async def main() -> None: + logger.info( + "Event lifecycle worker starting, poll_interval=%ds", + settings.EVENT_LIFECYCLE_POLL_INTERVAL_SECONDS, + ) + while True: + try: + await run_lifecycle_pass() + except Exception: + logger.exception("event_lifecycle: pass failed") + await asyncio.sleep(settings.EVENT_LIFECYCLE_POLL_INTERVAL_SECONDS) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/db/generated/event_participant.py b/db/generated/event_participant.py index 65129c77..2ad09031 100644 --- a/db/generated/event_participant.py +++ b/db/generated/event_participant.py @@ -40,10 +40,11 @@ class GetEventParticipantsRow: GET_USER_EVENTS = """-- name: get_user_events \\:many -SELECT - e.id, - e.name, - e.event_date, +SELECT + e.id, + e.name, + e.event_date, + e.end_date, e.status, ep.joined_at FROM events e @@ -58,6 +59,7 @@ class GetUserEventsRow: id: uuid.UUID name: str event_date: datetime.datetime + end_date: Optional[datetime.datetime] status: Any joined_at: datetime.datetime @@ -164,8 +166,9 @@ async def get_user_events(self, *, user_id: uuid.UUID) -> AsyncIterator[GetUserE id=row[0], name=row[1], event_date=row[2], - status=row[3], - joined_at=row[4], + end_date=row[3], + status=row[4], + joined_at=row[5], ) async def is_user_in_event(self, *, event_id: uuid.UUID, user_id: uuid.UUID) -> Optional[bool]: diff --git a/db/generated/events.py b/db/generated/events.py index 0395bfda..1cd52749 100644 --- a/db/generated/events.py +++ b/db/generated/events.py @@ -13,10 +13,30 @@ from db.generated import models +ACTIVATE_DUE_EVENTS = """-- name: activate_due_events \\:many +UPDATE events +SET status = 'scheduled'\\:\\:event_status +WHERE status = 'draft'\\:\\:event_status + AND event_date <= NOW() +RETURNING id +""" + + +ARCHIVE_ENDED_EVENTS = """-- name: archive_ended_events \\:many +UPDATE events +SET status = 'archived'\\:\\:event_status, + archived_at = NOW() +WHERE status = 'scheduled'\\:\\:event_status + AND end_date IS NOT NULL + AND end_date <= NOW() +RETURNING id +""" + + CREATE_EVENT = """-- name: create_event \\:one -INSERT INTO events (name, event_code, event_date, status, created_by) -VALUES (:p1, :p2, :p3, :p4, :p5) -RETURNING id, name, event_code, event_date, status, created_by, created_at, archived_at +INSERT INTO events (name, event_code, event_date, end_date, status, created_by) +VALUES (:p1, :p2, :p3, :p4, :p5, :p6) +RETURNING id, name, event_code, event_date, status, created_by, created_at, archived_at, end_date """ @@ -25,37 +45,38 @@ class CreateEventParams: name: str event_code: str event_date: datetime.datetime + end_date: Optional[datetime.datetime] status: Any created_by: uuid.UUID DELETE_EVENT = """-- name: delete_event \\:exec -DELETE FROM events +DELETE FROM events WHERE id = :p1 """ GET_EVENT_BY_CODE = """-- name: get_event_by_code \\:one -SELECT id, name, event_code, event_date, status, created_by, created_at, archived_at FROM events +SELECT id, name, event_code, event_date, status, created_by, created_at, archived_at, end_date FROM events WHERE event_code = :p1 """ GET_EVENT_BY_ID = """-- name: get_event_by_id \\:one -SELECT id, name, event_code, event_date, status, created_by, created_at, archived_at FROM events +SELECT id, name, event_code, event_date, status, created_by, created_at, archived_at, end_date FROM events WHERE id = :p1 """ GET_EVENTS_BY_NAME = """-- name: get_events_by_name \\:many -SELECT id, name, event_code, event_date, status, created_by, created_at, archived_at FROM events +SELECT id, name, event_code, event_date, status, created_by, created_at, archived_at, end_date FROM events WHERE name ILIKE '%' || :p1 || '%' ORDER BY event_date DESC """ LIST_EVENTS = """-- name: list_events \\:many -SELECT id, name, event_code, event_date, status, created_by, created_at, archived_at FROM events +SELECT id, name, event_code, event_date, status, created_by, created_at, archived_at, end_date FROM events WHERE -- Filter by Status (Optional) (:p3\\:\\:event_status IS NULL OR status = :p3) @@ -95,7 +116,7 @@ class ListEventsParams: SET status = :p2, archived_at = CASE WHEN :p2 = 'archived'\\:\\:event_status THEN NOW() ELSE archived_at END WHERE id = :p1 -RETURNING id, name, event_code, event_date, status, created_by, created_at, archived_at +RETURNING id, name, event_code, event_date, status, created_by, created_at, archived_at, end_date """ @@ -103,13 +124,24 @@ class AsyncQuerier: def __init__(self, conn: sqlalchemy.ext.asyncio.AsyncConnection): self._conn = conn + async def activate_due_events(self) -> AsyncIterator[uuid.UUID]: + result = await self._conn.stream(sqlalchemy.text(ACTIVATE_DUE_EVENTS)) + async for row in result: + yield row[0] + + async def archive_ended_events(self) -> AsyncIterator[uuid.UUID]: + result = await self._conn.stream(sqlalchemy.text(ARCHIVE_ENDED_EVENTS)) + async for row in result: + yield row[0] + async def create_event(self, arg: CreateEventParams) -> Optional[models.Event]: row = (await self._conn.execute(sqlalchemy.text(CREATE_EVENT), { "p1": arg.name, "p2": arg.event_code, "p3": arg.event_date, - "p4": arg.status, - "p5": arg.created_by, + "p4": arg.end_date, + "p5": arg.status, + "p6": arg.created_by, })).first() if row is None: return None @@ -122,6 +154,7 @@ async def create_event(self, arg: CreateEventParams) -> Optional[models.Event]: created_by=row[5], created_at=row[6], archived_at=row[7], + end_date=row[8], ) async def delete_event(self, *, id: uuid.UUID) -> None: @@ -140,6 +173,7 @@ async def get_event_by_code(self, *, event_code: str) -> Optional[models.Event]: created_by=row[5], created_at=row[6], archived_at=row[7], + end_date=row[8], ) async def get_event_by_id(self, *, id: uuid.UUID) -> Optional[models.Event]: @@ -155,6 +189,7 @@ async def get_event_by_id(self, *, id: uuid.UUID) -> Optional[models.Event]: created_by=row[5], created_at=row[6], archived_at=row[7], + end_date=row[8], ) async def get_events_by_name(self, *, dollar_1: Optional[str]) -> AsyncIterator[models.Event]: @@ -169,6 +204,7 @@ async def get_events_by_name(self, *, dollar_1: Optional[str]) -> AsyncIterator[ created_by=row[5], created_at=row[6], archived_at=row[7], + end_date=row[8], ) async def list_events(self, arg: ListEventsParams) -> AsyncIterator[models.Event]: @@ -191,6 +227,7 @@ async def list_events(self, arg: ListEventsParams) -> AsyncIterator[models.Event created_by=row[5], created_at=row[6], archived_at=row[7], + end_date=row[8], ) async def update_event_status(self, *, id: uuid.UUID, status: Any) -> Optional[models.Event]: @@ -206,4 +243,5 @@ async def update_event_status(self, *, id: uuid.UUID, status: Any) -> Optional[m created_by=row[5], created_at=row[6], archived_at=row[7], + end_date=row[8], ) diff --git a/db/generated/models.py b/db/generated/models.py index 21bac799..86617af5 100644 --- a/db/generated/models.py +++ b/db/generated/models.py @@ -74,6 +74,7 @@ class Event: created_by: uuid.UUID created_at: datetime.datetime archived_at: Optional[datetime.datetime] + end_date: Optional[datetime.datetime] @dataclasses.dataclass() diff --git a/db/queries/event_participant.sql b/db/queries/event_participant.sql index 9203c83d..42de3e66 100644 --- a/db/queries/event_participant.sql +++ b/db/queries/event_participant.sql @@ -6,10 +6,11 @@ RETURNING *; -- name: GetUserEvents :many -- Retrieves all events a specific user has successfully joined -SELECT - e.id, - e.name, - e.event_date, +SELECT + e.id, + e.name, + e.event_date, + e.end_date, e.status, ep.joined_at FROM events e diff --git a/db/queries/events.sql b/db/queries/events.sql index e7a7fd84..5e1fde41 100644 --- a/db/queries/events.sql +++ b/db/queries/events.sql @@ -1,6 +1,6 @@ -- name: CreateEvent :one -INSERT INTO events (name, event_code, event_date, status, created_by) -VALUES ($1, $2, $3, $4, $5) +INSERT INTO events (name, event_code, event_date, end_date, status, created_by) +VALUES ($1, $2, $3, $4, $5, $6) RETURNING *; -- name: GetEventById :one @@ -47,5 +47,21 @@ WHERE id = $1 RETURNING *; -- name: DeleteEvent :exec -DELETE FROM events -WHERE id = $1; \ No newline at end of file +DELETE FROM events +WHERE id = $1; + +-- name: ActivateDueEvents :many +UPDATE events +SET status = 'scheduled'::event_status +WHERE status = 'draft'::event_status + AND event_date <= NOW() +RETURNING id; + +-- name: ArchiveEndedEvents :many +UPDATE events +SET status = 'archived'::event_status, + archived_at = NOW() +WHERE status = 'scheduled'::event_status + AND end_date IS NOT NULL + AND end_date <= NOW() +RETURNING id; \ No newline at end of file diff --git a/makefile b/makefile index 241764ae..36690ed8 100644 --- a/makefile +++ b/makefile @@ -65,6 +65,7 @@ run-workers: uv run python -m app.worker.photo_worker.main & \ uv run python -m app.worker.storage_cleaner.main & \ uv run python -m app.worker.email_worker.main & \ + uv run python -m app.worker.event_lifecycle.main & \ wait lint: diff --git a/migrations/sql/down/add_end_date_to_events.sql b/migrations/sql/down/add_end_date_to_events.sql new file mode 100644 index 00000000..708425cc --- /dev/null +++ b/migrations/sql/down/add_end_date_to_events.sql @@ -0,0 +1 @@ +ALTER TABLE events DROP COLUMN end_date; diff --git a/migrations/sql/up/add_end_date_to_events.sql b/migrations/sql/up/add_end_date_to_events.sql new file mode 100644 index 00000000..5f2f8e72 --- /dev/null +++ b/migrations/sql/up/add_end_date_to_events.sql @@ -0,0 +1 @@ +ALTER TABLE events ADD COLUMN end_date timestamptz; diff --git a/migrations/versions/afb0a93aa21e_add_end_date_to_events.py b/migrations/versions/afb0a93aa21e_add_end_date_to_events.py new file mode 100644 index 00000000..9d7fdcd6 --- /dev/null +++ b/migrations/versions/afb0a93aa21e_add_end_date_to_events.py @@ -0,0 +1,27 @@ +"""add_end_date_to_events + +Revision ID: afb0a93aa21e +Revises: 9ec59aeb192b +Create Date: 2026-08-25 01:34:59.327934 + +""" +from typing import Sequence, Union + +from migrations.helper import run_sql_up, run_sql_down + + +# revision identifiers, used by Alembic. +revision: str = 'afb0a93aa21e' +down_revision: Union[str, Sequence[str], None] = '9ec59aeb192b' +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + """Upgrade schema.""" + run_sql_up("add_end_date_to_events") + + +def downgrade() -> None: + """Downgrade schema.""" + run_sql_down("add_end_date_to_events")