Commit 16f6f7fe14cab68bf51340f81df1857320beeeb4

Authored by totoro
1 parent b45f2413

Add support of domain in netservice

@@ -3,3 +3,4 @@ obj/Debug/xml_rpc_test.adm @@ -3,3 +3,4 @@ obj/Debug/xml_rpc_test.adm
3 obj/Debug/xml_rpc.unit_test.adm 3 obj/Debug/xml_rpc.unit_test.adm
4 obj/Debug/ami_test.adm 4 obj/Debug/ami_test.adm
5 bin/Debug/xml_rpc.unit_test.adm 5 bin/Debug/xml_rpc.unit_test.adm
  6 +*.bak
6 \ No newline at end of file 7 \ No newline at end of file
net_services/CXM_generic_client.anubis
@@ -136,14 +136,16 @@ define Maybe($T) @@ -136,14 +136,16 @@ define Maybe($T)
136 generic_request_for_service 136 generic_request_for_service
137 ( 137 (
138 MessageQueue queue, 138 MessageQueue queue,
139 - Word32 service_id,  
140 - Word32 service_version, 139 + Word32 service_id,
  140 + Word32 service_version,
  141 + String domain,
141 (MessageQueue, String) -> Maybe($T) handler, 142 (MessageQueue, String) -> Maybe($T) handler,
142 (String) -> One logger 143 (String) -> One logger
143 )= 144 )=
144 with test_msg = message(_CXM_REQUEST_FOR_SERVICE), 145 with test_msg = message(_CXM_REQUEST_FOR_SERVICE),
145 forget(add_int32(test_msg, "SERVICE", service_id)); 146 forget(add_int32(test_msg, "SERVICE", service_id));
146 forget(add_int32(test_msg, "VERSION", service_version)); 147 forget(add_int32(test_msg, "VERSION", service_version));
  148 + forget(add_string(test_msg, "DOMAIN", domain));
147 queue.add_Message_to_send(test_msg); 149 queue.add_Message_to_send(test_msg);
148 if queue.get_next_received_Message(10) is 150 if queue.get_next_received_Message(10) is
149 { 151 {
@@ -175,10 +177,11 @@ public define Maybe($T) @@ -175,10 +177,11 @@ public define Maybe($T)
175 generic_connect_to_net_service 177 generic_connect_to_net_service
176 ( 178 (
177 String queue_name, 179 String queue_name,
178 - Word32 server,  
179 - Word32 port,  
180 - Word32 service_id,  
181 - Word32 service_version, 180 + Word32 server,
  181 + Word32 port,
  182 + Word32 service_id,
  183 + Word32 service_version,
  184 + String domain,
182 (MessageQueue, String) -> Maybe($T) handler, 185 (MessageQueue, String) -> Maybe($T) handler,
183 (String) -> One logger 186 (String) -> One logger
184 ) 187 )
@@ -191,11 +194,27 @@ public define Maybe($T) @@ -191,11 +194,27 @@ public define Maybe($T)
191 with queue = create_MessageQueue(queue_name), 194 with queue = create_MessageQueue(queue_name),
192 message_transceiver(tcp(conn), queue); 195 message_transceiver(tcp(conn), queue);
193 // println("[" + virtual_machine_id + "] netservices generic_request_for_service()"); 196 // println("[" + virtual_machine_id + "] netservices generic_request_for_service()");
194 - with result = generic_request_for_service(queue, service_id, service_version, handler, logger), 197 + with result = generic_request_for_service(queue, service_id, service_version, domain, handler, logger),
195 // println("[" + virtual_machine_id + "] netservices client quit"); 198 // println("[" + virtual_machine_id + "] netservices client quit");
196 queue.quit(unique); 199 queue.quit(unique);
197 result 200 result
198 }. 201 }.
  202 +
  203 + //legacy version which not handle the domain
  204 +
  205 +public define Maybe($T)
  206 + generic_connect_to_net_service
  207 + (
  208 + String queue_name,
  209 + Word32 server,
  210 + Word32 port,
  211 + Word32 service_id,
  212 + Word32 service_version,
  213 + (MessageQueue, String) -> Maybe($T) handler,
  214 + (String) -> One logger
  215 + )
  216 + =
  217 + generic_connect_to_net_service(queue_name, server, port, service_id, service_version, "", handler, logger).
199 218
200 define Word32 219 define Word32
201 get_ip 220 get_ip
@@ -243,11 +262,24 @@ public define Maybe($T) @@ -243,11 +262,24 @@ public define Maybe($T)
243 Word32 port, 262 Word32 port,
244 Word32 service_id, 263 Word32 service_id,
245 Word32 service_version, 264 Word32 service_version,
  265 + String domain,
