package scheduler import ( "errors" "fmt" "sort" "sync" "time" ) // Priority ist die Dringlichkeit eines Auftrags. type Priority string const ( // PriorityCritical ist die höchste Stufe. PriorityCritical Priority = "critical" // PriorityHigh ist hohe Dringlichkeit. PriorityHigh Priority = "high" // PriorityNormal ist der Standard. PriorityNormal Priority = "normal" // PriorityLow ist niedrige Dringlichkeit. PriorityLow Priority = "low" ) // priorityRanks bilden die Stufen auf vergleichbare Zahlen ab. var priorityRanks = map[Priority]int{ PriorityCritical: 400, PriorityHigh: 300, PriorityNormal: 200, PriorityLow: 100, } // Rank liefert den Zahlenwert einer Stufe. func (priority Priority) Rank() int { if rankValue, isKnown := priorityRanks[priority]; isKnown { return rankValue } // Eine unbekannte Stufe gilt als normal statt als höchste. Ein Tippfehler // soll einen Auftrag nicht an die Spitze der Warteschlange befördern. return priorityRanks[PriorityNormal] } // IsValid meldet eine bekannte Stufe. func (priority Priority) IsValid() bool { _, isKnown := priorityRanks[priority] return isKnown } // agingBonusPerHour ist der Prioritätsgewinn je Wartestunde. // // Ohne ihn verhungerte ein Auftrag niedriger Priorität in einer Anlage mit // genug dringenden Aufträgen für immer — und niemand bemerkte es, weil er // technisch „wartet" statt zu scheitern. Mit 50 Punkten je Stunde überholt ein // wartender Auftrag nach zwei Stunden die nächsthöhere Stufe. const agingBonusPerHour = 50 // maximumAgingBonus begrenzt den Alterungsgewinn. // // Ohne Deckel überholte ein tagelang wartender Auftrag jede kritische // Sicherung. Der Deckel entspricht dem Abstand zweier Stufen: ein wartender // Auftrag steigt höchstens um eine Stufe auf. const maximumAgingBonus = 100 // QueuedJob ist ein zur Ausführung anstehender Auftrag. type QueuedJob struct { // JobID ist die Kennung des Auftrags. JobID string // Priority ist die Dringlichkeit. Priority Priority // ScheduledFor ist der geplante Zeitpunkt. ScheduledFor time.Time // EnqueuedAt ist der Zeitpunkt der Einreihung. EnqueuedAt time.Time // DependsOnJobIDs sind Aufträge, die zuerst gelingen müssen. DependsOnJobIDs []string // AttemptNumber ist die Nummer des Versuchs, beginnend bei 1. AttemptNumber int } // effectivePriority berechnet die Dringlichkeit einschließlich Alterung. func (queuedJob *QueuedJob) effectivePriority(currentTime time.Time) int { waitedHours := currentTime.Sub(queuedJob.EnqueuedAt).Hours() if waitedHours < 0 { waitedHours = 0 } agingBonus := int(waitedHours * agingBonusPerHour) if agingBonus > maximumAgingBonus { agingBonus = maximumAgingBonus } return queuedJob.Priority.Rank() + agingBonus } // ConcurrencyLimits begrenzen die gleichzeitige Ausführung. type ConcurrencyLimits struct { // GlobalLimit ist die Gesamtzahl gleichzeitiger Läufe; 0 bedeutet unbegrenzt. GlobalLimit int // PerRepositoryLimit begrenzt Läufe je Repository. // // Nötig, weil das Repository beim Schreiben eine Sperre hält: mehrere // gleichzeitige Läufe auf dasselbe Repository warten ohnehin aufeinander // und belegen dabei nur Arbeitsspeicher. PerRepositoryLimit int // PerSourceLimit begrenzt Läufe je Quelle. // // Zwei gleichzeitige Sicherungen derselben VM lasten den Wirt aus, ohne // zusätzlichen Nutzen. PerSourceLimit int } // RunningJob beschreibt einen laufenden Auftrag. type RunningJob struct { // JobID ist die Kennung des Auftrags. JobID string // RepositoryID ist das Ziel-Repository. RepositoryID string // SourceID ist die gesicherte Quelle. SourceID string // StartedAt ist der Beginn. StartedAt time.Time } // JobOutcome ist der Ausgang eines abgeschlossenen Laufs. type JobOutcome string const ( // OutcomeSucceeded ist ein vollständig gelungener Lauf. OutcomeSucceeded JobOutcome = "succeeded" // OutcomePartialFailure ist ein Lauf mit übergangenen Objekten. // // Er gilt ausdrücklich **nicht** als Erfolg (PROMPT.md §138). Abhängige // Aufträge laufen nach einem Teilfehler nicht an. OutcomePartialFailure JobOutcome = "partial_failure" // OutcomeFailed ist ein gescheiterter Lauf. OutcomeFailed JobOutcome = "failed" // OutcomeCancelled ist ein abgebrochener Lauf. OutcomeCancelled JobOutcome = "cancelled" ) // IsSuccess meldet einen vollständig gelungenen Ausgang. func (outcome JobOutcome) IsSuccess() bool { return outcome == OutcomeSucceeded } // ErrDependencyCycle meldet einen Ringschluss zwischen Aufträgen. var ErrDependencyCycle = errors.New("die abhängigkeiten der aufträge bilden einen ring") // ErrJobAlreadyQueued meldet einen bereits eingereihten Auftrag. var ErrJobAlreadyQueued = errors.New("der auftrag steht bereits in der warteschlange") // Queue verwaltet anstehende und laufende Aufträge. // // Sie ist bewusst zustandsbehaftet und nebenläufigkeitssicher, aber ohne // Datenbank: Was ausgeführt werden darf, ist eine Rechenfrage. Die Beständigkeit // über Neustarts hinweg liegt eine Schicht darüber. type Queue struct { // mutex schützt den Zustand. mutex sync.Mutex // pendingJobs sind die wartenden Aufträge. pendingJobs []QueuedJob // runningJobs sind die laufenden Aufträge, nach Auftragskennung. runningJobs map[string]RunningJob // lastOutcomes hält den letzten Ausgang je Auftrag für die Abhängigkeiten. lastOutcomes map[string]JobOutcome // limits sind die Nebenläufigkeitsgrenzen. limits ConcurrencyLimits // jobMetadata hält Repository und Quelle je Auftrag. jobMetadata map[string]RunningJob } // NewQueue erzeugt eine Warteschlange mit den angegebenen Grenzen. func NewQueue(limits ConcurrencyLimits) *Queue { return &Queue{ pendingJobs: make([]QueuedJob, 0, 16), runningJobs: make(map[string]RunningJob), lastOutcomes: make(map[string]JobOutcome), jobMetadata: make(map[string]RunningJob), limits: limits, } } // RegisterJobMetadata hinterlegt Repository und Quelle eines Auftrags. // // Ohne diese Angaben liessen sich die auftragsbezogenen Grenzen nicht prüfen. func (queue *Queue) RegisterJobMetadata(jobIdentifier string, repositoryIdentifier string, sourceIdentifier string) { queue.mutex.Lock() defer queue.mutex.Unlock() queue.jobMetadata[jobIdentifier] = RunningJob{ JobID: jobIdentifier, RepositoryID: repositoryIdentifier, SourceID: sourceIdentifier, } } // Enqueue reiht einen Auftrag ein. // // Ein bereits wartender oder laufender Auftrag wird abgelehnt. Andernfalls // entstünden bei einem überlasteten System immer mehr Einträge desselben // Auftrags, bis die Warteschlange den Speicher füllt. func (queue *Queue) Enqueue(queuedJob QueuedJob) error { queue.mutex.Lock() defer queue.mutex.Unlock() if _, isRunning := queue.runningJobs[queuedJob.JobID]; isRunning { return fmt.Errorf("%w: er läuft bereits", ErrJobAlreadyQueued) } for _, pendingJob := range queue.pendingJobs { if pendingJob.JobID == queuedJob.JobID { return ErrJobAlreadyQueued } } if queuedJob.EnqueuedAt.IsZero() { queuedJob.EnqueuedAt = time.Now().UTC() } if queuedJob.AttemptNumber <= 0 { queuedJob.AttemptNumber = 1 } queue.pendingJobs = append(queue.pendingJobs, queuedJob) return nil } // DispatchDecision erklärt, warum ein Auftrag läuft oder wartet. type DispatchDecision struct { // Job ist der betroffene Auftrag. Job QueuedJob // CanRun meldet die Ausführbarkeit. CanRun bool // Reason erklärt eine Ablehnung verständlich. // // Ein wartender Auftrag ohne Begründung ist für den Betrieb wertlos: Man // sieht, dass nichts geschieht, aber nicht warum. Reason string } // Dispatch wählt die als Nächstes auszuführenden Aufträge aus. // // Zurückgegeben werden alle Aufträge, die zum angegebenen Zeitpunkt starten // dürfen — in der Reihenfolge ihrer Dringlichkeit. Der Aufrufer startet sie und // meldet den Beginn über MarkRunning. func (queue *Queue) Dispatch(currentTime time.Time) []DispatchDecision { queue.mutex.Lock() defer queue.mutex.Unlock() dueJobs := make([]QueuedJob, 0, len(queue.pendingJobs)) for _, pendingJob := range queue.pendingJobs { if !pendingJob.ScheduledFor.After(currentTime) { dueJobs = append(dueJobs, pendingJob) } } queue.sortByEffectivePriority(dueJobs, currentTime) // Die Zählungen werden mitgeführt, damit mehrere in einem Durchgang // freigegebene Aufträge die Grenzen nicht gemeinsam überschreiten. runningPerRepository := make(map[string]int) runningPerSource := make(map[string]int) runningTotal := len(queue.runningJobs) for _, runningJob := range queue.runningJobs { runningPerRepository[runningJob.RepositoryID]++ runningPerSource[runningJob.SourceID]++ } decisions := make([]DispatchDecision, 0, len(dueJobs)) for _, dueJob := range dueJobs { jobMetadata := queue.jobMetadata[dueJob.JobID] if blockingReason := queue.checkDependencies(dueJob); blockingReason != "" { decisions = append(decisions, DispatchDecision{Job: dueJob, Reason: blockingReason}) continue } if queue.limits.GlobalLimit > 0 && runningTotal >= queue.limits.GlobalLimit { decisions = append(decisions, DispatchDecision{ Job: dueJob, Reason: fmt.Sprintf("Es laufen bereits %d Aufträge; mehr als %d sind nicht zugelassen.", runningTotal, queue.limits.GlobalLimit), }) continue } if queue.limits.PerRepositoryLimit > 0 && jobMetadata.RepositoryID != "" && runningPerRepository[jobMetadata.RepositoryID] >= queue.limits.PerRepositoryLimit { decisions = append(decisions, DispatchDecision{ Job: dueJob, Reason: fmt.Sprintf("Auf dem Repository %s laufen bereits %d Aufträge.", jobMetadata.RepositoryID, queue.limits.PerRepositoryLimit), }) continue } if queue.limits.PerSourceLimit > 0 && jobMetadata.SourceID != "" && runningPerSource[jobMetadata.SourceID] >= queue.limits.PerSourceLimit { decisions = append(decisions, DispatchDecision{ Job: dueJob, Reason: fmt.Sprintf("Die Quelle %s wird bereits gesichert.", jobMetadata.SourceID), }) continue } decisions = append(decisions, DispatchDecision{Job: dueJob, CanRun: true}) runningTotal++ runningPerRepository[jobMetadata.RepositoryID]++ runningPerSource[jobMetadata.SourceID]++ } return decisions } // sortByEffectivePriority ordnet Aufträge nach Dringlichkeit. // // Bei gleicher Dringlichkeit entscheidet der geplante Zeitpunkt, dann die // Kennung. Die letzte Stufe ist wichtig: Ohne sie wäre die Reihenfolge bei // gleichem Zeitpunkt zufällig und zwei Läufe derselben Anlage nicht // vergleichbar. func (queue *Queue) sortByEffectivePriority(jobsToSort []QueuedJob, currentTime time.Time) { sort.SliceStable(jobsToSort, func(firstIndex int, secondIndex int) bool { firstPriority := jobsToSort[firstIndex].effectivePriority(currentTime) secondPriority := jobsToSort[secondIndex].effectivePriority(currentTime) if firstPriority != secondPriority { return firstPriority > secondPriority } if !jobsToSort[firstIndex].ScheduledFor.Equal(jobsToSort[secondIndex].ScheduledFor) { return jobsToSort[firstIndex].ScheduledFor.Before(jobsToSort[secondIndex].ScheduledFor) } return jobsToSort[firstIndex].JobID < jobsToSort[secondIndex].JobID }) } // checkDependencies prüft die Abhängigkeiten eines Auftrags. // // Eine Abhängigkeit ist erfüllt, wenn der Vorgänger zuletzt **vollständig** // gelungen ist. Ein Teilfehler genügt nicht: Wer eine Datenbank sichert und // danach das Anwendungsverzeichnis, will nicht das Verzeichnis zu einer // halben Datenbank. func (queue *Queue) checkDependencies(queuedJob QueuedJob) string { for _, dependencyJobID := range queuedJob.DependsOnJobIDs { if _, isRunning := queue.runningJobs[dependencyJobID]; isRunning { return fmt.Sprintf("Der vorausgesetzte Auftrag %s läuft noch.", dependencyJobID) } lastOutcome, hasOutcome := queue.lastOutcomes[dependencyJobID] if !hasOutcome { return fmt.Sprintf("Der vorausgesetzte Auftrag %s ist noch nie gelaufen.", dependencyJobID) } if !lastOutcome.IsSuccess() { return fmt.Sprintf("Der vorausgesetzte Auftrag %s endete zuletzt mit %s.", dependencyJobID, lastOutcome) } } return "" } // MarkRunning vermerkt den Beginn eines Laufs. func (queue *Queue) MarkRunning(jobIdentifier string, startTime time.Time) { queue.mutex.Lock() defer queue.mutex.Unlock() remainingJobs := make([]QueuedJob, 0, len(queue.pendingJobs)) for _, pendingJob := range queue.pendingJobs { if pendingJob.JobID != jobIdentifier { remainingJobs = append(remainingJobs, pendingJob) } } queue.pendingJobs = remainingJobs jobMetadata := queue.jobMetadata[jobIdentifier] queue.runningJobs[jobIdentifier] = RunningJob{ JobID: jobIdentifier, RepositoryID: jobMetadata.RepositoryID, SourceID: jobMetadata.SourceID, StartedAt: startTime, } } // MarkFinished vermerkt das Ende eines Laufs. func (queue *Queue) MarkFinished(jobIdentifier string, outcome JobOutcome) { queue.mutex.Lock() defer queue.mutex.Unlock() delete(queue.runningJobs, jobIdentifier) queue.lastOutcomes[jobIdentifier] = outcome } // RunningCount liefert die Zahl laufender Aufträge. func (queue *Queue) RunningCount() int { queue.mutex.Lock() defer queue.mutex.Unlock() return len(queue.runningJobs) } // PendingCount liefert die Zahl wartender Aufträge. func (queue *Queue) PendingCount() int { queue.mutex.Lock() defer queue.mutex.Unlock() return len(queue.pendingJobs) } // Remove nimmt einen Auftrag aus der Warteschlange. func (queue *Queue) Remove(jobIdentifier string) bool { queue.mutex.Lock() defer queue.mutex.Unlock() remainingJobs := make([]QueuedJob, 0, len(queue.pendingJobs)) var wasRemoved bool for _, pendingJob := range queue.pendingJobs { if pendingJob.JobID == jobIdentifier { wasRemoved = true continue } remainingJobs = append(remainingJobs, pendingJob) } queue.pendingJobs = remainingJobs return wasRemoved } // ValidateDependencyGraph prüft eine Abhängigkeitsmenge auf Ringschlüsse. // // Ein Ring bliebe sonst unbemerkt: Jeder beteiligte Auftrag wartete auf einen // anderen, keiner liefe je an, und die Oberfläche zeigte lauter „wartende" // Aufträge ohne erkennbaren Grund. Deshalb wird beim Anlegen geprüft, nicht // beim Ausführen. func ValidateDependencyGraph(dependenciesByJob map[string][]string) error { // Tiefensuche mit drei Zuständen: unbesucht, in Bearbeitung, fertig. const ( stateUnvisited = 0 stateInProgress = 1 stateCompleted = 2 ) visitState := make(map[string]int, len(dependenciesByJob)) // Die feste Reihenfolge macht die Fehlermeldung bei mehreren Ringen // reproduzierbar. jobIdentifiers := make([]string, 0, len(dependenciesByJob)) for jobIdentifier := range dependenciesByJob { jobIdentifiers = append(jobIdentifiers, jobIdentifier) } sort.Strings(jobIdentifiers) var visitJob func(jobIdentifier string, visitPath []string) error visitJob = func(jobIdentifier string, visitPath []string) error { switch visitState[jobIdentifier] { case stateInProgress: return fmt.Errorf("%w: %s", ErrDependencyCycle, formatCyclePath(append(visitPath, jobIdentifier))) case stateCompleted: return nil } visitState[jobIdentifier] = stateInProgress for _, dependencyJobID := range dependenciesByJob[jobIdentifier] { if dependencyJobID == jobIdentifier { return fmt.Errorf("%w: der auftrag %s hängt von sich selbst ab", ErrDependencyCycle, jobIdentifier) } if visitError := visitJob(dependencyJobID, append(visitPath, jobIdentifier)); visitError != nil { return visitError } } visitState[jobIdentifier] = stateCompleted return nil } for _, jobIdentifier := range jobIdentifiers { if visitError := visitJob(jobIdentifier, nil); visitError != nil { return visitError } } return nil } // formatCyclePath beschreibt einen Ringschluss lesbar. func formatCyclePath(visitPath []string) string { pathText := "" for pathIndex, jobIdentifier := range visitPath { if pathIndex > 0 { pathText += " → " } pathText += jobIdentifier } return pathText }