큐 삽입 쪽(social_crud.decide, social_service.create_draft/publish_reused_text)이
JobType enum(SOCIAL_DRAFT=9, SOCIAL_POST=10)을 안 쓰고 숫자를 하드코딩(8, 9)해서,
워커 디스패처(worker/handlers.py)가 그 숫자로 엉뚱한 핸들러를 불렀다 — "게시" 잡(9)은
run_draft로, "초안" 잡(8)은 run_rollback으로. 둘 다 대상 상태 조건이 안 맞아 에러 없이
{"skipped": true}로 끝나 DONE 처리됐다 — 승인해도 실제로는 한 번도 게시되지 않는데
로그만 보면 정상으로 보이는 조용한 실패였다. SOCIAL_POSTING_ENABLED가 계속 꺼져 있어
지금까지 드러나지 않았다(2026-09-14부터 있던 버그).
- 세 호출부를 JobType.SOCIAL_DRAFT.value/SOCIAL_POST.value로 교체
- 기존 테스트의 job_type 기대값(8→9, 9→10)도 실제 enum에 맞게 수정
- 신규: 디스패처 매핑 정적 대조 + 실제 큐 삽입값으로 하는 엔드투엔드 회귀 테스트
(되돌려서 새 테스트가 실패하는 것까지 확인함)
전체 회귀 70 passed
68 lines
2.5 KiB
Python
68 lines
2.5 KiB
Python
"""승인 CAS와 큐 적재를 같은 트랜잭션으로 묶어 승인 후 잡 유실을 막는다."""
|
|
|
|
import json
|
|
from sqlalchemy import text
|
|
from common.database.db_session_manager import DB_SESSION_MNG
|
|
from common.database.model.models import place_social_posts as Post
|
|
from common.enums import JobType
|
|
|
|
|
|
async def transaction(fn):
|
|
return await DB_SESSION_MNG.execute_lambda_write(Post.DBType(), fn)
|
|
|
|
|
|
async def enqueue(s, post_id, job_type):
|
|
await s.execute(
|
|
text("""INSERT INTO jobs(job_type, payload, dedupe_key, max_attempts)
|
|
VALUES (:type, CAST(:payload AS jsonb), :key, 1)
|
|
ON CONFLICT (dedupe_key) WHERE status IN (1,2) AND dedupe_key IS NOT NULL DO NOTHING"""),
|
|
{
|
|
"type": job_type,
|
|
"payload": json.dumps({"post_id": str(post_id)}),
|
|
"key": f"social:{job_type}:{post_id}",
|
|
},
|
|
)
|
|
await s.execute(text("SELECT pg_notify('web4ai_job', '')"))
|
|
|
|
|
|
async def decide(s, post_id, sha, approve, via):
|
|
row = (
|
|
await s.execute(
|
|
text("""UPDATE place_social_posts
|
|
SET status=:status, decided_at=now(), decided_via=:via, updated_at=now()
|
|
WHERE post_id=:id AND deleted=false AND status='PENDING_APPROVAL'
|
|
AND approval_token_sha=:sha AND approval_expires_at>now()
|
|
RETURNING post_id, account_id"""),
|
|
{
|
|
"status": "APPROVED" if approve else "DECLINED",
|
|
"via": via,
|
|
"id": post_id,
|
|
"sha": sha,
|
|
},
|
|
)
|
|
).first()
|
|
if row and approve and row.account_id:
|
|
await enqueue(s, post_id, JobType.SOCIAL_POST.value)
|
|
return bool(row)
|
|
|
|
|
|
async def sweep():
|
|
async def run(s):
|
|
await s.execute(
|
|
text("""UPDATE place_social_posts SET status='EXPIRED', updated_at=now()
|
|
WHERE deleted=false AND status='PENDING_APPROVAL' AND approval_expires_at<=now()""")
|
|
)
|
|
# POSTING은 외부가 받았을 수 있다. 시간을 근거로 APPROVED로 돌리지 않는다.
|
|
await s.execute(
|
|
text("""UPDATE place_social_posts SET status='UNKNOWN',
|
|
last_error='POST_RESULT_UNKNOWN', updated_at=now()
|
|
WHERE deleted=false AND status='POSTING' AND updated_at<now()-interval '10 minutes'""")
|
|
)
|
|
await s.execute(
|
|
text("""UPDATE place_social_posts SET status='FAILED',
|
|
last_error='DRAFT_INTERRUPTED', updated_at=now()
|
|
WHERE deleted=false AND status='DRAFTING' AND updated_at<now()-interval '10 minutes'""")
|
|
)
|
|
|
|
await transaction(run)
|