From 221f09c9ba4e92be7dfb45b01ec132d0e4aecfc5 Mon Sep 17 00:00:00 2001 From: Cyril Galibern Date: Fri, 25 Sep 2026 15:38:05 +0200 Subject: [PATCH 1/2] [scheduler] Make task outdated thresholds configurable Tasks that decide a value is outdated now read their age limit from scheduler.task..max_age instead of hardcoded SQL intervals. The defaults keep the previous values. The value accepts time.ParseDuration units plus "d" (15m, 25h, 2d, 1d12h) and must be >= 1m. All scheduler.task.*.max_age keys are validated when the scheduler starts. scrub_object no longer uses the v_outdated_services view. Its query is now inlined with a configurable age. The view is left in the database. Tunable tasks: scrub_object, scrub_resources, scrub_instances, scrub_unfinished_actions, scrub_resmon, scrub_checks_live, scrub_diskinfo, scrub_svcdisks, scrub_stor_array, scrub_node_hba, scrub_packages, scrub_patches, scrub_comp_status, scrub_comp_status_unattached, scrub_static, scrub_tempviz, scrub_pdf, alert_instances_not_updated, alert_nodes_not_updated, alert_service_config_not_updated, alert_checks_not_updated, log_instances_not_updated --- cdb/db_actions.go | 12 ++-- cdb/db_checks.go | 7 ++- cdb/db_compliance.go | 31 +++++----- cdb/db_dashboard.go | 28 ++++----- cdb/db_diskinfo.go | 7 ++- cdb/db_instances.go | 16 ++--- cdb/db_nodes.go | 6 +- cdb/db_object.go | 12 ++-- cdb/db_packages.go | 17 +++--- cdb/db_resources.go | 14 +++-- cdb/db_storage.go | 7 ++- cdb/db_svcdisks.go | 7 ++- cdb/maxage.go | 52 ++++++++++++++++ cmd/conf.go | 22 +++++++ cmd/scheduler.go | 3 + scheduler/task.go | 37 +++++++++++ scheduler/task_alerts.go | 30 +++++++-- scheduler/task_scrub.go | 129 ++++++++++++++++++++++++++++++--------- 18 files changed, 329 insertions(+), 108 deletions(-) create mode 100644 cdb/maxage.go diff --git a/cdb/db_actions.go b/cdb/db_actions.go index 0dc2c6c2..e52ca148 100644 --- a/cdb/db_actions.go +++ b/cdb/db_actions.go @@ -221,17 +221,17 @@ func (oDb *DB) GetBActionErrors(ctx context.Context) (lines []BActionErrorCount, return } -func (oDb *DB) UpdateUnfinishedActions(ctx context.Context) error { +func (oDb *DB) UpdateUnfinishedActions(ctx context.Context, maxAge time.Duration) error { request := `UPDATE svcactions SET status = "err", end = "1000-01-01 00:00:00" WHERE - begin < DATE_SUB(NOW(), INTERVAL 120 MINUTE) + begin < DATE_SUB(NOW(), INTERVAL ? SECOND) AND end IS NULL AND status IS NULL AND action NOT LIKE "%#%"` - if count, err := oDb.execCountContext(ctx, request); err != nil { + if count, err := oDb.execCountContext(ctx, request, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("svcactions") @@ -239,15 +239,15 @@ func (oDb *DB) UpdateUnfinishedActions(ctx context.Context) error { return nil } -func (oDb *DB) GetUnfinishedActions(ctx context.Context) (lines []SvcAction, err error) { +func (oDb *DB) GetUnfinishedActions(ctx context.Context, maxAge time.Duration) (lines []SvcAction, err error) { query := `SELECT id, node_id, svc_id FROM svcactions WHERE - begin < DATE_SUB(NOW(), INTERVAL 120 MINUTE) + begin < DATE_SUB(NOW(), INTERVAL ? SECOND) AND end IS NULL AND status IS NULL AND action NOT LIKE "%#%"` var rows *sql.Rows - rows, err = oDb.DB.QueryContext(ctx, query) + rows, err = oDb.DB.QueryContext(ctx, query, maxAgeSeconds(maxAge)) if err != nil { return } diff --git a/cdb/db_checks.go b/cdb/db_checks.go index 4e6ea0f1..aa9975a5 100644 --- a/cdb/db_checks.go +++ b/cdb/db_checks.go @@ -4,11 +4,12 @@ import ( "context" "fmt" "log/slog" + "time" ) -func (oDb *DB) PurgeChecksOutdated(ctx context.Context) error { - request := fmt.Sprintf("DELETE FROM `checks_live` WHERE `chk_updated` < DATE_SUB(NOW(), INTERVAL 2 DAY)") - if count, err := oDb.execCountContext(ctx, request); err != nil { +func (oDb *DB) PurgeChecksOutdated(ctx context.Context, maxAge time.Duration) error { + request := fmt.Sprintf("DELETE FROM `checks_live` WHERE `chk_updated` < DATE_SUB(NOW(), INTERVAL ? SECOND)") + if count, err := oDb.execCountContext(ctx, request, maxAgeSeconds(maxAge)); err != nil { return fmt.Errorf("delete from checks_live: %w", err) } else if count > 0 { // TODO: add metrics about purged count diff --git a/cdb/db_compliance.go b/cdb/db_compliance.go index 62b2fd40..214d2dd2 100644 --- a/cdb/db_compliance.go +++ b/cdb/db_compliance.go @@ -5,6 +5,7 @@ import ( "database/sql" "fmt" "strings" + "time" ) type Moduleset struct { @@ -82,12 +83,12 @@ func (oDb *DB) PurgeCompModulesetsServices(ctx context.Context) error { } // purge entries older than 30 days -func (oDb *DB) PurgeCompStatusOutdated(ctx context.Context) error { +func (oDb *DB) PurgeCompStatusOutdated(ctx context.Context, maxAge time.Duration) error { var query = `DELETE FROM comp_status WHERE - run_date < DATE_SUB(NOW(), INTERVAL 31 DAY)` - if count, err := oDb.execCountContext(ctx, query); err != nil { + run_date < DATE_SUB(NOW(), INTERVAL ? SECOND)` + if count, err := oDb.execCountContext(ctx, query, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("comp_status") @@ -127,15 +128,15 @@ func (oDb *DB) PurgeCompStatusNodeOrphans(ctx context.Context) error { return nil } -// purge compliance status older than 7 days for modules in no moduleset, ie not schedulable -func (oDb *DB) PurgeCompStatusModulesetOrphans(ctx context.Context) error { +// purge compliance status older than maxAge for modules in no moduleset, ie not schedulable +func (oDb *DB) PurgeCompStatusModulesetOrphans(ctx context.Context, maxAge time.Duration) error { var query = `DELETE FROM comp_status WHERE - run_date < DATE_SUB(NOW(), INTERVAL 7 DAY) AND + run_date < DATE_SUB(NOW(), INTERVAL ? SECOND) AND run_module NOT IN ( SELECT modset_mod_name FROM comp_moduleset_modules )` - if count, err := oDb.execCountContext(ctx, query); err != nil { + if count, err := oDb.execCountContext(ctx, query, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("comp_status") @@ -143,11 +144,11 @@ func (oDb *DB) PurgeCompStatusModulesetOrphans(ctx context.Context) error { return nil } -// purge node compliance status older than 7 days for unattached modules -func (oDb *DB) PurgeCompStatusNodeUnattached(ctx context.Context) error { +// purge node compliance status older than maxAge for unattached modules +func (oDb *DB) PurgeCompStatusNodeUnattached(ctx context.Context, maxAge time.Duration) error { var query = `DELETE FROM comp_status WHERE - run_date < DATE_SUB(NOW(), INTERVAL 7 DAY) AND + run_date < DATE_SUB(NOW(), INTERVAL ? SECOND) AND svc_id = "" AND run_module NOT IN ( SELECT modset_mod_name @@ -157,7 +158,7 @@ func (oDb *DB) PurgeCompStatusNodeUnattached(ctx context.Context) error { FROM comp_node_moduleset ) )` - if count, err := oDb.execCountContext(ctx, query); err != nil { + if count, err := oDb.execCountContext(ctx, query, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("comp_status") @@ -165,11 +166,11 @@ func (oDb *DB) PurgeCompStatusNodeUnattached(ctx context.Context) error { return nil } -// purge svc compliance status older than 7 days for unattached modules -func (oDb *DB) PurgeCompStatusSvcUnattached(ctx context.Context) error { +// purge svc compliance status older than maxAge for unattached modules +func (oDb *DB) PurgeCompStatusSvcUnattached(ctx context.Context, maxAge time.Duration) error { var query = `DELETE FROM comp_status WHERE - run_date < DATE_SUB(NOW(), INTERVAL 7 DAY) AND + run_date < DATE_SUB(NOW(), INTERVAL ? SECOND) AND svc_id = "" AND run_module NOT IN ( SELECT modset_mod_name @@ -179,7 +180,7 @@ func (oDb *DB) PurgeCompStatusSvcUnattached(ctx context.Context) error { FROM comp_modulesets_services ) )` - if count, err := oDb.execCountContext(ctx, query); err != nil { + if count, err := oDb.execCountContext(ctx, query, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("comp_status") diff --git a/cdb/db_dashboard.go b/cdb/db_dashboard.go index 5a924f00..c08a9c31 100644 --- a/cdb/db_dashboard.go +++ b/cdb/db_dashboard.go @@ -636,7 +636,7 @@ func (oDb *DB) PurgeAlertsOnDeletedServices(ctx context.Context) error { return nil } -func (oDb *DB) DashboardUpdateNodesNotUpdated(ctx context.Context) error { +func (oDb *DB) DashboardUpdateNodesNotUpdated(ctx context.Context, maxAge time.Duration) error { request := `INSERT INTO dashboard SELECT NULL, @@ -653,10 +653,10 @@ func (oDb *DB) DashboardUpdateNodesNotUpdated(ctx context.Context) error { NULL, NULL FROM nodes - WHERE updated < date_sub(NOW(), interval 25 hour) + WHERE updated < date_sub(NOW(), INTERVAL ? SECOND) ON DUPLICATE KEY UPDATE dash_updated=NOW()` - if count, err := oDb.execCountContext(ctx, request); err != nil { + if count, err := oDb.execCountContext(ctx, request, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("dashboard") @@ -664,7 +664,7 @@ func (oDb *DB) DashboardUpdateNodesNotUpdated(ctx context.Context) error { return nil } -func (oDb *DB) DashboardUpdateChecksNotUpdated(ctx context.Context) error { +func (oDb *DB) DashboardUpdateChecksNotUpdated(ctx context.Context, maxAge time.Duration) error { request := ` DELETE FROM dashboard WHERE @@ -730,10 +730,10 @@ func (oDb *DB) DashboardUpdateChecksNotUpdated(ctx context.Context) error { CONCAT(chk_type, ":", chk_instance) FROM checks_live c JOIN nodes n ON c.node_id = n.node_id - WHERE chk_updated < DATE_SUB(NOW(), INTERVAL 1 DAY) + WHERE chk_updated < DATE_SUB(NOW(), INTERVAL ? SECOND) ON DUPLICATE KEY UPDATE dash_updated = NOW(); ` - if count, err := oDb.execCountContext(ctx, request); err != nil { + if count, err := oDb.execCountContext(ctx, request, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("dashboard") @@ -793,7 +793,7 @@ func (oDb *DB) DashboardDeleteActionErrorsWithNoError(ctx context.Context) error return nil } -func (oDb *DB) DashboardUpdateServiceConfigNotUpdated(ctx context.Context) error { +func (oDb *DB) DashboardUpdateServiceConfigNotUpdated(ctx context.Context, maxAge time.Duration) error { request := ` INSERT INTO dashboard SELECT @@ -811,11 +811,11 @@ func (oDb *DB) DashboardUpdateServiceConfigNotUpdated(ctx context.Context) error NULL, NULL FROM services - WHERE updated < DATE_SUB(NOW(), INTERVAL 25 HOUR) + WHERE updated < DATE_SUB(NOW(), INTERVAL ? SECOND) ON DUPLICATE KEY UPDATE dash_updated=NOW() ` - if count, err := oDb.execCountContext(ctx, request); err != nil { + if count, err := oDb.execCountContext(ctx, request, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("dashboard") @@ -823,7 +823,7 @@ func (oDb *DB) DashboardUpdateServiceConfigNotUpdated(ctx context.Context) error return nil } -func (oDb *DB) DashboardUpdateInstancesNotUpdated(ctx context.Context) error { +func (oDb *DB) DashboardUpdateInstancesNotUpdated(ctx context.Context, maxAge time.Duration) error { request := ` INSERT INTO dashboard SELECT @@ -841,11 +841,11 @@ func (oDb *DB) DashboardUpdateInstancesNotUpdated(ctx context.Context) error { NULL, NULL FROM svcmon - WHERE mon_updated < DATE_SUB(NOW(), INTERVAL 16 MINUTE) + WHERE mon_updated < DATE_SUB(NOW(), INTERVAL ? SECOND) ON DUPLICATE KEY UPDATE dash_updated=NOW() ` - if count, err := oDb.execCountContext(ctx, request); err != nil { + if count, err := oDb.execCountContext(ctx, request, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("dashboard") @@ -863,10 +863,10 @@ func (oDb *DB) DashboardUpdateInstancesNotUpdated(ctx context.Context) error { dashboard.dash_type = "service status not updated" AND dashboard.svc_id != "" AND dashboard.node_id != "" AND - (svcmon.id IS NULL OR svcmon.mon_updated >= DATE_SUB(NOW(), INTERVAL 16 MINUTE)) + (svcmon.id IS NULL OR svcmon.mon_updated >= DATE_SUB(NOW(), INTERVAL ? SECOND)) ) ` - if count, err := oDb.execCountContext(ctx, request); err != nil { + if count, err := oDb.execCountContext(ctx, request, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("dashboard") diff --git a/cdb/db_diskinfo.go b/cdb/db_diskinfo.go index cfb0885c..03f5e328 100644 --- a/cdb/db_diskinfo.go +++ b/cdb/db_diskinfo.go @@ -6,6 +6,7 @@ import ( "errors" "fmt" "strings" + "time" ) type ( @@ -223,12 +224,12 @@ func (oDb *DB) DeleteDiskinfoByDiskID(ctx context.Context, diskID string) (int64 return count, nil } -func (oDb *DB) PurgeDiskinfoOutdated(ctx context.Context) error { +func (oDb *DB) PurgeDiskinfoOutdated(ctx context.Context, maxAge time.Duration) error { var query = `DELETE FROM diskinfo WHERE - disk_updated < DATE_SUB(NOW(), INTERVAL 2 DAY)` - if count, err := oDb.execCountContext(ctx, query); err != nil { + disk_updated < DATE_SUB(NOW(), INTERVAL ? SECOND)` + if count, err := oDb.execCountContext(ctx, query, maxAgeSeconds(maxAge)); err != nil { return fmt.Errorf("purge diskinfo: %w", err) } else if count > 0 { oDb.SetChange("diskinfo") diff --git a/cdb/db_instances.go b/cdb/db_instances.go index 24b43432..1510a416 100644 --- a/cdb/db_instances.go +++ b/cdb/db_instances.go @@ -750,12 +750,13 @@ func (oDb *DB) PurgeInstance(ctx context.Context, id InstanceID) error { return err } -func (oDb *DB) InstancesOutdated(ctx context.Context) (instanceIDs []InstanceID, err error) { +// InstancesOutdated returns the svcmon instance ids not updated since maxAge. +func (oDb *DB) InstancesOutdated(ctx context.Context, maxAge time.Duration) (instanceIDs []InstanceID, err error) { var rows *sql.Rows query := "SELECT `svc_id`, `node_id` " + "FROM `svcmon` " + - "WHERE `mon_updated` < DATE_SUB(NOW(), INTERVAL 21 MINUTE)" - rows, err = oDb.DB.QueryContext(ctx, query) + "WHERE `mon_updated` < DATE_SUB(NOW(), INTERVAL ? SECOND)" + rows, err = oDb.DB.QueryContext(ctx, query, maxAgeSeconds(maxAge)) if err != nil { return } @@ -771,14 +772,13 @@ func (oDb *DB) InstancesOutdated(ctx context.Context) (instanceIDs []InstanceID, return } -func (oDb *DB) LogInstancesNotUpdated(ctx context.Context) error { - age := 2 +func (oDb *DB) LogInstancesNotUpdated(ctx context.Context, maxAge time.Duration) error { request := fmt.Sprintf(`INSERT IGNORE INTO log SELECT NULL, "service.status", "scheduler", - "instance status not updated for more than %dh (%%(date)s)", + "instance status not updated for more than %s (%%(date)s)", CONCAT('{"date": "', mon_updated, '"}'), NOW(), svc_id, @@ -788,8 +788,8 @@ func (oDb *DB) LogInstancesNotUpdated(ctx context.Context) error { "warning", node_id from svcmon - where mon_updated 0 { slog.Debug(fmt.Sprintf("alert: instance outdated: %d", count)) diff --git a/cdb/db_nodes.go b/cdb/db_nodes.go index 33d51156..5b80d435 100644 --- a/cdb/db_nodes.go +++ b/cdb/db_nodes.go @@ -433,9 +433,9 @@ func (oDb *DB) NodeUpdateClusterIDForNodeID(ctx context.Context, nodeID, cluster } } -func (oDb *DB) PurgeNodeHBAsOutdated(ctx context.Context) error { - request := fmt.Sprintf("DELETE FROM `node_hba` WHERE `updated` < DATE_SUB(NOW(), INTERVAL 7 DAY)") - if count, err := oDb.execCountContext(ctx, request); err != nil { +func (oDb *DB) PurgeNodeHBAsOutdated(ctx context.Context, maxAge time.Duration) error { + request := fmt.Sprintf("DELETE FROM `node_hba` WHERE `updated` < DATE_SUB(NOW(), INTERVAL ? SECOND)") + if count, err := oDb.execCountContext(ctx, request, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { slog.Debug(fmt.Sprintf("purged %d entries from table node_hba", count)) diff --git a/cdb/db_object.go b/cdb/db_object.go index 0e9e1773..98bfdc40 100644 --- a/cdb/db_object.go +++ b/cdb/db_object.go @@ -696,13 +696,17 @@ func (oDb *DB) PurgeTablesFromObjectID(ctx context.Context, id string) error { } // ObjectsOutdated return lists of ids, svc_ids and svcnames for objects that no -// longer have instances updated in the last 15 minutes and that don't have their object +// longer have instances updated since maxAge and that don't have their object // status set to "undef" yet. -func (oDb *DB) ObjectsOutdated(ctx context.Context) (objects []ObjectMeta, err error) { +func (oDb *DB) ObjectsOutdated(ctx context.Context, maxAge time.Duration) (objects []ObjectMeta, err error) { sql := `SELECT id, svc_id, svcname FROM services - WHERE svc_id IN (SELECT svc_id FROM v_outdated_services WHERE uptodate=0) + WHERE svc_id IN ( + SELECT svc_id FROM svcmon + GROUP BY svc_id + HAVING SUM(mon_updated >= DATE_SUB(NOW(), INTERVAL ? SECOND)) = 0 + ) AND (svc_status != "undef" OR svc_availstatus != "undef")` - rows, err := oDb.DB.QueryContext(ctx, sql) + rows, err := oDb.DB.QueryContext(ctx, sql, maxAgeSeconds(maxAge)) if err != nil { return } diff --git a/cdb/db_packages.go b/cdb/db_packages.go index 3a72af90..b37210df 100644 --- a/cdb/db_packages.go +++ b/cdb/db_packages.go @@ -1,13 +1,16 @@ package cdb -import "context" +import ( + "context" + "time" +) -func (oDb *DB) PurgePackagesOutdated(ctx context.Context) error { +func (oDb *DB) PurgePackagesOutdated(ctx context.Context, maxAge time.Duration) error { var query = `DELETE FROM packages WHERE - pkg_updated < DATE_SUB(NOW(), INTERVAL 100 DAY)` - if count, err := oDb.execCountContext(ctx, query); err != nil { + pkg_updated < DATE_SUB(NOW(), INTERVAL ? SECOND)` + if count, err := oDb.execCountContext(ctx, query, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("packages") @@ -15,12 +18,12 @@ func (oDb *DB) PurgePackagesOutdated(ctx context.Context) error { return nil } -func (oDb *DB) PurgePatchesOutdated(ctx context.Context) error { +func (oDb *DB) PurgePatchesOutdated(ctx context.Context, maxAge time.Duration) error { var query = `DELETE FROM patches WHERE - patch_updated < DATE_SUB(NOW(), INTERVAL 100 DAY)` - if count, err := oDb.execCountContext(ctx, query); err != nil { + patch_updated < DATE_SUB(NOW(), INTERVAL ? SECOND)` + if count, err := oDb.execCountContext(ctx, query, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("patches") diff --git a/cdb/db_resources.go b/cdb/db_resources.go index d092aabb..f162fcfe 100644 --- a/cdb/db_resources.go +++ b/cdb/db_resources.go @@ -177,11 +177,13 @@ func (oDb *DB) ResmonRefreshTimestamp(ctx context.Context, nodeID string, object return } -func (oDb *DB) ResourceOutdatedLists(ctx context.Context) (resources []ResourceMeta, err error) { +// ResourceOutdatedLists returns the resmon entries not in "undef" status and +// not updated since maxAge. +func (oDb *DB) ResourceOutdatedLists(ctx context.Context, maxAge time.Duration) (resources []ResourceMeta, err error) { sql := `SELECT id, rid, svc_id, node_id FROM resmon - WHERE updated < DATE_SUB(NOW(), INTERVAL 15 MINUTE) + WHERE updated < DATE_SUB(NOW(), INTERVAL ? SECOND) AND res_status != "undef"` - rows, err := oDb.DB.QueryContext(ctx, sql) + rows, err := oDb.DB.QueryContext(ctx, sql, maxAgeSeconds(maxAge)) if err != nil { return } @@ -280,12 +282,12 @@ func (oDb *DB) ResourceUpdateStatus(ctx context.Context, resources []ResourceMet return n, err } -func (oDb *DB) PurgeResmonOutdated(ctx context.Context) error { +func (oDb *DB) PurgeResmonOutdated(ctx context.Context, maxAge time.Duration) error { var query = `DELETE FROM resmon WHERE - updated < DATE_SUB(NOW(), INTERVAL 1 DAY)` - if count, err := oDb.execCountContext(ctx, query); err != nil { + updated < DATE_SUB(NOW(), INTERVAL ? SECOND)` + if count, err := oDb.execCountContext(ctx, query, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("resmon") diff --git a/cdb/db_storage.go b/cdb/db_storage.go index ab56417b..703eda72 100644 --- a/cdb/db_storage.go +++ b/cdb/db_storage.go @@ -2,14 +2,15 @@ package cdb import ( "context" + "time" ) -func (oDb *DB) PurgeStorArrayOutdated(ctx context.Context) error { +func (oDb *DB) PurgeStorArrayOutdated(ctx context.Context, maxAge time.Duration) error { var query = `DELETE FROM stor_array WHERE array_model LIKE "vdisk%" AND - array_updated < DATE_SUB(NOW(), INTERVAL 2 DAY)` - if count, err := oDb.execCountContext(ctx, query); err != nil { + array_updated < DATE_SUB(NOW(), INTERVAL ? SECOND)` + if count, err := oDb.execCountContext(ctx, query, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("stor_array") diff --git a/cdb/db_svcdisks.go b/cdb/db_svcdisks.go index e706d9c3..4b73e1b0 100644 --- a/cdb/db_svcdisks.go +++ b/cdb/db_svcdisks.go @@ -2,14 +2,15 @@ package cdb import ( "context" + "time" ) -func (oDb *DB) PurgeSvcdisksOutdated(ctx context.Context) error { +func (oDb *DB) PurgeSvcdisksOutdated(ctx context.Context, maxAge time.Duration) error { var query = `DELETE FROM svcdisks WHERE - disk_updated < DATE_SUB(NOW(), INTERVAL 2 DAY)` - if count, err := oDb.execCountContext(ctx, query); err != nil { + disk_updated < DATE_SUB(NOW(), INTERVAL ? SECOND)` + if count, err := oDb.execCountContext(ctx, query, maxAgeSeconds(maxAge)); err != nil { return err } else if count > 0 { oDb.SetChange("svcdisks") diff --git a/cdb/maxage.go b/cdb/maxage.go new file mode 100644 index 00000000..1ae09385 --- /dev/null +++ b/cdb/maxage.go @@ -0,0 +1,52 @@ +package cdb + +import ( + "fmt" + "strconv" + "strings" + "time" +) + +// ParseMaxAge parses a duration string accepting the time.ParseDuration +// units plus a "d" day unit, like "15m", "25h", "2d" or "1d12h". +func ParseMaxAge(s string) (time.Duration, error) { + input := s + s = strings.TrimSpace(s) + var days time.Duration + if i := strings.Index(s, "d"); i >= 0 { + n, err := strconv.ParseUint(s[:i], 10, 32) + if err != nil { + return 0, fmt.Errorf("invalid duration %q", input) + } + days = time.Duration(n) * 24 * time.Hour + s = s[i+1:] + if s == "" { + return days, nil + } + } + d, err := time.ParseDuration(s) + if err != nil { + return 0, fmt.Errorf("invalid duration %q", input) + } + return days + d, nil +} + +// FormatMaxAge formats a duration for human readable messages, using the +// largest of the d, h, m, s units that divides it. +func FormatMaxAge(d time.Duration) string { + switch { + case d%(24*time.Hour) == 0: + return fmt.Sprintf("%dd", d/(24*time.Hour)) + case d%time.Hour == 0: + return fmt.Sprintf("%dh", d/time.Hour) + case d%time.Minute == 0: + return fmt.Sprintf("%dm", d/time.Minute) + default: + return d.String() + } +} + +// maxAgeSeconds returns the maxAge duration as a SQL "INTERVAL ? SECOND" argument. +func maxAgeSeconds(maxAge time.Duration) int64 { + return int64(maxAge / time.Second) +} diff --git a/cmd/conf.go b/cmd/conf.go index a75f2e50..dbaa7e74 100644 --- a/cmd/conf.go +++ b/cmd/conf.go @@ -76,6 +76,28 @@ func setDefaultSchedulerConfig() { viper.SetDefault(s+".metrics.enable", false) viper.SetDefault(s+".task.trim.retention", 365) viper.SetDefault(s+".task.trim.batch_size", 1000) + viper.SetDefault(s+".task.alert_checks_not_updated.max_age", "1d") + viper.SetDefault(s+".task.alert_instances_not_updated.max_age", "16m") + viper.SetDefault(s+".task.alert_nodes_not_updated.max_age", "25h") + viper.SetDefault(s+".task.alert_service_config_not_updated.max_age", "25h") + viper.SetDefault(s+".task.log_instances_not_updated.max_age", "2h") + viper.SetDefault(s+".task.scrub_checks_live.max_age", "2d") + viper.SetDefault(s+".task.scrub_comp_status.max_age", "31d") + viper.SetDefault(s+".task.scrub_comp_status_unattached.max_age", "7d") + viper.SetDefault(s+".task.scrub_diskinfo.max_age", "2d") + viper.SetDefault(s+".task.scrub_instances.max_age", "21m") + viper.SetDefault(s+".task.scrub_node_hba.max_age", "7d") + viper.SetDefault(s+".task.scrub_object.max_age", "15m") + viper.SetDefault(s+".task.scrub_packages.max_age", "100d") + viper.SetDefault(s+".task.scrub_patches.max_age", "100d") + viper.SetDefault(s+".task.scrub_pdf.max_age", "1d") + viper.SetDefault(s+".task.scrub_resmon.max_age", "1d") + viper.SetDefault(s+".task.scrub_resources.max_age", "15m") + viper.SetDefault(s+".task.scrub_static.max_age", "1h") + viper.SetDefault(s+".task.scrub_stor_array.max_age", "2d") + viper.SetDefault(s+".task.scrub_svcdisks.max_age", "2d") + viper.SetDefault(s+".task.scrub_tempviz.max_age", "1h") + viper.SetDefault(s+".task.scrub_unfinished_actions.max_age", "2h") viper.SetDefault(s+".log.request.level", "none") } diff --git a/cmd/scheduler.go b/cmd/scheduler.go index ee299a62..043b1f21 100644 --- a/cmd/scheduler.go +++ b/cmd/scheduler.go @@ -32,6 +32,9 @@ func newScheduler() (*schedulerT, error) { if err := setup(sectionScheduler); err != nil { return nil, err } + if err := scheduler.ValidateMaxAges(); err != nil { + return nil, err + } if db, err := newDatabase(); err != nil { return nil, err } else { diff --git a/scheduler/task.go b/scheduler/task.go index 47dd54b1..a06bf84c 100644 --- a/scheduler/task.go +++ b/scheduler/task.go @@ -8,11 +8,13 @@ import ( "io" "log/slog" "os" + "strings" "time" "github.com/go-redis/redis/v8" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promauto" + "github.com/spf13/viper" "github.com/opensvc/oc3/cdb" ) @@ -296,3 +298,38 @@ func (t *Task) SetLastRunAt(ctx context.Context) error { } return nil } + +// maxAge returns the scheduler.task..max_age setting: the age after +// which the task considers a value outdated. +func maxAge(name string) (time.Duration, error) { + return parseMaxAgeKey("scheduler.task." + name + ".max_age") +} + +// parseMaxAgeKey parses the key value with cdb.ParseMaxAge, accepting the +// time.ParseDuration units plus "d", and rejects values < 1m. +func parseMaxAgeKey(key string) (time.Duration, error) { + s := viper.GetString(key) + d, err := cdb.ParseMaxAge(s) + if err != nil { + return 0, fmt.Errorf("invalid %s value %q: %w", key, s, err) + } + if d < time.Minute { + return 0, fmt.Errorf("invalid %s value %q: must be >= 1m", key, s) + } + return d, nil +} + +// ValidateMaxAges verifies all the scheduler.task..max_age settings, +// so a bad value is reported at startup instead of when the task runs. +func ValidateMaxAges() error { + var errs []error + for _, key := range viper.AllKeys() { + if !strings.HasPrefix(key, "scheduler.task.") || !strings.HasSuffix(key, ".max_age") { + continue + } + if _, err := parseMaxAgeKey(key); err != nil { + errs = append(errs, err) + } + } + return errors.Join(errs...) +} diff --git a/scheduler/task_alerts.go b/scheduler/task_alerts.go index e5e27110..14d06aae 100644 --- a/scheduler/task_alerts.go +++ b/scheduler/task_alerts.go @@ -270,7 +270,11 @@ func taskLogInstancesNotUpdated(ctx context.Context, task *Task) error { return err } defer odb.Rollback() - if err := odb.LogInstancesNotUpdated(ctx); err != nil { + age, err := maxAge("log_instances_not_updated") + if err != nil { + return err + } + if err := odb.LogInstancesNotUpdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -388,7 +392,11 @@ func taskAlertNodesNotUpdated(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.DashboardUpdateNodesNotUpdated(ctx); err != nil { + age, err := maxAge("alert_nodes_not_updated") + if err != nil { + return err + } + if err := odb.DashboardUpdateNodesNotUpdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -420,7 +428,11 @@ func taskAlertServiceConfigNotUpdated(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.DashboardUpdateServiceConfigNotUpdated(ctx); err != nil { + age, err := maxAge("alert_service_config_not_updated") + if err != nil { + return err + } + if err := odb.DashboardUpdateServiceConfigNotUpdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -436,7 +448,11 @@ func taskAlertInstancesNotUpdated(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.DashboardUpdateInstancesNotUpdated(ctx); err != nil { + age, err := maxAge("alert_instances_not_updated") + if err != nil { + return err + } + if err := odb.DashboardUpdateInstancesNotUpdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -500,7 +516,11 @@ func taskAlertChecksNotUpdated(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.DashboardUpdateChecksNotUpdated(ctx); err != nil { + age, err := maxAge("alert_checks_not_updated") + if err != nil { + return err + } + if err := odb.DashboardUpdateChecksNotUpdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { diff --git a/scheduler/task_scrub.go b/scheduler/task_scrub.go index ad5a7078..ad2a70e5 100644 --- a/scheduler/task_scrub.go +++ b/scheduler/task_scrub.go @@ -8,39 +8,44 @@ import ( "path/filepath" "time" - "github.com/opensvc/oc3/cdb" "github.com/spf13/viper" + + "github.com/opensvc/oc3/cdb" ) -// TaskScrubObjects marks services status "undef" if all instances have outdated or absent data. +// TaskScrubObjects marks services status "undef" if all instances have outdated data. // -// For testing, force a scrubable dataset with: +// For testing, force a scrubable dataset with (15 MINUTE being the default +// scheduler.task.scrub_object.max_age): // -// UPDATE services SET svc_status="up" WHERE svc_id IN (SELECT svc_id FROM v_outdated_services); +// UPDATE services SET svc_status="up" WHERE svc_id IN ( +// SELECT svc_id FROM svcmon GROUP BY svc_id +// HAVING SUM(mon_updated >= DATE_SUB(NOW(), INTERVAL 15 MINUTE)) = 0 +// ); var TaskScrubObjects = Task{ name: "scrub_object", - desc: "marks services status=undef if all instances have outdated (aged 15m) or absent data", + desc: "marks services status=undef if all instances have outdated (aged scheduler.task.scrub_object.max_age, default 15m) data", fn: taskScrubObjects, timeout: time.Minute, } var TaskScrubUnfinishedActions = Task{ name: "scrub_unfinished_actions", - desc: "set a end date and status=err on actions not finished after 2h running", + desc: "set a end date and status=err on actions not finished after scheduler.task.scrub_unfinished_actions.max_age (default 2h) running", fn: taskScrubUnfinishedActions, timeout: time.Minute, } var TaskScrubResources = Task{ name: "scrub_resources", - desc: "marks status=undef outdated (aged 15m) resources", + desc: "marks status=undef outdated (aged scheduler.task.scrub_resources.max_age, default 15m) resources", fn: taskScrubResources, timeout: time.Minute, } var TaskScrubInstances = Task{ name: "scrub_instances", - desc: "marks status=undef outdated (aged 15m) instances", + desc: "purges outdated (aged scheduler.task.scrub_instances.max_age, default 21m) instances", fn: taskScrubInstances, timeout: time.Minute, } @@ -206,7 +211,11 @@ func taskScrubInstances(ctx context.Context, task *Task) error { return err } defer odb.Rollback() - instanceIDs, err := odb.InstancesOutdated(ctx) + age, err := maxAge("scrub_instances") + if err != nil { + return err + } + instanceIDs, err := odb.InstancesOutdated(ctx, age) if err != nil { return err } @@ -230,7 +239,11 @@ func taskScrubResources(ctx context.Context, task *Task) error { defer odb.Rollback() // Fetch the outdated resources still not in "undef" availstatus - resources, err := odb.ResourceOutdatedLists(ctx) + age, err := maxAge("scrub_resources") + if err != nil { + return err + } + resources, err := odb.ResourceOutdatedLists(ctx, age) if err != nil { return err } @@ -294,7 +307,11 @@ func taskScrubObjects(ctx context.Context, task *Task) error { defer odb.Rollback() // Fetch the outdated services still not in "undef" availstatus - objects, err := odb.ObjectsOutdated(ctx) + age, err := maxAge("scrub_object") + if err != nil { + return err + } + objects, err := odb.ObjectsOutdated(ctx, age) if err != nil { return err } @@ -353,7 +370,11 @@ func taskScrubChecksLive(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.PurgeChecksOutdated(ctx); err != nil { + age, err := maxAge("scrub_checks_live") + if err != nil { + return err + } + if err := odb.PurgeChecksOutdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -369,7 +390,11 @@ func taskScrubNodeHBA(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.PurgeNodeHBAsOutdated(ctx); err != nil { + age, err := maxAge("scrub_node_hba") + if err != nil { + return err + } + if err := odb.PurgeNodeHBAsOutdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -385,7 +410,11 @@ func taskScrubPackages(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.PurgePackagesOutdated(ctx); err != nil { + age, err := maxAge("scrub_packages") + if err != nil { + return err + } + if err := odb.PurgePackagesOutdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -401,7 +430,11 @@ func taskScrubPatches(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.PurgePatchesOutdated(ctx); err != nil { + age, err := maxAge("scrub_patches") + if err != nil { + return err + } + if err := odb.PurgePatchesOutdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -417,7 +450,11 @@ func taskScrubResmon(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.PurgeResmonOutdated(ctx); err != nil { + age, err := maxAge("scrub_resmon") + if err != nil { + return err + } + if err := odb.PurgeResmonOutdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -433,7 +470,11 @@ func taskScrubDiskinfo(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.PurgeDiskinfoOutdated(ctx); err != nil { + age, err := maxAge("scrub_diskinfo") + if err != nil { + return err + } + if err := odb.PurgeDiskinfoOutdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -449,7 +490,11 @@ func taskScrubSvcdisks(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.PurgeSvcdisksOutdated(ctx); err != nil { + age, err := maxAge("scrub_svcdisks") + if err != nil { + return err + } + if err := odb.PurgeSvcdisksOutdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -465,7 +510,11 @@ func taskScrubStorArray(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.PurgeStorArrayOutdated(ctx); err != nil { + age, err := maxAge("scrub_stor_array") + if err != nil { + return err + } + if err := odb.PurgeStorArrayOutdated(ctx, age); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -545,7 +594,15 @@ func taskScrubCompStatus(ctx context.Context, task *Task) error { } defer odb.Rollback() - if err := odb.PurgeCompStatusOutdated(ctx); err != nil { + age, err := maxAge("scrub_comp_status") + if err != nil { + return err + } + unattachedAge, err := maxAge("scrub_comp_status_unattached") + if err != nil { + return err + } + if err := odb.PurgeCompStatusOutdated(ctx, age); err != nil { return err } if err := odb.PurgeCompStatusSvcOrphans(ctx); err != nil { @@ -554,13 +611,13 @@ func taskScrubCompStatus(ctx context.Context, task *Task) error { if err := odb.PurgeCompStatusNodeOrphans(ctx); err != nil { return err } - if err := odb.PurgeCompStatusModulesetOrphans(ctx); err != nil { + if err := odb.PurgeCompStatusModulesetOrphans(ctx, unattachedAge); err != nil { return err } - if err := odb.PurgeCompStatusNodeUnattached(ctx); err != nil { + if err := odb.PurgeCompStatusNodeUnattached(ctx, unattachedAge); err != nil { return err } - if err := odb.PurgeCompStatusSvcUnattached(ctx); err != nil { + if err := odb.PurgeCompStatusSvcUnattached(ctx, unattachedAge); err != nil { return err } if err := odb.Session.NotifyChanges(ctx); err != nil { @@ -609,7 +666,11 @@ func scrubFiles(pattern string, threshold time.Time) error { } func taskScrubStatic(ctx context.Context, task *Task) error { - threshold := time.Now().Add(-1 * time.Hour) + age, err := maxAge("scrub_static") + if err != nil { + return err + } + threshold := time.Now().Add(-age) directory := viper.GetString("scheduler.directories.static") if directory == "" { slog.Warn("skip: define scheduler.directories.static") @@ -638,7 +699,11 @@ func taskScrubStatic(ctx context.Context, task *Task) error { } func taskScrubTempviz(ctx context.Context, task *Task) error { - threshold := time.Now().Add(-1 * time.Hour) + age, err := maxAge("scrub_tempviz") + if err != nil { + return err + } + threshold := time.Now().Add(-age) directory := viper.GetString("scheduler.directories.static") if directory == "" { slog.Warn("skip: define scheduler.directories.static") @@ -648,7 +713,11 @@ func taskScrubTempviz(ctx context.Context, task *Task) error { } func taskScrubPdf(ctx context.Context, task *Task) error { - threshold := time.Now().Add(-24 * time.Hour) + age, err := maxAge("scrub_pdf") + if err != nil { + return err + } + threshold := time.Now().Add(-age) directory := viper.GetString("scheduler.directories.static") if directory == "" { slog.Warn("skip: define scheduler.directories.static") @@ -664,7 +733,11 @@ func taskScrubUnfinishedActions(ctx context.Context, task *Task) error { } defer odb.Rollback() - lines, err := odb.GetUnfinishedActions(ctx) + age, err := maxAge("scrub_unfinished_actions") + if err != nil { + return err + } + lines, err := odb.GetUnfinishedActions(ctx, age) if err != nil { return fmt.Errorf("get: %w", err) } @@ -689,7 +762,7 @@ func taskScrubUnfinishedActions(ctx context.Context, task *Task) error { if err := odb.Log(ctx, entries...); err != nil { return fmt.Errorf("log: %w", err) } - if err := odb.UpdateUnfinishedActions(ctx); err != nil { + if err := odb.UpdateUnfinishedActions(ctx, age); err != nil { return fmt.Errorf("update: %w", err) } if err := odb.Session.NotifyChanges(ctx); err != nil { From 7754f6f2b35bcd6f472344f704af1ceab647f003 Mon Sep 17 00:00:00 2001 From: Cyril Galibern Date: Fri, 25 Sep 2026 17:12:34 +0200 Subject: [PATCH 2/2] [worker] Fix resmon not refreshed on daemon ping InstancePingFromNodeID returned early when no svcmon row needed a refresh, so resmon.updated and resmon_log_last.res_end were skipped. When svcmon was kept fresh by the daemon status feed, resmon.updated aged, and scrub_resources flagged the resources of live instances as undef. The svcmon and resmon refreshes now run independently. --- cdb/db_instances.go | 23 +++++++++++++---------- 1 file changed, 13 insertions(+), 10 deletions(-) diff --git a/cdb/db_instances.go b/cdb/db_instances.go index 1510a416..50b6081f 100644 --- a/cdb/db_instances.go +++ b/cdb/db_instances.go @@ -232,9 +232,9 @@ func (oDb *DB) SvcmonRefreshTimestamp(ctx context.Context, nodeID string, object return } -// InstancePingFromNodeID updates match svcmon.mon_updated, svcmon_log_last.mon_end, -// resmon.updated and resmon_log_last.res_end when svcmon.mon_updated timestamp -// for node_id id older than 30s. +// InstancePingFromNodeID refreshes svcmon.mon_updated, svcmon_log_last.mon_end, +// resmon.updated and resmon_log_last.res_end of the node_id rows, when older +// than 30s. func (oDb *DB) InstancePingFromNodeID(ctx context.Context, nodeID string) (updates bool, err error) { defer logDuration("instancePing "+nodeID, time.Now()) const ( @@ -254,13 +254,16 @@ func (oDb *DB) InstancePingFromNodeID(ctx context.Context, nodeID string) (updat count int64 ) + // svcmon and resmon are refreshed independently: svcmon.mon_updated may + // have been refreshed by another path (daemon status) without refreshing + // all resmon rows, so skipping resmon when svcmon has no update would let + // resmon.updated age until the resources are scrubbed. if count, err = oDb.execCountContext(ctx, qUpdateSvcmon, nodeID); err != nil { return - } else if count == 0 { - return + } else if count > 0 { + updates = true + oDb.SetChange("svcmon") } - updates = true - oDb.SetChange("svcmon") if _, err = oDb.ExecContext(ctx, qUpdateSvcmonLogLast, nodeID); err != nil { return @@ -268,10 +271,10 @@ func (oDb *DB) InstancePingFromNodeID(ctx context.Context, nodeID string) (updat if count, err = oDb.execCountContext(ctx, qUpdateResmon, nodeID); err != nil { return - } else if count == 0 { - return + } else if count > 0 { + updates = true + oDb.SetChange("resmon") } - oDb.SetChange("resmon") _, err = oDb.ExecContext(ctx, qUpdateResmonLogLast, nodeID) return