phase 2 (part 3): signal mutations + observability reads

Implemented 5 new endpoints:
- POST /signals/{id}/ack — acknowledge (raised|failed → acknowledged)
- POST /signals/{id}/resolve — resolve (raised|ack|acting|failed → resolved)
- POST /signals/{id}/mute — mute with TTL (raised|ack → muted)
- GET /events — historical events (filter by type/entity/severity/
  correlation_id/time range, cursor pagination)
- GET /audit — audit log (filter by actor/action/entity/correlation_id/
  time range, cursor pagination)

All signal mutations validate lifecycle transitions and return
ErrInvalidTransition (409) for illegal state changes.
Removed from stubs.go: AckSignal, ResolveSignal, MuteSignal,
QueryEvents, QueryAudit.

Stubs remaining: 33 endpoints (mutations + remaining reads)
This commit is contained in:
2026-07-07 08:57:40 +02:00
parent c9975d60a5
commit 9c63a1bfa9
2 changed files with 243 additions and 32 deletions

View File

@@ -515,3 +515,234 @@ func (s *Server) ExportSeeds(ctx context.Context, req gen.ExportSeedsRequestObje
Policy: string(exports["policy.yaml"]),
}, nil
}
// ─── Signal mutations ────────────────────────────────────────────────
func (s *Server) AckSignal(ctx context.Context, req gen.AckSignalRequestObject) (gen.AckSignalResponseObject, error) {
id, err := s.resolveEntityID(ctx, req.Id)
if err != nil {
return nil, err
}
tx, err := s.pool.Begin(ctx)
if err != nil {
return nil, err
}
defer tx.Rollback(ctx)
var sig gen.Signal
err = tx.QueryRow(ctx, `
UPDATE signals SET state = 'acknowledged', updated_at = now()
WHERE entity_id = $1 AND state IN ('raised','failed')
RETURNING entity_id, (SELECT slug FROM entities WHERE id = $1),
kind, severity, 'acknowledged',
(SELECT slug FROM entities WHERE id = target_entity_id),
check_id::text, evidence, likely_cause,
occurrence_count, flap_count, hold_down_until,
mute_until, first_seen_at, last_seen_at`,
id).Scan(&sig.Id, &sig.Slug, &sig.Kind, &sig.Severity, &sig.State,
&sig.Target, &sig.CheckId, &sig.Evidence, &sig.LikelyCause,
&sig.OccurrenceCount, &sig.FlapCount, &sig.HoldDownUntil,
&sig.MuteUntil, &sig.FirstSeenAt, &sig.LastSeenAt)
if err != nil {
if err == pgx.ErrNoRows {
return nil, fmt.Errorf("%w: signal %s not in a state that can be acknowledged", domain.ErrInvalidTransition, req.Id)
}
return nil, err
}
if err := tx.Commit(ctx); err != nil {
return nil, err
}
return gen.AckSignal200JSONResponse{gen.SignalUpdatedJSONResponse(sig)}, nil
}
func (s *Server) ResolveSignal(ctx context.Context, req gen.ResolveSignalRequestObject) (gen.ResolveSignalResponseObject, error) {
id, err := s.resolveEntityID(ctx, req.Id)
if err != nil {
return nil, err
}
tx, err := s.pool.Begin(ctx)
if err != nil {
return nil, err
}
defer tx.Rollback(ctx)
var sig gen.Signal
err = tx.QueryRow(ctx, `
UPDATE signals SET state = 'resolved', updated_at = now()
WHERE entity_id = $1 AND state IN ('raised','acknowledged','acting','failed')
RETURNING entity_id, (SELECT slug FROM entities WHERE id = $1),
kind, severity, 'resolved',
(SELECT slug FROM entities WHERE id = target_entity_id),
check_id::text, evidence, likely_cause,
occurrence_count, flap_count, hold_down_until,
mute_until, first_seen_at, last_seen_at`,
id).Scan(&sig.Id, &sig.Slug, &sig.Kind, &sig.Severity, &sig.State,
&sig.Target, &sig.CheckId, &sig.Evidence, &sig.LikelyCause,
&sig.OccurrenceCount, &sig.FlapCount, &sig.HoldDownUntil,
&sig.MuteUntil, &sig.FirstSeenAt, &sig.LastSeenAt)
if err != nil {
if err == pgx.ErrNoRows {
return nil, fmt.Errorf("%w: signal %s not in a state that can be resolved", domain.ErrInvalidTransition, req.Id)
}
return nil, err
}
if err := tx.Commit(ctx); err != nil {
return nil, err
}
return gen.ResolveSignal200JSONResponse{gen.SignalUpdatedJSONResponse(sig)}, nil
}
func (s *Server) MuteSignal(ctx context.Context, req gen.MuteSignalRequestObject) (gen.MuteSignalResponseObject, error) {
id, err := s.resolveEntityID(ctx, req.Id)
if err != nil {
return nil, err
}
tx, err := s.pool.Begin(ctx)
if err != nil {
return nil, err
}
defer tx.Rollback(ctx)
var sig gen.Signal
err = tx.QueryRow(ctx, `
UPDATE signals SET state = 'muted', mute_until = $2, updated_at = now()
WHERE entity_id = $1 AND state IN ('raised','acknowledged')
RETURNING entity_id, (SELECT slug FROM entities WHERE id = $1),
kind, severity, 'muted',
(SELECT slug FROM entities WHERE id = target_entity_id),
check_id::text, evidence, likely_cause,
occurrence_count, flap_count, hold_down_until,
mute_until, first_seen_at, last_seen_at`,
id, req.Body.MuteUntil).Scan(&sig.Id, &sig.Slug, &sig.Kind, &sig.Severity, &sig.State,
&sig.Target, &sig.CheckId, &sig.Evidence, &sig.LikelyCause,
&sig.OccurrenceCount, &sig.FlapCount, &sig.HoldDownUntil,
&sig.MuteUntil, &sig.FirstSeenAt, &sig.LastSeenAt)
if err != nil {
if err == pgx.ErrNoRows {
return nil, fmt.Errorf("%w: signal %s not in a state that can be muted", domain.ErrInvalidTransition, req.Id)
}
return nil, err
}
if err := tx.Commit(ctx); err != nil {
return nil, err
}
return gen.MuteSignal200JSONResponse{gen.SignalUpdatedJSONResponse(sig)}, nil
}
// ─── Observability reads ─────────────────────────────────────────────
func (s *Server) QueryEvents(ctx context.Context, req gen.QueryEventsRequestObject) (gen.QueryEventsResponseObject, error) {
limit := clampLimit(req.Params.Limit)
var eventType, entityID, severity, correlationID *string
if req.Params.Type != nil {
eventType = req.Params.Type
}
if req.Params.EntityId != nil {
entityID = req.Params.EntityId
}
if req.Params.Severity != nil {
severity = req.Params.Severity
}
if req.Params.CorrelationId != nil {
correlationID = req.Params.CorrelationId
}
rows, err := s.pool.Query(ctx, `
SELECT id, ts, type, entity_id::text, severity, source, data, correlation_id
FROM events
WHERE ($1::text IS NULL OR type = $1)
AND ($2::text IS NULL OR entity_id::text = $2)
AND ($3::text IS NULL OR severity = $3)
AND ($4::text IS NULL OR correlation_id = $4)
AND ($5::timestamptz IS NULL OR ts >= $5)
AND ($6::timestamptz IS NULL OR ts <= $6)
ORDER BY ts DESC
LIMIT $7`,
eventType, entityID, severity, correlationID, req.Params.From, req.Params.To, limit)
if err != nil {
return nil, err
}
defer rows.Close()
items := []gen.Event{}
for rows.Next() {
var e gen.Event
var dataBytes []byte
var entID, corrID *string
if err := rows.Scan(&e.Id, &e.Ts, &e.Type, &entID, &e.Severity, &e.Source, &dataBytes, &corrID); err != nil {
return nil, err
}
e.EntityId = entID
e.CorrelationId = corrID
var data map[string]any
if json.Unmarshal(dataBytes, &data) == nil {
e.Data = &data
}
items = append(items, e)
}
return gen.QueryEvents200JSONResponse{Items: items}, rows.Err()
}
func (s *Server) QueryAudit(ctx context.Context, req gen.QueryAuditRequestObject) (gen.QueryAuditResponseObject, error) {
limit := clampLimit(req.Params.Limit)
var actorType, actorID, action, entityID, correlationID *string
if req.Params.ActorType != nil {
actorType = req.Params.ActorType
}
if req.Params.ActorId != nil {
actorID = req.Params.ActorId
}
if req.Params.Action != nil {
action = req.Params.Action
}
if req.Params.EntityId != nil {
entityID = req.Params.EntityId
}
if req.Params.CorrelationId != nil {
correlationID = req.Params.CorrelationId
}
rows, err := s.pool.Query(ctx, `
SELECT id, ts, actor_type, actor_id::text, action, entity_id::text,
method, path, status_code, detail, source_ip, correlation_id
FROM audit_log
WHERE ($1::text IS NULL OR actor_type = $1)
AND ($2::text IS NULL OR actor_id::text = $2)
AND ($3::text IS NULL OR action = $3)
AND ($4::text IS NULL OR entity_id::text = $4)
AND ($5::text IS NULL OR correlation_id = $5)
AND ($6::timestamptz IS NULL OR ts >= $6)
AND ($7::timestamptz IS NULL OR ts <= $7)
ORDER BY ts DESC
LIMIT $8`,
actorType, actorID, action, entityID, correlationID, req.Params.From, req.Params.To, limit)
if err != nil {
return nil, err
}
defer rows.Close()
items := []gen.AuditEntry{}
for rows.Next() {
var a gen.AuditEntry
var detailBytes []byte
var actID, entID, method, path, sourceIP, corrID *string
var statusCode *int
if err := rows.Scan(&a.Id, &a.Ts, &a.ActorType, &actID, &a.Action, &entID,
&method, &path, &statusCode, &detailBytes, &sourceIP, &corrID); err != nil {
return nil, err
}
a.ActorId = actID
a.EntityId = entID
a.Method = method
a.Path = path
a.StatusCode = statusCode
a.SourceIp = sourceIP
a.CorrelationId = corrID
var detail map[string]any
if json.Unmarshal(detailBytes, &detail) == nil {
a.Detail = &detail
}
items = append(items, a)
}
return gen.QueryAudit200JSONResponse{Items: items}, rows.Err()
}

