Hide keyboard shortcuts

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 

3 

4# -- stdlib -- 

5from typing import Optional, TYPE_CHECKING 

6from urllib.parse import urlparse 

7import logging 

8import socket 

9 

10# -- third party -- 

11from gevent import Greenlet 

12import gevent 

13 

14# -- own -- 

15from endpoint import Endpoint, EndpointDied 

16from utils.events import EventHub 

17import wire 

18 

19# -- typing -- 

20if TYPE_CHECKING: 

21 from client.core import Core # noqa: F401 

22 

23 

24# -- code -- 

25log = logging.getLogger('client.parts.Server') 

26STOP = EventHub.STOP_PROPAGATION 

27 

28 

29class Server(object): 

30 def __init__(self, core: Core): 

31 self.core = core 

32 

33 self.server_name = 'Unknown' 

34 self.state = 'initial' 

35 

36 self._ep: Optional[Endpoint] = None 

37 self._recv_gr: Optional[Greenlet] = None 

38 self._beater_gr: Optional[Greenlet] = None 

39 

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 

45 

46 def _greeting(self, ev: wire.Greeting) -> wire.Greeting: 

47 from settings import VERSION 

48 

49 core = self.core 

50 

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) 

57 

58 return ev 

59 

60 def _ping(self, ev: wire.Ping) -> wire.Ping: 

61 self.write(wire.Pong()) 

62 return ev 

63 

64 def _info(self, ev: wire.Info) -> wire.Info: 

65 core = self.core 

66 core.events.server_info.emit(ev.msg) 

67 return ev 

68 

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 

74 

75 # ----- Public Methods ----- 

76 def connect(self, uri: str) -> None: 

77 core = self.core 

78 

79 uri = urlparse(uri) 

80 assert uri.scheme == 'tcp' 

81 

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 

84 

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) 

97 

98 def disconnect(self) -> None: 

99 if self.state != 'connected': 

100 return 

101 

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' 

111 

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') 

118 

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') 

125 

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 

133 

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 

140 

141 core.events.server_dropped.emit(True) 

142 

143 def _dropped(self, _) -> None: 

144 core = self.core 

145 core.events.server_dropped.emit(True) 

146 

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 

151 

152 while self._ep: 

153 self._ep.write(wire.Beat()) 

154 core.runner.sleep(10)