99from threading import Thread , Lock
1010
1111from opcua import ua
12- from opcua import utils
12+ from opcua . uaprocessor import UAProcessor
1313
1414logger = 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-
0 commit comments