;;; -*- 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)))))))