summaryrefslogtreecommitdiff
path: root/lisp/queue
diff options
context:
space:
mode:
authorBret Koppel <bret@olo.com>2025-10-07 09:00:37 -0400
committerGitHub <noreply@github.com>2025-10-07 09:00:37 -0400
commit910d7fd656b9e957181732b36761449de5a970d1 (patch)
tree0f8ee46af3e93aa489785d64d54c5e155ef0057b /lisp/queue
parent9adba239b937df2830a5110adaa1dae17b7dd7d7 (diff)
parent0ad86ba37be273ee1fe4d9bdd5f7bb15c5701abf (diff)
Merge pull request #1 from ololabs/initial-commit
initial commit
Diffstat (limited to 'lisp/queue')
-rw-r--r--lisp/queue/fifo.lisp44
-rw-r--r--lisp/queue/queue.lisp87
2 files changed, 131 insertions, 0 deletions
diff --git a/lisp/queue/fifo.lisp b/lisp/queue/fifo.lisp
new file mode 100644
index 0000000..8bbd106
--- /dev/null
+++ b/lisp/queue/fifo.lisp
@@ -0,0 +1,44 @@
+;;; -*- Mode: Lisp; Syntax: ANSI-Common-Lisp; Base: 10 -*-
+(declaim (optimize (speed 0) (safety 3) (debug 3)))
+
+(in-package #:snow)
+
+(defclass fifo ()
+ ((buffer :initarg :buffer
+ :initform ()
+ :accessor buffer)
+ (mutex :initarg :mutex
+ :initform (sb-thread:make-mutex)
+ :accessor mutex)
+ (discard-preceding :initarg :discard-preceding
+ :initform nil
+ :accessor discard-preceding)
+ (wait-interval :initarg :wait-interval
+ :initform 0
+ :accessor wait-interval)
+ (timestamp :initarg :timestamp
+ :initform (get-universal-time)
+ :accessor timestamp))
+ (:documentation ""))
+
+(defmethod dequeue ((fifo fifo))
+ (sb-thread:with-mutex ((mutex fifo))
+ (when (or (= (wait-interval fifo) 0)
+ (> (get-universal-time) (+ (timestamp fifo) (wait-interval fifo))))
+ (setf (timestamp fifo) (get-universal-time))
+ (when (buffer fifo)
+ (pop (buffer fifo))))))
+
+(defmethod enqueue ((fifo fifo) obj)
+ (sb-thread:with-mutex ((mutex fifo))
+ (if (discard-preceding fifo)
+ (setf (buffer fifo) `(,obj))
+ (push obj (buffer fifo)))))
+
+(defmethod empty-p ((fifo fifo))
+ (sb-thread:with-mutex ((mutex fifo))
+ (endp (buffer fifo))))
+
+(defmethod len ((fifo fifo))
+ (sb-thread:with-mutex ((mutex fifo))
+ (length (buffer fifo))))
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)))))))