package workflow import ( "crypto/sha256" "encoding/binary" "errors" "fmt" "slices" "strconv" "strings" "time" _ "time/tzdata" "tangled.org/core/api/tangled" "github.com/bmatcuk/doublestar/v4" "github.com/go-git/go-git/v5/plumbing" "github.com/robfig/cron/v3" "gopkg.in/yaml.v3" ) // - when a repo is modified, it results in the trigger of a "Pipeline" // - a repo could consist of several workflow files // * .tangled/workflows/test.yml // * .tangled/workflows/lint.yml // - therefore a pipeline consists of several workflows, these execute in parallel // - each workflow consists of some execution steps, these execute serially type ( Pipeline []Workflow // this is simply a structural representation of the workflow file Workflow struct { Name string `yaml:"-"` // name of the workflow file Engine string `yaml:"engine"` RunsOn []string `yaml:"runs_on"` When []Constraint `yaml:"when"` CloneOpts CloneOpts `yaml:"clone"` Raw string `yaml:"-"` } Constraint struct { Event StringList `yaml:"event"` Types StringList `yaml:"types"` // optional; only applies to pull_request events. defaults to opened, reopened and synchronize Branch StringList `yaml:"branch"` // required for pull_request; for push, either branch or tag must be specified Tag StringList `yaml:"tag"` // optional; only applies to push events Paths StringList `yaml:"paths"` // optional; only run if any changed file matches a glob pattern Schedule []Schedule `yaml:"schedule"` // required for schedule events } Schedule struct { Cron string `yaml:"cron"` Timezone string `yaml:"timezone"` // optional; IANA timezone, defaults to UTC } CloneOpts struct { Skip bool `yaml:"skip"` Depth int `yaml:"depth"` IncludeSubmodules *bool `yaml:"submodules"` Tags *bool `yaml:"tags"` } StringList []string TriggerKind string ) const ( WorkflowDir = ".tangled/workflows" TriggerKindPush TriggerKind = "push" TriggerKindPullRequest TriggerKind = "pull_request" TriggerKindManual TriggerKind = "manual" TriggerKindSchedule TriggerKind = "schedule" // pull_request lifecycle actions, carried in the trigger metadata and // matched against a constraint's `types` list. PullRequestActionOpened = "opened" PullRequestActionReopened = "reopened" PullRequestActionClosed = "closed" PullRequestActionMerged = "merged" PullRequestActionSynchronize = "synchronize" ) // DefaultPullRequestActions is the set of pull_request actions a constraint // matches when it does not specify an explicit `types` list. This preserves // the historic behaviour of firing on PR creation and resubmission, plus // reopen, while leaving close/merge opt-in. var DefaultPullRequestActions = []string{ PullRequestActionOpened, PullRequestActionReopened, PullRequestActionSynchronize, } func (t TriggerKind) String() string { return strings.ReplaceAll(string(t), "_", " ") } var cronParser = cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow) // HasHashedCron reports whether a comma-separated cron term starts with // Jenkins-style hashed syntax. It intentionally reports malformed uses of H // such as Hfoo as well, so a caller can warn before ParseCron returns its error. func HasHashedCron(expression string) bool { for _, field := range strings.Fields(expression) { for _, term := range strings.Split(field, ",") { if strings.HasPrefix(term, "H") { return true } } } return false } // ParseCron parses a five-field cron expression in timezone using seed to // deterministically expand Jenkins-style hashed fields. func ParseCron(expression, timezone, seed string) (cron.Schedule, error) { return parseCron(expression, timezone, seed, true) } func parseCron(expression, timezone, seed string, validateOccurrence bool) (cron.Schedule, error) { if len(strings.Fields(expression)) != 5 { return nil, fmt.Errorf("expected exactly 5 fields") } if HasHashedCron(expression) { fields := strings.Fields(expression) for i, field := range fields { expanded, err := expandHashedCronField(field, i, seed) if err != nil { return nil, err } fields[i] = expanded } expression = strings.Join(fields, " ") } if timezone == "" { timezone = "UTC" } location, err := time.LoadLocation(timezone) if err != nil { return nil, fmt.Errorf("invalid timezone %q: %w", timezone, err) } schedule, err := cronParser.Parse("CRON_TZ=" + location.String() + " " + expression) if err != nil { return nil, err } if validateOccurrence { if err := validateCronHasOccurrence(schedule); err != nil { return nil, err } } return schedule, nil } // cronFieldBounds returns the bounds accepted for a hashed term. Jenkins // limits an unbounded day-of-month hash to 1..28 so it is safe in every month. func cronFieldBounds(field int, explicitRange bool) (int, int, error) { switch field { case 0: return 0, 59, nil case 1: return 0, 23, nil case 2: if !explicitRange { return 1, 28, nil } return 1, 31, nil case 3: return 1, 12, nil case 4: return 0, 6, nil default: return 0, 0, fmt.Errorf("invalid cron field index %d", field) } } func hashedCronValue(seed string, field int) uint64 { digest := sha256.Sum256([]byte(seed + "\x00" + strconv.Itoa(field))) return binary.BigEndian.Uint64(digest[:8]) } func parseHashedCronNumber(value string) (int, error) { if value == "" { return 0, fmt.Errorf("number is required") } for i := range value { if value[i] < '0' || value[i] > '9' { return 0, fmt.Errorf("%q is not a non-negative integer", value) } } number, err := strconv.Atoi(value) if err != nil { return 0, fmt.Errorf("invalid number %q: %w", value, err) } return number, nil } func expandHashedCronField(field string, fieldIndex int, seed string) (string, error) { terms := strings.Split(field, ",") hashed := false for i, term := range terms { if !strings.HasPrefix(term, "H") { continue } hashed = true expanded, err := expandHashedCronTerm(term, fieldIndex, seed) if err != nil { return "", err } terms[i] = expanded } if !hashed { return field, nil } return strings.Join(terms, ","), nil } func expandHashedCronTerm(term string, fieldIndex int, seed string) (string, error) { if !strings.HasPrefix(term, "H") { return "", fmt.Errorf("malformed hashed cron term %q", term) } rest := term[1:] explicitRange := false lo, hi, err := cronFieldBounds(fieldIndex, false) if err != nil { return "", err } if strings.HasPrefix(rest, "(") { close := strings.IndexByte(rest, ')') if close < 0 { return "", fmt.Errorf("malformed hashed cron term %q: missing ')'", term) } rangeParts := strings.Split(rest[1:close], "-") if len(rangeParts) != 2 { return "", fmt.Errorf("malformed hashed cron term %q: expected H(lo-hi)", term) } lo, err = parseHashedCronNumber(rangeParts[0]) if err != nil { return "", fmt.Errorf("malformed hashed cron term %q: %w", term, err) } hi, err = parseHashedCronNumber(rangeParts[1]) if err != nil { return "", fmt.Errorf("malformed hashed cron term %q: %w", term, err) } explicitRange = true boundLo, boundHi, err := cronFieldBounds(fieldIndex, explicitRange) if err != nil { return "", err } if lo < boundLo || hi > boundHi { return "", fmt.Errorf("hashed cron range %q is outside %d-%d", term, boundLo, boundHi) } if lo > hi { return "", fmt.Errorf("hashed cron range %q is reversed", term) } rest = rest[close+1:] } else if rest != "" && !strings.HasPrefix(rest, "/") { return "", fmt.Errorf("malformed hashed cron term %q", term) } step := 0 if rest != "" { if !strings.HasPrefix(rest, "/") || len(rest) == 1 { return "", fmt.Errorf("malformed hashed cron term %q: expected /step", term) } step, err = parseHashedCronNumber(rest[1:]) if err != nil { return "", fmt.Errorf("malformed hashed cron term %q: %w", term, err) } if step == 0 { return "", fmt.Errorf("hashed cron step must be positive") } } hash := hashedCronValue(seed, fieldIndex) width := hi - lo + 1 if step == 0 || step == 1 { return strconv.Itoa(lo + int(hash%uint64(width))), nil } if step > width { return "", fmt.Errorf("hashed cron step %d exceeds range %d-%d", step, lo, hi) } offset := int(hash % uint64(step)) values := make([]string, 0, (width+step-1)/step) for value := lo + offset; value <= hi; value += step { values = append(values, strconv.Itoa(value)) } return strings.Join(values, ","), nil } func validateCronHasOccurrence(schedule cron.Schedule) error { anchor := time.Now().Add(-time.Minute) if schedule.Next(anchor).IsZero() { return fmt.Errorf("expression has no satisfiable occurrence") } return nil } const ( // ValidateScheduleInterval examines a bounded future window of roughly 400 // days and at most 4096 merged occurrences total. This best-effort bound // covers dense and ordinary sparse calendars while preventing pathological // schedules from hanging validation. scheduleValidationHorizon = 400 * 24 * time.Hour scheduleValidationMaxOccurrences = 4096 ) type scheduleCursor struct { index int schedule cron.Schedule next time.Time } // ValidateScheduleInterval merges actual future occurrences from all schedules, // collapsing identical instants, and rejects distinct occurrences that are // closer than minimum. It scans roughly 400 days and at most 4096 merged // occurrences, returning nil when either bound is reached. func ValidateScheduleInterval(schedules []cron.Schedule, minimum time.Duration, after time.Time) error { if minimum <= 0 { return fmt.Errorf("minimum schedule interval must be positive") } horizon := after.Add(scheduleValidationHorizon) cursors := make([]scheduleCursor, 0, len(schedules)) for index, schedule := range schedules { if schedule == nil { return fmt.Errorf("schedule[%d] is nil", index) } next := schedule.Next(after) if next.IsZero() || next.After(horizon) { continue } if !next.After(after) { return fmt.Errorf("schedule[%d] did not advance beyond %s", index, after) } cursors = append(cursors, scheduleCursor{index: index, schedule: schedule, next: next}) } var previous time.Time havePrevious := false for range scheduleValidationMaxOccurrences { if len(cursors) == 0 { break } earliestIndex := 0 for index := 1; index < len(cursors); index++ { if cursors[index].next.Before(cursors[earliestIndex].next) { earliestIndex = index } } current := cursors[earliestIndex].next if !havePrevious || !current.Equal(previous) { if havePrevious { interval := current.Sub(previous) if interval < minimum { return fmt.Errorf("schedule occurrences %s and %s are %s apart, below minimum %s", previous, current, interval, minimum) } } previous = current havePrevious = true } for index := 0; index < len(cursors); { if !cursors[index].next.Equal(current) { index++ continue } prior := cursors[index].next next := cursors[index].schedule.Next(prior) if next.IsZero() || next.After(horizon) { cursors = append(cursors[:index], cursors[index+1:]...) continue } if !next.After(prior) { return fmt.Errorf("schedule[%d] did not advance beyond %s", cursors[index].index, prior) } cursors[index].next = next index++ } } return nil } // matchesPattern checks if a name matches any of the given patterns. // Patterns can be exact matches or glob patterns using * and **. // * matches any sequence of non-separator characters // ** matches any sequence of characters including separators func matchesPattern(name string, patterns []string) (bool, error) { for _, pattern := range patterns { matched, err := doublestar.Match(pattern, name) if err != nil { return false, err } if matched { return true, nil } } return false, nil } func FromFile(name string, contents []byte) (Workflow, error) { wf := Workflow{Name: name} if err := yaml.Unmarshal(contents, &wf); err != nil { return wf, err } if err := wf.validateSchedules(); err != nil { return wf, err } wf.Raw = string(contents) return wf, nil } func (w Workflow) validateSchedules() error { for i, constraint := range w.When { isSchedule := slices.Contains(constraint.Event, string(TriggerKindSchedule)) switch { case isSchedule && len(constraint.Schedule) == 0: return fmt.Errorf("when[%d]: schedule event requires schedule entries", i) case !isSchedule && len(constraint.Schedule) > 0: return fmt.Errorf("when[%d]: schedule entries require schedule event", i) } if isSchedule && (len(constraint.Branch) > 0 || len(constraint.Tag) > 0 || len(constraint.Types) > 0 || len(constraint.Paths) > 0) { return fmt.Errorf("when[%d]: schedule constraints cannot specify branch, tag, types, or paths", i) } for j, entry := range constraint.Schedule { if entry.Cron == "" { return fmt.Errorf("when[%d].schedule[%d]: cron is required", i, j) } if _, err := parseCron(entry.Cron, entry.Timezone, w.Name, !HasHashedCron(entry.Cron)); err != nil { return fmt.Errorf("when[%d].schedule[%d]: invalid cron expression %q: %w", i, j, entry.Cron, err) } } } return nil } // if any of the constraints on a workflow is true, return true func (w *Workflow) Match(trigger tangled.Pipeline_TriggerMetadata, changedFiles []string) (bool, error) { // manual dispatch skips matching constraints since selection is done by the caller if trigger.Manual != nil { return true, nil } // if not manual, run through the constraint list and see if any one matches for _, c := range w.When { matched, err := c.Match(trigger, changedFiles) if err != nil { return false, err } if matched { return true, nil } } // no constraints, always run this workflow if len(w.When) == 0 { return true, nil } return false, nil } func (c *Constraint) Match(trigger tangled.Pipeline_TriggerMetadata, changedFiles []string) (bool, error) { // manual triggers always pass this constraint if trigger.Manual != nil { return true, nil } // apply event constraints if !c.MatchEvent(trigger.Kind) { return false, nil } // apply branch and action constraints for PRs if trigger.PullRequest != nil { matched, err := c.MatchBranch(trigger.PullRequest.TargetBranch) if err != nil { return false, err } if !matched { return false, nil } action := "" if trigger.PullRequest.Action != nil { action = *trigger.PullRequest.Action } if !c.MatchTypes(action) { return false, nil } } // apply ref constraints for pushes if trigger.Push != nil { matched, err := c.MatchRef(trigger.Push.Ref) if err != nil { return false, err } if !matched { return false, nil } } // apply paths filter: if specified, at least one changed file must match. the // push trigger is the only one we compute a diff for, so a nil list there means // the diff could not be computed, and that must not skip the workflow quietly. // every other trigger keeps the old behaviour of a list that does not match if len(c.Paths) > 0 && (trigger.Push == nil || changedFiles != nil) { matched, err := matchesAnyFile(changedFiles, c.Paths) if err != nil { return false, err } if !matched { return false, nil } } return true, nil } // matchesAnyFile returns true if any file in files matches any of the glob patterns. func matchesAnyFile(files []string, patterns []string) (bool, error) { for _, f := range files { matched, err := matchesPattern(f, patterns) if err != nil { return false, err } if matched { return true, nil } } return false, nil } func (c *Constraint) MatchRef(ref string) (bool, error) { refName := plumbing.ReferenceName(ref) shortName := refName.Short() if refName.IsBranch() { return c.MatchBranch(shortName) } if refName.IsTag() { return c.MatchTag(shortName) } return false, nil } func (c *Constraint) MatchBranch(branch string) (bool, error) { return matchesPattern(branch, c.Branch) } func (c *Constraint) MatchTag(tag string) (bool, error) { return matchesPattern(tag, c.Tag) } func (c *Constraint) MatchEvent(event string) bool { return slices.Contains(c.Event, event) } // MatchTypes reports whether a pull_request action satisfies this constraint's // `types` filter. An empty `types` list falls back to DefaultPullRequestActions. // A missing action (e.g. legacy trigger metadata) is treated as "opened" so // existing pull_request workflows keep matching. func (c *Constraint) MatchTypes(action string) bool { if action == "" { action = PullRequestActionOpened } types := []string(c.Types) if len(types) == 0 { types = DefaultPullRequestActions } return slices.Contains(types, action) } // Custom unmarshaller for StringList func (s *StringList) UnmarshalYAML(unmarshal func(any) error) error { var stringType string if err := unmarshal(&stringType); err == nil { *s = []string{stringType} return nil } var sliceType []any if err := unmarshal(&sliceType); err == nil { if sliceType == nil { *s = nil return nil } parts := make([]string, len(sliceType)) for k, v := range sliceType { if sv, ok := v.(string); ok { parts[k] = sv } else { return fmt.Errorf("cannot unmarshal '%v' of type %T into a string value", v, v) } } *s = parts return nil } return errors.New("failed to unmarshal StringOrSlice") } func (c CloneOpts) AsRecord() tangled.Pipeline_CloneOpts { return tangled.Pipeline_CloneOpts{ Depth: int64(c.Depth), Skip: c.Skip, Submodules: c.IncludeSubmodules == nil || *c.IncludeSubmodules, Tags: c.Tags == nil || *c.Tags, } }