diff --git a/ModuleConfig.bx b/ModuleConfig.bx index 9e5b457d..e8f1a4f9 100644 --- a/ModuleConfig.bx +++ b/ModuleConfig.bx @@ -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 } } @@ -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" ) + } + } } /** diff --git a/docs/guides/visualizer.md b/docs/guides/visualizer.md index 3be31850..1cca4ea3 100644 --- a/docs/guides/visualizer.md +++ b/docs/guides/visualizer.md @@ -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 diff --git a/handlers/Visualizer.bx b/handlers/Visualizer.bx index cf57727b..65638dfa 100644 --- a/handlers/Visualizer.bx +++ b/handlers/Visualizer.bx @@ -20,6 +20,9 @@ class{ @inject( "RuleEventBus@rulebox" ) property name="eventBus"; + @inject( "LiveStreams@rulebox" ) + property name="liveStreams"; + this.layout = "Visualizer" /** @@ -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 + ) + } } /** diff --git a/models/metrics/LiveStreams.bx b/models/metrics/LiveStreams.bx new file mode 100644 index 00000000..7e817603 --- /dev/null +++ b/models/metrics/LiveStreams.bx @@ -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 + } + +} diff --git a/test-harness/tests/specs/RuleEventBusSpec.bx b/test-harness/tests/specs/RuleEventBusSpec.bx index 97ebd82b..65e807bd 100644 --- a/test-harness/tests/specs/RuleEventBusSpec.bx +++ b/test-harness/tests/specs/RuleEventBusSpec.bx @@ -118,6 +118,171 @@ class extends="tests.resources.BaseSpec"{ } ); } ); + + describe( "LiveStreams (the Live Tracker's stream cap and bounded queues)", function(){ + + function newLiveStreams(){ + return new rulebox.models.metrics.LiveStreams(); + } + + function newStreamBus(){ + var bus = new rulebox.models.metrics.RuleEventBus(); + bus.setWirebox( getWireBox() ); + bus.setLogger( getController().getLogBox().getRootLogger() ); + bus.setSettings( { visualizer: { enabled: true, metricsStore: "InMemoryMetricsStore@rulebox", datasourceName: "" } } ); + return bus; + } + + // Drain a queue into an array of its events' ruleName, oldest first + function drainRuleNames( required queue ){ + var names = []; + while( !arguments.queue.isEmpty() ){ + names.append( arguments.queue.poll().ruleName ); + } + return names; + } + + it( "rejects the N+1th stream and admits one again once a slot is released", function(){ + var streams = newLiveStreams(); + + expect( streams.tryAcquire( 2 ) ).toBeTrue(); + expect( streams.tryAcquire( 2 ) ).toBeTrue(); + expect( streams.tryAcquire( 2 ) ).toBeFalse(); + expect( streams.activeCount() ).toBe( 2 ); + + streams.release(); + + expect( streams.activeCount() ).toBe( 1 ); + expect( streams.tryAcquire( 2 ) ).toBeTrue(); + } ); + + it( "never lets release() drive the count below zero", function(){ + var streams = newLiveStreams(); + + streams.release(); + + expect( streams.activeCount() ).toBe( 0 ); + } ); + + it( "never admits more than the cap under concurrent attempts", function(){ + var streams = newLiveStreams(); + var ids = []; + for( var i = 1; i <= 100; i++ ){ + ids.append( i ); + } + + var admitted = ids.map( ( id ) => streams.tryAcquire( 10 ), true ); + + expect( admitted.filter( ( ok ) => ok ) ).toHaveLength( 10 ); + expect( streams.activeCount() ).toBe( 10 ); + } ); + + it( "serve() refuses to open a stream at the cap, without subscribing or running the opener", function(){ + var streams = newLiveStreams(); + var bus = newStreamBus(); + var opened = false; + streams.tryAcquire( 1 ); + + var served = streams.serve( bus, 1, ( queue ) => { opened = true } ); + + expect( served ).toBeFalse(); + expect( opened ).toBeFalse(); + expect( bus.getSubscribers() ).toBeEmpty(); + expect( streams.activeCount() ).toBe( 1 ); + } ); + + it( "serve() runs the opener against a queue the bus feeds, then unsubscribes and frees the slot", function(){ + var streams = newLiveStreams(); + var bus = newStreamBus(); + var seen = []; + var during = -1; + + var served = streams.serve( bus, 1, ( queue ) => { + bus.publish( { rulebookName: "rb", ruleName: "first", state: "EXECUTED", durationMs: 1, timestamp: "" } ); + seen = drainRuleNames( queue ); + during = streams.activeCount(); + } ); + + expect( served ).toBeTrue(); + expect( during ).toBe( 1 ); + expect( seen.toList() ).toBe( "first" ); + expect( bus.getSubscribers() ).toBeEmpty(); + expect( streams.activeCount() ).toBe( 0 ); + } ); + + it( "serve() releases the slot and unsubscribes when opening the stream throws", function(){ + var streams = newLiveStreams(); + var bus = newStreamBus(); + var threw = false; + + try{ + streams.serve( bus, 1, ( queue ) => { throw( type="Test.OpenFailed", message="SSE() blew up" ) } ); + } catch( any e ){ + threw = ( e.type == "Test.OpenFailed" ); + } + + expect( threw ).toBeTrue(); + expect( bus.getSubscribers() ).toBeEmpty(); + expect( streams.activeCount() ).toBe( 0 ); + // ...and the freed slot is genuinely usable + expect( streams.tryAcquire( 1 ) ).toBeTrue(); + } ); + + it( "serve() releases the slot when subscribing itself throws", function(){ + var streams = newLiveStreams(); + var failingBus = { + subscribe : ( listener ) => { throw( type="Test.SubscribeFailed", message="no" ) }, + unsubscribe: ( token ) => { throw( type="Test.ShouldNotUnsubscribe", message="nothing was subscribed" ) } + }; + var threw = false; + + try{ + streams.serve( failingBus, 1, ( queue ) => {} ); + } catch( any e ){ + threw = ( e.type == "Test.SubscribeFailed" ); + } + + expect( threw ).toBeTrue(); + expect( streams.activeCount() ).toBe( 0 ); + } ); + + it( "a full queue drops the OLDEST event instead of growing or blocking", function(){ + var streams = newLiveStreams(); + var queue = streams.newQueue( 3 ); + + for( var n = 1; n <= 5; n++ ){ + streams.offer( queue, { ruleName: "e#n#" } ); + } + + expect( queue.size() ).toBe( 3 ); + expect( drainRuleNames( queue ).toList() ).toBe( "e3,e4,e5" ); + } ); + + it( "bounds each stream's queue at 1000 events by default", function(){ + expect( newLiveStreams().newQueue().remainingCapacity() ).toBe( 1000 ); + } ); + + it( "keeps a stalled subscriber bounded while the bus keeps publishing", function(){ + var streams = newLiveStreams(); + var bus = newStreamBus(); + var queue = streams.newQueue( 10 ); + var token = bus.subscribe( ( evt ) => streams.offer( queue, evt ) ); + getInstance( "InMemoryMetricsStore@rulebox" ).reset(); + + // Nothing ever drains the queue: a stalled client + for( var n = 1; n <= 50; n++ ){ + bus.publish( { rulebookName: "rb", ruleName: "e#n#", state: "EXECUTED", durationMs: 1, timestamp: "" } ); + } + bus.unsubscribe( token ); + getInstance( "InMemoryMetricsStore@rulebox" ).reset(); + + var kept = drainRuleNames( queue ); + expect( kept ).toHaveLength( 10 ); + expect( kept[ 1 ] ).toBe( "e41" ); + expect( kept[ 10 ] ).toBe( "e50" ); + } ); + + } ); } } diff --git a/test-harness/tests/specs/VisualizerHandlerSpec.bx b/test-harness/tests/specs/VisualizerHandlerSpec.bx index 4c47b512..5ff70a69 100644 --- a/test-harness/tests/specs/VisualizerHandlerSpec.bx +++ b/test-harness/tests/specs/VisualizerHandlerSpec.bx @@ -171,6 +171,65 @@ class extends="tests.resources.BaseSpec"{ expect( summary ).toHaveKey( "countsByState" ); } ); + describe( "stream() connection cap", function(){ + + // Fresh request context per spec: execute() otherwise reuses the previous spec's + // (e.g. the disabled-visualizer spec's 404 renderData). + beforeEach( function(){ + setup(); + } ); + + // The exact settings struct and LiveStreams singleton the handler reads, so holding + // slots here is seen by the real stream() action. + function hold( required numeric count ){ + var streams = getWireBox().getInstance( "LiveStreams@rulebox" ); + for( var i = 1; i <= arguments.count; i++ ){ + streams.tryAcquire( 1000 ); + } + } + + function releaseAll(){ + var streams = getWireBox().getInstance( "LiveStreams@rulebox" ); + while( streams.activeCount() > 0 ){ + streams.release(); + } + } + + it( "answers 503 JSON instead of opening a stream once visualizer.maxStreams are open", function(){ + var settings = getWireBox().getInstance( dsl="coldbox:moduleSettings:rulebox" ); + var original = settings.visualizer.maxStreams; + settings.visualizer.maxStreams = 2; + hold( 2 ); + + try{ + var event = execute( event="rulebox:visualizer.stream", renderResults=true ); + + expect( event.getRenderData().statusCode ).toBe( 503 ); + expect( jsonDeserialize( event.getRenderedContent() ).error ).toBe( "Too many live tracker connections" ); + // A rejected request must not hold a slot of its own + expect( getWireBox().getInstance( "LiveStreams@rulebox" ).activeCount() ).toBe( 2 ); + } finally { + releaseAll(); + settings.visualizer.maxStreams = original; + } + } ); + + it( "uses the default cap of 25 that ModuleConfig fills in", function(){ + var settings = getWireBox().getInstance( dsl="coldbox:moduleSettings:rulebox" ); + expect( settings.visualizer.maxStreams ).toBe( 25 ); + hold( 25 ); + + try{ + var event = execute( event="rulebox:visualizer.stream", renderResults=true ); + + expect( event.getRenderData().statusCode ).toBe( 503 ); + } finally { + releaseAll(); + } + } ); + + } ); + it( "404s every action while the visualizer is disabled", function(){ // Same DSL the handler/RuleEventBus are injected with, so this is guaranteed to be // the exact same live settings struct they read from. diff --git a/test-harness/tests/specs/VisualizerSettingsSpec.bx b/test-harness/tests/specs/VisualizerSettingsSpec.bx index 46666497..0e70dc09 100644 --- a/test-harness/tests/specs/VisualizerSettingsSpec.bx +++ b/test-harness/tests/specs/VisualizerSettingsSpec.bx @@ -94,7 +94,10 @@ class extends="tests.resources.BaseSpec"{ var badShapes = [ { setting: { enabled: "maybe" }, key: "visualizer.enabled" }, { setting: { enabled: true, metricsStore: " " }, key: "visualizer.metricsStore" }, - { setting: { enabled: true, datasourceName: {} }, key: "visualizer.datasourceName" } + { setting: { enabled: true, datasourceName: {} }, key: "visualizer.datasourceName" }, + { setting: { enabled: true, maxStreams: "lots" }, key: "visualizer.maxStreams" }, + { setting: { enabled: true, maxStreams: 0 }, key: "visualizer.maxStreams" }, + { setting: { enabled: true, maxStreams: 2.5 }, key: "visualizer.maxStreams" } ]; try{ for( var bad in badShapes ){