feat: metrics agent for FreeSWITCH and Asterisk nodes

- pulse-lets-go-agent: collects metrics from PBX nodes via ESL (FreeSWITCH) or AMI (Asterisk)
- AMI client (raw TCP) for Asterisk Manager Interface
- Collector interface: collect, connect, subscribe hangup events, close
- FreeSWITCH collector: show channels, max_sessions, idle_cpu
- Asterisk collector: CoreShowChannels, CoreSettings
- System collector: /proc/loadavg, /proc/stat (CPU idle)
- Failure tracker: sliding window for call_failure_rate
- NATS publisher: pulse.metrics.<node_id> every 5 seconds
- Graceful shutdown, retry logic (5 failures → status=down)
- Template configs for UC nodes (Asterisk) and ses-sip (FreeSWITCH)
This commit is contained in:
Maksim Totmin
2026-06-25 16:59:10 +07:00
parent 3c698a8bdb
commit c713a4f142
12 changed files with 1170 additions and 0 deletions
+293
View File
@@ -0,0 +1,293 @@
// Package ami реализует клиент Asterisk Manager Interface (AMI).
// Использует сырой TCP — протокол AMI текстовый (Key: Value\r\n),
// аналогичен ESL по простоте, не нужна внешняя библиотека.
//
// Клиент обеспечивает:
// - Подключение и аутентификацию (Action: Login)
// - Синхронные команды (Action: Command)
// - Подписку на события (Action: Events)
// - Автоматический реконнект с exponential backoff
// - Потокобезопасность (mutex на отправку команд)
package ami
import (
"bufio"
"fmt"
"io"
"log"
"net"
"strings"
"sync"
"sync/atomic"
"time"
)
// Event — событие, полученное от Asterisk AMI.
type Event struct {
Headers map[string]string // все заголовки события (Event:, Channel:, Cause:, ...)
}
// Client — AMI-клиент для взаимодействия с Asterisk.
type Client struct {
host string
port int
username string
password string
conn net.Conn
reader *bufio.Reader
sendMu sync.Mutex // защищает запись в сокет
connected atomic.Bool
eventsCh chan *Event
closeCh chan struct{}
shouldReconn atomic.Bool
backoff time.Duration
}
// NewClient создаёт нового AMI-клиента.
func NewClient(host string, port int, username, password string) *Client {
return &Client{
host: host,
port: port,
username: username,
password: password,
eventsCh: make(chan *Event, 100),
closeCh: make(chan struct{}),
backoff: 1 * time.Second,
}
}
// Connect подключается к Asterisk AMI и выполняет аутентификацию.
func (c *Client) Connect() error {
c.shouldReconn.Store(true)
addr := fmt.Sprintf("%s:%d", c.host, c.port)
dialer := net.Dialer{Timeout: 5 * time.Second}
conn, err := dialer.Dial("tcp", addr)
if err != nil {
return fmt.Errorf("ami dial %s: %w", addr, err)
}
c.conn = conn
c.reader = bufio.NewReader(conn)
// 1. Читаем приветственное сообщение Asterisk (одна строка, без пустой строки)
line, err := c.reader.ReadString('\n')
if err != nil {
conn.Close()
return fmt.Errorf("чтение приветствия AMI: %w", err)
}
_ = line // "Asterisk Call Manager/5.0.5\r\n"
// 2. Отправляем логин
loginCmd := fmt.Sprintf("Action: Login\r\nUsername: %s\r\nSecret: %s\r\n\r\n", c.username, c.password)
if _, err := conn.Write([]byte(loginCmd)); err != nil {
conn.Close()
return fmt.Errorf("отправка Login: %w", err)
}
// 3. Читаем ответ
headers, err := c.readResponse()
if err != nil {
conn.Close()
return fmt.Errorf("чтение ответа Login: %w", err)
}
if resp := headers["Response"]; resp != "Success" {
conn.Close()
return fmt.Errorf("ami auth rejected: %s", headers["Message"])
}
c.connected.Store(true)
c.backoff = 1 * time.Second
log.Printf("[ami] подключён к %s", addr)
return nil
}
// Disconnect отключает клиент.
func (c *Client) Disconnect() {
c.shouldReconn.Store(false)
c.connected.Store(false)
if c.conn != nil {
c.conn.Close()
c.conn = nil
}
select {
case <-c.closeCh:
default:
close(c.closeCh)
}
}
// IsConnected возвращает статус подключения.
func (c *Client) IsConnected() bool {
return c.connected.Load()
}
// SendCommand отправляет синхронную команду и возвращает заголовки ответа.
func (c *Client) SendCommand(cmd string) (map[string]string, error) {
c.sendMu.Lock()
defer c.sendMu.Unlock()
if !c.connected.Load() || c.conn == nil {
return nil, fmt.Errorf("ami не подключён")
}
fullCmd := cmd + "\r\n\r\n"
if _, err := c.conn.Write([]byte(fullCmd)); err != nil {
return nil, fmt.Errorf("ami send: %w", err)
}
return c.readResponse()
}
// SendCommandBody отправляет команду и возвращает заголовки + тело ответа.
// Нужно для команд, возвращающих многострочный вывод (напр. Command).
func (c *Client) SendCommandBody(cmd string) (map[string]string, string, error) {
c.sendMu.Lock()
defer c.sendMu.Unlock()
if !c.connected.Load() || c.conn == nil {
return nil, "", fmt.Errorf("ami не подключён")
}
fullCmd := cmd + "\r\n\r\n"
if _, err := c.conn.Write([]byte(fullCmd)); err != nil {
return nil, "", fmt.Errorf("ami send: %w", err)
}
// Читаем заголовки до пустой строки
headers, err := c.readResponse()
if err != nil {
return nil, "", fmt.Errorf("ami read headers: %w", err)
}
// Если это ответ на Command — читаем тело до --END COMMAND--
if headers["Response"] == "Follows" {
var bodyLines []string
for {
line, err := c.reader.ReadString('\n')
if err != nil {
break
}
line = strings.TrimRight(line, "\r\n")
if line == "--END COMMAND--" {
break
}
bodyLines = append(bodyLines, line)
}
return headers, strings.Join(bodyLines, "\n"), nil
}
return headers, "", nil
}
// Subscribe подписывается на события AMI.
func (c *Client) Subscribe(eventMask string) error {
cmd := fmt.Sprintf("Action: Events\r\nEventMask: %s\r\n", eventMask)
_, _, err := c.SendCommandBody(cmd)
if err != nil {
return fmt.Errorf("подписка на события: %w", err)
}
log.Printf("[ami] подписан на события: %s", eventMask)
return nil
}
// Events возвращает канал с входящими событиями.
func (c *Client) Events() <-chan *Event {
return c.eventsCh
}
// StartReadingEvents запускает горутину чтения асинхронных событий.
func (c *Client) StartReadingEvents() {
go c.readEventsLoop()
}
// ConnectWithRetry подключается в фоне с бесконечными попытками.
func (c *Client) ConnectWithRetry() {
c.shouldReconn.Store(true)
for {
if err := c.Connect(); err == nil {
return
}
if !c.shouldReconn.Load() {
return
}
log.Printf("[ami] ошибка подключения, повтор через %v", c.backoff)
time.Sleep(c.backoff)
c.backoff *= 2
if c.backoff > 60*time.Second {
c.backoff = 60 * time.Second
}
}
}
// --- Приватные методы ---
// readResponse читает один AMI-ответ: заголовки Key: Value до пустой строки.
func (c *Client) readResponse() (map[string]string, error) {
headers := make(map[string]string)
for {
line, err := c.reader.ReadString('\n')
if err != nil {
if err == io.EOF {
return headers, nil
}
return headers, err
}
line = strings.TrimRight(line, "\r\n")
if line == "" {
return headers, nil
}
parts := strings.SplitN(line, ": ", 2)
if len(parts) == 2 {
headers[parts[0]] = parts[1]
}
}
}
// readEventsLoop читает асинхронные события от Asterisk.
func (c *Client) readEventsLoop() {
for {
select {
case <-c.closeCh:
return
default:
}
if !c.connected.Load() || c.conn == nil {
time.Sleep(500 * time.Millisecond)
continue
}
headers, err := c.readResponse()
if err != nil {
if err == io.EOF || strings.Contains(err.Error(), "use of closed network connection") {
c.connected.Store(false)
log.Printf("[ami] соединение разорвано")
if c.shouldReconn.Load() {
go c.ConnectWithRetry()
}
return
}
time.Sleep(100 * time.Millisecond)
continue
}
eventName := headers["Event"]
if eventName == "" {
continue
}
evt := &Event{Headers: headers}
select {
case c.eventsCh <- evt:
default:
log.Printf("[ami] канал событий переполнен, пропускаем %s", eventName)
}
}
}