mirror of
https://github.com/Tencent/WeKnora.git
synced 2026-09-19 02:18:25 +08:00
fix(multimodal): unblock processing when provider:// image read fails
When image bytes for a multimodal task cannot be read via FileService (e.g. tenant.StorageEngineConfig.MinIO is empty while the image was saved using the global MINIO_* env vars), the previous code fell back to the HTTP downloader, which rejected the provider:// URL with "unsupported URL scheme". The asynq handler then returned an error, asynq retried until exhaustion, and the per-knowledge "pending images" counter was never decremented — leaving the document stuck in "processing" indefinitely (issue #1282). Changes: - ImageMultimodalService now holds a default FileService and resolveFileServiceForPayload falls back to it when the tenant-scoped storage config cannot produce a usable service, mirroring the write-side fallback in knowledgeService.resolveFileService. - Extract readImageBytes: provider:// URLs are read exclusively via FileService and never handed to the HTTP downloader. - On unrecoverable read failure for a single image, log and skip that image but still call checkAndFinalizeAllImages so the parent knowledge can progress to post-processing. Fixes #1282
This commit is contained in:
@@ -63,6 +63,12 @@ type ImageMultimodalService struct {
|
||||
ollamaService *ollama.OllamaService
|
||||
taskEnqueuer interfaces.TaskEnqueuer
|
||||
redisClient *redis.Client
|
||||
// fileSvc is the globally configured default FileService used as a fallback
|
||||
// when the tenant-scoped storage config cannot produce a usable service
|
||||
// (e.g. images were saved using the global MINIO_* env vars while the
|
||||
// tenant's StorageEngineConfig.MinIO is empty). Mirrors the write-side
|
||||
// fallback in knowledgeService.resolveFileService.
|
||||
fileSvc interfaces.FileService
|
||||
}
|
||||
|
||||
func NewImageMultimodalService(
|
||||
@@ -75,6 +81,7 @@ func NewImageMultimodalService(
|
||||
ollamaService *ollama.OllamaService,
|
||||
taskEnqueuer interfaces.TaskEnqueuer,
|
||||
redisClient *redis.Client,
|
||||
fileSvc interfaces.FileService,
|
||||
) interfaces.TaskHandler {
|
||||
return &ImageMultimodalService{
|
||||
chunkService: chunkService,
|
||||
@@ -86,6 +93,7 @@ func NewImageMultimodalService(
|
||||
ollamaService: ollamaService,
|
||||
taskEnqueuer: taskEnqueuer,
|
||||
redisClient: redisClient,
|
||||
fileSvc: fileSvc,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -109,43 +117,16 @@ func (s *ImageMultimodalService) Handle(ctx context.Context, task *asynq.Task) e
|
||||
return fmt.Errorf("resolve VLM: %w", err)
|
||||
}
|
||||
|
||||
// Read image bytes: try provider:// via tenant-resolved FileService,
|
||||
// then legacy local path, then HTTP URL.
|
||||
var imgBytes []byte
|
||||
if types.ParseProviderScheme(payload.ImageURL) != "" {
|
||||
fileSvc := s.resolveFileServiceForPayload(ctx, payload)
|
||||
if fileSvc == nil {
|
||||
logger.Warnf(ctx, "[ImageMultimodal] Resolve tenant file service failed, fallback to URL/local: tenant=%d kb=%s",
|
||||
payload.TenantID, payload.KnowledgeBaseID)
|
||||
} else {
|
||||
// provider:// scheme — read via FileService
|
||||
reader, getErr := fileSvc.GetFile(ctx, payload.ImageURL)
|
||||
if getErr != nil {
|
||||
logger.Warnf(ctx, "[ImageMultimodal] FileService.GetFile(%s) failed: %v", payload.ImageURL, getErr)
|
||||
} else {
|
||||
imgBytes, err = io.ReadAll(reader)
|
||||
reader.Close()
|
||||
if err != nil {
|
||||
logger.Warnf(ctx, "[ImageMultimodal] Read provider file %s failed: %v", payload.ImageURL, err)
|
||||
imgBytes = nil
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if imgBytes == nil && payload.ImageLocalPath != "" {
|
||||
imgBytes, err = os.ReadFile(payload.ImageLocalPath)
|
||||
if err != nil {
|
||||
logger.Warnf(ctx, "[ImageMultimodal] Local file %s not available (%v), trying URL", payload.ImageLocalPath, err)
|
||||
imgBytes = nil
|
||||
}
|
||||
}
|
||||
if imgBytes == nil {
|
||||
imgBytes, err = downloadImageFromURL(payload.ImageURL)
|
||||
if err != nil {
|
||||
logger.Errorf(ctx, "[ImageMultimodal] Failed to download image from URL %s: %v", payload.ImageURL, err)
|
||||
return fmt.Errorf("read image from URL %s failed: %w", payload.ImageURL, err)
|
||||
}
|
||||
logger.Infof(ctx, "[ImageMultimodal] Image downloaded from URL, len=%d", len(imgBytes))
|
||||
// Read image bytes. A provider:// URL must be resolved via FileService —
|
||||
// it must NEVER be handed to the HTTP downloader (which would fail with
|
||||
// "unsupported URL scheme"). On unrecoverable read failure for a single
|
||||
// image, skip it and still trigger finalize so the parent knowledge
|
||||
// doesn't get stuck in "processing" forever (see issue #1282).
|
||||
imgBytes, readErr := s.readImageBytes(ctx, payload)
|
||||
if readErr != nil {
|
||||
logger.Errorf(ctx, "[ImageMultimodal] Skip unreadable image %s: %v", payload.ImageURL, readErr)
|
||||
s.checkAndFinalizeAllImages(ctx, payload)
|
||||
return nil
|
||||
}
|
||||
|
||||
imageInfo := types.ImageInfo{
|
||||
@@ -356,11 +337,16 @@ func (s *ImageMultimodalService) resolveVLM(ctx context.Context, kbID string) (v
|
||||
}
|
||||
|
||||
// resolveFileServiceForPayload resolves tenant/KB scoped file service for reading provider:// URLs.
|
||||
// Falls back to the globally configured default FileService when the tenant's
|
||||
// StorageEngineConfig does not carry a usable configuration for the URL's provider.
|
||||
// This mirrors the write-side fallback in knowledgeService.resolveFileService
|
||||
// and is required because images can be saved using global STORAGE_TYPE/MINIO_*
|
||||
// env vars while tenant.StorageEngineConfig.MinIO is left empty (issue #1282).
|
||||
func (s *ImageMultimodalService) resolveFileServiceForPayload(ctx context.Context, payload types.ImageMultimodalPayload) interfaces.FileService {
|
||||
tenant, err := s.tenantRepo.GetTenantByID(ctx, payload.TenantID)
|
||||
if err != nil || tenant == nil {
|
||||
logger.Warnf(ctx, "[ImageMultimodal] GetTenantByID failed: tenant=%d err=%v", payload.TenantID, err)
|
||||
return nil
|
||||
return s.fileSvc
|
||||
}
|
||||
|
||||
provider := types.ParseProviderScheme(payload.ImageURL)
|
||||
@@ -376,12 +362,54 @@ func (s *ImageMultimodalService) resolveFileServiceForPayload(ctx context.Contex
|
||||
baseDir := strings.TrimSpace(os.Getenv("LOCAL_STORAGE_BASE_DIR"))
|
||||
fileSvc, _, svcErr := filesvc.NewFileServiceFromStorageConfig(provider, tenant.StorageEngineConfig, baseDir)
|
||||
if svcErr != nil {
|
||||
logger.Warnf(ctx, "[ImageMultimodal] resolve file service failed: tenant=%d provider=%s err=%v", payload.TenantID, provider, svcErr)
|
||||
return nil
|
||||
logger.Warnf(ctx, "[ImageMultimodal] resolve file service failed (falling back to default): tenant=%d provider=%s err=%v",
|
||||
payload.TenantID, provider, svcErr)
|
||||
return s.fileSvc
|
||||
}
|
||||
return fileSvc
|
||||
}
|
||||
|
||||
// readImageBytes loads the image bytes for a multimodal payload.
|
||||
// - For provider:// URLs (local://, minio://, s3://, cos://, ...) it reads via
|
||||
// the resolved FileService and NEVER falls back to HTTP — handing a
|
||||
// provider:// URL to the HTTP downloader is what caused issue #1282.
|
||||
// - For legacy in-flight payloads with ImageLocalPath set, it tries the local
|
||||
// file before falling back to the URL.
|
||||
// - For plain http(s):// URLs it uses the SSRF-safe downloader.
|
||||
func (s *ImageMultimodalService) readImageBytes(ctx context.Context, payload types.ImageMultimodalPayload) ([]byte, error) {
|
||||
if types.ParseProviderScheme(payload.ImageURL) != "" {
|
||||
fileSvc := s.resolveFileServiceForPayload(ctx, payload)
|
||||
if fileSvc == nil {
|
||||
return nil, fmt.Errorf("no file service available for %s", payload.ImageURL)
|
||||
}
|
||||
reader, err := fileSvc.GetFile(ctx, payload.ImageURL)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("file service get %s: %w", payload.ImageURL, err)
|
||||
}
|
||||
defer reader.Close()
|
||||
data, err := io.ReadAll(reader)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read %s: %w", payload.ImageURL, err)
|
||||
}
|
||||
return data, nil
|
||||
}
|
||||
|
||||
if payload.ImageLocalPath != "" {
|
||||
if data, err := os.ReadFile(payload.ImageLocalPath); err == nil {
|
||||
return data, nil
|
||||
} else {
|
||||
logger.Warnf(ctx, "[ImageMultimodal] Local file %s not available (%v), falling back to URL", payload.ImageLocalPath, err)
|
||||
}
|
||||
}
|
||||
|
||||
data, err := downloadImageFromURL(payload.ImageURL)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("download %s: %w", payload.ImageURL, err)
|
||||
}
|
||||
logger.Infof(ctx, "[ImageMultimodal] Image downloaded from URL, len=%d", len(data))
|
||||
return data, nil
|
||||
}
|
||||
|
||||
// downloadImageFromURL downloads image bytes from an HTTP(S) URL.
|
||||
func downloadImageFromURL(imageURL string) ([]byte, error) {
|
||||
return secutils.DownloadBytes(imageURL)
|
||||
|
||||
Reference in New Issue
Block a user