diff --git a/docs/changelog/index.md b/docs/changelog/index.md index ecd99ece..b9ac2460 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -32,7 +32,8 @@ sidebar: false - 上传新部署后会按保留策略自动清理超出数量的历史部署包,减少磁盘占用。 - Agent 在同步 Pages 部署时信任控制面已完成的包校验,不再重复限制文件数与展开体积,仅校验下载完整性并安全解压到本地。 - 明确 Pages 历史部署保留语义为每个项目最多 N 条(激活必留、其余按从新到旧填充),裁剪失败会记录日志且不回滚已成功上传。 -- 主配置版本与 Pages 部署双轨管理:下发给 Agent 的配置会将 Pages 路由重绑定到项目当前激活部署,主配置回滚不再依赖已裁剪的旧部署包。 +- 主配置版本与 Pages 部署双轨管理:Agent 以 Pages 项目 ID 为锚点请求最新激活部署包,项目内切换激活版本无需重新发布主配置即可在边缘自动更新。 +- Agent 边缘仅保留每个 Pages 项目的最新部署包,新包就绪并切换成功后清理旧 release;拉取 latest 时增加 hash 二次校验与重试以降低激活竞态失败;单个项目同步失败不再阻塞同批其它项目。 ## [v3.3.0] - 2026-07-14 diff --git a/docs/design/agent-design.md b/docs/design/agent-design.md index 3957d831..ddf250f8 100644 --- a/docs/design/agent-design.md +++ b/docs/design/agent-design.md @@ -98,7 +98,7 @@ Agent 对数据面 OpenResty 的管控实现了端到端的闭环,包含配置 * `certs/`:证书存放目录(文件命名为 `{cert_id}.crt` 和 `{cert_id}.key`)。 * `waf/` 与 `pow/`:WAF 及防 CC 挑战所需的专用 Lua 运行时脚本。 * `waf_config.json` 与 `waf_ip_groups.json`:WAF 过滤引擎所需的结构化规则配置文件。 -* `pages_dir`:Pages 静态站点部署目录,默认位于 `data_dir/var/lib/openflare/pages`。当激活配置引用 Pages 部署时,Agent 会下载部署包(zip / tar.gz / tar.xz / tar.bz2 / tar / 7z)、校验下载 checksum、解压到部署 release 目录,并切换 `deployments/{deployment_id}/current` 供 OpenResty `root`/`try_files` 读取。业务侧体积/文件数校验由控制面完成,Agent 信任控制面结果,解压时仅保留本机路径安全防护。 +* `pages_dir`:Pages 静态站点部署目录,默认位于 `data_dir/var/lib/openflare/pages`。当激活配置引用 Pages **项目**时,Agent 按 `project_id` 请求控制面「最新激活包」(hash + package,下载后再校验 hash 防竞态),解压到 `projects/{project_id}/releases/{hash}`,切换 `current` 后**立即删除同项目其它历史 release**(仅保留最新)。项目内切换激活无需重发主配置;多项目对账时单项目失败不阻塞其它项目。 ### 2. 精细化的重载动作 1. **备份当前配置**:在写入新文件之前,Agent 会将现有的配置文件复制到 `.backup` 临时目录下,保留完整的现场快照。 diff --git a/docs/design/architecture.md b/docs/design/architecture.md index 4f83a43e..a04bfceb 100644 --- a/docs/design/architecture.md +++ b/docs/design/architecture.md @@ -124,7 +124,7 @@ OpenResty (Agent, TLS/WAF) * *同步与自愈的精细时序及回滚模型详见:[Agent 与发布模型设计](./agent-design.md)* ### 2. 静态托管与 API 代理流 -* 静态资源解压落地于 Agent 节点的 `deployments/{id}/current` 下,OpenResty 通过 `root`/`index`/`try_files` 指令在边缘直接向访客提供极低延迟的静态资源服务。 +* 静态资源解压落地于 Agent 节点的 `projects/{project_id}/current` 下(按项目 latest 拉取,仅保留最新包),OpenResty 通过 `root`/`index`/`try_files` 在边缘直接提供静态资源服务。 * 当启用 API 代理时,OpenResty 自动根据站点配置的 `api_proxy_path`(如 `/api`)将 API 请求重写并转发(`proxy_pass`)给后端动态接口。 * *部署包校验、解压逃逸防御及 Nginx 规则渲染详见:[Pages 静态托管设计文档](./pages-design.md)* diff --git a/docs/design/pages-design.md b/docs/design/pages-design.md index 781c2bad..799ba6ef 100644 --- a/docs/design/pages-design.md +++ b/docs/design/pages-design.md @@ -89,11 +89,16 @@ graph TD } ``` -### 3. 与主配置版本的双轨关系 -* **主配置版本**(`config_versions`)与 **Pages 部署** 是两套独立的版本体系。 -* 快照里记录的 `pages_deployment` 仅反映**发布当时**的激活部署,用于审计与当时渲染结果存档。 -* **运行时**:下发给 Agent 的配置会把各 Pages 路由**重绑定到该项目当前激活部署**(最新包)。主配置回滚/重新激活旧版本时,**不要求**仍能拉到旧 Pages 包,只保证指向当前最新激活部署。 -* 因此 Pages 历史裁剪可以安全删除非激活部署,无需为「主配置回滚到旧 Pages 包」预留存储。 +### 3. 与主配置版本的双轨关系(项目锚点 + latest 拉取) +* **主配置版本**与 **Pages 部署** 是两套独立的版本体系。 +* 主配置中 Pages 路由的稳定锚点是 **`pages_project_id`(项目 ID)**,不是某次部署 ID。 +* OpenResty `root` 使用项目级路径:`__OPENFLARE_PAGES_DIR__/projects/{project_id}/current`,激活切换时路径不变,无需为换包而重发主配置。 +* Agent 按项目请求「最新激活包」(类似 `github/release/latest`): + * `GET /api/v1/agent/pages/projects/:project_id/latest/hash` + * `GET /api/v1/agent/pages/projects/:project_id/latest/package` + * 控制面根据该项目**当前激活部署**返回哈希与压缩包;Agent 不关心具体 deployment_id。 +* 因此:在项目内切换激活部署后,**不必发布主配置**;Agent 在周期性对账时轮询 latest hash,发现变化即下载并切换 `current`。 +* 快照中的 `pages_deployment` 字段仍可记录发布时元数据(入口文件、SPA/API 代理等),但不作为 Agent 拉包的版本锁定。 --- @@ -119,19 +124,19 @@ graph TD Agent 运行在各边缘代理节点上,在应用配置版本前,必须先将 Pages 静态资源“原子”地拉取到节点本地。 -### 1. 校验式增量拉取 -1. Agent 解析激活配置中的 `SourceConfigJSON`,检索出所有 `UpstreamType == "pages"` 的路由引用的部署 `DeploymentID` 和 `Checksum`。 -2. 检查本地部署目录是否存在正确的版本标记文件 `.openflare-pages.json`,且 `Checksum` 匹配。 -3. 若不匹配,通过专属接口 `GET /api/agent/pages/deployments/:id/package` 下载对应的部署包。下载请求头必须携带节点独有的 `X-Agent-Token` 用于 Server 鉴权。 +### 1. 按项目拉取 latest +1. Agent 从激活主配置中解析 `UpstreamType == "pages"` 的路由,收集稳定锚点 **`pages_project_id`**。 +2. 对每个项目调用 `GET /api/v1/agent/pages/projects/:project_id/latest/hash` 获取控制面当前激活包哈希(类似 latest 指针)。 +3. 若本地 `projects/{project_id}/releases/{hash}` 尚未就绪,再下载 `.../latest/package`。下载后 **再次请求 hash** 与包内容 SHA-256 对齐,避免激活切换造成的竞态;不一致则有限次重试。 +4. 请求头携带节点 `X-Agent-Token`。 -### 2. 安全解压缩与原子切换 -为了保证配置应用过程的“无缝”且能在出错时立即回滚: -1. Agent 将下载的部署包数据写入临时目录,并重新计算 SHA-256 Checksum。如果与配置指明的 checksum 不符,立即报错并阻断发布流程。 -2. 解压部署包至临时目录 `releases/{checksum}.tmp`。解压支持 zip / tar.* / 7z。**Agent 默认信任控制面**:不再重复校验文件数/体积等业务限额(控制面上传时已完成);仅做本机落盘安全处理(路径逃逸、软链接拒绝)与下载 checksum 完整性校验。 -3. 解压成功后,写入标记文件 `.openflare-pages.json`。 -4. 清理 `releases/{checksum}` 目录,将整个临时目录重命名为 `releases/{checksum}`。 -5. **原子切换**:建立拷贝当前部署的物理副本到目标位置 `deployments/{deployment_id}/current`。切换前先备份上一版本的 `current`,一旦重载配置失败,Agent 能够快速恢复 `current` 目录并回滚 OpenResty。 -6. **定时清理**:每次配置成功应用后,Agent 自动比对本地部署目录,将所有不活跃的(即未被当前激活版本引用的)历史部署包和文件夹进行物理删除,释放磁盘空间。 +### 2. 安全解压缩、原子切换与只保留最新 +1. 下载字节计算 SHA-256,须与「下载后再次查询」的 latest hash 一致。 +2. 解压至 `projects/{project_id}/releases/{hash}.tmp`(支持 zip / tar.* / 7z)。Agent 信任控制面业务校验,仅做路径逃逸/软链防护。 +3. 写入 `.openflare-pages.json` 后 rename 为 `releases/{hash}`。 +4. **原子切换** `projects/{project_id}/current` 指向新 release(优先 symlink,失败则拷贝)。 +5. **仅当新包已就绪且 current 切换成功后**,删除该项目下其它 `releases/*`(含 `.tmp`),**不保留历史部署包**。边缘节点每个项目永远只保留一份最新内容。 +6. 多项目对账时 **隔离失败**:单个项目失败记日志并继续其它项目,最后汇总返回错误。 --- @@ -141,13 +146,13 @@ Agent 运行在各边缘代理节点上,在应用配置版本前,必须先 ### 1. 静态服务指令渲染 * **`root` 与 `index`**: - Server 根据配置将 `root` 指向 Agent 的 Pages 动态目录占位符 `__OPENFLARE_PAGES_DIR__/deployments/{deployment_id}/current`,并在此基础上追加项目的 `RootDir`。`index` 指向设置的入口文件。 + Server 将 `root` 指向项目级占位路径 `__OPENFLARE_PAGES_DIR__/projects/{project_id}/current`(可再追加 `RootDir`)。激活切换只换目录内容,路径不变,无需为换包重发主配置。 ```nginx server { listen 80; server_name myapp.example.com; - root "/var/lib/openflare/pages/deployments/12/current"; + root "/var/lib/openflare/pages/projects/3/current"; index "index.html"; ... } diff --git a/internal/apps/agent/httpclient/client.go b/internal/apps/agent/httpclient/client.go index ce6a3ca0..34e1398a 100644 --- a/internal/apps/agent/httpclient/client.go +++ b/internal/apps/agent/httpclient/client.go @@ -111,6 +111,31 @@ func (c *Client) DownloadPagesDeploymentPackage(ctx context.Context, deploymentI return io.ReadAll(res.Body) } +// GetPagesProjectLatestHash returns the active deployment package hash for a Pages project. +func (c *Client) GetPagesProjectLatestHash(ctx context.Context, projectID uint) (*protocol.PagesProjectLatestHashResponse, error) { + resp := protocol.APIResponse[protocol.PagesProjectLatestHashResponse]{} + if err := c.base.GetJSON(ctx, fmt.Sprintf("/api/v1/agent/pages/projects/%d/latest/hash", projectID), &resp); err != nil { + return nil, err + } + if err := edgehttp.APIError(resp.ErrorMsg); err != nil { + return nil, err + } + return &resp.Data, nil +} + +// DownloadPagesProjectLatestPackage downloads the active deployment package for a Pages project. +func (c *Client) DownloadPagesProjectLatestPackage(ctx context.Context, projectID uint) ([]byte, error) { + res, err := c.base.DoRaw(ctx, http.MethodGet, fmt.Sprintf("/api/v1/agent/pages/projects/%d/latest/package", projectID), nil) + if err != nil { + return nil, err + } + defer func() { _ = res.Body.Close() }() + if res.StatusCode != http.StatusOK { + return nil, edgehttp.ReadHTTPError(res) + } + return io.ReadAll(res.Body) +} + // SetToken updates the authentication token used for API requests. func (c *Client) SetToken(token string) { c.base.SetToken(token) diff --git a/internal/apps/agent/protocol/alias.go b/internal/apps/agent/protocol/alias.go index ae7753de..7853b753 100644 --- a/internal/apps/agent/protocol/alias.go +++ b/internal/apps/agent/protocol/alias.go @@ -75,6 +75,9 @@ type SupportFile = pkgprotocol.SupportFile // PagesDeploymentHashResponse is an alias for pkgprotocol.PagesDeploymentHashResponse. type PagesDeploymentHashResponse = pkgprotocol.PagesDeploymentHashResponse +// PagesProjectLatestHashResponse is an alias for pkgprotocol.PagesProjectLatestHashResponse. +type PagesProjectLatestHashResponse = pkgprotocol.PagesProjectLatestHashResponse + const ( // WSMessageTypeStatus is an alias for pkgprotocol.WSMessageTypeStatus. WSMessageTypeStatus = pkgprotocol.WSMessageTypeStatus diff --git a/internal/apps/agent/state/state.go b/internal/apps/agent/state/state.go index f460a0ea..e242cc31 100644 --- a/internal/apps/agent/state/state.go +++ b/internal/apps/agent/state/state.go @@ -15,9 +15,12 @@ const ( nodeIDRandomBytes = 8 ) -// PagesDeployment records a Pages release referenced by the active config. +// PagesDeployment records a Pages project tracked by the agent and the last +// applied package hash for that project's active deployment. +// ProjectID is the stable identity; DeploymentID/Hash follow control-plane "latest". type PagesDeployment struct { - DeploymentID uint `json:"deployment_id"` + ProjectID uint `json:"project_id"` + DeploymentID uint `json:"deployment_id,omitempty"` Hash string `json:"hash"` Checksum string `json:"checksum,omitempty"` } diff --git a/internal/apps/agent/sync/pages.go b/internal/apps/agent/sync/pages.go index 07adb259..67ab219a 100644 --- a/internal/apps/agent/sync/pages.go +++ b/internal/apps/agent/sync/pages.go @@ -9,6 +9,7 @@ import ( "errors" "fmt" "io" + "log/slog" "os" "path/filepath" "strings" @@ -22,6 +23,9 @@ const ( pagesDirPerm = 0o755 pagesFilePerm = 0o644 pagesManifestFilePerm = 0o644 + // pagesLatestPullAttempts covers a race where the active deployment changes + // between the hash probe and the package download. + pagesLatestPullAttempts = 2 ) type pagesSourceDocument struct { @@ -30,16 +34,20 @@ type pagesSourceDocument struct { type pagesSourceRoute struct { UpstreamType string `json:"upstream_type"` + PagesProjectID *uint `json:"pages_project_id"` PagesDeployment *pagesDeploymentSource `json:"pages_deployment"` } -type pagesDeploymentSource struct { - DeploymentID uint `json:"deployment_id"` - Checksum string `json:"checksum"` +// pagesProjectRef is the agent-side "latest" pointer for one Pages project. +type pagesProjectRef struct { + ProjectID uint + DeploymentID uint + Checksum string } type pagesDeploymentMarker struct { - DeploymentID uint `json:"deployment_id"` + ProjectID uint `json:"project_id"` + DeploymentID uint `json:"deployment_id,omitempty"` Checksum string `json:"checksum"` } @@ -50,13 +58,19 @@ func pagesDeploymentStateHash(item state.PagesDeployment) string { return strings.TrimSpace(item.Checksum) } -func snapshotPagesDeployments(snapshot *state.Snapshot) []pagesDeploymentSource { +func snapshotPagesProjects(snapshot *state.Snapshot) []pagesProjectRef { if snapshot == nil || snapshot.PagesDeployments == nil { return nil } - result := make([]pagesDeploymentSource, 0, len(snapshot.PagesDeployments)) + result := make([]pagesProjectRef, 0, len(snapshot.PagesDeployments)) for _, item := range snapshot.PagesDeployments { - result = append(result, pagesDeploymentSource{ + projectID := item.ProjectID + if projectID == 0 { + // Legacy agent state only stored deployment_id; skip until rediscovered from config. + continue + } + result = append(result, pagesProjectRef{ + ProjectID: projectID, DeploymentID: item.DeploymentID, Checksum: pagesDeploymentStateHash(item), }) @@ -64,32 +78,34 @@ func snapshotPagesDeployments(snapshot *state.Snapshot) []pagesDeploymentSource return result } -func setSnapshotPagesDeployments(snapshot *state.Snapshot, deployments []pagesDeploymentSource) { +func setSnapshotPagesProjects(snapshot *state.Snapshot, projects []pagesProjectRef) { if snapshot == nil { return } - if len(deployments) == 0 { + if len(projects) == 0 { snapshot.PagesDeployments = []state.PagesDeployment{} return } - snapshot.PagesDeployments = make([]state.PagesDeployment, len(deployments)) - for i, deployment := range deployments { + snapshot.PagesDeployments = make([]state.PagesDeployment, len(projects)) + for i, project := range projects { snapshot.PagesDeployments[i] = state.PagesDeployment{ - DeploymentID: deployment.DeploymentID, - Hash: strings.TrimSpace(deployment.Checksum), + ProjectID: project.ProjectID, + DeploymentID: project.DeploymentID, + Hash: strings.TrimSpace(project.Checksum), } } } -func updateSnapshotPagesDeploymentHash(snapshot *state.Snapshot, deployment pagesDeploymentSource) { +func updateSnapshotPagesProject(snapshot *state.Snapshot, project pagesProjectRef) { if snapshot == nil || snapshot.PagesDeployments == nil { return } - hash := strings.TrimSpace(deployment.Checksum) + hash := strings.TrimSpace(project.Checksum) for i := range snapshot.PagesDeployments { - if snapshot.PagesDeployments[i].DeploymentID != deployment.DeploymentID { + if snapshot.PagesDeployments[i].ProjectID != project.ProjectID { continue } + snapshot.PagesDeployments[i].DeploymentID = project.DeploymentID snapshot.PagesDeployments[i].Hash = hash snapshot.PagesDeployments[i].Checksum = "" return @@ -97,11 +113,19 @@ func updateSnapshotPagesDeploymentHash(snapshot *state.Snapshot, deployment page } func pagesDiscoveryNeeded(snapshot *state.Snapshot) bool { - return snapshot == nil || snapshot.PagesDeployments == nil + if snapshot == nil || snapshot.PagesDeployments == nil { + return true + } + // Legacy state rows may only have deployment_id (project_id == 0). Those + // cannot poll latest-by-project; force a full config rediscovery. + if len(snapshot.PagesDeployments) > 0 && len(snapshotPagesProjects(snapshot)) == 0 { + return true + } + return false } func pagesSyncNeeded(snapshot *state.Snapshot) bool { - return snapshot != nil && snapshot.PagesDeployments != nil && len(snapshot.PagesDeployments) > 0 + return snapshot != nil && len(snapshotPagesProjects(snapshot)) > 0 } func pagesReconcileNeeded(snapshot *state.Snapshot) bool { @@ -112,70 +136,173 @@ func pagesReconcileNeeded(snapshot *state.Snapshot) bool { } func (s *Service) syncPagesDeployments(ctx context.Context, snapshot *state.Snapshot, config *protocol.ActiveConfigResponse) error { - var deployments []pagesDeploymentSource + var projects []pagesProjectRef var err error if config != nil { - deployments, err = referencedPagesDeployments(config) + projects, err = referencedPagesProjects(config) if err != nil { return err } - setSnapshotPagesDeployments(snapshot, deployments) + setSnapshotPagesProjects(snapshot, projects) } else { - deployments = snapshotPagesDeployments(snapshot) + projects = snapshotPagesProjects(snapshot) } - if len(deployments) == 0 { + if len(projects) == 0 { return nil } if strings.TrimSpace(s.pagesDir) == "" { - return errors.New("pages_dir is required when active config references Pages deployments") + return errors.New("pages_dir is required when active config references Pages projects") } - for _, deployment := range deployments { - if err := s.ensurePagesDeployment(ctx, snapshot, deployment); err != nil { - return err + + // Isolate per-project failures so one bad project does not block others. + var failed []error + for _, project := range projects { + if ensureErr := s.ensurePagesProject(ctx, snapshot, project.ProjectID); ensureErr != nil { + slog.Error("ensure Pages project failed", + "project_id", project.ProjectID, + "error", ensureErr, + ) + failed = append(failed, fmt.Errorf("pages project %d: %w", project.ProjectID, ensureErr)) } } if s.nginxManager != nil { - if err := s.nginxManager.EnsureWorkerReadAccess(); err != nil { - return fmt.Errorf("ensure openresty worker read access: %w", err) + if accessErr := s.nginxManager.EnsureWorkerReadAccess(); accessErr != nil { + failed = append(failed, fmt.Errorf("ensure openresty worker read access: %w", accessErr)) } } - return nil + if len(failed) == 0 { + return nil + } + return errors.Join(failed...) } -func (s *Service) ensurePagesDeployment(ctx context.Context, snapshot *state.Snapshot, deployment pagesDeploymentSource) error { - serverHash, err := s.client.GetPagesDeploymentHash(ctx, deployment.DeploymentID) - if err != nil { - return fmt.Errorf("fetch Pages deployment %d hash: %w", deployment.DeploymentID, err) +// ensurePagesProject pulls the control-plane "latest" (active) package for a +// Pages project and switches local current to that release when needed. +// Only the latest release is retained on disk; older releases are removed after +// the new release is ready and current has been switched. +func (s *Service) ensurePagesProject(ctx context.Context, snapshot *state.Snapshot, projectID uint) error { + if projectID == 0 { + return errors.New("pages project id is required") } - serverHash = strings.TrimSpace(serverHash) - if serverHash == "" { - return fmt.Errorf("pages deployment %d hash is empty", deployment.DeploymentID) - } - effective := pagesDeploymentSource{ - DeploymentID: deployment.DeploymentID, - Checksum: serverHash, - } - updateSnapshotPagesDeploymentHash(snapshot, effective) - releaseDir := pagesReleaseDir(s.pagesDir, effective.DeploymentID, effective.Checksum) - if pagesReleaseReady(releaseDir, effective) { - return switchPagesCurrentDir(s.pagesDir, effective.DeploymentID, releaseDir) + var lastErr error + for attempt := 0; attempt < pagesLatestPullAttempts; attempt++ { + latest, err := s.client.GetPagesProjectLatestHash(ctx, projectID) + if err != nil { + return fmt.Errorf("fetch Pages project %d latest hash: %w", projectID, err) + } + hash := strings.TrimSpace(latest.Hash) + if hash == "" { + return fmt.Errorf("pages project %d latest hash is empty", projectID) + } + effective := pagesProjectRef{ + ProjectID: projectID, + DeploymentID: latest.DeploymentID, + Checksum: hash, + } + + releaseDir := pagesProjectReleaseDir(s.pagesDir, projectID, hash) + if pagesProjectReleaseReady(releaseDir, effective) { + if err := switchPagesProjectCurrentDir(s.pagesDir, projectID, releaseDir); err != nil { + return err + } + updateSnapshotPagesProject(snapshot, effective) + _ = cleanupPagesProjectStaleReleases(s.pagesDir, projectID, hash) + return nil + } + + packageBytes, err := s.client.DownloadPagesProjectLatestPackage(ctx, projectID) + if err != nil { + return fmt.Errorf("download Pages project %d latest package: %w", projectID, err) + } + got := checksumBytes(packageBytes) + + // Re-probe latest after download to detect activation races. + // Accept the package only when its content hash still matches latest. + verify, err := s.client.GetPagesProjectLatestHash(ctx, projectID) + if err != nil { + return fmt.Errorf("re-fetch Pages project %d latest hash: %w", projectID, err) + } + verifyHash := strings.TrimSpace(verify.Hash) + if verifyHash == "" { + return fmt.Errorf("pages project %d latest hash is empty", projectID) + } + if got != verifyHash { + lastErr = fmt.Errorf( + "pages project %d package/hash race: downloaded %s, latest now %s (attempt %d/%d)", + projectID, got, verifyHash, attempt+1, pagesLatestPullAttempts, + ) + slog.Warn("pages latest package race, retrying", + "project_id", projectID, + "downloaded_hash", got, + "latest_hash", verifyHash, + "attempt", attempt+1, + ) + continue + } + + effective = pagesProjectRef{ + ProjectID: projectID, + DeploymentID: verify.DeploymentID, + Checksum: got, + } + releaseDir = pagesProjectReleaseDir(s.pagesDir, projectID, got) + if err := extractPagesPackage(packageBytes, releaseDir, effective); err != nil { + return err + } + if err := switchPagesProjectCurrentDir(s.pagesDir, projectID, releaseDir); err != nil { + return err + } + updateSnapshotPagesProject(snapshot, effective) + // Only after the new release is ready and current switched: drop others. + _ = cleanupPagesProjectStaleReleases(s.pagesDir, projectID, got) + return nil } - packageBytes, err := s.client.DownloadPagesDeploymentPackage(ctx, effective.DeploymentID) + if lastErr != nil { + return lastErr + } + return fmt.Errorf("pages project %d latest pull failed", projectID) +} + +// cleanupPagesProjectStaleReleases keeps only keepHash under projects/{id}/releases. +// Must be called only after the keepHash release is ready and current points at it. +func cleanupPagesProjectStaleReleases(baseDir string, projectID uint, keepHash string) error { + keepHash = strings.TrimSpace(keepHash) + if projectID == 0 || keepHash == "" { + return nil + } + releasesRoot := filepath.Join(baseDir, "projects", fmt.Sprintf("%d", projectID), "releases") + entries, err := os.ReadDir(releasesRoot) //nolint:gosec // managed PagesDir if err != nil { - return fmt.Errorf("download Pages deployment %d: %w", effective.DeploymentID, err) - } - if got := checksumBytes(packageBytes); got != effective.Checksum { - return fmt.Errorf("pages deployment %d checksum mismatch: expected %s, got %s", effective.DeploymentID, effective.Checksum, got) - } - if err := extractPagesPackage(packageBytes, releaseDir, effective); err != nil { + if os.IsNotExist(err) { + return nil + } return err } - return switchPagesCurrentDir(s.pagesDir, effective.DeploymentID, releaseDir) + var firstErr error + for _, entry := range entries { + name := entry.Name() + if name == keepHash { + continue + } + // Drop partial extract leftovers as well (*.tmp). + target := filepath.Join(releasesRoot, name) + if removeErr := os.RemoveAll(target); removeErr != nil && firstErr == nil { + firstErr = removeErr + slog.Warn("failed to remove stale Pages release", + "project_id", projectID, + "path", target, + "error", removeErr, + ) + } + } + // Also remove legacy deployments/ tree leftovers if present (best-effort). + _ = os.RemoveAll(filepath.Join(baseDir, "deployments")) + return firstErr } -func pagesReleaseReady(dir string, deployment pagesDeploymentSource) bool { - if !markerMatches(dir, deployment) { +func pagesProjectReleaseReady(dir string, project pagesProjectRef) bool { + if !markerMatches(dir, project) { return false } entries, err := os.ReadDir(dir) //nolint:gosec // dir is managed PagesDir @@ -191,7 +318,7 @@ func pagesReleaseReady(dir string, deployment pagesDeploymentSource) bool { return false } -func referencedPagesDeployments(config *protocol.ActiveConfigResponse) ([]pagesDeploymentSource, error) { +func referencedPagesProjects(config *protocol.ActiveConfigResponse) ([]pagesProjectRef, error) { if config == nil || strings.TrimSpace(config.SourceConfigJSON) == "" { return nil, nil } @@ -200,26 +327,52 @@ func referencedPagesDeployments(config *protocol.ActiveConfigResponse) ([]pagesD return nil, fmt.Errorf("decode pages references: %w", err) } seen := make(map[uint]struct{}) - result := make([]pagesDeploymentSource, 0) + result := make([]pagesProjectRef, 0) for _, route := range doc.Routes { - if strings.ToLower(strings.TrimSpace(route.UpstreamType)) != "pages" || route.PagesDeployment == nil { + if strings.ToLower(strings.TrimSpace(route.UpstreamType)) != "pages" { continue } - deploymentID := route.PagesDeployment.DeploymentID - checksum := strings.TrimSpace(route.PagesDeployment.Checksum) - if deploymentID == 0 || checksum == "" { - return nil, errors.New("pages deployment snapshot is incomplete") + projectID := pagesProjectIDFromRoute(route) + if projectID == 0 { + return nil, errors.New("pages route is missing project_id") } - if _, ok := seen[deploymentID]; ok { + if _, ok := seen[projectID]; ok { continue } - seen[deploymentID] = struct{}{} - result = append(result, pagesDeploymentSource{DeploymentID: deploymentID, Checksum: checksum}) + seen[projectID] = struct{}{} + checksum := "" + deploymentID := uint(0) + if route.PagesDeployment != nil { + checksum = strings.TrimSpace(route.PagesDeployment.Checksum) + deploymentID = route.PagesDeployment.DeploymentID + } + result = append(result, pagesProjectRef{ + ProjectID: projectID, + DeploymentID: deploymentID, + Checksum: checksum, + }) } return result, nil } -func extractPagesPackage(packageBytes []byte, releaseDir string, deployment pagesDeploymentSource) error { +func pagesProjectIDFromRoute(route pagesSourceRoute) uint { + if route.PagesProjectID != nil && *route.PagesProjectID != 0 { + return *route.PagesProjectID + } + if route.PagesDeployment != nil && route.PagesDeployment.ProjectID != 0 { + return route.PagesDeployment.ProjectID + } + return 0 +} + +// pagesDeploymentSource is the subset of pages_deployment used when parsing config. +type pagesDeploymentSource struct { + ProjectID uint `json:"project_id"` + DeploymentID uint `json:"deployment_id"` + Checksum string `json:"checksum"` +} + +func extractPagesPackage(packageBytes []byte, releaseDir string, project pagesProjectRef) error { tmpDir := releaseDir + ".tmp" _ = os.RemoveAll(tmpDir) if err := os.MkdirAll(tmpDir, pagesDirPerm); err != nil { @@ -230,9 +383,7 @@ func extractPagesPackage(packageBytes []byte, releaseDir string, deployment page _ = os.RemoveAll(tmpDir) return fmt.Errorf("detect Pages package format: %w", err) } - // Control plane already inspected and accepted this package. Agent only - // verifies download integrity (checksum) and performs local-safe extract - // (path escape / symlink guards). Size and file-count limits are not re-applied. + // Control plane already inspected and accepted this package. if err := pagesarchive.ExtractBytes(packageBytes, format, tmpDir, pagesarchive.ExtractOptions{ StripCommonRoot: true, EnforceLimits: false, @@ -240,7 +391,7 @@ func extractPagesPackage(packageBytes []byte, releaseDir string, deployment page _ = os.RemoveAll(tmpDir) return fmt.Errorf("extract Pages package: %w", err) } - if err := writePagesMarker(tmpDir, deployment); err != nil { + if err := writePagesMarker(tmpDir, project); err != nil { _ = os.RemoveAll(tmpDir) return err } @@ -248,8 +399,8 @@ func extractPagesPackage(packageBytes []byte, releaseDir string, deployment page return os.Rename(tmpDir, releaseDir) } -func switchPagesCurrentDir(baseDir string, deploymentID uint, releaseDir string) error { - currentDir := pagesCurrentDir(baseDir, deploymentID) +func switchPagesProjectCurrentDir(baseDir string, projectID uint, releaseDir string) error { + currentDir := pagesProjectCurrentDir(baseDir, projectID) previousDir := currentDir + ".previous" _ = os.RemoveAll(previousDir) if err := os.MkdirAll(filepath.Dir(currentDir), pagesDirPerm); err != nil { @@ -261,7 +412,6 @@ func switchPagesCurrentDir(baseDir string, deploymentID uint, releaseDir string) relTarget = releaseDir } - // Try creating a temporary symlink first to check if symlinks are supported/feasible tmpSymlink := currentDir + ".tmp" _ = os.Remove(tmpSymlink) @@ -269,8 +419,6 @@ func switchPagesCurrentDir(baseDir string, deploymentID uint, releaseDir string) if symlinkErr != nil { return fallbackCopyPagesCurrentDir(currentDir, previousDir, releaseDir) } - - // Symlink is supported, proceed with symlink swap _ = os.Remove(tmpSymlink) if _, err := os.Lstat(currentDir); err == nil { @@ -337,7 +485,7 @@ func copyPagesDir(sourceDir string, targetDir string) error { }) } -func markerMatches(dir string, deployment pagesDeploymentSource) bool { +func markerMatches(dir string, project pagesProjectRef) bool { data, err := os.ReadFile(filepath.Join(dir, ".openflare-pages.json")) //nolint:gosec // dir is managed PagesDir if err != nil { return false @@ -346,23 +494,26 @@ func markerMatches(dir string, deployment pagesDeploymentSource) bool { if err := json.Unmarshal(data, &marker); err != nil { return false } - return marker.DeploymentID == deployment.DeploymentID && marker.Checksum == deployment.Checksum + if marker.ProjectID != 0 && marker.ProjectID != project.ProjectID { + return false + } + return marker.Checksum == project.Checksum } -func writePagesMarker(dir string, deployment pagesDeploymentSource) error { - data, err := json.Marshal(pagesDeploymentMarker(deployment)) +func writePagesMarker(dir string, project pagesProjectRef) error { + data, err := json.Marshal(pagesDeploymentMarker(project)) if err != nil { return err } return os.WriteFile(filepath.Join(dir, ".openflare-pages.json"), data, pagesManifestFilePerm) } -func pagesCurrentDir(baseDir string, deploymentID uint) string { - return filepath.Join(baseDir, "deployments", fmt.Sprintf("%d", deploymentID), "current") +func pagesProjectCurrentDir(baseDir string, projectID uint) string { + return filepath.Join(baseDir, "projects", fmt.Sprintf("%d", projectID), "current") } -func pagesReleaseDir(baseDir string, deploymentID uint, checksum string) string { - return filepath.Join(baseDir, "deployments", fmt.Sprintf("%d", deploymentID), "releases", checksum) +func pagesProjectReleaseDir(baseDir string, projectID uint, checksum string) string { + return filepath.Join(baseDir, "projects", fmt.Sprintf("%d", projectID), "releases", checksum) } func checksumBytes(data []byte) string { diff --git a/internal/apps/agent/sync/service.go b/internal/apps/agent/sync/service.go index 80ca9c6b..30562c91 100644 --- a/internal/apps/agent/sync/service.go +++ b/internal/apps/agent/sync/service.go @@ -32,6 +32,8 @@ type ConfigClient interface { GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigResponse, error) GetPagesDeploymentHash(ctx context.Context, deploymentID uint) (string, error) DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error) + GetPagesProjectLatestHash(ctx context.Context, projectID uint) (*protocol.PagesProjectLatestHashResponse, error) + DownloadPagesProjectLatestPackage(ctx context.Context, projectID uint) ([]byte, error) ReportApplyLog(ctx context.Context, payload protocol.ApplyLogPayload) error SyncWAFIPGroups(ctx context.Context, payload protocol.WAFIPGroupSyncRequest) (*protocol.WAFIPGroupSyncResponse, error) } diff --git a/internal/apps/agent/sync/service_test.go b/internal/apps/agent/sync/service_test.go index fdae7b7b..ecec6338 100644 --- a/internal/apps/agent/sync/service_test.go +++ b/internal/apps/agent/sync/service_test.go @@ -23,19 +23,20 @@ type fakeExecutor struct { reloadErr error } -func testPagesSourceConfigJSON(deploymentID uint, checksum string) string { - return fmt.Sprintf(`{"routes":[{"id":1,"site_name":"pages","domain":"pages.example.com","domains":["pages.example.com"],"origin_url":"openflare-pages://project/1","upstreams":["openflare-pages://project/1"],"enabled":true,"upstream_type":"pages","pages_deployment":{"project_id":1,"project_slug":"pages","deployment_id":%d,"deployment_number":1,"checksum":"%s","entry_file":"index.html","spa_fallback_enabled":true,"local_root":"__OPENFLARE_PAGES_DIR__/deployments/%d/current"}}],"openresty_config":{"worker_processes":"auto","worker_connections":1024,"worker_rlimit_nofile":65535,"events_multi_accept_enabled":true,"keepalive_timeout":20,"keepalive_requests":1000,"client_header_timeout":15,"client_body_timeout":15,"client_max_body_size":"64m","large_client_header_buffers":"4 16k","send_timeout":30,"proxy_connect_timeout":3,"proxy_send_timeout":60,"proxy_read_timeout":60,"websocket_enabled":true,"proxy_request_buffering":false,"proxy_buffering_enabled":true,"proxy_buffers":"16 16k","proxy_buffer_size":"8k","proxy_busy_buffers_size":"64k","gzip_enabled":true,"gzip_min_length":1024,"gzip_comp_level":5,"cache_enabled":false,"cache_levels":"1:2","cache_inactive":"30m","cache_max_size":"1g","cache_key_template":"$scheme$host$request_uri","cache_lock_enabled":true,"cache_lock_timeout":"5s","cache_use_stale":"error timeout updating http_500 http_502 http_503 http_504","main_config_template":"worker_processes {{OpenRestyWorkerProcesses}};"},"waf":{"rule_groups":[],"bindings":[]}}`, deploymentID, checksum, deploymentID) +func testPagesSourceConfigJSON(projectID, deploymentID uint, checksum string) string { + return fmt.Sprintf(`{"routes":[{"id":1,"site_name":"pages","domain":"pages.example.com","domains":["pages.example.com"],"origin_url":"openflare-pages://project/%d","upstreams":["openflare-pages://project/%d"],"enabled":true,"upstream_type":"pages","pages_project_id":%d,"pages_deployment":{"project_id":%d,"project_slug":"pages","deployment_id":%d,"deployment_number":1,"checksum":"%s","entry_file":"index.html","spa_fallback_enabled":true,"local_root":"__OPENFLARE_PAGES_DIR__/projects/%d/current"}}],"openresty_config":{"worker_processes":"auto","worker_connections":1024,"worker_rlimit_nofile":65535,"events_multi_accept_enabled":true,"keepalive_timeout":20,"keepalive_requests":1000,"client_header_timeout":15,"client_body_timeout":15,"client_max_body_size":"64m","large_client_header_buffers":"4 16k","send_timeout":30,"proxy_connect_timeout":3,"proxy_send_timeout":60,"proxy_read_timeout":60,"websocket_enabled":true,"proxy_request_buffering":false,"proxy_buffering_enabled":true,"proxy_buffers":"16 16k","proxy_buffer_size":"8k","proxy_busy_buffers_size":"64k","gzip_enabled":true,"gzip_min_length":1024,"gzip_comp_level":5,"cache_enabled":false,"cache_levels":"1:2","cache_inactive":"30m","cache_max_size":"1g","cache_key_template":"$scheme$host$request_uri","cache_lock_enabled":true,"cache_lock_timeout":"5s","cache_use_stale":"error timeout updating http_500 http_502 http_503 http_504","main_config_template":"worker_processes {{OpenRestyWorkerProcesses}};"},"waf":{"rule_groups":[],"bindings":[]}}`, projectID, projectID, projectID, projectID, deploymentID, checksum, projectID) } type fakeClient struct { - config protocol.ActiveConfigResponse - reports []protocol.ApplyLogPayload - wafSyncCalls []protocol.WAFIPGroupSyncRequest - pagesPackages map[uint][]byte - pagesHashes map[uint]string - wafSyncResult protocol.WAFIPGroupSyncResponse - fetchCalls int - hashCalls int + config protocol.ActiveConfigResponse + reports []protocol.ApplyLogPayload + wafSyncCalls []protocol.WAFIPGroupSyncRequest + pagesPackages map[uint][]byte // key: project_id (latest package) + pagesHashes map[uint]string // key: project_id + pagesLatestDeployIDs map[uint]uint // key: project_id → deployment_id + wafSyncResult protocol.WAFIPGroupSyncResponse + fetchCalls int + hashCalls int } type fakeManager struct { @@ -87,25 +88,71 @@ func (f *fakeClient) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfi } func (f *fakeClient) GetPagesDeploymentHash(ctx context.Context, deploymentID uint) (string, error) { + // Legacy path: map deployment id lookups via reverse of latest deploy ids when present. f.hashCalls++ + for projectID, depID := range f.pagesLatestDeployIDs { + if depID == deploymentID { + return f.projectHash(projectID) + } + } + return f.projectHash(deploymentID) +} + +func (f *fakeClient) DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error) { + for projectID, depID := range f.pagesLatestDeployIDs { + if depID == deploymentID { + return f.projectPackage(projectID) + } + } + return f.projectPackage(deploymentID) +} + +func (f *fakeClient) GetPagesProjectLatestHash(ctx context.Context, projectID uint) (*protocol.PagesProjectLatestHashResponse, error) { + f.hashCalls++ + hash, err := f.projectHash(projectID) + if err != nil { + return nil, err + } + deploymentID := projectID + if f.pagesLatestDeployIDs != nil { + if id, ok := f.pagesLatestDeployIDs[projectID]; ok { + deploymentID = id + } + } + return &protocol.PagesProjectLatestHashResponse{ + ProjectID: projectID, + DeploymentID: deploymentID, + Hash: hash, + }, nil +} + +func (f *fakeClient) DownloadPagesProjectLatestPackage(ctx context.Context, projectID uint) ([]byte, error) { + return f.projectPackage(projectID) +} + +func (f *fakeClient) projectHash(projectID uint) (string, error) { if f.pagesHashes != nil { - if hash, ok := f.pagesHashes[deploymentID]; ok { + if hash, ok := f.pagesHashes[projectID]; ok { return hash, nil } } if f.pagesPackages != nil { - if packageBytes, ok := f.pagesPackages[deploymentID]; ok { + if packageBytes, ok := f.pagesPackages[projectID]; ok { return testBytesChecksum(packageBytes), nil } } - return "", fmt.Errorf("missing Pages hash %d", deploymentID) + return "", fmt.Errorf("missing Pages project hash %d", projectID) } -func (f *fakeClient) DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error) { +func (f *fakeClient) projectPackage(projectID uint) ([]byte, error) { if f.pagesPackages == nil { - return nil, fmt.Errorf("missing Pages package %d", deploymentID) + return nil, fmt.Errorf("missing Pages project package %d", projectID) } - return f.pagesPackages[deploymentID], nil + packageBytes, ok := f.pagesPackages[projectID] + if !ok { + return nil, fmt.Errorf("missing Pages project package %d", projectID) + } + return packageBytes, nil } func (f *fakeClient) ReportApplyLog(ctx context.Context, payload protocol.ApplyLogPayload) error { @@ -378,7 +425,7 @@ func TestSyncOnceDownloadsPagesDeploymentBeforeApply(t *testing.T) { config: protocol.ActiveConfigResponse{ Version: "20260309-101", Checksum: "pages-config-checksum", - SourceConfigJSON: testPagesSourceConfigJSON(7, checksum), + SourceConfigJSON: testPagesSourceConfigJSON(7, 7, checksum), CreatedAt: time.Now().Format(time.RFC3339), }, pagesPackages: map[uint][]byte{7: packageBytes}, @@ -401,14 +448,14 @@ func TestSyncOnceDownloadsPagesDeploymentBeforeApply(t *testing.T) { if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{Version: "20260309-101", Checksum: "pages-config-checksum"}); err != nil { t.Fatalf("SyncOnce failed: %v", err) } - data, err := os.ReadFile(filepath.Join(pagesDir, "deployments", "7", "current", "index.html")) + data, err := os.ReadFile(filepath.Join(pagesDir, "projects", "7", "current", "index.html")) if err != nil { t.Fatalf("expected Pages file to be extracted: %v", err) } if string(data) != "hello" { t.Fatalf("unexpected Pages file content: %s", string(data)) } - if len(manager.applyRouteContents) != 1 || !strings.Contains(manager.applyRouteContents[0], "__OPENFLARE_PAGES_DIR__/deployments/7/current") { + if len(manager.applyRouteContents) != 1 || !strings.Contains(manager.applyRouteContents[0], "__OPENFLARE_PAGES_DIR__/projects/7/current") { t.Fatalf("expected Pages placeholder in rendered route config, got %#v", manager.applyRouteContents) } } @@ -426,7 +473,7 @@ func TestSyncPagesDeploymentEnsuresWorkerReadAccess(t *testing.T) { config := protocol.ActiveConfigResponse{ Version: "20260309-106", Checksum: "pages-config-checksum", - SourceConfigJSON: testPagesSourceConfigJSON(7, checksum), + SourceConfigJSON: testPagesSourceConfigJSON(7, 7, checksum), CreatedAt: time.Now().Format(time.RFC3339), } client := &fakeClient{ @@ -451,7 +498,7 @@ func TestSyncPagesDeploymentEnsuresWorkerReadAccess(t *testing.T) { if dataInfo.Mode().Perm()&0o005 == 0 { t.Fatalf("expected dataDir to be world-traversable, got %o", dataInfo.Mode().Perm()) } - indexPath := filepath.Join(pagesDir, "deployments", "7", "current", "index.html") + indexPath := filepath.Join(pagesDir, "projects", "7", "current", "index.html") indexInfo, err := os.Stat(indexPath) if err != nil { t.Fatalf("expected Pages file to be extracted: %v", err) @@ -471,7 +518,7 @@ func TestSyncOnceExtractsPagesPackageWithZeroByteFiles(t *testing.T) { config: protocol.ActiveConfigResponse{ Version: "20260309-103", Checksum: "pages-config-checksum", - SourceConfigJSON: testPagesSourceConfigJSON(9, checksum), + SourceConfigJSON: testPagesSourceConfigJSON(9, 9, checksum), CreatedAt: time.Now().Format(time.RFC3339), }, pagesPackages: map[uint][]byte{9: packageBytes}, @@ -494,7 +541,7 @@ func TestSyncOnceExtractsPagesPackageWithZeroByteFiles(t *testing.T) { if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{Version: "20260309-103", Checksum: "pages-config-checksum"}); err != nil { t.Fatalf("SyncOnce failed: %v", err) } - gitkeepPath := filepath.Join(pagesDir, "deployments", "9", "current", ".gitkeep") + gitkeepPath := filepath.Join(pagesDir, "projects", "9", "current", ".gitkeep") info, err := os.Stat(gitkeepPath) if err != nil { t.Fatalf("expected zero-byte Pages file to be extracted: %v", err) @@ -511,7 +558,7 @@ func TestSyncOnceRejectsPagesZipSlipBeforeApply(t *testing.T) { config: protocol.ActiveConfigResponse{ Version: "20260309-102", Checksum: "pages-config-checksum", - SourceConfigJSON: testPagesSourceConfigJSON(8, checksum), + SourceConfigJSON: testPagesSourceConfigJSON(8, 8, checksum), CreatedAt: time.Now().Format(time.RFC3339), }, pagesPackages: map[uint][]byte{8: packageBytes}, @@ -1169,7 +1216,7 @@ func TestSyncOnceDownloadsPagesDeploymentWhenChecksumMatches(t *testing.T) { config: protocol.ActiveConfigResponse{ Version: "20260309-101", Checksum: "pages-config-checksum", - SourceConfigJSON: testPagesSourceConfigJSON(7, checksum), + SourceConfigJSON: testPagesSourceConfigJSON(7, 7, checksum), CreatedAt: time.Now().Format(time.RFC3339), }, pagesPackages: map[uint][]byte{7: packageBytes}, @@ -1197,7 +1244,7 @@ func TestSyncOnceDownloadsPagesDeploymentWhenChecksumMatches(t *testing.T) { }); err != nil { t.Fatalf("SyncOnce failed: %v", err) } - data, err := os.ReadFile(filepath.Join(pagesDir, "deployments", "7", "current", "index.html")) + data, err := os.ReadFile(filepath.Join(pagesDir, "projects", "7", "current", "index.html")) if err != nil { t.Fatalf("expected Pages file to be extracted when checksum already matches: %v", err) } @@ -1225,8 +1272,10 @@ func TestSyncOnceDownloadsPagesDeploymentWhenChecksumMatches(t *testing.T) { if err != nil { t.Fatalf("failed to load state: %v", err) } - if len(snapshot.PagesDeployments) != 1 || snapshot.PagesDeployments[0].DeploymentID != 7 || snapshot.PagesDeployments[0].Hash != checksum { - t.Fatalf("expected Pages deployment refs to be cached in state, got %+v", snapshot.PagesDeployments) + if len(snapshot.PagesDeployments) != 1 || + snapshot.PagesDeployments[0].ProjectID != 7 || + snapshot.PagesDeployments[0].Hash != checksum { + t.Fatalf("expected Pages project refs to be cached in state, got %+v", snapshot.PagesDeployments) } if client.hashCalls == 0 { t.Fatal("expected Pages hash check during reconcile") @@ -1238,16 +1287,16 @@ func TestSyncOnceRedownloadsPagesDeploymentWhenServerHashChanges(t *testing.T) { updatedPackage := testPagesPackage(t, map[string]string{"index.html": "v2"}) initialHash := testBytesChecksum(initialPackage) updatedHash := testBytesChecksum(updatedPackage) - deploymentID := uint(12) + projectID := uint(12) client := &fakeClient{ config: protocol.ActiveConfigResponse{ Version: "20260309-107", Checksum: "pages-config-checksum", - SourceConfigJSON: testPagesSourceConfigJSON(deploymentID, initialHash), + SourceConfigJSON: testPagesSourceConfigJSON(projectID, projectID, initialHash), CreatedAt: time.Now().Format(time.RFC3339), }, - pagesPackages: map[uint][]byte{deploymentID: initialPackage}, - pagesHashes: map[uint]string{deploymentID: initialHash}, + pagesPackages: map[uint][]byte{projectID: initialPackage}, + pagesHashes: map[uint]string{projectID: initialHash}, } stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) nodeID, err := stateStore.EnsureNodeID() @@ -1259,26 +1308,28 @@ func TestSyncOnceRedownloadsPagesDeploymentWhenServerHashChanges(t *testing.T) { CurrentVersion: client.config.Version, CurrentChecksum: client.config.Checksum, PagesDeployments: []state.PagesDeployment{{ - DeploymentID: deploymentID, - Hash: initialHash, + ProjectID: projectID, + Hash: initialHash, }}, }); err != nil { t.Fatalf("save state failed: %v", err) } pagesDir := t.TempDir() - releaseDir := pagesReleaseDir(pagesDir, deploymentID, initialHash) - if err = extractPagesPackage(initialPackage, releaseDir, pagesDeploymentSource{ - DeploymentID: deploymentID, - Checksum: initialHash, + releaseDir := pagesProjectReleaseDir(pagesDir, projectID, initialHash) + if err = extractPagesPackage(initialPackage, releaseDir, pagesProjectRef{ + ProjectID: projectID, + Checksum: initialHash, }); err != nil { t.Fatalf("seed release failed: %v", err) } - if err = switchPagesCurrentDir(pagesDir, deploymentID, releaseDir); err != nil { + if err = switchPagesProjectCurrentDir(pagesDir, projectID, releaseDir); err != nil { t.Fatalf("seed current dir failed: %v", err) } - client.pagesPackages[deploymentID] = updatedPackage - client.pagesHashes[deploymentID] = updatedHash + // Control plane activates a new package; agent discovers via latest hash poll + // without main-config republish (checksum unchanged). + client.pagesPackages[projectID] = updatedPackage + client.pagesHashes[projectID] = updatedHash manager := &fakeManager{currentChecksum: client.config.Checksum} service := New(client, manager, stateStore) service.SetPagesDir(pagesDir) @@ -1288,7 +1339,7 @@ func TestSyncOnceRedownloadsPagesDeploymentWhenServerHashChanges(t *testing.T) { }); err != nil { t.Fatalf("SyncOnce failed: %v", err) } - data, err := os.ReadFile(filepath.Join(pagesDir, "deployments", "12", "current", "index.html")) + data, err := os.ReadFile(filepath.Join(pagesDir, "projects", "12", "current", "index.html")) if err != nil { t.Fatalf("expected updated Pages file: %v", err) } @@ -1298,20 +1349,236 @@ func TestSyncOnceRedownloadsPagesDeploymentWhenServerHashChanges(t *testing.T) { if client.fetchCalls != 0 { t.Fatalf("expected hash-only reconcile without active config fetch, got %d", client.fetchCalls) } + // Only the latest release must remain on disk. + releasesRoot := filepath.Join(pagesDir, "projects", "12", "releases") + entries, err := os.ReadDir(releasesRoot) + if err != nil { + t.Fatalf("read releases dir: %v", err) + } + if len(entries) != 1 || entries[0].Name() != updatedHash { + names := make([]string, 0, len(entries)) + for _, e := range entries { + names = append(names, e.Name()) + } + t.Fatalf("expected only latest release %s, got %v", updatedHash, names) + } + if _, err := os.Stat(filepath.Join(releasesRoot, initialHash)); !os.IsNotExist(err) { + t.Fatalf("expected stale release %s to be removed", initialHash) + } +} + +// racingLatestClient simulates control-plane activation changing between hash +// probe and package download (hash A → package B → verify hash B). +type racingLatestClient struct { + fakeClient + pkgA, pkgB []byte + hashA, hashB string + hashCall int + downloadCalls int +} + +func (r *racingLatestClient) GetPagesProjectLatestHash(ctx context.Context, projectID uint) (*protocol.PagesProjectLatestHashResponse, error) { + r.hashCall++ + r.hashCalls++ + // Sequence: + // 1: probe before download → A + // 2: verify after downloading B → B (race) + // 3+: stable on B for retry + hash, dep := r.hashA, uint(1) + if r.hashCall >= 2 { + hash, dep = r.hashB, 2 + } + return &protocol.PagesProjectLatestHashResponse{ + ProjectID: projectID, + DeploymentID: dep, + Hash: hash, + }, nil +} + +func (r *racingLatestClient) DownloadPagesProjectLatestPackage(ctx context.Context, projectID uint) ([]byte, error) { + r.downloadCalls++ + // Always return package B (what "latest download" would stream mid-race / after). + return r.pkgB, nil +} + +func TestEnsurePagesProjectSurvivesHashPackageRace(t *testing.T) { + pkgA := testPagesPackage(t, map[string]string{"index.html": "A"}) + pkgB := testPagesPackage(t, map[string]string{"index.html": "B"}) + hashA := testBytesChecksum(pkgA) + hashB := testBytesChecksum(pkgB) + projectID := uint(42) + client := &racingLatestClient{ + pkgA: pkgA, + pkgB: pkgB, + hashA: hashA, + hashB: hashB, + } + client.config = protocol.ActiveConfigResponse{ + Version: "v-race", + Checksum: "c-race", + SourceConfigJSON: testPagesSourceConfigJSON(projectID, 1, hashA), + CreatedAt: time.Now().Format(time.RFC3339), + } + stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) + snapshot, _ := stateStore.Load() + pagesDir := t.TempDir() + service := New(client, &fakeManager{}, stateStore) + service.SetPagesDir(pagesDir) + + if err := service.syncPagesDeployments(context.Background(), snapshot, &client.config); err != nil { + t.Fatalf("sync after race: %v", err) + } + data, err := os.ReadFile(filepath.Join(pagesDir, "projects", "42", "current", "index.html")) + if err != nil { + t.Fatalf("read current: %v", err) + } + if string(data) != "B" { + t.Fatalf("expected final content B after race recovery, got %q", string(data)) + } + // Only B release retained. + entries, err := os.ReadDir(filepath.Join(pagesDir, "projects", "42", "releases")) + if err != nil { + t.Fatal(err) + } + if len(entries) != 1 || entries[0].Name() != hashB { + t.Fatalf("expected only release %s, got %v", hashB, entries) + } + if client.downloadCalls < 1 { + t.Fatal("expected at least one package download") + } +} + +func TestCleanupDoesNotRunBeforeCurrentSwitch(t *testing.T) { + // If switch fails, stale release must remain so traffic can keep using it. + pagesDir := t.TempDir() + projectID := uint(9) + oldHash := "old-hash" + oldDir := pagesProjectReleaseDir(pagesDir, projectID, oldHash) + if err := os.MkdirAll(oldDir, 0o755); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(oldDir, "index.html"), []byte("old"), 0o644); err != nil { + t.Fatal(err) + } + if err := switchPagesProjectCurrentDir(pagesDir, projectID, oldDir); err != nil { + t.Fatal(err) + } + // Simulate failed upgrade path: new release extracted but we do NOT switch/cleanup. + newHash := "new-hash" + newDir := pagesProjectReleaseDir(pagesDir, projectID, newHash) + if err := os.MkdirAll(newDir, 0o755); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(newDir, "index.html"), []byte("new"), 0o644); err != nil { + t.Fatal(err) + } + // Old must still be readable via current. + data, err := os.ReadFile(filepath.Join(pagesDir, "projects", "9", "current", "index.html")) + if err != nil { + t.Fatal(err) + } + if string(data) != "old" { + t.Fatalf("current should still serve old until switch, got %q", data) + } + if _, err := os.Stat(oldDir); err != nil { + t.Fatal("old release must not be deleted before successful switch") + } +} + +func TestLegacyStateWithoutProjectIDForcesDiscovery(t *testing.T) { + snapshot := &state.Snapshot{ + PagesDeployments: []state.PagesDeployment{{ + DeploymentID: 99, + Hash: "abc", + }}, + } + if !pagesDiscoveryNeeded(snapshot) { + t.Fatal("legacy deployment-only state must force config rediscovery") + } + if pagesSyncNeeded(snapshot) { + t.Fatal("legacy rows without project_id must not enter hash-only sync") + } +} + +func TestCleanupPagesProjectStaleReleasesKeepsOnlyLatest(t *testing.T) { + pagesDir := t.TempDir() + projectID := uint(3) + keep := "keep-hash" + stale := "stale-hash" + keepDir := pagesProjectReleaseDir(pagesDir, projectID, keep) + staleDir := pagesProjectReleaseDir(pagesDir, projectID, stale) + if err := os.MkdirAll(filepath.Join(keepDir, "nested"), 0o755); err != nil { + t.Fatal(err) + } + if err := os.MkdirAll(staleDir, 0o755); err != nil { + t.Fatal(err) + } + tmpDir := staleDir + ".tmp" + if err := os.MkdirAll(tmpDir, 0o755); err != nil { + t.Fatal(err) + } + if err := cleanupPagesProjectStaleReleases(pagesDir, projectID, keep); err != nil { + t.Fatalf("cleanup: %v", err) + } + if _, err := os.Stat(keepDir); err != nil { + t.Fatalf("keep release missing: %v", err) + } + if _, err := os.Stat(staleDir); !os.IsNotExist(err) { + t.Fatal("stale release should be removed") + } + if _, err := os.Stat(tmpDir); !os.IsNotExist(err) { + t.Fatal("tmp leftover should be removed") + } +} + +func TestSyncPagesDeploymentsIsolatesProjectFailures(t *testing.T) { + okPackage := testPagesPackage(t, map[string]string{"index.html": "ok"}) + okHash := testBytesChecksum(okPackage) + client := &fakeClient{ + config: protocol.ActiveConfigResponse{ + Version: "20260309-200", + Checksum: "cfg", + SourceConfigJSON: fmt.Sprintf( + `{"routes":[{"upstream_type":"pages","pages_project_id":1,"pages_deployment":{"project_id":1}},{"upstream_type":"pages","pages_project_id":2,"pages_deployment":{"project_id":2}}]}`, + ), + CreatedAt: time.Now().Format(time.RFC3339), + }, + pagesPackages: map[uint][]byte{1: okPackage}, + pagesHashes: map[uint]string{1: okHash}, + // project 2 missing → ensure fails + } + stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) + snapshot, _ := stateStore.Load() + pagesDir := t.TempDir() + service := New(client, &fakeManager{}, stateStore) + service.SetPagesDir(pagesDir) + + err := service.syncPagesDeployments(context.Background(), snapshot, &client.config) + if err == nil { + t.Fatal("expected aggregated error when one project fails") + } + // Project 1 must still have been applied. + data, readErr := os.ReadFile(filepath.Join(pagesDir, "projects", "1", "current", "index.html")) + if readErr != nil { + t.Fatalf("project 1 should succeed despite project 2 failure: %v", readErr) + } + if string(data) != "ok" { + t.Fatalf("unexpected content: %s", data) + } } func TestSyncOnceRedownloadsPagesDeploymentWhenReleaseDirOnlyHasMarker(t *testing.T) { packageBytes := testPagesPackage(t, map[string]string{"index.html": "hello"}) checksum := testBytesChecksum(packageBytes) - deploymentID := uint(11) + projectID := uint(11) client := &fakeClient{ config: protocol.ActiveConfigResponse{ Version: "20260309-106", Checksum: "pages-config-checksum", - SourceConfigJSON: testPagesSourceConfigJSON(deploymentID, checksum), + SourceConfigJSON: testPagesSourceConfigJSON(projectID, projectID, checksum), CreatedAt: time.Now().Format(time.RFC3339), }, - pagesPackages: map[uint][]byte{deploymentID: packageBytes}, + pagesPackages: map[uint][]byte{projectID: packageBytes}, } stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) nodeID, err := stateStore.EnsureNodeID() @@ -1326,14 +1593,14 @@ func TestSyncOnceRedownloadsPagesDeploymentWhenReleaseDirOnlyHasMarker(t *testin t.Fatalf("save state failed: %v", err) } pagesDir := t.TempDir() - releaseDir := pagesReleaseDir(pagesDir, deploymentID, checksum) + releaseDir := pagesProjectReleaseDir(pagesDir, projectID, checksum) if err = os.MkdirAll(releaseDir, pagesDirPerm); err != nil { t.Fatalf("mkdir release dir failed: %v", err) } - if err = writePagesMarker(releaseDir, pagesDeploymentSource{DeploymentID: deploymentID, Checksum: checksum}); err != nil { + if err = writePagesMarker(releaseDir, pagesProjectRef{ProjectID: projectID, Checksum: checksum}); err != nil { t.Fatalf("write marker failed: %v", err) } - if pagesReleaseReady(releaseDir, pagesDeploymentSource{DeploymentID: deploymentID, Checksum: checksum}) { + if pagesProjectReleaseReady(releaseDir, pagesProjectRef{ProjectID: projectID, Checksum: checksum}) { t.Fatal("expected marker-only release dir to be treated as not ready") } @@ -1346,7 +1613,7 @@ func TestSyncOnceRedownloadsPagesDeploymentWhenReleaseDirOnlyHasMarker(t *testin }); err != nil { t.Fatalf("SyncOnce failed: %v", err) } - data, err := os.ReadFile(filepath.Join(pagesDir, "deployments", "11", "current", "index.html")) + data, err := os.ReadFile(filepath.Join(pagesDir, "projects", "11", "current", "index.html")) if err != nil { t.Fatalf("expected Pages file to be extracted after marker-only release dir: %v", err) } @@ -1389,7 +1656,7 @@ func TestSyncOnceDownloadsPagesDeploymentWithTopLevelFolder(t *testing.T) { config: protocol.ActiveConfigResponse{ Version: "20260309-105", Checksum: "pages-config-checksum", - SourceConfigJSON: testPagesSourceConfigJSON(77, checksum), + SourceConfigJSON: testPagesSourceConfigJSON(77, 77, checksum), CreatedAt: time.Now().Format(time.RFC3339), }, pagesPackages: map[uint][]byte{77: packageBytes}, @@ -1412,14 +1679,14 @@ func TestSyncOnceDownloadsPagesDeploymentWithTopLevelFolder(t *testing.T) { if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{Version: "20260309-105", Checksum: "pages-config-checksum"}); err != nil { t.Fatalf("SyncOnce failed: %v", err) } - data, err := os.ReadFile(filepath.Join(pagesDir, "deployments", "77", "current", "index.html")) + data, err := os.ReadFile(filepath.Join(pagesDir, "projects", "77", "current", "index.html")) if err != nil { t.Fatalf("expected Pages index.html file to be extracted: %v", err) } if string(data) != "hello html" { t.Fatalf("unexpected Pages index.html content: %s", string(data)) } - jsData, err := os.ReadFile(filepath.Join(pagesDir, "deployments", "77", "current", "assets", "app.js")) + jsData, err := os.ReadFile(filepath.Join(pagesDir, "projects", "77", "current", "assets", "app.js")) if err != nil { t.Fatalf("expected Pages assets/app.js file to be extracted: %v", err) } diff --git a/internal/apps/openflare/agent/routers.go b/internal/apps/openflare/agent/routers.go index f6433208..e39b27be 100644 --- a/internal/apps/openflare/agent/routers.go +++ b/internal/apps/openflare/agent/routers.go @@ -156,7 +156,7 @@ func ReportApplyLogHandler(c *gin.Context) { // GetPagesDeploymentHashHandler returns the upload SHA-256 hash for a Pages deployment package. // @Summary 查询 Pages 部署包哈希 -// @Description 返回 upload 框架记录的 SHA-256 哈希,供 Agent 对比本地缓存并按需拉取部署包 +// @Description 返回 upload 框架记录的 SHA-256 哈希,供 Agent 对比本地缓存并按需拉取部署包(兼容旧路径) // @Tags openflare-agent // @Produce json // @Security AgentTokenAuth @@ -166,7 +166,7 @@ func ReportApplyLogHandler(c *gin.Context) { // @Failure 401 {object} response.Any "Token 无效" // @Router /api/v1/agent/pages/deployments/{deployment_id}/hash [get] func GetPagesDeploymentHashHandler(c *gin.Context) { - deploymentID, ok := pagesDeploymentIDParam(c) + deploymentID, ok := pagesUintParam(c, "deployment_id") if !ok { return } @@ -182,7 +182,7 @@ func GetPagesDeploymentHashHandler(c *gin.Context) { // DownloadPagesPackageHandler streams the Pages deployment artifact to an authenticated agent. // @Summary 下载 Pages 部署包 -// @Description 流式下载指定部署的静态资源压缩包,供 Agent 边缘分发 +// @Description 流式下载指定部署的静态资源压缩包,供 Agent 边缘分发(兼容旧路径) // @Tags openflare-agent // @Produce application/octet-stream // @Security AgentTokenAuth @@ -192,7 +192,7 @@ func GetPagesDeploymentHashHandler(c *gin.Context) { // @Failure 401 {object} response.Any "Token 无效" // @Router /api/v1/agent/pages/deployments/{deployment_id}/package [get] func DownloadPagesPackageHandler(c *gin.Context) { - deploymentID, ok := pagesDeploymentIDParam(c) + deploymentID, ok := pagesUintParam(c, "deployment_id") if !ok { return } @@ -208,8 +208,63 @@ func DownloadPagesPackageHandler(c *gin.Context) { c.DataFromReader(http.StatusOK, packageObj.ContentLength, packageObj.ContentType, packageObj.Body, nil) } -func pagesDeploymentIDParam(c *gin.Context) (uint, bool) { - raw := c.Param("deployment_id") +// GetPagesProjectLatestHashHandler returns the hash of a project's currently active deployment. +// @Summary 查询 Pages 项目最新激活部署哈希 +// @Description 按项目 ID 返回当前激活部署的包哈希(类似 latest 指针),Agent 无需关心具体部署 ID +// @Tags openflare-agent +// @Produce json +// @Security AgentTokenAuth +// @Param project_id path int true "Pages 项目 ID" +// @Success 200 {object} response.Any{data=protocol.PagesProjectLatestHashResponse} +// @Failure 400 {object} response.Any "参数错误" +// @Failure 401 {object} response.Any "Token 无效" +// @Router /api/v1/agent/pages/projects/{project_id}/latest/hash [get] +func GetPagesProjectLatestHashHandler(c *gin.Context) { + projectID, ok := pagesUintParam(c, "project_id") + if !ok { + return + } + deploymentID, hash, err := pages.GetProjectLatestPackageHash(c.Request.Context(), projectID) + if apiutil.AbortBadRequestOnError(c, err) { + return + } + c.JSON(http.StatusOK, response.OK(protocol.PagesProjectLatestHashResponse{ + ProjectID: projectID, + DeploymentID: deploymentID, + Hash: hash, + })) +} + +// DownloadPagesProjectLatestPackageHandler streams the active deployment package for a project. +// @Summary 下载 Pages 项目最新激活部署包 +// @Description 按项目 ID 下载当前激活部署的压缩包,供 Agent 边缘分发 +// @Tags openflare-agent +// @Produce application/octet-stream +// @Security AgentTokenAuth +// @Param project_id path int true "Pages 项目 ID" +// @Success 200 {file} binary "部署包文件" +// @Failure 400 {object} response.Any "参数错误" +// @Failure 401 {object} response.Any "Token 无效" +// @Router /api/v1/agent/pages/projects/{project_id}/latest/package [get] +func DownloadPagesProjectLatestPackageHandler(c *gin.Context) { + projectID, ok := pagesUintParam(c, "project_id") + if !ok { + return + } + packageObj, err := pages.OpenProjectLatestPackage(c.Request.Context(), projectID) + if apiutil.AbortBadRequestOnError(c, err) { + return + } + defer func() { _ = packageObj.Body.Close() }() + c.Header("Content-Disposition", "attachment; filename="+packageObj.FileName) + if packageObj.ContentType != "" { + c.Header("Content-Type", packageObj.ContentType) + } + c.DataFromReader(http.StatusOK, packageObj.ContentLength, packageObj.ContentType, packageObj.Body, nil) +} + +func pagesUintParam(c *gin.Context, name string) (uint, bool) { + raw := c.Param(name) if raw == "" { response.AbortBadRequest(c, "无效的 ID") return 0, false diff --git a/internal/apps/openflare/config_version/pages_snapshot.go b/internal/apps/openflare/config_version/pages_snapshot.go index 456d99c3..d7012829 100644 --- a/internal/apps/openflare/config_version/pages_snapshot.go +++ b/internal/apps/openflare/config_version/pages_snapshot.go @@ -88,10 +88,8 @@ func buildSnapshotPagesDeployment(project *model.PagesProject, activeDeployment APIProxyPath: strings.TrimSpace(project.APIProxyPath), APIProxyPass: strings.TrimSpace(project.APIProxyPass), APIProxyRewrite: strings.TrimSpace(project.APIProxyRewrite), - LocalRoot: fmt.Sprintf( - "%s/deployments/%d/current", - openrestyrender.PagesDirPlaceholder, - activeDeployment.ID, - ), + // Root is project-scoped so Agents can swap active packages without + // re-publishing main config (nginx root stays stable). + LocalRoot: openrestyrender.PagesProjectLocalRoot(project.ID), } } diff --git a/internal/apps/openflare/config_version/pages_snapshot_test.go b/internal/apps/openflare/config_version/pages_snapshot_test.go index bff94153..9b6fb1a9 100644 --- a/internal/apps/openflare/config_version/pages_snapshot_test.go +++ b/internal/apps/openflare/config_version/pages_snapshot_test.go @@ -64,7 +64,7 @@ func TestBuildSnapshotRoutesPages(t *testing.T) { require.NotNil(t, snapshotRoute.PagesDeployment) assert.Equal(t, deployment.ID, snapshotRoute.PagesDeployment.DeploymentID) assert.Equal(t, deployment.Checksum, snapshotRoute.PagesDeployment.Checksum) - assert.Equal(t, "__OPENFLARE_PAGES_DIR__/deployments/1/current", snapshotRoute.PagesDeployment.LocalRoot) + assert.Equal(t, "__OPENFLARE_PAGES_DIR__/projects/1/current", snapshotRoute.PagesDeployment.LocalRoot) _, err = renderSnapshotConfig(bundle.SnapshotJSON, nil) require.NoError(t, err) diff --git a/internal/apps/openflare/pages/logics.go b/internal/apps/openflare/pages/logics.go index 883fc2f7..33bdf054 100644 --- a/internal/apps/openflare/pages/logics.go +++ b/internal/apps/openflare/pages/logics.go @@ -509,7 +509,8 @@ func ActivateDeployment(ctx context.Context, projectID uint, deploymentID uint) return GetProject(ctx, project.ID) } -// GetDeploymentPackageHash returns the SHA-256 hash of the deployment package from upload storage. +// GetDeploymentPackageHash returns the upload SHA-256 hash of the deployment package. +// Prefer GetProjectLatestPackageHash for Agent latest-pointer pulls. func GetDeploymentPackageHash(ctx context.Context, deploymentID uint) (string, error) { deployment, err := model.GetPagesDeploymentByID(ctx, deploymentID) if err != nil { @@ -518,14 +519,94 @@ func GetDeploymentPackageHash(ctx context.Context, deploymentID uint) (string, e if err = ensureDeploymentInActiveSnapshot(ctx, deployment.ID); err != nil { return "", err } + return deploymentPackageHash(ctx, deployment) +} + +// GetProjectLatestPackageHash returns the package hash of a project's active deployment. +// This is the Agent "latest" pointer: callers pass project_id only. +func GetProjectLatestPackageHash(ctx context.Context, projectID uint) (uint, string, error) { + deployment, err := resolveProjectActiveDeploymentForAgent(ctx, projectID) + if err != nil { + return 0, "", err + } + hash, err := deploymentPackageHash(ctx, deployment) + if err != nil { + return 0, "", err + } + return deployment.ID, hash, nil +} + +// OpenDeploymentPackage opens the deployment artifact from the upload storage framework. +func OpenDeploymentPackage(ctx context.Context, deploymentID uint) (DeploymentPackage, error) { + deployment, err := model.GetPagesDeploymentByID(ctx, deploymentID) + if err != nil { + return DeploymentPackage{}, err + } + if err = ensureDeploymentInActiveSnapshot(ctx, deployment.ID); err != nil { + return DeploymentPackage{}, err + } + if deployment.UploadID == 0 { + if err := ensureDeploymentUploadRecord(ctx, deployment); err != nil { + return DeploymentPackage{}, err + } + } + return openDeploymentPackageFromUpload(ctx, deployment.UploadID, deployment.ID) +} + +// OpenProjectLatestPackage opens the currently active deployment package for a project. +func OpenProjectLatestPackage(ctx context.Context, projectID uint) (DeploymentPackage, error) { + deployment, err := resolveProjectActiveDeploymentForAgent(ctx, projectID) + if err != nil { + return DeploymentPackage{}, err + } + if deployment.UploadID == 0 { + if err := ensureDeploymentUploadRecord(ctx, deployment); err != nil { + return DeploymentPackage{}, err + } + } + return openDeploymentPackageFromUpload(ctx, deployment.UploadID, deployment.ID) +} + +func resolveProjectActiveDeploymentForAgent(ctx context.Context, projectID uint) (*model.PagesDeployment, error) { + if projectID == 0 { + return nil, errors.New(errPagesProjectNotFound) + } + project, err := model.GetPagesProjectByID(ctx, projectID) + if err != nil { + return nil, err + } + if !project.Enabled { + return nil, errors.New(errPagesPackageNotInActiveConfig) + } + if project.ActiveDeploymentID == nil || *project.ActiveDeploymentID == 0 { + return nil, errors.New(errPagesPackageNotInActiveConfig) + } + if err := ensureProjectInActiveConfig(ctx, project.ID); err != nil { + return nil, err + } + deployment, err := model.GetPagesDeploymentByID(ctx, *project.ActiveDeploymentID) + if err != nil { + return nil, err + } + if deployment.ProjectID != project.ID { + return nil, errors.New(errPagesDeploymentMismatch) + } + return deployment, nil +} + +func deploymentPackageHash(ctx context.Context, deployment *model.PagesDeployment) (string, error) { + if deployment == nil { + return "", errors.New(errPagesDeploymentNotFound) + } if deployment.UploadID == 0 { if err := ensureDeploymentUploadRecord(ctx, deployment); err != nil { return "", err } - deployment, err = model.GetPagesDeploymentByID(ctx, deploymentID) + reloaded, err := model.GetPagesDeploymentByID(ctx, deployment.ID) if err != nil { return "", err } + deployment = reloaded } if deployment.UploadID == 0 { hash := strings.TrimSpace(deployment.Checksum) @@ -548,23 +629,6 @@ func GetDeploymentPackageHash(ctx context.Context, deploymentID uint) (string, e return hash, nil } -// OpenDeploymentPackage opens the deployment artifact from the upload storage framework. -func OpenDeploymentPackage(ctx context.Context, deploymentID uint) (DeploymentPackage, error) { - deployment, err := model.GetPagesDeploymentByID(ctx, deploymentID) - if err != nil { - return DeploymentPackage{}, err - } - if err = ensureDeploymentInActiveSnapshot(ctx, deployment.ID); err != nil { - return DeploymentPackage{}, err - } - if deployment.UploadID == 0 { - if err := ensureDeploymentUploadRecord(ctx, deployment); err != nil { - return DeploymentPackage{}, err - } - } - return openDeploymentPackageFromUpload(ctx, deployment.UploadID, deployment.ID) -} - func openDeploymentPackageFromUpload(ctx context.Context, uploadID uint64, deploymentID uint) (DeploymentPackage, error) { opened, err := upload.OpenStoredUpload(ctx, uploadID) if err != nil { @@ -648,13 +712,9 @@ func hydrateLegacyDeploymentUpload( return &ingestResult.Upload, nil } -// ensureDeploymentInActiveSnapshot allows Agent package download when the -// deployment is the project's current active deployment and that project is -// used by at least one pages route in the active main config. -// -// Main config versions pin historical pages_deployment ids for audit only. -// Runtime download always follows the live active Pages deployment (dual -// version control); rolling back main config must not require old packages. +// ensureDeploymentInActiveSnapshot allows download of a specific deployment when +// it is the project's current active deployment and the project is used by the +// active main config. func ensureDeploymentInActiveSnapshot(ctx context.Context, deploymentID uint) error { deployment, err := model.GetPagesDeploymentByID(ctx, deploymentID) if err != nil { @@ -667,7 +727,12 @@ func ensureDeploymentInActiveSnapshot(ctx context.Context, deploymentID uint) er if project.ActiveDeploymentID == nil || *project.ActiveDeploymentID != deployment.ID { return errors.New(errPagesPackageNotInActiveConfig) } + return ensureProjectInActiveConfig(ctx, project.ID) +} +// ensureProjectInActiveConfig checks that the Pages project is referenced by at +// least one pages route in the active main config snapshot. +func ensureProjectInActiveConfig(ctx context.Context, projectID uint) error { version, err := model.GetActiveConfigVersion(ctx) if err != nil { if errors.Is(err, gorm.ErrRecordNotFound) { @@ -683,28 +748,27 @@ func ensureDeploymentInActiveSnapshot(ctx context.Context, deploymentID uint) er if !strings.EqualFold(strings.TrimSpace(route.UpstreamType), "pages") { continue } - if route.PagesProjectID != nil && *route.PagesProjectID == project.ID { + if route.PagesProjectID != nil && *route.PagesProjectID == projectID { return nil } if route.PagesDeployment == nil { continue } - if route.PagesDeployment.ProjectID == project.ID { + if route.PagesDeployment.ProjectID == projectID { return nil } - // Frozen snapshot may only carry deployment_id; resolve project via that row. if route.PagesDeployment.DeploymentID == 0 { continue } - if route.PagesDeployment.DeploymentID == deployment.ID { - return nil - } - snapDeployment, snapErr := model.GetPagesDeploymentByID(ctx, route.PagesDeployment.DeploymentID) - if snapErr != nil { - continue - } - if snapDeployment.ProjectID == project.ID { - return nil + if route.PagesDeployment.DeploymentID != 0 { + // Historical snapshot may only pin deployment_id. + snapDeployment, snapErr := model.GetPagesDeploymentByID(ctx, route.PagesDeployment.DeploymentID) + if snapErr != nil { + continue + } + if snapDeployment.ProjectID == projectID { + return nil + } } } return errors.New(errPagesPackageNotInActiveConfig) diff --git a/internal/apps/openflare/pages/logics_test.go b/internal/apps/openflare/pages/logics_test.go index efed9f18..ba3d52ff 100644 --- a/internal/apps/openflare/pages/logics_test.go +++ b/internal/apps/openflare/pages/logics_test.go @@ -337,6 +337,16 @@ func TestOpenDeploymentPackageRequiresActiveConfigSnapshot(t *testing.T) { defer packageObj.Body.Close() assert.Equal(t, fmt.Sprintf("pages-deployment-%d.zip", deployment.ID), packageObj.FileName) + // Latest-by-project resolves the active package once the project is on active config. + depID, hash, err := GetProjectLatestPackageHash(ctx, project.ID) + require.NoError(t, err) + assert.Equal(t, deployment.ID, depID) + assert.NotEmpty(t, hash) + + latestPkg, err := OpenProjectLatestPackage(ctx, project.ID) + require.NoError(t, err) + defer latestPkg.Body.Close() + body, err := io.ReadAll(packageObj.Body) require.NoError(t, err) reader, err := zip.NewReader(bytes.NewReader(body), int64(len(body))) @@ -345,6 +355,53 @@ func TestOpenDeploymentPackageRequiresActiveConfigSnapshot(t *testing.T) { assert.Equal(t, "index.html", reader.File[0].Name) } +func TestProjectLatestRejectsWhenNotOnActiveConfigOrNotActive(t *testing.T) { + cleanup := setupPagesTestDB(t) + defer cleanup() + _, disableStorage := setupPagesStorageMock(t) + defer disableStorage() + ctx := context.Background() + + project, err := CreateProject(ctx, Input{Name: "Gate", Slug: "gate", Enabled: true}) + require.NoError(t, err) + d1, err := UploadDeployment(ctx, project.ID, testPagesMultipartFile(t, "a.zip", testPagesZip(t, map[string]string{"index.html": "a"})), "root") + require.NoError(t, err) + d2, err := UploadDeployment(ctx, project.ID, testPagesMultipartFile(t, "b.zip", testPagesZip(t, map[string]string{"index.html": "b"})), "root") + require.NoError(t, err) + _, err = ActivateDeployment(ctx, project.ID, d2.ID) + require.NoError(t, err) + + // No active main config → reject. + _, _, err = GetProjectLatestPackageHash(ctx, project.ID) + require.Error(t, err) + + require.NoError(t, db.DB(ctx).Create(&model.ConfigVersion{ + Version: "v-gate", + SnapshotJSON: fmt.Sprintf( + `{"routes":[{"upstream_type":"pages","pages_project_id":%d,"pages_deployment":{"project_id":%d,"deployment_id":%d}}]}`, + project.ID, project.ID, d2.ID, + ), + SupportFilesJSON: "[]", + Checksum: "c-gate", + IsActive: true, + CreatedBy: "test", + }).Error) + + // Non-active historical deployment must not be downloadable. + _, err = OpenDeploymentPackage(ctx, d1.ID) + require.Error(t, err) + + // Active latest works. + _, hash, err := GetProjectLatestPackageHash(ctx, project.ID) + require.NoError(t, err) + assert.NotEmpty(t, hash) + + // Disabled project rejects. + require.NoError(t, db.DB(ctx).Model(&model.PagesProject{}).Where("id = ?", project.ID).Update("enabled", false).Error) + _, _, err = GetProjectLatestPackageHash(ctx, project.ID) + require.Error(t, err) +} + func testPagesZip(t *testing.T, files map[string]string) []byte { t.Helper() diff --git a/internal/apps/openflare/pages/rebind.go b/internal/apps/openflare/pages/rebind.go index 1e177c98..d1950c9e 100644 --- a/internal/apps/openflare/pages/rebind.go +++ b/internal/apps/openflare/pages/rebind.go @@ -203,11 +203,7 @@ func buildLivePagesDeployment(project *model.PagesProject, active *model.PagesDe APIProxyPath: strings.TrimSpace(project.APIProxyPath), APIProxyPass: strings.TrimSpace(project.APIProxyPass), APIProxyRewrite: strings.TrimSpace(project.APIProxyRewrite), - LocalRoot: fmt.Sprintf( - "%s/deployments/%d/current", - openrestyrender.PagesDirPlaceholder, - active.ID, - ), + LocalRoot: openrestyrender.PagesProjectLocalRoot(project.ID), } } diff --git a/internal/apps/openflare/pages/rebind_test.go b/internal/apps/openflare/pages/rebind_test.go index 30b3f090..c37602d1 100644 --- a/internal/apps/openflare/pages/rebind_test.go +++ b/internal/apps/openflare/pages/rebind_test.go @@ -59,7 +59,7 @@ func TestRebindSnapshotPagesToCurrentActive(t *testing.T) { "project_id": project.ID, "deployment_id": old.ID, "checksum": "old-checksum", - "local_root": "__OPENFLARE_PAGES_DIR__/deployments/1/current", + "local_root": "__OPENFLARE_PAGES_DIR__/projects/1/current", }, "extra_keep_me": "yes", }, diff --git a/internal/router/v1/openflare/register_agent.go b/internal/router/v1/openflare/register_agent.go index fe59132a..28198836 100644 --- a/internal/router/v1/openflare/register_agent.go +++ b/internal/router/v1/openflare/register_agent.go @@ -25,6 +25,8 @@ func registerAgentRoutes(apiV1Router *gin.RouterGroup) { authorizedRoute.GET("/config-versions/active", agent.GetActiveConfigHandler) authorizedRoute.GET("/pages/deployments/:deployment_id/hash", agent.GetPagesDeploymentHashHandler) authorizedRoute.GET("/pages/deployments/:deployment_id/package", agent.DownloadPagesPackageHandler) + authorizedRoute.GET("/pages/projects/:project_id/latest/hash", agent.GetPagesProjectLatestHashHandler) + authorizedRoute.GET("/pages/projects/:project_id/latest/package", agent.DownloadPagesProjectLatestPackageHandler) authorizedRoute.POST("/waf/ip-groups/sync", agent.SyncWAFIPGroupsHandler) authorizedRoute.POST("/apply-logs", agent.ReportApplyLogHandler) } diff --git a/pkg/protocol/agent.go b/pkg/protocol/agent.go index b79311e4..93ba33a5 100644 --- a/pkg/protocol/agent.go +++ b/pkg/protocol/agent.go @@ -220,3 +220,11 @@ type PagesDeploymentHashResponse struct { DeploymentID uint `json:"deployment_id"` Hash string `json:"hash"` } + +// PagesProjectLatestHashResponse is the hash of a project's currently active Pages deployment. +// Agents poll this like a "latest" pointer without caring about historical deployment IDs. +type PagesProjectLatestHashResponse struct { + ProjectID uint `json:"project_id"` + DeploymentID uint `json:"deployment_id"` + Hash string `json:"hash"` +} diff --git a/pkg/render/openresty/render_test.go b/pkg/render/openresty/render_test.go index 445d2159..098a6200 100644 --- a/pkg/render/openresty/render_test.go +++ b/pkg/render/openresty/render_test.go @@ -345,7 +345,7 @@ func TestRenderRouteConfigPagesWithoutSPAFallbackServesRoot(t *testing.T) { UpstreamType: "pages", EnableHTTPS: false, PagesDeployment: &PagesDeployment{ - LocalRoot: "/data/var/lib/openflare/pages/deployments/1/current", + LocalRoot: "/data/var/lib/openflare/pages/projects/1/current", EntryFile: "index.html", SPAFallbackEnabled: false, }, @@ -378,7 +378,7 @@ func TestRenderRouteConfigPagesWithSPAFallbackServesRoot(t *testing.T) { UpstreamType: "pages", EnableHTTPS: false, PagesDeployment: &PagesDeployment{ - LocalRoot: "/data/var/lib/openflare/pages/deployments/1/current", + LocalRoot: "/data/var/lib/openflare/pages/projects/1/current", EntryFile: "index.html", SPAFallbackEnabled: true, SPAFallbackPath: "/index.html", diff --git a/pkg/render/openresty/types.go b/pkg/render/openresty/types.go index 1dee3df5..cb9256d5 100644 --- a/pkg/render/openresty/types.go +++ b/pkg/render/openresty/types.go @@ -1,6 +1,9 @@ package openresty -import "encoding/json" +import ( + "encoding/json" + "fmt" +) // Placeholder constants used as sentinel values in rendered OpenResty config // files; the deploy process replaces them with real paths before reload. @@ -163,6 +166,10 @@ type Route struct { // PagesDeployment holds the static-site deployment parameters for a Pages-type // route, including local root, entry file, SPA fallback, and API proxy options. +// +// LocalRoot is anchored on ProjectID (projects/{id}/current), not a specific +// deployment ID, so Agents can switch active packages without re-publishing +// main config / reloading OpenResty root paths. type PagesDeployment struct { ProjectID uint `json:"project_id"` ProjectSlug string `json:"project_slug"` @@ -179,6 +186,14 @@ type PagesDeployment struct { LocalRoot string `json:"local_root"` } +// PagesProjectLocalRoot returns the Agent-local root for a Pages project. +func PagesProjectLocalRoot(projectID uint) string { + if projectID == 0 { + return PagesDirPlaceholder + } + return fmt.Sprintf("%s/projects/%d/current", PagesDirPlaceholder, projectID) +} + // WAFRuleGraph is the compact graph executed by the OpenResty WAF runtime. type WAFRuleGraph struct { Entry string `json:"entry"`