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
This commit is contained in:
+18
-57
@@ -36,10 +36,10 @@ type Client struct {
|
||||
port int
|
||||
password string
|
||||
|
||||
conn net.Conn
|
||||
reader *bufio.Reader
|
||||
mu sync.Mutex // защищает отправку синхронных команд (auth, subscribe)
|
||||
sendMu sync.Mutex // защищает запись в сокет (async gateway commands)
|
||||
conn net.Conn
|
||||
reader *bufio.Reader
|
||||
mu sync.Mutex // защищает отправку синхронных команд (auth, subscribe)
|
||||
sendMu sync.Mutex // защищает запись в сокет (async gateway commands)
|
||||
|
||||
connected atomic.Bool
|
||||
eventsCh chan *Event
|
||||
@@ -109,7 +109,7 @@ func (c *Client) tryConnect() error {
|
||||
c.closeCh = make(chan struct{})
|
||||
}
|
||||
|
||||
addr := fmt.Sprintf("%s:%d", c.host, c.port)
|
||||
addr := net.JoinHostPort(c.host, fmt.Sprintf("%d", c.port))
|
||||
dialer := net.Dialer{Timeout: 5 * time.Second}
|
||||
conn, err := dialer.Dial("tcp", addr)
|
||||
if err != nil {
|
||||
@@ -119,7 +119,7 @@ func (c *Client) tryConnect() error {
|
||||
c.reader = bufio.NewReader(conn)
|
||||
|
||||
// 1. Читаем auth/request
|
||||
headers, _, err := c.readMessageLocked()
|
||||
headers, _, err := c.readMessage()
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return fmt.Errorf("чтение auth/request: %w", err)
|
||||
@@ -137,7 +137,7 @@ func (c *Client) tryConnect() error {
|
||||
}
|
||||
|
||||
// 3. Читаем auth response
|
||||
headers, _, err = c.readMessageLocked()
|
||||
headers, _, err = c.readMessage()
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return fmt.Errorf("чтение auth response: %w", err)
|
||||
@@ -199,7 +199,7 @@ func (c *Client) Send(command string) (map[string]string, string, error) {
|
||||
}
|
||||
|
||||
// Читаем ответ
|
||||
return c.readMessageLocked()
|
||||
return c.readMessage()
|
||||
}
|
||||
|
||||
// Subscribe подписывается на события ESL.
|
||||
@@ -263,7 +263,7 @@ func (c *Client) readEventsLoop() {
|
||||
continue
|
||||
}
|
||||
|
||||
headers, body, err := c.readMessageUnlocked()
|
||||
headers, body, err := c.readMessage()
|
||||
|
||||
if err != nil {
|
||||
if err == io.EOF || strings.Contains(err.Error(), "use of closed network connection") {
|
||||
@@ -330,48 +330,9 @@ func (c *Client) reconnect() {
|
||||
|
||||
// --- Приватные методы ---
|
||||
|
||||
// readMessageUnlocked читает одно ESL-сообщение: заголовки + тело.
|
||||
// Должна вызываться только из readEventsLoop (одиночный читатель).
|
||||
func (c *Client) readMessageUnlocked() (map[string]string, string, error) {
|
||||
headers := make(map[string]string)
|
||||
|
||||
for {
|
||||
line, err := c.reader.ReadString('\n')
|
||||
if err != nil {
|
||||
return nil, "", err
|
||||
}
|
||||
line = strings.TrimRight(line, "\r\n")
|
||||
if line == "" {
|
||||
break // конец заголовков
|
||||
}
|
||||
|
||||
parts := strings.SplitN(line, ":", 2)
|
||||
if len(parts) == 2 {
|
||||
headers[strings.TrimSpace(parts[0])] = strings.TrimSpace(parts[1])
|
||||
}
|
||||
}
|
||||
|
||||
// Читаем тело согласно Content-Length
|
||||
cl := headers["Content-Length"]
|
||||
if cl == "" {
|
||||
return headers, "", nil
|
||||
}
|
||||
length, err := strconv.Atoi(cl)
|
||||
if err != nil {
|
||||
return headers, "", fmt.Errorf("некорректный Content-Length: %s", cl)
|
||||
}
|
||||
|
||||
body := make([]byte, length)
|
||||
if _, err := io.ReadFull(c.reader, body); err != nil {
|
||||
return headers, "", fmt.Errorf("чтение тела (Content-Length=%d): %w", length, err)
|
||||
}
|
||||
|
||||
return headers, strings.TrimSpace(string(body)), nil
|
||||
}
|
||||
|
||||
// readMessageLocked читает одно ESL-сообщение: заголовки + тело.
|
||||
// Вызывающий держит c.mu (для синхронных команд Send).
|
||||
func (c *Client) readMessageLocked() (map[string]string, string, error) {
|
||||
// readMessage читает одно ESL-сообщение: заголовки + тело.
|
||||
// Вызывающий должен гарантировать монопольный доступ к c.reader.
|
||||
func (c *Client) readMessage() (map[string]string, string, error) {
|
||||
headers := make(map[string]string)
|
||||
|
||||
for {
|
||||
@@ -412,12 +373,12 @@ func (c *Client) readMessageLocked() (map[string]string, string, error) {
|
||||
|
||||
// Stats возвращает статистику ESL-клиента для health endpoint.
|
||||
type Stats struct {
|
||||
Status string `json:"status"`
|
||||
Host string `json:"host"`
|
||||
Reconnects int64 `json:"reconnects"`
|
||||
GatewayOps int64 `json:"gateway_ops_total"`
|
||||
EventsRecv int64 `json:"events_received"`
|
||||
UptimeSec int64 `json:"uptime_seconds"`
|
||||
Status string `json:"status"`
|
||||
Host string `json:"host"`
|
||||
Reconnects int64 `json:"reconnects"`
|
||||
GatewayOps int64 `json:"gateway_ops_total"`
|
||||
EventsRecv int64 `json:"events_received"`
|
||||
UptimeSec int64 `json:"uptime_seconds"`
|
||||
}
|
||||
|
||||
// GetStats возвращает текущую статистику.
|
||||
|
||||
@@ -16,6 +16,7 @@ type EventHandler func(eventName string, headers map[string]string, body string)
|
||||
// - Вызывает onConnect (например, для GatewaySyncAll)
|
||||
// - Подписывается на события SOFIA::gateway_register/unregister
|
||||
// - Запускает чтение асинхронных событий
|
||||
//
|
||||
// При разрыве соединения — автоматический реконнект.
|
||||
func (c *Client) StartEventLoop(ctx context.Context, profile string, onConnect func(), onEvent EventHandler) {
|
||||
const subscribeEvents = "SOFIA::gateway_register SOFIA::gateway_unregister SOFIA::gateway_expire"
|
||||
|
||||
@@ -32,8 +32,8 @@ func (c *Client) GatewayDelete(profile, name string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// GatewayList возвращает список имён gateway в заданном профиле.
|
||||
func (c *Client) GatewayList(profile string) ([]string, error) {
|
||||
// GatewayList возвращает список имён gateway в заданном профиле с указанным префиксом.
|
||||
func (c *Client) GatewayList(profile, prefix string) ([]string, error) {
|
||||
cmd := fmt.Sprintf("api sofia profile %s gwlist", profile)
|
||||
_, body, err := c.Send(cmd)
|
||||
|
||||
@@ -49,8 +49,8 @@ func (c *Client) GatewayList(profile string) ([]string, error) {
|
||||
if line == "" || strings.HasPrefix(line, "Name:") || strings.HasPrefix(line, "====") {
|
||||
continue
|
||||
}
|
||||
// Ищем строки вида "pulse-ingress-xxxxx sip:..."
|
||||
if strings.Contains(line, "pulse-") {
|
||||
// Ищем строки вида "<prefix>-ingress-xxxxx sip:..."
|
||||
if strings.Contains(line, prefix+"-ingress-") {
|
||||
parts := strings.Fields(line)
|
||||
if len(parts) > 0 {
|
||||
gateways = append(gateways, parts[0])
|
||||
|
||||
Reference in New Issue
Block a user