one logger per app and fixed some threading issues with newcron, thanks Paolo Pastori
This commit is contained in:
+57
-40
@@ -1,4 +1,5 @@
|
|||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
|
# vim: set ts=4 sw=4 et ai:
|
||||||
|
|
||||||
"""
|
"""
|
||||||
| This file is part of the web2py Web Framework
|
| This file is part of the web2py Web Framework
|
||||||
@@ -9,13 +10,15 @@
|
|||||||
Cron-style interface
|
Cron-style interface
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import sys
|
|
||||||
import os
|
|
||||||
import threading
|
import threading
|
||||||
import logging
|
import os
|
||||||
|
from logging import getLogger, DEBUG
|
||||||
import time
|
import time
|
||||||
import sched
|
import sched
|
||||||
|
import sys
|
||||||
import re
|
import re
|
||||||
|
import shlex
|
||||||
|
|
||||||
import datetime
|
import datetime
|
||||||
from functools import reduce
|
from functools import reduce
|
||||||
from gluon.settings import global_settings
|
from gluon.settings import global_settings
|
||||||
@@ -23,10 +26,18 @@ from gluon import fileutils
|
|||||||
from gluon._compat import to_bytes, pickle
|
from gluon._compat import to_bytes, pickle
|
||||||
from pydal.contrib import portalocker
|
from pydal.contrib import portalocker
|
||||||
|
|
||||||
logger = logging.getLogger("web2py.cron")
|
logger_name = 'web2py.cron'
|
||||||
|
|
||||||
|
|
||||||
_cron_stopping = False
|
_cron_stopping = False
|
||||||
|
|
||||||
|
_cron_subprocs_lock = threading.RLock()
|
||||||
_cron_subprocs = []
|
_cron_subprocs = []
|
||||||
|
|
||||||
|
def subprocess_count():
|
||||||
|
with _cron_subprocs_lock:
|
||||||
|
return len(_cron_subprocs)
|
||||||
|
|
||||||
|
|
||||||
def absolute_path_link(path):
|
def absolute_path_link(path):
|
||||||
"""
|
"""
|
||||||
@@ -45,14 +56,14 @@ def stopcron():
|
|||||||
"""Graceful shutdown of cron"""
|
"""Graceful shutdown of cron"""
|
||||||
global _cron_stopping
|
global _cron_stopping
|
||||||
_cron_stopping = True
|
_cron_stopping = True
|
||||||
while _cron_subprocs:
|
while subprocess_count():
|
||||||
proc = _cron_subprocs.pop()
|
with _cron_subprocs_lock:
|
||||||
|
proc = _cron_subprocs.pop()
|
||||||
if proc.poll() is None:
|
if proc.poll() is None:
|
||||||
try:
|
try:
|
||||||
proc.terminate()
|
proc.terminate()
|
||||||
except:
|
except Exception:
|
||||||
import traceback
|
getLogger(logger_name).exception('error in stopcron')
|
||||||
traceback.print_exc()
|
|
||||||
|
|
||||||
|
|
||||||
class extcron(threading.Thread):
|
class extcron(threading.Thread):
|
||||||
@@ -62,11 +73,10 @@ class extcron(threading.Thread):
|
|||||||
self.setDaemon(False)
|
self.setDaemon(False)
|
||||||
self.path = applications_parent
|
self.path = applications_parent
|
||||||
self.apps = apps
|
self.apps = apps
|
||||||
# crondance(self.path, 'external', startup=True, apps=self.apps)
|
|
||||||
|
|
||||||
def run(self):
|
def run(self):
|
||||||
if not _cron_stopping:
|
if not _cron_stopping:
|
||||||
logger.debug('external cron invocation')
|
getLogger(logger_name).debug('external cron invocation')
|
||||||
crondance(self.path, 'external', startup=False, apps=self.apps)
|
crondance(self.path, 'external', startup=False, apps=self.apps)
|
||||||
|
|
||||||
|
|
||||||
@@ -77,16 +87,19 @@ class hardcron(threading.Thread):
|
|||||||
self.setDaemon(True)
|
self.setDaemon(True)
|
||||||
self.path = applications_parent
|
self.path = applications_parent
|
||||||
self.apps = apps
|
self.apps = apps
|
||||||
|
# processing of '@reboot' entries in crontab (startup=True)
|
||||||
|
getLogger(logger_name).info('hard cron bootstrap')
|
||||||
crondance(self.path, 'hard', startup=True, apps=self.apps)
|
crondance(self.path, 'hard', startup=True, apps=self.apps)
|
||||||
|
|
||||||
def launch(self):
|
def launch(self):
|
||||||
if not _cron_stopping:
|
if not _cron_stopping:
|
||||||
logger.debug('hard cron invocation')
|
self.logger.debug('hard cron invocation')
|
||||||
crondance(self.path, 'hard', startup=False, apps=self.apps)
|
crondance(self.path, 'hard', startup=False, apps=self.apps)
|
||||||
|
|
||||||
def run(self):
|
def run(self):
|
||||||
|
self.logger = getLogger(logger_name)
|
||||||
|
self.logger.info('hard cron daemon started')
|
||||||
s = sched.scheduler(time.time, time.sleep)
|
s = sched.scheduler(time.time, time.sleep)
|
||||||
logger.info('hard cron daemon started')
|
|
||||||
while not _cron_stopping:
|
while not _cron_stopping:
|
||||||
now = time.time()
|
now = time.time()
|
||||||
s.enter(60 - now % 60, 1, self.launch, ())
|
s.enter(60 - now % 60, 1, self.launch, ())
|
||||||
@@ -99,14 +112,14 @@ class softcron(threading.Thread):
|
|||||||
threading.Thread.__init__(self)
|
threading.Thread.__init__(self)
|
||||||
self.path = applications_parent
|
self.path = applications_parent
|
||||||
self.apps = apps
|
self.apps = apps
|
||||||
# crondance(self.path, 'soft', startup=True, apps=self.apps)
|
|
||||||
|
|
||||||
def run(self):
|
def run(self):
|
||||||
if not _cron_stopping:
|
if not _cron_stopping:
|
||||||
logger.debug('soft cron invocation')
|
getLogger(logger_name).debug('soft cron invocation')
|
||||||
crondance(self.path, 'soft', startup=False, apps=self.apps)
|
crondance(self.path, 'soft', startup=False, apps=self.apps)
|
||||||
|
|
||||||
|
|
||||||
|
# TODO: context manager
|
||||||
class Token(object):
|
class Token(object):
|
||||||
|
|
||||||
def __init__(self, path):
|
def __init__(self, path):
|
||||||
@@ -115,6 +128,7 @@ class Token(object):
|
|||||||
fileutils.write_file(self.path, to_bytes(''), 'wb')
|
fileutils.write_file(self.path, to_bytes(''), 'wb')
|
||||||
self.master = None
|
self.master = None
|
||||||
self.now = time.time()
|
self.now = time.time()
|
||||||
|
self.logger = getLogger(logger_name)
|
||||||
|
|
||||||
def acquire(self, startup=False):
|
def acquire(self, startup=False):
|
||||||
"""
|
"""
|
||||||
@@ -133,7 +147,7 @@ class Token(object):
|
|||||||
else:
|
else:
|
||||||
locktime = 59.99
|
locktime = 59.99
|
||||||
if portalocker.LOCK_EX is None:
|
if portalocker.LOCK_EX is None:
|
||||||
logger.warning('cron disabled because no file locking')
|
self.logger.warning('cron disabled because no file locking')
|
||||||
return None
|
return None
|
||||||
self.master = fileutils.open_file(self.path, 'rb+')
|
self.master = fileutils.open_file(self.path, 'rb+')
|
||||||
try:
|
try:
|
||||||
@@ -148,8 +162,8 @@ class Token(object):
|
|||||||
ret = self.now
|
ret = self.now
|
||||||
if not stop:
|
if not stop:
|
||||||
# this happens if previous cron job longer than 1 minute
|
# this happens if previous cron job longer than 1 minute
|
||||||
logger.warning('stale cron.master detected')
|
self.logger.warning('stale cron.master detected')
|
||||||
logger.debug('acquiring lock')
|
self.logger.debug('acquiring lock')
|
||||||
self.master.seek(0)
|
self.master.seek(0)
|
||||||
pickle.dump((self.now, 0), self.master)
|
pickle.dump((self.now, 0), self.master)
|
||||||
self.master.flush()
|
self.master.flush()
|
||||||
@@ -167,7 +181,7 @@ class Token(object):
|
|||||||
ret = self.master.closed
|
ret = self.master.closed
|
||||||
if not self.master.closed:
|
if not self.master.closed:
|
||||||
portalocker.lock(self.master, portalocker.LOCK_EX)
|
portalocker.lock(self.master, portalocker.LOCK_EX)
|
||||||
logger.debug('releasing cron lock')
|
self.logger.debug('releasing cron lock')
|
||||||
self.master.seek(0)
|
self.master.seek(0)
|
||||||
(start, stop) = pickle.load(self.master)
|
(start, stop) = pickle.load(self.master)
|
||||||
if start == self.now: # if this is my lock
|
if start == self.now: # if this is my lock
|
||||||
@@ -248,26 +262,25 @@ class cronlauncher(threading.Thread):
|
|||||||
|
|
||||||
def run(self):
|
def run(self):
|
||||||
import subprocess
|
import subprocess
|
||||||
global _cron_subprocs
|
logger = getLogger(logger_name)
|
||||||
if isinstance(self.cmd, (list, tuple)):
|
proc = subprocess.Popen(self.cmd,
|
||||||
cmd = self.cmd
|
|
||||||
else:
|
|
||||||
cmd = self.cmd.split()
|
|
||||||
proc = subprocess.Popen(cmd,
|
|
||||||
stdin=subprocess.PIPE,
|
stdin=subprocess.PIPE,
|
||||||
stdout=subprocess.PIPE,
|
stdout=subprocess.PIPE,
|
||||||
stderr=subprocess.PIPE)
|
stderr=subprocess.PIPE)
|
||||||
_cron_subprocs.append(proc)
|
with _cron_subprocs_lock:
|
||||||
|
_cron_subprocs.append(proc)
|
||||||
(stdoutdata, stderrdata) = proc.communicate()
|
(stdoutdata, stderrdata) = proc.communicate()
|
||||||
try:
|
try:
|
||||||
_cron_subprocs.remove(proc)
|
with _cron_subprocs_lock:
|
||||||
|
_cron_subprocs.remove(proc)
|
||||||
except ValueError:
|
except ValueError:
|
||||||
pass
|
pass
|
||||||
if proc.returncode != 0:
|
if proc.returncode != 0:
|
||||||
logger.warning('call returned code %s:\n%s\n%s',
|
logger.warning('%r call returned code %s:\n%s\n%s',
|
||||||
proc.returncode, stdoutdata, stderrdata)
|
' '.join(self.cmd), proc.returncode, stdoutdata, stderrdata)
|
||||||
else:
|
elif logger.isEnabledFor(DEBUG):
|
||||||
logger.debug('call returned success:\n%s', stdoutdata)
|
logger.debug('%r call returned success:\n%s',
|
||||||
|
' '.join(self.cmd), stdoutdata)
|
||||||
|
|
||||||
|
|
||||||
def crondance(applications_parent, ctype='soft', startup=False, apps=None):
|
def crondance(applications_parent, ctype='soft', startup=False, apps=None):
|
||||||
@@ -287,6 +300,8 @@ def crondance(applications_parent, ctype='soft', startup=False, apps=None):
|
|||||||
('dom', now_s.tm_mday),
|
('dom', now_s.tm_mday),
|
||||||
('dow', (now_s.tm_wday + 1) % 7))
|
('dow', (now_s.tm_wday + 1) % 7))
|
||||||
|
|
||||||
|
logger = getLogger(logger_name)
|
||||||
|
|
||||||
if not apps:
|
if not apps:
|
||||||
apps = [x for x in os.listdir(apppath)
|
apps = [x for x in os.listdir(apppath)
|
||||||
if os.path.isdir(os.path.join(apppath, x))]
|
if os.path.isdir(os.path.join(apppath, x))]
|
||||||
@@ -332,15 +347,16 @@ def crondance(applications_parent, ctype='soft', startup=False, apps=None):
|
|||||||
for task in tasks:
|
for task in tasks:
|
||||||
if _cron_stopping:
|
if _cron_stopping:
|
||||||
break
|
break
|
||||||
citems = [(k in task and not v in task[k]) for k, v in checks]
|
|
||||||
task_min = task.get('min', [])
|
|
||||||
if not task:
|
if not task:
|
||||||
continue
|
continue
|
||||||
elif not startup and task_min == [-1]:
|
task_min = task.get('min', [])
|
||||||
|
if not startup and task_min == [-1]:
|
||||||
continue
|
continue
|
||||||
elif task_min != [-1] and reduce(lambda a, b: a or b, citems):
|
citems = [(k in task and not v in task[k]) for k, v in checks]
|
||||||
|
if task_min != [-1] and reduce(lambda a, b: a or b, citems):
|
||||||
continue
|
continue
|
||||||
logger.info('%s cron: %s executing %s in %s at %s',
|
|
||||||
|
logger.info('%s cron: %s executing %r in %s at %s',
|
||||||
ctype, app, task.get('cmd'),
|
ctype, app, task.get('cmd'),
|
||||||
os.getcwd(), datetime.datetime.now())
|
os.getcwd(), datetime.datetime.now())
|
||||||
action = models = False
|
action = models = False
|
||||||
@@ -361,11 +377,12 @@ def crondance(applications_parent, ctype='soft', startup=False, apps=None):
|
|||||||
if models:
|
if models:
|
||||||
commands.append('-M')
|
commands.append('-M')
|
||||||
else:
|
else:
|
||||||
commands = command
|
commands = shlex.split(command)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
|
# FIXME: using a new thread every time there is a task to
|
||||||
|
# launch is not a good idea in a long running process
|
||||||
cronlauncher(commands).start()
|
cronlauncher(commands).start()
|
||||||
except Exception as e:
|
except Exception:
|
||||||
logger.warning('execution error for %s: %s',
|
logger.exception('error starting %r', task['cmd'])
|
||||||
task.get('cmd'), e)
|
|
||||||
token.release()
|
token.release()
|
||||||
|
|||||||
+56
-11
@@ -1,22 +1,67 @@
|
|||||||
#!/usr/bin/env python
|
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
|
# vim: set ts=4 sw=4 et ai:
|
||||||
|
"""
|
||||||
|
Unit tests for cron
|
||||||
|
"""
|
||||||
|
|
||||||
""" Unit tests for cron """
|
import unittest
|
||||||
import unittest, os
|
import os
|
||||||
from gluon.newcron import Token, crondance
|
import shutil
|
||||||
|
import time
|
||||||
|
|
||||||
|
from gluon.newcron import Token, crondance, subprocess_count
|
||||||
|
from gluon.fileutils import create_app, write_file
|
||||||
|
|
||||||
|
|
||||||
|
test_app_name = '_test_cron'
|
||||||
|
appdir = os.path.join('applications', test_app_name)
|
||||||
|
|
||||||
|
def setUpModule():
|
||||||
|
if not os.path.exists(appdir):
|
||||||
|
os.mkdir(appdir)
|
||||||
|
create_app(appdir)
|
||||||
|
|
||||||
|
def tearDownModule():
|
||||||
|
if os.path.exists(appdir):
|
||||||
|
shutil.rmtree(appdir)
|
||||||
|
|
||||||
|
|
||||||
|
TEST_CRONTAB = """@reboot peppe **applications/%s/cron/test.py
|
||||||
|
""" % test_app_name
|
||||||
|
TEST_SCRIPT = """
|
||||||
|
from os.path import join as pjoin
|
||||||
|
|
||||||
|
with open(pjoin(request.folder, 'private', 'cron_req'), 'w') as f:
|
||||||
|
f.write(str(request))
|
||||||
|
"""
|
||||||
|
TARGET = os.path.join(appdir, 'private', 'cron_req')
|
||||||
|
|
||||||
class TestCron(unittest.TestCase):
|
class TestCron(unittest.TestCase):
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def tearDownClass(cls):
|
||||||
|
for master in (os.path.join(appdir, 'cron.master'), 'cron.master'):
|
||||||
|
if os.path.exists(master):
|
||||||
|
os.unlink(master)
|
||||||
|
|
||||||
def test_Token(self):
|
def test_Token(self):
|
||||||
appname_path = os.path.join(os.getcwd(), 'applications', 'welcome')
|
app_path = os.path.join(os.getcwd(), 'applications', test_app_name)
|
||||||
t = Token(path=appname_path)
|
t = Token(path=app_path)
|
||||||
self.assertNotEqual(t.acquire(), None)
|
self.assertEqual(t.acquire(), t.now)
|
||||||
self.assertFalse(t.release())
|
self.assertFalse(t.release())
|
||||||
self.assertEqual(t.acquire(), None)
|
self.assertIsNone(t.acquire())
|
||||||
self.assertTrue(t.release())
|
self.assertTrue(t.release())
|
||||||
return
|
|
||||||
|
|
||||||
def test_crondance(self):
|
def test_crondance(self):
|
||||||
#TODO update crondance to return something
|
base = os.path.join(appdir, 'cron')
|
||||||
crondance(os.getcwd())
|
write_file(os.path.join(base, 'crontab'), TEST_CRONTAB)
|
||||||
|
write_file(os.path.join(base, 'test.py'), TEST_SCRIPT)
|
||||||
|
if os.path.exists(TARGET):
|
||||||
|
os.unlink(TARGET)
|
||||||
|
crondance(os.getcwd(), 'hard', startup=True, apps=[test_app_name])
|
||||||
|
# must wait for the cron task, not very reliable
|
||||||
|
time.sleep(1)
|
||||||
|
while subprocess_count():
|
||||||
|
time.sleep(1)
|
||||||
|
time.sleep(1)
|
||||||
|
self.assertTrue(os.path.exists(TARGET))
|
||||||
|
|||||||
Reference in New Issue
Block a user