Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

twisted-2-asyncio #28

Closed
wants to merge 1 commit into from
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 24 additions & 15 deletions ctrader_open_api/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,56 +6,65 @@
from ctrader_open_api.factory import Factory
from twisted.internet import reactor, defer

class Client(ClientService):

import asyncio
from asyncio import Future

class Client(asyncio.AbstractServer):
def __init__(self, host, port, protocol, retryPolicy=None, clock=None, prepareConnection=None, numberOfMessagesToSendPerSecond=5):
self._runningReactor = reactor
self.numberOfMessagesToSendPerSecond = numberOfMessagesToSendPerSecond
endpoint = clientFromString(self._runningReactor, f"ssl:{host}:{port}")
# endpoint = clientFromString(self._runningReactor, f"ssl:{host}:{port}")
endpoint = asyncio.open_connection(host, port, ssl=True)
factory = Factory.forProtocol(protocol, client=self)
super().__init__(endpoint, factory, retryPolicy=retryPolicy, clock=clock, prepareConnection=prepareConnection)
self._events = dict()
self._responseDeferreds = dict()
self.isConnected = False

def startService(self):
async def startService(self):
if self.running:
return
ClientService.startService(self)
self.transport, self.protocol = await self.start_serving()

def stopService(self):
async def stopService(self):
if self.running and self.isConnected:
ClientService.stopService(self)
await self.close()

def _connected(self, protocol):
def connection_made(self, transport):
self.transport = transport
self.isConnected = True
if hasattr(self, "_connectedCallback"):
self._connectedCallback(self)

def _disconnected(self, reason):
def connection_lost(self, exc):
self.isConnected = False
self._responseDeferreds.clear()
if hasattr(self, "_disconnectedCallback"):
self._disconnectedCallback(self, reason)
self._disconnectedCallback(self, exc)

def _received(self, message):
def data_received(self, data):
message = self.factory.protocol.string_received(data)
if hasattr(self, "_messageReceivedCallback"):
self._messageReceivedCallback(self, message)
if (message.clientMsgId is not None and message.clientMsgId in self._responseDeferreds):
responseDeferred = self._responseDeferreds[message.clientMsgId]
self._responseDeferreds.pop(message.clientMsgId)
responseDeferred.callback(message)
responseDeferred.set_result(message)

def send(self, message, clientMsgId=None, responseTimeoutInSeconds=5, **params):
if type(message) in [str, int]:
message = Protobuf.get(message, **params)
responseDeferred = defer.Deferred(self._cancelMessageDiferred)
responseDeferred = Future()
if clientMsgId is None:
clientMsgId = str(id(responseDeferred))
if clientMsgId is not None:
self._responseDeferreds[clientMsgId] = responseDeferred
responseDeferred.addErrback(lambda failure: self._onResponseFailure(failure, clientMsgId))
responseDeferred.addTimeout(responseTimeoutInSeconds, self._runningReactor)
protocolDiferred = self.whenConnected(failAfterFailures=1)
# responseDeferred.addErrback(lambda failure: self._onResponseFailure(failure, clientMsgId))
responseDeferred.add_errback(lambda failure: self._onResponseFailure(failure, clientMsgId))
responseDeferred.add_done_callback(lambda fut: fut.result())
responseDeferred.add_timeout(responseTimeoutInSeconds, self._runningReactor)
protocolDiferred = self.whenConnected(failAfterFailures=1)
protocolDiferred.addCallbacks(lambda protocol: protocol.send(message, clientMsgId=clientMsgId, isCanceled=lambda: clientMsgId not in self._responseDeferreds), responseDeferred.errback)
return responseDeferred

Expand Down
58 changes: 28 additions & 30 deletions ctrader_open_api/tcpProtocol.py
Original file line number Diff line number Diff line change
@@ -1,30 +1,26 @@
#!/usr/bin/env python

import asyncio
from collections import deque
from twisted.protocols.basic import Int32StringReceiver
from twisted.internet import task
from ctrader_open_api.messages.OpenApiCommonMessages_pb2 import ProtoMessage, ProtoHeartbeatEvent
import datetime

class TcpProtocol(Int32StringReceiver):

class TcpProtocol(asyncio.Protocol):
MAX_LENGTH = 15000000
_send_queue = deque([])
_send_task = None
_lastSendMessageTime = None

def connectionMade(self):
super().connectionMade()
def connection_made(self, transport):
self.transport = transport

if not self._send_task:
self._send_task = task.LoopingCall(self._sendStrings)
self._send_task.start(1)
self._send_task = asyncio.create_task(self._send_strings())
self.factory.connected(self)

def connectionLost(self, reason):
super().connectionLost(reason)
if self._send_task.running:
self._send_task.stop()
self.factory.disconnected(reason)
def connection_lost(self, exc):
if self._send_task:
self._send_task.cancel()
self.factory.disconnected(exc)

def heartbeat(self):
self.send(ProtoHeartbeatEvent(), True)
Expand All @@ -45,27 +41,29 @@ def send(self, message, instant=False, clientMsgId=None, isCanceled = None):
data = msg.SerializeToString()

if instant:
self.sendString(data)
self.transport.write(data)
self._lastSendMessageTime = datetime.datetime.now()
else:
self._send_queue.append((isCanceled, data))

def _sendStrings(self):
size = len(self._send_queue)

if not size:
if self._lastSendMessageTime is None or (datetime.datetime.now() - self._lastSendMessageTime).total_seconds() > 20:
self.heartbeat()
return

for _ in range(min(size, self.factory.numberOfMessagesToSendPerSecond)):
isCanceled, data = self._send_queue.popleft()
if isCanceled is not None and isCanceled():
continue;
self.sendString(data)
self._lastSendMessageTime = datetime.datetime.now()
async def _send_strings(self):
while True:
await asyncio.sleep(1)
size = len(self._send_queue)

if not size:
if self._lastSendMessageTime is None or (datetime.datetime.now() - self._lastSendMessageTime).total_seconds() > 20:
self.heartbeat()
continue

for _ in range(min(size, self.factory.numberOfMessagesToSendPerSecond)):
isCanceled, data = self._send_queue.popleft()
if isCanceled is not None and isCanceled():
continue
self.transport.write(data)
self._lastSendMessageTime = datetime.datetime.now()

def stringReceived(self, data):
def data_received(self, data):
msg = ProtoMessage()
msg.ParseFromString(data)

Expand Down
Loading
Loading