Coverage for client/parts/server.py : 67%
Hot-keys on this page
r m x p toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
1# -*- coding: utf-8 -*-
2from __future__ import annotations
4# -- stdlib --
5from typing import Optional, TYPE_CHECKING, Any
6from urllib.parse import urlparse
7import logging
8import socket
10# -- third party --
11from gevent import Greenlet
12import gevent
14# -- own --
15from endpoint import Endpoint, EndpointDied
16from utils.events import EventHub
17import wire
19# -- typing --
20if TYPE_CHECKING:
21 from client.core import Core # noqa: F401
24# -- code --
25log = logging.getLogger('client.parts.Server')
26STOP = EventHub.STOP_PROPAGATION
29class Server(object):
30 def __init__(self, core: Core):
31 self.core = core
33 self.server_name = 'Unknown'
34 self.state = 'initial'
36 self._ep: Optional[Endpoint] = None
37 self._recv_gr: Optional[Greenlet] = None
38 self._beater_gr: Optional[Greenlet] = None
40 D = core.events.server_command
41 D[wire.Greeting] += self._greeting
42 D[wire.Ping] += self._ping
43 D[wire.Info] += self._info
44 D[wire.Error] += self._error
46 def _greeting(self, ev: wire.Greeting) -> wire.Greeting:
47 from settings import VERSION
49 core = self.core
51 if ev.version != VERSION: 51 ↛ 52line 51 didn't jump to line 52, because the condition on line 51 was never true
52 self.disconnect()
53 core.events.version_mismatch.emit(True)
54 else:
55 self.server_name = ev.node
56 core.events.server_connected.emit(True)
58 return ev
60 def _ping(self, ev: wire.Ping) -> wire.Ping:
61 self.write(wire.Pong())
62 return ev
64 def _info(self, ev: wire.Info) -> wire.Info:
65 core = self.core
66 core.events.server_info.emit(ev.msg)
67 return ev
69 def _error(self, ev: wire.Error) -> wire.Error:
70 log.warning('ServerError: %s', ev.msg)
71 core = self.core
72 core.events.server_error.emit(ev.msg)
73 return ev
75 # ----- Public Methods -----
76 def connect(self, uri: str) -> None:
77 core = self.core
79 uri = urlparse(uri)
80 assert uri.scheme == 'tcp'
82 if not self.state == 'initial': 82 ↛ 83line 82 didn't jump to line 83, because the condition on line 82 was never true
83 return
85 try:
86 self.state = 'connecting'
87 assert uri.port
88 addr = uri.hostname, uri.port
89 s = socket.create_connection(addr)
90 self._ep = Endpoint(s, addr)
91 self._recv_gr = core.runner.spawn(self._recv)
92 self._beater_gr = core.runner.spawn(self._beat)
93 self.state = 'connected'
94 except Exception:
95 self.state = 'initial'
96 log.exception('Error connecting server')
97 core.events.server_refused.emit(True)
99 def disconnect(self) -> None:
100 if self.state != 'connected':
101 return
103 self.state = 'dying'
104 ep, recv, beater = self._ep, self._recv_gr, self._beater_gr
105 self._ep = None
106 self._recv_gr = None
107 self._beater_gr = None
108 ep and ep.close()
109 recv and recv.kill()
110 beater and beater.kill()
111 self.state = 'initial'
113 def write(self, v: wire.ClientToServer) -> None:
114 ep = self._ep
115 if ep: 115 ↛ 118line 115 didn't jump to line 118, because the condition on line 115 was never false
116 ep.write(v)
117 else:
118 raise Exception('No endpoint present')
120 def raw_write(self, v: bytes) -> None:
121 ep = self._ep
122 if ep:
123 ep.raw_write(v)
124 else:
125 raise Exception('No endpoint present')
127 # ----- Methods -----
128 def _recv(self) -> None:
129 core = self.core
130 me = gevent.getcurrent()
131 me.link_exception(self._dropped)
132 me.gr_name = f'{core}::RECV'
133 D = core.events.server_command
135 assert self._ep
136 try:
137 for v in self._ep.messages(timeout=None): 137 ↛ 142line 137 didn't jump to line 142, because the loop on line 137 didn't complete
138 D[v.__class__].emit(v)
139 except EndpointDied:
140 pass
142 core.events.server_dropped.emit(True)
144 def _dropped(self, _: Any) -> None:
145 core = self.core
146 core.events.server_dropped.emit(True)
148 def _beat(self) -> None:
149 core = self.core
150 if core.options.testing: 150 ↛ 153line 150 didn't jump to line 153, because the condition on line 150 was never false
151 return
153 while self._ep:
154 self._ep.write(wire.Beat())
155 core.runner.sleep(10)