/* * * 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 system/logger.anubis read CXM_generic_protocol.anubis read xlib/CXM_message_constants.anubis public type NetService: net_service( Word32 version, Word32 id, String name, List(String) domains, (MessageQueue, String, String, (LogLevel, String) -> One ) -> One handler // Parameters are MessageQueue, peer IP and timestamp string, logger ) . define String dump_services ( List(NetService) net_services ) = join("\n", map((NetService net_s) |-> if net_s is net_service(version, id, name, domains, _) then " id : 0x"+ to_hexa(id)+"\n"+ " version : " + to_String(version)+"\n"+ " name : "+ name +"\n"+ " domains : "+"\n"+ join("\n",map((String domain) |-> " : "+domain, domains))+"\n"+ "----------------------------------------\n" ,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, //List of all available NetService Word32 service_id, //requested service ID Word32 service_version, //requested service Version String domain //requested domain for above resquested service ID/version )= 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, (LogLevel, String) -> One logger )= logger(logTrace, "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 logger(logError, "Unknown service"); 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, logger) } 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, (LogLevel, String) -> One logger ) = if queue.quit_requested(unique) then unique else logger(logTrace,"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, logger), closed then logger(logTrace,"PRE SERVICE message_receiver ["+virtual_machine_id + "] closed"); unique, msg(msg) then unique; //println("PRE negociation"); service_negociation(queue, msg, net_services, peer, logger); message_receiver(queue, net_services, peer, logger) }. define Server -> (RWStream) -> One net_services_handler ( List(NetService) net_services, (LogLevel, String) -> One logger ) = (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), logger(logInfo,"NET SERVICES Accepting connection with "+peer); //now managing the list of SERVICES with queue = create_MessageQueue("CXM Net Services", tcp(conn)), message_transceiver(queue); message_receiver(queue, net_services, peer, logger). public define Maybe(Server) start_net_services ( List(NetService) net_services, Word32 network_port, (LogLevel, String) -> One logger )= if start_server(0, network_port, net_services_handler(net_services, logger), (One u) |-> unique) is { cannot_create_the_socket then logger(logError, "Cannot create the listening socket."); failure, cannot_bind_to_port then logger(logError, "Cannot bind to port " + network_port); failure, cannot_listen_on_port then logger(logError, "Cannot listen on port " + network_port); failure, ok(server) then logger(logInfo, "Net services started on port " + network_port); logger(logInfo, "------ Available services ------"); logger(logInfo, dump_services(net_services)); success(server) } .