From 13a375e042488ca908d6e446a41d16a088aa35cc Mon Sep 17 00:00:00 2001 From: ryan Date: Sun, 21 Jun 2026 11:53:54 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=20Agent=20=E9=83=A8=E7=BD=B2?= =?UTF-8?q?=20Pages=20=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/changelog/index.md | 4 + docs/docs.go | 69 +++++ docs/swagger.json | 69 +++++ docs/swagger.yaml | 41 +++ internal/apps/agent/httpclient/client.go | 12 + internal/apps/agent/protocol/alias.go | 3 + internal/apps/agent/state/state.go | 36 +-- internal/apps/agent/sync/pages.go | 142 +++++++++-- internal/apps/agent/sync/service.go | 6 +- internal/apps/agent/sync/service_test.go | 236 +++++++++++++++++- internal/apps/agent/sync/sync_helpers.go | 32 +++ internal/apps/openflare/agent/routers.go | 27 ++ internal/apps/openflare/pages/errs.go | 1 + internal/apps/openflare/pages/logics.go | 53 +++- internal/apps/upload/exports.go | 11 +- internal/apps/upload/ingest/access.go | 49 ++++ internal/apps/upload/ingest/access_test.go | 16 ++ .../router/v1/openflare/register_agent.go | 1 + pkg/protocol/agent.go | 6 + 19 files changed, 755 insertions(+), 59 deletions(-) create mode 100644 internal/apps/upload/ingest/access.go create mode 100644 internal/apps/upload/ingest/access_test.go diff --git a/docs/changelog/index.md b/docs/changelog/index.md index b6b48cb0..da0cd5ba 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -30,6 +30,10 @@ sidebar: false - 修复 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。 - 修复仪表盘「24 小时请求趋势」摘要误显示当前小时请求量/错误量:改为汇总近 24 小时总量。 diff --git a/docs/docs.go b/docs/docs.go index a4ee2620..ce3b859f 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -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": { "get": { "security": [ @@ -16634,6 +16692,17 @@ const docTemplate = `{ } } }, + "protocol.PagesDeploymentHashResponse": { + "type": "object", + "properties": { + "deployment_id": { + "type": "integer" + }, + "hash": { + "type": "string" + } + } + }, "protocol.RelayConfig": { "type": "object", "properties": { diff --git a/docs/swagger.json b/docs/swagger.json index 01d47992..c65765e2 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -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": { "get": { "security": [ @@ -16627,6 +16685,17 @@ } } }, + "protocol.PagesDeploymentHashResponse": { + "type": "object", + "properties": { + "deployment_id": { + "type": "integer" + }, + "hash": { + "type": "string" + } + } + }, "protocol.RelayConfig": { "type": "object", "properties": { diff --git a/docs/swagger.yaml b/docs/swagger.yaml index 4f0a37fe..c66530c1 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -2588,6 +2588,13 @@ definitions: relay_node_id: type: string type: object + protocol.PagesDeploymentHashResponse: + properties: + deployment_id: + type: integer + hash: + type: string + type: object protocol.RelayConfig: properties: auth_token: @@ -6502,6 +6509,40 @@ paths: summary: 注册或发现 Agent 节点 tags: - 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: get: description: 流式下载指定部署的静态资源压缩包,供 Agent 边缘分发 diff --git a/internal/apps/agent/httpclient/client.go b/internal/apps/agent/httpclient/client.go index 45bf830e..ce6a3ca0 100644 --- a/internal/apps/agent/httpclient/client.go +++ b/internal/apps/agent/httpclient/client.go @@ -86,6 +86,18 @@ func (c *Client) SyncWAFIPGroups(ctx context.Context, payload protocol.WAFIPGrou 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. 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) diff --git a/internal/apps/agent/protocol/alias.go b/internal/apps/agent/protocol/alias.go index 9897513d..ae7753de 100644 --- a/internal/apps/agent/protocol/alias.go +++ b/internal/apps/agent/protocol/alias.go @@ -72,6 +72,9 @@ type WAFIPGroupSyncResponse = pkgprotocol.WAFIPGroupSyncResponse // SupportFile is an alias for pkgprotocol.SupportFile. type SupportFile = pkgprotocol.SupportFile +// PagesDeploymentHashResponse is an alias for pkgprotocol.PagesDeploymentHashResponse. +type PagesDeploymentHashResponse = pkgprotocol.PagesDeploymentHashResponse + 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 921fcb98..f460a0ea 100644 --- a/internal/apps/agent/state/state.go +++ b/internal/apps/agent/state/state.go @@ -15,22 +15,30 @@ const ( 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. type Snapshot struct { - NodeID string `json:"node_id"` - CurrentVersion string `json:"current_version"` - CurrentChecksum string `json:"current_checksum"` - BlockedVersion string `json:"blocked_version"` - BlockedChecksum string `json:"blocked_checksum"` - BlockedReason string `json:"blocked_reason"` - LastError string `json:"last_error"` - OpenrestyStatus string `json:"openresty_status"` - OpenrestyMessage string `json:"openresty_message"` - LastProfileFingerprint string `json:"last_profile_fingerprint"` - LastCPUStatTotal uint64 `json:"last_cpu_stat_total"` - LastCPUStatIdle uint64 `json:"last_cpu_stat_idle"` - LastMetricAtUnix int64 `json:"last_metric_at_unix"` - AccessLogOffset int64 `json:"access_log_offset"` + NodeID string `json:"node_id"` + CurrentVersion string `json:"current_version"` + CurrentChecksum string `json:"current_checksum"` + PagesDeployments []PagesDeployment `json:"pages_deployments"` + BlockedVersion string `json:"blocked_version"` + BlockedChecksum string `json:"blocked_checksum"` + BlockedReason string `json:"blocked_reason"` + LastError string `json:"last_error"` + OpenrestyStatus string `json:"openresty_status"` + OpenrestyMessage string `json:"openresty_message"` + LastProfileFingerprint string `json:"last_profile_fingerprint"` + LastCPUStatTotal uint64 `json:"last_cpu_stat_total"` + LastCPUStatIdle uint64 `json:"last_cpu_stat_idle"` + LastMetricAtUnix int64 `json:"last_metric_at_unix"` + AccessLogOffset int64 `json:"access_log_offset"` } // Store manages the storage and retrieval of the agent state snapshot. diff --git a/internal/apps/agent/sync/pages.go b/internal/apps/agent/sync/pages.go index 88cb8f84..eca3b9bb 100644 --- a/internal/apps/agent/sync/pages.go +++ b/internal/apps/agent/sync/pages.go @@ -18,6 +18,7 @@ import ( "strings" "github.com/Rain-kl/Wavelet/internal/apps/agent/protocol" + "github.com/Rain-kl/Wavelet/internal/apps/agent/state" ) const ( @@ -45,10 +46,85 @@ type pagesDeploymentMarker struct { Checksum string `json:"checksum"` } -func (s *Service) syncPagesDeployments(ctx context.Context, config *protocol.ActiveConfigResponse) error { - deployments, err := referencedPagesDeployments(config) - if err != nil { - return err +func pagesDeploymentStateHash(item state.PagesDeployment) string { + if hash := strings.TrimSpace(item.Hash); hash != "" { + return hash + } + 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 { 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") } for _, deployment := range deployments { - if err := s.ensurePagesDeployment(ctx, deployment); err != nil { + if err := s.ensurePagesDeployment(ctx, snapshot, deployment); err != nil { return err } } return nil } -func (s *Service) ensurePagesDeployment(ctx context.Context, deployment pagesDeploymentSource) error { - currentDir := pagesCurrentDir(s.pagesDir, deployment.DeploymentID) - if markerMatches(currentDir, deployment) { - return nil - } - packageBytes, err := s.client.DownloadPagesDeploymentPackage(ctx, deployment.DeploymentID) +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("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 { - return fmt.Errorf("pages deployment %d checksum mismatch: expected %s, got %s", deployment.DeploymentID, deployment.Checksum, got) + serverHash = strings.TrimSpace(serverHash) + if serverHash == "" { + return fmt.Errorf("pages deployment %d hash is empty", deployment.DeploymentID) } - releaseDir := pagesReleaseDir(s.pagesDir, deployment.DeploymentID, deployment.Checksum) - if !markerMatches(releaseDir, deployment) { - if err := extractPagesPackage(packageBytes, releaseDir, deployment); err != nil { - return err + 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) + } + 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) { diff --git a/internal/apps/agent/sync/service.go b/internal/apps/agent/sync/service.go index 774a87d4..e57c29b8 100644 --- a/internal/apps/agent/sync/service.go +++ b/internal/apps/agent/sync/service.go @@ -29,6 +29,7 @@ const ( // ConfigClient is the interface for communicating with the server control plane. 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) ReportApplyLog(ctx context.Context, payload protocol.ApplyLogPayload) 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) } 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) return s.stateStore.Save(snapshot) } @@ -163,7 +167,7 @@ func (s *Service) applyRenderedConfig(ctx context.Context, mode string, snapshot if err != nil { return err } - if err := s.syncPagesDeployments(ctx, config); err != nil { + if err := s.syncPagesDeployments(ctx, snapshot, config); err != nil { return err } mainConfigChecksum := checksumString(rendered.mainConfig) diff --git a/internal/apps/agent/sync/service_test.go b/internal/apps/agent/sync/service_test.go index 94d2dbee..d7f014fb 100644 --- a/internal/apps/agent/sync/service_test.go +++ b/internal/apps/agent/sync/service_test.go @@ -31,7 +31,9 @@ type fakeClient struct { config protocol.ActiveConfigResponse reports []protocol.ApplyLogPayload pagesPackages map[uint][]byte + pagesHashes map[uint]string fetchCalls int + hashCalls int } type fakeManager struct { @@ -76,6 +78,21 @@ func (f *fakeClient) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfi 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) { if f.pagesPackages == nil { return nil, fmt.Errorf("missing Pages package %d", deploymentID) @@ -501,8 +518,8 @@ func TestSyncOnceReportsNoopWhenVersionChangesButChecksumMatches(t *testing.T) { }); err != nil { t.Fatalf("SyncOnce failed: %v", err) } - if client.fetchCalls != 0 { - t.Fatalf("expected checksum match to skip config fetch, got %d", client.fetchCalls) + if client.fetchCalls != 1 { + t.Fatalf("expected checksum match to fetch active config once for Pages reconciliation, got %d", client.fetchCalls) } if len(manager.applyMainContents) != 0 { t.Fatal("expected checksum match to skip apply") @@ -568,9 +585,10 @@ func TestSyncOnceDoesNotRepeatNoopReportWhenStateAlreadyMatches(t *testing.T) { t.Fatalf("EnsureNodeID failed: %v", err) } if err = stateStore.Save(&state.Snapshot{ - NodeID: nodeID, - CurrentVersion: "20260309-003", - CurrentChecksum: "checksum-3", + NodeID: nodeID, + CurrentVersion: "20260309-003", + CurrentChecksum: "checksum-3", + PagesDeployments: []state.PagesDeployment{}, }); err != nil { t.Fatalf("failed to seed state: %v", err) } @@ -583,6 +601,9 @@ func TestSyncOnceDoesNotRepeatNoopReportWhenStateAlreadyMatches(t *testing.T) { }); err != nil { 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 { 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) } if err = stateStore.Save(&state.Snapshot{ - NodeID: nodeID, - CurrentVersion: client.config.Version, - CurrentChecksum: client.config.Checksum, + NodeID: nodeID, + CurrentVersion: client.config.Version, + CurrentChecksum: client.config.Checksum, + PagesDeployments: []state.PagesDeployment{}, }); err != nil { t.Fatalf("failed to seed state: %v", err) } @@ -923,13 +945,209 @@ func TestSyncOnceSkipsFetchWhenHeartbeatChecksumMatches(t *testing.T) { t.Fatalf("SyncOnce failed: %v", err) } 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 { 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 { t.Helper() var buffer bytes.Buffer diff --git a/internal/apps/agent/sync/sync_helpers.go b/internal/apps/agent/sync/sync_helpers.go index 85c9290e..7115d9b0 100644 --- a/internal/apps/agent/sync/sync_helpers.go +++ b/internal/apps/agent/sync/sync_helpers.go @@ -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 { + 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) if shouldReportNoopApply(snapshot, target.Version, target.Checksum) { 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) } 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) 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 { + 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) if shouldReportNoopApply(snapshot, config.Version, config.Checksum) { rendered, renderErr := renderActiveConfig(config) diff --git a/internal/apps/openflare/agent/routers.go b/internal/apps/openflare/agent/routers.go index bd3d8a10..54ddd1cc 100644 --- a/internal/apps/openflare/agent/routers.go +++ b/internal/apps/openflare/agent/routers.go @@ -11,6 +11,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/openflare/pages" "github.com/Rain-kl/Wavelet/internal/apps/openflare/websocket" "github.com/Rain-kl/Wavelet/internal/common/response" + "github.com/Rain-kl/Wavelet/pkg/protocol" "github.com/gin-gonic/gin" ) @@ -153,6 +154,32 @@ func ReportApplyLogHandler(c *gin.Context) { 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. // @Summary 下载 Pages 部署包 // @Description 流式下载指定部署的静态资源压缩包,供 Agent 边缘分发 diff --git a/internal/apps/openflare/pages/errs.go b/internal/apps/openflare/pages/errs.go index cdf26b3f..4de9edf9 100644 --- a/internal/apps/openflare/pages/errs.go +++ b/internal/apps/openflare/pages/errs.go @@ -24,5 +24,6 @@ const ( errPagesPackagePathEmpty = "pages 部署包路径为空" errPagesPackageUploadMissing = "pages 部署包上传记录不存在" errPagesPackageNotInActiveConfig = "pages 部署尚未进入激活配置" + errPagesDeploymentHashMissing = "pages 部署包哈希缺失" errPagesInvalidSnapshotFormat = "配置快照格式无效" ) diff --git a/internal/apps/openflare/pages/logics.go b/internal/apps/openflare/pages/logics.go index 1b76b99e..f7c80827 100644 --- a/internal/apps/openflare/pages/logics.go +++ b/internal/apps/openflare/pages/logics.go @@ -16,10 +16,8 @@ import ( "time" "github.com/Rain-kl/Wavelet/internal/apps/upload" - uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage" "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/model" - "github.com/Rain-kl/Wavelet/internal/repository" "github.com/Rain-kl/Wavelet/internal/storage" "gorm.io/gorm" ) @@ -354,6 +352,45 @@ 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. +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. func OpenDeploymentPackage(ctx context.Context, deploymentID uint) (*storage.Object, string, error) { 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) { - uploadRecord, err := repository.GetActiveUploadByID(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) + obj, _, err := upload.OpenStoredUpload(ctx, deployment.UploadID) if err != nil { return nil, "", fmt.Errorf("pages 部署包不存在: %w", err) } @@ -423,7 +452,7 @@ func hydrateLegacyDeploymentUpload( return nil, errors.New(errPagesPackagePathEmpty) } if deployment.UploadID > 0 { - uploadRecord, err := repository.GetActiveUploadByID(ctx, deployment.UploadID) + uploadRecord, err := upload.GetActiveUpload(ctx, deployment.UploadID) if err == nil { return &uploadRecord, nil } diff --git a/internal/apps/upload/exports.go b/internal/apps/upload/exports.go index 981f0c4c..206d1cce 100644 --- a/internal/apps/upload/exports.go +++ b/internal/apps/upload/exports.go @@ -31,10 +31,13 @@ var ( // Programmatic ingest API var ( - Ingest = ingest.Ingest - Remove = ingest.Remove - RemoveOwned = ingest.RemoveOwned - FindByHash = ingest.FindByHash + Ingest = ingest.Ingest + Remove = ingest.Remove + RemoveOwned = ingest.RemoveOwned + FindByHash = ingest.FindByHash + GetActiveUpload = ingest.GetActive + OpenStoredUpload = ingest.OpenActive + ActiveUploadHash = ingest.ActiveHash ) // Ingest policy constants diff --git a/internal/apps/upload/ingest/access.go b/internal/apps/upload/ingest/access.go new file mode 100644 index 00000000..8b5ee0d5 --- /dev/null +++ b/internal/apps/upload/ingest/access.go @@ -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 +} diff --git a/internal/apps/upload/ingest/access_test.go b/internal/apps/upload/ingest/access_test.go new file mode 100644 index 00000000..a9b4277e --- /dev/null +++ b/internal/apps/upload/ingest/access_test.go @@ -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") + } +} diff --git a/internal/router/v1/openflare/register_agent.go b/internal/router/v1/openflare/register_agent.go index 9d9d350d..fe59132a 100644 --- a/internal/router/v1/openflare/register_agent.go +++ b/internal/router/v1/openflare/register_agent.go @@ -23,6 +23,7 @@ func registerAgentRoutes(apiV1Router *gin.RouterGroup) { authorizedRoute.GET("/ws", agent.WebSocketHandler) authorizedRoute.POST("/nodes/heartbeat", agent.HeartbeatHandler) 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.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 be9241ec..33f529e4 100644 --- a/pkg/protocol/agent.go +++ b/pkg/protocol/agent.go @@ -213,3 +213,9 @@ type SupportFile struct { Path string `json:"path"` 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"` +}