Skip to content

Navigation Menu

Sign in
Appearance settings

Search code, repositories, users, issues, pull requests...

Provide feedback

We read every piece of feedback, and take your input very seriously.

Saved searches

Use saved searches to filter your results more quickly

Appearance settings

Commit e9746a4

Browse filesBrowse files
committed
cleanup, renaming and split some files
1 parent dcf7c80 commit e9746a4
Copy full SHA for e9746a4

13 files changed

+225-220Lines changed: 225 additions & 220 deletions
Expand file treeCollapse file tree
Open diff view settings
Collapse file

‎.gitignore‎

Copy file name to clipboard
+3Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,3 @@
1+
build*
2+
MANIFEST
3+
.idea*
Collapse file

‎example-server.py‎

Copy file name to clipboardExpand all lines: example-server.py
+1-3Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -23,9 +23,7 @@ def event(self, handle, event):
2323
logging.basicConfig(level=logging.WARN)
2424
logger = logging.getLogger("opcua.address_space")
2525
#logger = logging.getLogger("opcua.internal_server")
26-
logger = logging.getLogger("SubscriptionManager")
27-
logger.setLevel(logging.DEBUG)
28-
logger = logging.getLogger("Subscription")
26+
logger = logging.getLogger("opcua.subscription_server")
2927
logger.setLevel(logging.DEBUG)
3028

3129
# now setup our server and start it
Collapse file

‎opcua/binary_client.py‎

