최상단을 프로젝트 단위로 평평하게 둔다 — o2o-negosium 과 같은 규약이고, 이 레포만
다르게 갈 이유가 없다. negodata/{backend,front} 가 프로젝트 안에서 f/b 를 가르는 선례,
lps-admin/ 이 백엔드 없이 프론트만 가진 최상단 폴더의 선례다.
backend/ frontend/{admin,site,shared} → solution/{backend,front,site,shared} + admin/
## 왜
내부 라우트(/local-content, /places/:id/seo)의 이름과 화면 코드가 사장님 번들에
그대로 실려 나가고 있었다. UserRole.DEVELOPER 주석의 "고객사에 존재를 노출하지 않는다"를
번들이 깨고 있었다 — 라우트 가드는 화면을 가리지 번들은 못 가린다.
번들을 갈라 확인했다: 사장님 dist 에서 local-content · /places · SeoAudit 이 전부 0건이다.
그 과정에서 두 곳이 더 새고 있었다.
- AppShell 의 NAV 배열이 내부 메뉴를 하드코딩하고 있었다. 앱을 가른 뒤에도 dist 에
local-content 가 남아서 찾았다. 메뉴는 이제 앱이 prop 으로 들고 온다.
- EditorHeader·BuilderPage·LoginPage 가 /places 로 링크하고 있었다. 그 화면이 admin 으로
나갔으니 사장님 앱에서는 404 다. 링크를 걷어내고 LoginPage 기본 도착지는 '/' 로 바꿨다
(앱마다 홈이 다르고 각 라우터의 '/' 가 이미 그걸 안다).
## admin 에 백엔드를 두지 않았다
내부 화면이 부르는 훅이 전부 router/v1/{place,fact,local,validator} 에 이미 있다.
자체 백엔드를 두면 place·fact·link 를 같은 DB 에 대고 두 번 구현하게 된다.
대가는 solution/backend 가 죽으면 admin 도 멈추는 것 — 내부 도구라 감수한다.
## admin 의 `@` 는 solution/front/src 를 가리킨다
내부 화면이 쓰는 API 클라이언트·UI·수집 배선이 solution 에 한 벌만 있고 그 파일들끼리도
`@/...` 로 서로를 부른다. admin 에서 `@` 를 자기 src 로 잡으면 그 참조가 전부 깨진다
(실측 TS2307 14건). 복제하는 길도 있지만 RecollectPanel 주석이 금지한다 —
"수집 경로를 두 벌 만들면 확정 게이트"가 갈라진다.
admin 자기 파일만 `@admin` 이고, 의존 방향은 admin → solution 한 쪽뿐이다.
admin 이 여는 빌더는 다른 오리진이라 절대 URL + 새 탭이다(admin/src/lib/solutionUrl.ts).
react-router Link 로 두면 admin 안에서 라우트를 찾다 404 다.
## 그 밖
- npm 워크스페이스 루트를 레포 루트로 올렸다(admin 이 solution 밖이라).
- docker-compose 를 255→174줄로 줄이고 admin(:3002) 서비스를 넣었다. ADMIN_BIND 기본값은
127.0.0.1 — 0.0.0.0 으로 열면 앱을 가른 의미가 없다.
- 발행 호스트를 프론트 .env 에 따로 적지 않는다. compose 가 루트의 SITE_PUBLIC_HOST 를
VITE_PUBLISH_HOST 로 흘려보낸다 — 두 곳에 적으면 canonical 과 화면 주소가 조용히 갈라진다.
- nginx/site.conf 를 git 에서 빼고 .example 만 남겼다(.env·*.toml 과 같은 규약).
compose 가 bind mount 하므로 클론 직후 복사해야 한다 — 없으면 Docker 가 그 자리에
디렉토리를 만들어 nginx 가 설정 없이 뜬다.
- config.test.toml.example 을 추가했다. 없으면 클론한 사람이 pytest 를 아예 못 돌린다
(conftest import 단계에서 죽는다). 외부 API 키는 전부 빈값이다 —
APP_ENV=test 가 .env 를 안 읽는 이유를 여기서 우회하면 안 된다.
- 경로가 한 칸 깊어져 test_schema_ddl(parents[2]→[3]) 과 test_site_theme 을 고쳤다.
검증: front·admin·site 전부 lint 0 / build 0. 백엔드 514 passed.
남은 4건(test_build_publish 3 · test_snapshot 1)은 이 변경 전부터 실패하던 것으로,
손대지 않은 메인 체크아웃에서 같은 4건이 같게 실패하는 것을 확인했다.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_019uYhHQdssRubirPirrdJJC
275 lines
12 KiB
Python
275 lines
12 KiB
Python
"""작업 큐 CRUD — PostgreSQL 을 '제대로' 큐로 쓴다. (LPS `crud/job_crud.py` 이식)
|
|
|
|
- 할당은 **단일 문장 원자 claim**: FOR UPDATE SKIP LOCKED 서브쿼리 + 같은 UPDATE + RETURNING.
|
|
→ 워커 컨테이너가 몇 개든 같은 잡 이중 할당이 원천 불가. fetch 와 claim 을 분리하지 않는다.
|
|
- 모든 전이는 **조건부 CAS**(WHERE 에 status/worker_id 가드) + RETURNING.
|
|
- 복구는 timeout 추측이 아니라 **lease 만료 소유권**(reaper 가 회수).
|
|
- 재시도/백오프/dead-letter 를 큐에 내장.
|
|
|
|
전이가 조회/변경으로 나뉘지 않으므로(RETURNING) execute_lambda_write 로 실행한다.
|
|
큐 전이만 raw SQL 이다 — 다른 crud 는 전부 SQLAlchemy 표현식을 쓴다.
|
|
"""
|
|
|
|
import json
|
|
|
|
from sqlalchemy import text
|
|
|
|
from common.database.db_session_manager import DB_SESSION_MNG
|
|
from common.enums import DBType, DBWRType, JobStatus
|
|
|
|
|
|
# 잡 적재 시 워커를 즉시 깨우는 LISTEN/NOTIFY 채널(폴링 제거).
|
|
JOB_NOTIFY_CHANNEL = "web4ai_job"
|
|
|
|
|
|
def compute_backoff(attempts: int, base: float = 5.0, cap: float = 600.0) -> float:
|
|
"""지수 백오프(초). attempts 회 시도 후 다음 재시도까지 대기 = base * 2^(attempts-1), cap 상한."""
|
|
return min(cap, base * (2 ** max(0, attempts - 1)))
|
|
|
|
|
|
class JobQueue:
|
|
DB = DBType.MAIN.value
|
|
|
|
async def _tx(self, fn):
|
|
"""쓰기 트랜잭션 — 값 반환이 필요한 큐 전이 전용 진입점."""
|
|
return await DB_SESSION_MNG.execute_lambda_write(self.DB, fn)
|
|
|
|
# ---- 적재 ----
|
|
async def enqueue(
|
|
self,
|
|
job_type: int,
|
|
payload: dict,
|
|
priority: int = 100,
|
|
dedupe_key: str | None = None,
|
|
max_attempts: int = 3,
|
|
) -> str | None:
|
|
"""잡 적재. dedupe_key 가 활성(PENDING/RUNNING) 중복이면 삽입 없이 None 반환."""
|
|
sql = text("""
|
|
INSERT INTO job.jobs (job_type, priority, payload, dedupe_key, max_attempts)
|
|
VALUES (:t, :p, CAST(:payload AS jsonb), :dk, :ma)
|
|
ON CONFLICT (dedupe_key) WHERE status IN (1, 2) AND dedupe_key IS NOT NULL
|
|
DO NOTHING
|
|
RETURNING job_id
|
|
""")
|
|
|
|
async def run(s):
|
|
row = (await s.execute(sql, {
|
|
"t": job_type, "p": priority, "payload": json.dumps(payload, ensure_ascii=False),
|
|
"dk": dedupe_key, "ma": max_attempts,
|
|
})).first()
|
|
if row:
|
|
# 커밋 시 전달됨 → LISTEN 중인 유휴 워커를 즉시 깨운다(중복 스킵 시엔 알림 안 함).
|
|
await s.execute(text("SELECT pg_notify(:ch, '')"), {"ch": JOB_NOTIFY_CHANNEL})
|
|
return str(row[0]) if row else None
|
|
|
|
return await self._tx(run)
|
|
|
|
# ---- 원자적 claim ----
|
|
async def claim(self, worker_id: str, lease_sec: int = 120) -> dict | None:
|
|
"""대기 잡 1건을 원자적으로 점유. 없으면 None.
|
|
FOR UPDATE SKIP LOCKED 로 잠근 행을 같은 UPDATE 에서 RUNNING 으로 전이 → 이중 할당 불가."""
|
|
sql = text("""
|
|
UPDATE job.jobs SET
|
|
status = 2,
|
|
worker_id = :wid,
|
|
lease_until = now() + make_interval(secs => :lease),
|
|
run_started_at = now(),
|
|
attempts = attempts + 1,
|
|
updated_at = now()
|
|
WHERE job_id = (
|
|
SELECT job_id FROM job.jobs
|
|
WHERE status = 1 AND run_after <= now()
|
|
ORDER BY priority ASC, created_at ASC
|
|
FOR UPDATE SKIP LOCKED
|
|
LIMIT 1
|
|
)
|
|
RETURNING job_id, job_type, payload, attempts, max_attempts
|
|
""")
|
|
|
|
async def run(s):
|
|
row = (await s.execute(sql, {"wid": worker_id, "lease": lease_sec})).mappings().first()
|
|
if not row:
|
|
return None
|
|
d = dict(row)
|
|
d["job_id"] = str(d["job_id"])
|
|
if isinstance(d.get("payload"), str):
|
|
d["payload"] = json.loads(d["payload"])
|
|
return d
|
|
|
|
return await self._tx(run)
|
|
|
|
# ---- 완료/실패 (소유권 가드) ----
|
|
async def complete(self, job_id: str, worker_id: str, result: dict | None = None) -> bool:
|
|
sql = text("""
|
|
UPDATE job.jobs SET status = 3, result = CAST(:result AS jsonb),
|
|
lease_until = NULL, worker_id = NULL, updated_at = now()
|
|
WHERE job_id = CAST(:id AS uuid) AND status = 2 AND worker_id = :wid
|
|
RETURNING job_id
|
|
""")
|
|
|
|
async def run(s):
|
|
row = (await s.execute(sql, {
|
|
"id": job_id, "wid": worker_id,
|
|
"result": json.dumps(result, ensure_ascii=False) if result is not None else None,
|
|
})).first()
|
|
return row is not None
|
|
|
|
return await self._tx(run)
|
|
|
|
async def fail(self, job_id: str, worker_id: str, error: str, backoff_sec: float = 5.0) -> int | None:
|
|
"""실패 처리. 시도 남으면 PENDING(run_after=백오프)으로 재큐, 소진되면 DEAD(dead-letter).
|
|
전이 후 status(JobStatus 값)를 반환. 소유 불일치면 None."""
|
|
sql = text("""
|
|
UPDATE job.jobs SET
|
|
status = CASE WHEN attempts >= max_attempts THEN 4 ELSE 1 END,
|
|
run_after = CASE WHEN attempts >= max_attempts THEN run_after
|
|
ELSE now() + make_interval(secs => :backoff) END,
|
|
last_error = :err,
|
|
lease_until = NULL,
|
|
worker_id = NULL,
|
|
updated_at = now()
|
|
WHERE job_id = CAST(:id AS uuid) AND status = 2 AND worker_id = :wid
|
|
RETURNING status
|
|
""")
|
|
|
|
async def run(s):
|
|
row = (await s.execute(sql, {
|
|
"id": job_id, "wid": worker_id, "err": error[:2000], "backoff": backoff_sec,
|
|
})).first()
|
|
return int(row[0]) if row else None
|
|
|
|
return await self._tx(run)
|
|
|
|
# ---- lease 갱신(heartbeat) / 회수(reaper) ----
|
|
async def renew_lease(self, job_id: str, worker_id: str, lease_sec: int = 120) -> bool:
|
|
sql = text("""
|
|
UPDATE job.jobs SET lease_until = now() + make_interval(secs => :lease), updated_at = now()
|
|
WHERE job_id = CAST(:id AS uuid) AND worker_id = :wid AND status = 2
|
|
RETURNING job_id
|
|
""")
|
|
|
|
async def run(s):
|
|
row = (await s.execute(sql, {"id": job_id, "wid": worker_id, "lease": lease_sec})).first()
|
|
return row is not None
|
|
|
|
return await self._tx(run)
|
|
|
|
async def reap(self) -> list[str]:
|
|
"""만료된 lease(워커 사망 등)의 RUNNING 잡을 회수. 시도 남으면 즉시 재큐, 소진되면 DEAD.
|
|
회수된 job_id 목록 반환."""
|
|
sql = text("""
|
|
UPDATE job.jobs SET
|
|
status = CASE WHEN attempts >= max_attempts THEN 4 ELSE 1 END,
|
|
run_after = now(),
|
|
last_error = COALESCE(last_error, '') || ' [lease-expired reclaim]',
|
|
lease_until = NULL,
|
|
worker_id = NULL,
|
|
updated_at = now()
|
|
WHERE status = 2 AND lease_until IS NOT NULL AND lease_until < now()
|
|
RETURNING job_id
|
|
""")
|
|
|
|
async def run(s):
|
|
rows = (await s.execute(sql)).all()
|
|
return [str(r[0]) for r in rows]
|
|
|
|
return await self._tx(run)
|
|
|
|
# ---- 단건 조회 (상태 폴링) ----
|
|
async def get(self, job_id: str) -> dict | None:
|
|
"""잡 단건 조회(읽기). 없으면 None. status 는 정수(JobStatus 값)."""
|
|
sql = text("""
|
|
SELECT job_id, job_type, status, priority, attempts, max_attempts,
|
|
payload, result, last_error, run_after, run_started_at, created_at, updated_at
|
|
FROM job.jobs WHERE job_id = CAST(:id AS uuid)
|
|
""")
|
|
|
|
async def run(s):
|
|
row = (await s.execute(sql, {"id": job_id})).mappings().first()
|
|
if not row:
|
|
return None
|
|
d = dict(row)
|
|
d["job_id"] = str(d["job_id"])
|
|
for key in ("payload", "result"):
|
|
if isinstance(d.get(key), str):
|
|
d[key] = json.loads(d[key])
|
|
return d
|
|
|
|
return await DB_SESSION_MNG.execute_lambda(self.DB, DBWRType.DB_READ.value, run)
|
|
|
|
async def find_active(self, dedupe_key: str) -> dict | None:
|
|
"""dedupe_key 로 활성(PENDING/RUNNING) 잡을 찾는다.
|
|
enqueue 가 중복으로 None 을 돌려줬을 때, 이미 돌고 있는 잡의 id 를 알려주기 위함."""
|
|
sql = text("""
|
|
SELECT job_id, job_type, status FROM job.jobs
|
|
WHERE dedupe_key = :dk AND status IN (1, 2)
|
|
LIMIT 1
|
|
""")
|
|
|
|
async def run(s):
|
|
row = (await s.execute(sql, {"dk": dedupe_key})).mappings().first()
|
|
if not row:
|
|
return None
|
|
d = dict(row)
|
|
d["job_id"] = str(d["job_id"])
|
|
return d
|
|
|
|
return await DB_SESSION_MNG.execute_lambda(self.DB, DBWRType.DB_READ.value, run)
|
|
|
|
# ---- 관측(관리 API/알림용) ----
|
|
async def counts(self) -> dict[str, int]:
|
|
"""상태별 잡 개수."""
|
|
async def run(s):
|
|
rows = (await s.execute(text("SELECT status, count(*) FROM job.jobs GROUP BY status"))).all()
|
|
by_val = {int(st): int(c) for st, c in rows}
|
|
return {js.name: by_val.get(js.value, 0) for js in JobStatus}
|
|
|
|
return await DB_SESSION_MNG.execute_lambda(self.DB, DBWRType.DB_READ.value, run)
|
|
|
|
async def ops(self) -> dict:
|
|
"""운영 스냅샷(모니터링·알림용): 상태별 카운트 + 큐 지연(가장 오래된 PENDING 나이) +
|
|
최근 1시간 DEAD + stuck(좀비 신호).
|
|
|
|
stuck 은 두 축 — lease 만료(워커 사망인데 reaper 미회수) OR 실행 10분 초과(핸들러 행 —
|
|
heartbeat 가 lease 를 계속 갱신해 lease 축엔 안 잡히므로 run_started_at 으로 따로 본다)."""
|
|
sql = text("""
|
|
SELECT
|
|
count(*) FILTER (WHERE status = 1) AS pending,
|
|
count(*) FILTER (WHERE status = 2) AS running,
|
|
count(*) FILTER (WHERE status = 3) AS done,
|
|
count(*) FILTER (WHERE status = 4) AS dead,
|
|
count(*) FILTER (WHERE status = 4 AND updated_at > now() - interval '1 hour') AS dead_1h,
|
|
count(*) FILTER (WHERE status = 2 AND (
|
|
(lease_until IS NOT NULL AND lease_until < now())
|
|
OR run_started_at < now() - interval '10 minutes'
|
|
)) AS stuck_running,
|
|
COALESCE(EXTRACT(EPOCH FROM (now() - min(created_at) FILTER (WHERE status = 1)))::int, 0) AS oldest_pending_sec,
|
|
count(*) FILTER (WHERE last_error LIKE 'JobDeadlineExceeded%'
|
|
AND updated_at > now() - interval '1 hour') AS deadline_1h
|
|
FROM job.jobs
|
|
""")
|
|
|
|
async def run(s):
|
|
return {k: int(v) for k, v in dict((await s.execute(sql)).mappings().first()).items()}
|
|
|
|
return await DB_SESSION_MNG.execute_lambda(self.DB, DBWRType.DB_READ.value, run)
|
|
|
|
async def requeue(self, job_id: str) -> str | None:
|
|
"""DEAD 잡 재큐(관리자 액션): attempts 리셋 + PENDING 전이 + 워커 깨움.
|
|
DEAD 가 아니거나 없으면 None. 같은 dedupe_key 의 활성 잡이 있으면 부분 유니크 위반."""
|
|
sql = text("""
|
|
UPDATE job.jobs SET status = 1, attempts = 0, run_after = now(),
|
|
lease_until = NULL, worker_id = NULL, run_started_at = NULL,
|
|
last_error = NULL, updated_at = now()
|
|
WHERE job_id = CAST(:jid AS uuid) AND status = 4
|
|
RETURNING job_id
|
|
""")
|
|
|
|
async def run(s):
|
|
row = (await s.execute(sql, {"jid": job_id})).first()
|
|
if row:
|
|
await s.execute(text("SELECT pg_notify(:ch, '')"), {"ch": JOB_NOTIFY_CHANNEL})
|
|
return str(row[0]) if row else None
|
|
|
|
return await self._tx(run)
|