Просмотр исходного кода

dnet: add lilith support to dnet (mostly)

* calls spawns() rpc method to retrieve lilith info
* add data to model
* renders lilith name and name of spawned network

still TODO:

* display hosts and URLs when spawned network is selected
lunar-mining 2 лет назад
Родитель
Сommit
887ebff1f1
4 измененных файлов с 240 добавлено и 177 удалено
  1. 45 29
      bin/dnet/main.py
  2. 93 70
      bin/dnet/model.py
  3. 0 2
      bin/dnet/rpc.py
  4. 102 76
      bin/dnet/view.py

+ 45 - 29
bin/dnet/main.py

@@ -32,8 +32,13 @@ class Dnetview:
         self.model = Model()
         self.view = View(self.model)
 
-    async def subscribe(self, rpc, name, host, port):
+    async def subscribe(self, rpc, node):
+        name = node['name']
+        host = node['host']
+        port = node['port']
+        type = node['type']
         info = {}
+
         while True:
             try:
                 await rpc.start(host, port)
@@ -44,27 +49,34 @@ class Dnetview:
                 await self.queue.put(info)
                 continue
     
-        data = await rpc._make_request("p2p.get_info", [])
-        info[name] = data
-
-        await self.queue.put(info)
-        await rpc.dnet_switch(True)
-        await rpc.dnet_subscribe_events()
-
-        while True:
-            await asyncio.sleep(0.01)
-            data = await rpc.reader.readline()
-            try:
-                data = json.loads(data)
-                info[name] = data
-                await self.queue.put(info)
-            except:
-                info[name] = {}
-                await self.queue.put(info)
+        if type == "NORMAL":
+            data = await rpc._make_request("p2p.get_info", [])
+            info[name] = data
+
+            await self.queue.put(info)
+            await rpc.dnet_switch(True)
+            await rpc.dnet_subscribe_events()
+            
+            while True:
+                await asyncio.sleep(0.01)
+                data = await rpc.reader.readline()
+                try:
+                    data = json.loads(data)
+                    info[name] = data
+                    await self.queue.put(info)
+                except:
+                    info[name] = {}
+                    await self.queue.put(info)
+
+            await rpc.dnet_switch(False)
+    
+        if type == "LILITH":
+            data = await rpc._make_request("spawns", [])
+            info[name] = data
+            await self.queue.put(info)
 
-        await rpc.dnet_switch(False)
         await rpc.stop()
-    
+
     def get_config(self):
         with open("config.toml") as f:
             cfg = toml.load(f)
@@ -76,8 +88,7 @@ class Dnetview:
             for i, node in enumerate(nodes):
                 rpc = JsonRpc()
                 subscribe = tg.create_task(self.subscribe(
-                            rpc, node['name'], node['host'],
-                            node['port']))
+                            rpc, node))
                 nodes = tg.create_task(self.update_info())
 
     async def update_info(self):
@@ -85,15 +96,20 @@ class Dnetview:
             info = await self.queue.get()
             values = list(info.values())[0]
 
-            if values:
-                method = values.get("method")
-                if method == "dnet.subscribe_events":
-                    self.model.add_event(info)
-                else:
-                    self.model.add_node(info)
-            else:
+            if not values:
                 self.model.add_offline(info)
 
+            if 'result' in values:
+                result = values.get('result')
+                if 'spawns' in result:
+                    self.model.add_lilith(info)
+                if 'channels' in result:
+                    self.model.add_node(info)
+
+            if 'params' in values:
+                self.model.add_event(info)
+
+            logging.debug(self.model.liliths)
             self.queue.task_done()
 
     def main(self):

+ 93 - 70
bin/dnet/model.py

@@ -24,13 +24,14 @@ class Model:
 
     def __init__(self):
         self.nodes = {}
+        self.liliths = {}
 
     def add_node(self, node):
         channel_lookup = {}
         name = list(node.keys())[0]
         values = list(node.values())[0]
-        info = values["result"]
-        channels = info["channels"]
+        info = values['result']
+        channels = info['channels']
         
         self.nodes[name] = {}
         self.nodes[name]['outbound'] = {}