Copy file name to clipboardExpand all lines: opcua/binary_client.py
+1-1Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@ class BinaryClient(object):
2525
uaprotocol_auto.py and uaprotocol_hand.py
2626
"""
2727
def __init__(self):
28-
self.logger = logging.getLogger(self.__class__.__name__)
28+
self.logger = logging.getLogger(__name__)
2929
self._socket = None
3030
self._do_stop = False
3131
self._security_token = ua.ChannelSecurityToken()
Collapse file

‎opcua/binary_server.py‎

Copy file name to clipboardExpand all lines: opcua/binary_server.py
+1-195Lines changed: 1 addition & 195 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@
99
from threading import Thread, Lock
1010

1111
from opcua import ua
12-
from opcua import utils
12+
from opcua.uaprocessor import UAProcessor
1313

1414
logger = logging.getLogger(__name__)
1515

@@ -58,197 +58,3 @@ class ThreadingTCPServer(socketserver.ThreadingMixIn, socketserver.TCPServer):
5858
pass
5959

6060

61-
class UAProcessor(object):
62-
def __init__(self, internal_server, socket):
63-
self.logger = logging.getLogger(__name__)
64-
self.iserver = internal_server
65-
self.socket = socket
66-
self.channel = None
67-
self._lock = Lock()
68-
self.session = None
69-
70-
def loop(self):
71-
#first we want a hello message
72-
header = ua.Header.from_stream(self.socket)
73-
body = self.receive_body(header.body_size)
74-
if header.MessageType != ua.MessageType.Hello:
75-
self.logger.warning("received a message which is not a hello, sending back an error message %s", header)
76-
hdr = ua.Header(ua.MessageType.Error, ua.ChunkType.Single)
77-
self.write_socket(hdr)
78-
return
79-
hello = ua.Hello.from_binary(body)
80-
hdr = ua.Header(ua.MessageType.Acknowledge, ua.ChunkType.Single)
81-
ack = ua.Acknowledge()
82-
ack.ReceivebufferSize = hello.ReceiveBufferSize
83-
ack.SendbufferSize = hello.SendBufferSize
84-
self.write_socket(hdr, ack)
85-
86-
while True:
87-
header = ua.Header.from_stream(self.socket)
88-
if header is None:
89-
return
90-
if header.MessageType == ua.MessageType.Error:
91-
self.logger.warning("Received an error message type")
92-
return
93-
body = self.receive_body(header.body_size)
94-
if not self.process_body(header, body):
95-
break
96-
97-
def send_response(self, requesthandle, algohdr, seqhdr, response, msgtype=ua.MessageType.SecureMessage):
98-
with self._lock:
99-
response.ResponseHeader.RequestHandle = requesthandle
100-
seqhdr.SequenceNumber += 1
101-
hdr = ua.Header(msgtype, ua.ChunkType.Single, self.channel.SecurityToken.ChannelId)
102-
self.write_socket(hdr, algohdr, seqhdr, response)
103-
104-
def write_socket(self, hdr, *args):
105-
alle = []
106-
for arg in args:
107-
data = arg.to_binary()
108-
hdr.add_size(len(data))
109-
alle.append(data)
110-
alle.insert(0, hdr.to_binary())
111-
alle = b"".join(alle)
112-
self.logger.info("writting %s bytes to socket, with header %s ", len(alle), hdr)
113-
#self.logger.info("writting data %s", hdr, [i for i in args])
114-
#self.logger.debug("data: %s", alle)
115-
self.socket.send(alle)
116-
117-
def receive_body(self, size):
118-
self.logger.debug("reading body of message (%s bytes)", size)
119-
data = self.socket.recv(size)
120-
if size != len(data):
121-
raise Exception("Error, did not received expected number of bytes, got {}, asked for {}".format(len(data), size))
122-
return utils.Buffer(data)
123-
124-
def open_secure_channel(self, body):
125-
algohdr = ua.AsymmetricAlgorithmHeader.from_binary(body)
126-
seqhdr = ua.SequenceHeader.from_binary(body)
127-
request = ua.OpenSecureChannelRequest.from_binary(body)
128-
129-
self.channel = self.iserver.open_secure_channel(request.Parameters, self.channel)
130-
#send response
131-
response = ua.OpenSecureChannelResponse()
132-
response.Parameters = self.channel
133-
self.send_response(request.RequestHeader.RequestHandle, algohdr, seqhdr, response, ua.MessageType.SecureOpen)
134-
135-
def process_body(self, header, body):
136-
if header.MessageType == ua.MessageType.SecureOpen:
137-
self.open_secure_channel(body)
138-
139-
elif header.MessageType == ua.MessageType.SecureClose:
140-
if not self.channel or header.ChannelId != self.channel.SecurityToken.ChannelId:
141-
self.logger.warning("Request to close channel %s which was not issued, current channel is %s", header.ChannelId, self.channel)
142-
return False
143-
144-
elif header.MessageType == ua.MessageType.SecureMessage:
145-
algohdr = ua.SymmetricAlgorithmHeader.from_binary(body)
146-
seqhdr = ua.SequenceHeader.from_binary(body)
147-
self.process_message(algohdr, seqhdr, body)
148-
149-
else:
150-
self.logger.warning("Unsupported message type: %s", header.MessageType)
151-
return True
152-
153-
def process_message(self, algohdr, seqhdr, body):
154-
typeid = ua.NodeId.from_binary(body)
155-
requesthdr = ua.RequestHeader.from_binary(body)
156-
if typeid == ua.NodeId(ua.ObjectIds.CreateSessionRequest_Encoding_DefaultBinary):
157-
self.logger.info("Create session request")
158-
params = ua.CreateSessionParameters.from_binary(body)
159-
160-
self.session = self.iserver.create_session(params)
161-
162-
response = ua.CreateSessionResponse()
163-
response.Parameters = self.session
164-
self.send_response(requesthdr.RequestHandle, algohdr, seqhdr, response)
165-
166-
elif typeid == ua.NodeId(ua.ObjectIds.CloseSessionRequest_Encoding_DefaultBinary):
167-
self.logger.info("Close session request")
168-
deletesubs = ua.unpack_uatype('Boolean', body)
169-
170-
self.iserver.close_session(self.session, deletesubs)
171-
172-
response = ua.CloseSessionResponse()
173-
self.send_response(requesthdr.RequestHandle, algohdr, seqhdr, response)
174-
175-
elif typeid == ua.NodeId(ua.ObjectIds.ActivateSessionRequest_Encoding_DefaultBinary):
176-
self.logger.info("Activate session request")
177-
params = ua.ActivateSessionParameters.from_binary(body)
178-
179-
result = self.iserver.activate_session(self.session, params)
180-
181-
response = ua.ActivateSessionResponse()
182-
response.Parameters = result
183-
self.send_response(requesthdr.RequestHandle, algohdr, seqhdr, response)
184-
185-
elif typeid == ua.NodeId(ua.ObjectIds.ReadRequest_Encoding_DefaultBinary):
186-
self.logger.info("Read request")
187-
params = ua.ReadParameters.from_binary(body)
188-
189-
results = self.iserver.read(params)
190-
191-
response = ua.ReadResponse()
192-
response.Results = results
193-
self.send_response(requesthdr.RequestHandle, algohdr, seqhdr, response)
194-
195-
elif typeid == ua.NodeId(ua.ObjectIds.WriteRequest_Encoding_DefaultBinary):
196-
self.logger.info("Write request")
197-
params = ua.WriteParameters.from_binary(body)
198-
199-
results = self.iserver.write(params)
200-
201-
response = ua.WriteResponse()
202-
response.Results = results
203-
self.send_response(requesthdr.RequestHandle, algohdr, seqhdr, response)
204-
205-
elif typeid == ua.NodeId(ua.ObjectIds.BrowseRequest_Encoding_DefaultBinary):
206-
self.logger.info("Browse request")
207-
params = ua.BrowseParameters.from_binary(body)
208-
209-
results = self.iserver.browse(params)
210-
211-
response = ua.BrowseResponse()
212-
response.Results = results
213-
self.send_response(requesthdr.RequestHandle, algohdr, seqhdr, response)
214-
215-
elif typeid == ua.NodeId(ua.ObjectIds.GetEndpointsRequest_Encoding_DefaultBinary):
216-
self.logger.info("get endpoints request")
217-
params = ua.GetEndpointsParameters.from_binary(body)
218-
219-
endpoints = self.iserver.get_endpoints(params)
220-
221-
response = ua.GetEndpointsResponse()
222-
response.Endpoints = endpoints
223-
224-
self.send_response(requesthdr.RequestHandle, algohdr, seqhdr, response)
225-
226-
elif typeid == ua.NodeId(ua.ObjectIds.TranslateBrowsePathsToNodeIdsRequest_Encoding_DefaultBinary):
227-
self.logger.info("translate browsepaths to nodeids request")
228-
params = ua.TranslateBrowsePathsToNodeIdsParameters.from_binary(body)
229-
230-
paths = self.iserver.translate_browsepaths_to_nodeids(params.BrowsePaths)
231-
232-
response = ua.TranslateBrowsePathsToNodeIdsResponse()
233-
response.Results = paths
234-
235-
self.send_response(requesthdr.RequestHandle, algohdr, seqhdr, response)
236-
237-
elif typeid == ua.NodeId(ua.ObjectIds.AddNodesRequest_Encoding_DefaultBinary):
238-
self.logger.info("add nodes request")
239-
params = ua.AddNodesParameters.from_binary(body)
240-
241-
results = self.iserver.add_nodes(params.NodesToAdd)
242-
243-
response = ua.AddNodesResponse()
244-
response.Results = results
245-
246-
self.send_response(requesthdr.RequestHandle, algohdr, seqhdr, response)
247-
248-
else:
249-
self.logger.warning("Uknown message received %s", typeid)
250-
sf = ua.ServiceFault()
251-
sf.ResponseHeader.ServiceResult = ua.StatusCode(ua.StatusCodes.BadNotImplemented)
252-
self.send_response(requesthdr.RequestHandle, algohdr, seqhdr, sf)
253-
254-
Collapse file

‎opcua/client.py‎

Copy file name to clipboardExpand all lines: opcua/client.py
+2-2Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@ class KeepAlive(Thread):
1818
"""
1919
def __init__(self, client, timeout):
2020
Thread.__init__(self)
21-
self.logger = logging.getLogger(self.__class__.__name__)
21+
self.logger = logging.getLogger(__name__)
2222
if timeout == 0: # means no timeout bu we do not trust such servers
2323
timeout = 360000
2424
self.timeout = timeout
@@ -63,7 +63,7 @@ def __init__(self, url):
6363
if you are unsure of url, write at least hostname and port
6464
and call get_endpoints
6565
"""
66-
self.logger = logging.getLogger(self.__class__.__name__)
66+
self.logger = logging.getLogger(__name__)
6767
self.server_url = urlparse(url)
6868
self.name = "Pure Python Client"
6969
self.description = self.name
Collapse file

‎opcua/internal_server.py‎

Copy file name to clipboardExpand all lines: opcua/internal_server.py
+1-1Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,7 @@ def __str__(self):
3838

3939
class InternalServer(object):
4040
def __init__(self):
41-
self.logger = logging.getLogger(self.__class__.__name__)
41+
self.logger = logging.getLogger(__name__)
4242
self.endpoints = []
4343
self.sessions = {}
4444
self._channel_id_counter = 5
Collapse file

‎opcua/server.py‎

Copy file name to clipboardExpand all lines: opcua/server.py
-1Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,6 @@ def set_server_name(self, name):
5757
self.name = name
5858

5959
def start(self):
60-
print("START SERVER")
6160
self.iserver.start()
6261
self._set_endpoints()
6362
self.bserver = BinaryServer(self.iserver, self.endpoint.hostname, self.endpoint.port)
Collapse file

‎opcua/subscription.py‎

Copy file name to clipboardExpand all lines: opcua/subscription.py
+1-1Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ def __init__(self):
1616

1717
class Subscription(object):
1818
def __init__(self, server, params, handler):
19-
self.logger = logging.getLogger(self.__class__.__name__)
19+
self.logger = logging.getLogger(__name__)
2020
self.server = server
2121
self._client_handle = 200
2222
self._handler = handler
Collapse file

‎opcua/subscription_server.py‎

Copy file name to clipboardExpand all lines: opcua/subscription_server.py
+9-11Lines changed: 9 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -12,28 +12,26 @@
1212
class SubscriptionManager(Thread):
1313
def __init__(self, aspace):
1414
Thread.__init__(self)
15-
self.logger = logging.getLogger(self.__class__.__name__)
15+
self.logger = logging.getLogger(__name__)
1616
self.loop = None
1717
self.aspace = aspace
1818
self.subscriptions = {}
1919
self._sub_id_counter = 77
2020
self._cond = Condition()
2121

2222
def start(self):
23-
print("start internal")
2423
Thread.start(self)
2524
with self._cond:
2625
self._cond.wait()
27-
print("start internal finished")
2826

2927
def run(self):
30-
self.logger.warn("Starting subscription thread")
28+
self.logger.debug("Starting subscription thread")
3129
self.loop = asyncio.new_event_loop()
3230
asyncio.set_event_loop(self.loop)
3331
with self._cond:
3432
self._cond.notify_all()
3533
self.loop.run_forever()
36-
print("LOOP", self.loop)
34+
self.logger.debug("subsription thread ended")
3735

3836
def add_task(self, coro):
3937
"""
@@ -58,7 +56,7 @@ def create_subscription(self, params, callback):
5856
result.RevisedLifetimeCount = params.RequestedLifetimeCount
5957
result.RevisedMaxKeepAliveCount = params.RequestedMaxKeepAliveCount
6058

