refactored to work with django tasks as a backend

This commit is contained in:
Queue A
2026-08-27 08:55:08 +02:00
parent ba963d6367
commit 8028c2df76
15 changed files with 663 additions and 553 deletions
+60 -54
View File
@@ -1,9 +1,11 @@
from django.contrib import admin from django.contrib import admin
from django.utils import timezone from django.utils import timezone
from django.db.models import F from django.db.models import F
from django.tasks.base import TaskResultStatus
from .base.admin import BaseModelAdmin 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 os
import asyncio import asyncio
@@ -11,18 +13,18 @@ import humanize
@admin.register( Worker ) @admin.register( Worker )
class WorkerAdmin( BaseModelAdmin ): class WorkerAdmin( BaseModelAdmin ):
order = 4 order = 1
list_display = 'pid', 'thread_id', 'is_robust', 'is_master', 'is_running', 'health', list_display = 'process_id', 'thread_id', 'is_standalone', 'is_scheduler', 'is_running', # 'health',
def has_add_permission( self, request, obj = None ): return False 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 is_running.boolean = True
def health( self, obj ): #def health( self, obj ):
return (f"In Grace " if obj.in_grace else "") + humanize.naturaltime( obj.last_activity ) # 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" title = "app"
parameter_name = 'app_groups' parameter_name = 'app_groups'
@@ -30,33 +32,36 @@ class TaskAppFilter( admin.SimpleListFilter ): #TODO: Also check if it's an actu
return ( return (
( a.lower(), a ) ( a.lower(), a )
for a in sorted({ for a in sorted({
task.name.split(".", 1)[0] ts.task_path.split(".", 1)[0]
for task in Task.objects.all() for ts in TaskSchedule.objects.all()
}) if a }) if a
) )
def queryset( self, request, queryset ): def queryset( self, request, queryset ):
q = self.value() q = self.value()
if not q: return queryset if not q: return queryset
return queryset.filter( name__istartswith = q ) return queryset.filter( task_path__istartswith = q )
@admin.register( TaskSchedule )
@admin.register( Task ) class TaskScheduleAdmin( BaseModelAdmin ):
class TaskAdmin( BaseModelAdmin ): order = 2
order = 1 list_display = 'schedule_name', 'timeout', 'type', 'jitter', 'logged', 'last_execution', 'scheduled', 'is_enabled',
list_display = 'name', 'timeout', 'gracetime', 'jitter', 'type', 'worker_type', 'logged', 'last_execution', 'scheduled'
list_filter = TaskAppFilter, list_filter = TaskAppFilter,
fields = ["name", "description", "type", "on_model_change", "jitter", "self_aware"]
actions = 'schedule_execution', 'execution_now', 'delete_script_missing', 'update_details', fields = "name", "description", "type", #"on_model_change", "jitter",
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 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 ): def description( self, obj ):
try: try:
return obj.registered_tasks[obj.name].__doc__.strip("\n") return obj.task.func.__doc__.strip("\n")
except: except:
raise
return "N/A" return "N/A"
def jitter( self, obj ): def jitter( self, obj ):
@@ -69,38 +74,44 @@ class TaskAdmin( BaseModelAdmin ):
def type( self, obj ): def type( self, obj ):
results = [] results = []
if obj.timeout is None:
results.append( "Service" )
if obj.interval: if obj.interval:
delta = humanize.naturaldelta( obj.interval ) delta = humanize.naturaldelta( obj.interval )
delta = delta.replace("an ", "1 ").replace("a ", "1 ") 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: else:
results.append("Script Missing!") 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 ): def on_model_change( self, obj ):
try: return ", ".join( f"{m.__module__}.{m.__name__}" for m in obj.registered_tasks[obj.name].watching_models ) try: return ", ".join( f"{m.__module__}.{m.__name__}" for m in obj.registered_tasks[obj.name].watching_models )
except: return "N/A" except: return "N/A"
def logged( self, obj ): 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" if count == 1: return "1 trace"
return f"{count} traces" return f"{count} traces"
def last_execution( self, obj ): 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" 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 ): def scheduled( self, obj ):
return obj.trace_set.filter( status = "S" ).exists() return obj.trace_set.filter( status = TaskResultStatus.READY ).exists()
scheduled.boolean = True 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" ) @admin.action( description = "(Re)Schedule an execution for periodic tasks" )
def schedule_execution( self, request, qs ): def schedule_execution( self, request, qs ):
trace_ids = set() trace_ids = set()
@@ -132,12 +143,8 @@ class TaskAdmin( BaseModelAdmin ):
if task.name not in task.registered_tasks: if task.name not in task.registered_tasks:
task.delete() 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): class TraceAppFilter(admin.SimpleListFilter):
title = "app" title = "app"
@@ -178,33 +185,32 @@ class TraceNameFilter(admin.SimpleListFilter):
@admin.register( Trace ) @admin.register( Trace )
class TraceAdmin( BaseModelAdmin ): class TraceAdmin( BaseModelAdmin ):
order = 2 order = 3
list_display = 'task', 'execution', 'state', 'worker_lock' list_display = 'scheduled_datetime', 'task_path', 'execution', 'status', 'worker'
list_filter = TraceAppFilter, TraceNameFilter, 'task__worker_type', 'status', 'status_reason', list_filter = 'status', 'status_description', #'task__worker_type', TraceAppFilter, TraceNameFilter,
ordering = F('scheduled_datetime').desc(nulls_last=True), ordering = "-enqueued_datetime",
#readonly_fields = [ f.name for f in Trace._meta.fields ] #readonly_fields = [ f.name for f in Trace._meta.fields ]
def has_add_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 has_change_permission( self, request, obj = None ): return False
def execution( self, obj ): def execution( self, obj ):
if obj.last_run_datetime: if obj.finished_datetime:
return "- Ran " + humanize.naturaltime( obj.last_run_datetime ) return "- Ran " + humanize.naturaltime( obj.finished_datetime )
if obj.scheduled_datetime: if obj.scheduled_datetime < timezone.now():
if obj.scheduled_datetime < timezone.now(): return "- Should've run " + humanize.naturaltime( obj.scheduled_datetime )
return "- Should've run " + humanize.naturaltime( obj.scheduled_datetime ) else:
else: return "+ In " + humanize.naturaltime( obj.scheduled_datetime )
return "+ In " + humanize.naturaltime( obj.scheduled_datetime )
return "Never"
execution.admin_order_field = 'scheduled_datetime' execution.admin_order_field = 'scheduled_datetime'
def state( self, obj ): 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' state.admin_order_field = 'status'
actions = 'reschedule_to_now', actions = 'reschedule_to_now',
@admin.action( description = "Reschedule to run now" ) @admin.action( description = "Reschedule to run now" )
def reschedule_to_now( self, request, qs ): def reschedule_to_now( self, request, qs ):
+162
View File
@@ -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
@@ -12,7 +12,7 @@ from django.core.management.base import BaseCommand, CommandError
from django.conf import settings from django.conf import settings
from asyncron.workers import AsyncronWorker from asyncron.workers import AsyncronWorker
from asyncron.models import Task from asyncron.models import TaskSchedule
class bcolors: class bcolors:
HEADER = '\033[95m' HEADER = '\033[95m'
@@ -32,8 +32,8 @@ class Command(BaseCommand):
AsyncronWorker.IS_ACTIVE = True AsyncronWorker.IS_ACTIVE = True
while True: while True:
worker = AsyncronWorker() worker = AsyncronWorker()
worker.is_db_ready.set()
print( "Starting:", worker ) print( "Starting:", worker )
try: try:
+37 -52
View File
@@ -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 datetime
import django.db.models.deletion import django.db.models.deletion
import posix
import uuid
from django.db import migrations, models from django.db import migrations, models
@@ -10,83 +13,65 @@ class Migration(migrations.Migration):
initial = True initial = True
dependencies = [ dependencies = [
('contenttypes', '0002_remove_content_type_name'),
] ]
operations = [ operations = [
migrations.CreateModel( migrations.CreateModel(
name='Worker', name='TaskSchedule',
fields=[ fields=[
('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), ('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')),
('pid', models.IntegerField()), ('name', models.CharField(default='default', max_length=200)),
('thread_id', models.PositiveBigIntegerField()), ('is_enabled', models.BooleanField(default=True)),
('is_robust', models.BooleanField(default=False)), ('task_path', models.TextField()),
('is_master', models.BooleanField(default=False)), ('args', models.JSONField(blank=True, default=list)),
('in_grace', models.BooleanField(default=False)), ('kwargs', models.JSONField(blank=True, default=dict)),
('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))),
('interval', models.DurationField(blank=True, null=True)), ('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_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)), ('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={ options={
'abstract': False, 'unique_together': {('name', 'task_path')},
}, },
), ),
migrations.CreateModel( migrations.CreateModel(
name='Metadata', name='Worker',
fields=[ fields=[
('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), ('id', models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)),
('model_id', models.PositiveIntegerField()), ('process_id', models.IntegerField(default=posix.getpid)),
('name', models.CharField(max_length=256)), ('thread_id', models.PositiveBigIntegerField(default=_thread.get_ident)),
('data', models.JSONField(blank=True, null=True)), ('creation_datetime', models.DateTimeField(auto_now_add=True)),
('expiration_datetime', models.DateTimeField(blank=True, null=True)), ('is_standalone', models.BooleanField(default=False)),
('model_type', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, to='contenttypes.contenttype')), ('is_scheduler', models.BooleanField(default=False)),
], ],
options={ options={
'verbose_name': 'Metadata', 'constraints': [models.UniqueConstraint(condition=models.Q(('is_scheduler', True)), fields=('is_scheduler',), name='unique_scheduler')],
'verbose_name_plural': 'Metadata',
'indexes': [models.Index(fields=['model_type', 'model_id'], name='asyncron_me_model_t_d92186_idx')],
}, },
), ),
migrations.CreateModel( migrations.CreateModel(
name='Trace', name='Trace',
fields=[ fields=[
('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), ('id', models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)),
('status_reason', models.TextField(blank=True, default='')), ('task_path', models.TextField()),
('status', models.CharField(choices=[('S', 'Scheduled'), ('W', 'Waiting'), ('R', 'Running'), ('P', 'Paused'), ('C', 'Completed'), ('A', 'Aborted'), ('E', 'Error')], default='S', max_length=1)), ('status_description', models.TextField(blank=True, default='')),
('scheduled_datetime', models.DateTimeField(blank=True, null=True)), ('status', models.CharField(choices=[('READY', 'Ready'), ('RUNNING', 'Running'), ('FAILED', 'Failed'), ('SUCCESSFUL', 'Successful')], default='READY', max_length=10)),
('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)),
('args', models.JSONField(blank=True, default=list)), ('args', models.JSONField(blank=True, default=list)),
('kwargs', models.JSONField(blank=True, default=dict)), ('kwargs', models.JSONField(blank=True, default=dict)),
('stdout', models.TextField(blank=True, null=True)), ('enqueued_datetime', models.DateTimeField(auto_now_add=True)),
('stderr', models.TextField(blank=True, null=True)), ('scheduled_datetime', models.DateTimeField(blank=True, null=True)),
('returned', models.JSONField(blank=True, null=True)), ('started_datetime', models.DateTimeField(blank=True, null=True)),
('task', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, to='asyncron.task')), ('finished_datetime', models.DateTimeField(blank=True, null=True)),
('worker_lock', models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, to='asyncron.worker')), ('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={ 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')],
}, },
), ),
] ]
@@ -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."),
),
]
+19
View File
@@ -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'),
),
]
@@ -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."),
),
]
@@ -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,
),
]
@@ -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)),
),
]
@@ -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',
),
]
@@ -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',
),
]
+171 -152
View File
@@ -1,232 +1,251 @@
from django.utils import timezone from django.utils import timezone
from django.db import models from django.db import models
from django.db.models.constraints import UniqueConstraint, Q 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 from .base.models import BaseModel
#This is buggy, has print leakage and still not very good, from .utils import TIMEDELTA_PATTERN
#But better than nothing when a task is not self_aware or calls something that isn't
from unittest.mock import patch
import functools, traceback, io import functools, traceback, io
import random import random, uuid
import asyncio import asyncio
import os, threading
# Create your models here. # Create your models here.
from .base.models import BaseModel
class Worker( BaseModel ): class Worker( BaseModel ):
id = models.UUIDField( primary_key = True, default = uuid.uuid4, editable = False )
pid = models.IntegerField() process_id = models.IntegerField( default = os.getpid )
thread_id = models.PositiveBigIntegerField() thread_id = models.PositiveBigIntegerField( default = threading.get_ident )
is_robust = models.BooleanField( default = False ) creation_datetime = models.DateTimeField( auto_now_add = True )
is_master = 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 ) is_standalone = models.BooleanField( default = False )
is_scheduler = models.BooleanField( default = False )
#Variables with very feel good names! :) #in_grace = models.BooleanField( default = False ) #If the worker sees this as True, it should kill itself!
consumption_interval_seconds = models.IntegerField( default = 10 ) #last_activity = models.DateTimeField( null = True, blank = True )
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") 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: class Meta:
constraints = [ 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 ): class TaskSchedule( BaseModel ):
import os name = models.CharField( default = "default", max_length = 200 )
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 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 args = models.JSONField( default = list, blank = True )
worker_lock = models.ForeignKey( Worker, null = True, blank = True, on_delete = models.SET_NULL ) kwargs = models.JSONField( default = dict, blank = True )
worker_type = models.CharField( default = "A", choices = {
"A": "Any",
"R": "Robust", #Only seperate Robust workers
"D": "Dynamic", #Only on potentially reloadable workers
})
max_completed_traces = models.IntegerField( default = 10 ) #Distinguishes Periodic and Service like Tasks
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
interval = models.DurationField( null = True, blank = True ) 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_length = models.DurationField( default = timezone.timedelta( seconds = 0 ), blank = True )
jitter_pivot = models.CharField( default = "M", max_length = 1, choices = { jitter_pivot = models.CharField( default = "M", max_length = 1, choices = {
"S":"Start", "M":"Middle", "E":"End", "S":"Start", "M":"Middle", "E":"End",
}) })
def get_jitter( self ): def get_jitter( self ):
jitter = self.jitter_length * random.random() jitter = self.jitter_length * random.random()
match self.jitter_pivot: match self.jitter_pivot:
case "M": case "M": jitter -= self.jitter_length / 2
jitter -= self.jitter_length / 2 case "E": jitter *= -1
case "E":
jitter *= -1
return jitter 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 ): def __str__( self ):
type = "Callable" if self.interval is None else "Periodic" type_display = ( "Service Task" if self.type == "S" else "Periodic Task" )
mode = "Service" if self.timeout is None else "Task" short = self.task_path.rsplit('.')[-1]
short = self.name.rsplit('.')[-1] return f"{self.name} {type_display} {short}"
return " ".join([type, mode, short])
def register( self, f ): def as_trace( self ):
if not self.name: self.name = f"{f.__module__}.{f.__qualname__}" return Trace(
self.registered_tasks[self.name] = f schedule = self,
f.task = self task_path = self.task_path,
return f
def new_trace( self ): args = self.args,
trace = Trace( task_id = self.id ) kwargs = self.kwargs
trace.task = self #Less db hits )
return trace
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() now = timezone.now()
if await self.trace_set.filter( status = "W" ).aexists(): jitter_delta = self.get_jitter()
return 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 ): 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 ) task_path = models.TextField() #Path to the task function
status = models.CharField( default = "S", max_length = 1, choices = { schedule = models.ForeignKey( TaskSchedule, null = True, on_delete = models.SET_NULL )
"S":"Scheduled", @property
"W":"Waiting", def task( self ): return task_backends['default'].TASKS[self.task_path]
"R":"Running",
"P":"Paused", status_description = models.TextField( default = "", blank = True )
"C":"Completed", status = models.CharField(
"A":"Aborted", default = TaskResultStatus.READY,
"E":"Error", choices = TaskResultStatus.choices,
}) max_length = max( len(v) for v in TaskResultStatus.values ),
def set_status( self, status, reason = "" ): )
def set_status( self, status, desc = "" ):
self.status = status 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 ) args = models.JSONField( default = list, blank = True )
kwargs = models.JSONField( default = dict, 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: class Meta:
constraints = [ constraints = [
UniqueConstraint( UniqueConstraint(
fields = ['task_id'], fields = ['schedule_id'],
condition = models.Q(status = "S", scheduled_datetime = None), condition = models.Q(status = TaskResultStatus.READY),
name = "unique_unscheduled_for_task", 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 ): async def start( self ):
await self.eval_related('task') await self.eval_related('schedule')
assert self.status in "SPAWE", f"Cannot start a task that is in {self.get_status_display()} state!" schedule = self.schedule
task = self.task
self.last_run_datetime = timezone.now() #self.started_datetime = timezone.now()
self.last_end_datetime = None self.finished_datetime = None
self.returned = None self.return_value = None
self.stderr = "" self.exception_class_path = None
self.stdout = "" self.traceback = None
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()
#Runtime Bits #Runtime Bits
self.loop = asyncio.get_running_loop() loop = asyncio.get_running_loop()
self.new_print = asyncio.Event() 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: try:
async with asyncio.timeout( None ) as tmcm: async with asyncio.timeout( None ) as tmcm:
if self.task.timeout: if schedule and schedule.interval is not None and schedule.timeout:
tmcm.reschedule( self.loop.time() + self.task.timeout.total_seconds() ) tmcm.reschedule( loop.time() + self.schedule.timeout.total_seconds() )
if self.task.self_aware: if task.takes_context:
output = await func(self, *self.args, **self.kwargs ) output = await task.func( self, *self.args, **self.kwargs )
else: else:
with patch( 'builtins.print', self.print ): output = await task.func( *self.args, **self.kwargs )
output = await func( *self.args, **self.kwargs )
except TimeoutError: except TimeoutError as e:
self.set_status( "E", f"Timed out" ) self.set_status( TaskResultStatus.FAILED, f"Timed out" )
self.stderr = traceback.format_exc() self.traceback = traceback.format_exc()
self.exception_class_path = e.__qualname__
except Exception as e: except Exception as e:
self.set_status( "E", f"Exception: {e}" ) self.set_status( TaskResultStatus.FAILED, f"Exception: {e}" )
self.stderr = traceback.format_exc() self.traceback = traceback.format_exc()
self.exception_class_path = str(e) #.__qualname__
else: else:
self.set_status( "C" ) self.set_status( TaskResultStatus.SUCCESSFUL )
self.returned = output self.return_value = output
finally: finally:
self.commit_on_new_print_task.cancel() self.commit_on_new_print_task.cancel()
self.last_end_datetime = timezone.now() self.finished_datetime = timezone.now()
await self.asave() await self.asave()
+10 -43
View File
@@ -1,59 +1,26 @@
## ##
## decorators / functions to make the task calls easier ## decorators / functions to make the task calls easier
## ##
from django.utils.dateparse import parse_duration from django.utils.dateparse import parse_duration
from django.db import models from django.db import models
from django.utils import timezone from django.utils import timezone
from django.apps import apps from django.apps import apps
import re
# Regular expression pattern with named groups for "1w2d5h30m10s500ms1000us" without spaces
pattern = re.compile(
r'(\+|-)?'
r'(?:(?P<weeks>\d+)w)?'
r'(?:(?P<days>\d+)d)?'
r'(?:(?P<hours>\d+)h)?'
r'(?:(?P<minutes>\d+)m)?'
r'(?:(?P<seconds>\d+)s)?'
r'(?:(?P<milliseconds>\d+)ms)?'
r'(?:(?P<microseconds>\d+)us)?'
)
def task( *args, **kwargs ): def task( *args, **kwargs ):
from .models import Task from .backend import scheduled_task
kwargs.setdefault( 'schedule_name', 'scripted' )
jitter = kwargs.pop('jitter', "") return scheduled_task( *args, **kwargs )
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
def service( *args, **kwargs ): def service( *args, **kwargs ):
kwargs.setdefault( 'timeout', None ) from .backend import scheduled_task
kwargs.setdefault( 'worker_type', "R" ) kwargs.setdefault( 'interval', None )
kwargs.setdefault( 'self_aware', True ) kwargs.setdefault( 'is_sensitive', True )
return task( *args, **kwargs ) kwargs.setdefault( 'schedule_name', 'scripted' )
return scheduled_task( *args, **kwargs )
def run_on_model_change( *models ): def run_on_model_change( *models ):
return lambda x:x
models = [ models = [
apps.get_model(m) if isinstance(m, str) else m apps.get_model(m) if isinstance(m, str) else m
for m in models for m in models
+23
View File
@@ -1,4 +1,6 @@
import functools import functools
import os
import re
def rsetattr(obj, attr, val): def rsetattr(obj, attr, val):
pre, _, post = attr.rpartition('.') pre, _, post = attr.rpartition('.')
@@ -18,6 +20,27 @@ def rupdate(d, u): #https://stackoverflow.com/a/3233356/
d[k] = v d[k] = v
return d 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<weeks>\d+)w)?'
r'(?:(?P<days>\d+)d)?'
r'(?:(?P<hours>\d+)h)?'
r'(?:(?P<minutes>\d+)m)?'
r'(?:(?P<seconds>\d+)s)?'
r'(?:(?P<milliseconds>\d+)ms)?'
r'(?:(?P<microseconds>\d+)us)?'
)
#Django keeps giving: exception=OperationalError('the connection is closed') #Django keeps giving: exception=OperationalError('the connection is closed')
from django.db.utils import OperationalError from django.db.utils import OperationalError
+158 -156
View File
@@ -2,11 +2,12 @@
from django.db import IntegrityError, models, close_old_connections from django.db import IntegrityError, models, close_old_connections
from django.db.backends import signals as django_signals from django.db.backends import signals as django_signals
from django.db.utils import OperationalError from django.db.utils import OperationalError
from django.db import IntegrityError
from django.utils import timezone from django.utils import timezone
from django.tasks import task_backends
from django.tasks.base import TaskResultStatus
from asgiref.sync import sync_to_async from asgiref.sync import sync_to_async
import os, signal import os, signal
import time import time
import threading import threading
@@ -14,10 +15,14 @@ import logging, traceback
import asyncio import asyncio
import collections, functools import collections, functools
import random import random
import humanize
from .utils import retry_on_db_error, ignore_on_db_error from .utils import retry_on_db_error, ignore_on_db_error
from .asynctools import AsyncOriented from .asynctools import AsyncOriented
from .singleton import Singleton from .singleton import Singleton
from .models import Worker, TaskSchedule, Trace
class AsyncronWorker( Singleton, AsyncOriented ): class AsyncronWorker( Singleton, AsyncOriented ):
""" """
@@ -54,11 +59,9 @@ class AsyncronWorker( Singleton, AsyncOriented ):
def __init__( self ): def __init__( self ):
self.log #Evaluating the log property while we have the creation lock self.log #Evaluating the log property while we have the creation lock
self.is_db_ready_event = asyncio.Event() self.is_db_ready = asyncio.Event()
self.is_stopping_event = asyncio.Event() self.is_work_over = asyncio.Event()
self.model = Worker()
#Just so that the asyncron.apps.ready doens't trigger the django warning
self.start_after_db_ready = False
#Parallelism #Parallelism
self.thread = None #Once the worker starts, it'll be populated 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.clearing_dead_workers = False
self.watching_models = collections.defaultdict( set ) # Model -> Set of key name of the tasks 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 ) for callback in self.INIT_CALLBACKS: callback( self )
self.register_with_exit_signals() self.register_with_exit_signals()
@@ -107,14 +109,14 @@ class AsyncronWorker( Singleton, AsyncOriented ):
self.stop(f"Signal {signal.strsignal(signum)}") self.stop(f"Signal {signal.strsignal(signum)}")
def handle_new_db_connection( self, sender, **kwargs ): 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.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 ) 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. #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 ): 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()}" 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.thread = threading.current_thread()
self.start_working( is_robust = True ) self.start_working( is_standalone = True )
def stop( self, reason = None ): 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.log.info( f"Stopping Worker: {reason}" )
self.is_stopping_event.set() self.is_work_over.set()
#if not self.loop: return #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 ): async def startup( self ):
await super().startup() await super().startup()
self.backend = task_backends['default']
self.task_reason_jobs_queue = asyncio.Queue() #Run tasks from other threads, safely 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: with asyncio.Runner() as runner:
runner.run( self.startup() ) runner.run( self.startup() )
if self.start_after_db_ready: if not self.is_db_ready.is_set():
self.log.debug("Waiting on another module to create the first database connection...") #This Avoid's the django initialization warning, since waited for is_db_ready above!
runner.run( self.is_db_ready_event.wait() ) 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.is_standalone = is_standalone
self.model = Worker( pid = os.getpid(), thread_id = threading.get_ident(), is_robust = is_robust ) 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.maintain_scheduler() )
self.create_task( self.work_loop() )
self.model.save() #likley Avoid's the django initialization warning, since waited for is_db_ready_event above! self.create_task( self.run_tasks_on_schedule() )
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.attach_django_signals() self.attach_django_signals()
try: 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: except KeyboardInterrupt:
self.log.info(f"[W{self.model.id}] Worker Received KeyboardInterrupt, exiting...") self.log.info(f"[W{self.model.id}] Worker Received KeyboardInterrupt, exiting...")
@@ -203,15 +211,7 @@ class AsyncronWorker( Singleton, AsyncOriented ):
runner.run( self.cleanup() ) 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, # 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, # 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: # 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" # - 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.") 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 ): def attach_django_signals( self ):
django_name_to_signals = { django_name_to_signals = {
@@ -239,9 +325,10 @@ class AsyncronWorker( Singleton, AsyncOriented ):
and ( attr := getattr(models.signals, name) ) #Just an assignment and ( attr := getattr(models.signals, name) ) #Just an assignment
and isinstance( attr, models.signals.ModelSignal ) #Is a signal related to models! and isinstance( attr, models.signals.ModelSignal ) #Is a signal related to models!
} }
for name, signal in django_name_to_signals.items(): #for name, signal in django_name_to_signals.items():
signal.connect( functools.partial( self.model_changed, name ) ) # signal.connect( functools.partial( self.model_changed, name ) )
return
from .models import Task from .models import Task
for name, task in Task.registered_tasks.items(): for name, task in Task.registered_tasks.items():
if not hasattr(task, 'watching_models'): continue 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.") #print("Will not run another trace of the same task to reduce the change of an infinite cycle.")
continue 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( asyncio.run_coroutine_threadsafe(
task.ensure_quick_execution( reason = f"Change ({signal_name}) on {instance}" ), task.ensure_quick_execution( reason = f"Change ({signal_name}) on {instance}" ),
self.loop self.loop
@@ -286,12 +373,11 @@ class AsyncronWorker( Singleton, AsyncOriented ):
except: await asyncio.sleep( 1 ) except: await asyncio.sleep( 1 )
async def master_main( self ): async def maintain_scheduler( self ):
""" """
Fight over who's gonna be the master. Make sure at least one process is managing the scheduled tasks.
Prove your health in the process!
""" """
from .models import Worker, Task, Trace from .models import Worker, TaskSchedule, Trace
is_current_master = False is_current_master = False
next_overtake_attempt = time.time() + 1 + random.random() * 5 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() ): while await MyQs.aupdate( last_activity = timezone.now() ):
try: try:
await Worker.objects.filter( is_master = False ).aupdate( is_master = models.Q(id = self.model.id) ) await Worker.objects.filter( is_scheduler = False ).aupdate( is_scheduler = 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
except IntegrityError: # I'm not master! except IntegrityError: # I'm not master!
loop_wait = 5 + random.random() * 15 loop_wait = 5 + random.random() * 15
@@ -321,15 +400,15 @@ class AsyncronWorker( Singleton, AsyncOriented ):
next_overtake_attempt = time.time() + 60 next_overtake_attempt = time.time() + 60
took_master = False took_master = False
if self.model.is_robust: if self.model.is_standalone:
took_master = await Worker.objects.filter( is_master = True, is_robust = False ).aupdate( is_master = False ) took_master = await Worker.objects.filter( is_scheduler = True, is_standalone = False ).aupdate( is_scheduler = False )
loop_wait = 0 loop_wait = 0
else: 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 = None ) |
models.Q( last_activity__lte = timezone.now() - timezone.timedelta( minutes = 2 ) ) models.Q( last_activity__lte = timezone.now() - timezone.timedelta( minutes = 2 ) )
).aupdate( is_master = False ) ).aupdate( is_scheduler = False )
else: #I am Master! else: #I am Master!
@@ -352,7 +431,7 @@ class AsyncronWorker( Singleton, AsyncOriented ):
async def clear_orphaned_traces( self ): async def clear_orphaned_traces( self ):
from .models import Worker, Task, Trace 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 ): async def clear_dead_workers( self ):
self.clearing_dead_workers = True self.clearing_dead_workers = True
@@ -370,62 +449,8 @@ class AsyncronWorker( Singleton, AsyncOriented ):
await Worker.objects.filter( in_grace = True ).adelete() await Worker.objects.filter( in_grace = True ).adelete()
self.clearing_dead_workers = False 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 ): async def consume_task_reason_jobs_queue( self ):
while True: while True:
@@ -439,7 +464,7 @@ class AsyncronWorker( Singleton, AsyncOriented ):
#Schedule traces that aren't yet set. #Schedule traces that aren't yet set.
Ts = Task.objects.exclude( interval = None ).exclude( Ts = Task.objects.exclude( interval = None ).exclude(
trace__status = "S" 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: async for task in Ts:
trace = task.new_trace() 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 ) await Task.objects.filter( id = task.id, worker_lock = self.model ).aupdate( worker_lock = None )
early_seconds = 5 + self.check_interval * ( 1 + random.random() ) 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() await trace.eval_related()
#print(f"Checking {trace} to do now: {trace.scheduled_datetime - timezone.now()}") #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 except IntegrityError: count = 0
if not count: continue #Lost the race condition to another worker. 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. #start services that aren't running yet.
Ts = Task.objects.filter( interval = None, timeout = None ).exclude( Ts = Task.objects.filter( interval = None, timeout = None ).exclude(
trace__status__in = "WR" 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: async for task in Ts:
@@ -488,7 +513,7 @@ class AsyncronWorker( Singleton, AsyncOriented ):
trace = task.new_trace() trace = task.new_trace()
trace.set_status( "W", "Waiting to start the service ASAP." ) trace.set_status( "W", "Waiting to start the service ASAP." )
trace.worker_lock = self.model trace.worker = self.model
await trace.asave() await trace.asave()
#self.running_service_tasks[task.id] = #self.running_service_tasks[task.id] =
@@ -499,29 +524,6 @@ class AsyncronWorker( Singleton, AsyncOriented ):
trace = task.new_trace() trace = task.new_trace()
trace.set_status( "S", reason ) trace.set_status( "S", reason )
trace.scheduled_datetime = timezone.now() #So it runs instantly 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 ) ) 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()