Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions app/core/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions app/schema/request/web/event.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
class EventCreate(BaseModel):
name: str
event_date: datetime
end_date: Optional[datetime] = None
status: Optional[str] = "draft"

class JoinEventRequest(BaseModel):
Expand Down
2 changes: 2 additions & 0 deletions app/schema/response/web/event.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Expand Down
1 change: 1 addition & 0 deletions app/service/event.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
)
Expand Down
Empty file.
36 changes: 36 additions & 0 deletions app/worker/event_lifecycle/main.py
Original file line number Diff line number Diff line change
@@ -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())
15 changes: 9 additions & 6 deletions db/generated/event_participant.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Expand Down Expand Up @@ -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]:
Expand Down
60 changes: 49 additions & 11 deletions db/generated/events.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
"""


Expand All @@ -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)
Expand Down Expand Up @@ -95,21 +116,32 @@ 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
"""


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
Expand All @@ -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:
Expand All @@ -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]:
Expand All @@ -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]:
Expand All @@ -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]:
Expand All @@ -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]:
Expand All @@ -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],
)
1 change: 1 addition & 0 deletions db/generated/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
9 changes: 5 additions & 4 deletions db/queries/event_participant.sql
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
24 changes: 20 additions & 4 deletions db/queries/events.sql
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -47,5 +47,21 @@ WHERE id = $1
RETURNING *;

-- name: DeleteEvent :exec
DELETE FROM events
WHERE id = $1;
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;
1 change: 1 addition & 0 deletions makefile
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
1 change: 1 addition & 0 deletions migrations/sql/down/add_end_date_to_events.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
ALTER TABLE events DROP COLUMN end_date;
1 change: 1 addition & 0 deletions migrations/sql/up/add_end_date_to_events.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
ALTER TABLE events ADD COLUMN end_date timestamptz;
27 changes: 27 additions & 0 deletions migrations/versions/afb0a93aa21e_add_end_date_to_events.py
Original file line number Diff line number Diff line change
@@ -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")
Loading