Clone
1
Создание ВМ ASR (192.168.254.34)
yasusan edited this page 2026-06-02 13:30:48 +00:00
This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

Как использовать Сохраните скрипт как setup_asr.sh

Сделайте исполняемым: chmod +x setup_asr.sh

Запустите от root: ./setup_asr.sh

Проверка после установки bash

Проверка сервисов

systemctl status asr-command systemctl status nginx

Проверка cron

crontab -l

Ручной тест обработки

cd /root/asr_service source venv/bin/activate python transcribe_auto.py

#!/bin/bash set -e

--- НАСТРОЙКИ (ПРОВЕРЬТЕ ПЕРЕД ЗАПУСКОМ) ---

ASR_HOST="192.168.254.34" AUDIO_PORT="8080" CMD_PORT="8003" DB_HOST="192.168.254.35" DB_PORT="5432" DB_NAME="asr_db" DB_USER="asr_user" DB_PASSWORD="poluyan77" EMOTION_HOST="192.168.254.37" EMOTION_PORT="8001"

--- КОНЕЦ НАСТРОЕК ---

echo "=== 1. Установка базовых пакетов ===" apt update apt install -y python3 python3-pip python3-venv ffmpeg nginx

echo "=== 2. Создание проекта и виртуального окружения ===" mkdir -p /root/asr_service cd /root/asr_service python3 -m venv venv source venv/bin/activate

echo "=== 3. Установка Python-пакетов ===" pip install torch==2.5.1 torchaudio==2.5.1 --index-url https://download.pytorch.org/whl/cpu pip install kairos-asr[cpu] pydub numpy soundfile psycopg2-binary requests

echo "=== 4. Создание необходимых папок ===" mkdir -p /opt/asr/incoming_audio mkdir -p /opt/asr/output mkdir -p /opt/asr/processed mkdir -p /var/log/asr

echo "=== 5. Создание модуля извлечения имени (name_extractor.py) ===" cat > /root/asr_service/name_extractor.py << 'EOF' import re

def extract_name(text): patterns = [ r'Добрый день,\s*([А-ЯЁ][а-яё]+)\s*,\sслушаю', r'Здравствуйте,\s([А-ЯЁ][а-яё]+)\s*,\sслушаю', r'Добрый вечер,\s([А-ЯЁ][а-яё]+)\s*,\s*слушаю', ] for pattern in patterns: match = re.search(pattern, text) if match: return match.group(1) return None

def get_third_operator_segment(segments): operator_segments = [seg for seg in segments if seg[3] == 'оператор'] if len(operator_segments) >= 3: return operator_segments[2][2] return None EOF

echo "=== 6. Создание модуля парсинга имени файла (filename_parser.py) ===" cat > /root/asr_service/filename_parser.py << 'EOF' from datetime import datetime

def parse_filename(filename): name = filename.replace('.wav', '').replace('.flac', '') parts = name.split('-') if len(parts) >= 6 and parts[0] == 'in': incoming = parts[1] outgoing = parts[2] date_str = parts[3] time_str = parts[4] timestamp = parts[5] if len(parts) > 5 else None

    call_date = f"{date_str[:4]}-{date_str[4:6]}-{date_str[6:8]}" if len(date_str) == 8 else None
    call_time = f"{time_str[:2]}:{time_str[2:4]}:{time_str[4:6]}" if len(time_str) == 6 else None
    
    return {
        'incoming_number': incoming,
        'outgoing_number': outgoing,
        'call_date': call_date,
        'call_time': call_time,
        'call_timestamp': float(timestamp) if timestamp else None
    }
return None

EOF

echo "=== 7. Создание модуля тегирования (tagging_service.py) ===" cat > /root/asr_service/tagging_service.py << 'EOF' import psycopg2 from psycopg2.extras import RealDictCursor

DB_CONFIG = { "host": "192.168.254.35", "port": 5432, "database": "asr_db", "user": "asr_user", "password": "poluyan77" }

