mirror of
https://github.com/xcat2/confluent.git
synced 2026-09-29 08:41:00 +00:00
Compare commits
17 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 12bb5d583a | |||
| 03bdbfc8ed | |||
| 661b2ae815 | |||
| ddb8c4cce4 | |||
| 17fff4997b | |||
| 7b3129a1a2 | |||
| d183a3f99c | |||
| b3b3627bf9 | |||
| c1afc144cb | |||
| ac1f7c57b6 | |||
| 19e9c6910d | |||
| e38cd5d3e5 | |||
| f1d3e47439 | |||
| 9757cd1ae3 | |||
| c9452e65e8 | |||
| cd07e0e212 | |||
| f475d58955 |
@@ -102,9 +102,9 @@ def run():
|
||||
cmdv = ['ssh', sshnode] + cmdvbase + cmdstorun[0]
|
||||
if currprocs < concurrentprocs:
|
||||
currprocs += 1
|
||||
run_cmdv(node, cmdv, all, pipedesc)
|
||||
run_cmdv(sshnode, cmdv, all, pipedesc)
|
||||
else:
|
||||
pendingexecs.append((node, cmdv))
|
||||
pendingexecs.append((sshnode, cmdv))
|
||||
if not all or exitcode:
|
||||
sys.exit(exitcode)
|
||||
rdy, _, _ = select.select(all, [], [], 10)
|
||||
|
||||
@@ -90,17 +90,6 @@ def main(args):
|
||||
if 'error' in rsp:
|
||||
sys.stderr.write(rsp['error'] + '\n')
|
||||
sys.exit(1)
|
||||
if not args.clear and args.network and not args.prepareonly:
|
||||
rc = c.simple_noderange_command(args.noderange, '/boot/nextdevice', 'network',
|
||||
bootmode='uefi',
|
||||
persistent=False,
|
||||
errnodes=errnodes)
|
||||
if errnodes:
|
||||
sys.stderr.write(
|
||||
'Unable to set boot device for following nodes: {0}\n'.format(
|
||||
','.join(errnodes)))
|
||||
return 1
|
||||
rc |= c.simple_noderange_command(args.noderange, '/power/state', 'boot')
|
||||
if args.clear:
|
||||
cleararm(args.noderange, c)
|
||||
clearpending(args.noderange, c)
|
||||
@@ -120,7 +109,7 @@ def main(args):
|
||||
for profname in profnames:
|
||||
sys.stderr.write(' ' + profname + '\n')
|
||||
else:
|
||||
sys.stderr.write('No deployment profiles available, try osdeploy fiimport or imgutil capture\n')
|
||||
sys.stderr.write('No deployment profiles available, try osdeploy import or imgutil capture\n')
|
||||
sys.exit(1)
|
||||
armonce(args.noderange, c)
|
||||
setpending(args.noderange, args.profile, c)
|
||||
@@ -166,6 +155,17 @@ def main(args):
|
||||
else:
|
||||
print('{0}: {1}{2}'.format(node, profile, armed))
|
||||
sys.exit(0)
|
||||
if not args.clear and args.network and not args.prepareonly:
|
||||
rc = c.simple_noderange_command(args.noderange, '/boot/nextdevice', 'network',
|
||||
bootmode='uefi',
|
||||
persistent=False,
|
||||
errnodes=errnodes)
|
||||
if errnodes:
|
||||
sys.stderr.write(
|
||||
'Unable to set boot device for following nodes: {0}\n'.format(
|
||||
','.join(errnodes)))
|
||||
return 1
|
||||
rc |= c.simple_noderange_command(args.noderange, '/power/state', 'boot')
|
||||
if args.network and not args.prepareonly:
|
||||
return rc
|
||||
return 0
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -3,6 +3,8 @@ import os
|
||||
|
||||
class DiskInfo(object):
|
||||
def __init__(self, devname):
|
||||
if devname.startswith('nvme') and 'c' in devname:
|
||||
raise Exception("Skipping multipath devname")
|
||||
self.name = devname
|
||||
self.wwn = None
|
||||
self.path = None
|
||||
|
||||
@@ -72,6 +72,12 @@ def main(args):
|
||||
return rebase(cmdset.profile)
|
||||
ap.print_help()
|
||||
|
||||
def symlinkp(src, trg):
|
||||
try:
|
||||
os.symlink(src, trg)
|
||||
except Exception as e:
|
||||
if e.errno != 17:
|
||||
raise
|
||||
|
||||
def initialize_genesis():
|
||||
if not os.path.exists('/opt/confluent/genesis/x86_64/boot/kernel'):
|
||||
@@ -89,30 +95,33 @@ def initialize_genesis():
|
||||
return retval[1]
|
||||
retcode = 0
|
||||
try:
|
||||
util.mkdirp('/var/lib/confluent', 0o755)
|
||||
if hasconfluentuser:
|
||||
os.chown('/var/lib/confluent', hasconfluentuser.pw_uid, -1)
|
||||
os.setgid(hasconfluentuser.pw_gid)
|
||||
os.setuid(hasconfluentuser.pw_uid)
|
||||
os.umask(0o22)
|
||||
os.makedirs('/var/lib/confluent/public/os/genesis-x86_64/boot/efi/boot', 0o755)
|
||||
os.makedirs('/var/lib/confluent/public/os/genesis-x86_64/boot/initramfs', 0o755)
|
||||
os.symlink('/opt/confluent/genesis/x86_64/boot/efi/boot/BOOTX64.EFI',
|
||||
util.mkdirp('/var/lib/confluent/public/os/genesis-x86_64/boot/efi/boot', 0o755)
|
||||
util.mkdirp('/var/lib/confluent/public/os/genesis-x86_64/boot/initramfs', 0o755)
|
||||
symlinkp('/opt/confluent/genesis/x86_64/boot/efi/boot/BOOTX64.EFI',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/boot/efi/boot/BOOTX64.EFI')
|
||||
os.symlink('/opt/confluent/genesis/x86_64/boot/efi/boot/grubx64.efi',
|
||||
symlinkp('/opt/confluent/genesis/x86_64/boot/efi/boot/grubx64.efi',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/boot/efi/boot/grubx64.efi')
|
||||
os.symlink('/opt/confluent/genesis/x86_64/boot/initramfs/distribution',
|
||||
symlinkp('/opt/confluent/genesis/x86_64/boot/initramfs/distribution',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/boot/initramfs/distribution')
|
||||
os.symlink('/var/lib/confluent/public/site/initramfs.cpio',
|
||||
symlinkp('/var/lib/confluent/public/site/initramfs.cpio',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/boot/initramfs/site.cpio')
|
||||
os.symlink('/opt/confluent/lib/osdeploy/genesis/initramfs/addons.cpio',
|
||||
symlinkp('/opt/confluent/lib/osdeploy/genesis/initramfs/addons.cpio',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/boot/initramfs/addons.cpio')
|
||||
os.symlink('/opt/confluent/genesis/x86_64/boot/kernel',
|
||||
symlinkp('/opt/confluent/genesis/x86_64/boot/kernel',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/boot/kernel')
|
||||
shutil.copytree('/opt/confluent/lib/osdeploy/genesis/profiles/default/ansible/',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/ansible/')
|
||||
shutil.copytree('/opt/confluent/lib/osdeploy/genesis/profiles/default/scripts/',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/scripts/')
|
||||
shutil.copyfile('/opt/confluent/lib/osdeploy/genesis/profiles/default/profile.yaml',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/profile.yaml')
|
||||
if not os.path.exists('/var/lib/confluent/public/os/genesis-x86_64/ansible/'):
|
||||
shutil.copytree('/opt/confluent/lib/osdeploy/genesis/profiles/default/ansible/',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/ansible/')
|
||||
shutil.copytree('/opt/confluent/lib/osdeploy/genesis/profiles/default/scripts/',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/scripts/')
|
||||
shutil.copyfile('/opt/confluent/lib/osdeploy/genesis/profiles/default/profile.yaml',
|
||||
'/var/lib/confluent/public/os/genesis-x86_64/profile.yaml')
|
||||
except Exception as e:
|
||||
sys.stderr.write(str(e) + '\n')
|
||||
retcode = 1
|
||||
@@ -373,9 +382,14 @@ def initialize(cmdset):
|
||||
for rsp in c.read('/uuid'):
|
||||
uuid = rsp.get('uuid', {}).get('value', None)
|
||||
if uuid:
|
||||
with open('confluent_uuid', 'w') as uuidout:
|
||||
uuidout.write(uuid)
|
||||
uuidout.write('\n')
|
||||
oum = os.umask(0o11)
|
||||
try:
|
||||
with open('confluent_uuid', 'w') as uuidout:
|
||||
uuidout.write(uuid)
|
||||
uuidout.write('\n')
|
||||
os.chmod('confluent_uuid', 0o644)
|
||||
finally:
|
||||
os.umask(oum)
|
||||
totar.append('confluent_uuid')
|
||||
topack.append('confluent_uuid')
|
||||
if os.path.exists('ssh'):
|
||||
@@ -403,7 +417,17 @@ def initialize(cmdset):
|
||||
if res:
|
||||
sys.stderr.write('Error occurred while packing site initramfs')
|
||||
sys.exit(1)
|
||||
os.rename(tmpname, '/var/lib/confluent/public/site/initramfs.cpio')
|
||||
oum = os.umask(0o22)
|
||||
try:
|
||||
os.rename(tmpname, '/var/lib/confluent/public/site/initramfs.cpio')
|
||||
os.chmod('/var/lib/confluent/public/site/initramfs.cpio', 0o644)
|
||||
finally:
|
||||
os.umask(oum)
|
||||
oum = os.umask(0o22)
|
||||
try:
|
||||
os.chmod('/var/lib/confluent/public/site/initramfs.cpio', 0o644)
|
||||
finally:
|
||||
os.umask(oum)
|
||||
if cmdset.g:
|
||||
updateboot('genesis-x86_64')
|
||||
if totar:
|
||||
@@ -411,6 +435,11 @@ def initialize(cmdset):
|
||||
tarcmd = ['tar', '-czf', tmptarname] + totar
|
||||
subprocess.check_call(tarcmd)
|
||||
os.rename(tmptarname, '/var/lib/confluent/public/site/initramfs.tgz')
|
||||
oum = os.umask(0o22)
|
||||
try:
|
||||
os.chmod('/var/lib/confluent/public/site/initramfs.tgz', 0o644)
|
||||
finally:
|
||||
os.umask(0o22)
|
||||
os.chdir(opath)
|
||||
print('Site initramfs content packed successfully')
|
||||
|
||||
@@ -421,6 +450,9 @@ def initialize(cmdset):
|
||||
|
||||
|
||||
def updateboot(profilename):
|
||||
if not os.path.exists('/var/lib/confluent/public/site/initramfs.cpio'):
|
||||
emprint('Must generate site content first (TLS (-t) and/or SSH (-s))')
|
||||
return 1
|
||||
c = client.Command()
|
||||
for rsp in c.update('/deployment/profiles/{0}'.format(profilename),
|
||||
{'updateboot': 1}):
|
||||
|
||||
@@ -95,27 +95,29 @@ def assure_tls_ca():
|
||||
os.makedirs(os.path.dirname(fname))
|
||||
except OSError as e:
|
||||
if e.errno != 17:
|
||||
os.seteuid(ouid)
|
||||
raise
|
||||
try:
|
||||
shutil.copy2('/etc/confluent/tls/cacert.pem', fname)
|
||||
hv, _ = util.run(
|
||||
['openssl', 'x509', '-in', '/etc/confluent/tls/cacert.pem', '-hash', '-noout'])
|
||||
if not isinstance(hv, str):
|
||||
hv = hv.decode('utf8')
|
||||
hv = hv.strip()
|
||||
hashname = '/var/lib/confluent/public/site/tls/{0}.0'.format(hv)
|
||||
certname = '{0}.pem'.format(collective.get_myname())
|
||||
for currname in os.listdir('/var/lib/confluent/public/site/tls/'):
|
||||
currname = os.path.join('/var/lib/confluent/public/site/tls/', currname)
|
||||
if currname.endswith('.0'):
|
||||
try:
|
||||
realname = os.readlink(currname)
|
||||
if realname == certname:
|
||||
os.unlink(currname)
|
||||
except OSError:
|
||||
pass
|
||||
os.symlink(certname, hashname)
|
||||
finally:
|
||||
os.seteuid(ouid)
|
||||
shutil.copy2('/etc/confluent/tls/cacert.pem', fname)
|
||||
hv, _ = util.run(
|
||||
['openssl', 'x509', '-in', '/etc/confluent/tls/cacert.pem', '-hash', '-noout'])
|
||||
if not isinstance(hv, str):
|
||||
hv = hv.decode('utf8')
|
||||
hv = hv.strip()
|
||||
hashname = '/var/lib/confluent/public/site/tls/{0}.0'.format(hv)
|
||||
certname = '{0}.pem'.format(collective.get_myname())
|
||||
for currname in os.listdir('/var/lib/confluent/public/site/tls/'):
|
||||
currname = os.path.join('/var/lib/confluent/public/site/tls/', currname)
|
||||
if currname.endswith('.0'):
|
||||
try:
|
||||
realname = os.readlink(currname)
|
||||
if realname == certname:
|
||||
os.unlink(currname)
|
||||
except OSError:
|
||||
pass
|
||||
os.symlink(certname, hashname)
|
||||
|
||||
def substitute_cfg(setting, key, val, newval, cfgfile, line):
|
||||
if key.strip() == setting:
|
||||
|
||||
@@ -49,7 +49,6 @@ _handled_consoles = {}
|
||||
|
||||
_tracelog = None
|
||||
_bufferdaemon = None
|
||||
_bufferlock = None
|
||||
|
||||
try:
|
||||
range = xrange
|
||||
@@ -62,39 +61,38 @@ def chunk_output(output, n):
|
||||
yield output[i:i + n]
|
||||
|
||||
def get_buffer_output(nodename):
|
||||
out = _bufferdaemon.stdin
|
||||
instream = _bufferdaemon.stdout
|
||||
out = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
||||
out.setsockopt(socket.SOL_SOCKET, socket.SO_PASSCRED, 1)
|
||||
out.connect("\x00confluent-vtbuffer")
|
||||
if not isinstance(nodename, bytes):
|
||||
nodename = nodename.encode('utf8')
|
||||
outdata = bytearray()
|
||||
with _bufferlock:
|
||||
out.write(struct.pack('I', len(nodename)))
|
||||
out.write(nodename)
|
||||
out.flush()
|
||||
select.select((instream,), (), (), 30)
|
||||
while not outdata or outdata[-1]:
|
||||
try:
|
||||
chunk = os.read(instream.fileno(), 128)
|
||||
except IOError:
|
||||
chunk = None
|
||||
if chunk:
|
||||
outdata.extend(chunk)
|
||||
else:
|
||||
select.select((instream,), (), (), 0)
|
||||
return bytes(outdata[:-1])
|
||||
out.send(struct.pack('I', len(nodename)))
|
||||
out.send(nodename)
|
||||
select.select((out,), (), (), 30)
|
||||
while not outdata or outdata[-1]:
|
||||
try:
|
||||
chunk = os.read(out.fileno(), 128)
|
||||
except IOError:
|
||||
chunk = None
|
||||
if chunk:
|
||||
outdata.extend(chunk)
|
||||
else:
|
||||
select.select((out,), (), (), 0)
|
||||
return bytes(outdata[:-1])
|
||||
|
||||
|
||||
def send_output(nodename, output):
|
||||
if not isinstance(nodename, bytes):
|
||||
nodename = nodename.encode('utf8')
|
||||
with _bufferlock:
|
||||
_bufferdaemon.stdin.write(struct.pack('I', len(nodename) | (1 << 29)))
|
||||
_bufferdaemon.stdin.write(nodename)
|
||||
_bufferdaemon.stdin.flush()
|
||||
for chunk in chunk_output(output, 8192):
|
||||
_bufferdaemon.stdin.write(struct.pack('I', len(chunk) | (2 << 29)))
|
||||
_bufferdaemon.stdin.write(chunk)
|
||||
_bufferdaemon.stdin.flush()
|
||||
out = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
||||
out.setsockopt(socket.SOL_SOCKET, socket.SO_PASSCRED, 1)
|
||||
out.connect("\x00confluent-vtbuffer")
|
||||
out.send(struct.pack('I', len(nodename) | (1 << 29)))
|
||||
out.send(nodename)
|
||||
for chunk in chunk_output(output, 8192):
|
||||
out.send(struct.pack('I', len(chunk) | (2 << 29)))
|
||||
out.send(chunk)
|
||||
|
||||
def _utf8_normalize(data, decoder):
|
||||
# first we give the stateful decoder a crack at the byte stream,
|
||||
@@ -600,15 +598,10 @@ def _start_tenant_sessions(cfm):
|
||||
def initialize():
|
||||
global _tracelog
|
||||
global _bufferdaemon
|
||||
global _bufferlock
|
||||
_bufferlock = semaphore.Semaphore()
|
||||
_tracelog = log.Logger('trace')
|
||||
_bufferdaemon = subprocess.Popen(
|
||||
['/opt/confluent/bin/vtbufferd'], bufsize=0, stdin=subprocess.PIPE,
|
||||
stdout=subprocess.PIPE)
|
||||
fl = fcntl.fcntl(_bufferdaemon.stdout.fileno(), fcntl.F_GETFL)
|
||||
fcntl.fcntl(_bufferdaemon.stdout.fileno(),
|
||||
fcntl.F_SETFL, fl | os.O_NONBLOCK)
|
||||
['/opt/confluent/bin/vtbufferd', 'confluent-vtbuffer'], bufsize=0, stdin=subprocess.DEVNULL,
|
||||
stdout=subprocess.DEVNULL)
|
||||
|
||||
def start_console_sessions():
|
||||
configmodule.hook_new_configmanagers(_start_tenant_sessions)
|
||||
|
||||
@@ -247,6 +247,10 @@ class NodeHandler(immhandler.NodeHandler):
|
||||
if rsp.status == 200:
|
||||
pwdchanged = True
|
||||
password = newpassword
|
||||
wc.set_header('Authorization', 'Bearer ' + rspdata['access_token'])
|
||||
if '_csrf_token' in wc.cookies:
|
||||
wc.set_header('X-XSRF-TOKEN', wc.cookies['_csrf_token'])
|
||||
wc.grab_json_response_with_status('/api/providers/logout')
|
||||
else:
|
||||
if rspdata.get('locktime', 0) > 0:
|
||||
raise LockedUserException(
|
||||
@@ -280,6 +284,7 @@ class NodeHandler(immhandler.NodeHandler):
|
||||
rsp.read()
|
||||
if rsp.status != 200:
|
||||
return (None, None)
|
||||
wc.grab_json_response_with_status('/api/providers/logout')
|
||||
self._currcreds = (username, newpassword)
|
||||
wc.set_basic_credentials(username, newpassword)
|
||||
pwdchanged = True
|
||||
@@ -434,6 +439,7 @@ class NodeHandler(immhandler.NodeHandler):
|
||||
'/api/function',
|
||||
{'USER_UserModify': '{0},{1},,1,4,0,0,0,0,,8,,,'.format(uid, username)})
|
||||
if status == 200 and rsp.get('return', 0) == 13:
|
||||
wc.grab_json_response('/api/providers/logout')
|
||||
wc.set_basic_credentials(self._currcreds[0], self._currcreds[1])
|
||||
status = 503
|
||||
while status != 200:
|
||||
@@ -442,10 +448,13 @@ class NodeHandler(immhandler.NodeHandler):
|
||||
{'UserName': username}, method='PATCH')
|
||||
if status != 200:
|
||||
rsp = json.loads(rsp)
|
||||
if rsp.get('error', {}).get('code', 'Unknown') in ('Base.1.8.GeneralError', 'Base.1.12.GeneralError'):
|
||||
eventlet.sleep(10)
|
||||
if rsp.get('error', {}).get('code', 'Unknown') in ('Base.1.8.GeneralError', 'Base.1.12.GeneralError', 'Base.1.14.GeneralError'):
|
||||
eventlet.sleep(4)
|
||||
else:
|
||||
break
|
||||
self.tmppasswd = None
|
||||
self._currcreds = (username, passwd)
|
||||
return
|
||||
self.tmppasswd = None
|
||||
wc.grab_json_response('/api/providers/logout')
|
||||
self._currcreds = (username, passwd)
|
||||
@@ -632,3 +641,4 @@ def remote_nodecfg(nodename, cfm):
|
||||
info = {'addresses': [ipaddr]}
|
||||
nh = NodeHandler(info, cfm)
|
||||
nh.config(nodename)
|
||||
|
||||
|
||||
@@ -273,29 +273,6 @@ def opts_to_dict(rq, optidx, expectype=1):
|
||||
def ipfromint(numb):
|
||||
return socket.inet_ntoa(struct.pack('I', numb))
|
||||
|
||||
def get_pxe_bootfile(node, opts, disco, profile, cfg, myip):
|
||||
if opts.get(77, None) == b'iPXE':
|
||||
if not profile:
|
||||
profile = get_deployment_profile(node, cfg)
|
||||
if not profile:
|
||||
log.log({'info': 'No pending profile for {0}, skipping proxyDHCP reply'.format(node)})
|
||||
return None
|
||||
bootfile = 'http://{0}/confluent-public/os/{1}/boot.ipxe'.format(myip, profile).encode('utf8')
|
||||
elif disco['arch'] == 'uefi-x64':
|
||||
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 '
|
||||
'profile name "{1}" is {2} characters longer than is supported '
|
||||
'for this boot method.'.format(
|
||||
node, profile, len(bootfile) - 127)})
|
||||
return None
|
||||
return bootfile
|
||||
|
||||
def proxydhcp(handler, nodeguess):
|
||||
net4011 = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
||||
net4011.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
@@ -369,7 +346,6 @@ def proxydhcp(handler, nodeguess):
|
||||
profile = None
|
||||
if not myipn:
|
||||
myipn = socket.inet_aton(recv)
|
||||
myip = socket.inet_ntoa(myipn)
|
||||
profile = get_deployment_profile(node, cfg)
|
||||
if profile:
|
||||
log.log({
|
||||
@@ -378,9 +354,26 @@ def proxydhcp(handler, nodeguess):
|
||||
if not skiplogging:
|
||||
log.log({'info': 'No pending profile for {0}, skipping proxyDHCP reply'.format(node)})
|
||||
continue
|
||||
myip = socket.inet_ntoa(myipn)
|
||||
bootfile = get_pxe_bootfile(node, opts, disco, profile, cfg, myip)
|
||||
if not bootfile:
|
||||
if opts.get(77, None) == b'iPXE':
|
||||
if not profile:
|
||||
profile = get_deployment_profile(node, cfg)
|
||||
if not profile:
|
||||
log.log({'info': 'No pending profile for {0}, skipping proxyDHCP reply'.format(node)})
|
||||
continue
|
||||
myip = socket.inet_ntoa(myipn)
|
||||
bootfile = 'http://{0}/confluent-public/os/{1}/boot.ipxe'.format(myip, profile).encode('utf8')
|
||||
elif disco['arch'] == 'uefi-x64':
|
||||
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 '
|
||||
'profile name "{1}" is {2} characters longer than is supported '
|
||||
'for this boot method.'.format(
|
||||
node, profile, len(bootfile) - 127)})
|
||||
continue
|
||||
rpv[:240] = rqv[:240].tobytes()
|
||||
rpv[0:1] = b'\x02'
|
||||
@@ -405,7 +398,6 @@ def snoop(handler, protocol=None, nodeguess=None):
|
||||
#prominent
|
||||
#TODO(jjohnson2): enable unicast replies. This would suggest either
|
||||
# injection into the neigh table before OFFER or using SOCK_RAW.
|
||||
global tracelog
|
||||
start_proxydhcp(handler, nodeguess)
|
||||
tracelog = log.Logger('trace')
|
||||
global attribwatcher
|
||||
@@ -507,7 +499,7 @@ def process_dhcp6req(handler, rqv, addr, net, cfg, nodeguess):
|
||||
if ignoredisco.get(mac, 0) + 90 < time.time():
|
||||
ignoredisco[mac] = time.time()
|
||||
handler(info)
|
||||
consider_discover(info, req, net, cfg, None, nodeguess, disco, addr)
|
||||
consider_discover(info, req, net, cfg, None, nodeguess, addr)
|
||||
|
||||
def process_dhcp4req(handler, nodeguess, cfg, net4, idx, recv, rqv):
|
||||
rq = bytearray(rqv)
|
||||
@@ -547,7 +539,7 @@ def process_dhcp4req(handler, nodeguess, cfg, net4, idx, recv, rqv):
|
||||
and time.time() > ignoredisco.get(netaddr, 0) + 90):
|
||||
ignoredisco[netaddr] = time.time()
|
||||
handler(info)
|
||||
consider_discover(info, rqinfo, net4, cfg, rqv, nodeguess, disco)
|
||||
consider_discover(info, rqinfo, net4, cfg, rqv, nodeguess)
|
||||
|
||||
|
||||
|
||||
@@ -595,14 +587,17 @@ def get_deployment_profile(node, cfg, cfd=None):
|
||||
return None
|
||||
candmgrs = cfd.get(node, {}).get('collective.managercandidates', {}).get('value', None)
|
||||
if candmgrs:
|
||||
candmgrs = noderange.NodeRange(candmgrs, cfg).nodes
|
||||
try:
|
||||
candmgrs = noderange.NodeRange(candmgrs, cfg).nodes
|
||||
except Exception: # fallback to unverified noderange
|
||||
candmgrs = noderange.NodeRange(candmgrs).nodes
|
||||
if collective.get_myname() not in candmgrs:
|
||||
return None
|
||||
return profile
|
||||
|
||||
staticassigns = {}
|
||||
myipbypeer = {}
|
||||
def check_reply(node, info, packet, sock, cfg, reqview, addr, disco):
|
||||
def check_reply(node, info, packet, sock, cfg, reqview, addr):
|
||||
httpboot = info['architecture'] == 'uefi-httpboot'
|
||||
cfd = cfg.get_node_attributes(node, ('deployment.*', 'collective.managercandidates'))
|
||||
profile = get_deployment_profile(node, cfg, cfd)
|
||||
@@ -619,7 +614,7 @@ def check_reply(node, info, packet, sock, cfg, reqview, addr, disco):
|
||||
return
|
||||
return reply_dhcp6(node, addr, cfg, packet, cfd, profile, sock)
|
||||
else:
|
||||
return reply_dhcp4(node, info, packet, cfg, reqview, httpboot, cfd, profile, disco)
|
||||
return reply_dhcp4(node, info, packet, cfg, reqview, httpboot, cfd, profile)
|
||||
|
||||
def reply_dhcp6(node, addr, cfg, packet, cfd, profile, sock):
|
||||
myaddrs = netutil.get_my_addresses(addr[-1], socket.AF_INET6)
|
||||
@@ -703,7 +698,7 @@ def get_my_duid():
|
||||
return _myuuid
|
||||
|
||||
|
||||
def reply_dhcp4(node, info, packet, cfg, reqview, httpboot, cfd, profile, disco):
|
||||
def reply_dhcp4(node, info, packet, cfg, reqview, httpboot, cfd, profile):
|
||||
replen = 275 # default is going to be 286
|
||||
# while myipn is describing presumed destination, it's really
|
||||
# vague in the face of aliases, need to convert to ifidx and evaluate
|
||||
@@ -741,7 +736,7 @@ def reply_dhcp4(node, info, packet, cfg, reqview, httpboot, cfd, profile, disco)
|
||||
log.log({'error': nicerr})
|
||||
if niccfg.get('ipv4_broken', False):
|
||||
# Received a request over a nic with no ipv4 configured, ignore it
|
||||
log.log({'error': 'Skipping boot reply to {0} due to no viable IPv4 configuration on deployment system, over nic "{}"'.format(node, info['netinfo']['ifidx'])})
|
||||
log.log({'error': 'Skipping boot reply to {0} due to no viable IPv4 configuration on deployment system'.format(node)})
|
||||
return
|
||||
clipn = None
|
||||
if niccfg['ipv4_method'] == 'firmwarenone':
|
||||
@@ -794,7 +789,6 @@ def reply_dhcp4(node, info, packet, cfg, reqview, httpboot, cfd, profile, disco)
|
||||
repview[249:255] = b'\x33\x04\x00\x00\x00\xf0' # fixed short lease time
|
||||
repview[255:257] = b'\x61\x11'
|
||||
repview[257:274] = packet[97]
|
||||
|
||||
# Note that sending PXEClient kicks off the proxyDHCP procedure, ignoring
|
||||
# boot filename and such in the DHCP packet
|
||||
# we will simply always do it to provide the boot payload in a consistent
|
||||
@@ -803,9 +797,6 @@ def reply_dhcp4(node, info, packet, cfg, reqview, httpboot, cfd, profile, disco)
|
||||
repview[replen - 1:replen + 11] = b'\x3c\x0aHTTPClient'
|
||||
replen += 12
|
||||
else:
|
||||
bootfile = get_pxe_bootfile(node, packet, disco, profile, cfg, myip)
|
||||
if bootfile:
|
||||
repview[108:108 + len(bootfile)] = bootfile
|
||||
repview[replen - 1:replen + 10] = b'\x3c\x09PXEClient'
|
||||
replen += 11
|
||||
hwlen = bytearray(reqview[2:3].tobytes())[0]
|
||||
@@ -897,11 +888,11 @@ def ack_request(pkt, rq, info):
|
||||
repview[26:28] = struct.pack('!H', datasum)
|
||||
send_raw_packet(repview, len(rply), rq, info)
|
||||
|
||||
def consider_discover(info, packet, sock, cfg, reqview, nodeguess, disco, addr=None):
|
||||
def consider_discover(info, packet, sock, cfg, reqview, nodeguess, addr=None):
|
||||
if info.get('hwaddr', None) in macmap and info.get('uuid', None):
|
||||
check_reply(macmap[info['hwaddr']], info, packet, sock, cfg, reqview, addr, disco)
|
||||
check_reply(macmap[info['hwaddr']], info, packet, sock, cfg, reqview, addr)
|
||||
elif info.get('uuid', None) in uuidmap:
|
||||
check_reply(uuidmap[info['uuid']], info, packet, sock, cfg, reqview, addr, disco)
|
||||
check_reply(uuidmap[info['uuid']], info, packet, sock, cfg, reqview, addr)
|
||||
elif packet.get(53, None) == b'\x03':
|
||||
ack_request(packet, reqview, info)
|
||||
elif info.get('uuid', None) and info.get('hwaddr', None):
|
||||
|
||||
@@ -246,11 +246,11 @@ def _find_srvtype(net, net4, srvtype, addresses, xid):
|
||||
try:
|
||||
net4.sendto(data, ('239.255.255.253', 427))
|
||||
except socket.error as se:
|
||||
# On occasion, multicasting may be disabled
|
||||
# tolerate this scenario and move on
|
||||
if se.errno != 101:
|
||||
raise
|
||||
net4.sendto(data, (bcast, 427))
|
||||
pass
|
||||
try:
|
||||
net4.sendto(data, (bcast, 427))
|
||||
except socket.error as se:
|
||||
pass
|
||||
|
||||
|
||||
def _grab_rsps(socks, rsps, interval, xidmap, deferrals):
|
||||
|
||||
@@ -53,7 +53,7 @@ def execupdate(handler, filename, updateobj, type, owner, node, datfile):
|
||||
return
|
||||
if type == 'ffdc' and os.path.isdir(filename):
|
||||
filename += '/' + node
|
||||
if 'type' == 'ffdc':
|
||||
if type == 'ffdc':
|
||||
errstr = False
|
||||
if os.path.exists(filename):
|
||||
errstr = '{0} already exists on {1}, cannot overwrite'.format(
|
||||
|
||||
@@ -381,9 +381,10 @@ def list_info(parms, requestedparameter):
|
||||
break
|
||||
else:
|
||||
candidate = info[requestedparameter]
|
||||
candidate = candidate.strip()
|
||||
if candidate != '':
|
||||
results.add(_api_sanitize_string(candidate))
|
||||
if candidate:
|
||||
candidate = candidate.strip()
|
||||
if candidate != '':
|
||||
results.add(_api_sanitize_string(candidate))
|
||||
return [msg.ChildCollection(x + suffix) for x in util.natural_sort(results)]
|
||||
|
||||
def _handle_neighbor_query(pathcomponents, configmanager):
|
||||
|
||||
@@ -96,6 +96,7 @@ class Bracketer(object):
|
||||
txtnums = getnumbers_nodename(nodename)
|
||||
nums = [int(x) for x in txtnums]
|
||||
for n in range(self.count):
|
||||
# First pass to see if we have exactly one different number
|
||||
padto = len(txtnums[n])
|
||||
needpad = (padto != len('{}'.format(nums[n])))
|
||||
if self.sequences[n] is None:
|
||||
@@ -105,7 +106,24 @@ class Bracketer(object):
|
||||
elif self.sequences[n][2] == nums[n] and self.numlens[n][1] == padto:
|
||||
continue # new nodename has no new number, keep going
|
||||
else: # if self.sequences[n][2] != nums[n] or :
|
||||
if self.diffn is not None and (n != self.diffn or
|
||||
if self.diffn is not None and (n != self.diffn or
|
||||
(padto < self.numlens[n][1]) or
|
||||
(needpad and padto != self.numlens[n][1])):
|
||||
self.flush_current()
|
||||
self.sequences[n] = [[], nums[n], nums[n]]
|
||||
self.numlens[n] = [padto, padto]
|
||||
self.diffn = n
|
||||
for n in range(self.count):
|
||||
padto = len(txtnums[n])
|
||||
needpad = (padto != len('{}'.format(nums[n])))
|
||||
if self.sequences[n] is None:
|
||||
# We initialize to text pieces, 'currstart', and 'prev' number
|
||||
self.sequences[n] = [[], nums[n], nums[n]]
|
||||
self.numlens[n] = [len(txtnums[n]), len(txtnums[n])]
|
||||
elif self.sequences[n][2] == nums[n] and self.numlens[n][1] == padto:
|
||||
continue # new nodename has no new number, keep going
|
||||
else: # if self.sequences[n][2] != nums[n] or :
|
||||
if self.diffn is not None and (n != self.diffn or
|
||||
(padto < self.numlens[n][1]) or
|
||||
(needpad and padto != self.numlens[n][1])):
|
||||
self.flush_current()
|
||||
@@ -384,12 +402,16 @@ class NodeRange(object):
|
||||
def _expandstring(self, element, filternodes=None):
|
||||
prefix = ''
|
||||
if element[0][0] in ('/', '~'):
|
||||
if self.purenumeric:
|
||||
raise Exception('Regular expression not supported within "[]"')
|
||||
element = ''.join(element)
|
||||
nameexpression = element[1:]
|
||||
if self.cfm is None:
|
||||
raise Exception('Verification configmanager required')
|
||||
return set(self.cfm.filter_nodenames(nameexpression, filternodes))
|
||||
elif '=' in element[0] or '!~' in element[0]:
|
||||
if self.purenumeric:
|
||||
raise Exception('Equality/Inequality operators (=, !=, =~, !~) are invalid within "[]"')
|
||||
element = ''.join(element)
|
||||
if self.cfm is None:
|
||||
raise Exception('Verification configmanager required')
|
||||
@@ -449,3 +471,29 @@ class NodeRange(object):
|
||||
if self.cfm is None:
|
||||
return set([element])
|
||||
raise Exception(element + ' not a recognized node, group, or alias')
|
||||
|
||||
if __name__ == '__main__':
|
||||
cases = [
|
||||
(['r3u4', 'r5u6'], 'r3u4,r5u6'), # should not erroneously gather
|
||||
(['r3u4s1', 'r5u6s3'], 'r3u4s1,r5u6s3'), # should not erroneously gather
|
||||
(['r3u4s1', 'r3u4s2', 'r5u4s3'], 'r3u4s[1:2],r5u4s3'), # should not erroneously gather
|
||||
(['r3u4', 'r3u5', 'r3u6', 'r3u9', 'r4u1'], 'r3u[4:6,9],r4u1'),
|
||||
(['n01', 'n2', 'n03'], 'n01,n2,n03'),
|
||||
(['n7', 'n8', 'n09', 'n10', 'n11', 'n12', 'n13', 'n14', 'n15', 'n16',
|
||||
'n17', 'n18', 'n19', 'n20'], 'n[7:8],n[09:20]')
|
||||
]
|
||||
for case in cases:
|
||||
gc = case[0]
|
||||
bracketer = Bracketer(gc[0])
|
||||
for chnk in gc[1:]:
|
||||
bracketer.extend(chnk)
|
||||
br = bracketer.range
|
||||
resnodes = NodeRange(br).nodes
|
||||
if set(resnodes) != set(gc):
|
||||
print('FAILED: ' + repr(sorted(gc)))
|
||||
print('RESULT: ' + repr(sorted(resnodes)))
|
||||
print('EXPECTED: ' + repr(case[1]))
|
||||
print('ACTUAL: ' + br)
|
||||
|
||||
|
||||
|
||||
|
||||
@@ -98,14 +98,15 @@ def initialize_ca():
|
||||
preexec_fn=normalize_uid)
|
||||
ouid = normalize_uid()
|
||||
try:
|
||||
os.makedirs('/var/lib/confluent/public/site/ssh/', mode=0o755)
|
||||
except OSError as e:
|
||||
if e.errno != 17:
|
||||
raise
|
||||
try:
|
||||
os.makedirs('/var/lib/confluent/public/site/ssh/', mode=0o755)
|
||||
except OSError as e:
|
||||
if e.errno != 17:
|
||||
raise
|
||||
cafilename = '/var/lib/confluent/public/site/ssh/{0}.ca'.format(myname)
|
||||
shutil.copy('/etc/confluent/ssh/ca.pub', cafilename)
|
||||
finally:
|
||||
os.seteuid(ouid)
|
||||
cafilename = '/var/lib/confluent/public/site/ssh/{0}.ca'.format(myname)
|
||||
shutil.copy('/etc/confluent/ssh/ca.pub', cafilename)
|
||||
# newent = '@cert-authority * ' + capub.read()
|
||||
|
||||
|
||||
@@ -185,6 +186,14 @@ def initialize_root_key(generate, automation=False):
|
||||
if os.path.exists('/etc/confluent/ssh/automation'):
|
||||
alreadyexist = True
|
||||
else:
|
||||
ouid = normalize_uid()
|
||||
try:
|
||||
os.makedirs('/etc/confluent/ssh', mode=0o700)
|
||||
except OSError as e:
|
||||
if e.errno != 17:
|
||||
raise
|
||||
finally:
|
||||
os.seteuid(ouid)
|
||||
subprocess.check_call(
|
||||
['ssh-keygen', '-t', 'ed25519',
|
||||
'-f','/etc/confluent/ssh/automation', '-N', get_passphrase(),
|
||||
|
||||
@@ -29,9 +29,9 @@ import struct
|
||||
import eventlet.green.subprocess as subprocess
|
||||
|
||||
|
||||
def mkdirp(path):
|
||||
def mkdirp(path, mode=0o777):
|
||||
try:
|
||||
os.makedirs(path)
|
||||
os.makedirs(path, mode)
|
||||
except OSError as e:
|
||||
if e.errno != 17:
|
||||
raise
|
||||
|
||||
@@ -19,6 +19,7 @@ setup(
|
||||
'confluent/plugins/hardwaremanagement/',
|
||||
'confluent/plugins/deployment/',
|
||||
'confluent/plugins/console/',
|
||||
'confluent/plugins/info/',
|
||||
'confluent/plugins/shell/',
|
||||
'confluent/collective/',
|
||||
'confluent/plugins/configuration/'],
|
||||
|
||||
@@ -22,3 +22,16 @@ modification, are permitted provided that the following conditions are met:
|
||||
* Neither the name of the copyright holder nor the
|
||||
names of contributors may be used to endorse or promote products
|
||||
derived from this software without specific prior written permission.
|
||||
|
||||
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDER AND CONTRIBUTORS
|
||||
* "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
|
||||
* LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
|
||||
* A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE AUTHORS,
|
||||
* COPYRIGHT HOLDERS, OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT,
|
||||
* INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING,
|
||||
* BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF
|
||||
* USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON
|
||||
* ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
|
||||
* (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
|
||||
* OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
|
||||
|
||||
|
||||
+135
-44
@@ -1,8 +1,14 @@
|
||||
#include <asm-generic/socket.h>
|
||||
#define _GNU_SOURCE
|
||||
#include <stdio.h>
|
||||
#include <string.h>
|
||||
#include <stdlib.h>
|
||||
#include <locale.h>
|
||||
#include <unistd.h>
|
||||
#include <sys/socket.h>
|
||||
#include <sys/epoll.h>
|
||||
#include <sys/un.h>
|
||||
#include <fcntl.h>
|
||||
#include "tmt.h"
|
||||
#define HASHSIZE 2053
|
||||
#define MAXNAMELEN 256
|
||||
@@ -10,13 +16,17 @@
|
||||
struct terment {
|
||||
struct terment *next;
|
||||
char *name;
|
||||
int fd;
|
||||
TMT *vt;
|
||||
};
|
||||
|
||||
#define SETNODE 1
|
||||
#define WRITE 2
|
||||
#define READBUFF 0
|
||||
#define CLOSECONN 3
|
||||
#define MAXEVTS 16
|
||||
static struct terment *buffers[HASHSIZE];
|
||||
static char* nodenames[HASHSIZE];
|
||||
|
||||
unsigned long hash(char *str)
|
||||
/* djb2a */
|
||||
@@ -37,10 +47,13 @@ TMT *get_termentbyname(char *name) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
TMT *set_termentbyname(char *name) {
|
||||
TMT *set_termentbyname(char *name, int fd) {
|
||||
struct terment *ret;
|
||||
int idx;
|
||||
|
||||
if (nodenames[fd] == NULL) {
|
||||
nodenames[fd] = strdup(name);
|
||||
}
|
||||
idx = hash(name);
|
||||
for (ret = buffers[idx]; ret != NULL; ret = ret->next)
|
||||
if (strcmp(name, ret->name) == 0)
|
||||
@@ -48,12 +61,13 @@ TMT *set_termentbyname(char *name) {
|
||||
ret = (struct terment *)malloc(sizeof(*ret));
|
||||
ret->next = buffers[idx];
|
||||
ret->name = strdup(name);
|
||||
ret->fd = fd;
|
||||
ret->vt = tmt_open(31, 100, NULL, NULL, L"→←↑↓■◆▒°±▒┘┐┌└┼⎺───⎽├┤┴┬│≤≥π≠£•");
|
||||
buffers[idx] = ret;
|
||||
return ret->vt;
|
||||
}
|
||||
|
||||
void dump_vt(TMT* outvt) {
|
||||
void dump_vt(TMT* outvt, int outfd) {
|
||||
const TMTSCREEN *out = tmt_screen(outvt);
|
||||
const TMTPOINT *curs = tmt_cursor(outvt);
|
||||
int line, idx, maxcol, maxrow;
|
||||
@@ -67,9 +81,10 @@ void dump_vt(TMT* outvt) {
|
||||
tmt_color_t fg = TMT_COLOR_DEFAULT;
|
||||
tmt_color_t bg = TMT_COLOR_DEFAULT;
|
||||
wchar_t sgrline[30];
|
||||
char strbuffer[128];
|
||||
size_t srgidx = 0;
|
||||
char colorcode = 0;
|
||||
wprintf(L"\033c");
|
||||
write(outfd, "\033c", 2);
|
||||
maxcol = 0;
|
||||
maxrow = 0;
|
||||
for (line = out->nline - 1; line >= 0; --line) {
|
||||
@@ -148,60 +163,136 @@ void dump_vt(TMT* outvt) {
|
||||
}
|
||||
if (sgrline[0] != 0) {
|
||||
sgrline[wcslen(sgrline) - 1] = 0; // Trim last ;
|
||||
wprintf(L"\033[%lsm", sgrline);
|
||||
|
||||
snprintf(strbuffer, sizeof(strbuffer), "\033[%lsm", sgrline);
|
||||
write(outfd, strbuffer, strlen(strbuffer));
|
||||
write(outfd, "\033[]", 3);
|
||||
}
|
||||
wprintf(L"%lc", out->lines[line]->chars[idx].c);
|
||||
snprintf(strbuffer, sizeof(strbuffer), "%lc", out->lines[line]->chars[idx].c);
|
||||
write(outfd, strbuffer, strlen(strbuffer));
|
||||
}
|
||||
if (line < maxrow)
|
||||
wprintf(L"\r\n");
|
||||
write(outfd, "\r\n", 2);
|
||||
}
|
||||
fflush(stdout);
|
||||
wprintf(L"\x1b[%ld;%ldH", curs->r + 1, curs->c + 1);
|
||||
fflush(stdout);
|
||||
//fflush(stdout);
|
||||
snprintf(strbuffer, sizeof(strbuffer), "\x1b[%ld;%ldH", curs->r + 1, curs->c + 1);
|
||||
write(outfd, strbuffer, strlen(strbuffer));
|
||||
//fflush(stdout);
|
||||
}
|
||||
|
||||
int handle_traffic(int fd) {
|
||||
int cmd, length;
|
||||
char currnode[MAXNAMELEN];
|
||||
char cmdbuf[MAXDATALEN];
|
||||
char *nodename;
|
||||
TMT *currvt = NULL;
|
||||
TMT *outvt = NULL;
|
||||
length = read(fd, &cmd, 4);
|
||||
if (length <= 0) {
|
||||
return 0;
|
||||
}
|
||||
length = cmd & 536870911;
|
||||
cmd = cmd >> 29;
|
||||
if (cmd == SETNODE) {
|
||||
cmd = read(fd, currnode, length);
|
||||
currnode[length] = 0;
|
||||
if (cmd < 0)
|
||||
return 0;
|
||||
currvt = set_termentbyname(currnode, fd);
|
||||
} else if (cmd == WRITE) {
|
||||
if (currvt == NULL) {
|
||||
nodename = nodenames[fd];
|
||||
currvt = set_termentbyname(nodename, fd);
|
||||
}
|
||||
cmd = read(fd, cmdbuf, length);
|
||||
cmdbuf[length] = 0;
|
||||
if (cmd < 0)
|
||||
return 0;
|
||||
tmt_write(currvt, cmdbuf, length);
|
||||
} else if (cmd == READBUFF) {
|
||||
cmd = read(fd, cmdbuf, length);
|
||||
cmdbuf[length] = 0;
|
||||
if (cmd < 0)
|
||||
return 0;
|
||||
outvt = get_termentbyname(cmdbuf);
|
||||
if (outvt != NULL)
|
||||
dump_vt(outvt, fd);
|
||||
length = write(fd, "\x00", 1);
|
||||
if (length < 0)
|
||||
return 0;
|
||||
} else if (cmd == CLOSECONN) {
|
||||
return 0;
|
||||
}
|
||||
return 1;
|
||||
}
|
||||
|
||||
int main(int argc, char* argv[]) {
|
||||
int cmd, length;
|
||||
setlocale(LC_ALL, "");
|
||||
char cmdbuf[MAXDATALEN];
|
||||
char currnode[MAXNAMELEN];
|
||||
TMT *currvt = NULL;
|
||||
TMT *outvt = NULL;
|
||||
struct sockaddr_un addr;
|
||||
int numevts;
|
||||
int status;
|
||||
int poller;
|
||||
int n;
|
||||
socklen_t len;
|
||||
int ctlsock, currsock;
|
||||
socklen_t addrlen;
|
||||
struct ucred ucr;
|
||||
|
||||
struct epoll_event epvt, evts[MAXEVTS];
|
||||
stdin = freopen(NULL, "rb", stdin);
|
||||
if (stdin == NULL) {
|
||||
exit(1);
|
||||
}
|
||||
memset(&addr, 0, sizeof(struct sockaddr_un));
|
||||
addr.sun_family = AF_UNIX;
|
||||
strncpy(addr.sun_path + 1, argv[1], sizeof(addr.sun_path) - 2); // abstract namespace socket
|
||||
ctlsock = socket(AF_UNIX, SOCK_STREAM, 0);
|
||||
status = bind(ctlsock, (const struct sockaddr*)&addr, sizeof(sa_family_t) + strlen(argv[1]) + 1); //sizeof(struct sockaddr_un));
|
||||
if (status < 0) {
|
||||
perror("Unable to open unix socket - ");
|
||||
exit(1);
|
||||
}
|
||||
listen(ctlsock, 128);
|
||||
poller = epoll_create(1);
|
||||
memset(&epvt, 0, sizeof(struct epoll_event));
|
||||
epvt.events = EPOLLIN;
|
||||
epvt.data.fd = ctlsock;
|
||||
if (epoll_ctl(poller, EPOLL_CTL_ADD, ctlsock, &epvt) < 0) {
|
||||
perror("Unable to poll the socket");
|
||||
exit(1);
|
||||
}
|
||||
// create a unix domain socket for accepting, each connection is only allowed to either read or write, not both
|
||||
while (1) {
|
||||
length = fread(&cmd, 4, 1, stdin);
|
||||
if (length < 0)
|
||||
continue;
|
||||
length = cmd & 536870911;
|
||||
cmd = cmd >> 29;
|
||||
if (cmd == SETNODE) {
|
||||
cmd = fread(currnode, 1, length, stdin);
|
||||
currnode[length] = 0;
|
||||
if (cmd < 0)
|
||||
continue;
|
||||
currvt = set_termentbyname(currnode);
|
||||
} else if (cmd == WRITE) {
|
||||
if (currvt == NULL)
|
||||
currvt = set_termentbyname("");
|
||||
cmd = fread(cmdbuf, 1, length, stdin);
|
||||
cmdbuf[length] = 0;
|
||||
if (cmd < 0)
|
||||
continue;
|
||||
tmt_write(currvt, cmdbuf, length);
|
||||
} else if (cmd == READBUFF) {
|
||||
cmd = fread(cmdbuf, 1, length, stdin);
|
||||
cmdbuf[length] = 0;
|
||||
if (cmd < 0)
|
||||
continue;
|
||||
outvt = get_termentbyname(cmdbuf);
|
||||
if (outvt != NULL)
|
||||
dump_vt(outvt);
|
||||
length = write(1, "\x00", 1);
|
||||
if (length < 0)
|
||||
continue;
|
||||
numevts = epoll_wait(poller, evts, MAXEVTS, -1);
|
||||
if (numevts < 0) {
|
||||
perror("Failed wait");
|
||||
exit(1);
|
||||
}
|
||||
for (n = 0; n < numevts; ++n) {
|
||||
if (evts[n].data.fd == ctlsock) {
|
||||
currsock = accept(ctlsock, (struct sockaddr *) &addr, &addrlen);
|
||||
len = sizeof(ucr);
|
||||
getsockopt(currsock, SOL_SOCKET, SO_PEERCRED, &ucr, &len);
|
||||
if (ucr.uid != getuid()) { // block access for other users
|
||||
close(currsock);
|
||||
continue;
|
||||
}
|
||||
memset(&epvt, 0, sizeof(struct epoll_event));
|
||||
epvt.events = EPOLLIN;
|
||||
epvt.data.fd = currsock;
|
||||
epoll_ctl(poller, EPOLL_CTL_ADD, currsock, &epvt);
|
||||
} else {
|
||||
if (!handle_traffic(evts[n].data.fd)) {
|
||||
epoll_ctl(poller, EPOLL_CTL_DEL, evts[n].data.fd, NULL);
|
||||
close(evts[n].data.fd);
|
||||
if (nodenames[evts[n].data.fd] != NULL) {
|
||||
free(nodenames[evts[n].data.fd]);
|
||||
nodenames[evts[n].data.fd] = NULL;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user