Files
rogee b6d0af1a56
management-images / build-and-publish (push) Successful in 9m15s
docs: add one-click management deployment and image workflow
2026-09-16 17:57:05 +08:00

335 lines
10 KiB
Go

package management
import (
"context"
"database/sql"
"errors"
"net/http"
"strings"
"time"
"github.com/gofiber/fiber/v3"
)
func (s *Server) sipStatus(c fiber.Ctx) error {
p, authErr := s.authorize(c, "admin", "sip.status.read", "", "")
if authErr != nil {
return s.fail(c, authErr)
}
if err := s.checkFilterResources(c, p); err != nil {
return s.fail(c, err)
}
ctx, cancel := s.context(c)
defer cancel()
cells, err := s.store.ListCells(ctx)
if err != nil {
return s.fail(c, asAppError(err))
}
trunks, err := s.store.ListTrunks(ctx, s.cfg.PublicMode(), c.Query("provider_id"))
if err != nil {
return s.fail(c, asAppError(err))
}
visibleCells := cells[:0]
for _, cell := range cells {
if resourceAllowed(p.TokenSpec, "cell", cell.CellID) {
visibleCells = append(visibleCells, cell)
}
}
cells = visibleCells
visibleTrunks := trunks[:0]
for _, trunk := range trunks {
if resourceAllowed(p.TokenSpec, "trunk", trunk.TrunkID) && resourceAllowed(p.TokenSpec, "provider", trunk.ProviderID) {
visibleTrunks = append(visibleTrunks, trunk)
}
}
trunks = visibleTrunks
requested := parseCSV(c.Query("cell_ids"))
if len(requested) == 0 {
requested = parseCSV(c.Query("cell_id"))
}
items := make([]map[string]any, 0)
for _, cell := range cells {
if len(requested) > 0 && !containsString(requested, cell.CellID) {
continue
}
items = append(items, s.cellStatus(ctx, cell, "", trunks))
}
complete := true
for _, item := range items {
if value, _ := item["complete"].(bool); !value {
complete = false
}
}
coverage, dataAsOf := statusCoverage(items)
return c.JSON(map[string]any{"mode": s.cfg.PublicMode(), "complete": complete, "generated_at": utcString(time.Now()), "data_as_of": dataAsOf, "coverage": coverage, "cells": items})
}
func (s *Server) cellSIPStatus(c fiber.Ctx) error {
cellID := c.Params("cell_id")
p, authErr := s.authorize(c, "admin", "sip.status.read", "cell", cellID)
if authErr != nil {
return s.fail(c, authErr)
}
_ = p
ctx, cancel := s.context(c)
defer cancel()
cell, err := s.store.GetCell(ctx, cellID)
if errors.Is(err, sql.ErrNoRows) {
return s.fail(c, newAppError(404, "CELL_NOT_FOUND", "cell does not exist", nil))
}
if err != nil {
return s.fail(c, asAppError(err))
}
trunkID := c.Query("trunk_id")
if trunkID != "" && !resourceAllowed(p.TokenSpec, "trunk", trunkID) {
return s.fail(c, newAppError(http.StatusForbidden, "RESOURCE_FORBIDDEN", "requested trunk is outside the authorized scope", nil))
}
trunks, err := s.store.ListTrunks(ctx, s.cfg.PublicMode(), "")
if err != nil {
return s.fail(c, asAppError(err))
}
visibleTrunks := trunks[:0]
for _, trunk := range trunks {
if resourceAllowed(p.TokenSpec, "trunk", trunk.TrunkID) && resourceAllowed(p.TokenSpec, "provider", trunk.ProviderID) {
visibleTrunks = append(visibleTrunks, trunk)
}
}
return c.JSON(s.cellStatus(ctx, cell, trunkID, visibleTrunks))
}
func (s *Server) providerStatus(c fiber.Ctx) error {
providerID := c.Params("provider_id")
p, authErr := s.authorize(c, "admin", "sip.status.read", "provider", providerID)
if authErr != nil {
return s.fail(c, authErr)
}
ctx, cancel := s.context(c)
defer cancel()
provider, err := s.store.GetProvider(ctx, providerID)
if errors.Is(err, sql.ErrNoRows) {
return s.fail(c, newAppError(404, "PROVIDER_NOT_FOUND", "provider does not exist", nil))
}
if err != nil {
return s.fail(c, asAppError(err))
}
trunks, err := s.store.ListTrunks(ctx, s.cfg.PublicMode(), providerID)
if err != nil {
return s.fail(c, asAppError(err))
}
visibleTrunks := trunks[:0]
for _, trunk := range trunks {
if resourceAllowed(p.TokenSpec, "trunk", trunk.TrunkID) {
visibleTrunks = append(visibleTrunks, trunk)
}
}
trunks = visibleTrunks
cells, err := s.store.ListCells(ctx)
if err != nil {
return s.fail(c, asAppError(err))
}
visibleCells := cells[:0]
for _, cell := range cells {
if resourceAllowed(p.TokenSpec, "cell", cell.CellID) {
visibleCells = append(visibleCells, cell)
}
}
cells = visibleCells
cellViews := make([]map[string]any, 0, len(cells))
complete := true
for _, cell := range cells {
view := s.cellStatus(ctx, cell, "", trunks)
cellViews = append(cellViews, view)
if ok, _ := view["complete"].(bool); !ok {
complete = false
}
}
coverage, dataAsOf := statusCoverage(cellViews)
return c.JSON(map[string]any{"mode": s.cfg.PublicMode(), "complete": complete, "provider": provider, "trunks": trunks, "cells": cellViews, "generated_at": utcString(time.Now()), "data_as_of": dataAsOf, "coverage": coverage})
}
func (s *Server) cellStatus(ctx context.Context, cell Cell, trunkID string, trunks []AdminTrunk) map[string]any {
observations, err := s.store.LatestObservations(ctx, cell.CellID, trunkID)
view := map[string]any{"mode": s.cfg.PublicMode(), "cell_id": cell.CellID, "config_revision": cell.Revision, "config_status": cell.Config.Status, "egress_pool_id": cell.Config.EgressPoolID, "complete": false, "availability": "unknown", "eligibility": "unknown", "missing_sources": []string{"cell_agent", "asterisk", "ari", "registration", "media"}}
if cell.Config.Status == "disabled" {
view["availability"] = "disabled"
view["complete"] = true
view["missing_sources"] = []string{}
} else if err != nil {
view["error"] = "observation query failed"
return view
} else {
var selected *ObservationInput
if len(observations) > 0 {
selected = &observations[0]
}
if selected == nil {
view["reason"] = "no observation has been received for the current Cell boot"
} else {
age := time.Since(selected.ObservedAt).Seconds()
clockSkew := age < 0
view["clock_skew"] = clockSkew
view["observed_at"] = utcString(selected.ObservedAt)
view["received_at"] = utcString(selected.ReceivedAt)
if clockSkew {
view["observation_age_seconds"] = nil
} else {
view["observation_age_seconds"] = int64(age)
}
view["boot_id"] = selected.BootID
view["sequence"] = selected.Sequence
view["states"] = selected.States
view["occupancy"] = selected.Occupancy
if eligibility, ok := selected.States["eligibility"]; ok {
view["eligibility"] = eligibility
}
missing := missingSources(selected.States)
view["missing_sources"] = missing
if !clockSkew && age <= 15 && len(missing) == 0 {
view["availability"] = "healthy"
view["complete"] = true
} else if !clockSkew && age <= 30 {
view["availability"] = "stale"
view["reason"] = "observation is stale or incomplete"
} else if clockSkew {
view["availability"] = "unknown"
view["reason"] = "observation clock is ahead of the management server"
} else {
view["availability"] = "unknown"
view["reason"] = "observation is too old"
}
}
}
if trunkID != "" {
view["trunk_id"] = trunkID
view["publication"] = s.publicationState(ctx, cell.CellID, trunkID, trunks)
} else {
matrix := make([]map[string]any, 0, len(trunks))
for _, trunk := range trunks {
publication := s.publicationState(ctx, cell.CellID, trunk.TrunkID, trunks)
matrix = append(matrix, map[string]any{"provider_id": trunk.ProviderID, "trunk_id": trunk.TrunkID, "status": trunk.Status, "publication": publication})
}
view["trunks"] = matrix
}
return view
}
func statusCoverage(items []map[string]any) (map[string]any, *time.Time) {
covered := make([]string, 0, len(items))
missing := make([]string, 0)
var latest *time.Time
for _, item := range items {
cellID, _ := item["cell_id"].(string)
availability, _ := item["availability"].(string)
complete, _ := item["complete"].(bool)
if complete || availability == "disabled" {
covered = append(covered, cellID)
} else {
missing = append(missing, cellID)
}
if observed, ok := item["observed_at"].(string); ok {
if parsed, err := time.Parse(time.RFC3339Nano, observed); err == nil && (latest == nil || parsed.After(*latest)) {
value := parsed
latest = &value
}
}
}
return map[string]any{"covered_cell_ids": covered, "missing_cell_ids": missing}, latest
}
func missingSources(states map[string]any) []string {
if states == nil {
return []string{"cell_agent", "asterisk", "ari", "registration", "media"}
}
missing := make([]string, 0)
for _, key := range []string{"cell_agent", "asterisk", "ari", "registration", "media"} {
value, ok := states[key]
if !ok || value == nil || strings.EqualFold(toString(value), "unknown") || strings.EqualFold(toString(value), "missing") {
missing = append(missing, key)
}
}
return missing
}
func toString(v any) string {
switch value := v.(type) {
case string:
return value
case bool:
if value {
return "true"
}
return "false"
default:
return ""
}
}
func (s *Server) publicationState(ctx context.Context, cellID, trunkID string, trunks []AdminTrunk) map[string]any {
state := map[string]any{"state": "unknown", "cell_id": cellID, "trunk_id": trunkID}
var trunk *AdminTrunk
for i := range trunks {
if trunks[i].TrunkID == trunkID {
trunk = &trunks[i]
break
}
}
if trunk == nil {
return state
}
state["desired_revision"] = trunk.ActiveRevision
if trunk.ActiveRevision == 0 {
state["state"] = "not_published"
return state
}
pubs, err := s.store.GetPublicationRows(ctx, trunkID, &trunk.ActiveRevision)
if err != nil {
return state
}
for _, pub := range pubs {
if pub.CellID == cellID {
state["observed_revision"] = pub.LocalRevision
state["observed_digest"] = pub.LocalDigest
state["desired_digest"] = pub.TargetDigest
state["state"] = pub.Status
if pub.Status == "applied" && pub.LocalRevision == trunk.ActiveRevision && pub.LocalDigest == pub.TargetDigest {
state["state"] = "in_sync"
}
return state
}
}
state["state"] = "not_published"
return state
}
func (s *Server) ingestObservation(c fiber.Ctx) error {
cellID := c.Params("cell_id")
p, authErr := s.authorize(c, "admin", "sip.cell.write", "cell", cellID)
if authErr != nil {
return s.fail(c, authErr)
}
_ = p
var input ObservationInput
if err := decodeJSON(c, &input); err != nil {
return s.fail(c, err)
}
if input.CellID == "" {
input.CellID = cellID
}
if input.CellID != cellID {
return s.fail(c, newAppError(422, "CELL_ID_MISMATCH", "body cell_id does not match the path", nil))
}
input.ReceivedAt = time.Now()
if input.Source == "" {
input.Source = "mock"
}
if input.ObservationID == "" {
input.ObservationID = newID("obs")
}
ctx, cancel := s.context(c)
defer cancel()
if err := s.store.InsertObservation(ctx, input); err != nil {
return s.fail(c, asAppError(err))
}
return c.Status(http.StatusCreated).JSON(map[string]any{"mode": s.cfg.PublicMode(), "observation_id": input.ObservationID, "accepted": true})
}