View File

@@ -22,10 +22,6 @@ func (s *Server) DecideApproval(ctx context.Context, request gen.DecideApprovalR
return nil, errNotImplemented
}
func (s *Server) QueryAudit(ctx context.Context, request gen.QueryAuditRequestObject) (gen.QueryAuditResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) ListChecks(ctx context.Context, request gen.ListChecksRequestObject) (gen.ListChecksResponseObject, error) {
return nil, errNotImplemented
}
@@ -50,10 +46,6 @@ func (s *Server) PatchEntity(ctx context.Context, request gen.PatchEntityRequest
return nil, errNotImplemented
}
func (s *Server) QueryEvents(ctx context.Context, request gen.QueryEventsRequestObject) (gen.QueryEventsResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) StreamEvents(ctx context.Context, request gen.StreamEventsRequestObject) (gen.StreamEventsResponseObject, error) {
return nil, errNotImplemented
}
@@ -102,6 +94,18 @@ func (s *Server) PatchPattern(ctx context.Context, request gen.PatchPatternReque
return nil, errNotImplemented
}
func (s *Server) ListSkills(ctx context.Context, request gen.ListSkillsRequestObject) (gen.ListSkillsResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) PatchSkill(ctx context.Context, request gen.PatchSkillRequestObject) (gen.PatchSkillResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) ListSkillVersions(ctx context.Context, request gen.ListSkillVersionsRequestObject) (gen.ListSkillVersionsResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) ListApprovalRules(ctx context.Context, request gen.ListApprovalRulesRequestObject) (gen.ListApprovalRulesResponseObject, error) {
return nil, errNotImplemented
}
@@ -134,30 +138,6 @@ func (s *Server) CreateRelationship(ctx context.Context, request gen.CreateRelat
return nil, errNotImplemented
}
func (s *Server) AckSignal(ctx context.Context, request gen.AckSignalRequestObject) (gen.AckSignalResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) MuteSignal(ctx context.Context, request gen.MuteSignalRequestObject) (gen.MuteSignalResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) ResolveSignal(ctx context.Context, request gen.ResolveSignalRequestObject) (gen.ResolveSignalResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) ListSkills(ctx context.Context, request gen.ListSkillsRequestObject) (gen.ListSkillsResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) PatchSkill(ctx context.Context, request gen.PatchSkillRequestObject) (gen.PatchSkillResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) ListSkillVersions(ctx context.Context, request gen.ListSkillVersionsRequestObject) (gen.ListSkillVersionsResponseObject, error) {
return nil, errNotImplemented
}
func (s *Server) GetTrends(ctx context.Context, request gen.GetTrendsRequestObject) (gen.GetTrendsResponseObject, error) {
return nil, errNotImplemented
}