CXM_generic_client.anubis 16.6 KB
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442
/*
 * Created by PyramIDE.
 * User: ricard
 * Date: 02/02/2008
 * Time: 11:47
 * 
 *
 */

transmit system/message_transceiver.anubis
transmit system/logger.anubis
transmit calexium_lib/net_services_protocols/logger_service.anubis

transmit calexium_lib/CXM_message_constants.anubis
transmit calexium_lib/net_services/CXM_net_services.anubis
transmit 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(NetServiceAnswer)
/**
 * Sends the message to the server then parses the answer and returns the RESULT message on CMD success.
 * 
 */
  generic_send_message
  (
    MessageQueue              queue,
    Message                   msg_to_send,
    Int                       timeout,
    (LogLevel, String) -> One logger
  )=
  queue.add_Message_to_send(msg_to_send);
  if queue.get_next_received_Message(timeout) is
  {
    timeout then logger(logError, "["+queue.get_name(unique)+"]: receive timeout");failure,
    closed  then logger(logError, "["+queue.get_name(unique)+"]: socket closed");failure,
    msg(msg)  then
      if find_int32(msg, "CMD") is
      {
        failure       then logger(logError, "["+queue.get_name(unique)+"]: CMD field not found"); failure,
        success(cmd)  then 
          if find_int32(msg, "STATUS") is
          {
            failure     then logger(logError, "["+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,
    (LogLevel, 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(logError, "["+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(logError,"["+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,
    (LogLevel, 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(logError, "["+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(logError, "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,
    (LogLevel, 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,
    (LogLevel, 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(logError, "["+queue.get_name(unique)+"]: requesting service receive timeout");failure,
    closed  then logger(logError, "["+queue.get_name(unique)+"]: requesting service socket closed");failure,
    msg(msg)  then 
      if find_int32(msg, "STATUS") is
      {
        failure     then logger(logError, "["+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(logError, "["+queue.get_name(unique)+"]: requesting service RESULT not found");failure,
              success(result)  then 
                with timestamp =  if find_string(result, "TIMESTAMP") is
                                  {
                                    failure             then logger(logError, "["+queue.get_name(unique)+"]: requesting service TIMESTAMP not found"); "",
                                    success(timestamp)  then timestamp
                                  },
                handler(queue, timestamp)
            }
          else
            logger(logError, "["+queue.get_name(unique)+"]: the requested service is not available on server.");failure        
      }
  }.

// define Maybe($T)
//  generic_request_for_service
//  (
//    MessageQueue                queue,
//    Word32                      service_id,
//    Word32                      service_version,
//    String                      domain,
//    (MessageQueue, String) -> Maybe($T) handler,
//    (LogLevel, 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(logError, "["+queue.get_name(unique)+"]: requesting service receive timeout");failure,
//    closed  then logger(logError, "["+queue.get_name(unique)+"]: requesting service socket closed");failure,
//    msg(msg)  then 
//      if find_int32(msg, "STATUS") is
//      {
//        failure     then logger(logError, "["+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(logError, "["+queue.get_name(unique)+"]: requesting service RESULT not found");failure,
//              success(result) then 
//                with timestamp =  if find_string(result, "TIMESTAMP") is
//                                  {
//                                    failure             then logger(logError, "["+queue.get_name(unique)+"]: requesting service TIMESTAMP not found"); "",
//                                    success(timestamp)  then timestamp
//                                  },
//                handler(queue, timestamp)
//            }
//          else
//            logger(logError, "["+queue.get_name(unique)+"]: the requested service is not available on server.");failure        
//      }
//  }.
//  
public define Maybe(MessageQueue)
  get_message_queue_to_net_service
  (
    String                      queue_name,
    Word32                      server,
    Word32                      port,
    Word32                      service_id,
    Word32                      service_version,
    String                      domain,
    (MessageQueue, String) -> Maybe(One) handler, //1st function to apply if need (i.e authentication to remote service)
    (LogLevel, String) -> One   logger
  )
  =
  if connect( server, port) is
  {
    error(_)  then logger(logError, 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, 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 
      {
        failure     then  
          //we can't apply first function correctly, so we ask to Message Queue to quit and return failure
          logger(logError, "can't apply first function correctly");
          queue.quit(unique); 
          failure,
        success(_)  then  success(queue)
      }
  }.


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,
    (LogLevel, String) -> One   logger
  )
  =
  if connect( server, port) is
  {
    error(_)  then logger(logError, 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, 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");
      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,
    (LogLevel, 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,
    (LogLevel, String) -> One   logger
  ) =
  logger(logTrace, "Try to resolve URL [" + url_or_ip+ "].");
  if ip_address(url_or_ip) is success(ip) then
    logger(logTrace, "Try to resolve URL OK ip"+ip);
    ip
  else if dns(url_or_ip) is ok(ip_adr) then
    logger(logTrace, "Try to resolve URL OK dns"+ip_adr);
    ip_adr
  else
    logger(logError, "Can't resolve URL [" + url_or_ip+ "]. Using localhost (127.0.0.1) .");
    ip_address((127,0,0,1)).

  
public define Maybe(Word32)
  mb_get_ip
  ( 
    String                    url_or_ip,
    (LogLevel, String) -> One logger
  ) =
  logger(logTrace, "Try to resolve URL [" + url_or_ip+ "].");
  if ip_address(url_or_ip) is success(ip) then
    logger(logTrace, "Try to resolve URL OK ip"+ip);
    success(ip)
  else if dns(url_or_ip) is ok(ip_adr) then
    logger(logTrace, "Try to resolve URL OK dns"+ip_adr);
    success(ip_adr)
  else
    logger(logError, "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,
    (LogLevel, 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,
    (LogLevel, 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,
    (LogLevel, String) -> One   logger
  )
  =
  if open_SSL_connection( server_name, server_ip, port, accept_policy) is
  {
    error(_)  then logger(logError, 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, 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");
      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,
    (LogLevel, String) -> One   logger
  )
  = generic_connect_to_net_service_SSL(queue_name, server_name, server_ip, port, accept_policy, service_id, service_version, "", handler, logger).