@@ -41,37 +42,37 @@ class Model:
         self.nodes[name]['msgs'] = dd(list)
 
         for channel in channels:
-            id = channel["id"]
+            id = channel['id']
             channel_lookup[id] = channel
 
         for channel in channels:
-            if channel["session"] != "inbound":
+            if channel['session'] != 'inbound':
                 continue
-            id = channel["id"]
-            url = channel_lookup[id]["url"]
-            self.nodes[name]['inbound'][f"{id}"] = url
+            id = channel['id']
+            url = channel_lookup[id]['url']
+            self.nodes[name]['inbound'][f'{id}'] = url
 
-        for i, id in enumerate(info["outbound_slots"]):
+        for i, id in enumerate(info['outbound_slots']):
             if id == 0:
-                outbounds = self.nodes[name]['outbound'][f"{i}"] = ["none", 0]
+                outbounds = self.nodes[name]['outbound'][f'{i}'] = ['none', 0]
                 continue
             assert id in channel_lookup
-            url = channel_lookup[id]["url"]
-            outbounds = self.nodes[name]['outbound'][f"{i}"] = [url, id]
+            url = channel_lookup[id]['url']
+            outbounds = self.nodes[name]['outbound'][f'{i}'] = [url, id]
 
         for channel in channels:
-            if channel["session"] != "seed":
+            if channel['session'] != 'seed':
                 continue
-            id = channel["id"]
-            url = channel["url"]
-            self.nodes[name]['seed'][f"{id}"] = url
+            id = channel['id']
+            url = channel['url']
+            self.nodes[name]['seed'][f'{id}'] = url
 
         for channel in channels:
-            if channel["session"] != "manual":
+            if channel['session'] != 'manual':
                 continue
-            id = channel["id"]
-            url = channel["url"]
-            self.nodes[name]['manual'][f"{id}"] = url
+            id = channel['id']
+            url = channel['url']
+            self.nodes[name]['manual'][f'{id}'] = url
     
     def add_offline(self, node):
         name = list(node.keys())[0]
@@ -81,76 +82,98 @@ class Model:
     def add_event(self, event):
         name = list(event.keys())[0]
         values = list(event.values())[0]
-        params = values.get("params")
-        event = params[0].get("event")
-        info = params[0].get("info")
+        params = values.get('params')
+        event = params[0].get('event')
+        info = params[0].get('info')
 
         t = time.localtime()
-        current_time = time.strftime("%H:%M:%S", t)
+        current_time = time.strftime('%H:%M:%S', t)
 
         match event:                        
-            case "send":
-                nano = info.get("time")
-                cmd = info.get("cmd")
-                chan = info.get("chan")
-                addr = chan.get("addr")
+            case 'send':
+                nano = info.get('time')
+                cmd = info.get('cmd')
+                chan = info.get('chan')
+                addr = chan.get('addr')
                 t = (dt.datetime
                         .fromtimestamp(int(nano)/1000000000)
                         .strftime('%H:%M:%S'))
                 msgs = self.nodes[name]['msgs']
                 msgs[addr].append((t, event, cmd))
-            case "recv":
-                nano = info.get("time")
-                cmd = info.get("cmd")
-                chan = info.get("chan")
-                addr = chan.get("addr")
+            case 'recv':
+                nano = info.get('time')
+                cmd = info.get('cmd')
+                chan = info.get('chan')
+                addr = chan.get('addr')
                 t = (dt.datetime
                         .fromtimestamp(int(nano)/1000000000)
                         .strftime('%H:%M:%S'))
                 msgs = self.nodes[name]['msgs']
                 msgs[addr].append((t, event, cmd))
-            case "inbound_connected":
-                addr = info["addr"]
-                id = info.get("channel_id")
-                self.nodes[name]['inbound'][f"{id}"] = addr
-                logging.debug(f"{current_time}  inbound (connect):    {addr}")
-            case "inbound_disconnected":
-                addr = info["addr"]
-                id = info.get("channel_id")
+            case 'inbound_connected':
+                addr = info['addr']
+                id = info.get('channel_id')
+                self.nodes[name]['inbound'][f'{id}'] = addr
+                logging.debug(f'{current_time}  inbound (connect):    {addr}')
+            case 'inbound_disconnected':
+                addr = info['addr']
+                id = info.get('channel_id')
                 inbound = self.nodes[name]['inbound']
