diff --git a/Wavelet/internal/apps/openflare/agent/routers.go b/Wavelet/internal/apps/openflare/agent/routers.go index eabd176f..e44498e4 100644 --- a/Wavelet/internal/apps/openflare/agent/routers.go +++ b/Wavelet/internal/apps/openflare/agent/routers.go @@ -137,13 +137,17 @@ func DownloadPagesPackageHandler(c *gin.Context) { if !ok { return } - filePath, fileName, err := pages.GetDeploymentPackagePath(c.Request.Context(), deploymentID) + packageObj, fileName, err := pages.OpenDeploymentPackage(c.Request.Context(), deploymentID) if err != nil { compat.Fail(c, err.Error()) return } + defer packageObj.Body.Close() c.Header("Content-Disposition", "attachment; filename="+fileName) - c.File(filePath) + if packageObj.ContentType != "" { + c.Header("Content-Type", packageObj.ContentType) + } + c.DataFromReader(http.StatusOK, packageObj.ContentLength, packageObj.ContentType, packageObj.Body, nil) } func pagesDeploymentIDParam(c *gin.Context) (uint, bool) { diff --git a/Wavelet/internal/apps/openflare/pages/errs.go b/Wavelet/internal/apps/openflare/pages/errs.go index 80e05f01..412a81e7 100644 --- a/Wavelet/internal/apps/openflare/pages/errs.go +++ b/Wavelet/internal/apps/openflare/pages/errs.go @@ -21,6 +21,7 @@ const ( errPagesAPIProxyPassRequired = "启用 API 反代时,后端服务地址不能为空" errPagesAPIProxyPassInvalid = "API 反代后端服务地址必须是有效的 HTTP/HTTPS URL" errPagesPackagePathEmpty = "Pages 部署包路径为空" + errPagesPackageUploadMissing = "Pages 部署包上传记录不存在" errPagesPackageNotInActiveConfig = "Pages 部署尚未进入激活配置" errPagesInvalidSnapshotFormat = "配置快照格式无效" ) diff --git a/Wavelet/internal/apps/openflare/pages/helpers.go b/Wavelet/internal/apps/openflare/pages/helpers.go index 0b7a3832..7515f798 100644 --- a/Wavelet/internal/apps/openflare/pages/helpers.go +++ b/Wavelet/internal/apps/openflare/pages/helpers.go @@ -5,6 +5,7 @@ package pages import ( "archive/zip" + "context" "crypto/sha256" "encoding/hex" "errors" @@ -17,8 +18,9 @@ import ( "regexp" "strings" - "github.com/Rain-kl/Wavelet/internal/config" + "github.com/Rain-kl/Wavelet/internal/apps/upload" "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/repository" ) const ( @@ -26,6 +28,7 @@ const ( pagesMaxDeploymentBytes = 100 * 1024 * 1024 defaultPagesEntryFile = "index.html" defaultPagesFallbackPath = "/index.html" + pagesDeploymentUploadType = "openflare_pages_deployment" ) var pagesSlugPattern = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,126}[a-z0-9]$|^[a-z0-9]$`) @@ -144,15 +147,15 @@ func normalizePagesEntryFile(raw string) string { return strings.TrimPrefix(value, "/") } -func persistPagesUploadTemp(fileHeader *multipart.FileHeader) (string, string, error) { +func persistPagesUploadTemp(fileHeader *multipart.FileHeader) (string, string, int64, error) { file, err := fileHeader.Open() if err != nil { - return "", "", err + return "", "", 0, err } defer file.Close() temp, err := os.CreateTemp("", "openflare-pages-*.zip") if err != nil { - return "", "", err + return "", "", 0, err } defer temp.Close() hash := sha256.New() @@ -160,13 +163,64 @@ func persistPagesUploadTemp(fileHeader *multipart.FileHeader) (string, string, e written, err := io.Copy(io.MultiWriter(temp, hash), limited) if err != nil { _ = os.Remove(temp.Name()) - return "", "", err + return "", "", 0, err } if written > pagesMaxDeploymentBytes { _ = os.Remove(temp.Name()) - return "", "", fmt.Errorf("Pages 部署包不能超过 %d MiB", pagesMaxDeploymentBytes/1024/1024) + return "", "", 0, fmt.Errorf("Pages 部署包不能超过 %d MiB", pagesMaxDeploymentBytes/1024/1024) + } + return temp.Name(), hex.EncodeToString(hash.Sum(nil)), written, nil +} + +func ingestPagesDeploymentPackage( + ctx context.Context, + tempPath string, + checksum string, + size int64, + projectSlug string, + fileName string, +) (upload.IngestResult, error) { + file, err := os.Open(tempPath) + if err != nil { + return upload.IngestResult{}, err + } + defer file.Close() + + systemUser := repository.GetSystemUser(ctx) + accessMode := 0 + return upload.Ingest(ctx, upload.IngestRequest{ + UserID: systemUser.ID, + Reader: file, + Size: size, + FileName: fileName, + MimeType: "application/zip", + Extension: "zip", + Hash: checksum, + Type: pagesDeploymentUploadType, + AccessMode: &accessMode, + SkipExtensionCheck: true, + Policy: upload.PolicyDedupNewRecord, + Metadata: model.UploadMetadata{ + Extra: map[string]any{ + "project_slug": projectSlug, + }, + }, + }) +} + +func removeDeploymentArtifact(ctx context.Context, deployment *model.PagesDeployment) { + if deployment == nil { + return + } + if deployment.UploadID > 0 { + if _, err := upload.Remove(ctx, deployment.UploadID); err != nil { + return + } + return + } + if strings.TrimSpace(deployment.ArtifactPath) != "" { + _ = os.Remove(deployment.ArtifactPath) } - return temp.Name(), hex.EncodeToString(hash.Sum(nil)), nil } func findCommonRootPrefix(files []*zip.File) (string, error) { @@ -313,43 +367,4 @@ func checksumZipFile(item *zip.File) (string, error) { return hex.EncodeToString(hash.Sum(nil)), nil } -func pagesArtifactPath(projectSlug string, checksum string) (string, error) { - root, err := pagesStorageRoot() - if err != nil { - return "", err - } - return filepath.Join(root, "artifacts", projectSlug, checksum+".zip"), nil -} -func pagesStorageRoot() (string, error) { - cfg := config.Config.Database - if cfg.Enabled { - return filepath.Abs(filepath.Join("data", "pages")) - } - dbPath := strings.TrimSpace(cfg.SQLitePath) - if dbPath == "" || dbPath == ":memory:" { - return filepath.Abs(filepath.Join("data", "pages")) - } - dir := filepath.Dir(dbPath) - if dir == "." || dir == "" { - dir = "data" - } - return filepath.Abs(filepath.Join(dir, "pages")) -} - -func copyFile(src string, dst string) error { - input, err := os.Open(src) - if err != nil { - return err - } - defer input.Close() - output, err := os.OpenFile(dst, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o644) - if err != nil { - return err - } - defer output.Close() - if _, err = io.Copy(output, input); err != nil { - return err - } - return output.Sync() -} diff --git a/Wavelet/internal/apps/openflare/pages/logics.go b/Wavelet/internal/apps/openflare/pages/logics.go index 846667c1..5e9e38d2 100644 --- a/Wavelet/internal/apps/openflare/pages/logics.go +++ b/Wavelet/internal/apps/openflare/pages/logics.go @@ -15,8 +15,12 @@ import ( "strings" "time" + "github.com/Rain-kl/Wavelet/internal/apps/upload" + uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage" "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/internal/storage" "gorm.io/gorm" ) @@ -185,8 +189,8 @@ func DeleteProject(ctx context.Context, id uint) error { if err := tx.Delete(project).Error; err != nil { return err } - for _, deployment := range deployments { - _ = os.Remove(deployment.ArtifactPath) + for index := range deployments { + removeDeploymentArtifact(ctx, &deployments[index]) } return nil }) @@ -248,7 +252,7 @@ func UploadDeployment(ctx context.Context, projectID uint, fileHeader *multipart return nil, err } entryFile := normalizePagesEntryFile(project.EntryFile) - tempPath, checksum, err := persistPagesUploadTemp(fileHeader) + tempPath, checksum, packageSize, err := persistPagesUploadTemp(fileHeader) if err != nil { return nil, err } @@ -257,16 +261,23 @@ func UploadDeployment(ctx context.Context, projectID uint, fileHeader *multipart if err != nil { return nil, err } - artifactPath, err := pagesArtifactPath(project.Slug, checksum) + ingestResult, err := ingestPagesDeploymentPackage( + ctx, + tempPath, + checksum, + packageSize, + project.Slug, + fileHeader.Filename, + ) if err != nil { return nil, err } - if err = os.MkdirAll(filepath.Dir(artifactPath), 0o755); err != nil { - return nil, fmt.Errorf("创建 Pages 存储目录失败: %w", err) - } - if err = copyFile(tempPath, artifactPath); err != nil { - return nil, err - } + ingestCommitted := false + defer func() { + if !ingestCommitted && ingestResult.Created { + _, _ = upload.Remove(ctx, ingestResult.Upload.ID) + } + }() deployment := &model.PagesDeployment{} err = db.DB(ctx).Transaction(func(tx *gorm.DB) error { var maxNumber int @@ -281,7 +292,7 @@ func UploadDeployment(ctx context.Context, projectID uint, fileHeader *multipart DeploymentNumber: maxNumber + 1, Checksum: checksum, Status: model.PagesDeploymentStatusUploaded, - ArtifactPath: artifactPath, + UploadID: ingestResult.Upload.ID, FileCount: manifest.FileCount, TotalSize: manifest.TotalSize, CreatedBy: strings.TrimSpace(createdBy), @@ -300,9 +311,9 @@ func UploadDeployment(ctx context.Context, projectID uint, fileHeader *multipart return nil }) if err != nil { - _ = os.Remove(artifactPath) return nil, err } + ingestCommitted = true view := buildDeploymentView(deployment) return &view, nil } @@ -342,22 +353,47 @@ func ActivateDeployment(ctx context.Context, projectID uint, deploymentID uint) return GetProject(ctx, project.ID) } -// GetDeploymentPackagePath returns the on-disk artifact path and download filename for an agent package request. -func GetDeploymentPackagePath(ctx context.Context, deploymentID uint) (string, string, error) { +// OpenDeploymentPackage opens the deployment artifact from the upload storage framework. +func OpenDeploymentPackage(ctx context.Context, deploymentID uint) (*storage.Object, string, error) { deployment, err := model.GetPagesDeploymentByID(ctx, deploymentID) if err != nil { - return "", "", err + return nil, "", err } if err = ensureDeploymentInActiveSnapshot(ctx, deployment.ID); err != nil { - return "", "", err + return nil, "", err + } + fileName := fmt.Sprintf("pages-deployment-%d.zip", deployment.ID) + if deployment.UploadID > 0 { + uploadRecord, err := repository.GetActiveUploadByID(ctx, deployment.UploadID) + if err != nil { + return nil, "", errors.New(errPagesPackageUploadMissing) + } + obj, err := uploadstorage.OpenStoredObject(ctx, &uploadRecord) + if err != nil { + return nil, "", fmt.Errorf("Pages 部署包不存在: %w", err) + } + if obj.ContentType == "" { + obj.ContentType = "application/zip" + } + return obj, fileName, nil } if strings.TrimSpace(deployment.ArtifactPath) == "" { - return "", "", errors.New(errPagesPackagePathEmpty) + return nil, "", errors.New(errPagesPackagePathEmpty) } - if _, err = os.Stat(deployment.ArtifactPath); err != nil { - return "", "", fmt.Errorf("Pages 部署包不存在: %w", err) + file, err := os.Open(deployment.ArtifactPath) + if err != nil { + return nil, "", fmt.Errorf("Pages 部署包不存在: %w", err) } - return deployment.ArtifactPath, fmt.Sprintf("pages-deployment-%d.zip", deployment.ID), nil + info, err := file.Stat() + if err != nil { + _ = file.Close() + return nil, "", fmt.Errorf("Pages 部署包不存在: %w", err) + } + return &storage.Object{ + Body: file, + ContentLength: info.Size(), + ContentType: "application/zip", + }, fileName, nil } func ensureDeploymentInActiveSnapshot(ctx context.Context, deploymentID uint) error { @@ -436,7 +472,7 @@ func DeleteDeployment(ctx context.Context, projectID uint, deploymentID uint) er if err := tx.Delete(deployment).Error; err != nil { return err } - _ = os.Remove(deployment.ArtifactPath) + removeDeploymentArtifact(ctx, deployment) return nil }) } diff --git a/Wavelet/internal/apps/openflare/pages/logics_test.go b/Wavelet/internal/apps/openflare/pages/logics_test.go index f9791993..97b5bb50 100644 --- a/Wavelet/internal/apps/openflare/pages/logics_test.go +++ b/Wavelet/internal/apps/openflare/pages/logics_test.go @@ -8,13 +8,15 @@ import ( "bytes" "context" "fmt" + "io" "mime/multipart" "net/http/httptest" - "strconv" + "os" "testing" "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/storage" "github.com/glebarez/sqlite" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -29,11 +31,22 @@ func setupPagesTestDB(t *testing.T) func() { }) require.NoError(t, err) require.NoError(t, sqliteDB.AutoMigrate( + &model.User{}, + &model.Upload{}, + &model.UploadStat{}, + &model.TaskExecution{}, &model.PagesProject{}, &model.PagesDeployment{}, &model.PagesDeploymentFile{}, &model.ConfigVersion{}, )) + require.NoError(t, sqliteDB.Create(&model.User{ + ID: 999, + Username: "system", + Password: "*", + Nickname: "系统", + IsActive: true, + }).Error) db.SetDB(sqliteDB) return func() { @@ -41,6 +54,44 @@ func setupPagesTestDB(t *testing.T) func() { } } +func setupPagesStorageMock(t *testing.T) (restore func(), disable func()) { + t.Helper() + mockFiles := make(map[string][]byte) + restore = storage.MockStorage( + func(_ context.Context, key string, body io.Reader, _ int64, _ string) error { + data, err := io.ReadAll(body) + if err != nil { + return err + } + mockFiles[key] = data + return nil + }, + func(_ context.Context, key string) (*storage.Object, error) { + data, ok := mockFiles[key] + if !ok { + return nil, os.ErrNotExist + } + return &storage.Object{ + Body: io.NopCloser(bytes.NewReader(data)), + ContentLength: int64(len(data)), + ContentType: "application/zip", + }, nil + }, + func(_ context.Context, key string) error { + delete(mockFiles, key) + return nil + }, + ) + storage.IsEnabledFunc = func() bool { return true } + storage.ResetCache() + disable = func() { + storage.IsEnabledFunc = func() bool { return false } + storage.ResetCache() + restore() + } + return restore, disable +} + func TestCreateProject(t *testing.T) { cleanup := setupPagesTestDB(t) defer cleanup() @@ -90,9 +141,40 @@ func TestCreateProjectRejectsUnsafeFallbackPath(t *testing.T) { assert.Contains(t, err.Error(), "回退路径") } -func TestGetDeploymentPackagePathRequiresActiveConfigSnapshot(t *testing.T) { +func TestUploadDeploymentStoresPackageInUploadFramework(t *testing.T) { cleanup := setupPagesTestDB(t) defer cleanup() + _, disableStorage := setupPagesStorageMock(t) + defer disableStorage() + ctx := context.Background() + + project, err := CreateProject(ctx, Input{ + Name: "Upload Framework Site", + Slug: "upload-framework-site", + Enabled: true, + }) + require.NoError(t, err) + + deployment, err := UploadDeployment(ctx, project.ID, testPagesMultipartFile(t, "site.zip", testPagesZip(t, map[string]string{ + "index.html": "ok", + })), "root") + require.NoError(t, err) + + storedDeployment, err := model.GetPagesDeploymentByID(ctx, deployment.ID) + require.NoError(t, err) + assert.NotZero(t, storedDeployment.UploadID) + assert.Empty(t, storedDeployment.ArtifactPath) + + var uploadCount int64 + require.NoError(t, db.DB(ctx).Model(&model.Upload{}).Count(&uploadCount).Error) + assert.Equal(t, int64(1), uploadCount) +} + +func TestOpenDeploymentPackageRequiresActiveConfigSnapshot(t *testing.T) { + cleanup := setupPagesTestDB(t) + defer cleanup() + _, disableStorage := setupPagesStorageMock(t) + defer disableStorage() ctx := context.Background() project, err := CreateProject(ctx, Input{ @@ -110,7 +192,7 @@ func TestGetDeploymentPackagePathRequiresActiveConfigSnapshot(t *testing.T) { _, err = ActivateDeployment(ctx, project.ID, deployment.ID) require.NoError(t, err) - _, _, err = GetDeploymentPackagePath(ctx, deployment.ID) + _, _, err = OpenDeploymentPackage(ctx, deployment.ID) require.Error(t, err) assert.Contains(t, err.Error(), "激活配置") @@ -125,10 +207,17 @@ func TestGetDeploymentPackagePathRequiresActiveConfigSnapshot(t *testing.T) { CreatedBy: "test", }).Error) - filePath, fileName, err := GetDeploymentPackagePath(ctx, deployment.ID) + packageObj, fileName, err := OpenDeploymentPackage(ctx, deployment.ID) require.NoError(t, err) - assert.NotEmpty(t, filePath) - assert.Equal(t, "pages-deployment-"+strconv.FormatUint(uint64(deployment.ID), 10)+".zip", fileName) + defer packageObj.Body.Close() + assert.Equal(t, fmt.Sprintf("pages-deployment-%d.zip", deployment.ID), fileName) + + body, err := io.ReadAll(packageObj.Body) + require.NoError(t, err) + reader, err := zip.NewReader(bytes.NewReader(body), int64(len(body))) + require.NoError(t, err) + require.Len(t, reader.File, 1) + assert.Equal(t, "index.html", reader.File[0].Name) } func testPagesZip(t *testing.T, files map[string]string) []byte { @@ -165,4 +254,4 @@ func testPagesMultipartFile(t *testing.T, fileName string, content []byte) *mult require.NoError(t, err) file.Close() return header -} +} \ No newline at end of file diff --git a/Wavelet/internal/db/migrator/goose/postgres/202606190014_add_pages_deployment_upload_id.sql b/Wavelet/internal/db/migrator/goose/postgres/202606190014_add_pages_deployment_upload_id.sql new file mode 100644 index 00000000..8e791795 --- /dev/null +++ b/Wavelet/internal/db/migrator/goose/postgres/202606190014_add_pages_deployment_upload_id.sql @@ -0,0 +1,14 @@ +-- +goose Up +ALTER TABLE of_pages_deployments + ADD COLUMN IF NOT EXISTS upload_id BIGINT NOT NULL DEFAULT 0; + +CREATE INDEX IF NOT EXISTS idx_of_pages_deployments_upload_id ON of_pages_deployments (upload_id); + +ALTER TABLE of_pages_deployments + ALTER COLUMN artifact_path SET DEFAULT ''; + +-- +goose Down +DROP INDEX IF EXISTS idx_of_pages_deployments_upload_id; + +ALTER TABLE of_pages_deployments + DROP COLUMN IF EXISTS upload_id; \ No newline at end of file diff --git a/Wavelet/internal/db/migrator/goose/sqlite/202606190014_add_pages_deployment_upload_id.sql b/Wavelet/internal/db/migrator/goose/sqlite/202606190014_add_pages_deployment_upload_id.sql new file mode 100644 index 00000000..6804478c --- /dev/null +++ b/Wavelet/internal/db/migrator/goose/sqlite/202606190014_add_pages_deployment_upload_id.sql @@ -0,0 +1,10 @@ +-- +goose Up +ALTER TABLE of_pages_deployments + ADD COLUMN upload_id INTEGER NOT NULL DEFAULT 0; + +CREATE INDEX IF NOT EXISTS idx_of_pages_deployments_upload_id ON of_pages_deployments (upload_id); + +-- +goose Down +DROP INDEX IF EXISTS idx_of_pages_deployments_upload_id; + +-- SQLite cannot drop columns without table rebuild; keep upload_id on rollback. \ No newline at end of file diff --git a/Wavelet/internal/model/openflare_pages.go b/Wavelet/internal/model/openflare_pages.go index 3243324d..1643037b 100644 --- a/Wavelet/internal/model/openflare_pages.go +++ b/Wavelet/internal/model/openflare_pages.go @@ -47,7 +47,8 @@ type PagesDeployment struct { DeploymentNumber int `json:"deployment_number" gorm:"not null"` Checksum string `json:"checksum" gorm:"size:64;not null;index"` Status string `json:"status" gorm:"size:32;not null;default:'uploaded';index"` - ArtifactPath string `json:"artifact_path" gorm:"size:2048;not null"` + UploadID uint64 `json:"upload_id,string" gorm:"not null;default:0;index"` + ArtifactPath string `json:"artifact_path,omitempty" gorm:"size:2048;not null;default:''"` // legacy only FileCount int `json:"file_count" gorm:"not null;default:0"` TotalSize int64 `json:"total_size" gorm:"not null;default:0"` CreatedBy string `json:"created_by" gorm:"size:64;not null;default:''"`