246 (MessageQueue, String) -> Maybe($T) handler, 266 (MessageQueue, String) -> Maybe($T) handler,
247 (String) -> One logger 267 (String) -> One logger
248 ) 268 )
249 - = generic_connect_to_net_service(queue_name, get_ip(server, logger), port, service_id, service_version, handler, logger).  
250 - 269 + = generic_connect_to_net_service(queue_name, get_ip(server, logger), port, service_id, service_version, domain, handler, logger).
  270 +
  271 +public define Maybe($T)
  272 + generic_connect_to_net_service
  273 + (
  274 + String queue_name,
  275 + String server,
  276 + Word32 port,
  277 + Word32 service_id,
  278 + Word32 service_version,
  279 + (MessageQueue, String) -> Maybe($T) handler,
  280 + (String) -> One logger
  281 + )
  282 + = generic_connect_to_net_service(queue_name, server, port, service_id, service_version, "", handler, logger).
251 283
252 public define Maybe($T) 284 public define Maybe($T)
253 generic_connect_to_net_service_SSL 285 generic_connect_to_net_service_SSL
@@ -260,6 +292,7 @@ public define Maybe($T) @@ -260,6 +292,7 @@ public define Maybe($T)
260 // case of an invalid, non trusted or missing certificate 292 // case of an invalid, non trusted or missing certificate
261 Word32 service_id, 293 Word32 service_id,
262 Word32 service_version, 294 Word32 service_version,
  295 + String domain,
263 (MessageQueue, String) -> Maybe($T) handler, 296 (MessageQueue, String) -> Maybe($T) handler,
264 (String) -> One logger 297 (String) -> One logger
265 ) 298 )
@@ -270,8 +303,25 @@ public define Maybe($T) @@ -270,8 +303,25 @@ public define Maybe($T)
270 ok(conn) then 303 ok(conn) then
271 with queue = create_MessageQueue(queue_name), 304 with queue = create_MessageQueue(queue_name),
272 message_transceiver(ssl(conn), queue); 305 message_transceiver(ssl(conn), queue);
273 - with result = generic_request_for_service(queue, service_id, service_version, handler, logger), 306 + with result = generic_request_for_service(queue, service_id, service_version, domain, handler, logger),
274 queue.quit(unique); 307 queue.quit(unique);
275 //logInfo(debug_log,"domain_manager client quit"); 308 //logInfo(debug_log,"domain_manager client quit");
276 result 309 result
277 }. 310 }.
  311 +
  312 +public define Maybe($T)
  313 + generic_connect_to_net_service_SSL
  314 + (
  315 + String queue_name,
  316 + String server_name,
  317 + Word32 server_ip,
  318 + Word32 port,
  319 + (Maybe(X509)) -> Bool accept_policy, // your policy for accepting the server certificate in
  320 + // case of an invalid, non trusted or missing certificate
  321 + Word32 service_id,
  322 + Word32 service_version,
  323 + (MessageQueue, String) -> Maybe($T) handler,
  324 + (String) -> One logger
  325 + )
  326 + = generic_connect_to_net_service_SSL(queue_name, server_name, server_ip, port, accept_policy, service_id, service_version, "", handler, logger).
  327 +
