62 lines
1.7 KiB
Go
62 lines
1.7 KiB
Go
package reconciler
|
|
|
|
import (
|
|
"context"
|
|
"time"
|
|
)
|
|
|
|
// CycleTimeout — верхняя граница одного цикла reconcile.
|
|
const CycleTimeout = 5 * time.Minute
|
|
|
|
// Run запускает цикл: сразу, затем каждые interval, а также по сигналу trigger
|
|
// (после debounce, чтобы схлопнуть пачку событий). Блокируется до отмены ctx.
|
|
// report вызывается после каждого цикла.
|
|
func (r *Reconciler) Run(ctx context.Context, interval, debounce time.Duration, trigger <-chan struct{}, report func(Result, error)) {
|
|
cycle := func() {
|
|
cctx, cancel := context.WithTimeout(ctx, CycleTimeout)
|
|
defer cancel()
|
|
res, err := r.Reconcile(cctx)
|
|
if err != nil {
|
|
r.log.Error("цикл reconcile завершён с ошибками", "error", err)
|
|
} else {
|
|
r.log.Info("цикл reconcile завершён",
|
|
"created", res.Created, "updated", res.Updated, "deleted", res.Deleted,
|
|
"unchanged", res.Unchanged, "skippedForeign", res.SkippedForeign,
|
|
"skippedConflict", res.SkippedConflict, "pendingDelete", res.PendingDelete)
|
|
}
|
|
if report != nil {
|
|
report(res, err)
|
|
}
|
|
}
|
|
|
|
cycle()
|
|
ticker := time.NewTicker(interval)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
cycle()
|
|
case <-trigger:
|
|
if debounce > 0 {
|
|
t := time.NewTimer(debounce)
|
|
select {
|
|
case <-ctx.Done():
|
|
t.Stop()
|
|
return
|
|
case <-t.C:
|
|
}
|
|
}
|
|
// схлопываем накопившиеся сигналы
|
|
select {
|
|
case <-trigger:
|
|
default:
|
|
}
|
|
r.log.Debug("reconcile по событию Docker")
|
|
cycle()
|
|
ticker.Reset(interval)
|
|
}
|
|
}
|
|
}
|