도메인 간 수정 전파를 중앙에서 처리하는 로직이 필요했다. 찾아보면 이 메커니즘은 계층에 따라 서로 다른 이름으로 불린다. 설계 개념으로는 Observer, Pub/Sub, Mediator가 있고, 구현물로는 Domain Event와 Event Bus, 분산 환경까지 확장하면 Outbox와 Saga가 등장한다. 다만 대상 환경이 단일 프로세스, 단일 DB, 요청당 단일 세션의 FastAPI 모놀리식이라면 선택지는 자연스럽게 좁혀진다. 외부 브로커 없이 동일 트랜잭션 안에서 동작하는, 인프로세스 도메인 이벤트와 동기 디스패처의 조합이다.
본 글은 이 조합을 실제로 구현한 과정을 다룬다. 이벤트를 어떻게 정의하고, 어디에 수집하고, 어떻게 디스패치하며, 트랜잭션 경계에 어떻게 결합하는지, 그리고 양방향 전파에서 따라오는 왕복 문제를 어떻게 차단했는지 순서대로 살펴본다.
프로젝트와 도메인 구조
대상 프로젝트는 일정·타이머·할일을 통합 관리하는 생산성 백엔드다. 핵심 도메인은 셋으로, 할일(Todo), 캘린더 일정(Schedule), 타이머(Timer)가 있고 여기에 태그와 공유 설정(visibility)이 걸쳐 있다.
이 글에서 계속 등장할 관계는 Todo와 Schedule 사이의 연동이다. Todo에 마감(deadline)을 설정하면 그 마감을 캘린더에 표시하기 위한 Schedule이 자동 생성된다. 이렇게 만들어진 일정을 이하 "투영"이라 부르며,
Schedule.source_todo_id FK로 원본 Todo와 연결된다. 마감을 옮기면 투영도 옮겨져야 하고, 할일 제목을 바꾸면 투영의 제목도 바뀌어야 하며, 반대로 캘린더에서 투영을 옮기면 마감도 따라와야 한다. 즉 한 도메인의 쓰기가 다른 도메인의 쓰기를 유발하는 관계이고, 이 전파를 어디서 처리할지가 이 글의 주제다.배경
이 동기화는 원래 서비스 본문에 직접 들어 있었다.
# app/domain/todo/service.py — 이벤트 도입 전 실제 코드 if deadline_updated: existing_schedules = schedule_crud.get_schedules_by_source_todo_id( self.session, todo_id, self.owner_id) if old_deadline and not new_deadline: schedule_service = ScheduleService(self.session, self.current_user) for schedule in existing_schedules: schedule_service.delete_schedule(schedule.id) elif not old_deadline and new_deadline: from datetime import timedelta todo_tags = self.get_todo_tags(todo.id) tag_ids = [tag.id for tag in todo_tags] if todo_tags else None schedule_data = ScheduleCreate( title=todo.title, description=todo.description, start_time=new_deadline, end_time=new_deadline + timedelta(hours=1), source_todo_id=todo.id, tag_ids=tag_ids, ) schedule_service = ScheduleService(self.session, self.current_user) schedule_service.create_schedule(schedule_data) elif old_deadline and new_deadline: if existing_schedules: ... # ScheduleUpdate 조립 후 update_schedule 호출, 15줄가량
TodoService가 ScheduleService를 직접 생성해 호출하는 구조라 두 도메인이 강하게 결합되고, 순환 import를 피하기 위한 메서드 내부 지연 import가 계속 늘어났다. 역방향(일정 → 할일) 전파는 구현 자체가 없었다. 전파 규칙이 서비스 본문 곳곳에 흩어져 있으니 새 전파를 추가하려면 매번 다른 서비스의 내부를 열어야 했다.라이브러리 도입도 먼저 검토했다. blinker가 요구사항에 가장 가까웠지만, 따져보니 blinker가 대신해 주는 부분은 타입-핸들러 매핑과 순회 호출 정도로 직접 구현해도 130줄 남짓이다. 반면 정작 어려운 부분인 트랜잭션 경계 결합, 세션 컨텍스트 주입, 왕복 차단은 라이브러리를 도입해도 그대로 남는다. 여기에 kwargs 기반 비타입 payload 탓에 이벤트 스키마가 코드에 드러나지 않는 점, weakref 특성상 바운드 메서드 구독이 GC로 조용히 사라질 수 있는 점까지 더해져, 직접 구현하는 쪽이 낫다고 판단했다.
이벤트 정의
이벤트는 명령이 아니라 이미 일어난 사실이다. "일정을 갱신하라"가 아니라 "deadline이 변경되었다"를 발행하고, 누가 구독하는지 발행자는 모른다.
# app/shared/events.py @dataclass(frozen=True, kw_only=True) class DomainEvent: owner_id: str # handler는 변경 대상의 owner_id가 이와 일치할 때만 수정 origin: EventOrigin # 이벤트를 촉발한 도메인 def event_key(self) -> tuple: """중복 dispatch 억제 키: 타입 + 주요 식별자""" return (type(self).__qualname__, self.owner_id) # app/domain/todo/events.py @dataclass(frozen=True, kw_only=True) class TodoEvent(DomainEvent): todo_id: UUID def event_key(self) -> tuple: return (type(self).__qualname__, str(self.todo_id)) @dataclass(frozen=True, kw_only=True) class TodoDeadlineChanged(TodoEvent): """deadline 값이 변경됨""" # payload는 todo_id뿐
공통 베이스는 frozen dataclass로 정의했다. 도메인별 이벤트는 각 도메인의
events.py에 두고, 베이스가 요구하는 두 필드(owner_id, origin)와 중복 억제 키를 상속한다.payload에 식별자만 남긴 것은 초안이 실제로 틀렸기 때문이다. 처음에는
TodoDeadlineSet에 제목과 deadline을 스냅샷으로 실었는데, update_todo에서 이벤트를 만드는 위치가 필드 반영보다 앞이라 제목과 deadline을 한 요청에 같이 바꾸면 handler가 이전 제목으로 일정을 만들었다. 그래서 식별자만 싣고 핸들러가 동일 트랜잭션에서 현재 상태를 조회해 투영하는 방식으로 바꿨더니, 핸들러가 항상 최신 상태를 읽게 되면서 순서 문제 자체가 사라졌다.그래도 스냅샷이 필요한 자리가 있다. 삭제 이벤트를 처음 만들었을 때 handler가
source_todo_id로 연결 일정을 조회했는데 결과가 빈 리스트였다. Schedule.source_todo_id가 ondelete=SET NULL이라, Todo 삭제가 flush되는 순간 FK가 이미 NULL로 바뀐 뒤였던 것이다. dispatch 시점에 원본이 사라지는 정보는 이벤트가 직접 들고 가야 한다.@dataclass(frozen=True, kw_only=True) class TodoDeleted(TodoEvent): linked_schedule_ids: Tuple[UUID, ...] # 삭제 직전 스냅샷
이벤트 수집
발행한 이벤트를 다른 도메인으로 바로 들고 뛰면 그게 다시 직접 호출이다. 일단 우체통에 넣어두고 정해진 시점에 한꺼번에 비운다. cosmic python 책의 구현은 엔티티에
.events 리스트를 붙이고 Unit of Work가 수거하는데, 이 코드베이스는 도메인 모델이 SQLModel ORM 클래스의 alias라 엔티티에 리스트를 붙이면 로드 시 초기화가 꼬였다. 대신 SQLAlchemy 세션이 제공하는 세션 수명의 dict인 session.info를 우체통으로 썼다. 요청당 세션이 하나라 세션 스코프가 곧 요청 스코프다.# app/shared/uow.py def collect_event(session, event): session.info.setdefault("domain_events", []).append(event) def drain_events(session): events = session.info.get("domain_events", []) session.info["domain_events"] = [] return events def clear_events(session): session.info["domain_events"] = []
이 선택에는 대가가 하나 따라온다.
session.info는 트랜잭션과 연동되지 않으므로 DB rollback이 발생해도 이벤트 큐는 자동으로 비워지지 않는다. rollback 경계에서 clear_events를 호출하지 않으면 취소된 변경의 이벤트가 큐에 남았다가 다음 커밋에서 처리되는 사고가 난다. 이를 막기 위한 호출 위치는 트랜잭션 경계 절에서 함께 다룬다.Dispatch: drain loop
디스패치 방식으로는
emit() 즉시 핸들러를 재귀 실행하는 쪽이 가장 단순하다. 그러나 이 도메인에는 태그 그룹 삭제가 그룹 내 Todo 삭제로, 다시 투영 일정 삭제와 타이머 FK 해제로 이어지는 3홉 규모의 연쇄가 실제로 존재한다. 재귀 방식에서 이런 연쇄를 처리하면 핸들러가 핸들러를 중첩 호출하게 되고, 순환 가드가 재진입 문제와 얽힌다. 그래서 큐 방식을 택했다. 핸들러가 낳은 후속 이벤트를 큐 뒤에 붙여 같은 루프에서 순서대로 처리하면, 가드와 flush와 재수집이 루프 한 곳에 모인다.# app/shared/messagebus.py MAX_DISPATCH_DEPTH = 10 HANDLERS: Dict[Type[DomainEvent], List[Handler]] = {} def handle(events, session): queue = list(events) seen: set = set() depth = 0 while queue: event = queue.pop(0) key = event.event_key() if key in seen: continue seen.add(key) if depth >= MAX_DISPATCH_DEPTH: raise MaxDispatchDepthExceededError(event) depth += 1 for handler in HANDLERS.get(type(event), []): handler(event, session) # 예외는 삼키지 않음, 전체 rollback session.flush() # 후속 handler가 앞선 변경을 조회 가능 queue.extend(drain_events(session)) # handler가 collect한 이벤트 재수집
event_key의 억제 범위는 한 번의 drain loop 안이다. 요청 하나의 use-case에서 같은 Todo의 deadline 변경 이벤트가 두 번 만들어질 경로는 없으므로, 실사용에서 눌리는 것은 handler 연쇄가 만든 왕복 이벤트뿐이다. 별개 요청은 별개 loop라 서로 영향이 없다.솔직히 적어둘 것이 두 가지 있다.
MAX_DISPATCH_DEPTH는 이름과 달리 재귀 깊이가 아니라 loop가 처리하는 총 이벤트 수 상한이다. 지금은 요청당 이벤트가 많아야 서너 개라 10으로 잡았지만, 대량 발행이 생기면 이름과 한도를 같이 손봐야 한다. queue.pop(0)도 마찬가지로, 상한이 10인 리스트라 deque로 바꿀 이유가 없다고 판단하고 그대로 뒀다.Handler 등록과 배선
핸들러는 데코레이터로 등록한다. 같은 (타입, 함수) 쌍은 한 번만 등록되므로 모듈이 여러 번 import되어도 안전하다.
_KNOWN_REGISTRATIONS라는 별도 명세 리스트가 있는 이유는 모듈 import가 프로세스당 한 번만 데코레이터를 실행한다는 제약 때문이다. 테스트가 registry를 초기화하고 나면 모듈을 다시 import해도 재등록이 일어나지 않으므로, 데코레이터가 명세를 기록해 두고 register_handlers()가 이를 재반영하도록 했다. def register(event_type): def decorator(fn): if (event_type, fn) not in _KNOWN_REGISTRATIONS: _KNOWN_REGISTRATIONS.append((event_type, fn)) _add_handler(event_type, fn) return fn return decorator def register_handlers(): import app.domain.schedule.event_handlers # noqa: F401 import app.domain.todo.event_handlers # noqa: F401 for event_type, fn in _KNOWN_REGISTRATIONS: # 재-import는 데코레이터를 _add_handler(event_type, fn) # 재실행하지 않으므로 재반영
이 함수는 앱 lifespan과 테스트 conftest 양쪽에서 부른다. 서비스 단위 테스트는 앱 엔트리포인트를 import하지 않아서, conftest 배선이 없으면 handler 없이 테스트가 돌고 동기화 누락을 아무도 눈치채지 못한다.
트랜잭션 경계
이 설계에서 가장 중요한 제약은 디스패치 시점이다. 응답 DTO를 조립하기 전에 handler가 돌아야 deadline으로 자동 생성된 투영이 그 요청의 응답에 포함된다. 그래서 각 쓰기 use-case의 끝에서 dispatch하고, commit 지점에 안전망을 한 겹 더 둔다.
def flush_and_dispatch_events(session): session.flush() events = drain_events(session) if events: messagebus.handle(events, session) session.flush() def commit_with_events(session): flush_and_dispatch_events(session) session.commit() # app/db/session.py — FastAPI 쓰기 의존성이 곧 트랜잭션 경계 def get_db_transactional(): with _session_manager.get_session() as session: try: yield session commit_with_events(session) # 안전망 drain 후 commit except Exception: session.rollback() clear_events(session) # session.info는 rollback과 연동되지 않음 raise
Router -> Service use-case 1. 도메인 상태 변경 2. collect_event(session, event) 3. flush_and_dispatch_events(session) -> Router: 응답 DTO 조립 (handler가 만든 리소스 포함) -> 의존성 teardown: commit_with_events 예외 발생 시 -> rollback + clear_events
이전 작업은 배경에서 본 인라인 블록들을
collect_event로 바꾸는 것이었고, 교체 후 기존 테스트 947건이 수정 없이 통과했다. 발행 위치를 기존 호출 위치와 같게 유지했기 때문이다. 이후 새 행동을 테스트로 먼저 정의하며 얹은 결과 최종 1,010건이 되었는데, 추가분은 버스 단위 테스트와 전파 행동 테스트다. 제목 변경이 투영에 반영되는지, 투영이 없을 때 deadline 변경이 재생성으로 이어지는지, origin 가드가 실제로 작동하는지 같은 것들이다.Handler 작성
수신 handler는 상대 도메인 모듈에 둔다. Todo 이벤트의 handler는
app/domain/schedule/event_handlers.py에 있고, 이벤트를 자기 도메인의 변경으로 번역하는 책임은 받는 쪽이 진다.@register(TodoDeadlineChanged) def move_projection_schedule(event, session): """deadline 변경 -> 투영 이동. 투영이 없으면 재생성""" todo = _get_owned_todo(session, event) # owner_id 일치 확인 포함 if todo is None or todo.deadline is None: return schedules = schedule_crud.get_schedules_by_source_todo_id( session, event.todo_id, event.owner_id ) if schedules: _move_projection(session, schedules[0], todo) else: _create_projection(session, todo) def _move_projection(session, schedule, todo): """start만 이동하고 사용자가 조정한 길이는 보존""" duration = schedule.end_time - schedule.start_time schedule.start_time = todo.deadline schedule.end_time = todo.deadline + duration session.flush()
handler는 내부 전파 경로라 사용자 권한 검증을 거치지 않는 대신 변경 대상의
owner_id가 event.owner_id와 일치할 때만 수정하고, CurrentUser를 요구하는 서비스 계층 대신 crud와 모델 직접 조작으로 처리한다. 그리고 필드마다 원본 도메인을 정해두고 소유한 필드만 반영한다. 제목·설명·태그는 Todo가 원본이라 항상 덮어쓰지만 일정 길이는 Schedule 소유다. 위 코드가 duration을 보존하는 이유다.왕복 차단
역방향 전파를 추가하는 순간 걱정되는 시나리오가 있다.
- 사용자가 Todo 마감을 15일로 수정한다.
- Todo가 이벤트를 발행하고, Schedule handler가 투영을 15일로 옮긴다.
- 일정이 바뀌었으니 Schedule이 이벤트를 발행하고, Todo handler가 마감을 갱신한다.
- 마감이 바뀌었으니 다시 이벤트가 발행되고, 2와 3이 무한 반복된다.
이 구현에서 핑퐁은 2단계에서 끊긴다. handler는 대상 행을 직접 수정할 뿐 이벤트를 발행하지 않는다.
ScheduleTimeChanged는 사용자가 ScheduleService.update_schedule을 호출할 때만 발행되는데, Todo발 변경으로 일정을 옮기는 handler는 그 서비스를 경유하지 않기 때문이다.TodoService --TodoDeadlineChanged--> schedule/event_handlers --직접 수정(이벤트 미발행)--> schedule 행 ScheduleService --ScheduleTimeChanged--> todo/event_handlers --직접 수정(이벤트 미발행)--> todo 행
이벤트의
origin 필드로 자기 도메인발 이벤트를 거르는 가드도 함께 있다.@register(ScheduleTimeChanged) def backpropagate_deadline(event, session): if event.origin == EventOrigin.TODO: return schedule = schedule_crud.get_schedule(session, event.schedule_id, event.owner_id) if schedule is None or schedule.source_todo_id is None: return todo = todo_crud.get_todo(session, schedule.source_todo_id, event.owner_id) if todo is None: return todo.deadline = schedule.start_time session.flush()
고백하자면 이 가드는 현재 코드 경로에서 도달 불가능하다. handler가 이벤트를 발행하지 않는 구조가 유지되는 한
origin=TODO인 ScheduleTimeChanged는 만들어지지 않는다. 지우지 않고 남긴 이유는 이 구조가 리팩토링 한 번에 깨질 수 있는 암묵적 약속이기 때문인데, 도달 불가능한 코드는 테스트도 통과시키지 못하는 죽은 방어가 되기 쉽다. 그래서 origin=TODO 이벤트를 수동으로 발행해 역전파가 일어나지 않는 것을 확인하고, origin=SCHEDULE로는 일어나는 것을 대조하는 테스트를 따로 두어 가드를 계약으로 묶었다. drain loop의 event_key 억제와 MAX_DISPATCH_DEPTH는 이것과 별개의 역할로, 전자는 동일 이벤트의 재처리를 막고 후자는 폭주를 제한한다.마무리
디스패처 코어 자체는 30줄 남짓에 불과하지만, 실제로 시간이 들어간 곳은 그 주변이었다. payload에 무엇을 싣고 무엇을 조회할지, rollback 경계에서 큐를 언제 비울지, 왕복을 구조로 막을지 가드로 막을지 같은 문제들이다. 이벤트 버스 도입을 고려하고 있다면 버스 자체보다 이런 주변 규칙에 시간을 배정하는 편이 낫다.
남은 과제도 있다. WebSocket 경로는 아직 이 커밋 경계 바깥이라 이벤트를 발행해서는 안 되는 영역으로 남아 있고, 커밋 후 실행이 보장되어야 하는 부수효과(알림 등)는 transactional outbox로의 확장이 필요하다. 기회가 되면 이어서 다룰 예정이다.
부록
app/shared/events.py
""" Domain Event 기반 타입 이벤트는 발생한 도메인 사실만 담는 불변 값 객체다. payload에는 식별자와 도메인 사실만 담는다 (API DTO/ORM 인스턴스 금지, datetime은 naive-UTC, 컬렉션은 tuple — frozen 유지). """ from dataclasses import dataclass from enum import Enum class EventOrigin(str, Enum): """이벤트를 촉발한 source 도메인 (왕복 전파 차단용)""" TODO = "todo" SCHEDULE = "schedule" @dataclass(frozen=True, kw_only=True) class DomainEvent: """ 모든 도메인 이벤트의 베이스 :param owner_id: 리소스 소유자 (OIDC sub). 핸들러는 변경 대상의 owner_id가 이와 일치할 때만 수정한다. :param origin: 이벤트를 촉발한 source 도메인 """ owner_id: str origin: EventOrigin def event_key(self) -> tuple: """ 중복 dispatch 억제 키 — 타입 + 주요 식별자 같은 drain 루프 안에서 동일 키의 이벤트는 한 번만 처리된다. 서브클래스는 주요 식별자를 포함하도록 재정의한다. """ return (type(self).__qualname__, self.owner_id)
app/shared/uow.py
""" 세션 스코프 이벤트 수집 helper 범용 이벤트 버스/outbox가 아니라 동기 트랜잭션 내 도메인 이벤트 투영 큐 도메인 모델이 SQLModel ORM 클래스의 alias라 엔티티에 `.events` 리스트를 직접 붙일 수 없어, 세션 스코프 수집(`session.info`)을 사용한다. """ from typing import TYPE_CHECKING, List from sqlmodel import Session if TYPE_CHECKING: from app.shared.events import DomainEvent _EVENTS_KEY = "domain_events" def collect_event(session: Session, event: "DomainEvent") -> None: """이벤트를 세션 큐에 적재 (dispatch는 flush_and_dispatch_events에서)""" session.info.setdefault(_EVENTS_KEY, []).append(event) def drain_events(session: Session) -> List["DomainEvent"]: """수집된 이벤트를 전부 꺼내 반환 (큐는 비워짐)""" events = session.info.get(_EVENTS_KEY, []) session.info[_EVENTS_KEY] = [] return events def clear_events(session: Session) -> None: """이벤트 큐 비우기 — rollback 경계에서 필수 호출""" session.info[_EVENTS_KEY] = [] def flush_and_dispatch_events(session: Session) -> None: """ flush → drain/handle 루프 → flush (commit 없음) 쓰기 use-case가 응답 DTO를 조립하기 전에 호출해야 핸들러가 만든 자동 생성 리소스가 응답에 즉시 보인다. 핸들러 예외는 전파된다 → 호출자(트랜잭션 경계)에서 전체 rollback. """ from app.shared import messagebus session.flush() events = drain_events(session) if events: messagebus.handle(events, session) session.flush() def commit_with_events(session: Session) -> None: """미처리 이벤트 dispatch를 한 번 더 보장(안전망)한 뒤 commit""" flush_and_dispatch_events(session) session.commit()
app/shared/messagebus.py
""" In-process Message Bus 큐 방식 drain 루프로 이벤트를 핸들러에 전달한다. 모든 핸들러는 같은 세션/트랜잭션에서 실행되며, 하나라도 실패하면 예외가 전파되어 전체 rollback된다. cascade 무한루프는 event_key 중복 억제와 MAX_DISPATCH_DEPTH 상한으로 차단한다. """ import logging from typing import Callable, Dict, List, Sequence, Type from sqlmodel import Session from app.shared.events import DomainEvent from app.shared.uow import drain_events logger = logging.getLogger(__name__) Handler = Callable[[DomainEvent, Session], None] #: 한 drain 루프에서 처리하는 **총 이벤트 수** 상한 (폭주 cascade 차단). #: 재귀 깊이가 아니다 — 배치성 대량 발행이 필요해지면 깊이/총량 분리 필요. MAX_DISPATCH_DEPTH = 10 #: 중앙 핸들러 registry — register()로만 등록 HANDLERS: Dict[Type[DomainEvent], List[Handler]] = {} #: @register가 기록하는 등록 명세. 모듈 import는 한 번만 데코레이터를 #: 실행하므로, reset_handlers() 이후 register_handlers()가 이 명세를 #: replay해 registry를 복구한다. _KNOWN_REGISTRATIONS: List[tuple] = [] class MaxDispatchDepthExceededError(RuntimeError): """drain 루프가 MAX_DISPATCH_DEPTH를 초과 — 폭주 cascade 의심""" def __init__(self, event: DomainEvent): super().__init__( f"Max event dispatch depth ({MAX_DISPATCH_DEPTH}) exceeded " f"while handling {type(event).__qualname__}" ) self.event = event def register(event_type: Type[DomainEvent]): """ 핸들러 등록 데코레이터 — @register(TodoDeadlineSet) 같은 (event_type, handler) 쌍은 한 번만 등록된다(idempotent) — 모듈 재-import에 안전. """ def decorator(fn: Handler) -> Handler: if (event_type, fn) not in _KNOWN_REGISTRATIONS: _KNOWN_REGISTRATIONS.append((event_type, fn)) _add_handler(event_type, fn) return fn return decorator def _add_handler(event_type: Type[DomainEvent], fn: Handler) -> None: handlers = HANDLERS.setdefault(event_type, []) if fn not in handlers: handlers.append(fn) def handle(events: Sequence[DomainEvent], session: Session) -> None: """ 큐 방식 drain 루프 - event_key 중복은 skip (왕복 전파 종료) - 핸들러 실행 후 flush → 후속 핸들러가 앞선 변경을 조회 가능 - 핸들러가 collect한 후속 이벤트를 재수집해 같은 루프에서 처리 - 예외는 삼키지 않고 전파 → 호출자에서 전체 rollback """ queue = list(events) seen: set = set() depth = 0 while queue: event = queue.pop(0) key = event.event_key() if key in seen: logger.debug("Duplicate event suppressed in dispatch chain: %r", event) continue seen.add(key) if depth >= MAX_DISPATCH_DEPTH: raise MaxDispatchDepthExceededError(event) depth += 1 for handler in HANDLERS.get(type(event), []): handler(event, session) session.flush() queue.extend(drain_events(session)) def register_handlers() -> None: """ 도메인 핸들러 일괄 등록 (idempotent, reset 후 복구 가능) 앱 startup(lifespan)과 테스트 fixture 양쪽에서 호출되므로 재호출에 안전해야 한다. import는 첫 호출에서만 데코레이터를 실행하므로, 기록된 등록 명세를 replay해 reset_handlers() 이후에도 복구한다. 순환 import 방지를 위해 lazy import. """ import app.domain.schedule.event_handlers # noqa: F401 import app.domain.todo.event_handlers # noqa: F401 for event_type, fn in _KNOWN_REGISTRATIONS: _add_handler(event_type, fn) def reset_handlers() -> None: """registry 초기화 — 테스트 전용. register_handlers()로 복구할 수 있다.""" HANDLERS.clear()
app/domain/todo/enums.py
""" Todo Enums """ from enum import Enum class TodoStatus(str, Enum): """Todo 상태""" UNSCHEDULED = "UNSCHEDULED" SCHEDULED = "SCHEDULED" DONE = "DONE" CANCELLED = "CANCELLED" class LinkedSchedulePolicy(str, Enum): """마감 취소·Todo 삭제 시 연결 Schedule의 수명주기 처리 정책. 서버 기본값은 없다.""" DELETE = "delete" # 투영 삭제 + visibility cleanup UNLINK = "unlink" # 연결만 해제 — 일반 일정으로 생존 (데이터 보존)
app/domain/schedule/enums.py (발췌)
class SourceTodoPolicy(str, Enum): """ 투영 일정이 마감 표시이기를 그만둘 때(삭제·반복 전환) source Todo 처리 정책. 서버 기본값은 없다. """ CLEAR_DEADLINE = "clear_deadline" # Todo의 deadline도 제거 KEEP_DEADLINE = "keep_deadline" # deadline 유지 (연동 플래그만 해제)
app/domain/todo/events.py
""" Todo Domain Events payload는 식별자만 담고, 수신 핸들러가 같은 트랜잭션에서 현재 상태를 조회해 투영한다. 예외적으로 dispatch 시점에 원본이 사라지는 사실만 스냅샷으로 운반한다. """ from dataclasses import dataclass from typing import Tuple from uuid import UUID from app.domain.todo.enums import LinkedSchedulePolicy from app.shared.events import DomainEvent @dataclass(frozen=True, kw_only=True) class TodoEvent(DomainEvent): """Todo 이벤트 공통 베이스 — event_key에 todo_id 포함""" todo_id: UUID def event_key(self) -> tuple: return (type(self).__qualname__, str(self.todo_id)) @dataclass(frozen=True, kw_only=True) class TodoDeadlineSet(TodoEvent): """deadline이 새로 설정됨 (생성 시 포함)""" @dataclass(frozen=True, kw_only=True) class TodoDeadlineChanged(TodoEvent): """deadline 값이 변경됨""" @dataclass(frozen=True, kw_only=True) class TodoDeadlineRemoved(TodoEvent): """ deadline이 제거됨 :param policy: 클라이언트가 명시한 연결 투영 처리 정책 """ policy: LinkedSchedulePolicy @dataclass(frozen=True, kw_only=True) class TodoContentChanged(TodoEvent): """제목·설명이 변경됨""" @dataclass(frozen=True, kw_only=True) class TodoTagsChanged(TodoEvent): """태그가 변경됨""" @dataclass(frozen=True, kw_only=True) class TodoDeleted(TodoEvent): """ Todo가 삭제됨 :param linked_schedule_ids: 삭제 직전 스냅샷한 연결 Schedule ID. `Schedule.source_todo_id`가 `ondelete=SET NULL`이라 삭제 flush 후에는 조회할 수 없어 이벤트가 사실로서 운반한다. :param policy: 클라이언트가 명시한 연결 투영 처리 정책 """ linked_schedule_ids: Tuple[UUID, ...] policy: LinkedSchedulePolicy
app/domain/schedule/events.py
""" Schedule Domain Events 역전파(Schedule→Todo)는 `source_todo_id`가 있는 투영 일정에 한정된다. """ from dataclasses import dataclass from typing import Optional from uuid import UUID from app.domain.schedule.enums import SourceTodoPolicy from app.shared.events import DomainEvent @dataclass(frozen=True, kw_only=True) class ScheduleEvent(DomainEvent): """Schedule 이벤트 공통 베이스 — event_key에 schedule_id 포함""" schedule_id: UUID def event_key(self) -> tuple: return (type(self).__qualname__, str(self.schedule_id)) @dataclass(frozen=True, kw_only=True) class ScheduleTimeChanged(ScheduleEvent): """일정 시간이 변경됨. origin이 TODO면 재전파하지 않는다.""" @dataclass(frozen=True, kw_only=True) class ScheduleDeleted(ScheduleEvent): """ 일정이 삭제됨 :param source_todo_id: 삭제 직전 스냅샷 (row가 사라져 사후 조회 불가) :param policy: source-linked 일정이면 클라이언트가 명시한 source Todo 처리 정책. unlinked 일정 삭제면 None """ source_todo_id: Optional[UUID] policy: Optional[SourceTodoPolicy] @dataclass(frozen=True, kw_only=True) class ScheduleUnlinkedFromTodo(ScheduleEvent): """ 일정이 마감 투영이기를 그만둠 (반복 전환으로 연결 해제) :param source_todo_id: 해제 직전 스냅샷 :param policy: 클라이언트가 명시한 source Todo 처리 정책 """ source_todo_id: UUID policy: SourceTodoPolicy @dataclass(frozen=True, kw_only=True) class ScheduleCreatedWithTodoOption(ScheduleEvent): """ create_todo_options와 함께 일정이 생성됨 핸들러가 low-level로 Todo를 생성하므로 TodoDeadlineSet이 재발화하지 않는다 (A→B→A 순환의 구조적 차단). :param tag_group_id: 생성될 Todo가 속할 TagGroup """ tag_group_id: UUID
app/domain/schedule/event_handlers.py
""" Todo 이벤트를 받아 투영 Schedule을 동기화하는 핸들러. 핸들러는 내부 projection이라 사용자 권한 검증 대신 owner 불변조건 (변경 대상 owner_id == event.owner_id)을 강제하고, CurrentUser 없이 low-level 연산만 사용한다. 자동 생성 리소스에 visibility는 복사하지 않는다. """ import logging from datetime import timedelta from uuid import UUID from sqlmodel import Session from app.crud import schedule as schedule_crud from app.crud import tag as tag_crud from app.crud import todo as todo_crud from app.crud import visibility as visibility_crud from app.domain.todo.enums import LinkedSchedulePolicy from app.domain.todo.events import ( TodoContentChanged, TodoDeadlineChanged, TodoDeadlineRemoved, TodoDeadlineSet, TodoDeleted, TodoTagsChanged, ) from app.domain.visibility.enums import ResourceType from app.models.schedule import Schedule from app.models.todo import Todo from app.shared.messagebus import register logger = logging.getLogger(__name__) #: 투영 생성 시 초기 길이. 이후 길이는 사용자가 조정할 수 있고 #: deadline 이동 시 보존된다. INITIAL_PROJECTION_DURATION = timedelta(hours=1) # ============================================================ # 내부 helper # ============================================================ def _get_owned_todo(session: Session, event) -> Todo | None: """owner 불변조건을 통과한 Todo 반환 (없으면 None)""" todo = todo_crud.get_todo(session, event.todo_id, event.owner_id) if todo is None: logger.warning( "Todo %s not found for %s (owner=%s) — skipping projection", event.todo_id, type(event).__qualname__, event.owner_id, ) return todo def _sync_projection_tags(session: Session, schedule_id: UUID, todo: Todo) -> None: """투영 태그를 Todo 태그로 전체 교체. 소유자의 태그만 연결한다.""" tag_crud.delete_all_schedule_tags(session, schedule_id) for tag in todo.tags: if tag.owner_id == todo.owner_id: tag_crud.add_schedule_tag(session, schedule_id, tag.id) def _create_projection(session: Session, todo: Todo) -> Schedule: """투영 Schedule 생성 — 제목·설명·태그는 Todo에서, 길이는 초기값""" schedule = Schedule( owner_id=todo.owner_id, title=todo.title, description=todo.description, start_time=todo.deadline, end_time=todo.deadline + INITIAL_PROJECTION_DURATION, source_todo_id=todo.id, # state는 모델 기본값 PLANNED, visibility는 복사하지 않음 (기본 PRIVATE) ) session.add(schedule) session.flush() _sync_projection_tags(session, schedule.id, todo) return schedule def _move_projection(session: Session, schedule: Schedule, todo: Todo) -> None: """start만 이동하고 사용자가 조정한 길이는 보존한다.""" duration = schedule.end_time - schedule.start_time schedule.start_time = todo.deadline schedule.end_time = todo.deadline + duration session.flush() def _delete_schedule_with_cleanup(session: Session, schedule: Schedule) -> None: """Schedule 삭제 + 접근권한 설정 cleanup (orphan visibility 방지)""" visibility_crud.delete_visibility_by_resource( session, ResourceType.SCHEDULE, schedule.id ) schedule_crud.delete_schedule(session, schedule) def _unlink_schedule(session: Session, schedule: Schedule) -> None: """연결만 해제 — 일반 일정으로 생존 (메모·태그·공유 보존)""" schedule.source_todo_id = None session.flush() # ============================================================ # 핸들러 # ============================================================ @register(TodoDeadlineSet) def create_projection_schedule(event: TodoDeadlineSet, session: Session) -> None: """deadline 설정 → 투영 생성. 이미 있으면 이동 (단일 투영 불변조건)""" todo = _get_owned_todo(session, event) if todo is None or todo.deadline is None: return schedules = schedule_crud.get_schedules_by_source_todo_id( session, event.todo_id, event.owner_id ) if schedules: _move_projection(session, schedules[0], todo) else: _create_projection(session, todo) @register(TodoDeadlineChanged) def move_projection_schedule(event: TodoDeadlineChanged, session: Session) -> None: """deadline 변경 → 투영 이동. 투영이 없으면 재생성 (조용한 no-op 금지)""" todo = _get_owned_todo(session, event) if todo is None or todo.deadline is None: return schedules = schedule_crud.get_schedules_by_source_todo_id( session, event.todo_id, event.owner_id ) if schedules: _move_projection(session, schedules[0], todo) else: _create_projection(session, todo) @register(TodoDeadlineRemoved) def remove_projection_schedule(event: TodoDeadlineRemoved, session: Session) -> None: """deadline 제거 → 클라이언트 정책에 따라 삭제 또는 연결 해제""" schedules = schedule_crud.get_schedules_by_source_todo_id( session, event.todo_id, event.owner_id ) for schedule in schedules: if event.policy == LinkedSchedulePolicy.UNLINK: _unlink_schedule(session, schedule) else: _delete_schedule_with_cleanup(session, schedule) @register(TodoContentChanged) def sync_projection_content(event: TodoContentChanged, session: Session) -> None: """제목·설명 변경 → 투영에 최신 내용 반영""" todo = _get_owned_todo(session, event) if todo is None: return schedules = schedule_crud.get_schedules_by_source_todo_id( session, event.todo_id, event.owner_id ) for schedule in schedules: schedule.title = todo.title schedule.description = todo.description if schedules: session.flush() @register(TodoTagsChanged) def sync_projection_tags(event: TodoTagsChanged, session: Session) -> None: """태그 변경 → 투영 태그 전체 교체""" todo = _get_owned_todo(session, event) if todo is None: return schedules = schedule_crud.get_schedules_by_source_todo_id( session, event.todo_id, event.owner_id ) for schedule in schedules: _sync_projection_tags(session, schedule.id, todo) @register(TodoDeleted) def remove_projections_of_deleted_todo(event: TodoDeleted, session: Session) -> None: """Todo 삭제 → 스냅샷된 연결 투영을 클라이언트 정책에 따라 처리""" schedules = schedule_crud.get_schedules_by_ids( session, list(event.linked_schedule_ids) ) for schedule in schedules: if schedule.owner_id != event.owner_id: # owner 불변조건 logger.warning( "Schedule %s owner mismatch on TodoDeleted — skipping", schedule.id ) continue if event.policy == LinkedSchedulePolicy.UNLINK: _unlink_schedule(session, schedule) else: _delete_schedule_with_cleanup(session, schedule)
app/domain/todo/event_handlers.py
""" Schedule 이벤트를 받아 source Todo를 동기화하는 핸들러. 역투영은 source_todo_id가 있는 일정에 한정하며, origin이 TODO인 변경은 재전파하지 않는다. status는 직접 쓰기(이벤트 미발행)로 갱신해 drain 재수집을 막는다. """ import logging from sqlmodel import Session from app.crud import schedule as schedule_crud from app.crud import tag as tag_crud from app.crud import todo as todo_crud from app.domain.schedule.enums import SourceTodoPolicy from app.domain.schedule.events import ( ScheduleCreatedWithTodoOption, ScheduleDeleted, ScheduleTimeChanged, ScheduleUnlinkedFromTodo, ) from app.domain.todo.enums import TodoStatus from app.models.todo import Todo from app.shared.events import EventOrigin from app.shared.messagebus import register logger = logging.getLogger(__name__) def _apply_detach_policy(session: Session, todo: Todo, policy: SourceTodoPolicy) -> None: """투영이 사라진 Todo에 클라이언트 정책 적용 + 연동 플래그 해제""" if policy == SourceTodoPolicy.CLEAR_DEADLINE: todo.deadline = None # 터미널 상태(DONE/CANCELLED)는 덮지 않음 — 플래그 해제는 SCHEDULED에만 if todo.status == TodoStatus.SCHEDULED: todo.status = TodoStatus.UNSCHEDULED session.flush() @register(ScheduleTimeChanged) def backpropagate_deadline(event: ScheduleTimeChanged, session: Session) -> None: """투영 일정 시간 이동 → Todo.deadline 역투영""" if event.origin == EventOrigin.TODO: return # Todo발 변경의 재전파 차단 schedule = schedule_crud.get_schedule(session, event.schedule_id, event.owner_id) if schedule is None or schedule.source_todo_id is None: return if schedule.recurrence_rule: # 반복 일정이 된 투영은 시점 마감의 투영이 아님 — 역투영 제외 logger.debug( "Skipping deadline back-propagation for recurring schedule %s", schedule.id, ) return todo = todo_crud.get_todo(session, schedule.source_todo_id, event.owner_id) if todo is None: return todo.deadline = schedule.start_time session.flush() @register(ScheduleDeleted) def apply_policy_on_projection_deleted(event: ScheduleDeleted, session: Session) -> None: """투영 삭제 → 클라이언트 정책에 따라 deadline 처리 + 연동 플래그 해제""" if event.origin == EventOrigin.TODO: return # Todo발 삭제/정리는 Todo 쪽 흐름이 이미 상태를 관리 if event.source_todo_id is None: return todo = todo_crud.get_todo(session, event.source_todo_id, event.owner_id) if todo is None: return # 내부 자동 삭제 경로는 policy가 없을 수 있음 → 데이터 보존 쪽 선택 _apply_detach_policy(session, todo, event.policy or SourceTodoPolicy.KEEP_DEADLINE) @register(ScheduleUnlinkedFromTodo) def apply_policy_on_projection_unlinked( event: ScheduleUnlinkedFromTodo, session: Session ) -> None: """반복 전환으로 연결 해제 → 클라이언트 정책에 따라 deadline 처리 + 플래그 해제""" todo = todo_crud.get_todo(session, event.source_todo_id, event.owner_id) if todo is None: return _apply_detach_policy(session, todo, event.policy) @register(ScheduleCreatedWithTodoOption) def create_todo_for_schedule( event: ScheduleCreatedWithTodoOption, session: Session ) -> None: """ create_todo_options → 연결 Todo 생성 + 역연결 low-level 생성이므로 TodoDeadlineSet이 재발화하지 않는다 (A→B→A 순환의 구조적 차단). 태그는 Schedule에서 복사한다. """ schedule = schedule_crud.get_schedule(session, event.schedule_id, event.owner_id) if schedule is None or schedule.source_todo_id is not None: return todo = Todo( owner_id=event.owner_id, title=schedule.title, description=schedule.description, deadline=schedule.start_time, tag_group_id=event.tag_group_id, parent_id=None, status=TodoStatus.SCHEDULED, # 투영(이 Schedule)이 있으므로 SCHEDULED ) session.add(todo) session.flush() # 태그 복사 (owner 불변조건: 이벤트 소유자의 태그만) for tag in schedule.tags: if tag.owner_id == event.owner_id: tag_crud.add_todo_tag(session, todo.id, tag.id) schedule.source_todo_id = todo.id session.flush()
app/db/session.py (발췌)
def get_db() -> Generator[Session, None, None]: """ 읽기 전용 세션 (commit 없음) - GET 요청 등 읽기 작업에 사용 - FastAPI가 자동으로 세션 close 처리 """ from app.shared.uow import clear_events with _session_manager.get_session() as session: try: yield session except Exception: session.rollback() raise finally: if session.new or session.dirty or session.deleted: logger.warning( "Read-only session rolled back with pending changes. " "new=%d dirty=%d deleted=%d", len(session.new), len(session.dirty), len(session.deleted), ) session.rollback() # session.info는 트랜잭션 aware가 아님 → rollback 경계에서 clear clear_events(session) def get_db_transactional() -> Generator[Session, None, None]: """ 트랜잭션 자동 관리 세션 - POST/PUT/DELETE 등 쓰기 작업에 사용 - 함수 종료 성공: 미처리 도메인 이벤트 dispatch(안전망) 후 commit 자동 - 예외 발생: rollback 자동 + 이벤트 큐 clear """ from app.shared.uow import clear_events, commit_with_events with _session_manager.get_session() as session: try: yield session commit_with_events(session) except Exception: session.rollback() # session.info는 트랜잭션 aware가 아님 → rollback 경계에서 clear clear_events(session) raise
배선 (발췌)
# app/main.py — lifespan startup # 1-1. Domain Event 핸들러 등록 (idempotent) from app.shared.messagebus import register_handlers register_handlers() logger.info("✅ Domain event handlers registered")
# tests/conftest.py @pytest.fixture(autouse=True, scope="session") def _register_domain_event_handlers(): """ 도메인 이벤트 핸들러 일괄 등록 서비스 단위 테스트는 app.main(lifespan)을 import하지 않으므로, 여기서 등록하지 않으면 투영 동기화가 조용히 누락됩니다. register_handlers는 idempotent라 e2e(lifespan 경유)와 중복 호출에 안전합니다. """ from app.shared.messagebus import register_handlers register_handlers()
