diff --git a/net_services/CXM_generic_client.anubis b/net_services/CXM_generic_client.anubis index cd64a55..4ded637 100644 --- a/net_services/CXM_generic_client.anubis +++ b/net_services/CXM_generic_client.anubis @@ -260,8 +260,8 @@ public define Maybe(MessageQueue) error(_) then logError(logger, queue_name + ": Can't connect to service ["+ip_addr_to_string(server)+":"+port+"]");failure, ok(conn) then // println("[" + virtual_machine_id + "] netservices create queue"); - with queue = create_MessageQueue(queue_name), - message_transceiver(tcp(conn), queue); + with queue = create_MessageQueue(queue_name, tcp(conn)), + message_transceiver(/*tcp(conn),*/ queue); // println("[" + virtual_machine_id + "] netservices generic_request_for_service()"); if generic_request_for_service(queue, service_id, service_version, domain, handler, logger) is { @@ -293,8 +293,8 @@ public define Maybe($T) error(_) then logger(queue_name + ": Can't connect to service ["+ip_addr_to_string(server)+":"+port+"]");failure, ok(conn) then // println("[" + virtual_machine_id + "] netservices create queue"); - with queue = create_MessageQueue(queue_name), - message_transceiver(tcp(conn), queue); + with queue = create_MessageQueue(queue_name, tcp(conn)), + message_transceiver(queue); // println("[" + virtual_machine_id + "] netservices generic_request_for_service()"); with result = generic_request_for_service(queue, service_id, service_version, domain, handler, logger), // println("[" + virtual_machine_id + "] netservices client quit"); @@ -414,8 +414,8 @@ public define Maybe($T) { error(_) then logger(queue_name + ": Can't connect to domain manager ["+ip_addr_to_string(server_ip)+":"+port+"]");failure, ok(conn) then - with queue = create_MessageQueue(queue_name), - message_transceiver(ssl(conn), queue); + with queue = create_MessageQueue(queue_name, ssl(conn)), + message_transceiver(/*ssl(conn),*/ queue); with result = generic_request_for_service(queue, service_id, service_version, domain, handler, logger), queue.quit(unique); //logInfo(debug_log,"domain_manager client quit"); diff --git a/net_services/CXM_net_services.anubis b/net_services/CXM_net_services.anubis index 872d54b..451b016 100644 --- a/net_services/CXM_net_services.anubis +++ b/net_services/CXM_net_services.anubis @@ -1,213 +1,213 @@ -/* - * - * 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, - String domain - )= - if services_list is - { - [] then failure, - [h . t] then - if h.id = service_id & h.version >=+ service_version then - if domain = "" then - success(h) - else if domain:h.domains then //this writing (a:b) means, is a belonging to b where b is list of type a - success(h) - else - find_service(t, service_id, service_version, domain) - else - find_service(t, service_id, service_version, domain) - }. - - /** 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 - //if DOMAIN field exists, this mean we want to target only this domain - if find_string(msg, "DOMAIN") is - { - failure then find_service(net_services, service_id, service_version,""), - success(domain) then find_service(net_services, service_id, service_version, domain) - } - } - } - - 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 - //if DOMAIN field exists, this mean we want to target only this domain - if find_string(msg, "DOMAIN") is - { - failure then find_service(net_services, service_id, service_version,""), - success(domain) then find_service(net_services, service_id, service_version, domain) - } - } - }. - -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 - */ -public 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, + String domain + )= + if services_list is + { + [] then failure, + [h . t] then + if h.id = service_id & h.version >=+ service_version then + if domain = "" then + success(h) + else if domain:h.domains then //this writing (a:b) means, is a belonging to b where b is list of type a + success(h) + else + find_service(t, service_id, service_version, domain) + else + find_service(t, service_id, service_version, domain) + }. + + /** 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 + //if DOMAIN field exists, this mean we want to target only this domain + if find_string(msg, "DOMAIN") is + { + failure then find_service(net_services, service_id, service_version,""), + success(domain) then find_service(net_services, service_id, service_version, domain) + } + } + } + + 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 + //if DOMAIN field exists, this mean we want to target only this domain + if find_string(msg, "DOMAIN") is + { + failure then find_service(net_services, service_id, service_version,""), + success(domain) then find_service(net_services, service_id, service_version, domain) + } + } + }. + +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 + */ +public 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", tcp(conn)), + 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/net_services_protocols/ftp_client.anubis b/net_services_protocols/ftp_client.anubis index 010c06b..69aa469 100644 --- a/net_services_protocols/ftp_client.anubis +++ b/net_services_protocols/ftp_client.anubis @@ -154,8 +154,8 @@ public define Maybe(One) { error(_) then println("can't connect to ftp server"); failure, ok(conn) then - with queue = create_MessageQueue("ftp_get_file"), - message_transceiver(conn, queue); + with queue = create_MessageQueue("ftp_get_file", tcp(conn)), + message_transceiver(/*conn,*/ queue); println("request for ftp service"); if request_for_service(queue) then with result = get_file(queue, remote_file, local_file, ftp), -- libgit2 0.21.4