mirror of
https://github.com/xcat2/confluent.git
synced 2026-09-29 00:31:09 +00:00
Compare commits
13 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 12bb5d583a | |||
| 03bdbfc8ed | |||
| 661b2ae815 | |||
| ddb8c4cce4 | |||
| 17fff4997b | |||
| 7b3129a1a2 | |||
| d183a3f99c | |||
| b3b3627bf9 | |||
| c1afc144cb | |||
| ac1f7c57b6 | |||
| 19e9c6910d | |||
| e38cd5d3e5 | |||
| f1d3e47439 |
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -587,7 +587,10 @@ 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
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
|
||||
|
||||
@@ -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/'],
|
||||
|
||||
+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