server.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541
  1. import json
  2. import logging
  3. from io import BytesIO
  4. import attr
  5. from zope.interface import implementer
  6. from twisted.internet import address, threads, udp
  7. from twisted.internet._resolver import SimpleResolverComplexifier
  8. from twisted.internet.defer import Deferred, fail, succeed
  9. from twisted.internet.error import DNSLookupError
  10. from twisted.internet.interfaces import (
  11. IReactorPluggableNameResolver,
  12. IReactorTCP,
  13. IResolverSimple,
  14. )
  15. from twisted.python.failure import Failure
  16. from twisted.test.proto_helpers import AccumulatingProtocol, MemoryReactorClock
  17. from twisted.web.http import unquote
  18. from twisted.web.http_headers import Headers
  19. from twisted.web.server import Site
  20. from synapse.http.site import SynapseRequest
  21. from synapse.util import Clock
  22. from tests.utils import setup_test_homeserver as _sth
  23. logger = logging.getLogger(__name__)
  24. class TimedOutException(Exception):
  25. """
  26. A web query timed out.
  27. """
  28. @attr.s
  29. class FakeChannel(object):
  30. """
  31. A fake Twisted Web Channel (the part that interfaces with the
  32. wire).
  33. """
  34. site = attr.ib(type=Site)
  35. _reactor = attr.ib()
  36. result = attr.ib(default=attr.Factory(dict))
  37. _producer = None
  38. @property
  39. def json_body(self):
  40. if not self.result:
  41. raise Exception("No result yet.")
  42. return json.loads(self.result["body"].decode("utf8"))
  43. @property
  44. def code(self):
  45. if not self.result:
  46. raise Exception("No result yet.")
  47. return int(self.result["code"])
  48. @property
  49. def headers(self):
  50. if not self.result:
  51. raise Exception("No result yet.")
  52. h = Headers()
  53. for i in self.result["headers"]:
  54. h.addRawHeader(*i)
  55. return h
  56. def writeHeaders(self, version, code, reason, headers):
  57. self.result["version"] = version
  58. self.result["code"] = code
  59. self.result["reason"] = reason
  60. self.result["headers"] = headers
  61. def write(self, content):
  62. assert isinstance(content, bytes), "Should be bytes! " + repr(content)
  63. if "body" not in self.result:
  64. self.result["body"] = b""
  65. self.result["body"] += content
  66. def registerProducer(self, producer, streaming):
  67. self._producer = producer
  68. self.producerStreaming = streaming
  69. def _produce():
  70. if self._producer:
  71. self._producer.resumeProducing()
  72. self._reactor.callLater(0.1, _produce)
  73. if not streaming:
  74. self._reactor.callLater(0.0, _produce)
  75. def unregisterProducer(self):
  76. if self._producer is None:
  77. return
  78. self._producer = None
  79. def requestDone(self, _self):
  80. self.result["done"] = True
  81. def getPeer(self):
  82. # We give an address so that getClientIP returns a non null entry,
  83. # causing us to record the MAU
  84. return address.IPv4Address("TCP", "127.0.0.1", 3423)
  85. def getHost(self):
  86. return None
  87. @property
  88. def transport(self):
  89. return self
  90. class FakeSite:
  91. """
  92. A fake Twisted Web Site, with mocks of the extra things that
  93. Synapse adds.
  94. """
  95. server_version_string = b"1"
  96. site_tag = "test"
  97. access_logger = logging.getLogger("synapse.access.http.fake")
  98. def make_request(
  99. reactor,
  100. method,
  101. path,
  102. content=b"",
  103. access_token=None,
  104. request=SynapseRequest,
  105. shorthand=True,
  106. federation_auth_origin=None,
  107. ):
  108. """
  109. Make a web request using the given method and path, feed it the
  110. content, and return the Request and the Channel underneath.
  111. Args:
  112. method (bytes/unicode): The HTTP request method ("verb").
  113. path (bytes/unicode): The HTTP path, suitably URL encoded (e.g.
  114. escaped UTF-8 & spaces and such).
  115. content (bytes or dict): The body of the request. JSON-encoded, if
  116. a dict.
  117. shorthand: Whether to try and be helpful and prefix the given URL
  118. with the usual REST API path, if it doesn't contain it.
  119. federation_auth_origin (bytes|None): if set to not-None, we will add a fake
  120. Authorization header pretenting to be the given server name.
  121. Returns:
  122. Tuple[synapse.http.site.SynapseRequest, channel]
  123. """
  124. if not isinstance(method, bytes):
  125. method = method.encode("ascii")
  126. if not isinstance(path, bytes):
  127. path = path.encode("ascii")
  128. # Decorate it to be the full path, if we're using shorthand
  129. if (
  130. shorthand
  131. and not path.startswith(b"/_matrix")
  132. and not path.startswith(b"/_synapse")
  133. ):
  134. path = b"/_matrix/client/r0/" + path
  135. path = path.replace(b"//", b"/")
  136. if not path.startswith(b"/"):
  137. path = b"/" + path
  138. if isinstance(content, str):
  139. content = content.encode("utf8")
  140. site = FakeSite()
  141. channel = FakeChannel(site, reactor)
  142. req = request(channel)
  143. req.process = lambda: b""
  144. req.content = BytesIO(content)
  145. req.postpath = list(map(unquote, path[1:].split(b"/")))
  146. if access_token:
  147. req.requestHeaders.addRawHeader(
  148. b"Authorization", b"Bearer " + access_token.encode("ascii")
  149. )
  150. if federation_auth_origin is not None:
  151. req.requestHeaders.addRawHeader(
  152. b"Authorization",
  153. b"X-Matrix origin=%s,key=,sig=" % (federation_auth_origin,),
  154. )
  155. if content:
  156. req.requestHeaders.addRawHeader(b"Content-Type", b"application/json")
  157. req.requestReceived(method, path, b"1.1")
  158. return req, channel
  159. def wait_until_result(clock, request, timeout=100):
  160. """
  161. Wait until the request is finished.
  162. """
  163. clock.run()
  164. x = 0
  165. while not request.finished:
  166. # If there's a producer, tell it to resume producing so we get content
  167. if request._channel._producer:
  168. request._channel._producer.resumeProducing()
  169. x += 1
  170. if x > timeout:
  171. raise TimedOutException("Timed out waiting for request to finish.")
  172. clock.advance(0.1)
  173. def render(request, resource, clock):
  174. request.render(resource)
  175. wait_until_result(clock, request)
  176. @implementer(IReactorPluggableNameResolver)
  177. class ThreadedMemoryReactorClock(MemoryReactorClock):
  178. """
  179. A MemoryReactorClock that supports callFromThread.
  180. """
  181. def __init__(self):
  182. self.threadpool = ThreadPool(self)
  183. self._tcp_callbacks = {}
  184. self._udp = []
  185. lookups = self.lookups = {}
  186. @implementer(IResolverSimple)
  187. class FakeResolver(object):
  188. def getHostByName(self, name, timeout=None):
  189. if name not in lookups:
  190. return fail(DNSLookupError("OH NO: unknown %s" % (name,)))
  191. return succeed(lookups[name])
  192. self.nameResolver = SimpleResolverComplexifier(FakeResolver())
  193. super(ThreadedMemoryReactorClock, self).__init__()
  194. def listenUDP(self, port, protocol, interface="", maxPacketSize=8196):
  195. p = udp.Port(port, protocol, interface, maxPacketSize, self)
  196. p.startListening()
  197. self._udp.append(p)
  198. return p
  199. def callFromThread(self, callback, *args, **kwargs):
  200. """
  201. Make the callback fire in the next reactor iteration.
  202. """
  203. d = Deferred()
  204. d.addCallback(lambda x: callback(*args, **kwargs))
  205. self.callLater(0, d.callback, True)
  206. return d
  207. def getThreadPool(self):
  208. return self.threadpool
  209. def add_tcp_client_callback(self, host, port, callback):
  210. """Add a callback that will be invoked when we receive a connection
  211. attempt to the given IP/port using `connectTCP`.
  212. Note that the callback gets run before we return the connection to the
  213. client, which means callbacks cannot block while waiting for writes.
  214. """
  215. self._tcp_callbacks[(host, port)] = callback
  216. def connectTCP(self, host, port, factory, timeout=30, bindAddress=None):
  217. """Fake L{IReactorTCP.connectTCP}.
  218. """
  219. conn = super().connectTCP(
  220. host, port, factory, timeout=timeout, bindAddress=None
  221. )
  222. callback = self._tcp_callbacks.get((host, port))
  223. if callback:
  224. callback()
  225. return conn
  226. class ThreadPool:
  227. """
  228. Threadless thread pool.
  229. """
  230. def __init__(self, reactor):
  231. self._reactor = reactor
  232. def start(self):
  233. pass
  234. def stop(self):
  235. pass
  236. def callInThreadWithCallback(self, onResult, function, *args, **kwargs):
  237. def _(res):
  238. if isinstance(res, Failure):
  239. onResult(False, res)
  240. else:
  241. onResult(True, res)
  242. d = Deferred()
  243. d.addCallback(lambda x: function(*args, **kwargs))
  244. d.addBoth(_)
  245. self._reactor.callLater(0, d.callback, True)
  246. return d
  247. def setup_test_homeserver(cleanup_func, *args, **kwargs):
  248. """
  249. Set up a synchronous test server, driven by the reactor used by
  250. the homeserver.
  251. """
  252. server = _sth(cleanup_func, *args, **kwargs)
  253. database = server.config.database.get_single_database()
  254. # Make the thread pool synchronous.
  255. clock = server.get_clock()
  256. for database in server.get_datastores().databases:
  257. pool = database._db_pool
  258. def runWithConnection(func, *args, **kwargs):
  259. return threads.deferToThreadPool(
  260. pool._reactor,
  261. pool.threadpool,
  262. pool._runWithConnection,
  263. func,
  264. *args,
  265. **kwargs
  266. )
  267. def runInteraction(interaction, *args, **kwargs):
  268. return threads.deferToThreadPool(
  269. pool._reactor,
  270. pool.threadpool,
  271. pool._runInteraction,
  272. interaction,
  273. *args,
  274. **kwargs
  275. )
  276. pool.runWithConnection = runWithConnection
  277. pool.runInteraction = runInteraction
  278. pool.threadpool = ThreadPool(clock._reactor)
  279. pool.running = True
  280. return server
  281. def get_clock():
  282. clock = ThreadedMemoryReactorClock()
  283. hs_clock = Clock(clock)
  284. return clock, hs_clock
  285. @attr.s(cmp=False)
  286. class FakeTransport(object):
  287. """
  288. A twisted.internet.interfaces.ITransport implementation which sends all its data
  289. straight into an IProtocol object: it exists to connect two IProtocols together.
  290. To use it, instantiate it with the receiving IProtocol, and then pass it to the
  291. sending IProtocol's makeConnection method:
  292. server = HTTPChannel()
  293. client.makeConnection(FakeTransport(server, self.reactor))
  294. If you want bidirectional communication, you'll need two instances.
  295. """
  296. other = attr.ib()
  297. """The Protocol object which will receive any data written to this transport.
  298. :type: twisted.internet.interfaces.IProtocol
  299. """
  300. _reactor = attr.ib()
  301. """Test reactor
  302. :type: twisted.internet.interfaces.IReactorTime
  303. """
  304. _protocol = attr.ib(default=None)
  305. """The Protocol which is producing data for this transport. Optional, but if set
  306. will get called back for connectionLost() notifications etc.
  307. """
  308. disconnecting = False
  309. disconnected = False
  310. connected = True
  311. buffer = attr.ib(default=b"")
  312. producer = attr.ib(default=None)
  313. autoflush = attr.ib(default=True)
  314. def getPeer(self):
  315. return None
  316. def getHost(self):
  317. return None
  318. def loseConnection(self, reason=None):
  319. if not self.disconnecting:
  320. logger.info("FakeTransport: loseConnection(%s)", reason)
  321. self.disconnecting = True
  322. if self._protocol:
  323. self._protocol.connectionLost(reason)
  324. # if we still have data to write, delay until that is done
  325. if self.buffer:
  326. logger.info(
  327. "FakeTransport: Delaying disconnect until buffer is flushed"
  328. )
  329. else:
  330. self.connected = False
  331. self.disconnected = True
  332. def abortConnection(self):
  333. logger.info("FakeTransport: abortConnection()")
  334. if not self.disconnecting:
  335. self.disconnecting = True
  336. if self._protocol:
  337. self._protocol.connectionLost(None)
  338. self.disconnected = True
  339. def pauseProducing(self):
  340. if not self.producer:
  341. return
  342. self.producer.pauseProducing()
  343. def resumeProducing(self):
  344. if not self.producer:
  345. return
  346. self.producer.resumeProducing()
  347. def unregisterProducer(self):
  348. if not self.producer:
  349. return
  350. self.producer = None
  351. def registerProducer(self, producer, streaming):
  352. self.producer = producer
  353. self.producerStreaming = streaming
  354. def _produce():
  355. d = self.producer.resumeProducing()
  356. d.addCallback(lambda x: self._reactor.callLater(0.1, _produce))
  357. if not streaming:
  358. self._reactor.callLater(0.0, _produce)
  359. def write(self, byt):
  360. if self.disconnecting:
  361. raise Exception("Writing to disconnecting FakeTransport")
  362. self.buffer = self.buffer + byt
  363. # always actually do the write asynchronously. Some protocols (notably the
  364. # TLSMemoryBIOProtocol) get very confused if a read comes back while they are
  365. # still doing a write. Doing a callLater here breaks the cycle.
  366. if self.autoflush:
  367. self._reactor.callLater(0.0, self.flush)
  368. def writeSequence(self, seq):
  369. for x in seq:
  370. self.write(x)
  371. def flush(self, maxbytes=None):
  372. if not self.buffer:
  373. # nothing to do. Don't write empty buffers: it upsets the
  374. # TLSMemoryBIOProtocol
  375. return
  376. if self.disconnected:
  377. return
  378. if getattr(self.other, "transport") is None:
  379. # the other has no transport yet; reschedule
  380. if self.autoflush:
  381. self._reactor.callLater(0.0, self.flush)
  382. return
  383. if maxbytes is not None:
  384. to_write = self.buffer[:maxbytes]
  385. else:
  386. to_write = self.buffer
  387. logger.info("%s->%s: %s", self._protocol, self.other, to_write)
  388. try:
  389. self.other.dataReceived(to_write)
  390. except Exception as e:
  391. logger.exception("Exception writing to protocol: %s", e)
  392. return
  393. self.buffer = self.buffer[len(to_write) :]
  394. if self.buffer and self.autoflush:
  395. self._reactor.callLater(0.0, self.flush)
  396. if not self.buffer and self.disconnecting:
  397. logger.info("FakeTransport: Buffer now empty, completing disconnect")
  398. self.disconnected = True
  399. def connect_client(reactor: IReactorTCP, client_id: int) -> AccumulatingProtocol:
  400. """
  401. Connect a client to a fake TCP transport.
  402. Args:
  403. reactor
  404. factory: The connecting factory to build.
  405. """
  406. factory = reactor.tcpClients[client_id][2]
  407. client = factory.buildProtocol(None)
  408. server = AccumulatingProtocol()
  409. server.makeConnection(FakeTransport(client, reactor))
  410. client.makeConnection(FakeTransport(server, reactor))
  411. reactor.tcpClients.pop(client_id)
  412. return client, server