o2o-site-AEO/solution/backend/crud/social_crud.py
김성경 df1556f2b1 [fix] solution/backend: 쓰레드 게시 잡이 job_type 번호가 어긋나 워커에 조용히 씹혔다
큐 삽입 쪽(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
2026-09-22 11:36:21 +09:00

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)