You cannot select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
78 lines
2.4 KiB
Python
78 lines
2.4 KiB
Python
# backend/app/worker/main.py — 잡 등록 진입 (phase-7 소유, phase-13/14 가 잡 추가)
|
|
# WORKER_ENABLED=false(기본)면 스케줄러 미가동 — routers/worker.py 의 수동 트리거로 동작.
|
|
from sqlmodel import Session
|
|
|
|
from ..db import engine
|
|
|
|
|
|
def _wrap(fn, **kw):
|
|
def _job():
|
|
with Session(engine) as s:
|
|
fn(s, **kw)
|
|
|
|
return _job
|
|
|
|
|
|
def register(scheduler) -> None:
|
|
"""APScheduler 등 스케줄러가 주어질 때만 잡을 등록(있을 때). 없으면 수동 트리거."""
|
|
from .jobs.digest import run_digest
|
|
from .jobs.proactive import run_proactive
|
|
from .jobs.weekly_review import run_weekly_review
|
|
|
|
scheduler.add_job(
|
|
_wrap(run_proactive), "interval", minutes=30, id="proactive", replace_existing=True
|
|
)
|
|
for t in ["09:00", "13:00", "18:30"]:
|
|
h, m = t.split(":")
|
|
scheduler.add_job(
|
|
_wrap(run_digest, slot=t),
|
|
"cron",
|
|
hour=int(h),
|
|
minute=int(m),
|
|
id=f"digest-{t}",
|
|
replace_existing=True,
|
|
)
|
|
scheduler.add_job(
|
|
_wrap(run_weekly_review),
|
|
"cron",
|
|
day_of_week="sun",
|
|
hour=20,
|
|
id="weekly_review",
|
|
replace_existing=True,
|
|
)
|
|
# phase-13 sync_jobs.register(scheduler) 도 같은 진입에서 호출 가능
|
|
try:
|
|
from .sync_jobs import register as register_sync
|
|
|
|
register_sync(scheduler)
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
def register_event_subscribers(bus, session_factory) -> None:
|
|
"""런타임 event_bus 구독: approval.executed → 승인된 메일 결재 발송(phase-16)."""
|
|
from sqlmodel import select
|
|
|
|
from ..connectors.mail.outbound import send_outbound
|
|
from ..models import Approval, OutboundMail
|
|
|
|
def _on_approval_executed(ev):
|
|
# phase-16: 승인된 메일 결재 → 실제 발송(real) 또는 mock(Sent)
|
|
approval_id = (ev.payload or {}).get("approval_id")
|
|
if not approval_id:
|
|
return
|
|
with session_factory() as s:
|
|
appr = s.get(Approval, approval_id)
|
|
if not appr or appr.source != "mail":
|
|
return
|
|
ob = s.exec(
|
|
select(OutboundMail).where(
|
|
OutboundMail.approval_id == approval_id,
|
|
OutboundMail.status == "pending",
|
|
)
|
|
).first()
|
|
if ob:
|
|
send_outbound(s, ob)
|
|
|
|
bus.subscribe("approval.executed", _on_approval_executed)
|