/* * Created by PyramIDE. * User: ricard * Date: 02/02/2008 * Time: 11:47 * * To change this template use Tools | Options | Coding | Edit Standard Headers. */ read tools/basis.anubis read system/muscle.anubis read system/data_io.anubis read system/convert.anubis read system/string.anubis read system/message_queue.anubis read system/message_transceiver.anubis read tools/connections.anubis read calexium_lib/CXM_message_constants.anubis read calexium_lib/net_services/CXM_net_services.anubis read calexium_lib/net_services/CXM_generic_protocol.anubis // --Generic types--------------------------------------------------------------------- public type NetServiceAnswer: netservice_error (Word32 cmd, Word32 result_code, String result_string), netservice_ok (Word32 cmd, Maybe(Message) result_msg). // --Generic functions--------------------------------------------------------------------- /** * Sends the message to the server then parses the answer and returns the RESULT message on CMD success. */ public define Maybe(NetServiceAnswer) generic_send_message ( MessageQueue queue, Message msg_to_send, Int timeout, (String) -> One logger )= queue.add_Message_to_send(msg_to_send); if queue.get_next_received_Message(timeout) is { timeout then logger("["+queue.get_name(unique)+"]: receive timeout");failure, closed then logger("["+queue.get_name(unique)+"]: socket closed");failure, msg(msg) then if find_int32(msg, "CMD") is { failure then logger("["+queue.get_name(unique)+"]: CMD field not found"); failure, success(cmd) then if find_int32(msg, "STATUS") is { failure then logger("["+queue.get_name(unique)+"]: STATUS field not found"); failure, success(v) then if v = _CXM_OK then success(netservice_ok(cmd, find_message(msg, "RESULT"))) else with error_string = if find_string(msg, "STATUS_MSG") is success(s) then s else "", success(netservice_error(cmd, v, error_string)) } } }. public define Maybe($T) simple_handler ( MessageQueue queue, String timestamp, Int timeout, Message msg_to_send, (Message, MessageQueue, String) -> Maybe($T) handler, (String) -> One logger ) = if generic_send_message(queue, msg_to_send, timeout, logger) is { failure then failure, success(net_result) then if net_result is { netservice_error(cmd, err_code, err_str) then logger("["+queue.get_name(unique)+"]: message status ERROR [0x" + to_hexa(err_code) + ", '" + err_str + "']"); failure, netservice_ok(cmd, mb_msg) then if mb_msg is { failure then logger("["+queue.get_name(unique)+"]: can't find RESULT message."); failure, success(result) then handler(result, queue, timestamp) } } }. public define Maybe(One) no_result_handler ( MessageQueue queue, Int timeout, Message msg_to_send, (String) -> One logger ) = if generic_send_message(queue, msg_to_send, timeout, logger) is { failure then failure, success(net_result) then if net_result is { netservice_error(cmd, err_code, err_str) then logger("["+queue.get_name(unique)+"]: message status ERROR [0x" + to_hexa(err_code) + ", '" + err_str + "']"); failure, netservice_ok(cmd, mb_msg) then if mb_msg is { failure then unique, success(result) then logger("An unattended RESULT msg was found. Ignoring it...") }; success(unique) } }. public define (MessageQueue, String) -> Maybe($T) make_generic_handler ( Int timeout, Message msg_to_send, (Message, MessageQueue, String) -> Maybe($T) handler, (String) -> One logger ) = (MessageQueue queue, String timestamp) |-> simple_handler(queue, timestamp, timeout, msg_to_send, handler, logger). define Maybe($T) generic_request_for_service ( MessageQueue queue, Word32 service_id, Word32 service_version, (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)); queue.add_Message_to_send(test_msg); if queue.get_next_received_Message(10) is { timeout then logger("["+queue.get_name(unique)+"]: requesting service receive timeout");failure, closed then logger("["+queue.get_name(unique)+"]: requesting service socket closed");failure, msg(msg) then if find_int32(msg, "STATUS") is { failure then logger("["+queue.get_name(unique)+"]: requesting service STATUS not found");failure, success(v) then if v = _CXM_OK then if find_message(msg, "RESULT") is { failure then logger("["+queue.get_name(unique)+"]: requesting service RESULT not found");failure, success(result) then with timestamp = if find_string(result, "TIMESTAMP") is { failure then logger("["+queue.get_name(unique)+"]: requesting service TIMESTAMP not found"); "", success(timestamp) then timestamp }, handler(queue, timestamp) } else logger("["+queue.get_name(unique)+"]: the requested service is not available on server.");failure } }. 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 ) = if connect( server, port) is { error(_) then logger(queue_name + ": Can't connect to domain manager ["+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); // println("[" + virtual_machine_id + "] netservices generic_request_for_service()"); with result = generic_request_for_service(queue, service_id, service_version, handler, logger), // println("[" + virtual_machine_id + "] netservices client quit"); queue.quit(unique); 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 ) = if open_SSL_connection( server_name, server_ip, port, accept_policy) is { 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 result = generic_request_for_service(queue, service_id, service_version, handler, logger), queue.quit(unique); //logInfo(debug_log,"domain_manager client quit"); result }.