diff --git a/package.json b/package.json index a946eda1..e1cfa4a7 100644 --- a/package.json +++ b/package.json @@ -48,6 +48,7 @@ "ipc-rpc": "~0.1.3", "binary-search": "~1.2.0", "compression": "~1.4.0", + "prometheus-client": "~0.1.0", "debug": "~2.2.0" }, "devDependencies": { diff --git a/src/streammachine/index.coffee b/src/streammachine/index.coffee index cd40d2b9..34464748 100644 --- a/src/streammachine/index.coffee +++ b/src/streammachine/index.coffee @@ -27,6 +27,7 @@ module.exports = class StreamMachine behind_proxy: false + prometheus: true + admin: require_auth: false - diff --git a/src/streammachine/master/index.coffee b/src/streammachine/master/index.coffee index c3a010c0..c0eebdf4 100644 --- a/src/streammachine/master/index.coffee +++ b/src/streammachine/master/index.coffee @@ -1,7 +1,7 @@ -_ = require "underscore" -temp = require "temp" -net = require "net" -fs = require "fs" +_ = require "underscore" +temp = require "temp" +net = require "net" +fs = require "fs" express = require "express" Throttle = require "throttle" @@ -14,6 +14,7 @@ Alerts = require "../alerts" Analytics = require "./analytics" Monitoring = require "./monitoring" SlaveIO = require "./master_io" +Prometheus = require "./prometheus" RewindDumpRestore = require "../rewind/dump_restore" @@ -105,6 +106,11 @@ module.exports = class Master extends require("events").EventEmitter @monitoring = new Monitoring @, @log.child(module:"monitoring") + # -- Prometheus metrics -- # + + if @options.prometheus + @prometheus = new Prometheus @ + #---------- once_configured: (cb) -> diff --git a/src/streammachine/master/prometheus.coffee b/src/streammachine/master/prometheus.coffee new file mode 100644 index 00000000..d14559e0 --- /dev/null +++ b/src/streammachine/master/prometheus.coffee @@ -0,0 +1,57 @@ +Prometheus = require "prometheus-client" +express = require "express" + +module.exports = class PrometheusMaster + constructor: (@master) -> + @client = new Prometheus() + + # -- Register our Metrics -- # + + @connected_sources = @client.newGauge + namespace: "streammachine", + subsystem: "master" + name: "stream_sources" + help: "Number of sources connected for this stream." + + @source_latency = @client.newGauge + namespace: "streammachine", + subsystem: "master" + name: "stream_source_latency" + help: "How many milliseconds is the source chunk ts behind our clock time?" + + @connected_slaves = @client.newGauge + namespace: "streammachine", + subsystem: "master" + name: "slaves" + help: "Number of slaves that are connected." + + # -- Set our source loop -- # + + @_sourceInt = setInterval => + _process = (stream,idx,source) => + latency = + if source.last_ts + Number(new Date()) - Number(source.last_ts) + else + -1 + + @source_latency.set stream:stream, index:idx, latency + + for k,stream of @master.streams + _process(stream.key,idx,source) for source,idx in stream.sources + @connected_sources.set stream:stream.key, stream.sources.length + + for k,sg of @master.stream_groups + _process(sg._stream.key,idx,source) for source,idx in sg._stream.sources + @connected_sources.set stream:sg._stream.key, sg._stream.sources.length + + , 1000 + + @_slaveInt = setInterval => + @connected_slaves.set {}, Object.keys(@master.slaves?.slaves||{}).length + , 1000 + + # -- Attach metrics to the API -- # + + @app = express() + @app.get "/", @client.metricsFunc() diff --git a/src/streammachine/modes/master.coffee b/src/streammachine/modes/master.coffee index 42e1c9f1..f024ef51 100644 --- a/src/streammachine/modes/master.coffee +++ b/src/streammachine/modes/master.coffee @@ -32,6 +32,8 @@ module.exports = class MasterMode extends require("./base") @server.use "/s", @master.transport.app @server.use "/api", @master.api.app + @server.use "/metrics", @master.prometheus.app if @master.prometheus + if process.send? @_rpc = new RPC process, functions: OK: (msg,handle,cb) ->