summaryrefslogtreecommitdiff
path: root/lisp/queue/queue.lisp
blob: da86f6c7f22dfe5874cbe35f9544dce2de69faa5 (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
;;; -*- 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)
  "`level' is either :warning or :critical."
  (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)
  (with-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)))
      (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)))))))