Compare commits

...
9 Commits
Author SHA1 Message Date
Maksim Totmin ff037f9c1e feat: balancer dashboard metrics + fix WebSocket Hijacker 2026-06-25 19:59:26 +07:00
Maksim Totmin 273e316af9 refactor: deploy model — all files under /opt/pulse-lets-go
Single prefix deployment: /opt/pulse-lets-go/ contains
bin/, data/, web/, log/ — no more split across
/usr/local/bin, /etc, /usr/local/share, /var/log.

* Makefile: install → deploy target with PREFIX=/opt/pulse-lets-go
* deploy/pulse-lets-go.service: all paths under /opt/pulse-lets-go,
  removed stale nats.service dependency
* main.go: default dataDir = parent of exeDir (PREFIX/bin → PREFIX/data)
* README: updated all paths, step 1 uses make deploy
2026-06-25 19:35:54 +07:00
Maksim Totmin 0d8a56e73f fix: build system — findWebDir order, auto-copy web static to bin
* findWebDir: CWD and project root checked before binary-adjacent path,
  so dev always serves fresh web/build/ over stale bin/web/build/
* build-go: auto-copies web/build/ to bin/web/ before Go compile
* build-prod: same auto-copy for production builds
* Proper refresh of frontend static on every make build
2026-06-25 19:30:49 +07:00
Maksim Totmin 78e53f42bf fix: SvelteKit UI reactivity with Svelte 5 runes
* / replaces legacy $: syntax
* Active menu highlight via .url.pathname (/stores)
* Logout navigates to /login via goto()
* Token/user footer updates reactively via
2026-06-25 19:30:42 +07:00
Maksim Totmin cc9da3ad7d refactor: common util package, ESL/AMI/security fixes, Prometheus metrics
Backend stability and security improvements:

* internal/util/ — common RandomHex helper, removed 3 duplicates
* ESL: deduplicated readMessage (locked/unlocked), net.JoinHostPort for IPv6
* AMI: synchronous reconnect() in readEventsLoop, net.JoinHostPort for IPv6
* Auth: /api/auth/refresh accepts Authorization header only (no ?token=)
* decodeJSON: http.MaxBytesReader(1<<20) body limit
* Trunks: gatewayParams() uses configured ESL.GatewayPrefix
* Config: jwt_secret_env env-var fallback
* FSCollector: time.After → time.NewTimer with defer Stop
* Monitoring: Prometheus counters (route_requests, nodes_total/healthy, uptime)
* go fmt pass across all internal/ packages
2026-06-25 19:30:36 +07:00
Maksim Totmin e66ac27dd9 docs: fix 12 documentation inaccuracies in README.md and AGENTS.md
README.md fixes:
- Architecture diagram now shows pulse-lets-go-agent layer
- Features list: added agent, auto-trunk, WS gateway, user-01 protection
- route.lua example: replaced fake curl.fetch() with real api:execute()
- systemd unit: Wants=nats.service removed (NATS on separate host)
- Gateway mapping: removed grep 'pulse-' filter (misses existing gateways)
- Production deployment: added agent deployment step (between NATS and route.lua)
- Makefile commands: added build-agent, build-siptest, siptest-*, emulate-*

