Files
2026-08-22 21:19:43 +08:00

1203 lines
40 KiB
Go

package service
import (
"bytes"
"context"
"crypto/rand"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"mime"
"mime/multipart"
"net/http"
"net/url"
"os"
"path/filepath"
"strconv"
"strings"
"time"
"github.com/google/uuid"
"gorm.io/gorm"
"github.com/gochat/gochat/internal/config"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/repository"
"github.com/gochat/gochat/internal/security"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
)
const (
TaskTypeCleanupExpiredUploads = "storage:cleanup_expired_uploads"
uploadCleanupInterval = 24 * time.Hour
)
// UploadService handles file uploads for both account-level and widget direct uploads.
type UploadService struct {
directUploadRepo *repository.DirectUploadRepo
conversationRepo *repository.ConversationRepo
inboxRepo *repository.InboxRepo
contactInboxRepo *repository.ContactInboxRepo
accessDB *gorm.DB
fetchClient *security.SafeHTTPClient
cfg *config.Config
}
// NewUploadService creates a new UploadService.
func NewUploadService(directUploadRepo *repository.DirectUploadRepo, cfg *config.Config) *UploadService {
return &UploadService{
directUploadRepo: directUploadRepo,
fetchClient: security.NewSafeHTTPClient(security.DefaultSSRFConfig()),
cfg: cfg,
}
}
// WithWidgetAuth wires the widget session repositories used by Chatwoot's
// website_token + X-Auth-Token direct upload guard.
func (s *UploadService) WithWidgetAuth(inboxRepo *repository.InboxRepo, contactInboxRepo *repository.ContactInboxRepo) *UploadService {
s.inboxRepo = inboxRepo
s.contactInboxRepo = contactInboxRepo
return s
}
// WithConversationRepo wires the account-scoped conversation lookup needed by
// Chatwoot's nested conversation direct upload endpoint.
func (s *UploadService) WithConversationRepo(conversationRepo *repository.ConversationRepo) *UploadService {
s.conversationRepo = conversationRepo
return s
}
// WithAccessDB wires the scoped lookup used by the private local-upload handler.
func (s *UploadService) WithAccessDB(db *gorm.DB) *UploadService {
s.accessDB = db
return s
}
// --- DTOs ---
// AccountUploadRequest is the DTO for account-level file upload.
type AccountUploadRequest struct {
FileHeader *multipart.FileHeader `json:"-"`
}
// WidgetDirectUploadRequest is the DTO for widget direct file upload.
type WidgetDirectUploadRequest struct {
WebsiteToken string `json:"-"`
AuthToken string `json:"-"`
FileHeader *multipart.FileHeader `json:"-"`
}
// AccountDirectUploadRequest is the DTO for account-level direct file upload (staged for message attachment).
// Reference: Chatwoot POST /api/v1/accounts/:account_id/direct_uploads
type AccountDirectUploadRequest struct {
FileHeader *multipart.FileHeader `json:"-"`
}
type ActiveStorageDirectUploadRequest struct {
WebsiteToken string `json:"-"`
AuthToken string `json:"-"`
Blob ActiveStorageBlobParams `json:"blob"`
}
type ActiveStorageBlobParams struct {
Filename string `json:"filename"`
ByteSize int64 `json:"byte_size"`
Checksum string `json:"checksum"`
ContentType string `json:"content_type"`
Metadata map[string]any `json:"metadata"`
}
type ActiveStorageDirectUploadResponse struct {
ID uint `json:"id"`
Key string `json:"key"`
Filename string `json:"filename"`
ContentType string `json:"content_type"`
Metadata map[string]any `json:"metadata"`
ServiceName string `json:"service_name"`
ByteSize int64 `json:"byte_size"`
Checksum string `json:"checksum"`
CreatedAt time.Time `json:"created_at"`
SignedID string `json:"signed_id"`
DirectUpload ActiveStorageUploadURL `json:"direct_upload"`
}
type ActiveStorageUploadURL struct {
URL string `json:"url"`
Headers map[string]string `json:"headers"`
}
// UploadResponse is the unified response DTO for upload endpoints.
type UploadResponse struct {
UploadID uint `json:"upload_id"`
UploadUUID string `json:"upload_uuid"`
OriginalName string `json:"original_name"`
FileType string `json:"file_type"`
MimeType string `json:"mime_type"`
FileSize int64 `json:"file_size"`
FileURL string `json:"file_url"`
ThumbURL string `json:"thumb_url,omitempty"`
Status string `json:"status"`
ExpiresAt time.Time `json:"expires_at"`
}
type widgetUploadSession struct {
AccountID uint
ContactInboxID uint
}
// MaxRequestBodySize is the single hard ceiling for upload request bodies.
// Multipart metadata gets a small allowance while file-specific limits remain
// enforced from the actual content in the service.
func (s *UploadService) MaxRequestBodySize() int64 {
return s.maxUploadSize() + (1 << 20)
}
// --- Account Upload ---
// AccountUpload handles a file upload from the dashboard (account-scoped).
// Reference: Chatwoot api/v1/accounts/:account_id/upload
func (s *UploadService) AccountUpload(ctx context.Context, accountID uint, req AccountUploadRequest) (*UploadResponse, error) {
if accountID == 0 {
return nil, errors.New("account_id is required")
}
if req.FileHeader == nil {
return nil, errors.New("file is required")
}
return s.processUpload(ctx, accountID, req.FileHeader, model.DirectUploadSourceAccount)
}
// ProfileAvatarUpload stores a dashboard user's avatar as a durable account
// file. Profile avatars intentionally bypass the direct_uploads staging table:
// they are referenced immediately from users.avatar_url and do not expire.
func (s *UploadService) ProfileAvatarUpload(ctx context.Context, accountID uint, fileHeader *multipart.FileHeader) (*UploadResponse, error) {
if accountID == 0 {
return nil, errors.New("account_id is required")
}
if fileHeader == nil {
return nil, errors.New("avatar file is required")
}
maxSize := int64(s.cfg.Storage.MaxFileSize)
if maxSize <= 0 {
maxSize = 20 << 20
}
if fileHeader.Size > maxSize {
return nil, fmt.Errorf("avatar file size %d exceeds maximum %d", fileHeader.Size, maxSize)
}
src, err := fileHeader.Open()
if err != nil {
return nil, fmt.Errorf("failed to open avatar file: %w", err)
}
defer src.Close()
data, mimeType, err := readValidatedUpload(src, fileHeader.Filename, fileHeader.Header.Get("Content-Type"), fileHeader.Size, maxSize, true)
if err != nil {
return nil, err
}
if !strings.HasPrefix(mimeType, "image/") || !isUploadMIMEAllowed("image", mimeType) {
return nil, fmt.Errorf("avatar must be an image with a supported format, got %s", mimeType)
}
fileURL, thumbURL, err := s.saveUploadReader(accountID, model.DirectUploadSourceAccount, extensionForUploadMIME(mimeType), mimeType, bytes.NewReader(data))
if err != nil {
return nil, fmt.Errorf("failed to save avatar file: %w", err)
}
return &UploadResponse{
OriginalName: fileHeader.Filename,
FileType: "image",
MimeType: mimeType,
FileSize: int64(len(data)),
FileURL: fileURL,
ThumbURL: thumbURL,
Status: string(model.DirectUploadStatusCompleted),
}, nil
}
func (s *UploadService) AccountUploadFromURL(ctx context.Context, accountID uint, externalURL string) (*UploadResponse, error) {
if accountID == 0 {
return nil, errors.New("account_id is required")
}
parsed, err := url.Parse(externalURL)
ssrfCfg := security.DefaultSSRFConfig()
if err != nil || parsed.Hostname() == "" || (parsed.Scheme != "http" && parsed.Scheme != "https") {
return nil, errors.New("invalid url")
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, externalURL, nil)
if err != nil {
return nil, err
}
client := s.fetchClient
if client == nil {
client = security.NewSafeHTTPClient(ssrfCfg)
}
resp, err := client.Do(req)
if err != nil {
return nil, fmt.Errorf("failed to fetch external url: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode < http.StatusOK || resp.StatusCode >= http.StatusMultipleChoices {
return nil, fmt.Errorf("failed to fetch external url: status %d", resp.StatusCode)
}
maxSize := s.maxUploadSize()
if resp.ContentLength > maxSize {
return nil, errors.New("file too large")
}
data, err := io.ReadAll(io.LimitReader(resp.Body, maxSize+1))
if err != nil {
return nil, err
}
if int64(len(data)) > maxSize {
return nil, errors.New("file too large")
}
filename := filepath.Base(parsed.Path)
if filename == "." || filename == "/" || filename == "" {
filename = "upload"
}
contentType := resp.Header.Get("Content-Type")
if idx := strings.Index(contentType, ";"); idx >= 0 {
contentType = strings.TrimSpace(contentType[:idx])
}
return s.processUploadContent(ctx, accountID, model.DirectUploadSourceAccount, filename, contentType, int64(len(data)), bytes.NewReader(data))
}
// --- Account Direct Upload (staged for message attachment) ---
// AccountDirectUpload handles a staged file upload from the dashboard (account-scoped).
// Returns a blob/UUID that can be attached to a message later.
// Reference: Chatwoot POST /api/v1/accounts/:account_id/direct_uploads
func (s *UploadService) AccountDirectUpload(ctx context.Context, accountID uint, req AccountDirectUploadRequest) (*UploadResponse, error) {
if accountID == 0 {
return nil, errors.New("account_id is required")
}
if req.FileHeader == nil {
return nil, errors.New("file is required")
}
return s.processUpload(ctx, accountID, req.FileHeader, model.DirectUploadSourceAccount)
}
// --- Widget Direct Upload ---
// WidgetDirectUpload handles a file upload from the widget (visitor direct upload).
// Reference: Chatwoot POST /widget/direct_uploads
func (s *UploadService) WidgetDirectUpload(ctx context.Context, req WidgetDirectUploadRequest) (*UploadResponse, error) {
if req.FileHeader == nil {
return nil, errors.New("file is required")
}
session, err := s.validateWidgetUploadSession(ctx, req.WebsiteToken, req.AuthToken)
if err != nil {
return nil, err
}
result, err := s.processUpload(ctx, session.AccountID, req.FileHeader, model.DirectUploadSourceWidget)
if err != nil {
return nil, err
}
if err := s.addUploadSessionMetadata(ctx, result.UploadUUID, session.ContactInboxID); err != nil {
return nil, err
}
return result, nil
}
func (s *UploadService) CreateWidgetDirectUpload(ctx context.Context, req ActiveStorageDirectUploadRequest) (*ActiveStorageDirectUploadResponse, error) {
session, err := s.validateWidgetUploadSession(ctx, req.WebsiteToken, req.AuthToken)
if err != nil {
return nil, err
}
if req.Blob.Metadata == nil {
req.Blob.Metadata = map[string]any{}
}
req.Blob.Metadata["contact_inbox_id"] = session.ContactInboxID
return s.createActiveStorageDirectUpload(ctx, session.AccountID, 0, model.DirectUploadSourceWidget, req, "/api/v1/widget/direct_uploads/")
}
func (s *UploadService) CreateConversationDirectUpload(ctx context.Context, accountID, conversationID uint, req ActiveStorageDirectUploadRequest) (*ActiveStorageDirectUploadResponse, error) {
if accountID == 0 {
return nil, errors.New("account_id is required")
}
if conversationID == 0 {
return nil, errors.New("conversation_id is required")
}
if s.conversationRepo == nil {
return nil, errors.New("conversation repository is not configured")
}
if _, err := s.conversationRepo.FindByAccountAndDisplayIDOrID(ctx, accountID, conversationID); err != nil {
return nil, fmt.Errorf("conversation not found: %w", err)
}
urlPrefix := fmt.Sprintf("/api/v1/accounts/%d/conversations/%d/direct_uploads/", accountID, conversationID)
return s.createActiveStorageDirectUpload(ctx, accountID, accountID, model.DirectUploadSourceAccount, req, urlPrefix)
}
func (s *UploadService) validateWidgetUploadSession(ctx context.Context, websiteToken, authToken string) (widgetUploadSession, error) {
if websiteToken == "" {
return widgetUploadSession{}, errors.New("website_token is required")
}
if authToken == "" {
return widgetUploadSession{}, errors.New("widget auth token is required")
}
if s.inboxRepo == nil || s.contactInboxRepo == nil {
return widgetUploadSession{}, errors.New("widget upload authentication is not configured")
}
inbox, err := s.inboxRepo.FindByWebsiteToken(ctx, websiteToken)
if err != nil {
return widgetUploadSession{}, fmt.Errorf("invalid website_token: %w", err)
}
if !inbox.Enabled {
return widgetUploadSession{}, errors.New("inbox is disabled")
}
contactInbox, err := s.contactInboxRepo.FindByPubsubToken(ctx, authToken)
if err != nil {
return widgetUploadSession{}, fmt.Errorf("invalid widget auth token: %w", err)
}
if contactInbox.InboxID != inbox.ID {
return widgetUploadSession{}, errors.New("widget auth token does not belong to this inbox")
}
return widgetUploadSession{AccountID: inbox.AccountID, ContactInboxID: contactInbox.ID}, nil
}
func (s *UploadService) addUploadSessionMetadata(ctx context.Context, uploadUUID string, contactInboxID uint) error {
upload, err := s.directUploadRepo.FindByUUID(ctx, uploadUUID)
if err != nil {
return err
}
metadata := map[string]any{"contact_inbox_id": contactInboxID}
upload.Metadata, _ = json.Marshal(metadata)
return s.directUploadRepo.Update(ctx, upload)
}
func (s *UploadService) CompleteWidgetDirectUpload(ctx context.Context, uploadUUID string, body io.Reader, uploadTokens ...string) (*UploadResponse, error) {
if uploadUUID == "" {
return nil, errors.New("upload_uuid is required")
}
upload, err := s.directUploadRepo.FindByUUID(ctx, uploadUUID)
if err != nil {
return nil, fmt.Errorf("direct upload not found: %w", err)
}
if upload.Source != model.DirectUploadSourceWidget {
return nil, errors.New("direct upload source mismatch")
}
if time.Now().After(upload.ExpiresAt) {
upload.Status = model.DirectUploadStatusExpired
_ = s.directUploadRepo.Update(ctx, upload)
return nil, errors.New("direct upload has expired")
}
providedToken := ""
if len(uploadTokens) > 0 {
providedToken = uploadTokens[0]
}
if providedToken == "" || providedToken != directUploadToken(upload) {
return nil, errors.New("invalid direct upload token")
}
if err := s.completeDirectUpload(ctx, upload, body); err != nil {
return nil, fmt.Errorf("failed to save direct upload: %w", err)
}
return &UploadResponse{
UploadID: upload.ID,
UploadUUID: upload.UploadUUID,
OriginalName: upload.OriginalName,
FileType: upload.FileType,
MimeType: upload.MimeType,
FileSize: upload.FileSize,
FileURL: upload.FileURL,
ThumbURL: upload.ThumbURL,
Status: string(upload.Status),
ExpiresAt: upload.ExpiresAt,
}, nil
}
func (s *UploadService) CompleteConversationDirectUpload(ctx context.Context, accountID, conversationID uint, uploadUUID string, body io.Reader) (*UploadResponse, error) {
if accountID == 0 {
return nil, errors.New("account_id is required")
}
if conversationID == 0 {
return nil, errors.New("conversation_id is required")
}
if s.conversationRepo == nil {
return nil, errors.New("conversation repository is not configured")
}
if _, err := s.conversationRepo.FindByAccountAndDisplayIDOrID(ctx, accountID, conversationID); err != nil {
return nil, fmt.Errorf("conversation not found: %w", err)
}
if uploadUUID == "" {
return nil, errors.New("upload_uuid is required")
}
upload, err := s.directUploadRepo.FindByUUID(ctx, uploadUUID)
if err != nil {
return nil, fmt.Errorf("direct upload not found: %w", err)
}
if upload.Source != model.DirectUploadSourceAccount || upload.AccountID != accountID {
return nil, errors.New("direct upload source mismatch")
}
if time.Now().After(upload.ExpiresAt) {
upload.Status = model.DirectUploadStatusExpired
_ = s.directUploadRepo.Update(ctx, upload)
return nil, errors.New("direct upload has expired")
}
if err := s.completeDirectUpload(ctx, upload, body); err != nil {
return nil, fmt.Errorf("failed to save direct upload: %w", err)
}
return &UploadResponse{
UploadID: upload.ID,
UploadUUID: upload.UploadUUID,
OriginalName: upload.OriginalName,
FileType: upload.FileType,
MimeType: upload.MimeType,
FileSize: upload.FileSize,
FileURL: upload.FileURL,
ThumbURL: upload.ThumbURL,
Status: string(upload.Status),
ExpiresAt: upload.ExpiresAt,
}, nil
}
func directUploadToken(upload *model.DirectUpload) string {
if upload == nil || len(upload.Metadata) == 0 {
return ""
}
var metadata map[string]any
if json.Unmarshal(upload.Metadata, &metadata) != nil {
return ""
}
token, _ := metadata["active_storage_key"].(string)
return token
}
func (s *UploadService) completeDirectUpload(ctx context.Context, upload *model.DirectUpload, body io.Reader) error {
data, mimeType, err := readValidatedUpload(body, upload.OriginalName, upload.MimeType, upload.FileSize, upload.FileSize, true)
if err != nil {
return err
}
if err := s.saveReaderToDisk(upload.FileURL, bytes.NewReader(data), int64(len(data))); err != nil {
return err
}
upload.MimeType = mimeType
upload.FileType = categorizeUploadMIME(mimeType)
upload.FileSize = int64(len(data))
return s.directUploadRepo.Update(ctx, upload)
}
// --- Internal helpers ---
func (s *UploadService) createActiveStorageDirectUpload(ctx context.Context, accountID, storageAccountID uint, source model.DirectUploadSource, req ActiveStorageDirectUploadRequest, directUploadURLPrefix string) (*ActiveStorageDirectUploadResponse, error) {
if req.Blob.Filename == "" {
return nil, errors.New("filename is required")
}
if req.Blob.ByteSize <= 0 {
return nil, errors.New("byte_size is required")
}
mimeType := normalizeUploadMIME(req.Blob.ContentType)
if mimeType == "" || mimeType == "application/octet-stream" {
mimeType = detectUploadMIMEFromFilename(req.Blob.Filename)
}
fileCategory := categorizeUploadMIME(mimeType)
if fileCategory == "" {
return nil, fmt.Errorf("unsupported file type: %s", mimeType)
}
if !isUploadMIMEAllowed(fileCategory, mimeType) {
return nil, fmt.Errorf("MIME type %s is not allowed for category %s", mimeType, fileCategory)
}
maxSize := s.categoryUploadLimit(fileCategory)
if req.Blob.ByteSize > maxSize {
return nil, fmt.Errorf("file size %d exceeds maximum %d for type %s", req.Blob.ByteSize, maxSize, fileCategory)
}
uploadUUID := uuid.New().String()
fileURL, thumbURL := s.directUploadURL(source, storageAccountID, uploadUUID, mimeType)
metadata := req.Blob.Metadata
if metadata == nil {
metadata = map[string]any{}
}
metadata["checksum"] = req.Blob.Checksum
metadata["active_storage_key"] = randomStorageKey()
metadataJSON, _ := json.Marshal(metadata)
upload := &model.DirectUpload{
UploadUUID: uploadUUID,
AccountID: accountID,
Status: model.DirectUploadStatusPending,
Source: source,
OriginalName: req.Blob.Filename,
FileType: fileCategory,
MimeType: mimeType,
FileSize: req.Blob.ByteSize,
FileURL: fileURL,
ThumbURL: thumbURL,
Metadata: metadataJSON,
ExpiresAt: time.Now().Add(24 * time.Hour),
}
if err := s.directUploadRepo.Create(ctx, upload); err != nil {
return nil, fmt.Errorf("failed to create direct upload: %w", err)
}
return &ActiveStorageDirectUploadResponse{
ID: upload.ID,
Key: fmt.Sprint(metadata["active_storage_key"]),
Filename: upload.OriginalName,
ContentType: upload.MimeType,
Metadata: metadata,
ServiceName: "gochat_local",
ByteSize: upload.FileSize,
Checksum: req.Blob.Checksum,
CreatedAt: upload.CreatedAt,
SignedID: upload.UploadUUID,
DirectUpload: ActiveStorageUploadURL{
URL: directUploadURLPrefix + upload.UploadUUID + "?token=" + url.QueryEscape(fmt.Sprint(metadata["active_storage_key"])),
Headers: map[string]string{
"Content-Type": upload.MimeType,
},
},
}, nil
}
func (s *UploadService) processUpload(ctx context.Context, accountID uint, fileHeader *multipart.FileHeader, source model.DirectUploadSource) (*UploadResponse, error) {
mimeType := fileHeader.Header.Get("Content-Type")
src, err := fileHeader.Open()
if err != nil {
return nil, fmt.Errorf("failed to open uploaded file: %w", err)
}
defer src.Close()
return s.processUploadContent(ctx, accountID, source, fileHeader.Filename, mimeType, fileHeader.Size, src)
}
func (s *UploadService) processUploadContent(ctx context.Context, accountID uint, source model.DirectUploadSource, filename, mimeType string, size int64, reader io.Reader) (*UploadResponse, error) {
data, mimeType, err := readValidatedUpload(reader, filename, mimeType, size, s.maxUploadSize(), true)
if err != nil {
return nil, err
}
fileCategory := categorizeUploadMIME(mimeType)
if fileCategory == "" {
return nil, fmt.Errorf("unsupported file type: %s", mimeType)
}
if !isUploadMIMEAllowed(fileCategory, mimeType) {
return nil, fmt.Errorf("MIME type %s is not allowed for category %s", mimeType, fileCategory)
}
maxSize := s.categoryUploadLimit(fileCategory)
if int64(len(data)) > maxSize {
return nil, fmt.Errorf("file size %d exceeds maximum %d for type %s", len(data), maxSize, fileCategory)
}
// Step 2: Store file to disk
ext := extensionForUploadMIME(mimeType)
fileURL, thumbURL, err := s.saveUploadReader(accountID, source, ext, mimeType, bytes.NewReader(data))
if err != nil {
return nil, fmt.Errorf("failed to save file: %w", err)
}
// Step 3: Create upload record with UUID
expiryDuration := 24 * time.Hour
uploadUUID := uuid.New().String()
upload := &model.DirectUpload{
UploadUUID: uploadUUID,
AccountID: accountID,
Status: model.DirectUploadStatusPending,
Source: source,
OriginalName: filename,
FileType: fileCategory,
MimeType: mimeType,
FileSize: int64(len(data)),
FileURL: fileURL,
ThumbURL: thumbURL,
ExpiresAt: time.Now().Add(expiryDuration),
}
if err := s.directUploadRepo.Create(ctx, upload); err != nil {
// Clean up file on disk if DB insert fails
os.Remove(filepath.Join(s.cfg.Storage.LocalPath, fileURL))
return nil, fmt.Errorf("failed to create upload record: %w", err)
}
applogger.L().Infof("Direct file upload: account=%d source=%s uuid=%s file=%s size=%d",
accountID, source, uploadUUID, filename, size)
return &UploadResponse{
UploadID: upload.ID,
UploadUUID: uploadUUID,
OriginalName: upload.OriginalName,
FileType: upload.FileType,
MimeType: upload.MimeType,
FileSize: upload.FileSize,
FileURL: upload.FileURL,
ThumbURL: upload.ThumbURL,
Status: string(upload.Status),
ExpiresAt: upload.ExpiresAt,
}, nil
}
func (s *UploadService) saveUploadReader(accountID uint, source model.DirectUploadSource, ext, mimeType string, reader io.Reader) (string, string, error) {
dirPath := s.uploadDir(source, accountID)
if err := os.MkdirAll(dirPath, 0755); err != nil {
return "", "", fmt.Errorf("failed to create upload directory: %w", err)
}
// Generate unique filename
timestamp := time.Now().UnixMilli()
baseName := fmt.Sprintf("%d_%s", timestamp, uuid.New().String()[:8])
fileName := baseName + ext
fullPath := filepath.Join(dirPath, fileName)
// Create destination file
dst, err := os.Create(fullPath)
if err != nil {
return "", "", fmt.Errorf("failed to create destination file: %w", err)
}
defer dst.Close()
// Copy file content
if _, err := io.Copy(dst, reader); err != nil {
os.Remove(fullPath) // Clean up on failure
return "", "", fmt.Errorf("failed to copy file content: %w", err)
}
fileURL := s.uploadURL(source, accountID, fileName)
thumbURL := ""
// For images, we reference the same path (thumbnail generation can be added later)
if strings.HasPrefix(mimeType, "image/") {
thumbURL = fileURL // Placeholder: same as fileURL for now
}
return fileURL, thumbURL, nil
}
func (s *UploadService) saveReaderToDisk(fileURL string, body io.Reader, expectedSizes ...int64) error {
localPath := s.cfg.Storage.LocalPath
if localPath == "" {
localPath = "./uploads"
}
relative := strings.TrimPrefix(fileURL, "/uploads/")
fullPath := filepath.Join(localPath, relative)
if err := os.MkdirAll(filepath.Dir(fullPath), 0755); err != nil {
return err
}
dst, err := os.CreateTemp(filepath.Dir(fullPath), ".upload-*")
if err != nil {
return err
}
tmpPath := dst.Name()
defer os.Remove(tmpPath)
written, copyErr := io.Copy(dst, body)
closeErr := dst.Close()
if copyErr != nil {
return copyErr
}
if closeErr != nil {
return closeErr
}
if len(expectedSizes) > 0 && written != expectedSizes[0] {
return fmt.Errorf("file size mismatch: wrote %d bytes, expected %d", written, expectedSizes[0])
}
return os.Rename(tmpPath, fullPath)
}
// ResolveAuthorizedUpload maps a stored upload URL to disk only after the
// requesting account or widget conversation has been authorized.
func (s *UploadService) ResolveAuthorizedUpload(ctx context.Context, fileURL string, accountID uint, widgetToken string) (string, bool) {
if s.accessDB == nil || !strings.HasPrefix(fileURL, "/uploads/") {
return "", false
}
fileURL = "/uploads/" + strings.TrimPrefix(filepath.ToSlash(filepath.Clean(strings.TrimPrefix(fileURL, "/uploads/"))), "/")
contactInboxID, widgetInboxID, widgetContactID := s.widgetAccessIdentity(ctx, widgetToken)
var attachment struct {
AccountID uint
ContactInboxID *uint
}
if err := s.accessDB.WithContext(ctx).Raw(`
SELECT attachments.account_id, conversations.contact_inbox_id
FROM attachments
JOIN messages ON messages.id = attachments.message_id AND messages.deleted_at IS NULL
JOIN conversations ON conversations.id = messages.conversation_id AND conversations.deleted_at IS NULL
WHERE attachments.deleted_at IS NULL AND (attachments.file_url = ? OR attachments.thumb_url = ?)
LIMIT 1`, fileURL, fileURL).Scan(&attachment).Error; err == nil && attachment.AccountID != 0 {
allowed := accountID == attachment.AccountID ||
(contactInboxID != 0 && attachment.ContactInboxID != nil && *attachment.ContactInboxID == contactInboxID)
return s.localUploadPath(fileURL, allowed)
}
var directUpload model.DirectUpload
if err := s.accessDB.WithContext(ctx).
Where("file_url = ? OR thumb_url = ?", fileURL, fileURL).
First(&directUpload).Error; err == nil {
allowed := accountID != 0 && accountID == directUpload.AccountID
if directUpload.Source == model.DirectUploadSourceWidget && contactInboxID != 0 {
allowed = allowed || metadataUint(directUpload.Metadata, "contact_inbox_id") == contactInboxID
}
return s.localUploadPath(fileURL, allowed)
}
var widgetUpload model.WidgetFileUpload
if err := s.accessDB.WithContext(ctx).
Where("file_url = ? OR thumb_url = ?", fileURL, fileURL).
First(&widgetUpload).Error; err == nil {
var inbox model.Inbox
_ = s.accessDB.WithContext(ctx).Select("account_id").First(&inbox, widgetUpload.InboxID).Error
allowed := (accountID != 0 && accountID == inbox.AccountID) ||
(widgetInboxID == widgetUpload.InboxID && widgetContactID != 0 && widgetContactID == widgetUpload.ContactID)
return s.localUploadPath(fileURL, allowed)
}
// Durable account files such as profile avatars do not have attachment rows.
// Keep them account-private by the existing /account/<id>/ storage boundary.
parts := strings.Split(strings.TrimPrefix(fileURL, "/uploads/"), "/")
if len(parts) >= 3 && parts[0] == "account" {
pathAccountID, err := strconv.ParseUint(parts[1], 10, 32)
return s.localUploadPath(fileURL, err == nil && accountID != 0 && uint(pathAccountID) == accountID)
}
return "", false
}
func (s *UploadService) widgetAccessIdentity(ctx context.Context, token string) (uint, uint, uint) {
if token == "" || s.contactInboxRepo == nil {
return 0, 0, 0
}
contactInbox, err := s.contactInboxRepo.FindByPubsubToken(ctx, token)
if err != nil {
return 0, 0, 0
}
return contactInbox.ID, contactInbox.InboxID, contactInbox.ContactID
}
func (s *UploadService) localUploadPath(fileURL string, allowed bool) (string, bool) {
if !allowed {
return "", false
}
localPath := s.cfg.Storage.LocalPath
if localPath == "" {
localPath = "./uploads"
}
base, err := filepath.Abs(localPath)
if err != nil {
return "", false
}
fullPath, err := filepath.Abs(filepath.Join(base, filepath.FromSlash(strings.TrimPrefix(fileURL, "/uploads/"))))
if err != nil {
return "", false
}
relative, err := filepath.Rel(base, fullPath)
if err != nil || relative == ".." || strings.HasPrefix(relative, ".."+string(filepath.Separator)) {
return "", false
}
if info, err := os.Stat(fullPath); err != nil || info.IsDir() {
return "", false
}
return fullPath, true
}
func metadataUint(raw []byte, key string) uint {
var metadata map[string]any
if json.Unmarshal(raw, &metadata) != nil {
return 0
}
switch value := metadata[key].(type) {
case float64:
return uint(value)
case string:
parsed, _ := strconv.ParseUint(value, 10, 32)
return uint(parsed)
default:
return 0
}
}
func (s *UploadService) directUploadURL(source model.DirectUploadSource, accountID uint, uploadUUID, mimeType string) (string, string) {
ext := extensionForUploadMIME(mimeType)
fileName := uploadUUID + ext
fileURL := s.uploadURL(source, accountID, fileName)
thumbURL := ""
if strings.HasPrefix(mimeType, "image/") {
thumbURL = fileURL
}
return fileURL, thumbURL
}
func (s *UploadService) uploadDir(source model.DirectUploadSource, accountID uint) string {
localPath := s.cfg.Storage.LocalPath
if localPath == "" {
localPath = "./uploads"
}
subDir := uploadSubDir(source)
parts := []string{localPath, subDir}
if accountID > 0 {
parts = append(parts, fmt.Sprintf("%d", accountID))
}
return filepath.Join(parts...)
}
func (s *UploadService) uploadURL(source model.DirectUploadSource, accountID uint, fileName string) string {
subDir := uploadSubDir(source)
if accountID > 0 {
return fmt.Sprintf("/uploads/%s/%d/%s", subDir, accountID, fileName)
}
return fmt.Sprintf("/uploads/%s/%s", subDir, fileName)
}
func uploadSubDir(source model.DirectUploadSource) string {
if source == model.DirectUploadSourceWidget {
return "widget_direct"
}
return "account"
}
func randomStorageKey() string {
buf := make([]byte, 16)
if _, err := rand.Read(buf); err != nil {
return uuid.New().String()
}
return hex.EncodeToString(buf)
}
// CleanupExpiredUploads removes each object before its record. Failed object
// deletions leave the row intact so the durable job can retry safely.
func (s *UploadService) CleanupExpiredUploads(ctx context.Context) (int64, error) {
now := time.Now()
uploads, err := s.directUploadRepo.FindExpired(ctx, now)
if err != nil {
return 0, fmt.Errorf("find expired uploads: %w", err)
}
var count int64
var cleanupErr error
for i := range uploads {
if err := s.removeUploadFile(uploads[i].FileURL); err != nil {
cleanupErr = errors.Join(cleanupErr, fmt.Errorf("remove upload %s: %w", uploads[i].UploadUUID, err))
continue
}
if err := s.directUploadRepo.Delete(ctx, uploads[i].ID); err != nil {
cleanupErr = errors.Join(cleanupErr, fmt.Errorf("delete upload %s: %w", uploads[i].UploadUUID, err))
continue
}
count++
}
orphans, err := s.cleanupOrphanedUploadFiles(ctx, now.Add(-uploadCleanupInterval))
count += orphans
cleanupErr = errors.Join(cleanupErr, err)
applogger.L().Infof("Cleaned up %d expired or orphaned uploads", count)
return count, cleanupErr
}
// cleanupOrphanedUploadFiles reconciles only roots owned by UploadService.
// The grace period keeps an in-flight file-to-row write from racing cleanup.
func (s *UploadService) cleanupOrphanedUploadFiles(ctx context.Context, cutoff time.Time) (int64, error) {
if s.accessDB == nil {
return 0, nil
}
references := map[string]struct{}{}
var urls []string
for _, column := range []string{"file_url", "thumb_url"} {
urls = nil
if err := s.accessDB.WithContext(ctx).Model(&model.DirectUpload{}).
Where(column+" LIKE ?", "/uploads/%").Pluck(column, &urls).Error; err != nil {
return 0, fmt.Errorf("load direct upload references: %w", err)
}
for _, url := range urls {
references[url] = struct{}{}
}
}
if s.accessDB.Migrator().HasTable(&model.User{}) {
urls = nil
if err := s.accessDB.WithContext(ctx).Model(&model.User{}).
Where("avatar_url LIKE ?", "/uploads/account/%").Pluck("avatar_url", &urls).Error; err != nil {
return 0, fmt.Errorf("load avatar references: %w", err)
}
for _, url := range urls {
references[url] = struct{}{}
}
}
localPath := "./uploads"
if s.cfg != nil && s.cfg.Storage.LocalPath != "" {
localPath = s.cfg.Storage.LocalPath
}
var removed int64
var reconcileErr error
for _, subdir := range []string{"account", "widget_direct"} {
root := filepath.Join(localPath, subdir)
err := filepath.Walk(root, func(path string, info os.FileInfo, walkErr error) error {
if walkErr != nil {
if errors.Is(walkErr, os.ErrNotExist) {
return nil
}
return walkErr
}
if !info.Mode().IsRegular() || info.ModTime().After(cutoff) {
return nil
}
relative, err := filepath.Rel(localPath, path)
if err != nil {
return err
}
fileURL := "/uploads/" + filepath.ToSlash(relative)
if _, ok := references[fileURL]; ok {
return nil
}
if err := os.Remove(path); err != nil {
if errors.Is(err, os.ErrNotExist) {
return nil
}
return err
}
removed++
return nil
})
if err != nil && !errors.Is(err, os.ErrNotExist) {
reconcileErr = errors.Join(reconcileErr, fmt.Errorf("reconcile %s uploads: %w", subdir, err))
}
}
return removed, reconcileErr
}
func (s *UploadService) removeUploadFile(fileURL string) error {
relative := strings.TrimPrefix(fileURL, "/uploads/")
if relative == fileURL || relative == "" || filepath.IsAbs(relative) || strings.HasPrefix(filepath.Clean(relative), "..") {
return fmt.Errorf("invalid upload path %q", fileURL)
}
localPath := s.cfg.Storage.LocalPath
if localPath == "" {
localPath = "./uploads"
}
err := os.Remove(filepath.Join(localPath, relative))
if errors.Is(err, os.ErrNotExist) {
return nil
}
return err
}
// RegisterUploadCleanupJobs runs orphan cleanup daily through the existing
// durable, idempotent worker queue.
func RegisterUploadCleanupJobs(wp *worker.WorkerPool, svc *UploadService) {
wp.Register(TaskTypeCleanupExpiredUploads, func(ctx context.Context, _ *model.BackgroundJob) error {
_, cleanupErr := svc.CleanupExpiredUploads(ctx)
_, enqueueErr := EnqueueUploadCleanup(ctx, wp, time.Now().Add(uploadCleanupInterval))
return errors.Join(cleanupErr, enqueueErr)
})
}
func EnqueueUploadCleanup(ctx context.Context, wp *worker.WorkerPool, at time.Time) (*model.BackgroundJob, error) {
bucket := at.UTC().Truncate(uploadCleanupInterval).Unix()
return wp.Enqueue(ctx, TaskTypeCleanupExpiredUploads, nil,
worker.WithQueue("low"),
worker.WithScheduledAt(at),
worker.WithMaxAttempts(10),
worker.WithIdempotencyKey(fmt.Sprintf("storage:cleanup_expired_uploads:%d", bucket)),
)
}
// --- MIME detection helpers (reuse patterns from widget_theme_service.go) ---
func (s *UploadService) maxUploadSize() int64 {
if s.cfg != nil && s.cfg.Storage.MaxFileSize > 0 {
return s.cfg.Storage.MaxFileSize
}
return 20 << 20
}
func (s *UploadService) categoryUploadLimit(category string) int64 {
limit := s.maxUploadSize()
if categoryLimit := model.WidgetUploadMaxSizeByType[category]; categoryLimit > 0 && categoryLimit < limit {
return categoryLimit
}
return limit
}
func readValidatedUpload(reader io.Reader, filename, declaredMIME string, expectedSize, maxSize int64, requireExactSize bool) ([]byte, string, error) {
if reader == nil {
return nil, "", errors.New("file content is required")
}
if maxSize <= 0 {
maxSize = 20 << 20
}
if expectedSize > maxSize {
return nil, "", fmt.Errorf("file size %d exceeds maximum %d", expectedSize, maxSize)
}
data, err := io.ReadAll(io.LimitReader(reader, maxSize+1))
if err != nil {
return nil, "", err
}
if int64(len(data)) > maxSize {
return nil, "", fmt.Errorf("file size exceeds maximum %d", maxSize)
}
if requireExactSize && expectedSize >= 0 && int64(len(data)) != expectedSize {
return nil, "", fmt.Errorf("file size mismatch: got %d, expected %d", len(data), expectedSize)
}
if len(data) == 0 {
return nil, "", errors.New("empty file is not allowed")
}
detectedMIME := normalizeUploadMIME(http.DetectContentType(data[:min(len(data), 512)]))
declaredMIME = normalizeUploadMIME(declaredMIME)
if declaredMIME == "" || declaredMIME == "application/octet-stream" {
declaredMIME = detectUploadMIMEFromFilename(filename)
if declaredMIME == "application/octet-stream" {
return nil, "", fmt.Errorf("unsupported file type: extension %s", filepath.Ext(filename))
}
}
sample := bytes.TrimSpace(data[:min(len(data), 512)])
sample = bytes.TrimSpace(bytes.TrimPrefix(sample, []byte{0xef, 0xbb, 0xbf}))
if (declaredMIME == "text/plain" || declaredMIME == "text/csv") && bytes.HasPrefix(sample, []byte("<")) {
return nil, "", fmt.Errorf("active markup is not allowed for MIME type %s", declaredMIME)
}
if !uploadMIMEMatches(declaredMIME, detectedMIME) {
return nil, "", fmt.Errorf("file content type %s does not match declared type %s", detectedMIME, declaredMIME)
}
if declaredMIME == "text/csv" && detectedMIME == "text/plain" {
return data, declaredMIME, nil
}
if declaredMIME != detectedMIME {
return data, declaredMIME, nil
}
return data, detectedMIME, nil
}
func normalizeUploadMIME(value string) string {
value = strings.TrimSpace(strings.ToLower(value))
if value == "" {
return ""
}
mediaType, _, err := mime.ParseMediaType(value)
if err == nil {
return mediaType
}
return strings.TrimSpace(strings.Split(value, ";")[0])
}
func uploadMIMEMatches(declared, detected string) bool {
if declared == detected || (declared == "text/csv" && detected == "text/plain") {
return true
}
switch declared {
case "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
"application/vnd.openxmlformats-officedocument.wordprocessingml.document":
return detected == "application/zip" || detected == "application/octet-stream"
case "application/vnd.ms-excel", "application/msword":
return detected == "application/octet-stream"
}
return false
}
func extensionForUploadMIME(mimeType string) string {
switch normalizeUploadMIME(mimeType) {
case "image/png":
return ".png"
case "image/jpeg":
return ".jpg"
case "image/gif":
return ".gif"
case "image/webp":
return ".webp"
case "audio/mpeg", "audio/mp3":
return ".mp3"
case "audio/ogg":
return ".ogg"
case "audio/wav":
return ".wav"
case "audio/webm":
return ".webm"
case "video/mp4":
return ".mp4"
case "video/webm":
return ".webm"
case "video/ogg":
return ".ogv"
case "application/pdf":
return ".pdf"
case "text/csv":
return ".csv"
case "text/plain":
return ".txt"
case "application/vnd.ms-excel":
return ".xls"
case "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet":
return ".xlsx"
case "application/msword":
return ".doc"
case "application/vnd.openxmlformats-officedocument.wordprocessingml.document":
return ".docx"
default:
return ""
}
}
func detectUploadMIMEFromFilename(filename string) string {
ext := strings.ToLower(filepath.Ext(filename))
switch ext {
case ".png":
return "image/png"
case ".jpg", ".jpeg":
return "image/jpeg"
case ".gif":
return "image/gif"
case ".webp":
return "image/webp"
case ".svg":
return "image/svg+xml"
case ".mp3":
return "audio/mpeg"
case ".ogg":
return "audio/ogg"
case ".wav":
return "audio/wav"
case ".webm":
return "audio/webm"
case ".mp4":
return "video/mp4"
case ".pdf":
return "application/pdf"
case ".csv":
return "text/csv"
case ".txt":
return "text/plain"
case ".xls":
return "application/vnd.ms-excel"
case ".xlsx":
return "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet"
case ".doc":
return "application/msword"
case ".docx":
return "application/vnd.openxmlformats-officedocument.wordprocessingml.document"
default:
return "application/octet-stream"
}
}
func categorizeUploadMIME(mimeType string) string {
if strings.HasPrefix(mimeType, "image/") {
return "image"
}
if strings.HasPrefix(mimeType, "audio/") {
return "audio"
}
if strings.HasPrefix(mimeType, "video/") {
return "video"
}
// Check specific file MIME types
switch mimeType {
case "application/pdf",
"application/vnd.ms-excel",
"application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
"application/msword",
"application/vnd.openxmlformats-officedocument.wordprocessingml.document",
"text/plain",
"text/csv":
return "file"
}
return "" // unsupported
}
func isUploadMIMEAllowed(category string, mimeType string) bool {
allowed, ok := model.WidgetUploadAllowedTypes[category]
if !ok {
return false
}
for _, a := range allowed {
if a == mimeType {
return true
}
}
return false
}