-                self.nodes[name]['inbound'][f"{id}"] = {}
-                logging.debug(f"{current_time}  inbound (disconnect): {addr}")
-            case "outbound_slot_sleeping":
-                slot = info["slot"]
+                self.nodes[name]['inbound'][f'{id}'] = {}
+                logging.debug(f'{current_time}  inbound (disconnect): {addr}')
+            case 'outbound_slot_sleeping':
+                slot = info['slot']
                 event = self.nodes[name]['event']
-                event[(f"{name}", f"{slot}")] = ["sleeping", 0]
-                logging.debug(f"{current_time}  slot {slot}: sleeping")
-            case "outbound_slot_connecting":
-                slot = info["slot"]
-                addr = info["addr"]
+                event[(f'{name}', f'{slot}')] = ['sleeping', 0]
+                logging.debug(f'{current_time}  slot {slot}: sleeping')
+            case 'outbound_slot_connecting':
+                slot = info['slot']
+                addr = info['addr']
                 event = self.nodes[name]['event']
-                event[(f"{name}", f"{slot}")] = [f"connecting: addr={addr}", 0]
-                logging.debug(f"{current_time}  slot {slot}: connecting   addr={addr}")
-            case "outbound_slot_connected":
-                slot = info["slot"]
-                addr = info["addr"]
-                id = info["channel_id"]
-                self.nodes[name]['outbound'][f"{slot}"] = [addr, id]
-                logging.debug(f"{current_time}  slot {slot}: connected    addr={addr}")
-            case "outbound_slot_disconnected":
-                slot = info["slot"]
-                err = info["err"]
+                event[(f'{name}', f'{slot}')] = [f'connecting: addr={addr}', 0]
+                logging.debug(f'{current_time}  slot {slot}: connecting   addr={addr}')
+            case 'outbound_slot_connected':
+                slot = info['slot']
+                addr = info['addr']
+                id = info['channel_id']
+                self.nodes[name]['outbound'][f'{slot}'] = [addr, id]
+                logging.debug(f'{current_time}  slot {slot}: connected    addr={addr}')
+            case 'outbound_slot_disconnected':
+                slot = info['slot']
+                err = info['err']
                 event = self.nodes[name]['event']
-                event[(f"{name}", f"{slot}")] = [f"disconnected: {err}", 0]
-                logging.debug(f"{current_time}  slot {slot}: disconnected err='{err}'")
-            case "outbound_peer_discovery":
-                attempt = info["attempt"]
-                state = info["state"]
+                event[(f'{name}', f'{slot}')] = [f'disconnected: {err}', 0]
+                logging.debug(f'{current_time}  slot {slot}: disconnected err={err}')
+            case 'outbound_peer_discovery':
+                attempt = info['attempt']
+                state = info['state']
                 event = self.nodes[name]['event']
-                key = (f"{name}", "outbound")
-                event[key] = f"peer discovery: {state} (attempt {attempt})"
-                logging.debug(f"{current_time}  peer_discovery: {state} (attempt {attempt})")
+                key = (f'{name}', 'outbound')
+                event[key] = f'peer discovery: {state} (attempt {attempt})'
+                logging.debug(f'{current_time}  peer_discovery: {state} (attempt {attempt})')
 
 
+    def add_lilith(self, lilith):
+        key = list(lilith.keys())[0]
+        values = list(lilith.values())[0]
+        info = values['result']
+        spawns = info['spawns']
+
+        self.liliths[key] = {}
+        self.liliths[key]['name'] = {}
+        self.liliths[key]['hosts'] = {}
+        self.liliths[key]['urls'] = {}
+        
+        for spawn in spawns:
+            name = spawn['name']
+            urls = spawn['urls']
+            hosts = spawn['hosts']
+            
+            self.liliths[key]['name'] = name
+            self.liliths[key]['urls'] = urls
+            self.liliths[key]['hosts'] = hosts
+        
+
     def __repr__(self):
