2
0
mirror of https://github.com/xcat2/confluent.git synced 2026-09-29 08:41:00 +00:00

Compare commits

..

84 Commits

Author SHA1 Message Date
Jarrod Johnson 8b89232922 Do not get collective member when collective doesn't exist 2023-03-06 16:59:07 -05:00
Jarrod Johnson 22c464e092 Only add self to collective if self not yet in collective
Previously, it was safe to just do all the time, but now it may lose
the role.
2023-03-06 16:49:03 -05:00
Jarrod Johnson 4d9b11bc55 Fix quorum when there is no collective yet 2023-03-06 16:38:09 -05:00
Jarrod Johnson baa365fcac Implement non-voting collective members
Provide for applications
where only a small subset of collective
members should be
considered to count
toward whether the collective
can proceed.

Commonly, 'service' nodes may
be numerous to do work, but may all want to go offline
during a maintenance window.
2023-03-06 11:56:15 -05:00
Jarrod Johnson a385b1e93d Try strategy to have confignet run
confignet is special, it is designed
to work when networking
isn't right.  So have it run during firstboot
in case post fouled up
the network for firstboot.
2023-02-28 12:12:36 -05:00
Jarrod Johnson 733b6853dd Up newly added interfaces as a matter of course 2023-02-28 12:04:20 -05:00
Jarrod Johnson b4182cd4b5 Fix formation of error message
Use format to take in the parameters regardless of type
2023-02-27 14:55:01 -05:00
Jarrod Johnson 9f7e53701e Avoid latching onto USB nic in a vswitch as 'the nic'
In esxi, some builds may have USB nic brought up in a vswitch.

Detect and avoid that scenario.
2023-02-27 10:43:40 -05:00
Jarrod Johnson 70d8a1059c Consistently treat bytes as bytes in ssh
In Python3 systems,
there would be confusion
about bytes versus str.

Fix this so that ssh can work more consistently.
2023-02-24 15:47:20 -05:00
Jarrod Johnson 59b07665ab Modify float formatting again
Make sure at least one decimal is in a float.

Maximum precision of 5 past.
2023-02-24 12:03:43 -05:00
Jarrod Johnson 5ea214a726 Use eventlet subprocess
sshutil uses eventlet subprocess,
making calledprocesserror
hard to catch.

Adjust to consistently use same
subprocesss module.
2023-02-22 16:34:13 -05:00
Jarrod Johnson b99034f539 Improve reliability of collective join
While servicing an enrollment,
there's a window for a collective
member to be 'defined' but not
yet active, meaning quorum may transiently be lost as multiple enrollments progress.

Serialize enrollments by holding the enrollment process open.

Also, there is a chance that a transient transfer error may occur during loading
of the DB.  In such a case, restart
the connection rather thn aborting.
2023-02-22 16:11:38 -05:00
Jarrod Johnson 6df2e822a5 Correct api call in discovery 2023-02-22 09:34:32 -05:00
Jarrod Johnson 2379f6f90f Change nodesensors format of float
Floats are either unnecessarily long
in normal output, or too unconstrained in CSV output.

Normalize to as many digits as 'makes sense' up to 5 digits.

5 miight seem a bit much, but one common metric is kWh, which may need
that precision over short intervals.
2023-02-22 08:41:46 -05:00
Jarrod Johnson 77ba0acee6 Merge pull request #122 from Tkucherera/nodeconsole-kill
nodeconsole <noderange> kill: added functionality for closing open win…
2023-02-16 16:37:07 -05:00
Tinashe b2c773bb84 nodeconsole <noderange> kill:added functionality for closing open windowed consoles 2023-02-16 15:54:21 -05:00
Jarrod Johnson 241800b1c9 Restore filename-only import
The open file handle as implemented
could not pass to the subprocess.

Rather than figure out how to open
and pass the filehandle,
simply let the subprocess
independently open the file
if it isn't passed.
2023-02-16 09:13:05 -05:00
Jarrod Johnson abc639e32b Preferentially support HTTPS on Eaton PDU
While Eaton does not do HTTPS by default,
it can be configured to do so.

Support when available.

Mitigate downgrade attack by
stickying the cert fingerprint.
If fingerprint is present, then refuse
to even think about port 80.
2023-02-15 17:03:35 -05:00
Jarrod Johnson 90af99e864 Add more clear error on syncfile mistake
If a bad node was included in
a syncfile, the error was highly misleading.

Provide a more clear indicaiton of the problem on failure.
2023-02-14 14:53:40 -05:00
Jarrod Johnson 09ce824c85 Fix bad lookup attempts on slashed addr
While this should in theory be
harmless, it exacerbates some
DNS setups that would look
up the normal result quickly,
but would stall on
a bad lookup.
2023-02-14 14:53:40 -05:00
Jarrod Johnson 9c1e7a7142 Allow interfaces to supersede default
In some scenarios, the 'default'
interface is overlapped by another connection, either
identical or as a superset in a bond.

Whittle down the default
interface if superseded
to mitigate duplicate interface setup.
2023-02-14 14:53:40 -05:00
Jarrod Johnson 36195198a6 Add fallback for newer msgpack
Newer msgpack refuses the encoding argument, use raw=False instead.

Further, newer msgpack refuses to accept int as key by default.
Opt into it as the risk is hash collision due to msgpack int being used directly, and
we aren't dealing with untrusted
peer (we only talk to ourselves).
2023-02-14 14:53:40 -05:00
Jarrod Johnson 3798a33213 Merge pull request #121 from Tkucherera/nodeconsole
nodeconsole documentation: passthrough options
2023-02-14 08:45:43 -05:00
Tinashe 251b307bd7 nodeconsole documentation: passthrough options 2023-02-14 08:27:25 -05:00
Jarrod Johnson bb7a72db65 Fix for ipv6 deployment
Need to avoid double-bracketing of the server and also disable globbing
so curl does not mistake the ip address for a glob attempt.
2023-02-13 09:36:42 -05:00
Jarrod Johnson fcde113e08 Add a check of dns.domain to selfcheck for node 2023-02-08 14:45:16 -05:00
Jarrod Johnson a02f617b3d Add DDR5 dimm to nodeinventory CLI output 2023-02-07 14:01:18 -05:00
Jarrod Johnson 7f1ac92fc9 Store mgr from confluent= specificate 2023-02-01 16:51:21 -05:00
Jarrod Johnson 8cf97833ab Fixes for certificate directed discovery 2023-02-01 13:09:40 -05:00
Jarrod Johnson 3e747069d9 Try to get verified bay from SMMs
With V3 systems, we can now ask
the SMMs for the certificates
and use that for a verified
measurement, regardless of
whether the XCC is returning
the correct bay number.
2023-02-01 12:57:27 -05:00
Jarrod Johnson c687da4d5f Tweak architecture override on import 2023-01-31 15:57:41 -05:00
Jarrod Johnson 340ccc422c Specify check for arch override of addons.cpio
For now, keep using x86_64 as
default, but allow overrides
for other architectures.

One day it may be cleaner to move all addons.cpio to
arch specific subdirs.
2023-01-31 15:27:45 -05:00
Jarrod Johnson 2c3afac576 Restructure aarch64 addons
Avoid tripping over current copy over, prepare
for smarter selection by architecture.
2023-01-31 15:10:49 -05:00
Jarrod Johnson 8e1cc63ac0 Correct spelling of keyword argument in ipmi 2023-01-31 15:00:22 -05:00
Jarrod Johnson dc6c7c1acc Make sure both el8 and el9 binaries are packed 2023-01-31 13:29:24 -05:00
Jarrod Johnson 5c309db47c Further ARMv8 support
Handle aarch64 differences in
at least some distributions.
2023-01-31 11:20:40 -05:00
Jarrod Johnson 976e9ef563 Bump version on genesis 2023-01-31 09:10:23 -05:00
Jarrod Johnson 0efd2a4d74 Fix the amended license gathering 2023-01-31 08:58:56 -05:00
Jarrod Johnson 424830471d Note how to fetch srpms associated with genesis 2023-01-31 08:54:03 -05:00
Jarrod Johnson 23f33a8420 Revamp license gathering for genesis 2023-01-31 08:52:32 -05:00
Jarrod Johnson 521a58c1d9 Have a utility to generate NOTICE from tmux
tmux basically defers to the c files, so
generate a NOTICE file from c files.
2023-01-31 08:37:42 -05:00
Jarrod Johnson ce375a1162 Merge pull request #120 from Tkucherera/master
Adding missing imports
2023-01-30 14:48:05 -05:00
Tkucherera caee136012 Merge branch 'xcat2:master' into master 2023-01-30 14:15:53 -05:00
Tinashe 2e283f3442 nodeconsole: missing imports time and socket 2023-01-30 14:13:44 -05:00
Jarrod Johnson 2b01d9fbfa Properly store all candidate host ip addresses
This is needed to ensure that mis-detected primary ip
falls through to another viable ip
2023-01-30 12:40:40 -05:00
Jarrod Johnson 627bc9ffe3 Modify pkglist for aarch64 2023-01-27 12:14:37 -05:00
Jarrod Johnson 3e71e103b1 Fix unpacking of el8 and el9 built sources 2023-01-27 10:47:27 -05:00
Jarrod Johnson a90cd8515e Tweak osdeploy for ARM setup 2023-01-27 10:43:29 -05:00
Jarrod Johnson 284d042afe Merge pull request #119 from Tkucherera/master
nodeconsole windowed and tiled functionality
2023-01-27 09:33:45 -05:00
Tinashe 2cc134adeb nodeconsole: allow for passthrough args 2023-01-27 09:30:46 -05:00
Jarrod Johnson 02e242ec4e Restore link local cert in apiclient 2023-01-27 09:13:47 -05:00
Jarrod Johnson 1777223232 Fixes for osdeploy arm ipxe init 2023-01-27 08:40:31 -05:00
Jarrod Johnson 648290ffbc Begin implementing aarch64 deploy support 2023-01-27 08:00:38 -05:00
Tinashe c9b72225a9 nodeeventlog timeframe documentation 2023-01-26 10:37:30 -05:00
Tinashe 3433635e9b nodeeventlog: timeframe option 2023-01-26 10:16:33 -05:00
Tkucherera 7a6c4fc5b2 Merge branch 'xcat2:master' into master 2023-01-26 10:05:04 -05:00
Tinashe f176a836ae nodeecentlog: add timeframe option 2023-01-26 10:01:02 -05:00
Jarrod Johnson ce324e90f7 Draft spec to generate addons-aarch64 files 2023-01-25 12:54:03 -05:00
Tinashe 23ea53ab55 console geometry 100x31 2023-01-24 11:18:21 -05:00
Jarrod Johnson d14d28caf8 Confirm TLS connectivity when scanning hosts
In certain environments, Confluent may have an IP address that
is fake, but then there is elsewhere with that same IP for real.

To mitigate this, follow up basic connectivity with proof of having
an associated certificate.
2023-01-24 08:22:00 -05:00
Jarrod Johnson 0008998680 Add api method to request all mac data
This will provide easy way for
client to get FDB data, potentially
for use in conjunction with discovery data.

