diff --git a/cmd/app/main.go b/cmd/app/main.go index 22188c6..804a9f8 100644 --- a/cmd/app/main.go +++ b/cmd/app/main.go @@ -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" @@ -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 { @@ -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"` @@ -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") @@ -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) @@ -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) { diff --git a/cmd/simulator/main.go b/cmd/simulator/main.go index 3d84410..c815943 100644 --- a/cmd/simulator/main.go +++ b/cmd/simulator/main.go @@ -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) { diff --git a/config.example.yaml b/config.example.yaml index 0ff3bba..fd03744 100644 --- a/config.example.yaml +++ b/config.example.yaml @@ -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. # diff --git a/internal/queue/batchsize_test.go b/internal/queue/batchsize_test.go index 1090286..2ba4ba5 100644 --- a/internal/queue/batchsize_test.go +++ b/internal/queue/batchsize_test.go @@ -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 { diff --git a/internal/queue/metrics_init_test.go b/internal/queue/metrics_init_test.go index 5bf039b..dc050a6 100644 --- a/internal/queue/metrics_init_test.go +++ b/internal/queue/metrics_init_test.go @@ -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. // diff --git a/internal/queue/notification_method_test.go b/internal/queue/notification_method_test.go index f54a870..0e3666c 100644 --- a/internal/queue/notification_method_test.go +++ b/internal/queue/notification_method_test.go @@ -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. @@ -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, @@ -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) + } +} diff --git a/internal/queue/redelivery_test.go b/internal/queue/redelivery_test.go index b0c2525..987cdea 100644 --- a/internal/queue/redelivery_test.go +++ b/internal/queue/redelivery_test.go @@ -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() @@ -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 { diff --git a/internal/queue/registry.go b/internal/queue/registry.go index 4375fd7..f8df8b6 100644 --- a/internal/queue/registry.go +++ b/internal/queue/registry.go @@ -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, @@ -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 } @@ -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 @@ -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), diff --git a/internal/queue/registry_test.go b/internal/queue/registry_test.go index 725e3fd..584f779 100644 --- a/internal/queue/registry_test.go +++ b/internal/queue/registry_test.go @@ -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, @@ -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) @@ -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) diff --git a/internal/queue/stale_test.go b/internal/queue/stale_test.go index c1a9d14..3b79edf 100644 --- a/internal/queue/stale_test.go +++ b/internal/queue/stale_test.go @@ -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