-
Notifications
You must be signed in to change notification settings - Fork 24
Expand file tree
/
Copy pathchat_server.py
More file actions
267 lines (201 loc) · 10.7 KB
/
chat_server.py
File metadata and controls
267 lines (201 loc) · 10.7 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
import asyncio
import logging
import uuid
from asyncio import Queue
from collections import defaultdict
from dataclasses import dataclass, field
from typing import Dict, Optional, Set, Callable
from weakref import WeakValueDictionary, WeakSet
import reactivex
from reactivex import Observable, operators, Subject, Observer
from examples.tutorial.reactivex.shared import (Message, chat_filename_mimetype, ClientStatistics,
ServerStatisticsRequest, ServerStatistics, dataclass_to_payload,
decode_dataclass, decode_payload)
from rsocket.extensions.composite_metadata import CompositeMetadata
from rsocket.extensions.helpers import composite, metadata_item
from rsocket.frame_helpers import ensure_bytes
from rsocket.helpers import utf8_decode
from rsocket.payload import Payload
from rsocket.reactivex.back_pressure_publisher import (from_observable_with_backpressure, observable_from_queue,
observable_from_async_generator)
from rsocket.reactivex.reactivex_channel import ReactivexChannel
from rsocket.reactivex.reactivex_handler_adapter import reactivex_handler_factory
from rsocket.routing.request_router import RequestRouter
from rsocket.routing.routing_request_handler import RoutingRequestHandler
from rsocket.rsocket_server import RSocketServer
from rsocket.transports.tcp import TransportTCP
class SessionId(str): # allow weak reference
pass
@dataclass()
class UserSessionData:
username: str
session_id: SessionId
messages: Queue = field(default_factory=Queue)
statistics: Optional[ClientStatistics] = None
requested_statistics: ServerStatisticsRequest = field(default_factory=ServerStatisticsRequest)
@dataclass(frozen=True)
class ChatData:
channel_users: Dict[str, Set[SessionId]] = field(default_factory=lambda: defaultdict(WeakSet))
files: Dict[str, bytes] = field(default_factory=dict)
channel_messages: Dict[str, Queue] = field(default_factory=lambda: defaultdict(Queue))
user_session_by_id: Dict[str, UserSessionData] = field(default_factory=WeakValueDictionary)
chat_data = ChatData()
def ensure_channel_exists(channel_name):
if channel_name not in chat_data.channel_users:
chat_data.channel_users[channel_name] = WeakSet()
chat_data.channel_messages[channel_name] = Queue()
asyncio.create_task(channel_message_delivery(channel_name))
async def channel_message_delivery(channel_name: str):
logging.info('Starting channel delivery %s', channel_name)
while True:
try:
message = await chat_data.channel_messages[channel_name].get()
for session_id in chat_data.channel_users[channel_name]:
user_specific_message = Message(user=message.user,
content=message.content,
channel=channel_name)
chat_data.user_session_by_id[session_id].messages.put_nowait(user_specific_message)
except Exception as exception:
logging.error(str(exception), exc_info=True)
def get_file_name(composite_metadata):
return utf8_decode(composite_metadata.find_by_mimetype(chat_filename_mimetype)[0].content)
def find_session_by_username(username: str) -> Optional[UserSessionData]:
try:
return next(session for session in chat_data.user_session_by_id.values() if
session.username == username)
except StopIteration:
return None
def new_statistics_data(requested_statistics: ServerStatisticsRequest):
statistics_data = {}
if 'users' in requested_statistics.ids:
statistics_data['user_count'] = len(chat_data.user_session_by_id)
if 'channels' in requested_statistics.ids:
statistics_data['channel_count'] = len(chat_data.channel_messages)
return ServerStatistics(**statistics_data)
def find_username_by_session(session_id: SessionId) -> Optional[str]:
session = chat_data.user_session_by_id.get(session_id)
if session is None:
return None
return session.username
class ChatUserSession:
def __init__(self):
self._session: Optional[UserSessionData] = None
def remove(self):
if self._session is not None:
logging.info(f'Removing session: {self._session.session_id}')
del chat_data.user_session_by_id[self._session.session_id]
def router_factory(self):
router = RequestRouter(payload_deserializer=decode_payload)
@router.response('login')
async def login(username: str) -> Observable:
logging.info(f'New user: {username}')
session_id = SessionId(uuid.uuid4())
self._session = UserSessionData(username, session_id)
chat_data.user_session_by_id[session_id] = self._session
return reactivex.just(Payload(ensure_bytes(session_id)))
@router.response('channel.join')
async def join_channel(channel_name: str) -> Observable:
ensure_channel_exists(channel_name)
chat_data.channel_users[channel_name].add(self._session.session_id)
return reactivex.empty()
@router.response('channel.leave')
async def leave_channel(channel_name: str) -> Observable:
chat_data.channel_users[channel_name].discard(self._session.session_id)
return reactivex.empty()
@router.response('file.upload')
async def upload_file(payload: Payload, composite_metadata: CompositeMetadata) -> Observable:
chat_data.files[get_file_name(composite_metadata)] = payload.data
return reactivex.empty()
@router.response('file.download')
async def download_file(composite_metadata: CompositeMetadata) -> Observable:
file_name = get_file_name(composite_metadata)
return reactivex.just(Payload(chat_data.files[file_name],
composite(metadata_item(ensure_bytes(file_name), chat_filename_mimetype))))
@router.stream('files')
async def get_file_names() -> Callable[[Subject], Observable]:
async def generator():
for file_name in chat_data.files.keys():
yield file_name
return from_observable_with_backpressure(
lambda backpressure: observable_from_async_generator(generator(), backpressure).pipe(
operators.map(lambda file_name: Payload(ensure_bytes(file_name)))
))
@router.stream('channels')
async def get_channels() -> Observable:
return reactivex.from_iterable(
(Payload(ensure_bytes(channel)) for channel in chat_data.channel_messages.keys()))
@router.stream('channel.users')
async def get_channel_users(channel_name: str) -> Observable:
if channel_name not in chat_data.channel_users:
return reactivex.empty()
return reactivex.from_iterable(Payload(ensure_bytes(find_username_by_session(session_id))) for
session_id in
chat_data.channel_users[channel_name])
@router.fire_and_forget('statistics')
async def receive_statistics(statistics: ClientStatistics):
logging.info('Received client statistics. memory usage: %s', statistics.memory_usage)
self._session.statistics = statistics
@router.channel('statistics')
async def send_statistics() -> ReactivexChannel:
async def statistics_generator():
while True:
try:
await asyncio.sleep(self._session.requested_statistics.period_seconds)
yield new_statistics_data(self._session.requested_statistics)
except Exception:
logging.error('Statistics', exc_info=True)
def on_next(payload: Payload):
request = decode_dataclass(payload.data, ServerStatisticsRequest)
logging.info(f'Received statistics request {request.ids}, {request.period_seconds}')
if request.ids is not None:
self._session.requested_statistics.ids = request.ids
if request.period_seconds is not None:
self._session.requested_statistics.period_seconds = request.period_seconds
return ReactivexChannel(
from_observable_with_backpressure(
lambda backpressure: observable_from_async_generator(
statistics_generator(), backpressure
).pipe(
operators.map(dataclass_to_payload)
)),
Observer(on_next=on_next),
limit_rate=2)
@router.response('message')
async def send_message(message: Message) -> Observable:
logging.info('Received message for user: %s, channel: %s', message.user, message.channel)
target_message = Message(self._session.username, message.content, message.channel)
if message.channel is not None:
await chat_data.channel_messages[message.channel].put(target_message)
elif message.user is not None:
session = find_session_by_username(message.user)
await session.messages.put(target_message)
return reactivex.empty()
@router.stream('messages.incoming')
async def messages_incoming() -> Callable[[Subject], Observable]:
return from_observable_with_backpressure(
lambda backpressure: observable_from_queue(
self._session.messages, backpressure
).pipe(
operators.map(dataclass_to_payload)
)
)
return router
class CustomRoutingRequestHandler(RoutingRequestHandler):
def __init__(self, session: ChatUserSession):
super().__init__(session.router_factory())
self._session = session
async def on_close(self, rsocket, exception: Optional[Exception] = None):
self._session.remove()
return await super().on_close(rsocket, exception)
def handler_factory():
return CustomRoutingRequestHandler(ChatUserSession())
async def run_server():
def session(*connection):
RSocketServer(TransportTCP(*connection),
handler_factory=reactivex_handler_factory(handler_factory),
fragment_size_bytes=1_000_000)
async with await asyncio.start_server(session, 'localhost', 6565) as server:
await server.serve_forever()
if __name__ == '__main__':
logging.basicConfig(level=logging.INFO)
asyncio.run(run_server())