For now, leave LLDP out, as that isn't currently cached
at the confluent layer.
2023-01-23 13:37:29 -05:00
Jarrod Johnson 2e059b5887 Make an API for getting full discovery data in one fetch
This makes for faster nodediscover being possible, also
makes web management of the data easier
2023-01-23 11:47:33 -05:00
Jarrod Johnson 792e6472e4 Fix IPv6 addresses_match
fe80:: could be submitted during
collective startup, handle that problem appropriately.
2023-01-23 11:24:25 -05:00
Tinashe b965f9b758 nodeconsole windowed and tiled functionality 2023-01-20 16:41:56 -05:00
Jarrod Johnson a522e17a63 Merge pull request #118 from Tkucherera/master
nodeeventlog man page -l option
2023-01-20 14:49:11 -05:00
Tinashe 9f3b934ea4 nodeeventlog man page -l option 2023-01-20 14:24:32 -05:00
Jarrod Johnson 680ca2c4a2 Merge pull request #117 from Tkucherera/master
nodeeventlog: return last n entries
2023-01-20 12:51:16 -05:00
Tinashe c3d0d255d3 nodeeventlog: -l return last n lines for each node 2023-01-20 12:13:15 -05:00
Tinashe 46d0a8d222 nodeeventlog: return last n entries 2023-01-20 10:09:52 -05:00
Jarrod Johnson 75f020f53c Have apiarmed continuous be properly respected for shared secret
Remote media was erroneously being invalidated, despite user opting
out of the strict security.
2023-01-19 14:54:18 -05:00
Jarrod Johnson c09e8448c2 Change to POSIX compliant range
POSIX allows ., but does not allow +.  This was a problem with EL 8.4 libxcrypt,
though is not a problem otherwise.
2023-01-19 14:53:35 -05:00
Jarrod Johnson 01f939b871 Have SuSE path also not be bothered by inability to restart web service 2023-01-18 08:50:30 -05:00
Jarrod Johnson 1f23750356 Add affluent detection to confluent
Affluent agent will now have an SSDP
response.  Add support for at
least recognizing and presenting
this in the discovery data.
2023-01-17 15:11:12 -05:00
Jarrod Johnson d1265af828 Handle more errors
subprocess may throw other errors that aren't calledprocesserrors,
in newer python versions.  Handle the case more broadly.
2023-01-17 10:04:10 -05:00
Jarrod Johnson 0929f059e2 Increase track size in dir2img
Larger images still run afoul of track limits
in mtools.  Make tracks 4 times as big
to lower number of required tracks.
2023-01-17 10:02:51 -05:00
Jarrod Johnson 40c3f2da53 Actually display deployment state, when available 2023-01-13 13:02:17 -05:00
Jarrod Johnson 51e53405d8 Add attributes for profiles to report state
Profiles may want to report things
like success and error
2023-01-13 12:54:21 -05:00
Jarrod Johnson d644d34b60 Add '-s' to nodeinventory
This allows a quick command to
get attributes into confluent
for manually added nodes,
without having to go through 'discovery' process.
2023-01-13 12:29:36 -05:00
Jarrod Johnson 7f31ae5b57 Fix syntax error 2023-01-13 11:15:51 -05:00
Jarrod Johnson a09e1a3f8b Handle IPv6 not set on IPMI nodes 2023-01-13 11:07:13 -05:00
Jarrod Johnson 50c073670d Explicitly declare Textmode during autoconsole
This enables a workable console during text install,
while also allowing graphical to run
2023-01-13 10:54:29 -05:00
Jarrod Johnson bc452b9b9a Restore role-less group
If a group is missing a role,
coerce it to administrator
2023-01-13 10:01:52 -05:00
Jarrod Johnson 453d1f9ceb Add IPv6 configuration support
For redfish and IPMI devices,
support new IPv6 static configuration
controls
2023-01-13 10:01:28 -05:00
Jarrod Johnson feed125c86 Fix restoration of old confluent db
Old confluent DB may have None in role. This is no longer
allowed.  Restore such entries by coercing them to 'Administrator'
which is how old confluent treated such users.
2023-01-12 08:38:55 -05:00
59 changed files with 1118 additions and 215 deletions
+4 -4
View File
@@ -21,17 +21,17 @@ def create_image(directory, image, label=None):
currsz = (currsz // 512 +1) * 512
datasz += currsz
datasz += ents * 32768
datasz = datasz // 4096 + 1
datasz = datasz // 16384 + 1
with open(image, 'wb') as imgfile:
imgfile.seek(datasz * 4096 - 1)
imgfile.seek(datasz * 16384 - 1)
imgfile.write(b'\x00')
if label:
subprocess.check_call(['mformat', '-i', image, '-v', label,
'-r', '16', '-d', '1', '-t', str(datasz),
'-s', '4','-h', '2', '::'])
'-s', '16','-h', '2', '::'])
else:
subprocess.check_call(['mformat', '-i', image, '-r', '16', '-d', '1', '-t',
str(datasz), '-s', '4','-h', '2', '::'])
str(datasz), '-s', '16','-h', '2', '::'])
# Some clustered filesystems will have the lock from mformat
# linger after close (mformat doesn't unlock)
# do a blocking wait for shared lock and then explicitly
+6
View File
@@ -90,6 +90,12 @@ cfgpaths = {
'bmc.ipv4_gateway': (
'configuration/management_controller/net_interfaces/management',
'ipv4_gateway'),
'bmc.static_ipv6_addresses': (
'configuration/management_controller/net_interfaces/management',
'static_v6_addresses'),
'bmc.static_ipv6_gateway': (
'configuration/management_controller/net_interfaces/management',
'static_v6_gateway'),
'bmc.hostname': (
'configuration/management_controller/hostname', 'hostname'),
}
+176 -6
View File
@@ -26,10 +26,13 @@ if path.startswith('/opt'):
import confluent.client as client
import confluent.sortutil as sortutil
import confluent.logreader as logreader
import time
import socket
import re
confettypath = os.path.join(os.path.dirname(sys.argv[0]), 'confetty')
argparser = optparse.OptionParser(
usage="Usage: %prog [options] node",
usage="Usage: %prog [options] <noderange> [kill][-- [passthroughoptions]]",
epilog="Command sequences are available while connected to a console, hit "
"ctrl-'e', then release ctrl, then 'c', then '?' for a full list. "
"For example, ctrl-'e', then 'c', then '.' will exit the current "
@@ -58,10 +61,26 @@ argparser.add_option('-w','--windowed', action='store_true', default=False,
'--shell-type login". If the NODECONSOLE_WINDOWED_COMMAND '
'environment variable isn\'t set, xterm will be used by'
'default.')
(options, args) = argparser.parse_args()
pass_through_args = []
killcon = False
try:
noderange = args[0]
if len(args) > 1:
if args[1] == 'kill':
killcon = True
pass_through_args = args[1:]
args = args[:1]
except IndexError:
argparser.print_help()
sys.exit(1)
if len(args) != 1:
argparser.print_help()
sys.exit(1)
if options.log:
logname = args[0]
if not os.path.exists(logname) and logname[0] != '/':
@@ -71,14 +90,87 @@ if options.log:
sys.exit(1)
logreader.replay_to_console(logname)
sys.exit(0)
#added functionality for wcons
if options.windowed:
def kill(noderange):
sess = client.Command()
envstring=os.environ.get('NODECONSOLE_WINDOWED_COMMAND')
if not envstring:
envlist=["xterm", "-e"]
envstring = 'xterm'
nodes = []
for res in sess.read('/noderange/{0}/nodes/'.format(args[0])):
node = res.get('item', {}).get('href', '/').replace('/', '')
if not node:
sys.stderr.write(res.get('error', repr(res)) + '\n')
sys.exit(1)
nodes.append(node)
for node in nodes:
ps_data=subprocess.Popen(['ps', '-auxww' ], stdout=subprocess.PIPE)
wintr=ps_data.communicate()[0]
for line in wintr.decode('utf-8').split('\n'):
if confettypath in line and envstring in line and node in line:
pid_line = [x for x in line.split(' ') if x != '']
ps_data=subprocess.Popen(['kill', '-9', pid_line[1] ], stdout=subprocess.PIPE)
sys.exit(0)
def handle_geometry(envlist, sizegeometry, side_pad=0, top_pad=0, first=False):
if '-geometry' in envlist:
g_index = envlist.index('-geometry')
elif '-g' in envlist:
g_index = envlist.index('-g')
else:
g_index = 0
if g_index:
if first:
envlist[g_index+1] = '{0}+{1}+{2}'.format(envlist[g_index+1],side_pad, top_pad)
else:
envlist[g_index+1] = '{0}+{1}+{2}'.format(sizegeometry,side_pad, top_pad)
else:
envlist.insert(1, '-geometry')
envlist.insert(2, '{0}+{1}+{2}'.format(sizegeometry,side_pad, top_pad))
g_index = 1
return envlist
# add funcltionality to close/kill all open consoles
if killcon:
kill(noderange)
#added functionality for wcons
if options.windowed:
result=subprocess.Popen(['xwininfo', '-root'], stdout=subprocess.PIPE)
rootinfo=result.communicate()[0]
result.wait()
for line in rootinfo.decode('utf-8').split('\n'):
if 'Width' in line:
screenwidth = int(line.split(':')[1])
if 'Height' in line:
screenheight = int(line.split(':')[1])
envstring=os.environ.get('NODECONSOLE_WINDOWED_COMMAND')
if not envstring:
sizegeometry='100x31'
corrected_x, corrected_y = (13,84)
envlist = handle_geometry(['xterm'] + pass_through_args + ['-e'],sizegeometry, first=True)
#envlist=['xterm', '-bg', 'black', '-fg', 'white', '-geometry', '{sizegeometry}+0+0'.format(sizegeometry=sizegeometry), '-e']
else:
envlist=os.environ.get('NODECONSOLE_WINDOWED_COMMAND').split(' ')
if envlist[0] == 'xterm':
if '-geometry' in envlist:
g_index = envlist.index('-geometry')
elif '-g' in envlist:
g_index = envlist.index('-g')
else:
g_index = 0
if g_index:
envlist[g_index+1] = envlist[g_index+1] + '+0+0'
else:
envlist.insert(1, '-geometry')
envlist.insert(2, '100x31+0+0')
g_index = 1
nodes = []
sess = client.Command()
for res in sess.read('/noderange/{0}/nodes/'.format(args[0])):
@@ -87,9 +179,87 @@ if options.windowed:
sys.stderr.write(res.get('error', repr(res)) + '\n')
sys.exit(1)
nodes.append(node)
if options.tile and not envlist[0] == 'xterm':
sys.stderr.write('[ERROR] UNSUPPORTED OPTIONS. \nWindowed and tiled consoles are only supported when using xterm \n')
sys.exit(1)
firstnode=nodes[0]
nodes.pop(0)
with open(os.devnull, 'wb') as devnull:
xopen=subprocess.Popen(envlist + [confettypath, '-c', '/tmp/controlpath-{0}'.format(firstnode), '-m', '5', 'start', '/nodes/{0}/console/session'.format(firstnode) ] , stdin=devnull)
time.sleep(2)
s=socket.socket(socket.AF_UNIX)
winid=''
try:
s.connect('/tmp/controlpath-{firstnode}'.format(firstnode=firstnode))
s.recv(64)
s.send(b'GETWINID')
winid=s.recv(64).decode('utf-8')
except:
time.sleep(2)
# try to get id of first panel/xterm window using name
win=subprocess.Popen(['xwininfo', '-tree', '-root'], stdout=subprocess.PIPE)
wintr=win.communicate()[0]
for line in wintr.decode('utf-8').split('\n'):
if 'console: {firstnode}'.format(firstnode=firstnode) in line or 'confetty' in line:
win_obj = [ele for ele in line.split(' ') if ele.strip()]
winid = win_obj[0]
if winid:
firstnode_window=subprocess.Popen(['xwininfo', '-id', '{winid}'.format(winid=winid)], stdout=subprocess.PIPE)
xinfo=firstnode_window.communicate()[0]
xinfl = xinfo.decode('utf-8').split('\n')
for line in xinfl:
if 'Absolute upper-left X:' in line:
side_pad = int(line.split(':')[1])
elif 'Absolute upper-left Y:' in line:
top_pad = int(line.split(':')[1])
elif 'Width:' in line:
window_width = int(line.split(':')[1])
elif 'Height' in line:
window_height = int(line.split(':')[1])
elif '-geometry' in line:
l = re.split(' |x|\+', line)
l_nosp = [ele for ele in l if ele.strip()]
wmxo = int(l_nosp[1])
wmyo = int(l_nosp[2])
sizegeometry = str(wmxo) + 'x' + str(wmyo)
else:
pass
window_width += side_pad*2
window_height += side_pad+top_pad
screenwidth -= wmxo
screenheight -= wmyo
currx = window_width
curry = 0
maxcol = int(screenwidth/window_width)
for node in sortutil.natural_sort(nodes):
if options.tile and envlist[0] == 'xterm':
corrected_x = currx
corrected_y = curry
xgeometry = '{0}+{1}+{2}'.format(sizegeometry, corrected_x, corrected_y)
currx += window_width
if currx >= screenwidth:
currx=0
curry += window_height
if curry > screenheight:
curry =top_pad
if not envstring:
envlist= handle_geometry(envlist, sizegeometry, corrected_x, corrected_y)
else:
if g_index:
envlist[g_index+1] = xgeometry
elif envlist[0] == 'xterm':
envlist=handle_geometry(envlist, sizegeometry, side_pad, top_pad)
side_pad+=(side_pad+1)
top_pad+=(top_pad+30)
else:
pass
with open(os.devnull, 'wb') as devnull:
subprocess.Popen(envlist + [confettypath, '-m', '5', 'start', '/nodes/{0}/console/session'.format(node)], stdin=devnull)
xopen=subprocess.Popen(envlist + [confettypath, '-m', '5', 'start', '/nodes/{0}/console/session'.format(node)] , stdin=devnull)
sys.exit(0)
#end of wcons
if options.tile:
@@ -128,4 +298,4 @@ if options.tile:
os.execlp('tmux', 'tmux', 'attach', '-t', sessname)
else:
os.execl(confettypath, confettypath, 'start',
'/nodes/{0}/console/session'.format(args[0]))
'/nodes/{0}/console/session'.format(args[0]))
+14 -3
View File
@@ -49,7 +49,7 @@ def armonce(nr, cli):
def setpending(nr, profile, cli):
args = {'deployment.pendingprofile': profile}
args = {'deployment.pendingprofile': profile, 'deployment.state': '', 'deployment.state_detail': ''}
if not profile.startswith('genesis-'):
args['deployment.stagedprofile'] = ''
args['deployment.profile'] = ''
@@ -132,7 +132,7 @@ def main(args):
if node not in databynode:
databynode[node] = {}
for attr in dbn[node]:
if attr in ('deployment.pendingprofile', 'deployment.apiarmed', 'deployment.stagedprofile', 'deployment.profile'):
if attr in ('deployment.pendingprofile', 'deployment.apiarmed', 'deployment.stagedprofile', 'deployment.profile', 'deployment.state', 'deployment.state_detail'):
databynode[node][attr] = dbn[node][attr].get('value', '')
for node in sortutil.natural_sort(databynode):
profile = databynode[node].get('deployment.pendingprofile', '')
@@ -153,7 +153,18 @@ def main(args):
armed = ' (node authentication armed)'
else:
armed = ''
print('{0}: {1}{2}'.format(node, profile, armed))
stateinfo = ''
deploymentstate = databynode[node].get('deployment.state', '')
if deploymentstate:
statedetails = databynode[node].get('deployment.state_detail', '')
if statedetails:
stateinfo = '{}: {}'.format(deploymentstate, statedetails)
else:
stateinfo = deploymentstate
if stateinfo:
print('{0}: {1} ({2})'.format(node, profile, stateinfo))
else:
print('{0}: {1}{2}'.format(node, profile, armed))
sys.exit(0)
if args.network and not args.prepareonly:
return rc
+64 -5
View File
@@ -17,6 +17,7 @@
import codecs
from datetime import datetime as dt
from datetime import timedelta
import optparse
import os
import signal
@@ -41,7 +42,18 @@ argparser = optparse.OptionParser(
argparser.add_option('-m', '--maxnodes', type='int',
help='Specify a maximum number of '
'nodes to clear if clearing log, '
'prompting if over the threshold')
'prompting if over the threshold')
argparser.add_option('-l', '--lines', type='int',
help='return the last <n> entries '
'for each node in the eventlog. '
)
argparser.add_option('-t', '--timeframe', type='string',
help='return entries within a specified timeframe '
'for each node in the eventlog. This will return '
'entries from the last hours or days. '
'1h would be one hour, 4d would be four days. '
'format <num>h or <num>d'
)
(options, args) = argparser.parse_args()
try:
noderange = args[0]
@@ -89,12 +101,34 @@ def format_event(evt):
msg = ''
return ' '.join(retparts) + msg
if deletemode:
func = session.delete
session.stop_if_noderange_over(noderange, options.maxnodes)
else:
func = session.read
if options.timeframe:
try:
delta = int(options.timeframe[:-1])
except ValueError:
argparser.print_help()
sys.exit(1)
if options.timeframe[-1].lower() == 'd':
tdelta = timedelta(days=delta)
elif options.timeframe[-1].lower() == 'h':
tdelta = timedelta(hours=delta)
else:
argparser.print_help()
sys.exit(1)
timeframe = dt.now() - tdelta
event_dict = {}
nodes = []
for res in session.read('/noderange/{0}/nodes/'.format(args[0])):
node = res.get('item', {}).get('href', '/').replace('/', '')
nodes.append(node)
event_dict[node] = []
for rsp in func('/noderange/{0}/events/hardware/log'.format(noderange)):
if 'error' in rsp:
sys.stderr.write(rsp['error'] + '\n')
@@ -107,6 +141,31 @@ for rsp in func('/noderange/{0}/events/hardware/log'.format(noderange)):
sys.stderr.write('{0}: {1}\n'.format(node, thisdata['error']))
exitcode |= 1
if 'events' in thisdata:
evtdata = thisdata['events']
for evt in evtdata:
print('{0}: {1}'.format(node, format_event(evt)))
evtdata = thisdata['events']
if options.lines:
event_dict[node].extend(evtdata)
else:
for evt in evtdata:
if options.timeframe:
# check if line is in timeframe
if 'timestamp' in evt and evt['timestamp'] is not None:
display = dt.strptime(evt['timestamp'], '%Y-%m-%dT%H:%M:%S')
if display > timeframe:
print('{0}: {1}'.format(node, format_event(evt)))
else:
print('{0}: {1}'.format(node, format_event(evt)))
if options.lines:
for node in nodes:
evtdata_list = event_dict[node]
if len(evtdata_list) != 0:
if len(evtdata_list) > options.lines:
evtdata_list = evtdata_list[-abs(options.lines):]
for evt in evtdata_list:
if options.timeframe:
if 'timestamp' in evt and evt['timestamp'] is not None:
display = dt.strptime(evt['timestamp'], '%Y-%m-%dT%H:%M:%S')
if display > timeframe:
print('{0}: {1}'.format(node, format_event(evt)))
else:
print('{0}: {1}'.format(node, format_event(evt)))
+32 -4
View File
@@ -53,6 +53,8 @@ def print_mem_info(node, prefix, meminfo):
memdescfmt += '3-{1} '
elif 'DDR4' in meminfo['memory_type']:
memdescfmt += '4-{1} '
elif 'DDR5' in meminfo['memory_type']:
memdescfmt += '5-{1} '
elif 'DCPMM' in meminfo['memory_type']:
memdescfmt = '{0}GB {1} '
meminfo['module_type'] = 'DCPMM'
@@ -99,6 +101,7 @@ usedprefixes = set([])
argparser = optparse.OptionParser(
usage="Usage: %prog <noderange> [serial|model|uuid|mac]")
argparser.add_option('-j', '--json', action='store_true', help='Output JSON')
argparser.add_option('-s', '--store', action='store_true', help='Store serial, model, and uuid into id.serial, id.model, and id.uuid')
(options, args) = argparser.parse_args()
try:
noderange = args[0]
@@ -126,17 +129,36 @@ if len(args) > 1:
try:
if options.json:
databynode = {}
if options.store and len(args) <= 1:
url = '/noderange/{0}/inventory/hardware/all/system'
pushattribs = {}
session = client.Command()
for res in session.read(url.format(noderange)):
printerror(res)
if 'databynode' not in res:
continue
for node in res['databynode']:
if options.store and node not in pushattribs:
pushattribs[node] = {}
printerror(res['databynode'][node], node)
if 'inventory' not in res['databynode'][node]:
continue
for inv in res['databynode'][node]['inventory']:
prefix = inv['name']
if options.store and prefix == 'System':
currinfo = inv.get('information', {})
curruuid = currinfo.get('UUID', '')
if curruuid:
curruuid = curruuid.lower()
pushattribs[node]['id.uuid'] = curruuid
currserial = currinfo.get('Serial Number', '')
if currserial:
currserial = currserial.strip()
pushattribs[node]['id.serial'] = currserial
currmodelnum = currinfo.get('Model', '')
if currmodelnum:
currmodelnum = currmodelnum.strip()
pushattribs[node]['id.model'] = currmodelnum
idx = 2
while (node, prefix) in usedprefixes:
prefix = '{0} {1}'.format(inv['name'], idx)
@@ -148,7 +170,7 @@ try:
if node not in databynode:
databynode[node] = {}
databynode[node][prefix] = inv
else:
elif not options.store:
print('{0}: {1}: Not Present'.format(node, prefix))
continue
info = inv['information']
@@ -179,12 +201,18 @@ try:
databynode[node] = {}
databynode[node][prefix] = inv
break
print(u'{0}: {1} {2}: {3}'.format(node, prefix,
pretty(datum),
info[datum]))
elif not options.store:
print(u'{0}: {1} {2}: {3}'.format(node, prefix,
pretty(datum),
info[datum]))
if options.json:
print(json.dumps(databynode, sort_keys=True, indent=4,
separators=(',', ': ')))
if pushattribs:
for node in pushattribs:
for rsp in session.update('/nodes/{0}/attributes/current'.format(node), pushattribs[node]):
if 'error' in rsp:
sys.stderr.write(rsp['error'] + '\n')
except KeyboardInterrupt:
print('')
sys.exit(exitcode)
+9 -1
View File
@@ -37,6 +37,12 @@ class hybridcsv(csv.excel):
lineterminator = '\n'
def floatformat(num):
fm = u'{:.5f}'.format(num).rstrip('0')
if fm[-1:] == u'.':
return fm + u'0'
return fm
csv.register_dialect('hybrid', hybridcsv)
import confluent.client as client
@@ -135,7 +141,7 @@ def sensorpass(showout=True, appendtime=False):
if sensedata['value'] is None:
showval = ''
elif isinstance(sensedata['value'], float):
showval = u' {0:.5f} '.format(sensedata['value'])
showval = floatformat(sensedata['value'])
else:
showval = u' {0} '.format(sensedata['value'])
if sensedata['units'] not in (None, u''):
@@ -191,6 +197,8 @@ def format_csv(csvwriter, orderedsensors, resdata, showtime=True):
datum = ','.join([datum, healthstates])
else:
datum = healthstates
if isinstance(datum, float):
datum = floatformat(datum)
rowdata.append(datum)
except KeyError:
rowdata.append('N/A')
+15 -1
View File
@@ -2,7 +2,7 @@ nodeconsole(8) -- Open a console to a confluent node
=====================================================
## SYNOPSIS
`nodeconsole [options] <noderange>`
`nodeconsole [options] <noderange> [kill][-- [passthroughoptions]]`
## DESCRIPTION
@@ -16,6 +16,9 @@ will initiate an automatic retry interval that is randomized between 2 and 4 min
The reopen escape sequence below requests an immediate retry, as does connecting
a new session.
When a windowed console is open the `nodeconsole <noderange> kill` command will kill the
console process which will result in the console window closing.
## OPTIONS
* `-t`, `--tile`:
@@ -86,3 +89,14 @@ keystroke will be interpreted as a command. The following commands are availabl
Get a list of supported commands
* `<ent>`:
Hit enter to skip entering a command at the escape prompt.
## PASSTHROUGH OPTIONS
While opening a windowed console with xterm or any other console of choice. The
nodeconsole command gives capality to specify passthrough options targeted at
the console. All options after the -- will be parsed the console program. For
example, opening a windowed console using xterm with a black background.
`nodeconconsole -w n1 -- -bg black`
@@ -15,6 +15,15 @@ noderange.
* `-m MAXNODES`, `--maxnodes=MAXNODES`:
Specify a maximum number of nodes to clear if clearing log, prompting if
over the threshold
* `-l LINES`, `--lines=LINES`:
return the last <n> entries for each node in the eventlog.
* `-t TIMEFRAME`, `--timeframe=TIMEFRAME`:
return entries within a specified timeframe for each node's event log.
This will return entries from the last <n> hours or days. 1h would be
entries from with the last one hour.
* `-h`, `--help`:
Show help message and exit
+32
View File
@@ -0,0 +1,32 @@
VERSION=`git describe|cut -d- -f 1`
NUMCOMMITS=`git describe|cut -d- -f 2`
if [ "$NUMCOMMITS" != "$VERSION" ]; then
VERSION=$VERSION.dev$NUMCOMMITS.g`git describe|cut -d- -f 3`
fi
sed -e "s/#VERSION#/$VERSION/" confluent_osdeploy-aarch64.spec.tmpl > confluent_osdeploy-aarch64.spec
cd ..
cp ../LICENSE .
tar Jcvf confluent_osdeploy.tar.xz confluent_osdeploy
mv confluent_osdeploy.tar.xz ~/rpmbuild/SOURCES/
cd -
mkdir -p el9bin/opt/confluent/bin
mkdir -p el9bin/stateless-bin
mkdir -p el8bin/opt/confluent/bin
mkdir -p el8bin/stateless-bin
podman run --privileged --rm -v $(pwd)/utils:/buildutils -i -t el9builder make -C /buildutils
cd utils
mv confluent_imginfo copernicus clortho autocons ../el9bin/opt/confluent/bin
mv start_root urlmount ../el9bin/stateless-bin/
cd ..
podman run --privileged --rm -v $(pwd)/utils:/buildutils -i -t el8builder make -C /buildutils
cd utils
mv confluent_imginfo copernicus clortho autocons ../el8bin/opt/confluent/bin
mv start_root urlmount ../el8bin/stateless-bin/
cd ..
tar Jcvf confluent_el9bin.tar.xz el9bin/
tar Jcvf confluent_el8bin.tar.xz el8bin/
mv confluent_el8bin.tar.xz ~/rpmbuild/SOURCES/
mv confluent_el9bin.tar.xz ~/rpmbuild/SOURCES/
rm -rf el9bin
rm -rf el8bin
rpmbuild -ba confluent_osdeploy-aarch64.spec
@@ -253,6 +253,7 @@ class HTTPSClient(client.HTTPConnection, object):
self.stdheaders['CONFLUENT_NODENAME'] = node
if line.startswith('MANAGER:') and not host:
host = line.split(' ')[1]
self.hosts.append(host)
if not plainhost:
plainhost = host
if line.startswith('EXTMGRINFO:'):
@@ -304,6 +305,10 @@ class HTTPSClient(client.HTTPConnection, object):
def check_connections(self):
foundsrv = None
hosts = self.hosts
ctx = ssl.SSLContext(ssl.PROTOCOL_SSLv23)
ctx.load_verify_locations('/etc/confluent/ca.pem')
ctx.verify_mode = ssl.CERT_REQUIRED
ctx.check_hostname = True
for timeo in (0.1, 5):
for host in hosts:
try:
@@ -311,11 +316,17 @@ class HTTPSClient(client.HTTPConnection, object):
psock = socket.socket(addrinf[0])
psock.settimeout(timeo)
psock.connect(addrinf[4])
chost = host.split('%', 1)[0]
ctx.wrap_socket(psock, server_hostname=chost)
foundsrv = host
psock.close()
break
except OSError:
continue
except ssl.SSLError:
continue
except ssl.CertificateError:
continue
else:
continue
break
@@ -296,6 +296,10 @@ class NetworkManager(object):
subprocess.check_call(['nmcli', 'c', 'u', u])
else:
subprocess.check_call(['nmcli', 'c', 'add', 'type', self.devtypes[iname], 'con-name', cname, 'connection.interface-name', iname] + cargs)
self.read_connections()
u = self.uuidbyname.get(cname, None)
if u:
subprocess.check_call(['nmcli', 'c', 'u', u])
@@ -344,6 +348,13 @@ if __name__ == '__main__':
else:
netname_to_interfaces[uname] = {'interfaces': set([iname]), 'settings': nc['extranets'][netname]}
doneidxs.add(curridx)
if 'default' in netname_to_interfaces:
for netn in netname_to_interfaces:
if netn == 'default':
continue
netname_to_interfaces['default']['interfaces'] -= netname_to_interfaces[netn]['interfaces']
if not netname_to_interfaces['default']['interfaces']:
del netname_to_interfaces['default']
rm_tmp_llas(tmpllas)
if os.path.exists('/usr/bin/nmcli'):
nm = NetworkManager(devtypes)
@@ -0,0 +1,91 @@
Name: confluent_osdeploy-aarch64
Version: #VERSION#
Release: 1
Summary: OS Deployment support for confluent
License: Apache2
URL: https://hpc.lenovo.com/
Source0: confluent_osdeploy.tar.xz
Source1: confluent_el9bin.tar.xz
Source2: confluent_el8bin.tar.xz
BuildArch: noarch
Requires: confluent_ipxe mtools tar
BuildRoot: /tmp
%description
This contains support utilities for enabling deployment of aarch64 architecture systems
%define debug_package %{nil}
%prep
%setup -n confluent_osdeploy -a 2 -a 1
%build
mkdir -p opt/confluent/bin
mkdir -p stateless-bin
cp -a el8bin/* .
ln -s el8 el9
for os in rhvh4 el7 genesis el8 suse15 ubuntu20.04 ubuntu22.04 coreos el9; do
mkdir ${os}out
cd ${os}out
if [ -d ../${os}bin ]; then
cp -a ../${os}bin/opt .
else
cp -a ../opt .
fi
cp -a ../${os}/initramfs/* .
cp -a ../common/initramfs/* .
find . | cpio -H newc -o > ../addons.cpio
mv ../addons.cpio .
cd ..
done
for os in el7 el8 suse15 el9 ubuntu20.04; do
mkdir ${os}disklessout
cd ${os}disklessout
if [ -d ../${os}bin ]; then
cp -a ../${os}bin/opt .
else
cp -a ../opt .
fi
cp -a ../${os}-diskless/initramfs/* .
cp -a ../common/initramfs/* .
if [ -d ../${os}bin ]; then
cp -a ../${os}bin/stateless-bin/* opt/confluent/bin
else
cp -a ../stateless-bin/* opt/confluent/bin
fi
find . | cpio -H newc -o > ../addons.cpio
mv ../addons.cpio .
cd ..
done
mkdir esxi7out
cd esxi7out
cp -a ../opt .
cp -a ../esxi7/initramfs/* .
cp -a ../common/initramfs/* .
chmod +x bin/* opt/confluent/bin/*
tar zcvf ../addons.tgz *
mv ../addons.tgz .
cd ..
cp -a esxi7out esxi6out
cp -a esxi7 esxi6
cp -a esxi7out esxi8out
cp -a esxi7 esxi8
%install
mkdir -p %{buildroot}/opt/confluent/share/licenses/confluent_osdeploy/
#cp LICENSE %{buildroot}/opt/confluent/share/licenses/confluent_osdeploy/
for os in rhvh4 el7 el8 el9 genesis suse15 ubuntu20.04 ubuntu22.04 esxi6 esxi7 esxi8 coreos; do
mkdir -p %{buildroot}/opt/confluent/lib/osdeploy/$os/initramfs/aarch64/
cp ${os}out/addons.* %{buildroot}/opt/confluent/lib/osdeploy/$os/initramfs/aarch64/
if [ -d ${os}disklessout ]; then
mkdir -p %{buildroot}/opt/confluent/lib/osdeploy/${os}-diskless/initramfs/aarch64/
cp ${os}disklessout/addons.* %{buildroot}/opt/confluent/lib/osdeploy/${os}-diskless/initramfs/aarch64/
fi
done
find %{buildroot}/opt/confluent/lib/osdeploy/ -name .gitignore -exec rm -f {} +
%files
/opt/confluent/lib/osdeploy
#%license /opt/confluent/share/licenses/confluent_osdeploy/LICENSE
@@ -1,10 +1,10 @@
#!/bin/bash
function test_mgr() {
whost=$1
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
if curl -s https://${whost}/confluent-api/ > /dev/null; then
if curl -gs https://${whost}/confluent-api/ > /dev/null; then
return 0
fi
return 1
@@ -63,10 +63,10 @@ fetch_remote() {
set_confluent_vars
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo $1 failed to download; return 1; fi
}
@@ -174,10 +174,10 @@ run_remote_python() {
cd $confluentscripttmpdir
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo "'$*'" failed to download; return 1; fi
confluentpython $*
retcode=$?
@@ -1,10 +1,10 @@
#!/bin/bash
function test_mgr() {
whost=$1
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
if curl -s https://${whost}/confluent-api/ > /dev/null; then
if curl -gs https://${whost}/confluent-api/ > /dev/null; then
return 0
fi
return 1
@@ -63,10 +63,10 @@ fetch_remote() {
set_confluent_vars
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo $1 failed to download; return 1; fi
}
@@ -174,10 +174,10 @@ run_remote_python() {
cd $confluentscripttmpdir
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo "'$*'" failed to download; return 1; fi
confluentpython $*
retcode=$?
@@ -1,10 +1,10 @@
#!/bin/bash
function test_mgr() {
whost=$1
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
if curl -s https://${whost}/confluent-api/ > /dev/null; then
if curl -gs https://${whost}/confluent-api/ > /dev/null; then
return 0
fi
return 1
@@ -63,10 +63,10 @@ fetch_remote() {
set_confluent_vars
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo $1 failed to download; return 1; fi
}
@@ -174,10 +174,10 @@ run_remote_python() {
cd $confluentscripttmpdir
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo "'$*'" failed to download; return 1; fi
confluentpython $*
retcode=$?
@@ -100,7 +100,9 @@ cd /sys/class/net
if ! grep MANAGER: /etc/confluent/confluent.info; then
confluentsrv=$(getarg confluent)
if [ ! -z "$confluentsrv" ]; then
mgr=$confluentsrv
if [[ "$confluentsrv" = *":"* ]]; then
mgr="[$mgr]"
confluenthttpsrv=[$confluentsrv]
/usr/libexec/nm-initrd-generator ip=:dhcp6
else
@@ -146,7 +148,7 @@ fi
while ! confluentpython /opt/confluent/bin/apiclient $errout /confluent-api/self/deploycfg2 > /etc/confluent/confluent.deploycfg; do
sleep 10
done
ifidx=$(cat /tmp/confluent.ifidx)
ifidx=$(cat /tmp/confluent.ifidx 2> /dev/null)
if [ ! -z "$ifidx" ]; then
ifname=$(ip link |grep ^$ifidx:|awk '{print $2}')
ifname=${ifname%:}
@@ -2,6 +2,10 @@
sed -i 's/centos/CentOS/; s/rhel/Red Hat Enterprise Linux/; s/oraclelinux/Oracle Linux/; s/alma/AlmaLinux/' $2/profile.yaml
ln -s $1/images/pxeboot/vmlinuz $2/boot/kernel && \
ln -s $1/images/pxeboot/initrd.img $2/boot/initramfs/distribution
mkdir -p $2/boot/efi/boot && \
ln -s $1/EFI/BOOT/BOOTX64.EFI $1/EFI/BOOT/grubx64.efi $2/boot/efi/boot/
mkdir -p $2/boot/efi/boot
if [ -e $1/EFI/BOOT/BOOTAA64.EFI ]; then
ln -s $1/EFI/BOOT/BOOTAA64.EFI $1/EFI/BOOT/grubaa64.efi $2/boot/efi/boot/
else
ln -s $1/EFI/BOOT/BOOTX64.EFI $1/EFI/BOOT/grubx64.efi $2/boot/efi/boot/
fi
@@ -33,16 +33,21 @@ export nodename confluent_mgr confluent_profile
exec >> /var/log/confluent/confluent-firstboot.log
exec 2>> /var/log/confluent/confluent-firstboot.log
chmod 600 /var/log/confluent/confluent-firstboot.log
if [ ! -f /etc/confluent/firstboot.ran ]; then
cat /etc/confluent/tls/*.pem >> /etc/pki/tls/certs/ca-bundle.crt
confluentpython /root/confignet
rm /root/confignet
fi
while ! ping -c 1 $confluent_pingtarget >& /dev/null; do
sleep 1
done
if [ ! -f /etc/confluent/firstboot.ran ]; then
touch /etc/confluent/firstboot.ran
cat /etc/confluent/tls/*.pem >> /etc/pki/tls/certs/ca-bundle.crt
run_remote_python confignet
run_remote firstboot.custom
# Firstboot scripts may be placed into firstboot.d, e.g. firstboot.d/01-firstaction.sh, firstboot.d/02-secondaction.sh
run_remote_parts firstboot.d
@@ -1,10 +1,10 @@
#!/bin/bash
function test_mgr() {
whost=$1
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
if curl -s https://${whost}/confluent-api/ > /dev/null; then
if curl -gs https://${whost}/confluent-api/ > /dev/null; then
return 0
fi
return 1
@@ -63,10 +63,10 @@ fetch_remote() {
set_confluent_vars
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo $1 failed to download; return 1; fi
}
@@ -174,10 +174,10 @@ run_remote_python() {
cd $confluentscripttmpdir
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo "'$*'" failed to download; return 1; fi
confluentpython $*
retcode=$?
@@ -47,5 +47,8 @@ run_remote_parts post.d
# Induce execution of remote configuration, e.g. ansible plays in ansible/post.d/
run_remote_config post.d
cd /root
fetch_remote confignet
cd -
curl -sf -X POST -d 'status: staged' -H "CONFLUENT_NODENAME: $nodename" -H "CONFLUENT_APIKEY: $apikey" https://$confluent_mgr/confluent-api/self/updatestatus
kill $logshowpid
@@ -1,10 +1,10 @@
#!/bin/bash
function test_mgr() {
whost=$1
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
if curl -s https://${whost}/confluent-api/ > /dev/null; then
if curl -gs https://${whost}/confluent-api/ > /dev/null; then
return 0
fi
return 1
@@ -63,10 +63,10 @@ fetch_remote() {
set_confluent_vars
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo $1 failed to download; return 1; fi
}
@@ -174,10 +174,10 @@ run_remote_python() {
cd $confluentscripttmpdir
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo "'$*'" failed to download; return 1; fi
confluentpython $*
retcode=$?
@@ -33,6 +33,8 @@ for info in vswinfo.split('\n'):
upinfo = uplinkmatch.match(info)
if upinfo:
vmnic = upinfo.group(1)
if vmnic and 'vusb0' not in vmnic:
break
try:
with open('/tmp/confluentident/cnflnt.jsn') as identin:
identcfg = json.load(identin)
@@ -1,10 +1,10 @@
#!/bin/bash
function test_mgr() {
whost=$1
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
if curl -s https://${whost}/confluent-api/ > /dev/null; then
if curl -gs https://${whost}/confluent-api/ > /dev/null; then
return 0
fi
return 1
@@ -63,10 +63,10 @@ fetch_remote() {
set_confluent_vars
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo $1 failed to download; return 1; fi
}
@@ -174,10 +174,10 @@ run_remote_python() {
cd $confluentscripttmpdir
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo "'$*'" failed to download; return 1; fi
confluentpython $*
retcode=$?
@@ -77,6 +77,7 @@ if [ "$textconsole" = "true" ] && ! grep console= /proc/cmdline > /dev/null && [
sed -e s'/$/ 'console=${autocons#*/dev/}/ /proc/cmdline > /etc/fakecmdline
mount -o bind /etc/fakecmdline /proc/cmdline
echo "ConsoleDevice: ${autocons%,*}" >> /etc/linuxrc.d/01-confluent
echo "Textmode: 1" >> /etc/linuxrc.d/01-confluent
fi
tz=$(grep timezone: /etc/confluent/confluent.deploycfg | awk '{print $2}')
@@ -1,10 +1,10 @@
#!/bin/bash
function test_mgr() {
whost=$1
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
if curl -s https://${whost}/confluent-api/ > /dev/null; then
if curl -gs https://${whost}/confluent-api/ > /dev/null; then
return 0
fi
return 1
@@ -63,10 +63,10 @@ fetch_remote() {
set_confluent_vars
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo $1 failed to download; return 1; fi
}
@@ -174,10 +174,10 @@ run_remote_python() {
cd $confluentscripttmpdir
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo "'$*'" failed to download; return 1; fi
confluentpython $*
retcode=$?
@@ -1,10 +1,10 @@
#!/bin/bash
function test_mgr() {
whost=$1
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
if curl -s https://${whost}/confluent-api/ > /dev/null; then
if curl -gs https://${whost}/confluent-api/ > /dev/null; then
return 0
fi
return 1
@@ -63,10 +63,10 @@ fetch_remote() {
set_confluent_vars
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo $1 failed to download; return 1; fi
}
@@ -174,10 +174,10 @@ run_remote_python() {
cd $confluentscripttmpdir
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo "'$*'" failed to download; return 1; fi
confluentpython $*
retcode=$?
@@ -1,10 +1,10 @@
#!/bin/bash
function test_mgr() {
whost=$1
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
if curl -s https://${whost}/confluent-api/ > /dev/null; then
if curl -gs https://${whost}/confluent-api/ > /dev/null; then
return 0
fi
return 1
@@ -63,10 +63,10 @@ fetch_remote() {
set_confluent_vars
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo $1 failed to download; return 1; fi
}
@@ -174,10 +174,10 @@ run_remote_python() {
cd $confluentscripttmpdir
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo "'$*'" failed to download; return 1; fi
confluentpython $*
retcode=$?
@@ -1,10 +1,10 @@
#!/bin/bash
function test_mgr() {
whost=$1
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
if curl -s https://${whost}/confluent-api/ > /dev/null; then
if curl -gs https://${whost}/confluent-api/ > /dev/null; then
return 0
fi
return 1
@@ -63,10 +63,10 @@ fetch_remote() {
set_confluent_vars
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo $1 failed to download; return 1; fi
}
@@ -174,10 +174,10 @@ run_remote_python() {
cd $confluentscripttmpdir
mkdir -p $(dirname $1)
whost=$confluent_mgr
if [[ "$whost" == *:* ]]; then
if [[ "$whost" == *:* ]] && [[ "$whost" != *[* ]] ; then
whost="[$whost]"
fi
curl -f -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
curl -gf -sS $curlargs https://$whost/confluent-public/os/$confluent_profile/scripts/$1 > $1
if [ $? != 0 ]; then echo "'$*'" failed to download; return 1; fi
confluentpython $*
retcode=$?
+1 -1
View File
@@ -18,7 +18,7 @@
#define MAXPACKET 1024
static const char cryptalpha[] = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
static const char cryptalpha[] = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789./";
unsigned char* genpasswd(int len) {
unsigned char * passwd;
+11 -6
View File
@@ -52,11 +52,12 @@ def make_certificate():
os.umask(umask)
def show_invitation(name):
def show_invitation(name, nonvoting=False):
if not os.path.exists('/etc/confluent/srvcert.pem'):
make_certificate()
s = client.Command().connection
tlvdata.send(s, {'collective': {'operation': 'invite', 'name': name}})
role = 'nonvoting' if nonvoting else None
tlvdata.send(s, {'collective': {'operation': 'invite', 'name': name, 'role': role}})
invite = tlvdata.recv(s)['collective']
if 'error' in invite:
sys.stderr.write(invite['error'] + '\n')
@@ -104,16 +105,19 @@ def show_collective():
return
if 'quorum' in res['collective']:
print('Quorum: {0}'.format(res['collective']['quorum']))
print('Leader: {0}'.format(res['collective']['leader']))
leadernonvoting = res['collective']['leader'] in res['collective'].get('nonvoting', ())
print('Leader: {0}{1}'.format(res['collective']['leader'], ' (non-voting)' if leadernonvoting else ''))
if 'active' in res['collective']:
if res['collective']['active']:
print('Active collective members:')
for member in sortutil.natural_sort(res['collective']['active']):
print(' {0}'.format(member))
nonvoter = member in res['collective'].get('nonvoting', ())
print(' {0}{1}'.format(member, ' (nonvoting)' if nonvoter else ''))
if res['collective']['offline']:
print('Offline collective members:')
for member in sortutil.natural_sort(res['collective']['offline']):
print(' {0}'.format(member))
nonvoter = member in res['collective'].get('nonvoting', ())
print(' {0}{1}'.format(member, ' (nonvoting)' if nonvoter else ''))
else:
print('Run collective show on leader for more data')
@@ -128,6 +132,7 @@ def main():
'collective member. Run collective invite -h for more information')
ic.add_argument('name', help='Name of server to invite to join the '
'collective')
ic.add_argument('-n', help='Join as a non-voting member, do not have this member contribute to quorum', action='store_true')
dc = sp.add_parser('delete', help='Delete a member of a collective')
dc.add_argument('name', help='Name of server to delete from collective')
jc = sp.add_parser('join', help='Join a collective. Run collective join -h for more information')
@@ -138,7 +143,7 @@ def main():
if cmdset.command == 'gencert':
make_certificate()
elif cmdset.command == 'invite':
show_invitation(cmdset.name)
show_invitation(cmdset.name, cmdset.n)
elif cmdset.command == 'join':
join_collective(cmdset.server, cmdset.i)
elif cmdset.command == 'show':
+6 -1
View File
@@ -15,7 +15,7 @@ import confluent.sshutil as sshutil
import confluent.certutil as certutil
import confluent.client as client
import confluent.config.configmanager as configmanager
import subprocess
import eventlet.green.subprocess as subprocess
import tempfile
import shutil
import eventlet.green.socket as socket
@@ -242,6 +242,11 @@ if __name__ == '__main__':
emprint(f'There is no node named "{args.node}"')
allok = False
uuidok = True # not really, but suppress the spurious error
dnsdomain = rsp.get('dns.domain', {}).get('value', '')
if ',' in dnsdomain or ' ' in dnsdomain:
allok = False
emprint(f'{args.node} has a dns.domain that appears to be a search instead of singular domain')
uuidok = True # not really, but suppress the spurious error
uuid = rsp.get('id.uuid', {}).get('value', None)
if uuid:
uuidok = True
+12 -2
View File
@@ -217,14 +217,24 @@ def install_tftp_content():
else:
emprint(
'Detected {0} as tftp directory, but unable to determine tftp service, ensure that a tftp server is installed and enabled manually'.format(tftplocation))
otftplocation = tftplocation
tftplocation = '{0}/confluent/x86_64'.format(tftplocation)
try:
os.makedirs(tftplocation)
except OSError as e:
if e.errno != 17:
raise
armtftplocation = '{0}/confluent/aarch64'.format(otftplocation)
try:
os.makedirs(armtftplocation)
except OSError as e:
if e.errno != 17:
raise
shutil.copy('/opt/confluent/lib/ipxe/ipxe.efi', tftplocation)
shutil.copy('/opt/confluent/lib/ipxe/ipxe.kkpxe', tftplocation)
if os.path.exists('/opt/confluent/lib/ipxe/ipxe-aarch64.efi'):
shutil.copy('/opt/confluent/lib/ipxe/ipxe-aarch64.efi', os.path.join(armtftplocation, 'ipxe.efi'))
def initialize(cmdset):
@@ -285,13 +295,13 @@ def initialize(cmdset):
try:
subprocess.check_call(['systemctl', 'try-restart', 'httpd'])
print('HTTP server has been restarted if it was running')
except subprocess.CalledProcessError:
except Exception:
emprint('New HTTPS certificates generated, restart the web server manually')
elif os.path.exists('/usr/lib/systemd/system/apache2.service'):
try:
subprocess.check_call(['systemctl', 'try-restart', 'apache2'])
print('HTTP server has been restarted if it was running')
except subprocess.CalledProcessError:
except Exception:
emprint('New HTTPS certificates generated, restart the web server manually')
else:
emprint('New HTTPS certificates generated, restart the web server manually')
@@ -22,12 +22,13 @@ import hmac
import os
pending_invites = {}
def create_server_invitation(servername):
def create_server_invitation(servername, role):
servername = servername.encode('utf-8')
randbytes = (3 - ((len(servername) + 2) % 3)) % 3 + 64
invitation = os.urandom(randbytes)
pending_invites[servername] = invitation
return base64.b64encode(servername + b'@' + invitation)
pending_invites[servername] = {'invitation': invitation, 'role': role}
invite = servername + b'@' + invitation
return base64.b64encode(invite)
def create_client_proof(invitation, mycert, peercert):
return hmac.new(invitation, peercert + mycert, hashlib.sha256).digest()
@@ -42,6 +43,8 @@ def check_client_proof(servername, mycert, peercert, proof):
if servername not in pending_invites:
return False
invitation = pending_invites[servername]
role = invitation['role']
invitation = invitation['invitation']
validproof = hmac.new(invitation, mycert + peercert, hashlib.sha256
).digest()
if proof == validproof:
@@ -54,7 +57,7 @@ def check_client_proof(servername, mycert, peercert, proof):
# Now to generate an answer...., reverse the cert order so our answer
# is different, but still proving things
return hmac.new(invitation, peercert + mycert, hashlib.sha256
).digest()
).digest(), role
# The given proof did not verify the invitation
return False
return False, None
@@ -61,10 +61,15 @@ class ContextBool(object):
connecting = ContextBool()
leader_init = ContextBool()
enrolling = ContextBool()
def connect_to_leader(cert=None, name=None, leader=None, remote=None):
global currentleader
global follower
ocert = cert
oname = name
oleader = leader
oremote = remote
if leader is None:
leader = currentleader
log.log({'info': 'Attempting connection to leader {0}'.format(leader),
@@ -138,7 +143,9 @@ def connect_to_leader(cert=None, name=None, leader=None, remote=None):
remote.close()
except Exception:
pass
raise Exception("Error doing initial DB transfer")
log.log({'error': 'Retrying connection, error during initial sync', 'subsystem': 'collective'})
return connect_to_leader(ocert, oname, oleader, oremote)
raise Exception("Error doing initial DB transfer") # bad ssl write retry
dbjson += ndata
cfm.clear_configuration()
try:
@@ -146,7 +153,7 @@ def connect_to_leader(cert=None, name=None, leader=None, remote=None):
for c in colldata:
cfm._true_add_collective_member(c, colldata[c]['address'],
colldata[c]['fingerprint'],
sync=False)
sync=False, role=colldata[c].get('role', None))
for globvar in globaldata:
cfm.set_global(globvar, globaldata[globvar], False)
cfm._txcount = dbi.get('txcount', 0)
@@ -321,7 +328,8 @@ def handle_connection(connection, cert, request, local=False):
#TODO(jjohnson2): Cannot do the invitation if not the head node, the certificate hand-carrying
#can't work in such a case.
name = request['name']
invitation = invites.create_server_invitation(name)
role = request.get('role', '')
invitation = invites.create_server_invitation(name, role)
tlvdata.send(connection,
{'collective': {'invitation': invitation}})
connection.close()
@@ -390,34 +398,44 @@ def handle_connection(connection, cert, request, local=False):
eventlet.spawn_n(connect_to_leader, rsp['collective'][
'fingerprint'], name)
if 'enroll' == operation:
#TODO(jjohnson2): error appropriately when asked to enroll, but the master is elsewhere
mycert = util.get_certificate_from_file('/etc/confluent/srvcert.pem')
proof = base64.b64decode(request['hmac'])
myrsp = invites.check_client_proof(request['name'], mycert,
cert, proof)
if not myrsp:
tlvdata.send(connection, {'error': 'Invalid token'})
connection.close()
return
if not list(cfm.list_collective()):
# First enrollment of a collective, since the collective doesn't
# quite exist, then set initting false to let the enrollment action
# drive this particular initialization
initting = False
myrsp = base64.b64encode(myrsp)
fprint = util.get_fingerprint(cert)
myfprint = util.get_fingerprint(mycert)
cfm.add_collective_member(get_myname(),
connection.getsockname()[0], myfprint)
cfm.add_collective_member(request['name'],
connection.getpeername()[0], fprint)
myleader = get_leader(connection)
ldrfprint = cfm.get_collective_member_by_address(
myleader)['fingerprint']
tlvdata.send(connection,
{'collective': {'approval': myrsp,
'fingerprint': ldrfprint,
'leader': get_leader(connection)}})
with enrolling:
cfm.check_quorum()
mycert = util.get_certificate_from_file('/etc/confluent/srvcert.pem')
proof = base64.b64decode(request['hmac'])
myrsp, role = invites.check_client_proof(request['name'], mycert,
cert, proof)
if not myrsp:
tlvdata.send(connection, {'error': 'Invalid token'})
connection.close()
return
if not list(cfm.list_collective()):
# First enrollment of a collective, since the collective doesn't
# quite exist, then set initting false to let the enrollment action
# drive this particular initialization
initting = False
myrsp = base64.b64encode(myrsp)
fprint = util.get_fingerprint(cert)
myfprint = util.get_fingerprint(mycert)
iam = cfm.get_collective_member(get_myname())
if not iam:
cfm.add_collective_member(get_myname(),
connection.getsockname()[0], myfprint)
cfm.add_collective_member(request['name'],
connection.getpeername()[0], fprint, role)
myleader = get_leader(connection)
ldrfprint = cfm.get_collective_member_by_address(
myleader)['fingerprint']
tlvdata.send(connection,
{'collective': {'approval': myrsp,
'fingerprint': ldrfprint,
'leader': get_leader(connection)}})
havequorum = False
while not havequorum:
try:
cfm.check_quorum()
havequorum = True
except exc.DegradedCollective:
eventlet.sleep(0.1)
if 'assimilate' == operation:
drone = request['name']
droneinfo = cfm.get_collective_member(drone)
@@ -575,9 +593,12 @@ def populate_collinfo(collinfo):
activemembers = set(cfm.cfgstreams)
activemembers.add(iam)
collinfo['offline'] = []
collinfo['nonvoting'] = []
for member in cfm.list_collective():
if member not in activemembers:
collinfo['offline'].append(member)
if cfm.get_collective_member(member).get('role', None) == 'nonvoting':
collinfo['nonvoting'].append(member)
def try_assimilate(drone, followcount, remote):
@@ -233,6 +233,12 @@ node = {
'a stateless profile would be both after first boot.')
},
'deployment.state': {
'description': ('Profiles may push more specific state, for example, it may set the state to "failed" or "succeded"'),
},
'deployment.state_detail': {
'description': ('Detailed state information as reported by an OS profile, when available'),
},
'deployment.useinsecureprotocols': {
'description': ('What phase(s) of boot are permitted to use insecure protocols '
'(TFTP and HTTP without TLS. By default, only HTTPS is used. However '
@@ -89,6 +89,7 @@ import string
import struct
import sys
import threading
import time
import traceback
try:
unicode
@@ -320,7 +321,7 @@ def _rpc_set_group_attributes(tenant, attribmap, autocreate):
def check_quorum():
if isinstance(cfgleader, bool):
raise exc.DegradedCollective()
if (not cfgleader) and len(cfgstreams) < (len(_cfgstore.get('collective', {})) // 2):
if (not cfgleader) and (not has_quorum()):
# the leader counts in addition to registered streams
raise exc.DegradedCollective()
if cfgleader and not _hasquorum:
@@ -348,9 +349,9 @@ def exec_on_followers(fnname, *args):
pushes = eventlet.GreenPool()
# Check health of collective prior to attempting
for _ in pushes.starmap(
_push_rpc, [(cfgstreams[s], b'') for s in cfgstreams]):
_push_rpc, [(cfgstreams[s]['stream'], b'') for s in cfgstreams]):
pass
if len(cfgstreams) < (len(_cfgstore['collective']) // 2):
if not has_quorum():
# the leader counts in addition to registered streams
raise exc.DegradedCollective()
exec_on_followers_unconditional(fnname, *args)
@@ -363,7 +364,7 @@ def exec_on_followers_unconditional(fnname, *args):
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]):
_push_rpc, [(cfgstreams[s]['stream'], payload) for s in cfgstreams]):
pass
@@ -615,6 +616,40 @@ def set_global(globalname, value, sync=True):
if sync:
ConfigManager._bg_sync_to_file()
mycachedname = [None, 0]
def get_myname():
if mycachedname[1] > time.time() - 15:
return mycachedname[0]
try:
with open('/etc/confluent/cfg/myname', 'r') as f:
mycachedname[0] = f.read().strip()
mycachedname[1] = time.time()
return mycachedname[0]
except IOError:
myname = socket.gethostname().split('.')[0]
with open('/etc/confluent/cfg/myname', 'w') as f:
f.write(myname)
mycachedname[0] = myname
mycachedname[1] = time.time()
return myname
def has_quorum():
voters = 0
for follower in cfgstreams:
if cfgstreams[follower].get('role', None) != 'nonvoting':
voters += 1
iam = get_collective_member(get_myname())
myrole = None
if iam:
myrole = iam.get('role', None)
if myrole != 'nonvoting':
voters += 1
allvoters = 0
for ghost in list_collective():
if get_collective_member(ghost).get('role', None) != 'nonvoting':
allvoters += 1
return voters > allvoters // 2
cfgstreams = {}
def relay_slaved_requests(name, listener):
global cfgleader
@@ -622,21 +657,21 @@ def relay_slaved_requests(name, listener):
pushes = eventlet.GreenPool()
if name not in _followerlocks:
_followerlocks[name] = gthread.RLock()
meminfo = get_collective_member(name)
with _followerlocks[name]:
try:
stop_following()
if name in cfgstreams:
try:
cfgstreams[name].close()
cfgstreams[name]['stream'].close()
except Exception:
pass
del cfgstreams[name]
if membership_callback:
membership_callback()
cfgstreams[name] = listener
cfgstreams[name] = {'stream': listener, 'role': meminfo.get('role', None)}
lh = StreamHandler(listener)
_hasquorum = len(cfgstreams) >= (
len(_cfgstore['collective']) // 2)
_hasquorum = has_quorum()
_newquorum = None
while _hasquorum != _newquorum:
if _newquorum is not None:
@@ -644,10 +679,9 @@ def relay_slaved_requests(name, listener):
payload = msgpack.packb({'quorum': _hasquorum}, use_bin_type=False)
for _ in pushes.starmap(
_push_rpc,
[(cfgstreams[s], payload) for s in cfgstreams]):
[(cfgstreams[s]['stream'], payload) for s in cfgstreams]):
pass
_newquorum = len(cfgstreams) >= (
len(_cfgstore['collective']) // 2)
_newquorum = has_quorum()
_hasquorum = _newquorum
if _hasquorum and _pending_collective_updates:
apply_pending_collective_updates()
@@ -693,12 +727,11 @@ def relay_slaved_requests(name, listener):
except KeyError:
pass # May have already been closed/deleted...
if cfgstreams:
_hasquorum = len(cfgstreams) >= (
len(_cfgstore['collective']) // 2)
_hasquorum = has_quorum()
payload = msgpack.packb({'quorum': _hasquorum}, use_bin_type=False)
for _ in pushes.starmap(
_push_rpc,
[(cfgstreams[s], payload) for s in cfgstreams]):
[(cfgstreams[s]['stream'], payload) for s in cfgstreams]):
pass
if membership_callback:
membership_callback()
@@ -756,8 +789,8 @@ def stop_leading(newleader=None):
for stream in list(cfgstreams):
try:
if rpcpayload is not None:
_push_rpc(cfgstreams[stream], rpcpayload)
cfgstreams[stream].close()
_push_rpc(cfgstreams[stream]['stream'], rpcpayload)
cfgstreams[stream]['stream'].close()
except Exception:
pass
try:
@@ -876,12 +909,12 @@ def follow_channel(channel):
return {}
def add_collective_member(name, address, fingerprint):
def add_collective_member(name, address, fingerprint, role=None):
if cfgleader:
return exec_on_leader('add_collective_member', name, address, fingerprint)
return exec_on_leader('add_collective_member', name, address, fingerprint, role)
if cfgstreams:
exec_on_followers('_true_add_collective_member', name, address, fingerprint)
_true_add_collective_member(name, address, fingerprint)
exec_on_followers('_true_add_collective_member', name, address, fingerprint, True, role)
_true_add_collective_member(name, address, fingerprint, role=role)
def del_collective_member(name):
if cfgleader and not isinstance(cfgleader, bool):
@@ -934,7 +967,7 @@ def apply_pending_collective_updates():
del _pending_collective_updates[name]
def _true_add_collective_member(name, address, fingerprint, sync=True):
def _true_add_collective_member(name, address, fingerprint, sync=True, role=None):
name = confluent.util.stringify(name)
if _cfgstore is None:
init(not sync) # use not sync to avoid read from disk
@@ -942,6 +975,8 @@ def _true_add_collective_member(name, address, fingerprint, sync=True):
_cfgstore['collective'] = {}
_cfgstore['collective'][name] = {'name': name, 'address': address,
'fingerprint': fingerprint}
if role:
_cfgstore['collective'][name]['role'] = role
with _dirtylock:
if 'collectivedirty' not in _cfgstore:
_cfgstore['collectivedirty'] = set([])
@@ -1515,6 +1550,8 @@ class ConfigManager(object):
groupname = confluent.util.stringify(groupname)
if groupname in self._cfgstore['usergroups']:
raise Exception("Duplicate groupname requested")
if role is None:
role = 'Administrator'
for candrole in _validroles:
if candrole.lower().startswith(role.lower()):
role = candrole
@@ -1625,6 +1662,8 @@ class ConfigManager(object):
name = confluent.util.stringify(name)
if name in self._cfgstore['users']:
raise Exception("Duplicate username requested")
if role is None:
role = 'Administrator'
for candrole in _validroles:
if candrole.lower().startswith(role.lower()):
role = candrole
+4 -4
View File
@@ -79,14 +79,14 @@ class CredServer(object):
hmackey = hmackey.get(nodename, {}).get('secret.selfapiarmtoken', {}).get('value', None)
elif tlv[1]:
client.recv(tlv[1])
apimats = self.cfm.get_node_attributes(nodename,
['deployment.apiarmed', 'deployment.sealedapikey'])
apiarmed = apimats.get(nodename, {}).get('deployment.apiarmed', {}).get(
'value', None)
if not hmackey:
if not address_is_somewhat_trusted(peer[0], nodename, self.cfm):
client.close()
return
apimats = self.cfm.get_node_attributes(nodename,
['deployment.apiarmed', 'deployment.sealedapikey'])
apiarmed = apimats.get(nodename, {}).get('deployment.apiarmed', {}).get(
'value', None)
if not apiarmed:
if apimats.get(nodename, {}).get(
'deployment.sealedapikey', {}).get('value', None):
@@ -116,6 +116,7 @@ nodehandlers = {
'pxe-client': pxeh,
'onie-switch': None,
'cumulus-switch': None,
'affluent-switch': None,
'service:io-device.Lenovo:management-module': None,
'service:thinkagile-storage': cpstorage,
'service:lenovo-tsm': tsm,
@@ -127,6 +128,7 @@ servicenames = {
'cumulus-switch': 'cumulus-switch',
'service:lenovo-smm': 'lenovo-smm',
'service:lenovo-smm2': 'lenovo-smm2',
'affluent-switch': 'affluent-switch',
'lenovo-xcc': 'lenovo-xcc',
'service:management-hardware.IBM:integrated-management-module2': 'lenovo-imm2',
'service:io-device.Lenovo:management-module': 'lenovo-switch',
@@ -140,6 +142,7 @@ servicebyname = {
'cumulus-switch': 'cumulus-switch',
'lenovo-smm': 'service:lenovo-smm',
'lenovo-smm2': 'service:lenovo-smm2',
'affluent-switch': 'affluent-switch',
'lenovo-xcc': 'lenovo-xcc',
'lenovo-imm2': 'service:management-hardware.IBM:integrated-management-module2',
'lenovo-switch': 'service:io-device.Lenovo:management-module',
@@ -356,6 +359,14 @@ def show_info(mac):
for i in send_discovery_datum(known_info[mac]):
yield i
def dump_discovery():
infobymac = {}
for mac in known_info:
infobymac[mac] = {}
for i in send_discovery_datum(known_info[mac]):
for kn in i.kvpairs:
infobymac[mac][kn] = i.kvpairs[kn]
yield msg.KeyValueData(infobymac)
list_info = {
'by-node': list_matching_nodes,
@@ -599,11 +610,14 @@ def handle_read_api_request(pathcomponents):
# starting at 2 are parameters to previous index
if pathcomponents == ['discovery', 'rescan']:
return (msg.KeyValueData({'scanning': bool(scanner)}),)
if pathcomponents == ['discovery', 'alldata']:
return dump_discovery()
subcats, queryparms, indexof, coll = _parameterize_path(pathcomponents[1:])
if len(pathcomponents) == 1:
dirlist = [msg.ChildCollection(x + '/') for x in sorted(list(subcats))]
dirlist.append(msg.ChildCollection('rescan'))
dirlist.append(msg.ChildCollection('autosense'))
dirlist.append(msg.ChildCollection('alldata'))
dirlist.append(msg.ChildCollection('subscriptions/'))
return dirlist
if not coll:
@@ -1143,6 +1157,38 @@ def get_nodename_from_enclosures(cfg, info):
return nodename
def search_smms_by_cert(currsmm, cert, cfg):
neighs = []
cv = util.TLSCertVerifier(
cfg, currsmm, 'pubkeys.tls_hardwaremanager').verify_cert
try:
cd = cfg.get_node_attributes(currsmm, ['hardwaremanagement.manager',
'pubkeys.tls_hardwaremanager'])
smmaddr = cd.get(currsmm, {}).get('hardwaremanagement.manager', {}).get('value', None)
wc = webclient.SecureHTTPConnection(currsmm, verifycallback=cv)
neighs = wc.grab_json_response('/scripts/neighdata.json')
except Exception:
return None
for neigh in neighs:
fprint = neigh.get('sha384', None)
if fprint and fprint.endswith('AA=='):
fprint = fprint[:-4]
if fprint and util.cert_matches(fprint, cert):
port = neigh.get('port', None)
if port is not None:
bay = port + 1
nl = list(
cfg.filter_node_attributes('enclosure.manager=' + currsmm))
nl = list(
cfg.filter_node_attributes('enclosure.bay={}'.format(bay), nl))
if len(nl) == 1:
return currsmm, bay, nl[0]
return currsmm, bay, None
exnl = list(cfg.filter_node_attributes('enclosure.extends=' + currsmm))
if len(exnl) == 1:
return search_smms_by_cert(exnl[0], cert, cfg)
def eval_node(cfg, handler, info, nodename, manual=False):
try:
handler.probe() # unicast interrogation as possible to get more data
@@ -1177,6 +1223,14 @@ def eval_node(cfg, handler, info, nodename, manual=False):
# The specified node is an enclosure (has nodes mapped to it), but
# what we are talking to is *not* an enclosure
# might be ambiguous, need to match chassis-uuid as well..
match = search_smms_by_cert(nodename, handler.https_cert, cfg)
if match:
info['verfied'] = True
info['enclosure.bay'] = match[1]
if match[2]:
if not discover_node(cfg, handler, info, match[2], manual):
pending_nodes[match[2]] = nodename
return
if 'enclosure.bay' not in info:
unknown_info[info['hwaddr']] = info
info['discostatus'] = 'unidentified'
@@ -365,6 +365,8 @@ def proxydhcp(handler, nodeguess):
bootfile = b'confluent/x86_64/ipxe.efi'
elif disco['arch'] == 'bios-x86':
bootfile = b'confluent/x86_64/ipxe.kkpxe'
elif disco['arch'] == 'uefi-aarch64':
bootfile = b'confluent/aarch64/ipxe.efi'
if len(bootfile) > 127:
log.log(
{'info': 'Boot offer cannot be made to {0} as the '
@@ -58,7 +58,7 @@ smsg = ('M-SEARCH * HTTP/1.1\r\n'
def active_scan(handler, protocol=None):
known_peers = set([])
for scanned in scan(['urn:dmtf-org:service:redfish-rest:1']):
for scanned in scan(['urn:dmtf-org:service:redfish-rest:1', 'urn::service:affluent']):
for addr in scanned['addresses']:
if addr in known_peers:
break
@@ -381,6 +381,24 @@ def _find_service(service, target):
querypool = gp.GreenPool()
pooltargs = []
for nid in peerdata:
if peerdata[nid].get('services', [None])[0] == 'urn::service:affluent:1':
peerdata[nid]['attributes'] = {
'type': 'affluent-switch',
}
peerdata[nid]['services'] = ['affluent-switch']
mya = peerdata[nid]['attributes']
usn = peerdata[nid]['usn']
idinfo = usn.split('::')
for idi in idinfo:
key, val = idi.split(':', 1)
if key == 'uuid':
peerdata[nid]['uuid'] = val
elif key == 'serial':
mya['enclosure-serial-number'] = [val]
elif key == 'model':
mya['enclosure-machinetype-model'] = [val]
yield peerdata[nid]
continue
if '/redfish/v1/' not in peerdata[nid].get('urls', ()) and '/redfish/v1' not in peerdata[nid].get('urls', ()):
continue
if '/DeviceDescription.json' in peerdata[nid]['urls']:
@@ -466,6 +484,10 @@ def _parse_ssdp(peer, rsp, peerdata):
peerdatum['services'] = [value]
elif value not in peerdatum['services']:
peerdatum['services'].append(value)
elif header == 'USN':
peerdatum['usn'] = value
elif header == 'MODELNAME':
peerdatum['modelname'] = value
+4 -1
View File
@@ -1682,16 +1682,19 @@ class NetworkConfiguration(ConfluentMessage):
desc = 'Network configuration'
def __init__(self, name=None, ipv4addr=None, ipv4gateway=None,
ipv4cfgmethod=None, hwaddr=None):
ipv4cfgmethod=None, hwaddr=None, staticv6addrs=(), staticv6gateway=None):
self.myargs = (name, ipv4addr, ipv4gateway, ipv4cfgmethod, hwaddr)
self.notnode = name is None
self.stripped = False
v6addrs = ','.join(staticv6addrs)
kvpairs = {
'ipv4_address': {'value': ipv4addr},
'ipv4_gateway': {'value': ipv4gateway},
'ipv4_configuration': {'value': ipv4cfgmethod},
'hw_addr': {'value': hwaddr},
'static_v6_addresses': {'value': v6addrs},
'static_v6_gateway': {'value': staticv6gateway}
}
if self.notnode:
self.kvpairs = kvpairs
+5 -2
View File
@@ -198,8 +198,9 @@ class NetManager(object):
ipv4addr = attribs.get('ipv4_address', None)
if ipv4addr:
try:
for ai in socket.getaddrinfo(ipv4addr, 0, socket.AF_INET, socket.SOCK_STREAM):
ipv4addr = ai[-1][0]
luaddr = ipv4addr.split('/', 1)[0]
for ai in socket.getaddrinfo(luaddr, 0, socket.AF_INET, socket.SOCK_STREAM):
ipv4addr.replace(luaddr, ai[-1][0])
except socket.gaierror:
pass
else:
@@ -661,6 +662,8 @@ def addresses_match(addr1, addr2):
:param addr2:
:return: True if the given addresses refer to the same thing
"""
if '%' in addr1 or '%' in addr2:
return False
for addrinfo in socket.getaddrinfo(addr1, 0, 0, socket.SOCK_STREAM):
rootaddr1 = socket.inet_pton(addrinfo[0], addrinfo[4][0])
if addrinfo[0] == socket.AF_INET6 and rootaddr1[:12] == b'\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\xff\xff':
@@ -220,7 +220,10 @@ def _start_offloader():
def _recv_offload():
upacker = msgpack.Unpacker(encoding='utf8')
try:
upacker = msgpack.Unpacker(encoding='utf8')
except TypeError:
upacker = msgpack.Unpacker(raw=False, strict_map_key=False)
instream = _offloader.stdout.fileno()
while True:
select.select([_offloader.stdout], [], [])
@@ -601,7 +604,7 @@ def handle_read_api_request(pathcomponents, configmanager):
elif len(pathcomponents) == 2:
if pathcomponents[-1] == 'macs':
return [msg.ChildCollection(x) for x in (# 'by-node/',
'by-mac/', 'by-switch/',
'alldata', 'by-mac/', 'by-switch/',
'rescan')]
elif pathcomponents[-1] == 'neighbors':
return [msg.ChildCollection('by-switch/')]
@@ -616,6 +619,8 @@ def handle_read_api_request(pathcomponents, configmanager):
elif len(pathcomponents) == 4:
macaddr = pathcomponents[-1].replace('-', ':')
return dump_macinfo(macaddr)
elif pathcomponents[2] == 'alldata':
return [msg.KeyValueData(_apimacmap)]
elif pathcomponents[2] == 'by-mac':
if len(pathcomponents) == 3:
return [msg.ChildCollection(x.replace(':', '-'))
@@ -687,7 +692,10 @@ def rescan(cfg):
if __name__ == '__main__':
if len(sys.argv) > 1 and sys.argv[1] == '-o':
upacker = msgpack.Unpacker(encoding='utf8')
try:
upacker = msgpack.Unpacker(encoding='utf8')
except TypeError:
upacker = msgpack.Unpacker(raw=False, strict_map_key=False)
currfl = fcntl.fcntl(sys.stdin.fileno(), fcntl.F_GETFL)
fcntl.fcntl(sys.stdin.fileno(), fcntl.F_SETFL, currfl | os.O_NONBLOCK)
@@ -712,4 +720,4 @@ if __name__ == '__main__':
print("Mac to location lookup table: -------------------")
print(repr(_macmap))
print("switch to fdb lookup table: -------------------")
print(repr(_macsbyswitch))
print(repr(_macsbyswitch))
+25 -6
View File
@@ -376,6 +376,14 @@ def check_ubuntu(isoinfo):
if not isinstance(arch, str):
arch = arch.decode('utf8')
major = '.'.join(ver.split('.', 2)[:2])
if 'efi/boot/bootaa64.efi' in isoinfo[0]:
exlist = ['casper/vmlinuz', 'casper/initrd',
'efi/boot/bootaa64.efi', 'efi/boot/grubaa64.efi'
]
else:
exlist = ['casper/vmlinuz', 'casper/initrd',
'efi/boot/bootx64.efi', 'efi/boot/grubx64.efi'
]
return {'name': 'ubuntu-{0}-{1}'.format(ver, arch),
'method': EXTRACT|COPY,
'extractlist': ['casper/vmlinuz', 'casper/initrd',
@@ -765,10 +773,15 @@ def generate_stock_profiles(defprofile, distpath, targpath, osname,
yout.write('# This manifest enables rebase to know original source of profile data and if any customizations have been done\n')
manifestdata = {'distdir': srcname, 'disthashes': hmap}
yout.write(yaml.dump(manifestdata, default_flow_style=False))
for initrd in os.listdir('{0}/initramfs'.format(defprofile)):
fullpath = '{0}/initramfs/{1}'.format(defprofile, initrd)
initrds = ['{0}/initramfs/{1}'.format(defprofile, initrd) for initrd in os.listdir('{0}/initramfs'.format(defprofile))]
if os.path.exists('{0}/initramfs/{1}'.format(defprofile, arch)):
initrds.extend(['{0}/initramfs/{1}/{2}'.format(defprofile, arch, initrd) for initrd in os.listdir('{0}/initramfs/{1}'.format(defprofile, arch))])
for fullpath in initrds:
initrd = os.path.basename(fullpath)
if os.path.isdir(fullpath):
continue
if os.path.exists('{0}/boot/initramfs/{1}'.format(dirname, initrd)):
os.remove('{0}/boot/initramfs/{1}'.format(dirname, initrd))
os.symlink(fullpath,
'{0}/boot/initramfs/{1}'.format(dirname, initrd))
os.symlink(
@@ -792,11 +805,17 @@ class MediaImporter(object):
raise Exception('`osdeploy initialize` must be executed before importing any media')
self.profiles = []
medfile = None
self.medfile = None
if cfm and media in cfm.clientfiles:
medfile = cfm.clientfiles[media]
self.medfile = cfm.clientfiles[media]
medfile = self.medfile
else:
medfile = open(media, 'rb')
identity = fingerprint(medfile)
try:
identity = fingerprint(medfile)
finally:
if not self.medfile:
medfile.close()
if not identity:
raise exc.InvalidArgumentException('Unsupported Media')
self.percent = 0.0
@@ -824,7 +843,6 @@ class MediaImporter(object):
del importing[importkey]
raise Exception('{0} already exists'.format(self.targpath))
self.filename = os.path.abspath(media)
self.medfile = medfile
self.error = ''
self.importer = eventlet.spawn(self.importmedia)
@@ -838,7 +856,8 @@ class MediaImporter(object):
def importmedia(self):
os.environ['PYTHONPATH'] = ':'.join(sys.path)
os.environ['CONFLUENT_MEDIAFD'] = '{0}'.format(self.medfile.fileno())
if self.medfile:
os.environ['CONFLUENT_MEDIAFD'] = '{0}'.format(self.medfile.fileno())
with open(os.devnull, 'w') as devnull:
self.worker = subprocess.Popen(
[sys.executable, __file__, self.filename, '-b'],
@@ -49,7 +49,8 @@ class WebClient(object):
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))
#must be str not bytes
results.put(msg.ConfluentNodeError(self.node, 'Unknown error: {} while retrieving {}'.format(rsp, url)))
return {}
return rsp
@@ -192,4 +193,4 @@ def retrieve_health(configmanager, creds, node, results, element):
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))
results.put(msg.SensorReadings(hinfo.get('sensors', []), name=node))
@@ -17,6 +17,9 @@ import confluent.util as util
import confluent.messages as msg
import confluent.exceptions as exc
import eventlet
import eventlet.green.socket as socket
import pyghmi.util.webclient as wc
import confluent.util as util
import re
import hashlib
import json
@@ -59,12 +62,32 @@ class WebResponse(httplib.HTTPResponse):
def _check_close(self):
return True
class WebConnection(httplib.HTTPConnection):
class WebConnection(wc.SecureHTTPConnection):
response_class = WebResponse
def __init__(self, host):
httplib.HTTPConnection.__init__(self, host, 80)
def __init__(self, host, secure, verifycallback):
if secure:
port = 443
else:
port = 80
wc.SecureHTTPConnection.__init__(self, host, port, verifycallback=verifycallback)
self.secure = secure
self.cookies = {}
def connect(self):
if self.secure:
return super(WebConnection, self).connect()
addrinfo = socket.getaddrinfo(self.host, self.port)[0]
# workaround problems of too large mtu, moderately frequent occurance
# in this space
plainsock = socket.socket(addrinfo[0])
plainsock.settimeout(self.mytimeout)
try:
plainsock.setsockopt(socket.IPPROTO_TCP, socket.TCP_MAXSEG, 1456)
except socket.error:
pass
plainsock.connect(addrinfo[4])
self.sock = plainsock
def getresponse(self):
try:
rsp = super(WebConnection, self).getresponse()
@@ -132,8 +155,18 @@ class PDUClient(object):
if not target:
target = self.node
target = target.split('/', 1)[0]
self._wc = WebConnection(target)
self.login(self.configmanager)
verifier = util.TLSCertVerifier(
self.configmanager, self.node, 'pubkeys.tls_hardwaremanager')
try:
self._wc = WebConnection(target, secure=True, verifycallback=verifier.verify_cert)
self.login(self.configmanager)
except socket.error as e:
pkey = self.configmanager.get_node_attributes(self.node, 'pubkeys.tls_hardwaremanager')
pkey = pkey.get(self.node, {}).get('pubkeys.tls_hardwaremanager', {}).get('value', None)
if pkey:
raise
self._wc = WebConnection(target, secure=False, verifycallback=verifier.verify_cert)
self.login(self.configmanager)
return self._wc
def login(self, configmanager):
@@ -153,18 +186,27 @@ class PDUClient(object):
if not username or not passwd:
raise Exception('Missing username or password')
b64user = base64.b64encode(username.encode('utf8')).decode('utf8')
b64pass = base64.b64encode(passwd.encode('utf8')).decode('utf8')
rsp = self.wc.grab_response('/config/gateway?page=cgi_authentication&login={}&_dc={}'.format(b64user, int(time.time())))
rsp = json.loads(sanitize_json(rsp[0]))
parms = answer_challenge(username, passwd, rsp['data'][-1])
self.sessid = rsp['data'][0]
url = '/config/gateway?page=cgi_authenticationChallenge&sessionId={}&login={}&sessionKey={}&szResponse={}&szResponseValue={}&dc={}'.format(
rsp['data'][0],
b64user,
parms['sessionKey'],
parms['szResponse'],
parms['szResponseValue'],
int(time.time()),
)
if rsp['data'][-1] == 'password':
url = '/config/gateway?page=cgi_authenticationPassword&login={}&sessionId={}&password={}&dc={}'.format(
b64user,
rsp['data'][0],
b64pass,
int(time.time()),
)
else:
parms = answer_challenge(username, passwd, rsp['data'][-1])
url = '/config/gateway?page=cgi_authenticationChallenge&sessionId={}&login={}&sessionKey={}&szResponse={}&szResponseValue={}&dc={}'.format(
rsp['data'][0],
b64user,
parms['sessionKey'],
parms['szResponse'],
parms['szResponseValue'],
int(time.time()),
)
rsp = self.wc.grab_response(url)
rsp = json.loads(sanitize_json(rsp[0]))
if rsp['success'] != True:
@@ -177,7 +219,7 @@ class PDUClient(object):
return wc.grab_response(url)
def logout(self):
print(repr(self.do_request('cgi_logout')))
self.do_request('cgi_logout')
def get_outlet(self, outlet):
rsp = self.do_request('cgi_pdu_outlets')
@@ -735,11 +735,14 @@ class IpmiHandler(object):
elif len(self.element) == 4 and self.element[-1] == 'management':
if self.op == 'read':
lancfg = self.ipmicmd.get_net_configuration()
v6cfg = self.ipmicmd.get_net6_configuration()
self.output.put(msg.NetworkConfiguration(
self.node, ipv4addr=lancfg['ipv4_address'],
ipv4gateway=lancfg['ipv4_gateway'],
ipv4cfgmethod=lancfg['ipv4_configuration'],
hwaddr=lancfg['mac_address']
hwaddr=lancfg['mac_address'],
staticv6addrs=v6cfg.get('static_addrs', ''),
staticv6gateway=v6cfg.get('static_gateway', ''),
))
elif self.op == 'update':
config = self.inputdata.netconfig(self.node)
@@ -748,6 +751,11 @@ class IpmiHandler(object):
ipv4_address=config['ipv4_address'],
ipv4_configuration=config['ipv4_configuration'],
ipv4_gateway=config['ipv4_gateway'])
v6addrs = config.get('static_v6_addresses', None)
if v6addrs is not None:
v6addrs = v6addrs.split(',')
v6gw = config.get('static_v6_gateway', None)
self.ipmicmd.set_net6_configuration(static_addresses=v6addrs, static_gateway=v6gw)
except socket.error as se:
self.output.put(msg.ConfluentNodeError(self.node,
se.message))
@@ -589,11 +589,14 @@ class IpmiHandler(object):
elif len(self.element) == 4 and self.element[-1] == 'management':
if self.op == 'read':
lancfg = self.ipmicmd.get_net_configuration()
v6cfg = self.ipmicmd.get_net6_configuration()
self.output.put(msg.NetworkConfiguration(
self.node, ipv4addr=lancfg['ipv4_address'],
ipv4gateway=lancfg['ipv4_gateway'],
ipv4cfgmethod=lancfg['ipv4_configuration'],
hwaddr=lancfg['mac_address']
hwaddr=lancfg['mac_address'],
staticv6addrs=v6cfg['static_addrs'],
staticv6gateway=v6cfg['static_gateway']
))
elif self.op == 'update':
config = self.inputdata.netconfig(self.node)
@@ -602,6 +605,11 @@ class IpmiHandler(object):
ipv4_address=config['ipv4_address'],
ipv4_configuration=config['ipv4_configuration'],
ipv4_gateway=config['ipv4_gateway'])
v6addrs = config.get('static_v6_addresses', None)
if v6addrs is not None:
v6addrs = v6addrs.split(',')
v6gw = config.get('static_v6_gateway', None)
self.ipmicmd.set_net6_configuration(static_addresses=v6addrs, static_gateway=v6gw)
except socket.error as se:
self.output.put(msg.ConfluentNodeError(self.node,
se.message))
@@ -146,7 +146,7 @@ class SshShell(conapi.Console):
return
except cexc.PubkeyInvalid as pi:
self.ssh.close()
self.keyaction = ''
self.keyaction = b''
self.candidatefprint = pi.fingerprint
self.datacallback(pi.message)
self.keyattrname = pi.attrname
@@ -197,18 +197,18 @@ class SshShell(conapi.Console):
delidx = data.index(b'\x7f')
data = data[:delidx - 1] + data[delidx + 1:]
self.keyaction += data
if '\r' in self.keyaction:
action = self.keyaction.split('\r')[0]
if action.lower() == 'accept':
if b'\r' in self.keyaction:
action = self.keyaction.split(b'\r')[0]
if action.lower() == b'accept':
self.nodeconfig.set_node_attributes(
{self.node:
{self.keyattrname: self.candidatefprint}})
self.datacallback('\r\n')
self.logon()
elif action.lower() == 'disconnect':
elif action.lower() == b'disconnect':
self.datacallback(conapi.ConsoleEvent.Disconnect)
else:
self.keyaction = ''
self.keyaction = b''
self.datacallback('\r\nEnter "disconnect" or "accept": ')
elif len(data) > 0:
self.datacallback(data)
+13
View File
@@ -420,6 +420,19 @@ def handle_request(env, start_response):
yield 'complete'
elif env['PATH_INFO'] == '/self/updatestatus' and reqbody:
update = yaml.safe_load(reqbody)
statusstr = update.get('state', None)
statusdetail = update.get('state_detail', None)
didstateupdate = False
if statusstr:
cfg.set_node_attributes({nodename: {'deployment.state': statusstr}})
didstateupdate = True
if statusdetail:
cfg.set_node_attributes({nodename: {'deployment.state_detail': statusdetail}})
didstateupdate = True
if 'status' not in update and didstateupdate:
start_response('200 Ok', ())
yield 'Accepted'
return
if update['status'] == 'staged':
targattr = 'deployment.stagedprofile'
elif update['status'] == 'complete':
+8 -5
View File
@@ -98,11 +98,14 @@ class SyncList(object):
v = v.strip()
if ':' in v:
nr, v = v.split(':', 1)
for candidate in noderange.NodeRange(nr, cfg).nodes:
if candidate == nodename:
break
else:
continue
try:
for candidate in noderange.NodeRange(nr, cfg).nodes:
if candidate == nodename:
break
else:
continue
except Exception as e:
raise Exception('Error on syncfile line "{}": {}'.format(ent, str(e)))
optparts = v.split()
v = optparts[0]
optparts = optparts[1:]
+7 -2
View File
@@ -18,7 +18,11 @@ popd
rm -rf $tdir
cp $tfile rpmlist
cp confluent-genesis.spec confluent-genesis-out.spec
for lic in $(python3 getlicenses.py rpmlist); do
python3 getlicenses.py rpmlist > /tmp/tmpliclist
if [ $? -ne 0 ]; then
exit 1
fi
for lic in $(cat /tmp/tmpliclist); do
lo=${lic#/usr/share/}
lo=${lo#licenses/}
fname=$(basename $lo)
@@ -45,12 +49,13 @@ cp /usr/share/doc/ipmitool/COPYING licenses/ipmitool
echo %license /opt/confluent/genesis/%{arch}/licenses/ipmitool/COPYING >> confluent-genesis-out.spec
cp -f /boot/vmlinuz-$(uname -r) boot/kernel
cp /boot/efi/EFI/BOOT/BOOTX64.EFI boot/efi/boot
cp /boot/efi/EFI/centos/grubx64.efi boot/efi/boot/grubx64.efi
find /boot/efi -name grubx64.efi -exec cp {} boot/efi/boot/grubx64.efi \;
mkdir -p ~/rpmbuild/SOURCES/
tar cf ~/rpmbuild/SOURCES/confluent-genesis.tar boot rpmlist licenses
rpmbuild -bb confluent-genesis-out.spec
rm -rf /usr/lib/dracut/modules.d/97genesis
popd
# for rpm in $(cat ../rpmlist); do dnf download --source $rpm; done
# getting src rpms would be nice, but centos isn't consistent..
# /usr/lib/dracut/skipcpio /opt/confluent/genesis/x86_64/boot/initramfs/distribution | xzcat | cpio -dumiv
# rpm -qf $(find . -type f | sed -e 's/^.//') |sort -u|grep -v 'not owned' > ../rpmlist
+2 -2
View File
@@ -1,6 +1,6 @@
%define arch x86_64
Version: 3.5.0
Release: 4
Version: 3.6.2
Release: 1
Name: confluent-genesis-%{arch}
BuildArch: noarch
Summary: Genesis servicing image for confluent
+59
View File
@@ -0,0 +1,59 @@
import glob
yearsbyname = {}
namesbylicense = {}
filesbylicense = {}
for source in glob.glob('*.c'):
with open(source, 'r') as sourcein:
cap = False
thelicense = ''
currnames = set([])
for line in sourcein.readlines():
if '$OpenBSD$' in line:
continue
if '/*' in line:
cap = True
elif '*/' in line:
cap = False
break
elif cap:
line = line[3:]
if line.startswith('Author: '):
continue
if line.startswith('Copyright'):
_, _, years, name = line.split(maxsplit=3)
name = name.split('>', 1)[0] + '>'
currnames.add(name)
if name not in yearsbyname:
yearsbyname[name] = set([])
yearsbyname[name].add(years)
continue
thelicense += line
if thelicense not in namesbylicense:
namesbylicense[thelicense] = set([])
namesbylicense[thelicense].update(currnames)
if thelicense not in filesbylicense:
filesbylicense[thelicense] = set([])
filesbylicense[thelicense].add(source)
# with open(source + '.license', 'w') as liceout:
# liceout.write(thelicense)
for license in namesbylicense:
for file in sorted(filesbylicense[license]):
print('File: ' + file)
print('')
for author in namesbylicense[license]:
years = []
for year in sorted(yearsbyname[author]):
if not years:
years.append(year)
continue
if int(years[-1].split('-')[-1]) == int(year) - 1:
if '-' in years[-1]:
years[-1] = years[-1].split('-', 1)[0] + '-' + year
else:
years[-1] = years[-1] + '-' + year
else:
years.append(year)
authline = 'Copyright (c) {} {}'.format(','.join(years), author)
print(authline)
print("\n" + license + "\n\n")
+38 -1
View File
@@ -28,6 +28,7 @@ for rpm in allrpmlist:
with open(sys.argv[1]) as rpmlist:
rpmlist = rpmlist.read().split('\n')
licenses = set([])
licensesbyrpm = {}
for rpm in rpmlist:
if not rpm:
continue
@@ -37,11 +38,47 @@ for rpm in rpmlist:
for relrpm in srpmtorpm[srpm]:
liclist = runcmd(f'rpm -qL {relrpm}')
for lic in liclist:
if not lic:
continue
if lic == '(contains no files)':
continue
licensesbyrpm[rpm] = lic
licenses.add(lic)
for lic in sorted(licenses):
print(lic)
manualrpms = [
'ipmitool',
'almalinux-release',
'libaio',
'hwdata',
'snmp',
'libnl3',
'libbpf', # this is covered by kernel
'sqlite', # public domain
'linux-firmware', #all pertinent licenses are stripped out
'xfsprogs', # manually added by hand below (not in rpm)
'tmux', # use the extracttmuxlicenses on the source to generate NOTICE below
]
manuallicenses = [
'/usr/share/doc/ipmitool/COPYING',
'/usr/share/doc/libaio/COPYING',
'/usr/share/doc/net-snmp/COPYING',
'/usr/share/doc/libnl3/COPYING',
'/usr/share/licenses/xfsprogs/GPL-2.0',
'/usr/share/licenses/xfsprogs/LGPL-2.1',
'/usr/share/licenses/tmux/NOTICE',
]
for lic in manuallicenses:
print(lic)
for rpm in rpmlist:
if not rpm:
continue
for manualrpm in manualrpms:
if manualrpm in rpm:
break
else:
if rpm not in licensesbyrpm:
raise Exception('Unresolved license info for ' + rpm)
print("UH OH: " + rpm)
+20
View File
@@ -0,0 +1,20 @@
dnf
hostname
irqbalance
less
sssd-client
NetworkManager
nfs-utils
numactl-libs
passwd
rootfiles
sudo
tuned
yum
initscripts
tpm2-tools
xfsprogs
e2fsprogs
fuse-libs
libnl3
chrony kernel net-tools nfs-utils openssh-server rsync tar util-linux python3 tar dracut dracut-network ethtool parted openssl dhclient openssh-clients bash vim-minimal rpm iputils lvm2 efibootmgr shim-aa64 grub2-efi-aa64 attr
+20
View File
@@ -0,0 +1,20 @@
dnf
hostname
irqbalance
less
sssd-client
NetworkManager
nfs-utils
numactl-libs
passwd
rootfiles
sudo
tuned
yum
initscripts
tpm2-tools
xfsprogs
e2fsprogs
fuse-libs
libnl3
chrony kernel net-tools nfs-utils openssh-server rsync tar util-linux python3 tar dracut dracut-network ethtool parted openssl dhclient openssh-clients bash vim-minimal rpm iputils lvm2 efibootmgr shim-aa64 grub2-efi-aa64 attr
+25 -5
View File
@@ -8,6 +8,7 @@ import glob
import json
import argparse
import os
import platform
import pwd
import re
import shutil
@@ -202,8 +203,12 @@ def capture_remote(args):
os.symlink('/var/lib/confluent/public/site/initramfs.cpio',
os.path.join(outdir, 'boot/initramfs/site.cpio'))
confdir = '/opt/confluent/lib/osdeploy/{}-diskless'.format(oscat)
os.symlink('{}/initramfs/addons.cpio'.format(confdir),
os.path.join(outdir, 'boot/initramfs/addons.cpio'))
archaddon = '/opt/confluent/lib/osdeploy/{}-diskless/initramfs/{}/addons.cpio'.format(oscat, platform.machine())
if os.path.exists(archaddon):
os.symlink(archaddon, os.path.join(outdir, 'boot/initramfs/addons.cpio'))
else:
os.symlink('{}/initramfs/addons.cpio'.format(confdir),
os.path.join(outdir, 'boot/initramfs/addons.cpio'))
indir = '{}/profiles/default'.format(confdir)
if os.path.exists(indir):
copy_tree(indir, outdir)
@@ -671,6 +676,11 @@ def version_sort(iterable):
def get_kern_version(filename):
with open(filename, 'rb') as kernfile:
checkgzip = kernfile.read(2)
if checkgzip == b'\x1f\x8b':
# gzipped... this would probably be aarch64
# assume the filename has the version embedded
return os.path.basename(filename).replace('vmlinuz-', '')
kernfile.seek(0x20e)
offset = struct.unpack('<H', kernfile.read(2))[0] + 0x200
kernfile.seek(offset)
@@ -1276,8 +1286,12 @@ def pack_image(args):
profiley.write('label: {0}\nkernelargs: quiet # confluent_imagemethod=untethered|tethered # tethered is default when unspecified to save on memory, untethered will use more ram, but will not have any ongoing runtime root fs dependency on the http servers.\n'.format(label))
oscat = oshandler.oscategory
confdir = '/opt/confluent/lib/osdeploy/{}-diskless'.format(oscat)
os.symlink('{}/initramfs/addons.cpio'.format(confdir),
os.path.join(outdir, 'boot/initramfs/addons.cpio'))
archaddon = '/opt/confluent/lib/osdeploy/{}-diskless/initramfs/{}/addons.cpio'.format(oscat, platform.machine())
if os.path.exists(archaddon):
os.symlink(archaddon, os.path.join(outdir, 'boot/initramfs/addons.cpio'))
else:
os.symlink('{}/initramfs/addons.cpio'.format(confdir),
os.path.join(outdir, 'boot/initramfs/addons.cpio'))
indir = '{}/profiles/default'.format(confdir)
if os.path.exists(indir):
copy_tree(indir, outdir)
@@ -1298,6 +1312,8 @@ def pack_image(args):
def gather_bootloader(outdir, rootpath='/'):
shimlocation = os.path.join(rootpath, 'boot/efi/EFI/BOOT/BOOTX64.EFI')
if not os.path.exists(shimlocation):
shimlocation = os.path.join(rootpath, 'boot/efi/EFI/BOOT/BOOTAA64.EFI')
if not os.path.exists(shimlocation):
shimlocation = os.path.join(rootpath, 'usr/lib64/efi/shim.efi')
if not os.path.exists(shimlocation):
@@ -1308,7 +1324,11 @@ def gather_bootloader(outdir, rootpath='/'):
for candidate in glob.glob(os.path.join(rootpath, 'boot/efi/EFI/*')):
if 'BOOT' not in candidate:
grubbin = os.path.join(candidate, 'grubx64.efi')
break
if os.path.exists(grubbin):
break
grubbin = os.path.join(candidate, 'grubaa64.efi')
if os.path.exists(grubbin):
break
if not grubbin:
grubbin = os.path.join(rootpath, 'usr/lib64/efi/grub.efi')
if not os.path.exists(grubbin):