From f6fc2ccaeabaf06f6fc11ffdef30d15701eec46b Mon Sep 17 00:00:00 2001 From: ponmeloco Date: Fri, 18 Sep 2026 08:39:57 +0200 Subject: [PATCH] Store the notification START event where the END never arrives A broker module that distributes notifications, such as mod_gearman, answers NEBTYPE_CONTACTNOTIFICATIONMETHOD_START with NEBERROR_CALLBACKOVERRIDE. Naemon then continues with the next notification command without running this one, so it never brokers NEBTYPE_CONTACTNOTIFICATIONMETHOD_END (notifications.c: the override at the method level is a `continue`, and the END call comes after the command). This worker stores only the END event, so in that setup statusengine_host_notifications and statusengine_service_notifications stay empty for good while the notifications themselves go out. The START event does arrive and carries everything but the end time: state, output, contact, command and start time. store_notification_start, or STATUSENGINE_STORE_NOTIFICATION_START, or -store-notification-start, makes the worker store it instead. It defaults to off, so a core that brokers the END event is stored exactly as before, and with it on the END event is dropped so no notification is recorded twice. end_time is filled from start_time. The column is NOT NULL with no sub-second part, so a zero would be rendered as 1970, and a duration of exactly zero reads as the placeholder it is rather than as a measurement. The same is done before the event is published on the WebSocket, so both see one value. This is the same option as in the PHP worker. Measured with it on in an openITCOCKPIT Docker stack, which distributes notifications by default: 3489 service and 2325 host notifications stored, where there had been none. NewRouter takes the flag as a new last parameter, as statusMaxAge and mysqlBatchSize were added. The failing tests are the six that need a Gearman or RabbitMQ broker on localhost; the same six fail on main. --- cmd/app/main.go | 16 ++++-- cmd/simulator/main.go | 2 +- config.example.yaml | 11 ++++ internal/queue/batchsize_test.go | 2 +- internal/queue/metrics_init_test.go | 2 +- internal/queue/notification_method_test.go | 67 +++++++++++++++++++++- internal/queue/redelivery_test.go | 4 +- internal/queue/registry.go | 50 ++++++++++++---- internal/queue/registry_test.go | 6 +- internal/queue/stale_test.go | 2 +- 10 files changed, 135 insertions(+), 27 deletions(-) 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