moved lots of things around, and returned most features using cleaner code

This commit is contained in:
Queue A
2026-08-27 18:08:10 +02:00
parent 8028c2df76
commit 2dac1e31f8
13 changed files with 568 additions and 467 deletions
+104 -81
View File
@@ -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 )
if took.total_seconds() < 1:
bits.append(f"Ran and finished {humanize.naturaltime(obj.finished_datetime)}")
else:
return "+ In " + humanize.naturaltime( obj.scheduled_datetime )
execution.admin_order_field = 'scheduled_datetime'
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 )
-162
View File
@@ -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
+8
View File
@@ -0,0 +1,8 @@
from .backend import AsyncronBackend
from .decorators import scheduled_task
__all__ = [
AsyncronBackend,
scheduled_task
]
+93
View File
@@ -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
+48
View File
@@ -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
+27
View File
@@ -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 ) ),
)
@@ -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),
),
]
@@ -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',
),
]
@@ -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'),
),
]
@@ -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),
),
]
+43 -71
View File
@@ -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()
+25
View File
@@ -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 ):
+127 -149
View File
@@ -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
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_scheduler = True ).filter(
#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_scheduler = False )
).aupdate( is_coordinator = False )
if coordinator_removed: continue #Try being the coordinator right now!
else: #I am the coordinator!
loop_wait_seconds = 2 + 3 * random.random()
else: #I am Master!
loop_wait = 2 + random.random() * 3
if not is_current_coordinator: self.log.info(f"[W{self.model.id}] Running as coordinator.")
is_current_coordinator = True
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 )():
#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()
await asyncio.sleep( 30 )
await Worker.objects.filter( in_grace = True ).adelete()
self.clearing_dead_workers = False
#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...")
@@ -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()