비동기 작업 시스템
문서 파싱, 인덱스 구축, 요약과 질문 생성, 그래프 추출, Wiki 생성, 데이터 소스 동기화, 일괄 작업은 작업 시스템이 스케줄링합니다. 표준 배포는 asynq와 Redis를 사용하고, Lite 모드는 Redis 없는 실행기를 사용합니다. 두 모드는 작업 처리 로직을 공유하지만 실행 방식과 운영 기능은 다릅니다.
전체 아키텍처: 두 가지 실행 모드
WeKnora에는 배포 형태에 따라 선택하는 두 가지 작업 실행 모드가 있습니다.
- asynq 모드(표준 배포): 작업은
asynq.Client를 통해 JSON payload로 직렬화되어 Redis 큐에 기록되고, 여러 독립asynq.Server(worker pool)가 소비합니다.internal/router/task.go의RunAsynqServer()는 공통asynq.ServeMux를 구성하여 6개 pool에서 실행합니다. - Lite 모드(단일 머신 / macOS App, Redis 없음):
internal/router/sync_task.go의SyncTaskExecutor는 동일한interfaces.TaskEnqueuer인터페이스를 구현합니다.Enqueue는 작업을 goroutine에 바로 전달해 실행하며ProcessIn(지연)과MaxRetry옵션을 지원합니다. 재시도는 선형 백오프(attempt * 5s, 상한 30s)를 사용합니다.
// internal/router/sync_task.go
// SyncTaskExecutor는 Redis 없이 작업을 (goroutine에서) 동기적으로 실행합니다.
// Lite 모드에서 *asynq.Client를 그대로 대체하는 용도로 사용합니다.두 모드에 등록되는 handler 집합은 완전히 동일합니다(RunAsynqServer와 RegisterSyncHandlers 비교). 따라서 배포 형태에 따라 작업 의미가 달라지지 않습니다.
시스템에서 Redis의 역할
| 역할 | 설명 | 소스 코드 위치 |
|---|---|---|
| asynq broker | 모든 작업 큐(pending list, scheduled/retry ZSET, archived ZSET)를 Redis에 저장합니다. dequeue는 원자적이며(BRPOPLPUSH), 하나의 작업을 하나의 worker만 실행하도록 보장합니다 | internal/router/task.go getAsynqRedisClientOpt() |
| 작업 점검 데이터 소스 | asynq.Inspector + 직접적인 LPos/ZRank/ZRevRank 페이지 조회 | internal/router/task_inspector.go |
| Wiki ingest 뮤텍스 | wiki:active:<kbID>, finalize 잠금, slug 잠금 모두 SetNX + TTL 사용 | internal/application/service/wiki_ingest.go, wiki_ingest_batch.go |
| 멀티모달 하위 작업 카운터 | 이미지 하위 작업 완료 카운트(DECR). 마지막 attempt가 finalize 트리거 | image_multimodal 관련 서비스 |
| 요청 제한 | 슬라이딩 윈도우 요청 제한 ZSET(관측성 문서 참조) | internal/ratelimit/limiter.go |
Redis 연결 매개변수는 환경 변수 REDIS_ADDR / REDIS_USERNAME / REDIS_PASSWORD / REDIS_DB / TLS 설정에서 가져옵니다. 읽기·쓰기 타임아웃은 WEKNORA_REDIS_OP_TIMEOUT_MS로 제어하며 기본값은 500ms입니다(쓰기 타임아웃은 선두 차단을 흡수하기 위해 그 2배).
// internal/router/task.go
const defaultRedisOpTimeoutMs = 500
opt := &asynq.RedisClientOpt{
Addr: os.Getenv("REDIS_ADDR"),
ReadTimeout: time.Duration(timeoutMs) * time.Millisecond,
WriteTimeout: time.Duration(timeoutMs*2) * time.Millisecond,
...
}작업 유형 목록
작업 유형 상수는 internal/types/task.go에 정의되어 있습니다.
| 작업 유형 | 상수 | 용도 | 큐 |
|---|---|---|---|
document:process | TypeDocumentProcess | 문서 파싱 진입점(DocReader / 분할 / 벡터화) | default |
manual:process | TypeManualProcess | 수동 지식 업데이트(cleanup + 재인덱싱) | default |
temporary_document:process | TypeTemporaryDocumentProcess | 세션 임시 문서(채팅 첨부 파일) 파싱 | chat_attachment |
knowledge:post_process | TypeKnowledgePostProcess | 지식 후처리 통합 스케줄링(보강 하위 작업 fan-out) | postprocess |
knowledge:auto_tag | TypeKnowledgeAutoTag | 기존 태그를 문서에 자동 연결 | summary |
memory:extract | TypeMemoryExtract | 개인 메모리 백그라운드 추출 | memory |
summary:generation | TypeSummaryGeneration | 요약 생성 | summary |
datatable:summary | TypeDataTableSummary | 표 요약 | summary |
image:multimodal | TypeImageMultimodal | 이미지 OCR + VLM Caption | multimodal |
chunk:extract | TypeChunkExtract | 그래프 엔티티/관계 추출(chunk별) | graph |
question:generation | TypeQuestionGeneration | 질문 생성(chunk 배치별 fan-out) | question |
datasource:sync | TypeDataSourceSync | 데이터 소스 동기화 | sync |
faq:import | TypeFAQImport | FAQ 가져오기(dry run 포함) | low(maintenance) |
kb:clone | TypeKBClone | 지식 베이스 복사 | low |
kb:delete | TypeKBDelete | 지식 베이스 삭제 | low |
index:delete | TypeIndexDelete | 인덱스 삭제 | low |
knowledge:list_delete | TypeKnowledgeListDelete | 지식 일괄 삭제 | low |
knowledge:list_reparse | TypeKnowledgeListReparse | 일괄 재파싱 | low |
knowledge:move | TypeKnowledgeMove | 지식 이동 | low |
wiki:ingest | TypeWikiIngest | Wiki 페이지 생성/동기화 | wiki |
wiki:finalize | TypeWikiFinalize | Wiki KB 수준 마무리(디바운스: 인덱스 재구축/깨진 링크 정리/교차 링크) | wiki |
모든 payload 구조체(예: DocumentProcessPayload, ImageMultimodalPayload)는 types.TracingContext를 임베딩하여 Langfuse/W3C traceparent를 프로세스 간 전달합니다(관측성 문서 참조). 또한 tenant_id / knowledge_id / knowledge_base_id 등의 라우팅 필드를 공통으로 포함하며, 데드레터 보관과 취소 대상 매칭에 사용합니다.
자동 태그와 요약은 summary 큐를 공유하는 선택적 보강 작업입니다. 문서 지식 베이스에서 auto_tag_config를 켰을 때만 큐에 추가하며, 실패해도 완료된 파싱에는 영향을 주지 않습니다. 메모리 추출은 독립 memory 큐를 사용하고 enrichment pool이 소비합니다. 가중치 1로 shared pool의 탄력적 용량 사용에도 참여하며, worker pool 총수는 여전히 6개입니다.
메모리 작업은 개인 주체별로 중복 제거하고 지연 집계합니다. memory_subjects의 extract_cursor, pending_sessions, 예약 시간을 사용해 처리를 이어 가므로 질문할 때마다 모델 추출을 즉시 시작하지 않습니다. 공간의 메모리가 꺼져 있거나 write_mode가 auto가 아니면 백그라운드 증류를 실행하지 않습니다. Lite의 동기 실행기도 자동 태그와 메모리 작업을 등록하며 같은 비즈니스 기능 설정을 따릅니다.
Worker Pool 토폴로지와 제어 전략
internal/types/task.go의 queueDefinitions는 큐 토폴로지의 단일 진실 공급원(single source of truth)입니다. worker server 구성(QueueWeightsForPool)과 운영 대시보드 표시(QueueStats)가 이 레지스트리를 공유하여 가중치 불일치를 방지합니다.
6개의 독립 worker pool
각 pool은 독립적인 asynq.Server이며 동시 실행 용량이 엄격히 격리됩니다(가중치에 따른 선호가 아님). 기본 동시 실행 수와 설정 키(system_settings 키 / 환경 변수, types.ResolveWorkerPoolConcurrency 참조)는 다음과 같습니다.
| Pool | 기본 동시 실행 수 | 소비 큐(가중치) | 설정 키 / 환경 변수 |
|---|---|---|---|
core | 8 | default(1), chat_attachment(3) | asynq.core_concurrency / WEKNORA_ASYNQ_CORE_CONCURRENCY |
postprocess | 2 | postprocess(1) | asynq.postprocess_concurrency / WEKNORA_ASYNQ_POSTPROCESS_CONCURRENCY |
enrichment | 12 | summary(2), multimodal(1), graph(1), question(1), memory(1) | asynq.enrichment_concurrency / WEKNORA_ASYNQ_ENRICHMENT_CONCURRENCY |
maintenance | 4 | sync(2), low(1) | asynq.maintenance_concurrency / WEKNORA_ASYNQ_MAINTENANCE_CONCURRENCY |
shared(탄력 계층) | 6 | core + enrichment 중 SharedWeight > 0인 큐 | asynq.shared_concurrency / WEKNORA_ASYNQ_SHARED_CONCURRENCY |
wiki | 8 | wiki(1) | asynq.wiki_concurrency / WEKNORA_WIKI_ASYNQ_CONCURRENCY |
설계 핵심 사항(소스 주석에서도 확인 가능):
- 보장 용량 + 탄력적 차용: core/postprocess/enrichment/maintenance가 최소 용량을 보장합니다.
sharedpool은 core와 enrichment 큐를 동시에 구독하므로 어느 단계든 유휴 용량을 빌려 쓸 수 있습니다(NewSharedAsynqServer: Redis dequeue는 원자적이므로 여러 server가 같은 큐를 구독해도 각 작업은 한 번만 실행). post-process와 maintenance는 shared pool에 참여하지 않습니다. 전자는 지연 시간을 독립적으로 보장해야 하고, 후자는 실행 시간이 길어 대화형 작업의 순간 부하용 용량을 점유할 수 있기 때문입니다. 범위는QueueWeightsForSharedPool이 정의합니다. - Wiki의 엄격한 격리:
wikipool은wiki큐만 가져와 파싱 파이프라인과 Wiki 생성이 서로 기아 상태를 일으키지 않도록 합니다(NewWikiAsynqServer주석). - 채팅 첨부 파일 우선: core pool에서
chat_attachment의 가중치는 3으로default의 1보다 높습니다. 따라서 대규모 KB 가져오기가 대화형 채팅 업로드를 대기시키지 않습니다. - 롤링 업그레이드 호환성:
QueueMaintenance상수의 물리 Redis 큐 이름은 이전 버전의"low"를 유지합니다. 롤링 배포 중에도 이전 버전이 큐에 넣은 작업을 소비할 수 있습니다.
Worker Pool 아키텍처 다이어그램
미들웨어 제어
RunAsynqServer(internal/router/task.go)는 같은 mux에 다음 세 미들웨어를 순서대로 설치합니다.
asynqdl.MiddlewareWithCallback(데드레터) — handler가 반환한 원래 오류를 확인할 수 있도록 반드시 가장 먼저 설치해야 합니다(이후 미들웨어가 오류를 변환할 수 있음). 실패 재시도와 데드레터 처리를 참고하세요.backgroundTaskMiddleware— 각 작업 context에types.WithBackgroundTask표시를 붙입니다. 모델별 채팅 동시성 제어기(chat concurrency governor)는 이를 통해 ingestion/enrichment의 LLM 호출을 제한하되 대화형 사용자 채팅에는 영향을 주지 않습니다.langfuse.AsynqMiddleware— Langfuse가 꺼져 있으면 그대로 통과시킵니다. 켜져 있으면 상위 HTTP trace를 이어 받거나 독립 trace를 새로 시작하고 handler 실행을 SPAN으로 감쌉니다.
재시도 백오프 전략
기본적으로 asynq의 지수 백오프(약 10s, 40s, 90s, 2.5m…)를 사용하지만 Wiki ingest 잠금 충돌은 별도로 처리합니다(asynqRetryDelayFunc).
// internal/router/task.go
func asynqRetryDelayFunc(n int, e error, t *asynq.Task) time.Duration {
if errors.Is(e, service.ErrWikiIngestConcurrent) {
return wikiIngestRetryDelay // 고정 15s
}
return asynq.DefaultRetryDelayFunc(n, e, t)
}이유: 고아 잠금의 TTL은 ≤ 60s이므로 고정 15s 재시도는 거의 확실하게 성공합니다. 반면 지수 백오프는 프로세스 비정상 종료 후 재시작했을 때 KB를 7–10분 동안 멈춰 있게 할 수 있습니다.
작업 생명주기 상태 머신
asynq 측 런타임 상태(internal/router/task_inspector.go의 runtimeTaskState가 types.RuntimeTaskState로 매핑)는 pending, active, scheduled, retry, archived, completed입니다. 비즈니스 측 지식 행의 parse_status(internal/types/knowledge.go)는 pending → processing → finalizing → completed와 failed / deleting / cancelled입니다.
이에 대응하는 지식 행 상태(작업이 구동):
작업 점검, 취소, 운영 대시보드(TaskInspector)
internal/router/task_inspector.go는 interfaces.TaskInspector를 구현합니다. asynq 모드에서는 asynq.Inspector + 네이티브 Redis client를 기반으로 하고, Lite 모드에서는 noopTaskInspector를 사용합니다(goroutine은 시작 전에 제거할 수 없으므로 체크포인트 방식의 중단이 유일한 중지 신호).
지식 / 지식 베이스별 취소
CancelTasksForKnowledge(ctx, knowledgeID): 등록된 모든 큐(queuesScanned는types.QueueDefinitions()에서 가져옴)의 pending/scheduled/retry/active 네 상태를 스캔하여 payload의knowledge_id가 일치하면 처리합니다. 취소 가능한 작업 유형 허용 목록taskTypesForKnowledgeCancel은document:process,manual:process,image:multimodal,knowledge:post_process,question:generation,summary:generation,chunk:extract입니다(FAQ 가져오기와 지식 베이스 수준 작업은 제외).- 취소 절차는 세 단계입니다(
cancelMatchingTasks). ① 먼저 대기 상태를 모두 삭제합니다. ② active 작업의 스냅샷을 만든 뒤Inspector.CancelProcessing으로 신호를 보내고, 1s 안정화 구간에서 25ms 간격으로 폴링하여context.Canceled때문에 retry로 이동한 레코드를 삭제합니다(deleteCancelledTransitions). ③ 대기 상태를 다시 스캔해 취소 도중 새로 큐에 들어온 후속 작업까지 처리합니다. CancelTasksForKnowledgeBase: KB 삭제 후 고아 작업을 정리합니다.kb:delete와index:delete는 명시적으로 제외합니다(스냅샷을 가지고 실제 저장소 정리를 담당하므로 삭제하면 리소스가 누수됨). clone/move의 의미상 KB 필드(source_id/target_id/source_kb_id/target_kb_id)도 매칭에 참여합니다.- 모두 best-effort로 동작합니다. Redis가 불안정하면 Warn 로그를 남기고 오류를 삼키며, 취소 API는 여전히 성공을 반환합니다.
HasQueuedTasksForKnowledge: 읽기 전용 탐지입니다. housekeeping 정리는 이를 사용해 "적체되어 있지만 고아는 아닌" 행을 구분하여 failed로 잘못 표시하지 않도록 합니다.
운영 대시보드(SystemAdmin Runtime Dashboard)
QueueStats(): 큐마다GetQueueInfo를 호출해types.QueueStat(size/pending/active/scheduled/retry/archived/completed, 당일 processed/failed, paused,latency_ms(가장 오래된 pending 작업의 경과 시간), 메모리 사용량)을 출력하고 정적 pool/weight 메타데이터를 덧붙입니다. 한 번도 생성되지 않은 큐는 0값 행을 반환합니다(isAsynqQueueNotFound는 asynq v0.26에서 노출되는 내부NOT_FOUND오류 문자열도 지원.task_inspector_errors.go참조).ListRuntimeTasks(): Redis 키asynq:{<queue>}:<state>를 기반으로 직접 페이지를 나눕니다. pending/active는 LIST(최신 항목 우선), scheduled/retry는NextProcessAt오름차순 ZSET, archived/completed는 점수 내림차순입니다. 커서는 base64로 인코딩한 앵커 윈도우(최대 32개 앵커,runtimeTaskCursorMaxAnchors)이며, 앵커가 사라져도(작업 완료/재시도/삭제) 페이지 조회를 이어 갈 수 있습니다. payload는 허용 목록의 라우팅 메타데이터(tenant/kb/knowledge/task/sync 등의 ID)만 투영하고 문서 내용이나 키는 절대 노출하지 않습니다.- 작업 동작은
runtimeTaskActions의 상태 검사로 제한합니다.cancel(pending/active/scheduled/retry이면서 취소 가능한 유형),run_now(scheduled/retry/archived, asynq가 재시도 횟수 유지),delete(archived만)입니다.PurgeArchivedRuntimeTasks로 단일 큐의 archived 집합을 한 번에 비울 수도 있습니다. WorkerServerStats(): asynq server 하트비트(동시 실행 수, 활성 worker 수, 상태, 큐 가중치)를 읽고 복제본 전체를 집계하여 "설정된 단일 인스턴스 용량"과 "실제 클러스터 용량"을 구분합니다.
대응 HTTP API(internal/router/router.go, SystemAdmin + 플랫폼 API Key capability로 접근 제어):
| 메서드 | 경로 | 설명 |
|---|---|---|
| GET | /api/v1/system/admin/runtime/queues | 큐 깊이 스냅샷 + worker 하트비트 |
| GET | /api/v1/system/admin/runtime/queues/:queue/tasks | 상태별 커서 페이지 조회 작업 목록 |
| POST | /api/v1/system/admin/runtime/queues/:queue/tasks/:task_id/actions/:action | cancel / run_now / delete(플랫폼 감사 기록) |
| DELETE | /api/v1/system/admin/runtime/queues/:queue/archived | archived 비우기(플랫폼 감사 system.queue_archived_purged 기록) |
실패 재시도와 데드레터 처리
asynq 데드레터 미들웨어(internal/middleware/asynqdl/asynqdl.go)
- 마지막 시도가 실패했을 때만(
isFinalAttempt:retried >= max_retry)task_dead_letters에 한 행을 기록하여 일시적 불안정이 발생할 때마다 행이 생기지 않도록 합니다. buildDeadLetter는 허용적인payloadProbe로 임의의 payload에서tenant_id/knowledge_base_id/kb_id/knowledge_id/source_kb_id를 추출합니다.inferScope는 "영향 범위"에 따라 scope를 추론합니다(knowledge_base>knowledge>tenant>unknown). payload는 원형 그대로 유지하고(향후 재실행에 사용 가능),last_error는 8KB로 제한합니다.- 삽입은 best-effort입니다. DB 실패는 로그만 남기며 원래 작업 오류는 항상 그대로 asynq로 전파합니다(archived로 이동).
OnDeadLetter콜백(internal/router/task.go의newDeadLetterKnowledgeFailer):document:process/knowledge:post_process/manual:process가 재시도를 소진하면 단일 UPDATE로 지식 행의parse_status=failed+error_message를 함께 기록합니다(부분 업데이트 방지). 이어서SpanTracker.FinalizeAttempt로 해당 attempt의 루트 span을 닫아 타임라인에 더 이상 "진행 중"으로 표시되지 않게 합니다.knowledge:list_delete에는 전용 분기markKnowledgeListDeleteFailed가 있습니다.image:multimodal은 부모 지식을 실패로 표시하지 않습니다(finalize-on-last-attempt가 이미 진행을 보장). 콜백은context.Background()로 실행하며 panic을 포착하고 원래 작업 오류를 절대 바꾸지 않습니다.
영구 저장 작업 큐와 서비스 수준 데드레터(internal/application/repository/task_queue.go)
task_pending_ops 테이블은 Redis list 큐의 영구 저장 대안입니다(재시작 시 손실 없음, TTL 축출 없음). 큐 식별자는 (task_type, scope, scope_id) 세 값의 조합이며, 현재 주 소비자는 Wiki ingest입니다.
Enqueue/EnqueueIfKnowledgeBaseActive: 후자는 트랜잭션 안에서 PostgresSHARE행 잠금으로 KB가 여전히 유효한지 검증하여 KB 소프트 삭제 후 새로운 영구 작업이 기록되지 않도록 합니다.ClaimBatch:dedup_key(=문서)별 전체 그룹을 원자적으로 선점합니다. 핵심 불변식은 같은 문서의 여러 op(예: ingest 뒤 retract)를 절대 두 병렬 배치로 나누지 않는 것입니다. 유효한 claim(claimed_at >= staleBefore)이 있는 key는 그룹 전체를 건너뛰며, 늦게 도착한 형제 op는 소유자의 완료 또는 claim 만료를 기다립니다. Postgres에서는 각 key의 anchor 행에FOR UPDATE SKIP LOCKED를 적용해 병렬 선점자가 서로 겹치지 않는 key 집합을 얻도록 보장합니다. SQLite(Lite/테스트)는 단일 쓰기 실행자 엔진에 의존합니다.IncrFailCount(UPDATE ... RETURNING으로 단일 왕복 원자적 증가)는 서비스 측 상한(wiki의wikiMaxFailRetries)과 함께 사용합니다. 상한을 넘으면 해당 op를task_pending_ops에서task_dead_letters로 옮깁니다(internal/application/service/wiki_ingest.go가deadLetterRepo.Insert직접 호출).ReleaseByIDs/DeleteByIDs/DeleteByScope/DeleteByDedupKey/PendingCount는 선점 해제, 소비 확인, KB 생명주기 정리, 적체 관측을 제공합니다.
데드레터 저장소 taskDeadLetterRepository는 ListByScope / ListByTaskType(id 내림차순 커서 페이지 조회, limit 1–200)와 DeleteByID를 제공합니다. 운영자는 로그를 뒤지지 않고 SQL로 작업 유형 / scope / 테넌트별 실패를 직접 조회할 수 있습니다.
최종 안전장치: housekeeping 정리
internal/application/service/knowledge_housekeeping.go: cron이 5분마다(0 */5 * * * *) stale 임계값보다 오래 pending/processing/finalizing에 머문 지식 행을 스캔해 failed로 표시합니다. 이는 asynq 재시도, 데드레터 콜백, multimodal finalize 외의 마지막 방어선입니다(handler 실행 중 worker가 kill되어 defer가 실행되지 않은 경우 등). 정리 시 span 하트비트, updated_at, TaskInspector.HasQueuedTasksForKnowledge를 함께 확인하여 "적체되어 있지만 고아는 아닌" 행을 잘못 종료하지 않도록 합니다. WEKNORA_HOUSEKEEPING_ENABLED=false로 끌 수 있습니다.
이벤트 버스(internal/event)
이벤트 버스는 프로세스 내부의 세션/에이전트 스트리밍 이벤트 배포(SSE 전송, IM 콜백 등)에 사용하며, asynq(프로세스 간 영구 작업)와 상호 보완적입니다.
구조와 전달 보장
// internal/event/event.go
type Event struct {
ID string // 이벤트 ID(자동 생성 UUID, 스트리밍 업데이트 추적용)
Type EventType
SessionID string
Data interface{}
Metadata map[string]interface{}
RequestID string
}EventBus.On(type, handler)로 등록하고(같은 유형에 여러 handler 가능)Off/Clear로 제거합니다.HasHandlers/GetHandlerCount로 조회합니다.- 동기 모드(
NewEventBus, 기본값):Emit이 handler를 순서대로 실행하고 하나라도 오류가 나면 즉시 오류를 반환합니다(at-most-once, 오류 시 후속 handler 중단). - 비동기 모드(
NewAsyncEventBus):Emit은 handler마다 goroutine을 시작하는 fire-and-forget 방식입니다. 오류는 버리고 panic은 recover하여 로그에 남깁니다. EmitAndWait: 두 모드 모두에서 전체 handler를 병렬 실행하고 완료를 기다리며 오류와 panic을 수집합니다.- 전달 보장은 프로세스 내부에 한정되며 영구 저장되지 않습니다. 등록된 handler가 없으면 이벤트를 조용히 버립니다(nil 반환). 프로세스가 충돌하면 진행 중인 이벤트가 손실됩니다. 영구 저장이 필요하면 asynq 또는
task_pending_ops를 사용해야 합니다. global.go는 전역 싱글턴(event.On/event.Emit)을 제공합니다. 실제 세션 수준 스트리밍 처리는 독립 bus 인스턴스를 사용합니다(구독자 참조).middleware.go는 handler 미들웨어를 제공합니다.WithLogging(트리거/실패 로그),WithTiming(소요 시간을 metadata에 기록),WithRecovery(panic을PanicError로 변환)를Chain/ApplyMiddleware로 조합합니다.adapter.go의EventBusAdapter는 순환 의존성을 피하기 위해*EventBus를types.EventBusInterface에 맞춥니다.
이벤트 유형 목록(internal/event/event.go)
| 그룹 | 이벤트 유형 |
|---|---|
| 쿼리 처리 | query.received, query.validated, query.preprocess, query.rewrite, query.rewritten |
| 검색 | retrieval.start, retrieval.vector, retrieval.keyword, retrieval.entity, retrieval.complete |
| 리랭킹 | rerank.start, rerank.complete |
| 병합 | merge.start, merge.complete |
| 채팅 생성 | chat.start, chat.complete, chat.stream |
| 에이전트 생명주기 | agent.query, agent.plan, agent.step, agent.tool, agent.complete |
| 에이전트 스트리밍(실시간 피드백) | thought, tool_call, tool_result, reflection, references, final_answer |
| MCP 도구 수동 승인 | tool_approval_required, tool_approval_resolved |
| MCP OAuth 세션 내 인증 | mcp_oauth_required, mcp_oauth_resolved |
| 오류 / 세션 / 제어 | error, session_title, stop |
각 이벤트의 데이터 구조는 internal/event/event_data.go에 정의되어 있습니다(예: AgentToolCallData는 tool_call_id/tool_name/arguments/hint를 포함하고, AgentFinalAnswerData는 content/done/is_fallback 등을 포함).
주요 구독자
| 구독자 | 소스 코드 | 구독 내용 |
|---|---|---|
| SSE 에이전트 스트리밍 handler | internal/handler/session/agent_stream_handler.go | thought, tool_call, tool_result, references, final_answer, reflection, error, session_title, agent.complete, tool approval 및 MCP OAuth의 네 유형 |
| 지식 질의응답 handler | internal/handler/session/qa.go, helpers.go | thought, final_answer, stop |
| IM 통합(WeCom 등) | internal/im/service.go | final_answer, error, references, agent.complete, thought, tool_call, tool_result, mcp_oauth_required 등을 각 IM 플랫폼 메시지로 변환 |
internal/runtime 패키지
이 패키지는 작으며 worker 로직이 아닌 런타임 인프라를 담당합니다.
container.go:init()가 전역*dig.Container(uber dig)를 생성하고GetContainer()가 각 패키지의 의존성 등록/해석에 사용됩니다. 모든 asynq server, handler, repository를 이를 통해 구성합니다(실제 대규모 구성은internal/container/container.go에서 수행).server.go:MarkServerStarted()/ServerStartedAt()/ServerUptime()— 프로세스 시작 시각을 기록하여 운영 대시보드에서 uptime을 표시하는 데 사용합니다.startup.go:SilenceGinRouteSpam()은 약 150줄의 Gin 라우트 등록 로그를 억제하고 한 줄로 요약합니다(LogGinRouteCount).LogStartupEnv()는 선별한 환경 변수 배너를 출력하고(민감 값은set (N chars)만 표시), 흔한 위험 설정에 명시적으로 경고합니다(예:SYSTEM_AES_KEY길이가 32가 아니면 암호화가 실제로 비활성화됨,REDIS_TLS_INSECURE_SKIP_VERIFY=true).
작업 모니터링 방법
- 운영 대시보드 / Runtime API(6.2절): 큐 깊이, 가장 오래된 pending 지연(
latency_ms), 당일 processed/failed, worker 하트비트를 확인합니다. 상태별 작업 탐색,last_error,retried/max_retry확인,run_now/cancel/delete실행을 지원합니다. - 데드레터 테이블 SQL:
SELECT * FROM task_dead_letters WHERE scope='knowledge_base' AND scope_id='<kbID>' ORDER BY id DESC;또는task_type별 실패율을 집계합니다.task_pending_ops의PendingCount/enqueued_at로 한 번도 해소되지 않은 적체를 찾을 수 있습니다. - 로그: worker 측은 모두
internal/logger를 사용합니다. 주요 접두사는[TaskInspector](취소/점검),asynq dead-letter,[SyncTask](Lite 모드),[Housekeeping]입니다. 시작 시 각 pool은asynq <pool> server starting with concurrency=...를 출력합니다. - Langfuse trace: 활성화하면 각 asynq 작업은
asynq.<task_type>SPAN(queue, retry, payload 크기 메타데이터 포함)이 되며, 이를 트리거한 HTTP 요청과 같은 trace에 속합니다(관측성 문서 참조). - 플랫폼 감사: archived 작업의
run_now/delete/purge 동작은audit_logs에 기록되어(system.queue_task_*동작) 책임을 추적할 수 있습니다.
구현 참고
| 모듈 | 소스 코드 경로 |
|---|---|
| 작업 등록과 worker pool 구성 | internal/router/task.go |
| Lite 모드 동기 실행기(Redis 없음) | internal/router/sync_task.go |
| 작업 점검 / 취소 / 운영 대시보드 | internal/router/task_inspector.go, internal/router/task_inspector_errors.go |
| 큐 토폴로지와 작업 유형 정의 | internal/types/task.go |
| 데드레터 미들웨어 | internal/middleware/asynqdl/asynqdl.go |
| 영구 저장 작업 큐 / 데드레터 저장소 | internal/application/repository/task_queue.go |
| 데드레터 / 대기 작업 모델 | internal/types/task_dead_letter.go, internal/types/task_pending_op.go |
| 이벤트 버스 | internal/event/(event.go, event_data.go, global.go, middleware.go, adapter.go) |
| 런타임 보조(DI 컨테이너, 시작 배너, uptime) | internal/runtime/(container.go, server.go, startup.go) |
| 멈춘 작업의 최종 안전장치 정리 | internal/application/service/knowledge_housekeeping.go |