net_services/CXM_net_services.anubis
1 -๏ปฟ/*  
2 - *  
3 - * User: David RENE  
4 - * Date: 25/04/2007  
5 - * Time: 11:01  
6 - * (c) Calexium  
7 - *  
8 - */  
9 -read tools/basis.anubis  
10 -read system/convert.anubis  
11 -read system/string.anubis  
12 -read system/muscle.anubis  
13 -read system/data_io.anubis  
14 -read system/message_queue.anubis  
15 -read system/message_transceiver.anubis  
16 -read CXM_generic_protocol.anubis  
17 -read calexium_lib/CXM_message_constants.anubis  
18 -  
19 -public type NetService:  
20 - net_service(  
21 - Word32 version,  
22 - Word32 id,  
23 - String name,  
24 - (MessageQueue, String, String) -> One handler // Parameters are MessageQueue, peer IP and timestamp string  
25 - ).  
26 -  
27 -define One  
28 - print_services  
29 - (  
30 - List(NetService) net_services  
31 - ) =  
32 - map_forget((NetService net_s)|->  
33 - println(" id : 0x"+ to_hexa(net_s.id));  
34 - println(" version : " + to_String(net_s.version));  
35 - println(" name : "+ net_s.name );  
36 - println("----------------------------------------")  
37 - ,net_services).  
38 -  
39 - /** Try to find the service_id in services_list. If the service is found in that list  
40 - * the corresponding NetService object is return  
41 - */  
42 -define Maybe(NetService)  
43 - find_service  
44 - (  
45 - List(NetService) services_list,  
46 - Word32 service_id,  
47 - Word32 service_version  
48 - )=  
49 - if services_list is  
50 - {  
51 - [] then failure,  
52 - [h . t] then  
53 - if h.id = service_id & h.version >=+ service_version then  
54 - success(h)  
55 - else  
56 - find_service(t, service_id, service_version)  
57 - }.  
58 -  
59 - /** Check if the muscle message msg has the correct fields for requesting a net_services  
60 - * if we found "service" and "version" fields on the message, we try to find if the service  
61 - * referenced in "service" is available in net_services list  
62 - */  
63 -define Maybe(NetService)  
64 - has_service  
65 - (  
66 - MessageQueue queue,  
67 - Message msg,  
68 - List(NetService) net_services  
69 - )=  
70 - if find_int32(msg, "SERVICE") is  
71 - {  
72 - failure then //send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE); failure,  
73 - // old names... should be removed soon  
74 - if find_int32(msg, "service") is  
75 - {  
76 - failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE); failure,  
77 - success(service_id) then  
78 - if find_int32(msg, "version") is  
79 - {  
80 - failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE);failure,  
81 - success(service_version) then find_service(net_services, service_id, service_version)  
82 - }  
83 - }  
84 -  
85 - success(service_id) then  
86 - if find_int32(msg, "VERSION") is  
87 - {  
88 - failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE);failure,  
89 - success(service_version) then find_service(net_services, service_id, service_version)  
90 - }  
91 - }.  
92 -  
93 -define String  
94 - get_time_stamp  
95 - =  
96 - with time = (UTime) unow,  
97 - "<"+virtual_machine_id+"@"+time.seconds+">".  
98 -  
99 - /** This message_received function just handle the negociation process the available net_services.  
100 - * In other words, it only recognize the _CXM_REQUEST_FOR_SERVICE message and try to launch the  
101 - * corresponding servcice  
102 - */  
103 -  
104 -define One  
105 - service_negociation  
106 - (  
107 - MessageQueue queue,  
108 - Message msg,  
109 - List(NetService) net_services,  
110 - String peer  
111 - )=  
112 - //println("Service NEGOCIATION [" + to_hexa(*msg.what) + "] received");  
113 - if * msg.what = _CXM_REQUEST_FOR_SERVICE then  
114 - if has_service(queue, msg, net_services) is  
115 - {  
116 - failure then  
117 - send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE, _CXM_UNKNOW_SERVICE, "Unknown service")  
118 - success(net_service) then  
119 - with result = message(0),  
120 - timestamp = get_time_stamp,  
121 - forget(add_string(result, "TIMESTAMP", timestamp));  
122 - send_ACK_ok(queue, _CXM_REQUEST_FOR_SERVICE, result);  
123 - net_service.handler(queue, peer, timestamp)  
124 - }  
125 - else  
126 - send_ACK_error(queue, *msg.what, _CXM_UNKNOW_CMD, "Unknown command [" + (*msg.what) + "]")  
127 - .  
128 -  
129 - /**  
130 - * this function unflatten muscle message and give the correct message to service_negociation function  
131 - */  
132 -define One  
133 - message_receiver  
134 - (  
135 - MessageQueue queue,  
136 - List(NetService) net_services,  
137 - String peer  
138 - ) =  
139 - if queue.quit_requested(unique) then  
140 - unique  
141 - else  
142 - //println("PRE SERVICE message_receiver "+"["+virtual_machine_id + "]");  
143 - if queue.get_next_received_Message(1) is  
144 - {  
145 - timeout then //println("PRE timeout");  
146 - message_receiver(queue, net_services, peer),  
147 - closed then //println("PRE closed");  
148 - unique,  
149 - msg(msg) then unique; //println("PRE negociation");  
150 - service_negociation(queue, msg, net_services, peer);  
151 - message_receiver(queue, net_services, peer)  
152 - }.  
153 -  
154 -define Server -> (RWStream) -> One  
155 - net_services_handler  
156 - (  
157 - List(NetService) net_services,  
158 - ) =  
159 - (Server server) |-> (RWStream conn) |->  
160 - if remote_IP_address_and_port(conn) is (num_peer,_) then  
161 - //convert IP address of the client to string  
162 - with peer = ip_addr_to_string(num_peer),  
163 - //println("NET SERVICES Accepting connection with "+peer);  
164 -  
165 - //now managing the list of SERVICES  
166 - with queue = create_MessageQueue("CXM Net Services"),  
167 - message_transceiver(conn, queue);  
168 - message_receiver(queue, net_services, peer).  
169 -  
170 -  
171 -public define Maybe(Server)  
172 - start_net_services  
173 - (  
174 - List(NetService) net_services,  
175 - Word32 network_port,  
176 - )=  
177 - if start_server(0,  
178 - network_port,  
179 - net_services_handler(net_services),  
180 - (One u) |-> unique) is  
181 - {  
182 - cannot_create_the_socket then println("Cannot create the listening socket."); failure,  
183 - cannot_bind_to_port then println("Cannot bind to port " + network_port ); failure,  
184 - cannot_listen_on_port then println("Cannot listen on port " + network_port); failure,  
185 - ok(server) then  
186 - println("Net services started on port " + network_port);  
187 - println("------ Available services ------");  
188 - print_services(net_services);  
189 - success(server)  
190 - }. 1 +๏ปฟ/*
  2 + *
  3 + * User: David RENE
  4 + * Date: 25/04/2007
  5 + * Time: 11:01
  6 + * (c) Calexium
  7 + *
  8 + */
  9 +read tools/basis.anubis
  10 +read system/convert.anubis
  11 +read system/string.anubis
  12 +read system/muscle.anubis
  13 +read system/data_io.anubis
  14 +read system/message_queue.anubis
  15 +read system/message_transceiver.anubis
  16 +read CXM_generic_protocol.anubis
  17 +read calexium_lib/CXM_message_constants.anubis
  18 +
  19 +public type NetService:
  20 + net_service(
  21 + Word32 version,
  22 + Word32 id,
  23 + String name,
  24 + List(String) domains,
  25 + (MessageQueue, String, String) -> One handler // Parameters are MessageQueue, peer IP and timestamp string
  26 + ).
  27 +
  28 +define One
  29 + print_services
  30 + (
  31 + List(NetService) net_services
  32 + ) =
  33 +
  34 + map_forget((NetService net_s)|->
  35 + if net_s is net_service(version, id, name, domains, _) then
  36 + println(" id : 0x"+ to_hexa(id));
  37 + println(" version : " + to_String(version));
  38 + println(" name : "+ name );
  39 + println(" domains : ");
  40 + map_forget((String domain) |-> println(" : "+domain), domains);
  41 + println("----------------------------------------")
  42 + ,net_services).
  43 +
  44 + /** Try to find the service_id in services_list. If the service is found in that list
  45 + * the corresponding NetService object is return
  46 + */
  47 +define Maybe(NetService)
  48 + find_service
  49 + (
  50 + List(NetService) services_list,
  51 + Word32 service_id,
  52 + Word32 service_version
  53 + )=
  54 + if services_list is
  55 + {
  56 + [] then failure,
  57 + [h . t] then
  58 + if h.id = service_id & h.version >=+ service_version then
  59 + success(h)
  60 + else
  61 + find_service(t, service_id, service_version)
  62 + }.
  63 +
  64 + /** Check if the muscle message msg has the correct fields for requesting a net_services
  65 + * if we found "service" and "version" fields on the message, we try to find if the service
  66 + * referenced in "service" is available in net_services list
  67 + */
  68 +define Maybe(NetService)
  69 + has_service
  70 + (
  71 + MessageQueue queue,
  72 + Message msg,
  73 + List(NetService) net_services
  74 + )=
  75 + if find_int32(msg, "SERVICE") is
  76 + {
  77 + failure then //send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE); failure,
  78 + // old names... should be removed soon
  79 + if find_int32(msg, "service") is
  80 + {
  81 + failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE); failure,
  82 + success(service_id) then
  83 + if find_int32(msg, "version") is
  84 + {
  85 + failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE);failure,
  86 + success(service_version) then find_service(net_services, service_id, service_version)
  87 + }
  88 + }
  89 +
  90 + success(service_id) then
  91 + if find_int32(msg, "VERSION") is
  92 + {
  93 + failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE);failure,
  94 + success(service_version) then find_service(net_services, service_id, service_version)
  95 + }
  96 + }.
  97 +
  98 +define String
  99 + get_time_stamp
  100 + =
  101 + with time = (UTime) unow,
  102 + "<"+virtual_machine_id+"@"+time.seconds+">".
  103 +
  104 + /** This message_received function just handle the negociation process the available net_services.
  105 + * In other words, it only recognize the _CXM_REQUEST_FOR_SERVICE message and try to launch the
  106 + * corresponding servcice
  107 + */
  108 +
  109 +define One
  110 + service_negociation
  111 + (
  112 + MessageQueue queue,
  113 + Message msg,
  114 + List(NetService) net_services,
  115 + String peer
  116 + )=
  117 + //println("Service NEGOCIATION [" + to_hexa(*msg.what) + "] received");
  118 + if * msg.what = _CXM_REQUEST_FOR_SERVICE then
  119 + if has_service(queue, msg, net_services) is
  120 + {
  121 + failure then
  122 + send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE, _CXM_UNKNOW_SERVICE, "Unknown service")
  123 + success(net_service) then
  124 + with result = message(0),
  125 + timestamp = get_time_stamp,
  126 + forget(add_string(result, "TIMESTAMP", timestamp));
  127 + send_ACK_ok(queue, _CXM_REQUEST_FOR_SERVICE, result);
  128 + net_service.handler(queue, peer, timestamp)
  129 + }
  130 + else
  131 + send_ACK_error(queue, *msg.what, _CXM_UNKNOW_CMD, "Unknown command [" + (*msg.what) + "]")
  132 + .
  133 +
  134 + /**
  135 + * this function unflatten muscle message and give the correct message to service_negociation function
  136 + */
  137 +define One
  138 + message_receiver
  139 + (
  140 + MessageQueue queue,
  141 + List(NetService) net_services,
  142 + String peer
  143 + ) =
  144 + if queue.quit_requested(unique) then
  145 + unique
  146 + else
  147 + //println("PRE SERVICE message_receiver "+"["+virtual_machine_id + "]");
  148 + if queue.get_next_received_Message(1) is
  149 + {
  150 + timeout then //println("PRE timeout");
  151 + message_receiver(queue, net_services, peer),
  152 + closed then //println("PRE closed");
  153 + unique,
  154 + msg(msg) then unique; //println("PRE negociation");
  155 + service_negociation(queue, msg, net_services, peer);
  156 + message_receiver(queue, net_services, peer)
  157 + }.
  158 +
  159 +define Server -> (RWStream) -> One
  160 + net_services_handler
  161 + (
  162 + List(NetService) net_services,
  163 + ) =
  164 + (Server server) |-> (RWStream conn) |->
  165 + if remote_IP_address_and_port(conn) is (num_peer,_) then
  166 + //convert IP address of the client to string
  167 + with peer = ip_addr_to_string(num_peer),
  168 + //println("NET SERVICES Accepting connection with "+peer);
  169 +
  170 + //now managing the list of SERVICES
  171 + with queue = create_MessageQueue("CXM Net Services"),
  172 + message_transceiver(conn, queue);
  173 + message_receiver(queue, net_services, peer).
  174 +
  175 +
  176 +public define Maybe(Server)
  177 + start_net_services
  178 + (
  179 + List(NetService) net_services,
  180 + Word32 network_port,
  181 + )=
  182 + if start_server(0,
  183 + network_port,
  184 + net_services_handler(net_services),
  185 + (One u) |-> unique) is
  186 + {
  187 + cannot_create_the_socket then println("Cannot create the listening socket."); failure,
  188 + cannot_bind_to_port then println("Cannot bind to port " + network_port ); failure,
  189 + cannot_listen_on_port then println("Cannot listen on port " + network_port); failure,
  190 + ok(server) then
  191 + println("Net services started on port " + network_port);
  192 + println("------ Available services ------");
  193 + print_services(net_services);
  194 + success(server)
  195 + }.
web/CXM_multihost_http_server.anubis
@@ -857,7 +857,7 @@ define ReadResult @@ -857,7 +857,7 @@ define ReadResult
857 //if unow > dead_line then record_dubious_connection(connection,dead_line,dos) else 857 //if unow > dead_line then record_dubious_connection(connection,dead_line,dos) else
858 if read(connection.conn, 16384, time_out) is // the connection is closed after 10 minutes of inactivity 858 if read(connection.conn, 16384, time_out) is // the connection is closed after 10 minutes of inactivity
859 { 859 {
860 - error then println(pid + "read failed)"); error, 860 + error then println(pid + "read failed ["+to_string(result_buffer)+"]"); error,
861 timeout then timeout, 861 timeout then timeout,
862 ok(ba) then 862 ok(ba) then
863 // println(pid + "ba = " + length(ba)); 863 // println(pid + "ba = " + length(ba));