diff --git a/.hgignore b/.hgignore index eb804ee..6337528 100644 --- a/.hgignore +++ b/.hgignore @@ -3,3 +3,4 @@ obj/Debug/xml_rpc_test.adm obj/Debug/xml_rpc.unit_test.adm obj/Debug/ami_test.adm bin/Debug/xml_rpc.unit_test.adm +*.bak \ No newline at end of file diff --git a/net_services/CXM_generic_client.anubis b/net_services/CXM_generic_client.anubis index 8a4d531..d8d369d 100644 --- a/net_services/CXM_generic_client.anubis +++ b/net_services/CXM_generic_client.anubis @@ -136,14 +136,16 @@ define Maybe($T) generic_request_for_service ( MessageQueue queue, - Word32 service_id, - Word32 service_version, + Word32 service_id, + Word32 service_version, + String domain, (MessageQueue, String) -> Maybe($T) handler, (String) -> One logger )= with test_msg = message(_CXM_REQUEST_FOR_SERVICE), forget(add_int32(test_msg, "SERVICE", service_id)); forget(add_int32(test_msg, "VERSION", service_version)); + forget(add_string(test_msg, "DOMAIN", domain)); queue.add_Message_to_send(test_msg); if queue.get_next_received_Message(10) is { @@ -175,10 +177,11 @@ public define Maybe($T) generic_connect_to_net_service ( String queue_name, - Word32 server, - Word32 port, - Word32 service_id, - Word32 service_version, + Word32 server, + Word32 port, + Word32 service_id, + Word32 service_version, + String domain, (MessageQueue, String) -> Maybe($T) handler, (String) -> One logger ) @@ -191,11 +194,27 @@ public define Maybe($T) with queue = create_MessageQueue(queue_name), message_transceiver(tcp(conn), queue); // println("[" + virtual_machine_id + "] netservices generic_request_for_service()"); - with result = generic_request_for_service(queue, service_id, service_version, handler, logger), + with result = generic_request_for_service(queue, service_id, service_version, domain, handler, logger), // println("[" + virtual_machine_id + "] netservices client quit"); queue.quit(unique); result }. + + //legacy version which not handle the domain + +public define Maybe($T) + generic_connect_to_net_service + ( + String queue_name, + Word32 server, + Word32 port, + Word32 service_id, + Word32 service_version, + (MessageQueue, String) -> Maybe($T) handler, + (String) -> One logger + ) + = + generic_connect_to_net_service(queue_name, server, port, service_id, service_version, "", handler, logger). define Word32 get_ip @@ -243,11 +262,24 @@ public define Maybe($T) Word32 port, Word32 service_id, Word32 service_version, + String domain, (MessageQueue, String) -> Maybe($T) handler, (String) -> One logger ) - = generic_connect_to_net_service(queue_name, get_ip(server, logger), port, service_id, service_version, handler, logger). - + = generic_connect_to_net_service(queue_name, get_ip(server, logger), port, service_id, service_version, domain, handler, logger). + +public define Maybe($T) + generic_connect_to_net_service + ( + String queue_name, + String server, + Word32 port, + Word32 service_id, + Word32 service_version, + (MessageQueue, String) -> Maybe($T) handler, + (String) -> One logger + ) + = generic_connect_to_net_service(queue_name, server, port, service_id, service_version, "", handler, logger). public define Maybe($T) generic_connect_to_net_service_SSL @@ -260,6 +292,7 @@ public define Maybe($T) // case of an invalid, non trusted or missing certificate Word32 service_id, Word32 service_version, + String domain, (MessageQueue, String) -> Maybe($T) handler, (String) -> One logger ) @@ -270,8 +303,25 @@ public define Maybe($T) ok(conn) then with queue = create_MessageQueue(queue_name), message_transceiver(ssl(conn), queue); - with result = generic_request_for_service(queue, service_id, service_version, handler, logger), + with result = generic_request_for_service(queue, service_id, service_version, domain, handler, logger), queue.quit(unique); //logInfo(debug_log,"domain_manager client quit"); result }. + +public define Maybe($T) + generic_connect_to_net_service_SSL + ( + String queue_name, + String server_name, + Word32 server_ip, + Word32 port, + (Maybe(X509)) -> Bool accept_policy, // your policy for accepting the server certificate in + // case of an invalid, non trusted or missing certificate + Word32 service_id, + Word32 service_version, + (MessageQueue, String) -> Maybe($T) handler, + (String) -> One logger + ) + = generic_connect_to_net_service_SSL(queue_name, server_name, server_ip, port, accept_policy, service_id, service_version, "", handler, logger). + diff --git a/net_services/CXM_net_services.anubis b/net_services/CXM_net_services.anubis index dcf687b..d797a3d 100644 --- a/net_services/CXM_net_services.anubis +++ b/net_services/CXM_net_services.anubis @@ -1,190 +1,195 @@ -/* - * - * User: David RENE - * Date: 25/04/2007 - * Time: 11:01 - * (c) Calexium - * - */ -read tools/basis.anubis -read system/convert.anubis -read system/string.anubis -read system/muscle.anubis -read system/data_io.anubis -read system/message_queue.anubis -read system/message_transceiver.anubis -read CXM_generic_protocol.anubis -read calexium_lib/CXM_message_constants.anubis - -public type NetService: - net_service( - Word32 version, - Word32 id, - String name, - (MessageQueue, String, String) -> One handler // Parameters are MessageQueue, peer IP and timestamp string - ). - -define One - print_services - ( - List(NetService) net_services - ) = - map_forget((NetService net_s)|-> - println(" id : 0x"+ to_hexa(net_s.id)); - println(" version : " + to_String(net_s.version)); - println(" name : "+ net_s.name ); - println("----------------------------------------") - ,net_services). - - /** Try to find the service_id in services_list. If the service is found in that list - * the corresponding NetService object is return - */ -define Maybe(NetService) - find_service - ( - List(NetService) services_list, - Word32 service_id, - Word32 service_version - )= - if services_list is - { - [] then failure, - [h . t] then - if h.id = service_id & h.version >=+ service_version then - success(h) - else - find_service(t, service_id, service_version) - }. - - /** Check if the muscle message msg has the correct fields for requesting a net_services - * if we found "service" and "version" fields on the message, we try to find if the service - * referenced in "service" is available in net_services list - */ -define Maybe(NetService) - has_service - ( - MessageQueue queue, - Message msg, - List(NetService) net_services - )= - if find_int32(msg, "SERVICE") is - { - failure then //send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE); failure, - // old names... should be removed soon - if find_int32(msg, "service") is - { - failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE); failure, - success(service_id) then - if find_int32(msg, "version") is - { - failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE);failure, - success(service_version) then find_service(net_services, service_id, service_version) - } - } - - success(service_id) then - if find_int32(msg, "VERSION") is - { - failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE);failure, - success(service_version) then find_service(net_services, service_id, service_version) - } - }. - -define String - get_time_stamp - = - with time = (UTime) unow, - "<"+virtual_machine_id+"@"+time.seconds+">". - - /** This message_received function just handle the negociation process the available net_services. - * In other words, it only recognize the _CXM_REQUEST_FOR_SERVICE message and try to launch the - * corresponding servcice - */ - -define One - service_negociation - ( - MessageQueue queue, - Message msg, - List(NetService) net_services, - String peer - )= - //println("Service NEGOCIATION [" + to_hexa(*msg.what) + "] received"); - if * msg.what = _CXM_REQUEST_FOR_SERVICE then - if has_service(queue, msg, net_services) is - { - failure then - send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE, _CXM_UNKNOW_SERVICE, "Unknown service") - success(net_service) then - with result = message(0), - timestamp = get_time_stamp, - forget(add_string(result, "TIMESTAMP", timestamp)); - send_ACK_ok(queue, _CXM_REQUEST_FOR_SERVICE, result); - net_service.handler(queue, peer, timestamp) - } - else - send_ACK_error(queue, *msg.what, _CXM_UNKNOW_CMD, "Unknown command [" + (*msg.what) + "]") - . - - /** - * this function unflatten muscle message and give the correct message to service_negociation function - */ -define One - message_receiver - ( - MessageQueue queue, - List(NetService) net_services, - String peer - ) = - if queue.quit_requested(unique) then - unique - else - //println("PRE SERVICE message_receiver "+"["+virtual_machine_id + "]"); - if queue.get_next_received_Message(1) is - { - timeout then //println("PRE timeout"); - message_receiver(queue, net_services, peer), - closed then //println("PRE closed"); - unique, - msg(msg) then unique; //println("PRE negociation"); - service_negociation(queue, msg, net_services, peer); - message_receiver(queue, net_services, peer) - }. - -define Server -> (RWStream) -> One - net_services_handler - ( - List(NetService) net_services, - ) = - (Server server) |-> (RWStream conn) |-> - if remote_IP_address_and_port(conn) is (num_peer,_) then - //convert IP address of the client to string - with peer = ip_addr_to_string(num_peer), - //println("NET SERVICES Accepting connection with "+peer); - - //now managing the list of SERVICES - with queue = create_MessageQueue("CXM Net Services"), - message_transceiver(conn, queue); - message_receiver(queue, net_services, peer). - - -public define Maybe(Server) - start_net_services - ( - List(NetService) net_services, - Word32 network_port, - )= - if start_server(0, - network_port, - net_services_handler(net_services), - (One u) |-> unique) is - { - cannot_create_the_socket then println("Cannot create the listening socket."); failure, - cannot_bind_to_port then println("Cannot bind to port " + network_port ); failure, - cannot_listen_on_port then println("Cannot listen on port " + network_port); failure, - ok(server) then - println("Net services started on port " + network_port); - println("------ Available services ------"); - print_services(net_services); - success(server) - }. +/* + * + * User: David RENE + * Date: 25/04/2007 + * Time: 11:01 + * (c) Calexium + * + */ +read tools/basis.anubis +read system/convert.anubis +read system/string.anubis +read system/muscle.anubis +read system/data_io.anubis +read system/message_queue.anubis +read system/message_transceiver.anubis +read CXM_generic_protocol.anubis +read calexium_lib/CXM_message_constants.anubis + +public type NetService: + net_service( + Word32 version, + Word32 id, + String name, + List(String) domains, + (MessageQueue, String, String) -> One handler // Parameters are MessageQueue, peer IP and timestamp string + ). + +define One + print_services + ( + List(NetService) net_services + ) = + + map_forget((NetService net_s)|-> + if net_s is net_service(version, id, name, domains, _) then + println(" id : 0x"+ to_hexa(id)); + println(" version : " + to_String(version)); + println(" name : "+ name ); + println(" domains : "); + map_forget((String domain) |-> println(" : "+domain), domains); + println("----------------------------------------") + ,net_services). + + /** Try to find the service_id in services_list. If the service is found in that list + * the corresponding NetService object is return + */ +define Maybe(NetService) + find_service + ( + List(NetService) services_list, + Word32 service_id, + Word32 service_version + )= + if services_list is + { + [] then failure, + [h . t] then + if h.id = service_id & h.version >=+ service_version then + success(h) + else + find_service(t, service_id, service_version) + }. + + /** Check if the muscle message msg has the correct fields for requesting a net_services + * if we found "service" and "version" fields on the message, we try to find if the service + * referenced in "service" is available in net_services list + */ +define Maybe(NetService) + has_service + ( + MessageQueue queue, + Message msg, + List(NetService) net_services + )= + if find_int32(msg, "SERVICE") is + { + failure then //send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE); failure, + // old names... should be removed soon + if find_int32(msg, "service") is + { + failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE); failure, + success(service_id) then + if find_int32(msg, "version") is + { + failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE);failure, + success(service_version) then find_service(net_services, service_id, service_version) + } + } + + success(service_id) then + if find_int32(msg, "VERSION") is + { + failure then send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE);failure, + success(service_version) then find_service(net_services, service_id, service_version) + } + }. + +define String + get_time_stamp + = + with time = (UTime) unow, + "<"+virtual_machine_id+"@"+time.seconds+">". + + /** This message_received function just handle the negociation process the available net_services. + * In other words, it only recognize the _CXM_REQUEST_FOR_SERVICE message and try to launch the + * corresponding servcice + */ + +define One + service_negociation + ( + MessageQueue queue, + Message msg, + List(NetService) net_services, + String peer + )= + //println("Service NEGOCIATION [" + to_hexa(*msg.what) + "] received"); + if * msg.what = _CXM_REQUEST_FOR_SERVICE then + if has_service(queue, msg, net_services) is + { + failure then + send_ACK_error(queue, _CXM_REQUEST_FOR_SERVICE, _CXM_UNKNOW_SERVICE, "Unknown service") + success(net_service) then + with result = message(0), + timestamp = get_time_stamp, + forget(add_string(result, "TIMESTAMP", timestamp)); + send_ACK_ok(queue, _CXM_REQUEST_FOR_SERVICE, result); + net_service.handler(queue, peer, timestamp) + } + else + send_ACK_error(queue, *msg.what, _CXM_UNKNOW_CMD, "Unknown command [" + (*msg.what) + "]") + . + + /** + * this function unflatten muscle message and give the correct message to service_negociation function + */ +define One + message_receiver + ( + MessageQueue queue, + List(NetService) net_services, + String peer + ) = + if queue.quit_requested(unique) then + unique + else + //println("PRE SERVICE message_receiver "+"["+virtual_machine_id + "]"); + if queue.get_next_received_Message(1) is + { + timeout then //println("PRE timeout"); + message_receiver(queue, net_services, peer), + closed then //println("PRE closed"); + unique, + msg(msg) then unique; //println("PRE negociation"); + service_negociation(queue, msg, net_services, peer); + message_receiver(queue, net_services, peer) + }. + +define Server -> (RWStream) -> One + net_services_handler + ( + List(NetService) net_services, + ) = + (Server server) |-> (RWStream conn) |-> + if remote_IP_address_and_port(conn) is (num_peer,_) then + //convert IP address of the client to string + with peer = ip_addr_to_string(num_peer), + //println("NET SERVICES Accepting connection with "+peer); + + //now managing the list of SERVICES + with queue = create_MessageQueue("CXM Net Services"), + message_transceiver(conn, queue); + message_receiver(queue, net_services, peer). + + +public define Maybe(Server) + start_net_services + ( + List(NetService) net_services, + Word32 network_port, + )= + if start_server(0, + network_port, + net_services_handler(net_services), + (One u) |-> unique) is + { + cannot_create_the_socket then println("Cannot create the listening socket."); failure, + cannot_bind_to_port then println("Cannot bind to port " + network_port ); failure, + cannot_listen_on_port then println("Cannot listen on port " + network_port); failure, + ok(server) then + println("Net services started on port " + network_port); + println("------ Available services ------"); + print_services(net_services); + success(server) + }. diff --git a/web/CXM_multihost_http_server.anubis b/web/CXM_multihost_http_server.anubis index 4570c25..b4d8a1e 100644 --- a/web/CXM_multihost_http_server.anubis +++ b/web/CXM_multihost_http_server.anubis @@ -857,7 +857,7 @@ define ReadResult //if unow > dead_line then record_dubious_connection(connection,dead_line,dos) else if read(connection.conn, 16384, time_out) is // the connection is closed after 10 minutes of inactivity { - error then println(pid + "read failed)"); error, + error then println(pid + "read failed ["+to_string(result_buffer)+"]"); error, timeout then timeout, ok(ba) then // println(pid + "ba = " + length(ba)); -- libgit2 0.21.4