mirror of
https://github.com/xcat2/confluent.git
synced 2026-09-29 00:31:09 +00:00
Compare commits
116 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 20457bbe37 | |||
| 206a79a976 | |||
| 99f9e852da | |||
| 6ab176218e | |||
| 5ccfa36da6 | |||
| 5b1e144d32 | |||
| 4722c3ec92 | |||
| d75867050c | |||
| 213d440052 | |||
| 0800290c8e | |||
| c5c5b36536 | |||
| 87a7e65b42 | |||
| 51c09d844f | |||
| 2c4f8dfceb | |||
| 3f53c55a66 | |||
| 598ec4a294 | |||
| 501ab64e18 | |||
| 2936c7e8fd | |||
| 4f85ba2bff | |||
| 5232b7c9c4 | |||
| f964fd8ce1 | |||
| f97fd3105f | |||
| bc03da47af | |||
| bd39171611 | |||
| ed050b37e1 | |||
| 8d1d19d9a8 | |||
| 017f3fb372 | |||
| 46518f890b | |||
| 7e86a72872 | |||
| 2567503662 | |||
| 6b56181a52 | |||
| c18ce50138 | |||
| a0684520d8 | |||
| 374aa49016 | |||
| 0b95daa30d | |||
| d33365195b | |||
| 3429173c27 | |||
| f6c44922f8 | |||
| a86d962984 | |||
| 9ee29aabe1 | |||
| a413f321fe | |||
| f2bd796c2a | |||
| bf31c4872f | |||
| 634e5a8944 | |||
| 67e3530d16 | |||
| 3c26beda1d | |||
| e2d0e49fc7 | |||
| da5a34c2e4 | |||
| 3629cb8ee7 | |||
| 8233e0a5bd | |||
| eae7b3bd80 | |||
| 868367e052 | |||
| 6289cfaac4 | |||
| f6d4fef5e6 | |||
| b1b7ec4d50 | |||
| c0cd6de4f7 | |||
| 4437e81e04 | |||
| 6a12af1242 | |||
| 9879a83a10 | |||
| cce6b824de | |||
| ce1cb952e8 | |||
| c6812274e4 | |||
| 7cd7068dd7 | |||
| 48f0330568 | |||
| 66e1d17d28 | |||
| 7480494432 | |||
| 49c00bfbb7 | |||
| 201985dd0e | |||
| 1aee19997a | |||
| 3bc366bef4 | |||
| 4c83a1a04e | |||
| cfae28a869 | |||
| 44e6a72847 | |||
| 006fdc8280 | |||
| 895b5264f6 | |||
| 0b577af1ca | |||
| ff0b1bba7f | |||
| 0badd9e5b4 | |||
| c02064f0a5 | |||
| c1b82d8163 | |||
| 0d5fa7a98a | |||
| 968efe719a | |||
| 7a63ca8759 | |||
| a24866c2df | |||
| c666b11138 | |||
| 22f6198f60 | |||
| c99d01dffc | |||
| 8d0028a1de | |||
| bb9c2297c9 | |||
| 91fa5bd1eb | |||
| ac9609c40d | |||
| 0c4cb49c20 | |||
| 4be4100014 | |||
| 9f7c8c69f2 | |||
| c35f7d99f7 | |||
| 445950d02a | |||
| cf72cf2d8c | |||
| 0652a7321b | |||
| 4c8ba92856 | |||
| 09582d7597 | |||
| 8a9e9aa7b3 | |||
| b766e7b0ee | |||
| 92699e47f2 | |||
| 2aa9910d83 | |||
| 18b6398c64 | |||
| 47b68e4258 | |||
| 79b6d099ab | |||
| 604ebcde3b | |||
| b4b733a573 | |||
| 3bf083deb3 | |||
| 9d770632ce | |||
| 79afd174c9 | |||
| 13a0bf4fbe | |||
| f1e1d9804a | |||
| 546296ce71 | |||
| 954b2dd15c |
+4
-1
@@ -1 +1,4 @@
|
||||
.idea
|
||||
*.pyc
|
||||
.*.
|
||||
confluent_client/man/man*
|
||||
.*.sw*
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
#!/usr/bin/python
|
||||
import os
|
||||
import sys
|
||||
path = os.path.dirname(os.path.realpath(__file__))
|
||||
try:
|
||||
sys.path.remove(path)
|
||||
except Exception:
|
||||
pass
|
||||
path = os.path.realpath(os.path.join(path, '..', 'confluent_server'))
|
||||
sys.path.append(path)
|
||||
|
||||
import confluent.config.attributes as attr
|
||||
import shutil
|
||||
|
||||
shutil.copyfile('doc/man/nodeattrib.ronn.tmpl', 'doc/man/nodeattrib.ronn')
|
||||
shutil.copyfile('doc/man/nodegroupattrib.ronn.tmpl', 'doc/man/nodegroupattrib.ronn')
|
||||
with open('doc/man/nodeattrib.ronn', 'a') as outf:
|
||||
for field in sorted(attr.node):
|
||||
outf.write('\n* `{0}`:\n {1}\n'.format(field, attr.node[field]['description']))
|
||||
with open('doc/man/nodegroupattrib.ronn', 'a') as outf:
|
||||
for field in sorted(attr.node):
|
||||
outf.write('\n* `{0}`:\n {1}\n'.format(field, attr.node[field]['description']))
|
||||
|
||||
|
||||
@@ -43,6 +43,8 @@ argparser.add_option('-d', '--diff', action='store_true',
|
||||
'output group and others')
|
||||
argparser.add_option('-w', '--watch', action='store_true',
|
||||
help='Show intermediate results while running')
|
||||
argparser.add_option('-g', '--groupcount', action='store_true',
|
||||
help='Show count of output groups rather than the actual output')
|
||||
argparser.add_option('-s', '--skipcommon', action='store_true',
|
||||
help='Do not print most common result, only non modal '
|
||||
'groups, useful when combined with -d')
|
||||
@@ -69,6 +71,9 @@ def print_current():
|
||||
if options.diff:
|
||||
grouped.print_deviants(skipmodal=options.skipcommon, count=options.count,
|
||||
reverse=options.reverse, basenode=options.base)
|
||||
elif options.groupcount:
|
||||
grouped.generate_byoutput()
|
||||
print(len(grouped.byoutput))
|
||||
else:
|
||||
grouped.print_all(skipmodal=options.skipcommon,
|
||||
count=options.count,
|
||||
|
||||
@@ -76,6 +76,11 @@ import confluent.termhandler as termhandler
|
||||
import confluent.tlvdata as tlvdata
|
||||
import confluent.client as client
|
||||
|
||||
try:
|
||||
unicode
|
||||
except NameError:
|
||||
unicode = str
|
||||
|
||||
conserversequence = '\x05c' # ctrl-e, c
|
||||
clearpowermessage = False
|
||||
|
||||
@@ -299,8 +304,11 @@ currchildren = None
|
||||
|
||||
|
||||
def print_result(res):
|
||||
global exitcode
|
||||
if 'errorcode' in res or 'error' in res:
|
||||
print(res['error'])
|
||||
if 'errorcode' in res:
|
||||
exitcode |= res['errorcode']
|
||||
return
|
||||
if 'databynode' in res:
|
||||
print_result(res['databynode'])
|
||||
@@ -957,7 +965,7 @@ def main():
|
||||
try:
|
||||
server_connect()
|
||||
connected = True
|
||||
except (socket.gaierror, socket.error):
|
||||
except Exception:
|
||||
pass
|
||||
if not connected:
|
||||
time.sleep(1)
|
||||
|
||||
@@ -47,6 +47,8 @@ exitcode = 0
|
||||
errorNodes = set([])
|
||||
session.stop_if_noderange_over(noderange, options.maxnodes)
|
||||
success = session.simple_noderange_command(noderange, 'configuration/management_controller/reset', 'reset', key='state', errnodes=errorNodes) # = 0 if successful
|
||||
if success != 0:
|
||||
sys.exit(success)
|
||||
|
||||
# Determine which nodes were successful and print them
|
||||
|
||||
|
||||
@@ -58,8 +58,8 @@ except IndexError:
|
||||
sys.exit(1)
|
||||
client.check_globbing(noderange)
|
||||
bootdev = None
|
||||
if len(sys.argv) > 2:
|
||||
bootdev = sys.argv[2]
|
||||
if len(args) > 1:
|
||||
bootdev = args[1]
|
||||
if bootdev in ('net', 'pxe'):
|
||||
bootdev = 'network'
|
||||
session = client.Command()
|
||||
|
||||
@@ -54,6 +54,12 @@ argparser.add_option('-d', '--detail', dest='detail',
|
||||
action='store_true', default=False,
|
||||
help='Provide verbose information as available, such as '
|
||||
'help text and possible valid values')
|
||||
argparser.add_option('-e', '--extra', dest='extra',
|
||||
action='store_true', default=False,
|
||||
help='Access extra configuration. Extra configuration is generally '
|
||||
'reserved for unpopular or redundant options that may be slow to '
|
||||
'read. Notably the IMM category on Lenovo settings is considered '
|
||||
'to be extra configuration')
|
||||
argparser.add_option('-x', '--exclude', dest='exclude',
|
||||
action='store_true', default=False,
|
||||
help='Treat positional arguments as items to not '
|
||||
@@ -71,7 +77,7 @@ argparser.add_option('-r', '--restoredefault', default=False,
|
||||
argparser.add_option('-m', '--maxnodes', type='int',
|
||||
help='Specify a maximum number of '
|
||||
'nodes to configure, '
|
||||
'prompting if over the threshold')
|
||||
'prompting if over the threshold')
|
||||
(options, args) = argparser.parse_args()
|
||||
|
||||
cfgpaths = {
|
||||
@@ -103,6 +109,7 @@ assignment = {}
|
||||
queryparms = {}
|
||||
printsys = []
|
||||
printbmc = []
|
||||
printextbmc = []
|
||||
printallbmc = False
|
||||
setsys = {}
|
||||
forceset = False
|
||||
@@ -177,11 +184,17 @@ def parse_config_line(arguments):
|
||||
del queryparms[path]
|
||||
except KeyError:
|
||||
pass
|
||||
if not matchedparms:
|
||||
if param.lower() == 'imm':
|
||||
printextbmc.append(param)
|
||||
options.extra = True
|
||||
elif not matchedparms:
|
||||
printsys.append(param)
|
||||
elif param not in cfgpaths:
|
||||
if param.startswith('bmc.'):
|
||||
printbmc.append(param.replace('bmc.', ''))
|
||||
elif param.lower().startswith('imm'):
|
||||
options.extra = True
|
||||
printextbmc.append(param)
|
||||
else:
|
||||
printsys.append(param)
|
||||
else:
|
||||
@@ -275,10 +288,15 @@ else:
|
||||
NullOpt(), queryparms[path])
|
||||
if rc:
|
||||
sys.exit(rc)
|
||||
if printbmc or printallbmc:
|
||||
rcode = client.print_attrib_path(
|
||||
'/noderange/{0}/configuration/management_controller/extended/all'.format(noderange),
|
||||
session, printbmc, options, attrprefix='bmc.')
|
||||
if printsys == 'all' or printextbmc or printbmc or printallbmc:
|
||||
if printbmc or not printextbmc:
|
||||
rcode = client.print_attrib_path(
|
||||
'/noderange/{0}/configuration/management_controller/extended/all'.format(noderange),
|
||||
session, printbmc, options, attrprefix='bmc.')
|
||||
if options.extra:
|
||||
rcode |= client.print_attrib_path(
|
||||
'/noderange/{0}/configuration/management_controller/extended/extra'.format(noderange),
|
||||
session, printextbmc, options)
|
||||
if printsys or options.exclude:
|
||||
if printsys == 'all':
|
||||
printsys = []
|
||||
@@ -287,7 +305,6 @@ else:
|
||||
else:
|
||||
path = '/noderange/{0}/configuration/system/advanced'.format(
|
||||
noderange)
|
||||
|
||||
rcode = client.print_attrib_path(path, session, printsys,
|
||||
options)
|
||||
sys.exit(rcode)
|
||||
|
||||
@@ -51,27 +51,28 @@ if options.tile:
|
||||
nodes.append(node)
|
||||
initial = True
|
||||
pane = 0
|
||||
sessname = 'nodeconsole_{0}'.format(os.getpid())
|
||||
for node in sortutil.natural_sort(nodes):
|
||||
panename = '{0}:{1}'.format(sessname, pane)
|
||||
if initial:
|
||||
initial = False
|
||||
subprocess.call(
|
||||
['tmux', 'new-session', '-d', '-s',
|
||||
'nodeconsole_{0}'.format(os.getpid()), '-x', '800', '-y',
|
||||
sessname, '-x', '800', '-y',
|
||||
'800', '{0} -m 5 start /nodes/{1}/console/session'.format(
|
||||
confettypath, node)])
|
||||
else:
|
||||
subprocess.call(['tmux', 'select-pane', '-t', str(pane)])
|
||||
subprocess.call(['tmux', 'set-option', 'pane-border-status', 'top'], stderr=null)
|
||||
pane += 1
|
||||
subprocess.call(['tmux', 'select-pane', '-t', sessname])
|
||||
subprocess.call(['tmux', 'set-option', '-t', panename, 'pane-border-status', 'top'], stderr=null)
|
||||
subprocess.call(
|
||||
['tmux', 'split', '-h',
|
||||
['tmux', 'split', '-h', '-t', sessname,
|
||||
'{0} -m 5 start /nodes/{1}/console/session'.format(
|
||||
confettypath, node)])
|
||||
subprocess.call(['tmux', 'select-layout', 'tiled'], stdout=null)
|
||||
subprocess.call(['tmux', 'select-pane', '-t', '0'])
|
||||
subprocess.call(['tmux', 'set-option', 'pane-border-status', 'top'], stderr=null)
|
||||
os.execlp('tmux', 'tmux', 'attach', '-t', 'nodeconsole_{0}'.format(
|
||||
os.getpid()))
|
||||
subprocess.call(['tmux', 'select-layout', '-t', sessname, 'tiled'], stdout=null)
|
||||
pane += 1
|
||||
subprocess.call(['tmux', 'select-pane', '-t', sessname])
|
||||
subprocess.call(['tmux', 'set-option', '-t', panename, 'pane-border-status', 'top'], stderr=null)
|
||||
os.execlp('tmux', 'tmux', 'attach', '-t', sessname)
|
||||
else:
|
||||
os.execl(confettypath, confettypath, 'start',
|
||||
'/nodes/{0}/console/session'.format(args[0]))
|
||||
|
||||
@@ -237,10 +237,12 @@ def clear_discovery(options, session):
|
||||
else:
|
||||
print(repr(res))
|
||||
|
||||
def list_matching_macs(options, session, node=None):
|
||||
def list_matching_macs(options, session, node=None, checknode=True):
|
||||
path = '/discovery/'
|
||||
if node:
|
||||
path += 'by-node/{0}/'.format(node)
|
||||
elif checknode and options.node:
|
||||
path += 'by-node/{0}/'.format(options.node)
|
||||
if options.model:
|
||||
path += 'by-model/{0}/'.format(options.model)
|
||||
if options.serial:
|
||||
@@ -277,11 +279,11 @@ def assign_discovery(options, session, needid=True):
|
||||
abort = True
|
||||
if abort:
|
||||
sys.exit(1)
|
||||
matches = list_matching_macs(options, session, None if needid else options.node)
|
||||
matches = list_matching_macs(options, session, None if needid else options.node, False)
|
||||
if not matches:
|
||||
# Do a rescan to catch missing requested data
|
||||
blocking_scan(session)
|
||||
matches = list_matching_macs(options, session, None if needid else options.node)
|
||||
matches = list_matching_macs(options, session, None if needid else options.node, False)
|
||||
if not matches:
|
||||
sys.stderr.write("No matching discovery candidates found\n")
|
||||
sys.exit(1)
|
||||
|
||||
@@ -65,10 +65,6 @@ client.check_globbing(noderange)
|
||||
|
||||
def install_license(session, filename):
|
||||
global exitcode
|
||||
if not os.path.exists(filename):
|
||||
sys.stderr.write('Unable to locate requested file {0}\n'.format(
|
||||
filename))
|
||||
sys.exit(404)
|
||||
resource = '/noderange/{0}/configuration/' \
|
||||
'management_controller/licenses/'.format(noderange)
|
||||
filename = os.path.abspath(filename)
|
||||
@@ -87,18 +83,16 @@ def save_licenses(session, dirname):
|
||||
resource = '/noderange/{0}/configuration/' \
|
||||
'management_controller/save_licenses'.format(noderange)
|
||||
filename = os.path.abspath(dirname)
|
||||
if not os.path.exists(filename):
|
||||
sys.stderr.write('Unable to locate specified directory {0}\n'.format(
|
||||
filename))
|
||||
sys.exit(404)
|
||||
instargs = {'dirname': filename}
|
||||
for res in session.create(resource, instargs):
|
||||
for node in res.get('databynode', {}):
|
||||
fname = res['databynode'][node].get('filename', None)
|
||||
if fname:
|
||||
print('{0}: Saved license to {1}'.format(node, fname))
|
||||
elif 'error' in res['databynode'][node]:
|
||||
sys.stderr.write('{0}: {1}\n'.format(node, res['databynode'][node]['error']))
|
||||
else:
|
||||
sys.stderr.write('{0}: {1}', node, repr(res['databynode'][node]))
|
||||
sys.stderr.write('{0}: {1}\n'.format(node, repr(res['databynode'][node])))
|
||||
|
||||
|
||||
def show_licenses(session):
|
||||
|
||||
@@ -46,6 +46,8 @@ def run():
|
||||
help='Number of commands to run at a time')
|
||||
argparser.add_option('-n', '--nonodeprefix', action='store_true',
|
||||
help='Do not prefix output with node names')
|
||||
argparser.add_option('-p', '--port', type='int', default=0,
|
||||
help='Specify a custom port for ssh')
|
||||
argparser.add_option('-m', '--maxnodes', type='int',
|
||||
help='Specify a maximum number of '
|
||||
'nodes to run remote ssh command to, '
|
||||
@@ -59,7 +61,7 @@ def run():
|
||||
sys.exit(1)
|
||||
client.check_globbing(args[0])
|
||||
concurrentprocs = options.count
|
||||
c = client.Command()
|
||||
c = client.Command()
|
||||
cmdstr = " ".join(args[1:])
|
||||
|
||||
currprocs = 0
|
||||
@@ -79,7 +81,10 @@ def run():
|
||||
cmd = ex[node]['value']
|
||||
if not isinstance(cmd, str) and not isinstance(cmd, bytes):
|
||||
cmd = cmd.encode('utf-8')
|
||||
cmdv = ['ssh', node, cmd]
|
||||
if options.port:
|
||||
cmdv = ['ssh', '-p', '{0}'.format(options.port), node, cmd]
|
||||
else:
|
||||
cmdv = ['ssh', node, cmd]
|
||||
if currprocs < concurrentprocs:
|
||||
currprocs += 1
|
||||
run_cmdv(node, cmdv, all, pipedesc)
|
||||
|
||||
@@ -429,6 +429,10 @@ def printattributes(session, requestargs, showtype, nodetype, noderange, options
|
||||
path = '/{0}/{1}/attributes/{2}'.format(nodetype, noderange, showtype)
|
||||
return print_attrib_path(path, session, requestargs, options)
|
||||
|
||||
def _sort_attrib(k):
|
||||
if isinstance(k[1], dict) and k[1].get('sortid', None) is not None:
|
||||
return k[1]['sortid']
|
||||
return k[0]
|
||||
|
||||
def print_attrib_path(path, session, requestargs, options, rename=None, attrprefix=None):
|
||||
exitcode = 0
|
||||
@@ -439,9 +443,7 @@ def print_attrib_path(path, session, requestargs, options, rename=None, attrpref
|
||||
exitcode = 1
|
||||
continue
|
||||
for node in sorted(res['databynode']):
|
||||
for attr, val in sorted(
|
||||
res['databynode'][node].items(),
|
||||
key=lambda k: k[1].get('sortid', k[0]) if isinstance(k[1], dict) else k[0]):
|
||||
for attr, val in sorted(res['databynode'][node].items(), key=_sort_attrib):
|
||||
if attr == 'error':
|
||||
sys.stderr.write('{0}: Error: {1}\n'.format(node, val))
|
||||
continue
|
||||
|
||||
+2
@@ -100,3 +100,5 @@ See nodegroupattrib(8) command on how to manage attributes on a group level.
|
||||
## SEE ALSO
|
||||
|
||||
nodegroupattrib(8), nodeattribexpressions(5)
|
||||
|
||||
## ATTRIBUTES
|
||||
@@ -22,6 +22,11 @@ given as a node expression, as documented in the man page for nodeattribexpressi
|
||||
If combined with `-x`, will show all differing values except those indicated
|
||||
by `-x`
|
||||
|
||||
* `-e`, `--extra`:
|
||||
Read settings that are generally not needed, but may be slow to retrieve.
|
||||
Notably this includes the IMM category of Lenovo systems. The most popular
|
||||
IMM settings are available through faster 'bmc' attributes.
|
||||
|
||||
* `-x`, `--exclude`:
|
||||
Rather than listing only the specified configuration parameters, list all
|
||||
attributes except for the specified ones
|
||||
|
||||
@@ -23,6 +23,9 @@ data may be filtered by various parameters, as denoted in the options below.
|
||||
**nodediscover assign** performs manual discovery, assigning an entry to a node
|
||||
identity or, using `-i`, using a csv file to assign nodes all at once. For
|
||||
example, a spreadsheet of serial numbers to desired node names could be used.
|
||||
Note that if you see that the host is unreachable, it may be due to the IP
|
||||
address on the endpoint having changed since last detected. In such a case, it
|
||||
may help to **clear** and try **assign** again.
|
||||
|
||||
**nodediscover rescan** requests the server to do an active sweep for new
|
||||
devices. Generally every effort is made to passively detect devices as they
|
||||
|
||||
+2
@@ -41,3 +41,5 @@ the attributes are set on the node versus a group to which a node belongs.
|
||||
## SEE ALSO
|
||||
|
||||
nodeattrib(8), nodeattribexpressions(5)
|
||||
|
||||
## ATTRIBUTES
|
||||
@@ -0,0 +1,21 @@
|
||||
nodersync(8) -- Run rsync in parallel against a noderange
|
||||
=========================================================================
|
||||
|
||||
## SYNOPSIS
|
||||
|
||||
`nodersync <file/directorylist> <noderange>:<destination>`
|
||||
|
||||
## DESCRIPTION
|
||||
|
||||
Supervises execution of rsync to push files or a directory tree to the specified
|
||||
noderange. This will present progress as percentage for all nodes.
|
||||
|
||||
## OPTIONS
|
||||
|
||||
* `-m`:
|
||||
Specify maximum number of nodes for noderange max.
|
||||
|
||||
* `-c`:
|
||||
Specify how many rsync executions to do concurrently. If noderange
|
||||
exceeds the count, then excess nodes will wait until one of the
|
||||
active count completes.
|
||||
@@ -1,4 +1,6 @@
|
||||
#!/bin/sh
|
||||
cd `dirname $0`
|
||||
python3 addattribs.py || python2 addattribs.py
|
||||
cd `dirname $0`/doc/man
|
||||
mkdir -p ../../man/man1
|
||||
mkdir -p ../../man/man5
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
import argparse
|
||||
import errno
|
||||
import os
|
||||
import pwd
|
||||
import socket
|
||||
import subprocess
|
||||
import sys
|
||||
@@ -40,6 +41,12 @@ def make_certificate():
|
||||
'/etc/confluent/srvcert.pem -subj /CN='
|
||||
'{0}'.format(socket.gethostname()).split(' ')):
|
||||
raise Exception('Error generating certificate')
|
||||
try:
|
||||
uid = pwd.getpwnam('confluent').pw_uid
|
||||
os.chown('/etc/confluent/privkey.pem', uid, -1)
|
||||
os.chown('/etc/confluent/srvcert.pem', uid, -1)
|
||||
except KeyError:
|
||||
pass
|
||||
print('Certificate generated successfully')
|
||||
os.umask(umask)
|
||||
|
||||
@@ -88,6 +95,9 @@ def show_collective():
|
||||
s = client.Command().connection
|
||||
tlvdata.send(s, {'collective': {'operation': 'show'}})
|
||||
res = tlvdata.recv(s)
|
||||
if 'error' in res:
|
||||
print(res['error'])
|
||||
return
|
||||
if 'error' in res['collective']:
|
||||
print(res['collective']['error'])
|
||||
return
|
||||
|
||||
@@ -32,7 +32,7 @@ import confluent.main
|
||||
import multiprocessing
|
||||
if __name__ == '__main__':
|
||||
multiprocessing.freeze_support()
|
||||
confluent.main.run()
|
||||
confluent.main.run(sys.argv)
|
||||
#except:
|
||||
# pass
|
||||
#p.disable()
|
||||
|
||||
@@ -7,7 +7,7 @@ import tempfile
|
||||
def get_openssl_conf_location():
|
||||
if exists('/etc/pki/tls/openssl.cnf'):
|
||||
return '/etc/pki/tls/openssl.cnf'
|
||||
elif exists('/etc/ssl/openssl.cnf');
|
||||
elif exists('/etc/ssl/openssl.cnf'):
|
||||
return '/etc/ssl/openssl.cnf'
|
||||
else:
|
||||
raise Exception("Cannot find openssl config file")
|
||||
|
||||
@@ -62,8 +62,11 @@ if args[0] == 'restore':
|
||||
if pid is not None:
|
||||
print("Confluent is running, must shut down to restore db")
|
||||
sys.exit(1)
|
||||
password = options.password
|
||||
if options.interactivepassword:
|
||||
password = getpass.getpass('Enter password to restore backup: ')
|
||||
try:
|
||||
cfm.restore_db_from_directory(dumpdir, options.password)
|
||||
cfm.restore_db_from_directory(dumpdir, password)
|
||||
except Exception as e:
|
||||
print(str(e))
|
||||
sys.exit(1)
|
||||
@@ -86,7 +89,7 @@ elif args[0] == 'dump':
|
||||
main._initsecurity(conf.get_config())
|
||||
if not os.path.exists(dumpdir):
|
||||
os.makedirs(dumpdir)
|
||||
cfm.dump_db_to_directory(dumpdir, options.password, options.redact,
|
||||
cfm.dump_db_to_directory(dumpdir, password, options.redact,
|
||||
options.skipkeys)
|
||||
|
||||
|
||||
|
||||
@@ -3,8 +3,12 @@ cd `dirname $0`
|
||||
PKGNAME=$(basename $(pwd))
|
||||
DPKGNAME=$(basename $(pwd) | sed -e s/_/-/)
|
||||
OPKGNAME=$(basename $(pwd) | sed -e s/_/-/)
|
||||
PYEXEC=python3
|
||||
DSCARGS="--with-python3=True --with-python2=False"
|
||||
if grep wheezy /etc/os-release; then
|
||||
DPKGNAME=python-$DPKGNAME
|
||||
PYEXEC=python
|
||||
DSCARGS=""
|
||||
fi
|
||||
cd ..
|
||||
mkdir -p /tmp/confluent # $DPKGNAME
|
||||
@@ -24,15 +28,15 @@ install-scripts=/opt/confluent/bin
|
||||
package=$DPKGNAME
|
||||
EOF
|
||||
|
||||
python setup.py sdist > /dev/null 2>&1
|
||||
py2dsc dist/*.tar.gz
|
||||
$PYEXEC setup.py sdist > /dev/null 2>&1
|
||||
py2dsc $DSCARGS dist/*.tar.gz
|
||||
shopt -s extglob
|
||||
cd deb_dist/!(*.orig)/
|
||||
if [ "$OPKGNAME" = "confluent-server" ]; then
|
||||
if grep wheezy /etc/os-release; then
|
||||
sed -i 's/^\(Depends:.*\)/\1, python-confluent-client, python-lxml, python-eficompressor, python-pycryptodomex, python-dateutil, python-pyopenssl/' debian/control
|
||||
sed -i 's/^\(Depends:.*\)/\1, python-confluent-client, python-lxml, python-eficompressor, python-pycryptodomex, python-dateutil, python-pyopenssl, python-msgpack/' debian/control
|
||||
else
|
||||
sed -i 's/^\(Depends:.*\)/\1, confluent-client, python-lxml, python-eficompressor, python-pycryptodome, python-dateutil, python-websocket/' debian/control
|
||||
sed -i 's/^\(Depends:.*\)/\1, confluent-client, python3-lxml, python3-eficompressor, python3-pycryptodome, python3-websocket, python3-msgpack, python3-eventlet, python3-pyparsing, python3-pyte, python3-pyghmi, python3-paramiko/' debian/control
|
||||
fi
|
||||
if grep wheezy /etc/os-release; then
|
||||
echo 'confluent_client python-confluent-client' >> debian/pydist-overrides
|
||||
@@ -40,6 +44,9 @@ if [ "$OPKGNAME" = "confluent-server" ]; then
|
||||
echo 'confluent_client confluent-client' >> debian/pydist-overrides
|
||||
fi
|
||||
fi
|
||||
if ! grep wheezy /etc/os-release; then
|
||||
sed -i 's/^Package: python3-/Package: /' debian/control
|
||||
fi
|
||||
head -n -1 debian/control > debian/control1
|
||||
mv debian/control1 debian/control
|
||||
echo 'export PYBUILD_INSTALL_ARGS=--install-lib=/opt/confluent/lib/python' >> debian/rules
|
||||
|
||||
@@ -64,12 +64,12 @@ class AsyncTermRelation(object):
|
||||
# Need to keep an association of term object to async
|
||||
# This allows the async handler to know the context of
|
||||
# outgoing data to provide to calling code
|
||||
def __init__(self, termid, async):
|
||||
self.async = async
|
||||
def __init__(self, termid, asynchdl):
|
||||
self.asynchdl = asynchdl
|
||||
self.termid = termid
|
||||
|
||||
def got_data(self, data):
|
||||
self.async.add(self.termid, data)
|
||||
self.asynchdl.add(self.termid, data)
|
||||
|
||||
|
||||
class AsyncSession(object):
|
||||
|
||||
@@ -27,6 +27,8 @@ from fnmatch import fnmatch
|
||||
import hashlib
|
||||
import hmac
|
||||
import multiprocessing
|
||||
import os
|
||||
import pwd
|
||||
import confluent.userutil as userutil
|
||||
import confluent.util as util
|
||||
pam = None
|
||||
@@ -58,9 +60,9 @@ _allowedbyrole = {
|
||||
'/node*/configuration/*',
|
||||
],
|
||||
'update': [
|
||||
'/discovery/*',
|
||||
'/discovery/*',
|
||||
'/networking/macs/rescan',
|
||||
'/node*/power/state',
|
||||
'/node*/power/state',
|
||||
'/node*/power/reseat',
|
||||
'/node*/attributes/*',
|
||||
'/node*/media/*tach',
|
||||
@@ -98,17 +100,6 @@ _deniedbyrole = {
|
||||
}
|
||||
|
||||
|
||||
def _prune_passcache():
|
||||
# This function makes sure we don't remember a passphrase in memory more
|
||||
# than 10 seconds
|
||||
while True:
|
||||
curtime = time.time()
|
||||
for passent in _passcache.iterkeys():
|
||||
if passent[2] < curtime - 90:
|
||||
del _passcache[passent]
|
||||
eventlet.sleep(90)
|
||||
|
||||
|
||||
def _get_usertenant(name, tenant=False):
|
||||
"""_get_usertenant
|
||||
|
||||
@@ -268,12 +259,36 @@ def check_user_passphrase(name, passphrase, operation=None, element=None, tenant
|
||||
_passcache[(user, tenant)] = hashlib.sha256(passphrase).digest()
|
||||
return authorize(user, element, tenant, operation)
|
||||
if pam:
|
||||
pammy = pam.pam()
|
||||
usergood = pammy.authenticate(user, passphrase)
|
||||
del pammy
|
||||
pwe = None
|
||||
try:
|
||||
pwe = pwd.getpwnam(user)
|
||||
except KeyError:
|
||||
#pam won't work if the user doesn't exist, don't go further
|
||||
eventlet.sleep(0.05) # stall even on test for existence of a username
|
||||
return None
|
||||
if os.getuid() != 0:
|
||||
# confluent is running with reduced privilege, however, pam_unix refuses
|
||||
# to let a non-0 user check anothers password.
|
||||
# We will fork and the child will assume elevated privilege to
|
||||
# get unix_chkpwd helper to enable checking /etc/shadow
|
||||
pid = os.fork()
|
||||
if not pid:
|
||||
usergood = False
|
||||
try:
|
||||
# we change to the uid we are trying to authenticate as, because
|
||||
# pam_unix uses unix_chkpwd which reque
|
||||
os.setuid(pwe.pw_uid)
|
||||
usergood = pam.authenticate(user, passphrase, service=_pamservice)
|
||||
finally:
|
||||
os._exit(0 if usergood else 1)
|
||||
usergood = os.waitpid(pid, 0)[1] == 0
|
||||
else:
|
||||
# We are running as root, we don't need to fork in order to authenticate the
|
||||
# user
|
||||
usergood = pam.authenticate(user, passphrase, service=_pamservice)
|
||||
if usergood:
|
||||
_passcache[(user, tenant)] = hashlib.sha256(passphrase).digest()
|
||||
return authorize(user, element, tenant, operation, skipuserobj=False)
|
||||
return authorize(user, element, tenant, operation, skipuserobj=False)
|
||||
eventlet.sleep(0.05) # stall even on test for existence of a username
|
||||
return None
|
||||
|
||||
|
||||
@@ -73,20 +73,12 @@ def connect_to_leader(cert=None, name=None, leader=None):
|
||||
with cfm._initlock:
|
||||
banner = tlvdata.recv(remote) # the banner
|
||||
vers = banner.split()[2]
|
||||
pvers = 0
|
||||
reqver = 4
|
||||
if vers == b'v0':
|
||||
pvers = 2
|
||||
elif vers == b'v1':
|
||||
pvers = 4
|
||||
if sys.version_info[0] < 3:
|
||||
pvers = 2
|
||||
reqver = 2
|
||||
if vers != b'v2':
|
||||
raise Exception('This instance only supports protocol 2, synchronize versions between collective members')
|
||||
tlvdata.recv(remote) # authpassed... 0..
|
||||
if name is None:
|
||||
name = get_myname()
|
||||
tlvdata.send(remote, {'collective': {'operation': 'connect',
|
||||
'protover': reqver,
|
||||
'name': name,
|
||||
'txcount': cfm._txcount}})
|
||||
keydata = tlvdata.recv(remote)
|
||||
@@ -160,15 +152,15 @@ def connect_to_leader(cert=None, name=None, leader=None):
|
||||
raise
|
||||
currentleader = leader
|
||||
#spawn this as a thread...
|
||||
follower = eventlet.spawn(follow_leader, remote, pvers, leader)
|
||||
follower = eventlet.spawn(follow_leader, remote, leader)
|
||||
return True
|
||||
|
||||
|
||||
def follow_leader(remote, proto, leader):
|
||||
def follow_leader(remote, leader):
|
||||
global currentleader
|
||||
cleanexit = False
|
||||
try:
|
||||
cfm.follow_channel(remote, proto)
|
||||
cfm.follow_channel(remote)
|
||||
except greenlet.GreenletExit:
|
||||
cleanexit = True
|
||||
finally:
|
||||
@@ -430,7 +422,6 @@ def handle_connection(connection, cert, request, local=False):
|
||||
tlvdata.send(connection, collinfo)
|
||||
if 'connect' == operation:
|
||||
drone = request['name']
|
||||
folver = request.get('protover', 2)
|
||||
droneinfo = cfm.get_collective_member(drone)
|
||||
if not (droneinfo and util.cert_matches(droneinfo['fingerprint'],
|
||||
cert)):
|
||||
@@ -479,7 +470,7 @@ def handle_connection(connection, cert, request, local=False):
|
||||
connection.sendall(cfgdata)
|
||||
#tlvdata.send(connection, {'tenants': 0}) # skip the tenants for now,
|
||||
# so far unused anyway
|
||||
if not cfm.relay_slaved_requests(drone, connection, folver):
|
||||
if not cfm.relay_slaved_requests(drone, connection):
|
||||
if not retrythread: # start a recovery if everyone else seems
|
||||
# to have disappeared
|
||||
retrythread = eventlet.spawn_after(30 + random.random(),
|
||||
|
||||
@@ -270,7 +270,8 @@ node = {
|
||||
},
|
||||
'console.method': {
|
||||
'description': ('Indicate the method used to access the console of '
|
||||
'the managed node.'),
|
||||
'the managed node. If not specified, then console '
|
||||
'is disabled'),
|
||||
'validvalues': ('ssh', 'ipmi', 'tsmsol'),
|
||||
},
|
||||
# 'virtualization.host': {
|
||||
@@ -301,6 +302,7 @@ node = {
|
||||
'hardwaremanagement.method': {
|
||||
'description': 'The method used to perform operations such as power '
|
||||
'control, get sensor data, get inventory, and so on. '
|
||||
'ipmi is used if not specified.'
|
||||
},
|
||||
'enclosure.bay': {
|
||||
'description': 'The bay in the enclosure, if any',
|
||||
|
||||
@@ -71,6 +71,7 @@ import eventlet.green.select as select
|
||||
import eventlet.green.threading as gthread
|
||||
import fnmatch
|
||||
import json
|
||||
import msgpack
|
||||
import operator
|
||||
import os
|
||||
import random
|
||||
@@ -101,10 +102,6 @@ _cfgstore = None
|
||||
_pendingchangesets = {}
|
||||
_txcount = 0
|
||||
_hasquorum = True
|
||||
if sys.version_info[0] >= 3:
|
||||
lowestver = 4
|
||||
else:
|
||||
lowestver = 2
|
||||
|
||||
_attraliases = {
|
||||
'bmc': 'hardwaremanagement.manager',
|
||||
@@ -115,6 +112,14 @@ _attraliases = {
|
||||
}
|
||||
_validroles = ('Administrator', 'Operator', 'Monitor')
|
||||
|
||||
|
||||
def attrib_supports_expression(attrib):
|
||||
attrib = _attraliases.get(attrib, attrib)
|
||||
if attrib.startswith('secret.') or attrib.startswith('crypted.'):
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def _mkpath(pathname):
|
||||
try:
|
||||
os.makedirs(pathname)
|
||||
@@ -317,8 +322,8 @@ def exec_on_leader(function, *args):
|
||||
while xid in _pendingchangesets:
|
||||
xid = confluent.util.stringify(base64.b64encode(os.urandom(8)))
|
||||
_pendingchangesets[xid] = event.Event()
|
||||
rpcpayload = cPickle.dumps({'function': function, 'args': args,
|
||||
'xid': xid}, protocol=cfgproto)
|
||||
rpcpayload = msgpack.packb({'function': function, 'args': args,
|
||||
'xid': xid}, use_bin_type=False)
|
||||
rpclen = len(rpcpayload)
|
||||
cfgleader.sendall(struct.pack('!Q', rpclen))
|
||||
cfgleader.sendall(rpcpayload)
|
||||
@@ -343,8 +348,8 @@ def exec_on_followers_unconditional(fnname, *args):
|
||||
global _txcount
|
||||
pushes = eventlet.GreenPool()
|
||||
_txcount += 1
|
||||
payload = cPickle.dumps({'function': fnname, 'args': args,
|
||||
'txcount': _txcount}, protocol=lowestver)
|
||||
payload = msgpack.packb({'function': fnname, 'args': args,
|
||||
'txcount': _txcount}, use_bin_type=False)
|
||||
for _ in pushes.starmap(
|
||||
_push_rpc, [(cfgstreams[s], payload) for s in cfgstreams]):
|
||||
pass
|
||||
@@ -491,7 +496,8 @@ def crypt_value(value,
|
||||
key = _masterkey
|
||||
iv = os.urandom(12)
|
||||
crypter = AES.new(key, AES.MODE_GCM, nonce=iv)
|
||||
value = confluent.util.stringify(value).encode('utf-8')
|
||||
if not isinstance(value, bytes):
|
||||
value = value.encode('utf-8')
|
||||
cryptval, hmac = crypter.encrypt_and_digest(value)
|
||||
return iv, cryptval, hmac, b'\x02'
|
||||
|
||||
@@ -565,14 +571,9 @@ def set_global(globalname, value, sync=True):
|
||||
ConfigManager._bg_sync_to_file()
|
||||
|
||||
cfgstreams = {}
|
||||
def relay_slaved_requests(name, listener, vers):
|
||||
def relay_slaved_requests(name, listener):
|
||||
global cfgleader
|
||||
global _hasquorum
|
||||
global lowestver
|
||||
if vers > 2 and sys.version_info[0] < 3:
|
||||
vers = 2
|
||||
if vers < lowestver:
|
||||
lowestver = vers
|
||||
pushes = eventlet.GreenPool()
|
||||
if name not in _followerlocks:
|
||||
_followerlocks[name] = gthread.RLock()
|
||||
@@ -593,7 +594,7 @@ def relay_slaved_requests(name, listener, vers):
|
||||
while _hasquorum != _newquorum:
|
||||
if _newquorum is not None:
|
||||
_hasquorum = _newquorum
|
||||
payload = cPickle.dumps({'quorum': _hasquorum}, protocol=lowestver)
|
||||
payload = msgpack.packb({'quorum': _hasquorum}, use_bin_type=False)
|
||||
for _ in pushes.starmap(
|
||||
_push_rpc,
|
||||
[(cfgstreams[s], payload) for s in cfgstreams]):
|
||||
@@ -615,15 +616,19 @@ def relay_slaved_requests(name, listener, vers):
|
||||
if not nrpc:
|
||||
raise Exception('Truncated client error')
|
||||
rpc += nrpc
|
||||
rpc = cPickle.loads(rpc)
|
||||
rpc = msgpack.unpackb(rpc, raw=False)
|
||||
exc = None
|
||||
if not (rpc['function'].startswith('_rpc_') or rpc['function'].endswith('_collective_member')):
|
||||
raise Exception('Unsupported function {0} called'.format(rpc['function']))
|
||||
try:
|
||||
globals()[rpc['function']](*rpc['args'])
|
||||
except ValueError as ve:
|
||||
exc = ['ValueError', str(ve)]
|
||||
except Exception as e:
|
||||
exc = e
|
||||
exc = ['Exception', str(e)]
|
||||
if 'xid' in rpc:
|
||||
res = _push_rpc(listener, cPickle.dumps({'xid': rpc['xid'],
|
||||
'exc': exc}, protocol=vers))
|
||||
res = _push_rpc(listener, msgpack.packb({'xid': rpc['xid'],
|
||||
'exc': exc}, use_bin_type=False))
|
||||
if not res:
|
||||
break
|
||||
try:
|
||||
@@ -642,7 +647,7 @@ def relay_slaved_requests(name, listener, vers):
|
||||
if cfgstreams:
|
||||
_hasquorum = len(cfgstreams) >= (
|
||||
len(_cfgstore['collective']) // 2)
|
||||
payload = cPickle.dumps({'quorum': _hasquorum}, protocol=lowestver)
|
||||
payload = msgpack.packb({'quorum': _hasquorum}, use_bin_type=False)
|
||||
for _ in pushes.starmap(
|
||||
_push_rpc,
|
||||
[(cfgstreams[s], payload) for s in cfgstreams]):
|
||||
@@ -684,19 +689,15 @@ class StreamHandler(object):
|
||||
self.sock = None
|
||||
|
||||
|
||||
def stop_following(replacement=None, proto=2):
|
||||
def stop_following(replacement=None):
|
||||
with _leaderlock:
|
||||
global cfgleader
|
||||
global cfgproto
|
||||
if cfgleader and not isinstance(cfgleader, bool):
|
||||
try:
|
||||
cfgleader.close()
|
||||
except Exception:
|
||||
pass
|
||||
cfgleader = replacement
|
||||
if proto > 2 and sys.version_info[0] < 3:
|
||||
proto = 2
|
||||
cfgproto = proto
|
||||
|
||||
def stop_leading():
|
||||
for stream in list(cfgstreams):
|
||||
@@ -754,15 +755,14 @@ def commit_clear():
|
||||
ConfigManager._bg_sync_to_file()
|
||||
|
||||
cfgleader = None
|
||||
cfgproto = 2
|
||||
|
||||
|
||||
def follow_channel(channel, proto=2):
|
||||
def follow_channel(channel):
|
||||
global _txcount
|
||||
global _hasquorum
|
||||
try:
|
||||
stop_leading()
|
||||
stop_following(channel, proto)
|
||||
stop_following(channel)
|
||||
lh = StreamHandler(channel)
|
||||
msg = lh.get_next_msg()
|
||||
while msg:
|
||||
@@ -774,17 +774,24 @@ def follow_channel(channel, proto=2):
|
||||
if not nrpc:
|
||||
raise Exception('Truncated message error')
|
||||
rpc += nrpc
|
||||
rpc = cPickle.loads(rpc)
|
||||
rpc = msgpack.unpackb(rpc, raw=False)
|
||||
if 'txcount' in rpc:
|
||||
_txcount = rpc['txcount']
|
||||
if 'function' in rpc:
|
||||
if not (rpc['function'].startswith('_true') or rpc['function'].startswith('_rpc')):
|
||||
raise Exception("Received unsupported function call: {0}".format(rpc['function']))
|
||||
try:
|
||||
globals()[rpc['function']](*rpc['args'])
|
||||
except Exception as e:
|
||||
print(repr(e))
|
||||
if 'xid' in rpc and rpc['xid']:
|
||||
if rpc.get('exc', None):
|
||||
_pendingchangesets[rpc['xid']].send_exception(rpc['exc'])
|
||||
exctype, excstr = rpc['exc']
|
||||
if exctype == 'ValueError':
|
||||
exc = ValueError(excstr)
|
||||
else:
|
||||
exc = Exception(excstr)
|
||||
_pendingchangesets[rpc['xid']].send_exception(exc)
|
||||
else:
|
||||
_pendingchangesets[rpc['xid']].send()
|
||||
if 'quorum' in rpc:
|
||||
@@ -1407,7 +1414,7 @@ class ConfigManager(object):
|
||||
return exec_on_leader('_rpc_master_set_user', self.tenant, name,
|
||||
attributemap)
|
||||
if cfgstreams:
|
||||
exec_on_followers('_rpc_set_user', self.tenant, name)
|
||||
exec_on_followers('_rpc_set_user', self.tenant, name, attributemap)
|
||||
self._true_set_user(name, attributemap)
|
||||
|
||||
def _true_set_user(self, name, attributemap):
|
||||
@@ -1487,9 +1494,10 @@ class ConfigManager(object):
|
||||
self._cfgstore['users'][name]['displayname'] = displayname
|
||||
_cfgstore['main']['idmap'][uid] = {
|
||||
'tenant': self.tenant,
|
||||
'username': name
|
||||
'username': name,
|
||||
'role': role,
|
||||
}
|
||||
if attributemap is not None:
|
||||
if attributemap:
|
||||
self._true_set_user(name, attributemap)
|
||||
_mark_dirtykey('users', name, self.tenant)
|
||||
_mark_dirtykey('idmap', uid)
|
||||
@@ -1894,6 +1902,8 @@ class ConfigManager(object):
|
||||
eventlet.spawn_n(_do_notifier, self, watcher, callback)
|
||||
|
||||
def del_nodes(self, nodes):
|
||||
if isinstance(nodes, set):
|
||||
nodes = list(nodes) # msgpack can't handle set
|
||||
if cfgleader: # slaved to a collective
|
||||
return exec_on_leader('_rpc_master_del_nodes', self.tenant,
|
||||
nodes)
|
||||
|
||||
@@ -108,10 +108,10 @@ def pytechars2line(chars, maxlen=None):
|
||||
char = chars[charidx]
|
||||
csi = bytearray([])
|
||||
if char.fg != lfg:
|
||||
csi.append(30 + pytecolors2ansi[char.fg])
|
||||
csi.append(30 + pytecolors2ansi.get(char.fg, 9))
|
||||
lfg = char.fg
|
||||
if char.bg != lbg:
|
||||
csi.append(40 + pytecolors2ansi[char.bg])
|
||||
csi.append(40 + pytecolors2ansi.get(char.bg, 9))
|
||||
lbg = char.bg
|
||||
if char.bold != lb:
|
||||
lb = char.bold
|
||||
@@ -243,7 +243,7 @@ class ConsoleHandler(object):
|
||||
def check_collective(self, attrvalue):
|
||||
myc = attrvalue.get(self.node, {}).get('collective.manager', {}).get(
|
||||
'value', None)
|
||||
if configmodule.list_collective() and not myc:
|
||||
if list(configmodule.list_collective()) and not myc:
|
||||
self._is_local = False
|
||||
self._detach()
|
||||
self._disconnect()
|
||||
|
||||
@@ -60,20 +60,19 @@ import eventlet.greenpool as greenpool
|
||||
import eventlet.green.ssl as ssl
|
||||
import eventlet.queue as queue
|
||||
import itertools
|
||||
import msgpack
|
||||
import os
|
||||
try:
|
||||
import cPickle as pickle
|
||||
pargs = {}
|
||||
except ImportError:
|
||||
import pickle
|
||||
pargs = {'encoding': 'utf-8'}
|
||||
import socket
|
||||
import struct
|
||||
import sys
|
||||
|
||||
pluginmap = {}
|
||||
dispatch_plugins = (b'ipmi', u'ipmi')
|
||||
dispatch_plugins = (b'ipmi', u'ipmi', b'redfish', u'redfish', b'tsmsol', u'tsmsol')
|
||||
|
||||
try:
|
||||
unicode
|
||||
except NameError:
|
||||
unicode = str
|
||||
|
||||
def seek_element(currplace, currkey):
|
||||
try:
|
||||
@@ -221,10 +220,14 @@ def _init_core():
|
||||
'pluginattrs': ['hardwaremanagement.method'],
|
||||
'default': 'ipmi',
|
||||
}),
|
||||
'extra': PluginRoute({
|
||||
'pluginattrs': ['hardwaremanagement.method'],
|
||||
'default': 'ipmi',
|
||||
}),
|
||||
'advanced': PluginRoute({
|
||||
'pluginattrs': ['hardwaremanagement.method'],
|
||||
'default': 'ipmi',
|
||||
}),
|
||||
}),
|
||||
},
|
||||
},
|
||||
'storage': {
|
||||
@@ -417,18 +420,22 @@ def create_user(inputdata, configmanager):
|
||||
try:
|
||||
username = inputdata['name']
|
||||
del inputdata['name']
|
||||
role = inputdata['role']
|
||||
del inputdata['role']
|
||||
except (KeyError, ValueError):
|
||||
raise exc.InvalidArgumentException()
|
||||
configmanager.create_user(username, attributemap=inputdata)
|
||||
raise exc.InvalidArgumentException('Missing user name or role')
|
||||
configmanager.create_user(username, role, attributemap=inputdata)
|
||||
|
||||
|
||||
def create_usergroup(inputdata, configmanager):
|
||||
try:
|
||||
groupname = inputdata['name']
|
||||
role = inputdata['role']
|
||||
del inputdata['name']
|
||||
del inputdata['role']
|
||||
except (KeyError, ValueError):
|
||||
raise exc.InvalidArgumentException()
|
||||
configmanager.create_usergroup(groupname)
|
||||
raise exc.InvalidArgumentException("Missing user name or role")
|
||||
configmanager.create_usergroup(groupname, role)
|
||||
|
||||
|
||||
def update_usergroup(groupname, attribmap, configmanager):
|
||||
@@ -689,16 +696,22 @@ def handle_dispatch(connection, cert, dispatch, peername):
|
||||
cfm.get_collective_member(peername)['fingerprint'], cert):
|
||||
connection.close()
|
||||
return
|
||||
pversion = 0
|
||||
if bytearray(dispatch)[0] == 0x80:
|
||||
pversion = bytearray(dispatch)[1]
|
||||
dispatch = pickle.loads(dispatch, **pargs)
|
||||
if dispatch[0:2] != b'\x01\x03': # magic value to indicate msgpack
|
||||
# We only support msgpack now
|
||||
# The magic should preclude any pickle, as the first byte can never be
|
||||
# under 0x20 or so.
|
||||
connection.close()
|
||||
return
|
||||
dispatch = msgpack.unpackb(dispatch[2:], raw=False)
|
||||
configmanager = cfm.ConfigManager(dispatch['tenant'])
|
||||
nodes = dispatch['nodes']
|
||||
inputdata = dispatch['inputdata']
|
||||
operation = dispatch['operation']
|
||||
pathcomponents = dispatch['path']
|
||||
routespec = nested_lookup(noderesources, pathcomponents)
|
||||
inputdata = msg.get_input_message(
|
||||
pathcomponents, operation, inputdata, nodes, dispatch['isnoderange'],
|
||||
configmanager)
|
||||
plugroute = routespec.routeinfo
|
||||
plugpath = None
|
||||
nodesbyhandler = {}
|
||||
@@ -728,18 +741,26 @@ def handle_dispatch(connection, cert, dispatch, peername):
|
||||
configmanager=configmanager,
|
||||
inputdata=inputdata))
|
||||
for res in itertools.chain(*passvalues):
|
||||
_forward_rsp(connection, res, pversion)
|
||||
_forward_rsp(connection, res)
|
||||
except Exception as res:
|
||||
_forward_rsp(connection, res, pversion)
|
||||
_forward_rsp(connection, res)
|
||||
connection.sendall('\x00\x00\x00\x00\x00\x00\x00\x00')
|
||||
|
||||
|
||||
def _forward_rsp(connection, res, pversion):
|
||||
def _forward_rsp(connection, res):
|
||||
try:
|
||||
r = pickle.dumps(res, protocol=pversion)
|
||||
except TypeError:
|
||||
r = pickle.dumps(Exception(
|
||||
'Cannot serialize error, check collective.manager error logs for details' + str(res)), protocol=pversion)
|
||||
r = res.serialize()
|
||||
except AttributeError:
|
||||
if isinstance(res, Exception):
|
||||
r = msgpack.packb(['Exception', str(res)], use_bin_type=False)
|
||||
else:
|
||||
r = msgpack.packb(
|
||||
['Exception', 'Unable to serialize response ' + repr(res)],
|
||||
use_bin_type=False)
|
||||
except Exception as e:
|
||||
r = msgpack.packb(
|
||||
['Exception', 'Unable to serialize response ' + repr(res) + ' due to ' + str(e)],
|
||||
use_bin_type=False)
|
||||
rlen = len(r)
|
||||
if not rlen:
|
||||
return
|
||||
@@ -830,7 +851,7 @@ def handle_node_request(configmanager, inputdata, operation,
|
||||
del pathcomponents[0:2]
|
||||
passvalues = queue.Queue()
|
||||
plugroute = routespec.routeinfo
|
||||
inputdata = msg.get_input_message(
|
||||
msginputdata = msg.get_input_message(
|
||||
pathcomponents, operation, inputdata, nodes, isnoderange,
|
||||
configmanager)
|
||||
if 'handler' in plugroute: # fixed handler definition, easy enough
|
||||
@@ -841,7 +862,7 @@ def handle_node_request(configmanager, inputdata, operation,
|
||||
passvalue = hfunc(
|
||||
nodes=nodes, element=pathcomponents,
|
||||
configmanager=configmanager,
|
||||
inputdata=inputdata)
|
||||
inputdata=msginputdata)
|
||||
if isnoderange:
|
||||
return passvalue
|
||||
elif isinstance(passvalue, console.Console):
|
||||
@@ -894,13 +915,13 @@ def handle_node_request(configmanager, inputdata, operation,
|
||||
workers.spawn(addtoqueue, passvalues, hfunc, {'nodes': nodesbyhandler[hfunc],
|
||||
'element': pathcomponents,
|
||||
'configmanager': configmanager,
|
||||
'inputdata': inputdata})
|
||||
'inputdata': msginputdata})
|
||||
for manager in nodesbymanager:
|
||||
numworkers += 1
|
||||
workers.spawn(addtoqueue, passvalues, dispatch_request, {
|
||||
'nodes': nodesbymanager[manager], 'manager': manager,
|
||||
'element': pathcomponents, 'configmanager': configmanager,
|
||||
'inputdata': inputdata, 'operation': operation})
|
||||
'inputdata': inputdata, 'operation': operation, 'isnoderange': isnoderange})
|
||||
if isnoderange or not autostrip:
|
||||
return iterate_queue(numworkers, passvalues)
|
||||
else:
|
||||
@@ -944,7 +965,7 @@ def addtoqueue(theq, fun, kwargs):
|
||||
|
||||
|
||||
def dispatch_request(nodes, manager, element, configmanager, inputdata,
|
||||
operation):
|
||||
operation, isnoderange):
|
||||
a = configmanager.get_collective_member(manager)
|
||||
try:
|
||||
remote = socket.create_connection((a['address'], 13001))
|
||||
@@ -978,10 +999,10 @@ def dispatch_request(nodes, manager, element, configmanager, inputdata,
|
||||
pvers = 2
|
||||
tlvdata.recv(remote)
|
||||
myname = collective.get_myname()
|
||||
dreq = pickle.dumps({'name': myname, 'nodes': list(nodes),
|
||||
'path': element,'tenant': configmanager.tenant,
|
||||
'operation': operation, 'inputdata': inputdata},
|
||||
protocol=pvers)
|
||||
dreq = b'\x01\x03' + msgpack.packb(
|
||||
{'name': myname, 'nodes': list(nodes),
|
||||
'path': element,'tenant': configmanager.tenant,
|
||||
'operation': operation, 'inputdata': inputdata, 'isnoderange': isnoderange}, use_bin_type=False)
|
||||
tlvdata.send(remote, {'dispatch': {'name': myname, 'length': len(dreq)}})
|
||||
remote.sendall(dreq)
|
||||
while True:
|
||||
@@ -1029,11 +1050,13 @@ def dispatch_request(nodes, manager, element, configmanager, inputdata,
|
||||
return
|
||||
rsp += nrsp
|
||||
try:
|
||||
rsp = pickle.loads(rsp, **pargs)
|
||||
except UnicodeDecodeError:
|
||||
rsp = pickle.loads(rsp, encoding='latin1')
|
||||
rsp = msg.msg_deserialize(rsp)
|
||||
except Exception:
|
||||
rsp = exc.deserialize_exc(rsp)
|
||||
if isinstance(rsp, Exception):
|
||||
raise rsp
|
||||
if not rsp:
|
||||
raise Exception('Error in cross-collective serialize/deserialze, see remote logs')
|
||||
yield rsp
|
||||
|
||||
|
||||
|
||||
@@ -65,7 +65,7 @@ import base64
|
||||
import confluent.config.configmanager as cfm
|
||||
import confluent.collective.manager as collective
|
||||
import confluent.discovery.protocols.pxe as pxe
|
||||
#import confluent.discovery.protocols.ssdp as ssdp
|
||||
import confluent.discovery.protocols.ssdp as ssdp
|
||||
import confluent.discovery.protocols.slp as slp
|
||||
import confluent.discovery.handlers.imm as imm
|
||||
import confluent.discovery.handlers.cpstorage as cpstorage
|
||||
@@ -107,6 +107,8 @@ nodehandlers = {
|
||||
'service:management-hardware.Lenovo:lenovo-xclarity-controller': xcc,
|
||||
'service:management-hardware.IBM:integrated-management-module2': imm,
|
||||
'pxe-client': pxeh,
|
||||
'onie-switch': None,
|
||||
'cumulus-switch': None,
|
||||
'service:io-device.Lenovo:management-module': None,
|
||||
'service:thinkagile-storage': cpstorage,
|
||||
'service:lenovo-tsm': tsm,
|
||||
@@ -114,6 +116,8 @@ nodehandlers = {
|
||||
|
||||
servicenames = {
|
||||
'pxe-client': 'pxe-client',
|
||||
'onie-switch': 'onie-switch',
|
||||
'cumulus-switch': 'cumulus-switch',
|
||||
'service:lenovo-smm': 'lenovo-smm',
|
||||
'service:management-hardware.Lenovo:lenovo-xclarity-controller': 'lenovo-xcc',
|
||||
'service:management-hardware.IBM:integrated-management-module2': 'lenovo-imm2',
|
||||
@@ -124,6 +128,8 @@ servicenames = {
|
||||
|
||||
servicebyname = {
|
||||
'pxe-client': 'pxe-client',
|
||||
'onie-switch': 'onie-switch',
|
||||
'cumulus-switch': 'cumulus-switch',
|
||||
'lenovo-smm': 'service:lenovo-smm',
|
||||
'lenovo-xcc': 'service:management-hardware.Lenovo:lenovo-xclarity-controller',
|
||||
'lenovo-imm2': 'service:management-hardware.IBM:integrated-management-module2',
|
||||
@@ -1070,7 +1076,7 @@ def discover_node(cfg, handler, info, nodename, manual):
|
||||
traceback.print_exc()
|
||||
return False
|
||||
newnodeattribs = {}
|
||||
if cfm.list_collective():
|
||||
if list(cfm.list_collective()):
|
||||
# We are in a collective, check collective.manager
|
||||
cmc = cfg.get_node_attributes(nodename, 'collective.manager')
|
||||
cm = cmc.get(nodename, {}).get('collective.manager', {}).get('value', None)
|
||||
@@ -1211,9 +1217,17 @@ def rescan():
|
||||
if scanner:
|
||||
return
|
||||
else:
|
||||
scanner = eventlet.spawn(slp.active_scan, safe_detected, slp)
|
||||
scanner = eventlet.spawn(blocking_scan)
|
||||
|
||||
|
||||
def blocking_scan():
|
||||
global scanner
|
||||
slpscan = eventlet.spawn(slp.active_scan, safe_detected, slp)
|
||||
ssdpscan = eventlet.spawn(ssdp.active_scan, safe_detected, ssdp)
|
||||
slpscan.wait()
|
||||
ssdpscan.wait()
|
||||
scanner = None
|
||||
|
||||
def start_detection():
|
||||
global attribwatcher
|
||||
global rechecker
|
||||
|
||||
@@ -13,9 +13,58 @@
|
||||
# limitations under the License.
|
||||
|
||||
import confluent.discovery.handlers.bmc as bmchandler
|
||||
import eventlet
|
||||
import confluent.util as util
|
||||
try:
|
||||
from urllib import urlencode
|
||||
except ImportError:
|
||||
from urllib.parse import urlencode
|
||||
webclient = eventlet.import_patched('pyghmi.util.webclient')
|
||||
|
||||
class NodeHandler(bmchandler.NodeHandler):
|
||||
DEFAULT_USER = 'admin'
|
||||
DEFAULT_PASS = 'admin'
|
||||
devname = 'BMC'
|
||||
maxmacs = 2
|
||||
|
||||
def validate_cert(self, certificate):
|
||||
# broadly speaking, merely checks consistency moment to moment,
|
||||
# but if https_cert gets stricter, this check means something
|
||||
fprint = util.get_fingerprint(self.https_cert)
|
||||
return util.cert_matches(fprint, certificate)
|
||||
|
||||
def get_webclient(self, user, passwd, newuser, newpass):
|
||||
wc = webclient.SecureHTTPConnection(self.ipaddr, 443,
|
||||
verifycallback=self.validate_cert)
|
||||
wc.connect()
|
||||
authdata = urlencode({'username': user, 'password': passwd,
|
||||
'weblogsign': 1})
|
||||
res = wc.grab_json_response_with_status('/api/session', authdata)
|
||||
if res[1] == 200:
|
||||
if res[0].get('force_password', 1) == 0:
|
||||
# Need to handle password change
|
||||
passchange = {
|
||||
'Password': newpass,
|
||||
'RetypePassword': newpass,
|
||||
'param': 4,
|
||||
'username': 'admin',
|
||||
'privilege': 4,
|
||||
}
|
||||
passchange = urlencode(passchange)
|
||||
rsp = wc.grab_json_response_with_status('/api/reset-pass',
|
||||
passchange)
|
||||
rsp = wc.grab_json_response_with_status('/api/session',
|
||||
method='DELETE')
|
||||
|
||||
def config(self, nodename, reset=False):
|
||||
self.nodename = nodename
|
||||
creds = self.configmanager.get_node_attributes(
|
||||
self.nodename, ['secret.hardwaremanagementuser',
|
||||
'secret.hardwaremanagementpassword'],
|
||||
decrypt=True)
|
||||
user, passwd, isdefault = self.get_node_credentials(
|
||||
nodename, creds, 'admin', 'admin')
|
||||
if not isdefault:
|
||||
self.get_webclient(self.DEFAULT_USER, self.DEFAULT_PASS, user,
|
||||
passwd)
|
||||
self._bmcconfig(nodename, False)
|
||||
|
||||
@@ -90,6 +90,13 @@ class NodeHandler(bmchandler.NodeHandler):
|
||||
smmip = smmip[-1][0]
|
||||
if smmip and ':' in smmip:
|
||||
raise exc.NotImplementedException('IPv6 not supported')
|
||||
wc.request('POST', '/data', 'get=hostname')
|
||||
rsp = wc.getresponse()
|
||||
rspdata = fromstring(util.stringify(rsp.read()))
|
||||
currip = rspdata.find('netConfig').find('ifConfigEntries').find(
|
||||
'ifConfig').find('v4IPAddr').text
|
||||
if currip == smmip:
|
||||
return
|
||||
netconfig = netutil.get_nic_config(cfg, nodename, ip=smmip)
|
||||
netmask = netutil.cidr_to_mask(netconfig['prefix'])
|
||||
setdata = 'set=ifIndex:0,v4DHCPEnabled:0,v4IPAddr:{0},v4NetMask:{1}'.format(smmip, netmask)
|
||||
|
||||
@@ -15,6 +15,7 @@
|
||||
import base64
|
||||
import codecs
|
||||
import confluent.discovery.handlers.imm as immhandler
|
||||
import confluent.exceptions as exc
|
||||
import confluent.netutil as netutil
|
||||
import confluent.util as util
|
||||
import errno
|
||||
@@ -95,7 +96,8 @@ class NodeHandler(immhandler.NodeHandler):
|
||||
ipmicmd.xraw_command(netfn=0x3a, command=0xf1, data=(1,))
|
||||
except pygexc.IpmiException as e:
|
||||
if (e.ipmicode != 193 and 'Unauthorized name' not in str(e) and
|
||||
'Incorrect password' not in str(e)):
|
||||
'Incorrect password' not in str(e) and
|
||||
str(e) != 'Session no longer connected'):
|
||||
# raise an issue if anything other than to be expected
|
||||
if disableipmi:
|
||||
_, _ = wc.grab_json_response_with_status(
|
||||
@@ -173,8 +175,10 @@ class NodeHandler(immhandler.NodeHandler):
|
||||
if pwdchanged:
|
||||
# Remove the minimum change interval, to allow sane
|
||||
# password changes after provisional changes
|
||||
self.set_password_policy('')
|
||||
wc = self.wc
|
||||
self.set_password_policy('', wc)
|
||||
return (wc, pwdchanged)
|
||||
return (None, None)
|
||||
|
||||
@property
|
||||
def wc(self):
|
||||
@@ -230,7 +234,7 @@ class NodeHandler(immhandler.NodeHandler):
|
||||
if wc:
|
||||
return wc
|
||||
|
||||
def set_password_policy(self, strruleset):
|
||||
def set_password_policy(self, strruleset, wc):
|
||||
ruleset = {'USER_GlobalMinPassChgInt': '0'}
|
||||
for rule in strruleset.split(','):
|
||||
if '=' not in rule:
|
||||
@@ -251,9 +255,7 @@ class NodeHandler(immhandler.NodeHandler):
|
||||
if name.lower() == 'reuse':
|
||||
ruleset['USER_GlobalMinPassReuseCycle'] = value
|
||||
try:
|
||||
wc = self.wc
|
||||
wc.grab_json_response('/api/dataset', ruleset)
|
||||
wc.grab_json_response('/api/providers/logout')
|
||||
except Exception as e:
|
||||
print(repr(e))
|
||||
pass
|
||||
@@ -361,7 +363,7 @@ class NodeHandler(immhandler.NodeHandler):
|
||||
self.nodename, ['secret.hardwaremanagementuser',
|
||||
'secret.hardwaremanagementpassword'], decrypt=True)
|
||||
user, passwd, isdefault = self.get_node_credentials(nodename, creds, 'USERID', 'PASSW0RD')
|
||||
self.set_password_policy(strruleset)
|
||||
self.set_password_policy(strruleset, wc)
|
||||
if self._atdefaultcreds:
|
||||
if isdefault and self.tmppasswd:
|
||||
raise Exception(
|
||||
|
||||
@@ -30,9 +30,16 @@ pxearchs = {
|
||||
'\x00\x07': 'uefi-x64',
|
||||
'\x00\x09': 'uefi-x64',
|
||||
'\x00\x0b': 'uefi-aarch64',
|
||||
'\x00\x10': 'uefi-httpboot',
|
||||
}
|
||||
|
||||
|
||||
def stringify(value):
|
||||
string = bytes(value)
|
||||
if not isinstance(string, str):
|
||||
string = string.decode('utf8')
|
||||
return string
|
||||
|
||||
def decode_uuid(rawguid):
|
||||
lebytes = struct.unpack_from('<IHH', rawguid[:8])
|
||||
bebytes = struct.unpack_from('>HHI', rawguid[8:])
|
||||
@@ -40,23 +47,50 @@ def decode_uuid(rawguid):
|
||||
lebytes[0], lebytes[1], lebytes[2], bebytes[0], bebytes[1], bebytes[2]).lower()
|
||||
|
||||
|
||||
def _decode_ocp_vivso(rq, idx, size):
|
||||
end = idx + size
|
||||
vivso = {'service-type': 'onie-switch'}
|
||||
while idx < end:
|
||||
if rq[idx] == 3:
|
||||
vivso['machine'] = stringify(rq[idx + 2:idx + 2 + rq[idx + 1]])
|
||||
elif rq[idx] == 4:
|
||||
vivso['arch'] = stringify(rq[idx + 2:idx + 2 + rq[idx + 1]])
|
||||
elif rq[idx] == 5:
|
||||
vivso['revision'] = stringify(rq[idx + 2:idx + 2 + rq[idx + 1]])
|
||||
idx += rq[idx + 1] + 2
|
||||
return '', None, vivso
|
||||
|
||||
|
||||
def find_info_in_options(rq, optidx):
|
||||
uuid = None
|
||||
arch = None
|
||||
vivso = None
|
||||
ztpurlrequested = False
|
||||
iscumulus = False
|
||||
try:
|
||||
while uuid is None or arch is None:
|
||||
if rq[optidx] == 53: # DHCP message type
|
||||
# we want only length 1 and only discover (type 1)
|
||||
if rq[optidx + 1] != 1 or rq[optidx + 2] != 1:
|
||||
return uuid, arch
|
||||
return uuid, arch, vivso
|
||||
optidx += 3
|
||||
elif rq[optidx] == 55:
|
||||
if 239 in rq[optidx + 2:optidx + 2 + rq[optidx + 1]]:
|
||||
ztpurlrequested = True
|
||||
optidx += rq[optidx + 1] + 2
|
||||
elif rq[optidx] == 60:
|
||||
vci = stringify(rq[optidx + 2:optidx + 2 + rq[optidx + 1]])
|
||||
if vci.startswith('cumulus-linux'):
|
||||
iscumulus = True
|
||||
arch = vci.replace('cumulus-linux', '').strip()
|
||||
optidx += rq[optidx + 1] + 2
|
||||
elif rq[optidx] == 97:
|
||||
if rq[optidx + 1] != 17:
|
||||
# 16 bytes of uuid and one reserved byte
|
||||
return uuid, arch
|
||||
return uuid, arch, vivso
|
||||
if rq[optidx + 2] != 0: # the reserved byte should be zero,
|
||||
# anything else would be a new spec that we don't know yet
|
||||
return uuid, arch
|
||||
return uuid, arch, vivso
|
||||
uuid = decode_uuid(rq[optidx + 3:optidx + 19])
|
||||
optidx += 19
|
||||
elif rq[optidx] == 93:
|
||||
@@ -66,11 +100,20 @@ def find_info_in_options(rq, optidx):
|
||||
if archraw in pxearchs:
|
||||
arch = pxearchs[archraw]
|
||||
optidx += 4
|
||||
elif rq[optidx] == 125:
|
||||
#vivso = rq[optidx + 2:optidx + 2 + rq[optidx + 1]]
|
||||
if rq[optidx + 2:optidx + 6] == b'\x00\x00\xa6\x7f': # OCP
|
||||
return _decode_ocp_vivso(rq, optidx + 7, rq[optidx + 6])
|
||||
optidx += rq[optidx + 1] + 2
|
||||
else:
|
||||
optidx += rq[optidx + 1] + 2
|
||||
except IndexError:
|
||||
return uuid, arch
|
||||
return uuid, arch
|
||||
pass
|
||||
if not vivso and iscumulus and ztpurlrequested:
|
||||
if not uuid:
|
||||
uuid = ''
|
||||
vivso = {'service-type': 'cumulus-switch', 'arch': arch}
|
||||
return uuid, arch, vivso
|
||||
|
||||
def snoop(handler, protocol=None):
|
||||
#TODO(jjohnson2): ipv6 socket and multicast for DHCPv6, should that be
|
||||
@@ -101,7 +144,14 @@ def snoop(handler, protocol=None):
|
||||
optidx = rq.index(b'\x63\x82\x53\x63') + 4
|
||||
except ValueError:
|
||||
continue
|
||||
uuid, arch = find_info_in_options(rq, optidx)
|
||||
uuid, arch, vivso = find_info_in_options(rq, optidx)
|
||||
if vivso:
|
||||
# info['modelnumber'] = info['attributes']['enclosure-machinetype-model'][0]
|
||||
handler({'hwaddr': netaddr, 'uuid': uuid,
|
||||
'architecture': vivso.get('arch', ''),
|
||||
'services': (vivso['service-type'],),
|
||||
'attributes': {'enclosure-machinetype-model': [vivso.get('machine', '')]}})
|
||||
continue
|
||||
if uuid is None:
|
||||
continue
|
||||
# We will fill out service to have something to byte into,
|
||||
|
||||
@@ -493,11 +493,13 @@ def snoop(handler, protocol=None):
|
||||
_add_attributes(peerbymacaddress[mac])
|
||||
peerbymacaddress[mac]['hwaddr'] = mac
|
||||
peerbymacaddress[mac]['protocol'] = protocol
|
||||
if 'service:ipmi' in peerbymacaddress[mac]['services']:
|
||||
if 'service:ipmi//Athena:623' in peerbymacaddress[mac].get('urls', ()):
|
||||
peerbymacaddress[mac]['services'] = ['service:thinkagile-storage']
|
||||
else:
|
||||
for srvurl in peerbymacaddress[mac].get('urls', ()):
|
||||
if len(srvurl) > 4:
|
||||
srvurl = srvurl[:-3]
|
||||
if srvurl.endswith('://Athena:'):
|
||||
continue
|
||||
if 'service:ipmi' in peerbymacaddress[mac]['services']:
|
||||
continue
|
||||
if 'service:lightttpd' in peerbymacaddress[mac]['services']:
|
||||
currinf = peerbymacaddress[mac]
|
||||
curratt = currinf.get('attributes', {})
|
||||
@@ -568,11 +570,13 @@ def scan(srvtypes=_slp_services, addresses=None, localonly=False):
|
||||
_grab_rsps((net, net4), rsps, 1, xidmap)
|
||||
# now to analyze and flesh out the responses
|
||||
for id in rsps:
|
||||
if 'service:ipmi' in rsps[id]['services']:
|
||||
if 'service:ipmi://Athena:623' in rsps[id].get('urls', ''):
|
||||
rsps[id]['services'] = ['service:thinkagile-storage']
|
||||
else:
|
||||
for srvurl in rsps[id].get('urls', ()):
|
||||
if len(srvurl) > 4:
|
||||
srvurl = srvurl[:-3]
|
||||
if srvurl.endswith('://Athena:'):
|
||||
continue
|
||||
if 'service:ipmi' in rsps[id]['services']:
|
||||
continue
|
||||
if localonly:
|
||||
for addr in rsps[id]['addresses']:
|
||||
if 'fe80' in addr[0]:
|
||||
|
||||
@@ -32,6 +32,10 @@ import confluent.neighutil as neighutil
|
||||
import confluent.util as util
|
||||
import eventlet.green.select as select
|
||||
import eventlet.green.socket as socket
|
||||
try:
|
||||
from eventlet.green.urllib.request import urlopen
|
||||
except (ImportError, AssertionError):
|
||||
from eventlet.green.urllib2 import urlopen
|
||||
import struct
|
||||
|
||||
mcastv4addr = '239.255.255.250'
|
||||
@@ -45,6 +49,20 @@ smsg = ('M-SEARCH * HTTP/1.1\r\n'
|
||||
'MX: 3\r\n\r\n')
|
||||
|
||||
|
||||
def active_scan(handler, protocol=None):
|
||||
known_peers = set([])
|
||||
for scanned in scan(['urn:dmtf-org:service:redfish-rest:1']):
|
||||
for addr in scanned['addresses']:
|
||||
ip = addr[0].partition('%')[0] # discard scope if present
|
||||
if ip not in neighutil.neightable:
|
||||
continue
|
||||
if addr in known_peers:
|
||||
break
|
||||
known_peers.add(addr)
|
||||
else:
|
||||
scanned['protocol'] = protocol
|
||||
handler(scanned)
|
||||
|
||||
def scan(services, target=None):
|
||||
for service in services:
|
||||
for rply in _find_service(service, target):
|
||||
@@ -109,11 +127,11 @@ def snoop(handler, byehandler=None):
|
||||
known_peers.add(peer)
|
||||
newmacs.add(mac)
|
||||
if mac in peerbymacaddress:
|
||||
peerbymacaddress[mac]['peers'].append(peer)
|
||||
peerbymacaddress[mac]['addresses'].append(peer)
|
||||
else:
|
||||
peerbymacaddress[mac] = {
|
||||
'hwaddr': mac,
|
||||
'peers': [peer],
|
||||
'addresses': [peer],
|
||||
}
|
||||
peerdata = peerbymacaddress[mac]
|
||||
for headline in rsp[1:]:
|
||||
@@ -185,7 +203,12 @@ def _find_service(service, target):
|
||||
timeout = 0
|
||||
r, _, _ = select.select((net4, net6), (), (), timeout)
|
||||
for nid in peerdata:
|
||||
yield peerdata[nid]
|
||||
for url in peerdata[nid].get('urls', ()):
|
||||
if url.endswith('/desc.tmpl'):
|
||||
info = urlopen(url).read()
|
||||
if '<friendlyName>Athena</friendlyName>' in info:
|
||||
peerdata[nid]['services'] = ['service:thinkagile-storage']
|
||||
yield peerdata[nid]
|
||||
|
||||
|
||||
def _parse_ssdp(peer, rsp, peerdata):
|
||||
@@ -203,11 +226,11 @@ def _parse_ssdp(peer, rsp, peerdata):
|
||||
if code == '200':
|
||||
if nid in peerdata:
|
||||
peerdatum = peerdata[nid]
|
||||
if peer not in peerdatum['peers']:
|
||||
peerdatum['peers'].append(peer)
|
||||
if peer not in peerdatum['addresses']:
|
||||
peerdatum['addresses'].append(peer)
|
||||
else:
|
||||
peerdatum = {
|
||||
'peers': [peer],
|
||||
'addresses': [peer],
|
||||
'hwaddr': mac,
|
||||
}
|
||||
peerdata[nid] = peerdatum
|
||||
|
||||
@@ -17,7 +17,17 @@
|
||||
|
||||
import base64
|
||||
import json
|
||||
import msgpack
|
||||
|
||||
def deserialize_exc(msg):
|
||||
excd = msgpack.unpackb(msg, raw=False)
|
||||
if excd[0] == 'Exception':
|
||||
return Exception(excd[1])
|
||||
if excd[0] not in globals():
|
||||
return Exception('Cannot deserialize: {0}'.format(repr(excd)))
|
||||
if not issubclass(excd[0], ConfluentException):
|
||||
return Exception('Cannot deserialize: {0}'.format(repr(excd)))
|
||||
return globals(excd[0])(*excd[1])
|
||||
|
||||
class ConfluentException(Exception):
|
||||
apierrorcode = 500
|
||||
@@ -27,6 +37,10 @@ class ConfluentException(Exception):
|
||||
errstr = ' - '.join((self._apierrorstr, str(self)))
|
||||
return json.dumps({'error': errstr })
|
||||
|
||||
def serialize(self):
|
||||
return msgpack.packb([self.__class__.__name__, [str(self)]],
|
||||
use_bin_type=False)
|
||||
|
||||
@property
|
||||
def apierrorstr(self):
|
||||
if str(self):
|
||||
@@ -104,6 +118,7 @@ class PubkeyInvalid(ConfluentException):
|
||||
|
||||
def __init__(self, text, certificate, fingerprint, attribname, event):
|
||||
super(PubkeyInvalid, self).__init__(self, text)
|
||||
self.myargs = (text, certificate, fingerprint, attribname, event)
|
||||
self.fingerprint = fingerprint
|
||||
self.attrname = attribname
|
||||
self.message = text
|
||||
@@ -117,6 +132,10 @@ class PubkeyInvalid(ConfluentException):
|
||||
'certificate': certtxt}
|
||||
self.errorbody = json.dumps(bodydata)
|
||||
|
||||
def serialize(self):
|
||||
return msgpack.packb([self.__class__.__name__, self.myargs],
|
||||
use_bin_type=False)
|
||||
|
||||
def get_error_body(self):
|
||||
return self.errorbody
|
||||
|
||||
|
||||
@@ -36,14 +36,31 @@ _tracelog = None
|
||||
|
||||
def execupdate(handler, filename, updateobj, type, owner, node):
|
||||
global _tracelog
|
||||
if type != 'ffdc' and not os.path.exists(filename):
|
||||
errstr = '{0} does not appear to exist on {1}'.format(
|
||||
filename, socket.gethostname())
|
||||
updateobj.handle_progress({'phase': 'error', 'progress': 0.0,
|
||||
'detail': errstr})
|
||||
return
|
||||
if type != 'ffdc':
|
||||
errstr = False
|
||||
if not os.path.exists(filename):
|
||||
errstr = '{0} does not appear to exist on {1}'.format(
|
||||
filename, socket.gethostname())
|
||||
elif not os.access(filename, os.R_OK):
|
||||
errstr = '{0} is not readable by confluent on {1} (ensure confluent user or group can access file and parent directories)'.format(
|
||||
filename, socket.gethostname())
|
||||
if errstr:
|
||||
updateobj.handle_progress({'phase': 'error', 'progress': 0.0,
|
||||
'detail': errstr})
|
||||
return
|
||||
if type == 'ffdc' and os.path.isdir(filename):
|
||||
filename += '/' + node
|
||||
if 'type' == 'ffdc':
|
||||
errstr = False
|
||||
if os.path.exists(filename):
|
||||
errstr = '{0} already exists on {1}, cannot overwrite'.format(
|
||||
filename, socket.gethostname())
|
||||
elif not os.access(os.path.dirname(filename), os.W_OK):
|
||||
errstr = '{0} directory not writable by confluent user/group on {1}, check the directory and parent directory ownership and permissions'.format(filename, socket.gethostname())
|
||||
if errstr:
|
||||
updateobj.handle_progress({'phase': 'error', 'progress': 0.0,
|
||||
'detail': errstr})
|
||||
return
|
||||
try:
|
||||
if type == 'firmware':
|
||||
completion = handler(filename, progress=updateobj.handle_progress,
|
||||
|
||||
@@ -416,7 +416,7 @@ def resourcehandler_backend(env, start_response):
|
||||
reqtype = env['CONTENT_TYPE']
|
||||
operation = opmap[env['REQUEST_METHOD']]
|
||||
querydict = _get_query_dict(env, reqbody, reqtype)
|
||||
if 'restexplorerop' in querydict:
|
||||
if operation != 'retrieve' and 'restexplorerop' in querydict:
|
||||
operation = querydict['restexplorerop']
|
||||
del querydict['restexplorerop']
|
||||
authorized = _authorize_request(env, operation)
|
||||
@@ -511,10 +511,10 @@ def resourcehandler_backend(env, start_response):
|
||||
width = querydict.get('width', 80)
|
||||
height = querydict.get('height', 24)
|
||||
datacallback = None
|
||||
async = None
|
||||
asynchdl = None
|
||||
if 'HTTP_CONFLUENTASYNCID' in env:
|
||||
async = confluent.asynchttp.get_async(env, querydict)
|
||||
termrel = async.set_term_relation(env)
|
||||
asynchdl = confluent.asynchttp.get_async(env, querydict)
|
||||
termrel = asynchdl.set_term_relation(env)
|
||||
datacallback = termrel.got_data
|
||||
try:
|
||||
if shellsession:
|
||||
@@ -537,8 +537,8 @@ def resourcehandler_backend(env, start_response):
|
||||
start_response("500 Internal Server Error", headers)
|
||||
return
|
||||
sessid = _assign_consessionid(consession)
|
||||
if async:
|
||||
async.add_console_session(sessid)
|
||||
if asynchdl:
|
||||
asynchdl.add_console_session(sessid)
|
||||
start_response('200 OK', headers)
|
||||
yield '{"session":"%s","data":""}' % sessid
|
||||
return
|
||||
|
||||
@@ -77,13 +77,16 @@ def _daemonize():
|
||||
print('confluent server starting as pid {0}'.format(thispid))
|
||||
os._exit(0)
|
||||
os.closerange(0, 2)
|
||||
os.umask(63)
|
||||
os.open(os.devnull, os.O_RDWR)
|
||||
os.dup2(0, 1)
|
||||
os.dup2(0, 2)
|
||||
log.daemonized = True
|
||||
|
||||
|
||||
def _redirectoutput():
|
||||
os.umask(63)
|
||||
sys.stdout = log.Logger('stdout', buffered=False)
|
||||
sys.stderr = log.Logger('stderr', buffered=False)
|
||||
log.daemonized = True
|
||||
|
||||
|
||||
def _updatepidfile():
|
||||
@@ -206,7 +209,7 @@ def setlimits():
|
||||
pass
|
||||
|
||||
|
||||
def run():
|
||||
def run(args):
|
||||
setlimits()
|
||||
try:
|
||||
signal.signal(signal.SIGUSR1, dumptrace)
|
||||
@@ -232,7 +235,10 @@ def run():
|
||||
except (OSError, IOError) as e:
|
||||
print(repr(e))
|
||||
sys.exit(1)
|
||||
_daemonize()
|
||||
if '-f' not in args:
|
||||
_daemonize()
|
||||
if '-o' not in args:
|
||||
_redirectoutput()
|
||||
if havefcntl:
|
||||
_updatepidfile()
|
||||
signal.signal(signal.SIGINT, terminate)
|
||||
|
||||
@@ -24,6 +24,7 @@ import confluent.config.conf as cfgfile
|
||||
from copy import deepcopy
|
||||
from datetime import datetime
|
||||
import confluent.util as util
|
||||
import msgpack
|
||||
import json
|
||||
|
||||
try:
|
||||
@@ -84,6 +85,13 @@ def _htmlify_structure(indict):
|
||||
return ret + '</ul>'
|
||||
|
||||
|
||||
def msg_deserialize(packed):
|
||||
m = msgpack.unpackb(packed, raw=False)
|
||||
cls = globals()[m[0]]
|
||||
if issubclass(cls, ConfluentMessage) or issubclass(cls, ConfluentNodeError):
|
||||
return cls(*m[1:])
|
||||
raise Exception("Unknown shenanigans")
|
||||
|
||||
class ConfluentMessage(object):
|
||||
apicode = 200
|
||||
readonly = False
|
||||
@@ -105,6 +113,15 @@ class ConfluentMessage(object):
|
||||
jsonsnippet = json.dumps(datasource, sort_keys=True, separators=(',', ':'))[1:-1]
|
||||
return jsonsnippet
|
||||
|
||||
def serialize(self):
|
||||
msg = [self.__class__.__name__]
|
||||
msg.extend(self.myargs)
|
||||
return msgpack.packb(msg, use_bin_type=False)
|
||||
|
||||
@classmethod
|
||||
def deserialize(cls, data):
|
||||
return cls(*data)
|
||||
|
||||
def raw(self):
|
||||
"""Return pythonic representation of the response.
|
||||
|
||||
@@ -211,6 +228,15 @@ class ConfluentNodeError(object):
|
||||
self.node = node
|
||||
self.error = errorstr
|
||||
|
||||
def serialize(self):
|
||||
return msgpack.packb(
|
||||
[self.__class__.__name__, self.node, self.error],
|
||||
use_bin_type=False)
|
||||
|
||||
@classmethod
|
||||
def deserialize(cls, data):
|
||||
return cls(*data)
|
||||
|
||||
def raw(self):
|
||||
return {'databynode': {self.node: {'errorcode': self.apicode,
|
||||
'error': self.error}}}
|
||||
@@ -221,7 +247,7 @@ class ConfluentNodeError(object):
|
||||
def strip_node(self, node):
|
||||
# NOTE(jjohnson2): For single node errors, raise exception to
|
||||
# trigger what a developer of that medium would expect
|
||||
raise Exception(self.error)
|
||||
raise Exception('{0}: {1}'.format(self.node, self.error))
|
||||
|
||||
|
||||
class ConfluentResourceUnavailable(ConfluentNodeError):
|
||||
@@ -260,9 +286,9 @@ class ConfluentTargetNotFound(ConfluentNodeError):
|
||||
|
||||
class ConfluentTargetInvalidCredentials(ConfluentNodeError):
|
||||
apicode = 502
|
||||
def __init__(self, node):
|
||||
def __init__(self, node, errstr='bad credentials'):
|
||||
self.node = node
|
||||
self.error = 'bad credentials'
|
||||
self.error = errstr
|
||||
|
||||
def strip_node(self, node):
|
||||
raise exc.TargetEndpointBadCredentials
|
||||
@@ -271,6 +297,7 @@ class ConfluentTargetInvalidCredentials(ConfluentNodeError):
|
||||
class DeletedResource(ConfluentMessage):
|
||||
notnode = True
|
||||
def __init__(self, resource):
|
||||
self.myargs = [resource]
|
||||
self.kvpairs = {'deleted': resource}
|
||||
|
||||
def strip_node(self, node):
|
||||
@@ -282,6 +309,7 @@ class CreatedResource(ConfluentMessage):
|
||||
readonly = True
|
||||
|
||||
def __init__(self, resource):
|
||||
self.myargs = [resource]
|
||||
self.kvpairs = {'created': resource}
|
||||
|
||||
def strip_node(self, node):
|
||||
@@ -293,6 +321,7 @@ class RenamedResource(ConfluentMessage):
|
||||
readonly = True
|
||||
|
||||
def __init__(self, oldname, newname):
|
||||
self.myargs = (oldname, newname)
|
||||
self.kvpairs = {'oldname': oldname, 'newname': newname}
|
||||
|
||||
def strip_node(self, node):
|
||||
@@ -301,6 +330,7 @@ class RenamedResource(ConfluentMessage):
|
||||
|
||||
class RenamedNode(ConfluentMessage):
|
||||
def __init__(self, name, rename):
|
||||
self.myargs = (name, rename)
|
||||
self.desc = 'New Name'
|
||||
kv = {'rename': {'value': rename}}
|
||||
self.kvpairs = {name: kv}
|
||||
@@ -311,13 +341,16 @@ class AssignedResource(ConfluentMessage):
|
||||
readonly = True
|
||||
|
||||
def __init__(self, resource):
|
||||
self.myargs = [resource]
|
||||
self.kvpairs = {'assigned': resource}
|
||||
|
||||
|
||||
class ConfluentChoiceMessage(ConfluentMessage):
|
||||
valid_values = set()
|
||||
valid_paramset = {}
|
||||
|
||||
def __init__(self, node, state):
|
||||
self.myargs = (node, state)
|
||||
self.stripped = False
|
||||
self.kvpairs = {
|
||||
node: {
|
||||
@@ -391,6 +424,7 @@ class LinkRelation(ConfluentMessage):
|
||||
|
||||
class ChildCollection(LinkRelation):
|
||||
def __init__(self, collname, candelete=False):
|
||||
self.myargs = (collname, candelete)
|
||||
self.rel = 'item'
|
||||
self.href = collname
|
||||
self.candelete = candelete
|
||||
@@ -511,9 +545,21 @@ class InputFirmwareUpdate(ConfluentMessage):
|
||||
raise Exception('User requested substitutions, but code is '
|
||||
'written against old api, code must be fixed or '
|
||||
'skip {} expansion')
|
||||
if self.filebynode[node].startswith('/etc/confluent'):
|
||||
raise Exception(
|
||||
'File transfer with /etc/confluent is not supported')
|
||||
if self.filebynode[node].startswith('/var/log/confluent'):
|
||||
raise Exception(
|
||||
'File transfer with /var/log/confluent is not supported')
|
||||
return self._filename
|
||||
|
||||
def nodefile(self, node):
|
||||
if self.filebynode[node].startswith('/etc/confluent'):
|
||||
raise Exception(
|
||||
'File transfer with /etc/confluent is not supported')
|
||||
if self.filebynode[node].startswith('/var/log/confluent'):
|
||||
raise Exception(
|
||||
'File transfer with /var/log/confluent is not supported')
|
||||
return self.filebynode[node]
|
||||
|
||||
class InputMedia(InputFirmwareUpdate):
|
||||
@@ -532,11 +578,15 @@ class DetachMedia(ConfluentMessage):
|
||||
|
||||
|
||||
class Media(ConfluentMessage):
|
||||
def __init__(self, node, media):
|
||||
self.kvpairs = {node: {'name': media.name, 'url': media.url}}
|
||||
def __init__(self, node, media=None, rawmedia=None):
|
||||
if media:
|
||||
rawmedia = {'name': media.name, 'url': media.url}
|
||||
self.myargs = (node, None, rawmedia)
|
||||
self.kvpairs = {node: rawmedia}
|
||||
|
||||
class SavedFile(ConfluentMessage):
|
||||
def __init__(self, node, file):
|
||||
self.myargs = (node, file)
|
||||
self.kvpairs = {node: {'filename': file}}
|
||||
|
||||
class InputAlertData(ConfluentMessage):
|
||||
@@ -624,6 +674,8 @@ class InputAttributes(ConfluentMessage):
|
||||
if nodes is None:
|
||||
self.attribs = inputdata
|
||||
for attrib in self.attribs:
|
||||
if not cfm.attrib_supports_expression(attrib):
|
||||
continue
|
||||
if type(self.attribs[attrib]) in (bytes, unicode):
|
||||
try:
|
||||
# ok, try to use format against the string
|
||||
@@ -650,7 +702,7 @@ class InputAttributes(ConfluentMessage):
|
||||
return {}
|
||||
nodeattr = deepcopy(self.nodeattribs[node])
|
||||
for attr in nodeattr:
|
||||
if type(nodeattr[attr]) in (bytes, unicode):
|
||||
if type(nodeattr[attr]) in (bytes, unicode) and cfm.attrib_supports_expression(attr):
|
||||
try:
|
||||
# as above, use format() to see if string follows
|
||||
# expression, store value back in case of escapes
|
||||
@@ -1108,6 +1160,7 @@ class BootDevice(ConfluentChoiceMessage):
|
||||
}
|
||||
|
||||
def __init__(self, node, device, bootmode='unspecified', persistent=False):
|
||||
self.myargs = (node, device, bootmode, persistent)
|
||||
if device not in self.valid_values:
|
||||
raise Exception("Invalid boot device argument passed in:" +
|
||||
repr(device))
|
||||
@@ -1206,10 +1259,10 @@ class PowerState(ConfluentChoiceMessage):
|
||||
|
||||
def __init__(self, node, state, oldstate=None):
|
||||
super(PowerState, self).__init__(node, state)
|
||||
self.myargs = (node, state, oldstate)
|
||||
if oldstate is not None:
|
||||
self.kvpairs[node]['oldstate'] = {'value': oldstate}
|
||||
|
||||
|
||||
class BMCReset(ConfluentChoiceMessage):
|
||||
valid_values = set([
|
||||
'reset',
|
||||
@@ -1225,13 +1278,13 @@ class NTPEnabled(ConfluentChoiceMessage):
|
||||
|
||||
def __init__(self, node, enabled):
|
||||
self.stripped = False
|
||||
self.myargs = (node, enabled)
|
||||
self.kvpairs = {
|
||||
node: {
|
||||
'state': {'value': str(enabled)},
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
class EventCollection(ConfluentMessage):
|
||||
"""A collection of events
|
||||
|
||||
@@ -1251,6 +1304,8 @@ class EventCollection(ConfluentMessage):
|
||||
def __init__(self, events=(), name=None):
|
||||
eventdata = []
|
||||
self.notnode = name is None
|
||||
self.myname = name
|
||||
self.myargs = (eventdata, name)
|
||||
for event in events:
|
||||
entry = {
|
||||
'id': event.get('id', None),
|
||||
@@ -1278,6 +1333,10 @@ class AsyncCompletion(ConfluentMessage):
|
||||
self.stripped = True
|
||||
self.notnode = True
|
||||
|
||||
@classmethod
|
||||
def deserialize(cls):
|
||||
raise Exception("Not supported")
|
||||
|
||||
def raw(self):
|
||||
return {'_requestdone': True}
|
||||
|
||||
@@ -1288,6 +1347,10 @@ class AsyncMessage(ConfluentMessage):
|
||||
self.notnode = True
|
||||
self.msgpair = pair
|
||||
|
||||
@classmethod
|
||||
def deserialize(cls):
|
||||
raise Exception("Not supported")
|
||||
|
||||
def raw(self):
|
||||
rsp = self.msgpair[1]
|
||||
rspdict = None
|
||||
@@ -1319,6 +1382,7 @@ class User(ConfluentMessage):
|
||||
self.desc = 'foo'
|
||||
self.stripped = False
|
||||
self.notnode = name is None
|
||||
self.myargs = (uid, username, privilege_level, name, expiration)
|
||||
kvpairs = {'username': {'value': username},
|
||||
'password': {'value': '', 'type': 'password'},
|
||||
'privilege_level': {'value': privilege_level},
|
||||
@@ -1338,7 +1402,11 @@ class UserCollection(ConfluentMessage):
|
||||
self.notnode = name is None
|
||||
self.desc = 'list of users'
|
||||
userlist = []
|
||||
self.myargs = (userlist, name)
|
||||
for user in users:
|
||||
if 'username' in user: # processing an already translated dict
|
||||
userlist.append(user)
|
||||
continue
|
||||
entry = {
|
||||
'uid': user['uid'],
|
||||
'username': user['name'],
|
||||
@@ -1352,8 +1420,10 @@ class UserCollection(ConfluentMessage):
|
||||
self.kvpairs = {name: {'users': userlist}}
|
||||
|
||||
|
||||
|
||||
class AlertDestination(ConfluentMessage):
|
||||
def __init__(self, ip, acknowledge=False, acknowledge_timeout=None, retries=0, name=None):
|
||||
self.myargs = (ip, acknowledge, acknowledge_timeout, retries, name)
|
||||
self.desc = 'foo'
|
||||
self.stripped = False
|
||||
self.notnode = name is None
|
||||
@@ -1418,7 +1488,11 @@ class SensorReadings(ConfluentMessage):
|
||||
def __init__(self, sensors=(), name=None):
|
||||
readings = []
|
||||
self.notnode = name is None
|
||||
self.myargs = (readings, name)
|
||||
for sensor in sensors:
|
||||
if isinstance(sensor, dict):
|
||||
readings.append(sensor)
|
||||
continue
|
||||
sensordict = {'name': sensor.name}
|
||||
if hasattr(sensor, 'value'):
|
||||
sensordict['value'] = sensor.value
|
||||
@@ -1443,6 +1517,13 @@ class Firmware(ConfluentMessage):
|
||||
readonly = True
|
||||
|
||||
def __init__(self, data, name):
|
||||
for datum in data:
|
||||
for component in datum:
|
||||
for field in datum[component]:
|
||||
tdatum = datum[component]
|
||||
if isinstance(tdatum[field], datetime):
|
||||
tdatum[field] = tdatum[field].strftime('%Y-%m-%dT%H:%M:%S')
|
||||
self.myargs = (data, name)
|
||||
self.notnode = name is None
|
||||
self.desc = 'Firmware information'
|
||||
if self.notnode:
|
||||
@@ -1455,6 +1536,7 @@ class KeyValueData(ConfluentMessage):
|
||||
readonly = True
|
||||
|
||||
def __init__(self, kvdata, name=None):
|
||||
self.myargs = (kvdata, name)
|
||||
self.notnode = name is None
|
||||
if self.notnode:
|
||||
self.kvpairs = kvdata
|
||||
@@ -1464,6 +1546,7 @@ class KeyValueData(ConfluentMessage):
|
||||
class Array(ConfluentMessage):
|
||||
def __init__(self, name, disks=None, raid=None, volumes=None,
|
||||
id=None, capacity=None, available=None):
|
||||
self.myargs = (name, disks, raid, volumes, id, capacity, available)
|
||||
self.kvpairs = {
|
||||
name: {
|
||||
'type': 'array',
|
||||
@@ -1478,6 +1561,7 @@ class Array(ConfluentMessage):
|
||||
|
||||
class Volume(ConfluentMessage):
|
||||
def __init__(self, name, volname, size, state, array, stripsize=None):
|
||||
self.myargs = (name, volname, size, state, array, stripsize)
|
||||
self.kvpairs = {
|
||||
name: {
|
||||
'type': 'volume',
|
||||
@@ -1518,6 +1602,8 @@ class Disk(ConfluentMessage):
|
||||
def __init__(self, name, label=None, description=None,
|
||||
diskid=None, state=None, serial=None, fru=None,
|
||||
array=None):
|
||||
self.myargs = (name, label, description, diskid, state,
|
||||
serial, fru, array)
|
||||
state = self._normalize_state(state)
|
||||
self.kvpairs = {
|
||||
name: {
|
||||
@@ -1539,6 +1625,7 @@ class LEDStatus(ConfluentMessage):
|
||||
readonly = True
|
||||
|
||||
def __init__(self, data, name):
|
||||
self.myargs = (data, name)
|
||||
self.notnode = name is None
|
||||
self.desc = 'led status'
|
||||
|
||||
@@ -1553,6 +1640,7 @@ class NetworkConfiguration(ConfluentMessage):
|
||||
|
||||
def __init__(self, name=None, ipv4addr=None, ipv4gateway=None,
|
||||
ipv4cfgmethod=None, hwaddr=None):
|
||||
self.myargs = (name, ipv4addr, ipv4gateway, ipv4cfgmethod, hwaddr)
|
||||
self.notnode = name is None
|
||||
self.stripped = False
|
||||
|
||||
@@ -1573,6 +1661,7 @@ class HealthSummary(ConfluentMessage):
|
||||
valid_values = valid_health_values
|
||||
|
||||
def __init__(self, health, name=None):
|
||||
self.myargs = (health, name)
|
||||
self.stripped = False
|
||||
self.notnode = name is None
|
||||
if health not in self.valid_values:
|
||||
@@ -1585,6 +1674,7 @@ class HealthSummary(ConfluentMessage):
|
||||
|
||||
class Attributes(ConfluentMessage):
|
||||
def __init__(self, name=None, kv=None, desc=''):
|
||||
self.myargs = (name, kv, desc)
|
||||
self.desc = desc
|
||||
nkv = {}
|
||||
self.notnode = name is None
|
||||
@@ -1605,6 +1695,7 @@ class ConfigSet(Attributes):
|
||||
|
||||
class ListAttributes(ConfluentMessage):
|
||||
def __init__(self, name=None, kv=None, desc=''):
|
||||
self.myargs = (name, kv, desc)
|
||||
self.desc = desc
|
||||
self.notnode = name is None
|
||||
if self.notnode:
|
||||
@@ -1615,6 +1706,7 @@ class ListAttributes(ConfluentMessage):
|
||||
|
||||
class MCI(ConfluentMessage):
|
||||
def __init__(self, name=None, mci=None):
|
||||
self.myargs = (name, mci)
|
||||
self.notnode = name is None
|
||||
self.desc = 'BMC identifier'
|
||||
|
||||
@@ -1627,6 +1719,7 @@ class MCI(ConfluentMessage):
|
||||
|
||||
class Hostname(ConfluentMessage):
|
||||
def __init__(self, name=None, hostname=None):
|
||||
self.myargs = (name, hostname)
|
||||
self.notnode = name is None
|
||||
self.desc = 'BMC hostname'
|
||||
|
||||
@@ -1638,6 +1731,7 @@ class Hostname(ConfluentMessage):
|
||||
|
||||
class DomainName(ConfluentMessage):
|
||||
def __init__(self, name=None, dn=None):
|
||||
self.myargs = (name, dn)
|
||||
self.notnode = name is None
|
||||
self.desc = 'BMC domain name'
|
||||
|
||||
@@ -1652,6 +1746,7 @@ class NTPServers(ConfluentMessage):
|
||||
readonly = True
|
||||
|
||||
def __init__(self, name=None, servers=None):
|
||||
self.myargs = (name, servers)
|
||||
self.notnode = name is None
|
||||
self.desc = 'NTP Server'
|
||||
|
||||
@@ -1666,6 +1761,7 @@ class NTPServers(ConfluentMessage):
|
||||
|
||||
class NTPServer(ConfluentMessage):
|
||||
def __init__(self, name=None, server=None):
|
||||
self.myargs = (name, server)
|
||||
self.notnode = name is None
|
||||
self.desc = 'NTP Server'
|
||||
|
||||
@@ -1682,6 +1778,7 @@ class License(ConfluentMessage):
|
||||
readonly = True
|
||||
|
||||
def __init__(self, name=None, kvm=None, feature=None, state=None):
|
||||
self.myargs = (name, kvm, feature, state)
|
||||
self.notnode = name is None
|
||||
self.desc = 'License'
|
||||
|
||||
@@ -1697,6 +1794,7 @@ class CryptedAttributes(Attributes):
|
||||
defaulttype = 'password'
|
||||
|
||||
def __init__(self, name=None, kv=None, desc=''):
|
||||
self.myargs = (name, kv, desc)
|
||||
# for now, just keep the dictionary keys and discard crypt value
|
||||
self.desc = desc
|
||||
nkv = {}
|
||||
|
||||
@@ -44,7 +44,7 @@ import eventlet
|
||||
from eventlet.greenpool import GreenPool
|
||||
import eventlet.semaphore
|
||||
import re
|
||||
|
||||
webclient = eventlet.import_patched('pyghmi.util.webclient')
|
||||
# The interesting OIDs are:
|
||||
# lldpLocChassisId - to cross reference (1.0.8802.1.1.2.1.3.2.0)
|
||||
# lldpLocPortId - for cross referencing.. (1.0.8802.1.1.2.1.3.7.1.3)
|
||||
@@ -85,6 +85,7 @@ _neighdata = {}
|
||||
_neighbypeerid = {}
|
||||
_updatelocks = {}
|
||||
_chassisidbyswitch = {}
|
||||
_noaffluent = set([])
|
||||
|
||||
def lenovoname(idx, desc):
|
||||
if desc.isdigit():
|
||||
@@ -171,21 +172,58 @@ def _init_lldp(data, iname, idx, idxtoportid, switch):
|
||||
data[iname] = {'port': iname, 'portid': str(idxtoportid[idx]),
|
||||
'chassisid': _chassisidbyswitch[switch]}
|
||||
|
||||
def _extract_neighbor_data_affluent(switch, user, password, cfm, lldpdata):
|
||||
kv = util.TLSCertVerifier(cfm, switch,
|
||||
'pubkeys.tls_hardwaremanager').verify_cert
|
||||
wc = webclient.SecureHTTPConnection(
|
||||
switch, 443, verifycallback=kv, timeout=5)
|
||||
wc.set_basic_credentials(user, password)
|
||||
neighdata = wc.grab_json_response('/affluent/lldp/all')
|
||||
chassisid = neighdata['chassis']['id']
|
||||
_chassisidbyswitch[switch] = chassisid,
|
||||
for record in neighdata['neighbors']:
|
||||
localport = record['localport']
|
||||
peerid = '{0}.{1}'.format(
|
||||
record.get('peerchassisid', '').replace(':', '-').replace('/', '-'),
|
||||
record.get('peerportid', '').replace(':', '-').replace('/', '-'),
|
||||
)
|
||||
portdata = {
|
||||
'verified': True, # It is over TLS after all
|
||||
'peerdescription': record.get('peerdescription', None),
|
||||
'peerchassisid': record['peerchassisid'],
|
||||
'peername': record['peername'],
|
||||
'switch': switch,
|
||||
'chassisid': chassisid,
|
||||
'portid': record['localport'],
|
||||
'peerportid': record['peerportid'],
|
||||
'port': record['localport'],
|
||||
'peerid': peerid,
|
||||
}
|
||||
_neighbypeerid[peerid] = portdata
|
||||
lldpdata[localport] = portdata
|
||||
neighdata[switch] = lldpdata
|
||||
|
||||
|
||||
def _extract_neighbor_data_b(args):
|
||||
"""Build LLDP data about elements connected to switch
|
||||
|
||||
args are carried as a tuple, because of eventlet convenience
|
||||
"""
|
||||
switch, password, user, force = args[:4]
|
||||
switch, password, user, cfm, force = args[:5]
|
||||
vintage = _neighdata.get(switch, {}).get('!!vintage', 0)
|
||||
now = util.monotonic_time()
|
||||
if vintage > (now - 60) and not force:
|
||||
return
|
||||
lldpdata = {'!!vintage': now}
|
||||
try:
|
||||
return _extract_neighbor_data_affluent(switch, user, password, cfm, lldpdata)
|
||||
except Exception:
|
||||
pass
|
||||
conn = snmp.Session(switch, password, user)
|
||||
sid = None
|
||||
lldpdata = {'!!vintage': now}
|
||||
for sysid in conn.walk('1.3.6.1.2.1.1.2'):
|
||||
sid = str(sysid[1][6:])
|
||||
_noaffluent.add(switch)
|
||||
idxtoifname = {}
|
||||
idxtoportid = {}
|
||||
_chassisidbyswitch[switch] = sanitize(list(
|
||||
@@ -268,8 +306,8 @@ def _extract_neighbor_data(args):
|
||||
return _extract_neighbor_data_b(args)
|
||||
except Exception as e:
|
||||
yieldexc = False
|
||||
if len(args) >= 5:
|
||||
yieldexc = args[4]
|
||||
if len(args) >= 6:
|
||||
yieldexc = args[5]
|
||||
if yieldexc:
|
||||
return e
|
||||
else:
|
||||
@@ -358,10 +396,3 @@ def _handle_neighbor_query(pathcomponents, configmanager):
|
||||
raise x
|
||||
return list_info(parms, listrequested)
|
||||
|
||||
|
||||
def _list_interfaces(switchname, configmanager):
|
||||
switchcreds = get_switchcreds(configmanager, (switchname,))
|
||||
switchcreds = switchcreds[0]
|
||||
conn = snmp.Session(*switchcreds)
|
||||
ifnames = netutil.get_portnamemap(conn)
|
||||
return util.natural_sort(ifnames.values())
|
||||
@@ -45,13 +45,16 @@ from eventlet.greenpool import GreenPool
|
||||
import eventlet
|
||||
import eventlet.semaphore
|
||||
import re
|
||||
webclient = eventlet.import_patched('pyghmi.util.webclient')
|
||||
|
||||
|
||||
noaffluent = set([])
|
||||
|
||||
_macmap = {}
|
||||
_apimacmap = {}
|
||||
_macsbyswitch = {}
|
||||
_nodesbymac = {}
|
||||
_switchportmap = {}
|
||||
_neighdata = {}
|
||||
vintage = None
|
||||
|
||||
|
||||
@@ -127,6 +130,36 @@ def _nodelookup(switch, ifname):
|
||||
return None
|
||||
|
||||
|
||||
def _affluent_map_switch(args):
|
||||
switch, password, user, cfm = args
|
||||
kv = util.TLSCertVerifier(cfm, switch,
|
||||
'pubkeys.tls_hardwaremanager').verify_cert
|
||||
wc = webclient.SecureHTTPConnection(
|
||||
switch, 443, verifycallback=kv, timeout=5)
|
||||
wc.set_basic_credentials(user, password)
|
||||
macs = wc.grab_json_response('/affluent/macs/by-port')
|
||||
_macsbyswitch[switch] = macs
|
||||
|
||||
for iface in macs:
|
||||
nummacs = len(macs[iface])
|
||||
for mac in macs[iface]:
|
||||
if mac in _macmap:
|
||||
_macmap[mac].append((switch, iface, nummacs))
|
||||
else:
|
||||
_macmap[mac] = [(switch, iface, nummacs)]
|
||||
nodename = _nodelookup(switch, iface)
|
||||
if nodename is not None:
|
||||
if mac in _nodesbymac and _nodesbymac[mac][0] != nodename:
|
||||
# For example, listed on both a real edge port
|
||||
# and by accident a trunk port
|
||||
log.log({'error': '{0} and {1} described by ambiguous'
|
||||
' switch topology values'.format(
|
||||
nodename, _nodesbymac[mac][0])})
|
||||
_nodesbymac[mac] = (None, None)
|
||||
else:
|
||||
_nodesbymac[mac] = (nodename, nummacs)
|
||||
|
||||
|
||||
def _map_switch_backend(args):
|
||||
"""Manipulate portions of mac address map relevant to a given switch
|
||||
"""
|
||||
@@ -144,13 +177,18 @@ def _map_switch_backend(args):
|
||||
# fallback if ifName is empty
|
||||
#
|
||||
global _macmap
|
||||
if len(args) == 3:
|
||||
switch, password, user = args
|
||||
if len(args) == 4:
|
||||
switch, password, user, cfm = args
|
||||
if not user:
|
||||
user = None
|
||||
else:
|
||||
switch, password = args
|
||||
user = None
|
||||
if switch not in noaffluent:
|
||||
try:
|
||||
return _affluent_map_switch(args)
|
||||
except Exception:
|
||||
pass
|
||||
haveqbridge = False
|
||||
mactobridge = {}
|
||||
conn = snmp.Session(switch, password, user)
|
||||
@@ -164,6 +202,7 @@ def _map_switch_backend(args):
|
||||
*([int(x) for x in oid[-6:]])
|
||||
)
|
||||
mactobridge[macaddr] = int(bridgeport)
|
||||
noaffluent.add(switch)
|
||||
if not haveqbridge:
|
||||
for vb in conn.walk('1.3.6.1.2.1.17.4.3.1.2'):
|
||||
oid, bridgeport = vb
|
||||
|
||||
@@ -36,7 +36,7 @@ def get_switchcreds(configmanager, switches):
|
||||
'secret.hardwaremanagementuser', {}).get('value', None)
|
||||
if not user:
|
||||
user = None
|
||||
switchauth.append((switch, password, user))
|
||||
switchauth.append((switch, password, user, configmanager))
|
||||
return switchauth
|
||||
|
||||
|
||||
|
||||
@@ -250,6 +250,8 @@ class NodeRange(object):
|
||||
return nodes
|
||||
if ':' in element: # : range for less ambiguity
|
||||
return self.expandrange(element, ':')
|
||||
elif '..' in element:
|
||||
return self.expandrange(element, '..')
|
||||
elif '-' in element:
|
||||
return self.expandrange(element, '-')
|
||||
elif '+' in element:
|
||||
|
||||
@@ -0,0 +1,153 @@
|
||||
|
||||
# Copyright 2019-2020 Lenovo
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
|
||||
import eventlet
|
||||
import eventlet.queue as queue
|
||||
import confluent.exceptions as exc
|
||||
webclient = eventlet.import_patched('pyghmi.util.webclient')
|
||||
import confluent.messages as msg
|
||||
import confluent.util as util
|
||||
|
||||
class SwitchSensor(object):
|
||||
def __init__(self, name, states, value=None, health=None):
|
||||
self.name = name
|
||||
self.value = value
|
||||
self.states = states
|
||||
self.health = health
|
||||
|
||||
class WebClient(object):
|
||||
def __init__(self, node, configmanager, creds):
|
||||
self.node = node
|
||||
self.wc = webclient.SecureHTTPConnection(node, port=443, verifycallback=util.TLSCertVerifier(
|
||||
configmanager, node, 'pubkeys.tls_hardwaremanager').verify_cert)
|
||||
self.wc.set_basic_credentials(creds[node]['secret.hardwaremanagementuser']['value'], creds[node]['secret.hardwaremanagementpassword']['value'])
|
||||
|
||||
def fetch(self, url, results):
|
||||
rsp, status = self.wc.grab_json_response_with_status(url)
|
||||
if status == 401:
|
||||
results.put(msg.ConfluentTargetInvalidCredentials(self.node, 'Unable to authenticate'))
|
||||
return {}
|
||||
elif status != 200:
|
||||
results.put(msg.ConfluentNodeError(self.node, 'Unknown error: ' + rsp + ' while retrieving ' + url))
|
||||
return {}
|
||||
return rsp
|
||||
|
||||
|
||||
def update(nodes, element, configmanager, inputdata):
|
||||
for node in nodes:
|
||||
yield msg.ConfluentNodeError(node, 'Not Implemented')
|
||||
|
||||
|
||||
def delete(nodes, element, configmanager, inputdata):
|
||||
for node in nodes:
|
||||
yield msg.ConfluentNodeError(node, 'Not Implemented')
|
||||
|
||||
|
||||
def create(nodes, element, configmanager, inputdata):
|
||||
for node in nodes:
|
||||
yield msg.ConfluentNodeError(node, 'Not Implemented')
|
||||
|
||||
|
||||
def _run_method(method, workers, results, configmanager, nodes, element):
|
||||
creds = configmanager.get_node_attributes(
|
||||
nodes, ['secret.hardwaremanagementuser', 'secret.hardwaremanagementpassword'], decrypt=True)
|
||||
for node in nodes:
|
||||
workers.add(eventlet.spawn(method, configmanager, creds,
|
||||
node, results, element))
|
||||
|
||||
def retrieve(nodes, element, configmanager, inputdata):
|
||||
results = queue.LightQueue()
|
||||
workers = set([])
|
||||
if element == ['power', 'state']:
|
||||
for node in nodes:
|
||||
yield msg.PowerState(node=node, state='on')
|
||||
return
|
||||
elif element == ['health', 'hardware']:
|
||||
_run_method(retrieve_health, workers, results, configmanager, nodes, element)
|
||||
elif element[:3] == ['inventory', 'hardware', 'all']:
|
||||
_run_method(retrieve_inventory, workers, results, configmanager, nodes, element)
|
||||
elif element[:3] == ['inventory', 'firmware', 'all']:
|
||||
_run_method(retrieve_firmware, workers, results, configmanager, nodes, element)
|
||||
elif element == ['sensors', 'hardware', 'all']:
|
||||
_run_method(list_sensors, workers, results, configmanager, nodes, element)
|
||||
elif element[:3] == ['sensors', 'hardware', 'all']:
|
||||
_run_method(retrieve_sensors, workers, results, configmanager, nodes, element)
|
||||
else:
|
||||
for node in nodes:
|
||||
yield msg.ConfluentNodeError(node, 'Not Implemented')
|
||||
return
|
||||
while workers:
|
||||
try:
|
||||
datum = results.get(10)
|
||||
while datum:
|
||||
if datum:
|
||||
yield datum
|
||||
datum = results.get_nowait()
|
||||
except queue.Empty:
|
||||
pass
|
||||
eventlet.sleep(0.001)
|
||||
for t in list(workers):
|
||||
if t.dead:
|
||||
workers.discard(t)
|
||||
try:
|
||||
while True:
|
||||
datum = results.get_nowait()
|
||||
if datum:
|
||||
yield datum
|
||||
except queue.Empty:
|
||||
pass
|
||||
|
||||
|
||||
def retrieve_inventory(configmanager, creds, node, results, element):
|
||||
if len(element) == 3:
|
||||
results.put(msg.ChildCollection('all'))
|
||||
results.put(msg.ChildCollection('system'))
|
||||
return
|
||||
wc = WebClient(node, configmanager, creds)
|
||||
invinfo = wc.fetch('/affluent/inventory/hardware/all', results)
|
||||
if invinfo:
|
||||
results.put(msg.KeyValueData(invinfo, node))
|
||||
|
||||
|
||||
def retrieve_firmware(configmanager, creds, node, results, element):
|
||||
if len(element) == 3:
|
||||
results.put(msg.ChildCollection('all'))
|
||||
return
|
||||
wc = WebClient(node, configmanager, creds)
|
||||
fwinfo = wc.fetch('/affluent/inventory/firmware/all', results)
|
||||
if fwinfo:
|
||||
results.put(msg.Firmware(fwinfo, node))
|
||||
|
||||
def list_sensors(configmanager, creds, node, results, element):
|
||||
wc = WebClient(node, configmanager, creds)
|
||||
sensors = wc.fetch('/affluent/sensors/hardware/all', results)
|
||||
for sensor in sensors['item']:
|
||||
results.put(msg.ChildCollection(sensor))
|
||||
|
||||
def retrieve_sensors(configmanager, creds, node, results, element):
|
||||
wc = WebClient(node, configmanager, creds)
|
||||
sensors = wc.fetch('/affluent/sensors/hardware/all/{0}'.format(element[-1]), results)
|
||||
if sensors:
|
||||
results.put(msg.SensorReadings(sensors['sensors'], node))
|
||||
|
||||
|
||||
|
||||
def retrieve_health(configmanager, creds, node, results, element):
|
||||
wc = WebClient(node, configmanager, creds)
|
||||
hinfo = wc.fetch('/affluent/health', results)
|
||||
if hinfo:
|
||||
results.put(msg.HealthSummary(hinfo.get('health', 'unknown'), name=node))
|
||||
results.put(msg.SensorReadings(hinfo.get('sensors', []), name=node))
|
||||
@@ -29,6 +29,7 @@ import eventlet.queue as queue
|
||||
import eventlet.support.greendns
|
||||
from fnmatch import fnmatch
|
||||
import os
|
||||
import pwd
|
||||
import pyghmi.constants as pygconstants
|
||||
import pyghmi.exceptions as pygexc
|
||||
import pyghmi.storage as storage
|
||||
@@ -646,8 +647,10 @@ class IpmiHandler(object):
|
||||
return self.handle_ntp()
|
||||
elif self.element[1:4] == ['management_controller', 'extended', 'all']:
|
||||
return self.handle_bmcconfig()
|
||||
elif self.element[1:4] == ['management_controller', 'extended', 'all']:
|
||||
elif self.element[1:4] == ['management_controller', 'extended', 'advanced']:
|
||||
return self.handle_bmcconfig(True)
|
||||
elif self.element[1:4] == ['management_controller', 'extended', 'extra']:
|
||||
return self.handle_bmcconfig(True, extended=True)
|
||||
elif self.element[1:3] == ['system', 'all']:
|
||||
return self.handle_sysconfig()
|
||||
elif self.element[1:3] == ['system', 'advanced']:
|
||||
@@ -853,6 +856,10 @@ class IpmiHandler(object):
|
||||
raise
|
||||
if hasattr(reading, 'health'):
|
||||
reading.health = _str_health(reading.health)
|
||||
if hasattr(reading, 'unavailable') and reading.unavailable:
|
||||
self.output.put(msg.SensorReadings([EmptySensor(
|
||||
reading.name)], name=self.node))
|
||||
continue
|
||||
readings.append(reading)
|
||||
self.output.put(msg.SensorReadings(readings, name=self.node))
|
||||
else:
|
||||
@@ -868,9 +875,13 @@ class IpmiHandler(object):
|
||||
self.ipmicmd.sensormap[sensorname])
|
||||
if hasattr(reading, 'health'):
|
||||
reading.health = _str_health(reading.health)
|
||||
self.output.put(
|
||||
msg.SensorReadings([reading],
|
||||
name=self.node))
|
||||
if hasattr(reading, 'unavailable') and reading.unavailable:
|
||||
self.output.put(msg.SensorReadings([EmptySensor(
|
||||
reading.name)], name=self.node))
|
||||
else:
|
||||
self.output.put(
|
||||
msg.SensorReadings([reading],
|
||||
name=self.node))
|
||||
except pygexc.IpmiException as ie:
|
||||
if ie.ipmicode == 203:
|
||||
self.output.put(msg.ConfluentResourceUnavailable(
|
||||
@@ -1024,6 +1035,8 @@ class IpmiHandler(object):
|
||||
return self._create_storage(storelem)
|
||||
|
||||
def _delete_storage(self, storelem):
|
||||
if len(storelem) < 2:
|
||||
storelem.append('')
|
||||
if len(storelem) < 2 or storelem[0] != 'volumes':
|
||||
raise exc.InvalidArgumentException('Must target a specific volume')
|
||||
volname = storelem[-1]
|
||||
@@ -1412,12 +1425,14 @@ class IpmiHandler(object):
|
||||
'Cannot read the "clear" resource')
|
||||
self.ipmicmd.clear_system_configuration()
|
||||
|
||||
def handle_bmcconfig(self, advanced=False):
|
||||
def handle_bmcconfig(self, advanced=False, extended=False):
|
||||
if 'read' == self.op:
|
||||
try:
|
||||
self.output.put(msg.ConfigSet(
|
||||
self.node,
|
||||
self.ipmicmd.get_bmc_configuration()))
|
||||
if extended:
|
||||
bmccfg = self.ipmicmd.get_extended_bmc_configuration()
|
||||
else:
|
||||
bmccfg = self.ipmicmd.get_bmc_configuration()
|
||||
self.output.put(msg.ConfigSet(self.node, bmccfg))
|
||||
except Exception as e:
|
||||
self.output.put(
|
||||
msg.ConfluentNodeError(self.node, str(e)))
|
||||
@@ -1490,14 +1505,33 @@ class IpmiHandler(object):
|
||||
|
||||
def save_licenses(self):
|
||||
directory = self.inputdata.nodefile(self.node)
|
||||
checkdir = directory
|
||||
if not os.access(directory, os.W_OK):
|
||||
raise exc.InvalidArgumentException(
|
||||
'The confluent system user/group is unable to write to '
|
||||
'directory {0}, check ownership and permissions'.format(
|
||||
checkdir))
|
||||
for saved in self.ipmicmd.save_licenses(directory):
|
||||
try:
|
||||
pwent = pwd.getpwnam(self.current_user)
|
||||
os.chown(saved, pwent.pw_uid, pwent.pw_gid)
|
||||
except KeyError:
|
||||
pass
|
||||
self.output.put(msg.SavedFile(self.node, saved))
|
||||
|
||||
def handle_licenses(self):
|
||||
if self.element[-1] == '':
|
||||
self.element = self.element[:-1]
|
||||
if self.op in ('create', 'update'):
|
||||
self.ipmicmd.apply_license(self.inputdata.nodefile(self.node))
|
||||
filename = self.inputdata.nodefile(self.node)
|
||||
if not os.access(filename, os.R_OK):
|
||||
errstr = ('{0} is not readable by confluent on {1} '
|
||||
'(ensure confluent user or group can access file '
|
||||
'and parent directories)').format(
|
||||
filename, socket.gethostname())
|
||||
self.output.put(msg.ConfluentNodeError(self.node, errstr))
|
||||
return
|
||||
self.ipmicmd.apply_license(filename)
|
||||
if len(self.element) == 3:
|
||||
self.output.put(msg.ChildCollection('all'))
|
||||
i = 1
|
||||
|
||||
@@ -27,6 +27,7 @@ import eventlet.queue as queue
|
||||
import eventlet.support.greendns
|
||||
from fnmatch import fnmatch
|
||||
import os
|
||||
import pwd
|
||||
import pyghmi.constants as pygconstants
|
||||
import pyghmi.exceptions as pygexc
|
||||
import pyghmi.storage as storage
|
||||
@@ -893,6 +894,8 @@ class IpmiHandler(object):
|
||||
return self._create_storage(storelem)
|
||||
|
||||
def _delete_storage(self, storelem):
|
||||
if len(storelem) < 2:
|
||||
storelem.append('')
|
||||
if len(storelem) < 2 or storelem[0] != 'volumes':
|
||||
raise exc.InvalidArgumentException('Must target a specific volume')
|
||||
volname = storelem[-1]
|
||||
@@ -1347,13 +1350,31 @@ class IpmiHandler(object):
|
||||
|
||||
def save_licenses(self):
|
||||
directory = self.inputdata.nodefile(self.node)
|
||||
if not os.access(directory, os.W_OK):
|
||||
raise exc.InvalidArgumentException(
|
||||
'The confluent system user/group is unable to write to '
|
||||
'directory {0}, check ownership and permissions'.format(
|
||||
directory))
|
||||
for saved in self.ipmicmd.save_licenses(directory):
|
||||
try:
|
||||
pwent = pwd.getpwnam(self.current_user)
|
||||
os.chown(saved, pwent.pw_uid, pwent.pw_gid)
|
||||
except KeyError:
|
||||
pass
|
||||
self.output.put(msg.SavedFile(self.node, saved))
|
||||
|
||||
def handle_licenses(self):
|
||||
if self.element[-1] == '':
|
||||
self.element = self.element[:-1]
|
||||
if self.op in ('create', 'update'):
|
||||
filename = self.inputdata.nodefile(self.node)
|
||||
if not os.access(filename, os.R_OK):
|
||||
errstr = ('{0} is not readable by confluent on {1} '
|
||||
'(ensure confluent user or group can access file '
|
||||
'and parent directories)').format(
|
||||
filename, socket.gethostname())
|
||||
self.output.put(msg.ConfluentNodeError(self.node, errstr))
|
||||
return
|
||||
self.ipmicmd.apply_license(self.inputdata.nodefile(self.node))
|
||||
if len(self.element) == 3:
|
||||
self.output.put(msg.ChildCollection('all'))
|
||||
|
||||
@@ -99,8 +99,10 @@ class SshShell(conapi.Console):
|
||||
def recvdata(self):
|
||||
while self.connected:
|
||||
pendingdata = self.shell.recv(8192)
|
||||
if pendingdata == '':
|
||||
self.datacallback(conapi.ConsoleEvent.Disconnect)
|
||||
if not pendingdata:
|
||||
self.ssh.close()
|
||||
if self.datacallback:
|
||||
self.datacallback(conapi.ConsoleEvent.Disconnect)
|
||||
return
|
||||
self.datacallback(pendingdata)
|
||||
|
||||
@@ -110,7 +112,7 @@ class SshShell(conapi.Console):
|
||||
# that would rather not use the nodename as anything but an opaque
|
||||
# identifier
|
||||
self.datacallback = callback
|
||||
if self.username is not '':
|
||||
if self.username is not b'':
|
||||
self.logon()
|
||||
else:
|
||||
self.inputmode = 0
|
||||
@@ -126,12 +128,14 @@ class SshShell(conapi.Console):
|
||||
password=self.password, allow_agent=False,
|
||||
look_for_keys=False)
|
||||
except paramiko.AuthenticationException:
|
||||
self.ssh.close()
|
||||
self.inputmode = 0
|
||||
self.username = b''
|
||||
self.password = b''
|
||||
self.datacallback('\r\nlogin as: ')
|
||||
return
|
||||
except paramiko.ssh_exception.NoValidConnectionsError as e:
|
||||
self.ssh.close()
|
||||
self.datacallback(str(e))
|
||||
self.inputmode = 0
|
||||
self.username = b''
|
||||
@@ -139,6 +143,7 @@ class SshShell(conapi.Console):
|
||||
self.datacallback('\r\nlogin as: ')
|
||||
return
|
||||
except cexc.PubkeyInvalid as pi:
|
||||
self.ssh.close()
|
||||
self.keyaction = ''
|
||||
self.candidatefprint = pi.fingerprint
|
||||
self.datacallback(pi.message)
|
||||
@@ -148,6 +153,7 @@ class SshShell(conapi.Console):
|
||||
self.datacallback('\r\nEnter "disconnect" or "accept": ')
|
||||
return
|
||||
except paramiko.SSHException as pi:
|
||||
self.ssh.close()
|
||||
self.inputmode = -2
|
||||
warn = str(pi)
|
||||
if warnhostkey:
|
||||
@@ -169,11 +175,11 @@ class SshShell(conapi.Console):
|
||||
self.datacallback(conapi.ConsoleEvent.Disconnect)
|
||||
return
|
||||
elif self.inputmode == -1:
|
||||
while len(data) and data[0] == b'\x7f' and len(self.keyaction):
|
||||
while len(data) and data[0:1] == b'\x7f' and len(self.keyaction):
|
||||
self.datacallback('\b \b') # erase previously echoed value
|
||||
self.keyaction = self.keyaction[:-1]
|
||||
data = data[1:]
|
||||
while len(data) and data[0] == b'\x7f':
|
||||
while len(data) and data[0:1] == b'\x7f':
|
||||
data = data[1:]
|
||||
while b'\x7f' in data:
|
||||
delidx = data.index(b'\x7f')
|
||||
@@ -195,11 +201,11 @@ class SshShell(conapi.Console):
|
||||
elif len(data) > 0:
|
||||
self.datacallback(data)
|
||||
elif self.inputmode == 0:
|
||||
while len(data) and data[0] == b'\x7f' and len(self.username):
|
||||
while len(data) and data[0:1] == b'\x7f' and len(self.username):
|
||||
self.datacallback('\b \b') # erase previously echoed value
|
||||
self.username = self.username[:-1]
|
||||
data = data[1:]
|
||||
while len(data) and data[0] == b'\x7f':
|
||||
while len(data) and data[0:1] == b'\x7f':
|
||||
data = data[1:]
|
||||
while b'\x7f' in data:
|
||||
delidx = data.index(b'\x7f')
|
||||
@@ -216,7 +222,7 @@ class SshShell(conapi.Console):
|
||||
# echo back typed data
|
||||
self.datacallback(data)
|
||||
elif self.inputmode == 1:
|
||||
while len(data) > 0 and data[0] == b'\x7f':
|
||||
while len(data) > 0 and data[0:1] == b'\x7f':
|
||||
self.password = self.password[:-1]
|
||||
data = data[1:]
|
||||
while b'\x7f' in data:
|
||||
@@ -237,4 +243,4 @@ class SshShell(conapi.Console):
|
||||
|
||||
def create(nodes, element, configmanager, inputdata):
|
||||
if len(nodes) == 1:
|
||||
return SshShell(nodes[0], configmanager)
|
||||
return SshShell(nodes[0], configmanager)
|
||||
|
||||
@@ -111,6 +111,8 @@ class ShellSession(consoleserver.ConsoleSession):
|
||||
|
||||
def destroy(self):
|
||||
try:
|
||||
activesessions[(self.configmanager.tenant, self.node,
|
||||
self.username)][self.sessionid].close()
|
||||
del activesessions[(self.configmanager.tenant, self.node,
|
||||
self.username)][self.sessionid]
|
||||
except KeyError:
|
||||
|
||||
@@ -123,8 +123,8 @@ def sessionhdl(connection, authname, skipauth=False, cert=None):
|
||||
if authdata:
|
||||
cfm = authdata[1]
|
||||
authenticated = True
|
||||
# version 0 == original, version 1 == pickle3 allowed
|
||||
send_data(connection, "Confluent -- v{0} --".format(sys.version_info[0] - 2))
|
||||
# version 0 == original, version 1 == pickle3 allowed, 2 = pickle forbidden, msgpack allowed
|
||||
send_data(connection, "Confluent -- v2 --")
|
||||
while not authenticated: # prompt for name and passphrase
|
||||
send_data(connection, {'authpassed': 0})
|
||||
response = tlvdata.recv(connection)
|
||||
@@ -154,20 +154,25 @@ def sessionhdl(connection, authname, skipauth=False, cert=None):
|
||||
cfm = authdata[1]
|
||||
send_data(connection, {'authpassed': 1})
|
||||
request = tlvdata.recv(connection)
|
||||
if request and 'collective' in request and skipauth:
|
||||
if not libssl:
|
||||
if request and 'collective' in request:
|
||||
if skipauth:
|
||||
if not libssl:
|
||||
tlvdata.send(
|
||||
connection,
|
||||
{'collective': {'error': 'Server either does not have '
|
||||
'python-pyopenssl installed or has an '
|
||||
'incorrect version installed '
|
||||
'(e.g. pyOpenSSL would need to be '
|
||||
'replaced with python-pyopenssl). '
|
||||
'Restart confluent after updating '
|
||||
'the dependency.'}})
|
||||
return
|
||||
return collective.handle_connection(connection, None, request['collective'],
|
||||
local=True)
|
||||
else:
|
||||
tlvdata.send(
|
||||
connection,
|
||||
{'collective': {'error': 'Server either does not have '
|
||||
'python-pyopenssl installed or has an '
|
||||
'incorrect version installed '
|
||||
'(e.g. pyOpenSSL would need to be '
|
||||
'replaced with python-pyopenssl). '
|
||||
'Restart confluent after updating '
|
||||
'the dependency.'}})
|
||||
return
|
||||
return collective.handle_connection(connection, None, request['collective'],
|
||||
local=True)
|
||||
connection,
|
||||
{'collective': {'error': 'collective management commands may only be used by root'}})
|
||||
while request is not None:
|
||||
try:
|
||||
process_request(
|
||||
@@ -467,10 +472,14 @@ class SockApi(object):
|
||||
|
||||
def watch_for_cert(self):
|
||||
libc = ctypes.CDLL(ctypes.util.find_library('c'))
|
||||
watcher = libc.inotify_init()
|
||||
if libc.inotify_add_watch(watcher, '/etc/confluent/', 0x100) > -1:
|
||||
watcher = libc.inotify_init1(os.O_NONBLOCK)
|
||||
if libc.inotify_add_watch(watcher, b'/etc/confluent/', 0x100) > -1:
|
||||
while True:
|
||||
select.select((watcher,), (), (), 86400)
|
||||
try:
|
||||
os.read(watcher, 1024)
|
||||
except Exception:
|
||||
pass
|
||||
if self.should_run_remoteapi():
|
||||
os.close(watcher)
|
||||
self.start_remoteapi()
|
||||
|
||||
@@ -13,9 +13,9 @@ BuildRoot: %{_tmppath}/%{name}-%{version}-%{release}-buildroot
|
||||
Prefix: %{_prefix}
|
||||
BuildArch: noarch
|
||||
%if "%{dist}" == ".el8"
|
||||
Requires: python3-pyghmi >= 1.0.34, python3-eventlet, python3-greenlet, python3-pycryptodomex >= 3.4.7, confluent_client, python3-pyparsing, python3-paramiko, python3-dns, python3-netifaces, python3-pyasn1 >= 0.2.3, python3-pysnmp >= 4.3.4, python3-pyte, python3-lxml, python3-eficompressor, python3-setuptools, python3-dateutil, python3-enum34, python3-asn1crypto, python3-cffi, python3-pyOpenSSL, python3-monotonic, python3-websocket-client
|
||||
Requires: python3-pyghmi >= 1.0.34, python3-eventlet, python3-greenlet, python3-pycryptodomex >= 3.4.7, confluent_client, python3-pyparsing, python3-paramiko, python3-dns, python3-netifaces, python3-pyasn1 >= 0.2.3, python3-pysnmp >= 4.3.4, python3-pyte, python3-lxml, python3-eficompressor, python3-setuptools, python3-dateutil, python3-enum34, python3-asn1crypto, python3-cffi, python3-pyOpenSSL, python3-monotonic, python3-websocket-client python3-msgpack
|
||||
%else
|
||||
Requires: python-pyghmi >= 1.0.34, python-eventlet, python-greenlet, python-pycryptodomex >= 3.4.7, confluent_client, python-pyparsing, python-paramiko, python-dns, python-netifaces, python2-pyasn1 >= 0.2.3, python-pysnmp >= 4.3.4, python-pyte, python-lxml, python-eficompressor, python-setuptools, python-dateutil, python2-websocket-client
|
||||
Requires: python-pyghmi >= 1.0.34, python-eventlet, python-greenlet, python-pycryptodomex >= 3.4.7, confluent_client, python-pyparsing, python-paramiko, python-dns, python-netifaces, python2-pyasn1 >= 0.2.3, python-pysnmp >= 4.3.4, python-pyte, python-lxml, python-eficompressor, python-setuptools, python-dateutil, python2-websocket-client python2-msgpack
|
||||
%endif
|
||||
Vendor: Jarrod Johnson <jjohnson2@lenovo.com>
|
||||
Url: http://xcat.sf.net/
|
||||
@@ -46,15 +46,37 @@ grep -v confluent/__init__.py INSTALLED_FILES.bare | grep -v etc/init.d/confluen
|
||||
rm $RPM_BUILD_ROOT/etc/init.d/confluent
|
||||
rmdir $RPM_BUILD_ROOT/etc/init.d
|
||||
rmdir $RPM_BUILD_ROOT/etc
|
||||
# Only do non-root confluent if systemd of the platform supports it
|
||||
systemd-analyze verify $RPM_BUILD_ROOT/usr/lib/systemd/system/confluent.service 2>&1 | grep "'AmbientCapabilities'" > /dev/null && sed -e 's/User=.*//' -e 's/Group=.*//' -e 's/AmbientCapabilities=.*//' -i $RPM_BUILD_ROOT/usr/lib/systemd/system/confluent.service
|
||||
cat INSTALLED_FILES
|
||||
|
||||
%triggerin -- python-pyghmi
|
||||
%triggerin -- python-pyghmi, python3-pyghmi, python2-pyghmi
|
||||
if [ -x /usr/bin/systemctl ]; then /usr/bin/systemctl try-restart confluent >& /dev/null; fi
|
||||
true
|
||||
|
||||
%pre
|
||||
getent group confluent > /dev/null || /usr/sbin/groupadd -r confluent
|
||||
getent passwd confluent > /dev/null || /usr/sbin/useradd -r -g confluent -d /var/lib/confluent -s /sbin/nologin confluent
|
||||
mkdir -p /etc/confluent /var/lib/confluent /var/log/confluent /var/cache/confluent
|
||||
chown -R confluent:confluent /etc/confluent /var/lib/confluent /var/log/confluent /var/cache/confluent
|
||||
|
||||
%post
|
||||
sysctl -p /usr/lib/sysctl.d/confluent.conf >& /dev/null
|
||||
if [ -x /usr/bin/systemctl ]; then /usr/bin/systemctl try-restart confluent >& /dev/null; fi
|
||||
NEEDCHOWN=0
|
||||
NEEDSTART=0
|
||||
find /etc/confluent -uid 0 | egrep '.*' > /dev/null && NEEDCHOWN=1
|
||||
find /var/log/confluent -uid 0 | egrep '.*' > /dev/null && NEEDCHOWN=1
|
||||
find /var/run/confluent -uid 0 | egrep '.*' > /dev/null && NEEDCHOWN=1
|
||||
find /var/cache/confluent -uid 0 | egrep '.*' > /dev/null && NEEDCHOWN=1
|
||||
if [ $NEEDCHOWN = 1 ]; then
|
||||
if systemctl is-active confluent > /dev/null; then
|
||||
NEEDSTART=1
|
||||
systemctl stop confluent
|
||||
fi
|
||||
chown -R confluent:confluent /etc/confluent /var/lib/confluent /var/log/confluent /var/cache/confluent
|
||||
fi
|
||||
systemctl daemon-reload
|
||||
if systemctl is-active confluent > /dev/null || [ $NEEDSTART = 1 ]; then /usr/bin/systemctl restart confluent >& /dev/null; fi
|
||||
if [ ! -e /etc/pam.d/confluent ]; then
|
||||
ln -s /etc/pam.d/sshd /etc/pam.d/confluent
|
||||
fi
|
||||
|
||||
@@ -1,13 +1,25 @@
|
||||
# IBM(c) 2015 Apache 2.0
|
||||
# Lenovo(c) 2020 Apache 2.0
|
||||
[Unit]
|
||||
Description=Confluent hardware manager
|
||||
Description=Confluent hardware manager
|
||||
|
||||
[Service]
|
||||
Type=forking
|
||||
#PIDFile=/var/run/confluent/pid
|
||||
RuntimeDirectory=confluent
|
||||
StateDirectory=confluent
|
||||
CacheDirectory=confluent
|
||||
LogsDirectory=confluent
|
||||
ConfigurationDirectory=confluent
|
||||
ExecStart=/opt/confluent/bin/confluent
|
||||
ExecStop=/opt/confluent/bin/confetty shutdown /
|
||||
Restart=on-failure
|
||||
AmbientCapabilities=CAP_NET_BIND_SERVICE CAP_SETUID CAP_SETGID CAP_CHOWN
|
||||
User=confluent
|
||||
Group=confluent
|
||||
DevicePolicy=closed
|
||||
ProtectControlGroups=true
|
||||
ProtectSystem=true
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
|
||||
Reference in New Issue
Block a user