def get_current_tag_version(): conn = psycopg2.connect(**DB_CONFIG, cursor_factory=RealDictCursor) cursor = conn.cursor() cursor.execute("SELECT version FROM tag_version WHERE id = 1") row = cursor.fetchone() cursor.close() conn.close() return row['version'] if row else 1

def get_rules(): conn = psycopg2.connect(**DB_CONFIG, cursor_factory=RealDictCursor) cursor = conn.cursor() cursor.execute(""" SELECT tr.tag_name, tk.keyword FROM tag_rules tr JOIN tag_keywords tk ON tk.tag_id = tr.id WHERE tr.is_active = TRUE """) rows = cursor.fetchall() cursor.close() conn.close()

rules = {}
for row in rows:
    tag = row['tag_name']
    keyword = row['keyword']
    if tag not in rules:
        rules[tag] = []
    rules[tag].append(keyword)
return rules

def tag_text(text): rules = get_rules() found_tags = set() text_lower = text.lower()

for tag, keywords in rules.items():
    for keyword in keywords:
        if keyword in text_lower:
            found_tags.add(tag)
            break
return list(found_tags)

def tag_call(call_id): conn = psycopg2.connect(**DB_CONFIG) cursor = conn.cursor()

cursor.execute("SELECT full_text FROM transcriptions WHERE call_id = %s", (call_id,))
row = cursor.fetchone()

if not row or not row[0]:
    cursor.close()
    conn.close()
    return []

full_text = row[0]
tags = tag_text(full_text)
current_version = get_current_tag_version()

cursor.execute("DELETE FROM auto_tags WHERE call_id = %s", (call_id,))

for tag in tags:
    cursor.execute("INSERT INTO auto_tags (call_id, tag_name, confidence) VALUES (%s, %s, 0.8)", (call_id, tag))

cursor.execute("UPDATE calls SET tag_version = %s WHERE id = %s", (current_version, call_id))

conn.commit()
cursor.close()
conn.close()
return tags

EOF

echo "=== 8. Создание основного скрипта обработки (transcribe_auto.py) ===" cat > /root/asr_service/transcribe_auto.py << 'EOF' #!/usr/bin/env python3 import warnings warnings.filterwarnings("ignore") import logging logging.getLogger().setLevel(logging.ERROR)

import os import sys import json import tempfile import shutil import fcntl from datetime import datetime from pydub import AudioSegment from pydub.effects import normalize from kairos_asr import KairosASR import psycopg2 from psycopg2.extras import execute_values from name_extractor import extract_name, get_third_operator_segment from filename_parser import parse_filename from tagging_service import tag_call

========== НАСТРОЙКИ ==========

INPUT_DIR = "/opt/asr/incoming_audio" OUTPUT_DIR = "/opt/asr/output" PROCESSED_DIR = "/opt/asr/processed" LOG_FILE = "/var/log/asr/processor.log"

DB_CONFIG = { "host": "192.168.254.35", "port": 5432, "database": "asr_db", "user": "asr_user", "password": "poluyan77" }

================================

os.makedirs(INPUT_DIR, exist_ok=True) os.makedirs(OUTPUT_DIR, exist_ok=True) os.makedirs(PROCESSED_DIR, exist_ok=True)

Кэш имён операторов

_operator_names_cache = None _operator_names_cache_time = None

def log(msg): timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") with open(LOG_FILE, "a") as f: f.write(f"[{timestamp}] {msg}\n") print(msg)

