feat: Phase 5bc — SignalService + ObservationService + MetricsRepo/SignalRepo
Problem: signal lifecycle (upsert, resolve, health aggregation) and observe-pass orchestration (load checks, resolve targets, run probes, aggregate health) were embedded in scheduler/scheduler.go — 1095 lines of monolith with no port abstraction. Change: - app/signals.go: SignalService — ProcessCheckResult evaluates probe outcomes (upserts signals on critical/warning, resolves on ok), records metrics, computes health changes (ok/degraded/down/stale). WorstHealthForTarget aggregates open signals into entity health. - app/observation.go: ObservationService — RunPass loads enabled checks via CheckRepository, resolves targets via TargetResolver, dispatches probes through CheckerLookup (probes.Registry) with bounded concurrency (default 10), sends results through SignalService. - adapters/postgres/signals.go: MetricsRepo (InsertSamples via sqlcgen InsertMetricSample), SignalRepo (Open/UpsertWithTriggers/ Transition with inline SQL matching the scheduler's patterns). Verification: go build/vet, full test suite (18 pkgs green), DB integration (postgres + mcp — green).
This commit is contained in:
88
internal/adapters/postgres/signals.go
Normal file
88
internal/adapters/postgres/signals.go
Normal file
@@ -0,0 +1,88 @@
|
|||||||
|
package db
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
|
||||||
|
"github.com/dtoro/oikos/internal/adapters/postgres/sqlcgen"
|
||||||
|
"github.com/dtoro/oikos/internal/core/domain"
|
||||||
|
"github.com/dtoro/oikos/internal/core/ports"
|
||||||
|
"github.com/google/uuid"
|
||||||
|
)
|
||||||
|
|
||||||
|
// MetricsRepo implements ports.MetricsRepository.
|
||||||
|
type MetricsRepo struct {
|
||||||
|
pool *Pool
|
||||||
|
}
|
||||||
|
|
||||||
|
var _ ports.MetricsRepository = (*MetricsRepo)(nil)
|
||||||
|
|
||||||
|
func NewMetricsRepo(pool *Pool) *MetricsRepo { return &MetricsRepo{pool: pool} }
|
||||||
|
|
||||||
|
func (m *MetricsRepo) InsertSamples(ctx context.Context, entityID domain.UUID, samples []ports.MetricSample) error {
|
||||||
|
for _, s := range samples {
|
||||||
|
if err := sqlcgen.New(m.pool).InsertMetricSample(ctx, sqlcgen.InsertMetricSampleParams{
|
||||||
|
EntityID: mustUUID(entityID), Metric: s.Metric, Value: s.Value,
|
||||||
|
}); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// SignalRepo implements ports.SignalRepository with inline SQL.
|
||||||
|
type SignalRepo struct {
|
||||||
|
pool *Pool
|
||||||
|
}
|
||||||
|
|
||||||
|
var _ ports.SignalRepository = (*SignalRepo)(nil)
|
||||||
|
|
||||||
|
func NewSignalRepo(pool *Pool) *SignalRepo { return &SignalRepo{pool: pool} }
|
||||||
|
|
||||||
|
func (r *SignalRepo) Open(ctx context.Context) ([]domain.Signal, error) {
|
||||||
|
rows, err := r.pool.Query(ctx,
|
||||||
|
`SELECT entity_id, kind, severity, state FROM signals WHERE state NOT IN ('resolved','failed') ORDER BY severity DESC`)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
defer rows.Close()
|
||||||
|
var items []domain.Signal
|
||||||
|
for rows.Next() {
|
||||||
|
var s domain.Signal
|
||||||
|
var id uuid.UUID
|
||||||
|
if err := rows.Scan(&id, &s.Kind, &s.Severity, &s.State); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
s.EntityID = domain.UUID(id.String())
|
||||||
|
items = append(items, s)
|
||||||
|
}
|
||||||
|
return items, rows.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *SignalRepo) History(ctx context.Context, entityID domain.UUID, limit int) ([]domain.Signal, error) {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *SignalRepo) UpsertWithTriggers(ctx context.Context, input ports.SignalUpsertInput) error {
|
||||||
|
eid := mustUUID(input.Signal.EntityID)
|
||||||
|
_, err := r.pool.Exec(ctx,
|
||||||
|
`INSERT INTO signals (entity_id, kind, severity, target_entity_id, state)
|
||||||
|
VALUES ($1, $2, $3, $4, 'raised')
|
||||||
|
ON CONFLICT (target_entity_id, kind) WHERE state NOT IN ('resolved','failed')
|
||||||
|
DO UPDATE SET occurrence_count = signals.occurrence_count + 1,
|
||||||
|
last_seen_at = now(), updated_at = now()`,
|
||||||
|
eid, input.Signal.Kind, input.Signal.Severity, eid)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *SignalRepo) Transition(ctx context.Context, input ports.SignalTransitionInput) (domain.Signal, error) {
|
||||||
|
tag, err := r.pool.Exec(ctx,
|
||||||
|
`UPDATE signals SET state = 'resolved', updated_at = now()
|
||||||
|
WHERE entity_id = $1 AND state IN ('raised', 'acknowledged')`, mustUUID(input.SignalID))
|
||||||
|
if err != nil {
|
||||||
|
return domain.Signal{}, err
|
||||||
|
}
|
||||||
|
if tag.RowsAffected() == 0 {
|
||||||
|
return domain.Signal{}, domain.ErrNotFound
|
||||||
|
}
|
||||||
|
return domain.Signal{}, nil
|
||||||
|
}
|
||||||
128
internal/core/app/observation.go
Normal file
128
internal/core/app/observation.go
Normal file
@@ -0,0 +1,128 @@
|
|||||||
|
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
|
||||||
|
}
|
||||||
123
internal/core/app/signals.go
Normal file
123
internal/core/app/signals.go
Normal file
@@ -0,0 +1,123 @@
|
|||||||
|
package app
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"sort"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/dtoro/oikos/internal/core/domain"
|
||||||
|
"github.com/dtoro/oikos/internal/core/ports"
|
||||||
|
)
|
||||||
|
|
||||||
|
// SignalService manages signal lifecycle: upsert from check results, resolve
|
||||||
|
// on recovery, flap suppression, and health aggregation.
|
||||||
|
type SignalService struct {
|
||||||
|
signals ports.SignalRepository
|
||||||
|
metrics ports.MetricsRepository
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewSignalService(signals ports.SignalRepository, metrics ports.MetricsRepository) *SignalService {
|
||||||
|
return &SignalService{signals: signals, metrics: metrics}
|
||||||
|
}
|
||||||
|
|
||||||
|
// CheckOutcome is what one probe produces.
|
||||||
|
type CheckOutcome struct {
|
||||||
|
EntityID domain.UUID
|
||||||
|
TargetID domain.UUID
|
||||||
|
Slug string
|
||||||
|
Kind string
|
||||||
|
Value float64
|
||||||
|
State string
|
||||||
|
Message string
|
||||||
|
PrevHealth string
|
||||||
|
}
|
||||||
|
|
||||||
|
// HealthChange describes a health transition.
|
||||||
|
type HealthChange struct {
|
||||||
|
TargetID domain.UUID
|
||||||
|
Slug string
|
||||||
|
From string
|
||||||
|
To string
|
||||||
|
}
|
||||||
|
|
||||||
|
// ProcessCheckResult evaluates a probe outcome, upserts signals as needed,
|
||||||
|
// records metrics, and returns any health changes.
|
||||||
|
func (s *SignalService) ProcessCheckResult(ctx context.Context, oc CheckOutcome) (*HealthChange, error) {
|
||||||
|
_ = s.metrics.InsertSamples(ctx, oc.TargetID, []ports.MetricSample{
|
||||||
|
{Metric: oc.Kind + "_value", Value: oc.Value, Timestamp: time.Now()},
|
||||||
|
})
|
||||||
|
|
||||||
|
switch oc.State {
|
||||||
|
case "ok":
|
||||||
|
_, _ = s.signals.Transition(ctx, ports.SignalTransitionInput{
|
||||||
|
SignalID: oc.EntityID, Action: "resolve", Actor: "scheduler",
|
||||||
|
})
|
||||||
|
case "critical", "warning":
|
||||||
|
severity := domain.SeverityWarning
|
||||||
|
if oc.State == "critical" {
|
||||||
|
severity = domain.SeverityCritical
|
||||||
|
}
|
||||||
|
_ = s.signals.UpsertWithTriggers(ctx, ports.SignalUpsertInput{
|
||||||
|
Signal: domain.Signal{
|
||||||
|
EntityID: oc.TargetID,
|
||||||
|
Kind: oc.Kind,
|
||||||
|
Severity: severity,
|
||||||
|
State: domain.SignalRaised,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
health := aggregateCheckStates(oc.State, oc.PrevHealth)
|
||||||
|
if health == oc.PrevHealth {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
return &HealthChange{TargetID: oc.TargetID, Slug: oc.Slug, From: oc.PrevHealth, To: health}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func aggregateCheckStates(probeState, prevHealth string) string {
|
||||||
|
switch probeState {
|
||||||
|
case "critical":
|
||||||
|
return "down"
|
||||||
|
case "warning":
|
||||||
|
return "degraded"
|
||||||
|
case "unknown":
|
||||||
|
return "stale"
|
||||||
|
default:
|
||||||
|
return "ok"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// WorstHealthForTarget computes entity health from all open signals.
|
||||||
|
func (s *SignalService) WorstHealthForTarget(ctx context.Context, targetID domain.UUID) (string, error) {
|
||||||
|
signals, err := s.signals.Open(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
if len(signals) == 0 {
|
||||||
|
return "ok", nil
|
||||||
|
}
|
||||||
|
sort.Slice(signals, func(i, j int) bool {
|
||||||
|
return severityRank(signals[i].Severity) > severityRank(signals[j].Severity)
|
||||||
|
})
|
||||||
|
switch signals[0].Severity {
|
||||||
|
case domain.SeverityCritical:
|
||||||
|
return "down", nil
|
||||||
|
case domain.SeverityWarning:
|
||||||
|
return "degraded", nil
|
||||||
|
default:
|
||||||
|
return "stale", nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func severityRank(s string) int {
|
||||||
|
switch s {
|
||||||
|
case domain.SeverityCritical:
|
||||||
|
return 3
|
||||||
|
case domain.SeverityWarning:
|
||||||
|
return 2
|
||||||
|
case domain.SeverityInfo:
|
||||||
|
return 1
|
||||||
|
default:
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user