blob: 47b9e658e7b2b0327f864f4eec64a5177a292e16 (
plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
|
;;; -*- Mode: Lisp; Syntax: ANSI-Common-Lisp; Base: 10 -*-
(declaim (optimize (speed 0) (safety 3) (debug 3)))
(in-package :snow)
(defvar *queue-git* nil)
(defvar *queue-aws* nil)
(defvar *queue-tfcloud* nil)
(defvar *queue-octopus-slow* nil)
(defvar *queue-octopus-fast* nil)
(defvar *queue-awx* nil)
(defvar *queue-mutex* (sb-thread:make-mutex))
(defmacro with-error-handled-thread ((thread-name level) &body body)
(declare (ignore level))
(let ((thread-fn (intern (string-upcase thread-name) *package*)))
`(progn
(defun ,thread-fn ()
(handler-case
,@body
(error (e)
(let ((message (format nil "The ~a thread died unexpectedly.~% exception = [~a]" ,thread-name e)))
(org-ckons-core::logger message)
;; sleep so we don't get an error shitstorm
(sleep 10)
(sb-thread:make-thread (lambda () (funcall #',thread-fn))
:name ,thread-name)))))
(sb-thread:make-thread (lambda () (funcall #',thread-fn))
:name ,thread-name))))
(defmacro without-error-handled-thread ((thread-name level) &body body)
"For development. Falls into the debugger and dies on error."
(declare (ignore level))
(let ((thread-fn (intern (string-upcase thread-name) *package*)))
`(progn
(defun ,thread-fn ()
,@body)
(sb-thread:make-thread (lambda () (funcall #',thread-fn))
:name ,thread-name))))
(defun generate-thread-id ()
"Generates a unique random string to use in the process thread
names. The string is a SHA1 hash."
(subseq (generate-sessionid) 0 32))
(defun queue-generator (queue-type)
(let* ((queue-name (symbol-name queue-type))
(symbol-queue-name (intern (format nil "*QUEUE-~a*" queue-name) (package-name #.*package*)))
(queue-keyword (intern queue-name "KEYWORD"))
(symbol-sleep-interval (intern "SLEEP-INTERVAL" "KEYWORD"))
(symbol-num-process-threads (intern "NUM-PROCESS-THREADS" "KEYWORD"))
(symbol-wait-interval (intern "WAIT-INTERVAL" "KEYWORD")))
(when (null (eval symbol-queue-name))
(set symbol-queue-name (make-instance 'fifo :wait-interval (getf (getf (queue *webapp*) queue-keyword) symbol-wait-interval)))
(process-thread-generator queue-name queue-keyword symbol-queue-name symbol-sleep-interval symbol-num-process-threads))))
(defun process-thread-generator (queue-name queue-keyword symbol-queue-name symbol-sleep-interval symbol-num-process-threads)
(let ((sleep-interval (getf (queue *webapp*) symbol-sleep-interval))
(num-threads (getf (getf (queue *webapp*) queue-keyword) symbol-num-process-threads)))
(loop for i from 1 to num-threads
do (let ((process-thread-name (format nil "~a-PROCESS-THREAD-~a" queue-name (generate-thread-id))))
(process-thread-function symbol-queue-name process-thread-name sleep-interval)))))
(defun process-thread-function (symbol-queue-name process-thread-name sleep-interval)
(without-error-handled-thread (process-thread-name :warning)
(labels ((do-dequeue ()
(dequeue (eval symbol-queue-name)))
(do-process (object)
(bg-perform object))
(do-sleep ()
(sleep sleep-interval))
(do-loop ()
(handler-case
(loop
(let (object
do-process-p
do-sleep-p)
(sb-thread:with-mutex (*queue-mutex*)
(cond ((empty-p (eval symbol-queue-name))
(setf do-sleep-p t))
(t
(setf object (do-dequeue))
(if object
(setf do-process-p t)
(setf do-sleep-p t)))))
(when do-sleep-p
(do-sleep))
(when (and do-process-p object)
(do-process object))))
(error (e)
(org-ckons-core::logger e)
(unless hunchentoot:*catch-errors-p*
(invoke-debugger e))
(do-loop)))))
(do-loop))))
|