Create schedules for new pipelines and delete schedules for removed pipelines.
()
| 85 | |
| 86 | |
| 87 | def update_pipeline_schedule(): |
| 88 | """Create schedules for new pipelines and delete schedules for removed pipelines.""" |
| 89 | |
| 90 | from vulnerabilities.importers import IMPORTERS_REGISTRY |
| 91 | from vulnerabilities.improvers import IMPROVERS_REGISTRY |
| 92 | from vulnerabilities.models import PipelineSchedule |
| 93 | from vulnerabilities.pipelines.exporters import EXPORTERS_REGISTRY |
| 94 | |
| 95 | pipelines = IMPORTERS_REGISTRY | IMPROVERS_REGISTRY | EXPORTERS_REGISTRY |
| 96 | |
| 97 | PipelineSchedule.objects.exclude(pipeline_id__in=pipelines.keys()).delete() |
| 98 | for id, pipeline_class in pipelines.items(): |
| 99 | run_once = getattr(pipeline_class, "run_once", False) |
| 100 | run_interval = getattr(pipeline_class, "run_interval", 1440) |
| 101 | run_priority = getattr( |
| 102 | pipeline_class, "run_priority", PipelineSchedule.ExecutionPriority.DEFAULT |
| 103 | ) |
| 104 | |
| 105 | pipeline, created = PipelineSchedule.objects.get_or_create( |
| 106 | pipeline_id=id, |
| 107 | defaults={ |
| 108 | "is_run_once": run_once, |
| 109 | "run_interval": run_interval, |
| 110 | "run_priority": run_priority, |
| 111 | }, |
| 112 | ) |
| 113 | |
| 114 | if not created: |
| 115 | if pipeline.run_priority != run_priority or pipeline.run_interval != run_interval: |
| 116 | pipeline.run_priority = run_priority |
| 117 | pipeline.run_interval = run_interval |
| 118 | pipeline.save() |