def format_time(seconds): minutes = int(seconds // 60) secs = seconds % 60 return f"{minutes:02d}:{secs:06.3f}"

def load_operator_names(force=False): global _operator_names_cache, _operator_names_cache_time

if not force and _operator_names_cache_time is not None:
    age = (datetime.now() - _operator_names_cache_time).total_seconds()
    if age < 86400:
        return _operator_names_cache

try:
    conn = psycopg2.connect(**DB_CONFIG)
    cursor = conn.cursor()
    cursor.execute("SELECT name FROM operator_names ORDER BY name")
    rows = cursor.fetchall()
    cursor.close()
    conn.close()
    
    _operator_names_cache = set(row[0].lower() for row in rows)
    _operator_names_cache_time = datetime.now()
    log(f"📇 Загружено имён операторов: {len(_operator_names_cache)}")
    return _operator_names_cache
except Exception as e:
    log(f"❌ Ошибка загрузки имён: {e}")
    return set()

def extract_name_from_third_segment(segments): operator_segments = [seg for seg in segments if seg[3] == 'оператор'] if len(operator_segments) < 3: return None

third_phrase = operator_segments[2][2]
if not third_phrase:
    return None

names_set = load_operator_names()
words = third_phrase.lower().split()

for word in words:
    clean_word = word.strip('.,!?;:()"\'')
    if clean_word in names_set:
        return clean_word.capitalize()
return None

def analyze_emotion_text(text, max_retries=3, timeout=10): import requests import time

for attempt in range(max_retries):
    try:
        response = requests.post(f'http://192.168.254.37:8001/analyze', data={'text': text}, timeout=timeout)
        if response.status_code == 200:
            result = response.json()
            return result.get('emotion'), result.get('confidence')
        else:
            log(f"⚠️ Эмоции: HTTP {response.status_code}, попытка {attempt + 1}/{max_retries}")
    except requests.exceptions.Timeout:
        log(f"⚠️ Эмоции: таймаут, попытка {attempt + 1}/{max_retries}")
    except requests.exceptions.ConnectionError:
        log(f"⚠️ Эмоции: нет соединения, попытка {attempt + 1}/{max_retries}")
    except Exception as e:
        log(f"⚠️ Эмоции: ошибка {e}, попытка {attempt + 1}/{max_retries}")
    
    if attempt < max_retries - 1:
        time.sleep(2)

log("❌ Эмоции: не удалось получить ответ")
return None, None

def save_to_db(call_data, segments): conn = None try: conn = psycopg2.connect(**DB_CONFIG) cursor = conn.cursor()

    cursor.execute("""
        INSERT INTO calls 
        (incoming_number, outgoing_number, call_date, call_time, call_timestamp,
         audio_filename, audio_filepath, audio_duration_sec, audio_size_bytes, status)
        VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, 'completed')
        RETURNING id
    """, (
        call_data.get('incoming_number'),
        call_data.get('outgoing_number'),
        call_data.get('call_date'),
        call_data.get('call_time'),
        datetime.fromtimestamp(call_data['call_timestamp']) if call_data.get('call_timestamp') else None,
        call_data['audio_filename'],
        call_data['audio_filepath'],
        call_data['audio_duration_sec'],
        call_data['audio_size_bytes']
    ))
    call_id = cursor.fetchone()[0]
    
    full_text = ' '.join([seg[2] for seg in segments])
    cursor.execute("INSERT INTO transcriptions (call_id, full_text) VALUES (%s, %s) RETURNING id", (call_id, full_text))
    transcript_id = cursor.fetchone()[0]
    
    segments_data = [(transcript_id, idx, seg[3], seg[0], seg[1], seg[2]) for idx, seg in enumerate(segments)]
    execute_values(cursor, """
        INSERT INTO transcription_segments 
        (transcript_id, segment_index, speaker, start_time_sec, end_time_sec, text)
        VALUES %s
    """, segments_data)
    
    operator_name = extract_name_from_third_segment(segments)
    if operator_name:
        cursor.execute("UPDATE calls SET operator_name = %s WHERE id = %s", (operator_name, call_id))
        log(f"👤 Извлечено имя оператора: {operator_name}")

    log("🔍 Запускаем тегирование...")
    tags = tag_call(call_id)
    if tags:
        log(f"🏷️ Найдены теги: {', '.join(tags)}")
    
    emotion, confidence = analyze_emotion_text(full_text)
    if emotion:
        cursor.execute("INSERT INTO auto_tags (call_id, tag_name, confidence) VALUES (%s, %s, %s)", (call_id, f"эмоция_{emotion}", confidence))
        log(f"🎭 Эмоция звонка: {emotion} ({confidence:.2f})")
    
    conn.commit()
    log(f"💾 Сохранено в БД: call_id={call_id}, segments={len(segments)}")
    return call_id
    
except Exception as e:
    log(f"❌ Ошибка БД: {e}")
    if conn:
        conn.rollback()
    return None
finally:
    if conn:
        conn.close()

def process_file(filepath): log(f"🔄 Обработка: {filepath}")

filename = os.path.basename(filepath)
file_size = os.path.getsize(filepath)

metadata = parse_filename(filename)

audio = AudioSegment.from_wav(filepath)
duration = len(audio) / 1000.0
log(f"📊 Длительность: {duration:.1f} сек, размер: {file_size/1024/1024:.2f} МБ")

audio_asr = audio.set_frame_rate(16000).set_sample_width(2)
audio_asr = normalize(audio_asr, headroom=0.5)

if audio_asr.channels < 2:
    left_audio = audio_asr
    right_audio = None
else:
    mono = audio_asr.split_to_mono()
    left_audio = mono[0]
    right_audio = mono[1]

asr = KairosASR(device="cpu")
all_segments = []

if left_audio:
    with tempfile.NamedTemporaryFile(suffix='.wav', delete=False) as tmp:
        left_audio.export(tmp.name, format='wav')
        result = asr.transcribe(wav_file=tmp.name)
        for sent in result.sentences:
            all_segments.append((sent.start, sent.end, sent.text, "абонент"))
        os.unlink(tmp.name)

if right_audio:
    with tempfile.NamedTemporaryFile(suffix='.wav', delete=False) as tmp:
        right_audio.export(tmp.name, format='wav')
        result = asr.transcribe(wav_file=tmp.name)
        for sent in result.sentences:
            all_segments.append((sent.start, sent.end, sent.text, "оператор"))
        os.unlink(tmp.name)

all_segments.sort(key=lambda x: x[0])

output_file = os.path.join(OUTPUT_DIR, filename.replace('.wav', '.json'))
with open(output_file, 'w', encoding='utf-8') as f:
    json.dump({"file": filepath, "segments": [[s, e, t, sp] for s, e, t, sp in all_segments]}, f, ensure_ascii=False, indent=2)

text_file = os.path.join(OUTPUT_DIR, filename.replace('.wav', '.txt'))
with open(text_file, 'w', encoding='utf-8') as f:
    for start, end, text, speaker in all_segments:
        f.write(f"[{format_time(start)} → {format_time(end)}] {speaker}: {text}\n")

call_data = {
    'audio_filename': filename,
    'audio_filepath': filepath,
    'audio_duration_sec': duration,
    'audio_size_bytes': file_size,
}
if metadata:
    call_data.update(metadata)

call_id = save_to_db(call_data, all_segments)

flac_filename = filename.replace('.wav', '.flac')
flac_path = os.path.join(PROCESSED_DIR, flac_filename)
audio.export(flac_path, format='flac')
os.remove(filepath)

flac_size = os.path.getsize(flac_path) / 1024 / 1024
log(f"✅ Готово: {flac_filename} ({flac_size:.2f} МБ), call_id={call_id}")
return True

if name == "main": lock_fd = open('/tmp/asr_processor.lock', 'w') try: fcntl.flock(lock_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) except BlockingIOError: print("Скрипт уже запущен, выхожу") sys.exit(0)

log("="*60)
log("ASR Processor started")

files = [f for f in os.listdir(INPUT_DIR) if f.endswith('.wav')]

if not files:
    log("Нет новых файлов")
else:
    log(f"Найдено файлов: {len(files)}")
    for filename in files:
        filepath = os.path.join(INPUT_DIR, filename)
        try:
            process_file(filepath)
        except Exception as e:
            log(f"❌ Ошибка при обработке {filename}: {e}")

log("ASR Processor finished")
fcntl.flock(lock_fd, fcntl.LOCK_UN)
lock_fd.close()

EOF

chmod +x /root/asr_service/transcribe_auto.py

echo "=== 9. Создание командного сервера (command_server.py) ===" cat > /root/asr_service/command_server.py << 'EOF' #!/usr/bin/env python3 import sys import os import json import http.server

sys.path.insert(0, '/root/asr_service')

def load_operator_names(): import psycopg2 DB_CONFIG = {"host": "192.168.254.35", "port": 5432, "database": "asr_db", "user": "asr_user", "password": "poluyan77"} try: conn = psycopg2.connect(**DB_CONFIG) cursor = conn.cursor() cursor.execute("SELECT name FROM operator_names ORDER BY name") rows = cursor.fetchall() cursor.close() conn.close() return [row[0] for row in rows] except Exception as e: print(f"Ошибка: {e}") return []

def reload_operator_names(): global _operator_names_cache _operator_names_cache = load_operator_names() return _operator_names_cache

_operator_names_cache = load_operator_names()

class CommandHandler(http.server.SimpleHTTPRequestHandler): def do_POST(self): if self.path == '/reload-names': try: names = reload_operator_names() self.send_response(200) self.send_header('Content-type', 'application/json') self.end_headers() self.wfile.write(json.dumps({"status": "ok", "count": len(names)}).encode()) print(f"🔄 Кэш имён операторов обновлён, загружено {len(names)} имён") except Exception as e: self.send_response(500) self.end_headers() self.wfile.write(json.dumps({"status": "error", "error": str(e)}).encode()) else: self.send_response(404) self.end_headers()

def log_message(self, format, *args):
    pass

def start_command_server(port=8003): server = http.server.HTTPServer(('0.0.0.0', port), CommandHandler) print(f"🚀 Командный сервер запущен на порту {port}") server.serve_forever()

if name == "main": start_command_server() EOF

chmod +x /root/asr_service/command_server.py

echo "=== 10. Настройка Nginx для аудио ===" cat > /etc/nginx/sites-available/audio-server << EOF server { listen $AUDIO_PORT; server_name $ASR_HOST;

location /audio/ {
    alias /opt/asr/processed/;
    add_header Access-Control-Allow-Origin *;
    add_header Access-Control-Allow-Methods "GET, HEAD, OPTIONS";
    add_header Access-Control-Allow-Headers "Range";
    add_header Content-Type audio/flac;
}

} EOF

ln -sf /etc/nginx/sites-available/audio-server /etc/nginx/sites-enabled/ rm -f /etc/nginx/sites-enabled/default nginx -t systemctl restart nginx

echo "=== 11. Создание systemd сервиса для командного сервера ===" cat > /etc/systemd/system/asr-command.service << EOF [Unit] Description=ASR Command Server for operator names cache After=network.target

[Service] Type=simple User=root WorkingDirectory=/root/asr_service Environment="PATH=/root/asr_service/venv/bin" ExecStart=/root/asr_service/venv/bin/python /root/asr_service/command_server.py Restart=always RestartSec=10

[Install] WantedBy=multi-user.target EOF

echo "=== 12. Настройка cron для запуска воркера ===" cat > /etc/cron.d/asr-processor << EOF */5 * * * * root cd /root/asr_service && /root/asr_service/venv/bin/python /root/asr_service/transcribe_auto.py >> /var/log/asr/cron.log 2>&1 EOF

echo "=== 13. Запуск сервисов ===" systemctl daemon-reload systemctl enable asr-command systemctl start asr-command systemctl restart nginx

echo "=== ASR ВМ готова ===" echo "Аудио порт: http://$ASR_HOST:$AUDIO_PORT" echo "Командный порт: $CMD_PORT" echo "Cron настроен на запуск каждые 5 минут"