mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-09 00:56:37 +08:00
修复 Agent 部署 Pages 问题
This commit is contained in:
@@ -30,6 +30,10 @@ sidebar: false
|
|||||||
|
|
||||||
- 修复 Pages 上传或节点同步时报 `pages file size out of bounds`:允许 ZIP 包内的 0 字节文件,并兼容未声明解压大小的 ZIP 条目。
|
- 修复 Pages 上传或节点同步时报 `pages file size out of bounds`:允许 ZIP 包内的 0 字节文件,并兼容未声明解压大小的 ZIP 条目。
|
||||||
|
|
||||||
|
- 修复 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`。
|
||||||
|
|
||||||
- 修复节点详情 OpenResty 连接数与吞吐显示为「—」:节点可观测 API 将 OpenResty 观测数据合并进 `metric_snapshots`;指标文案改为「请求/分钟」(近 60 秒窗口),连接数为 0 时正常显示 0。
|
- 修复节点详情 OpenResty 连接数与吞吐显示为「—」:节点可观测 API 将 OpenResty 观测数据合并进 `metric_snapshots`;指标文案改为「请求/分钟」(近 60 秒窗口),连接数为 0 时正常显示 0。
|
||||||
|
|
||||||
- 修复仪表盘「24 小时请求趋势」摘要误显示当前小时请求量/错误量:改为汇总近 24 小时总量。
|
- 修复仪表盘「24 小时请求趋势」摘要误显示当前小时请求量/错误量:改为汇总近 24 小时总量。
|
||||||
|
|||||||
@@ -4316,6 +4316,64 @@ const docTemplate = `{
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"/api/v1/agent/pages/deployments/{deployment_id}/hash": {
|
||||||
|
"get": {
|
||||||
|
"security": [
|
||||||
|
{
|
||||||
|
"AgentTokenAuth": []
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"description": "返回 upload 框架记录的 SHA-256 哈希,供 Agent 对比本地缓存并按需拉取部署包",
|
||||||
|
"produces": [
|
||||||
|
"application/json"
|
||||||
|
],
|
||||||
|
"tags": [
|
||||||
|
"openflare-agent"
|
||||||
|
],
|
||||||
|
"summary": "查询 Pages 部署包哈希",
|
||||||
|
"parameters": [
|
||||||
|
{
|
||||||
|
"type": "integer",
|
||||||
|
"description": "部署 ID",
|
||||||
|
"name": "deployment_id",
|
||||||
|
"in": "path",
|
||||||
|
"required": true
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"responses": {
|
||||||
|
"200": {
|
||||||
|
"description": "OK",
|
||||||
|
"schema": {
|
||||||
|
"allOf": [
|
||||||
|
{
|
||||||
|
"$ref": "#/definitions/response.Any"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"data": {
|
||||||
|
"$ref": "#/definitions/protocol.PagesDeploymentHashResponse"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"400": {
|
||||||
|
"description": "参数错误",
|
||||||
|
"schema": {
|
||||||
|
"$ref": "#/definitions/response.Any"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"401": {
|
||||||
|
"description": "Token 无效",
|
||||||
|
"schema": {
|
||||||
|
"$ref": "#/definitions/response.Any"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
"/api/v1/agent/pages/deployments/{deployment_id}/package": {
|
"/api/v1/agent/pages/deployments/{deployment_id}/package": {
|
||||||
"get": {
|
"get": {
|
||||||
"security": [
|
"security": [
|
||||||
@@ -16634,6 +16692,17 @@ const docTemplate = `{
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"protocol.PagesDeploymentHashResponse": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"deployment_id": {
|
||||||
|
"type": "integer"
|
||||||
|
},
|
||||||
|
"hash": {
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
"protocol.RelayConfig": {
|
"protocol.RelayConfig": {
|
||||||
"type": "object",
|
"type": "object",
|
||||||
"properties": {
|
"properties": {
|
||||||
|
|||||||
@@ -4309,6 +4309,64 @@
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"/api/v1/agent/pages/deployments/{deployment_id}/hash": {
|
||||||
|
"get": {
|
||||||
|
"security": [
|
||||||
|
{
|
||||||
|
"AgentTokenAuth": []
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"description": "返回 upload 框架记录的 SHA-256 哈希,供 Agent 对比本地缓存并按需拉取部署包",
|
||||||
|
"produces": [
|
||||||
|
"application/json"
|
||||||
|
],
|
||||||
|
"tags": [
|
||||||
|
"openflare-agent"
|
||||||
|
],
|
||||||
|
"summary": "查询 Pages 部署包哈希",
|
||||||
|
"parameters": [
|
||||||
|
{
|
||||||
|
"type": "integer",
|
||||||
|
"description": "部署 ID",
|
||||||
|
"name": "deployment_id",
|
||||||
|
"in": "path",
|
||||||
|
"required": true
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"responses": {
|
||||||
|
"200": {
|
||||||
|
"description": "OK",
|
||||||
|
"schema": {
|
||||||
|
"allOf": [
|
||||||
|
{
|
||||||
|
"$ref": "#/definitions/response.Any"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"data": {
|
||||||
|
"$ref": "#/definitions/protocol.PagesDeploymentHashResponse"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"400": {
|
||||||
|
"description": "参数错误",
|
||||||
|
"schema": {
|
||||||
|
"$ref": "#/definitions/response.Any"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"401": {
|
||||||
|
"description": "Token 无效",
|
||||||
|
"schema": {
|
||||||
|
"$ref": "#/definitions/response.Any"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
"/api/v1/agent/pages/deployments/{deployment_id}/package": {
|
"/api/v1/agent/pages/deployments/{deployment_id}/package": {
|
||||||
"get": {
|
"get": {
|
||||||
"security": [
|
"security": [
|
||||||
@@ -16627,6 +16685,17 @@
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"protocol.PagesDeploymentHashResponse": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"deployment_id": {
|
||||||
|
"type": "integer"
|
||||||
|
},
|
||||||
|
"hash": {
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
"protocol.RelayConfig": {
|
"protocol.RelayConfig": {
|
||||||
"type": "object",
|
"type": "object",
|
||||||
"properties": {
|
"properties": {
|
||||||
|
|||||||
@@ -2588,6 +2588,13 @@ definitions:
|
|||||||
relay_node_id:
|
relay_node_id:
|
||||||
type: string
|
type: string
|
||||||
type: object
|
type: object
|
||||||
|
protocol.PagesDeploymentHashResponse:
|
||||||
|
properties:
|
||||||
|
deployment_id:
|
||||||
|
type: integer
|
||||||
|
hash:
|
||||||
|
type: string
|
||||||
|
type: object
|
||||||
protocol.RelayConfig:
|
protocol.RelayConfig:
|
||||||
properties:
|
properties:
|
||||||
auth_token:
|
auth_token:
|
||||||
@@ -6502,6 +6509,40 @@ paths:
|
|||||||
summary: 注册或发现 Agent 节点
|
summary: 注册或发现 Agent 节点
|
||||||
tags:
|
tags:
|
||||||
- openflare-agent
|
- openflare-agent
|
||||||
|
/api/v1/agent/pages/deployments/{deployment_id}/hash:
|
||||||
|
get:
|
||||||
|
description: 返回 upload 框架记录的 SHA-256 哈希,供 Agent 对比本地缓存并按需拉取部署包
|
||||||
|
parameters:
|
||||||
|
- description: 部署 ID
|
||||||
|
in: path
|
||||||
|
name: deployment_id
|
||||||
|
required: true
|
||||||
|
type: integer
|
||||||
|
produces:
|
||||||
|
- application/json
|
||||||
|
responses:
|
||||||
|
"200":
|
||||||
|
description: OK
|
||||||
|
schema:
|
||||||
|
allOf:
|
||||||
|
- $ref: '#/definitions/response.Any'
|
||||||
|
- properties:
|
||||||
|
data:
|
||||||
|
$ref: '#/definitions/protocol.PagesDeploymentHashResponse'
|
||||||
|
type: object
|
||||||
|
"400":
|
||||||
|
description: 参数错误
|
||||||
|
schema:
|
||||||
|
$ref: '#/definitions/response.Any'
|
||||||
|
"401":
|
||||||
|
description: Token 无效
|
||||||
|
schema:
|
||||||
|
$ref: '#/definitions/response.Any'
|
||||||
|
security:
|
||||||
|
- AgentTokenAuth: []
|
||||||
|
summary: 查询 Pages 部署包哈希
|
||||||
|
tags:
|
||||||
|
- openflare-agent
|
||||||
/api/v1/agent/pages/deployments/{deployment_id}/package:
|
/api/v1/agent/pages/deployments/{deployment_id}/package:
|
||||||
get:
|
get:
|
||||||
description: 流式下载指定部署的静态资源压缩包,供 Agent 边缘分发
|
description: 流式下载指定部署的静态资源压缩包,供 Agent 边缘分发
|
||||||
|
|||||||
@@ -86,6 +86,18 @@ func (c *Client) SyncWAFIPGroups(ctx context.Context, payload protocol.WAFIPGrou
|
|||||||
return &resp.Data, nil
|
return &resp.Data, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// GetPagesDeploymentHash returns the upload SHA-256 hash for the given Pages deployment ID.
|
||||||
|
func (c *Client) GetPagesDeploymentHash(ctx context.Context, deploymentID uint) (string, error) {
|
||||||
|
resp := protocol.APIResponse[protocol.PagesDeploymentHashResponse]{}
|
||||||
|
if err := c.base.GetJSON(ctx, fmt.Sprintf("/api/v1/agent/pages/deployments/%d/hash", deploymentID), &resp); err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
if err := edgehttp.APIError(resp.ErrorMsg); err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
return resp.Data.Hash, nil
|
||||||
|
}
|
||||||
|
|
||||||
// DownloadPagesDeploymentPackage downloads the deployment package for the given Pages deployment ID.
|
// DownloadPagesDeploymentPackage downloads the deployment package for the given Pages deployment ID.
|
||||||
func (c *Client) DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error) {
|
func (c *Client) DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error) {
|
||||||
res, err := c.base.DoRaw(ctx, http.MethodGet, fmt.Sprintf("/api/v1/agent/pages/deployments/%d/package", deploymentID), nil)
|
res, err := c.base.DoRaw(ctx, http.MethodGet, fmt.Sprintf("/api/v1/agent/pages/deployments/%d/package", deploymentID), nil)
|
||||||
|
|||||||
@@ -72,6 +72,9 @@ type WAFIPGroupSyncResponse = pkgprotocol.WAFIPGroupSyncResponse
|
|||||||
// SupportFile is an alias for pkgprotocol.SupportFile.
|
// SupportFile is an alias for pkgprotocol.SupportFile.
|
||||||
type SupportFile = pkgprotocol.SupportFile
|
type SupportFile = pkgprotocol.SupportFile
|
||||||
|
|
||||||
|
// PagesDeploymentHashResponse is an alias for pkgprotocol.PagesDeploymentHashResponse.
|
||||||
|
type PagesDeploymentHashResponse = pkgprotocol.PagesDeploymentHashResponse
|
||||||
|
|
||||||
const (
|
const (
|
||||||
// WSMessageTypeStatus is an alias for pkgprotocol.WSMessageTypeStatus.
|
// WSMessageTypeStatus is an alias for pkgprotocol.WSMessageTypeStatus.
|
||||||
WSMessageTypeStatus = pkgprotocol.WSMessageTypeStatus
|
WSMessageTypeStatus = pkgprotocol.WSMessageTypeStatus
|
||||||
|
|||||||
@@ -15,22 +15,30 @@ const (
|
|||||||
nodeIDRandomBytes = 8
|
nodeIDRandomBytes = 8
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// PagesDeployment records a Pages release referenced by the active config.
|
||||||
|
type PagesDeployment struct {
|
||||||
|
DeploymentID uint `json:"deployment_id"`
|
||||||
|
Hash string `json:"hash"`
|
||||||
|
Checksum string `json:"checksum,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
// Snapshot represents the state of the agent at a given point in time.
|
// Snapshot represents the state of the agent at a given point in time.
|
||||||
type Snapshot struct {
|
type Snapshot struct {
|
||||||
NodeID string `json:"node_id"`
|
NodeID string `json:"node_id"`
|
||||||
CurrentVersion string `json:"current_version"`
|
CurrentVersion string `json:"current_version"`
|
||||||
CurrentChecksum string `json:"current_checksum"`
|
CurrentChecksum string `json:"current_checksum"`
|
||||||
BlockedVersion string `json:"blocked_version"`
|
PagesDeployments []PagesDeployment `json:"pages_deployments"`
|
||||||
BlockedChecksum string `json:"blocked_checksum"`
|
BlockedVersion string `json:"blocked_version"`
|
||||||
BlockedReason string `json:"blocked_reason"`
|
BlockedChecksum string `json:"blocked_checksum"`
|
||||||
LastError string `json:"last_error"`
|
BlockedReason string `json:"blocked_reason"`
|
||||||
OpenrestyStatus string `json:"openresty_status"`
|
LastError string `json:"last_error"`
|
||||||
OpenrestyMessage string `json:"openresty_message"`
|
OpenrestyStatus string `json:"openresty_status"`
|
||||||
LastProfileFingerprint string `json:"last_profile_fingerprint"`
|
OpenrestyMessage string `json:"openresty_message"`
|
||||||
LastCPUStatTotal uint64 `json:"last_cpu_stat_total"`
|
LastProfileFingerprint string `json:"last_profile_fingerprint"`
|
||||||
LastCPUStatIdle uint64 `json:"last_cpu_stat_idle"`
|
LastCPUStatTotal uint64 `json:"last_cpu_stat_total"`
|
||||||
LastMetricAtUnix int64 `json:"last_metric_at_unix"`
|
LastCPUStatIdle uint64 `json:"last_cpu_stat_idle"`
|
||||||
AccessLogOffset int64 `json:"access_log_offset"`
|
LastMetricAtUnix int64 `json:"last_metric_at_unix"`
|
||||||
|
AccessLogOffset int64 `json:"access_log_offset"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// Store manages the storage and retrieval of the agent state snapshot.
|
// Store manages the storage and retrieval of the agent state snapshot.
|
||||||
|
|||||||
@@ -18,6 +18,7 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/agent/protocol"
|
"github.com/Rain-kl/Wavelet/internal/apps/agent/protocol"
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/apps/agent/state"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -45,10 +46,85 @@ type pagesDeploymentMarker struct {
|
|||||||
Checksum string `json:"checksum"`
|
Checksum string `json:"checksum"`
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) syncPagesDeployments(ctx context.Context, config *protocol.ActiveConfigResponse) error {
|
func pagesDeploymentStateHash(item state.PagesDeployment) string {
|
||||||
deployments, err := referencedPagesDeployments(config)
|
if hash := strings.TrimSpace(item.Hash); hash != "" {
|
||||||
if err != nil {
|
return hash
|
||||||
return err
|
}
|
||||||
|
return strings.TrimSpace(item.Checksum)
|
||||||
|
}
|
||||||
|
|
||||||
|
func snapshotPagesDeployments(snapshot *state.Snapshot) []pagesDeploymentSource {
|
||||||
|
if snapshot == nil || snapshot.PagesDeployments == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
result := make([]pagesDeploymentSource, 0, len(snapshot.PagesDeployments))
|
||||||
|
for _, item := range snapshot.PagesDeployments {
|
||||||
|
result = append(result, pagesDeploymentSource{
|
||||||
|
DeploymentID: item.DeploymentID,
|
||||||
|
Checksum: pagesDeploymentStateHash(item),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
return result
|
||||||
|
}
|
||||||
|
|
||||||
|
func setSnapshotPagesDeployments(snapshot *state.Snapshot, deployments []pagesDeploymentSource) {
|
||||||
|
if snapshot == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if len(deployments) == 0 {
|
||||||
|
snapshot.PagesDeployments = []state.PagesDeployment{}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
snapshot.PagesDeployments = make([]state.PagesDeployment, len(deployments))
|
||||||
|
for i, deployment := range deployments {
|
||||||
|
snapshot.PagesDeployments[i] = state.PagesDeployment{
|
||||||
|
DeploymentID: deployment.DeploymentID,
|
||||||
|
Hash: strings.TrimSpace(deployment.Checksum),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func updateSnapshotPagesDeploymentHash(snapshot *state.Snapshot, deployment pagesDeploymentSource) {
|
||||||
|
if snapshot == nil || snapshot.PagesDeployments == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
hash := strings.TrimSpace(deployment.Checksum)
|
||||||
|
for i := range snapshot.PagesDeployments {
|
||||||
|
if snapshot.PagesDeployments[i].DeploymentID != deployment.DeploymentID {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
snapshot.PagesDeployments[i].Hash = hash
|
||||||
|
snapshot.PagesDeployments[i].Checksum = ""
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func pagesDiscoveryNeeded(snapshot *state.Snapshot) bool {
|
||||||
|
return snapshot == nil || snapshot.PagesDeployments == nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func pagesSyncNeeded(snapshot *state.Snapshot) bool {
|
||||||
|
return snapshot != nil && snapshot.PagesDeployments != nil && len(snapshot.PagesDeployments) > 0
|
||||||
|
}
|
||||||
|
|
||||||
|
func pagesReconcileNeeded(snapshot *state.Snapshot) bool {
|
||||||
|
if pagesDiscoveryNeeded(snapshot) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return pagesSyncNeeded(snapshot)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Service) syncPagesDeployments(ctx context.Context, snapshot *state.Snapshot, config *protocol.ActiveConfigResponse) error {
|
||||||
|
var deployments []pagesDeploymentSource
|
||||||
|
var err error
|
||||||
|
if config != nil {
|
||||||
|
deployments, err = referencedPagesDeployments(config)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
setSnapshotPagesDeployments(snapshot, deployments)
|
||||||
|
} else {
|
||||||
|
deployments = snapshotPagesDeployments(snapshot)
|
||||||
}
|
}
|
||||||
if len(deployments) == 0 {
|
if len(deployments) == 0 {
|
||||||
return nil
|
return nil
|
||||||
@@ -57,32 +133,60 @@ func (s *Service) syncPagesDeployments(ctx context.Context, config *protocol.Act
|
|||||||
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 deployments")
|
||||||
}
|
}
|
||||||
for _, deployment := range deployments {
|
for _, deployment := range deployments {
|
||||||
if err := s.ensurePagesDeployment(ctx, deployment); err != nil {
|
if err := s.ensurePagesDeployment(ctx, snapshot, deployment); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) ensurePagesDeployment(ctx context.Context, deployment pagesDeploymentSource) error {
|
func (s *Service) ensurePagesDeployment(ctx context.Context, snapshot *state.Snapshot, deployment pagesDeploymentSource) error {
|
||||||
currentDir := pagesCurrentDir(s.pagesDir, deployment.DeploymentID)
|
serverHash, err := s.client.GetPagesDeploymentHash(ctx, deployment.DeploymentID)
|
||||||
if markerMatches(currentDir, deployment) {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
packageBytes, err := s.client.DownloadPagesDeploymentPackage(ctx, deployment.DeploymentID)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("download Pages deployment %d: %w", deployment.DeploymentID, err)
|
return fmt.Errorf("fetch Pages deployment %d hash: %w", deployment.DeploymentID, err)
|
||||||
}
|
}
|
||||||
if got := checksumBytes(packageBytes); got != deployment.Checksum {
|
serverHash = strings.TrimSpace(serverHash)
|
||||||
return fmt.Errorf("pages deployment %d checksum mismatch: expected %s, got %s", deployment.DeploymentID, deployment.Checksum, got)
|
if serverHash == "" {
|
||||||
|
return fmt.Errorf("pages deployment %d hash is empty", deployment.DeploymentID)
|
||||||
}
|
}
|
||||||
releaseDir := pagesReleaseDir(s.pagesDir, deployment.DeploymentID, deployment.Checksum)
|
effective := pagesDeploymentSource{
|
||||||
if !markerMatches(releaseDir, deployment) {
|
DeploymentID: deployment.DeploymentID,
|
||||||
if err := extractPagesPackage(packageBytes, releaseDir, deployment); err != nil {
|
Checksum: serverHash,
|
||||||
return err
|
}
|
||||||
|
updateSnapshotPagesDeploymentHash(snapshot, effective)
|
||||||
|
|
||||||
|
releaseDir := pagesReleaseDir(s.pagesDir, effective.DeploymentID, effective.Checksum)
|
||||||
|
if pagesReleaseReady(releaseDir, effective) {
|
||||||
|
return switchPagesCurrentDir(s.pagesDir, effective.DeploymentID, releaseDir)
|
||||||
|
}
|
||||||
|
packageBytes, err := s.client.DownloadPagesDeploymentPackage(ctx, effective.DeploymentID)
|
||||||
|
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 {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return switchPagesCurrentDir(s.pagesDir, effective.DeploymentID, releaseDir)
|
||||||
|
}
|
||||||
|
|
||||||
|
func pagesReleaseReady(dir string, deployment pagesDeploymentSource) bool {
|
||||||
|
if !markerMatches(dir, deployment) {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
entries, err := os.ReadDir(dir) //nolint:gosec // dir is managed PagesDir
|
||||||
|
if err != nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
for _, entry := range entries {
|
||||||
|
if entry.Name() == ".openflare-pages.json" {
|
||||||
|
continue
|
||||||
}
|
}
|
||||||
|
return true
|
||||||
}
|
}
|
||||||
return switchPagesCurrentDir(s.pagesDir, deployment.DeploymentID, releaseDir)
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
func referencedPagesDeployments(config *protocol.ActiveConfigResponse) ([]pagesDeploymentSource, error) {
|
func referencedPagesDeployments(config *protocol.ActiveConfigResponse) ([]pagesDeploymentSource, error) {
|
||||||
|
|||||||
@@ -29,6 +29,7 @@ const (
|
|||||||
// ConfigClient is the interface for communicating with the server control plane.
|
// ConfigClient is the interface for communicating with the server control plane.
|
||||||
type ConfigClient interface {
|
type ConfigClient interface {
|
||||||
GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigResponse, error)
|
GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigResponse, error)
|
||||||
|
GetPagesDeploymentHash(ctx context.Context, deploymentID uint) (string, error)
|
||||||
DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error)
|
DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error)
|
||||||
ReportApplyLog(ctx context.Context, payload protocol.ApplyLogPayload) error
|
ReportApplyLog(ctx context.Context, payload protocol.ApplyLogPayload) error
|
||||||
SyncWAFIPGroups(ctx context.Context, payload protocol.WAFIPGroupSyncRequest) (*protocol.WAFIPGroupSyncResponse, error)
|
SyncWAFIPGroups(ctx context.Context, payload protocol.WAFIPGroupSyncRequest) (*protocol.WAFIPGroupSyncResponse, error)
|
||||||
@@ -152,6 +153,9 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
|
|||||||
clearBlockedTarget(snapshot)
|
clearBlockedTarget(snapshot)
|
||||||
}
|
}
|
||||||
if snapshot.CurrentVersion == config.Version && snapshot.CurrentChecksum == config.Checksum && !startup {
|
if snapshot.CurrentVersion == config.Version && snapshot.CurrentChecksum == config.Checksum && !startup {
|
||||||
|
if err := s.syncPagesDeployments(ctx, snapshot, config); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
slog.Debug("skipping apply because state already records target version/checksum", "version", config.Version, "checksum", config.Checksum)
|
slog.Debug("skipping apply because state already records target version/checksum", "version", config.Version, "checksum", config.Checksum)
|
||||||
return s.stateStore.Save(snapshot)
|
return s.stateStore.Save(snapshot)
|
||||||
}
|
}
|
||||||
@@ -163,7 +167,7 @@ func (s *Service) applyRenderedConfig(ctx context.Context, mode string, snapshot
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := s.syncPagesDeployments(ctx, config); err != nil {
|
if err := s.syncPagesDeployments(ctx, snapshot, config); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
mainConfigChecksum := checksumString(rendered.mainConfig)
|
mainConfigChecksum := checksumString(rendered.mainConfig)
|
||||||
|
|||||||
@@ -31,7 +31,9 @@ type fakeClient struct {
|
|||||||
config protocol.ActiveConfigResponse
|
config protocol.ActiveConfigResponse
|
||||||
reports []protocol.ApplyLogPayload
|
reports []protocol.ApplyLogPayload
|
||||||
pagesPackages map[uint][]byte
|
pagesPackages map[uint][]byte
|
||||||
|
pagesHashes map[uint]string
|
||||||
fetchCalls int
|
fetchCalls int
|
||||||
|
hashCalls int
|
||||||
}
|
}
|
||||||
|
|
||||||
type fakeManager struct {
|
type fakeManager struct {
|
||||||
@@ -76,6 +78,21 @@ func (f *fakeClient) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfi
|
|||||||
return &f.config, nil
|
return &f.config, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (f *fakeClient) GetPagesDeploymentHash(ctx context.Context, deploymentID uint) (string, error) {
|
||||||
|
f.hashCalls++
|
||||||
|
if f.pagesHashes != nil {
|
||||||
|
if hash, ok := f.pagesHashes[deploymentID]; ok {
|
||||||
|
return hash, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if f.pagesPackages != nil {
|
||||||
|
if packageBytes, ok := f.pagesPackages[deploymentID]; ok {
|
||||||
|
return testBytesChecksum(packageBytes), nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return "", fmt.Errorf("missing Pages hash %d", deploymentID)
|
||||||
|
}
|
||||||
|
|
||||||
func (f *fakeClient) DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error) {
|
func (f *fakeClient) DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error) {
|
||||||
if f.pagesPackages == nil {
|
if f.pagesPackages == nil {
|
||||||
return nil, fmt.Errorf("missing Pages package %d", deploymentID)
|
return nil, fmt.Errorf("missing Pages package %d", deploymentID)
|
||||||
@@ -501,8 +518,8 @@ func TestSyncOnceReportsNoopWhenVersionChangesButChecksumMatches(t *testing.T) {
|
|||||||
}); err != nil {
|
}); err != nil {
|
||||||
t.Fatalf("SyncOnce failed: %v", err)
|
t.Fatalf("SyncOnce failed: %v", err)
|
||||||
}
|
}
|
||||||
if client.fetchCalls != 0 {
|
if client.fetchCalls != 1 {
|
||||||
t.Fatalf("expected checksum match to skip config fetch, got %d", client.fetchCalls)
|
t.Fatalf("expected checksum match to fetch active config once for Pages reconciliation, got %d", client.fetchCalls)
|
||||||
}
|
}
|
||||||
if len(manager.applyMainContents) != 0 {
|
if len(manager.applyMainContents) != 0 {
|
||||||
t.Fatal("expected checksum match to skip apply")
|
t.Fatal("expected checksum match to skip apply")
|
||||||
@@ -568,9 +585,10 @@ func TestSyncOnceDoesNotRepeatNoopReportWhenStateAlreadyMatches(t *testing.T) {
|
|||||||
t.Fatalf("EnsureNodeID failed: %v", err)
|
t.Fatalf("EnsureNodeID failed: %v", err)
|
||||||
}
|
}
|
||||||
if err = stateStore.Save(&state.Snapshot{
|
if err = stateStore.Save(&state.Snapshot{
|
||||||
NodeID: nodeID,
|
NodeID: nodeID,
|
||||||
CurrentVersion: "20260309-003",
|
CurrentVersion: "20260309-003",
|
||||||
CurrentChecksum: "checksum-3",
|
CurrentChecksum: "checksum-3",
|
||||||
|
PagesDeployments: []state.PagesDeployment{},
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
t.Fatalf("failed to seed state: %v", err)
|
t.Fatalf("failed to seed state: %v", err)
|
||||||
}
|
}
|
||||||
@@ -583,6 +601,9 @@ func TestSyncOnceDoesNotRepeatNoopReportWhenStateAlreadyMatches(t *testing.T) {
|
|||||||
}); err != nil {
|
}); err != nil {
|
||||||
t.Fatalf("SyncOnce failed: %v", err)
|
t.Fatalf("SyncOnce failed: %v", err)
|
||||||
}
|
}
|
||||||
|
if client.fetchCalls != 0 {
|
||||||
|
t.Fatalf("expected no active config fetch when Pages releases are already reconciled, got %d", client.fetchCalls)
|
||||||
|
}
|
||||||
if len(client.reports) != 0 {
|
if len(client.reports) != 0 {
|
||||||
t.Fatalf("expected matching state to skip duplicate noop report, got %+v", client.reports)
|
t.Fatalf("expected matching state to skip duplicate noop report, got %+v", client.reports)
|
||||||
}
|
}
|
||||||
@@ -907,9 +928,10 @@ func TestSyncOnceSkipsFetchWhenHeartbeatChecksumMatches(t *testing.T) {
|
|||||||
t.Fatalf("EnsureNodeID failed: %v", err)
|
t.Fatalf("EnsureNodeID failed: %v", err)
|
||||||
}
|
}
|
||||||
if err = stateStore.Save(&state.Snapshot{
|
if err = stateStore.Save(&state.Snapshot{
|
||||||
NodeID: nodeID,
|
NodeID: nodeID,
|
||||||
CurrentVersion: client.config.Version,
|
CurrentVersion: client.config.Version,
|
||||||
CurrentChecksum: client.config.Checksum,
|
CurrentChecksum: client.config.Checksum,
|
||||||
|
PagesDeployments: []state.PagesDeployment{},
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
t.Fatalf("failed to seed state: %v", err)
|
t.Fatalf("failed to seed state: %v", err)
|
||||||
}
|
}
|
||||||
@@ -923,13 +945,209 @@ func TestSyncOnceSkipsFetchWhenHeartbeatChecksumMatches(t *testing.T) {
|
|||||||
t.Fatalf("SyncOnce failed: %v", err)
|
t.Fatalf("SyncOnce failed: %v", err)
|
||||||
}
|
}
|
||||||
if client.fetchCalls != 0 {
|
if client.fetchCalls != 0 {
|
||||||
t.Fatalf("expected no active config fetch when heartbeat checksum matches, got %d", client.fetchCalls)
|
t.Fatalf("expected no active config fetch when Pages releases are already reconciled, got %d", client.fetchCalls)
|
||||||
|
}
|
||||||
|
if len(manager.applyMainContents) != 0 {
|
||||||
|
t.Fatal("expected checksum match to skip apply")
|
||||||
}
|
}
|
||||||
if len(client.reports) != 0 {
|
if len(client.reports) != 0 {
|
||||||
t.Fatal("expected no apply log when no config change is needed")
|
t.Fatal("expected no apply log when no config change is needed")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestSyncOnceDownloadsPagesDeploymentWhenChecksumMatches(t *testing.T) {
|
||||||
|
packageBytes := testPagesPackage(t, map[string]string{"index.html": "hello"})
|
||||||
|
checksum := testBytesChecksum(packageBytes)
|
||||||
|
client := &fakeClient{
|
||||||
|
config: protocol.ActiveConfigResponse{
|
||||||
|
Version: "20260309-101",
|
||||||
|
Checksum: "pages-config-checksum",
|
||||||
|
SourceConfigJSON: testPagesSourceConfigJSON(7, checksum),
|
||||||
|
CreatedAt: time.Now().Format(time.RFC3339),
|
||||||
|
},
|
||||||
|
pagesPackages: map[uint][]byte{7: packageBytes},
|
||||||
|
}
|
||||||
|
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
|
||||||
|
nodeID, err := stateStore.EnsureNodeID()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("EnsureNodeID failed: %v", err)
|
||||||
|
}
|
||||||
|
if err = stateStore.Save(&state.Snapshot{
|
||||||
|
NodeID: nodeID,
|
||||||
|
CurrentVersion: client.config.Version,
|
||||||
|
CurrentChecksum: client.config.Checksum,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("save state failed: %v", err)
|
||||||
|
}
|
||||||
|
manager := &fakeManager{currentChecksum: client.config.Checksum}
|
||||||
|
service := New(client, manager, stateStore)
|
||||||
|
pagesDir := t.TempDir()
|
||||||
|
service.SetPagesDir(pagesDir)
|
||||||
|
|
||||||
|
if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{
|
||||||
|
Version: client.config.Version,
|
||||||
|
Checksum: client.config.Checksum,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("SyncOnce failed: %v", err)
|
||||||
|
}
|
||||||
|
data, err := os.ReadFile(filepath.Join(pagesDir, "deployments", "7", "current", "index.html"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("expected Pages file to be extracted when checksum already matches: %v", err)
|
||||||
|
}
|
||||||
|
if string(data) != "hello" {
|
||||||
|
t.Fatalf("unexpected Pages file content: %s", string(data))
|
||||||
|
}
|
||||||
|
if len(manager.applyMainContents) != 0 {
|
||||||
|
t.Fatal("expected checksum match to skip OpenResty apply")
|
||||||
|
}
|
||||||
|
if client.fetchCalls != 1 {
|
||||||
|
t.Fatalf("expected one active config fetch on first reconcile, got %d", client.fetchCalls)
|
||||||
|
}
|
||||||
|
|
||||||
|
client.fetchCalls = 0
|
||||||
|
if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{
|
||||||
|
Version: client.config.Version,
|
||||||
|
Checksum: client.config.Checksum,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("second SyncOnce failed: %v", err)
|
||||||
|
}
|
||||||
|
if client.fetchCalls != 0 {
|
||||||
|
t.Fatalf("expected no active config fetch after Pages release is ready, got %d", client.fetchCalls)
|
||||||
|
}
|
||||||
|
snapshot, err := stateStore.Load()
|
||||||
|
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 client.hashCalls == 0 {
|
||||||
|
t.Fatal("expected Pages hash check during reconcile")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSyncOnceRedownloadsPagesDeploymentWhenServerHashChanges(t *testing.T) {
|
||||||
|
initialPackage := testPagesPackage(t, map[string]string{"index.html": "v1"})
|
||||||
|
updatedPackage := testPagesPackage(t, map[string]string{"index.html": "v2"})
|
||||||
|
initialHash := testBytesChecksum(initialPackage)
|
||||||
|
updatedHash := testBytesChecksum(updatedPackage)
|
||||||
|
deploymentID := uint(12)
|
||||||
|
client := &fakeClient{
|
||||||
|
config: protocol.ActiveConfigResponse{
|
||||||
|
Version: "20260309-107",
|
||||||
|
Checksum: "pages-config-checksum",
|
||||||
|
SourceConfigJSON: testPagesSourceConfigJSON(deploymentID, initialHash),
|
||||||
|
CreatedAt: time.Now().Format(time.RFC3339),
|
||||||
|
},
|
||||||
|
pagesPackages: map[uint][]byte{deploymentID: initialPackage},
|
||||||
|
pagesHashes: map[uint]string{deploymentID: initialHash},
|
||||||
|
}
|
||||||
|
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
|
||||||
|
nodeID, err := stateStore.EnsureNodeID()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("EnsureNodeID failed: %v", err)
|
||||||
|
}
|
||||||
|
if err = stateStore.Save(&state.Snapshot{
|
||||||
|
NodeID: nodeID,
|
||||||
|
CurrentVersion: client.config.Version,
|
||||||
|
CurrentChecksum: client.config.Checksum,
|
||||||
|
PagesDeployments: []state.PagesDeployment{{
|
||||||
|
DeploymentID: deploymentID,
|
||||||
|
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,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("seed release failed: %v", err)
|
||||||
|
}
|
||||||
|
if err = switchPagesCurrentDir(pagesDir, deploymentID, releaseDir); err != nil {
|
||||||
|
t.Fatalf("seed current dir failed: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
client.pagesPackages[deploymentID] = updatedPackage
|
||||||
|
client.pagesHashes[deploymentID] = updatedHash
|
||||||
|
manager := &fakeManager{currentChecksum: client.config.Checksum}
|
||||||
|
service := New(client, manager, stateStore)
|
||||||
|
service.SetPagesDir(pagesDir)
|
||||||
|
if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{
|
||||||
|
Version: client.config.Version,
|
||||||
|
Checksum: client.config.Checksum,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("SyncOnce failed: %v", err)
|
||||||
|
}
|
||||||
|
data, err := os.ReadFile(filepath.Join(pagesDir, "deployments", "12", "current", "index.html"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("expected updated Pages file: %v", err)
|
||||||
|
}
|
||||||
|
if string(data) != "v2" {
|
||||||
|
t.Fatalf("unexpected Pages file content after hash change: %s", string(data))
|
||||||
|
}
|
||||||
|
if client.fetchCalls != 0 {
|
||||||
|
t.Fatalf("expected hash-only reconcile without active config fetch, got %d", client.fetchCalls)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSyncOnceRedownloadsPagesDeploymentWhenReleaseDirOnlyHasMarker(t *testing.T) {
|
||||||
|
packageBytes := testPagesPackage(t, map[string]string{"index.html": "hello"})
|
||||||
|
checksum := testBytesChecksum(packageBytes)
|
||||||
|
deploymentID := uint(11)
|
||||||
|
client := &fakeClient{
|
||||||
|
config: protocol.ActiveConfigResponse{
|
||||||
|
Version: "20260309-106",
|
||||||
|
Checksum: "pages-config-checksum",
|
||||||
|
SourceConfigJSON: testPagesSourceConfigJSON(deploymentID, checksum),
|
||||||
|
CreatedAt: time.Now().Format(time.RFC3339),
|
||||||
|
},
|
||||||
|
pagesPackages: map[uint][]byte{deploymentID: packageBytes},
|
||||||
|
}
|
||||||
|
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
|
||||||
|
nodeID, err := stateStore.EnsureNodeID()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("EnsureNodeID failed: %v", err)
|
||||||
|
}
|
||||||
|
if err = stateStore.Save(&state.Snapshot{
|
||||||
|
NodeID: nodeID,
|
||||||
|
CurrentVersion: client.config.Version,
|
||||||
|
CurrentChecksum: client.config.Checksum,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("save state failed: %v", err)
|
||||||
|
}
|
||||||
|
pagesDir := t.TempDir()
|
||||||
|
releaseDir := pagesReleaseDir(pagesDir, deploymentID, 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 {
|
||||||
|
t.Fatalf("write marker failed: %v", err)
|
||||||
|
}
|
||||||
|
if pagesReleaseReady(releaseDir, pagesDeploymentSource{DeploymentID: deploymentID, Checksum: checksum}) {
|
||||||
|
t.Fatal("expected marker-only release dir to be treated as not ready")
|
||||||
|
}
|
||||||
|
|
||||||
|
manager := &fakeManager{currentChecksum: client.config.Checksum}
|
||||||
|
service := New(client, manager, stateStore)
|
||||||
|
service.SetPagesDir(pagesDir)
|
||||||
|
if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{
|
||||||
|
Version: client.config.Version,
|
||||||
|
Checksum: client.config.Checksum,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("SyncOnce failed: %v", err)
|
||||||
|
}
|
||||||
|
data, err := os.ReadFile(filepath.Join(pagesDir, "deployments", "11", "current", "index.html"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("expected Pages file to be extracted after marker-only release dir: %v", err)
|
||||||
|
}
|
||||||
|
if string(data) != "hello" {
|
||||||
|
t.Fatalf("unexpected Pages file content: %s", string(data))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func testPagesPackage(t *testing.T, files map[string]string) []byte {
|
func testPagesPackage(t *testing.T, files map[string]string) []byte {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
var buffer bytes.Buffer
|
var buffer bytes.Buffer
|
||||||
|
|||||||
@@ -68,6 +68,20 @@ func (s *Service) syncMatchingChecksum(ctx context.Context, mode string, startup
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) finishUpToDateSync(ctx context.Context, mode string, snapshot *state.Snapshot, target *protocol.ActiveConfigMeta) error {
|
func (s *Service) finishUpToDateSync(ctx context.Context, mode string, snapshot *state.Snapshot, target *protocol.ActiveConfigMeta) error {
|
||||||
|
if pagesReconcileNeeded(snapshot) {
|
||||||
|
if pagesDiscoveryNeeded(snapshot) {
|
||||||
|
config, err := s.client.GetActiveConfig(ctx)
|
||||||
|
if err != nil {
|
||||||
|
slog.Error("fetch active config failed", "mode", mode, "error", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := s.syncPagesDeployments(ctx, snapshot, config); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
} else if err := s.syncPagesDeployments(ctx, snapshot, nil); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
}
|
||||||
slog.Debug("local openresty config already up to date", "mode", mode, "version", target.Version)
|
slog.Debug("local openresty config already up to date", "mode", mode, "version", target.Version)
|
||||||
if shouldReportNoopApply(snapshot, target.Version, target.Checksum) {
|
if shouldReportNoopApply(snapshot, target.Version, target.Checksum) {
|
||||||
if err := s.reportNoopApply(ctx, snapshot.NodeID, target.Version, target.Checksum, "", "", 0); err != nil {
|
if err := s.reportNoopApply(ctx, snapshot.NodeID, target.Version, target.Checksum, "", "", 0); err != nil {
|
||||||
@@ -97,6 +111,21 @@ func (s *Service) syncMismatchedChecksum(ctx context.Context, mode string, start
|
|||||||
clearBlockedTarget(snapshot)
|
clearBlockedTarget(snapshot)
|
||||||
}
|
}
|
||||||
if snapshot.CurrentVersion == target.Version && snapshot.CurrentChecksum == target.Checksum && !startup {
|
if snapshot.CurrentVersion == target.Version && snapshot.CurrentChecksum == target.Checksum && !startup {
|
||||||
|
if pagesReconcileNeeded(snapshot) {
|
||||||
|
if pagesDiscoveryNeeded(snapshot) {
|
||||||
|
config, fetchErr := s.client.GetActiveConfig(ctx)
|
||||||
|
if fetchErr != nil {
|
||||||
|
slog.Error("fetch active config failed", "mode", mode, "error", fetchErr)
|
||||||
|
return fetchErr
|
||||||
|
}
|
||||||
|
if err := s.syncPagesDeployments(ctx, snapshot, config); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
} else if err := s.syncPagesDeployments(ctx, snapshot, nil); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return s.stateStore.Save(snapshot)
|
||||||
|
}
|
||||||
slog.Debug("skipping config fetch because state already records target version/checksum", "version", target.Version, "checksum", target.Checksum)
|
slog.Debug("skipping config fetch because state already records target version/checksum", "version", target.Version, "checksum", target.Checksum)
|
||||||
return s.stateStore.Save(snapshot)
|
return s.stateStore.Save(snapshot)
|
||||||
}
|
}
|
||||||
@@ -110,6 +139,9 @@ func (s *Service) syncMismatchedChecksum(ctx context.Context, mode string, start
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) handleUpToDateConfig(ctx context.Context, mode string, snapshot *state.Snapshot, config *protocol.ActiveConfigResponse) error {
|
func (s *Service) handleUpToDateConfig(ctx context.Context, mode string, snapshot *state.Snapshot, config *protocol.ActiveConfigResponse) error {
|
||||||
|
if err := s.syncPagesDeployments(ctx, snapshot, config); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
slog.Debug("local openresty config already up to date", "mode", mode, "version", config.Version)
|
slog.Debug("local openresty config already up to date", "mode", mode, "version", config.Version)
|
||||||
if shouldReportNoopApply(snapshot, config.Version, config.Checksum) {
|
if shouldReportNoopApply(snapshot, config.Version, config.Checksum) {
|
||||||
rendered, renderErr := renderActiveConfig(config)
|
rendered, renderErr := renderActiveConfig(config)
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ import (
|
|||||||
"github.com/Rain-kl/Wavelet/internal/apps/openflare/pages"
|
"github.com/Rain-kl/Wavelet/internal/apps/openflare/pages"
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/openflare/websocket"
|
"github.com/Rain-kl/Wavelet/internal/apps/openflare/websocket"
|
||||||
"github.com/Rain-kl/Wavelet/internal/common/response"
|
"github.com/Rain-kl/Wavelet/internal/common/response"
|
||||||
|
"github.com/Rain-kl/Wavelet/pkg/protocol"
|
||||||
"github.com/gin-gonic/gin"
|
"github.com/gin-gonic/gin"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -153,6 +154,32 @@ func ReportApplyLogHandler(c *gin.Context) {
|
|||||||
c.JSON(http.StatusOK, response.OK(log))
|
c.JSON(http.StatusOK, response.OK(log))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// GetPagesDeploymentHashHandler returns the upload SHA-256 hash for a Pages deployment package.
|
||||||
|
// @Summary 查询 Pages 部署包哈希
|
||||||
|
// @Description 返回 upload 框架记录的 SHA-256 哈希,供 Agent 对比本地缓存并按需拉取部署包
|
||||||
|
// @Tags openflare-agent
|
||||||
|
// @Produce json
|
||||||
|
// @Security AgentTokenAuth
|
||||||
|
// @Param deployment_id path int true "部署 ID"
|
||||||
|
// @Success 200 {object} response.Any{data=protocol.PagesDeploymentHashResponse}
|
||||||
|
// @Failure 400 {object} response.Any "参数错误"
|
||||||
|
// @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)
|
||||||
|
if !ok {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
hash, err := pages.GetDeploymentPackageHash(c.Request.Context(), deploymentID)
|
||||||
|
if apiutil.AbortBadRequestOnError(c, err) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
c.JSON(http.StatusOK, response.OK(protocol.PagesDeploymentHashResponse{
|
||||||
|
DeploymentID: deploymentID,
|
||||||
|
Hash: hash,
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
// DownloadPagesPackageHandler streams the Pages deployment artifact to an authenticated agent.
|
// DownloadPagesPackageHandler streams the Pages deployment artifact to an authenticated agent.
|
||||||
// @Summary 下载 Pages 部署包
|
// @Summary 下载 Pages 部署包
|
||||||
// @Description 流式下载指定部署的静态资源压缩包,供 Agent 边缘分发
|
// @Description 流式下载指定部署的静态资源压缩包,供 Agent 边缘分发
|
||||||
|
|||||||
@@ -24,5 +24,6 @@ const (
|
|||||||
errPagesPackagePathEmpty = "pages 部署包路径为空"
|
errPagesPackagePathEmpty = "pages 部署包路径为空"
|
||||||
errPagesPackageUploadMissing = "pages 部署包上传记录不存在"
|
errPagesPackageUploadMissing = "pages 部署包上传记录不存在"
|
||||||
errPagesPackageNotInActiveConfig = "pages 部署尚未进入激活配置"
|
errPagesPackageNotInActiveConfig = "pages 部署尚未进入激活配置"
|
||||||
|
errPagesDeploymentHashMissing = "pages 部署包哈希缺失"
|
||||||
errPagesInvalidSnapshotFormat = "配置快照格式无效"
|
errPagesInvalidSnapshotFormat = "配置快照格式无效"
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -16,10 +16,8 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/upload"
|
"github.com/Rain-kl/Wavelet/internal/apps/upload"
|
||||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/db"
|
"github.com/Rain-kl/Wavelet/internal/db"
|
||||||
"github.com/Rain-kl/Wavelet/internal/model"
|
"github.com/Rain-kl/Wavelet/internal/model"
|
||||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
)
|
)
|
||||||
@@ -354,6 +352,45 @@ func ActivateDeployment(ctx context.Context, projectID uint, deploymentID uint)
|
|||||||
return GetProject(ctx, project.ID)
|
return GetProject(ctx, project.ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// GetDeploymentPackageHash returns the SHA-256 hash of the deployment package from upload storage.
|
||||||
|
func GetDeploymentPackageHash(ctx context.Context, deploymentID uint) (string, error) {
|
||||||
|
deployment, err := model.GetPagesDeploymentByID(ctx, deploymentID)
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
if err = ensureDeploymentInActiveSnapshot(ctx, deployment.ID); err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
if deployment.UploadID == 0 {
|
||||||
|
if err := ensureDeploymentUploadRecord(ctx, deployment); err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
deployment, err = model.GetPagesDeploymentByID(ctx, deploymentID)
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if deployment.UploadID == 0 {
|
||||||
|
hash := strings.TrimSpace(deployment.Checksum)
|
||||||
|
if hash == "" {
|
||||||
|
return "", errors.New(errPagesDeploymentHashMissing)
|
||||||
|
}
|
||||||
|
return hash, nil
|
||||||
|
}
|
||||||
|
uploadRecord, err := upload.GetActiveUpload(ctx, deployment.UploadID)
|
||||||
|
if err != nil {
|
||||||
|
return "", fmt.Errorf("pages 部署包不存在: %w", err)
|
||||||
|
}
|
||||||
|
hash := strings.TrimSpace(uploadRecord.Hash)
|
||||||
|
if hash == "" {
|
||||||
|
hash = strings.TrimSpace(deployment.Checksum)
|
||||||
|
}
|
||||||
|
if hash == "" {
|
||||||
|
return "", errors.New(errPagesDeploymentHashMissing)
|
||||||
|
}
|
||||||
|
return hash, nil
|
||||||
|
}
|
||||||
|
|
||||||
// OpenDeploymentPackage opens the deployment artifact from the upload storage framework.
|
// 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) (*storage.Object, string, error) {
|
||||||
deployment, err := model.GetPagesDeploymentByID(ctx, deploymentID)
|
deployment, err := model.GetPagesDeploymentByID(ctx, deploymentID)
|
||||||
@@ -373,15 +410,7 @@ func OpenDeploymentPackage(ctx context.Context, deploymentID uint) (*storage.Obj
|
|||||||
}
|
}
|
||||||
|
|
||||||
func openDeploymentPackageFromUpload(ctx context.Context, deployment *model.PagesDeployment, fileName string) (*storage.Object, string, error) {
|
func openDeploymentPackageFromUpload(ctx context.Context, deployment *model.PagesDeployment, fileName string) (*storage.Object, string, error) {
|
||||||
uploadRecord, err := repository.GetActiveUploadByID(ctx, deployment.UploadID)
|
obj, _, err := upload.OpenStoredUpload(ctx, deployment.UploadID)
|
||||||
if err != nil {
|
|
||||||
return nil, "", fmt.Errorf("pages 部署包不存在: %w", err)
|
|
||||||
}
|
|
||||||
return openDeploymentPackageFromUploadRecord(ctx, &uploadRecord, fileName)
|
|
||||||
}
|
|
||||||
|
|
||||||
func openDeploymentPackageFromUploadRecord(ctx context.Context, uploadRecord *model.Upload, fileName string) (*storage.Object, string, error) {
|
|
||||||
obj, err := uploadstorage.OpenStoredObject(ctx, uploadRecord)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, "", fmt.Errorf("pages 部署包不存在: %w", err)
|
return nil, "", fmt.Errorf("pages 部署包不存在: %w", err)
|
||||||
}
|
}
|
||||||
@@ -423,7 +452,7 @@ func hydrateLegacyDeploymentUpload(
|
|||||||
return nil, errors.New(errPagesPackagePathEmpty)
|
return nil, errors.New(errPagesPackagePathEmpty)
|
||||||
}
|
}
|
||||||
if deployment.UploadID > 0 {
|
if deployment.UploadID > 0 {
|
||||||
uploadRecord, err := repository.GetActiveUploadByID(ctx, deployment.UploadID)
|
uploadRecord, err := upload.GetActiveUpload(ctx, deployment.UploadID)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
return &uploadRecord, nil
|
return &uploadRecord, nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -31,10 +31,13 @@ var (
|
|||||||
|
|
||||||
// Programmatic ingest API
|
// Programmatic ingest API
|
||||||
var (
|
var (
|
||||||
Ingest = ingest.Ingest
|
Ingest = ingest.Ingest
|
||||||
Remove = ingest.Remove
|
Remove = ingest.Remove
|
||||||
RemoveOwned = ingest.RemoveOwned
|
RemoveOwned = ingest.RemoveOwned
|
||||||
FindByHash = ingest.FindByHash
|
FindByHash = ingest.FindByHash
|
||||||
|
GetActiveUpload = ingest.GetActive
|
||||||
|
OpenStoredUpload = ingest.OpenActive
|
||||||
|
ActiveUploadHash = ingest.ActiveHash
|
||||||
)
|
)
|
||||||
|
|
||||||
// Ingest policy constants
|
// Ingest policy constants
|
||||||
|
|||||||
@@ -0,0 +1,49 @@
|
|||||||
|
// Copyright 2026 Arctel.net
|
||||||
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
|
|
||||||
|
package ingest
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"strings"
|
||||||
|
|
||||||
|
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.
|
||||||
|
func GetActive(ctx context.Context, uploadID uint64) (model.Upload, error) {
|
||||||
|
if uploadID == 0 {
|
||||||
|
return model.Upload{}, errors.New("upload id is required")
|
||||||
|
}
|
||||||
|
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) {
|
||||||
|
record, err := GetActive(ctx, uploadID)
|
||||||
|
if err != nil {
|
||||||
|
return nil, model.Upload{}, err
|
||||||
|
}
|
||||||
|
obj, err := uploadstorage.OpenStoredObject(ctx, &record)
|
||||||
|
if err != nil {
|
||||||
|
return nil, model.Upload{}, err
|
||||||
|
}
|
||||||
|
return obj, record, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ActiveHash returns the SHA-256 hash recorded for an active upload.
|
||||||
|
func ActiveHash(ctx context.Context, uploadID uint64) (string, error) {
|
||||||
|
record, err := GetActive(ctx, uploadID)
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
hash := strings.TrimSpace(record.Hash)
|
||||||
|
if hash == "" {
|
||||||
|
return "", errors.New("upload hash is empty")
|
||||||
|
}
|
||||||
|
return hash, nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,16 @@
|
|||||||
|
// Copyright 2026 Arctel.net
|
||||||
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
|
|
||||||
|
package ingest
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestGetActiveRequiresUploadID(t *testing.T) {
|
||||||
|
_, err := GetActive(context.Background(), 0)
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("expected error for empty upload id")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -23,6 +23,7 @@ func registerAgentRoutes(apiV1Router *gin.RouterGroup) {
|
|||||||
authorizedRoute.GET("/ws", agent.WebSocketHandler)
|
authorizedRoute.GET("/ws", agent.WebSocketHandler)
|
||||||
authorizedRoute.POST("/nodes/heartbeat", agent.HeartbeatHandler)
|
authorizedRoute.POST("/nodes/heartbeat", agent.HeartbeatHandler)
|
||||||
authorizedRoute.GET("/config-versions/active", agent.GetActiveConfigHandler)
|
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/deployments/:deployment_id/package", agent.DownloadPagesPackageHandler)
|
||||||
authorizedRoute.POST("/waf/ip-groups/sync", agent.SyncWAFIPGroupsHandler)
|
authorizedRoute.POST("/waf/ip-groups/sync", agent.SyncWAFIPGroupsHandler)
|
||||||
authorizedRoute.POST("/apply-logs", agent.ReportApplyLogHandler)
|
authorizedRoute.POST("/apply-logs", agent.ReportApplyLogHandler)
|
||||||
|
|||||||
@@ -213,3 +213,9 @@ type SupportFile struct {
|
|||||||
Path string `json:"path"`
|
Path string `json:"path"`
|
||||||
Content string `json:"content"`
|
Content string `json:"content"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// PagesDeploymentHashResponse is the upload SHA-256 hash for a Pages deployment package.
|
||||||
|
type PagesDeploymentHashResponse struct {
|
||||||
|
DeploymentID uint `json:"deployment_id"`
|
||||||
|
Hash string `json:"hash"`
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user