/* * * User: フランスのトトロ aka (David RENÉ) * Date: 25/04/2007 * Time: 16:20 * © David RENÉ * */ read system/muscle.anubis read system/data_io.anubis read system/string.anubis read system/files.anubis read tools/basis.anubis read system/message_queue.anubis read xlib/message_constants.anubis read xlib/types/generated/file_ref.anubis public define Word32 _CXM_OK = 0. public define Word32 _CXM_ERROR = 1. public define Word32 _CXM_UNKNOW_CMD = 2. public define Word32 _CXM_UNKNOW_SERVICE = 3. public define Word32 _CXM_MISSING_REQUIRED_FIELD = 4. public define Word32 _CXM_FORBIDDEN = 5. public define Word32 _CXM_BAD_AUTHENTICATION = 6. public define Word32 _CXM_TEMPORARY_ERROR = 7. // When received, the client should try later public type ProtocolResult: failure, timeout, unknow_cmd, error, error(Word32, String), ok, ok_msg(Message). public define One send_ACK_error ( MessageQueue queue, Word32 cmd_id, Word32 error_code, String error_string, )= with err_msg = message(_CXM_ACK), forget(add_int32(err_msg, "CMD", cmd_id)); forget(add_int32(err_msg, "STATUS", error_code)); (if error_string /= "" then forget(add_string(err_msg, "STATUS_MSG", error_string)) else unique); forget(queue.add_Message_to_send(err_msg)). public define One send_ACK_error ( MessageQueue queue, Word32 cmd_id )= send_ACK_error(queue, cmd_id, _CXM_ERROR, ""). public define One send_ACK_ok ( MessageQueue queue, Word32 cmd_id )= with ok_msg = message(_CXM_ACK), forget(add_int32(ok_msg, "CMD", cmd_id)); forget(add_int32(ok_msg, "STATUS", _CXM_OK)); forget(queue.add_Message_to_send(ok_msg)) . public define One send_ACK_ok ( MessageQueue queue, Word32 cmd_id, Message result )= with ok_msg = message(_CXM_ACK), forget(add_int32(ok_msg, "CMD", cmd_id)); forget(add_int32(ok_msg, "STATUS", _CXM_OK)); forget(add_message(ok_msg, "RESULT", result)); forget(queue.add_Message_to_send(ok_msg)) . public define One send_result ( MessageQueue queue, Word32 cmd_id, Maybe(Message) mb_msg )= if mb_msg is { failure then send_ACK_error(queue, cmd_id), success(msg) then send_ACK_ok(queue, cmd_id, msg) }. public define One send_result ( MessageQueue queue, Word32 cmd_id, Result((Word32, String), Message) mb_msg )= if mb_msg is { error(err) then if err is (err_code, err_string) then send_ACK_error(queue, cmd_id, err_code, err_string), ok(msg) then send_ACK_ok(queue, cmd_id, msg) }. public define One send_result ( MessageQueue queue, Word32 cmd_id, Bool result )= if result then send_ACK_ok(queue, cmd_id) else send_ACK_error(queue, cmd_id). public define ProtocolResult wait_for_reply ( MessageQueue mQ, Word32 wait_cmd, Int t_out ) = if mQ.get_next_received_Message(t_out) is { timeout then timeout, closed then failure, //println("wait_for_reply closed"); msg(_msg) then if *_msg.what = _CXM_ACK then if find_int32(_msg, "CMD") is { failure then failure, //println("wait_for_reply CMD"); success(cmd) then // println("wait_for_reply CMD="+to_hexa(cmd)); if find_int32(_msg, "STATUS") is { failure then failure, //println("wait_for_reply STATUS"); success(status) then if cmd = wait_cmd then ( if status = _CXM_OK then if find_message(_msg, "RESULT") is { failure then ok, success(ok_message) then ok_msg(ok_message) } else if status = _CXM_ERROR then error else if status = _CXM_UNKNOW_CMD then unknow_cmd else error(status, if find_string(_msg, "STATUS_MSG") is success(txt) then txt else "") ) else failure //println("wait_for_reply "); } } else failure //println("wait_for_reply not ACK"); }. public define Bool simple_wait_for_reply ( MessageQueue mQ, Word32 wait_cmd, Int t_out ) = if wait_for_reply(mQ, wait_cmd, t_out) is { failure then false, timeout then false, unknow_cmd then false, error then false, error(_, _) then false, ok then true, ok_msg(msg) then true }. public define Bool get_file_ref ( MessageQueue mQ, File_ref f_ref, String tmp_path )= //create the target with target_file = tmp_path + "/" + f_ref.name, //println("NET_SERVICE get_file_ref for "+target_file); with msg = message(_CXM_GET_FILE_REF), forget(add_message(msg, "GET_FILE_REF", to_Message(f_ref))); forget(mQ.add_Message_to_send(msg)); if wait_for_reply(mQ, _CXM_GET_FILE_REF, 30) is { failure then false, timeout then false, unknow_cmd then false, error then false, error(_, _) then false, ok then false, //false because OK without result message is not allowed ok_msg(rmsg) then //println("get_file_ref wait_for_reply _CXM_GET_FILE_REF OK"); with size = if find_string(rmsg, "SIZE") is {failure then f_ref.size, success(size_str) then if decimal_scan(size_str) is { failure then should_not_happen(0), success(_size_) then _size_}}, with mode = find_string(rmsg, "MODE", "NEW"), if mode = "APPEND" then if (Maybe(RWStream))file(target_file, append) is { failure then println("can't create target file"+target_file);false, //nothing to write success(target) then with buffer = mQ.raw_mode_on(size), println("raw_mode ON buffer len "+length(buffer)+" buffer ["+to_string(buffer)+"]"); forget(flush(buffer, weaken(target))); if copy_file_to_Connection(mQ.get_connection(unique), file(target), size - length(buffer)) is { failure then mQ.raw_mode_off(unique);false, success(_) then mQ.raw_mode_off(unique);true } } else //by default the mode is new if (Maybe(RWStream))file(target_file, new) is { failure then println("can't create target file"+target_file);false, //nothing to write success(target) then with buffer = mQ.raw_mode_on(size), println("raw_mode ON buffer len "+length(buffer)+" buffer ["+to_string(buffer)+"]"); forget(flush(buffer, weaken(target))); if copy_file_to_Connection(mQ.get_connection(unique), file(target), size - length(buffer)) is { failure then mQ.raw_mode_off(unique);false, success(_) then mQ.raw_mode_off(unique);true } } } .