Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 11 additions & 5 deletions cmd/app/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ type config struct {
commandAPIKeys string // comma-separated; empty leaves the command endpoint unregistered
commandListenAddr string
enableOpenITCockpitTweaks bool // selects the core-restart hoststatus/servicestatus cleanup query
storeNotificationStart bool // store NEBTYPE_CONTACTNOTIFICATIONMETHOD_START instead of END, for distributed notifications
statusMaxAge string // max age of a hoststatus/servicestatus event before it is discarded; "0" disables
logLevel string // "debug", "info", "warn" or "error"
logFormat string // "text" or "json"
Expand All @@ -69,10 +70,10 @@ type config struct {
// fileConfig mirrors config's fields for -config's optional YAML file (see
// config.example.yaml for every key, its default and a description). Every
// key is optional: a zero value (empty string, nil for APIKeys/
// EnableOpenITCockpitTweaks) means "not set in the file", so it never
// overrides an environment variable or hardcoded default - see resolveString/
// resolveBool. EnableOpenITCockpitTweaks is a *bool (rather than bool) for
// exactly this reason: unlike a missing string, Go can't otherwise tell
// EnableOpenITCockpitTweaks/StoreNotificationStart) means "not set in the
// file", so it never overrides an environment variable or hardcoded default -
// see resolveString/resolveBool. The two bool keys are *bool (rather than
// bool) for exactly this reason: unlike a missing string, Go can't otherwise tell
// "the file didn't mention this key" apart from "the file explicitly set
// it to false".
type fileConfig struct {
Expand Down Expand Up @@ -103,6 +104,7 @@ type fileConfig struct {
CommandAPIKeys []string `yaml:"command_api_keys"`
CommandListenAddr string `yaml:"command_listen_addr"`
EnableOpenITCockpitTweaks *bool `yaml:"enable_openitcockpit_tweaks"`
StoreNotificationStart *bool `yaml:"store_notification_start"`
StatusMaxAge string `yaml:"status_max_age"`
LogLevel string `yaml:"log_level"`
LogFormat string `yaml:"log_format"`
Expand Down Expand Up @@ -282,6 +284,9 @@ func loadConfig() config {
flag.BoolVar(&cfg.enableOpenITCockpitTweaks, "enable-openitcockpit-tweaks", false,
"on a core restart, delete only hoststatus/servicestatus rows for objects openITCockpit no longer "+
"knows about instead of truncating both tables outright")
flag.BoolVar(&cfg.storeNotificationStart, "store-notification-start", false,
"store the START event of a notification method instead of its END; for a core whose notifications "+
"a broker module such as mod_gearman distributes, which never brokers the END event")
flag.StringVar(&cfg.statusMaxAge, "status-max-age", "5m",
"discard statusngin_hoststatus/statusngin_servicestatus events older than this Go duration (e.g. \"5m\", \"90s\"); "+
"they are superseded snapshots, so a backlog of them is not worth draining after downtime. \"0\" processes every event regardless of age")
Expand Down Expand Up @@ -331,6 +336,7 @@ func loadConfig() config {
cfg.commandAPIKeys = resolveString(explicit, "command-api-keys", cfg.commandAPIKeys, "STATUSENGINE_API_COMMAND_KEYS", strings.Join(fc.CommandAPIKeys, ","))
cfg.commandListenAddr = resolveString(explicit, "command-listen-addr", cfg.commandListenAddr, "STATUSENGINE_COMMAND_LISTEN_ADDR", fc.CommandListenAddr)
cfg.enableOpenITCockpitTweaks = resolveBool(explicit, "enable-openitcockpit-tweaks", cfg.enableOpenITCockpitTweaks, "ENABLE_OPENITCOCKPIT_TWEAKS", fc.EnableOpenITCockpitTweaks)
cfg.storeNotificationStart = resolveBool(explicit, "store-notification-start", cfg.storeNotificationStart, "STATUSENGINE_STORE_NOTIFICATION_START", fc.StoreNotificationStart)
cfg.statusMaxAge = resolveString(explicit, "status-max-age", cfg.statusMaxAge, "STATUSENGINE_STATUS_MAX_AGE", fc.StatusMaxAge)
cfg.logLevel = resolveString(explicit, "log-level", cfg.logLevel, "STATUSENGINE_LOG_LEVEL", fc.LogLevel)
cfg.logFormat = resolveString(explicit, "log-format", cfg.logFormat, "STATUSENGINE_LOG_FORMAT", fc.LogFormat)
Expand Down Expand Up @@ -790,7 +796,7 @@ func main() {
// connection is ever dialed (CLAUDE.md rule 5).
gc := graphite.NewClient(cfg.graphiteAddr, graphite.WithMaxBatchSize(cfg.graphiteBatchSize))

router, runners := queue.NewRouter(sqlDB, hub, gc, perfdataRoute, cfg.graphitePrefix, cfg.nodeName, cfg.enableOpenITCockpitTweaks, statusMaxAge, cfg.mysqlBatchSize)
router, runners := queue.NewRouter(sqlDB, hub, gc, perfdataRoute, cfg.graphitePrefix, cfg.nodeName, cfg.enableOpenITCockpitTweaks, statusMaxAge, cfg.mysqlBatchSize, cfg.storeNotificationStart)
for _, r := range runners {
wg.Add(1)
go func(r queue.Runner) {
Expand Down
2 changes: 1 addition & 1 deletion cmd/simulator/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -158,7 +158,7 @@ func main() {
// statusMaxAge 0: the simulator shifts every fixture timestamp to keep
// primary keys unique (see withUniqueTimestamps), so ages here are
// synthetic and an age filter would only make its output unpredictable.
router, runners := queue.NewRouter(sqlDB, hub, gc, queue.PerfdataRouteMySQL, "statusengine-simulator", "statusengine-simulator", false, 0, db.DefaultMaxBatchSize)
router, runners := queue.NewRouter(sqlDB, hub, gc, queue.PerfdataRouteMySQL, "statusengine-simulator", "statusengine-simulator", false, 0, db.DefaultMaxBatchSize, false)
for _, r := range runners {
wg.Add(1)
go func(r queue.Runner) {
Expand Down
11 changes: 11 additions & 0 deletions config.example.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,17 @@ command_listen_addr: 127.0.0.1:8081
# exist survives the restart instead of needing to be rebuilt.
enable_openitcockpit_tweaks: false

# Which notification method event to store in statusengine_host_notifications
# and statusengine_service_notifications?
# false (default): the END event, as the core brokers it once a notification
# command has run.
# true: the START event. For a core whose notifications a broker module such
# as mod_gearman distributes: it answers the START event with
# NEBERROR_CALLBACKOVERRIDE, the core then never runs the command itself
# and never brokers the END event, so with the default nothing is ever
# stored. end_time is set to start_time, since the end is not known.
store_notification_start: false

# Discard statusngin_hoststatus and statusngin_servicestatus events older
# than this, instead of writing them to MySQL and broadcasting them.
#
Expand Down
2 changes: 1 addition & 1 deletion internal/queue/batchsize_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ func TestBatchSizeStaysUnderPlaceholderLimit(t *testing.T) {
// the default of 100 and would sail past a table that has grown too wide.
hub := websocket.NewHub()
_, runners := NewRouter(mockDB, hub, graphite.NewClient("127.0.0.1:2003"), PerfdataRouteMySQL,
"statusengine-test", "statusengine-test", false, noAgeFilter, db.MaxConfigurableBatchSize)
"statusengine-test", "statusengine-test", false, noAgeFilter, db.MaxConfigurableBatchSize, false)

var checked, widest int
for _, r := range runners {
Expand Down
2 changes: 1 addition & 1 deletion internal/queue/metrics_init_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,7 @@ func TestNewRouterPreCreatesMetricSeries(t *testing.T) {
hub := websocket.NewHub()

router, _ := NewRouter(sqlDB, hub, graphite.NewClient("127.0.0.1:2003"),
PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize)
PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize, false)

// Every queue in the router, on all four per-queue metrics.
//
Expand Down
67 changes: 65 additions & 2 deletions internal/queue/notification_method_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ func TestContactNotificationMethodHandlerDiscardsNonEndType(t *testing.T) {

hostIns := &fakeEnqueuer[notificationMethodEvent]{}
serviceIns := &fakeEnqueuer[notificationMethodEvent]{}
handler := newContactNotificationMethodHandler(hub, QueueContactNotificationMethod, hostIns, serviceIns)
handler := newContactNotificationMethodHandler(hub, QueueContactNotificationMethod, hostIns, serviceIns, false)

// The real fixture carries type 605 (NEBTYPE_CONTACTNOTIFICATIONMETHOD_END)
// and a service_description, so it must land in serviceIns.
Expand Down Expand Up @@ -60,7 +60,7 @@ func TestContactNotificationMethodHandlerRoutesHostVsService(t *testing.T) {

hostIns := &fakeEnqueuer[notificationMethodEvent]{}
serviceIns := &fakeEnqueuer[notificationMethodEvent]{}
handler := newContactNotificationMethodHandler(hub, QueueContactNotificationMethod, hostIns, serviceIns)
handler := newContactNotificationMethodHandler(hub, QueueContactNotificationMethod, hostIns, serviceIns, false)

hostEvent := []byte(`{
"type": 605,
Expand Down Expand Up @@ -114,3 +114,66 @@ func TestHostNotificationRowAndServiceNotificationRowColumns(t *testing.T) {
t.Fatalf("serviceNotificationRow[3] (hostname) = %v, want %v", serviceRow[3], ev.HostName)
}
}

func TestContactNotificationMethodHandlerStoresStartEventWhenAsked(t *testing.T) {
hub := websocket.NewHub()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go hub.Run(ctx)

hostIns := &fakeEnqueuer[notificationMethodEvent]{}
serviceIns := &fakeEnqueuer[notificationMethodEvent]{}
handler := newContactNotificationMethodHandler(hub, QueueContactNotificationMethod, hostIns, serviceIns, true)

// What Naemon brokers when mod_gearman distributes notifications: the
// START event, without an end time, and no END event after it.
start := []byte(`{
"type": 604,
"timestamp": 1785517089,
"timestamp_usec": 927284,
"contactnotificationmethod": {
"host_name": "localhost",
"service_description": "Swap Usage",
"contact_name": "someone",
"start_time": 1785517089,
"end_time": 0
}
}`)
if err := handler(ctx, start); err != nil {
t.Fatalf("handler: %v", err)
}
got := serviceIns.snapshot()
if len(got) != 1 {
t.Fatalf("serviceIns got %d items, want 1", len(got))
}
if got[0].StartTime != 1785517089 || got[0].EndTime != got[0].StartTime {
t.Fatalf("start_time/end_time = %d/%d, want 1785517089 for both", got[0].StartTime, got[0].EndTime)
}
if got[0].ContactName != "someone" {
t.Fatalf("contact_name = %q, want the one from the START event", got[0].ContactName)
}

// The END event is not stored as well: where it does arrive, keeping both
// would record every notification twice.
end := []byte(`{
"type": 605,
"timestamp": 1785517090,
"timestamp_usec": 11,
"contactnotificationmethod": {
"host_name": "localhost",
"service_description": "Swap Usage",
"contact_name": "someone",
"start_time": 1785517089,
"end_time": 1785517090
}
}`)
if err := handler(ctx, end); err != nil {
t.Fatalf("handler: %v", err)
}
if got := len(serviceIns.snapshot()); got != 1 {
t.Fatalf("serviceIns got %d items after the END event, want still 1", got)
}
if got := len(hostIns.snapshot()); got != 0 {
t.Fatalf("hostIns got %d items, want 0", got)
}
}
4 changes: 2 additions & 2 deletions internal/queue/redelivery_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -213,7 +213,7 @@ func TestNewRouterEmitsUpsertForCheckTables(t *testing.T) {

hub := websocket.NewHub()
router, runners := NewRouter(mockDB, hub, graphite.NewClient("127.0.0.1:2003"),
PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize)
PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize, false)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()
Expand Down Expand Up @@ -312,7 +312,7 @@ func TestUpsertTablesHaveNoSecondaryUniqueIndex(t *testing.T) {

hub := websocket.NewHub()
_, runners := NewRouter(mockDB, hub, graphite.NewClient("127.0.0.1:2003"),
PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize)
PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize, false)

tables := make([]string, 0, 16)
for _, r := range runners {
Expand Down
50 changes: 39 additions & 11 deletions internal/queue/registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -245,13 +245,22 @@ func serviceAcknowledgementRow(ev acknowledgementEvent, dst []any) []any {
}

// notificationTypeContactNotificationMethodEnd is the Nagios/Icinga/Naemon
// NEBTYPE_CONTACTNOTIFICATIONMETHOD_END event type: the only
// NEBTYPE_CONTACTNOTIFICATIONMETHOD_END event type: the
// contactnotificationmethod event that represents a completed notification
// method delivery, and therefore the only one persisted to
// method delivery, and by default the only one persisted to
// statusengine_host_notifications/statusengine_service_notifications. Every
// other type value on this queue is discarded immediately.
const notificationTypeContactNotificationMethodEnd = 605

// notificationTypeContactNotificationMethodStart is
// NEBTYPE_CONTACTNOTIFICATIONMETHOD_START, persisted instead of the END event
// when storeNotificationStart is set. A broker module that distributes
// notifications, such as mod_gearman, answers this event with
// NEBERROR_CALLBACKOVERRIDE. Naemon then continues with the next notification
// command without running this one, so it never brokers the END event, and
// the START is all there is to store. It carries everything but the end time.
const notificationTypeContactNotificationMethodStart = 604

func hostNotificationRow(ev notificationMethodEvent, dst []any) []any {
return append(dst,
ev.HostName, ev.StartTime, ev.TimestampUsec, ev.ContactName, ev.CommandName, ev.CommandArgs,
Expand All @@ -267,23 +276,42 @@ func serviceNotificationRow(ev notificationMethodEvent, dst []any) []any {
}

// newContactNotificationMethodHandler filters out every event whose type
// isn't notificationTypeContactNotificationMethodEnd, then routes the rest
// to hostIns or serviceIns depending on whether service_description is set
// - mirroring newStateChangeHandler/newAcknowledgementHandler's host-vs-
// service split.
func newContactNotificationMethodHandler(hub *websocket.Hub, topic string, hostIns, serviceIns enqueuer[notificationMethodEvent]) Handler {
// isn't the one this worker stores - notificationTypeContactNotificationMethodEnd,
// or notificationTypeContactNotificationMethodStart with storeNotificationStart
// - then routes the rest to hostIns or serviceIns depending on whether
// service_description is set, mirroring newStateChangeHandler/
// newAcknowledgementHandler's host-vs-service split.
func newContactNotificationMethodHandler(hub *websocket.Hub, topic string, hostIns, serviceIns enqueuer[notificationMethodEvent], storeNotificationStart bool) Handler {
stored := notificationTypeContactNotificationMethodEnd
if storeNotificationStart {
stored = notificationTypeContactNotificationMethodStart
}

return func(ctx context.Context, payload []byte) error {
events, err := decodeContactNotificationMethod(payload)
if err != nil {
return decodeError(topic, err)
}

if storeNotificationStart {
// A START event has no end: the module that took the notification
// over does the sending and never reports back. end_time is NOT
// NULL with no sub-second part, so a zero would read as 1970. The
// start time makes the duration exactly zero, which reads as the
// placeholder it is rather than as a measurement.
for i := range events {
if events[i].Type == notificationTypeContactNotificationMethodStart {
events[i].EndTime = events[i].StartTime
}
}
}

publishFiltered(hub, topic, events, func(ev notificationMethodEvent) bool {
return ev.Type == notificationTypeContactNotificationMethodEnd
return ev.Type == stored
})

for _, ev := range events {
if ev.Type != notificationTypeContactNotificationMethodEnd {
if ev.Type != stored {
continue
}

Expand Down Expand Up @@ -900,7 +928,7 @@ var redeliverySafePKColumn = map[string]string{
"statusengine_service_notifications_log": "hostname",
}

func NewRouter(sqlDB *sql.DB, hub *websocket.Hub, gc *graphite.Client, perfdataRoute PerfdataRoute, graphitePrefix, nodeName string, enableOpenITCockpitTweaks bool, statusMaxAge time.Duration, mysqlBatchSize int) (Router, []Runner) {
func NewRouter(sqlDB *sql.DB, hub *websocket.Hub, gc *graphite.Client, perfdataRoute PerfdataRoute, graphitePrefix, nodeName string, enableOpenITCockpitTweaks bool, statusMaxAge time.Duration, mysqlBatchSize int, storeNotificationStart bool) (Router, []Runner) {
// Every table shares one batch size, built once here and spread into
// each constructor below, so a table added later cannot quietly keep
// the default. db.WithMaxBatchSize clamps; cmd/app is what rejects an
Expand Down Expand Up @@ -994,7 +1022,7 @@ func NewRouter(sqlDB *sql.DB, hub *websocket.Hub, gc *graphite.Client, perfdataR
// NewStaleDroppingHandler for why that is safe here and nowhere else.
QueueHostStatus: NewStaleDroppingHandler(hub, QueueHostStatus, hostStatus, decodeHostStatus, statusMaxAge),
QueueServiceStatus: NewStaleDroppingHandler(hub, QueueServiceStatus, serviceStatus, decodeServiceStatus, statusMaxAge),
QueueContactNotificationMethod: newContactNotificationMethodHandler(hub, QueueContactNotificationMethod, hostNotifications, serviceNotifications),
QueueContactNotificationMethod: newContactNotificationMethodHandler(hub, QueueContactNotificationMethod, hostNotifications, serviceNotifications, storeNotificationStart),
QueueNotifications: newNotificationHandler(hub, QueueNotifications, hostNotificationsLog, serviceNotificationsLog),

QueueDowntimes: newDowntimeHandler(hub, QueueDowntimes, sqlDB, nodeName),
Expand Down
6 changes: 3 additions & 3 deletions internal/queue/registry_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ const testBatchSize = db.DefaultMaxBatchSize
func TestNewRouterCoversAllQueues(t *testing.T) {
sqlDB := openTestDB(t)
hub := websocket.NewHub()
router, runners := NewRouter(sqlDB, hub, graphite.NewClient("127.0.0.1:2003"), PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize)
router, runners := NewRouter(sqlDB, hub, graphite.NewClient("127.0.0.1:2003"), PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize, false)

want := []string{
QueueHostStatus, QueueServiceStatus, QueueHostChecks, QueueServiceChecks,
Expand Down Expand Up @@ -98,7 +98,7 @@ func runAllAndFlush(t *testing.T, runners []Runner) context.Context {
func TestHostCheckHandlerPersistsToMySQL(t *testing.T) {
sqlDB := openTestDB(t)
hub := websocket.NewHub()
router, runners := NewRouter(sqlDB, hub, graphite.NewClient("127.0.0.1:2003"), PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize)
router, runners := NewRouter(sqlDB, hub, graphite.NewClient("127.0.0.1:2003"), PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize, false)
ctx := runAllAndFlush(t, runners)
go hub.Run(ctx)

Expand Down Expand Up @@ -160,7 +160,7 @@ func TestHostCheckHandlerPersistsToMySQL(t *testing.T) {
func TestAcknowledgementHandlerRoutesToHostAndServiceTables(t *testing.T) {
sqlDB := openTestDB(t)
hub := websocket.NewHub()
router, runners := NewRouter(sqlDB, hub, graphite.NewClient("127.0.0.1:2003"), PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize)
router, runners := NewRouter(sqlDB, hub, graphite.NewClient("127.0.0.1:2003"), PerfdataRouteMySQL, "statusengine-test", "statusengine-test", false, noAgeFilter, testBatchSize, false)
ctx := runAllAndFlush(t, runners)
go hub.Run(ctx)

Expand Down
2 changes: 1 addition & 1 deletion internal/queue/stale_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -258,7 +258,7 @@ func TestOnlyStatusQueuesDiscardOnAge(t *testing.T) {
go hub.Run(ctx)

router, _ := NewRouter(sqlDB, hub, graphite.NewClient("127.0.0.1:2003"), PerfdataRouteMySQL,
"statusengine-test", "statusengine-test", false, 5*time.Minute, testBatchSize)
"statusengine-test", "statusengine-test", false, 5*time.Minute, testBatchSize, false)

// NewRouter must have pre-created both series at zero, so a dashboard
// panel reads 0 rather than "No data" on a worker that has not
Expand Down
Loading