-        return f"{self.nodes}"
+        return f'{self.nodes}'
+        return f'{self.liliths}'

+ 0 - 2
bin/dnet/rpc.py

@@ -73,5 +73,3 @@ class JsonRpc:
 
     async def dnet_subscribe_events(self):
         return await self._subscribe("dnet.subscribe_events", [])
-
-

+ 102 - 76
bin/dnet/view.py

@@ -56,6 +56,9 @@ class Session(DnetWidget):
         txt = urwid.Text(f"  {self.session}")
         super().update(txt)
 
+    def is_lilith(self):
+        return True
+
 
 class Slot(DnetWidget):
     def set_txt(self, i, addr):
@@ -68,7 +71,7 @@ class Slot(DnetWidget):
             self.addr = addr
             txt = urwid.Text(f"    {self.addr}")
         super().update(txt)
-
+    
 
 class View():
     palette = [
@@ -88,6 +91,13 @@ class View():
         leftbox = urwid.LineBox(self.list)
         columns = urwid.Columns([leftbox, rightbox], focus_column=0)
         self.ui = urwid.Frame(urwid.AttrWrap( columns, 'body' ))
+        self.known_outbound = []
+        self.known_inbound = []
+        self.known_nodes = []
+        self.live_nodes = []
+        self.dead_nodes = []
+        self.refresh = False
+
 
     #-----------------------------------------------------------------
     # Render get_info()
@@ -135,6 +145,13 @@ class View():
                slot.set_txt(i, addr)
                self.listwalker.contents.append(slot)
 
+       if 'name' in info and info['name']:
+           spawn = info.get('name')
+           session = Session(node_name, spawn)
+           session.is_lilith()
+           session.set_txt()  
+           self.listwalker.contents.append(session)
+
     def draw_empty(self, node_name, info):
        node = Node(node_name, "node")
        node.set_txt(True)
@@ -189,94 +206,103 @@ class View():
                     self.pile.contents.append((urwid.Text(
                             f"{time}: {event}: {msg}"),
                             self.pile.options()))
-
+        
+    # Sort nodes into lists.
+    def sort(self, nodes):
+        for name, info in nodes:
+            if bool(info) and name not in self.live_nodes:
+                self.live_nodes.append(name)
+            if not bool(info) and name not in self.dead_nodes:
+                self.dead_nodes.append(name)
+            if bool(info) and name in self.dead_nodes:
+                logging.debug("Refresh: dead node online.")
+                self.refresh = True
+            if not bool(info) and name in self.live_nodes:
+                logging.debug("Refresh: online node offline.")
+                self.refresh = True
+
+    # Display nodes according to list.
+    async def draw(self, nodes):
+        for name, info in nodes:
+            if name in self.live_nodes and name not in self.known_nodes:
+                self.draw_info(name, info)
+            if name in self.dead_nodes and name not in self.known_nodes:
+                self.draw_empty(name, info)
+            if self.refresh:
+                logging.debug("Refresh initiated.")
+                await asyncio.sleep(0.1)
+                self.known_outbound.clear()
+                self.known_inbound.clear()
+                self.known_nodes.clear()
+                self.live_nodes.clear()
+                self.dead_nodes.clear()
+                refresh = False
+                listw.clear()
+                logging.debug("Refresh complete.")
+
+    # Handle events.
+    def events(self, nodes):
+        for name, info in nodes:
+            if bool(info) and name in self.known_nodes:
+                self.fill_left_box()
+                self.fill_right_box()
+
+                if 'inbound' in info:
+                    for key in info['inbound'].keys():
+                        # New inbound online.
+                        if key not in self.known_inbound:
+                            addr = info['inbound'].get(key)
+                            if bool(addr):
+                                logging.debug(f"Refresh: inbound {key} online")
+                                refresh = True
+                        # Known inbound offline.
+                        for key in self.known_inbound:
+                            addr = info['inbound'].get(key)
+                            if bool(addr):
+                                continue
+                            logging.debug(f"Refresh: inbound {key} offline")
+                            self.refresh = True
+
+                # New outbound online.
+                if 'outbound' in info:
+                    for i, info in info['outbound'].items():
+                        addr = info[0]
+                        id = info[1]
+                        if id == 0:
+                            continue
+                        if id in self.known_outbound:
+                            continue
+                        logging.debug(f"Outbound {i}, {addr} came online.")
+                        self.refresh = True
+    
     async def update_view(self, evloop: asyncio.AbstractEventLoop,
                           loop: urwid.MainLoop):
-        live_nodes = []
-        dead_nodes = []
-        known_nodes = []
-        known_inbound = []
-        known_outbound = []
-        refresh = False
-
         while True:
             await asyncio.sleep(0.1)
 
             nodes = self.model.nodes.items()
+            liliths = self.model.liliths.items()
             listw = self.listwalker.contents
             evloop.call_soon(loop.draw_screen)
 
             for index, item in enumerate(listw):
                 # Keep track of known nodes.
-                if item.node_name not in known_nodes:
-                    known_nodes.append(item.node_name)
+                if item.node_name not in self.known_nodes:
+                    self.known_nodes.append(item.node_name)
                 # Keep track of known inbounds.
                 if (item.session == "inbound-slot"
-                        and item.i not in known_inbound):
-                    known_inbound.append(item.i)
+                        and item.i not in self.known_inbound):
+                    self.known_inbound.append(item.i)
                 # Keep track of known outbounds.
                 if (item.session == "outbound-slot"
-                        and item.id not in known_outbound
+                        and item.id not in self.known_outbound
                         and not item.id == 0):
