package app import ( "context" "log/slog" "time" "github.com/dtoro/oikos/internal/core/ports" ) // ObserveConfig controls the observe pass. type ObserveConfig struct { MaxConcurrency int } // ObserveResult summarizes one pass. type ObserveResult struct { ChecksRun int HealthChanges int Duration time.Duration } // CheckerLookup resolves a check kind to its probe implementation. type CheckerLookup interface { Get(kind string) ports.Checker } // ObservationService runs observe passes: load enabled checks, resolve // targets, run probes via the Checker registry, process results through // SignalService, and aggregate health. type ObservationService struct { checks ports.CheckRepository signals *SignalService targets ports.TargetResolver reg CheckerLookup } func NewObservationService( checks ports.CheckRepository, signals *SignalService, targets ports.TargetResolver, reg CheckerLookup, ) *ObservationService { return &ObservationService{checks: checks, signals: signals, targets: targets, reg: reg} } // RunPass loads enabled checks, resolves targets, runs probes with bounded // concurrency, and processes results through SignalService. func (o *ObservationService) RunPass(ctx context.Context, cfg ObserveConfig) (ObserveResult, error) { start := time.Now() var res ObserveResult defs, err := o.checks.ListEnabled(ctx) if err != nil { return res, err } if cfg.MaxConcurrency <= 0 { cfg.MaxConcurrency = 10 } type work struct { def ports.CheckDef target ports.Target health string } var jobs []work for _, d := range defs { checker := o.reg.Get(d.Kind) if checker == nil { slog.Debug("observation: no checker for kind", "kind", d.Kind, "check", d.ID) continue } target, err := o.targets.ResolveForCheck(ctx, d.EntityID, d.Kind) if err != nil { slog.Warn("observation: resolve target", "check", d.ID, "error", err) continue } jobs = append(jobs, work{def: d, target: target}) } sem := make(chan struct{}, cfg.MaxConcurrency) type result struct { j work out ports.CheckResult } results := make(chan result, len(jobs)) for _, j := range jobs { j := j go func() { sem <- struct{}{} defer func() { <-sem }() checker := o.reg.Get(j.def.Kind) out := checker.Check(ctx, j.def, j.target) results <- result{j: j, out: out} }() } var healthChanges int for i := 0; i < len(jobs); i++ { r := <-results res.ChecksRun++ hc, err := o.signals.ProcessCheckResult(ctx, CheckOutcome{ EntityID: r.j.def.ID, TargetID: r.j.def.EntityID, Slug: r.j.def.Name, Kind: r.j.def.Kind, Value: r.out.Value, State: r.out.State, Message: r.out.Message, PrevHealth: "", }) if err != nil { slog.Warn("observation: signal processing failed", "check", r.j.def.ID, "error", err) continue } if hc != nil { healthChanges++ } } res.HealthChanges = healthChanges res.Duration = time.Since(start) return res, nil }