#!/usr/bin/python3 # -*- coding: utf-8 -*- import asyncio import json import re import traceback from datetime import datetime import aiohttp from panoramisk import Manager # Settings IP = '127.0.0.1' # Default USERNAME = 'Your_user_name' # AMI Name SECRET = 'Your_user_password' # AMI Secret URL = 'Your_planfix_endpoint_url' # Planfix endpoint URL DEBUG = False DEBUG_CALL_FILE_PATH = "/etc/asterisk/scripts/ami_listener_debug_call.txt" # Path to the file with calls logs DEBUG_LOG_FILE_PATH = "/etc/asterisk/scripts/ami_listener_debug_log.txt" # Path to the file with script logs def to_json(message): pairs = re.findall(r"(\S+)='([^']+)'", message) return {key: value for key, value in pairs} class AMIListener: def __init__(self, ip, username, secret, url, debug, debug_call_file_path, debug_log_file_path): self.ip = ip self.username = username self.secret = secret self.url = url self.debug = debug self.debug_call_file_path = debug_call_file_path self.debug_log_file_path = debug_log_file_path self.manager = None async def handle_newchannel(self, manager, message): await self.handle_debug_all_events(manager, message) data = to_json(str(message)) unique_id = data.get('Uniqueid', '') linked_id = data.get('Linkedid', '') channel_state = data.get('ChannelState', '') caller_id_num = data.get('CallerIDNum', '') if not unique_id: await self.handle_log_if_needed("Missing 'Uniqueid' in message.", "WARNING") return if not linked_id: await self.handle_log_if_needed("Missing 'Linkedid' in message.", "WARNING") return if unique_id != linked_id: await self.handle_log_if_needed(f"'Uniqueid' ({unique_id}) doesn't match 'Linkedid' ({linked_id}).", "INFO") return if not channel_state: await self.handle_log_if_needed("Missing 'ChannelState' in message.", "INFO") return if channel_state != '4': await self.handle_log_if_needed(f"'ChannelState' is not '4': received {channel_state}.", "INFO") #return if not caller_id_num: await self.handle_log_if_needed("Missing 'CallerIDNum' in message.", "WARNING") return req_mes = { 'event': 'newchannel', 'uniqueId': unique_id, 'callerNum': caller_id_num } await self.send_to_crm(req_mes) async def handle_dialbegin(self, manager, message): await self.handle_debug_all_events(manager, message) data = to_json(str(message)) linked_id = data.get('Linkedid', '') dest_caller_id_num = data.get('DestCallerIDNum', '') channel = data.get('Channel', '') if not linked_id: await self.handle_log_if_needed("Missing 'Linkedid' in message.", "WARNING") return if not dest_caller_id_num: await self.handle_log_if_needed("Missing 'DestCallerIDNum' in message.", "WARNING") return if not channel: await self.handle_log_if_needed("Missing 'Channel' in message.", "WARNING") return direction = await self.get_direction(data) if not direction: await self.handle_log_if_needed("Call direction - unknown", "WARNING") return rec_name = await self.get_recording_file(channel) if not rec_name: await self.handle_log_if_needed("Call recording name - unknown", "WARNING") line_number = await self.get_line_number(data, channel, direction, dest_caller_id_num) req_mes = { 'event': 'dialbegin', 'uniqueId': linked_id, 'direction': direction, 'destCallerNum': dest_caller_id_num, 'callFileName': rec_name } if line_number: req_mes['lineNumber'] = line_number await self.send_to_crm(req_mes) async def handle_dialend(self, manager, message): await self.handle_debug_all_events(manager, message) data = to_json(str(message)) linked_id = data.get('Linkedid', '') dest_caller_id_num = data.get('DestCallerIDNum', '') channel = data.get('Channel', '') dial_status = data.get('DialStatus', '') if not linked_id: await self.handle_log_if_needed("Missing 'Linkedid' in message.", "WARNING") return if not dest_caller_id_num: await self.handle_log_if_needed("Missing 'DestCallerIDNum' in message.", "WARNING") return if not channel: await self.handle_log_if_needed("Missing 'Channel' in message.", "WARNING") return if not dial_status: await self.handle_log_if_needed("Missing 'DialStatus' in message.", "WARNING") return direction = await self.get_direction(data) if not direction: await self.handle_log_if_needed("Call direction - unknown", "WARNING") return rec_name = await self.get_recording_file(channel) if not rec_name: await self.handle_log_if_needed("Call recording name - unknown", "WARNING") line_number = await self.get_line_number(data, channel, direction, dest_caller_id_num) req_mes = { 'event': 'dialend', 'uniqueId': linked_id, 'direction': direction, 'dialStatus': dial_status, 'destCallerNum': dest_caller_id_num, 'callFileName': rec_name } if line_number: req_mes['lineNumber'] = line_number await self.send_to_crm(req_mes) async def handle_cdr(self, manager, message): await self.handle_debug_all_events(manager, message) data = to_json(str(message)) unique_id = data.get('UniqueID', '') linked_id = data.get('Linkedid', '') disposition = data.get('Disposition', '') answer_time = data.get('AnswerTime', data.get('StartTime', '')) duration = data.get('BillableSeconds', '') rec_name = await self.try_to_compile_rec_name(data) if not unique_id: await self.handle_log_if_needed("Missing 'UniqueID' in message.", "WARNING") return if disposition != 'ANSWERED': await self.handle_log_if_needed("'Disposition' is not 'ANSWERED'.", "INFO") return if not answer_time: await self.handle_log_if_needed("Missing 'AnswerTime' in message.", "WARNING") return if not duration: await self.handle_log_if_needed("Missing 'BillableSeconds' in message.", "WARNING") return req_mes = { 'event': 'cdr', 'uniqueId': unique_id, 'disposition': disposition, 'answerTime': answer_time, 'billableSeconds': duration, 'callFileName': rec_name, } if linked_id: req_mes['linkedId'] = linked_id await self.send_to_crm(req_mes) async def send_to_crm(self, data): await self.handle_log_if_needed(f"Sending to CRM: {data}", "INFO") headers = {'Content-type': 'application/json', 'Accept': 'text/plain'} body = json.dumps(data) async with aiohttp.ClientSession() as session: async with session.post(self.url, data=body, headers=headers) as response: response_text = await response.text() await self.handle_log_if_needed(f"CRM response: {response_text}", "INFO") async def get_ami_variable(self, channel, variable_name): response = await self.manager.send_action({ 'Action': 'GetVar', 'Channel': channel, 'Variable': variable_name }) return response.get('Value', '') async def try_to_compile_rec_name(self, data): unique_id = data.get('UniqueID', '') destination = data.get('Destination', '') cnum = data.get('cnum', '') start_time = data.get('StartTime', '') dt_obj = datetime.strptime(start_time, "%Y-%m-%d %H:%M:%S") formatted_str = dt_obj.strftime("%Y%m%d-%H%M%S") rec_name = "" if destination and destination != '' and cnum and cnum != '': if destination != cnum: if len(destination) <= 5: rec_name = f"in-{destination}-{cnum}-{formatted_str}-{unique_id}" await self.handle_log_if_needed(f"Variable 'destination' is {destination}, setting rec_name - {rec_name}", "INFO") else: rec_name = f"out-{destination}-{cnum}-{formatted_str}-{unique_id}" await self.handle_log_if_needed(f"Variable 'destination' is {destination}, setting rec_name - {rec_name}", "INFO") else: await self.handle_log_if_needed(f"Variable 'destination' equal 'cnum' = {destination}.", "INFO") else: await self.handle_log_if_needed("Variables 'destination' or 'cnum' is empty.", "INFO") return rec_name async def get_recording_file(self, channel): rec_name = await self.get_ami_variable(channel, 'CALLFILENAME') if not rec_name: await self.handle_log_if_needed("Variable 'CALLFILENAME' is empty.", "INFO") if not rec_name: rec_name = await self.get_ami_variable(channel, 'CDR(recordingfile)') if not rec_name: await self.handle_log_if_needed("Variable 'CDR(recordingfile)' is empty.", "INFO") return rec_name async def get_line_number(self, data, channel, direction, dest_caller_id_num): if direction != 'OUTBOUND': return '' line_number = data.get('CallerIDNum', '') if line_number == '': line_number = '' if not line_number or len(line_number) <= 3: current_caller_id_num = await self.get_ami_variable(channel, 'CALLERID(num)') if current_caller_id_num and current_caller_id_num != '': line_number = current_caller_id_num if not line_number: await self.handle_log_if_needed("Outgoing line number is empty.", "INFO") return '' if len(line_number) <= 3: await self.handle_log_if_needed(f"Outgoing line number looks internal = {line_number}.", "INFO") return '' if dest_caller_id_num and line_number == dest_caller_id_num: await self.handle_log_if_needed(f"Outgoing line number equals destination number = {line_number}.", "INFO") return '' return line_number async def get_direction(self, data): channel = data.get('Channel', '') direction = await self.get_ami_variable(channel, 'CRM_DIRECTION') if not direction: await self.handle_log_if_needed("Variable 'CRM_DIRECTION' is empty.", "INFO") if not direction: direction = await self.get_ami_variable(channel, 'DIRECTION') if not direction: await self.handle_log_if_needed("Variable 'DIRECTION' is empty.", "INFO") if not direction: caller_id_num = data.get('CallerIDNum', '') dest_caller_id_num = data.get('DestCallerIDNum', '') if caller_id_num and caller_id_num != '' and dest_caller_id_num and dest_caller_id_num != '': if caller_id_num != dest_caller_id_num: if len(dest_caller_id_num) <= 5: direction = "INBOUND" await self.handle_log_if_needed(f"Variable 'DestCallerIDNum' is {dest_caller_id_num}, setting direction - INBOUND", "INFO") else: direction = "OUTBOUND" await self.handle_log_if_needed(f"Variable 'DestCallerIDNum' is {dest_caller_id_num}, setting direction - OUTBOUND", "INFO") else: await self.handle_log_if_needed(f"Variable 'CallerIDNum' equal 'DestCallerIDNum' = {dest_caller_id_num}.", "INFO") else: await self.handle_log_if_needed("Variables 'CallerIDNum' or 'DestCallerIDNum' is empty.", "INFO") return direction async def handle_debug_all_events(self, manager, message): if not self.debug: return await self.handle_log_if_needed(f"Received event: {to_json(str(message))}", "INFO") log_entry = { "date": datetime.now().strftime('%Y-%m-%d'), "time": datetime.now().strftime('%H:%M:%S'), "data": to_json(str(message)) } try: with open(self.debug_call_file_path, "a", encoding="utf-8") as log_file: log_file.write(json.dumps(log_entry, ensure_ascii=False) + "\n") except Exception as e: await self.handle_log_if_needed(f"Error writing to log file: {e}", "ERROR") async def handle_log_if_needed(self, message, log_level="INFO"): if not self.debug: return log_entry = { "date": datetime.now().strftime('%Y-%m-%d'), "time": datetime.now().strftime('%H:%M:%S'), "level": log_level, "message": message } try: with open(self.debug_log_file_path, "a", encoding="utf-8") as log_file: log_file.write(json.dumps(log_entry, ensure_ascii=False) + "\n") except Exception as e: print(f"Error writing to log file: {e}") async def check_eventmask(self): response = await self.manager.send_action({ 'Action': 'Command', 'Command': 'manager show connected' }) event_mask_on = 'EventMask: on' in response.get('Output', '') await self.handle_log_if_needed( f"AMI EventMask is {'ON' if event_mask_on else 'OFF'}", "INFO") return event_mask_on async def monitor_eventmask(self): while True: try: await asyncio.sleep(60) if not await self.check_eventmask(): await self.manager.send_action( {'Action': 'Events', 'EventMask': 'on'}) except Exception: await self.handle_log_if_needed( f"monitor_eventmask error:\n{traceback.format_exc()}", "ERROR") async def cleanup_filters(self): TARGET_EVENTS = {'Newchannel', 'DialBegin', 'DialEnd', 'Cdr'} resp = await self.manager.send_action({ 'Action': 'Command', 'Command': 'manager show connected' }) output = resp.get('Output', '') if isinstance(output, list): output = '\n'.join(output) for flt in re.findall(r'EventFilter:\s+(.+)', output): m = re.match(r'!Event:(\w+)', flt) if m and m.group(1) in TARGET_EVENTS: await self.manager.send_action({ 'Action': 'Filter', 'Operation': 'del', 'Filter': flt }) await self.handle_log_if_needed(f'Removed filter {flt}', 'INFO') async def filter_watcher(self): while True: try: await asyncio.sleep(300) await self.cleanup_filters() except Exception: await self.handle_log_if_needed( f'filter_watcher error:\n{traceback.format_exc()}', 'ERROR') async def start(self): loop = asyncio.get_event_loop() self.manager = Manager(loop=loop, host=self.ip, username=self.username, secret=self.secret) self.manager.register_event('Newchannel', self.handle_newchannel) self.manager.register_event('DialBegin', self.handle_dialbegin) self.manager.register_event('DialEnd', self.handle_dialend) self.manager.register_event('Cdr', self.handle_cdr) while True: try: await self.manager.connect() await self.manager.send_action({ 'Action': 'Events', 'EventMask': 'on' }) await self.check_eventmask() await self.cleanup_filters() break except ConnectionRefusedError as e: print(f"[AMI] Connection refused, retrying in 5 seconds: {e}") await asyncio.sleep(5) except Exception: print(f"[AMI] Unexpected error, retrying in 5 seconds:\n{traceback.format_exc()}") await asyncio.sleep(5) loop.create_task(self.monitor_eventmask()) loop.create_task(self.filter_watcher()) try: await asyncio.Event().wait() finally: self.manager.close() def main(): ami_listener = AMIListener(IP, USERNAME, SECRET, URL, DEBUG, DEBUG_CALL_FILE_PATH, DEBUG_LOG_FILE_PATH) loop = asyncio.get_event_loop() loop.run_until_complete(ami_listener.start()) if __name__ == '__main__': main()