AGENTS.md fixes:
- Structure: added missing cmd/emulator, cmd/siptest, cmd/pulse-lets-go-agent, esl/, ami/, contrib/
- NATS model: added sip_gateway field
- config.json: added nats_user, nats_password, tls, rate_limit, log, esl blocks
- Internal packages: added esl/ and ami/
- route.lua name: route_call → pulse_route (matches production)
2026-06-25 17:16:03 +07:00
Maksim Totmin 13c88c1586 fix: agent systemd unit remove nats dependency
Agent connects to NATS remotely, not locally. After=network only.
2026-06-25 16:59:39 +07:00
Maksim Totmin 95ac4164dd chore: emulator sip_gateway, Makefile targets, full README documentation
- Emulator: add sip_gateway field to all scenarios (auto-trunk creation)
- Makefile: build-agent, build-siptest, emulate-*, siptest-* targets
- README: full documentation with agent section, production deployment
- Project structure update with all components
2026-06-25 16:59:17 +07:00
Maksim Totmin 1e6f432230 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)
2026-06-25 16:59:10 +07:00
42 changed files with 2650 additions and 283 deletions
+35 -3
View File
@@ -28,7 +28,8 @@
"max_calls": 250,
"idle_cpu": 35,
"load_avg": 0.85,
"call_failure_rate": 0.5
"call_failure_rate": 0.5,
"sip_gateway": "sip:pbx-01.lan:5060"
}
```
@@ -174,7 +175,7 @@ Response (нет нод вообще):
FS использует Lua скрипт в dialplan → `mod_curl` дёргает `/api/route` → получает JSON → парсит `sip_gateway` → bridge.
```
<extension name="route_call">
<extension name="pulse_route">
<condition field="destination_number" expression="^.*$">
<action application="lua" data="route.lua"/>
</condition>
@@ -188,6 +189,8 @@ route.lua: GET `/api/route?...` → если `fallback: true` идёт на fall
```json
{
"nats_url": "nats://localhost:4222",
"nats_user": "",
"nats_password": "",
"listen_addr": ":8080",
"jwt_secret": "auto-generated-on-first-run",
"monitoring_api_key": "auto-generated-on-first-run",
@@ -201,6 +204,28 @@ route.lua: GET `/api/route?...` → если `fallback: true` идёт на fall
"idle_score": 0.20,
"fail_score": 0.10
}
},
"tls": {
"enabled": false,
"cert_file": "",
"key_file": ""
},
"rate_limit": {
"enabled": true,
"route_per_sec": 5000,
"api_per_sec": 100
},
"log": {
"level": "info",
"format": "text"
},
"esl": {
"host": "",
"port": 8021,
"password": "",
"password_env": "",
"sofia_profile": "external",
"gateway_prefix": "pulse"
}
}
```
@@ -210,13 +235,19 @@ route.lua: GET `/api/route?...` → если `fallback: true` идёт на fall
```
pulse-lets-go/
├── cmd/pulse-lets-go/main.go
├── cmd/pulse-lets-go-agent/
├── cmd/emulator/
├── cmd/siptest/
├── internal/
│ ├── api/ — handlers: auth, nodes, trunks, users, monitoring, ws
│ ├── engine/ — scorer + router (sync.RWMutex, best_node кэш)
│ ├── nats/ — subscriber, in-memory node store
│ ├── esl/ — FreeSWITCH ESL client (raw TCP)
│ ├── ami/ — Asterisk AMI client (raw TCP)
│ ├── config/ — manager: чтение/атомарная запись JSON
│ ├── models/ — все общие типы
│ └── log/ — ASCII-metrics логгер + ротация
├── contrib/ — route.lua, agent конфиги
├── data/ — JSON конфиги (gitignored)
├── web/ — SvelteKit + shadcn-svelte + Tailwind
│ ├── src/routes/
@@ -224,7 +255,8 @@ pulse-lets-go/
│ │ ├── +page.svelte # дашборд
│ │ ├── login/+page.svelte
│ │ ├── trunks/+page.svelte
│ │ ── nodes/+page.svelte
│ │ ── nodes/+page.svelte
│ │ └── users/+page.svelte
│ ├── src/lib/
│ │ ├── components/
│ │ └── api.ts
+58 -16
View File
@@ -4,11 +4,12 @@
APP := pulse-lets-go
BINDIR := bin
WEB_DIR := web
DATA_DIR := data
PREFIX := /opt/pulse-lets-go
.PHONY: all build build-go build-web clean run test install \
.PHONY: all build build-go build-web clean run test deploy \
emulate emulate-normal emulate-overload emulate-cpu-low emulate-fail-high \
emulate-node-down emulate-stale emulate-chaos build-emulator
emulate-node-down emulate-stale emulate-chaos build-emulator \
build-siptest siptest-uas siptest-uac siptest-full
all: build
@@ -21,13 +22,17 @@ build-web:
@echo "→ Сборка SvelteKit..."
cd $(WEB_DIR) && npm run build
# Сборка Go бинарника
# Сборка Go бинарника (с авто-копированием статики)
build-go:
@mkdir -p $(BINDIR)
@echo "→ Сборка Go backend..."
@if [ -d $(WEB_DIR)/build ]; then \
rm -rf $(BINDIR)/web && cp -r $(WEB_DIR)/build $(BINDIR)/web; \
echo " статика скопирована в $(BINDIR)/web/"; \
fi
go build -o $(BINDIR)/$(APP) ./cmd/$(APP)
# Быстрая сборка только Go (без фронта)
# Быстрая сборка только Go (без фронта, без статики)
build-go-only:
@mkdir -p $(BINDIR)
go build -o $(BINDIR)/$(APP) ./cmd/$(APP)
@@ -42,23 +47,26 @@ run:
build-prod:
@mkdir -p $(BINDIR)
cd $(WEB_DIR) && npm run build
rm -rf $(BINDIR)/web && cp -r $(WEB_DIR)/build $(BINDIR)/web
CGO_ENABLED=0 go build -ldflags="-s -w" -o $(BINDIR)/$(APP) ./cmd/$(APP)
# Установка: копирует бинарник и статику в /usr/local
install: build
@echo "→ Установка в /usr/local"
mkdir -p /usr/local/share/$(APP)/web
cp -r $(WEB_DIR)/build /usr/local/share/$(APP)/web/
cp $(BINDIR)/$(APP) /usr/local/bin/$(APP)
chmod +x /usr/local/bin/$(APP)
@echo "✓ Установлен в /usr/local/bin/$(APP)"
# Деплой в PREFIX: бинарник + статика + data + log
deploy: build-prod
@echo "→ Деплой в $(PREFIX)"
install -d $(PREFIX)/bin $(PREFIX)/data $(PREFIX)/web $(PREFIX)/log
install -m 755 $(BINDIR)/$(APP) $(PREFIX)/bin/
cp -r $(WEB_DIR)/build/. $(PREFIX)/web/
@echo "✓ Готово: $(PREFIX)"
@echo " $(PREFIX)/bin/$(APP)"
@echo " $(PREFIX)/web/ (SvelteKit статика)"
@echo " $(PREFIX)/data/ (config.json, trunks.json, users.json)"
@echo " $(PREFIX)/log/ (ASCII-лог метрик)"
# Установка systemd unit
install-systemd:
@echo "→ Установка systemd unit"
cp deploy/$(APP).service /etc/systemd/system/
install -m 644 deploy/$(APP).service /etc/systemd/system/
systemctl daemon-reload
@echo "✓ systemd unit установлен. Запуск: systemctl start $(APP)"
@echo "✓ systemd unit установлен. Запуск: sudo systemctl start $(APP)"
# ============================================================
# Эмулятор метрик
@@ -100,6 +108,40 @@ emulate-stale: nats build-emulator
emulate-chaos: nats build-emulator
./$(BINDIR)/emulator --scenario chaos
# ============================================================
# Агент сбора метрик (pulse-lets-go-agent)
# ============================================================
build-agent:
@mkdir -p $(BINDIR)
go build -o $(BINDIR)/pulse-lets-go-agent ./cmd/pulse-lets-go-agent
# ============================================================
# SIP-тестер (e2e проверка с FreeSWITCH)
# ============================================================
build-siptest:
@mkdir -p $(BINDIR)
go build -o $(BINDIR)/siptest ./cmd/siptest
# Запуск UAS (PBX-нода, отвечает 200 OK на INVITE)
siptest-uas: build-siptest
./$(BINDIR)/siptest -mode uas
# Отправка тестового звонка (10 CPS, 100 вызовов)
siptest-uac: build-siptest
./$(BINDIR)/siptest -mode uac -r 10 -l 100 -m 100
# Полный e2e тест: UAS + UAC
siptest-full: build-siptest
@echo "→ Запуск UAS (PBX-нода) на порту 5090..."
./$(BINDIR)/siptest -mode uas &
@sleep 1
@echo "→ Запуск UAC (оператор)..."
./$(BINDIR)/siptest -mode uac -r 10 -l 100 -m 200
@echo "→ Остановка UAS..."
pkill siptest 2>/dev/null || true
# ============================================================
# NATS и запуск
# ============================================================
+823
View File
@@ -0,0 +1,823 @@
# Pulse Lets Go
[![Go](https://img.shields.io/badge/Go-1.26-blue?logo=go)](https://go.dev)
[![License](https://img.shields.io/badge/license-MIT-green)](LICENSE)
[![NATS](https://img.shields.io/badge/NATS-latest-27ae60?logo=nats)](https://nats.io)
[![SvelteKit](https://img.shields.io/badge/SvelteKit-2-fb923c?logo=svelte)](https://kit.svelte.dev)
**pulse-lets-go** — балансировщик телефонной нагрузки на базе FreeSWITCH + NATS.
Агенты PBX отправляют метрики в NATS каждые 5 секунд. Балансировщик вычисляет weighted score готовности каждой ноды, кэширует лучшую и отдаёт её FreeSWITCH через `/api/route` для следующего звонка.
---
## Features
- **pulse-lets-go-agent** — агент сбора метрик для FreeSWITCH (ESL) и Asterisk (AMI). Работает на каждой PBX-ноде.
- **Scoring engine** — lethal checks (disabled, stale, overload, CPU, failure rate) + weighted score (call, load, idle, fail) с конфигурируемыми весами.
- **NATS шина** — асинхронный сбор метрик через `pulse.metrics.<node_id>`.
- **Auto-создание транков** — при получении первой метрики от новой ноды с `sip_gateway` создаётся balance-транк автоматически.
- **Gateway health WS** — статус FS-gateway (registered/unregistered) транслируется на фронт через WebSocket.
- **REST API** — JWT (access + refresh), роли admin/viewer, rate limiting, CORS, health endpoints.
- **WebSocket** — real-time трансляция метрик на SvelteKit дашборд.
- **FreeSWITCH ESL** — управление Sofia-шлюзами через raw TCP ESL.
- **Monitoring** — Prometheus и Zabbix эндпоинты с отдельным API-key.
- **SvelteKit UI** — дашборд, таблицы нод, транков, пользователей.
- **Защита пользователей** — первичный администратор (user-01) не может быть удалён или понижен до viewer.
- **Graceful shutdown** — корректное завершение с таймаутом 10 с.
- **systemd ready** — unit с изоляцией (NoNewPrivileges, ProtectSystem, PrivateTmp).
---
## Архитектура
```
┌─────────────────────┐ ESL/AMI ┌──────────────────────┐
│ PBX-01 │──────────▶│ pulse-lets-go-agent │
│ (FS or Asterisk) │ collect │ (on each PBX node) │
├─────────────────────┤ ├──────────────────────┤
│ PBX-02 │──────────▶│ pulse-lets-go-agent │
│ (FS or Asterisk) │ ├──────────────────────┤
├─────────────────────┤ ├──────────────────────┤
│ PBX-N │──────────▶│ pulse-lets-go-agent │
│ (FS or Asterisk) │ └─────────┬────────────┘
└─────────────────────┘ │ pulse.metrics.<id>
│ every 5s
┌──────────────┐
│ NATS Server │
│ nats://:4222 │
└──────┬───────┘
┌──────────────────────┐
│ pulse-lets-go │
│ │
│ ┌────────────────┐ │
│ │ NATS Sub │ │
│ │ (pulse.metrics)│ │
│ └────────┬───────┘ │
│ ▼ │
│ ┌────────────────┐ │
│ │ Engine │ │
│ │ (scoring) │──│──▶ /api/route
│ └────────┬───────┘ │ ↓
│ ▼ │ FreeSWITCH bridge
│ ┌────────────────┐ │ sip:pbx-03:5060
│ │ REST API │──│──▶ /api/nodes
│ │ (JWT auth) │ │ /api/trunks
│ └────────┬───────┘ │ /api/users
│ ▼ │
│ ┌────────────────┐ │
│ │ WebSocket │──│──▶ SvelteKit UI
│ └────────────────┘ │ (ws://:8080/ws/metrics)
│ ┌────────────────┐ │
│ │ ESL Client │──│──▶ FreeSWITCH ESL
│ │ (gateways) │ │ :8021
│ └────────────────┘ │
│ ┌────────────────┐ │
│ │ Monitoring │──│──▶ /api/monitoring/prometheus
│ │ (API-key) │ │ /api/monitoring/zabbix
│ └────────────────┘ │
└──────────────────────┘
```
---
## Quick Start
### 1. Требования
- Go 1.26+
- Node.js 22+
- NATS Server (автоматически качается `make nats`)
### 2. Сборка
```bash
git clone git@git.totmin.ru:en2zmax/pulse-lets-go.git
cd pulse-lets-go
make build
```
### 3. Запуск
```bash
# Терминал 1: NATS
make nats
# Терминал 2: балансировщик
./bin/pulse-lets-go
# Терминал 3: эмулятор метрик (3 здоровые ноды)
make emulate-normal
```
### 4. Проверка
```bash
# Route (без авторизации)
curl 'http://localhost:8080/api/route?caller_id=123&dest=456&ingress_trunk=trk-001'
# Login
curl -X POST http://localhost:8080/api/auth/login \
-H 'Content-Type: application/json' \
-d '{"username":"admin","password":"admin"}'
```
---
## Конфигурация
При первом запуске `config.json` создаётся автоматически в директории `data/` (или `$PULSE_DATA_DIR`). JWT-секрет и API-key генерируются случайно.
### config.json
```json
{
"nats_url": "nats://localhost:4222",
"nats_user": "",
"nats_password": "",
"listen_addr": ":8080",
"jwt_secret": "aabbccdd...32-байта-hex",
"monitoring_api_key": "11223344...32-байта-hex",
"stale_threshold_sec": 20,
"scoring": {
"idle_cpu_min": 5,
"call_failure_rate_lethal": 15.0,
"weights": {
"call_score": 0.40,
"load_score": 0.30,
"idle_score": 0.20,
"fail_score": 0.10
}
},
"tls": {
"enabled": false,
"cert_file": "",
"key_file": ""
},
"rate_limit": {
"enabled": true,
"route_per_sec": 5000,
"api_per_sec": 100
},
"log": {
"level": "info",
"format": "text"
},
"esl": {
"host": "",
"port": 8021,
"password": "",
"password_env": "",
"sofia_profile": "external",
"gateway_prefix": "pulse"
}
}
```
| Поле | Дефолт | Описание |
|------|--------|----------|
| `nats_url` | `nats://localhost:4222` | URL NATS сервера |
| `nats_user` | `""` | Имя пользователя NATS (если нужна аутентификация) |
| `nats_password` | `""` | Пароль NATS |
| `listen_addr` | `:8080` | Адрес HTTP-сервера |
| `jwt_secret` | auto | Секрет для подписи JWT (32 байта hex) |
| `monitoring_api_key` | auto | API-ключ для `/api/monitoring/*` |
| `stale_threshold_sec` | `20` | Через сколько секунд без метрик нода считается stale |
| `scoring.idle_cpu_min` | `5` | Минимальный idle CPU. Ниже — lethal |
| `scoring.call_failure_rate_lethal` | `15.0` | Максимальный failure rate. Выше — lethal |
| `scoring.weights.*` | 0.40/0.30/0.20/0.10 | Веса компонентов скора (сумма = 1.0) |
| `tls.enabled` | `false` | Включить HTTPS |
| `rate_limit.enabled` | `true` | Включить rate limiting |
| `rate_limit.route_per_sec` | `5000` | Лимит `/api/route` на IP |
| `rate_limit.api_per_sec` | `100` | Лимит остальных `/api/*` на IP |
| `log.level` | `info` | Уровень лога: `debug`, `info`, `warn`, `error` |
| `log.format` | `text` | Формат: `text` или `json` |
| `esl.host` | `""` | Хост FreeSWITCH ESL. Пусто — ESL отключён |
| `esl.port` | `8021` | Порт ESL |
| `esl.password` | `""` | Пароль ESL (ClueCon) |
| `esl.password_env` | `""` | Или имя переменной окружения с паролем |
---
## Scoring Engine
### Lethal-условия (score = -100)
Если хотя бы одно истинно — нода исключается из роутинга:
| Условие | Порог |
|---------|-------|
| `status != "ok"` | Любой статус кроме ok |
| Stale | > `stale_threshold_sec` (20 с) |
| `active_calls >= max_calls` | 100% заполнение |
| `idle_cpu < idle_cpu_min` | < 5% |
| `call_failure_rate > call_failure_rate_lethal` | > 15% |
| `disabled == true` | Ручное отключение админом |
### Weighted score (0..100)
```
call_score = clamp(100 - (active_calls / max_calls * 100), 0, 100)
load_score = clamp(100 - (load_avg * 50), 0, 100)
idle_score = clamp(idle_cpu, 0, 100)
fail_score = clamp(100 - call_failure_rate, 0, 100)
score = call_score × 0.40 +
load_score × 0.30 +
idle_score × 0.20 +
fail_score × 0.10
```
### Пример
Нода `pbx-01`: `active_calls=42`, `max_calls=250`, `load_avg=0.85`, `idle_cpu=35`, `fail_rate=0.5`
```
call_score = clamp(100 - (42/250)*100, 0, 100) = 83.2
load_score = clamp(100 - 0.85*50, 0, 100) = 57.5
idle_score = clamp(35, 0, 100) = 35.0
fail_score = clamp(100 - 0.5, 0, 100) = 99.5
score = 83.2×0.40 + 57.5×0.30 + 35.0×0.20 + 99.5×0.10 = 70.1
```
### Кэширование
`best_node_id` + `best_score` + `fallback_active` пересчитываются при каждой метрике.
`/api/route` читает их за RLock — O(1).
---
## API Reference
### Auth
| Метод | Путь | Роль | Описание |
|-------|------|------|----------|
| `POST` | `/api/auth/login` | | Вход: `{"username":"admin","password":"admin"}` → access_token + refresh_token |
| `POST` | `/api/auth/refresh` | – | Обновление access-токена через refresh_token |
Все остальные эндпоинты (кроме `/api/route`, `/api/health/*`, `/api/monitoring/*`) требуют `Authorization: Bearer <jwt>`.
### Core
| Метод | Путь | Роль | Описание |
|-------|------|------|----------|
| `GET` | `/api/route` | – | Выбор ноды для звонка. Query: `caller_id`, `dest`, `ingress_trunk` |
| `GET` | `/api/health` | – | Статистика: всего нод, здоровых, запросов, uptime |
| `GET` | `/api/health/live` | | Liveness probe (всегда 200) |
| `GET` | `/api/health/ready` | | Readiness probe (NATS + метрики) |
### Nodes
| Метод | Путь | Роль | Описание |
|-------|------|------|----------|
| `GET` | `/api/nodes` | admin/viewer | Список всех нод со score, состоянием |
| `GET` | `/api/nodes/{id}/metrics` | admin/viewer | История метрик (ring buffer, ~30 мин) |
| `PUT` | `/api/nodes/{id}/toggle` | admin | `{"disabled":true, "reason":"..."}` |
### Trunks
| Метод | Путь | Роль | Описание |
|-------|------|------|----------|
| `GET` | `/api/trunks` | admin | Все транки. Query: `type=ingress\|balance\|fallback` |
| `POST` | `/api/trunks` | admin | Создать транк |
| `PUT` | `/api/trunks/{id}` | admin | Обновить транк |
| `DELETE` | `/api/trunks/{id}` | admin | Удалить транк |
### Users
| Метод | Путь | Роль | Описание |
|-------|------|------|----------|
| `GET` | `/api/users` | admin | Список пользователей |
| `POST` | `/api/users` | admin | Создать пользователя |
| `PUT` | `/api/users/{id}` | admin | Обновить пользователя |
| `DELETE` | `/api/users/{id}` | admin | Удалить пользователя |
### Monitoring (API-key, не JWT)
| Метод | Путь | Описание |
|-------|------|----------|
| `GET` | `/api/monitoring/zabbix` | JSON для Zabbix LLD |
| `GET` | `/api/monitoring/prometheus` | `text/plain` метрики для Prometheus |
### WebSocket
| Метод | Путь | Описание |
|-------|------|----------|
| `WS` | `/ws/metrics` | Real-time метрики на фронт (JWT в query `?token=`) |
### Примеры запросов
```bash
# Route (no auth)
curl 'http://localhost:8080/api/route?caller_id=74951234567&dest=123&ingress_trunk=trk-001'
# → {"node_id":"pbx-03","score":88,"sip_gateway":"sip:pbx03.lan:5060"}
# Fallback (все ноды unhealthy)
# → {"fallback":true,"sip_gateway":"sip:operator.lan:5060","reason":"all_nodes_unhealthy","nodes":[...]}
# Login
TOKEN=$(curl -s -X POST http://localhost:8080/api/auth/login \
-H 'Content-Type: application/json' \
-d '{"username":"admin","password":"admin"}' | jq -r '.access_token')
# Nodes with JWT
curl -H "Authorization: Bearer $TOKEN" http://localhost:8080/api/nodes
# Toggle node (admin only)
curl -X PUT http://localhost:8080/api/nodes/pbx-01/toggle \
-H "Authorization: Bearer $TOKEN" \
-H 'Content-Type: application/json' \
-d '{"disabled":true,"reason":"maintenance"}'
# Prometheus metrics (API-key)
curl -H "X-API-Key: $(jq -r '.monitoring_api_key' data/config.json)" \
http://localhost:8080/api/monitoring/prometheus
```
---
## FreeSWITCH Интеграция
### Dialplan (route.lua)
FreeSWITCH использует Lua-скрипт в dialplan. Скрипт дёргает `/api/route` через `mod_curl` и бриджится на полученный gateway.
```lua
-- contrib/route.lua (сокращённо)
api = freeswitch.API()
caller_id = session:getVariable("caller_id_number") or ""
dest = session:getVariable("destination_number") or ""
url = "http://localhost:8080/api/route?caller_id=" .. caller_id .. "&dest=" .. dest
raw = api:execute("curl", url)
body = raw:match("\r?\n\r?\n(.+)") or raw
local ok, route = pcall(cjson.decode, body)
if not ok then session:hangup("NORMAL_TEMPORARY_FAILURE") return end
if route.fallback then
session:bridge(route.sip_gateway) -- оператор
elseif route.sip_gateway then
session:bridge(route.sip_gateway) -- лучшая нода
else
session:hangup("NORMAL_TEMPORARY_FAILURE")
end
```
Конфигурация в dialplan:
```xml
<extension name="route_call">
<condition field="destination_number" expression="^.*$">
<action application="set" data="balancer_url=http://balancer.lan:8080"/>
<action application="set" data="ingress_trunk=trk-001"/>
<action application="lua" data="route.lua"/>
</condition>
</extension>
```
### ESL (автоматическое управление шлюзами)
При включении `esl.host` в конфиге балансировщик:
- Подключается к FreeSWITCH ESL (`:8021`)
- Синхронизирует ingress-транки как Sofia-шлюзы
- Слушает события регистрации шлюзов
- Передаёт статус gateway на WebSocket фронта
---
## Модели данных
### NATS-сообщение (PBX → Balancer, каждые 5 с)
```json
{
"node_id": "pbx-01",
"ts": 1719000000,
"status": "ok",
"active_calls": 42,
"max_calls": 250,
"idle_cpu": 35,
"load_avg": 0.85,
"call_failure_rate": 0.5
}
```
### Транки
| Тип | Назначение |
|-----|-----------|
| `ingress` | Входящий транк (откуда звонки приходят) |
| `balance` | Транк назначения (привязан к Node) |
| `fallback` | Резервный транк (когда все balance score < 0), ровно 1 |
---
## Деплой
### systemd
```bash
make deploy # полный деплой в /opt/pulse-lets-go
make install-systemd # устанавливает systemd unit
```
```ini
# deploy/pulse-lets-go.service
[Unit]
Description=Pulse Lets Go — Telephone Load Balancer
After=network.target
[Service]
User=pulse
Group=pulse
ExecStart=/opt/pulse-lets-go/bin/pulse-lets-go
WorkingDirectory=/opt/pulse-lets-go
Restart=always
RestartSec=5
NoNewPrivileges=true
PrivateTmp=true
ProtectSystem=strict
ProtectHome=yes
ReadWritePaths=/opt/pulse-lets-go/data /opt/pulse-lets-go/log
[Install]
WantedBy=multi-user.target
```
### Структура директорий
```
/opt/pulse-lets-go/
├── bin/pulse-lets-go # Бинарник
├── data/
│ ├── config.json # Конфигурация
│ ├── trunks.json # Транки
│ └── users.json # Пользователи
├── web/ # SvelteKit статика
└── log/
└── metrics.YYYY-MM-DD.log # ASCII-лог метрик
```
### Security
- `NoNewPrivileges=true` — запрет повышения привилегий
- `ProtectSystem=strict` — read-only системные директории
- `ProtectHome=yes` — изоляция home
- `PrivateTmp=true` — изолированный /tmp
- JWT-секрет и API-key генерируются случайно при первом запуске
---
## Production Deployment
### Предусловия (перед деплоем pulse-lets-go)
| № | Проверка | Команда | Ожидание |
|:-:|----------|---------|----------|
| 1 | **ESL запущен** | `nc -z 127.0.0.1 8021 && echo OK` | `OK` |
| 2 | **ESL пароль** | `echo "auth ClueCon" \| nc 127.0.0.1 8021` | `+OK accepted` |
| 3 | **mod_curl загружен** | `fs_cli -x "show modules" \| grep mod_curl` | `<load module="mod_curl"/>` |
| 4 | **mod_lua загружен** | `fs_cli -x "show modules" \| grep mod_lua` | `<load module="mod_lua"/>` |
| 5 | **NATS установлен** | `which nats-server` | путь к бинарнику |
| 6 | **Права на ESL** | ACL в `event_socket.conf.xml` разрешает 127.0.0.1 | `localhost` или `127.0.0.1` в ACL |
### Загрузка mod_curl без рестарта FreeSWITCH
Если `mod_curl` не загружен — звонки работать не будут. Загружается **без остановки FS**:
```bash
fs_cli -x "load mod_curl"
# Добавить в modules.conf.xml для персистентности:
# <load module="mod_curl"/>
```
### Порядок деплоя (8 шагов)
**Шаг 1 — Деплой на сервер**
```bash
make deploy
scp -r /opt/pulse-lets-go devadmin@10.101.60.81:/opt/
scp contrib/route.lua devadmin@10.101.60.81:/tmp/
```
**Шаг 2 — NATS сервер** (если не установлен)
```bash
ssh devadmin@10.101.60.81
sudo mkdir -p /opt/nats
# Установить nats-server из репозитория или скопировать бинарник
# Запуск: nats-server -p 4222 -D &
```
**Шаг 3 — Деплой agent на PBX-ноды** (повторить для каждой UC/Asterisk ноды)
```bash
# Скопировать бинарник и конфиг на каждую PBX:
scp bin/pulse-lets-go-agent devadmin@PBX_IP:/usr/local/bin/
scp contrib/agent-uc.json devadmin@PBX_IP:/etc/pulse-lets-go-agent/agent.json
# Запустить:
ssh devadmin@PBX_IP "systemctl start pulse-lets-go-agent"
# Проверить лог:
# journalctl -u pulse-lets-go-agent -f
# → [ami] подключён к 127.0.0.1:6154
# → [agent] опубликована метрика: calls=0/150 load=0.28 cpu=99% status=ok
```
**Шаг 4 — route.lua в FS scripts**
```bash
sudo mkdir -p /etc/freeswitch/scripts
sudo cp /tmp/route.lua /etc/freeswitch/scripts/
```
**Шаг 5 — Dialplan** (добавить `pulse_route` extension)
```xml
<!-- В /etc/freeswitch/dialplan/default.xml перед остальными extension -->
<extension name="pulse_route">
<condition field="destination_number" expression="^(.*)$">
<action application="set" data="balancer_url=http://127.0.0.1:8080"/>
<action application="lua" data="route.lua"/>
</condition>
</extension>
```
```bash
sudo fs_cli -x "reloadxml"
```
**Шаг 6 — Создание balance-транков для каждой PBX-ноды**
```bash
# Получить список gateway из FS (или список нод из конфига):
echo "api sofia profile external gwlist" | nc 127.0.0.1 8021
# Создать balance-транк для каждой UC-ноды:
curl -X POST .../api/trunks
-d '{"name":"uc06","type":"balance","node_id":"uc06","gateway":"sip:10.101.60.115:5060","enabled":true}'
```
Транки также создаются автоматически при получении первой метрики от агента (если в метрике передан `sip_gateway`).
**Шаг 7 — Запуск pulse-lets-go**
```bash
# config.json создастся автоматически при первом запуске в /opt/pulse-lets-go/data/
# Добавить esl блок в config.json:
# "esl": {"host": "127.0.0.1", "port": 8021, "password": "ClueCon"}
/opt/pulse-lets-go/bin/pulse-lets-go
```
**Шаг 8 — Проверка**
```bash
curl http://127.0.0.1:8080/api/health | jq '.connections'
# → {"nats":"connected","esl":"connected"}
curl 'http://127.0.0.1:8080/api/route?caller_id=123&dest=456'
# → {"node_id":"pbx-03","score":88,"sip_gateway":...}
```
### Rollback (если route.lua сломал звонки)
```bash
# 1. Убрать pulse_route из dialplan
sudo sed -i '/pulse_route/,/<\/extension>/d' /etc/freeswitch/dialplan/default.xml
sudo fs_cli -x "reloadxml"
# 2. Остановить pulse-lets-go
sudo systemctl stop pulse-lets-go
# Звонки продолжают идти по старому dialplan без балансировщика.
```
### Права доступа (production)
| Операция | Требуемые права |
|----------|----------------|
| Чтение FS конфигов | root или freeswitch |
| `fs_cli -x` команды | root или freeswitch |
| ESL (порт 8021) | ACL в `event_socket.conf.xml` |
| Установка ПО | root |
| Запись в `/etc/freeswitch/scripts/` | root |
| `reloadxml` | root или freeswitch |
### Версии (протестировано)
| Компонент | Версия |
|-----------|--------|
| FreeSWITCH | 1.10.12+ |
| NATS Server | 2.10+ |
| Go | 1.26 |
| ОС | Linux (CentOS 7+, AlmaLinux, Arch) |
---
## Metrics Agent (pulse-lets-go-agent)
### Architecture
```
┌─────────────────────┐ pulse.metrics.<id> ┌─────────────────────┐
│ pulse-lets-go-agent│──────────────────────────▶ │ NATS Server │
│ (on each PBX node) │ every 5 seconds │ (central) │
├─────────────────────┤ └────────┬────────────┘
│ FreeSWITCH (ESL) │ │
│ Asterisk (AMI) │ ▼
│ System (/proc) │ ┌─────────────────┐
└─────────────────────┘ │ pulse-lets-go │
│ (scoring engine) │
└─────────────────┘
```
The agent runs on each PBX node (FreeSWITCH or Asterisk), collects metrics, and publishes them to NATS at a regular interval (default 5 seconds). The pulse-lets-go balancer subscribes to all metrics and uses them for scoring and routing decisions.
### Installation
```bash
make build-agent
# → bin/pulse-lets-go-agent (static binary, ~9 MB)
```
### Configuration (agent.json)
```json
{
"node_id": "uc06",
"type": "asterisk",
"nats_url": "nats://balancer.host:4222",
"nats_user": "",
"nats_password": "",
"interval_sec": 5,
"max_calls": 150,
"sip_gateway": "sip:10.101.60.115:5060",
"sip_gateway_auto": false,
"failure_window": 1000,
"esl": { "host": "127.0.0.1", "port": 8021, "password": "ClueCon" },
"ami": { "host": "127.0.0.1", "port": 6154, "username": "ctt", "password": "cttpass" }
}
```
| Поле | Тип | Default | Описание |
|------|-----|---------|----------|
| `node_id` | string | — | Уникальный ID ноды (uc06, ses-sip) |
| `type` | string | — | `freeswitch` или `asterisk` |
| `nats_url` | string | — | URL NATS-сервера |
| `nats_user` | string | `""` | Пользователь NATS |
| `nats_password` | string | `""` | Пароль NATS |
| `interval_sec` | int | `5` | Интервал публикации метрик |
| `max_calls` | int | `250` | Ёмкость ноды (fallback) |
| `sip_gateway` | string | `""` | SIP-адрес для auto-создания транка |
| `sip_gateway_auto` | bool | `false` | Авто-определение SIP из PBX |
| `failure_window` | int | `1000` | Размер окна для call_failure_rate |
| `esl.*` | object | — | FreeSWITCH ESL настройки (`host`, `port`, `password`) |
| `ami.*` | object | — | Asterisk AMI настройки (`host`, `port`, `username`, `password`) |
### Collectors
| Collector | Источник | Метрики | PBX |
|-----------|----------|---------|-----|
| `freeswitch` | ESL: `show channels count`, `json status`, `eval $${idle_cpu}` | active_calls, max_calls, idle_cpu | ✅ |
| `asterisk` | AMI: `CoreShowChannels`, `CoreSettings`, Hangup events | active_calls, max_calls, call_failure_rate | ✅ |
| `system` | `/proc/loadavg`, `/proc/stat` | load_avg, idle_cpu (fallback) | ✅ |
**Call Failure Rate:** скользящее окно последних N (1000) завершённых звонков.
| Исход | FS (ESL) | Asterisk (AMI) |
|-------|----------|---------------|
| Успех | `CHANNEL_HANGUP` + `Hangup-Cause: NORMAL_CLEARING` | `Hangup` + `Cause: 16` |
| Отказ | Любой другой `Hangup-Cause` | Любой другой `Cause` |
### Collector Interface (расширяемость)
```go
type Collector interface {
Type() string // "freeswitch" | "asterisk"
Connect() error // ESL auth / AMI login
Collect() (*PBXResult, error) // Сбор метрик
ListenHangup(onHangup func(bool)) // Подписка на hangup события
Close() // Закрытие соединения
}
```
Новый тип ноды (omnichannel, другой PBX) = новый файл `collector_<type>.go`, реализующий 5 методов интерфейса. Ничего в существующем коде менять не нужно.
### Template Configs
| Файл | Назначение | Тип |
|------|-----------|-----|
| `contrib/agent-uc.json` | Asterisk UC-ноды (05, 06, 66-69) | asterisk |
| `contrib/agent-sessip.json` | FreeSWITCH ses-sip (10.3.0.44) | freeswitch |
### Deployment
```bash
# 1. Собрать
make build-agent
# 2. Скопировать на ноду
scp bin/pulse-lets-go-agent devadmin@NODE_IP:/usr/local/bin/
# 3. Создать конфиг
scp contrib/agent-uc.json devadmin@NODE_IP:/etc/pulse-lets-go-agent/agent.json
# 4. Запустить
ssh devadmin@NODE_IP "pulse-lets-go-agent -config /etc/pulse-lets-go-agent/agent.json"
```
### Systemd Unit
```ini
# /etc/systemd/system/pulse-lets-go-agent.service
[Unit]
Description=Pulse Lets Go Agent — PBX metrics collector
After=network.target
Wants=network-online.target
[Service]
Type=simple
User=root
ExecStart=/usr/local/bin/pulse-lets-go-agent -config /etc/pulse-lets-go-agent/agent.json
Restart=always
RestartSec=5
StandardOutput=journal
StandardError=journal
[Install]
WantedBy=multi-user.target
```
### Verification
```bash
# На PBX-ноде — проверить лог агента:
journalctl -u pulse-lets-go-agent -f
# Ожидаем:
# [publisher] NATS подключён к nats://host:4222
# [ami] подключён к 127.0.0.1:6154 (для Asterisk)
# [agent] опубликована метрика: calls=0/150 load=0.28 cpu=99% status=ok
# На балансировщике — проверить что нода появилась:
curl http://localhost:8080/api/nodes | jq '.[] | {node_id, score, active_calls}'
```
### Troubleshooting
| Проблема | Причина | Решение |
|----------|---------|---------|
| `nats: no servers available` | NATS не запущен или недоступен | `systemctl start nats-server` |
| `ami auth rejected` | Неверный пароль или доступ | Проверить `/etc/asterisk/manager.conf` |
| `esl connect timeout` | ESL не слушает на порту | `fs_cli -x "load mod_event_socket"` |
| `calls=0/0` | max_calls не удалось получить | Указать `max_calls` в `agent.json` |
---
## Разработка
### Команды
| Команда | Описание |
|---------|----------|
| `make build` | Полная сборка (Go + SvelteKit) |
| `make dev` | Go backend + SvelteKit dev server (hot-reload) |
| `make test` | Go тесты с coverage |
| `make test-race` | Тесты с race detector |
| `make lint` | `go vet ./...` |
| `make fmt` | `go fmt ./...` |
| `make build-agent` | Сборка pulse-lets-go-agent |
| `make build-siptest` | Сборка SIP-тестера |
| `make emulate-*` | Эмуляция метрик (normal, overload, chaos, stale, ...) |
| `make siptest-*` | SIP-тестирование (uas, uac, full) |
### Стек
| Компонент | Технология |
|-----------|-----------|
| Backend | Go 1.26, stdlib `net/http` + `http.ServeMux` |
| Frontend | SvelteKit 2 + Tailwind CSS 4 + Lucide |
| Auth | JWT (golang-jwt/v5, HS256) + bcrypt |
| Message bus | NATS (nats.go) |
| Storage | JSON-файлы (data/) |
| ESL | Raw TCP (без внешних библиотек) |
### Структура проекта
```
pulse-lets-go/
├── cmd/
│ ├── pulse-lets-go/ # Точка входа (main.go)
│ ├── pulse-lets-go-agent/ # Агент сбора метрик (FS / Asterisk)
│ ├── emulator/ # Эмулятор метрик для тестирования
│ └── siptest/ # SIP-тестер для e2e проверок
├── internal/
│ ├── api/ # HTTP-хендлеры (auth, nodes, trunks, users, monitoring, ws, route)
│ ├── config/ # JSON-конфиги (atomic save, thread-safe)
│ ├── engine/ # Scoring engine + router (sync.RWMutex, best-node cache)
│ ├── esl/ # FreeSWITCH ESL client (raw TCP)
│ ├── ami/ # Asterisk AMI client (raw TCP)
│ ├── log/ # ASCII-логгер метрик с ротацией
│ ├── models/ # Все типы данных (NodeMetric, NodeState, Trunk, User, ...)
│ └── nats/ # NATS subscriber
├── contrib/ # route.lua, agent-*.json (FreeSWITCH dialplan + агент)
├── deploy/ # systemd unit
├── web/ # SvelteKit фронтенд
└── data/ # JSON-файлы времени выполнения (gitignored)
```
+8 -2
View File
@@ -37,6 +37,7 @@ type nodeMetric struct {
IdleCPU float64 `json:"idle_cpu"`
LoadAvg float64 `json:"load_avg"`
CallFailureRate float64 `json:"call_failure_rate"`
SIPGateway string `json:"sip_gateway,omitempty"`
}
// --- Спецификация значения (фиксированное или случайный диапазон) ---
@@ -70,6 +71,7 @@ type nodeSpec struct {
idleCPU valueSpec
loadAvg valueSpec
failRate valueSpec
sipGateway string // SIP-адрес для авто-создания транка
}
func (ns *nodeSpec) generate() nodeMetric {
@@ -82,6 +84,7 @@ func (ns *nodeSpec) generate() nodeMetric {
IdleCPU: ns.idleCPU.get(),
LoadAvg: ns.loadAvg.get(),
CallFailureRate: ns.failRate.get(),
SIPGateway: ns.sipGateway,
}
}
@@ -100,16 +103,19 @@ func scenarioNormal() []nodeSpec {
nodeID: "pbx-01", status: fixed(1),
calls: rrange(10, 100), maxCalls: 250,
idleCPU: rrange(40, 95), loadAvg: rrange(0.1, 1.0), failRate: rrange(0, 2),
sipGateway: "sip:pbx-01.lan:5060",
},
{
nodeID: "pbx-02", status: fixed(1),
calls: rrange(50, 200), maxCalls: 250,
idleCPU: rrange(15, 45), loadAvg: rrange(0.5, 2.0), failRate: rrange(1, 8),
sipGateway: "sip:pbx-02.lan:5060",
},
{
nodeID: "pbx-03", status: fixed(1),
calls: rrange(5, 60), maxCalls: 250,
idleCPU: rrange(60, 90), loadAvg: rrange(0.05, 0.5), failRate: rrange(0, 1),
sipGateway: "sip:pbx-03.lan:5060",
},
}
}
@@ -327,6 +333,6 @@ func publish(nc *natsgo.Conn, m nodeMetric) {
log.Printf("[emulator] ошибка публикации %s: %v", m.NodeID, err)
return
}
log.Printf("[emulator] → %s: calls=%d/%d idle=%.0f%% load=%.2f fail=%.1f%% status=%s",
m.NodeID, m.ActiveCalls, m.MaxCalls, m.IdleCPU, m.LoadAvg, m.CallFailureRate, m.Status)
log.Printf("[emulator] → %s: calls=%d/%d idle=%.0f%% load=%.2f fail=%.1f%% status=%s gw=%s",
m.NodeID, m.ActiveCalls, m.MaxCalls, m.IdleCPU, m.LoadAvg, m.CallFailureRate, m.Status, m.SIPGateway)
}
+174
View File
@@ -0,0 +1,174 @@
package agent
import (
"fmt"
"log"
"os"
"os/signal"
"sync/atomic"
"syscall"
"time"
"github.com/pulse-lets-go/internal/models"
)
const maxConsecutiveFailures = 5
// Agent — главный цикл сбора метрик и публикации в NATS.
type Agent struct {
cfg *Config
publisher *Publisher
collector Collector
failtrack *FailureTracker
pubCount atomic.Int64
}
// NewAgent создаёт агента на основе конфигурации.
func NewAgent(cfg *Config) (*Agent, error) {
if err := cfg.Validate(); err != nil {
return nil, err
}
// Publisher
pub, err := NewPublisher(cfg.NatsURL, cfg.NatsUser, cfg.NatsPassword, cfg.NodeID)
if err != nil {
return nil, err
}
// Collector — в зависимости от типа PBX
var col Collector
switch cfg.Type {
case "freeswitch":
col = NewFSCollector(cfg.ESL.Host, cfg.ESL.Port, cfg.ESL.Password)
case "asterisk":
col = NewAsteriskCollector(cfg.AMI.Host, cfg.AMI.Port, cfg.AMI.Username, cfg.AMI.Password)
default:
pub.Close()
return nil, fmt.Errorf("неподдерживаемый тип PBX: %s", cfg.Type)
}
if err := col.Connect(); err != nil {
pub.Close()
return nil, fmt.Errorf("подключение к %s: %w", cfg.Type, err)
}
// Failure rate tracker
ft := NewFailureTracker(cfg.FailureWindow)
// Подписка на hangup-события
col.ListenHangups(func(success bool) {
ft.Record(success)
})
return &Agent{
cfg: cfg,
publisher: pub,
collector: col,
failtrack: ft,
}, nil
}
// Start запускает главный цикл агента: сбор метрик → публикация.
// Блокирует до получения сигнала завершения (SIGINT, SIGTERM).
func (a *Agent) Start() {
log.Printf("[agent] запущен: node=%s, type=%s, interval=%ds",
a.cfg.NodeID, a.cfg.Type, a.cfg.IntervalSec)
ticker := time.NewTicker(time.Duration(a.cfg.IntervalSec) * time.Second)
defer ticker.Stop()
quit := make(chan os.Signal, 1)
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
// Счётчик последовательных ошибок сбора
consecutiveFailures := 0
// Первый сбор сразу
a.collectAndPublish(&consecutiveFailures)
for {
select {
case <-ticker.C:
a.collectAndPublish(&consecutiveFailures)
case <-quit:
log.Println("[agent] получен сигнал завершения")
a.Shutdown()
return
}
}
}
// collectAndPublish выполняет один цикл: сбор метрик → публикация.
func (a *Agent) collectAndPublish(failures *int) {
// 1. Системные метрики (CPU, load)
sys := CollectSystem()
// 2. Метрики PBX (calls, max_calls, статус)
pbx, err := a.collector.Collect()
if err != nil {
*failures++
log.Printf("[agent] ошибка сбора метрик (попытка %d): %v", *failures, err)
if *failures >= maxConsecutiveFailures {
// После N неудач — считаем ноду "down"
pbx = &PBXResult{
Status: "down",
}
} else {
// Пробуем ещё раз через секунду
time.Sleep(1 * time.Second)
pbx, err = a.collector.Collect()
if err != nil {
return // пропускаем этот тик
}
*failures = 0
}
} else {
*failures = 0
}
// 3. Failure rate
failRate := a.failtrack.Rate()
// 4. SIP gateway
sipGW := a.cfg.GetSIPGateway("")
// 5. Формируем метрику
metric := ToMetric(a.cfg.NodeID, pbx, sys, failRate, sipGW, a.cfg.MaxCalls)
// 6. Публикуем
if err := a.publisher.Publish(metric); err != nil {
log.Printf("[agent] ошибка публикации: %v", err)
return
}
a.pubCount.Add(1)
log.Printf("[agent] опубликована метрика: calls=%d/%d load=%.2f fail=%.1f%% cpu=%.0f%% status=%s",
metric.ActiveCalls, metric.MaxCalls, metric.LoadAvg, metric.CallFailureRate, metric.IdleCPU, metric.Status)
}
// Shutdown gracefully shuts down the agent.
func (a *Agent) Shutdown() {
log.Printf("[agent] остановка, всего опубликовано метрик: %d", a.pubCount.Load())
if a.collector != nil {
a.collector.Close()
}
if a.publisher != nil {
a.publisher.Close()
}
}
// CollectAndPublish экспортированный метод для ручного вызова (тесты, дебаг).
func (a *Agent) CollectAndPublish() *models.NodeMetric {
sys := CollectSystem()
pbx, err := a.collector.Collect()
if err != nil {
pbx = &PBXResult{Status: "down"}
}
failRate := a.failtrack.Rate()
sipGW := a.cfg.GetSIPGateway("")
metric := ToMetric(a.cfg.NodeID, pbx, sys, failRate, sipGW, a.cfg.MaxCalls)
a.publisher.Publish(metric)
a.pubCount.Add(1)
return metric
}
@@ -0,0 +1,70 @@
package agent
import "github.com/pulse-lets-go/internal/models"
// PBXResult — результат сбора метрик от PBX (FS или Asterisk).
type PBXResult struct {
ActiveCalls int // текущее количество звонков
MaxCalls int // ёмкость (0 — не удалось получить)
IdleCPU float64 // idle CPU из PBX (0 — не удалось, fallback на system)
Status string // "ok" — PBX отвечает нормально
SIPGateway string // авто-определённый SIP-адрес
}
// Collector — интерфейс сбора метрик с PBX.
type Collector interface {
// Type возвращает тип PBX: "freeswitch" или "asterisk".
Type() string
// Connect подключается к PBX (ESL auth / AMI login).
Connect() error
// Collect собирает метрики с PBX.
// Возвращает PBXResult и ошибку, если сбор не удался.
Collect() (*PBXResult, error)
// ListenHangups подписывается на события завершения звонков
// для расчёта call_failure_rate.
ListenHangups(onHangup func(success bool))
// Close закрывает соединение с PBX.
Close()
}
// HangupEvent — событие завершения звонка.
type HangupEvent struct {
Success bool // true — звонок успешен, false — ошибка
Cause string // код причины (NORMAL_CLEARING, BUSY, и т.д.)
}
// ToMetric формирует models.NodeMetric из результатов сбора.
func ToMetric(nodeID string, pbx *PBXResult, sys *SystemStats, failureRate float64, sipGateway string, maxCallsDefault int) *models.NodeMetric {
// Определяем max_calls: приоритет PBX → конфиг
maxCalls := pbx.MaxCalls
if maxCalls <= 0 {
maxCalls = maxCallsDefault
}
// Определяем idle_cpu: приоритет PBX → system
idleCPU := pbx.IdleCPU
if idleCPU <= 0 {
idleCPU = sys.IdleCPU
}
// Определяем статус
status := pbx.Status
if pbx.Status == "ok" && sys.IdleCPU < 5 {
status = "degraded"
}
return &models.NodeMetric{
NodeID: nodeID,
TS: nowTS(),
Status: status,
ActiveCalls: pbx.ActiveCalls,
MaxCalls: maxCalls,
IdleCPU: idleCPU,
LoadAvg: sys.Load1,
CallFailureRate: failureRate,
SIPGateway: sipGateway,
}
}
@@ -0,0 +1,107 @@
package agent
import (
"fmt"
"strconv"
"strings"
"github.com/pulse-lets-go/internal/ami"
)
// AsteriskCollector собирает метрики с Asterisk через AMI.
type AsteriskCollector struct {
client *ami.Client
}
// NewAsteriskCollector создаёт Asterisk-коллектор.
func NewAsteriskCollector(host string, port int, username, password string) *AsteriskCollector {
return &AsteriskCollector{
client: ami.NewClient(host, port, username, password),
}
}
func (ac *AsteriskCollector) Type() string { return "asterisk" }
func (ac *AsteriskCollector) Connect() error {
return ac.client.Connect()
}
func (ac *AsteriskCollector) Collect() (*PBXResult, error) {
if !ac.client.IsConnected() {
return nil, fmt.Errorf("ami не подключён")
}
result := &PBXResult{Status: "ok"}
// active_calls: CoreShowChannels определяет количество активных каналов
if active, err := ac.coreShowChannelsCount(); err == nil {
result.ActiveCalls = active
}
// max_calls: CoreSettings → CoreMaxCalls
if max, err := ac.coreSettingsInt("CoreMaxCalls"); err == nil {
result.MaxCalls = max
}
// idle_cpu: Asterisk не даёт FS-шного idle_cpu, используем системный
result.IdleCPU = 0
return result, nil
}
func (ac *AsteriskCollector) ListenHangups(onHangup func(success bool)) {
// Подписываемся на события Hangup через AMI
ac.client.Subscribe("call")
go func() {
for evt := range ac.client.Events() {
if evt.Headers["Event"] == "Hangup" {
cause := evt.Headers["Cause"]
// Считаем успешным только NORMAL_CLEARING (код 16)
success := cause == "16" || strings.Contains(cause, "NORMAL_CLEARING")
onHangup(success)
}
}
}()
}
func (ac *AsteriskCollector) Close() {
ac.client.Disconnect()
}
// --- helpers ---
// coreShowChannelsCount получает количество активных каналов.
func (ac *AsteriskCollector) coreShowChannelsCount() (int, error) {
cmd := "Action: Command\r\nCommand: core show channels count\r\n"
headers, body, err := ac.client.SendCommandBody(cmd)
if err != nil {
return 0, err
}
_ = headers
// Ответ выглядит как: "1 active call" или "5 active calls"
body = strings.TrimSpace(body)
fields := strings.Fields(body)
if len(fields) >= 3 {
// Первое слово — число
n, err := strconv.Atoi(fields[0])
if err == nil {
return n, nil
}
}
return 0, fmt.Errorf("не удалось распарсить количество каналов из: %s", body)
}
// coreSettingsInt получает целочисленное значение из CoreSettings.
func (ac *AsteriskCollector) coreSettingsInt(key string) (int, error) {
headers, _, err := ac.client.SendCommandBody("Action: CoreSettings\r\n")
if err != nil {
return 0, err
}
val := headers[key]
if val == "" {
return 0, fmt.Errorf("поле %s не найдено", key)
}
return strconv.Atoi(strings.TrimSpace(val))
}
@@ -0,0 +1,137 @@
package agent
import (
"context"
"fmt"
"strconv"
"strings"
"time"
"github.com/pulse-lets-go/internal/esl"
)
// FSCollector собирает метрики с FreeSWITCH через ESL.
type FSCollector struct {
client *esl.Client
}
// NewFSCollector создаёт FreeSWITCH-коллектор.
func NewFSCollector(host string, port int, password string) *FSCollector {
return &FSCollector{
client: esl.NewClient(host, port, password),
}
}
func (fc *FSCollector) Type() string { return "freeswitch" }
func (fc *FSCollector) Connect() error {
// ESL client подключается синхронно через tryConnect
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
fc.client.ConnectWithRetry()
// Ждём подключения с таймаутом
timer := time.NewTimer(100 * time.Millisecond)
defer timer.Stop()
for !fc.client.IsConnected() {
timer.Reset(100 * time.Millisecond)
select {
case <-ctx.Done():
return fmt.Errorf("esl connect timeout")
case <-timer.C:
return nil
}
}
return nil
}
func (fc *FSCollector) Collect() (*PBXResult, error) {
if !fc.client.IsConnected() {
return nil, fmt.Errorf("esl не подключён")
}
result := &PBXResult{Status: "ok"}
// active_calls: api show channels count
if active, err := fc.apiInt("show channels count"); err == nil {
result.ActiveCalls = active
}
// max_calls: api json status → max_sessions
if max, err := fc.apiJSONInt("json", "status", "max_sessions"); err == nil {
result.MaxCalls = max
}
// idle_cpu: eval $${idle_cpu} — FS-специфичный
if idle, err := fc.apiFloat("eval $${idle_cpu}"); err == nil && idle > 0 {
result.IdleCPU = idle
}
return result, nil
}
func (fc *FSCollector) ListenHangups(onHangup func(success bool)) {
// Подписываемся на CHANNEL_HANGUP через ESL
fc.client.Subscribe("CHANNEL_HANGUP")
go func() {
// Для FS используем обработку событий через eventsCh
// Упрощённо: не реализуем полноценную подписку на первом этапе
// failure rate остаётся на нуле
time.Sleep(1 * time.Second)
}()
}
func (fc *FSCollector) Close() {
fc.client.Disconnect()
}
// --- helpers ---
func (fc *FSCollector) apiInt(cmd string) (int, error) {
headers, body, err := fc.client.Send("api " + cmd)
if err != nil {
return 0, err
}
_ = headers
val := strings.TrimSpace(body)
n, err := strconv.Atoi(val)
if err != nil {
// иногда ответ содержит текст помимо числа
n, _ = strconv.Atoi(strings.Fields(val)[0])
}
return n, nil
}
func (fc *FSCollector) apiFloat(cmd string) (float64, error) {
_, body, err := fc.client.Send("api " + cmd)
if err != nil {
return 0, err
}
val := strings.TrimSpace(body)
f, _ := strconv.ParseFloat(val, 64)
return f, nil
}
// apiJSONInt парсит JSON-ответ от FS api: api json status → field1.field2
func (fc *FSCollector) apiJSONInt(cmd string, field string, subfield string) (int, error) {
_, body, err := fc.client.Send("api " + cmd)
if err != nil {
return 0, err
}
// Быстрый парсинг: ищем "subfield": число
search := fmt.Sprintf(`"%s":`, subfield)
idx := strings.Index(body, search)
if idx < 0 {
return 0, fmt.Errorf("поле %s не найдено в ответе", subfield)
}
rest := body[idx+len(search):]
rest = strings.TrimLeft(rest, " \t")
end := strings.IndexAny(rest, ",\n\r}")
if end > 0 {
rest = rest[:end]
}
return strconv.Atoi(strings.TrimSpace(rest))
}
+124
View File
@@ -0,0 +1,124 @@
// Package agent реализует агент сбора метрик для pulse-lets-go.
// Агент работает на каждой PBX-ноде (FreeSWITCH или Asterisk),
// собирает метрики и публикует их в NATS каждые 5 секунд.
//
// Поддерживаемые типы PBX:
// - "freeswitch" — сбор через ESL (Event Socket Library)
// - "asterisk" — сбор через AMI (Asterisk Manager Interface)
//
// Пример agent.json:
//
// {
// "node_id": "uc06",
// "type": "asterisk",
// "nats_url": "nats://10.101.60.81:4222",
// "sip_gateway": "sip:10.101.60.115:5060",
// "interval_sec": 5,
// "max_calls": 150,
// "ami": { "host": "127.0.0.1", "port": 6154, "username": "ctt", "password": "cttpass" }
// }
package agent
import (
"encoding/json"
"fmt"
"os"
)
// Config — конфигурация агента (agent.json).
type Config struct {
NodeID string `json:"node_id"` // уникальный ID ноды (uc06, ses-sip)
Type string `json:"type"` // "freeswitch" или "asterisk"
NatsURL string `json:"nats_url"` // NATS URL
NatsUser string `json:"nats_user,omitempty"`
NatsPassword string `json:"nats_password,omitempty"`
IntervalSec int `json:"interval_sec"` // интервал сбора (default: 5)
MaxCalls int `json:"max_calls"` // ёмкость ноды (default: 250)
SIPGateway string `json:"sip_gateway"` // SIP-адрес для auto-транка
SIPGatewayAuto bool `json:"sip_gateway_auto"` // авто-определить из PBX
FailureWindow int `json:"failure_window"` // окно для call_failure_rate (default: 1000)
ESL ESLCfg `json:"esl,omitempty"`
AMI AMICfg `json:"ami,omitempty"`
}
// ESLCfg — настройки подключения к FreeSWITCH ESL.
type ESLCfg struct {
Host string `json:"host"`
Port int `json:"port"`
Password string `json:"password"`
}
// AMICfg — настройки подключения к Asterisk AMI.
type AMICfg struct {
Host string `json:"host"`
Port int `json:"port"`
Username string `json:"username"`
Password string `json:"password"`
}
// LoadConfig загружает конфигурацию агента из JSON-файла.
func LoadConfig(path string) (*Config, error) {
data, err := os.ReadFile(path)
if err != nil {
return nil, fmt.Errorf("чтение %s: %w", path, err)
}
var cfg Config
if err := json.Unmarshal(data, &cfg); err != nil {
return nil, fmt.Errorf("парсинг %s: %w", path, err)
}
// Defaults
if cfg.IntervalSec <= 0 {
cfg.IntervalSec = 5
}
if cfg.MaxCalls <= 0 {
cfg.MaxCalls = 250
}
if cfg.FailureWindow <= 0 {
cfg.FailureWindow = 1000
}
if cfg.ESL.Port == 0 {
cfg.ESL.Port = 8021
}
if cfg.AMI.Port == 0 {
cfg.AMI.Port = 5038
}
return &cfg, nil
}
// Validate проверяет корректность конфигурации.
func (c *Config) Validate() error {
if c.NodeID == "" {
return fmt.Errorf("node_id обязателен")
}
if c.NatsURL == "" {
return fmt.Errorf("nats_url обязателен")
}
switch c.Type {
case "freeswitch":
if c.ESL.Host == "" {
return fmt.Errorf("esl.host обязателен для type=freeswitch")
}
case "asterisk":
if c.AMI.Host == "" {
return fmt.Errorf("ami.host обязателен для type=asterisk")
}
default:
return fmt.Errorf("type должен быть freeswitch или asterisk, получен: %s", c.Type)
}
return nil
}
// GetSIPGateway возвращает SIP-адрес для авто-транка.
// Если SIPGatewayAuto = true — вернётся "" (collector определит сам).
func (c *Config) GetSIPGateway(autoDetected string) string {
if c.SIPGateway != "" {
return c.SIPGateway
}
if c.SIPGatewayAuto {
return autoDetected
}
return ""
}
@@ -0,0 +1,52 @@
package agent
import "sync"
// FailureTracker — скользящее окно для расчёта call_failure_rate.
// Хранит последние N результатов звонков (success/failure),
// вычисляет процент неудачных звонков в реальном времени.
type FailureTracker struct {
mu sync.Mutex
window []bool // true=success, false=failure
position int // позиция следующей записи
size int // текущий размер окна
capacity int // максимальный размер окна
}
// NewFailureTracker создаёт трекер с окном заданного размера.
func NewFailureTracker(capacity int) *FailureTracker {
return &FailureTracker{
window: make([]bool, capacity),
capacity: capacity,
}
}
// Record записывает результат звонка (true=success, false=failure).
func (ft *FailureTracker) Record(success bool) {
ft.mu.Lock()
defer ft.mu.Unlock()
ft.window[ft.position] = success
ft.position = (ft.position + 1) % ft.capacity
if ft.size < ft.capacity {
ft.size++
}
}
// Rate возвращает процент неудачных звонков (0.0..100.0).
func (ft *FailureTracker) Rate() float64 {
ft.mu.Lock()
defer ft.mu.Unlock()
if ft.size == 0 {
return 0
}
failures := 0
for i := 0; i < ft.size; i++ {
if !ft.window[i] {
failures++
}
}
return float64(failures) / float64(ft.size) * 100.0
}
@@ -0,0 +1,57 @@
package agent
import (
"encoding/json"
"fmt"
"log"
natsgo "github.com/nats-io/nats.go"
"github.com/pulse-lets-go/internal/models"
)
// Publisher публикует метрики в NATS.
type Publisher struct {
nc *natsgo.Conn
nodeID string
}
// NewPublisher создаёт NATS publisher.
func NewPublisher(natsURL, user, password, nodeID string) (*Publisher, error) {
opts := []natsgo.Option{
natsgo.ReconnectWait(natsgo.DefaultReconnectWait),
natsgo.MaxReconnects(-1),
}
if user != "" && password != "" {
opts = append(opts, natsgo.UserInfo(user, password))
}
nc, err := natsgo.Connect(natsURL, opts...)
if err != nil {
return nil, fmt.Errorf("nats connect: %w", err)
}
log.Printf("[publisher] NATS подключён к %s", natsURL)
return &Publisher{nc: nc, nodeID: nodeID}, nil
}
// Publish отправляет метрику в NATS: pulse.metrics.<node_id>.
func (p *Publisher) Publish(m *models.NodeMetric) error {
subject := fmt.Sprintf("pulse.metrics.%s", p.nodeID)
data, err := json.Marshal(m)
if err != nil {
return fmt.Errorf("marshal metric: %w", err)
}
if err := p.nc.Publish(subject, data); err != nil {
return fmt.Errorf("publish: %w", err)
}
return nil
}
// Close закрывает NATS-соединение.
func (p *Publisher) Close() {
if p.nc != nil {
p.nc.Close()
}
}
+84
View File
@@ -0,0 +1,84 @@
package agent
import (
"bufio"
"os"
"strconv"
"strings"
"time"
)
// nowTS возвращает текущий unix timestamp.
func nowTS() int64 {
return time.Now().Unix()
}
// SystemStats содержит системные метрики (CPU, load).
type SystemStats struct {
Load1 float64 // load average 1m
Load5 float64 // load average 5m
Load15 float64 // load average 15m
IdleCPU float64 // процент простоя CPU
}
// CollectSystem собирает системные метрики из /proc.
func CollectSystem() *SystemStats {
load1, load5, load15 := loadAvg()
return &SystemStats{
Load1: load1,
Load5: load5,
Load15: load15,
IdleCPU: cpuIdle(),
}
}
// loadAvg читает /proc/loadavg.
func loadAvg() (float64, float64, float64) {
data, err := os.ReadFile("/proc/loadavg")
if err != nil {
return 0, 0, 0
}
parts := strings.Fields(string(data))
if len(parts) < 3 {
return 0, 0, 0
}
load1, _ := strconv.ParseFloat(parts[0], 64)
load5, _ := strconv.ParseFloat(parts[1], 64)
load15, _ := strconv.ParseFloat(parts[2], 64)
return load1, load5, load15
}
// cpuIdle вычисляет процент простоя CPU из /proc/stat.
func cpuIdle() float64 {
file, err := os.Open("/proc/stat")
if err != nil {
return 0
}
defer file.Close()
scanner := bufio.NewScanner(file)
for scanner.Scan() {
line := scanner.Text()
if !strings.HasPrefix(line, "cpu ") {
continue
}
fields := strings.Fields(line)
if len(fields) < 8 {
return 0
}
// cpu user nice system idle iowait irq softirq steal ...
var total, idle float64
for i, f := range fields[1:] {
val, _ := strconv.ParseFloat(f, 64)
total += val
if i == 3 || i == 4 { // idle + iowait
idle += val
}
}
if total > 0 {
return (idle / total) * 100.0
}
return 0
}
return 0
}
+44
View File
@@ -0,0 +1,44 @@
// pulse-lets-go-agent — агент сбора метрик для pulse-lets-go.
// Работает на каждой PBX-ноде (FreeSWITCH или Asterisk),
// собирает метрики и публикует их в NATS каждые 5 секунд.
//
// Использование:
//
// pulse-lets-go-agent -config /etc/pulse-lets-go-agent/agent.json
//
// Пример agent.json:
//
// {
// "node_id": "uc06",
// "type": "asterisk",
// "nats_url": "nats://10.101.60.81:4222",
// "sip_gateway": "sip:10.101.60.115:5060",
// "interval_sec": 5,
// "max_calls": 150,
// "ami": { "host": "127.0.0.1", "port": 6154, "username": "ctt", "password": "cttpass" }
// }
package main
import (
"flag"
"log"
"github.com/pulse-lets-go/cmd/pulse-lets-go-agent/agent"
)
func main() {
configPath := flag.String("config", "agent.json", "путь к agent.json")
flag.Parse()
cfg, err := agent.LoadConfig(*configPath)
if err != nil {
log.Fatalf("ошибка загрузки конфига: %v", err)
}
a, err := agent.NewAgent(cfg)
if err != nil {
log.Fatalf("ошибка создания агента: %v", err)
}
a.Start()
}
+23 -20
View File
@@ -5,8 +5,6 @@ package main
import (
"context"
"crypto/rand"
"encoding/hex"
"fmt"
"log"
"net/http"
@@ -24,15 +22,17 @@ import (
filelog "github.com/pulse-lets-go/internal/log"
"github.com/pulse-lets-go/internal/models"
"github.com/pulse-lets-go/internal/nats"
"github.com/pulse-lets-go/internal/util"
)
func main() {
// Определяем рабочую директорию
// По умолчанию data/ — родительская папка от бинарника.
// Пример: /opt/pulse-lets-go/bin/pulse-lets-go → /opt/pulse-lets-go/data
exe, _ := os.Executable()
exeDir := filepath.Dir(exe)
dataDir := filepath.Join(exeDir, "data")
dataDir := filepath.Join(filepath.Dir(exeDir), "data")
// Разрешаем переопределение через переменную окружения
if envDir := os.Getenv("PULSE_DATA_DIR"); envDir != "" {
dataDir = envDir
}
@@ -91,7 +91,7 @@ func main() {
}
// Создаём новый balance-транк
now := time.Now().UTC()
trunkID := "trk-" + randomHex(6)
trunkID := "trk-" + util.RandomHex(6)
trunk := models.Trunk{
ID: trunkID,
Name: fmt.Sprintf("%s баланс", nodeID),
@@ -113,7 +113,7 @@ func main() {
})
// 6. HTTP API
apiHandler := api.NewAPI(eng, cfgMgr, cfg.JWTSecret, cfg.MonitoringAPIKey, sub.IsConnected, cfg)
apiHandler := api.NewAPI(eng, cfgMgr, cfg.GetJWTSecret(), cfg.MonitoringAPIKey, sub.IsConnected, cfg)
handler := apiHandler.Handler()
// 6.1 ESL-клиент (опционально — если сконфигурирован)
@@ -177,6 +177,19 @@ func main() {
}
}()
// 6.2 Periodic WS broadcast + FS stats poll
go func() {
ticker := time.NewTicker(5 * time.Second)
for range ticker.C {
if eslClient != nil && eslClient.IsConnected() {
if stats, err := eslClient.FetchFsStats(); err == nil {
apiHandler.UpdateFsStats(stats)
}
}
apiHandler.BroadcastAll()
}
}()
// Подключаем раздачу SvelteKit статики, если директория существует
webDir := findWebDir(exeDir)
if webDir != "" {
@@ -230,20 +243,19 @@ func main() {
log.Printf("[main] ошибка остановки HTTP сервера: %v", err)
}
sub.Close()
metricsLogger.Close()
log.Println("[main] pulse-lets-go остановлен")
}
// findWebDir ищет директорию со SvelteKit сборкой (web/build/).
// Пробует по порядку: переменная окружения, рядом с бинарником, в корне проекта, CWD.
// Пробует по порядку: PULSE_WEB_DIR, CWD, ../web/build, рядом с бинарником.
func findWebDir(exeDir string) string {
candidates := []string{
os.Getenv("PULSE_WEB_DIR"),
filepath.Join(exeDir, "web", "build"), // bin/web/build
filepath.Join(filepath.Dir(exeDir), "web", "build"), // ../web/build (корень проекта)
"web/build", // ./web/build (CWD)
"web/build", // ./web/build (CWD — dev, systemd working dir)
filepath.Join(filepath.Dir(exeDir), "web", "build"), // ../web/build (корень проекта, dev с bin/)
filepath.Join(exeDir, "web", "build"), // рядом с бинарником (production install, make deploy)
}
for _, d := range candidates {
if d == "" {
@@ -294,15 +306,6 @@ func withSPA(webDir string, apiHandler http.Handler) http.Handler {
})
}
// randomHex генерирует случайную hex-строку заданной длины.
func randomHex(n int) string {
b := make([]byte, n)
if _, err := rand.Read(b); err != nil {
return fmt.Sprintf("%x", time.Now().UnixNano())[:n]
}
return hex.EncodeToString(b)[:n]
}
// extractHostFromGateway извлекает хост из SIP URI.
// "sip:mts-gw.lan:5060" → "mts-gw.lan"
func extractHostFromGateway(gateway string) string {
+7 -11
View File
@@ -1,16 +1,16 @@
// SIP тестер для e2e проверки pulse-lets-go + FreeSWITCH.
// Два режима:
//
// uas — отвечает 200 OK на SIP INVITE, симулирует PBX-ноду
// uac — шлёт SIP INVITE в FreeSWITCH, симулирует оператора
//
// Примеры:
//
// ./siptest -mode uas (PBX-нода на порту 5090)
// ./siptest -mode uac -r 10 -l 100 (оператор, 10 CPS, до 100 одновременных)
package main
import (
"crypto/rand"
"encoding/hex"
"flag"
"fmt"
"log"
@@ -19,6 +19,8 @@ import (
"sync"
"sync/atomic"
"time"
"github.com/pulse-lets-go/internal/util"
)
const (
@@ -103,7 +105,7 @@ func buildSIPResponse(invite, localAddr string) string {
// Добавляем tag к To если его нет
to = strings.TrimSpace(to)
if !strings.Contains(to, ";tag=") {
to += ";tag=uas-" + randomHex(4)
to += ";tag=uas-" + util.RandomHex(4)
}
return fmt.Sprintf("SIP/2.0 200 OK\r\n"+
@@ -213,8 +215,8 @@ func runUAC(rate, limit, maxCalls int, dest, caller, host string) {
}
func buildSIPInvite(dest, caller string, callNum int, srcAddr, fsAddr string) string {
callID := fmt.Sprintf("call-%d-%s@%s", callNum, randomHex(4), "127.0.0.1")
branch := "z9hG4bK-" + randomHex(8)
callID := fmt.Sprintf("call-%d-%s@%s", callNum, util.RandomHex(4), "127.0.0.1")
branch := "z9hG4bK-" + util.RandomHex(8)
return fmt.Sprintf("INVITE sip:%s@%s SIP/2.0\r\n"+
"Via: SIP/2.0/UDP %s;branch=%s\r\n"+
@@ -227,9 +229,3 @@ func buildSIPInvite(dest, caller string, callNum int, srcAddr, fsAddr string) st
"\r\n",
dest, fsAddr, srcAddr, branch, caller, callNum, dest, callID, caller, srcAddr)
}
func randomHex(n int) string {
b := make([]byte, n)
rand.Read(b)
return hex.EncodeToString(b)[:n]
}
+15
View File
@@ -0,0 +1,15 @@
{
"node_id": "ses-sip",
"type": "freeswitch",
"nats_url": "nats://10.101.60.81:4222",
"interval_sec": 5,
"max_calls": 250,
"sip_gateway": "sip:10.3.0.44:5060",
"sip_gateway_auto": false,
"failure_window": 1000,
"esl": {
"host": "127.0.0.1",
"port": 8021,
"password": "ClueCon"
}
}
+16
View File
@@ -0,0 +1,16 @@
{
"node_id": "uc06",
"type": "asterisk",
"nats_url": "nats://10.101.60.81:4222",
"interval_sec": 5,
"max_calls": 150,
"sip_gateway": "sip:10.101.60.115:5060",
"sip_gateway_auto": false,
"failure_window": 1000,
"ami": {
"host": "127.0.0.1",
"port": 6154,
"username": "ctt",
"password": "cttpass"
}
}
+4 -6
View File
@@ -1,16 +1,14 @@
[Unit]
Description=pulse-lets-go — Telephone Load Balancer
Documentation=https://github.com/pulse-lets-go
After=network.target nats.service
Wants=nats.service
After=network.target
[Service]
Type=simple
User=pulse
Group=pulse
ExecStart=/usr/local/bin/pulse-lets-go
WorkingDirectory=/usr/local/share/pulse-lets-go
Environment=PULSE_DATA_DIR=/etc/pulse-lets-go
ExecStart=/opt/pulse-lets-go/bin/pulse-lets-go
WorkingDirectory=/opt/pulse-lets-go
Restart=always
RestartSec=5
@@ -19,7 +17,7 @@ NoNewPrivileges=yes
PrivateTmp=yes
ProtectSystem=strict
ProtectHome=yes
ReadWritePaths=/etc/pulse-lets-go /var/log/pulse-lets-go
ReadWritePaths=/opt/pulse-lets-go/data /opt/pulse-lets-go/log
# Логирование
StandardOutput=journal
+317
View File
@@ -0,0 +1,317 @@
// 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 := net.JoinHostPort(c.host, fmt.Sprintf("%d", 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] соединение разорвано, реконнект...")
c.reconnect()
continue
}
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)
}
}
}
// reconnect выполняет реконнект с backoff (без горутин — синхронный вызов).
func (c *Client) reconnect() {
if !c.shouldReconn.Load() {
return
}
log.Printf("[ami] реконнект через %v...", c.backoff)
time.Sleep(c.backoff)
if c.conn != nil {
c.conn.Close()
c.conn = nil
}
if err := c.Connect(); err != nil {
log.Printf("[ami] реконнект не удался: %v", err)
c.backoff *= 2
if c.backoff > 60*time.Second {
c.backoff = 60 * time.Second
}
return
}
c.backoff = 1 * time.Second
log.Printf("[ami] реконнект успешен")
}
+1 -5
View File
@@ -139,7 +139,7 @@ func (a *API) apiKeyMiddleware(next http.Handler) http.Handler {
// handleLogin обрабатывает POST /api/auth/login.
func (a *API) handleLogin(w http.ResponseWriter, r *http.Request) {
var req models.LoginRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
if err := decodeJSON(r, &req); err != nil {
writeError(w, http.StatusBadRequest, "некорректный JSON")
return
}
@@ -199,10 +199,6 @@ func (a *API) handleLogin(w http.ResponseWriter, r *http.Request) {
// handleRefresh выдаёт новый access токен по refresh токену.
func (a *API) handleRefresh(w http.ResponseWriter, r *http.Request) {
tokenStr := extractBearerToken(r)
if tokenStr == "" {
// пробуем получить из query param (для WebSocket)
tokenStr = r.URL.Query().Get("token")
}
if tokenStr == "" {
writeError(w, http.StatusUnauthorized, "требуется токен")
return
+11
View File
@@ -1,10 +1,13 @@
package api
import (
"bufio"
"crypto/rand"
"encoding/hex"
"encoding/json"
"fmt"
"log"
"net"
"net/http"
"time"
)
@@ -82,3 +85,11 @@ func (w *loggingResponseWriter) WriteHeader(code int) {
func (w *loggingResponseWriter) Write(b []byte) (int, error) {
return w.ResponseWriter.Write(b)
}
func (w *loggingResponseWriter) Hijack() (net.Conn, *bufio.ReadWriter, error) {
hijacker, ok := w.ResponseWriter.(http.Hijacker)
if !ok {
return nil, nil, fmt.Errorf("loggingResponseWriter: underlying ResponseWriter does not implement http.Hijacker")
}
return hijacker.Hijack()
}
+20
View File
@@ -38,6 +38,26 @@ func (a *API) handlePrometheus(w http.ResponseWriter, r *http.Request) {
var sb strings.Builder
// Общие метрики (counter gauge)
stats := a.engine.GetHealthStats()
fmt.Fprintf(&sb, "# HELP pulse_route_requests_total Total number of route requests.\n")
fmt.Fprintf(&sb, "# TYPE pulse_route_requests_total counter\n")
fmt.Fprintf(&sb, "pulse_route_requests_total %d\n", stats.RouteRequests)
fmt.Fprintf(&sb, "# HELP pulse_nodes_total Total nodes registered.\n")
fmt.Fprintf(&sb, "# TYPE pulse_nodes_total gauge\n")
fmt.Fprintf(&sb, "pulse_nodes_total %d\n", stats.TotalNodes)
fmt.Fprintf(&sb, "# HELP pulse_nodes_healthy Healthy nodes count.\n")
fmt.Fprintf(&sb, "# TYPE pulse_nodes_healthy gauge\n")
fmt.Fprintf(&sb, "pulse_nodes_healthy %d\n", stats.HealthyNodes)
fmt.Fprintf(&sb, "# HELP pulse_uptime_seconds Uptime in seconds.\n")
fmt.Fprintf(&sb, "# TYPE pulse_uptime_seconds gauge\n")
fmt.Fprintf(&sb, "pulse_uptime_seconds %d\n", stats.UptimeSeconds)
fmt.Fprintf(&sb, "\n")
for _, n := range nodes {
id := sanitizePromLabel(n.NodeID)
+2
View File
@@ -14,6 +14,7 @@ func (a *API) handleRoute(w http.ResponseWriter, r *http.Request) {
// Нет зарегистрированных нод
if nodeID == "" {
a.engine.IncrementRouteFallbacks()
writeJSON(w, http.StatusServiceUnavailable, models.RouteResponse{
Error: "no_nodes_registered",
})
@@ -22,6 +23,7 @@ func (a *API) handleRoute(w http.ResponseWriter, r *http.Request) {
// Все ноды unhealthy — fallback
if fallback || score < 0 {
a.engine.IncrementRouteFallbacks()
fallbackGW, _ := a.findFallbackGateway()
nodes := a.getRouteNodeInfo()
+38 -1
View File
@@ -4,6 +4,7 @@ import (
"encoding/json"
"fmt"
"net/http"
"sync"
"github.com/pulse-lets-go/internal/config"
"github.com/pulse-lets-go/internal/engine"
@@ -23,6 +24,11 @@ type API struct {
logFormat string
routeLimiter *rateLimiter
apiLimiter *rateLimiter
gatewayMu sync.RWMutex
gatewayStates map[string]string // trunkID → "up"/"down"
fsStats esl.FsStats
fsStatsMu sync.RWMutex
}
// NewAPI создаёт новый HTTP API с заданными зависимостями.
@@ -35,6 +41,7 @@ func NewAPI(eng *engine.Engine, cfgMgr *config.Manager, jwtSecret, monitoringAPI
wsHub: newWSHub(),
natsConnected: natsFn,
logFormat: cfg.Log.Format,
gatewayStates: make(map[string]string),
}
if cfg.RateLimit.Enabled {
a.routeLimiter = newRateLimiter(cfg.RateLimit.RoutePerSec, cfg.RateLimit.RoutePerSec)
@@ -130,6 +137,10 @@ func corsMiddleware(next http.Handler) http.Handler {
// BroadcastGatewayEvent транслирует статус gateway всем WS-клиентам.
func (a *API) BroadcastGatewayEvent(gatewayName, trunkID, status string) {
a.gatewayMu.Lock()
a.gatewayStates[trunkID] = status
a.gatewayMu.Unlock()
msg := &models.WsMetricsMessage{
Type: "gateway_update",
Payload: map[string]string{
@@ -141,9 +152,35 @@ func (a *API) BroadcastGatewayEvent(gatewayName, trunkID, status string) {
a.wsHub.broadcast(msg)
}
// decodeJSON декодирует тело запроса в структуру v.
func (a *API) gatewayStats() (total, up, down int) {
a.gatewayMu.RLock()
defer a.gatewayMu.RUnlock()
for _, s := range a.gatewayStates {
total++
switch s {
case "up":
up++
case "down":
down++
}
}
return
}
// UpdateFsStats обновляет кэш метрик FreeSWITCH (вызывается из ticker).
func (a *API) UpdateFsStats(stats esl.FsStats) {
a.fsStatsMu.Lock()
a.fsStats = stats
a.fsStatsMu.Unlock()
}
// maxBodySize — максимальный размер тела запроса в байтах (1 MiB).
const maxBodySize = 1 << 20
// decodeJSON декодирует тело запроса в структуру v с ограничением размера.
func decodeJSON(r *http.Request, v interface{}) error {
defer r.Body.Close()
r.Body = http.MaxBytesReader(nil, r.Body, maxBodySize)
if err := json.NewDecoder(r.Body).Decode(v); err != nil {
return fmt.Errorf("декодирование JSON: %w", err)
}
+3 -2
View File
@@ -268,8 +268,9 @@ func (a *API) eslPushGatewayDelete(trunk models.Trunk) {
// gatewayParams возвращает параметры для создания FS gateway на основе транка.
func (a *API) gatewayParams(trunk models.Trunk) (profile, name, proxy string) {
profile = a.readConfig().ESL.SofiaProfile
name = fmt.Sprintf("pulse-ingress-%s", trunk.ID)
cfg := a.readConfig()
profile = cfg.ESL.SofiaProfile
name = fmt.Sprintf("%s-ingress-%s", cfg.ESL.GatewayPrefix, trunk.ID)
proxy = extractHost(trunk.Gateway)
return
}
+2 -13
View File
@@ -1,13 +1,12 @@
package api
import (
"crypto/rand"
"encoding/hex"
"fmt"
"net/http"
"time"
"github.com/pulse-lets-go/internal/models"
"github.com/pulse-lets-go/internal/util"
"golang.org/x/crypto/bcrypt"
)
@@ -227,15 +226,5 @@ func findUserIndex(users []models.User, id string) int {
// generateID генерирует ID вида "prefix-XXXXX".
func generateID(prefix string) string {
return fmt.Sprintf("%s-%s", prefix, randomHex(6))
}
// randomHex генерирует случайную hex-строку заданной длины.
func randomHex(n int) string {
b := make([]byte, n)
if _, err := rand.Read(b); err != nil {
// fallback на time-based, если crypto/rand недоступен
return fmt.Sprintf("%x", time.Now().UnixNano())
}
return hex.EncodeToString(b)[:n]
return fmt.Sprintf("%s-%s", prefix, util.RandomHex(6))
}
+42 -7
View File
@@ -104,12 +104,7 @@ func (a *API) handleWS(w http.ResponseWriter, r *http.Request) {
// Отправляем текущее состояние при подключении
go func() {
nodes := a.engine.GetAllNodes()
msg := &models.WsMetricsMessage{
Type: "nodes_update",
Payload: nodes,
}
a.wsHub.broadcast(msg)
a.BroadcastAll()
}()
// Читаем из вебсокета (ping/pong + закрытие)
@@ -125,7 +120,7 @@ func (a *API) handleWS(w http.ResponseWriter, r *http.Request) {
}()
}
// BroadcastMetrics рассылает метрики всем WS-клиентам (вызывается engine при обновлении).
// BroadcastMetrics рассылает метрики нод всем WS-клиентам.
func (a *API) BroadcastMetrics() {
nodes := a.engine.GetAllNodes()
msg := &models.WsMetricsMessage{
@@ -135,6 +130,46 @@ func (a *API) BroadcastMetrics() {
a.wsHub.broadcast(msg)
}
// BroadcastAll рассылает ноды и состояние балансировщика.
func (a *API) BroadcastAll() {
a.BroadcastMetrics()
a.broadcastBalancerUpdate()
}
func (a *API) broadcastBalancerUpdate() {
stats := a.engine.GetHealthStats()
gwTotal, gwUp, gwDown := a.gatewayStats()
info := models.BalancerInfo{
NatsConnected: a.natsConnected != nil && a.natsConnected(),
GatewaysTotal: gwTotal,
GatewaysUp: gwUp,
GatewaysDown: gwDown,
RouteTotal: stats.RouteRequests,
RouteFallbacks: stats.RouteFallbacks,
UptimeSec: stats.UptimeSeconds,
}
if a.eslClient != nil {
eslStats := a.eslClient.GetStats()
info.EslConnected = eslStats.Status == "connected"
info.EslUptimeSec = eslStats.UptimeSec
info.EslReconnects = eslStats.Reconnects
a.fsStatsMu.RLock()
info.FsActiveCalls = a.fsStats.ActiveCalls
info.FsMaxSessions = a.fsStats.MaxSessions
info.FsUptime = a.fsStats.FsUptime
a.fsStatsMu.RUnlock()
}
msg := &models.WsMetricsMessage{
Type: "balancer_update",
Payload: info,
}
a.wsHub.broadcast(msg)
}
// validateToken проверяет JWT и возвращает UserInfo.
func (a *API) validateToken(tokenStr string) (*UserInfo, error) {
token, err := parseJWT(tokenStr, a.jwtSecret)
+11
View File
@@ -72,6 +72,16 @@ func (e *ESLConfig) GetPassword() string {
return e.Password
}
// GetJWTSecret возвращает JWT-секрет: из переменной окружения, если задан jwt_secret_env, иначе из поля jwt_secret.
func (c *Config) GetJWTSecret() string {
if c.JWTSecretEnv != "" {
if val := os.Getenv(c.JWTSecretEnv); val != "" {
return val
}
}
return c.JWTSecret
}
// Config — корневая конфигурация приложения (config.json).
type Config struct {
NatsURL string `json:"nats_url"`
@@ -79,6 +89,7 @@ type Config struct {
NatsPassword string `json:"nats_password"`
ListenAddr string `json:"listen_addr"`
JWTSecret string `json:"jwt_secret"`
JWTSecretEnv string `json:"jwt_secret_env"` // имя переменной окружения с JWT-секретом
MonitoringAPIKey string `json:"monitoring_api_key"`
StaleThresholdSec int `json:"stale_threshold_sec"`
Scoring ScoringConfig `json:"scoring"`
+8
View File
@@ -29,6 +29,7 @@ type Engine struct {
// Статистика для health endpoint
startTime time.Time
routeRequests atomic.Int64
routeFallbacks atomic.Int64
lastMetricTime atomic.Int64 // unix ts последней полученной метрики
}
@@ -188,11 +189,17 @@ func (e *Engine) IncrementRouteRequests() {
e.routeRequests.Add(1)
}
// IncrementRouteFallbacks увеличивает счётчик fallback-запросов.
func (e *Engine) IncrementRouteFallbacks() {
e.routeFallbacks.Add(1)
}
// HealthStats возвращает статистику для health endpoint.
type HealthStats struct {
TotalNodes int `json:"total_nodes"`
HealthyNodes int `json:"healthy_nodes"`
RouteRequests int64 `json:"route_requests_total"`
RouteFallbacks int64 `json:"route_fallbacks_total"`
UptimeSeconds int64 `json:"uptime_seconds"`
LastMetricTS int64 `json:"last_metric_ts"`
}
@@ -212,6 +219,7 @@ func (e *Engine) GetHealthStats() HealthStats {
TotalNodes: len(e.nodes),
HealthyNodes: healthy,
RouteRequests: e.routeRequests.Load(),
RouteFallbacks: e.routeFallbacks.Load(),
UptimeSeconds: int64(time.Since(e.startTime).Seconds()),
LastMetricTS: e.lastMetricTime.Load(),
}
+75 -47
View File
@@ -109,7 +109,7 @@ func (c *Client) tryConnect() error {
c.closeCh = make(chan struct{})
}
addr := fmt.Sprintf("%s:%d", c.host, c.port)
addr := net.JoinHostPort(c.host, fmt.Sprintf("%d", c.port))
dialer := net.Dialer{Timeout: 5 * time.Second}
conn, err := dialer.Dial("tcp", addr)
if err != nil {
@@ -119,7 +119,7 @@ func (c *Client) tryConnect() error {
c.reader = bufio.NewReader(conn)
// 1. Читаем auth/request
headers, _, err := c.readMessageLocked()
headers, _, err := c.readMessage()
if err != nil {
conn.Close()
return fmt.Errorf("чтение auth/request: %w", err)
@@ -137,7 +137,7 @@ func (c *Client) tryConnect() error {
}
// 3. Читаем auth response
headers, _, err = c.readMessageLocked()
headers, _, err = c.readMessage()
if err != nil {
conn.Close()
return fmt.Errorf("чтение auth response: %w", err)
@@ -199,7 +199,7 @@ func (c *Client) Send(command string) (map[string]string, string, error) {
}
// Читаем ответ
return c.readMessageLocked()
return c.readMessage()
}
// Subscribe подписывается на события ESL.
@@ -263,7 +263,7 @@ func (c *Client) readEventsLoop() {
continue
}
headers, body, err := c.readMessageUnlocked()
headers, body, err := c.readMessage()
if err != nil {
if err == io.EOF || strings.Contains(err.Error(), "use of closed network connection") {
@@ -330,48 +330,9 @@ func (c *Client) reconnect() {
// --- Приватные методы ---
// readMessageUnlocked читает одно ESL-сообщение: заголовки + тело.
// Должна вызываться только из readEventsLoop (одиночный читатель).
func (c *Client) readMessageUnlocked() (map[string]string, string, error) {
headers := make(map[string]string)
for {
line, err := c.reader.ReadString('\n')
if err != nil {
return nil, "", err
}
line = strings.TrimRight(line, "\r\n")
if line == "" {
break // конец заголовков
}
parts := strings.SplitN(line, ":", 2)
if len(parts) == 2 {
headers[strings.TrimSpace(parts[0])] = strings.TrimSpace(parts[1])
}
}
// Читаем тело согласно Content-Length
cl := headers["Content-Length"]
if cl == "" {
return headers, "", nil
}
length, err := strconv.Atoi(cl)
if err != nil {
return headers, "", fmt.Errorf("некорректный Content-Length: %s", cl)
}
body := make([]byte, length)
if _, err := io.ReadFull(c.reader, body); err != nil {
return headers, "", fmt.Errorf("чтение тела (Content-Length=%d): %w", length, err)
}
return headers, strings.TrimSpace(string(body)), nil
}
// readMessageLocked читает одно ESL-сообщение: заголовки + тело.
// Вызывающий держит c.mu (для синхронных команд Send).
func (c *Client) readMessageLocked() (map[string]string, string, error) {
// readMessage читает одно ESL-сообщение: заголовки + тело.
// Вызывающий должен гарантировать монопольный доступ к c.reader.
func (c *Client) readMessage() (map[string]string, string, error) {
headers := make(map[string]string)
for {
@@ -440,3 +401,70 @@ func (c *Client) GetStats() Stats {
func (c *Client) IncGatewayOps() {
c.gatewayOps.Add(1)
}
// --- FS stats ---
// FsStats — показатели FreeSWITCH балансировщика.
type FsStats struct {
ActiveCalls int
MaxSessions int
FsUptime string
}
// FetchFsStats запрашивает метрики FreeSWITCH через api команды.
func (c *Client) FetchFsStats() (FsStats, error) {
if !c.connected.Load() {
return FsStats{}, fmt.Errorf("esl не подключён")
}
var stats FsStats
_, body, err := c.Send("api show calls count")
if err != nil {
return FsStats{}, err
}
stats.ActiveCalls = parseFirstInt(body)
_, body, err = c.Send("api status")
if err != nil {
return stats, err
}
stats.MaxSessions = parseMaxSessions(body)
stats.FsUptime = parseUptime(body)
return stats, nil
}
func parseFirstInt(s string) int {
s = strings.TrimSpace(s)
for i, c := range s {
if c < '0' || c > '9' {
if i == 0 {
return 0
}
v, err := strconv.Atoi(s[:i])
if err != nil {
return 0
}
return v
}
}
return 0
}
func parseMaxSessions(body string) int {
for _, line := range strings.Split(body, "\n") {
if strings.Contains(line, "session(s) max") {
return parseFirstInt(line)
}
}
return 0
}
func parseUptime(body string) string {
for _, line := range strings.Split(body, "\n") {
if strings.HasPrefix(line, "UP ") {
return strings.TrimPrefix(line, "UP ")
}
}
return ""
}
+1
View File
@@ -16,6 +16,7 @@ type EventHandler func(eventName string, headers map[string]string, body string)
// - Вызывает onConnect (например, для GatewaySyncAll)
// - Подписывается на события SOFIA::gateway_register/unregister
// - Запускает чтение асинхронных событий
//
// При разрыве соединения — автоматический реконнект.
func (c *Client) StartEventLoop(ctx context.Context, profile string, onConnect func(), onEvent EventHandler) {
const subscribeEvents = "SOFIA::gateway_register SOFIA::gateway_unregister SOFIA::gateway_expire"
+4 -4
View File
@@ -32,8 +32,8 @@ func (c *Client) GatewayDelete(profile, name string) error {
return nil
}
// GatewayList возвращает список имён gateway в заданном профиле.
func (c *Client) GatewayList(profile string) ([]string, error) {
// GatewayList возвращает список имён gateway в заданном профиле с указанным префиксом.
func (c *Client) GatewayList(profile, prefix string) ([]string, error) {
cmd := fmt.Sprintf("api sofia profile %s gwlist", profile)
_, body, err := c.Send(cmd)
@@ -49,8 +49,8 @@ func (c *Client) GatewayList(profile string) ([]string, error) {
if line == "" || strings.HasPrefix(line, "Name:") || strings.HasPrefix(line, "====") {
continue
}
// Ищем строки вида "pulse-ingress-xxxxx sip:..."
if strings.Contains(line, "pulse-") {
// Ищем строки вида "<prefix>-ingress-xxxxx sip:..."
if strings.Contains(line, prefix+"-ingress-") {
parts := strings.Fields(line)
if len(parts) > 0 {
gateways = append(gateways, parts[0])
+20 -1
View File
@@ -209,10 +209,29 @@ type ZabbixResponse struct {
Data []ZabbixNode `json:"data"`
}
// --- Balancer ---
// BalancerInfo — состояние балансировщика для дашборда.
type BalancerInfo struct {
EslConnected bool `json:"esl_connected"`
EslUptimeSec int64 `json:"esl_uptime_sec,omitempty"`
EslReconnects int64 `json:"esl_reconnects,omitempty"`
NatsConnected bool `json:"nats_connected"`
GatewaysTotal int `json:"gateways_total"`
GatewaysUp int `json:"gateways_up"`
GatewaysDown int `json:"gateways_down"`
FsActiveCalls int `json:"fs_active_calls"`
FsMaxSessions int `json:"fs_max_sessions"`
FsUptime string `json:"fs_uptime"`
RouteTotal int64 `json:"route_total"`
RouteFallbacks int64 `json:"route_fallbacks"`
UptimeSec int64 `json:"uptime_sec"`
}
// --- WebSocket ---
// WsMetricsMessage — сообщение, отправляемое по WebSocket.
type WsMetricsMessage struct {
Type string `json:"type"` // "nodes_update", "node_toggle", etc
Type string `json:"type"` // "nodes_update", "balancer_update", "gateway_update", etc
Payload interface{} `json:"payload"`
}
+1 -1
View File
@@ -9,9 +9,9 @@ import (
natsgo "github.com/nats-io/nats.go"
"github.com/pulse-lets-go/internal/models"
"github.com/pulse-lets-go/internal/engine"
filelog "github.com/pulse-lets-go/internal/log"
"github.com/pulse-lets-go/internal/models"
)
const subject = "pulse.metrics.>"
+18
View File
@@ -0,0 +1,18 @@
package util
import (
"crypto/rand"
"encoding/hex"
"fmt"
"time"
)
// RandomHex генерирует случайную hex-строку из n байт.
// При ошибке crypto/rand использует time-based fallback.
func RandomHex(n int) string {
b := make([]byte, n)
if _, err := rand.Read(b); err != nil {
return fmt.Sprintf("%x", time.Now().UnixNano())[:n*2]
}
return hex.EncodeToString(b)[:n*2]
}
+16
View File
@@ -184,6 +184,22 @@ export async function deleteUser(id: string) {
}
// --- WebSocket ---
export interface BalancerInfo {
esl_connected: boolean;
esl_uptime_sec: number;
esl_reconnects: number;
nats_connected: boolean;
gateways_total: number;
gateways_up: number;
gateways_down: number;
fs_active_calls: number;
fs_max_sessions: number;
fs_uptime: string;
route_total: number;
route_fallbacks: number;
uptime_sec: number;
}
export function connectWebSocket(onMessage: (data: any) => void): WebSocket {
const token = getToken();
const proto = window.location.protocol === 'https:' ? 'wss' : 'ws';
+24 -15
View File
@@ -1,31 +1,40 @@
<script lang="ts">
import { browser } from '$app/environment';
import { page } from '$app/stores';
import { goto } from '$app/navigation';
import '../app.css';
import { Toaster } from 'svelte-sonner';
import { LogIn, LayoutDashboard, Phone, Cable, Users } from 'lucide-svelte';
import { LayoutDashboard, Phone, Cable, Users } from 'lucide-svelte';
let token: string | null = browser ? localStorage.getItem('token') : null;
let user: { username: string; role: string } | null = token
? (() => { try { return JSON.parse(atob(token.split('.')[1])); } catch { return null; } })()
: null;
let token: string | null = null;
let user: { username: string; role: string } | null = null;
function logout() {
if (browser) {
localStorage.removeItem('token');
localStorage.removeItem('refreshToken');
}
token = null;
function updateAuth() {
token = typeof localStorage !== 'undefined' ? localStorage.getItem('token') : null;
if (token) {
try { user = JSON.parse(atob(token.split('.')[1])); } catch { user = null; }
} else {
user = null;
}
}
$: pathname = browser ? window.location.pathname : '';
$: isLogin = pathname === '/login';
function logout() {
localStorage.removeItem('token');
localStorage.removeItem('refreshToken');
updateAuth();
goto('/login');
}
let pathname = $derived($page.url.pathname);
$effect(() => {
$page.url.pathname;
updateAuth();
});
</script>
<Toaster position="top-right" />
{#if !isLogin}
{#if pathname !== '/login'}
<div class="flex min-h-screen">
<!-- Sidebar -->
<aside class="w-64 border-r bg-card flex flex-col">
+71 -2
View File
@@ -4,11 +4,12 @@
import { onMount, onDestroy } from 'svelte';
import {
getNodes, connectWebSocket, isAuthenticated,
type NodeInfo
type NodeInfo, type BalancerInfo
} from '$lib/api';
import { Activity, Phone, Zap, AlertTriangle, Wifi, Server } from 'lucide-svelte';
import { Activity, Phone, Zap, AlertTriangle, Wifi, Server, Plug, Globe, ArrowLeftRight, Network } from 'lucide-svelte';
let nodes: NodeInfo[] = [];
let balancer: BalancerInfo | null = null;
let ws: WebSocket | null = null;
let error: string | null = null;
@@ -21,6 +22,8 @@
ws = connectWebSocket((msg) => {
if (msg.type === 'nodes_update') {
nodes = msg.payload;
} else if (msg.type === 'balancer_update') {
balancer = msg.payload;
}
});
});
@@ -67,10 +70,31 @@
}
}
function dotColor(ok: boolean, degrade?: boolean): string {
if (ok) return 'bg-green-500';
if (degrade) return 'bg-yellow-500';
return 'bg-red-500';
}
function fmtUptime(sec: number): string {
if (sec < 60) return `${sec}с`;
if (sec < 3600) return `${Math.floor(sec / 60)}м`;
if (sec < 86400) return `${Math.floor(sec / 3600)}ч`;
return `${Math.floor(sec / 86400)}д`;
}
function fmtNum(n: number): string {
if (n >= 1000) return (n / 1000).toFixed(1) + 'K';
return String(n);
}
$: healthyCount = nodes.filter(n => n.score >= 0 && !n.disabled && !n.is_stale).length;
$: totalCalls = nodes.reduce((s, n) => s + n.active_calls, 0);
$: avgScore = nodes.length > 0 ? nodes.reduce((s, n) => s + n.score, 0) / nodes.length : 0;
$: unhealthyList = nodes.filter(n => n.score < 0 || n.disabled || n.is_stale);
$: fallbackPct = balancer && balancer.route_total > 0
? (balancer.route_fallbacks / balancer.route_total * 100).toFixed(1)
: '0';
</script>
{#if error}
@@ -79,6 +103,51 @@
</div>
{/if}
<!-- Balancer Status Bar -->
{#if balancer}
<div class="flex flex-wrap items-center gap-1 text-xs text-muted-foreground mb-4 px-3 py-2 border rounded-md bg-muted/30">
<span class="flex items-center gap-1 mr-2">
<span class="w-1.5 h-1.5 rounded-full {dotColor(balancer.esl_connected)}"></span>
<Plug class="h-3 w-3" />
ESL {balancer.esl_connected ? 'Online' : 'Offline'}
{#if balancer.esl_uptime_sec > 0}{fmtUptime(balancer.esl_uptime_sec)}{/if}
</span>
<span class="opacity-30">|</span>
<span class="flex items-center gap-1 mr-2">
<span class="w-1.5 h-1.5 rounded-full {dotColor(balancer.nats_connected)}"></span>
<Network class="h-3 w-3" />
NATS {balancer.nats_connected ? 'Online' : 'Offline'}
</span>
<span class="opacity-30">|</span>
<span class="flex items-center gap-1 mr-2">
<span class="w-1.5 h-1.5 rounded-full {balancer.gateways_down > 0 ? 'bg-red-500' : 'bg-green-500'}"></span>
<Globe class="h-3 w-3" />
GW {balancer.gateways_up}/{balancer.gateways_total}
{#if balancer.gateways_down > 0}
<span class="text-red-500">({balancer.gateways_down} down)</span>
{/if}
</span>
<span class="opacity-30">|</span>
<span class="flex items-center gap-1">
<ArrowLeftRight class="h-3 w-3" />
R {fmtNum(balancer.route_total)}
{#if balancer.route_fallbacks > 0}
<span class="text-orange-500">({fallbackPct}% fb)</span>
{/if}
</span>
{#if balancer.esl_connected && balancer.fs_max_sessions > 0}
<span class="opacity-30">|</span>
<span class="flex items-center gap-1">
<Phone class="h-3 w-3" />
FS {balancer.fs_active_calls}·{balancer.fs_max_sessions}
{#if balancer.fs_uptime}
<span class="opacity-50">{balancer.fs_uptime}</span>
{/if}
</span>
{/if}
</div>
{/if}
<!-- Summary Cards -->
<div class="grid grid-cols-1 sm:grid-cols-2 lg:grid-cols-4 gap-4 mb-6">
<div class="border rounded-lg p-4 bg-card">