61-
sub = Subscription(self, result, self.aspace, callback)
59+
sub = InternalSubscription(self, result, self.aspace, callback)
6260
sub.start()
6361
self.subscriptions[result.SubscriptionId] = sub
6462

@@ -74,7 +72,7 @@ def delete_subscriptions(self, ids):
7472
return res
7573

7674
def publish(self, acks):
77-
self.logger.warn("publish request with acks %s", acks)
75+
self.logger.info("publish request with acks %s", acks)
7876

7977
def create_monitored_items(self, params):
8078
self.logger.info("create monitored items")
@@ -97,9 +95,9 @@ def __init__(self):
9795
self.mode = None
9896

9997

100-
class Subscription(object):
98+
class InternalSubscription(object):
10199
def __init__(self, manager, data, addressspace, callback):
102-
self.logger = logging.getLogger(self.__class__.__name__)
100+
self.logger = logging.getLogger(__name__)
103101
self.aspace = addressspace
104102
self.manager = manager
105103
self.data = data
@@ -123,7 +121,7 @@ def loop(self):
123121
yield from asyncio.sleep(1)
124122

125123
def publish_results(self):
126-
print("looking for results and publishing")
124+
self.logger.debug("looking for results and publishing")
127125

128126
def create_monitored_items(self, params):
129127
results = []
@@ -158,7 +156,7 @@ def _create_monitored_item(self, params):
158156
return result
159157

160158
def datachange_callback(self, handle, value):
161-
self.logger.warn("subscription %s: datachange callback called with %s, %s", self, handle, value)
159+
self.logger.info("subscription %s: datachange callback called with %s, %s", self, handle, value)
162160

163161

164162

0 commit comments

Comments
0 (0)
Morty Proxy This is a proxified and sanitized view of the page, visit original site.