diff --git a/cmd/server/main.go b/cmd/server/main.go index ad8e33e..f8a84c6 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -17,6 +17,8 @@ import ( "github.com/CeruleanFlow/cerulean/internal/search" "github.com/CeruleanFlow/cerulean/internal/storage" "github.com/CeruleanFlow/cerulean/internal/task" + + docparser "github.com/CeruleanFlow/cerulean/internal/parser" ) func main() { @@ -31,9 +33,10 @@ func main() { log.Fatal(err) } taskManager := task.NewMemoryManager() + documentParser := docparser.NewPDFTextParser() searchBackend, err := buildSearchBackend(cfg) ragService := rag.NewService(paperRepo, searchBackend) - ingestService := ingest.NewService(paperRepo, chunkRepo, objectStore, taskManager, searchBackend) + ingestService := ingest.NewService(paperRepo, chunkRepo, objectStore, taskManager, searchBackend, documentParser) router := api.NewRouter(api.RouterOptions{ Config: cfg, diff --git a/go.mod b/go.mod index 99d5dd4..4b5a179 100644 --- a/go.mod +++ b/go.mod @@ -5,6 +5,7 @@ go 1.26.0 require ( github.com/gin-gonic/gin v1.12.0 github.com/joho/godotenv v1.5.1 + github.com/ledongthuc/pdf v0.0.0-20250511090121-5959a4027728 github.com/minio/minio-go/v7 v7.2.1 gorm.io/datatypes v1.2.7 gorm.io/driver/mysql v1.6.0 diff --git a/go.sum b/go.sum index 8926e00..df595ad 100644 --- a/go.sum +++ b/go.sum @@ -71,6 +71,8 @@ github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/ledongthuc/pdf v0.0.0-20250511090121-5959a4027728 h1:QwWKgMY28TAXaDl+ExRDqGQltzXqN/xypdKP86niVn8= +github.com/ledongthuc/pdf v0.0.0-20250511090121-5959a4027728/go.mod h1:1fEHWurg7pvf5SG6XNE5Q8UZmOwex51Mkx3SLhrW5B4= github.com/leodido/go-urn v1.4.0 h1:WT9HwE9SGECu3lg4d/dIA+jxlljEa1/ffXKmRjqdmIQ= github.com/leodido/go-urn v1.4.0/go.mod h1:bvxc+MVxLKB4z00jd1z+Dvzr47oO32F/QSNjSBOlFxI= github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= diff --git a/internal/ingest/service.go b/internal/ingest/service.go index dcc4754..9e86bfb 100644 --- a/internal/ingest/service.go +++ b/internal/ingest/service.go @@ -4,6 +4,8 @@ import ( "bytes" "context" "fmt" + "io" + "os" "strings" "time" @@ -12,6 +14,8 @@ import ( "github.com/CeruleanFlow/cerulean/internal/search" "github.com/CeruleanFlow/cerulean/internal/storage" "github.com/CeruleanFlow/cerulean/internal/task" + + docparser "github.com/CeruleanFlow/cerulean/internal/parser" ) type Service struct { @@ -20,15 +24,24 @@ type Service struct { store storage.ObjectStorage search search.Backend tasks task.Manager + parser docparser.Parser } -func NewService(papers repository.PaperRepository, chunks repository.ChunkRepository, store storage.ObjectStorage, tasks task.Manager, searchBackend search.Backend) *Service { +func NewService( + papers repository.PaperRepository, + chunks repository.ChunkRepository, + store storage.ObjectStorage, + tasks task.Manager, + searchBackend search.Backend, + parser docparser.Parser, +) *Service { return &Service{ papers: papers, chunks: chunks, - tasks: tasks, store: store, + tasks: tasks, search: searchBackend, + parser: parser, } } @@ -37,119 +50,210 @@ func (s *Service) StartPaperIngest(ctx context.Context, paperID string) (task.Ta if err != nil { return task.Task{}, err } + + optCtx, cancel := context.WithTimeout(ctx, 30*time.Second) + defer cancel() + now := time.Now() job := task.Task{ ID: fmt.Sprintf("task_%d", now.UnixNano()), PaperID: paperID, Type: "paper_ingest", Status: task.Queued, - Message: "queued local placeholder ingestion; PaddleOCR worker will replace this stage", + Message: "queued PDF text ingestion; PaddleOCR will be used later for scanned PDFs", CreatedAt: now, UpdatedAt: now, } paper.Status = domain.PaperProcessing paper.Error = "" paper.UpdatedAt = now - if err := s.papers.Update(ctx, paper); err != nil { + + if err := s.papers.Update(optCtx, paper); err != nil { return task.Task{}, err } - if err := s.tasks.Create(ctx, job); err != nil { + if err := s.tasks.Create(optCtx, job); err != nil { return task.Task{}, err } // Keep this asynchronous so the public API shape is already compatible with // the future PaddleOCR worker and queue based pipeline. - go s.runPlaceholderIngest(context.Background(), job, paper) + go func(job task.Task, paper domain.Paper) { + defer func() { + if r := recover(); r != nil { + s.fail(context.Background(), job, paper, fmt.Errorf("panic during paper ingest: %v", r)) + } + }() + + taskCtx, cancel := context.WithTimeout(context.Background(), 10*time.Minute) + defer cancel() + + if err := s.runPDFTextIngest(taskCtx, job, paper); err != nil { + s.fail(context.Background(), job, paper, err) + } + }(job, paper) return job, nil } -func (s *Service) runPlaceholderIngest(ctx context.Context, job task.Task, paper domain.Paper) error { +func (s *Service) runPDFTextIngest(ctx context.Context, job task.Task, paper domain.Paper) error { + // precheck + if s.parser == nil { + return fmt.Errorf("document parser is not initialized") + } + now := time.Now() job.Status = task.Running - job.Message = "building placeholder chunks" + job.Message = "parsing PDF text" job.UpdatedAt = now _ = s.tasks.Update(ctx, job) - markdown := placeholderMarkdown(paper) + // download the file to local + pdfPath, cleanup, err := s.downloadOriginalPDF(ctx, paper) + if err != nil { + return err + } + defer cleanup() + + // parse the file + doc, err := s.parser.ParseFile(ctx, pdfPath) + if err != nil { + return fmt.Errorf("parse pdf text: %w", err) + } + + // generate the markdown + markdown := parsedMarkdown(paper, doc) artifactKey := fmt.Sprintf("papers/%s/parsed/document.md", paper.ID) - if _, err := s.store.Put(ctx, artifactKey, bytes.NewReader([]byte(markdown)), int64(len(markdown)), storage.PutOptions{ContentType: "text/markdown; charset=utf-8"}); err != nil { - s.fail(ctx, job, paper, err) - return fmt.Errorf("store.Put: %w", err) + + // save the markdown to MinIO + if _, err := s.store.Put( + ctx, + artifactKey, + bytes.NewReader([]byte(markdown)), + int64(len(markdown)), + storage.PutOptions{ContentType: "text/markdown; charset=utf-8"}, + ); err != nil { + return fmt.Errorf("store parsed markdown: %w", err) } - chunks := makePlaceholderChunks(paper, artifactKey) + // generate the chunk + chunks := docparser.BuildChunks(paper, artifactKey, doc, 1200, 150) + if len(chunks) == 0 { + return fmt.Errorf("pdf text parser produced no chunks") + } + + // delete old chunks in mysql if err := s.chunks.DeleteByPaperID(ctx, paper.ID); err != nil { - s.fail(ctx, job, paper, err) - return fmt.Errorf("s.chunks.DeleteByPaperID: %w", err) + return fmt.Errorf("delete old chunks from mysql: %w", err) } + // delete old chunk index in es + if s.search != nil { + if err := s.search.DeleteByPaperID(ctx, paper.ID); err != nil { + return fmt.Errorf("delete old chunks from elasticsearch: %w", err) + } + } + // save new chunks in mysql if err := s.chunks.UpsertMany(ctx, chunks); err != nil { return fmt.Errorf("save chunks to mysql: %w", err) } + // save new chunks index in es if s.search != nil { if err := s.search.IndexChunks(ctx, chunks); err != nil { return fmt.Errorf("index chunks to elasticsearch: %w", err) } } + // update paper / task status now = time.Now() + paper.Status = domain.PaperParsed - paper.PageCount = 1 - paper.Error = "" paper.UpdatedAt = now - _ = s.papers.Update(ctx, paper) + paper.Error = "" + paper.PageCount = doc.PageCount + + if err := s.papers.Update(ctx, paper); err != nil { + return fmt.Errorf("update paper status: %w", err) + } + job.Status = task.Succeeded - job.Message = "placeholder ingestion completed; next step is wiring PaddleOCR output into the same chunk path" + job.Message = fmt.Sprintf("PDF text ingestion completed: %d pages, %d chunks", doc.PageCount, len(chunks)) job.UpdatedAt = now - _ = s.tasks.Update(ctx, job) + + if err := s.tasks.Update(ctx, job); err != nil { + return fmt.Errorf("update task status: %w", err) + } return nil } func (s *Service) fail(ctx context.Context, job task.Task, paper domain.Paper, err error) { + opCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + now := time.Now() + paper.Status = domain.PaperFailed paper.Error = err.Error() paper.UpdatedAt = now - _ = s.papers.Update(ctx, paper) + _ = s.papers.Update(opCtx, paper) + job.Status = task.Failed job.Message = err.Error() job.UpdatedAt = now - _ = s.tasks.Update(ctx, job) + _ = s.tasks.Update(opCtx, job) } -func placeholderMarkdown(paper domain.Paper) string { - return fmt.Sprintf("# %s\n\n"+ - "This is a placeholder parsed document for **%s**.\n\n"+ - "The current Cerulean pipeline has stored the original PDF object at `%s` and created searchable placeholder chunks. The next implementation step is to replace this placeholder with PaddleOCR output, page-level layout JSON, chunking, embedding, reranking, and hybrid retrieval.\n\n"+ - "Suggested future metadata:\n\n"+ - "- sha256: `%s`\n"+ - "- content_type: `%s`\n"+ - "- file_size: `%d`\n", paper.Title, paper.Filename, paper.ObjectKey, paper.SHA256, paper.ContentType, paper.Size) +// downloadOriginalPDF download original PDF to tmp +func (s *Service) downloadOriginalPDF(ctx context.Context, paper domain.Paper) (string, func(), error) { + if s.store == nil { + return "", nil, fmt.Errorf("object storage is not initialized") + } + + if err := os.MkdirAll(".var/temp", 0755); err != nil { + return "", nil, fmt.Errorf("create tmp dir: %w", err) + } + + reader, _, err := s.store.Get(ctx, paper.ObjectKey) + if err != nil { + return "", nil, fmt.Errorf("get original pdf from storage: %w", err) + } + defer reader.Close() + + tmp, err := os.CreateTemp(".var/tmp", "ingest-*.pdf") + if err != nil { + return "", nil, fmt.Errorf("create ingest temp file: %w", err) + } + + tmpPath := tmp.Name() + cleanup := func() { + _ = os.Remove(tmpPath) + } + + if _, err := io.Copy(tmp, reader); err != nil { + _ = os.Remove(tmpPath) + cleanup() + return "", nil, fmt.Errorf("copy pdf to temp file: %w", err) + } + + if err := tmp.Close(); err != nil { + cleanup() + return "", nil, fmt.Errorf("close temp file: %w", err) + } + return tmpPath, cleanup, nil } -func makePlaceholderChunks(paper domain.Paper, artifactKey string) []domain.Chunk { - now := time.Now() - base := []string{ - fmt.Sprintf("Paper title: %s. Original filename: %s. This paper has been uploaded into Cerulean and is ready for OCR parsing.", paper.Title, paper.Filename), - "Cerulean is a paper-oriented RAG system. The planned pipeline is PDF upload, MinIO artifact storage, PaddleOCR document parsing, chunking, AstraFlow embedding, Amaranth vector retrieval, Elasticsearch lexical retrieval, AstraFlow reranking, and DeepSeek answer generation.", - "This placeholder chunk exists so search and RAG screens can be tested before PaddleOCR is connected. After OCR is wired, this chunk will be replaced by page-aware content with page numbers, section labels, and citation information.", - } - chunks := make([]domain.Chunk, 0, len(base)) - for i, text := range base { - chunks = append(chunks, domain.Chunk{ - ID: fmt.Sprintf("%s_chunk_%03d", paper.ID, i+1), - PaperID: paper.ID, - PageNo: 1, - Index: i, - Text: strings.TrimSpace(text), - ObjectKey: artifactKey, - Metadata: map[string]string{ - "source": "placeholder_ingest", - "title": paper.Title, - }, - CreatedAt: now, - UpdatedAt: now, - }) - } - return chunks +func parsedMarkdown(paper domain.Paper, doc docparser.Document) string { + var b strings.Builder + + b.WriteString("# ") + b.WriteString(paper.Title) + b.WriteString("\n\n") + + b.WriteString("\n\n") + + for _, page := range doc.Pages { + b.WriteString(fmt.Sprintf("## Page %d\n\n", page.PageNo)) + b.WriteString(strings.TrimSpace(page.Text)) + b.WriteString("\n\n") + } + + return b.String() } diff --git a/internal/parser/chunker.go b/internal/parser/chunker.go new file mode 100644 index 0000000..8609432 --- /dev/null +++ b/internal/parser/chunker.go @@ -0,0 +1,104 @@ +package parser + +import ( + "fmt" + "strings" + "time" + + "github.com/CeruleanFlow/cerulean/internal/domain" +) + +func BuildChunks(paper domain.Paper, artifactKey string, doc Document, maxRunes, overlap int) []domain.Chunk { + if maxRunes <= 0 { + maxRunes = 1200 + } + if overlap < 0 { + overlap = 0 + } + + now := time.Now() + chunks := make([]domain.Chunk, 0) + + chunkIndex := 0 + + for _, page := range doc.Pages { + pageText := strings.TrimSpace(page.Text) + if pageText == "" { + continue + } + + pieces := splitTextByRunes(pageText, maxRunes, overlap) + + for _, piece := range pieces { + piece = strings.TrimSpace(piece) + if piece == "" { + continue + } + + chunks = append(chunks, domain.Chunk{ + ID: fmt.Sprintf("%s_chunk_%04d", paper.ID, chunkIndex+1), + PaperID: paper.ID, + PageNo: page.PageNo, + Index: chunkIndex, + Text: piece, + ObjectKey: artifactKey, + Metadata: map[string]string{ + "source": "pdf_text_parse", + "title": paper.Title, + }, + CreatedAt: now, + UpdatedAt: now, + }) + + chunkIndex++ + } + } + return chunks +} + +// splitTextByRunes split the text by runes +func splitTextByRunes(text string, maxRunes int, overlap int) []string { + runes := []rune(text) + if len(runes) == 0 { + return nil + } + + result := make([]string, 0) + + start := 0 + for start < len(runes) { + end := start + maxRunes + if end >= len(runes) { + end = len(runes) + } else { + end = chooseChunkEnd(runes, start, end) + } + + result = append(result, string(runes[start:end])) + + if end >= len(runes) { + break + } + + next := end - overlap + if next <= start { + next = end + } + start = next + } + + return result +} + +func chooseChunkEnd(runes []rune, start int, end int) int { + minEnd := start + int(float64(end-start)*0.7) + + for i := end; i > minEnd; i-- { + switch runes[i-1] { + case '\n', '.', '。', ';', ';', '!', '!', '?', '?': + return i + } + } + + return end +} diff --git a/internal/parser/document.go b/internal/parser/document.go new file mode 100644 index 0000000..7ca1ca8 --- /dev/null +++ b/internal/parser/document.go @@ -0,0 +1,18 @@ +package parser + +import "context" + +type Page struct { + PageNo int `json:"page_no"` + Text string `json:"text"` +} + +type Document struct { + PageCount int `json:"page_count"` + Pages []Page `json:"pages"` + Text string `json:"text"` +} + +type Parser interface { + ParseFile(ctx context.Context, path string) (Document, error) +} diff --git a/internal/parser/pdf_text.go b/internal/parser/pdf_text.go new file mode 100644 index 0000000..6d53291 --- /dev/null +++ b/internal/parser/pdf_text.go @@ -0,0 +1,74 @@ +package parser + +import ( + "context" + "errors" + "strings" + + "github.com/ledongthuc/pdf" +) + +var ( + ErrNoExtractableText = errors.New("no extractable text from PDF; OCR is required") +) + +type PDFTextParser struct{} + +func NewPDFTextParser() *PDFTextParser { + return &PDFTextParser{} +} + +func (p *PDFTextParser) ParseFile(ctx context.Context, path string) (Document, error) { + f, reader, err := pdf.Open(path) + if err != nil { + return Document{}, err + } + defer f.Close() + + pageCount := reader.NumPage() + pages := make([]Page, 0, pageCount) + + var full strings.Builder + + for i := 1; i <= pageCount; i++ { + if err := ctx.Err(); err != nil { + return Document{}, err + } + + page := reader.Page(i) + if page.V.IsNull() { + continue + } + + text, err := page.GetPlainText(nil) + if err != nil { + return Document{}, err + } + + text = strings.TrimSpace(text) + if text == "" { + continue + } + + pages = append(pages, Page{ + PageNo: i, + Text: text, + }) + + if full.Len() > 0 { + full.WriteString("\n\n") + } + full.WriteString(text) + } + + documentText := strings.TrimSpace(full.String()) + if documentText == "" { + return Document{}, ErrNoExtractableText + } + + return Document{ + PageCount: pageCount, + Pages: pages, + Text: documentText, + }, nil +}