-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathdemo_server.py
More file actions
112 lines (84 loc) · 3.23 KB
/
Copy pathdemo_server.py
File metadata and controls
112 lines (84 loc) · 3.23 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
import asyncio
from asyncio.queues import Queue
from typing import Any
from starlette.applications import Starlette
from starlette.endpoints import WebSocketEndpoint
from starlette.routing import WebSocketRoute
from starlette.websockets import WebSocket
from olink.core.types import Name
from olink.remote import IObjectSource, RemoteNode
class Counter:
count = 0
_node: RemoteNode
def increment(self):
self.count += 1
# notify all registered clients
RemoteNode.notify_property_change("demo.Counter/count", self.count)
class CounterAdapter(IObjectSource):
node: RemoteNode = None
def __init__(self, impl):
self.impl = impl
self._methods = {"increment": impl.increment}
self._properties = {"count": lambda v: setattr(impl, "count", v)}
# need to register this source with the registry
RemoteNode.register_source(self)
def olink_object_name(self):
# name this source is registered under
return "demo.Counter"
def olink_invoke(self, name: str, args: list[Any]) -> Any:
# called on incoming invoke message
path = Name.path_from_name(name)
func = self._methods.get(path)
if func is None:
return None
return func(*args)
def olink_set_property(self, name: str, value: Any):
# called on incoming set property message
path = Name.path_from_name(name)
setter = self._properties.get(path)
if setter is not None:
setter(value)
def olink_linked(self, name: str, node: "RemoteNode"):
# called when a remote node is linked to this node
self.impl._node = node
def olink_unlinked(self, name: str, node: "RemoteNode"):
# called when a remote node is unlinked from this node
self.impl._node = None
def olink_collect_properties(self) -> object:
return {k: getattr(self.impl, k) for k in ["count"]}
counter = Counter()
adapter = CounterAdapter(counter)
class RemoteEndpoint(WebSocketEndpoint):
encoding = "text"
def __init__(self, scope, receive, send):
super().__init__(scope, receive, send)
self.node = RemoteNode()
self.queue = Queue()
async def sender(self, ws):
print("start sender")
while True:
msg = await self.queue.get()
print("send", msg)
await ws.send_text(msg)
self.queue.task_done()
async def on_connect(self, ws: WebSocket):
print("on_connect")
self._sender_task = asyncio.create_task(self.sender(ws))
def writer(msg: str):
print("writer", msg)
self.queue.put_nowait(msg)
self.node.on_write(writer)
await super().on_connect(ws)
async def on_receive(self, ws: WebSocket, data: Any) -> None:
print("on_receive", data)
self.node.handle_message(data)
async def on_disconnect(self, websocket: WebSocket, close_code: int) -> None:
await super().on_disconnect(websocket, close_code)
self.node.on_write(None)
self._sender_task.cancel()
try:
await self._sender_task
except asyncio.CancelledError:
pass
routes = [WebSocketRoute("/ws", RemoteEndpoint)]
app = Starlette(routes=routes)