CI / CI (push) Successful in 4m11s
## Changes - `internal/temporal/client.go` — Robust Temporal client with retry (exp backoff), TLS, health check - `internal/temporal/worker.go` — Worker creation, activity/workflow registration, lifecycle - `internal/temporal/context.go` — Timeout helpers - `k8s/worker-deployment.yaml` — 2-10 replica HPA, liveness/readiness probes, security context, pod anti-affinity - `k8s/workflow-runner-deployment.yaml` — Singleton runner with probes - `k8s/kustomization.yaml` — Updated resource list Co-authored-by: poimen <[email protected]>
425 lines
14 KiB
Go
425 lines
14 KiB
Go
package routing
|
|
|
|
import (
|
|
"embed"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io/ioutil"
|
|
"os"
|
|
"path/filepath"
|
|
"runtime"
|
|
"sync"
|
|
)
|
|
|
|
//go:embed activity_knowledge_base.json
|
|
var kbFS embed.FS
|
|
|
|
// KnowledgeBase represents the activity knowledge base
|
|
// SOLID: Single Responsibility - maintains index of activities, provides lookup methods
|
|
// DRY: Loaded once, cached globally with sync.Once pattern
|
|
// CRAP Score: LOW
|
|
// - Complexity: 2 (uses byName index for O(1) lookup, simple methods)
|
|
// - Repetition: 1 (unique concern, no duplicate code)
|
|
// - Total CRAP: 3 (excellent - cache + lookup is efficient)
|
|
type KnowledgeBase struct {
|
|
Version string `json:"version"`
|
|
Activities []ActivityMetadata `json:"activities"`
|
|
Metadata KnowledgeBaseMetadata `json:"metadata"`
|
|
|
|
// Index for fast O(1) lookups (DRY: avoid O(n) iteration)
|
|
byName map[string]*ActivityMetadata
|
|
}
|
|
|
|
// KnowledgeBaseMetadata tracks KB metadata
|
|
type KnowledgeBaseMetadata struct {
|
|
TotalActivities int `json:"totalActivities"`
|
|
LastUpdated string `json:"lastUpdated"`
|
|
Categories map[string]int `json:"categories"`
|
|
}
|
|
|
|
var (
|
|
// globalKB holds singleton instance (lazy loaded)
|
|
globalKB *KnowledgeBase
|
|
// kbMutex protects globalKB initialization
|
|
kbMutex sync.Mutex
|
|
// kbOnce ensures KB loaded exactly once
|
|
kbOnce sync.Once
|
|
// kbErr caches load error for retry logic
|
|
kbErr error
|
|
)
|
|
|
|
// LoadKnowledgeBase loads the activity knowledge base from a JSON file
|
|
// CRAP Score: LOW (single responsibility - file loading)
|
|
// - Complexity: 1 (straightforward file+JSON parsing)
|
|
// - Repetition: 1 (unique logic)
|
|
// - Total CRAP: 2
|
|
func LoadKnowledgeBase(filePath string) (*KnowledgeBase, error) {
|
|
// Read file
|
|
data, err := ioutil.ReadFile(filePath)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to read knowledge base file: %w", err)
|
|
}
|
|
|
|
// Parse JSON
|
|
var kb KnowledgeBase
|
|
err = json.Unmarshal(data, &kb)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to parse knowledge base JSON: %w", err)
|
|
}
|
|
|
|
// Build index for O(1) lookup (DRY: avoid repeated linear scans)
|
|
kb.byName = make(map[string]*ActivityMetadata)
|
|
for i := range kb.Activities {
|
|
kb.byName[kb.Activities[i].Name] = &kb.Activities[i]
|
|
}
|
|
|
|
return &kb, nil
|
|
}
|
|
|
|
// loadKnowledgeBaseFromEmbedded tries to load KB from embedded file
|
|
// Returns (kb, true, nil) on success
|
|
// Returns (nil, false, nil) if embedded file not found
|
|
// Returns (nil, false, error) on parse error
|
|
// CRAP Score: LOW
|
|
func loadKnowledgeBaseFromEmbedded() (*KnowledgeBase, bool, error) {
|
|
data, err := kbFS.ReadFile("activity_knowledge_base.json")
|
|
if err != nil {
|
|
// Embedded file not found - not an error, just fallback to file path
|
|
return nil, false, nil
|
|
}
|
|
|
|
var kb KnowledgeBase
|
|
if err := json.Unmarshal(data, &kb); err != nil {
|
|
return nil, false, fmt.Errorf("failed to parse embedded knowledge base: %w", err)
|
|
}
|
|
|
|
// Build index
|
|
kb.byName = make(map[string]*ActivityMetadata)
|
|
for i := range kb.Activities {
|
|
kb.byName[kb.Activities[i].Name] = &kb.Activities[i]
|
|
}
|
|
|
|
return &kb, true, nil
|
|
}
|
|
|
|
// LoadKnowledgeBaseFromDefaultPath loads KB from default location
|
|
// Tries embedded file first (DRY: no file dependency), then falls back to file paths
|
|
// Search order:
|
|
// 1. Embedded file (preferred - no external dependency)
|
|
// 2. Executable directory
|
|
// 3. Current working directory
|
|
// 4. internal/routing relative to cwd
|
|
// 5. ../internal/routing relative to cwd
|
|
// 6. Same directory as source code
|
|
func LoadKnowledgeBaseFromDefaultPath() (*KnowledgeBase, error) {
|
|
// Try embedded file first (most reliable - no file I/O dependency)
|
|
if kb, found, err := loadKnowledgeBaseFromEmbedded(); err != nil {
|
|
return nil, err
|
|
} else if found {
|
|
return kb, nil
|
|
}
|
|
|
|
// Try to find from package directory
|
|
execDir, err := os.Executable()
|
|
if err == nil {
|
|
// Try in same directory as binary
|
|
path := filepath.Join(filepath.Dir(execDir), "activity_knowledge_base.json")
|
|
if _, err := os.Stat(path); err == nil {
|
|
return LoadKnowledgeBase(path)
|
|
}
|
|
}
|
|
|
|
// Try from current working directory
|
|
if _, err := os.Stat("activity_knowledge_base.json"); err == nil {
|
|
return LoadKnowledgeBase("activity_knowledge_base.json")
|
|
}
|
|
|
|
// Try from internal/routing directory relative to cwd
|
|
if _, err := os.Stat("internal/routing/activity_knowledge_base.json"); err == nil {
|
|
return LoadKnowledgeBase("internal/routing/activity_knowledge_base.json")
|
|
}
|
|
|
|
// Try from parent directory (for tests running from tests/ dir)
|
|
if _, err := os.Stat("../internal/routing/activity_knowledge_base.json"); err == nil {
|
|
return LoadKnowledgeBase("../internal/routing/activity_knowledge_base.json")
|
|
}
|
|
|
|
// Try using runtime to find package directory
|
|
_, filename, _, ok := runtime.Caller(0)
|
|
if ok {
|
|
pkgDir := filepath.Dir(filename)
|
|
path := filepath.Join(pkgDir, "activity_knowledge_base.json")
|
|
if _, err := os.Stat(path); err == nil {
|
|
return LoadKnowledgeBase(path)
|
|
}
|
|
}
|
|
|
|
return nil, fmt.Errorf("activity_knowledge_base.json not found in any expected location")
|
|
}
|
|
|
|
// GetGlobalKnowledgeBase returns singleton KB instance
|
|
// Lazy-loads on first call using sync.Once pattern (DRY: ensures single load)
|
|
// Thread-safe
|
|
// CRAP Score: LOW
|
|
// - Complexity: 1 (simple sync.Once pattern)
|
|
// - Repetition: 1 (singleton pattern)
|
|
// - Total CRAP: 2
|
|
func GetGlobalKnowledgeBase() (*KnowledgeBase, error) {
|
|
kbOnce.Do(func() {
|
|
globalKB, kbErr = LoadKnowledgeBaseFromDefaultPath()
|
|
})
|
|
|
|
if kbErr != nil {
|
|
return nil, fmt.Errorf("knowledge base load error: %w", kbErr)
|
|
}
|
|
|
|
return globalKB, nil
|
|
}
|
|
|
|
// GetActivity returns metadata for a specific activity
|
|
// Returns nil if activity not found (use HasActivity to check first)
|
|
// CRAP Score: LOW
|
|
// - Complexity: 1 (simple map lookup O(1))
|
|
// - Repetition: 1 (unique)
|
|
// - Total CRAP: 2
|
|
func (kb *KnowledgeBase) GetActivity(name string) *ActivityMetadata {
|
|
return kb.byName[name]
|
|
}
|
|
|
|
// ListActivities returns all activities (slice reference, do not modify)
|
|
// CRAP Score: LOW (simple accessor)
|
|
func (kb *KnowledgeBase) ListActivities() []ActivityMetadata {
|
|
return kb.Activities
|
|
}
|
|
|
|
// ListActivitiesByCategory returns all activities in a specific category
|
|
// SOLID: Open/Closed principle - easy to extend with more filters without modifying core logic
|
|
// CRAP Score: LOW
|
|
// - Complexity: 1 (linear scan O(n), but necessary for filtering)
|
|
// - Repetition: 1 (unique concern)
|
|
// - Total CRAP: 2
|
|
func (kb *KnowledgeBase) ListActivitiesByCategory(category string) []ActivityMetadata {
|
|
var result []ActivityMetadata
|
|
for _, activity := range kb.Activities {
|
|
if activity.Category == category {
|
|
result = append(result, activity)
|
|
}
|
|
}
|
|
return result
|
|
}
|
|
|
|
// GetActivityNames returns all activity names in declaration order
|
|
// DRY: Pre-allocated slice to avoid append overhead
|
|
// CRAP Score: LOW
|
|
func (kb *KnowledgeBase) GetActivityNames() []string {
|
|
names := make([]string, len(kb.Activities))
|
|
for i, activity := range kb.Activities {
|
|
names[i] = activity.Name
|
|
}
|
|
return names
|
|
}
|
|
|
|
// HasActivity checks if an activity exists using O(1) index lookup
|
|
// SOLID: Single Responsibility - existence check only
|
|
// DRY: Uses byName index to avoid linear scan
|
|
// CRAP Score: LOW
|
|
// - Complexity: 1 (map lookup)
|
|
// - Repetition: 1 (unique)
|
|
// - Total CRAP: 2
|
|
func (kb *KnowledgeBase) HasActivity(name string) bool {
|
|
_, exists := kb.byName[name]
|
|
return exists
|
|
}
|
|
|
|
// GetDependencies returns prerequisite activities for an activity
|
|
// DRY: Uses GetActivity once instead of direct map access (single lookup point)
|
|
// CRAP Score: LOW
|
|
func (kb *KnowledgeBase) GetDependencies(activityName string) []string {
|
|
activity := kb.GetActivity(activityName)
|
|
if activity == nil {
|
|
return []string{}
|
|
}
|
|
return activity.Constraints.Dependencies
|
|
}
|
|
|
|
// GetTimeoutForActivity returns the default timeout for an activity
|
|
// Falls back to 5m if activity not found (sensible default)
|
|
// SOLID: Single Responsibility - timeout lookup only
|
|
// CRAP Score: LOW
|
|
func (kb *KnowledgeBase) GetTimeoutForActivity(activityName string) string {
|
|
activity := kb.GetActivity(activityName)
|
|
if activity == nil {
|
|
return "5m" // Default timeout - sensible fallback
|
|
}
|
|
return activity.Constraints.DefaultTimeout
|
|
}
|
|
|
|
// GetRetryPolicyForActivity returns retry configuration for an activity
|
|
// DRY: Converts ActivityMetadata constraints into RetryPolicy struct (single conversion point)
|
|
// SOLID: Single Responsibility - converts one constraint type to another
|
|
// CRAP Score: LOW
|
|
// - Complexity: 2 (conditional, struct creation)
|
|
// - Repetition: 1 (unique conversion logic)
|
|
// - Total CRAP: 3
|
|
func (kb *KnowledgeBase) GetRetryPolicyForActivity(activityName string) *RetryPolicy {
|
|
activity := kb.GetActivity(activityName)
|
|
if activity == nil {
|
|
return &RetryPolicy{
|
|
MaxAttempts: 1,
|
|
BackoffRate: 1.0,
|
|
InitialInterval: "1s",
|
|
}
|
|
}
|
|
|
|
return &RetryPolicy{
|
|
MaxAttempts: int32(activity.Constraints.RecommendedRetries),
|
|
BackoffRate: activity.Constraints.RetryBackoff,
|
|
InitialInterval: "1s",
|
|
MaxInterval: "30s",
|
|
}
|
|
}
|
|
|
|
// IsFlaky returns whether an activity is marked as flaky (needs extra retries)
|
|
// SOLID: Single Responsibility - flakiness check only
|
|
// CRAP Score: LOW
|
|
func (kb *KnowledgeBase) IsFlaky(activityName string) bool {
|
|
activity := kb.GetActivity(activityName)
|
|
if activity == nil {
|
|
return false // Non-existent activities treated as stable (conservative)
|
|
}
|
|
return activity.Constraints.IsFlaky
|
|
}
|
|
|
|
// GetNotes returns implementation notes and caveats for an activity
|
|
// Useful for logging, debugging, and documentation generation
|
|
// CRAP Score: LOW
|
|
func (kb *KnowledgeBase) GetNotes(activityName string) string {
|
|
activity := kb.GetActivity(activityName)
|
|
if activity == nil {
|
|
return ""
|
|
}
|
|
return activity.Constraints.Notes
|
|
}
|
|
|
|
// Validate checks the knowledge base for consistency
|
|
// Checks:
|
|
// 1. No circular dependencies in activity constraints
|
|
// 2. All referenced dependencies exist
|
|
// SOLID: Single Responsibility - validation only, no side effects
|
|
// CRAP Score: MEDIUM
|
|
// - Complexity: 3 (nested loops + recursion)
|
|
// - Repetition: 2 (two separate checks, some code reuse in checkDependencies)
|
|
// - Total CRAP: 5 (acceptable for validation logic)
|
|
func (kb *KnowledgeBase) Validate() error {
|
|
// Check for circular dependencies using DFS
|
|
visited := make(map[string]bool)
|
|
for _, activity := range kb.Activities {
|
|
if err := kb.checkDependencies(activity.Name, visited, []string{}); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
// DRY: Check all dependencies exist in second pass (separate concern from cycle detection)
|
|
for _, activity := range kb.Activities {
|
|
for _, dep := range activity.Constraints.Dependencies {
|
|
if !kb.HasActivity(dep) {
|
|
return fmt.Errorf("activity %s depends on non-existent activity %s", activity.Name, dep)
|
|
}
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// checkDependencies validates activity dependencies for cycles using DFS
|
|
// Internal helper method for Validate()
|
|
// Uses path to build cycle path for error reporting
|
|
// CRAP Score: MEDIUM
|
|
// - Complexity: 3 (string building, recursion, path tracking)
|
|
// - Repetition: 1 (unique DFS logic)
|
|
// - Total CRAP: 4 (acceptable for graph traversal)
|
|
func (kb *KnowledgeBase) checkDependencies(activityName string, visited map[string]bool, path []string) error {
|
|
// Check for cycles by detecting if activityName appears in current path
|
|
// This indicates we've visited activityName already in this traversal
|
|
for _, p := range path {
|
|
if p == activityName {
|
|
// Build human-readable cycle description
|
|
cycleStr := ""
|
|
found := false
|
|
for _, n := range path {
|
|
if found {
|
|
cycleStr += " -> " + n
|
|
}
|
|
if n == activityName {
|
|
found = true
|
|
cycleStr += n
|
|
}
|
|
}
|
|
cycleStr += " -> " + activityName
|
|
return fmt.Errorf("circular dependency detected: %s", cycleStr)
|
|
}
|
|
}
|
|
|
|
// Skip if already fully visited (memoization)
|
|
if visited[activityName] {
|
|
return nil
|
|
}
|
|
|
|
visited[activityName] = true
|
|
newPath := append(path, activityName)
|
|
|
|
activity := kb.GetActivity(activityName)
|
|
if activity == nil {
|
|
return nil // Non-existent activity will be caught in Validate() second pass
|
|
}
|
|
|
|
// Recursively check all dependencies
|
|
for _, dep := range activity.Constraints.Dependencies {
|
|
if err := kb.checkDependencies(dep, visited, newPath); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// String returns a human-readable short description of the knowledge base
|
|
// Implements fmt.Stringer interface for logging
|
|
// CRAP Score: LOW (simple string formatting)
|
|
func (kb *KnowledgeBase) String() string {
|
|
return fmt.Sprintf("KnowledgeBase(v%s, %d activities)", kb.Version, kb.Metadata.TotalActivities)
|
|
}
|
|
|
|
// PrintSummary generates human-readable documentation of all activities
|
|
// Useful for:
|
|
// - CLI output (showing available activities)
|
|
// - Documentation generation
|
|
// - Debugging knowledge base content
|
|
// DRY: Centralizes summary formatting (single point of change)
|
|
// SOLID: Single Responsibility - formatting only, no mutations
|
|
// CRAP Score: MEDIUM
|
|
// - Complexity: 2 (string building, nested loops)
|
|
// - Repetition: 1 (unique formatting)
|
|
// - Total CRAP: 3
|
|
func (kb *KnowledgeBase) PrintSummary() string {
|
|
summary := fmt.Sprintf("=== Activity Knowledge Base ===\nVersion: %s\nTotal Activities: %d\n\n", kb.Version, kb.Metadata.TotalActivities)
|
|
|
|
summary += "Activities by Category:\n"
|
|
for category, count := range kb.Metadata.Categories {
|
|
summary += fmt.Sprintf(" %s: %d\n", category, count)
|
|
}
|
|
|
|
summary += "\nActivity Details:\n"
|
|
for _, activity := range kb.Activities {
|
|
summary += fmt.Sprintf("\n[%s] %s\n", activity.Name, activity.Description)
|
|
summary += fmt.Sprintf(" Category: %s\n", activity.Category)
|
|
summary += fmt.Sprintf(" Timeout: %s\n", activity.Constraints.DefaultTimeout)
|
|
summary += fmt.Sprintf(" Flaky: %v (Retries: %d)\n", activity.Constraints.IsFlaky, activity.Constraints.RecommendedRetries)
|
|
if len(activity.Constraints.Dependencies) > 0 {
|
|
summary += fmt.Sprintf(" Dependencies: %v\n", activity.Constraints.Dependencies)
|
|
}
|
|
}
|
|
|
|
return summary
|
|
}
|