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, Any 

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

98 

99 def disconnect(self) -> None: 

100 if self.state != 'connected': 

101 return 

102 

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' 

112 

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

119 

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

126 

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 

134 

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 

141 

142 core.events.server_dropped.emit(True) 

143 

144 def _dropped(self, _: Any) -> None: 

145 core = self.core 

146 core.events.server_dropped.emit(True) 

147 

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 

152 

153 while self._ep: 

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

155 core.runner.sleep(10)