문서 수집 파이프라인(Document Ingestion Pipeline)
문서는 파일 업로드, URL 가져오기 또는 수동 생성을 통해 시스템에 들어옵니다. 서버는 콘텐츠를 저장하고 처리 작업을 제출한 뒤 파싱, 청크 분할, 인덱싱을 차례로 완료하며, 지식 베이스 설정에 따라 요약, 질문, 그래프 및 Wiki를 생성합니다. 상태와 진행률은 처리 단계에 따라 갱신되며, 실패한 작업은 재시도 및 점검 메커니즘으로 복구됩니다.
전체 아키텍처
WeKnora의 수집 경로는 Asynq(Redis)기반 분산 비동기 파이프라인입니다. HTTP Handler는 DB 저장과 큐 등록만 담당하고, 시간이 오래 걸리는 모든 작업(파싱, 벡터화, LLM 보강)은 독립 Worker 풀이 큐에서 가져와 처리합니다.
진입 계층: 세 가지 생성 방식
라우트는 internal/router/routes_knowledge.go에 등록됩니다.
kb.POST("/file", g.OwnedKBOrAdmin(), g.KBAccessWrite("id"), handler.CreateKnowledgeFromFile)
kb.POST("/url", g.OwnedKBOrAdmin(), g.KBAccessWrite("id"), handler.CreateKnowledgeFromURL)
kb.POST("/manual", g.OwnedKBOrAdmin(), g.KBAccessWrite("id"), handler.CreateManualKnowledge)관련 관리 엔드포인트는 POST /knowledge/:id/reparse(다시 파싱), POST /knowledge/:id/cancel-parse(파싱 취소), POST /knowledge/batch-reparse, POST /knowledge/batch-delete, POST /knowledge/move(지식 베이스 간 이동)입니다.
파일 업로드(CreateKnowledgeFromFile)
- 폼 매개변수:
file,fileName,metadata,enable_multimodel,tag_ids,process_config(업로드할 때마다 KB 수준 처리 설정을 재정의할 수 있습니다. 처리 설정: KB 기본값 + 개별 업로드 재정의 참고). - 흐름: 확장자 검증 → MD5 중복 제거 →
FileService.SaveFile저장 →Knowledge레코드 생성 → 큐 등록.
통합 확장자 게이트
internal/application/service/knowledge_util.go의 supportedImportFileExtensions는 모든 가져오기 경로의 단일 진실 공급원입니다. 직접 업로드, 파일 URL 다운로드, worker 다운로드 완료 후 재검사 모두 같은 표를 확인합니다.
pdf txt docx doc epub html htm mhtml md markdown
png jpg jpeg gif csv xlsx xls pptx ppt json
mp3 wav m4a flac ogg이전에는 URL 가져오기가 더 짧은 별도 허용 목록을 관리하여 "직접 업로드한 xlsx는 가능하지만 URL로 가져온 xlsx는 거부됨"과 같은 불일치가 있었습니다(#2447). 이제 isSupportedImportExtension() / validateImportFileType()으로 통합 판정하며, 동영상 유형에는 "현재 동영상 파일 업로드는 지원되지 않습니다"라는 명확한 안내를 제공합니다.
표 형식 확장자(csv / xlsx / xls, dataTableFileExtensions)는 문서 처리 작업 뒤에 표 요약 작업(enqueueDataTableSummaryIfNeeded)을 추가로 연결합니다.
이미지 및 오디오 파일의 추가 사전 검증(객체 스토리지 설정 완전성, VLM / ASR 모델 설정 여부)은 process_config 검증과 함께 resolveFileImportProcessConfig()로 통합되며, 업로드와 URL 가져오기가 이를 공용으로 사용합니다.
URL 가져오기(CreateKnowledgeFromURL)
- JSON Body:
{url, file_name?, file_type?, enable_multimodel?, title?, tag_ids?, channel?, process_config?}. isFileURL()은 위의 통합 확장자 집합을 기준으로 이것이 "파일 다운로드"인지 "웹페이지 수집"인지 판정합니다.- Handler와 Service 양쪽에서 SSRF를 방어합니다(
internal/handler/knowledge.go와knowledge_create.go가 모두 호출).
if err := secutils.ValidateURLForSSRF(req.URL); err != nil {
c.Error(errors.NewBadRequestError(secutils.FormatSSRFError("URL", req.URL, err)))
return
}Worker 측은 실제 수집 전에 다시 검증하며(knowledge_process.go의 convert()), 삼중 방어선으로 TOCTOU를 방지합니다.
수동 생성(CreateManualKnowledge)
- JSON Body는
types.ManualKnowledgePayload{Title, Content, Status, TagIDs, Channel, ProcessConfig}이며 초안(Draft)상태를 지원합니다. 게시할 때는triggerManualProcessing()을 통해 파일과 동일한 청크 분할/인덱싱 파이프라인으로 진입합니다(DocReader 단계는 건너뜀).
중복 제거 메커니즘
knowledge_create.go는 업로드 파일의 MD5를 계산하고 네 항목의 조합으로 DB를 조회합니다.
hash, err := calculateFileHash(file) // MD5
exists, existingKnowledge, err := s.repo.CheckKnowledgeExists(ctx, tenantID, kbID,
&types.KnowledgeCheckParams{
Type: "file",
FileName: fileName,
FileType: getFileType(fileName),
FileSize: file.Size,
FileHash: hash,
})
if exists {
return existingKnowledge, types.NewDuplicateFileError(existingKnowledge)
}일치하면 중복 수집하지 않고 기존 Knowledge를 DuplicateFileError와 함께 반환합니다(프런트엔드는 이를 바탕으로 "파일이 이미 존재합니다"라고 안내). FileType도 판정에 참여합니다. 즉, 중복은 동일한 파일 유형 안에서만 성립하므로 콘텐츠가 완전히 같은 notes.md와 notes.txt는 별개의 지식 두 건으로 공존합니다(CheckKnowledgeExists는 해시 및 "파일명 + 크기" 두 분기 모두에 LOWER(file_type) 조건을 추가함).
초기 상태
새 Knowledge 레코드의 주요 초기 필드(knowledge_create.go)는 다음과 같습니다.
knowledge := &types.Knowledge{
ID: uuid.New().String(),
Type: "file", // 또는 "url" / "manual"
ParseStatus: "pending", // 초기 파싱 상태
EnableStatus: "disabled", // 인덱싱 완료 전에는 검색 불가
FileHash: hash,
...
}CSV/Excel 데이터 표 지식은 생성 후 TypeDataTableSummary(datatable:summary)작업도 추가로 큐에 등록하여 표 질의응답에 사용할 table_summary / table_column 유형 Chunk를 생성합니다.
파일 스토리지 계층(FileService와 스토리지 백엔드)
인터페이스 정의
internal/types/interfaces/file.go:
type FileService interface {
CheckConnectivity(ctx context.Context) error
SaveFile(ctx context.Context, file *multipart.FileHeader, tenantID uint64, knowledgeID string) (string, error)
SaveBytes(ctx context.Context, data []byte, tenantID uint64, fileName string, temp bool) (string, error)
GetFile(ctx context.Context, filePath string) (io.ReadCloser, error)
GetFileURL(ctx context.Context, filePath string) (string, error)
DeleteFile(ctx context.Context, filePath string) error
CopyFile(ctx context.Context, srcPath string, tenantID uint64, knowledgeID string) (string, error)
}지원하는 스토리지 백엔드
팩터리 함수 NewFileServiceFromStorageConfig()(internal/application/service/file/factory.go)는 types.StorageEngineConfig.DefaultProvider에 따라 백엔드를 선택합니다. 실제 지원 백엔드는 다음과 같습니다.
| Provider | 경로 접두사 | 구현 파일 | 설명 | 주요 설정 |
|---|---|---|---|---|
local | local:// | file/local.go | 단일 머신 로컬 디스크 | LocalEngineConfig.PathPrefix, 기본 디렉터리는 LOCAL_STORAGE_BASE_DIR, 외부 링크 서명은 APP_EXTERNAL_URL 사용 |
minio | minio:// | file/minio.go | MinIO / S3 호환 | MinIOEngineConfig(mode: docker이면 환경 변수 MINIO_ENDPOINT / MINIO_ACCESS_KEY_ID / MINIO_SECRET_ACCESS_KEY / MINIO_BUCKET_NAME을 읽고, mode: remote이면 설정 필드를 읽음) |
cos | cos:// | file/cos.go | Tencent Cloud COS | SecretID/SecretKey/Region/BucketName/AppID, 별도 임시 버킷 TempBucketName/TempRegion 지원 |
oss | oss:// | file/oss.go | Alibaba Cloud OSS | Endpoint/Region/AccessKey/SecretKey/BucketName, 임시 버킷 지원 |
s3 | s3:// | file/s3.go | AWS S3 / 호환 프로토콜 | Endpoint/Region/AccessKey/SecretKey/BucketName/UseSSL/ForcePathStyle |
tos | tos:// | file/tos.go | Volcano Engine TOS | 위와 동일하며 임시 버킷 지원 |
obs | obs:// | file/obs.go | Huawei Cloud OBS | Endpoint/Region/AccessKey/SecretKey/BucketName/UseSSL |
ks3 | ks3:// | file/ks3.go | Kingsoft Cloud KS3 | Endpoint/Region/AccessKey/SecretKey/BucketName |
dummy | dummy:// | file/dummy.go | 테스트용 빈 구현 | 없음 |
객체 Key 구성 규칙
- 정식 파일:
{tenantID}/{knowledgeID}/{uuid 또는 나노초 타임스탬프}{ext}. 예:local://12345/kb-001/1722045600000000000.pdf. - 내보내기/임시/복제 산출물:
{tenantID}/exports/{fileName}_{timestamp}{ext}. - 경로 보안:
secutils.SafePathUnderBase(디렉터리 순회 방지),secutils.SafeFileName, 객체 스토리지 측utils.SafeObjectKey.
두 래핑 계층
backend_scoped.go: 여러 스토리지 백엔드를 배포할 때 경로에storage://{backendID}/{innerPath}형태의 인스턴스 접두사를 붙입니다.wrap/unwrap으로 인코딩/디코딩하고 백엔드 간 작업을 거부합니다. KB는StorageBackendID로 특정 백엔드 인스턴스에 바인딩할 수 있습니다.resource_catalog.go: 물리 경로를 안정적인resource://{uuid}참조로 등록하며Bind(리소스를 knowledge 등의 owner와 연결),MarkDeleted,CreateAccessGrant(임시 접근 토큰을 생성하여/r/{token}형태 URL 산출)를 지원합니다. 애플리케이션 계층은resource://참조만 보유하면 하위 스토리지를 투명하게 이전할 수 있습니다.
비동기 작업 메커니즘(Asynq + Redis)
큐 등록
knowledge_create.go는 types.DocumentProcessPayload(TenantID/KnowledgeID/KnowledgeBaseID/FilePath/FileName/FileType/EnableMultimodel/EnableQuestionGeneration/QuestionCount/Language/Attempt 등 포함)를 구성하며, 작업 옵션은 knowledge_task_options.go에서 가져옵니다.
opts := []asynq.Option{
asynq.Queue(types.QueueDefault),
asynq.Timeout(config.DocumentProcessTimeout(cfg)), // 기본 30분
asynq.MaxRetry(3), // 실패 시 최대 3회 재시도
}
task := asynq.NewTask(types.TypeDocumentProcess, payloadBytes, opts...)
info, err := s.task.Enqueue(task)큐 등록에 실패하면 ParseStatus를 failed로 설정합니다(파일은 이미 저장되어 있으므로 reparse로 다시 실행할 수 있음).
큐 토폴로지와 Worker 풀
internal/types/task.go에 정의된 큐는 다음과 같습니다.
| 큐 상수 | 이름 | 용도 |
|---|---|---|
QueueDefault | default | 핵심 문서 처리(파싱/청크 분할/임베딩/인덱싱) |
QueuePostProcess | postprocess | 후처리 오케스트레이션 작업 |
QueueSummary | summary | 요약 / 질문 생성 계열 LLM 작업 |
QueueMultimodal | multimodal | 이미지 OCR / VLM 캡션 |
QueueMaintenance | low | 유지보수 작업(FAQ 일괄 가져오기 등) |
기본 동시 실행 수(internal/types/task.go)는 핵심 풀 DefaultCoreWorkerConcurrency = 8, 후처리 풀 2, 보강 풀 12, 유지보수 풀 4입니다.
실패 재시도 의미 체계
TypeDocumentProcess:MaxRetry(3)→ 최초 실행 + 재시도 3회로 총 4회 시도하며, 각 시도에는DocumentProcessTimeout(기본 30분)이 적용됩니다.- Payload는
Attempt를 포함합니다(다시 파싱할 때 과거 최대 attempt+1). Span Tracker는 attempt로 각 처리 회차의 진행률 트리를 격리하며, 새 attempt는 이전 작업의 마무리 동작을 "대체"(supersede)합니다. - 처리 함수는 "마지막 asynq 시도인지"(
isLastRetry)를 구분합니다. 마지막이 아닌 시도의 실패는 오류를 그대로 반환하여 asynq가 재시도하게 하고, 마지막 시도에서만ParseStatus를failed로 저장하고ErrorMessage를 기록합니다.
처리 설정: KB 기본값 + 개별 업로드 재정의
knowledge_process_config.go의 ResolveProcessConfig(kb, overrides)는 KB 기본 설정과 업로드 시 전달된 process_config(types.KnowledgeProcessOverrides)를 병합하여 types.EffectiveProcessConfig를 만듭니다.
- 재정의 가능 항목:
ChunkingConfig(chunk 크기/겹침/전략/부모-자식 청크 분할 등),EnableMultimodel,VLMConfig,ASRConfig,QuestionGenerationConfig,GraphEnabled,ExtractConfig,ParserEngineRules. - 제약:
eff.GraphEnabled = eff.GraphEnabled && eff.ExtractConfig.Enabled(그래프는 추출 설정 활성화에 의존). ValidateProcessOverrides는 파일 유형에 따라 사전 검증합니다. 이미지 업로드에는 VLM 모델, 오디오 업로드에는 ASR 모델이 반드시 설정되어야 하며, 멀티모달에는 완전한 객체 스토리지 설정도 필요합니다(validateImageMultimodalConfig).- 재정의 설정은
knowledge.SetProcessOverrides를 통해 Knowledge 행에 영속화되며 reparse 때도 그대로 사용됩니다.
어떤 파이프라인을 실행할지는 KB의 IndexingStrategy(internal/types/indexing_strategy.go)가 결정합니다.
type IndexingStrategy struct {
VectorEnabled bool // 의미 벡터 인덱스
KeywordEnabled bool // BM25 키워드 인덱스
WikiEnabled bool // Wiki 페이지 자동 생성
GraphEnabled bool // 지식 그래프 추출
}NeedsEmbedding() = Vector || Keyword, NeedsChunks() = 하나라도 활성화입니다. 기본값은 vector+keyword 활성화입니다.
핵심 처리 파이프라인(knowledge_process.go)
Worker는 TypeDocumentProcess를 가져온 뒤 표준화된 다섯 단계로 진행하며, 각 단계는 하나의 Span에 대응합니다(Housekeeping 자가 복구(knowledge_housekeeping.go) 참고).
docreader → chunking → embedding → multimodal → postprocess
파싱(convert, Stage: docreader)
beginStage(StageDocReader)가 입력(file_name/file_type/is_url)을 기록합니다.- URL 모드에서는
ValidateURLForSSRF를 다시 실행하며, 실패하면 즉시failStage+ParseStatus=failed로 처리합니다. - 엔진 선택:
eff.ChunkingConfig.ResolveParserEngine(fileType)(URL은 가상 유형"url"사용)으로 KB에 설정된ParserEngineRules(파일 유형 → 엔진)에 따라 라우팅합니다.MergeParserEngineOverrides는 테넌트 수준 및 업로드 수준 엔진 매개변수 재정의를 병합합니다. resolveDocReader는interfaces.DocReader를 반환합니다.- builtin: gRPC(
docparser/grpc_parser.go)또는 HTTP(http_parser.go)를 통해 Python docreader 서비스를 호출합니다. - simple: Go 네이티브로 md/txt/csv/json/이미지/오디오를 파싱합니다(
builtin_converter.go, CSV→Markdown 표, JSON→재귀 분할 코드 블록, 이미지/오디오→자리표시자 참조). - anydoc: Go 프로세스 안에서 docx/doc/pptx/ppt/xlsx/xls/odf/rtf/epub/csv/pdf를 파싱합니다(
anydoc_reader.go). 내부적으로 cgo로 연결한 anydoc Rust 라이브러리를 사용합니다. office 문서의 포함 이미지는 문서 모델에 따라 Markdown의 원래 위치에 다시 삽입되며, 텍스트 레이어가 없는 스캔 PDF는 DocReader를 사용할 수 있을 때 builtin 전체 페이지 렌더링으로 폴백합니다.anydoc빌드 태그가 있는 바이너리에서만 사용할 수 있고, 나머지 빌드에서는 엔진 목록에 사용할 수 없는 것으로 표시됩니다. - weknoracloud / mineru / mineru_cloud / paddleocr_vl / paddleocr_vl_cloud: HTTP 변환기입니다(
engines.go에 등록되며mineru_endpoint,mineru_api_key,paddleocr_vl_endpoint등의 설정으로 가용성 판정).
- builtin: gRPC(
엔진 카탈로그는 internal/infrastructure/docparser/engines.go에 모여 있습니다. 각 엔진은 메타데이터(이름, 설명, 파일 유형, 가용성 프로브)와 NewReader 팩터리를 함께 선언합니다. docparser.NewReader는 이름에 따라 분기하며, 등록되지 않은 이름(예: docreader에만 있는 markitdown)은 docreader 클라이언트로 전달됩니다. 5. 파일 모드에서는 FileService.GetFile(payload.FilePath)에서 바이트를 읽어 ReadRequest.FileContent에 채웁니다.
docreader 서비스 측(docreader/, Python gRPC): proto는 docreader/proto/docreader.proto에 정의되며, 서비스 메서드는 Read / ReadStream(스트리밍: 첫 프레임은 meta, 이후 이미지마다 한 프레임을 전송하여 대형 스캔 PDF가 gRPC 메시지 한도를 초과하지 않게 함)/ ListEngines입니다. 내장 parser는 docx/doc/pdf/md/xlsx/xls/epub/html/htm/mhtml/이미지/웹페이지(URL은 WebParser가 처리)를 지원하며, markitdown(Microsoft MarkItDown)과 opendataloader(PDF 레이아웃 분석, Java 11+ 필요)엔진을 선택적으로 등록할 수 있습니다. Go 측은 ListEngines로 원격 엔진을 자동 탐색합니다. 반환 형식은 ReadResult{MarkdownContent, ImageRefs, Metadata, IsAudio, AudioData}로 통일됩니다. 즉, 파싱 산출물은 항상 Markdown 텍스트 + 이미지 바이트이며 이미지 영속화는 Go 측이 담당합니다.
ASR 전사(오디오 파일)
convertResult.IsAudio가 참일 때(오디오 파일이 자리표시자 + 원본 바이트로 파싱됨):
asrModel, err := s.modelService.GetASRModel(ctx, eff.ASRConfig.ModelID)
transcriptionResult, err := asrModel.Transcribe(ctx, convertResult.AudioData, knowledge.FileName)전사 텍스트로 MarkdownContent를 교체한 뒤 일반 텍스트 파이프라인을 계속 실행합니다. ASR이 설정되지 않았으면 즉시 실패합니다.
이미지 추출 및 업로드
docparser/image_resolver.go의 ImageResolver.ResolveAndStore는 다음을 수행합니다.
<!link>로 감싼 이미지,data:URI, HTML 인라인 base64, 순수 base64, docreader가 반환한ImageRefs인라인 바이트를 차례로 처리합니다.- 아이콘 수준의 작은 이미지(가로·세로 < 64px 또는 < 512바이트,
IsOriginal=true인 원본 업로드 파일 제외)를 필터링합니다. SaveBytes로 현재 KB의 스토리지 백엔드에 업로드하고savedRefs캐시로 중복을 제거합니다.- Markdown의 참조를 스토리지 URL로 다시 씁니다(
markdown_image_scanner.go가위치를 정확히 찾음).
이후 ResolveRemoteImages가 Markdown의 외부 http(s) 이미지도 다운로드하여 저장합니다(마찬가지로 SSRF 방어 적용). 산출된 storedImages []docparser.StoredImage는 멀티모달 단계에서 사용됩니다.
청크 분할(Stage: chunking)
청크 분할은 Go 측에서 수행됩니다(internal/infrastructure/chunker, 자세한 내용은 《청크 분할 메커니즘》 장 참고).
chunkCfg := buildSplitterConfigFromChunking(eff.ChunkingConfig)
if eff.ChunkingConfig.EnableParentChild {
parentCfg, childCfg := buildParentChildConfigs(eff.ChunkingConfig, chunkCfg)
pcResult := chunker.SplitParentChild(convertResult.MarkdownContent, parentCfg, childCfg)
// children → types.ParsedChunk(ParentIndex 포함); parents → ParsedParentChunk
} else {
splitChunks := chunker.Split(convertResult.MarkdownContent, chunkCfg)
}DB 기록 및 인덱싱(processChunks, Stage: chunking + embedding)
processChunks는 핵심 조립 함수입니다.
- 부모 블록(부모-자식 청크 분할 모드): 각 parent에 대해
ChunkTypeParentText레코드를 만들고PreChunkID/NextChunkID연결 리스트를 구성합니다. 부모 블록은 DB에만 저장하고 벡터 인덱스에는 넣지 않습니다(검색에서 자식 블록이 일치한 뒤 부모 블록 콘텐츠를 다시 가져옴). - 텍스트 블록: 각
ParsedChunk에 대해ChunkTypeText레코드를 만들며StartAt/EndAt(원문의 rune 오프셋으로 복원/강조에 사용 가능)과 메모리 상태의ContextHeader(제목 이동 경로, DB에는 저장하지 않음)를 포함합니다. 부모-자식 모드에서는ParentChunkID도 기록합니다. chunkService.CreateChunks(ctx, insertChunks)가 DB에 일괄 기록합니다. 실패하면ParseStatus=failed+failStage(StageChunking)로 처리합니다.- 벡터화 및 인덱싱(
kb.NeedsEmbeddingModel()일 때):
indexContent := titlePrefix + chunk.EmbeddingContent() // 제목 + 이동 경로 + 콘텐츠
indexInfoList = append(indexInfoList, &types.IndexInfo{
Content: indexContent, SourceID: chunk.ID, SourceType: types.ChunkSourceType,
ChunkID: chunk.ID, KnowledgeID: knowledge.ID, KnowledgeBaseID: ..., IsEnabled: true,
})
err = retrieveEngine.BatchIndex(ctx, embeddingModel, indexInfoList)인덱싱 실패 시 보상 롤백을 실행합니다. 이미 기록한 chunks를 삭제하고(DeleteChunksByKnowledgeID)벡터 인덱스를 정리한 뒤(DeleteByKnowledgeIDList)failed로 설정하여 불완전한 산출물이 남지 않게 합니다. 5. 이미지 멀티모달 작업 팬아웃: enableMultimodel && len(storedImages) > 0이면 enqueueImageMultimodalTasks가 이미지마다 TypeImageMultimodal 작업 하나를 QueueMultimodal에 등록합니다. payload에는 ImageURL/EnableOCR/EnableCaption/Attempt/ImageIndex가 포함됩니다. 6. finalizeIndexedKnowledgeState: 실행할 멀티모달/후처리가 남아 있으면 processing을 유지하고, 없으면 바로 completed로 전환합니다. 동시에 EnableStatus="enabled"로 설정하고(이때부터 문서를 검색할 수 있음)테넌트 스토리지 사용량을 누적합니다.
후처리 오케스트레이션(knowledge_post_process.go, Stage: postprocess)
모든 멀티모달 작업이 완료되면(또는 멀티모달이 없으면)TypeKnowledgePostProcess를 큐에 등록합니다. 이 작업은 보강 하위 작업의 오케스트레이터이며, 원자적 카운터로 최종 상태 수렴을 보장합니다.
willSpawnSummary := len(textChunks) > 0
willSpawnQuestion := willSpawnSummary && kb.NeedsEmbeddingModel() && eff.QuestionGenerationConfig.Enabled
willSpawnWiki := kb.IndexingStrategy.WikiEnabled && len(textChunks) > 0
willSpawnGraph := eff.GraphEnabled && len(textChunks) > 0
// questionGenChunkBatchSize = 20: 질문 생성은 chunk 20개씩 배치 처리
expectedSubtasks = summary(0/1) + questionBatchCount + wiki(0/1) + graphChunkCount
// parse_status를 processing에서 finalizing으로 원자적으로 승격하고 pending_subtasks_count 기록
promoted, err := s.knowledgeRepo.SetFinalizing(ctx, payload.KnowledgeID, expectedSubtasks)expectedSubtasks == 0이면 빠른 경로로 바로completed가 됩니다.- 각 하위 작업은 최종 상태로 종료할 때
FinalizeSubtask를 호출하여pending_subtasks_count를 원자적으로 감소시키며, 0이 되면 자동으로completed로 승격합니다. - 부족분 조정: 실제 큐 등록 수가 계획보다 적으면(예: 특정 큐 등록 실패)즉시 차이만큼 보상 감소하여
finalizing에 영원히 멈추지 않게 합니다. finalizeSubtaskDetached(knowledge.go): 감소 동작은context.WithoutCancel+ 10초 타임아웃의 분리된 컨텍스트에서 실행됩니다. worker가 정상 종료할 때 ctx 취소로 카운트가 유실되어 지식이finalizing에 영구 정체되는 것을 방지합니다.
네 종류의 보강 하위 작업은 다음과 같습니다.
| 작업 | 큐 | 단위 | 설명 |
|---|---|---|---|
TypeSummaryGeneration | summary | 지식당 1개 | 문서 요약 생성, summary_status 독립 상태 머신 |
TypeQuestionGeneration | question 큐 | chunk 20개당 한 배치 | chunk의 검색 질문 생성 |
TypeChunkExtract | graph 큐 | chunk당 1개 | 엔터티/관계 추출 후 그래프 엔진에 기록 |
TypeWikiIngest | wiki 큐 | 디바운스 일괄 처리 | Wiki 페이지 생성/갱신 |
문서 자동 태그
파싱 후 후처리는 지식 베이스의 auto_tag_config에 따라 knowledge:auto_tag(summary 큐)등록 여부를 결정합니다. 처리기는 현재 KB의 기존 태그 중에서 선택하고 모델은 summary_model_id로 폴백합니다. max_tags 기본값은 3, 상한은 10이며 기본적으로 태그가 이미 있는 문서는 건너뜁니다. 태그 연결만 증분 추가하고 새 분류를 만들거나 수동 태그를 삭제하지 않으며, 실패해도 문서 수집을 막지 않습니다. 설정 및 새 파싱/다시 파싱에 적용되는 범위는 지식 베이스 관리를 참고하세요.
요약 새로고침(knowledge_summary_refresh.go)
최초 수집 외에도 청크 콘텐츠 편집, 청크 활성화/비활성화, 사용자 정의 메타데이터 변경으로 기존 요약이 오래된 상태가 되면 요약 새로고침을 한 번 큐에 등록합니다(POST /knowledge/:id/regenerate-summary로 수동 실행할 수도 있음).
- 작업 시작 시 각 원본 청크의
content_revision/is_enabled와custom_metadata버전으로 입력 스냅샷을 기록합니다. - 생성 완료 후
summarySourceChanged()로 스냅샷을 다시 확인합니다. 그사이 다시 편집되었다면ErrSummaryRefreshStale을 반환하고 이번 결과를 버리며summary_status도 변경하지 않습니다. 더 최신 새로고침이 마무리하도록 하여 이전 요약이 새 요약을 덮어쓰지 않게 합니다. - 데이터베이스 읽기 실패와 "입력이 변경됨"을 별도로 처리하여 일시적인 읽기 오류를 오래된 작업으로 오인해 조용히 버리지 않게 합니다.
- 새로고침 작업은 HTTP 미들웨어가 테넌트 컨텍스트를 주입하지 않는 Asynq worker에서 실행되므로
restoreSummaryRefreshTenantInfo()가 완전한 테넌트 설정을 복원합니다. 검색 엔진 팩터리에 이 설정이 필요합니다.
이미지 멀티모달(image_multimodal.go)
ImageMultimodalService.Handle은 단일 이미지 작업을 처리합니다.
readImageBytes가 스토리지/URL에서 이미지를 가져오고,resolveVLM이 KB의 VLM 설정을 가져옵니다.- 캡션(VLM, prompt는
buildVLMCaptionPrompt가DescriptionLanguage/CustomInstructions에 따라 구성)과 OCR 텍스트를 생성합니다. - 결과를 소속 텍스트 Chunk의
ImageInfo(JSON)에 다시 기록하고 자식 Chunk 두 개, 즉ChunkTypeImageCaption과ChunkTypeImageOCR을 생성/갱신합니다.ParentChunkID는 텍스트 블록을 가리키며, 이후indexChunks로 벡터 인덱스에 별도 등록합니다. 따라서 "이미지 설명을 검색해도 원문 블록을 후보 검색"할 수 있습니다. shouldDropOrphanedMultimodal은 부모 블록이 삭제되거나 대체되었는지 확인하며, 고아 작업은 즉시 버립니다.checkAndFinalizeAllImages: 모든 이미지 처리가 완료되면enqueueKnowledgePostProcessTask가 후처리 오케스트레이션(knowledge_post_process.go, Stage: postprocess)을 실행합니다.
상태 머신
Knowledge 기본 상태(ParseStatus)
internal/types/knowledge.go에 정의된 전체 값은 다음과 같습니다.
| 값 | 의미 |
|---|---|
pending | 생성 완료, worker가 가져가기를 기다리는 중 |
processing | 파싱/청크 분할/임베딩/멀티모달 실행 중 |
finalizing | 기본 흐름 완료, 보강 하위 작업 대기 중(pending_subtasks_count > 0) |
completed | 모두 완료 |
failed | 처리 실패(원인은 ErrorMessage에 기록) |
deleting | 삭제 중(동시 실행 방지 표시) |
cancelled | 사용자가 파싱 취소 |
보조 상태: EnableStatus ∈ {enabled, disabled}(검색 가능 여부. 인덱싱 성공 즉시 enabled이며 보강을 기다리지 않음), SummaryStatus ∈ {none, pending, processing, completed, failed}.
단계별 진행률(Span Tracker)
knowledge_span_tracker.go + internal/types/knowledge_span.go는 단계별 진행률 트리를 제공합니다(프런트엔드 타임라인이 이를 바탕으로 렌더링됨).
- 다섯 표준 단계:
StageDocReader / StageChunking / StageEmbedding / StageMultimodal / StagePostProcess(types.AllStages). - Span 상태:
pending / running / done / failed / skipped / cancelled.skipped는 의도적인 건너뛰기(예: 멀티모달 미활성화),cancelled는 상위 단계 실패에 따른 연쇄 취소에 사용합니다. - 각 처리 회차는 독립된
Attempt(repo.NextAttempt)를 가지며, 루트 Span은name="knowledge_processing",Kind=SpanKindRoot입니다. 단계는beginStage / endStage / failStage / skipStage로 기록하고, 입출력은JSONMap(예:chunks_planned/chunks_written/total_text_chars)에 기록합니다. - 진행률을 기록할 때마다
touchKnowledgeHeartbeat가 하트비트도 갱신하여 Housekeeping이 여전히 처리 중인 작업과 갱신이 멈춘 작업을 구분할 수 있게 합니다.
Housekeeping 자가 복구(knowledge_housekeeping.go)
백그라운드에서 5분마다 실행되며(WEKNORA_HOUSEKEEPING_ENABLED로 비활성화 가능), worker 충돌 / Redis 작업 유실로 생긴 좀비 상태를 복구합니다.
Sweep A —— 정체된 지식 복구는 세 단계로 필터링합니다.
- 1차 선별:
parse_status IN (pending, processing, finalizing) AND updated_at < cutoff. filterByLastSpanActivity:knowledge_processing_spans의MAX(updated_at)하트비트를 조회합니다. 하트비트가 여전히 임계값 안이면 유지하고(아직 처리 중), span이 전혀 없는 경우도 정체로 판정합니다.filterOutQueued: asynq TaskInspector로 대기 작업이 남아 있는지 확인하고, 있으면 유지합니다(단순히 대기 중).
정체로 판정된 지식은 다음과 같이 갱신됩니다.
UPDATE knowledge SET parse_status = 'failed',
error_message = 'task stuck in processing > [threshold], recovered by housekeeping',
pending_subtasks_count = 0
WHERE id IN (stuck_ids)임계값은 staleThreshold() = max(1h, DocumentProcessTimeout) + 10min입니다.
Sweep B —— 정체된 요약 복구: summary_status = 'processing' AND updated_at < 1시간 전 → failed로 설정합니다.
삭제 정리 경로(knowledge_delete.go)
DeleteKnowledge(ctx, id)의 순서는 신중하게 설계되었습니다(행을 먼저 삭제하고 파일은 나중에 삭제, 실패 시 재시도 가능).
ParseStatus = deleting으로 표시하여 동시 실행 작업의 기록을 막습니다.pending/processing상태의 지식에는dequeueKnowledgeTasks()를 실행하여 큐의 하위 작업을 취소합니다.- errgroup으로 네 종류의 리소스를 병렬 정리합니다.
- 벡터/키워드 인덱스:
retrieveEngine.DeleteByKnowledgeIDList(embedding 차원과 KB 유형에 따라 라우팅). - Wiki:
cleanupWikiOnKnowledgeDelete(Redis tombstone 기록 → pending ingest 정리 → 기존 페이지 reconcile → WikiRetract 큐 등록). - Chunks:
chunkService.DeleteChunksByKnowledgeID. - 그래프:
graphEngine.DelGraph.
- 벡터/키워드 인덱스:
- Tag 연결을 삭제한 뒤 Knowledge 데이터베이스 행을 삭제합니다.
- 마지막으로 best-effort 방식으로 물리 파일을 정리합니다. 원본 파일 +
chunk_image_info에서 수집한 모든 추출 이미지(collectImageURLs+deleteExtractedImages)를 삭제하고 테넌트 스토리지 통계를 차감합니다.
일괄 버전 DeleteKnowledgeList는 각 KB의 FileService를 미리 로드하고, KB별로 이미지 URL을 그룹화하며, embedding 모델별로 인덱스를 그룹 삭제하여 goroutine 안의 반복 조회를 피합니다.
FAQ 유형 지식 가져오기(knowledge_faq.go / knowledge_faq_import.go)
FAQ 지식 베이스는 문서 파싱 파이프라인을 사용하지 않습니다. 각 FAQ KB에는 Knowledge 인스턴스가 하나만 있으며(ensureFAQKnowledge), 각 질의응답 쌍은 ChunkTypeFAQ 유형 Chunk 하나이고 메타데이터는 Chunk.Metadata에 저장됩니다.
type FAQChunkMetadata struct {
StandardQuestion string // 표준 질문
SimilarQuestions []string // 유사 질문
NegativeQuestions []string // 반례 질문(부정 예 필터링, 인덱싱에는 참여하지 않음)
Answers []string
AnswerStrategy AnswerStrategy // "all" | "random"
...
}- 단일 생성
CreateFAQEntry: 정제 및 검증 → 중복 확인(checkFAQQuestionDuplicate)→ Chunk 구성(buildFAQChunkContent가FAQIndexMode에 따라 답변을 Content에 기록할지 결정)→indexFAQChunks동기 인덱싱 →ChunkStatusIndexed. - 인덱스 모드(KB 수준 설정):
FAQIndexModeQuestionOnly(question_only, 질문만 인덱싱)/FAQIndexModeQuestionAnswer(question_answer, 질문+답변). 질문 인덱싱은 다시FAQQuestionIndexModeCombined(표준 질문+유사 질문을 벡터 하나로 병합)와FAQQuestionIndexModeSeparate(각 유사 질문을 독립 벡터로 구성, source_id는{chunkID}-{index}형태이며 증분 인덱싱incrementalIndexFAQEntry지원)로 나뉩니다. - 일괄 가져오기
UpsertFAQEntries:- 모드는
append(추가/병합)또는replace(전체 교체)이며,DryRun으로 검증만 수행할 수 있습니다. - 200개 또는 50KB를 초과하면 항목을 먼저
SaveBytes로 객체 스토리지에 업로드하고 payload에는EntriesURL만 넣습니다. TypeFAQImport→QueueMaintenance에 등록하며MaxRetry 5(dry-run은 3), Timeout은 2시간입니다. 같은 KB에서는 동시에 하나의 가져오기 작업만 허용합니다(Redis 잠금).- 중복 제거는
CalculateFAQContentHash를 기반으로 합니다. 표준 질문/유사 질문/반례/답변을 정규화(URL 제거, 소문자 변환, 번체→간체, 전각→반각, 지능형 공백)한 뒤 SHA256을 계산합니다. - append 모드는 네 단계 검증을 수행합니다(표준 질문 충돌→항목 전체 실패, 유사 질문/반례 충돌→충돌 항목만 제거하는 부분 실패, 표준 질문이 이미 존재→합집합으로 병합).
- 진행률은 Redis에 기록합니다(
FAQImportProgress:pending/processing/completed/failed, 성공/실패/부분 실패/건너뜀 수). 실패 항목은 다운로드할 수 있도록 UTF-8 BOM이 포함된 CSV로 내보냅니다.
- 모드는
지식 복제 및 이동(knowledge_clone_move.go)
복제(CloneKnowledgeBase / CloneChunk)
- KB 수준 복제는 먼저 KB 설정을 복사하고, 그다음 집합 차(
AminusB)에 따라 Knowledge를 추가/삭제하며 병렬 처리합니다(삭제 배치 10개, 복제는 하나씩). - Chunk 수준 복제(배치 100개)는
Text/ParentText/Summary/ImageCaption/ImageOCR다섯 종류 chunk를 복사합니다.- 이미지 깊은 복사:
cloneChunkImageInfo가 원본 스토리지에서 바이트를 읽어 → 대상 테넌트의exports/네임스페이스에 기록하고urlCache로 중복을 제거합니다.rewriteContentImageURLs는 Content의 기존 URL을 모두 교체합니다(부분 일치를 피하려고 긴 URL부터 처리). - 태그 매핑
getOrCreateTagInTarget(같은 이름이 있으면 재사용, 없으면 생성). PreChunkID/NextChunkID/ParentChunkID매핑을 재구성한 뒤 일괄 삽입합니다.- 벡터 인덱스는
retrieveEngine.CopyIndices()로 직접 복사하며 embedding을 다시 계산하지 않습니다.
- 이미지 깊은 복사:
- FAQ KB 복제는 차등 동기화를 사용합니다.
chunkRepo.FAQChunkDiff가content_hash를 기준으로 추가/삭제/일치 세 그룹을 계산하고, 일치 쌍은 상태(IsEnabled/Flags/TagID/AnswerStrategy)만 동기화합니다. - 진행률은 Redis에 기록합니다(
KBCloneProgress).
이동(ProcessKnowledgeMove)
사전 조건은 원본/대상 KB의 유형이 같고 EmbeddingModelID도 같아야 한다는 것입니다. 모드는 두 가지입니다.
reuse_vectors:sourceKB.SharesStoreWith(targetKB)(동일한 벡터 스토리지 인스턴스)가 필요합니다.CopyIndices로 인덱스 복사 → 원본 인덱스 삭제 →MoveChunksByKnowledgeID로 chunk 소속 변경 → Tag 연결 정리 → Knowledge의 KB ID 갱신 순서로 처리합니다.reparse: 서로 다른 벡터 스토리지 사이에서 사용합니다.cleanupKnowledgeResources(인덱스/chunks/그래프 삭제/스토리지 통계 차감)→ Knowledge를pending으로 초기화하여 대상 KB에 연결 →TypeDocumentProcess를 다시 큐에 등록합니다(manual 유형은triggerManualProcessing사용).
엔드투엔드 시퀀스 요약
멀티모달, 질문 생성, 그래프가 활성화된 PDF 한 편의 전체 여정은 다음과 같습니다.
POST /knowledge-bases/:id/knowledge/file→ MD5 중복 제거 →cos://tenant/kb/uuid.pdf→ Knowledge(pending) → asynqdocument:process.- Worker: Span attempt=1 루트 생성 →
docreader단계에서 gRPC로 Python 서비스를 호출하여 Markdown+이미지 바이트 획득 → 이미지를 스토리지에 업로드하고 URL 다시 쓰기 →chunking단계에서 Go chunker가 블록 분할 → chunks 기록 →embedding단계에서 BatchIndex →EnableStatus=enabled(이때 이미 검색 가능)→ 이미지마다 multimodal 작업 큐 등록. - 멀티모달 worker가 이미지별로 OCR+Caption을 수행하여 image_caption/image_ocr 자식 chunk를 생성하고 인덱싱합니다. 모두 완료되면 post-process를 실행합니다.
- 오케스트레이터가
expectedSubtasks(요약 1개 + 질문 배치 N/20개 + 그래프 M개 + Wiki 0/1개)를 계산 →SetFinalizing→ 팬아웃합니다. 각 하위 작업의 최종 상태에서FinalizeSubtask가 감소하고 0이 되면 →completed. - 그사이 어느 단계든 정체되면 Housekeeping이 5분마다 "updated_at + span 하트비트 + 큐 검사"의 삼중 기준으로
failed로 회수하며, 사용자는 reparse(attempt+1)로 다시 시작할 수 있습니다.
구현 참고 자료
각 단계에 해당하는 소스 위치는 다음과 같습니다.
| 단계 | 소스 위치 |
|---|---|
| HTTP 진입점 | internal/handler/knowledge.go, internal/router/router.go |
| 생성 및 큐 등록 | internal/application/service/knowledge_create.go, knowledge_task_options.go |
| 파일 스토리지 | internal/application/service/file/(factory.go, 각 백엔드 구현) |
| 파싱 인프라 | internal/infrastructure/docparser/, docreader/(Python 서비스) |
| 기본 처리 파이프라인 | internal/application/service/knowledge_process.go |
| 처리 설정 병합 | internal/application/service/knowledge_process_config.go |
| 후처리 | internal/application/service/knowledge_post_process.go, image_multimodal.go |
| 진행률 추적 | internal/application/service/knowledge_span_tracker.go, internal/types/knowledge_span.go |
| 자가 복구 | internal/application/service/knowledge_housekeeping.go |
| 삭제 | internal/application/service/knowledge_delete.go |
| FAQ | internal/application/service/knowledge_faq.go, knowledge_faq_import.go |
| 복제/이동 | internal/application/service/knowledge_clone_move.go |