diff --git a/asyncron/admin.py b/asyncron/admin.py index d216226..4f51ba4 100644 --- a/asyncron/admin.py +++ b/asyncron/admin.py @@ -14,39 +14,90 @@ import humanize @admin.register( Worker ) class WorkerAdmin( BaseModelAdmin ): order = 1 - list_display = 'process_id', 'thread_id', 'is_standalone', 'is_scheduler', 'is_running', # 'health', + list_display = 'process_id', 'thread_id', 'is_standalone', 'is_coordinator', 'is_running', 'last_activity', def has_add_permission( self, request, obj = None ): return False 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 ) + @admin.display( ordering = "last_activity" ) + def last_activity( self, obj ): return humanize.naturaltime( obj.last_activity ) -class TaskAppFilter( admin.SimpleListFilter ): - title = "app" - parameter_name = 'app_groups' + +class TaskPathFilter( admin.SimpleListFilter ): + #Vibe Coded nested list_filter usable inside two models + + title = "task path" + parameter_name = "task_path_prefix" def lookups( self, request, model_admin ): - return ( - ( a.lower(), a ) - for a in sorted({ - ts.task_path.split(".", 1)[0] - for ts in TaskSchedule.objects.all() - }) if a - ) + # Only needs to return *something* non-empty so Django considers + # this filter "active" (has_output()). Actual UI is in choices(). + return [("", "")] - def queryset( self, request, queryset ): - q = self.value() - if not q: return queryset - return queryset.filter( task_path__istartswith = q ) + def queryset(self, request, queryset): + value = self.value() + if value: return queryset.filter( task_path__startswith = value ) + return queryset + + def choices( self, changelist ): + + current_prefix = self.value() or "" + current_prefix_dotted = f"{current_prefix}." if current_prefix else "" + + # Only look at task_paths matching the current prefix, to keep + # the "next segment" query cheap-ish on large tables. + qs = changelist.queryset.model.objects.all() + if current_prefix: + qs = qs.filter( task_path__startswith = current_prefix + "." ) + + paths = qs.values_list( "task_path", flat = True ).distinct() + + next_segments = set() + for path in paths: + remainder = path[len(current_prefix_dotted):] + if not remainder: + continue + next_segments.add(remainder.split(".", 1)[0]) + + # "All" — full reset + yield { + "selected": not current_prefix, + "query_string": changelist.get_query_string(remove=[self.parameter_name]), + "display": "All", + } + + # Breadcrumb trail — click any ancestor to jump back up + if current_prefix: + parts = current_prefix.split(".") + for i in range(len(parts)): + ancestor = ".".join(parts[: i + 1]) + yield { + "selected": ancestor == current_prefix, + "query_string": changelist.get_query_string( + {self.parameter_name: ancestor} + ), + "display": f" {'-' * i} {parts[i]}", + } + + # Next-level children under the current prefix + depth = current_prefix.count(".") + 1 if current_prefix else 0 + for segment in sorted(next_segments): + child_prefix = f"{current_prefix_dotted}{segment}" + yield { + "selected": False, + "query_string": changelist.get_query_string( + {self.parameter_name: child_prefix} + ), + "display": f" {'-' * depth} {segment}", + } @admin.register( TaskSchedule ) class TaskScheduleAdmin( BaseModelAdmin ): order = 2 list_display = 'schedule_name', 'timeout', 'type', 'jitter', 'logged', 'last_execution', 'scheduled', 'is_enabled', - list_filter = TaskAppFilter, + list_filter = TaskPathFilter, fields = "name", "description", "type", #"on_model_change", "jitter", @@ -146,84 +197,56 @@ class TaskScheduleAdmin( BaseModelAdmin ): -class TraceAppFilter(admin.SimpleListFilter): - title = "app" - parameter_name = 'app_groups' - - def lookups( self, request, model_admin ): - return ( - ( a.lower(), a ) - for a in sorted({ - task.name.split(".", 1)[0] - for task in Task.objects.all() - }) if a - ) - - def queryset( self, request, queryset ): - q = self.value() - if not q: return queryset - return queryset.filter( task__name__istartswith = q ) - -class TraceNameFilter(admin.SimpleListFilter): - title = "task name" - parameter_name = 'short_name' - - def lookups( self, request, model_admin ): - return ( - ( a.lower(), a ) - for a in sorted({ - task.name.rsplit(".", 1)[-1] - for task in Task.objects.all() - }) if a - ) - - def queryset( self, request, queryset ): - q = self.value() - if not q: return queryset - return queryset.filter( task__name__iendswith = q ) - - @admin.register( Trace ) class TraceAdmin( BaseModelAdmin ): order = 3 - list_display = 'scheduled_datetime', 'task_path', 'execution', 'status', 'worker' - list_filter = 'status', 'status_description', #'task__worker_type', TraceAppFilter, TraceNameFilter, + list_display = 'scheduled_datetime', 'task_path', 'execution', 'status_desc', 'worker' + list_filter = 'status', TaskPathFilter 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 + @admin.display( description = "Status description", ordering = "status" ) + def status_desc( self, obj ): + desc = obj.lifetime.get( obj.status, "" ) + return f"{obj.get_status_display()}: {desc}" if desc else obj.get_status_display() + + + @admin.display( description = "Execution Summary", ordering = F("started_datetime").asc(nulls_first=True) ) def execution( self, obj ): + now = timezone.now() + if obj.finished_datetime: - return "- Ran " + humanize.naturaltime( obj.finished_datetime ) + if not obj.started_datetime: return "Malformed Trace" + took = obj.finished_datetime - obj.started_datetime + bits = [] - 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' + if took.total_seconds() < 1: + bits.append(f"Ran and finished {humanize.naturaltime(obj.finished_datetime)}") + else: + bits.append(f"Ran {humanize.naturaltime(obj.finished_datetime)} for {humanize.naturaldelta(took)}") - def state( self, obj ): - if obj.status_description: return f"{obj.status}: {obj.status_description}" - return obj.get_status_display() - state.admin_order_field = 'status' + delay = obj.started_datetime - obj.scheduled_datetime + if delay.total_seconds() > 1: bits.append(f"delayed {humanize.naturaldelta(delay)}") + return ", ".join(bits) + + if obj.started_datetime: + return f"Started {humanize.naturaltime(obj.started_datetime)}" + + if obj.scheduled_datetime < now: + overdue = now - obj.scheduled_datetime + return f"Overdue by {humanize.naturaldelta(overdue)}" + + return f"In {humanize.naturaltime(obj.scheduled_datetime)}" - actions = 'reschedule_to_now', - @admin.action( description = "Reschedule to run now" ) - def reschedule_to_now( self, request, qs ): - results = asyncio.run( - qs.exclude( task__interval = None ).filter( status = "S" ).gather_method( - 'reschedule', - reason = "Manually Rescheduled", - target_datetime = timezone.now(), - ) - ) - self.explain_gather_results( request, results, 5 ) - - + actions = 'duplicate_and_enqueue', + @admin.action( description = "Duplicate and Enqueue" ) + def duplicate_and_enqueue( self, request, qs ): + for trace in qs: trace.task.enqueue( *trace.args, **trace.kwargs ) diff --git a/asyncron/backend.py b/asyncron/backend.py deleted file mode 100644 index 3cd5a3e..0000000 --- a/asyncron/backend.py +++ /dev/null @@ -1,162 +0,0 @@ - -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/backend/__init__.py b/asyncron/backend/__init__.py new file mode 100644 index 0000000..e1f6385 --- /dev/null +++ b/asyncron/backend/__init__.py @@ -0,0 +1,8 @@ +from .backend import AsyncronBackend +from .decorators import scheduled_task + + +__all__ = [ + AsyncronBackend, + scheduled_task +] diff --git a/asyncron/backend/backend.py b/asyncron/backend/backend.py new file mode 100644 index 0000000..94c690c --- /dev/null +++ b/asyncron/backend/backend.py @@ -0,0 +1,93 @@ +from django.conf import settings +from django.tasks.backends.base import BaseTaskBackend +from django.tasks.exceptions import InvalidTask +from django.tasks.base import ( + DEFAULT_TASK_QUEUE_NAME, + DEFAULT_TASK_PRIORITY, + TASK_MAX_PRIORITY, + TASK_MIN_PRIORITY, + TaskResultStatus +) + +from django.utils import timezone +from django.utils.inspect import get_func_args, is_module_level_function + +from asgiref.sync import sync_to_async +from inspect import iscoroutinefunction + +from ..utils import get_caller_name +from .task import AsyncronTask + + +class AsyncronBackend( BaseTaskBackend ): + task_class = AsyncronTask + + supports_defer = True + 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 validate_task( self, task ): + """ + Determine whether the provided Task can be executed by the backend. + changed version of: django/tasks/backends/base.py + """ + if not is_module_level_function( task.func ): + raise InvalidTask("Task function must be defined at a module level.") + + if not iscoroutinefunction( task.func ): + raise InvalidTask("Backend currently does not support sync Tasks.") + + task_func_args = get_func_args( task.func ) + if task.takes_context and ( + not task_func_args or task_func_args[0] not in ("context", "trace") + ): + raise InvalidTask( + "Task takes context but does not have a first argument of 'context' or 'trace'." + ) + + if task.priority != DEFAULT_TASK_PRIORITY: + raise InvalidTask("Backend does not support setting priority of tasks yet.") + + if not self.supports_defer and task.run_after is not None: + raise InvalidTask("Backend does not support run_after.") + + if ( + settings.USE_TZ + and task.run_after is not None + and not timezone.is_aware(task.run_after) + ): + raise InvalidTask("run_after must be an aware datetime.") + + if self.queues and task.queue_name not in self.queues: + raise InvalidTask(f"Queue '{task.queue_name}' is not valid for backend.") + + + def enqueue( self, task, args, kwargs ): + self.validate_task( task ) + + trace = task.as_trace() + trace.set_status( TaskResultStatus.READY, f"Sync enqueue from: {get_caller_name()}" ) + trace.args = args + trace.kwargs = kwargs + trace.save() + return trace.task_result + + async def aenqueue( self, task, args, kwargs ): + self.validate_task( task ) + + trace = task.as_trace() + trace.set_status( TaskResultStatus.READY, f"Async enqueue from: {get_caller_name()}" ) + trace.args = args + trace.kwargs = kwargs + await trace.asave() + return trace diff --git a/asyncron/backend/decorators.py b/asyncron/backend/decorators.py new file mode 100644 index 0000000..6c76c70 --- /dev/null +++ b/asyncron/backend/decorators.py @@ -0,0 +1,48 @@ +from django.tasks.base import DEFAULT_TASK_PRIORITY, DEFAULT_TASK_QUEUE_NAME, DEFAULT_TASK_BACKEND_ALIAS +from django.tasks import task_backends +from asyncron.models import TaskSchedule + +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 diff --git a/asyncron/backend/task.py b/asyncron/backend/task.py new file mode 100644 index 0000000..9759589 --- /dev/null +++ b/asyncron/backend/task.py @@ -0,0 +1,27 @@ +from django.tasks.base import TaskResultStatus +from django.tasks import Task +from django.utils import timezone + +from asyncron.models import Worker, Trace +import dataclasses + +@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 ) + + def allowed_workers( self ): + qs = Worker.objects.all() + if self.is_sensitive: qs = qs.filter( is_standalone = True ) + return qs + + def as_trace( self ): + return Trace( + task_path = self.module_path, + scheduled_datetime = timezone.now() + ( self.run_after or timezone.timedelta( seconds = 0 ) ), + ) diff --git a/asyncron/migrations/0004_remove_trace_status_description_trace_lifetime.py b/asyncron/migrations/0004_remove_trace_status_description_trace_lifetime.py new file mode 100644 index 0000000..377034f --- /dev/null +++ b/asyncron/migrations/0004_remove_trace_status_description_trace_lifetime.py @@ -0,0 +1,22 @@ +# Generated by Django 6.1 on 2026-08-27 09:40 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('asyncron', '0003_alter_trace_scheduled_datetime'), + ] + + operations = [ + migrations.RemoveField( + model_name='trace', + name='status_description', + ), + migrations.AddField( + model_name='trace', + name='lifetime', + field=models.JSONField(blank=True, default=dict), + ), + ] diff --git a/asyncron/migrations/0005_rename_prune_failed_after_taskschedule_prune_failed_over_and_more.py b/asyncron/migrations/0005_rename_prune_failed_after_taskschedule_prune_failed_over_and_more.py new file mode 100644 index 0000000..710233a --- /dev/null +++ b/asyncron/migrations/0005_rename_prune_failed_after_taskschedule_prune_failed_over_and_more.py @@ -0,0 +1,23 @@ +# Generated by Django 6.1 on 2026-08-27 12:04 + +from django.db import migrations + + +class Migration(migrations.Migration): + + dependencies = [ + ('asyncron', '0004_remove_trace_status_description_trace_lifetime'), + ] + + operations = [ + migrations.RenameField( + model_name='taskschedule', + old_name='prune_failed_after', + new_name='prune_failed_over', + ), + migrations.RenameField( + model_name='taskschedule', + old_name='prune_success_after', + new_name='prune_success_over', + ), + ] diff --git a/asyncron/migrations/0006_remove_worker_unique_scheduler_and_more.py b/asyncron/migrations/0006_remove_worker_unique_scheduler_and_more.py new file mode 100644 index 0000000..a559968 --- /dev/null +++ b/asyncron/migrations/0006_remove_worker_unique_scheduler_and_more.py @@ -0,0 +1,26 @@ +# Generated by Django 6.1 on 2026-08-27 13:10 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('asyncron', '0005_rename_prune_failed_after_taskschedule_prune_failed_over_and_more'), + ] + + operations = [ + migrations.RemoveConstraint( + model_name='worker', + name='unique_scheduler', + ), + migrations.RenameField( + model_name='worker', + old_name='is_scheduler', + new_name='is_coordinator', + ), + migrations.AddConstraint( + model_name='worker', + constraint=models.UniqueConstraint(condition=models.Q(('is_coordinator', True)), fields=('is_coordinator',), name='unique_coordinator'), + ), + ] diff --git a/asyncron/migrations/0007_worker_last_activity.py b/asyncron/migrations/0007_worker_last_activity.py new file mode 100644 index 0000000..25035a3 --- /dev/null +++ b/asyncron/migrations/0007_worker_last_activity.py @@ -0,0 +1,18 @@ +# Generated by Django 6.1 on 2026-08-27 13:21 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('asyncron', '0006_remove_worker_unique_scheduler_and_more'), + ] + + operations = [ + migrations.AddField( + model_name='worker', + name='last_activity', + field=models.DateTimeField(blank=True, null=True), + ), + ] diff --git a/asyncron/models.py b/asyncron/models.py index 7ff54b0..33fb2b2 100644 --- a/asyncron/models.py +++ b/asyncron/models.py @@ -1,7 +1,7 @@ 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.base import TaskResult, TaskResultStatus from django.tasks import task_backends from .base.models import BaseModel @@ -12,7 +12,7 @@ import random, uuid import asyncio import os, threading -# Create your models here. + class Worker( BaseModel ): id = models.UUIDField( primary_key = True, default = uuid.uuid4, editable = False ) @@ -22,22 +22,19 @@ class Worker( BaseModel ): creation_datetime = models.DateTimeField( auto_now_add = True ) is_standalone = models.BooleanField( default = False ) - is_scheduler = models.BooleanField( default = False ) - + is_coordinator = models.BooleanField( default = False ) #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") + last_activity = models.DateTimeField( null = True, blank = True ) - def is_in_main( self ): #not cross process - return self.thread_id == threading.main_thread().ident + def __str__( self ): return f"P{self.process_id}W{self.thread_id}" class Meta: constraints = [ - UniqueConstraint( fields = ['is_scheduler'], condition = Q( is_scheduler = True ), name='unique_scheduler'), + UniqueConstraint( fields = ['is_coordinator'], condition = Q( is_coordinator = True ), name = 'unique_coordinator' ), ] + class TaskSchedule( BaseModel ): name = models.CharField( default = "default", max_length = 200 ) @@ -94,8 +91,8 @@ class TaskSchedule( BaseModel ): "+": "E", }[match.group(1)] - prune_success_after = models.IntegerField( default = 10 ) - prune_failed_after = models.IntegerField( default = 1000 ) + prune_success_over = models.IntegerField( default = 10 ) + prune_failed_over = models.IntegerField( default = 1000 ) class Meta: unique_together = [ @@ -144,19 +141,13 @@ class Trace( BaseModel ): 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_description = desc args = models.JSONField( default = list, blank = True ) kwargs = models.JSONField( default = dict, blank = True ) @@ -166,16 +157,45 @@ class Trace( BaseModel ): started_datetime = models.DateTimeField( null = True, blank = True ) finished_datetime = models.DateTimeField( null = True, blank = True ) + lifetime = models.JSONField( default = dict, blank = True ) stdout = models.TextField( default = "", blank = True ) exception_class_path = models.TextField( null = True ) traceback = models.TextField( null = True ) + def set_status( self, status, desc = "" ): + + if isinstance( status, BaseException ): + exception_type = type(status) + self.status = TaskResultStatus.FAILED + self.lifetime[self.status] = desc + self.traceback = traceback.format_exception( status ) + self.exception_class_path = f"{exception_type.__module__}.{exception_type.__qualname__}" + return + + assert status in TaskResultStatus, f"Unkown status '{status}' provided for trace!" + self.status = status + self.lifetime[self.status] = desc + 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 ): + + @property + def task( self ): return task_backends['default'].TASKS[self.task_path] + + @property + def task_error( self ): + if not self.exception_class_path: return + return TaskError( + exception_class_path = self.exception_class_path, + traceback = self.traceback, + ) + + @property + def task_result( self ): + error = self.task_error return TaskResult( task = self.task, id = self.id, @@ -187,8 +207,8 @@ class Trace( BaseModel ): args = self.args, kwargs = self.kwargs, backend = self.task.backend, - errors = [], - worker_ids = [], + errors = [error] if error else [], + worker_ids = [self.worker_id] if self.worker_id else [], ) class Meta: @@ -201,57 +221,9 @@ class Trace( BaseModel ): ] - async def start( self ): - await self.eval_related('schedule') - schedule = self.schedule - task = self.task - - #self.started_datetime = timezone.now() - self.finished_datetime = None - self.return_value = None - self.exception_class_path = None - self.traceback = None - - #Runtime Bits - loop = asyncio.get_running_loop() - self.new_print = asyncio.Event() - self.commit_on_new_print_task = loop.create_task( self.commit_on_new_print() ) - - try: - async with asyncio.timeout( None ) as tmcm: - - if schedule and schedule.interval is not None and schedule.timeout: - tmcm.reschedule( loop.time() + self.schedule.timeout.total_seconds() ) - - if task.takes_context: - output = await task.func( self, *self.args, **self.kwargs ) - else: - output = await task.func( *self.args, **self.kwargs ) - - 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( TaskResultStatus.FAILED, f"Exception: {e}" ) - self.traceback = traceback.format_exc() - self.exception_class_path = str(e) #.__qualname__ - - else: - self.set_status( TaskResultStatus.SUCCESSFUL ) - self.return_value = output - - finally: - self.commit_on_new_print_task.cancel() - - self.finished_datetime = timezone.now() - await self.asave() - - - ### Runtime methods for self aware tasks + ### Runtime methods for context aware tasks async def commit_on_new_print( self ): - while True: + while True: #new_print event needs to be created in the worker. await self.new_print.wait() await self.asave( update_fields = ['stdout'] ) self.new_print.clear() diff --git a/asyncron/utils.py b/asyncron/utils.py index e833384..18bc914 100644 --- a/asyncron/utils.py +++ b/asyncron/utils.py @@ -1,4 +1,5 @@ import functools +import inspect import os import re @@ -28,6 +29,10 @@ def is_proc_alive( pid ): 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 +def get_caller_name(): + #filename, line_number, function_name, lines, index + try: return inspect.getframeinfo(inspect.currentframe().f_back.f_back.f_back)[2] + except: return "N/A" # Regular expression pattern with named groups for "1w2d5h30m10s500ms1000us" without spaces TIMEDELTA_PATTERN = re.compile( @@ -42,6 +47,26 @@ TIMEDELTA_PATTERN = re.compile( ) +import json +from django.db.models import Func, Value, TextField, JSONField +from django.contrib.postgres.fields import ArrayField +class JSONSet(Func): #Vibe Coded Class + """Set a single key (or nested path) in a JSONField via jsonb_set().""" + function = "jsonb_set" + output_field = JSONField() + + def __init__(self, field_name, path, value, create_missing=True): + # path e.g. ['key'] or ['nested', 'key'] + super().__init__( + field_name, + Value(path, output_field=ArrayField(TextField())), + Value(json.dumps(value)), + Value(create_missing), + ) + + + + #Django keeps giving: exception=OperationalError('the connection is closed') from django.db.utils import OperationalError def retry_on_db_error( f ): diff --git a/asyncron/workers.py b/asyncron/workers.py index fbc0ea5..b48f4be 100644 --- a/asyncron/workers.py +++ b/asyncron/workers.py @@ -17,7 +17,7 @@ import collections, functools import random import humanize -from .utils import retry_on_db_error, ignore_on_db_error +from .utils import is_proc_alive, JSONSet from .asynctools import AsyncOriented from .singleton import Singleton from .models import Worker, TaskSchedule, Trace @@ -66,6 +66,7 @@ class AsyncronWorker( Singleton, AsyncOriented ): #Parallelism self.thread = None #Once the worker starts, it'll be populated + self.compatible_task_paths = [] #We have to conver and use it in database queries a lot, might as well not use a set. self.clearing_dead_workers = False self.watching_models = collections.defaultdict( set ) # Model -> Set of key name of the tasks @@ -158,7 +159,8 @@ class AsyncronWorker( Singleton, AsyncOriented ): async def cleanup( self ): try: count = await Trace.objects.filter( status = TaskResultStatus.RUNNING, worker = self.model ).aupdate( - status_description = "Worker died during execution", + lifetime = JSONSet( "lifetime", [TaskResultStatus.FAILED], "Worker died during execution" ), + finished_datetime = timezone.now(), status = TaskResultStatus.FAILED, worker = None ) await self.model.adelete() @@ -170,6 +172,7 @@ class AsyncronWorker( Singleton, AsyncOriented ): await super().cleanup() def start_working( self, is_standalone = False ): + with asyncio.Runner() as runner: runner.run( self.startup() ) @@ -181,13 +184,16 @@ class AsyncronWorker( Singleton, AsyncOriented ): self.log.debug(f"DB Connection ready, continuing...") self.model.is_standalone = is_standalone + for path, task in self.backend.TASKS.items(): + if task.is_sensitive and not is_standalone: continue + self.compatible_task_paths.append( path ) + runner.run( self.model.asave() ) runner.run( self.sync_task_schedules() ) + #self.create_task( self.consume_task_reason_jobs_queue() ) - - #self.create_task( self.maintain_scheduler() ) - + self.create_task( self.maintain_coordinator() ) self.create_task( self.run_tasks_on_schedule() ) self.attach_django_signals() @@ -227,7 +233,6 @@ class AsyncronWorker( Singleton, AsyncOriented ): try: await ts.asave() except IntegrityError: pass - async def run_tasks_on_schedule( self ): self.check_interval = 0 @@ -248,25 +253,20 @@ class AsyncronWorker( Singleton, AsyncOriented ): except IntegrityError: continue #else: print("Created a new trace:", trace, humanize.naturaltime( trace.scheduled_datetime ) ) + early_seconds = 1 + self.check_interval * ( 1 + random.random() ) Ts = Trace.objects.filter( + task_path__in = self.compatible_task_paths, status = TaskResultStatus.READY, - scheduled_datetime__lte = timezone.now(), + scheduled_datetime__lte = timezone.now() + timezone.timedelta( seconds = early_seconds ), worker = None - ).order_by('-scheduled_datetime') + ).prefetch_related("schedule").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 - ) + locked = await Trace.objects.filter( id = trace.id, worker = None ).aupdate( 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 )() @@ -275,27 +275,84 @@ class AsyncronWorker( Singleton, AsyncOriented ): self.is_work_over.set() async def start_trace_on_time( self, trace ): + await trace.eval_related("schedule") if timezone.now() < trace.scheduled_datetime: #If this is a periodic task - await asyncio.sleep( ( timezone.now() - trace.scheduled_datetime ).total_seconds() ) + await asyncio.sleep( ( trace.scheduled_datetime - timezone.now() ).total_seconds() ) await trace.arefresh_from_db() - await trace.start() + await self.start_trace_now( trace ) - trace.worker = None - await trace.asave( update_fields = ['worker'] ) + #If this was not a scheduled tasked, return here, other wise, go prune old logs. + if not trace.schedule: return + + prune_over = -1 + if trace.status == TaskResultStatus.SUCCESSFUL: + prune_over = trace.schedule.prune_success_over + if trace.status == TaskResultStatus.FAILED: + prune_over = trace.schedule.prune_failed_over + + if prune_over < 0: return #nothing to prune - 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') + prune_candidates = Trace.objects.filter( + schedule_id = trace.schedule_id, + status = trace.status, #Prune only the same types of traces + enqueued_datetime__lt = trace.enqueued_datetime, #Do not prune traces that got added after this trace + prune_protected = False, worker = None, + ).order_by('-enqueued_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 ) + await prune_candidates.exclude( + id__in = prune_candidates[:prune_over].values_list( 'id', flat = True ) ).adelete() + + async def start_trace_now( self, trace, info = None ): + await trace.eval_related('schedule') + schedule = trace.schedule + task = trace.task + + trace.status = TaskResultStatus.RUNNING + if info: trace.lifetime[trace.status] = info + + trace.started_datetime = timezone.now() + trace.finished_datetime = None + + trace.return_value = None + trace.exception_class_path = None + trace.traceback = None + await trace.asave() + + trace.new_print = asyncio.Event() + trace.commit_on_new_print_task = self.loop.create_task( trace.commit_on_new_print() ) + + try: + async with asyncio.timeout( None ) as tmcm: + + if schedule and schedule.interval is not None and schedule.timeout: + tmcm.reschedule( self.loop.time() + trace.schedule.timeout.total_seconds() ) + + if task.takes_context: + trace.return_value = await task.func( trace, *trace.args, **trace.kwargs ) + else: + trace.return_value = await task.func( *trace.args, **trace.kwargs ) + + except TimeoutError as e: + trace.set_status( e, f"Timed out" ) + + except Exception as e: + trace.set_status( e, f"Exception: {e}" ) + + else: + trace.set_status( TaskResultStatus.SUCCESSFUL ) + + finally: + trace.commit_on_new_print_task.cancel() + del trace.commit_on_new_print_task + del trace.new_print + + trace.finished_datetime = timezone.now() + trace.worker = None + await trace.asave() @@ -373,81 +430,63 @@ class AsyncronWorker( Singleton, AsyncOriented ): except: await asyncio.sleep( 1 ) - async def maintain_scheduler( self ): + async def maintain_coordinator( self ): """ - Make sure at least one process is managing the scheduled tasks. + Make sure at exactly one process is managing the scheduled tasks, and asycron's housekeeping """ - from .models import Worker, TaskSchedule, Trace - is_current_master = False - next_overtake_attempt = time.time() + 1 + random.random() * 5 - loop_wait = 5 + is_current_coordinator = False + this_worker_as_queryset = Worker.objects.filter( id = self.model.id ) + loop_wait_seconds = 5 - MyQs = Worker.objects.filter( id = self.model.id ) - while await MyQs.aupdate( last_activity = timezone.now() ): + while await this_worker_as_queryset.aupdate( last_activity = timezone.now() ): try: - await Worker.objects.filter( is_scheduler = False ).aupdate( is_scheduler = models.Q(id = self.model.id) ) + #Try to set yourself as the coordinator, if no other workers are currently coordinating + await Worker.objects.filter( is_coordinator = False ).aupdate( is_coordinator = models.Q( id = self.model.id ) ) - except IntegrityError: # I'm not master! - loop_wait = 5 + random.random() * 15 + except IntegrityError: # Not Coordinating + loop_wait_seconds = 5 + 15 * random.random() - if is_current_master: self.log.info(f"[W{self.model.id}] No longer master.") - is_current_master = False + if is_current_coordinator: self.log.info(f"[W{self.model.id}] No longer coordinating.") + is_current_coordinator = False - #Deletes dead masters every now and then! - if next_overtake_attempt <= time.time(): - next_overtake_attempt = time.time() + 60 - took_master = False + #Remove the coordinator tag of any worker that has not been coordinating for 2 minutes! + coordinator_removed = await Worker.objects.filter( is_coordinator = True ).filter( + models.Q( last_activity = None ) | + models.Q( last_activity__lte = timezone.now() - timezone.timedelta( minutes = 2 ) ) + ).aupdate( is_coordinator = False ) + if coordinator_removed: continue #Try being the coordinator right now! - 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: #I am the coordinator! + loop_wait_seconds = 2 + 3 * random.random() - else: - 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_scheduler = False ) + if not is_current_coordinator: self.log.info(f"[W{self.model.id}] Running as coordinator.") + is_current_coordinator = True + + #Delete inactive workers if the proccess does not exist. + inactive_workers_queryset = Worker.objects.filter( + models.Q( last_activity = None ) | + models.Q( last_activity__lte = timezone.now() - timezone.timedelta( seconds = 30 ) ) + ) + async for worker in inactive_workers_queryset: + if not is_proc_alive( worker.process_id ): + await worker.adelete() + + #Set Orphaned tasks as FAILED + await Trace.objects.filter( + status = TaskResultStatus.RUNNING, + worker = None, + ).aupdate( + status = TaskResultStatus.FAILED, + lifetime = JSONSet( "lifetime", [TaskResultStatus.FAILED], "Worker lost during execution" ), + ) + + await asyncio.sleep( loop_wait_seconds ) + + self.log.warning(f"[W{self.model.id}] Cannot find the corresponding database model, exiting maintain_coordinator...") - else: #I am Master! - loop_wait = 2 + random.random() * 3 - - if not is_current_master: self.log.info(f"[W{self.model.id}] Running as master.") - is_current_master = True - - if not self.clearing_dead_workers: - self.create_task( self.clear_dead_workers() ) - - await self.sync_tasks() - await self.clear_orphaned_traces() - - if loop_wait: - await asyncio.sleep( loop_wait ) - - else: - self.log.warning(f"[W{self.model.id}] Worker in Master Loop Cannot find it's corresponding database model!") - - async def clear_orphaned_traces( self ): - from .models import Worker, Task, Trace - await Trace.objects.filter( worker = None, status = "RUNNING" ).adelete() - - async def clear_dead_workers( self ): - self.clearing_dead_workers = True - from .models import Worker, Task, Trace - await Worker.objects.filter( - last_activity__lte = timezone.now() - timezone.timedelta( seconds = 30 ), - in_grace = False - ).aupdate( in_grace = True ) - - async for worker in Worker.objects.filter( in_grace = False, last_activity = None ): - if not await sync_to_async( worker.is_proc_alive )(): - await worker.adelete() - - await asyncio.sleep( 30 ) - await Worker.objects.filter( in_grace = True ).adelete() - self.clearing_dead_workers = False @@ -458,67 +497,6 @@ class AsyncronWorker( Singleton, AsyncOriented ): await self.start_task_now( task, reason ) self.task_reason_jobs_queue.task_done() - async def check_scheduled( self ): - from .models import Task, Trace - - #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_standalone else "R" ) - - async for task in Ts: - trace = task.new_trace() - await trace.reschedule( reason = "Auto Scheduled" ) - - try: - locked = await Task.objects.filter( id = task.id ).filter( - models.Q(worker_lock = None) | - models.Q(worker_lock = self.model) #This is incase the lock has been aquired for some reason before. - ).aupdate( worker_lock = self.model ) - except IntegrityError: - locked = False - - if locked: - await trace.asave() - 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 = 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 = self.model ) - except IntegrityError: count = 0 - if not count: continue #Lost the race condition to another worker. - - self.create_task( self.start_trace_on_time( trace ) ) - - async def check_services( self ): - from .models import Task, Trace - - #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_standalone else "R" ) - - async for task in Ts: - - try: - locked = await Task.objects.filter( id = task.id ).filter( - models.Q(worker_lock = None) | - models.Q(worker_lock = self.model) #This is incase the lock has been aquired for some reason before. - ).aupdate( worker_lock = self.model ) - except IntegrityError: continue - if not locked: continue - - trace = task.new_trace() - trace.set_status( "W", "Waiting to start the service ASAP." ) - trace.worker = self.model - await trace.asave() - - #self.running_service_tasks[task.id] = - self.create_task( self.start_trace_on_time( trace ) ) - await Task.objects.filter( id = task.id, worker_lock = self.model ).aupdate( worker_lock = None ) async def start_task_now( self, task, reason = "" ): trace = task.new_trace()