From 8028c2df765a3cbecde073985e66c1a6baff308f Mon Sep 17 00:00:00 2001 From: Queue A Date: Thu, 27 Aug 2026 08:55:08 +0200 Subject: [PATCH] refactored to work with django tasks as a backend --- asyncron/admin.py | 114 ++++--- asyncron/backend.py | 162 +++++++++ .../commands/run_asyncron_worker.py | 4 +- asyncron/migrations/0001_initial.py | 89 ++--- asyncron/migrations/0002_task_self_aware.py | 18 - asyncron/migrations/0002_trace_worker.py | 19 ++ .../migrations/0003_alter_task_self_aware.py | 18 - .../0003_alter_trace_scheduled_datetime.py | 20 ++ .../migrations/0004_alter_task_gracetime.py | 19 -- ...t_crowning_attempt_worker_last_activity.py | 18 - ...asyncron_me_model_t_d92186_idx_and_more.py | 20 -- asyncron/models.py | 323 +++++++++--------- asyncron/shortcuts.py | 53 +-- asyncron/utils.py | 25 +- asyncron/workers.py | 314 ++++++++--------- 15 files changed, 663 insertions(+), 553 deletions(-) create mode 100644 asyncron/backend.py delete mode 100644 asyncron/migrations/0002_task_self_aware.py create mode 100644 asyncron/migrations/0002_trace_worker.py delete mode 100644 asyncron/migrations/0003_alter_task_self_aware.py create mode 100644 asyncron/migrations/0003_alter_trace_scheduled_datetime.py delete mode 100644 asyncron/migrations/0004_alter_task_gracetime.py delete mode 100644 asyncron/migrations/0005_rename_last_crowning_attempt_worker_last_activity.py delete mode 100644 asyncron/migrations/0006_remove_metadata_asyncron_me_model_t_d92186_idx_and_more.py diff --git a/asyncron/admin.py b/asyncron/admin.py index 9f918e4..d216226 100644 --- a/asyncron/admin.py +++ b/asyncron/admin.py @@ -1,9 +1,11 @@ from django.contrib import admin from django.utils import timezone from django.db.models import F +from django.tasks.base import TaskResultStatus from .base.admin import BaseModelAdmin -from .models import Worker, Task, Trace +from .models import Worker, TaskSchedule, Trace +from .utils import is_proc_alive import os import asyncio @@ -11,18 +13,18 @@ import humanize @admin.register( Worker ) class WorkerAdmin( BaseModelAdmin ): - order = 4 - list_display = 'pid', 'thread_id', 'is_robust', 'is_master', 'is_running', 'health', + order = 1 + list_display = 'process_id', 'thread_id', 'is_standalone', 'is_scheduler', 'is_running', # 'health', def has_add_permission( self, request, obj = None ): return False - def is_running( self, obj ): return obj.is_proc_alive() + def is_running( self, obj ): return is_proc_alive(obj.process_id) is_running.boolean = True - def health( self, obj ): - return (f"In Grace " if obj.in_grace else "") + humanize.naturaltime( obj.last_activity ) + #def health( self, obj ): + # return (f"In Grace " if obj.in_grace else "") + humanize.naturaltime( obj.last_activity ) -class TaskAppFilter( admin.SimpleListFilter ): #TODO: Also check if it's an actuall app, same for trace +class TaskAppFilter( admin.SimpleListFilter ): title = "app" parameter_name = 'app_groups' @@ -30,33 +32,36 @@ class TaskAppFilter( admin.SimpleListFilter ): #TODO: Also check if it's an actu return ( ( a.lower(), a ) for a in sorted({ - task.name.split(".", 1)[0] - for task in Task.objects.all() + ts.task_path.split(".", 1)[0] + for ts in TaskSchedule.objects.all() }) if a ) def queryset( self, request, queryset ): q = self.value() if not q: return queryset - return queryset.filter( name__istartswith = q ) + return queryset.filter( task_path__istartswith = q ) - -@admin.register( Task ) -class TaskAdmin( BaseModelAdmin ): - order = 1 - list_display = 'name', 'timeout', 'gracetime', 'jitter', 'type', 'worker_type', 'logged', 'last_execution', 'scheduled' +@admin.register( TaskSchedule ) +class TaskScheduleAdmin( BaseModelAdmin ): + order = 2 + list_display = 'schedule_name', 'timeout', 'type', 'jitter', 'logged', 'last_execution', 'scheduled', 'is_enabled', list_filter = TaskAppFilter, - fields = ["name", "description", "type", "on_model_change", "jitter", "self_aware"] - actions = 'schedule_execution', 'execution_now', 'delete_script_missing', 'update_details', - def has_add_permission( self, request, obj = None ): return False - def has_delete_permission( self, request, obj = None ): return False - def has_change_permission( self, request, obj = None ): return False + fields = "name", "description", "type", #"on_model_change", "jitter", + + + def has_add_permission( self, request, obj = None ): return True + #def has_delete_permission( self, request, obj = None ): return not obj or obj.name != "scripted" + def has_change_permission( self, request, obj = None ): return not obj or obj.name != "scripted" + + def schedule_name( self, obj ): return f"{obj.task_path} ({obj.name})" def description( self, obj ): try: - return obj.registered_tasks[obj.name].__doc__.strip("\n") + return obj.task.func.__doc__.strip("\n") except: + raise return "N/A" def jitter( self, obj ): @@ -69,38 +74,44 @@ class TaskAdmin( BaseModelAdmin ): def type( self, obj ): results = [] - if obj.timeout is None: - results.append( "Service" ) - if obj.interval: delta = humanize.naturaldelta( obj.interval ) delta = delta.replace("an ", "1 ").replace("a ", "1 ") - results.append( f"Periodic, every {delta}" ) + results.append( f"Every {delta}" ) - if obj.name not in obj.registered_tasks: - results.append("Script Missing!") + else: + results.append( "Service" ) - return ", ".join( results ) if results else "Callable" + #if obj.name not in obj.registered_tasks: + # results.append("Script Missing!") + + return ", ".join( results ) def on_model_change( self, obj ): try: return ", ".join( f"{m.__module__}.{m.__name__}" for m in obj.registered_tasks[obj.name].watching_models ) except: return "N/A" def logged( self, obj ): - count = obj.trace_set.exclude( status = "S" ).count() + count = obj.trace_set.exclude( status = TaskResultStatus.READY ).count() if count == 1: return "1 trace" return f"{count} traces" def last_execution( self, obj ): - last_trace = obj.trace_set.exclude( status = "S" ).exclude( last_run_datetime = None ).order_by('last_run_datetime').last() + last_trace = obj.trace_set.exclude( status = TaskResultStatus.READY ).exclude( finished_datetime = None ).order_by('finished_datetime').last() if not last_trace: return "Never" - return "Ongoing..." if last_trace.status == "R" else humanize.naturaltime( last_trace.last_run_datetime ) + return "Ongoing..." if last_trace.status == TaskResultStatus.RUNNING else humanize.naturaltime( last_trace.finished_datetime ) def scheduled( self, obj ): - return obj.trace_set.filter( status = "S" ).exists() + return obj.trace_set.filter( status = TaskResultStatus.READY ).exists() scheduled.boolean = True + + actions = 'toggle_is_enabled', # 'schedule_execution', 'execution_now', 'delete_script_missing', + @admin.action( description = "Toggle Enable/Disable" ) + def toggle_is_enabled( self, request, qs ): + qs.update( is_enabled = ~F('is_enabled') ) + @admin.action( description = "(Re)Schedule an execution for periodic tasks" ) def schedule_execution( self, request, qs ): trace_ids = set() @@ -132,12 +143,8 @@ class TaskAdmin( BaseModelAdmin ): if task.name not in task.registered_tasks: task.delete() - @admin.action( description = "Check for updates in details of this tasks" ) - def update_details( self, request, qs ): - for task in qs: - if task.name not in task.registered_tasks: continue - task_f = task.registered_tasks[task.name] - print( "T:", task, task_f ) + + class TraceAppFilter(admin.SimpleListFilter): title = "app" @@ -178,33 +185,32 @@ class TraceNameFilter(admin.SimpleListFilter): @admin.register( Trace ) class TraceAdmin( BaseModelAdmin ): - order = 2 - list_display = 'task', 'execution', 'state', 'worker_lock' - list_filter = TraceAppFilter, TraceNameFilter, 'task__worker_type', 'status', 'status_reason', - ordering = F('scheduled_datetime').desc(nulls_last=True), + order = 3 + list_display = 'scheduled_datetime', 'task_path', 'execution', 'status', 'worker' + list_filter = 'status', 'status_description', #'task__worker_type', TraceAppFilter, TraceNameFilter, + ordering = "-enqueued_datetime", #readonly_fields = [ f.name for f in Trace._meta.fields ] - def has_add_permission( self, request, obj = None ): return False - def has_change_permission( self, request, obj = None ): return False + #def has_add_permission( self, request, obj = None ): return False + #def has_change_permission( self, request, obj = None ): return False def execution( self, obj ): - if obj.last_run_datetime: - return "- Ran " + humanize.naturaltime( obj.last_run_datetime ) + if obj.finished_datetime: + return "- Ran " + humanize.naturaltime( obj.finished_datetime ) - if obj.scheduled_datetime: - if obj.scheduled_datetime < timezone.now(): - return "- Should've run " + humanize.naturaltime( obj.scheduled_datetime ) - else: - return "+ In " + humanize.naturaltime( obj.scheduled_datetime ) - - return "Never" + if obj.scheduled_datetime < timezone.now(): + return "- Should've run " + humanize.naturaltime( obj.scheduled_datetime ) + else: + return "+ In " + humanize.naturaltime( obj.scheduled_datetime ) execution.admin_order_field = 'scheduled_datetime' def state( self, obj ): - return f"{obj.status}: {obj.status_reason}" if obj.status_reason else f"{obj.get_status_display()}" + if obj.status_description: return f"{obj.status}: {obj.status_description}" + return obj.get_status_display() state.admin_order_field = 'status' + actions = 'reschedule_to_now', @admin.action( description = "Reschedule to run now" ) def reschedule_to_now( self, request, qs ): diff --git a/asyncron/backend.py b/asyncron/backend.py new file mode 100644 index 0000000..3cd5a3e --- /dev/null +++ b/asyncron/backend.py @@ -0,0 +1,162 @@ + +from django.db import IntegrityError +from django.tasks.backends.base import BaseTaskBackend +from django.tasks.exceptions import InvalidTask +from django.tasks import task_backends, Task +from django.tasks.base import DEFAULT_TASK_PRIORITY, DEFAULT_TASK_QUEUE_NAME, DEFAULT_TASK_BACKEND_ALIAS + +from .models import Worker, TaskSchedule + +import dataclasses +import asyncio + + + +def scheduled_task( + function = None, *, + + schedule_name = "default", + interval = None, jitter = None, + + priority = DEFAULT_TASK_PRIORITY, + queue_name = DEFAULT_TASK_QUEUE_NAME, + backend = DEFAULT_TASK_BACKEND_ALIAS, + takes_context = True, + **kwargs +): + + if "run_after" in kwargs: + raise TypeError( + "run_after cannot be defined statically with the @task decorator. " + "Use .using(run_after=...) to set it dynamically." + ) + + def wrapper(f): + back = task_backends[backend] + assert hasattr( back, "TASK_SCHEDULES" ), f"Task Backend {back} ({backend=}) Does not support task schedules!" + + task = back.task_class( + func=f, + priority=priority, + queue_name=queue_name, + backend=backend, + takes_context=takes_context, + **kwargs, + ) + + ts = TaskSchedule( + name = schedule_name, + task_path = task.module_path, + interval = interval, + ) + ts.set_jitter( jitter ) + back.TASK_SCHEDULES.append( ts ) + + return task + + if function: return wrapper(function) + return wrapper + + +@dataclasses.dataclass(frozen=True, slots=True, kw_only=True) +class AsyncronTask( Task ): + timeout: int | None = None + is_sensitive: bool = False + + def __post_init__(self): + backend = self.get_backend() + #backend.validate_task(self) + backend.register_task(self) + + +class AsyncronBackend( BaseTaskBackend ): + task_class = AsyncronTask + supports_defer = False + supports_async_task = True + supports_get_result = False + supports_priority = False + + TASKS = {} + TASK_SCHEDULES = [] + #def __init__(self, alias, params, **kwargs): + # super().__init__(alias, params) + + + def register_task( self, task ): + self.TASKS[ task.module_path ] = task + + + + def _execute_task(self, task_result): + """ + Execute the Task for the given TaskResult, mutating it with the + outcome. + """ + object.__setattr__(task_result, "enqueued_at", timezone.now()) + task_enqueued.send(type(self), task_result=task_result) + + task = task_result.task + task_start_time = timezone.now() + object.__setattr__(task_result, "status", TaskResultStatus.RUNNING) + object.__setattr__(task_result, "started_at", task_start_time) + object.__setattr__(task_result, "last_attempted_at", task_start_time) + task_result.worker_ids.append(self.worker_id) + task_started.send(sender=type(self), task_result=task_result) + + try: + if task.takes_context: + raw_return_value = task.call( + TaskContext(task_result=task_result), + *task_result.args, + **task_result.kwargs, + ) + else: + raw_return_value = task.call(*task_result.args, **task_result.kwargs) + + object.__setattr__( + task_result, + "_return_value", + normalize_json(raw_return_value), + ) + except KeyboardInterrupt: + # If the user tried to terminate, let them + raise + except BaseException as e: + object.__setattr__(task_result, "finished_at", timezone.now()) + exception_type = type(e) + task_result.errors.append( + TaskError( + exception_class_path=( + f"{exception_type.__module__}.{exception_type.__qualname__}" + ), + traceback="".join(format_exception(e)), + ) + ) + object.__setattr__(task_result, "status", TaskResultStatus.FAILED) + task_finished.send(type(self), task_result=task_result) + else: + object.__setattr__(task_result, "finished_at", timezone.now()) + object.__setattr__(task_result, "status", TaskResultStatus.SUCCESSFUL) + task_finished.send(type(self), task_result=task_result) + + def enqueue( self, task, args, kwargs ): + self.validate_task(task) + + task_result = TaskResult( + task=task, + id=get_random_string(32), + status=TaskResultStatus.READY, + enqueued_at=None, + started_at=None, + last_attempted_at=None, + finished_at=None, + args=args, + kwargs=kwargs, + backend=self.alias, + errors=[], + worker_ids=[], + ) + + self._execute_task(task_result) + + return task_result diff --git a/asyncron/management/commands/run_asyncron_worker.py b/asyncron/management/commands/run_asyncron_worker.py index 68bc06d..42d4194 100644 --- a/asyncron/management/commands/run_asyncron_worker.py +++ b/asyncron/management/commands/run_asyncron_worker.py @@ -12,7 +12,7 @@ from django.core.management.base import BaseCommand, CommandError from django.conf import settings from asyncron.workers import AsyncronWorker -from asyncron.models import Task +from asyncron.models import TaskSchedule class bcolors: HEADER = '\033[95m' @@ -32,8 +32,8 @@ class Command(BaseCommand): AsyncronWorker.IS_ACTIVE = True while True: - worker = AsyncronWorker() + worker.is_db_ready.set() print( "Starting:", worker ) try: diff --git a/asyncron/migrations/0001_initial.py b/asyncron/migrations/0001_initial.py index 666cef5..e16bb63 100644 --- a/asyncron/migrations/0001_initial.py +++ b/asyncron/migrations/0001_initial.py @@ -1,7 +1,10 @@ -# Generated by Django 5.1.2 +# Generated by Django 6.1 on 2026-08-24 10:48 +import _thread import datetime import django.db.models.deletion +import posix +import uuid from django.db import migrations, models @@ -10,83 +13,65 @@ class Migration(migrations.Migration): initial = True dependencies = [ - ('contenttypes', '0002_remove_content_type_name'), ] operations = [ migrations.CreateModel( - name='Worker', + name='TaskSchedule', fields=[ ('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), - ('pid', models.IntegerField()), - ('thread_id', models.PositiveBigIntegerField()), - ('is_robust', models.BooleanField(default=False)), - ('is_master', models.BooleanField(default=False)), - ('in_grace', models.BooleanField(default=False)), - ('last_crowning_attempt', models.DateTimeField(blank=True, null=True)), - ('consumption_interval_seconds', models.IntegerField(default=10)), - ('consumption_total_active', models.IntegerField(default=0)), - ], - options={ - 'constraints': [models.UniqueConstraint(condition=models.Q(('is_master', True)), fields=('is_master',), name='only_one_master')], - }, - ), - migrations.CreateModel( - name='Task', - fields=[ - ('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), - ('name', models.TextField(unique=True)), - ('worker_type', models.CharField(choices=[('A', 'Any'), ('R', 'Robust'), ('D', 'Dynamic')], default='A')), - ('max_completed_traces', models.IntegerField(default=10)), - ('max_failed_traces', models.IntegerField(default=1000)), - ('timeout', models.DurationField(blank=True, default=datetime.timedelta(seconds=300), null=True)), - ('gracetime', models.DurationField(default=datetime.timedelta(seconds=60))), + ('name', models.CharField(default='default', max_length=200)), + ('is_enabled', models.BooleanField(default=True)), + ('task_path', models.TextField()), + ('args', models.JSONField(blank=True, default=list)), + ('kwargs', models.JSONField(blank=True, default=dict)), ('interval', models.DurationField(blank=True, null=True)), + ('timeout', models.DurationField(default=datetime.timedelta(seconds=300))), ('jitter_length', models.DurationField(blank=True, default=datetime.timedelta(0))), ('jitter_pivot', models.CharField(choices=[('S', 'Start'), ('M', 'Middle'), ('E', 'End')], default='M', max_length=1)), - ('worker_lock', models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, to='asyncron.worker')), + ('prune_success_after', models.IntegerField(default=10)), + ('prune_failed_after', models.IntegerField(default=1000)), ], options={ - 'abstract': False, + 'unique_together': {('name', 'task_path')}, }, ), migrations.CreateModel( - name='Metadata', + name='Worker', fields=[ - ('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), - ('model_id', models.PositiveIntegerField()), - ('name', models.CharField(max_length=256)), - ('data', models.JSONField(blank=True, null=True)), - ('expiration_datetime', models.DateTimeField(blank=True, null=True)), - ('model_type', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, to='contenttypes.contenttype')), + ('id', models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ('process_id', models.IntegerField(default=posix.getpid)), + ('thread_id', models.PositiveBigIntegerField(default=_thread.get_ident)), + ('creation_datetime', models.DateTimeField(auto_now_add=True)), + ('is_standalone', models.BooleanField(default=False)), + ('is_scheduler', models.BooleanField(default=False)), ], options={ - 'verbose_name': 'Metadata', - 'verbose_name_plural': 'Metadata', - 'indexes': [models.Index(fields=['model_type', 'model_id'], name='asyncron_me_model_t_d92186_idx')], + 'constraints': [models.UniqueConstraint(condition=models.Q(('is_scheduler', True)), fields=('is_scheduler',), name='unique_scheduler')], }, ), migrations.CreateModel( name='Trace', fields=[ - ('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), - ('status_reason', models.TextField(blank=True, default='')), - ('status', models.CharField(choices=[('S', 'Scheduled'), ('W', 'Waiting'), ('R', 'Running'), ('P', 'Paused'), ('C', 'Completed'), ('A', 'Aborted'), ('E', 'Error')], default='S', max_length=1)), - ('scheduled_datetime', models.DateTimeField(blank=True, null=True)), - ('register_datetime', models.DateTimeField(auto_now_add=True)), - ('last_run_datetime', models.DateTimeField(blank=True, null=True)), - ('last_end_datetime', models.DateTimeField(blank=True, null=True)), - ('protected', models.BooleanField(default=False)), + ('id', models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ('task_path', models.TextField()), + ('status_description', models.TextField(blank=True, default='')), + ('status', models.CharField(choices=[('READY', 'Ready'), ('RUNNING', 'Running'), ('FAILED', 'Failed'), ('SUCCESSFUL', 'Successful')], default='READY', max_length=10)), ('args', models.JSONField(blank=True, default=list)), ('kwargs', models.JSONField(blank=True, default=dict)), - ('stdout', models.TextField(blank=True, null=True)), - ('stderr', models.TextField(blank=True, null=True)), - ('returned', models.JSONField(blank=True, null=True)), - ('task', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, to='asyncron.task')), - ('worker_lock', models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, to='asyncron.worker')), + ('enqueued_datetime', models.DateTimeField(auto_now_add=True)), + ('scheduled_datetime', models.DateTimeField(blank=True, null=True)), + ('started_datetime', models.DateTimeField(blank=True, null=True)), + ('finished_datetime', models.DateTimeField(blank=True, null=True)), + ('stdout', models.TextField(blank=True, default='')), + ('exception_class_path', models.TextField(null=True)), + ('traceback', models.TextField(null=True)), + ('return_value', models.JSONField(blank=True, null=True)), + ('prune_protected', models.BooleanField(default=False)), + ('schedule', models.ForeignKey(null=True, on_delete=django.db.models.deletion.SET_NULL, to='asyncron.taskschedule')), ], options={ - 'constraints': [models.UniqueConstraint(condition=models.Q(('scheduled_datetime', None), ('status', 'S')), fields=('task_id',), name='unique_unscheduled_for_task')], + 'constraints': [models.UniqueConstraint(condition=models.Q(('status', 'READY')), fields=('schedule_id',), name='unique_ready_for_each_schedule')], }, ), ] diff --git a/asyncron/migrations/0002_task_self_aware.py b/asyncron/migrations/0002_task_self_aware.py deleted file mode 100644 index 8e35e49..0000000 --- a/asyncron/migrations/0002_task_self_aware.py +++ /dev/null @@ -1,18 +0,0 @@ -# Generated by Django 5.1.2 on 2025-01-21 22:21 - -from django.db import migrations, models - - -class Migration(migrations.Migration): - - dependencies = [ - ('asyncron', '0001_initial'), - ] - - operations = [ - migrations.AddField( - model_name='task', - name='self_aware', - field=models.BooleanField(default=False, help_text="Whether It's first argument is 'self', being a trace instance."), - ), - ] diff --git a/asyncron/migrations/0002_trace_worker.py b/asyncron/migrations/0002_trace_worker.py new file mode 100644 index 0000000..bf4f881 --- /dev/null +++ b/asyncron/migrations/0002_trace_worker.py @@ -0,0 +1,19 @@ +# Generated by Django 6.1 on 2026-08-24 11:15 + +import django.db.models.deletion +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('asyncron', '0001_initial'), + ] + + operations = [ + migrations.AddField( + model_name='trace', + name='worker', + field=models.ForeignKey(null=True, on_delete=django.db.models.deletion.SET_NULL, to='asyncron.worker'), + ), + ] diff --git a/asyncron/migrations/0003_alter_task_self_aware.py b/asyncron/migrations/0003_alter_task_self_aware.py deleted file mode 100644 index f973d44..0000000 --- a/asyncron/migrations/0003_alter_task_self_aware.py +++ /dev/null @@ -1,18 +0,0 @@ -# Generated by Django 5.1.2 on 2025-08-09 09:55 - -from django.db import migrations, models - - -class Migration(migrations.Migration): - - dependencies = [ - ('asyncron', '0002_task_self_aware'), - ] - - operations = [ - migrations.AlterField( - model_name='task', - name='self_aware', - field=models.BooleanField(default=True, help_text="Whether It's first argument is 'self', being a trace instance."), - ), - ] diff --git a/asyncron/migrations/0003_alter_trace_scheduled_datetime.py b/asyncron/migrations/0003_alter_trace_scheduled_datetime.py new file mode 100644 index 0000000..8171f10 --- /dev/null +++ b/asyncron/migrations/0003_alter_trace_scheduled_datetime.py @@ -0,0 +1,20 @@ +# Generated by Django 6.1 on 2026-08-24 15:58 + +import django.utils.timezone +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('asyncron', '0002_trace_worker'), + ] + + operations = [ + migrations.AlterField( + model_name='trace', + name='scheduled_datetime', + field=models.DateTimeField(default=django.utils.timezone.now), + preserve_default=False, + ), + ] diff --git a/asyncron/migrations/0004_alter_task_gracetime.py b/asyncron/migrations/0004_alter_task_gracetime.py deleted file mode 100644 index 98fad57..0000000 --- a/asyncron/migrations/0004_alter_task_gracetime.py +++ /dev/null @@ -1,19 +0,0 @@ -# Generated by Django 5.2.5 on 2025-11-22 16:54 - -import datetime -from django.db import migrations, models - - -class Migration(migrations.Migration): - - dependencies = [ - ('asyncron', '0003_alter_task_self_aware'), - ] - - operations = [ - migrations.AlterField( - model_name='task', - name='gracetime', - field=models.DurationField(default=datetime.timedelta(seconds=1)), - ), - ] diff --git a/asyncron/migrations/0005_rename_last_crowning_attempt_worker_last_activity.py b/asyncron/migrations/0005_rename_last_crowning_attempt_worker_last_activity.py deleted file mode 100644 index aa9d273..0000000 --- a/asyncron/migrations/0005_rename_last_crowning_attempt_worker_last_activity.py +++ /dev/null @@ -1,18 +0,0 @@ -# Generated by Django 5.2.5 on 2025-11-22 19:09 - -from django.db import migrations - - -class Migration(migrations.Migration): - - dependencies = [ - ('asyncron', '0004_alter_task_gracetime'), - ] - - operations = [ - migrations.RenameField( - model_name='worker', - old_name='last_crowning_attempt', - new_name='last_activity', - ), - ] diff --git a/asyncron/migrations/0006_remove_metadata_asyncron_me_model_t_d92186_idx_and_more.py b/asyncron/migrations/0006_remove_metadata_asyncron_me_model_t_d92186_idx_and_more.py deleted file mode 100644 index da79f53..0000000 --- a/asyncron/migrations/0006_remove_metadata_asyncron_me_model_t_d92186_idx_and_more.py +++ /dev/null @@ -1,20 +0,0 @@ -# Generated by Django 6.0.6 on 2026-07-25 15:34 - -from django.db import migrations - - -class Migration(migrations.Migration): - - dependencies = [ - ('asyncron', '0005_rename_last_crowning_attempt_worker_last_activity'), - ] - - operations = [ - migrations.RemoveIndex( - model_name='metadata', - name='asyncron_me_model_t_d92186_idx', - ), - migrations.DeleteModel( - name='Metadata', - ), - ] diff --git a/asyncron/models.py b/asyncron/models.py index 157b76f..7ff54b0 100644 --- a/asyncron/models.py +++ b/asyncron/models.py @@ -1,232 +1,251 @@ from django.utils import timezone from django.db import models from django.db.models.constraints import UniqueConstraint, Q +from django.tasks.base import TaskResultStatus +from django.tasks import task_backends -#To mock print, can't use redirect_stdout in async code -#This is buggy, has print leakage and still not very good, -#But better than nothing when a task is not self_aware or calls something that isn't -from unittest.mock import patch - +from .base.models import BaseModel +from .utils import TIMEDELTA_PATTERN import functools, traceback, io -import random +import random, uuid import asyncio - +import os, threading # Create your models here. -from .base.models import BaseModel class Worker( BaseModel ): + id = models.UUIDField( primary_key = True, default = uuid.uuid4, editable = False ) - pid = models.IntegerField() - thread_id = models.PositiveBigIntegerField() + process_id = models.IntegerField( default = os.getpid ) + thread_id = models.PositiveBigIntegerField( default = threading.get_ident ) - is_robust = models.BooleanField( default = False ) - is_master = models.BooleanField( default = False ) - in_grace = models.BooleanField( default = False ) #If the worker sees this as True, it should kill itself! + creation_datetime = models.DateTimeField( auto_now_add = True ) - last_activity = models.DateTimeField( null = True, blank = True ) + is_standalone = models.BooleanField( default = False ) + is_scheduler = models.BooleanField( default = False ) - #Variables with very feel good names! :) - consumption_interval_seconds = models.IntegerField( default = 10 ) - consumption_total_active = models.IntegerField( default = 0 ) - def __str__( self ): return f"P{self.pid}W{self.thread_id}" + ("R" if self.is_robust else "D") + #in_grace = models.BooleanField( default = False ) #If the worker sees this as True, it should kill itself! + #last_activity = models.DateTimeField( null = True, blank = True ) + + def __str__( self ): + return f"P{self.process_id}W{self.thread_id}" + ("M" if self.is_in_main() else "D") + + def is_in_main( self ): #not cross process + return self.thread_id == threading.main_thread().ident class Meta: constraints = [ - UniqueConstraint( fields = ('is_master',), condition = Q( is_master = True ), name='only_one_master'), + UniqueConstraint( fields = ['is_scheduler'], condition = Q( is_scheduler = True ), name='unique_scheduler'), ] - def is_proc_alive( self ): - import os - pid = self.pid #Slightly Altered: https://stackoverflow.com/a/20186516 - if pid < 0: return False #NOTE: pid == 0 returns True - try: os.kill(pid, 0) - except ProcessLookupError: return False # errno.ESRCH: No such process - except PermissionError: return True # errno.EPERM: Operation not permitted (i.e., process exists) - else: return True # no error, we can send a signal to the process +class TaskSchedule( BaseModel ): + name = models.CharField( default = "default", max_length = 200 ) -class Task( BaseModel ): + is_enabled = models.BooleanField( default = True ) - registered_tasks = {} #Name -> the actual function + task_path = models.TextField() #Path to the task function + @property + def task( self ): return task_backends['default'].TASKS[self.task_path] - name = models.TextField( unique = True ) #Path to the function - worker_lock = models.ForeignKey( Worker, null = True, blank = True, on_delete = models.SET_NULL ) - worker_type = models.CharField( default = "A", choices = { - "A": "Any", - "R": "Robust", #Only seperate Robust workers - "D": "Dynamic", #Only on potentially reloadable workers - }) + args = models.JSONField( default = list, blank = True ) + kwargs = models.JSONField( default = dict, blank = True ) - max_completed_traces = models.IntegerField( default = 10 ) - max_failed_traces = models.IntegerField( default = 1000 ) - - timeout = models.DurationField( - default = timezone.timedelta( minutes = 5 ), - null = True, blank = True - ) #None will mean it's a "service" like task - gracetime = models.DurationField( default = timezone.timedelta( seconds = 1 ) ) - self_aware = models.BooleanField( default = True, help_text = "Whether It's first argument is 'self', being a trace instance." ) - - #Periodic Tasks + #Distinguishes Periodic and Service like Tasks interval = models.DurationField( null = True, blank = True ) + @property + def type( self ): + if self.interval is None: return "S" #Service + return "P" #Periodic + + #For Periodic tasks, it's the whole execution, + #For Service like tasks, it's the gracetime after the first exit signal. + timeout = models.DurationField( default = timezone.timedelta( minutes = 5 ) ) + + #delay before execution for Both task types jitter_length = models.DurationField( default = timezone.timedelta( seconds = 0 ), blank = True ) jitter_pivot = models.CharField( default = "M", max_length = 1, choices = { "S":"Start", "M":"Middle", "E":"End", }) + def get_jitter( self ): jitter = self.jitter_length * random.random() match self.jitter_pivot: - case "M": - jitter -= self.jitter_length / 2 - case "E": - jitter *= -1 + case "M": jitter -= self.jitter_length / 2 + case "E": jitter *= -1 return jitter + def set_jitter( self, jitter ): + if not jitter: #Set the default values + self.jitter_length = timezone.timedelta( seconds = 0 ) + self.jitter_pivot = "M" + return + + match = TIMEDELTA_PATTERN.match( jitter ) + assert match, "Provided jitter value does not match the correct timedelta pattern! (ex: 1w2d5h30m10s500ms1000us)" + + self.jitter_length = timezone.timedelta( **{ + k: int(v) + for k, v in match.groupdict().items() + if v is not None + } ) + self.jitter_pivot = { + "-": "S", + None: "M", + "+": "E", + }[match.group(1)] + + prune_success_after = models.IntegerField( default = 10 ) + prune_failed_after = models.IntegerField( default = 1000 ) + + class Meta: + unique_together = [ + ('name', 'task_path'), + ] def __str__( self ): - type = "Callable" if self.interval is None else "Periodic" - mode = "Service" if self.timeout is None else "Task" - short = self.name.rsplit('.')[-1] - return " ".join([type, mode, short]) + type_display = ( "Service Task" if self.type == "S" else "Periodic Task" ) + short = self.task_path.rsplit('.')[-1] + return f"{self.name} {type_display} {short}" - def register( self, f ): - if not self.name: self.name = f"{f.__module__}.{f.__qualname__}" - self.registered_tasks[self.name] = f - f.task = self - return f + def as_trace( self ): + return Trace( + schedule = self, + task_path = self.task_path, - def new_trace( self ): - trace = Trace( task_id = self.id ) - trace.task = self #Less db hits - return trace + args = self.args, + kwargs = self.kwargs + ) + async def set_traces_scheduled_datetime( self, trace ): + assert trace.schedule_id == self.id, "This trace does not belong to this schedule!" - async def ensure_quick_execution( self, reason = "Quick Exec" ): now = timezone.now() - if await self.trace_set.filter( status = "W" ).aexists(): - return + jitter_delta = self.get_jitter() + last_trace = await self.trace_set.exclude( + status = TaskResultStatus.READY + ).order_by( '-started_datetime' ).afirst() + + if last_trace: #Execute now + jitter + trace.scheduled_datetime = last_trace.started_datetime + jitter_delta + if self.interval is not None: #Add the "Period" in periodic tasks + trace.scheduled_datetime += self.interval + + else: + trace.scheduled_datetime = now + jitter_delta + + if trace.scheduled_datetime < now: #So, in case jitter is negative + trace.scheduled_datetime = now - if await self.trace_set.filter( status = "S", scheduled_datetime__lte = now ).aexists(): - return - trace = await self.trace_set.filter( status = "S" ).order_by('scheduled_datetime').afirst() - if not trace: trace = self.new_trace() - await trace.reschedule( reason = reason, target_datetime = now ) - await trace.asave() class Trace( BaseModel ): - task = models.ForeignKey( Task, on_delete = models.CASCADE ) + id = models.UUIDField( primary_key = True, default = uuid.uuid4, editable = False ) + worker = models.ForeignKey( Worker, null = True, on_delete = models.SET_NULL ) - status_reason = models.TextField( default = "", blank = True ) - status = models.CharField( default = "S", max_length = 1, choices = { - "S":"Scheduled", - "W":"Waiting", - "R":"Running", - "P":"Paused", - "C":"Completed", - "A":"Aborted", - "E":"Error", - }) - def set_status( self, status, reason = "" ): + task_path = models.TextField() #Path to the task function + schedule = models.ForeignKey( TaskSchedule, null = True, on_delete = models.SET_NULL ) + @property + def task( self ): return task_backends['default'].TASKS[self.task_path] + + status_description = models.TextField( default = "", blank = True ) + status = models.CharField( + default = TaskResultStatus.READY, + choices = TaskResultStatus.choices, + max_length = max( len(v) for v in TaskResultStatus.values ), + ) + + def set_status( self, status, desc = "" ): self.status = status - self.status_reason = reason + self.status_description = desc - scheduled_datetime = models.DateTimeField( null = True, blank = True ) - register_datetime = models.DateTimeField( auto_now_add = True ) - last_run_datetime = models.DateTimeField( null = True, blank = True ) - last_end_datetime = models.DateTimeField( null = True, blank = True ) - - worker_lock = models.ForeignKey( Worker, null = True, blank = True, on_delete = models.SET_NULL ) - protected = models.BooleanField( default = False ) #Do not delete these. args = models.JSONField( default = list, blank = True ) kwargs = models.JSONField( default = dict, blank = True ) - stdout = models.TextField( null = True, blank = True ) - stderr = models.TextField( null = True, blank = True ) - returned = models.JSONField( null = True, blank = True ) - def __str__( self ): return f"Trace of Task {self.task}" + enqueued_datetime = models.DateTimeField( auto_now_add = True ) + scheduled_datetime = models.DateTimeField() + started_datetime = models.DateTimeField( null = True, blank = True ) + finished_datetime = models.DateTimeField( null = True, blank = True ) + + stdout = models.TextField( default = "", blank = True ) + exception_class_path = models.TextField( null = True ) + traceback = models.TextField( null = True ) + + return_value = models.JSONField( null = True, blank = True ) + prune_protected = models.BooleanField( default = False ) #Do not delete traces with this flag. + + def __str__( self ): return f"{self.get_status_display()} Trace of {self.task_path}" + + def as_task_result( self ): + return TaskResult( + task = self.task, + id = self.id, + status = self.status, + enqueued_at = self.enqueued_datetime, + started_at = self.started_datetime, + last_attempted_at = self.started_datetime, + finished_at = self.finished_datetime, + args = self.args, + kwargs = self.kwargs, + backend = self.task.backend, + errors = [], + worker_ids = [], + ) class Meta: constraints = [ UniqueConstraint( - fields = ['task_id'], - condition = models.Q(status = "S", scheduled_datetime = None), - name = "unique_unscheduled_for_task", - ) + fields = ['schedule_id'], + condition = models.Q(status = TaskResultStatus.READY), + name = "unique_ready_for_each_schedule", + ), ] - async def reschedule( self, reason = "", target_datetime = None ): - assert self.status in "SAE", f"Cannot reschedule a task that is in {self.get_status_display()} state!" - - await self.eval_related('task') - assert self.task.interval or target_datetime, "This is not a periodic task! Nothing to reschedule." - - self.set_status( "S", reason ) - if target_datetime: - self.scheduled_datetime = target_datetime - - else: - base_time = self.last_run_datetime or timezone.now() - jitter = self.task.get_jitter() - self.scheduled_datetime = base_time + self.task.interval + jitter - - if self.id: await self.asave( update_fields = ["status", "status_reason", "scheduled_datetime"] ) - - async def start( self ): - await self.eval_related('task') - assert self.status in "SPAWE", f"Cannot start a task that is in {self.get_status_display()} state!" + await self.eval_related('schedule') + schedule = self.schedule + task = self.task - self.last_run_datetime = timezone.now() - self.last_end_datetime = None - self.returned = None - self.stderr = "" - self.stdout = "" - - try: - func = Task.registered_tasks[self.task.name] - except KeyError: - self.set_status( "E", "Script Missing!" ) - await self.asave() - return - - self.set_status( "R" ) - await self.asave() + #self.started_datetime = timezone.now() + self.finished_datetime = None + self.return_value = None + self.exception_class_path = None + self.traceback = None #Runtime Bits - self.loop = asyncio.get_running_loop() + loop = asyncio.get_running_loop() self.new_print = asyncio.Event() - self.commit_on_new_print_task = self.loop.create_task( self.commit_on_new_print() ) + self.commit_on_new_print_task = loop.create_task( self.commit_on_new_print() ) try: async with asyncio.timeout( None ) as tmcm: - if self.task.timeout: - tmcm.reschedule( self.loop.time() + self.task.timeout.total_seconds() ) + if schedule and schedule.interval is not None and schedule.timeout: + tmcm.reschedule( loop.time() + self.schedule.timeout.total_seconds() ) - if self.task.self_aware: - output = await func(self, *self.args, **self.kwargs ) + if task.takes_context: + output = await task.func( self, *self.args, **self.kwargs ) else: - with patch( 'builtins.print', self.print ): - output = await func( *self.args, **self.kwargs ) + output = await task.func( *self.args, **self.kwargs ) - except TimeoutError: - self.set_status( "E", f"Timed out" ) - self.stderr = traceback.format_exc() + except TimeoutError as e: + self.set_status( TaskResultStatus.FAILED, f"Timed out" ) + self.traceback = traceback.format_exc() + self.exception_class_path = e.__qualname__ except Exception as e: - self.set_status( "E", f"Exception: {e}" ) - self.stderr = traceback.format_exc() + self.set_status( TaskResultStatus.FAILED, f"Exception: {e}" ) + self.traceback = traceback.format_exc() + self.exception_class_path = str(e) #.__qualname__ else: - self.set_status( "C" ) - self.returned = output + self.set_status( TaskResultStatus.SUCCESSFUL ) + self.return_value = output finally: self.commit_on_new_print_task.cancel() - self.last_end_datetime = timezone.now() + self.finished_datetime = timezone.now() await self.asave() diff --git a/asyncron/shortcuts.py b/asyncron/shortcuts.py index 8b5a4a7..35a56c5 100644 --- a/asyncron/shortcuts.py +++ b/asyncron/shortcuts.py @@ -1,59 +1,26 @@ ## ## decorators / functions to make the task calls easier ## + from django.utils.dateparse import parse_duration from django.db import models from django.utils import timezone from django.apps import apps -import re - -# Regular expression pattern with named groups for "1w2d5h30m10s500ms1000us" without spaces -pattern = re.compile( - r'(\+|-)?' - r'(?:(?P\d+)w)?' - r'(?:(?P\d+)d)?' - r'(?:(?P\d+)h)?' - r'(?:(?P\d+)m)?' - r'(?:(?P\d+)s)?' - r'(?:(?P\d+)ms)?' - r'(?:(?P\d+)us)?' -) def task( *args, **kwargs ): - from .models import Task - - jitter = kwargs.pop('jitter', "") - match = pattern.match( jitter ) - if jitter and match: - kwargs['jitter_length'] = timezone.timedelta( **{ - k: int(v) - for k, v in match.groupdict().items() - if v is not None - } ) - kwargs['jitter_pivot'] = { - "-": "S", - None: "M", - "+": "E", - }[match.group(1)] - - for f in Task._meta.fields: - if not isinstance(f, models.DurationField): continue - if f not in kwargs: continue - if not kwargs[f]: continue - - kwargs[f] = parse_duration( kwargs[f] ) - - return Task( *args, **kwargs ).register - + from .backend import scheduled_task + kwargs.setdefault( 'schedule_name', 'scripted' ) + return scheduled_task( *args, **kwargs ) def service( *args, **kwargs ): - kwargs.setdefault( 'timeout', None ) - kwargs.setdefault( 'worker_type', "R" ) - kwargs.setdefault( 'self_aware', True ) - return task( *args, **kwargs ) - + from .backend import scheduled_task + kwargs.setdefault( 'interval', None ) + kwargs.setdefault( 'is_sensitive', True ) + kwargs.setdefault( 'schedule_name', 'scripted' ) + return scheduled_task( *args, **kwargs ) def run_on_model_change( *models ): + return lambda x:x models = [ apps.get_model(m) if isinstance(m, str) else m for m in models diff --git a/asyncron/utils.py b/asyncron/utils.py index 84815ab..e833384 100644 --- a/asyncron/utils.py +++ b/asyncron/utils.py @@ -1,4 +1,6 @@ import functools +import os +import re def rsetattr(obj, attr, val): pre, _, post = attr.rpartition('.') @@ -18,6 +20,27 @@ def rupdate(d, u): #https://stackoverflow.com/a/3233356/ d[k] = v return d +def is_proc_alive( pid ): + #Slightly Altered: https://stackoverflow.com/a/20186516 + if pid < 0: return False #NOTE: pid == 0 returns True + try: os.kill(pid, 0) + except ProcessLookupError: return False # errno.ESRCH: No such process + except PermissionError: return True # errno.EPERM: Operation not permitted (i.e., process exists) + else: return True # no error, we can send a signal to the process + + +# Regular expression pattern with named groups for "1w2d5h30m10s500ms1000us" without spaces +TIMEDELTA_PATTERN = re.compile( + r'(\+|-)?' + r'(?:(?P\d+)w)?' + r'(?:(?P\d+)d)?' + r'(?:(?P\d+)h)?' + r'(?:(?P\d+)m)?' + r'(?:(?P\d+)s)?' + r'(?:(?P\d+)ms)?' + r'(?:(?P\d+)us)?' +) + #Django keeps giving: exception=OperationalError('the connection is closed') from django.db.utils import OperationalError @@ -37,5 +60,5 @@ def ignore_on_db_error( f ): return await f( *args, **kwargs ) except OperationalError as e: raise #return For DEV - + return decorator diff --git a/asyncron/workers.py b/asyncron/workers.py index 54d485b..fbc0ea5 100644 --- a/asyncron/workers.py +++ b/asyncron/workers.py @@ -2,11 +2,12 @@ from django.db import IntegrityError, models, close_old_connections from django.db.backends import signals as django_signals from django.db.utils import OperationalError +from django.db import IntegrityError from django.utils import timezone - +from django.tasks import task_backends +from django.tasks.base import TaskResultStatus from asgiref.sync import sync_to_async - import os, signal import time import threading @@ -14,10 +15,14 @@ import logging, traceback import asyncio import collections, functools import random +import humanize from .utils import retry_on_db_error, ignore_on_db_error from .asynctools import AsyncOriented from .singleton import Singleton +from .models import Worker, TaskSchedule, Trace + + class AsyncronWorker( Singleton, AsyncOriented ): """ @@ -54,11 +59,9 @@ class AsyncronWorker( Singleton, AsyncOriented ): def __init__( self ): self.log #Evaluating the log property while we have the creation lock - self.is_db_ready_event = asyncio.Event() - self.is_stopping_event = asyncio.Event() - - #Just so that the asyncron.apps.ready doens't trigger the django warning - self.start_after_db_ready = False + self.is_db_ready = asyncio.Event() + self.is_work_over = asyncio.Event() + self.model = Worker() #Parallelism self.thread = None #Once the worker starts, it'll be populated @@ -66,7 +69,6 @@ class AsyncronWorker( Singleton, AsyncOriented ): self.clearing_dead_workers = False self.watching_models = collections.defaultdict( set ) # Model -> Set of key name of the tasks - self.database_unreachable = False for callback in self.INIT_CALLBACKS: callback( self ) self.register_with_exit_signals() @@ -107,14 +109,14 @@ class AsyncronWorker( Singleton, AsyncOriented ): self.stop(f"Signal {signal.strsignal(signum)}") def handle_new_db_connection( self, sender, **kwargs ): - if self.is_db_ready_event.is_set(): return + if self.is_db_ready.is_set(): return self.log.debug(f"First DB connection: {sender}") - self.is_db_ready_event.set() + self.is_db_ready.set() django_signals.connection_created.disconnect( self.handle_new_db_connection ) #if not self.loop: return #We're not inside an async runner in this(?) or another thread. - #self.loop.call_soon_threadsafe( self.is_db_ready_event.set ) + #self.loop.call_soon_threadsafe( self.is_db_ready.set ) def start( self, daemon = False ): @@ -127,17 +129,17 @@ class AsyncronWorker( Singleton, AsyncOriented ): assert threading.main_thread() == threading.current_thread(), f"Cannot run a non daemon worker, in a thread other than main! Current Thread: {threading.current_thread()}" self.thread = threading.current_thread() - self.start_working( is_robust = True ) + self.start_working( is_standalone = True ) def stop( self, reason = None ): - if self.is_stopping_event.is_set(): return #TODO: Insisting on exiting faster should probably be managed in the signal handler + if self.is_work_over.is_set(): return #TODO: Insisting on exiting faster should probably be managed in the signal handler self.log.info( f"Stopping Worker: {reason}" ) - self.is_stopping_event.set() + self.is_work_over.set() #if not self.loop: return - #self.loop.call_soon_threadsafe( self.is_stopping_event.set ) + #self.loop.call_soon_threadsafe( self.is_work_over.set ) ## @@ -150,42 +152,48 @@ class AsyncronWorker( Singleton, AsyncOriented ): async def startup( self ): await super().startup() + self.backend = task_backends['default'] self.task_reason_jobs_queue = asyncio.Queue() #Run tasks from other threads, safely + async def cleanup( self ): + try: + count = await Trace.objects.filter( status = TaskResultStatus.RUNNING, worker = self.model ).aupdate( + status_description = "Worker died during execution", + status = TaskResultStatus.FAILED, worker = None + ) + await self.model.adelete() + except Exception as e: + self.log.info(f"Marking traces as 'Aborted' Failed: {e}") + else: + if count: self.log.info(f"Marked {count} trace(s) as 'Aborted'.") - def start_working( self, is_robust = False ): + await super().cleanup() + + def start_working( self, is_standalone = False ): with asyncio.Runner() as runner: runner.run( self.startup() ) - if self.start_after_db_ready: - self.log.debug("Waiting on another module to create the first database connection...") - runner.run( self.is_db_ready_event.wait() ) + if not self.is_db_ready.is_set(): + #This Avoid's the django initialization warning, since waited for is_db_ready above! + self.log.debug(f"Waiting on another module to create the first database connection...") + runner.run( self.is_db_ready.wait() ) + self.log.debug(f"DB Connection ready, continuing...") - from .models import Worker, Task, Trace - self.model = Worker( pid = os.getpid(), thread_id = threading.get_ident(), is_robust = is_robust ) + self.model.is_standalone = is_standalone + runner.run( self.model.asave() ) - self.create_task( self.consume_task_reason_jobs_queue() ) + runner.run( self.sync_task_schedules() ) + #self.create_task( self.consume_task_reason_jobs_queue() ) - self.create_task( self.master_main() ) - self.create_task( self.work_loop() ) + #self.create_task( self.maintain_scheduler() ) - self.model.save() #likley Avoid's the django initialization warning, since waited for is_db_ready_event above! - self.model.refresh_from_db() - - - #Fill in the ID fields of the tasks we didn't dare to check with db until now - from .models import Task - for func in Task.registered_tasks.values(): - task = func.task - if not task.pk: - try: task.pk = Task.objects.get( name = task.name ).pk - except Task.DoesNotExist: pass #It's a new one, it's fine. + self.create_task( self.run_tasks_on_schedule() ) self.attach_django_signals() try: - runner.run( self.is_stopping_event.wait() ) #This is the lifetime of this worker + runner.run( self.is_work_over.wait() ) #This is the lifetime of this worker except KeyboardInterrupt: self.log.info(f"[W{self.model.id}] Worker Received KeyboardInterrupt, exiting...") @@ -203,15 +211,7 @@ class AsyncronWorker( Singleton, AsyncOriented ): runner.run( self.cleanup() ) - try: - count = Trace.objects.filter( status__in = "SWRP", worker_lock = self.model ).update( - status_reason = "Worker died during execution", - status = "A", worker_lock = None - ) - except Exception as e: - self.log.info(f"Marking traces as 'Aborted' Failed: {e}") - else: - if count: self.log.info(f"Marked {count} trace(s) as 'Aborted'.") + # Almost always, a worker can delete it's model from the db, # But it seems that there is a sitation where despite Task.worker_lock being delete=SET_NULL, @@ -220,16 +220,102 @@ class AsyncronWorker( Singleton, AsyncOriented ): # Error example: # - django.db.utils.IntegrityError: update or delete on table "asyncron_worker" violates foreign key constraint "asyncron_task_worker_lock_id_0bb55026_fk_asyncron_worker_id" on table "asyncron_task" - for attempt in range(3): - try: - self.model.delete() - except IntegrityError as e: - self.log.warning(f"Deleting worker (self) failed {attempt+1} time(s): {e}") - time.sleep( 0.1 ) - else: break - self.log.debug("Worker stopped working.") + async def sync_task_schedules( self ): + for ts in self.backend.TASK_SCHEDULES: + try: await ts.asave() + except IntegrityError: pass + + + async def run_tasks_on_schedule( self ): + self.check_interval = 0 + + while True: + + if not await Worker.objects.filter( id = self.model.id ).aexists(): break + + await asyncio.sleep( self.check_interval ) + self.check_interval = 10 + + async for ts in TaskSchedule.objects.filter( is_enabled = True ): + if await ts.trace_set.filter( status = TaskResultStatus.RUNNING ).aexists(): continue + + trace = ts.as_trace() + await ts.set_traces_scheduled_datetime( trace ) + + try: await trace.asave() + except IntegrityError: continue + #else: print("Created a new trace:", trace, humanize.naturaltime( trace.scheduled_datetime ) ) + + Ts = Trace.objects.filter( + status = TaskResultStatus.READY, + scheduled_datetime__lte = timezone.now(), + worker = None + ).order_by('-scheduled_datetime') + + async for trace in Ts: + locked = await Trace.objects.filter( id = trace.id, worker = None ).aupdate( + started_datetime = timezone.now(), + status = TaskResultStatus.RUNNING, + status_description = "Running scheduled task", + worker = self.model.id + ) + if not locked: continue + + await trace.arefresh_from_db() + await trace.eval_related("schedule") + schedule = trace.schedule + + self.create_task( self.start_trace_on_time( trace ) ) + + await sync_to_async( close_old_connections )() + + + self.is_work_over.set() + + async def start_trace_on_time( self, trace ): + + if timezone.now() < trace.scheduled_datetime: #If this is a periodic task + await asyncio.sleep( ( timezone.now() - trace.scheduled_datetime ).total_seconds() ) + await trace.arefresh_from_db() + + await trace.start() + + trace.worker = None + await trace.asave( update_fields = ['worker'] ) + + return + #Traces for the same task that we are done with (Completed, Aborted, Errored) + QuerySet = Trace.objects.filter( + task_id = trace.task_id, status__in = "CAE", protected = False, worker = None + ).order_by('-register_datetime') + + #Should be deleted after the threashold + max_count = trace.task.max_completed_traces if trace.status == "C" else trace.task.max_failed_traces + await QuerySet.exclude( + id__in = QuerySet[:max_count].values_list( 'id', flat = True ) + ).adelete() + + + + + + + + + + + + + + + + + + + + def attach_django_signals( self ): django_name_to_signals = { @@ -239,9 +325,10 @@ class AsyncronWorker( Singleton, AsyncOriented ): and ( attr := getattr(models.signals, name) ) #Just an assignment and isinstance( attr, models.signals.ModelSignal ) #Is a signal related to models! } - for name, signal in django_name_to_signals.items(): - signal.connect( functools.partial( self.model_changed, name ) ) + #for name, signal in django_name_to_signals.items(): + # signal.connect( functools.partial( self.model_changed, name ) ) + return from .models import Task for name, task in Task.registered_tasks.items(): if not hasattr(task, 'watching_models'): continue @@ -258,7 +345,7 @@ class AsyncronWorker( Singleton, AsyncOriented ): #print("Will not run another trace of the same task to reduce the change of an infinite cycle.") continue - if task.worker_type not in ("AR" if self.model.is_robust else "AD"): #If we can't run this trace + if task.worker_type not in ("AR" if self.model.is_standalone else "AD"): #If we can't run this trace asyncio.run_coroutine_threadsafe( task.ensure_quick_execution( reason = f"Change ({signal_name}) on {instance}" ), self.loop @@ -286,12 +373,11 @@ class AsyncronWorker( Singleton, AsyncOriented ): except: await asyncio.sleep( 1 ) - async def master_main( self ): + async def maintain_scheduler( self ): """ - Fight over who's gonna be the master. - Prove your health in the process! + Make sure at least one process is managing the scheduled tasks. """ - from .models import Worker, Task, Trace + from .models import Worker, TaskSchedule, Trace is_current_master = False next_overtake_attempt = time.time() + 1 + random.random() * 5 @@ -301,14 +387,7 @@ class AsyncronWorker( Singleton, AsyncOriented ): while await MyQs.aupdate( last_activity = timezone.now() ): try: - await Worker.objects.filter( is_master = False ).aupdate( is_master = models.Q(id = self.model.id) ) - - except RuntimeError as e: #Syntax Error cause: cannot schedule new futures after interpreter shutdown - self.log.critical(f"[W{self.model.id}] Master Loop Runtime Error:", e ) - if "interpreter shutdown" not in e.args[0]: raise - self.database_unreachable = True - loop_wait = 0 - break + await Worker.objects.filter( is_scheduler = False ).aupdate( is_scheduler = models.Q(id = self.model.id) ) except IntegrityError: # I'm not master! loop_wait = 5 + random.random() * 15 @@ -321,15 +400,15 @@ class AsyncronWorker( Singleton, AsyncOriented ): next_overtake_attempt = time.time() + 60 took_master = False - if self.model.is_robust: - took_master = await Worker.objects.filter( is_master = True, is_robust = False ).aupdate( is_master = False ) + if self.model.is_standalone: + took_master = await Worker.objects.filter( is_scheduler = True, is_standalone = False ).aupdate( is_scheduler = False ) loop_wait = 0 else: - await Worker.objects.filter( is_master = True ).filter( + await Worker.objects.filter( is_scheduler = True ).filter( models.Q( last_activity = None ) | models.Q( last_activity__lte = timezone.now() - timezone.timedelta( minutes = 2 ) ) - ).aupdate( is_master = False ) + ).aupdate( is_scheduler = False ) else: #I am Master! @@ -352,7 +431,7 @@ class AsyncronWorker( Singleton, AsyncOriented ): async def clear_orphaned_traces( self ): from .models import Worker, Task, Trace - await Trace.objects.filter( worker_lock = None, status__in = "RPW" ).adelete() + await Trace.objects.filter( worker = None, status = "RUNNING" ).adelete() async def clear_dead_workers( self ): self.clearing_dead_workers = True @@ -370,62 +449,8 @@ class AsyncronWorker( Singleton, AsyncOriented ): await Worker.objects.filter( in_grace = True ).adelete() self.clearing_dead_workers = False - async def sync_tasks( self ): - from .models import Task - for name, func in Task.registered_tasks.items(): - - init_task = func.task - try: - func.task = await Task.objects.aget( name = name ) - except Task.DoesNotExist: - await func.task.asave() - await func.task.arefresh_from_db() - else: #For now, to commit changes to db - init_task.id = func.task.id - - #DEBUG this: - #django.db.utils.IntegrityError: insert or update on table "asyncron_task" violates foreign key constraint "asyncron_task_worker_lock_id_0bb55026_fk_asyncron_worker_id" - #DETAIL: Key (worker_lock_id)=(6035) is not present in table "asyncron_worker". - await init_task.asave() - #END of DEBUG! - - await func.task.arefresh_from_db() - async def work_loop( self ): - from .models import Worker, Task, Trace - - self.check_interval = 0 - - while not self.database_unreachable: - - if not await Worker.objects.filter( id = self.model.id ).aexists(): break - - await asyncio.sleep( self.check_interval ) - self.check_interval = 10 - - try: - await self.check_services() - await self.check_scheduled() - except OperationalError as e: - self.log.warning(f"DB Connection Error: {e}") - self.log.warning(f"Traceback:\n{traceback.format_exc()}" ) - self.check_interval = 60 #break - - except Exception as e: - self.log.warning(f"check_scheduled failed: {e}") - self.log.warning(f"Traceback:\n{traceback.format_exc()}" ) - self.check_interval = 20 - - try: - await sync_to_async( close_old_connections )() - except Exception as e: - self.log.warning(f"close_old_connections failed: {e}") - self.log.warning(f"Traceback:\n{traceback.format_exc()}" ) - #break - - - self.is_stopping_event.set() async def consume_task_reason_jobs_queue( self ): while True: @@ -439,7 +464,7 @@ class AsyncronWorker( Singleton, AsyncOriented ): #Schedule traces that aren't yet set. Ts = Task.objects.exclude( interval = None ).exclude( trace__status = "S" - ).exclude( worker_type = "D" if self.model.is_robust else "R" ) + ).exclude( worker_type = "D" if self.model.is_standalone else "R" ) async for task in Ts: trace = task.new_trace() @@ -458,11 +483,11 @@ class AsyncronWorker( Singleton, AsyncOriented ): await Task.objects.filter( id = task.id, worker_lock = self.model ).aupdate( worker_lock = None ) early_seconds = 5 + self.check_interval * ( 1 + random.random() ) - async for trace in Trace.objects.filter( status = "S", worker_lock = None, scheduled_datetime__lte = timezone.now() + timezone.timedelta( seconds = early_seconds ) ): + async for trace in Trace.objects.filter( status = "S", worker = None, scheduled_datetime__lte = timezone.now() + timezone.timedelta( seconds = early_seconds ) ): await trace.eval_related() #print(f"Checking {trace} to do now: {trace.scheduled_datetime - timezone.now()}") - try: count = await Trace.objects.filter( id = trace.id, status = "S" ).aupdate( status = "W", worker_lock = self.model ) + try: count = await Trace.objects.filter( id = trace.id, status = "S" ).aupdate( status = "W", worker = self.model ) except IntegrityError: count = 0 if not count: continue #Lost the race condition to another worker. @@ -474,7 +499,7 @@ class AsyncronWorker( Singleton, AsyncOriented ): #start services that aren't running yet. Ts = Task.objects.filter( interval = None, timeout = None ).exclude( trace__status__in = "WR" - ).exclude( worker_type = "D" if self.model.is_robust else "R" ) + ).exclude( worker_type = "D" if self.model.is_standalone else "R" ) async for task in Ts: @@ -488,7 +513,7 @@ class AsyncronWorker( Singleton, AsyncOriented ): trace = task.new_trace() trace.set_status( "W", "Waiting to start the service ASAP." ) - trace.worker_lock = self.model + trace.worker = self.model await trace.asave() #self.running_service_tasks[task.id] = @@ -499,29 +524,6 @@ class AsyncronWorker( Singleton, AsyncOriented ): trace = task.new_trace() trace.set_status( "S", reason ) trace.scheduled_datetime = timezone.now() #So it runs instantly - trace.worker_lock_id = self.model.id + trace.worker_id = self.model.id self.create_task( self.start_trace_on_time( trace ) ) - - async def start_trace_on_time( self, trace ): - from .models import Trace - - if trace.scheduled_datetime and timezone.now() < trace.scheduled_datetime: #If this is a periodic task - await asyncio.sleep( ( timezone.now() - trace.scheduled_datetime ).total_seconds() ) - await trace.arefresh_from_db() - - await trace.start() - - trace.worker_lock = None - await trace.asave( update_fields = ['worker_lock'] ) - - #Traces for the same task that we are done with (Completed, Aborted, Errored) - QuerySet = Trace.objects.filter( - task_id = trace.task_id, status__in = "CAE", protected = False, worker_lock = None - ).order_by('-register_datetime') - - #Should be deleted after the threashold - max_count = trace.task.max_completed_traces if trace.status == "C" else trace.task.max_failed_traces - await QuerySet.exclude( - id__in = QuerySet[:max_count].values_list( 'id', flat = True ) - ).adelete()