CXM_generic_client.anubis 11.9 KB
/*
 * Created by PyramIDE.
 * User: ricard
 * Date: 02/02/2008
 * Time: 11:47
 * 
 *
 */


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,
    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
  {
    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,
    String                      domain,
    (MessageQueue, String) -> Maybe($T) handler,
    (String) -> One             logger
  )
  =
  if connect( server, port) is
  {
    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);
//      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");
      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
  ( 
    String          url_or_ip,
    (String) -> One logger
  ) =
  //logInfo(debug_log,"Try to resolve URL [" + url_or_ip+ "].");
  if ip_address(url_or_ip) is success(ip) then
    //logInfo(debug_log,"Try to resolve URL OK ip"+ip);
    ip
  else if dns(url_or_ip) is ok(ip_adr) then
    //logInfo(debug_log,"Try to resolve URL OK dns"+ip_adr);
    ip_adr
  else
    logger("Can't resolve URL [" + url_or_ip+ "]. Using localhost.");
    ip_address((127,0,0,1)).

  
public define Maybe(Word32)
  mb_get_ip
  ( 
    String          url_or_ip
  ) =
  //logInfo(debug_log,"Try to resolve URL [" + url_or_ip+ "].");
  if ip_address(url_or_ip) is success(ip) then
    //logInfo(debug_log,"Try to resolve URL OK ip"+ip);
    success(ip)
  else if dns(url_or_ip) is ok(ip_adr) then
    //logInfo(debug_log,"Try to resolve URL OK dns"+ip_adr);
    success(ip_adr)
  else
    println("Can't resolve URL [" + url_or_ip+ "].");
    failure.

public define Word32
  get_ip
  ( 
    String          url_or_ip
  ) =
  if ip_address(url_or_ip) is success(ip) then  ip
  else if dns(url_or_ip) is ok(ip_adr)    then  ip_adr
  else
    println("Can't resolve URL [" + url_or_ip+ "]. Using localhost.");
    ip_address((127,0,0,1)).
    
  /* Same version as above, but server is string containing IP or URL
   * It's resolve by get_ip and call generic_connect_to_net_service with IP
   */
   
public define Maybe($T)
  generic_connect_to_net_service
  (
    String                        queue_name,
    String                        server,
    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, 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
  (
    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,
    String                      domain,
    (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, 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).