Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 13 additions & 1 deletion ModuleConfig.bx
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,9 @@ class {
// IMetricsStore@rulebox yourself and point this at its mapping.
metricsStore = "InMemoryMetricsStore@rulebox",
// Datasource name SQLiteMetricsStore reads/writes, if you opt into it above.
datasourceName = "rulebox_visualizer"
datasourceName = "rulebox_visualizer",
// Max live-tracker (SSE) connections open at once; each pins two server threads. Over the cap, stream() answers 503.
maxStreams = 25
}
}

Expand Down Expand Up @@ -97,6 +99,16 @@ class {
invalidSetting( "visualizer.#key#", value, "a non-blank WireBox mapping or name" )
}
}
// Whole-number settings and the smallest value each accepts
var minimums = {
maxStreams : 1
}
for( var key, minimum in minimums ){
var value = arguments.visualizer[ key ] ?: ""
if( !isNumeric( value ) || value < minimum || value != int( value ) ){
invalidSetting( "visualizer.#key#", value, "a whole number of #minimum# or more" )
}
}
}

/**
Expand Down
25 changes: 25 additions & 0 deletions docs/guides/visualizer.md
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,31 @@ Each row shows the time, rulebook, rule, outcome state, and duration. Use
> `whitespaceCompressionEnabled` to `false` in `boxlang.json`, keeping in mind
> that it applies to all of your app's output.

#### Limiting live connections

Each open Live Tracker tab holds a stream open, and that pins two server
threads for as long as it stays connected. RuleBox therefore caps how many
streams can be open at once with `visualizer.maxStreams` (default `25`):

```cfc
moduleSettings = {
rulebox = {
visualizer = {
enabled = true,
maxStreams = 10
}
}
}
```

Once the cap is reached, further requests to `stream` get an HTTP `503` with
`{ "error": "Too many live tracker connections" }` instead of a new stream, and
a slot frees up as soon as a tab closes or its connection drops. Keep the cap
comfortably below your servlet container's worker thread count so the tracker
can never starve the rest of your application. Each stream also buffers at
most 1000 events; if a browser stalls and falls behind, its oldest unsent
events are dropped rather than letting memory grow.

## Metrics persistence

Every rule evaluation is recorded through `RuleEventBus@rulebox`, which fans
Expand Down
49 changes: 30 additions & 19 deletions handlers/Visualizer.bx
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,9 @@ class{
@inject( "RuleEventBus@rulebox" )
property name="eventBus";

@inject( "LiveStreams@rulebox" )
property name="liveStreams";

this.layout = "Visualizer"

/**
Expand Down Expand Up @@ -159,27 +162,35 @@ class{
* SSE stream of every rule-evaluation event, as it happens, across all rulebooks.
*/
function stream( event, rc, prc ){
var queue = createObject( "java", "java.util.concurrent.LinkedBlockingQueue" ).init()
var bus = variables.eventBus
var token = bus.subscribe( ( evt ) => queue.offer( evt ) )

SSE(
callback: ( emitter ) => {
try{
while( !emitter.isClosed() ){
var evt = queue.poll( 1000, createObject( "java", "java.util.concurrent.TimeUnit" ).MILLISECONDS )
if( !isNull( evt ) ){
emitter.send( evt, "rule" )
// Slot accounting, the bounded per-stream queue and the subscribe/unsubscribe pairing live
// in LiveStreams; serve() releases the slot and unsubscribes even if SSE() itself throws.
var served = variables.liveStreams.serve(
bus : variables.eventBus,
maxStreams : variables.settings.visualizer.maxStreams,
opener : ( queue ) => {
SSE(
callback: ( emitter ) => {
while( !emitter.isClosed() ){
var evt = queue.poll( 1000, createObject( "java", "java.util.concurrent.TimeUnit" ).MILLISECONDS )
if( !isNull( evt ) ){
emitter.send( evt, "rule" )
}
}
}
} finally {
bus.unsubscribe( token )
}
},
async: true,
keepAliveInterval: 15000,
timeout: 0
},
async: true,
keepAliveInterval: 15000,
timeout: 0
)
}
)

if( !served ){
event.renderData(
type = "json",
data = { "error": "Too many live tracker connections" },
statusCode = 503
)
}
}

/**
Expand Down
137 changes: 137 additions & 0 deletions models/metrics/LiveStreams.bx
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
/**
* Bookkeeping for the Rule Visualizer's live SSE streams (handlers/Visualizer.bx stream()):
*
* - caps how many streams may be open at once (each open stream pins two server threads),
* - hands out a bounded per-stream event queue that never blocks the event bus and never grows,
* - guarantees a stream's slot is released and its bus subscription removed however it ends.
*
* Kept out of the handler so the limit logic is unit-testable without opening a real stream.
*/
@singleton
@threadsafe
class{

/**
* Fixed limits, the same for every instance. Static so they are not copied per instance.
*/
static {
// Events held per stream before the oldest start being dropped
QUEUE_CAPACITY = 1000
}

/**
* Constructor: no streams open.
*
* @return This LiveStreams instance
*/
function init(){
variables.active = createObject( "java", "java.util.concurrent.atomic.AtomicInteger" ).init( 0 )
return this
}

/**
* How many streams currently hold a slot.
*/
numeric function activeCount(){
return variables.active.get()
}

/**
* Atomically take a stream slot unless maxStreams are already in use.
*
* @maxStreams The cap on concurrent streams
*
* @return true if a slot was taken (the caller must release() it), false if at the cap
*/
boolean function tryAcquire( required numeric maxStreams ){
while( true ){
var current = variables.active.get()
if( current >= arguments.maxStreams ){
return false
}
if( variables.active.compareAndSet( current, current + 1 ) ){
return true
}
}
}

/**
* Give back a slot taken by tryAcquire(). Never drops the count below zero.
*/
void function release(){
while( true ){
var current = variables.active.get()
if( current <= 0 || variables.active.compareAndSet( current, current - 1 ) ){
return
}
}
}

/**
* A fresh bounded queue for one stream's events.
*
* @capacity Max events held; defaults to QUEUE_CAPACITY
*/
any function newQueue( numeric capacity=static.QUEUE_CAPACITY ){
return createObject( "java", "java.util.concurrent.LinkedBlockingQueue" ).init( javaCast( "int", arguments.capacity ) )
}

/**
* Add an event to a bounded queue without ever blocking the publisher. When the queue is full
* (a slow or stalled client) the OLDEST queued event is dropped to make room: a live tracker
* is most useful showing what just happened.
*
* @queue A queue from newQueue()
* @evt The event to add
*
* @return true if nothing had to be dropped
*/
boolean function offer( required queue, required evt ){
if( arguments.queue.offer( arguments.evt ) ){
return true
}
// Full: evict the head and retry. Another producer can win the freed slot, so retry a few
// times; if we still lose, dropping this (newest) event is the same bounded behavior.
for( var attempt = 1; attempt <= 3; attempt++ ){
arguments.queue.poll()
if( arguments.queue.offer( arguments.evt ) ){
return false
}
}
return false
}

/**
* Open one live stream under the cap. Takes a slot, subscribes a bounded queue to the bus, runs
* the opener, and in every outcome (normal end, client gone, opener throws) unsubscribes and
* releases the slot. The opener normally blocks until the stream closes.
*
* @bus The RuleEventBus to subscribe to
* @maxStreams The cap on concurrent streams
* @opener A ( queue ) => void closure that opens the stream and drains the queue
*
* @return false (and runs nothing) if the cap is reached, true once a stream ran
*/
boolean function serve( required bus, required numeric maxStreams, required opener ){
if( !tryAcquire( arguments.maxStreams ) ){
return false
}

var token = 0
var subscribed = false
try{
var queue = newQueue()
var self = this
token = arguments.bus.subscribe( ( evt ) => self.offer( queue, evt ) )
subscribed = true
arguments.opener( queue )
} finally {
if( subscribed ){
arguments.bus.unsubscribe( token )
}
release()
}
return true
}

}
Loading
Loading