-
-
Notifications
You must be signed in to change notification settings - Fork 56
/
bulkhead.clj
195 lines (168 loc) · 5.62 KB
/
bulkhead.clj
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
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
;; This Source Code Form is subject to the terms of the Mozilla Public
;; License, v. 2.0. If a copy of the MPL was not distributed with this
;; file, You can obtain one at http://mozilla.org/MPL/2.0/.
;;
;; Copyright (c) Andrey Antukh <niwi@niwi.nz>
(ns promesa.exec.bulkhead
"Bulkhead pattern: limiter of concurrent executions."
(:refer-clojure :exclude [run! prn])
(:require
[promesa.exec :as px]
[promesa.impl :as pi]
[promesa.protocols :as pt]
[promesa.util :as pu])
(:import
clojure.lang.IObj
clojure.lang.IFn
clojure.lang.ILookup
java.util.concurrent.LinkedBlockingQueue
java.util.concurrent.BlockingQueue
java.util.concurrent.CompletableFuture
java.util.concurrent.ExecutorService
java.util.concurrent.Semaphore))
(set! *warn-on-reflection* true)
(declare ^:private instant)
(declare ^:privare run-hook!)
(defn- prn
[& args]
(locking prn
(apply clojure.core/prn args)))
(defn- current-thread
[]
(.getName (Thread/currentThread)))
(defprotocol IQueue
(-poll! [_])
(-offer! [_ _]))
(deftype Task [bulkhead task-fn result inst]
pt/ICancellable
(-cancelled? [_]
(pt/-cancelled? result))
(-cancel! [_]
(pt/-cancel! result))
clojure.core/Inst
(inst-ms* [_] inst)
IFn
(invoke [this]
;; (prn "Task.invoke0" (current-thread))
(let [semaphore (::semaphore bulkhead)
executor (::executor bulkhead)]
(run-hook! bulkhead ::on-run this)
(try
(if-let [rval (task-fn)]
(do
;; (prn "Task.invoke0" "rval=" rval)
(-> (pt/-promise rval)
(pt/-handle (fn [v e]
(pt/-release! semaphore)
(pt/-submit! executor bulkhead)
(if e
(pt/-reject! result e)
(pt/-resolve! result v))))))
(pt/-resolve! result nil))
(catch Throwable cause
(pt/-release! semaphore)
(pt/-submit! executor bulkhead)
(pt/-reject! result cause))))))
(deftype Bulkhead [queue semaphore executor metadata]
IObj
(meta [_] metadata)
(withMeta [this mdata]
(Bulkhead. queue semaphore executor mdata))
ILookup
(valAt [this key]
(case key
::executor executor
::semaphore semaphore
::queue-size (-> this meta ::queue-size)
::concurrency (-> this meta ::concurrency)
::current-queue-size (.size ^BlockingQueue queue)
::current-concurrency (- (-> this meta ::concurrency)
(.availablePermits ^Semaphore semaphore))
nil))
(valAt [this key default]
(or (.valAt ^ILookup this key) nil))
pt/IExecutor
(-run! [this task-fn]
(-> (-offer! this task-fn)
(pt/-map px/noop)))
(-submit! [this task-fn]
(-offer! this task-fn))
IFn
(invoke [this]
;; (prn "Bulkhead.invoke0" (current-thread) "start")
(loop []
(when-let [task (-poll! this)]
(if (pt/-cancelled? task)
(pt/-release! semaphore)
(pt/-submit! executor task))
(recur)))
;; (prn "Bulkhead.invoke0" (current-thread) "end")
)
IQueue
(-offer! [this task-fn]
(let [result (pi/deferred)
task (Task. this task-fn result (instant))]
(if (-offer! queue task)
(do
;; (prn "Bulkhead.offer!" (current-thread) "success")
(run-hook! this ::on-queue)
(this))
(let [size (-> this meta ::queue-size)
msg (str "Queue max capacity reached: " size)
props {:type :bulkhead-error
:code :capacity-limit-reached
::instance this}]
;; (prn "Bulkhead.offer!" (current-thread) "fail")
(pt/-reject! result (ex-info msg props))))
result))
(-poll! [this]
(when (pt/-try-acquire! semaphore)
(if-let [task (-poll! queue)]
(do
;; (prn "Bulkhead.poll!" (current-thread))
task)
(.release ^Semaphore semaphore)))))
(defn bulkhead?
"Check if the provided object is instance of Bulkhead type."
[o]
(instance? Bulkhead o))
(extend-type BlockingQueue
IQueue
(-poll! [this] (.poll ^BlockingQueue this))
(-offer! [this o] (.offer ^BlockingQueue this o)))
(extend-type Semaphore
pt/ISemaphore
(-try-acquire!
([this] (.tryAcquire ^Semaphore this))
([this permits] (.tryAcquire ^Semaphore this (int permits))))
(-acquire!
([this] (.acquire ^Semaphore this))
([this permits] (.acquire ^Semaphore this (int permits))))
(-release!
([this] (.release ^Semaphore this))
([this permits] (.release ^Semaphore this (int permits)))))
(defn- instant
[]
(System/currentTimeMillis))
(defn- run-hook!
([instance key-fn]
(when-let [hook-fn (-> instance meta key-fn)]
(let [executor (::executor instance)]
(pt/-submit! executor (partial hook-fn instance)))))
([instance key-fn param1]
(when-let [hook-fn (-> instance meta key-fn)]
(let [executor (::executor instance)]
(pt/-submit! executor (partial hook-fn instance param1))))))
(ns-unmap *ns* '->Bulkhead)
(defn create
[& {:keys [executor concurrency queue-size on-run on-queue]
:or {concurrency 1 queue-size Integer/MAX_VALUE}
:as params}]
(let [executor (px/resolve-executor (or executor px/*default-executor*))
queue (LinkedBlockingQueue. (int queue-size))
semaphore (Semaphore. (int concurrency))
metadata {::concurrency concurrency
::queue-size queue-size
::on-run on-run
::on-queue on-queue}]
(Bulkhead. queue semaphore executor metadata)))