Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 34 additions & 10 deletions kolibri/core/auth/test/sync_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ def __init__(
self.port = get_free_tcp_port()
self.baseurl = "http://127.0.0.1:{}/".format(self.port)
self.enable_automatic_download = enable_automatic_download
self._instance = None
if seeded_kolibri_home is not None:
shutil.rmtree(self.env["KOLIBRI_HOME"])
shutil.copytree(seeded_kolibri_home, self.env["KOLIBRI_HOME"])
Expand Down Expand Up @@ -132,7 +133,8 @@ def _wait_for_server_start(self, timeout=20):
def kill(self):
try:
subprocess.Popen("kolibri stop", env=self.env, shell=True)
self._instance.kill()
if self._instance is not None:
self._instance.kill()
shutil.rmtree(self.env["KOLIBRI_HOME"])
except OSError:
pass
Expand Down Expand Up @@ -188,6 +190,7 @@ def generate_base_data(self):
class multiple_kolibri_servers:
def __init__(self, count=2, **server_kwargs):
self.server_count = count
self.servers = []
self.server_kwargs = [
{
key: value[i] if isinstance(value, (list, tuple)) else value
Expand All @@ -197,6 +200,19 @@ def __init__(self, count=2, **server_kwargs):
]

def __enter__(self):
try:
self._start_servers()
except BaseException:
# Servers left running hold connections that stop every later test
# creating its own database. BaseException: Django exits rather than
# raises when it cannot clobber one.
self.__exit__(None, None, None)
raise
return self.servers

def _start_servers(self):
# the same instance is reused for every invocation, so start from scratch
self.servers = []
# spin up the servers
if "sqlite" in connection.vendor:
tempserver = KolibriServer(
Expand All @@ -208,12 +224,15 @@ def __enter__(self):
tempserver.delete_model(DatabaseIDModel)
preseeded_home = tempserver.env["KOLIBRI_HOME"]

self.servers = [
KolibriServer(
seeded_kolibri_home=preseeded_home, **self.server_kwargs[i]
for i in range(self.server_count):
# track before starting, so a failed start is still shut down
server = KolibriServer(
autostart=False,
seeded_kolibri_home=preseeded_home,
**self.server_kwargs[i],
)
for i in range(self.server_count)
]
self.servers.append(server)
server.start()

# calculate the DATABASE settings
for server in self.servers:
Expand Down Expand Up @@ -259,17 +278,22 @@ def __enter__(self):
server_conn.close()
server.start()

return self.servers

def __exit__(self, typ, val, traceback):
# make sure all the servers are shut down
# kill every server before touching any database, so that a database that
# refuses to drop cannot abort the loop and leave later servers running
for server in self.servers:
server.kill()
for server in self.servers:
# a server abandoned before its alias was registered has no database
if server.db_alias not in connections.databases:
continue
# destroy the test databases
server_conn = connections[server.db_alias]
try:
server_conn.creation.destroy_test_db()
except OSError:
except Exception:
# Nothing narrower will do: Django surfaces a missing database
# as a RuntimeError from _nodb_cursor.
pass
server_conn.close()
# Remove the database alias from settings to prevent subsequent tests
Expand Down
2 changes: 2 additions & 0 deletions kolibri/core/auth/test/test_auth_tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,8 @@ class dummy_orm_job_data:
interval = 8600
retry_interval = 5
max_retries = 3
last_finished_state = None
last_finished_time = None


@patch("kolibri.core.tasks.viewsets.tasks.job_storage")
Expand Down
18 changes: 18 additions & 0 deletions kolibri/core/tasks/job.py
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,12 @@ class State:
CANCELING,
}

FINISHED_STATES = {
FAILED,
CANCELED,
COMPLETED,
}


JobStatus = namedtuple("Status", ("title", "text"))

Expand Down Expand Up @@ -263,6 +269,18 @@ def __init__(
self._supervisor_id = NO_VALUE
self.func = callable_to_import_path(func)

def reset_for_new_run(self):
"""
Clear the fields describing a single run, so a repeating job does not
inherit the previous run's progress or outcome. extra_metadata is kept:
it is what the task manager renders for the run that just finished.
"""
self.exception = None
self.traceback = ""
self.progress = 0
self.total_progress = 0
self.result = None

def _check_storage_attached(self):
if self._storage is None:
raise ReferenceError(
Expand Down
21 changes: 21 additions & 0 deletions kolibri/core/tasks/migrations/0005_add_last_finished_fields.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
from django.db import migrations
from django.db import models


class Migration(migrations.Migration):
dependencies = [
("kolibritasks", "0004_add_supervisor_registry"),
]

operations = [
migrations.AddField(
model_name="job",
name="last_finished_state",
field=models.CharField(blank=True, max_length=20, null=True),
),
migrations.AddField(
model_name="job",
name="last_finished_time",
field=models.DateTimeField(blank=True, null=True),
),
]
5 changes: 5 additions & 0 deletions kolibri/core/tasks/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,11 @@ class Job(models.Model):
# Maximum number of retries allowed for the job
max_retries = models.IntegerField(null=True, blank=True)

# The outcome of the last finished run, which a repeating job's own state
# stops describing as soon as the job is re-queued.
last_finished_state = models.CharField(max_length=20, null=True, blank=True)
last_finished_time = models.DateTimeField(null=True, blank=True)

# References the supervisor currently responsible for this job.
# No FK constraint to avoid complexity with supervisor cleanup.
supervisor_id = UUIDField(null=True, blank=True)
Expand Down
19 changes: 9 additions & 10 deletions kolibri/core/tasks/storage.py
Original file line number Diff line number Diff line change
Expand Up @@ -475,11 +475,7 @@ def clear(self, queue=None, job_id=None, force=False):

# filter only by the finished jobs, if we are not specified to force
if not force:
queryset = queryset.filter(
Q(state=State.COMPLETED)
| Q(state=State.FAILED)
| Q(state=State.CANCELED)
)
queryset = queryset.filter(state__in=State.FINISHED_STATES)

if self._hooks:
for orm_job in queryset:
Expand Down Expand Up @@ -673,7 +669,7 @@ def reschedule_finished_job_if_needed( # noqa: C901
orm_job = self.get_orm_job(job_id)

# Only allow this function to be run on a job that is in a finished state.
if orm_job.state not in {State.COMPLETED, State.FAILED, State.CANCELED}:
if orm_job.state not in State.FINISHED_STATES:
raise JobNotRestartable(
"Cannot reschedule job with state={}".format(orm_job.state)
)
Expand Down Expand Up @@ -712,10 +708,7 @@ def reschedule_finished_job_if_needed( # noqa: C901
current_retries = orm_job.retries if orm_job.retries is not None else 0
kwargs["retries"] = current_retries + 1

elif (
orm_job.state in {State.COMPLETED, State.FAILED, State.CANCELED}
and kwargs["repeat"] != 0
):
elif orm_job.state in State.FINISHED_STATES and kwargs["repeat"] != 0:
# Otherwise, if we are in a finished state and repeat is not 0, then we can reschedule, either because
# repeat is None, or because repeat is not None and is greater than 0.
if kwargs["repeat"] is not None:
Expand Down Expand Up @@ -847,6 +840,9 @@ def _write_job_update(
self, job, orm_job, state, supervisor_id, kwargs, repeat=NO_VALUE
):
if state is not None:
# orm_job.state is the state being left, until the assignment below.
if state == State.RUNNING and orm_job.state != State.RUNNING:
job.reset_for_new_run()
orm_job.state = job.state = state
# Ownership exists only in supervised states; a bare re-mark
# preserves the owner, a terminal state clears it.
Expand All @@ -855,6 +851,9 @@ def _write_job_update(
orm_job.supervisor_id = supervisor_id
else:
orm_job.supervisor_id = None
if state in State.FINISHED_STATES:
orm_job.last_finished_state = state
orm_job.last_finished_time = self._now()
# repeat is nullable, so None is a real value; NO_VALUE means "leave it".
if repeat is not NO_VALUE:
orm_job.repeat = repeat
Expand Down
100 changes: 100 additions & 0 deletions kolibri/core/tasks/test/taskrunner/test_storage.py
Original file line number Diff line number Diff line change
Expand Up @@ -484,6 +484,106 @@ def test_reschedule_finished_job_failed(self, defaultbackend, simplejob):
assert requeued_job.state == State.QUEUED
assert requeued_orm_job.scheduled_time > previous_scheduled_time

def test_reschedule_recurring_job_keeps_last_finished(
self, defaultbackend, simplejob
):
job_id = defaultbackend.enqueue_at(
local_now(), simplejob, QUEUE, interval=10, repeat=None
)
assert defaultbackend.get_orm_job(job_id).last_finished_state is None
defaultbackend.complete_job(job_id)

defaultbackend.reschedule_finished_job_if_needed(job_id)

requeued_orm_job = defaultbackend.get_orm_job(job_id)

assert requeued_orm_job.state == State.QUEUED
assert requeued_orm_job.last_finished_state == State.COMPLETED
assert requeued_orm_job.last_finished_time is not None

def test_reschedule_with_delay_keeps_last_finished(self, defaultbackend, simplejob):
job_id = defaultbackend.enqueue_job(simplejob, QUEUE)
defaultbackend.complete_job(job_id)

defaultbackend.reschedule_finished_job_if_needed(
job_id, delay=datetime.timedelta(seconds=5)
)

requeued_orm_job = defaultbackend.get_orm_job(job_id)

assert requeued_orm_job.state == State.QUEUED
assert requeued_orm_job.last_finished_state == State.COMPLETED
assert requeued_orm_job.last_finished_time is not None

def test_last_finished_describes_the_most_recent_run(
self, defaultbackend, simplejob
):
exception = ValueError("Error")
job_id = defaultbackend.enqueue_job(
simplejob, QUEUE, retry_interval=5, max_retries=3
)
defaultbackend.mark_job_as_failed(job_id, exception, "Traceback")
defaultbackend.reschedule_finished_job_if_needed(job_id, exception=exception)
assert defaultbackend.get_orm_job(job_id).last_finished_state == State.FAILED
defaultbackend.mark_job_as_running(job_id)

defaultbackend.complete_job(job_id)

assert defaultbackend.get_orm_job(job_id).last_finished_state == State.COMPLETED

def test_clear_leaves_rescheduled_recurring_job(self, defaultbackend, simplejob):
job_id = defaultbackend.enqueue_at(
local_now(), simplejob, QUEUE, interval=10, repeat=None
)
defaultbackend.complete_job(job_id)
defaultbackend.reschedule_finished_job_if_needed(job_id)

defaultbackend.clear(force=False)

assert defaultbackend.get_orm_job(job_id).state == State.QUEUED

def test_marking_running_clears_completed_run_fields(
self, defaultbackend, simplejob
):
job_id = defaultbackend.enqueue_job(simplejob, QUEUE)
defaultbackend.mark_job_as_running(job_id)
defaultbackend.update_job_progress(job_id, 5, 10)
defaultbackend.complete_job(job_id, result="run one")

defaultbackend.mark_job_as_running(job_id)

job = defaultbackend.get_job(job_id)

assert job.progress == 0
assert job.total_progress == 0
assert job.result is None

def test_marking_running_clears_failed_run_fields(self, defaultbackend, simplejob):
job_id = defaultbackend.enqueue_job(simplejob, QUEUE)
defaultbackend.mark_job_as_running(job_id)
defaultbackend.mark_job_as_failed(job_id, ValueError("Error"), "Traceback")

defaultbackend.mark_job_as_running(job_id)

job = defaultbackend.get_job(job_id)

assert job.exception is None
assert job.traceback == ""

def test_marking_running_keeps_last_finished(self, defaultbackend, simplejob):
job_id = defaultbackend.enqueue_at(
local_now(), simplejob, QUEUE, interval=10, repeat=None
)
defaultbackend.complete_job(job_id)
defaultbackend.reschedule_finished_job_if_needed(job_id)

defaultbackend.mark_job_as_running(job_id)

orm_job = defaultbackend.get_orm_job(job_id)

assert orm_job.last_finished_state == State.COMPLETED
assert orm_job.last_finished_time is not None

def test_reschedule_finished_job_invalid_state_queued(
self, defaultbackend, simplejob
):
Expand Down
Loading
Loading