-
Notifications
You must be signed in to change notification settings - Fork 87
Expand file tree
/
Copy pathrelay_manager.py
More file actions
107 lines (86 loc) · 3.5 KB
/
Copy pathrelay_manager.py
File metadata and controls
107 lines (86 loc) · 3.5 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
import json
import time
import threading
from typing import Optional
from dataclasses import dataclass
from threading import Lock
from .event import Event
from .filter import Filters
from .message_pool import MessagePool
from .message_type import ClientMessageType
from .relay import Relay, RelayPolicy, RelayProxyConnectionConfig
from .request import Request
class RelayException(Exception):
pass
@dataclass
class RelayManager:
def __post_init__(self):
self.relays: dict[str, Relay] = {}
self.message_pool: MessagePool = MessagePool()
self.lock: Lock = Lock()
def add_relay(
self,
url: str,
policy: RelayPolicy = RelayPolicy(),
ssl_options: Optional[dict] = None,
proxy_config: Optional[RelayProxyConnectionConfig] = None):
relay = Relay(url, self.message_pool, policy, proxy_config, ssl_options)
with self.lock:
self.relays[url] = relay
threading.Thread(
target=relay.connect,
name=f"{relay.url}-thread",
daemon=True,
).start()
time.sleep(1)
def remove_relay(self, url: str):
with self.lock:
if url in self.relays:
relay = self.relays.pop(url)
relay.close()
def add_subscription_on_relay(self, url: str, id: str, filters: Filters):
with self.lock:
if url in self.relays:
relay = self.relays[url]
if not relay.policy.should_read:
raise RelayException(f"Could not send request: {url} is not configured to read from")
relay.add_subscription(id, filters)
request = Request(id, filters)
relay.publish(request.to_message())
else:
raise RelayException(f"Invalid relay url: no connection to {url}")
def add_subscription_on_all_relays(self, id: str, filters: Filters):
with self.lock:
for relay in self.relays.values():
if relay.policy.should_read:
relay.add_subscription(id, filters)
request = Request(id, filters)
relay.publish(request.to_message())
def close_subscription_on_relay(self, url: str, id: str):
with self.lock:
if url in self.relays:
relay = self.relays[url]
relay.close_subscription(id)
relay.publish(json.dumps(["CLOSE", id]))
else:
raise RelayException(f"Invalid relay url: no connection to {url}")
def close_subscription_on_all_relays(self, id: str):
with self.lock:
for relay in self.relays.values():
relay.close_subscription(id)
relay.publish(json.dumps(["CLOSE", id]))
def close_all_relay_connections(self):
with self.lock:
for url in self.relays:
relay = self.relays[url]
relay.close()
def publish_event(self, event: Event):
""" Verifies that the Event is publishable before submitting it to relays """
if event.signature is None:
raise RelayException(f"Could not publish {event.id}: must be signed")
if not event.verify():
raise RelayException(f"Could not publish {event.id}: failed to verify signature {event.signature}")
with self.lock:
for relay in self.relays.values():
if relay.policy.should_write:
relay.publish(event.to_message())