summaryrefslogtreecommitdiff
path: root/qpid/tools/src/py/qpid-stat
diff options
context:
space:
mode:
Diffstat (limited to 'qpid/tools/src/py/qpid-stat')
-rwxr-xr-xqpid/tools/src/py/qpid-stat452
1 files changed, 198 insertions, 254 deletions
diff --git a/qpid/tools/src/py/qpid-stat b/qpid/tools/src/py/qpid-stat
index a7272da3f1..bb094554e6 100755
--- a/qpid/tools/src/py/qpid-stat
+++ b/qpid/tools/src/py/qpid-stat
@@ -21,13 +21,18 @@
import os
from optparse import OptionParser, OptionGroup
-from time import sleep ### debug
import sys
import locale
import socket
import re
-from qmf.console import Session, Console
-from qpid.disp import Display, Header, Sorter
+from qpid.messaging import Connection
+
+home = os.environ.get("QPID_TOOLS_HOME", os.path.normpath("/usr/share/qpid-tools"))
+sys.path.append(os.path.join(home, "python"))
+
+from qpidtoollibs.broker import BrokerAgent
+from qpidtoollibs.disp import Display, Header, Sorter
+
class Config:
def __init__(self):
@@ -37,7 +42,7 @@ class Config:
self._limit = 50
self._increasing = False
self._sortcol = None
- self._cluster_detail = False
+ self._details = None
self._sasl_mechanism = None
config = Config()
@@ -56,24 +61,16 @@ def OptionsAndArguments(argv):
parser.add_option_group(group1)
group2 = OptionGroup(parser, "Display Options")
- group2.add_option("-b", "--broker", help="Show Brokers",
- action="store_const", const="b", dest="show")
- group2.add_option("-c", "--connections", help="Show Connections",
- action="store_const", const="c", dest="show")
- group2.add_option("-e", "--exchanges", help="Show Exchanges",
- action="store_const", const="e", dest="show")
- group2.add_option("-q", "--queues", help="Show Queues",
- action="store_const", const="q", dest="show")
- group2.add_option("-u", "--subscriptions", help="Show Subscriptions",
- action="store_const", const="u", dest="show")
- group2.add_option("-S", "--sort-by", metavar="<colname>",
- help="Sort by column name")
- group2.add_option("-I", "--increasing", action="store_true", default=False,
- help="Sort by increasing value (default = decreasing)")
- group2.add_option("-L", "--limit", type="int", default=50, metavar="<n>",
- help="Limit output to n rows")
- group2.add_option("-C", "--cluster", action="store_true", default=False,
- help="Display per-broker cluster detail.")
+ group2.add_option("-b", "--broker", help="Show Brokers", action="store_const", const="b", dest="show")
+ group2.add_option("-c", "--connections", help="Show Connections", action="store_const", const="c", dest="show")
+ group2.add_option("-e", "--exchanges", help="Show Exchanges", action="store_const", const="e", dest="show")
+ group2.add_option("-q", "--queues", help="Show Queues", action="store_const", const="q", dest="show")
+ group2.add_option("-u", "--subscriptions", help="Show Subscriptions", action="store_const", const="u", dest="show")
+ group2.add_option("-m", "--memory", help="Show Broker Memory Stats", action="store_const", const="m", dest="show")
+ group2.add_option("-S", "--sort-by", metavar="<colname>", help="Sort by column name")
+ group2.add_option("-I", "--increasing", action="store_true", default=False, help="Sort by increasing value (default = decreasing)")
+ group2.add_option("-L", "--limit", type="int", default=50, metavar="<n>", help="Limit output to n rows")
+ group2.add_option("-D", "--details", action="store", metavar="<name>", dest="detail", default=None, help="Display details on a single object.")
parser.add_option_group(group2)
opts, args = parser.parse_args(args=argv)
@@ -86,8 +83,8 @@ def OptionsAndArguments(argv):
config._connTimeout = opts.timeout
config._increasing = opts.increasing
config._limit = opts.limit
- config._cluster_detail = opts.cluster
config._sasl_mechanism = opts.sasl_mechanism
+ config._detail = opts.detail
if args:
config._host = args[0]
@@ -119,86 +116,26 @@ class IpAddr:
bestAddr = addrPort
return bestAddr
-class Broker(object):
- def __init__(self, qmf, broker):
- self.broker = broker
-
- agents = qmf.getAgents()
- for a in agents:
- if a.getAgentBank() == '0':
- self.brokerAgent = a
-
- bobj = qmf.getObjects(_class="broker", _package="org.apache.qpid.broker", _agent=self.brokerAgent)[0]
- self.currentTime = bobj.getTimestamps()[0]
- try:
- self.uptime = bobj.uptime
- except:
- self.uptime = 0
- self.connections = {}
- self.sessions = {}
- self.exchanges = {}
- self.queues = {}
- self.subscriptions = {}
- package = "org.apache.qpid.broker"
-
- list = qmf.getObjects(_class="connection", _package=package, _agent=self.brokerAgent)
- for conn in list:
- if not conn.shadow:
- self.connections[conn.getObjectId()] = conn
-
- list = qmf.getObjects(_class="session", _package=package, _agent=self.brokerAgent)
- for sess in list:
- if sess.connectionRef in self.connections:
- self.sessions[sess.getObjectId()] = sess
-
- list = qmf.getObjects(_class="exchange", _package=package, _agent=self.brokerAgent)
- for exchange in list:
- self.exchanges[exchange.getObjectId()] = exchange
-
- list = qmf.getObjects(_class="queue", _package=package, _agent=self.brokerAgent)
- for queue in list:
- self.queues[queue.getObjectId()] = queue
-
- list = qmf.getObjects(_class="subscription", _package=package, _agent=self.brokerAgent)
- for subscription in list:
- self.subscriptions[subscription.getObjectId()] = subscription
-
- def getName(self):
- return self.broker.getUrl()
-
- def getCurrentTime(self):
- return self.currentTime
-
- def getUptime(self):
- return self.uptime
-
-class BrokerManager(Console):
+class BrokerManager:
def __init__(self):
- self.brokerName = None
- self.qmf = None
- self.broker = None
- self.brokers = []
- self.cluster = None
+ self.brokerName = None
+ self.connections = []
+ self.brokers = []
+ self.cluster = None
def SetBroker(self, brokerUrl, mechanism):
self.url = brokerUrl
- self.qmf = Session()
- self.mechanism = mechanism
- self.broker = self.qmf.addBroker(brokerUrl, config._connTimeout, mechanism)
- agents = self.qmf.getAgents()
- for a in agents:
- if a.getAgentBank() == '0':
- self.brokerAgent = a
+ self.connections.append(Connection(self.url, sasl_mechanism=mechanism))
+ self.connections[0].open()
+ self.brokers.append(BrokerAgent(self.connections[0]))
def Disconnect(self):
""" Release any allocated brokers. Ignore any failures as the tool is
shutting down.
"""
try:
- if self.broker:
- self.qmf.delBroker(self.broker)
- else:
- for b in self.brokers: self.qmf.delBroker(b.broker)
+ for conn in self.connections:
+ conn.close()
except:
pass
@@ -238,62 +175,63 @@ class BrokerManager(Console):
hosts.append(bestUrl)
return hosts
- def displaySubs(self, subs, indent, broker=None, conn=None, sess=None, exchange=None, queue=None):
- if len(subs) == 0:
- return
- this = subs[0]
- remaining = subs[1:]
- newindent = indent + " "
- if this == 'b':
- pass
- elif this == 'c':
- if broker:
- for oid in broker.connections:
- iconn = broker.connections[oid]
- self.printConnSub(indent, broker.getName(), iconn)
- self.displaySubs(remaining, newindent, broker=broker, conn=iconn,
- sess=sess, exchange=exchange, queue=queue)
- elif this == 's':
- pass
- elif this == 'e':
- pass
- elif this == 'q':
- pass
- print
-
def displayBroker(self, subs):
disp = Display(prefix=" ")
heads = []
- heads.append(Header('broker'))
- heads.append(Header('cluster'))
heads.append(Header('uptime', Header.DURATION))
- heads.append(Header('conn', Header.KMG))
- heads.append(Header('sess', Header.KMG))
- heads.append(Header('exch', Header.KMG))
- heads.append(Header('queue', Header.KMG))
+ heads.append(Header('connections', Header.COMMAS))
+ heads.append(Header('sessions', Header.COMMAS))
+ heads.append(Header('exchanges', Header.COMMAS))
+ heads.append(Header('queues', Header.COMMAS))
rows = []
- for broker in self.brokers:
- if self.cluster:
- ctext = "%s(%s)" % (self.cluster.clusterName, self.cluster.status)
- else:
- ctext = "<standalone>"
- row = (broker.getName(), ctext, broker.getUptime(),
- len(broker.connections), len(broker.sessions),
- len(broker.exchanges), len(broker.queues))
- rows.append(row)
- title = "Brokers"
- if config._sortcol:
- sorter = Sorter(heads, rows, config._sortcol, config._limit, config._increasing)
- dispRows = sorter.getSorted()
- else:
- dispRows = rows
- disp.formattedTable(title, heads, dispRows)
+ broker = self.brokers[0].getBroker()
+ connections = self.getConnectionMap()
+ sessions = self.getSessionMap()
+ exchanges = self.getExchangeMap()
+ queues = self.getQueueMap()
+ row = (broker.getUpdateTime() - broker.getCreateTime(),
+ len(connections), len(sessions),
+ len(exchanges), len(queues))
+ rows.append(row)
+ disp.formattedTable('Broker Summary:', heads, rows)
+
+ if 'queueCount' not in broker.values:
+ return
+
+ print
+ heads = []
+ heads.append(Header('Statistic'))
+ heads.append(Header('Messages', Header.COMMAS))
+ heads.append(Header('Bytes', Header.COMMAS))
+ rows = []
+ rows.append(['queue-depth', broker.msgDepth, broker.byteDepth])
+ rows.append(['total-enqueues', broker.msgTotalEnqueues, broker.byteTotalEnqueues])
+ rows.append(['total-dequeues', broker.msgTotalDequeues, broker.byteTotalDequeues])
+ rows.append(['persistent-enqueues', broker.msgPersistEnqueues, broker.bytePersistEnqueues])
+ rows.append(['persistent-dequeues', broker.msgPersistDequeues, broker.bytePersistDequeues])
+ rows.append(['transactional-enqueues', broker.msgTxnEnqueues, broker.byteTxnEnqueues])
+ rows.append(['transactional-dequeues', broker.msgTxnDequeues, broker.byteTxnDequeues])
+ rows.append(['flow-to-disk-depth', broker.msgFtdDepth, broker.byteFtdDepth])
+ rows.append(['flow-to-disk-enqueues', broker.msgFtdEnqueues, broker.byteFtdEnqueues])
+ rows.append(['flow-to-disk-dequeues', broker.msgFtdDequeues, broker.byteFtdDequeues])
+ rows.append(['acquires', broker.acquires, None])
+ rows.append(['releases', broker.releases, None])
+ rows.append(['discards-no-route', broker.discardsNoRoute, None])
+ rows.append(['discards-ttl-expired', broker.discardsTtl, None])
+ rows.append(['discards-limit-overflow', broker.discardsOverflow, None])
+ rows.append(['discards-ring-overflow', broker.discardsRing, None])
+ rows.append(['discards-lvq-replace', broker.discardsLvq, None])
+ rows.append(['discards-subscriber-reject', broker.discardsSubscriber, None])
+ rows.append(['discards-purged', broker.discardsPurge, None])
+ rows.append(['reroutes', broker.reroutes, None])
+ rows.append(['abandoned', broker.abandoned, None])
+ rows.append(['abandoned-via-alt', broker.abandonedViaAlt, None])
+ disp.formattedTable('Aggregate Broker Statistics:', heads, rows)
+
def displayConn(self, subs):
disp = Display(prefix=" ")
heads = []
- if self.cluster:
- heads.append(Header('broker'))
heads.append(Header('client-addr'))
heads.append(Header('cproc'))
heads.append(Header('cpid'))
@@ -303,25 +241,20 @@ class BrokerManager(Console):
heads.append(Header('msgIn', Header.KMG))
heads.append(Header('msgOut', Header.KMG))
rows = []
- for broker in self.brokers:
- for oid in broker.connections:
- conn = broker.connections[oid]
- row = []
- if self.cluster:
- row.append(broker.getName())
- row.append(conn.address)
- row.append(conn.remoteProcessName)
- row.append(conn.remotePid)
- row.append(conn.authIdentity)
- row.append(broker.getCurrentTime() - conn.getTimestamps()[1])
- idle = broker.getCurrentTime() - conn.getTimestamps()[0]
- row.append(broker.getCurrentTime() - conn.getTimestamps()[0])
- row.append(conn.msgsFromClient)
- row.append(conn.msgsToClient)
- rows.append(row)
+ connections = self.brokers[0].getAllConnections()
+ broker = self.brokers[0].getBroker()
+ for conn in connections:
+ row = []
+ row.append(conn.address)
+ row.append(conn.remoteProcessName)
+ row.append(conn.remotePid)
+ row.append(conn.authIdentity)
+ row.append(broker.getUpdateTime() - conn.getCreateTime())
+ row.append(broker.getUpdateTime() - conn.getUpdateTime())
+ row.append(conn.msgsFromClient)
+ row.append(conn.msgsToClient)
+ rows.append(row)
title = "Connections"
- if self.cluster:
- title += " for cluster '%s'" % self.cluster.clusterName
if config._sortcol:
sorter = Sorter(heads, rows, config._sortcol, config._limit, config._increasing)
dispRows = sorter.getSorted()
@@ -335,8 +268,6 @@ class BrokerManager(Console):
def displayExchange(self, subs):
disp = Display(prefix=" ")
heads = []
- if self.cluster:
- heads.append(Header('broker'))
heads.append(Header("exchange"))
heads.append(Header("type"))
heads.append(Header("dur", Header.Y))
@@ -348,26 +279,21 @@ class BrokerManager(Console):
heads.append(Header("byteOut", Header.KMG))
heads.append(Header("byteDrop", Header.KMG))
rows = []
- for broker in self.brokers:
- for oid in broker.exchanges:
- ex = broker.exchanges[oid]
- row = []
- if self.cluster:
- row.append(broker.getName())
- row.append(ex.name)
- row.append(ex.type)
- row.append(ex.durable)
- row.append(ex.bindingCount)
- row.append(ex.msgReceives)
- row.append(ex.msgRoutes)
- row.append(ex.msgDrops)
- row.append(ex.byteReceives)
- row.append(ex.byteRoutes)
- row.append(ex.byteDrops)
- rows.append(row)
+ exchanges = self.brokers[0].getAllExchanges()
+ for ex in exchanges:
+ row = []
+ row.append(ex.name)
+ row.append(ex.type)
+ row.append(ex.durable)
+ row.append(ex.bindingCount)
+ row.append(ex.msgReceives)
+ row.append(ex.msgRoutes)
+ row.append(ex.msgDrops)
+ row.append(ex.byteReceives)
+ row.append(ex.byteRoutes)
+ row.append(ex.byteDrops)
+ rows.append(row)
title = "Exchanges"
- if self.cluster:
- title += " for cluster '%s'" % self.cluster.clusterName
if config._sortcol:
sorter = Sorter(heads, rows, config._sortcol, config._limit, config._increasing)
dispRows = sorter.getSorted()
@@ -375,11 +301,9 @@ class BrokerManager(Console):
dispRows = rows
disp.formattedTable(title, heads, dispRows)
- def displayQueue(self, subs):
+ def displayQueues(self, subs):
disp = Display(prefix=" ")
heads = []
- if self.cluster:
- heads.append(Header('broker'))
heads.append(Header("queue"))
heads.append(Header("dur", Header.Y))
heads.append(Header("autoDel", Header.Y))
@@ -393,28 +317,23 @@ class BrokerManager(Console):
heads.append(Header("cons", Header.KMG))
heads.append(Header("bind", Header.KMG))
rows = []
- for broker in self.brokers:
- for oid in broker.queues:
- q = broker.queues[oid]
- row = []
- if self.cluster:
- row.append(broker.getName())
- row.append(q.name)
- row.append(q.durable)
- row.append(q.autoDelete)
- row.append(q.exclusive)
- row.append(q.msgDepth)
- row.append(q.msgTotalEnqueues)
- row.append(q.msgTotalDequeues)
- row.append(q.byteDepth)
- row.append(q.byteTotalEnqueues)
- row.append(q.byteTotalDequeues)
- row.append(q.consumerCount)
- row.append(q.bindingCount)
- rows.append(row)
+ queues = self.brokers[0].getAllQueues()
+ for q in queues:
+ row = []
+ row.append(q.name)
+ row.append(q.durable)
+ row.append(q.autoDelete)
+ row.append(q.exclusive)
+ row.append(q.msgDepth)
+ row.append(q.msgTotalEnqueues)
+ row.append(q.msgTotalDequeues)
+ row.append(q.byteDepth)
+ row.append(q.byteTotalEnqueues)
+ row.append(q.byteTotalDequeues)
+ row.append(q.consumerCount)
+ row.append(q.bindingCount)
+ rows.append(row)
title = "Queues"
- if self.cluster:
- title += " for cluster '%s'" % self.cluster.clusterName
if config._sortcol:
sorter = Sorter(heads, rows, config._sortcol, config._limit, config._increasing)
dispRows = sorter.getSorted()
@@ -422,46 +341,46 @@ class BrokerManager(Console):
dispRows = rows
disp.formattedTable(title, heads, dispRows)
+ def displayQueue(self, subs):
+ disp = Display(prefix=" ")
+ heads = []
+
def displaySubscriptions(self, subs):
disp = Display(prefix=" ")
heads = []
- if self.cluster:
- heads.append(Header('broker'))
- heads.append(Header("subscription"))
+ heads.append(Header("subscr"))
heads.append(Header("queue"))
- heads.append(Header("connection"))
- heads.append(Header("processName"))
- heads.append(Header("processId"))
- heads.append(Header("browsing", Header.Y))
- heads.append(Header("acknowledged", Header.Y))
- heads.append(Header("exclusive", Header.Y))
+ heads.append(Header("conn"))
+ heads.append(Header("procName"))
+ heads.append(Header("procId"))
+ heads.append(Header("browse", Header.Y))
+ heads.append(Header("acked", Header.Y))
+ heads.append(Header("excl", Header.Y))
heads.append(Header("creditMode"))
heads.append(Header("delivered", Header.KMG))
rows = []
- for broker in self.brokers:
- for oid in broker.subscriptions:
- s = broker.subscriptions[oid]
- row = []
- try:
- if self.cluster:
- row.append(broker.getName())
- row.append(s.name)
- row.append(self.qmf.getObjects(_objectId=s.queueRef)[0].name)
- connectionRef = self.qmf.getObjects(_objectId=s.sessionRef)[0].connectionRef
- row.append(self.qmf.getObjects(_objectId=connectionRef)[0].address)
- row.append(self.qmf.getObjects(_objectId=connectionRef)[0].remoteProcessName)
- row.append(self.qmf.getObjects(_objectId=connectionRef)[0].remotePid)
- row.append(s.browsing)
- row.append(s.acknowledged)
- row.append(s.exclusive)
- row.append(s.creditMode)
- row.append(s.delivered)
- rows.append(row)
- except:
- pass
+ subscriptions = self.brokers[0].getAllSubscriptions()
+ sessions = self.getSessionMap()
+ connections = self.getConnectionMap()
+ for s in subscriptions:
+ row = []
+ try:
+ row.append(s.name)
+ row.append(s.queueRef)
+ session = sessions[s.sessionRef]
+ connection = connections[session.connectionRef]
+ row.append(connection.address)
+ row.append(connection.remoteProcessName)
+ row.append(connection.remotePid)
+ row.append(s.browsing)
+ row.append(s.acknowledged)
+ row.append(s.exclusive)
+ row.append(s.creditMode)
+ row.append(s.delivered)
+ rows.append(row)
+ except:
+ pass
title = "Subscriptions"
- if self.cluster:
- title += " for cluster '%s'" % self.cluster.clusterName
if config._sortcol:
sorter = Sorter(heads, rows, config._sortcol, config._limit, config._increasing)
dispRows = sorter.getSorted()
@@ -469,33 +388,58 @@ class BrokerManager(Console):
dispRows = rows
disp.formattedTable(title, heads, dispRows)
+ def displayMemory(self, unused):
+ disp = Display(prefix=" ")
+ heads = [Header('Statistic'), Header('Value', Header.COMMAS)]
+ rows = []
+ memory = self.brokers[0].getMemory()
+ for k,v in memory.values.items():
+ if k != 'name':
+ rows.append([k, v])
+ disp.formattedTable('Broker Memory Statistics:', heads, rows)
+
+ def getExchangeMap(self):
+ exchanges = self.brokers[0].getAllExchanges()
+ emap = {}
+ for e in exchanges:
+ emap[e.name] = e
+ return emap
+
+ def getQueueMap(self):
+ queues = self.brokers[0].getAllQueues()
+ qmap = {}
+ for q in queues:
+ qmap[q.name] = q
+ return qmap
+
+ def getSessionMap(self):
+ sessions = self.brokers[0].getAllSessions()
+ smap = {}
+ for s in sessions:
+ smap[s.name] = s
+ return smap
+
+ def getConnectionMap(self):
+ connections = self.brokers[0].getAllConnections()
+ cmap = {}
+ for c in connections:
+ cmap[c.address] = c
+ return cmap
+
def displayMain(self, main, subs):
if main == 'b': self.displayBroker(subs)
elif main == 'c': self.displayConn(subs)
elif main == 's': self.displaySession(subs)
elif main == 'e': self.displayExchange(subs)
- elif main == 'q': self.displayQueue(subs)
+ elif main == 'q':
+ if config._detail:
+ self.displayQueue(subs, config._detail)
+ else:
+ self.displayQueues(subs)
elif main == 'u': self.displaySubscriptions(subs)
+ elif main == 'm': self.displayMemory(subs)
def display(self):
- if config._cluster_detail or config._types[0] == 'b':
- # always show cluster detail when dumping broker stats
- self._getCluster()
- if self.cluster:
- memberList = self.cluster.members.split(";")
- hostList = self._getHostList(memberList)
- self.qmf.delBroker(self.broker)
- self.broker = None
- if config._host.find("@") > 0:
- authString = config._host.split("@")[0] + "@"
- else:
- authString = ""
- for host in hostList:
- b = self.qmf.addBroker(authString + host, config._connTimeout)
- self.brokers.append(Broker(self.qmf, b))
- else:
- self.brokers.append(Broker(self.qmf, self.broker))
-
self.displayMain(config._types[0], config._types[1:])