mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-05 23:26:38 +08:00
[新增] OpenFlare Pages
This commit is contained in:
@@ -71,6 +71,7 @@ func main() {
|
||||
LuaDir: cfg.LuaDir,
|
||||
NginxLuaDir: cfg.OpenrestyLuaDir,
|
||||
RuntimeConfigDir: cfg.RuntimeConfigDir,
|
||||
PagesDir: cfg.PagesDir,
|
||||
OpenrestyObservabilityListen: nginx.ObservabilityListenAddress(cfg.OpenrestyObservabilityPort),
|
||||
OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort,
|
||||
OpenrestyResolverDirective: "",
|
||||
@@ -89,12 +90,14 @@ func main() {
|
||||
slog.Error("ensure managed lua assets failed", "error", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
syncService := syncservice.New(client, runtimeManager, stateStore)
|
||||
syncService.SetPagesDir(cfg.PagesDir)
|
||||
runner := &agent.Runner{
|
||||
Config: cfg,
|
||||
StateStore: stateStore,
|
||||
ObservabilityBuffer: observabilityBuffer,
|
||||
HeartbeatService: heartbeat.New(client),
|
||||
SyncService: syncservice.New(client, runtimeManager, stateStore),
|
||||
SyncService: syncService,
|
||||
Updater: updater.New(),
|
||||
RuntimeManager: runtimeManager,
|
||||
WebSocketService: wsClient,
|
||||
|
||||
@@ -23,6 +23,7 @@ const (
|
||||
defaultCertDirRelativePath = "etc/nginx/certs"
|
||||
defaultLuaDirRelativePath = "etc/nginx/lua"
|
||||
defaultRuntimeConfigDirRelativePath = "etc/openflare"
|
||||
defaultPagesDirRelativePath = "var/lib/openflare/pages"
|
||||
defaultMMDBRelativePath = "etc/openflare/GeoLite2-Country.mmdb"
|
||||
defaultAccessLogRelativePath = "var/log/openflare/access.log"
|
||||
defaultStateRelativePath = "var/lib/openflare/agent-state.json"
|
||||
@@ -57,6 +58,7 @@ type Config struct {
|
||||
LuaDir string `json:"lua_dir"`
|
||||
OpenrestyLuaDir string `json:"openresty_lua_dir"`
|
||||
RuntimeConfigDir string `json:"runtime_config_dir"`
|
||||
PagesDir string `json:"pages_dir"`
|
||||
MMDBPath string `json:"mmdb_path"`
|
||||
MMDBUpdateInterval MillisecondDuration `json:"mmdb_update_interval"`
|
||||
MMDBDownloadURL string `json:"mmdb_download_url"`
|
||||
@@ -86,6 +88,7 @@ type configFile struct {
|
||||
LuaDir string `json:"lua_dir"`
|
||||
OpenrestyLuaDir string `json:"openresty_lua_dir"`
|
||||
RuntimeConfigDir string `json:"runtime_config_dir"`
|
||||
PagesDir string `json:"pages_dir"`
|
||||
MMDBPath string `json:"mmdb_path"`
|
||||
MMDBUpdateInterval MillisecondDuration `json:"mmdb_update_interval"`
|
||||
MMDBDownloadURL string `json:"mmdb_download_url"`
|
||||
@@ -128,6 +131,7 @@ func Load(path string) (*Config, error) {
|
||||
LuaDir: file.LuaDir,
|
||||
OpenrestyLuaDir: file.OpenrestyLuaDir,
|
||||
RuntimeConfigDir: file.RuntimeConfigDir,
|
||||
PagesDir: file.PagesDir,
|
||||
MMDBPath: file.MMDBPath,
|
||||
MMDBUpdateInterval: file.MMDBUpdateInterval,
|
||||
MMDBDownloadURL: file.MMDBDownloadURL,
|
||||
@@ -190,6 +194,9 @@ func applyDefaults(cfg *Config, baseDir string) {
|
||||
if cfg.RuntimeConfigDir == "" {
|
||||
cfg.RuntimeConfigDir = joinManagedPath(cfg.DataDir, defaultRuntimeConfigDirRelativePath)
|
||||
}
|
||||
if cfg.PagesDir == "" {
|
||||
cfg.PagesDir = joinManagedPath(cfg.DataDir, defaultPagesDirRelativePath)
|
||||
}
|
||||
if cfg.MMDBPath == "" {
|
||||
cfg.MMDBPath = joinManagedPath(cfg.DataDir, defaultMMDBRelativePath)
|
||||
}
|
||||
@@ -231,6 +238,7 @@ func normalizeManagedPaths(cfg *Config) {
|
||||
&cfg.LuaDir,
|
||||
&cfg.OpenrestyLuaDir,
|
||||
&cfg.RuntimeConfigDir,
|
||||
&cfg.PagesDir,
|
||||
&cfg.StatePath,
|
||||
&cfg.ObservabilityBufferPath,
|
||||
&cfg.MMDBPath,
|
||||
@@ -251,6 +259,7 @@ func hasEnvConfig() bool {
|
||||
"OPENFLARE_NODE_IP",
|
||||
"OPENFLARE_DATA_DIR",
|
||||
"OPENFLARE_OPENRESTY_PATH",
|
||||
"OPENFLARE_PAGES_DIR",
|
||||
"OPENFLARE_HEARTBEAT_INTERVAL",
|
||||
"OPENFLARE_REQUEST_TIMEOUT",
|
||||
"OPENFLARE_OPENRESTY_OBSERVABILITY_PORT",
|
||||
@@ -281,6 +290,7 @@ func applyEnvOverrides(cfg *Config) {
|
||||
overrideString("OPENFLARE_NODE_IP", &cfg.NodeIP)
|
||||
overrideString("OPENFLARE_DATA_DIR", &cfg.DataDir)
|
||||
overrideString("OPENFLARE_OPENRESTY_PATH", &cfg.OpenrestyPath)
|
||||
overrideString("OPENFLARE_PAGES_DIR", &cfg.PagesDir)
|
||||
overrideString("OPENFLARE_MMDB_PATH", &cfg.MMDBPath)
|
||||
overrideString("OPENFLARE_MMDB_DOWNLOAD_URL", &cfg.MMDBDownloadURL)
|
||||
if value := strings.TrimSpace(os.Getenv("OPENFLARE_HEARTBEAT_INTERVAL")); value != "" {
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
@@ -86,6 +87,23 @@ func (c *Client) SyncWAFIPGroups(ctx context.Context, payload protocol.WAFIPGrou
|
||||
return &resp.Data, nil
|
||||
}
|
||||
|
||||
func (c *Client) DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error) {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.baseURL+fmt.Sprintf("/api/agent/pages/deployments/%d/package", deploymentID), nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
req.Header.Set("X-Agent-Token", c.token)
|
||||
res, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer res.Body.Close()
|
||||
if res.StatusCode != http.StatusOK {
|
||||
return nil, errors.New(res.Status)
|
||||
}
|
||||
return io.ReadAll(res.Body)
|
||||
}
|
||||
|
||||
func (c *Client) SetToken(token string) {
|
||||
c.token = strings.TrimSpace(token)
|
||||
slog.Debug("http client token updated")
|
||||
|
||||
@@ -141,6 +141,7 @@ type Manager struct {
|
||||
LuaDir string
|
||||
NginxLuaDir string
|
||||
RuntimeConfigDir string
|
||||
PagesDir string
|
||||
OpenrestyObservabilityListen string
|
||||
OpenrestyObservabilityPort int
|
||||
OpenrestyResolverDirective string
|
||||
@@ -422,6 +423,9 @@ func (m *Manager) CurrentChecksum() (string, error) {
|
||||
normalizedRoute = strings.ReplaceAll(normalizedRoute, luaDir+"/pow/static", openrestyrender.PowStaticDirPlaceholder)
|
||||
normalizedRoute = strings.ReplaceAll(normalizedRoute, luaDir, openrestyrender.LuaDirPlaceholder)
|
||||
}
|
||||
if pagesDir := m.pagesRuntimePath(); pagesDir != "" {
|
||||
normalizedRoute = strings.ReplaceAll(normalizedRoute, pagesDir, openrestyrender.PagesDirPlaceholder)
|
||||
}
|
||||
files, err := m.readManagedSupportFiles()
|
||||
if err != nil {
|
||||
return "", err
|
||||
@@ -1130,6 +1134,9 @@ func (m *Manager) renderRouteConfig(content string) string {
|
||||
rendered = strings.ReplaceAll(rendered, openrestyrender.LuaDirPlaceholder, luaDir)
|
||||
rendered = strings.ReplaceAll(rendered, openrestyrender.PowStaticDirPlaceholder, luaDir+"/pow/static")
|
||||
}
|
||||
if pagesDir := m.pagesRuntimePath(); pagesDir != "" {
|
||||
rendered = strings.ReplaceAll(rendered, openrestyrender.PagesDirPlaceholder, pagesDir)
|
||||
}
|
||||
return rendered
|
||||
}
|
||||
|
||||
@@ -1258,6 +1265,10 @@ func (m *Manager) luaRuntimePath() string {
|
||||
return filepath.ToSlash(m.NginxLuaDir)
|
||||
}
|
||||
|
||||
func (m *Manager) pagesRuntimePath() string {
|
||||
return filepath.ToSlash(strings.TrimSpace(m.PagesDir))
|
||||
}
|
||||
|
||||
func checksum(content string) string {
|
||||
sum := sha256.Sum256([]byte(content))
|
||||
return hex.EncodeToString(sum[:])
|
||||
|
||||
@@ -0,0 +1,273 @@
|
||||
package sync
|
||||
|
||||
import (
|
||||
"archive/zip"
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"path"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
||||
"openflare-agent/internal/protocol"
|
||||
)
|
||||
|
||||
type pagesSourceDocument struct {
|
||||
Routes []pagesSourceRoute `json:"routes"`
|
||||
}
|
||||
|
||||
type pagesSourceRoute struct {
|
||||
UpstreamType string `json:"upstream_type"`
|
||||
PagesDeployment *pagesDeploymentSource `json:"pages_deployment"`
|
||||
}
|
||||
|
||||
type pagesDeploymentSource struct {
|
||||
DeploymentID uint `json:"deployment_id"`
|
||||
Checksum string `json:"checksum"`
|
||||
}
|
||||
|
||||
type pagesDeploymentMarker struct {
|
||||
DeploymentID uint `json:"deployment_id"`
|
||||
Checksum string `json:"checksum"`
|
||||
}
|
||||
|
||||
func (s *Service) syncPagesDeployments(ctx context.Context, config *protocol.ActiveConfigResponse) error {
|
||||
deployments, err := referencedPagesDeployments(config)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(deployments) == 0 {
|
||||
return nil
|
||||
}
|
||||
if strings.TrimSpace(s.pagesDir) == "" {
|
||||
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 {
|
||||
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)
|
||||
if err != nil {
|
||||
return fmt.Errorf("download Pages deployment %d: %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)
|
||||
}
|
||||
releaseDir := pagesReleaseDir(s.pagesDir, deployment.DeploymentID, deployment.Checksum)
|
||||
if !markerMatches(releaseDir, deployment) {
|
||||
if err := extractPagesPackage(packageBytes, releaseDir, deployment); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return switchPagesCurrentDir(s.pagesDir, deployment.DeploymentID, releaseDir)
|
||||
}
|
||||
|
||||
func referencedPagesDeployments(config *protocol.ActiveConfigResponse) ([]pagesDeploymentSource, error) {
|
||||
if config == nil || strings.TrimSpace(config.SourceConfigJSON) == "" {
|
||||
return nil, nil
|
||||
}
|
||||
var doc pagesSourceDocument
|
||||
if err := json.Unmarshal([]byte(config.SourceConfigJSON), &doc); err != nil {
|
||||
return nil, fmt.Errorf("decode Pages references: %w", err)
|
||||
}
|
||||
seen := make(map[uint]struct{})
|
||||
result := make([]pagesDeploymentSource, 0)
|
||||
for _, route := range doc.Routes {
|
||||
if strings.ToLower(strings.TrimSpace(route.UpstreamType)) != "pages" || route.PagesDeployment == nil {
|
||||
continue
|
||||
}
|
||||
deploymentID := route.PagesDeployment.DeploymentID
|
||||
checksum := strings.TrimSpace(route.PagesDeployment.Checksum)
|
||||
if deploymentID == 0 || checksum == "" {
|
||||
return nil, errors.New("Pages deployment snapshot is incomplete")
|
||||
}
|
||||
if _, ok := seen[deploymentID]; ok {
|
||||
continue
|
||||
}
|
||||
seen[deploymentID] = struct{}{}
|
||||
result = append(result, pagesDeploymentSource{DeploymentID: deploymentID, Checksum: checksum})
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func extractPagesPackage(packageBytes []byte, releaseDir string, deployment pagesDeploymentSource) error {
|
||||
tmpDir := releaseDir + ".tmp"
|
||||
_ = os.RemoveAll(tmpDir)
|
||||
if err := os.MkdirAll(tmpDir, 0o755); err != nil {
|
||||
return err
|
||||
}
|
||||
reader, err := zip.NewReader(bytes.NewReader(packageBytes), int64(len(packageBytes)))
|
||||
if err != nil {
|
||||
_ = os.RemoveAll(tmpDir)
|
||||
return fmt.Errorf("open Pages zip: %w", err)
|
||||
}
|
||||
for _, item := range reader.File {
|
||||
relativePath, skip, err := normalizePagesArchivePath(item.Name)
|
||||
if err != nil {
|
||||
_ = os.RemoveAll(tmpDir)
|
||||
return err
|
||||
}
|
||||
if skip {
|
||||
continue
|
||||
}
|
||||
if item.FileInfo().Mode()&os.ModeSymlink != 0 {
|
||||
_ = os.RemoveAll(tmpDir)
|
||||
return fmt.Errorf("Pages package contains unsupported symlink: %s", relativePath)
|
||||
}
|
||||
if err := extractPagesFile(item, filepath.Join(tmpDir, relativePath)); err != nil {
|
||||
_ = os.RemoveAll(tmpDir)
|
||||
return err
|
||||
}
|
||||
}
|
||||
if err := writePagesMarker(tmpDir, deployment); err != nil {
|
||||
_ = os.RemoveAll(tmpDir)
|
||||
return err
|
||||
}
|
||||
_ = os.RemoveAll(releaseDir)
|
||||
return os.Rename(tmpDir, releaseDir)
|
||||
}
|
||||
|
||||
func extractPagesFile(item *zip.File, targetPath string) error {
|
||||
if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil {
|
||||
return err
|
||||
}
|
||||
source, err := item.Open()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer source.Close()
|
||||
target, err := os.OpenFile(targetPath, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, item.FileInfo().Mode().Perm())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer target.Close()
|
||||
_, err = io.Copy(target, source)
|
||||
return err
|
||||
}
|
||||
|
||||
func switchPagesCurrentDir(baseDir string, deploymentID uint, releaseDir string) error {
|
||||
currentDir := pagesCurrentDir(baseDir, deploymentID)
|
||||
previousDir := currentDir + ".previous"
|
||||
_ = os.RemoveAll(previousDir)
|
||||
if err := os.MkdirAll(filepath.Dir(currentDir), 0o755); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := os.Stat(currentDir); err == nil {
|
||||
if err := os.Rename(currentDir, previousDir); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if err := copyPagesDir(releaseDir, currentDir); err != nil {
|
||||
_ = os.RemoveAll(currentDir)
|
||||
if _, restoreErr := os.Stat(previousDir); restoreErr == nil {
|
||||
_ = os.Rename(previousDir, currentDir)
|
||||
}
|
||||
return err
|
||||
}
|
||||
_ = os.RemoveAll(previousDir)
|
||||
return nil
|
||||
}
|
||||
|
||||
func copyPagesDir(sourceDir string, targetDir string) error {
|
||||
return filepath.WalkDir(sourceDir, func(sourcePath string, entry os.DirEntry, err error) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
relativePath, err := filepath.Rel(sourceDir, sourcePath)
|
||||
if err != nil || relativePath == "." {
|
||||
return err
|
||||
}
|
||||
targetPath := filepath.Join(targetDir, relativePath)
|
||||
if entry.IsDir() {
|
||||
return os.MkdirAll(targetPath, 0o755)
|
||||
}
|
||||
info, err := entry.Info()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
input, err := os.Open(sourcePath)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer input.Close()
|
||||
if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil {
|
||||
return err
|
||||
}
|
||||
output, err := os.OpenFile(targetPath, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, info.Mode().Perm())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer output.Close()
|
||||
_, err = io.Copy(output, input)
|
||||
return err
|
||||
})
|
||||
}
|
||||
|
||||
func normalizePagesArchivePath(raw string) (string, bool, error) {
|
||||
name := strings.TrimSpace(filepath.ToSlash(raw))
|
||||
if name == "" || strings.HasSuffix(name, "/") {
|
||||
return "", true, nil
|
||||
}
|
||||
if strings.HasPrefix(name, "/") {
|
||||
return "", false, fmt.Errorf("Pages package contains absolute path: %s", raw)
|
||||
}
|
||||
cleaned := path.Clean(name)
|
||||
if cleaned == "." {
|
||||
return "", true, nil
|
||||
}
|
||||
if cleaned == ".." || strings.HasPrefix(cleaned, "../") || strings.Contains(cleaned, "/../") {
|
||||
return "", false, fmt.Errorf("Pages package path escapes deployment root: %s", raw)
|
||||
}
|
||||
return filepath.FromSlash(cleaned), false, nil
|
||||
}
|
||||
|
||||
func markerMatches(dir string, deployment pagesDeploymentSource) bool {
|
||||
data, err := os.ReadFile(filepath.Join(dir, ".openflare-pages.json"))
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
var marker pagesDeploymentMarker
|
||||
if err := json.Unmarshal(data, &marker); err != nil {
|
||||
return false
|
||||
}
|
||||
return marker.DeploymentID == deployment.DeploymentID && marker.Checksum == deployment.Checksum
|
||||
}
|
||||
|
||||
func writePagesMarker(dir string, deployment pagesDeploymentSource) error {
|
||||
data, err := json.Marshal(pagesDeploymentMarker{
|
||||
DeploymentID: deployment.DeploymentID,
|
||||
Checksum: deployment.Checksum,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return os.WriteFile(filepath.Join(dir, ".openflare-pages.json"), data, 0o644)
|
||||
}
|
||||
|
||||
func pagesCurrentDir(baseDir string, deploymentID uint) string {
|
||||
return filepath.Join(baseDir, "deployments", fmt.Sprintf("%d", deploymentID), "current")
|
||||
}
|
||||
|
||||
func pagesReleaseDir(baseDir string, deploymentID uint, checksum string) string {
|
||||
return filepath.Join(baseDir, "deployments", fmt.Sprintf("%d", deploymentID), "releases", checksum)
|
||||
}
|
||||
|
||||
func checksumBytes(data []byte) string {
|
||||
sum := sha256.Sum256(data)
|
||||
return hex.EncodeToString(sum[:])
|
||||
}
|
||||
@@ -25,6 +25,7 @@ const (
|
||||
|
||||
type ConfigClient interface {
|
||||
GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigResponse, 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)
|
||||
}
|
||||
@@ -42,6 +43,11 @@ type Service struct {
|
||||
client ConfigClient
|
||||
nginxManager NginxManager
|
||||
stateStore *state.Store
|
||||
pagesDir string
|
||||
}
|
||||
|
||||
func (s *Service) SetPagesDir(path string) {
|
||||
s.pagesDir = strings.TrimSpace(path)
|
||||
}
|
||||
|
||||
func New(client ConfigClient, nginxManager NginxManager, stateStore *state.Store) *Service {
|
||||
@@ -225,6 +231,9 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := s.syncPagesDeployments(ctx, config); err != nil {
|
||||
return err
|
||||
}
|
||||
mainConfigChecksum := checksumString(rendered.mainConfig)
|
||||
routeConfigChecksum := checksumString(rendered.routeConfig)
|
||||
slog.Info("applying new openresty config", "mode", mode, "from_version", snapshot.CurrentVersion, "to_version", config.Version, "old_checksum", currentChecksum, "new_checksum", config.Checksum)
|
||||
|
||||
@@ -1,7 +1,11 @@
|
||||
package sync
|
||||
|
||||
import (
|
||||
"archive/zip"
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
@@ -19,10 +23,15 @@ type fakeExecutor struct {
|
||||
reloadErr error
|
||||
}
|
||||
|
||||
func testPagesSourceConfigJSON(deploymentID uint, checksum string) string {
|
||||
return fmt.Sprintf(`{"routes":[{"id":1,"site_name":"pages","domain":"pages.example.com","domains":["pages.example.com"],"origin_url":"openflare-pages://project/1","upstreams":["openflare-pages://project/1"],"enabled":true,"upstream_type":"pages","pages_deployment":{"project_id":1,"project_slug":"pages","deployment_id":%d,"deployment_number":1,"checksum":"%s","entry_file":"index.html","spa_fallback_enabled":true,"local_root":"__OPENFLARE_PAGES_DIR__/deployments/%d/current"}}],"openresty_config":{"worker_processes":"auto","worker_connections":1024,"worker_rlimit_nofile":65535,"events_multi_accept_enabled":true,"keepalive_timeout":20,"keepalive_requests":1000,"client_header_timeout":15,"client_body_timeout":15,"client_max_body_size":"64m","large_client_header_buffers":"4 16k","send_timeout":30,"proxy_connect_timeout":3,"proxy_send_timeout":60,"proxy_read_timeout":60,"websocket_enabled":true,"proxy_request_buffering":false,"proxy_buffering_enabled":true,"proxy_buffers":"16 16k","proxy_buffer_size":"8k","proxy_busy_buffers_size":"64k","gzip_enabled":true,"gzip_min_length":1024,"gzip_comp_level":5,"cache_enabled":false,"cache_levels":"1:2","cache_inactive":"30m","cache_max_size":"1g","cache_key_template":"$scheme$host$request_uri","cache_lock_enabled":true,"cache_lock_timeout":"5s","cache_use_stale":"error timeout updating http_500 http_502 http_503 http_504","main_config_template":"worker_processes {{OpenRestyWorkerProcesses}};"},"waf":{"rule_groups":[],"bindings":[]}}`, deploymentID, checksum, deploymentID)
|
||||
}
|
||||
|
||||
type fakeClient struct {
|
||||
config protocol.ActiveConfigResponse
|
||||
reports []protocol.ApplyLogPayload
|
||||
fetchCalls int
|
||||
config protocol.ActiveConfigResponse
|
||||
reports []protocol.ApplyLogPayload
|
||||
pagesPackages map[uint][]byte
|
||||
fetchCalls int
|
||||
}
|
||||
|
||||
type fakeManager struct {
|
||||
@@ -67,6 +76,13 @@ func (f *fakeClient) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfi
|
||||
return &f.config, nil
|
||||
}
|
||||
|
||||
func (f *fakeClient) DownloadPagesDeploymentPackage(ctx context.Context, deploymentID uint) ([]byte, error) {
|
||||
if f.pagesPackages == nil {
|
||||
return nil, fmt.Errorf("missing Pages package %d", deploymentID)
|
||||
}
|
||||
return f.pagesPackages[deploymentID], nil
|
||||
}
|
||||
|
||||
func (f *fakeClient) ReportApplyLog(ctx context.Context, payload protocol.ApplyLogPayload) error {
|
||||
f.reports = append(f.reports, payload)
|
||||
return nil
|
||||
@@ -179,6 +195,77 @@ func TestSyncOnceSuccess(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncOnceDownloadsPagesDeploymentBeforeApply(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)
|
||||
}
|
||||
snapshot, _ := stateStore.Load()
|
||||
snapshot.NodeID = nodeID
|
||||
if err = stateStore.Save(snapshot); err != nil {
|
||||
t.Fatalf("save state failed: %v", err)
|
||||
}
|
||||
manager := &fakeManager{currentChecksum: "old-checksum"}
|
||||
service := New(client, manager, stateStore)
|
||||
pagesDir := t.TempDir()
|
||||
service.SetPagesDir(pagesDir)
|
||||
|
||||
if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{Version: "20260309-101", Checksum: "pages-config-checksum"}); err != nil {
|
||||
t.Fatalf("SyncOnce failed: %v", err)
|
||||
}
|
||||
data, err := os.ReadFile(filepath.Join(pagesDir, "deployments", "7", "current", "index.html"))
|
||||
if err != nil {
|
||||
t.Fatalf("expected Pages file to be extracted: %v", err)
|
||||
}
|
||||
if string(data) != "hello" {
|
||||
t.Fatalf("unexpected Pages file content: %s", string(data))
|
||||
}
|
||||
if len(manager.applyRouteContents) != 1 || !strings.Contains(manager.applyRouteContents[0], "__OPENFLARE_PAGES_DIR__/deployments/7/current") {
|
||||
t.Fatalf("expected Pages placeholder in rendered route config, got %#v", manager.applyRouteContents)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncOnceRejectsPagesZipSlipBeforeApply(t *testing.T) {
|
||||
packageBytes := testPagesPackage(t, map[string]string{"../escape.html": "bad", "index.html": "ok"})
|
||||
checksum := testBytesChecksum(packageBytes)
|
||||
client := &fakeClient{
|
||||
config: protocol.ActiveConfigResponse{
|
||||
Version: "20260309-102",
|
||||
Checksum: "pages-config-checksum",
|
||||
SourceConfigJSON: testPagesSourceConfigJSON(8, checksum),
|
||||
CreatedAt: time.Now().Format(time.RFC3339),
|
||||
},
|
||||
pagesPackages: map[uint][]byte{8: packageBytes},
|
||||
}
|
||||
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
|
||||
if _, err := stateStore.EnsureNodeID(); err != nil {
|
||||
t.Fatalf("EnsureNodeID failed: %v", err)
|
||||
}
|
||||
manager := &fakeManager{currentChecksum: "old-checksum"}
|
||||
service := New(client, manager, stateStore)
|
||||
service.SetPagesDir(t.TempDir())
|
||||
|
||||
err := service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{Version: "20260309-102", Checksum: "pages-config-checksum"})
|
||||
if err == nil || !strings.Contains(err.Error(), "escapes deployment root") {
|
||||
t.Fatalf("expected zip-slip rejection, got %v", err)
|
||||
}
|
||||
if len(manager.applyRouteContents) != 0 {
|
||||
t.Fatalf("OpenResty apply must not run after Pages package rejection")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncOnceRollbackOnNginxFailure(t *testing.T) {
|
||||
client := &fakeClient{
|
||||
config: protocol.ActiveConfigResponse{
|
||||
@@ -761,3 +848,27 @@ func TestSyncOnceSkipsFetchWhenHeartbeatChecksumMatches(t *testing.T) {
|
||||
t.Fatal("expected no apply log when no config change is needed")
|
||||
}
|
||||
}
|
||||
|
||||
func testPagesPackage(t *testing.T, files map[string]string) []byte {
|
||||
t.Helper()
|
||||
var buffer bytes.Buffer
|
||||
writer := zip.NewWriter(&buffer)
|
||||
for name, content := range files {
|
||||
file, err := writer.Create(name)
|
||||
if err != nil {
|
||||
t.Fatalf("create zip file failed: %v", err)
|
||||
}
|
||||
if _, err := file.Write([]byte(content)); err != nil {
|
||||
t.Fatalf("write zip file failed: %v", err)
|
||||
}
|
||||
}
|
||||
if err := writer.Close(); err != nil {
|
||||
t.Fatalf("close zip failed: %v", err)
|
||||
}
|
||||
return buffer.Bytes()
|
||||
}
|
||||
|
||||
func testBytesChecksum(data []byte) string {
|
||||
sum := sha256.Sum256(data)
|
||||
return hex.EncodeToString(sum[:])
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user