feat(openflare): Pages 部署包改用 upload 本地文件存储框架

通过 upload.Ingest 摄取 zip 部署包并记录 upload_id,删除部署时调用
upload.Remove;Agent 下载改为 OpenDeploymentPackage 从存储后端流式输出。
新增 of_pages_deployments.upload_id 迁移,保留 artifact_path 兼容旧数据。
This commit is contained in:
ryan
2026-06-18 19:52:59 +08:00
parent 9ac5ff6925
commit 29176da35f
8 changed files with 247 additions and 77 deletions
@@ -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) {
@@ -21,6 +21,7 @@ const (
errPagesAPIProxyPassRequired = "启用 API 反代时,后端服务地址不能为空"
errPagesAPIProxyPassInvalid = "API 反代后端服务地址必须是有效的 HTTP/HTTPS URL"
errPagesPackagePathEmpty = "Pages 部署包路径为空"
errPagesPackageUploadMissing = "Pages 部署包上传记录不存在"
errPagesPackageNotInActiveConfig = "Pages 部署尚未进入激活配置"
errPagesInvalidSnapshotFormat = "配置快照格式无效"
)
@@ -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()
}
+57 -21
View File
@@ -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
})
}
@@ -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
}
}
@@ -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;
@@ -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.
+2 -1
View File
@@ -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:''"`