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
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 addr = uri.hostname, uri.port
88 s = socket.create_connection(addr)
89 self._ep = Endpoint(s, addr)
90 self._recv_gr = core.runner.spawn(self._recv)
91 self._beater_gr = core.runner.spawn(self._beat)
92 self.state = 'connected'
93 except Exception:
94 self.state = 'initial'
95 log.exception('Error connecting server')
96 core.events.server_refused.emit(True)
98 def disconnect(self) -> None:
99 if self.state != 'connected':
100 return
102 self.state = 'dying'
103 ep, recv, beater = self._ep, self._recv_gr, self._beater_gr
104 self._ep = None
105 self._recv_gr = None
106 self._beater_gr = None
107 ep and ep.close()
108 recv and recv.kill()
109 beater and beater.kill()
110 self.state = 'initial'
112 def write(self, v: wire.ClientToServer) -> None:
113 ep = self._ep
114 if ep: 114 ↛ 117line 114 didn't jump to line 117, because the condition on line 114 was never false
115 ep.write(v)
116 else:
117 raise Exception('No endpoint present')
119 def raw_write(self, v: bytes) -> None:
120 ep = self._ep
121 if ep:
122 ep.raw_write(v)
123 else:
124 raise Exception('No endpoint present')
126 # ----- Methods -----
127 def _recv(self) -> None:
128 core = self.core
129 me = gevent.getcurrent()
130 me.link_exception(self._dropped)
131 me.gr_name = f'{core}::RECV'
132 D = core.events.server_command
134 assert self._ep
135 try:
136 for v in self._ep.messages(timeout=None): 136 ↛ 141line 136 didn't jump to line 141, because the loop on line 136 didn't complete
137 D[v.__class__].emit(v)
138 except EndpointDied:
139 pass
141 core.events.server_dropped.emit(True)
143 def _dropped(self, _) -> None:
144 core = self.core
145 core.events.server_dropped.emit(True)
147 def _beat(self) -> None:
148 core = self.core
149 if core.options.testing: 149 ↛ 152line 149 didn't jump to line 152, because the condition on line 149 was never false
150 return
152 while self._ep:
153 self._ep.write(wire.Beat())
154 core.runner.sleep(10)