diff --git a/framework/clclient/client.go b/framework/clclient/client.go index 1f1b63258..90d1aa4bc 100644 --- a/framework/clclient/client.go +++ b/framework/clclient/client.go @@ -9,6 +9,7 @@ import ( "net/http" "os" "regexp" + "strconv" "strings" "sync" "time" @@ -1431,3 +1432,38 @@ func ImportP2PKeys(cl []*ChainlinkClient, keys [][]byte) error { } return eg.Wait() } + +func (c *ChainlinkClient) ReadWorkflowEvents(workflowID string, sequence int64, limit int) (*WorkflowDebugEvents, *http.Response, error) { + specObj := &WorkflowDebugEvents{} + framework.L.Info().Str(NodeURL, c.Config.URL).Str("ID", workflowID).Int64("sequence", sequence).Int("limit", limit).Msg("Reading Workflow Events") + resp, err := c.APIClient.R(). + SetResult(&specObj). + SetPathParams(map[string]string{ + "id": workflowID, + }). + SetQueryParams(map[string]string{ + "sequence": strconv.FormatInt(sequence, 10), + "limit": strconv.Itoa(limit), + }). + Get("/v2/debug/workflow/{id}/events?sequence={sequence}&limit={limit}") + if err != nil { + return nil, nil, err + } + return specObj, resp.RawResponse, err +} + +func (c *ChainlinkClient) ReadOrphanEvents(sequence int64, limit int) (*WorkflowOrphanEvents, *http.Response, error) { + specObj := &WorkflowOrphanEvents{} + framework.L.Info().Str(NodeURL, c.Config.URL).Int64("sequence", sequence).Int("limit", limit).Msg("Reading Workflow Orphan Events") + resp, err := c.APIClient.R(). + SetResult(&specObj). + SetQueryParams(map[string]string{ + "sequence": strconv.FormatInt(sequence, 10), + "limit": strconv.Itoa(limit), + }). + Get("/v2/debug/workflow/orphan_events?sequence={sequence}&limit={limit}") + if err != nil { + return nil, nil, err + } + return specObj, resp.RawResponse, err +} diff --git a/framework/clclient/models.go b/framework/clclient/models.go index 69f9d9d93..5d8eaa794 100644 --- a/framework/clclient/models.go +++ b/framework/clclient/models.go @@ -1419,3 +1419,45 @@ type ForwarderAttributes struct { CreatedAt time.Time `json:"createdAt"` UpdatedAt time.Time `json:"updatedAt"` } + +type WorkflowDebugEvents struct { + Data WorkflowDebugEventsData `json:"data"` +} + +type WorkflowDebugEventsData struct { + Type string `json:"type"` + ID string `json:"id"` + Attributes WorkflowDebugEventsAttributes `json:"attributes"` +} + +type WorkflowDebugEventsAttributes struct { + Events []WorkflowDebugEvent `json:"events"` +} + +type WorkflowDebugEvent struct { + Timestamp time.Time `json:"timestamp"` + Sequence int64 `json:"sequence"` + Message []byte `json:"message"` // protobuf encoded + Type string `json:"type"` // protobuf type name +} + +type WorkflowOrphanEvents struct { + Data WorkflowOrphanEventsData `json:"data"` +} + +type WorkflowOrphanEventsData struct { + Type string `json:"type"` + ID string `json:"id"` + Attributes WorkflowDebugEventsAttributes `json:"attributes"` +} + +type WorkflowOrphanEventsAttributes struct { + Events []WorkflowOrphanEvent `json:"events"` +} + +type WorkflowOrphanEvent struct { + Timestamp time.Time `json:"timestamp"` + Sequence int64 `json:"sequence"` + Message []byte `json:"message"` // protobuf encoded + Type string `json:"type"` // protobuf type name +} diff --git a/framework/go.mod b/framework/go.mod index 5e0900c27..c6b0bf298 100644 --- a/framework/go.mod +++ b/framework/go.mod @@ -125,6 +125,8 @@ require ( github.com/golang/snappy v0.0.5-0.20220116011046-fa5810519dcb // indirect github.com/google/btree v1.1.3 // indirect github.com/google/gnostic-models v0.6.8 // indirect + github.com/google/go-github/v72 v72.0.0 // indirect + github.com/google/go-querystring v1.1.0 // indirect github.com/google/gofuzz v1.2.0 // indirect github.com/google/s2a-go v0.1.9 // indirect github.com/googleapis/enterprise-certificate-proxy v0.3.4 // indirect diff --git a/framework/go.sum b/framework/go.sum index b9a7e1b05..509afd770 100644 --- a/framework/go.sum +++ b/framework/go.sum @@ -332,12 +332,15 @@ github.com/google/go-cmp v0.2.0/go.mod h1:oXzfMopK8JAjlY9xF4vHSVASa0yLyX7SntLO5a github.com/google/go-cmp v0.3.0/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= github.com/google/go-cmp v0.3.1/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= github.com/google/go-cmp v0.4.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= +github.com/google/go-cmp v0.5.2/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.5.4/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.5.9/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/go-github/v72 v72.0.0 h1:FcIO37BLoVPBO9igQQ6tStsv2asG4IPcYFi655PPvBM= +github.com/google/go-github/v72 v72.0.0/go.mod h1:WWtw8GMRiL62mvIquf1kO3onRHeWWKmK01qdCY8c5fg= github.com/google/go-querystring v1.1.0 h1:AnCroh3fv4ZBgVIf1Iwtovgjaw/GiKJo8M8yD/fhyJ8= github.com/google/go-querystring v1.1.0/go.mod h1:Kcdr2DB4koayq7X8pmAG4sNG59So17icRSOU623lUBU= github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= diff --git a/framework/observability.go b/framework/observability.go index 9ff510f2c..ee8fb812c 100644 --- a/framework/observability.go +++ b/framework/observability.go @@ -1,12 +1,16 @@ package framework import ( + "context" "embed" "fmt" "io/fs" "os" "path/filepath" "strings" + "time" + + "github.com/google/go-github/v72/github" ) //go:embed observability/* @@ -22,18 +26,189 @@ const ( LocalPrometheusURL = "http://localhost:3000/explore?panes=%7B%22qZw%22:%7B%22datasource%22:%22PBFA97CFB590B2093%22,%22queries%22:%5B%7B%22refId%22:%22A%22,%22expr%22:%22%22,%22range%22:true,%22datasource%22:%7B%22type%22:%22prometheus%22,%22uid%22:%22PBFA97CFB590B2093%22%7D%7D%5D,%22range%22:%7B%22from%22:%22now-15m%22,%22to%22:%22now%22%7D%7D%7D&schemaVersion=1&orgId=1" LocalPostgresDebugURL = "http://localhost:3000/d/000000039/postgresql-database?orgId=1&refresh=5s&var-DS_PROMETHEUS=PBFA97CFB590B2093&var-interval=$__auto_interval_interval&var-namespace=&var-release=&var-instance=postgres_exporter_0:9187&var-datname=All&var-mode=All&from=now-15m&to=now" LocalPyroScopeURL = "http://localhost:4040/?query=process_cpu%3Acpu%3Ananoseconds%3Acpu%3Ananoseconds%7Bservice_name%3D%22chainlink-node%22%7D&from=now-15m" + + CTFCacheDir = ".local/share/ctf" + DefaultGitHubOwner = "smartcontractkit" + DefaultGitHubRepo = "chainlink-testing-framework" + DefaultObservabilityPath = "framework/observability" ) +// resolveObservabilitySource determines where to load observability files from based on the source parameter +// - Empty string: use embedded files +// - file:// prefix: use local filesystem +// - http(s):// prefix: download and cache from remote URL +func resolveObservabilitySource(source string) (fs.FS, string, error) { + if source == "" { + // Default: use embedded files + return EmbeddedObservabilityFiles, "observability", nil + } + + if strings.HasPrefix(source, "file://") { + // Local filesystem path + localPath := strings.TrimPrefix(source, "file://") + if _, err := os.Stat(localPath); err != nil { + return nil, "", fmt.Errorf("local observability path does not exist: %s: %w", localPath, err) + } + return os.DirFS(localPath), ".", nil + } + + if strings.HasPrefix(source, "http://") || strings.HasPrefix(source, "https://") { + // Remote URL: download and cache + cachePath, err := downloadAndCacheObservabilityFiles(source) + if err != nil { + return nil, "", fmt.Errorf("failed to download observability files: %w", err) + } + return os.DirFS(cachePath), ".", nil + } + + return nil, "", fmt.Errorf("invalid source format: %s (must be empty, file://, or http(s)://)", source) +} + +// downloadAndCacheObservabilityFiles downloads observability files from a GitHub URL and caches them +func downloadAndCacheObservabilityFiles(url string) (string, error) { + // Parse URL to extract repo info and create cache key + // Expected format: https://github.com/owner/repo/tree/ref/path/to/observability + owner, repo, ref, path, err := parseGitHubURL(url) + if err != nil { + return "", err + } + + // Create cache directory using just the ref + homeDir, err := os.UserHomeDir() + if err != nil { + return "", fmt.Errorf("failed to get home directory: %w", err) + } + cachedPath := filepath.Join(homeDir, CTFCacheDir, "observability", ref) + + // Check if already cached and has content + if info, err := os.Stat(cachedPath); err == nil && info.IsDir() { + // Verify the cache directory has files + entries, err := os.ReadDir(cachedPath) + if err == nil && len(entries) > 0 { + L.Debug().Msgf("Using cached observability files from %s", cachedPath) + return cachedPath, nil + } + L.Debug().Msg("Cache directory exists but is empty, re-downloading") + } + + L.Info().Msgf("Downloading observability files from GitHub: %s/%s@%s (path: %s)", owner, repo, ref, path) + + // Create GitHub client with optional authentication and timeout context + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute) + defer cancel() + + var client *github.Client + if token := os.Getenv("GITHUB_TOKEN"); token != "" { + L.Debug().Msg("Using authenticated GitHub client") + client = github.NewClient(nil).WithAuthToken(token) + } else { + L.Debug().Msg("Using unauthenticated GitHub client") + client = github.NewClient(nil) + } + + // Download directory contents recursively + if err := downloadDirectoryRecursive(ctx, client, owner, repo, ref, path, cachedPath); err != nil { + return "", fmt.Errorf("failed to download directory: %w", err) + } + + L.Info().Msgf("Observability files cached at: %s", cachedPath) + return cachedPath, nil +} + +// parseGitHubURL parses a GitHub URL and extracts owner, repo, ref, and path +func parseGitHubURL(url string) (owner, repo, ref, path string, err error) { + // Expected format: https://github.com/owner/repo/tree/ref/path/to/observability + url = strings.TrimPrefix(url, "https://github.com/") + url = strings.TrimPrefix(url, "http://github.com/") + parts := strings.Split(url, "/") + + if len(parts) < 5 { + return "", "", "", "", fmt.Errorf("invalid GitHub URL format (expected: https://github.com/owner/repo/tree|blob/ref/path)") + } + + owner = parts[0] + repo = parts[1] + treeOrBlob := parts[2] + ref = parts[3] + path = strings.Join(parts[4:], "/") + + if treeOrBlob != "tree" && treeOrBlob != "blob" { + return "", "", "", "", fmt.Errorf("unsupported GitHub URL type: %s (expected 'tree' or 'blob')", treeOrBlob) + } + + return owner, repo, ref, path, nil +} + +// downloadDirectoryRecursive recursively downloads a directory from GitHub +func downloadDirectoryRecursive(ctx context.Context, client *github.Client, owner, repo, ref, path, destPath string) error { + // Get directory contents + _, directoryContent, _, err := client.Repositories.GetContents(ctx, owner, repo, path, &github.RepositoryContentGetOptions{ + Ref: ref, + }) + if err != nil { + return fmt.Errorf("failed to get directory contents: %w", err) + } + + // Create destination directory + if err := os.MkdirAll(destPath, 0o755); err != nil { + return fmt.Errorf("failed to create directory: %w", err) + } + + // Process each item in the directory + for _, item := range directoryContent { + if item.GetName() == "README.md" { + continue + } + + itemPath := item.GetPath() + itemName := item.GetName() + targetPath := filepath.Join(destPath, itemName) + + switch item.GetType() { + case "file": + // Download file + fileContent, _, _, err := client.Repositories.GetContents(ctx, owner, repo, itemPath, &github.RepositoryContentGetOptions{ + Ref: ref, + }) + if err != nil { + return fmt.Errorf("failed to get file %s: %w", itemPath, err) + } + + content, err := fileContent.GetContent() + if err != nil { + return fmt.Errorf("failed to decode file %s: %w", itemPath, err) + } + + if err := os.WriteFile(targetPath, []byte(content), 0o644); err != nil { + return fmt.Errorf("failed to write file %s: %w", targetPath, err) + } + + case "dir": + // Recursively download subdirectory + if err := downloadDirectoryRecursive(ctx, client, owner, repo, ref, itemPath, targetPath); err != nil { + return err + } + } + } + + return nil +} + // extractAllFiles goes through the embedded directory and extracts all files to the current directory func extractAllFiles(embeddedDir string) error { + return extractAllFilesFromFS(EmbeddedObservabilityFiles, embeddedDir) +} + +// extractAllFilesFromFS goes through a filesystem and extracts all files to the current directory +func extractAllFilesFromFS(fsys fs.FS, embeddedDir string) error { // Get current working directory where CLI is running currentDir, err := os.Getwd() if err != nil { return fmt.Errorf("failed to get current directory: %w", err) } - // Walk through the embedded files - err = fs.WalkDir(EmbeddedObservabilityFiles, embeddedDir, func(path string, d fs.DirEntry, err error) error { + // Walk through the files + err = fs.WalkDir(fsys, embeddedDir, func(path string, d fs.DirEntry, err error) error { if err != nil { return fmt.Errorf("error walking the directory: %w", err) } @@ -46,8 +221,8 @@ func extractAllFiles(embeddedDir string) error { return nil } - // Read file content from embedded file system - content, err := EmbeddedObservabilityFiles.ReadFile(path) + // Read file content from file system + content, err := fs.ReadFile(fsys, path) if err != nil { return fmt.Errorf("failed to read file %s: %w", path, err) } @@ -123,15 +298,28 @@ func BlockScoutDown(url string) error { // ObservabilityUpOnlyLoki slim stack with only Loki to verify specific logs of CL nodes or services in tests func ObservabilityUpOnlyLoki() error { + return ObservabilityUpOnlyLokiWithSource("") +} + +// ObservabilityUpOnlyLokiWithSource slim stack with only Loki using custom observability file source +// source can be: +// - "" (empty): use embedded files (default) +// - "file:///path/to/observability": use local filesystem +// - "https://github.com/owner/repo/tree/tag/framework/observability": download from GitHub +func ObservabilityUpOnlyLokiWithSource(source string) error { L.Info().Msg("Creating local observability stack") - if err := extractAllFiles("observability"); err != nil { + fsys, dir, err := resolveObservabilitySource(source) + if err != nil { + return err + } + if err := extractAllFilesFromFS(fsys, dir); err != nil { return err } _ = DefaultNetwork(nil) if err := NewPromtail(); err != nil { return err } - err := RunCommand("bash", "-c", fmt.Sprintf(` + err = RunCommand("bash", "-c", fmt.Sprintf(` cd %s && \ docker compose up -d loki grafana `, "compose")) @@ -145,15 +333,28 @@ func ObservabilityUpOnlyLoki() error { // ObservabilityUp standard stack with logs/metrics for load testing and observability func ObservabilityUp() error { + return ObservabilityUpWithSource("") +} + +// ObservabilityUpWithSource standard stack with logs/metrics using custom observability file source +// source can be: +// - "" (empty): use embedded files (default) +// - "file:///path/to/observability": use local filesystem +// - "https://github.com/owner/repo/tree/tag/framework/observability": download from GitHub +func ObservabilityUpWithSource(source string) error { L.Info().Msg("Creating local observability stack") - if err := extractAllFiles("observability"); err != nil { + fsys, dir, err := resolveObservabilitySource(source) + if err != nil { + return err + } + if err := extractAllFilesFromFS(fsys, dir); err != nil { return err } _ = DefaultNetwork(nil) if err := NewPromtail(); err != nil { return err } - err := RunCommand("bash", "-c", fmt.Sprintf(` + err = RunCommand("bash", "-c", fmt.Sprintf(` cd %s && \ docker compose up -d otel-collector prometheus loki grafana `, "compose")) @@ -170,15 +371,28 @@ func ObservabilityUp() error { // ObservabilityUpFull full stack for load testing and performance investigations func ObservabilityUpFull() error { + return ObservabilityUpFullWithSource("") +} + +// ObservabilityUpFullWithSource full stack for load testing using custom observability file source +// source can be: +// - "" (empty): use embedded files (default) +// - "file:///path/to/observability": use local filesystem +// - "https://github.com/owner/repo/tree/tag/framework/observability": download from GitHub +func ObservabilityUpFullWithSource(source string) error { L.Info().Msg("Creating full local observability stack") - if err := extractAllFiles("observability"); err != nil { + fsys, dir, err := resolveObservabilitySource(source) + if err != nil { + return err + } + if err := extractAllFilesFromFS(fsys, dir); err != nil { return err } _ = DefaultNetwork(nil) if err := NewPromtail(); err != nil { return err } - err := RunCommand("bash", "-c", fmt.Sprintf(` + err = RunCommand("bash", "-c", fmt.Sprintf(` cd %s && \ docker compose up -d `, "compose"))