-                    known_outbound.append(item.id)
-
-            for name, info in nodes:
-                # 1. Sort nodes into lists.
-                if bool(info) and name not in live_nodes:
-                    live_nodes.append(name)
-                if not bool(info) and name not in dead_nodes:
-                    dead_nodes.append(name)
-                if bool(info) and name in dead_nodes:
-                    logging.debug("Refresh: dead node online.")
-                    refresh = True
-                if not bool(info) and name in live_nodes:
-                    logging.debug("Refresh: online node offline.")
-                    refresh = True
-
-                # 2. Display nodes according to list.
-                if name in live_nodes and name not in known_nodes:
-                    self.draw_info(name, info)
-                if name in dead_nodes and name not in known_nodes:
-                    self.draw_empty(name, info)
-                if refresh:
-                    logging.debug("Refresh initiated.")
-                    await asyncio.sleep(0.1)
-                    known_outbound.clear()
-                    known_inbound.clear()
-                    known_nodes.clear()
-                    live_nodes.clear()
-                    dead_nodes.clear()
-                    refresh = False
-                    listw.clear()
-                    logging.debug("Refresh complete.")
-
-                # 3. Handle events on nodes we know.
-                if bool(info) and name in known_nodes:
-                    self.fill_left_box()
-                    self.fill_right_box()
-                    if 'inbound' in info:
-                        for key in info['inbound'].keys():
-                            # New inbound online.
-                            if key not in known_inbound:
-                                addr = info['inbound'].get(key)
-                                if bool(addr):
-                                    logging.debug(f"Refresh: inbound {key} online")
-                                    refresh = True
-                            # Known inbound offline.
-                            for key in known_inbound:
-                                addr = info['inbound'].get(key)
-                                if bool(addr):
-                                    continue
-                                logging.debug(f"Refresh: inbound {key} offline")
-                                refresh = True
-                    # New outbound online.
-                    if 'outbound' in info:
-                        for i, info in info['outbound'].items():
-                            addr = info[0]
-                            id = info[1]
-                            if id == 0:
-                                continue
-                            if id in known_outbound:
-                                continue
-                            logging.debug(f"Outbound {i}, {addr} came online.")
-                            refresh = True
+                    self.known_outbound.append(item.id)
+
+            self.sort(nodes)
+            self.sort(liliths)
+            
+            await self.draw(nodes)
+            await self.draw(liliths)
+
+            self.events(nodes)