From 6ed3c0c81f260761e68edca651e4b40653c6fc9d Mon Sep 17 00:00:00 2001 From: ryan Date: Sun, 21 Jun 2026 12:04:54 +0800 Subject: [PATCH] =?UTF-8?q?=E6=94=B6=E6=95=9B=20Pages=20=E9=83=A8=E7=BD=B2?= =?UTF-8?q?=E5=8C=85=E8=AF=BB=E5=8F=96=E8=B7=AF=E5=BE=84?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/changelog/index.md | 2 +- internal/apps/openflare/agent/routers.go | 4 +- internal/apps/openflare/pages/helpers.go | 140 ++++-------------- internal/apps/openflare/pages/logics.go | 63 ++++---- internal/apps/openflare/pages/logics_test.go | 12 +- internal/apps/upload/exports.go | 23 ++- internal/apps/upload/ingest/access.go | 16 +- internal/apps/upload/ingest/local_file.go | 111 ++++++++++++++ .../apps/upload/ingest/local_file_test.go | 36 +++++ internal/apps/upload/ingest/object.go | 26 ++++ 10 files changed, 274 insertions(+), 159 deletions(-) create mode 100644 internal/apps/upload/ingest/local_file.go create mode 100644 internal/apps/upload/ingest/local_file_test.go create mode 100644 internal/apps/upload/ingest/object.go diff --git a/docs/changelog/index.md b/docs/changelog/index.md index da0cd5ba..566d32d8 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -32,7 +32,7 @@ sidebar: false - 修复 Agent 在 OpenResty 配置 checksum 已一致时跳过 Pages 部署包下载,导致 `deployments/{id}/releases` 为空、站点文件未落地:在 state 中缓存 Pages 部署引用;周期同步通过 `GET /api/v1/agent/pages/deployments/:id/hash` 对比 upload SHA-256,仅在哈希变化或本地 release 未就绪时下载 ZIP,避免重复拉取完整配置与部署包。 -- 收敛 Pages 部署包读取路径:`upload` 域新增 `GetActiveUpload` / `OpenStoredUpload` / `ActiveUploadHash` 门面,Pages 业务不再直接调用 `repository.GetActiveUploadByID` 或 `upload/storage.OpenStoredObject`。 +- 收敛 Pages 部署包读取路径:`upload` 域新增 `GetActiveUpload` / `OpenStoredUpload` / `ActiveUploadHash` / `ResolveLocalFile` / `IngestFromLocalPath` 门面;遗留 `artifact_path` 回填与本地路径解析迁入 upload 域;`OpenDeploymentPackage` 改为返回 `DeploymentPackage`(`io.ReadCloser` + 元数据),不再向 Agent Handler 泄漏 `storage.Object`。 - 修复节点详情 OpenResty 连接数与吞吐显示为「—」:节点可观测 API 将 OpenResty 观测数据合并进 `metric_snapshots`;指标文案改为「请求/分钟」(近 60 秒窗口),连接数为 0 时正常显示 0。 diff --git a/internal/apps/openflare/agent/routers.go b/internal/apps/openflare/agent/routers.go index 54ddd1cc..f6433208 100644 --- a/internal/apps/openflare/agent/routers.go +++ b/internal/apps/openflare/agent/routers.go @@ -196,12 +196,12 @@ func DownloadPagesPackageHandler(c *gin.Context) { if !ok { return } - packageObj, fileName, err := pages.OpenDeploymentPackage(c.Request.Context(), deploymentID) + packageObj, err := pages.OpenDeploymentPackage(c.Request.Context(), deploymentID) if apiutil.AbortBadRequestOnError(c, err) { return } defer func() { _ = packageObj.Body.Close() }() - c.Header("Content-Disposition", "attachment; filename="+fileName) + c.Header("Content-Disposition", "attachment; filename="+packageObj.FileName) if packageObj.ContentType != "" { c.Header("Content-Type", packageObj.ContentType) } diff --git a/internal/apps/openflare/pages/helpers.go b/internal/apps/openflare/pages/helpers.go index 5306c0b9..086178bd 100644 --- a/internal/apps/openflare/pages/helpers.go +++ b/internal/apps/openflare/pages/helpers.go @@ -22,20 +22,17 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/upload" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" - platformstorage "github.com/Rain-kl/Wavelet/internal/storage" ) const ( - pagesLegacyArtifactCandidateCapacity = 8 - pagesLegacyArtifactRootCapacity = 4 - pagesMaxDeploymentFiles = 1000 - pagesMaxDeploymentBytes = 100 * 1024 * 1024 - defaultPagesEntryFile = "index.html" - defaultPagesFallbackPath = "/index.html" - pagesDeploymentUploadType = "openflare_pages_deployment" - mimeTypeApplicationZip = "application/zip" - pagesMaxPathLength = 512 - bytesPerKiB = 1024 + pagesMaxDeploymentFiles = 1000 + pagesMaxDeploymentBytes = 100 * 1024 * 1024 + defaultPagesEntryFile = "index.html" + defaultPagesFallbackPath = "/index.html" + pagesDeploymentUploadType = "openflare_pages_deployment" + mimeTypeApplicationZip = "application/zip" + pagesMaxPathLength = 512 + bytesPerKiB = 1024 ) var pagesSlugPattern = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,126}[a-z0-9]$|^[a-z0-9]$`) @@ -179,26 +176,34 @@ func persistPagesUploadTemp(fileHeader *multipart.FileHeader) (string, string, i return temp.Name(), hex.EncodeToString(hash.Sum(nil)), written, nil } +func pagesLegacyRelativeCandidates(project *model.PagesProject, deployment *model.PagesDeployment) []string { + if project == nil || deployment == nil { + return nil + } + slug := strings.TrimSpace(project.Slug) + checksum := strings.TrimSpace(deployment.Checksum) + if slug == "" || checksum == "" { + return nil + } + fileName := checksum + ".zip" + return []string{ + filepath.Join("artifacts", slug, fileName), + filepath.Join("pages", "artifacts", slug, fileName), + filepath.Join("data", "pages", "artifacts", slug, fileName), + } +} + func ingestPagesDeploymentPackage( ctx context.Context, - tempPath string, + localPath string, checksum string, - size int64, projectSlug string, fileName string, ) (upload.IngestResult, error) { - file, err := os.Open(tempPath) //nolint:gosec // tempPath is a validated pages deployment staging file - if err != nil { - return upload.IngestResult{}, err - } - defer func() { _ = file.Close() }() - systemUser := repository.GetSystemUser(ctx) accessMode := 0 - return upload.Ingest(ctx, upload.IngestRequest{ + return upload.IngestFromLocalPath(ctx, localPath, upload.IngestRequest{ UserID: systemUser.ID, - Reader: file, - Size: size, FileName: fileName, MimeType: mimeTypeApplicationZip, Extension: "zip", @@ -215,101 +220,12 @@ func ingestPagesDeploymentPackage( }) } -func legacyArtifactCandidatePaths(ctx context.Context, project *model.PagesProject, deployment *model.PagesDeployment) []string { - if deployment == nil { - return nil - } - - seen := make(map[string]struct{}) - candidates := make([]string, 0, pagesLegacyArtifactCandidateCapacity) - add := func(raw string) { - value := strings.TrimSpace(raw) - if value == "" { - return - } - if _, ok := seen[value]; ok { - return - } - info, statErr := os.Stat(value) - if statErr != nil || info.IsDir() { - return - } - seen[value] = struct{}{} - candidates = append(candidates, value) - } - - storedPath := strings.TrimSpace(deployment.ArtifactPath) - add(storedPath) - if storedPath != "" { - add(filepath.Clean(storedPath)) - add(strings.ReplaceAll(storedPath, "/data/data/", "/data/")) - add(strings.ReplaceAll(filepath.Clean(storedPath), string(filepath.Separator)+string(filepath.Separator), string(filepath.Separator))) - } - - slug := "" - if project != nil { - slug = strings.TrimSpace(project.Slug) - } - checksum := strings.TrimSpace(deployment.Checksum) - if slug != "" && checksum != "" { - add(filepath.Join("artifacts", slug, checksum+".zip")) - add(filepath.Join("pages", "artifacts", slug, checksum+".zip")) - add(filepath.Join("data", "pages", "artifacts", slug, checksum+".zip")) - } - - roots := make([]string, 0, pagesLegacyArtifactRootCapacity) - if cfg, err := platformstorage.LoadConfig(ctx); err == nil { - root := strings.TrimSpace(cfg.Local.Root) - if root != "" { - roots = append(roots, root) - } - } - for _, root := range roots { - if storedPath != "" && !filepath.IsAbs(storedPath) { - add(filepath.Join(root, storedPath)) - } - if slug != "" && checksum != "" { - add(filepath.Join(root, "artifacts", slug, checksum+".zip")) - add(filepath.Join(root, "pages", "artifacts", slug, checksum+".zip")) - add(filepath.Join(root, "data", "pages", "artifacts", slug, checksum+".zip")) - } - } - - return candidates -} - -func openLegacyDeploymentArtifact(ctx context.Context, project *model.PagesProject, deployment *model.PagesDeployment) (string, *os.File, os.FileInfo, error) { - for _, candidate := range legacyArtifactCandidatePaths(ctx, project, deployment) { - file, err := os.Open(candidate) //nolint:gosec // candidate is resolved from managed legacy artifact metadata - if err != nil { - continue - } - info, statErr := file.Stat() - if statErr != nil { - _ = file.Close() - continue - } - if info.IsDir() { - _ = file.Close() - continue - } - return candidate, file, info, nil - } - return "", nil, nil, os.ErrNotExist -} - 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) + _, _ = upload.Remove(ctx, deployment.UploadID) } } diff --git a/internal/apps/openflare/pages/logics.go b/internal/apps/openflare/pages/logics.go index f7c80827..6011241b 100644 --- a/internal/apps/openflare/pages/logics.go +++ b/internal/apps/openflare/pages/logics.go @@ -8,6 +8,7 @@ import ( "encoding/json" "errors" "fmt" + "io" "mime/multipart" "net/url" "os" @@ -18,10 +19,17 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/upload" "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/model" - "github.com/Rain-kl/Wavelet/internal/storage" "gorm.io/gorm" ) +// DeploymentPackage is a streamable Pages deployment artifact for agent download. +type DeploymentPackage struct { + FileName string + ContentType string + ContentLength int64 + Body io.ReadCloser +} + // Input Pages 项目创建/更新请求。 type Input struct { Name string `json:"name"` @@ -251,7 +259,7 @@ func UploadDeployment(ctx context.Context, projectID uint, fileHeader *multipart return nil, err } entryFile := normalizePagesEntryFile(project.EntryFile) - tempPath, checksum, packageSize, err := persistPagesUploadTemp(fileHeader) + tempPath, checksum, _, err := persistPagesUploadTemp(fileHeader) if err != nil { return nil, err } @@ -264,7 +272,6 @@ func UploadDeployment(ctx context.Context, projectID uint, fileHeader *multipart ctx, tempPath, checksum, - packageSize, project.Slug, fileHeader.Filename, ) @@ -392,32 +399,37 @@ func GetDeploymentPackageHash(ctx context.Context, deploymentID uint) (string, e } // OpenDeploymentPackage opens the deployment artifact from the upload storage framework. -func OpenDeploymentPackage(ctx context.Context, deploymentID uint) (*storage.Object, string, error) { +func OpenDeploymentPackage(ctx context.Context, deploymentID uint) (DeploymentPackage, error) { deployment, err := model.GetPagesDeploymentByID(ctx, deploymentID) if err != nil { - return nil, "", err + return DeploymentPackage{}, err } if err = ensureDeploymentInActiveSnapshot(ctx, deployment.ID); err != nil { - return nil, "", err + return DeploymentPackage{}, err } - fileName := fmt.Sprintf("pages-deployment-%d.zip", deployment.ID) if deployment.UploadID == 0 { if err := ensureDeploymentUploadRecord(ctx, deployment); err != nil { - return nil, "", err + return DeploymentPackage{}, err } } - return openDeploymentPackageFromUpload(ctx, deployment, fileName) + return openDeploymentPackageFromUpload(ctx, deployment.UploadID, deployment.ID) } -func openDeploymentPackageFromUpload(ctx context.Context, deployment *model.PagesDeployment, fileName string) (*storage.Object, string, error) { - obj, _, err := upload.OpenStoredUpload(ctx, deployment.UploadID) +func openDeploymentPackageFromUpload(ctx context.Context, uploadID uint64, deploymentID uint) (DeploymentPackage, error) { + opened, err := upload.OpenStoredUpload(ctx, uploadID) if err != nil { - return nil, "", fmt.Errorf("pages 部署包不存在: %w", err) + return DeploymentPackage{}, fmt.Errorf("pages 部署包不存在: %w", err) } - if obj.ContentType == "" { - obj.ContentType = mimeTypeApplicationZip + contentType := opened.ContentType + if contentType == "" { + contentType = mimeTypeApplicationZip } - return obj, fileName, nil + return DeploymentPackage{ + FileName: fmt.Sprintf("pages-deployment-%d.zip", deploymentID), + ContentType: contentType, + ContentLength: opened.ContentLength, + Body: opened.Body, + }, nil } func ensureDeploymentUploadRecord(ctx context.Context, deployment *model.PagesDeployment) error { @@ -431,11 +443,7 @@ func ensureDeploymentUploadRecord(ctx context.Context, deployment *model.PagesDe if err != nil { return err } - artifactPath, _, info, err := openLegacyDeploymentArtifact(ctx, project, deployment) - if err != nil { - return fmt.Errorf("pages 部署包不存在: %w", err) - } - if _, err := hydrateLegacyDeploymentUpload(ctx, deployment, project, artifactPath, info.Size()); err != nil { + if _, err := hydrateLegacyDeploymentUpload(ctx, deployment, project); err != nil { return err } return nil @@ -445,11 +453,9 @@ func hydrateLegacyDeploymentUpload( ctx context.Context, deployment *model.PagesDeployment, project *model.PagesProject, - artifactPath string, - size int64, ) (*model.Upload, error) { - if deployment == nil || project == nil || strings.TrimSpace(artifactPath) == "" { - return nil, errors.New(errPagesPackagePathEmpty) + if deployment == nil || project == nil { + return nil, errors.New(errPagesDeploymentNotFound) } if deployment.UploadID > 0 { uploadRecord, err := upload.GetActiveUpload(ctx, deployment.UploadID) @@ -458,11 +464,18 @@ func hydrateLegacyDeploymentUpload( } } + artifactPath, _, err := upload.ResolveLocalFile(ctx, upload.LocalFileCandidateRequest{ + StoredPath: deployment.ArtifactPath, + RelativePaths: pagesLegacyRelativeCandidates(project, deployment), + }) + if err != nil { + return nil, fmt.Errorf("pages 部署包不存在: %w", err) + } + ingestResult, err := ingestPagesDeploymentPackage( ctx, artifactPath, deployment.Checksum, - size, project.Slug, fmt.Sprintf("pages-deployment-%d.zip", deployment.ID), ) diff --git a/internal/apps/openflare/pages/logics_test.go b/internal/apps/openflare/pages/logics_test.go index 3d8f96ba..ccef6a54 100644 --- a/internal/apps/openflare/pages/logics_test.go +++ b/internal/apps/openflare/pages/logics_test.go @@ -245,10 +245,10 @@ func TestOpenDeploymentPackageHydratesLegacyArtifactPath(t *testing.T) { CreatedBy: "test", }).Error) - packageObj, fileName, err := OpenDeploymentPackage(ctx, deployment.ID) + packageObj, err := OpenDeploymentPackage(ctx, deployment.ID) require.NoError(t, err) defer packageObj.Body.Close() - assert.Equal(t, fmt.Sprintf("pages-deployment-%d.zip", deployment.ID), fileName) + assert.Equal(t, fmt.Sprintf("pages-deployment-%d.zip", deployment.ID), packageObj.FileName) body, err := io.ReadAll(packageObj.Body) require.NoError(t, err) @@ -266,7 +266,7 @@ func TestOpenDeploymentPackageHydratesLegacyArtifactPath(t *testing.T) { require.NoError(t, db.DB(ctx).Model(&model.Upload{}).Count(&uploadCount).Error) assert.Equal(t, int64(1), uploadCount) - packageObj2, _, err := OpenDeploymentPackage(ctx, deployment.ID) + packageObj2, err := OpenDeploymentPackage(ctx, deployment.ID) require.NoError(t, err) defer packageObj2.Body.Close() body2, err := io.ReadAll(packageObj2.Body) @@ -296,7 +296,7 @@ func TestOpenDeploymentPackageRequiresActiveConfigSnapshot(t *testing.T) { _, err = ActivateDeployment(ctx, project.ID, deployment.ID) require.NoError(t, err) - _, _, err = OpenDeploymentPackage(ctx, deployment.ID) + _, err = OpenDeploymentPackage(ctx, deployment.ID) require.Error(t, err) assert.Contains(t, err.Error(), "激活配置") @@ -311,10 +311,10 @@ func TestOpenDeploymentPackageRequiresActiveConfigSnapshot(t *testing.T) { CreatedBy: "test", }).Error) - packageObj, fileName, err := OpenDeploymentPackage(ctx, deployment.ID) + packageObj, err := OpenDeploymentPackage(ctx, deployment.ID) require.NoError(t, err) defer packageObj.Body.Close() - assert.Equal(t, fmt.Sprintf("pages-deployment-%d.zip", deployment.ID), fileName) + assert.Equal(t, fmt.Sprintf("pages-deployment-%d.zip", deployment.ID), packageObj.FileName) body, err := io.ReadAll(packageObj.Body) require.NoError(t, err) diff --git a/internal/apps/upload/exports.go b/internal/apps/upload/exports.go index 206d1cce..dd60af22 100644 --- a/internal/apps/upload/exports.go +++ b/internal/apps/upload/exports.go @@ -31,13 +31,22 @@ var ( // Programmatic ingest API var ( - Ingest = ingest.Ingest - Remove = ingest.Remove - RemoveOwned = ingest.RemoveOwned - FindByHash = ingest.FindByHash - GetActiveUpload = ingest.GetActive - OpenStoredUpload = ingest.OpenActive - ActiveUploadHash = ingest.ActiveHash + Ingest = ingest.Ingest + Remove = ingest.Remove + RemoveOwned = ingest.RemoveOwned + FindByHash = ingest.FindByHash + GetActiveUpload = ingest.GetActive + OpenStoredUpload = ingest.OpenActiveObject + ActiveUploadHash = ingest.ActiveHash + ResolveLocalFile = ingest.ResolveLocalFile + IngestFromLocalPath = ingest.IngestFromLocalPath +) + +type ( + // OpenedUploadObject is the upload-domain view of a stored object stream. + OpenedUploadObject = ingest.OpenedObject + // LocalFileCandidateRequest describes filesystem locations that may host a legacy blob. + LocalFileCandidateRequest = ingest.LocalFileCandidateRequest ) // Ingest policy constants diff --git a/internal/apps/upload/ingest/access.go b/internal/apps/upload/ingest/access.go index 8b5ee0d5..e7b8944f 100644 --- a/internal/apps/upload/ingest/access.go +++ b/internal/apps/upload/ingest/access.go @@ -11,7 +11,6 @@ import ( uploadcache "github.com/Rain-kl/Wavelet/internal/apps/upload/cache" uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage" "github.com/Rain-kl/Wavelet/internal/model" - "github.com/Rain-kl/Wavelet/internal/storage" ) // GetActive loads an active upload record by ID through the upload metadata cache path. @@ -22,17 +21,22 @@ func GetActive(ctx context.Context, uploadID uint64) (model.Upload, error) { return uploadcache.GetUploadByID(ctx, uploadID) } -// OpenActive opens the stored object for an active upload record. -func OpenActive(ctx context.Context, uploadID uint64) (*storage.Object, model.Upload, error) { +// OpenActiveObject opens the stored object for an active upload record. +func OpenActiveObject(ctx context.Context, uploadID uint64) (OpenedObject, error) { record, err := GetActive(ctx, uploadID) if err != nil { - return nil, model.Upload{}, err + return OpenedObject{}, err } obj, err := uploadstorage.OpenStoredObject(ctx, &record) if err != nil { - return nil, model.Upload{}, err + return OpenedObject{}, err } - return obj, record, nil + return OpenedObject{ + Body: obj.Body, + ContentType: obj.ContentType, + ContentLength: obj.ContentLength, + Upload: record, + }, nil } // ActiveHash returns the SHA-256 hash recorded for an active upload. diff --git a/internal/apps/upload/ingest/local_file.go b/internal/apps/upload/ingest/local_file.go new file mode 100644 index 00000000..f14173fd --- /dev/null +++ b/internal/apps/upload/ingest/local_file.go @@ -0,0 +1,111 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package ingest + +import ( + "context" + "errors" + "os" + "path/filepath" + "strings" + + "github.com/Rain-kl/Wavelet/internal/storage" +) + +// LocalFileCandidateRequest describes filesystem locations that may host a legacy blob. +type LocalFileCandidateRequest struct { + StoredPath string + RelativePaths []string +} + +// ResolveLocalFile returns the first existing regular file among managed candidate paths. +func ResolveLocalFile(ctx context.Context, req LocalFileCandidateRequest) (string, int64, error) { + for _, candidate := range buildLocalFileCandidates(ctx, req) { + info, err := os.Stat(candidate) //nolint:gosec // candidate is resolved from managed legacy metadata + if err != nil || info.IsDir() { + continue + } + return candidate, info.Size(), nil + } + return "", 0, os.ErrNotExist +} + +// IngestFromLocalPath ingests a local regular file through the standard upload ingest path. +func IngestFromLocalPath(ctx context.Context, localPath string, req Request) (Result, error) { + localPath = strings.TrimSpace(localPath) + if localPath == "" { + return Result{}, errors.New("local path is required") + } + file, err := os.Open(localPath) //nolint:gosec // localPath is resolved from managed legacy metadata + if err != nil { + return Result{}, err + } + defer func() { _ = file.Close() }() + + info, err := file.Stat() + if err != nil { + return Result{}, err + } + if info.IsDir() { + return Result{}, errors.New("local path must be a regular file") + } + if req.Size <= 0 { + req.Size = info.Size() + } + req.Reader = file + return Ingest(ctx, req) +} + +func buildLocalFileCandidates(ctx context.Context, req LocalFileCandidateRequest) []string { + seen := make(map[string]struct{}) + candidates := make([]string, 0, 8+len(req.RelativePaths)*4) + add := func(raw string) { + value := strings.TrimSpace(raw) + if value == "" { + return + } + if _, ok := seen[value]; ok { + return + } + seen[value] = struct{}{} + candidates = append(candidates, value) + } + + storedPath := strings.TrimSpace(req.StoredPath) + add(storedPath) + if storedPath != "" { + add(filepath.Clean(storedPath)) + add(strings.ReplaceAll(storedPath, "/data/data/", "/data/")) + add(strings.ReplaceAll( + filepath.Clean(storedPath), + string(filepath.Separator)+string(filepath.Separator), + string(filepath.Separator), + )) + } + for _, relativePath := range req.RelativePaths { + add(relativePath) + } + + for _, root := range localStorageRoots(ctx) { + if storedPath != "" && !filepath.IsAbs(storedPath) { + add(filepath.Join(root, storedPath)) + } + for _, relativePath := range req.RelativePaths { + add(filepath.Join(root, relativePath)) + } + } + return candidates +} + +func localStorageRoots(ctx context.Context) []string { + cfg, err := storage.LoadConfig(ctx) + if err != nil { + return nil + } + root := strings.TrimSpace(cfg.Local.Root) + if root == "" { + return nil + } + return []string{root} +} diff --git a/internal/apps/upload/ingest/local_file_test.go b/internal/apps/upload/ingest/local_file_test.go new file mode 100644 index 00000000..08fbd1ae --- /dev/null +++ b/internal/apps/upload/ingest/local_file_test.go @@ -0,0 +1,36 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package ingest + +import ( + "context" + "os" + "path/filepath" + "testing" +) + +func TestResolveLocalFileFindsStoredPath(t *testing.T) { + dir := t.TempDir() + artifactPath := filepath.Join(dir, "legacy.zip") + if err := os.WriteFile(artifactPath, []byte("legacy"), 0o644); err != nil { + t.Fatalf("write artifact: %v", err) + } + + path, size, err := ResolveLocalFile(context.Background(), LocalFileCandidateRequest{ + StoredPath: artifactPath, + }) + if err != nil { + t.Fatalf("ResolveLocalFile failed: %v", err) + } + if path != artifactPath || size != int64(len("legacy")) { + t.Fatalf("unexpected resolve result: path=%q size=%d", path, size) + } +} + +func TestIngestFromLocalPathRequiresPath(t *testing.T) { + _, err := IngestFromLocalPath(context.Background(), "", Request{}) + if err == nil { + t.Fatal("expected error for empty local path") + } +} diff --git a/internal/apps/upload/ingest/object.go b/internal/apps/upload/ingest/object.go new file mode 100644 index 00000000..02fed0a1 --- /dev/null +++ b/internal/apps/upload/ingest/object.go @@ -0,0 +1,26 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package ingest + +import ( + "io" + + "github.com/Rain-kl/Wavelet/internal/model" +) + +// OpenedObject is the upload-domain view of a stored object stream. +type OpenedObject struct { + Body io.ReadCloser + ContentType string + ContentLength int64 + Upload model.Upload +} + +// Close closes the object body when present. +func (o OpenedObject) Close() error { + if o.Body == nil { + return nil + } + return o.Body.Close() +}