From 50273c110856e18e3ba563b0f62b7273a4dc022b Mon Sep 17 00:00:00 2001 From: ckonstanski Date: Tue, 7 Oct 2025 06:53:09 -0600 Subject: initial commit --- lisp/queue/queue.lisp | 87 +++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 87 insertions(+) create mode 100644 lisp/queue/queue.lisp (limited to 'lisp/queue/queue.lisp') diff --git a/lisp/queue/queue.lisp b/lisp/queue/queue.lisp new file mode 100644 index 0000000..da86f6c --- /dev/null +++ b/lisp/queue/queue.lisp @@ -0,0 +1,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))))))) -- cgit v1.3