fix(ui): 自定义工具表单 overflow-y-scroll→auto + 认证方式NONE翻译改为'无' + Captain文档爬虫/同步后端
This commit is contained in:
@@ -580,6 +580,9 @@ func Bootstrap(env string) (*App, error) {
|
||||
captainAssistantService := service.NewCaptainAssistantService(captainAssistantRepo, captainInboxRepo, captainDocumentRepo, captainAssistantResponseRepo, llmProvider, rdb)
|
||||
captainDocumentService := service.NewCaptainDocumentService(captainDocumentRepo, llmProvider, captainAssistantRepo)
|
||||
captainDocumentService.SetResponseRepo(captainAssistantResponseRepo)
|
||||
captainDocumentService.SetSyncBackend(service.NewCaptainDocumentSyncBackend())
|
||||
captainDocumentService.SetCrawlBackend(service.NewCaptainDocumentCrawlBackend())
|
||||
captainDocumentService.SetPageParserBackend(service.NewCaptainDocumentPageParserBackend())
|
||||
captainDocumentService.SetWorkerPool(workerPool)
|
||||
if _, err := service.EnqueueCaptainDocumentScheduleSyncs(context.Background(), workerPool, time.Now()); err != nil {
|
||||
applogger.L().Warnf("failed to enqueue initial Captain document sync scheduler: %v", err)
|
||||
|
||||
@@ -3,6 +3,7 @@ package v1
|
||||
import (
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
@@ -24,6 +25,7 @@ func NewCaptainDocumentHandler(svc *service.CaptainDocumentService) *CaptainDocu
|
||||
|
||||
// Create creates a new captain document.
|
||||
// POST /api/v1/accounts/:account_id/captain_assistants/:assistant_id/documents
|
||||
// Supports both JSON (for URL/content input) and multipart/form-data (for PDF upload).
|
||||
func (h *CaptainDocumentHandler) Create(c *gin.Context) {
|
||||
accountID := parseAccountIDParam(c)
|
||||
if accountID == 0 {
|
||||
@@ -32,10 +34,31 @@ func (h *CaptainDocumentHandler) Create(c *gin.Context) {
|
||||
}
|
||||
|
||||
var req service.CreateDocumentRequest
|
||||
if err := bindNestedJSONPayload(c, "document", &req); err != nil {
|
||||
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrValidation, err.Error())
|
||||
return
|
||||
|
||||
contentType := c.GetHeader("Content-Type")
|
||||
if strings.Contains(contentType, "multipart/form-data") {
|
||||
// Multipart form-data: PDF file upload or form fields
|
||||
req.Name = c.PostForm("document[name]")
|
||||
req.ExternalLink = c.PostForm("document[external_link]")
|
||||
req.Content = c.PostForm("document[content]")
|
||||
if assistantIDStr := c.PostForm("document[assistant_id]"); assistantIDStr != "" {
|
||||
if id, err := strconv.ParseUint(assistantIDStr, 10, 64); err == nil {
|
||||
req.AssistantID = uint(id)
|
||||
}
|
||||
}
|
||||
|
||||
// Check for PDF file
|
||||
if fileHeader, err := c.FormFile("document[pdf_file]"); err == nil {
|
||||
req.PdfFile = fileHeader
|
||||
}
|
||||
} else {
|
||||
// JSON request
|
||||
if err := bindNestedJSONPayload(c, "document", &req); err != nil {
|
||||
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrValidation, err.Error())
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
assistantID := req.AssistantID
|
||||
if assistantID == 0 {
|
||||
assistantID, _ = parseUintParam(c, "assistant_id")
|
||||
@@ -48,7 +71,7 @@ func (h *CaptainDocumentHandler) Create(c *gin.Context) {
|
||||
doc, err := h.svc.Create(c.Request.Context(), assistantID, accountID, &req)
|
||||
if err != nil {
|
||||
applogger.L().Errorf("Create captain document: %v", err)
|
||||
response.AbortWithStatusError(c, http.StatusUnprocessableEntity, response.ErrValidation, "failed to create document")
|
||||
response.AbortWithStatusError(c, http.StatusUnprocessableEntity, response.ErrValidation, err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
@@ -212,16 +235,18 @@ func captainDocumentPayload(doc *model.CaptainDocument) gin.H {
|
||||
if syncStatus == model.DocumentSyncStatusPending {
|
||||
syncStatus = "syncing"
|
||||
}
|
||||
pdfDocument := doc.ContentType == "application/pdf" || doc.FileURL != ""
|
||||
payload := gin.H{
|
||||
"account_id": doc.AccountID,
|
||||
"assistant": captainDocumentAssistantPayload(doc),
|
||||
"content": doc.Content,
|
||||
"content_type": "text/html",
|
||||
"content_type": doc.ContentType,
|
||||
"created_at": doc.CreatedAt.Unix(),
|
||||
"external_link": doc.ExternalLink,
|
||||
"display_url": doc.ExternalLink,
|
||||
"file_size": 0,
|
||||
"pdf_document": false,
|
||||
"file_size": doc.FileSize,
|
||||
"file_url": doc.FileURL,
|
||||
"pdf_document": pdfDocument,
|
||||
"id": doc.ID,
|
||||
"name": doc.Name,
|
||||
"status": status,
|
||||
|
||||
@@ -182,7 +182,7 @@ type CaptainDocument struct {
|
||||
AccountID uint `gorm:"index;not null" json:"account_id"`
|
||||
AssistantID uint `gorm:"index;not null" json:"assistant_id"`
|
||||
Name string `gorm:"size:255" json:"name"`
|
||||
ExternalLink string `gorm:"size:2048;not null" json:"external_link"`
|
||||
ExternalLink string `gorm:"size:2048" json:"external_link"`
|
||||
Content string `gorm:"type:text" json:"content,omitempty"`
|
||||
ContentFingerprint string `gorm:"size:64" json:"content_fingerprint,omitempty"`
|
||||
Status DocumentStatus `gorm:"size:50;default:in_progress;not null" json:"status"`
|
||||
@@ -192,6 +192,11 @@ type CaptainDocument struct {
|
||||
LastSyncErrorCode string `gorm:"size:50" json:"last_sync_error_code,omitempty"`
|
||||
Metadata json.RawMessage `gorm:"type:jsonb;serializer:json" json:"metadata,omitempty"`
|
||||
|
||||
// File upload metadata (PDF support)
|
||||
FileSize int64 `gorm:"default:0" json:"file_size,omitempty"`
|
||||
ContentType string `gorm:"size:100" json:"content_type,omitempty"`
|
||||
FileURL string `gorm:"size:2048" json:"file_url,omitempty"`
|
||||
|
||||
// Relationships
|
||||
Assistant CaptainAssistant `gorm:"foreignKey:AssistantID" json:"assistant,omitempty"`
|
||||
Responses []CaptainAssistantResponse `gorm:"foreignKey:DocumentableID" json:"responses,omitempty"`
|
||||
|
||||
@@ -1,5 +1,14 @@
|
||||
package pgvector
|
||||
|
||||
import (
|
||||
"database/sql/driver"
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"math"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// Vector represents a vector for similarity search.
|
||||
// Stub implementation for building without PostgreSQL/pgvector extension.
|
||||
type Vector []float32
|
||||
@@ -17,4 +26,89 @@ func (v Vector) String() string {
|
||||
// Dimensions returns the number of dimensions in the vector.
|
||||
func (v Vector) Dimensions() int {
|
||||
return len(v)
|
||||
}
|
||||
}
|
||||
|
||||
// Scan implements sql.Scanner so GORM/database/sql can scan vector columns
|
||||
// from PostgreSQL into the Vector type.
|
||||
//
|
||||
// PostgreSQL pgvector returns the vector as a string like "[0.1,0.2,0.3]"
|
||||
// via text protocol, or as raw bytes via binary protocol.
|
||||
// We handle both cases.
|
||||
func (v *Vector) Scan(src any) error {
|
||||
if src == nil {
|
||||
*v = nil
|
||||
return nil
|
||||
}
|
||||
|
||||
switch val := src.(type) {
|
||||
case string:
|
||||
return v.parseVectorString(val)
|
||||
case []byte:
|
||||
// Try text protocol first: "[0.1,0.2,...]"
|
||||
s := string(val)
|
||||
if strings.HasPrefix(s, "[") {
|
||||
return v.parseVectorString(s)
|
||||
}
|
||||
// Binary protocol: pgvector binary format
|
||||
// Header: 2 bytes unused (version), 2 bytes dim count, then dim*4 bytes float32 big-endian
|
||||
return v.parseVectorBinary(val)
|
||||
default:
|
||||
return fmt.Errorf("pgvector: cannot scan %T into Vector", src)
|
||||
}
|
||||
}
|
||||
|
||||
// Value implements driver.Valuer so GORM/database/sql can use Vector
|
||||
// in parameter binding. Returns the PostgreSQL text representation.
|
||||
func (v Vector) Value() (driver.Value, error) {
|
||||
if len(v) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
strs := make([]string, len(v))
|
||||
for i, f := range v {
|
||||
strs[i] = strconv.FormatFloat(float64(f), 'f', -1, 32)
|
||||
}
|
||||
return "[" + strings.Join(strs, ",") + "]", nil
|
||||
}
|
||||
|
||||
// parseVectorString parses a PostgreSQL vector text representation "[0.1,0.2,0.3]".
|
||||
func (v *Vector) parseVectorString(s string) error {
|
||||
s = strings.TrimSpace(s)
|
||||
s = strings.TrimPrefix(s, "[")
|
||||
s = strings.TrimSuffix(s, "]")
|
||||
if s == "" {
|
||||
*v = Vector{}
|
||||
return nil
|
||||
}
|
||||
parts := strings.Split(s, ",")
|
||||
result := make(Vector, len(parts))
|
||||
for i, part := range parts {
|
||||
f, err := strconv.ParseFloat(strings.TrimSpace(part), 32)
|
||||
if err != nil {
|
||||
return fmt.Errorf("pgvector: parse float %q: %w", part, err)
|
||||
}
|
||||
result[i] = float32(f)
|
||||
}
|
||||
*v = result
|
||||
return nil
|
||||
}
|
||||
|
||||
// parseVectorBinary parses pgvector binary format.
|
||||
// Format: 2 bytes (version, currently 0), 2 bytes (dims), then dims * 4 bytes float32.
|
||||
func (v *Vector) parseVectorBinary(data []byte) error {
|
||||
if len(data) < 4 {
|
||||
return fmt.Errorf("pgvector: binary vector too short (%d bytes)", len(data))
|
||||
}
|
||||
dims := int(binary.BigEndian.Uint16(data[2:4]))
|
||||
expected := 4 + dims*4
|
||||
if len(data) < expected {
|
||||
return fmt.Errorf("pgvector: binary vector expected %d bytes, got %d", expected, len(data))
|
||||
}
|
||||
result := make(Vector, dims)
|
||||
for i := 0; i < dims; i++ {
|
||||
offset := 4 + i*4
|
||||
bits := binary.BigEndian.Uint32(data[offset : offset+4])
|
||||
result[i] = math.Float32frombits(bits)
|
||||
}
|
||||
*v = result
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -84,9 +84,14 @@ func (r *CaptainDocumentRepo) ListByAccount(ctx context.Context, accountID uint,
|
||||
}
|
||||
switch filters.Source {
|
||||
case "web":
|
||||
db = db.Where("external_link <> ''")
|
||||
// Web documents: have external_link but not a PDF.
|
||||
db = db.Where("(content_type IS NULL OR content_type = '' OR content_type != 'application/pdf') AND (file_url = '' OR file_url IS NULL) AND external_link <> ''")
|
||||
case "pdf":
|
||||
db = db.Where("external_link = ''")
|
||||
// PDF documents: identified by content_type or file_url.
|
||||
db = db.Where("content_type = 'application/pdf' OR (file_url IS NOT NULL AND file_url <> '')")
|
||||
case "text":
|
||||
// Text documents: inline content, no external link and no file.
|
||||
db = db.Where("(external_link = '' OR external_link IS NULL) AND (file_url = '' OR file_url IS NULL) AND content <> ''")
|
||||
}
|
||||
switch filters.Filter {
|
||||
case "syncing":
|
||||
@@ -95,6 +100,11 @@ func (r *CaptainDocumentRepo) ListByAccount(ctx context.Context, accountID uint,
|
||||
db = db.Where("sync_status = ?", model.DocumentSyncStatusFailed)
|
||||
case "synced":
|
||||
db = db.Where("sync_status = ?", model.DocumentSyncStatusSynced)
|
||||
case "stale":
|
||||
// Documents that need updating: never synced or explicitly marked stale,
|
||||
// excluding those currently syncing.
|
||||
db = db.Where("(last_synced_at IS NULL OR sync_status = ?) AND sync_status != ?",
|
||||
model.DocumentSyncStatusStale, model.DocumentSyncStatusPending)
|
||||
}
|
||||
if filters.SearchKey != "" {
|
||||
query := "%" + strings.ToLower(filters.SearchKey) + "%"
|
||||
|
||||
@@ -0,0 +1,139 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"golang.org/x/net/html"
|
||||
)
|
||||
|
||||
// captainDocumentCrawlBackendImpl is the production implementation of
|
||||
// CaptainDocumentCrawlBackend. It fetches the document's external URL,
|
||||
// extracts page links for crawl fan-out.
|
||||
type captainDocumentCrawlBackendImpl struct {
|
||||
httpClient *http.Client
|
||||
}
|
||||
|
||||
// NewCaptainDocumentCrawlBackend creates the production crawl backend.
|
||||
func NewCaptainDocumentCrawlBackend() CaptainDocumentCrawlBackend {
|
||||
return &captainDocumentCrawlBackendImpl{
|
||||
httpClient: &http.Client{Timeout: 30 * time.Second},
|
||||
}
|
||||
}
|
||||
|
||||
func (b *captainDocumentCrawlBackendImpl) CrawlCaptainDocument(ctx context.Context, doc *model.CaptainDocument) (*CaptainDocumentCrawlResult, error) {
|
||||
url := doc.ExternalLink
|
||||
if url == "" {
|
||||
return &CaptainDocumentCrawlResult{ErrorCode: "not_found"}, nil
|
||||
}
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
||||
if err != nil {
|
||||
return &CaptainDocumentCrawlResult{ErrorCode: "fetch_failed"}, fmt.Errorf("create request: %w", err)
|
||||
}
|
||||
req.Header.Set("User-Agent", "GoChat-Captain/1.0")
|
||||
|
||||
resp, err := b.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return &CaptainDocumentCrawlResult{ErrorCode: "fetch_failed"}, fmt.Errorf("fetch url: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return &CaptainDocumentCrawlResult{ErrorCode: "fetch_failed"}, fmt.Errorf("HTTP %d for %s", resp.StatusCode, url)
|
||||
}
|
||||
|
||||
body, err := io.ReadAll(io.LimitReader(resp.Body, 5<<20)) // 5MB max
|
||||
if err != nil {
|
||||
return &CaptainDocumentCrawlResult{ErrorCode: "fetch_failed"}, fmt.Errorf("read body: %w", err)
|
||||
}
|
||||
|
||||
links := extractPageLinks(string(body), url)
|
||||
return &CaptainDocumentCrawlResult{
|
||||
PageLinks: links,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// extractPageLinks parses an HTML document and returns all unique absolute
|
||||
// href links found in <a> tags.
|
||||
func extractPageLinks(htmlStr, baseURL string) []string {
|
||||
doc, err := html.Parse(strings.NewReader(htmlStr))
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
base := baseURL
|
||||
// Remove trailing slash for consistent prefix matching
|
||||
base = strings.TrimSuffix(base, "/")
|
||||
|
||||
var links []string
|
||||
var walk func(*html.Node)
|
||||
walk = func(n *html.Node) {
|
||||
if n.Type == html.ElementNode && n.Data == "a" {
|
||||
for _, attr := range n.Attr {
|
||||
if attr.Key == "href" {
|
||||
href := strings.TrimSpace(attr.Val)
|
||||
if href == "" || strings.HasPrefix(href, "#") {
|
||||
continue
|
||||
}
|
||||
absolute := resolveURL(href, base)
|
||||
if absolute != "" {
|
||||
links = append(links, absolute)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
for c := n.FirstChild; c != nil; c = c.NextSibling {
|
||||
walk(c)
|
||||
}
|
||||
}
|
||||
walk(doc)
|
||||
|
||||
// Deduplicate
|
||||
seen := make(map[string]struct{}, len(links))
|
||||
unique := make([]string, 0, len(links))
|
||||
for _, l := range links {
|
||||
if _, ok := seen[l]; ok {
|
||||
continue
|
||||
}
|
||||
seen[l] = struct{}{}
|
||||
unique = append(unique, l)
|
||||
}
|
||||
return unique
|
||||
}
|
||||
|
||||
// resolveURL converts a relative URL to an absolute URL using the base URL.
|
||||
func resolveURL(href, base string) string {
|
||||
href = strings.TrimSpace(href)
|
||||
if href == "" {
|
||||
return ""
|
||||
}
|
||||
// Already absolute
|
||||
if strings.HasPrefix(href, "http://") || strings.HasPrefix(href, "https://") {
|
||||
return href
|
||||
}
|
||||
// Protocol-relative
|
||||
if strings.HasPrefix(href, "//") {
|
||||
return "https:" + href
|
||||
}
|
||||
// Absolute path
|
||||
if strings.HasPrefix(href, "/") {
|
||||
// Extract scheme://host from base
|
||||
idx := strings.Index(base, "://")
|
||||
if idx < 0 {
|
||||
return ""
|
||||
}
|
||||
hostEnd := strings.IndexByte(base[idx+3:], '/')
|
||||
if hostEnd < 0 {
|
||||
return base + href
|
||||
}
|
||||
return base[:idx+3+hostEnd] + href
|
||||
}
|
||||
// Relative path
|
||||
return base + "/" + href
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// captainDocumentPageParserBackendImpl is the production implementation of
|
||||
// CaptainDocumentPageParserBackend. It fetches a single web page and
|
||||
// extracts its title and visible text content — the same logic used by
|
||||
// the sync backend for individual page URLs.
|
||||
type captainDocumentPageParserBackendImpl struct {
|
||||
httpClient *http.Client
|
||||
}
|
||||
|
||||
// NewCaptainDocumentPageParserBackend creates the production page parser backend.
|
||||
func NewCaptainDocumentPageParserBackend() CaptainDocumentPageParserBackend {
|
||||
return &captainDocumentPageParserBackendImpl{
|
||||
httpClient: &http.Client{Timeout: 30 * time.Second},
|
||||
}
|
||||
}
|
||||
|
||||
func (b *captainDocumentPageParserBackendImpl) ParseCaptainDocumentPage(ctx context.Context, pageLink string) (*CaptainDocumentSyncResult, error) {
|
||||
if pageLink == "" {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "not_found"}, nil
|
||||
}
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, pageLink, nil)
|
||||
if err != nil {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "fetch_failed"}, fmt.Errorf("create request: %w", err)
|
||||
}
|
||||
req.Header.Set("User-Agent", "GoChat-Captain/1.0")
|
||||
|
||||
resp, err := b.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "fetch_failed"}, fmt.Errorf("fetch url: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "fetch_failed"}, fmt.Errorf("HTTP %d for %s", resp.StatusCode, pageLink)
|
||||
}
|
||||
|
||||
body, err := io.ReadAll(io.LimitReader(resp.Body, 5<<20)) // 5MB max
|
||||
if err != nil {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "fetch_failed"}, fmt.Errorf("read body: %w", err)
|
||||
}
|
||||
|
||||
title, content := extractHTMLText(string(body))
|
||||
content = strings.TrimSpace(content)
|
||||
if content == "" {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "content_empty"}, nil
|
||||
}
|
||||
|
||||
if title == "" {
|
||||
title = pageLink
|
||||
}
|
||||
return &CaptainDocumentSyncResult{
|
||||
Content: content,
|
||||
Title: title,
|
||||
}, nil
|
||||
}
|
||||
@@ -6,7 +6,10 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"mime/multipart"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -120,8 +123,11 @@ func (s *CaptainDocumentService) SetWorkerPool(wp *worker.WorkerPool) {
|
||||
// CreateDocumentRequest is the DTO for creating a document.
|
||||
type CreateDocumentRequest struct {
|
||||
Name string `json:"name" validate:"required"`
|
||||
ExternalLink string `json:"external_link" validate:"required"`
|
||||
ExternalLink string `json:"external_link"`
|
||||
Content string `json:"content"`
|
||||
AssistantID uint `json:"assistant_id"`
|
||||
// File upload fields (set by handler when multipart/form-data)
|
||||
PdfFile *multipart.FileHeader `json:"-"`
|
||||
}
|
||||
|
||||
// UpdateDocumentRequest is the DTO for updating a document.
|
||||
@@ -149,30 +155,146 @@ func (s *CaptainDocumentService) Create(ctx context.Context, assistantID, accoun
|
||||
return nil, fmt.Errorf("assistant not found: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Validate: at least one source must be provided
|
||||
if req.ExternalLink == "" && req.PdfFile == nil && strings.TrimSpace(req.Content) == "" {
|
||||
return nil, errors.New("at least one of external_link, content, or pdf_file is required")
|
||||
}
|
||||
|
||||
doc := &model.CaptainDocument{
|
||||
AccountID: accountID,
|
||||
AssistantID: assistantID,
|
||||
Name: req.Name,
|
||||
ExternalLink: req.ExternalLink,
|
||||
Content: strings.TrimSpace(req.Content),
|
||||
Status: model.DocumentStatusPending,
|
||||
}
|
||||
|
||||
// Handle PDF file upload
|
||||
if req.PdfFile != nil {
|
||||
fileURL, err := s.saveUploadedFile(ctx, accountID, req.PdfFile)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("save pdf file: %w", err)
|
||||
}
|
||||
doc.FileURL = fileURL
|
||||
doc.FileSize = req.PdfFile.Size
|
||||
doc.ContentType = req.PdfFile.Header.Get("Content-Type")
|
||||
if doc.ContentType == "" {
|
||||
doc.ContentType = "application/pdf"
|
||||
}
|
||||
// If no external link was provided, use the uploaded file URL
|
||||
if doc.ExternalLink == "" {
|
||||
doc.ExternalLink = fileURL
|
||||
}
|
||||
}
|
||||
|
||||
if err := s.documentRepo.Create(ctx, doc); err != nil {
|
||||
applogger.L().Errorf("Create captain document: %v", err)
|
||||
return nil, fmt.Errorf("create document: %w", err)
|
||||
}
|
||||
if created, err := s.documentRepo.GetByAccountAndID(ctx, accountID, doc.ID); err == nil {
|
||||
if enqueueErr := s.enqueueDocumentCrawl(ctx, accountID, created.ID); enqueueErr != nil {
|
||||
|
||||
// Route the document to the correct processing pipeline:
|
||||
// - PDF upload → sync job (PDF text extraction is the sync backend's job)
|
||||
// - URL → crawl job (fetch page, discover links, parse)
|
||||
// - Text content → sync job (process content directly, no fetch needed)
|
||||
docStatus := doc.Status
|
||||
if doc.Content != "" && doc.ExternalLink == "" {
|
||||
// Direct content input: mark as in_progress so the sync pipeline
|
||||
// can pick it up and generate FAQ responses.
|
||||
docStatus = model.DocumentStatusInProgress
|
||||
doc.SyncStatus = model.DocumentSyncStatusPending
|
||||
if err := s.documentRepo.Update(ctx, doc); err != nil {
|
||||
applogger.L().Errorf("Update captain document status after create: %v", err)
|
||||
return nil, fmt.Errorf("update document status: %w", err)
|
||||
}
|
||||
if s.worker != nil {
|
||||
if _, err := s.worker.Enqueue(ctx, TaskTypeCaptainDocumentSync,
|
||||
captainDocumentSyncJob{AccountID: accountID, DocumentID: doc.ID},
|
||||
worker.WithQueue("low"), worker.WithMaxAttempts(3)); err != nil {
|
||||
applogger.L().Warnf("enqueue document sync for content: %v", err)
|
||||
}
|
||||
}
|
||||
} else if doc.FileURL != "" {
|
||||
// PDF upload: enqueue sync job to extract text from the PDF file.
|
||||
if s.worker != nil {
|
||||
if _, err := s.worker.Enqueue(ctx, TaskTypeCaptainDocumentSync,
|
||||
captainDocumentSyncJob{AccountID: accountID, DocumentID: doc.ID},
|
||||
worker.WithQueue("low"), worker.WithMaxAttempts(3)); err != nil {
|
||||
applogger.L().Warnf("enqueue document sync for pdf: %v", err)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// URL: enqueue crawl job to fetch/extract content
|
||||
if created, err := s.documentRepo.GetByAccountAndID(ctx, accountID, doc.ID); err == nil {
|
||||
if enqueueErr := s.enqueueDocumentCrawl(ctx, accountID, created.ID); enqueueErr != nil {
|
||||
return nil, enqueueErr
|
||||
}
|
||||
return created, nil
|
||||
}
|
||||
if enqueueErr := s.enqueueDocumentCrawl(ctx, accountID, doc.ID); enqueueErr != nil {
|
||||
return nil, enqueueErr
|
||||
}
|
||||
return created, nil
|
||||
}
|
||||
if enqueueErr := s.enqueueDocumentCrawl(ctx, accountID, doc.ID); enqueueErr != nil {
|
||||
return nil, enqueueErr
|
||||
}
|
||||
_ = docStatus
|
||||
return doc, nil
|
||||
}
|
||||
|
||||
// saveUploadedFile saves an uploaded PDF file to local storage and returns the file URL.
|
||||
func (s *CaptainDocumentService) saveUploadedFile(ctx context.Context, accountID uint, fileHeader *multipart.FileHeader) (string, error) {
|
||||
// Read the file content
|
||||
src, err := fileHeader.Open()
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("open uploaded file: %w", err)
|
||||
}
|
||||
defer src.Close()
|
||||
|
||||
// Generate a unique filename
|
||||
ext := filepath.Ext(fileHeader.Filename)
|
||||
if ext == "" {
|
||||
ext = ".pdf"
|
||||
}
|
||||
timestamp := time.Now().UnixNano()
|
||||
filename := fmt.Sprintf("captain_docs/%d/%d%s", accountID, timestamp, ext)
|
||||
destPath := filepath.Join(s.uploadDir(), filename)
|
||||
|
||||
// Create directory if needed
|
||||
if err := os.MkdirAll(filepath.Dir(destPath), 0o755); err != nil {
|
||||
return "", fmt.Errorf("create upload directory: %w", err)
|
||||
}
|
||||
|
||||
// Write the file
|
||||
dst, err := os.Create(destPath)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("create destination file: %w", err)
|
||||
}
|
||||
defer dst.Close()
|
||||
|
||||
if _, err := io.Copy(dst, src); err != nil {
|
||||
return "", fmt.Errorf("write uploaded file: %w", err)
|
||||
}
|
||||
|
||||
// Return the relative URL path (served by StaticFS at /uploads)
|
||||
return fmt.Sprintf("/uploads/captain_docs/%d/%d%s", accountID, timestamp, ext), nil
|
||||
}
|
||||
|
||||
// uploadDir returns the base directory for file uploads.
|
||||
// This should match the configured storage.local_path (default: ./uploads).
|
||||
func (s *CaptainDocumentService) uploadDir() string {
|
||||
return "./uploads"
|
||||
}
|
||||
|
||||
// enqueueDocumentProcess enqueues a job to process document content directly (for inline content input).
|
||||
func (s *CaptainDocumentService) enqueueDocumentProcess(ctx context.Context, accountID, docID uint) error {
|
||||
if s.worker == nil {
|
||||
return nil
|
||||
}
|
||||
_, err := s.worker.Enqueue(ctx, TaskTypeCaptainDocumentResponseBuilder, captainDocumentResponseBuilderJob{AccountID: accountID, DocumentID: docID},
|
||||
worker.WithQueue("low"),
|
||||
worker.WithMaxAttempts(3),
|
||||
)
|
||||
return err
|
||||
}
|
||||
|
||||
// Get retrieves a document by ID.
|
||||
func (s *CaptainDocumentService) Get(ctx context.Context, id uint) (*model.CaptainDocument, error) {
|
||||
doc, err := s.documentRepo.GetByID(ctx, id)
|
||||
@@ -468,6 +590,28 @@ func (s *CaptainDocumentService) SyncDocumentByAccount(ctx context.Context, acco
|
||||
if err := s.markDocumentSyncStarted(ctx, doc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Direct content input: content is already present, skip fetch/extraction.
|
||||
// Go straight to FAQ response building.
|
||||
if strings.TrimSpace(doc.Content) != "" && doc.FileURL == "" && doc.ExternalLink == "" {
|
||||
doc.Content = strings.TrimSpace(doc.Content)
|
||||
doc.ContentFingerprint = computeFingerprint(doc.Content)
|
||||
doc.Status = model.DocumentStatusCompleted
|
||||
doc.SyncStatus = model.DocumentSyncStatusSynced
|
||||
doc.LastSyncErrorCode = ""
|
||||
now := time.Now().Unix()
|
||||
doc.LastSyncedAt = &now
|
||||
doc.LastSyncAttemptedAt = &now
|
||||
if err := s.documentRepo.Update(ctx, doc); err != nil {
|
||||
return nil, fmt.Errorf("update synced document: %w", err)
|
||||
}
|
||||
if err := s.enqueueDocumentResponseBuilder(ctx, accountID, id); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return s.documentRepo.GetByAccountAndID(ctx, accountID, id)
|
||||
}
|
||||
|
||||
// PDF and URL documents need a sync backend to fetch/extract content.
|
||||
if s.syncBackend == nil {
|
||||
return s.markDocumentSyncFailed(ctx, accountID, id, "sync_disabled")
|
||||
}
|
||||
|
||||
@@ -0,0 +1,175 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"golang.org/x/net/html"
|
||||
)
|
||||
|
||||
// captainDocumentSyncBackendImpl is the production implementation of
|
||||
// CaptainDocumentSyncBackend. It handles two document types:
|
||||
// - PDF uploads: extracts text via the `pdftotext` CLI (poppler-utils),
|
||||
// which correctly handles CJK/CID fonts that pure-Go PDF libraries cannot.
|
||||
// - Web URLs: fetches the page and extracts visible text from HTML.
|
||||
type captainDocumentSyncBackendImpl struct {
|
||||
httpClient *http.Client
|
||||
}
|
||||
|
||||
// NewCaptainDocumentSyncBackend creates the production sync backend.
|
||||
func NewCaptainDocumentSyncBackend() CaptainDocumentSyncBackend {
|
||||
return &captainDocumentSyncBackendImpl{
|
||||
httpClient: &http.Client{Timeout: 30 * time.Second},
|
||||
}
|
||||
}
|
||||
|
||||
func (b *captainDocumentSyncBackendImpl) SyncCaptainDocument(ctx context.Context, doc *model.CaptainDocument) (*CaptainDocumentSyncResult, error) {
|
||||
// PDF document: extract text from the uploaded file.
|
||||
if doc.ContentType == "application/pdf" || doc.FileURL != "" {
|
||||
return b.syncPDFDocument(ctx, doc)
|
||||
}
|
||||
// Web URL document: fetch and extract text from the page.
|
||||
if doc.ExternalLink != "" {
|
||||
return b.syncWebDocument(ctx, doc)
|
||||
}
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "content_empty"}, nil
|
||||
}
|
||||
|
||||
// syncPDFDocument extracts text content from a locally stored PDF file
|
||||
// using the `pdftotext` CLI (poppler-utils). This handles CJK fonts and
|
||||
// complex PDF encodings that pure-Go libraries like ledongthuc/pdf cannot.
|
||||
func (b *captainDocumentSyncBackendImpl) syncPDFDocument(ctx context.Context, doc *model.CaptainDocument) (*CaptainDocumentSyncResult, error) {
|
||||
filePath := b.resolvePDFPath(doc)
|
||||
if filePath == "" {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "not_found"}, nil
|
||||
}
|
||||
|
||||
// Use pdftotext CLI for robust text extraction (supports CJK).
|
||||
cmd := exec.CommandContext(ctx, "pdftotext", "-enc", "UTF-8", filePath, "-")
|
||||
var stdout, stderr bytes.Buffer
|
||||
cmd.Stdout = &stdout
|
||||
cmd.Stderr = &stderr
|
||||
if err := cmd.Run(); err != nil {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "fetch_failed"}, fmt.Errorf("pdftotext: %w, stderr: %s", err, stderr.String())
|
||||
}
|
||||
|
||||
content := strings.TrimSpace(stdout.String())
|
||||
if content == "" {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "content_empty"}, nil
|
||||
}
|
||||
|
||||
title := doc.Name
|
||||
return &CaptainDocumentSyncResult{
|
||||
Content: content,
|
||||
Title: title,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// resolvePDFPath converts the doc's file_url to a local filesystem path.
|
||||
// file_url is stored as "/uploads/captain_docs/<account>/<timestamp>.pdf"
|
||||
// and served from the local "./uploads" directory.
|
||||
func (b *captainDocumentSyncBackendImpl) resolvePDFPath(doc *model.CaptainDocument) string {
|
||||
if doc.FileURL == "" {
|
||||
return ""
|
||||
}
|
||||
// file_url is a relative URL path like "/uploads/captain_docs/1/123.pdf"
|
||||
// Map it to the local filesystem path.
|
||||
path := doc.FileURL
|
||||
if strings.HasPrefix(path, "/uploads/") {
|
||||
return "." + path
|
||||
}
|
||||
// If it's already a filesystem path, use it directly.
|
||||
if _, err := os.Stat(path); err == nil {
|
||||
return path
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// syncWebDocument fetches a web page and extracts its visible text content.
|
||||
func (b *captainDocumentSyncBackendImpl) syncWebDocument(ctx context.Context, doc *model.CaptainDocument) (*CaptainDocumentSyncResult, error) {
|
||||
url := doc.ExternalLink
|
||||
if url == "" {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "not_found"}, nil
|
||||
}
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
||||
if err != nil {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "fetch_failed"}, fmt.Errorf("create request: %w", err)
|
||||
}
|
||||
req.Header.Set("User-Agent", "GoChat-Captain/1.0")
|
||||
|
||||
resp, err := b.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "fetch_failed"}, fmt.Errorf("fetch url: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "fetch_failed"}, fmt.Errorf("HTTP %d for %s", resp.StatusCode, url)
|
||||
}
|
||||
|
||||
body, err := io.ReadAll(io.LimitReader(resp.Body, 5<<20)) // 5MB max
|
||||
if err != nil {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "fetch_failed"}, fmt.Errorf("read body: %w", err)
|
||||
}
|
||||
|
||||
title, content := extractHTMLText(string(body))
|
||||
content = strings.TrimSpace(content)
|
||||
if content == "" {
|
||||
return &CaptainDocumentSyncResult{ErrorCode: "content_empty"}, nil
|
||||
}
|
||||
|
||||
if title == "" {
|
||||
title = doc.Name
|
||||
}
|
||||
return &CaptainDocumentSyncResult{
|
||||
Content: content,
|
||||
Title: title,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// extractHTMLText parses an HTML document and returns the page title and
|
||||
// visible text content (stripping scripts, styles, and HTML tags).
|
||||
func extractHTMLText(htmlStr string) (title string, content string) {
|
||||
doc, err := html.Parse(strings.NewReader(htmlStr))
|
||||
if err != nil {
|
||||
// Fall back to raw text if HTML parsing fails.
|
||||
return "", strings.TrimSpace(htmlStr)
|
||||
}
|
||||
|
||||
var extract func(*html.Node)
|
||||
extract = func(n *html.Node) {
|
||||
if n.Type == html.ElementNode {
|
||||
switch n.Data {
|
||||
case "script", "style", "noscript", "head":
|
||||
return
|
||||
case "title":
|
||||
if n.FirstChild != nil {
|
||||
title = strings.TrimSpace(n.FirstChild.Data)
|
||||
}
|
||||
return
|
||||
}
|
||||
}
|
||||
if n.Type == html.TextNode {
|
||||
text := strings.TrimSpace(n.Data)
|
||||
if text != "" {
|
||||
content += text + " "
|
||||
}
|
||||
}
|
||||
for c := n.FirstChild; c != nil; c = c.NextSibling {
|
||||
extract(c)
|
||||
}
|
||||
}
|
||||
extract(doc)
|
||||
|
||||
content = strings.TrimSpace(content)
|
||||
return title, content
|
||||
}
|
||||
Reference in New Issue
Block a user