Source code for duct.sources.linux.basic

"""
.. module:: basic
   :platform: unix
   :synopsis: Basic linux host checks

.. moduleauthor:: Colin Alston <colin@tamvera.com>
"""
from zope.interface import implementer

from duct.interfaces import IDuctSource
from duct.objects import Source
from duct.aggregators import Counter64


[docs] @implementer(IDuctSource) class LoadAverage(Source): """Reports system load average for the current host **Metrics:** :(service name): Load average """ def _parse_loadaverage(self, data): la1 = data.split()[0] return self.createEvent('ok', 'Load average', float(la1))
[docs] async def get(self): with open('/proc/loadavg', 'rt', encoding='utf-8') as f: return self._parse_loadaverage(f.read())
[docs] async def sshGet(self): loadavg, err, code = await self.fork('/bin/cat /proc/loadavg') if code == 0: return self._parse_loadaverage(loadavg) else: raise RuntimeError(err)
[docs] @implementer(IDuctSource) class DiskIO(Source): """Reports disk IO statistics per device :param devices: List of devices to check (optional) :type devices: list. **Metrics:** :(service name).(device name).reads: Number of completed reads :(service name).(device name).read_bytes: Bytes read per second :(service name).(device name).read_latency: Disk read latency :(service name).(device name).writes: Number of completed writes :(service name).(device name).write_bytes: Bytes written per second :(service name).(device name).write_latency: Disk write latency """ def __init__(self, *a, **kw): Source.__init__(self, *a, **kw) self.devices = self.config.get('devices') self.tcache = {} self.trc = {} self.twc = {} def _parse_stats(self, stats): events = [] for s in stats: parts = s.strip().split() n = parts[2] # Filter things we don't care about if (n[:4] != 'loop') and (n[:3] != 'ram'): dname = "/dev/" + n if self.devices and (dname not in self.devices): continue nums = [int(i) for i in parts[3:]] reads, _readm, read_sec, read_t = nums[:4] writes, _writem, write_sec, write_t = nums[4:8] # Calculate the average latency of read/write ops if n in self.tcache: (last_r, last_w, last_rt, last_wt) = self.tcache[n] r_delta = float(reads - last_r) w_delta = float(writes - last_w) if r_delta > 0: read_lat = (read_t - last_rt)/float(reads - last_r) self.trc[n] = read_lat else: read_lat = self.trc.get(n, None) if w_delta > 0: write_lat = (write_t - last_wt)/float(writes - last_w) self.twc[n] = write_lat else: write_lat = self.twc.get(n, None) else: if reads > 0: read_lat = read_t / float(reads) self.trc[n] = read_lat else: read_lat = None if writes > 0: write_lat = write_t / float(writes) self.twc[n] = write_lat else: write_lat = None self.tcache[n] = (reads, writes, read_t, write_t) if read_lat: events.append(self.createEvent( 'ok', 'Read latency (ms)', read_lat, prefix=f'{dname}.read_latency' )) if write_lat: events.append(self.createEvent( 'ok', 'Write latency (ms)', write_lat, prefix=f'{dname}.write_latency' )) events.extend([ self.createEvent('ok', 'Reads', reads, prefix=f'{dname}.reads', aggregation=Counter64), self.createEvent('ok', 'Read Bps', read_sec * 512, prefix=f'{dname}.read_bytes', aggregation=Counter64), self.createEvent('ok', 'Writes', writes, prefix=f'{dname}.writes', aggregation=Counter64), self.createEvent('ok', 'Write Bps', write_sec * 512, prefix=f'{dname}.write_bytes', aggregation=Counter64), ]) return events
[docs] async def sshGet(self): diskstats, err, code = await self.fork('/bin/cat /proc/diskstats') if code == 0: stats = diskstats.strip('\n').split('\n') return self._parse_stats(stats) else: raise RuntimeError(err)
def _getstats(self): with open('/proc/diskstats', 'rt', encoding='utf-8') as f: stats = f.read() return stats.strip('\n').split('\n')
[docs] async def get(self): stats = self._getstats() return self._parse_stats(stats)
[docs] @implementer(IDuctSource) class CPU(Source): """Reports system CPU utilisation as a percentage/100 **Metrics:** :(service name): Percentage CPU utilisation :(service name).(type): Percentage CPU utilisation by type :(service name).coreX: Percentage CPU utilisation per core :(service name).coreX.(type): Percentage CPU utilisation per core by type """ cols = ['user', 'nice', 'system', 'idle', 'iowait', 'irq', 'softirq', 'steal', 'guest', 'guest_nice'] def __init__(self, *a): Source.__init__(self, *a) self.cpu = {} self.prev_total = {} self.prev_usage = {} def _read_proc_stat(self): cpus = [] with open('/proc/stat', 'rt', encoding='utf-8') as procstat: for l in procstat: if l.startswith('cpu'): cpus.append(l.strip('\n')) return cpus def _calculate_metrics(self, stat): cpuid = stat.split()[0] cpu = [int(i) for i in stat.split()[1:]] # We might not have all the virt-related numbers, so zero-pad. cpu = (cpu + [0, 0, 0])[:10] (user, nice, system, idle, iowait, irq, softirq, steal, _guest, _guestnice) = cpu usage = user + nice + system + irq + softirq + steal total = usage + iowait + idle if not self.cpu.get(cpuid): # No initial values, so set them and return no events. self.cpu[cpuid] = cpu self.prev_total[cpuid] = total self.prev_usage[cpuid] = usage return None total_diff = total - self.prev_total[cpuid] if total_diff != 0: metrics = [(None, (usage - self.prev_usage[cpuid]) / float(total_diff))] for i, name in enumerate(self.cols): prev = self.cpu[cpuid][i] cpu_m = (cpu[i] - prev) / float(total_diff) metrics.append((name, cpu_m)) self.cpu[cpuid] = cpu self.prev_total[cpuid] = total self.prev_usage[cpuid] = usage return metrics return None def _transpose_metrics(self, metrics, prefix): if metrics: events = [ self.createEvent('ok', f'CPU {name} {int(cpu_m * 100)}%', cpu_m, prefix=prefix+name) for name, cpu_m in metrics[1:] ] events.append(self.createEvent( 'ok', f'CPU {int(metrics[0][1] * 100)}%', metrics[0][1], prefix=prefix.rstrip('.'))) return events return None
[docs] async def sshGet(self): procstat, err, code = await self.fork('cat /proc/stat') if code == 0: metrics = [] stat = procstat.strip('\n').split('\n') for cpu in stat: if not cpu.startswith('cpu'): continue if cpu.split()[0] == 'cpu': prefix = "" else: prefix = 'core' + cpu.split()[0].strip('cpu') + '.' stats = self._calculate_metrics(cpu) if stats: metrics.extend(self._transpose_metrics(stats, prefix)) return metrics or None else: raise RuntimeError(err)
[docs] async def get(self): stat = self._read_proc_stat() metrics = [] for cpu in stat: if not cpu.startswith('cpu'): continue if cpu.split()[0] == 'cpu': prefix = "" else: prefix = 'core' + cpu.split()[0].strip('cpu') + '.' stats = self._calculate_metrics(cpu) if stats: metrics.extend(self._transpose_metrics(stats, prefix)) return metrics or None
[docs] @implementer(IDuctSource) class Memory(Source): """Reports system memory utilisation as a percentage/100 **Metrics:** :(service name): Percentage memory utilisation """ def _parse_stats(self, mem): dat = {} for l in mem: if ':' not in l: continue k, v = l.replace(':', '').split()[:2] dat[k] = int(v) if 'MemAvailable' in dat: free = dat['MemAvailable'] else: free = dat['MemFree'] + dat['Buffers'] + dat['Cached'] total = dat['MemTotal'] used = total - free return self.createEvent('ok', f'Memory {used}/{total}', used/float(total))
[docs] async def get(self): with open('/proc/meminfo', 'rt', encoding='utf-8') as f: return self._parse_stats(f)
[docs] async def sshGet(self): mem, err, code = await self.fork('/bin/cat /proc/meminfo') if code == 0: return self._parse_stats(mem.strip('\n').split('\n')) else: raise RuntimeError(err)
[docs] @implementer(IDuctSource) class DiskFree(Source): """Returns the free space for all mounted filesystems :param disks: List of devices to check (optional) :type disks: list. **Metrics:** :(service name).(device).used: Used space (%) :(service name).(device).bytes: Used space (kbytes) :(service name).(device).free: Free space (kbytes) """ ssh = True
[docs] async def get(self): disks = self.config.get('disks') out, _, _ = await self.fork('/bin/df', args=('-lPx', 'tmpfs',)) out = [i.split() for i in out.strip('\n').split('\n')[1:]] events = [] for disk, _size, used, free, util, _mnt in out: if disks and (disk not in disks): continue if disk != "udev": util = int(util.strip('%')) used = int(used) free = int(free) events.extend([ self.createEvent('ok', f'Disk percent used {util}%', util, prefix=f'{disk}.used'), self.createEvent('ok', f'Disk bytes used {used} kB', used, prefix=f'{disk}.bytes'), self.createEvent('ok', f'Disk free {free} kB', free, prefix=f'{disk}.free') ]) return events
[docs] @implementer(IDuctSource) class Network(Source): """Returns all network interface statistics :param interfaces: List of interfaces to check (optional) :type interfaces: list. **Metrics:** :(service name).(device).tx_bytes: Bytes transmitted :(service name).(device).tx_packets: Packets transmitted :(service name).(device).tx_errors: Errors :(service name).(device).rx_bytes: Bytes received :(service name).(device).rx_packets: Packets received :(service name).(device).rx_errors: Errors """ def _parse_stats(self, stats): ifaces = self.config.get('interfaces') ev = [] for stat in stats: items = stat.split() iface = items[0].strip(':') if ifaces and (iface not in ifaces): continue tx_bytes = int(items[1]) tx_packets = int(items[2]) tx_err = int(items[3]) rx_bytes = int(items[9]) rx_packets = int(items[10]) rx_err = int(items[11]) ev.extend([ self.createEvent('ok', f'Network {iface} TX bytes/sec', tx_bytes, prefix=f'{iface}.tx_bytes', aggregation=Counter64), self.createEvent('ok', f'Network {iface} TX packets/sec', tx_packets, prefix=f'{iface}.tx_packets', aggregation=Counter64), self.createEvent('ok', f'Network {iface} TX errors/sec', tx_err, prefix=f'{iface}.tx_errors', aggregation=Counter64), self.createEvent('ok', f'Network {iface} RX bytes/sec', rx_bytes, prefix=f'{iface}.rx_bytes', aggregation=Counter64), self.createEvent('ok', f'Network {iface} RX packets/sec', rx_packets, prefix=f'{iface}.rx_packets', aggregation=Counter64), self.createEvent('ok', f'Network {iface} RX errors/sec', rx_err, prefix=f'{iface}.rx_errors', aggregation=Counter64), ]) return ev def _readStats(self): with open('/proc/net/dev', 'rt', encoding='utf-8') as f: proc_dev = f.read() return proc_dev.strip('\n').split('\n')[2:]
[docs] async def sshGet(self): net, err, code = await self.fork('/bin/cat /proc/net/dev') if code == 0: return self._parse_stats(net.strip('\n').split('\n')[2:]) else: raise RuntimeError(err)
[docs] async def get(self): return self._parse_stats(self._readStats())