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, "max_calls": 250,
"idle_cpu": 35, "idle_cpu": 35,
"load_avg": 0.85, "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. FS использует Lua скрипт в dialplan → `mod_curl` дёргает `/api/route` → получает JSON → парсит `sip_gateway` → bridge.
``` ```
<extension name="route_call"> <extension name="pulse_route">
<condition field="destination_number" expression="^.*$"> <condition field="destination_number" expression="^.*$">
<action application="lua" data="route.lua"/> <action application="lua" data="route.lua"/>
</condition> </condition>
@@ -188,6 +189,8 @@ route.lua: GET `/api/route?...` → если `fallback: true` идёт на fall
```json ```json
{ {
"nats_url": "nats://localhost:4222", "nats_url": "nats://localhost:4222",
"nats_user": "",
"nats_password": "",
"listen_addr": ":8080", "listen_addr": ":8080",
"jwt_secret": "auto-generated-on-first-run", "jwt_secret": "auto-generated-on-first-run",
"monitoring_api_key": "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, "idle_score": 0.20,
"fail_score": 0.10 "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/ pulse-lets-go/
├── cmd/pulse-lets-go/main.go ├── cmd/pulse-lets-go/main.go
├── cmd/pulse-lets-go-agent/
├── cmd/emulator/
├── cmd/siptest/
├── internal/ ├── internal/
│ ├── api/ — handlers: auth, nodes, trunks, users, monitoring, ws │ ├── api/ — handlers: auth, nodes, trunks, users, monitoring, ws
│ ├── engine/ — scorer + router (sync.RWMutex, best_node кэш) │ ├── engine/ — scorer + router (sync.RWMutex, best_node кэш)
│ ├── nats/ — subscriber, in-memory node store │ ├── nats/ — subscriber, in-memory node store
│ ├── esl/ — FreeSWITCH ESL client (raw TCP)
│ ├── ami/ — Asterisk AMI client (raw TCP)
│ ├── config/ — manager: чтение/атомарная запись JSON │ ├── config/ — manager: чтение/атомарная запись JSON
│ ├── models/ — все общие типы │ ├── models/ — все общие типы
│ └── log/ — ASCII-metrics логгер + ротация │ └── log/ — ASCII-metrics логгер + ротация
├── contrib/ — route.lua, agent конфиги
├── data/ — JSON конфиги (gitignored) ├── data/ — JSON конфиги (gitignored)
├── web/ — SvelteKit + shadcn-svelte + Tailwind ├── web/ — SvelteKit + shadcn-svelte + Tailwind
│ ├── src/routes/ │ ├── src/routes/
@@ -224,7 +255,8 @@ pulse-lets-go/
│ │ ├── +page.svelte # дашборд │ │ ├── +page.svelte # дашборд
│ │ ├── login/+page.svelte │ │ ├── login/+page.svelte
│ │ ├── trunks/+page.svelte │ │ ├── trunks/+page.svelte
│ │ ── nodes/+page.svelte │ │ ── nodes/+page.svelte
│ │ └── users/+page.svelte
│ ├── src/lib/ │ ├── src/lib/
│ │ ├── components/ │ │ ├── components/
│ │ └── api.ts │ │ └── api.ts
+58 -16
View File
@@ -4,11 +4,12 @@
APP := pulse-lets-go APP := pulse-lets-go
BINDIR := bin BINDIR := bin
WEB_DIR := web 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 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 all: build
@@ -21,13 +22,17 @@ build-web:
@echo "→ Сборка SvelteKit..." @echo "→ Сборка SvelteKit..."
cd $(WEB_DIR) && npm run build cd $(WEB_DIR) && npm run build
# Сборка Go бинарника # Сборка Go бинарника (с авто-копированием статики)
build-go: build-go:
@mkdir -p $(BINDIR) @mkdir -p $(BINDIR)
@echo "→ Сборка Go backend..." @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 build -o $(BINDIR)/$(APP) ./cmd/$(APP)
# Быстрая сборка только Go (без фронта) # Быстрая сборка только Go (без фронта, без статики)
build-go-only: build-go-only:
@mkdir -p $(BINDIR) @mkdir -p $(BINDIR)
go build -o $(BINDIR)/$(APP) ./cmd/$(APP) go build -o $(BINDIR)/$(APP) ./cmd/$(APP)
@@ -42,23 +47,26 @@ run:
build-prod: build-prod:
@mkdir -p $(BINDIR) @mkdir -p $(BINDIR)
cd $(WEB_DIR) && npm run build 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) CGO_ENABLED=0 go build -ldflags="-s -w" -o $(BINDIR)/$(APP) ./cmd/$(APP)
# Установка: копирует бинарник и статику в /usr/local # Деплой в PREFIX: бинарник + статика + data + log
install: build deploy: build-prod
@echo "→ Установка в /usr/local" @echo "→ Деплой в $(PREFIX)"
mkdir -p /usr/local/share/$(APP)/web install -d $(PREFIX)/bin $(PREFIX)/data $(PREFIX)/web $(PREFIX)/log
cp -r $(WEB_DIR)/build /usr/local/share/$(APP)/web/ install -m 755 $(BINDIR)/$(APP) $(PREFIX)/bin/
cp $(BINDIR)/$(APP) /usr/local/bin/$(APP) cp -r $(WEB_DIR)/build/. $(PREFIX)/web/
chmod +x /usr/local/bin/$(APP) @echo "✓ Готово: $(PREFIX)"
@echo "✓ Установлен в /usr/local/bin/$(APP)" @echo " $(PREFIX)/bin/$(APP)"
@echo " $(PREFIX)/web/ (SvelteKit статика)"
@echo " $(PREFIX)/data/ (config.json, trunks.json, users.json)"
@echo " $(PREFIX)/log/ (ASCII-лог метрик)"
# Установка systemd unit # Установка systemd unit
install-systemd: install-systemd:
@echo "→ Установка systemd unit" install -m 644 deploy/$(APP).service /etc/systemd/system/
cp deploy/$(APP).service /etc/systemd/system/
systemctl daemon-reload 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 emulate-chaos: nats build-emulator
./$(BINDIR)/emulator --scenario chaos ./$(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 и запуск # 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"` IdleCPU float64 `json:"idle_cpu"`
LoadAvg float64 `json:"load_avg"` LoadAvg float64 `json:"load_avg"`
CallFailureRate float64 `json:"call_failure_rate"` CallFailureRate float64 `json:"call_failure_rate"`
SIPGateway string `json:"sip_gateway,omitempty"`
} }
// --- Спецификация значения (фиксированное или случайный диапазон) --- // --- Спецификация значения (фиксированное или случайный диапазон) ---
@@ -70,6 +71,7 @@ type nodeSpec struct {
idleCPU valueSpec idleCPU valueSpec
loadAvg valueSpec loadAvg valueSpec
failRate valueSpec failRate valueSpec
sipGateway string // SIP-адрес для авто-создания транка
} }
func (ns *nodeSpec) generate() nodeMetric { func (ns *nodeSpec) generate() nodeMetric {
@@ -82,6 +84,7 @@ func (ns *nodeSpec) generate() nodeMetric {
IdleCPU: ns.idleCPU.get(), IdleCPU: ns.idleCPU.get(),
LoadAvg: ns.loadAvg.get(), LoadAvg: ns.loadAvg.get(),
CallFailureRate: ns.failRate.get(), CallFailureRate: ns.failRate.get(),
SIPGateway: ns.sipGateway,
} }
} }
@@ -100,16 +103,19 @@ func scenarioNormal() []nodeSpec {
nodeID: "pbx-01", status: fixed(1), nodeID: "pbx-01", status: fixed(1),
calls: rrange(10, 100), maxCalls: 250, calls: rrange(10, 100), maxCalls: 250,
idleCPU: rrange(40, 95), loadAvg: rrange(0.1, 1.0), failRate: rrange(0, 2), 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), nodeID: "pbx-02", status: fixed(1),
calls: rrange(50, 200), maxCalls: 250, calls: rrange(50, 200), maxCalls: 250,
idleCPU: rrange(15, 45), loadAvg: rrange(0.5, 2.0), failRate: rrange(1, 8), 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), nodeID: "pbx-03", status: fixed(1),
calls: rrange(5, 60), maxCalls: 250, calls: rrange(5, 60), maxCalls: 250,
idleCPU: rrange(60, 90), loadAvg: rrange(0.05, 0.5), failRate: rrange(0, 1), 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) log.Printf("[emulator] ошибка публикации %s: %v", m.NodeID, err)
return return
} }
log.Printf("[emulator] → %s: calls=%d/%d idle=%.0f%% load=%.2f fail=%.1f%% status=%s", 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.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 ( import (
"context" "context"
"crypto/rand"
"encoding/hex"
"fmt" "fmt"
"log" "log"
"net/http" "net/http"
@@ -24,15 +22,17 @@ import (
filelog "github.com/pulse-lets-go/internal/log" filelog "github.com/pulse-lets-go/internal/log"
"github.com/pulse-lets-go/internal/models" "github.com/pulse-lets-go/internal/models"
"github.com/pulse-lets-go/internal/nats" "github.com/pulse-lets-go/internal/nats"
"github.com/pulse-lets-go/internal/util"
) )
func main() { func main() {
// Определяем рабочую директорию // Определяем рабочую директорию
// По умолчанию data/ — родительская папка от бинарника.
// Пример: /opt/pulse-lets-go/bin/pulse-lets-go → /opt/pulse-lets-go/data
exe, _ := os.Executable() exe, _ := os.Executable()
exeDir := filepath.Dir(exe) exeDir := filepath.Dir(exe)
dataDir := filepath.Join(exeDir, "data") dataDir := filepath.Join(filepath.Dir(exeDir), "data")
// Разрешаем переопределение через переменную окружения
if envDir := os.Getenv("PULSE_DATA_DIR"); envDir != "" { if envDir := os.Getenv("PULSE_DATA_DIR"); envDir != "" {
dataDir = envDir dataDir = envDir
} }
@@ -91,7 +91,7 @@ func main() {
} }
// Создаём новый balance-транк // Создаём новый balance-транк
now := time.Now().UTC() now := time.Now().UTC()
trunkID := "trk-" + randomHex(6) trunkID := "trk-" + util.RandomHex(6)
trunk := models.Trunk{ trunk := models.Trunk{
ID: trunkID, ID: trunkID,
Name: fmt.Sprintf("%s баланс", nodeID), Name: fmt.Sprintf("%s баланс", nodeID),
@@ -113,7 +113,7 @@ func main() {
}) })
// 6. HTTP API // 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() handler := apiHandler.Handler()
// 6.1 ESL-клиент (опционально — если сконфигурирован) // 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 статики, если директория существует // Подключаем раздачу SvelteKit статики, если директория существует
webDir := findWebDir(exeDir) webDir := findWebDir(exeDir)
if webDir != "" { if webDir != "" {
@@ -230,20 +243,19 @@ func main() {
log.Printf("[main] ошибка остановки HTTP сервера: %v", err) log.Printf("[main] ошибка остановки HTTP сервера: %v", err)
} }
sub.Close()
metricsLogger.Close() metricsLogger.Close()
log.Println("[main] pulse-lets-go остановлен") log.Println("[main] pulse-lets-go остановлен")
} }
// findWebDir ищет директорию со SvelteKit сборкой (web/build/). // findWebDir ищет директорию со SvelteKit сборкой (web/build/).
// Пробует по порядку: переменная окружения, рядом с бинарником, в корне проекта, CWD. // Пробует по порядку: PULSE_WEB_DIR, CWD, ../web/build, рядом с бинарником.
func findWebDir(exeDir string) string { func findWebDir(exeDir string) string {
candidates := []string{ candidates := []string{
os.Getenv("PULSE_WEB_DIR"), os.Getenv("PULSE_WEB_DIR"),
filepath.Join(exeDir, "web", "build"), // bin/web/build "web/build", // ./web/build (CWD — dev, systemd working dir)
filepath.Join(filepath.Dir(exeDir), "web", "build"), // ../web/build (корень проекта) filepath.Join(filepath.Dir(exeDir), "web", "build"), // ../web/build (корень проекта, dev с bin/)
"web/build", // ./web/build (CWD) filepath.Join(exeDir, "web", "build"), // рядом с бинарником (production install, make deploy)
} }
for _, d := range candidates { for _, d := range candidates {
if d == "" { 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. // extractHostFromGateway извлекает хост из SIP URI.
// "sip:mts-gw.lan:5060" → "mts-gw.lan" // "sip:mts-gw.lan:5060" → "mts-gw.lan"
func extractHostFromGateway(gateway string) string { func extractHostFromGateway(gateway string) string {
+7 -11
View File
@@ -1,16 +1,16 @@
// SIP тестер для e2e проверки pulse-lets-go + FreeSWITCH. // SIP тестер для e2e проверки pulse-lets-go + FreeSWITCH.
// Два режима: // Два режима:
//
// uas — отвечает 200 OK на SIP INVITE, симулирует PBX-ноду // uas — отвечает 200 OK на SIP INVITE, симулирует PBX-ноду
// uac — шлёт SIP INVITE в FreeSWITCH, симулирует оператора // uac — шлёт SIP INVITE в FreeSWITCH, симулирует оператора
// //
// Примеры: // Примеры:
//
// ./siptest -mode uas (PBX-нода на порту 5090) // ./siptest -mode uas (PBX-нода на порту 5090)
// ./siptest -mode uac -r 10 -l 100 (оператор, 10 CPS, до 100 одновременных) // ./siptest -mode uac -r 10 -l 100 (оператор, 10 CPS, до 100 одновременных)
package main package main
import ( import (
"crypto/rand"
"encoding/hex"
"flag" "flag"
"fmt" "fmt"
"log" "log"
@@ -19,6 +19,8 @@ import (
"sync" "sync"
"sync/atomic" "sync/atomic"
"time" "time"
"github.com/pulse-lets-go/internal/util"
) )
const ( const (
@@ -103,7 +105,7 @@ func buildSIPResponse(invite, localAddr string) string {
// Добавляем tag к To если его нет // Добавляем tag к To если его нет
to = strings.TrimSpace(to) to = strings.TrimSpace(to)
if !strings.Contains(to, ";tag=") { 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"+ 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 { 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") callID := fmt.Sprintf("call-%d-%s@%s", callNum, util.RandomHex(4), "127.0.0.1")
branch := "z9hG4bK-" + randomHex(8) branch := "z9hG4bK-" + util.RandomHex(8)
return fmt.Sprintf("INVITE sip:%s@%s SIP/2.0\r\n"+ return fmt.Sprintf("INVITE sip:%s@%s SIP/2.0\r\n"+
"Via: SIP/2.0/UDP %s;branch=%s\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", "\r\n",
dest, fsAddr, srcAddr, branch, caller, callNum, dest, callID, caller, srcAddr) 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] [Unit]
Description=pulse-lets-go — Telephone Load Balancer Description=pulse-lets-go — Telephone Load Balancer
Documentation=https://github.com/pulse-lets-go Documentation=https://github.com/pulse-lets-go
After=network.target nats.service After=network.target
Wants=nats.service
[Service] [Service]
Type=simple Type=simple
User=pulse User=pulse
Group=pulse Group=pulse
ExecStart=/usr/local/bin/pulse-lets-go ExecStart=/opt/pulse-lets-go/bin/pulse-lets-go
WorkingDirectory=/usr/local/share/pulse-lets-go WorkingDirectory=/opt/pulse-lets-go
Environment=PULSE_DATA_DIR=/etc/pulse-lets-go
Restart=always Restart=always
RestartSec=5 RestartSec=5
@@ -19,7 +17,7 @@ NoNewPrivileges=yes
PrivateTmp=yes PrivateTmp=yes
ProtectSystem=strict ProtectSystem=strict
ProtectHome=yes 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 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. // handleLogin обрабатывает POST /api/auth/login.
func (a *API) handleLogin(w http.ResponseWriter, r *http.Request) { func (a *API) handleLogin(w http.ResponseWriter, r *http.Request) {
var req models.LoginRequest 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") writeError(w, http.StatusBadRequest, "некорректный JSON")
return return
} }
@@ -199,10 +199,6 @@ func (a *API) handleLogin(w http.ResponseWriter, r *http.Request) {
// handleRefresh выдаёт новый access токен по refresh токену. // handleRefresh выдаёт новый access токен по refresh токену.
func (a *API) handleRefresh(w http.ResponseWriter, r *http.Request) { func (a *API) handleRefresh(w http.ResponseWriter, r *http.Request) {
tokenStr := extractBearerToken(r) tokenStr := extractBearerToken(r)
if tokenStr == "" {
// пробуем получить из query param (для WebSocket)
tokenStr = r.URL.Query().Get("token")
}
if tokenStr == "" { if tokenStr == "" {
writeError(w, http.StatusUnauthorized, "требуется токен") writeError(w, http.StatusUnauthorized, "требуется токен")
return return
+11
View File
@@ -1,10 +1,13 @@
package api package api
import ( import (
"bufio"
"crypto/rand" "crypto/rand"
"encoding/hex" "encoding/hex"
"encoding/json" "encoding/json"
"fmt"
"log" "log"
"net"
"net/http" "net/http"
"time" "time"
) )
@@ -82,3 +85,11 @@ func (w *loggingResponseWriter) WriteHeader(code int) {
func (w *loggingResponseWriter) Write(b []byte) (int, error) { func (w *loggingResponseWriter) Write(b []byte) (int, error) {
return w.ResponseWriter.Write(b) 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 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 { for _, n := range nodes {
id := sanitizePromLabel(n.NodeID) id := sanitizePromLabel(n.NodeID)
+2
View File
@@ -14,6 +14,7 @@ func (a *API) handleRoute(w http.ResponseWriter, r *http.Request) {
// Нет зарегистрированных нод // Нет зарегистрированных нод
if nodeID == "" { if nodeID == "" {
a.engine.IncrementRouteFallbacks()
writeJSON(w, http.StatusServiceUnavailable, models.RouteResponse{ writeJSON(w, http.StatusServiceUnavailable, models.RouteResponse{
Error: "no_nodes_registered", Error: "no_nodes_registered",
}) })
@@ -22,6 +23,7 @@ func (a *API) handleRoute(w http.ResponseWriter, r *http.Request) {
// Все ноды unhealthy — fallback // Все ноды unhealthy — fallback
if fallback || score < 0 { if fallback || score < 0 {
a.engine.IncrementRouteFallbacks()
fallbackGW, _ := a.findFallbackGateway() fallbackGW, _ := a.findFallbackGateway()
nodes := a.getRouteNodeInfo() nodes := a.getRouteNodeInfo()
+38 -1
View File
@@ -4,6 +4,7 @@ import (
"encoding/json" "encoding/json"
"fmt" "fmt"
"net/http" "net/http"
"sync"
"github.com/pulse-lets-go/internal/config" "github.com/pulse-lets-go/internal/config"
"github.com/pulse-lets-go/internal/engine" "github.com/pulse-lets-go/internal/engine"
@@ -23,6 +24,11 @@ type API struct {
logFormat string logFormat string
routeLimiter *rateLimiter routeLimiter *rateLimiter
apiLimiter *rateLimiter apiLimiter *rateLimiter
gatewayMu sync.RWMutex
gatewayStates map[string]string // trunkID → "up"/"down"
fsStats esl.FsStats
fsStatsMu sync.RWMutex
} }
// NewAPI создаёт новый HTTP API с заданными зависимостями. // NewAPI создаёт новый HTTP API с заданными зависимостями.
@@ -35,6 +41,7 @@ func NewAPI(eng *engine.Engine, cfgMgr *config.Manager, jwtSecret, monitoringAPI
wsHub: newWSHub(), wsHub: newWSHub(),
natsConnected: natsFn, natsConnected: natsFn,
logFormat: cfg.Log.Format, logFormat: cfg.Log.Format,
gatewayStates: make(map[string]string),
} }
if cfg.RateLimit.Enabled { if cfg.RateLimit.Enabled {
a.routeLimiter = newRateLimiter(cfg.RateLimit.RoutePerSec, cfg.RateLimit.RoutePerSec) a.routeLimiter = newRateLimiter(cfg.RateLimit.RoutePerSec, cfg.RateLimit.RoutePerSec)
@@ -130,6 +137,10 @@ func corsMiddleware(next http.Handler) http.Handler {
// BroadcastGatewayEvent транслирует статус gateway всем WS-клиентам. // BroadcastGatewayEvent транслирует статус gateway всем WS-клиентам.
func (a *API) BroadcastGatewayEvent(gatewayName, trunkID, status string) { func (a *API) BroadcastGatewayEvent(gatewayName, trunkID, status string) {
a.gatewayMu.Lock()
a.gatewayStates[trunkID] = status
a.gatewayMu.Unlock()
msg := &models.WsMetricsMessage{ msg := &models.WsMetricsMessage{
Type: "gateway_update", Type: "gateway_update",
Payload: map[string]string{ Payload: map[string]string{
@@ -141,9 +152,35 @@ func (a *API) BroadcastGatewayEvent(gatewayName, trunkID, status string) {
a.wsHub.broadcast(msg) 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 { func decodeJSON(r *http.Request, v interface{}) error {
defer r.Body.Close() defer r.Body.Close()
r.Body = http.MaxBytesReader(nil, r.Body, maxBodySize)
if err := json.NewDecoder(r.Body).Decode(v); err != nil { if err := json.NewDecoder(r.Body).Decode(v); err != nil {
return fmt.Errorf("декодирование JSON: %w", err) return fmt.Errorf("декодирование JSON: %w", err)
} }
+3 -2
View File
@@ -268,8 +268,9 @@ func (a *API) eslPushGatewayDelete(trunk models.Trunk) {
// gatewayParams возвращает параметры для создания FS gateway на основе транка. // gatewayParams возвращает параметры для создания FS gateway на основе транка.
func (a *API) gatewayParams(trunk models.Trunk) (profile, name, proxy string) { func (a *API) gatewayParams(trunk models.Trunk) (profile, name, proxy string) {
profile = a.readConfig().ESL.SofiaProfile cfg := a.readConfig()
name = fmt.Sprintf("pulse-ingress-%s", trunk.ID) profile = cfg.ESL.SofiaProfile
name = fmt.Sprintf("%s-ingress-%s", cfg.ESL.GatewayPrefix, trunk.ID)
proxy = extractHost(trunk.Gateway) proxy = extractHost(trunk.Gateway)
return return
} }
+2 -13
View File
@@ -1,13 +1,12 @@
package api package api
import ( import (
"crypto/rand"
"encoding/hex"
"fmt" "fmt"
"net/http" "net/http"
"time" "time"
"github.com/pulse-lets-go/internal/models" "github.com/pulse-lets-go/internal/models"
"github.com/pulse-lets-go/internal/util"
"golang.org/x/crypto/bcrypt" "golang.org/x/crypto/bcrypt"
) )
@@ -227,15 +226,5 @@ func findUserIndex(users []models.User, id string) int {
// generateID генерирует ID вида "prefix-XXXXX". // generateID генерирует ID вида "prefix-XXXXX".
func generateID(prefix string) string { func generateID(prefix string) string {
return fmt.Sprintf("%s-%s", prefix, randomHex(6)) return fmt.Sprintf("%s-%s", prefix, util.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]
} }
+42 -7
View File
@@ -104,12 +104,7 @@ func (a *API) handleWS(w http.ResponseWriter, r *http.Request) {
// Отправляем текущее состояние при подключении // Отправляем текущее состояние при подключении
go func() { go func() {
nodes := a.engine.GetAllNodes() a.BroadcastAll()
msg := &models.WsMetricsMessage{
Type: "nodes_update",
Payload: nodes,
}
a.wsHub.broadcast(msg)
}() }()
// Читаем из вебсокета (ping/pong + закрытие) // Читаем из вебсокета (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() { func (a *API) BroadcastMetrics() {
nodes := a.engine.GetAllNodes() nodes := a.engine.GetAllNodes()
msg := &models.WsMetricsMessage{ msg := &models.WsMetricsMessage{
@@ -135,6 +130,46 @@ func (a *API) BroadcastMetrics() {
a.wsHub.broadcast(msg) 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. // validateToken проверяет JWT и возвращает UserInfo.
func (a *API) validateToken(tokenStr string) (*UserInfo, error) { func (a *API) validateToken(tokenStr string) (*UserInfo, error) {
token, err := parseJWT(tokenStr, a.jwtSecret) token, err := parseJWT(tokenStr, a.jwtSecret)
+11
View File
@@ -72,6 +72,16 @@ func (e *ESLConfig) GetPassword() string {
return e.Password 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). // Config — корневая конфигурация приложения (config.json).
type Config struct { type Config struct {
NatsURL string `json:"nats_url"` NatsURL string `json:"nats_url"`
@@ -79,6 +89,7 @@ type Config struct {
NatsPassword string `json:"nats_password"` NatsPassword string `json:"nats_password"`
ListenAddr string `json:"listen_addr"` ListenAddr string `json:"listen_addr"`
JWTSecret string `json:"jwt_secret"` JWTSecret string `json:"jwt_secret"`
JWTSecretEnv string `json:"jwt_secret_env"` // имя переменной окружения с JWT-секретом
MonitoringAPIKey string `json:"monitoring_api_key"` MonitoringAPIKey string `json:"monitoring_api_key"`
StaleThresholdSec int `json:"stale_threshold_sec"` StaleThresholdSec int `json:"stale_threshold_sec"`
Scoring ScoringConfig `json:"scoring"` Scoring ScoringConfig `json:"scoring"`
+8
View File
@@ -29,6 +29,7 @@ type Engine struct {
// Статистика для health endpoint // Статистика для health endpoint
startTime time.Time startTime time.Time
routeRequests atomic.Int64 routeRequests atomic.Int64
routeFallbacks atomic.Int64
lastMetricTime atomic.Int64 // unix ts последней полученной метрики lastMetricTime atomic.Int64 // unix ts последней полученной метрики
} }
@@ -188,11 +189,17 @@ func (e *Engine) IncrementRouteRequests() {
e.routeRequests.Add(1) e.routeRequests.Add(1)
} }
// IncrementRouteFallbacks увеличивает счётчик fallback-запросов.
func (e *Engine) IncrementRouteFallbacks() {
e.routeFallbacks.Add(1)
}
// HealthStats возвращает статистику для health endpoint. // HealthStats возвращает статистику для health endpoint.
type HealthStats struct { type HealthStats struct {
TotalNodes int `json:"total_nodes"` TotalNodes int `json:"total_nodes"`
HealthyNodes int `json:"healthy_nodes"` HealthyNodes int `json:"healthy_nodes"`
RouteRequests int64 `json:"route_requests_total"` RouteRequests int64 `json:"route_requests_total"`
RouteFallbacks int64 `json:"route_fallbacks_total"`
UptimeSeconds int64 `json:"uptime_seconds"` UptimeSeconds int64 `json:"uptime_seconds"`
LastMetricTS int64 `json:"last_metric_ts"` LastMetricTS int64 `json:"last_metric_ts"`
} }
@@ -212,6 +219,7 @@ func (e *Engine) GetHealthStats() HealthStats {
TotalNodes: len(e.nodes), TotalNodes: len(e.nodes),
HealthyNodes: healthy, HealthyNodes: healthy,
RouteRequests: e.routeRequests.Load(), RouteRequests: e.routeRequests.Load(),
RouteFallbacks: e.routeFallbacks.Load(),
UptimeSeconds: int64(time.Since(e.startTime).Seconds()), UptimeSeconds: int64(time.Since(e.startTime).Seconds()),
LastMetricTS: e.lastMetricTime.Load(), LastMetricTS: e.lastMetricTime.Load(),
} }
+75 -47
View File
@@ -109,7 +109,7 @@ func (c *Client) tryConnect() error {
c.closeCh = make(chan struct{}) 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} dialer := net.Dialer{Timeout: 5 * time.Second}
conn, err := dialer.Dial("tcp", addr) conn, err := dialer.Dial("tcp", addr)
if err != nil { if err != nil {
@@ -119,7 +119,7 @@ func (c *Client) tryConnect() error {
c.reader = bufio.NewReader(conn) c.reader = bufio.NewReader(conn)
// 1. Читаем auth/request // 1. Читаем auth/request
headers, _, err := c.readMessageLocked() headers, _, err := c.readMessage()
if err != nil { if err != nil {
conn.Close() conn.Close()
return fmt.Errorf("чтение auth/request: %w", err) return fmt.Errorf("чтение auth/request: %w", err)
@@ -137,7 +137,7 @@ func (c *Client) tryConnect() error {
} }
// 3. Читаем auth response // 3. Читаем auth response
headers, _, err = c.readMessageLocked() headers, _, err = c.readMessage()
if err != nil { if err != nil {
conn.Close() conn.Close()
return fmt.Errorf("чтение auth response: %w", err) 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. // Subscribe подписывается на события ESL.
@@ -263,7 +263,7 @@ func (c *Client) readEventsLoop() {
continue continue
} }
headers, body, err := c.readMessageUnlocked() headers, body, err := c.readMessage()
if err != nil { if err != nil {
if err == io.EOF || strings.Contains(err.Error(), "use of closed network connection") { if err == io.EOF || strings.Contains(err.Error(), "use of closed network connection") {
@@ -330,48 +330,9 @@ func (c *Client) reconnect() {
// --- Приватные методы --- // --- Приватные методы ---
// readMessageUnlocked читает одно ESL-сообщение: заголовки + тело. // readMessage читает одно ESL-сообщение: заголовки + тело.
// Должна вызываться только из readEventsLoop (одиночный читатель). // Вызывающий должен гарантировать монопольный доступ к c.reader.
func (c *Client) readMessageUnlocked() (map[string]string, string, error) { func (c *Client) readMessage() (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) {
headers := make(map[string]string) headers := make(map[string]string)
for { for {
@@ -440,3 +401,70 @@ func (c *Client) GetStats() Stats {
func (c *Client) IncGatewayOps() { func (c *Client) IncGatewayOps() {
c.gatewayOps.Add(1) 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) // - Вызывает onConnect (например, для GatewaySyncAll)
// - Подписывается на события SOFIA::gateway_register/unregister // - Подписывается на события SOFIA::gateway_register/unregister
// - Запускает чтение асинхронных событий // - Запускает чтение асинхронных событий
//
// При разрыве соединения — автоматический реконнект. // При разрыве соединения — автоматический реконнект.
func (c *Client) StartEventLoop(ctx context.Context, profile string, onConnect func(), onEvent EventHandler) { func (c *Client) StartEventLoop(ctx context.Context, profile string, onConnect func(), onEvent EventHandler) {
const subscribeEvents = "SOFIA::gateway_register SOFIA::gateway_unregister SOFIA::gateway_expire" 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 return nil
} }
// GatewayList возвращает список имён gateway в заданном профиле. // GatewayList возвращает список имён gateway в заданном профиле с указанным префиксом.
func (c *Client) GatewayList(profile string) ([]string, error) { func (c *Client) GatewayList(profile, prefix string) ([]string, error) {
cmd := fmt.Sprintf("api sofia profile %s gwlist", profile) cmd := fmt.Sprintf("api sofia profile %s gwlist", profile)
_, body, err := c.Send(cmd) _, 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, "====") { if line == "" || strings.HasPrefix(line, "Name:") || strings.HasPrefix(line, "====") {
continue continue
} }
// Ищем строки вида "pulse-ingress-xxxxx sip:..." // Ищем строки вида "<prefix>-ingress-xxxxx sip:..."
if strings.Contains(line, "pulse-") { if strings.Contains(line, prefix+"-ingress-") {
parts := strings.Fields(line) parts := strings.Fields(line)
if len(parts) > 0 { if len(parts) > 0 {
gateways = append(gateways, parts[0]) gateways = append(gateways, parts[0])
+20 -1
View File
@@ -209,10 +209,29 @@ type ZabbixResponse struct {
Data []ZabbixNode `json:"data"` 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 --- // --- WebSocket ---
// WsMetricsMessage — сообщение, отправляемое по WebSocket. // WsMetricsMessage — сообщение, отправляемое по WebSocket.
type WsMetricsMessage struct { 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"` Payload interface{} `json:"payload"`
} }
+1 -1
View File
@@ -9,9 +9,9 @@ import (
natsgo "github.com/nats-io/nats.go" natsgo "github.com/nats-io/nats.go"
"github.com/pulse-lets-go/internal/models"
"github.com/pulse-lets-go/internal/engine" "github.com/pulse-lets-go/internal/engine"
filelog "github.com/pulse-lets-go/internal/log" filelog "github.com/pulse-lets-go/internal/log"
"github.com/pulse-lets-go/internal/models"
) )
const subject = "pulse.metrics.>" 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 --- // --- 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 { export function connectWebSocket(onMessage: (data: any) => void): WebSocket {
const token = getToken(); const token = getToken();
const proto = window.location.protocol === 'https:' ? 'wss' : 'ws'; const proto = window.location.protocol === 'https:' ? 'wss' : 'ws';
+24 -15
View File
@@ -1,31 +1,40 @@
<script lang="ts"> <script lang="ts">
import { browser } from '$app/environment';
import { page } from '$app/stores'; import { page } from '$app/stores';
import { goto } from '$app/navigation';
import '../app.css'; import '../app.css';
import { Toaster } from 'svelte-sonner'; 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 token: string | null = null;
let user: { username: string; role: string } | null = token let user: { username: string; role: string } | null = null;
? (() => { try { return JSON.parse(atob(token.split('.')[1])); } catch { return null; } })()
: null;
function logout() { function updateAuth() {
if (browser) { token = typeof localStorage !== 'undefined' ? localStorage.getItem('token') : null;
localStorage.removeItem('token'); if (token) {
localStorage.removeItem('refreshToken'); try { user = JSON.parse(atob(token.split('.')[1])); } catch { user = null; }
} } else {
token = null;
user = null; user = null;
} }
}
$: pathname = browser ? window.location.pathname : ''; function logout() {
$: isLogin = pathname === '/login'; localStorage.removeItem('token');
localStorage.removeItem('refreshToken');
updateAuth();
goto('/login');
}
let pathname = $derived($page.url.pathname);
$effect(() => {
$page.url.pathname;
updateAuth();
});
</script> </script>
<Toaster position="top-right" /> <Toaster position="top-right" />
{#if !isLogin} {#if pathname !== '/login'}
<div class="flex min-h-screen"> <div class="flex min-h-screen">
<!-- Sidebar --> <!-- Sidebar -->
<aside class="w-64 border-r bg-card flex flex-col"> <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 { onMount, onDestroy } from 'svelte';
import { import {
getNodes, connectWebSocket, isAuthenticated, getNodes, connectWebSocket, isAuthenticated,
type NodeInfo type NodeInfo, type BalancerInfo
} from '$lib/api'; } 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 nodes: NodeInfo[] = [];
let balancer: BalancerInfo | null = null;
let ws: WebSocket | null = null; let ws: WebSocket | null = null;
let error: string | null = null; let error: string | null = null;
@@ -21,6 +22,8 @@
ws = connectWebSocket((msg) => { ws = connectWebSocket((msg) => {
if (msg.type === 'nodes_update') { if (msg.type === 'nodes_update') {
nodes = msg.payload; 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; $: healthyCount = nodes.filter(n => n.score >= 0 && !n.disabled && !n.is_stale).length;
$: totalCalls = nodes.reduce((s, n) => s + n.active_calls, 0); $: 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; $: 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); $: 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> </script>
{#if error} {#if error}
@@ -79,6 +103,51 @@
</div> </div>
{/if} {/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 --> <!-- Summary Cards -->
<div class="grid grid-cols-1 sm:grid-cols-2 lg:grid-cols-4 gap-4 mb-6"> <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"> <div class="border rounded-lg p-4 bg-card">