mirror of
https://github.com/MeshCore-Beacon/beacon-server.git
synced 2026-09-25 13:13:39 +00:00
refactor db package
layout and organization only
This commit is contained in:
+217
@@ -0,0 +1,217 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
sqlc "github.com/MeshCore-Tower/tower-server/db/sqlc"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/api"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/ingest"
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
)
|
||||
|
||||
func (s *Store) UpsertChannel(ctx context.Context, channelHash []byte, keyFingerprint []byte, name string, hashtag string) (int, error) {
|
||||
var namePtr, hashtagPtr *string
|
||||
if name != "" {
|
||||
namePtr = &name
|
||||
}
|
||||
if hashtag != "" {
|
||||
hashtagPtr = &hashtag
|
||||
}
|
||||
isHashtag := hashtag != ""
|
||||
row, err := s.q.UpsertChannel(ctx, sqlc.UpsertChannelParams{
|
||||
ChannelHash: channelHash,
|
||||
Column2: keyFingerprint, // key_fingerprint
|
||||
Name: namePtr,
|
||||
Hashtag: hashtagPtr,
|
||||
IsHashtag: &isHashtag,
|
||||
MessageCount: nil, // message count bumped separately by InsertChannelMessage
|
||||
})
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return int(row.ID), nil
|
||||
}
|
||||
|
||||
func (s *Store) UpsertChannelHashOnly(ctx context.Context, channelHash []byte) (int, error) {
|
||||
rowID, err := s.q.UpsertChannelHashOnly(ctx, channelHash)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return int(rowID), nil
|
||||
}
|
||||
|
||||
func (s *Store) ListChannels(ctx context.Context, limit int32, hash []byte, iata string, cursor int64) (api.Page[api.ChannelSummary], error) {
|
||||
var cursorTS pgtype.Timestamptz
|
||||
if cursor > 0 {
|
||||
cursorTS = pgtype.Timestamptz{Time: time.UnixMilli(cursor), Valid: true}
|
||||
}
|
||||
rows, err := s.q.ListChannels(ctx, sqlc.ListChannelsParams{
|
||||
Column1: hash,
|
||||
Column2: iata,
|
||||
Column3: cursorTS,
|
||||
Limit: limit + 1,
|
||||
})
|
||||
if err != nil {
|
||||
return api.Page[api.ChannelSummary]{}, err
|
||||
}
|
||||
hasMore := len(rows) > int(limit)
|
||||
if hasMore {
|
||||
rows = rows[:limit]
|
||||
}
|
||||
items := make([]api.ChannelSummary, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
items = append(items, api.ChannelSummary{
|
||||
ID: int(v.ID),
|
||||
Name: v.Name,
|
||||
ChannelHash: hex.EncodeToString(v.ChannelHash),
|
||||
LastSeen: v.LastSeen.Time.UnixMilli(),
|
||||
IsHashtag: v.IsHashtag != nil && *v.IsHashtag,
|
||||
KeyKnown: v.KeyKnown != nil && *v.KeyKnown,
|
||||
})
|
||||
}
|
||||
var nextCursor *int64
|
||||
if hasMore {
|
||||
last := items[len(items)-1].LastSeen
|
||||
nextCursor = &last
|
||||
}
|
||||
return api.Page[api.ChannelSummary]{
|
||||
Items: items,
|
||||
NextCursor: nextCursor,
|
||||
HasMore: hasMore,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetChannel(ctx context.Context, channelID int32) (*api.Channel, error) {
|
||||
row, err := s.q.GetChannelByID(ctx, channelID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
channel := api.Channel{
|
||||
ChannelSummary: api.ChannelSummary{
|
||||
ID: int(row.ID),
|
||||
Name: row.Name,
|
||||
ChannelHash: hex.EncodeToString(row.ChannelHash),
|
||||
LastSeen: row.LastSeen.Time.UnixMilli(),
|
||||
IsHashtag: row.IsHashtag != nil && *row.IsHashtag,
|
||||
KeyKnown: row.KeyKnown != nil && *row.KeyKnown,
|
||||
},
|
||||
Hashtag: row.Hashtag,
|
||||
MessageCount: 0,
|
||||
}
|
||||
if row.MessageCount != nil {
|
||||
channel.MessageCount = *row.MessageCount
|
||||
}
|
||||
if row.IsHashtag != nil && *row.IsHashtag && row.KeyFingerprint != nil {
|
||||
fp := hex.EncodeToString(row.KeyFingerprint)
|
||||
channel.KeyFingerprint = &fp
|
||||
}
|
||||
return &channel, nil
|
||||
}
|
||||
|
||||
func (s *Store) InsertChannelMessage(ctx context.Context, m ingest.InsertChannelMessageParams) (bool, error) {
|
||||
params := sqlc.InsertChannelMessageParams{ChannelID: int32(m.ChannelID), PacketHash: m.PacketHash, SenderName: &m.SenderName, Content: &m.Content, SentAt: pgtype.Timestamptz{Time: m.SentAt, Valid: true}}
|
||||
_, err := s.q.InsertChannelMessage(ctx, params)
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return false, nil // duplicate
|
||||
}
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
func (s *Store) ListChannelMessages(ctx context.Context, channelID *int32, since time.Time, limit int32, iatas []string, scope string, cursor int64) (api.Page[api.ChannelMessage], error) {
|
||||
ts := pgtype.Timestamptz{Time: since, Valid: !since.IsZero()}
|
||||
var messages []api.ChannelMessage
|
||||
var hasMore bool
|
||||
iataFilter := strings.Join(iatas, ",")
|
||||
if channelID == nil {
|
||||
rows, err := s.q.ListAllChannelMessages(ctx, sqlc.ListAllChannelMessagesParams{
|
||||
Column1: ts,
|
||||
Column2: iataFilter,
|
||||
Column3: scope,
|
||||
Column4: cursor,
|
||||
Limit: limit + 1,
|
||||
})
|
||||
if err != nil {
|
||||
return api.Page[api.ChannelMessage]{}, err
|
||||
}
|
||||
hasMore = len(rows) > int(limit)
|
||||
if hasMore {
|
||||
rows = rows[:limit]
|
||||
}
|
||||
messages = make([]api.ChannelMessage, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
messages = append(messages, toChannelMessage(v.ID, v.PacketHashHex, v.ChannelHash, v.SenderName, v.Content, v.SentAt, v.ObservationCount))
|
||||
}
|
||||
} else {
|
||||
rows, err := s.q.ListChannelMessages(ctx, sqlc.ListChannelMessagesParams{
|
||||
ChannelID: *channelID,
|
||||
Column2: ts,
|
||||
Column3: iataFilter,
|
||||
Column4: scope,
|
||||
Column5: cursor,
|
||||
Limit: limit + 1,
|
||||
})
|
||||
if err != nil {
|
||||
return api.Page[api.ChannelMessage]{}, err
|
||||
}
|
||||
hasMore = len(rows) > int(limit)
|
||||
if hasMore {
|
||||
rows = rows[:limit]
|
||||
}
|
||||
messages = make([]api.ChannelMessage, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
messages = append(messages, toChannelMessage(v.ID, v.PacketHashHex, v.ChannelHash, v.SenderName, v.Content, v.SentAt, v.ObservationCount))
|
||||
}
|
||||
}
|
||||
|
||||
var nextCursor *int64
|
||||
if hasMore && len(messages) > 0 {
|
||||
last := messages[len(messages)-1].ID
|
||||
nextCursor = &last
|
||||
}
|
||||
return api.Page[api.ChannelMessage]{
|
||||
Items: messages,
|
||||
NextCursor: nextCursor,
|
||||
HasMore: hasMore,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *Store) ListChannelMessagesByHash(ctx context.Context, hash []byte, since time.Time, limit int32, iatas []string, scope string, cursor int64) (api.Page[api.ChannelMessage], error) {
|
||||
iataFilter := strings.Join(iatas, ",")
|
||||
rows, err := s.q.ListChannelMessagesByHash(ctx, sqlc.ListChannelMessagesByHashParams{
|
||||
ChannelHash: hash,
|
||||
Column2: pgtype.Timestamptz{Time: since, Valid: !since.IsZero()},
|
||||
Column3: iataFilter,
|
||||
Column4: scope,
|
||||
Column5: cursor,
|
||||
Limit: limit + 1,
|
||||
})
|
||||
if err != nil {
|
||||
return api.Page[api.ChannelMessage]{}, err
|
||||
}
|
||||
hasMore := len(rows) > int(limit)
|
||||
if hasMore {
|
||||
rows = rows[:limit]
|
||||
}
|
||||
messages := make([]api.ChannelMessage, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
messages = append(messages, toChannelMessage(v.ID, hex.EncodeToString(v.PacketHash), v.ChannelHash, v.SenderName, v.Content, v.SentAt, v.ObservationCount))
|
||||
}
|
||||
var nextCursor *int64
|
||||
if hasMore && len(messages) > 0 {
|
||||
last := messages[len(messages)-1].ID
|
||||
nextCursor = &last
|
||||
}
|
||||
return api.Page[api.ChannelMessage]{
|
||||
Items: messages,
|
||||
NextCursor: nextCursor,
|
||||
HasMore: hasMore,
|
||||
}, nil
|
||||
}
|
||||
+146
@@ -0,0 +1,146 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
sqlc "github.com/MeshCore-Tower/tower-server/db/sqlc"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/api"
|
||||
)
|
||||
|
||||
func (s *Store) UpsertIATADetails(ctx context.Context, iata string, name string, lat, lng *float64) error {
|
||||
return s.q.UpsertIATADetails(ctx, sqlc.UpsertIATADetailsParams{
|
||||
Iata: iata,
|
||||
DisplayName: &name,
|
||||
ApproxLat: lat,
|
||||
ApproxLng: lng,
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Store) ListIATAs(ctx context.Context) ([]api.IATA, error) {
|
||||
rows, err := s.q.ListIATAs(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
iatas := make([]api.IATA, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
iatas = append(iatas, api.IATA{
|
||||
IATA: v.Iata,
|
||||
DisplayName: v.DisplayName,
|
||||
Lat: v.ApproxLat,
|
||||
Lng: v.ApproxLng,
|
||||
})
|
||||
}
|
||||
return iatas, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetIATA(ctx context.Context, iata string) (*api.IATA, error) {
|
||||
i, err := s.q.GetIATA(ctx, iata)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &api.IATA{
|
||||
IATA: i.Iata,
|
||||
DisplayName: i.DisplayName,
|
||||
Lat: i.ApproxLat,
|
||||
Lng: i.ApproxLng,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *Store) UpsertRegion(ctx context.Context, slug, name, description string, displayOrder int, centerLat, centerLng *float64, zoomLevel *int) (int32, error) {
|
||||
var zl *int32
|
||||
if zoomLevel != nil {
|
||||
z := int32(*zoomLevel)
|
||||
zl = &z
|
||||
}
|
||||
do := int32(displayOrder)
|
||||
return s.q.UpsertRegion(ctx, sqlc.UpsertRegionParams{
|
||||
Slug: slug,
|
||||
Name: name,
|
||||
Description: &description,
|
||||
DisplayOrder: &do,
|
||||
CenterLat: centerLat,
|
||||
CenterLng: centerLng,
|
||||
ZoomLevel: zl,
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Store) ListRegions(ctx context.Context) ([]api.RegionSummary, error) {
|
||||
rows, err := s.q.ListRegions(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
regions := make([]api.RegionSummary, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
regions = append(regions, api.RegionSummary{
|
||||
ID: int(v.ID),
|
||||
Slug: v.Slug,
|
||||
Name: v.Name,
|
||||
})
|
||||
}
|
||||
return regions, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetRegion(ctx context.Context, regionID int32) (*api.Region, error) {
|
||||
region, err := s.q.GetRegion(ctx, regionID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
result := api.Region{
|
||||
RegionSummary: api.RegionSummary{
|
||||
ID: int(region.ID),
|
||||
Slug: region.Slug,
|
||||
Name: region.Name,
|
||||
},
|
||||
Description: region.Description,
|
||||
CenterLat: region.CenterLat,
|
||||
CenterLng: region.CenterLng,
|
||||
}
|
||||
var zoomLevel *int
|
||||
if region.ZoomLevel != nil {
|
||||
z := int(*region.ZoomLevel)
|
||||
zoomLevel = &z
|
||||
}
|
||||
result.ZoomLevel = zoomLevel
|
||||
iatas, err := s.q.GetRegionIATAs(ctx, regionID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
result.IATAs = iatas
|
||||
return &result, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetRegionBySlug(ctx context.Context, slug string) (*api.Region, error) {
|
||||
region, err := s.q.GetRegionBySlug(ctx, slug)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
result := api.Region{
|
||||
RegionSummary: api.RegionSummary{
|
||||
ID: int(region.ID),
|
||||
Slug: region.Slug,
|
||||
Name: region.Name,
|
||||
},
|
||||
Description: region.Description,
|
||||
CenterLat: region.CenterLat,
|
||||
CenterLng: region.CenterLng,
|
||||
}
|
||||
var zoomLevel *int
|
||||
if region.ZoomLevel != nil {
|
||||
z := int(*region.ZoomLevel)
|
||||
zoomLevel = &z
|
||||
}
|
||||
result.ZoomLevel = zoomLevel
|
||||
iatas, err := s.q.GetRegionIATAs(ctx, region.ID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
result.IATAs = iatas
|
||||
return &result, nil
|
||||
}
|
||||
|
||||
func (s *Store) UpsertRegionIATA(ctx context.Context, regionID int32, iata string) error {
|
||||
return s.q.UpsertRegionIATA(ctx, sqlc.UpsertRegionIATAParams{
|
||||
RegionID: regionID,
|
||||
Iata: iata,
|
||||
})
|
||||
}
|
||||
+173
@@ -0,0 +1,173 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
sqlc "github.com/MeshCore-Tower/tower-server/db/sqlc"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/api"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/ingest"
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
)
|
||||
|
||||
func (s *Store) UpsertNode(ctx context.Context, n ingest.UpsertNodeParams, radio ingest.RadioSettings) (uuid.UUID, error) {
|
||||
params := sqlc.UpsertNodeParams{
|
||||
PublicKey: n.PublicKey,
|
||||
NodeType: int16(n.NodeType),
|
||||
Name: &n.Name,
|
||||
Latitude: n.Latitude,
|
||||
Longitude: n.Longitude,
|
||||
}
|
||||
if radio.FreqMHz != 0 {
|
||||
params.RadioFreqMhz = &radio.FreqMHz
|
||||
params.RadioSf = &radio.SF
|
||||
params.RadioBwKhz = &radio.BWKHz
|
||||
}
|
||||
row, err := s.q.UpsertNode(ctx, params)
|
||||
if err != nil {
|
||||
return uuid.Nil, err
|
||||
}
|
||||
return row.ID, nil
|
||||
}
|
||||
|
||||
func (s *Store) UpsertNodeIATA(ctx context.Context, nodeID uuid.UUID, iata string) error {
|
||||
params := sqlc.UpsertNodeIATAParams{NodeID: nodeID, Iata: iata}
|
||||
return s.q.UpsertNodeIATA(ctx, params)
|
||||
}
|
||||
|
||||
func (s *Store) UpsertNodeShortID(ctx context.Context, nodeID uuid.UUID, iata string, prefix4 []byte) error {
|
||||
return s.q.UpsertNodeShortID(ctx, sqlc.UpsertNodeShortIDParams{
|
||||
NodeID: nodeID,
|
||||
Iata: iata,
|
||||
Prefix4: prefix4,
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Store) SetNodeCapability(ctx context.Context, nodeID uuid.UUID, paths, traces bool) error {
|
||||
var errs []error
|
||||
if paths {
|
||||
errs = append(errs, s.q.SetNodeMultibytePaths(ctx, nodeID))
|
||||
}
|
||||
if traces {
|
||||
errs = append(errs, s.q.SetNodeMultibyteTraces(ctx, nodeID))
|
||||
}
|
||||
return errors.Join(errs...)
|
||||
}
|
||||
|
||||
func (s *Store) SetNodeDefaultScope(ctx context.Context, nodeID uuid.UUID, scopeID int32) error {
|
||||
return s.q.SetNodeDefaultScope(ctx, sqlc.SetNodeDefaultScopeParams{
|
||||
ID: nodeID,
|
||||
DefaultScopeID: &scopeID,
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Store) ListNodes(ctx context.Context, nodeType int16, iatas []string, supportsMultibytePaths, supportsMultibyteTraces *bool, pubkey []byte, name, scope string, cursor int64, limit int32) (api.Page[api.NodeSummary], error) {
|
||||
var cursorTS pgtype.Timestamptz
|
||||
if cursor > 0 {
|
||||
cursorTS = pgtype.Timestamptz{Time: time.UnixMilli(cursor), Valid: true}
|
||||
}
|
||||
iataFilter := strings.Join(iatas, ",")
|
||||
rows, err := s.q.ListNodes(ctx, sqlc.ListNodesParams{
|
||||
Column1: nodeType,
|
||||
Column2: iataFilter,
|
||||
Column3: tristate(supportsMultibytePaths),
|
||||
Column4: tristate(supportsMultibyteTraces),
|
||||
Column5: pubkey,
|
||||
Column6: name,
|
||||
Column7: cursorTS,
|
||||
Limit: limit + 1,
|
||||
Column9: scope,
|
||||
})
|
||||
if err != nil {
|
||||
return api.Page[api.NodeSummary]{}, err
|
||||
}
|
||||
hasMore := len(rows) > int(limit)
|
||||
if hasMore {
|
||||
rows = rows[:limit]
|
||||
}
|
||||
items := make([]api.NodeSummary, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
node := api.NodeSummary{
|
||||
ID: v.ID,
|
||||
PublicKey: hex.EncodeToString(v.PublicKey),
|
||||
NodeType: v.NodeType,
|
||||
NodeTypeName: api.NodeTypeName(v.NodeType),
|
||||
Name: v.Name,
|
||||
Latitude: v.Latitude,
|
||||
Longitude: v.Longitude,
|
||||
IsObserver: v.IsObserver,
|
||||
ObvserverID: nullableUUID(v.ObserverID),
|
||||
}
|
||||
if len(v.Iatas) > 0 {
|
||||
if err := json.Unmarshal(v.Iatas, &node.IATAs); err != nil {
|
||||
log.Printf("store: failed to unmarshal node iatas: %v", err)
|
||||
node.IATAs = []api.NodeIATA{}
|
||||
}
|
||||
}
|
||||
if v.RadioFreqMhz != nil && v.RadioSf != nil && v.RadioBwKhz != nil {
|
||||
s := fmt.Sprintf("%.1f,%g,%d", *v.RadioFreqMhz, *v.RadioBwKhz, *v.RadioSf)
|
||||
node.Radio = &s
|
||||
}
|
||||
items = append(items, node)
|
||||
}
|
||||
var nextCursor *int64
|
||||
if hasMore && len(items) > 0 {
|
||||
ms := rows[len(rows)-1].LastSeen.Time.UnixMilli()
|
||||
nextCursor = &ms
|
||||
}
|
||||
return api.Page[api.NodeSummary]{
|
||||
Items: items,
|
||||
NextCursor: nextCursor,
|
||||
HasMore: hasMore,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetNode(ctx context.Context, nodeID uuid.UUID) (*api.Node, error) {
|
||||
row, err := s.q.GetNodeByID(ctx, nodeID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
node := &api.Node{
|
||||
NodeSummary: api.NodeSummary{
|
||||
ID: row.ID,
|
||||
PublicKey: hex.EncodeToString(row.PublicKey),
|
||||
NodeType: row.NodeType,
|
||||
NodeTypeName: api.NodeTypeName(row.NodeType),
|
||||
Name: row.Name,
|
||||
Latitude: row.Latitude,
|
||||
Longitude: row.Longitude,
|
||||
IsObserver: row.IsObserver,
|
||||
ObvserverID: nullableUUID(row.ObserverID),
|
||||
DefaultScope: row.DefaultScopeName,
|
||||
},
|
||||
LocationSource: row.LocationSource,
|
||||
SupportsMultibytePaths: row.SupportsMultibytePaths,
|
||||
SupportsMultibyteTraces: row.SupportsMultibyteTraces,
|
||||
MinFirmwareVersion: row.MinFirmwareVersion,
|
||||
FirstSeen: row.FirstSeen.Time.UnixMilli(),
|
||||
LastSeen: row.LastSeen.Time.UnixMilli(),
|
||||
Metadata: row.Metadata,
|
||||
}
|
||||
if len(row.Iatas) > 0 {
|
||||
if err := json.Unmarshal(row.Iatas, &node.IATAs); err != nil {
|
||||
log.Printf("store: failed to unmarshal node iatas: %v", err)
|
||||
node.IATAs = []api.NodeIATA{}
|
||||
}
|
||||
}
|
||||
if row.RadioFreqMhz != nil && row.RadioSf != nil && row.RadioBwKhz != nil {
|
||||
s := fmt.Sprintf("%.1f,%g,%d", *row.RadioFreqMhz, *row.RadioBwKhz, *row.RadioSf)
|
||||
node.Radio = &s
|
||||
}
|
||||
if row.LastAdvertAt.Valid {
|
||||
ms := row.LastAdvertAt.Time.UnixMilli()
|
||||
node.LastAdvertAt = &ms
|
||||
}
|
||||
return node, nil
|
||||
}
|
||||
+291
@@ -0,0 +1,291 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
"log"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
sqlc "github.com/MeshCore-Tower/tower-server/db/sqlc"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/api"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/ingest"
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
)
|
||||
|
||||
func (s *Store) UpsertObserver(ctx context.Context, pubkey []byte) (uuid.UUID, string, error) {
|
||||
row, err := s.q.UpsertObserver(ctx, pubkey)
|
||||
if err != nil {
|
||||
return uuid.Nil, "", err
|
||||
}
|
||||
displayName := ""
|
||||
if row.DisplayName != nil {
|
||||
displayName = *row.DisplayName
|
||||
}
|
||||
return row.ID, displayName, err
|
||||
}
|
||||
|
||||
func (s *Store) ListObservers(ctx context.Context, iatas []string, observerType, broker, status, name, scope string, cursor int64, limit int32) (api.Page[api.ObserverSummary], error) {
|
||||
var cursorTS pgtype.Timestamptz
|
||||
if cursor > 0 {
|
||||
cursorTS = pgtype.Timestamptz{Time: time.UnixMilli(cursor), Valid: true}
|
||||
}
|
||||
iataFilter := strings.Join(iatas, ",")
|
||||
params := sqlc.ListObserversParams{
|
||||
Column1: iataFilter,
|
||||
Column2: observerType,
|
||||
Column3: broker,
|
||||
Column4: status,
|
||||
Column5: name,
|
||||
Column6: cursorTS,
|
||||
Limit: limit + 1,
|
||||
Column8: scope,
|
||||
}
|
||||
rows, err := s.q.ListObservers(ctx, params)
|
||||
if err != nil {
|
||||
return api.Page[api.ObserverSummary]{}, err
|
||||
}
|
||||
hasMore := len(rows) > int(limit)
|
||||
if hasMore {
|
||||
rows = rows[:limit]
|
||||
}
|
||||
items := make([]api.ObserverSummary, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
observer := api.ObserverSummary{
|
||||
ID: v.ID,
|
||||
IATA: v.Iata,
|
||||
Status: v.Status,
|
||||
Scopes: v.Scopes,
|
||||
}
|
||||
if v.RadioFreqMhz != nil && v.RadioSf != nil && v.RadioBwKhz != nil {
|
||||
s := fmt.Sprintf("%.1f,%g,%d", *v.RadioFreqMhz, *v.RadioBwKhz, *v.RadioSf)
|
||||
observer.Radio = &s
|
||||
}
|
||||
if v.DisplayName != nil {
|
||||
observer.DisplayName = v.DisplayName
|
||||
}
|
||||
if v.ObserverType != nil {
|
||||
observer.ObserverType = v.ObserverType
|
||||
}
|
||||
items = append(items, observer)
|
||||
}
|
||||
var nextCursor *int64
|
||||
if hasMore {
|
||||
// observers use UUID so encode last_seen as cursor
|
||||
if rows[len(rows)-1].LastStatusAt.Valid {
|
||||
ms := rows[len(rows)-1].LastStatusAt.Time.UnixMilli()
|
||||
nextCursor = &ms
|
||||
}
|
||||
}
|
||||
return api.Page[api.ObserverSummary]{
|
||||
Items: items,
|
||||
NextCursor: nextCursor,
|
||||
HasMore: hasMore,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetObserver(ctx context.Context, observerID uuid.UUID) (*api.Observer, error) {
|
||||
obs, err := s.q.GetObserverByID(ctx, observerID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
brokerRows, err := s.q.GetObserverBrokers(ctx, observerID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
observer := api.Observer{
|
||||
ObserverSummary: api.ObserverSummary{
|
||||
ID: obs.ID,
|
||||
DisplayName: obs.DisplayName,
|
||||
ObserverType: obs.ObserverType,
|
||||
Status: "offline",
|
||||
},
|
||||
PublicKey: hex.EncodeToString(obs.PublicKey),
|
||||
SoftwareVersion: obs.SoftwareVersion,
|
||||
HardwareModel: obs.HardwareModel,
|
||||
FirmwareVersion: obs.FirmwareVersion,
|
||||
FirmwareBuild: obs.FirmwareBuild,
|
||||
RadioFreqMHz: obs.RadioFreqMhz,
|
||||
RadioSF: obs.RadioSf,
|
||||
RadioBWKHz: obs.RadioBwKhz,
|
||||
RadioCR: obs.RadioCr,
|
||||
BatteryLevel: obs.BatteryLevel,
|
||||
UptimeSeconds: obs.UptimeSeconds,
|
||||
StatusMetadata: obs.StatusMetadata,
|
||||
FirstSeen: obs.FirstSeen.Time.UnixMilli(),
|
||||
LastSeen: obs.LastSeen.Time.UnixMilli(),
|
||||
ObservationCount: *obs.ObservationCount,
|
||||
}
|
||||
scopes, err := s.GetObserverScopes(ctx, observerID)
|
||||
if err != nil {
|
||||
log.Printf("store: GetObserverScopes failed for %s: %v", observerID, err)
|
||||
scopes = []string{}
|
||||
}
|
||||
observer.Scopes = scopes
|
||||
brokers := make([]api.ObserverBroker, 0, len(brokerRows))
|
||||
for _, v := range brokerRows {
|
||||
var lastPacketAt int64
|
||||
if v.LastPacketAt.Valid {
|
||||
lastPacketAt = v.LastPacketAt.Time.UnixMilli()
|
||||
}
|
||||
brokers = append(brokers, api.ObserverBroker{
|
||||
Name: v.BrokerName,
|
||||
LastPacketAt: lastPacketAt,
|
||||
LastSeenAt: v.LastSeen.Time.UnixMilli(),
|
||||
})
|
||||
}
|
||||
observer.Brokers = brokers
|
||||
if obs.LastStatusAt.Valid && time.Since(obs.LastStatusAt.Time) < 5*time.Minute {
|
||||
observer.Status = "online"
|
||||
}
|
||||
var lastStatusAt *int64
|
||||
if obs.LastStatusAt.Valid {
|
||||
ms := obs.LastStatusAt.Time.UnixMilli()
|
||||
lastStatusAt = &ms
|
||||
}
|
||||
observer.LastStatusAt = lastStatusAt
|
||||
observer.IATA, _ = s.GetObserverLastIATA(ctx, observerID)
|
||||
return &observer, nil
|
||||
}
|
||||
|
||||
func (s *Store) InsertObserverTelemetry(ctx context.Context, observerID uuid.UUID, reportedAt time.Time, batteryMV *int32, txAirSecs, rxAirSecs *float32, noiseFloor float32, uptimeSeconds int64, queueLen, debugFlags, recvErrors *int32) error {
|
||||
return s.q.InsertObserverTelemetry(ctx, sqlc.InsertObserverTelemetryParams{
|
||||
ObserverID: observerID,
|
||||
ReportedAt: pgtype.Timestamptz{Time: reportedAt, Valid: true},
|
||||
BatteryVoltageMv: batteryMV,
|
||||
AirtimeTxPct: txAirSecs,
|
||||
AirtimeRxPct: rxAirSecs,
|
||||
NoiseFloorDb: &noiseFloor,
|
||||
UptimeSeconds: &uptimeSeconds,
|
||||
QueueLength: queueLen,
|
||||
DebugFlags: debugFlags,
|
||||
ReceiveErrors: recvErrors,
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Store) GetObserverTelemetry(ctx context.Context, observerID uuid.UUID, since, until time.Time, afterID int64) (*api.ObserverTelemetry, error) {
|
||||
// TODO: implement server-side bucketing by interval when needed.
|
||||
// Currently returns all points in the range at stored resolution.
|
||||
rows, err := s.q.GetObserverTelemetry(ctx, sqlc.GetObserverTelemetryParams{
|
||||
ObserverID: observerID,
|
||||
Column2: pgtype.Timestamptz{Time: since, Valid: !since.IsZero()},
|
||||
Column3: pgtype.Timestamptz{Time: until, Valid: !until.IsZero()},
|
||||
Column4: afterID,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
points := make([]api.ObserverTelemetryPoint, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
points = append(points, api.ObserverTelemetryPoint{
|
||||
T: v.ReportedAt.Time.Unix(),
|
||||
BatteryMV: v.BatteryVoltageMv,
|
||||
AirtimeTxPct: v.AirtimeTxPct,
|
||||
AirtimeRxPct: v.AirtimeRxPct,
|
||||
NoiseFloorDB: v.NoiseFloorDb,
|
||||
UptimeSeconds: v.UptimeSeconds,
|
||||
QueueLength: v.QueueLength,
|
||||
ReceiveErrors: v.ReceiveErrors,
|
||||
})
|
||||
}
|
||||
return &api.ObserverTelemetry{Points: points}, nil
|
||||
}
|
||||
|
||||
func (s *Store) ListObserverAdverts(ctx context.Context, observerID uuid.UUID, cursor int64, limit int32) (api.Page[api.AdvertObservation], error) {
|
||||
rows, err := s.q.ListObserverAdverts(ctx, sqlc.ListObserverAdvertsParams{
|
||||
ObserverID: observerID,
|
||||
Column2: cursor,
|
||||
Limit: limit + 1, // fetch one extra to detect hasMore
|
||||
})
|
||||
if err != nil {
|
||||
log.Printf("api: ListObserverAdverts failed: %v", err)
|
||||
return api.Page[api.AdvertObservation]{}, err
|
||||
}
|
||||
hasMore := len(rows) > int(limit)
|
||||
if hasMore {
|
||||
rows = rows[:limit]
|
||||
}
|
||||
items := make([]api.AdvertObservation, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
items = append(items, api.AdvertObservation{
|
||||
PacketObservationSummary: api.PacketObservationSummary{
|
||||
ID: v.ID,
|
||||
PacketHash: v.PacketHashHex,
|
||||
PayloadType: v.PayloadType,
|
||||
PayloadTypeName: api.PayloadTypeName(v.PayloadType),
|
||||
IATA: v.Iata,
|
||||
HeardAt: v.HeardAt.Time.UnixMilli(),
|
||||
RSSI: v.Rssi,
|
||||
SNR: v.Snr,
|
||||
HopCount: &v.HopCount,
|
||||
},
|
||||
NodeName: v.NodeName,
|
||||
NodePublicKey: &v.NodePublicKey,
|
||||
})
|
||||
}
|
||||
var nextCursor *int64
|
||||
if hasMore {
|
||||
last := items[len(items)-1].ID
|
||||
nextCursor = &last
|
||||
}
|
||||
return api.Page[api.AdvertObservation]{
|
||||
Items: items,
|
||||
NextCursor: nextCursor,
|
||||
HasMore: hasMore,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *Store) UpdateObserverStatus(ctx context.Context, p ingest.UpdateObserverStatusParams) (uuid.UUID, error) {
|
||||
params := sqlc.UpdateObserverStatusParams{PublicKey: p.PublicKey, Column2: p.DisplayName, Column3: p.ObserverType, SoftwareVersion: &p.SoftwareVersion, HardwareModel: &p.HardwareModel, FirmwareVersion: &p.FirmwareVersion, FirmwareBuild: &p.FirmwareBuild, RadioFreqMhz: &p.RadioFreqMHz, RadioSf: &p.RadioSF, RadioBwKhz: &p.RadioBWKHz, RadioCr: &p.RadioCR, BatteryLevel: p.BatteryLevel, UptimeSeconds: p.UptimeSeconds, StatusMetadata: p.StatusMetadata}
|
||||
return s.q.UpdateObserverStatus(ctx, params)
|
||||
}
|
||||
|
||||
func (s *Store) GetObserverLastIATA(ctx context.Context, observerID uuid.UUID) (string, error) {
|
||||
return s.q.GetObserverLastIATA(ctx, observerID)
|
||||
}
|
||||
|
||||
func (s *Store) GetObserverRadio(ctx context.Context, observerID uuid.UUID) (ingest.RadioSettings, error) {
|
||||
row, err := s.q.GetObserverRadio(ctx, observerID)
|
||||
if err != nil {
|
||||
return ingest.RadioSettings{}, err
|
||||
}
|
||||
var settings ingest.RadioSettings
|
||||
if row.RadioFreqMhz != nil {
|
||||
settings.FreqMHz = *row.RadioFreqMhz
|
||||
}
|
||||
if row.RadioSf != nil {
|
||||
settings.SF = *row.RadioSf
|
||||
}
|
||||
if row.RadioBwKhz != nil {
|
||||
settings.BWKHz = *row.RadioBwKhz
|
||||
}
|
||||
if row.RadioCr != nil {
|
||||
settings.CR = *row.RadioCr
|
||||
}
|
||||
return settings, nil
|
||||
}
|
||||
|
||||
func (s *Store) UpsertObserverBroker(ctx context.Context, observerID uuid.UUID, brokerName string) error {
|
||||
params := sqlc.UpsertObserverBrokerParams{
|
||||
ObserverID: observerID,
|
||||
BrokerName: brokerName,
|
||||
}
|
||||
return s.q.UpsertObserverBroker(ctx, params)
|
||||
}
|
||||
|
||||
func (s *Store) UpsertObserverScope(ctx context.Context, observerID uuid.UUID, scopeID int32) error {
|
||||
return s.q.UpsertObserverScope(ctx, sqlc.UpsertObserverScopeParams{
|
||||
ObserverID: observerID,
|
||||
ScopeID: scopeID,
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Store) GetObserverScopes(ctx context.Context, observerID uuid.UUID) ([]string, error) {
|
||||
return s.q.GetObserverScopes(ctx, observerID)
|
||||
}
|
||||
|
||||
func (s *Store) DeleteOldTelemetry(ctx context.Context, cutoff time.Time) error {
|
||||
return s.q.DeleteOldTelemetry(ctx, pgtype.Timestamptz{Time: cutoff, Valid: true})
|
||||
}
|
||||
+407
@@ -0,0 +1,407 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/binary"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
sqlc "github.com/MeshCore-Tower/tower-server/db/sqlc"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/api"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/ingest"
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
)
|
||||
|
||||
func (s *Store) UpsertPacket(ctx context.Context, p ingest.UpsertPacketParams) (bool, error) {
|
||||
var regionCode, subRegionCode *int32
|
||||
hasTransportCodes := len(p.TransportCodes) == 4
|
||||
if hasTransportCodes {
|
||||
r := int32(binary.LittleEndian.Uint16(p.TransportCodes[0:2]))
|
||||
s := int32(binary.LittleEndian.Uint16(p.TransportCodes[2:4]))
|
||||
regionCode = &r
|
||||
subRegionCode = &s
|
||||
}
|
||||
params := sqlc.UpsertPacketParams{
|
||||
PacketHash: p.PacketHash,
|
||||
PayloadType: int16(p.PayloadType),
|
||||
PayloadVersion: int16(p.PayloadVersion),
|
||||
RouteType: int16(p.RouteType),
|
||||
TransportCodesPresent: &hasTransportCodes,
|
||||
RegionCode: regionCode,
|
||||
SubRegionCode: subRegionCode,
|
||||
OriginPubkey: p.OriginPubkey,
|
||||
RawPayload: p.RawPayload,
|
||||
RawHeader: p.RawHeader,
|
||||
ParsedPayload: p.ParsedPayload,
|
||||
ChannelHash: p.ChannelHash,
|
||||
ScopeID: p.ScopeID,
|
||||
}
|
||||
row, err := s.q.UpsertPacket(ctx, params)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return row.Inserted, nil
|
||||
}
|
||||
|
||||
func (s *Store) ListPackets(ctx context.Context, payloadType, routeType int16, iatas []string, scope string, since, until time.Time, cursor int64, limit int32) (api.Page[api.PacketSummary], error) {
|
||||
var cursorTS pgtype.Timestamptz
|
||||
if cursor > 0 {
|
||||
cursorTS = pgtype.Timestamptz{Time: time.UnixMilli(cursor), Valid: true}
|
||||
}
|
||||
var sinceTS pgtype.Timestamptz
|
||||
if !since.IsZero() {
|
||||
sinceTS = pgtype.Timestamptz{Time: since, Valid: true}
|
||||
}
|
||||
var untilTS pgtype.Timestamptz
|
||||
if !until.IsZero() {
|
||||
untilTS = pgtype.Timestamptz{Time: until, Valid: true}
|
||||
}
|
||||
iataFilter := strings.Join(iatas, ",")
|
||||
rows, err := s.q.ListPackets(ctx, sqlc.ListPacketsParams{
|
||||
Column1: payloadType,
|
||||
Column2: routeType,
|
||||
Column3: iataFilter,
|
||||
Column4: sinceTS,
|
||||
Column5: untilTS,
|
||||
Column6: cursorTS,
|
||||
Limit: limit + 1,
|
||||
Column8: scope,
|
||||
})
|
||||
if err != nil {
|
||||
return api.Page[api.PacketSummary]{}, err
|
||||
}
|
||||
hasMore := len(rows) > int(limit)
|
||||
if hasMore {
|
||||
rows = rows[:limit]
|
||||
}
|
||||
items := make([]api.PacketSummary, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
item := api.PacketSummary{
|
||||
PacketHash: hex.EncodeToString(v.PacketHash),
|
||||
PayloadType: v.PayloadType,
|
||||
PayloadTypeName: api.PayloadTypeName(v.PayloadType),
|
||||
RouteType: v.RouteType,
|
||||
RouteTypeName: api.RouteTypeName(v.RouteType),
|
||||
Scope: v.ScopeName,
|
||||
FirstHeardAt: v.FirstHeardAt.Time.UnixMilli(),
|
||||
LastHeardAt: v.LastHeardAt.Time.UnixMilli(),
|
||||
ObservationCount: int32(v.ObservationCount),
|
||||
}
|
||||
if v.LatestObserverID != (uuid.UUID{}) {
|
||||
item.LatestObserver = &api.PacketLatestObserver{
|
||||
ID: v.LatestObserverID,
|
||||
DisplayName: v.LatestObserverName,
|
||||
IATA: v.LatestObserverIata,
|
||||
}
|
||||
}
|
||||
items = append(items, item)
|
||||
}
|
||||
var nextCursor *int64
|
||||
if hasMore && len(items) > 0 {
|
||||
last := items[len(items)-1].LastHeardAt
|
||||
nextCursor = &last
|
||||
}
|
||||
return api.Page[api.PacketSummary]{
|
||||
Items: items,
|
||||
NextCursor: nextCursor,
|
||||
HasMore: hasMore,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetPacket(ctx context.Context, packetHash []byte) (*api.Packet, error) {
|
||||
row, err := s.q.GetPacketByHash(ctx, packetHash)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
obsRows, err := s.q.ListObservationsForPacket(ctx, packetHash)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
p := &api.Packet{
|
||||
PacketHash: hex.EncodeToString(row.PacketHash),
|
||||
Header: api.PacketHeader{
|
||||
Raw: hex.EncodeToString(row.RawHeader),
|
||||
RouteType: row.RouteType,
|
||||
RouteTypeName: api.RouteTypeName(row.RouteType),
|
||||
PayloadType: row.PayloadType,
|
||||
PayloadTypeName: api.PayloadTypeName(row.PayloadType),
|
||||
PayloadVersion: row.PayloadVersion,
|
||||
},
|
||||
ParsedPayload: row.ParsedPayload,
|
||||
RawPayload: hex.EncodeToString(row.RawPayload),
|
||||
Decrypted: row.Decrypted != nil && *row.Decrypted,
|
||||
Scope: row.ScopeName,
|
||||
FirstHeardAt: row.FirstHeardAt.Time.UnixMilli(),
|
||||
LastHeardAt: row.LastHeardAt.Time.UnixMilli(),
|
||||
ObservationCount: int32(len(obsRows)),
|
||||
Observations: make([]api.PacketObservationDetail, 0, len(obsRows)),
|
||||
}
|
||||
minHeardAt := obsRows[0].HeardAt.Time
|
||||
if len(obsRows) > 1 {
|
||||
maxHeardAt := obsRows[0].HeardAt.Time
|
||||
for _, v := range obsRows[1:] {
|
||||
if v.HeardAt.Time.Before(minHeardAt) {
|
||||
minHeardAt = v.HeardAt.Time
|
||||
}
|
||||
if v.HeardAt.Time.After(maxHeardAt) {
|
||||
maxHeardAt = v.HeardAt.Time
|
||||
}
|
||||
}
|
||||
p.FirstToLastMs = maxHeardAt.Sub(minHeardAt).Milliseconds()
|
||||
}
|
||||
if row.OriginPubkey != nil {
|
||||
s := hex.EncodeToString(row.OriginPubkey)
|
||||
p.OriginPubkey = &s
|
||||
}
|
||||
if row.ChannelHash != nil {
|
||||
ch := hex.EncodeToString(row.ChannelHash)
|
||||
p.ChannelHash = &ch
|
||||
}
|
||||
if row.TransportCodesPresent != nil && *row.TransportCodesPresent {
|
||||
tc := &api.PacketTransportCodes{}
|
||||
if row.RegionCode != nil {
|
||||
tc.RegionCode = *row.RegionCode
|
||||
}
|
||||
if row.SubRegionCode != nil {
|
||||
tc.SubRegionCode = *row.SubRegionCode
|
||||
}
|
||||
p.TransportCodes = tc
|
||||
}
|
||||
for _, v := range obsRows {
|
||||
obs := api.PacketObservationDetail{
|
||||
ID: v.ID,
|
||||
ObserverID: v.ObserverID,
|
||||
ObserverName: v.ObserverName,
|
||||
IATA: v.Iata,
|
||||
HeardAt: v.HeardAt.Time.UnixMilli(),
|
||||
PathLength: api.PacketPathLength{
|
||||
Raw: fmt.Sprintf("%02x", v.PathLengthByte),
|
||||
HashSize: v.HashSize,
|
||||
HopCount: v.HopCount,
|
||||
},
|
||||
RSSI: v.Rssi,
|
||||
SNR: v.Snr,
|
||||
SourceBroker: *v.SourceBroker,
|
||||
}
|
||||
prop := int32(v.HeardAt.Time.Sub(minHeardAt).Milliseconds())
|
||||
obs.PropagationTimeMs = &prop
|
||||
resolvedPath := []api.ResolvedHop{}
|
||||
if v.PathBytes != nil && v.HashSize > 0 {
|
||||
hashSize := int(v.HashSize)
|
||||
hashes := make([][]byte, 0, len(v.PathBytes)/hashSize)
|
||||
for i := 0; i+hashSize <= len(v.PathBytes); i += hashSize {
|
||||
hashes = append(hashes, v.PathBytes[i:i+hashSize])
|
||||
}
|
||||
resolved, err := s.ResolvePathHashes(ctx, v.Iata, hashes)
|
||||
if err != nil {
|
||||
log.Printf("store: path resolution failed for observation %d: %v", v.ID, err)
|
||||
} else {
|
||||
for _, hash := range hashes {
|
||||
key := hex.EncodeToString(hash)
|
||||
entries := resolved[key]
|
||||
hop := api.ResolvedHop{
|
||||
Nodes: make([]api.ResolvedNode, 0, len(entries)),
|
||||
}
|
||||
switch len(entries) {
|
||||
case 0:
|
||||
hop.Confidence = "none"
|
||||
case 1:
|
||||
hop.Confidence = "high"
|
||||
default:
|
||||
hop.Confidence = "ambiguous"
|
||||
}
|
||||
for _, e := range entries {
|
||||
hop.Nodes = append(hop.Nodes, api.ResolvedNode{
|
||||
ID: e.NodeID,
|
||||
Name: e.Name,
|
||||
Latitude: e.Latitude,
|
||||
Longitude: e.Longitude,
|
||||
PublicKey: hex.EncodeToString(e.PublicKey),
|
||||
})
|
||||
}
|
||||
resolvedPath = append(resolvedPath, hop)
|
||||
}
|
||||
}
|
||||
}
|
||||
obs.ResolvedPath = resolvedPath
|
||||
if v.PathBytes != nil {
|
||||
pb := hex.EncodeToString(v.PathBytes)
|
||||
obs.PathBytes = &pb
|
||||
}
|
||||
if v.RadioFreqMhz != nil || v.SpreadFactor != nil || v.BandwidthKhz != nil || v.CodingRate != nil {
|
||||
obs.Radio = &api.PacketRadio{
|
||||
FreqMHz: v.RadioFreqMhz,
|
||||
SpreadFactor: v.SpreadFactor,
|
||||
BandwidthKHz: v.BandwidthKhz,
|
||||
CodingRate: v.CodingRate,
|
||||
}
|
||||
}
|
||||
p.Observations = append(p.Observations, obs)
|
||||
}
|
||||
if row.PayloadType == 9 && len(obsRows) > 0 {
|
||||
var tracePayload struct {
|
||||
PathHashes []string `json:"pathHashes"`
|
||||
Flags byte `json:"flags"`
|
||||
}
|
||||
if err := json.Unmarshal(row.ParsedPayload, &tracePayload); err == nil && len(tracePayload.PathHashes) > 0 {
|
||||
hashSize := int(1 << (tracePayload.Flags & 0x03))
|
||||
hashes := make([][]byte, 0, len(tracePayload.PathHashes))
|
||||
for _, h := range tracePayload.PathHashes {
|
||||
b, err := hex.DecodeString(h)
|
||||
if err == nil {
|
||||
hashes = append(hashes, b)
|
||||
}
|
||||
}
|
||||
|
||||
// collect unique IATAs from observations
|
||||
seenIATAs := make(map[string]struct{})
|
||||
for _, v := range obsRows {
|
||||
seenIATAs[v.Iata] = struct{}{}
|
||||
}
|
||||
|
||||
// merge results across all IATAs — best confidence per hop wins
|
||||
// confidence priority: high > ambiguous > none
|
||||
type hopResult struct {
|
||||
confidence string
|
||||
entries []api.ResolvedPathEntry
|
||||
}
|
||||
merged := make([]hopResult, len(hashes))
|
||||
for i := range merged {
|
||||
merged[i] = hopResult{confidence: "none"}
|
||||
}
|
||||
|
||||
confidenceRank := map[string]int{"none": 0, "ambiguous": 1, "high": 2}
|
||||
|
||||
for iata := range seenIATAs {
|
||||
resolved, err := s.ResolvePathHashes(ctx, iata, hashes)
|
||||
if err != nil {
|
||||
log.Printf("store: trace route resolution failed for iata=%s: %v", iata, err)
|
||||
continue
|
||||
}
|
||||
for i, hash := range hashes {
|
||||
key := hex.EncodeToString(hash[:hashSize])
|
||||
entries := resolved[key]
|
||||
var confidence string
|
||||
switch len(entries) {
|
||||
case 0:
|
||||
confidence = "none"
|
||||
case 1:
|
||||
confidence = "high"
|
||||
default:
|
||||
confidence = "ambiguous"
|
||||
}
|
||||
if confidenceRank[confidence] > confidenceRank[merged[i].confidence] {
|
||||
merged[i] = hopResult{confidence: confidence, entries: entries}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
route := make([]api.ResolvedHop, 0, len(hashes))
|
||||
for _, hr := range merged {
|
||||
hop := api.ResolvedHop{
|
||||
Confidence: hr.confidence,
|
||||
Nodes: make([]api.ResolvedNode, 0, len(hr.entries)),
|
||||
}
|
||||
for _, e := range hr.entries {
|
||||
hop.Nodes = append(hop.Nodes, api.ResolvedNode{
|
||||
ID: e.NodeID,
|
||||
Name: e.Name,
|
||||
Latitude: e.Latitude,
|
||||
Longitude: e.Longitude,
|
||||
PublicKey: hex.EncodeToString(e.PublicKey),
|
||||
})
|
||||
}
|
||||
route = append(route, hop)
|
||||
}
|
||||
p.ResolvedRoute = route
|
||||
}
|
||||
}
|
||||
return p, nil
|
||||
}
|
||||
|
||||
func (s *Store) UpsertIATA(ctx context.Context, iata string) error {
|
||||
return s.q.UpsertIATA(ctx, iata)
|
||||
}
|
||||
|
||||
func (s *Store) InsertObservation(ctx context.Context, o ingest.InsertObservationParams) (bool, error) {
|
||||
params := sqlc.InsertObservationParams{
|
||||
PacketHash: o.PacketHash,
|
||||
ObserverID: o.ObserverID,
|
||||
Iata: o.IATA,
|
||||
HeardAt: pgtype.Timestamptz{Time: o.HeardAt, Valid: true},
|
||||
PathLengthByte: int16(o.PathLengthByte),
|
||||
HashSize: int16(o.HashSize),
|
||||
HopCount: int16(o.HopCount),
|
||||
PathBytes: o.PathBytes,
|
||||
Rssi: &o.RSSI,
|
||||
Snr: &o.SNR,
|
||||
PropagationTimeMs: &o.PropagationTimeMs,
|
||||
RadioFreqMhz: &o.RadioFreqMHz,
|
||||
SpreadFactor: &o.SpreadFactor,
|
||||
BandwidthKhz: &o.BandwidthKHz,
|
||||
CodingRate: &o.CodingRate,
|
||||
SourceBroker: &o.SourceBroker,
|
||||
}
|
||||
row, err := s.q.InsertObservation(ctx, params)
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return false, nil // conflict, not an error
|
||||
}
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return row.ID != 0, nil
|
||||
}
|
||||
|
||||
func (s *Store) ListNodeObservations(ctx context.Context, nodeID uuid.UUID, cursor int64, limit int32) (api.Page[api.PacketObservationSummary], error) {
|
||||
rows, err := s.q.ListNodeObservations(ctx, sqlc.ListNodeObservationsParams{
|
||||
ID: nodeID,
|
||||
Column2: cursor,
|
||||
Limit: limit + 1,
|
||||
})
|
||||
if err != nil {
|
||||
return api.Page[api.PacketObservationSummary]{}, err
|
||||
}
|
||||
hasMore := len(rows) > int(limit)
|
||||
if hasMore {
|
||||
rows = rows[:limit]
|
||||
}
|
||||
items := make([]api.PacketObservationSummary, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
items = append(items, api.PacketObservationSummary{
|
||||
ID: v.ID,
|
||||
PacketHash: v.PacketHashHex,
|
||||
PayloadType: v.PayloadType,
|
||||
PayloadTypeName: api.PayloadTypeName(v.PayloadType),
|
||||
IATA: v.Iata,
|
||||
HeardAt: v.HeardAt.Time.UnixMilli(),
|
||||
RSSI: v.Rssi,
|
||||
SNR: v.Snr,
|
||||
HopCount: &v.HopCount,
|
||||
})
|
||||
}
|
||||
var nextCursor *int64
|
||||
if hasMore && len(items) > 0 {
|
||||
last := items[len(items)-1].ID
|
||||
nextCursor = &last
|
||||
}
|
||||
return api.Page[api.PacketObservationSummary]{
|
||||
Items: items,
|
||||
NextCursor: nextCursor,
|
||||
HasMore: hasMore,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetPacketObservationCount(ctx context.Context, packetHash []byte) (int64, error) {
|
||||
return s.q.GetPacketObservationCount(ctx, packetHash)
|
||||
}
|
||||
|
||||
func (s *Store) DeleteOldPackets(ctx context.Context, cutoff time.Time) error {
|
||||
return s.q.DeleteOldPackets(ctx, pgtype.Timestamptz{Time: cutoff, Valid: true})
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
sqlc "github.com/MeshCore-Tower/tower-server/db/sqlc"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/scopestore"
|
||||
)
|
||||
|
||||
func (s *Store) UpsertTransportScope(ctx context.Context, name, displayName string, transportKey, keyFingerprint []byte) error {
|
||||
var dn *string
|
||||
if displayName != "" {
|
||||
dn = &displayName
|
||||
}
|
||||
return s.q.UpsertTransportScope(ctx, sqlc.UpsertTransportScopeParams{
|
||||
Name: name,
|
||||
DisplayName: dn,
|
||||
TransportKey: transportKey,
|
||||
KeyFingerprint: keyFingerprint,
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Store) GetTransportScopes(ctx context.Context) ([]scopestore.Entry, error) {
|
||||
rows, err := s.q.GetTransportScopes(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
entries := make([]scopestore.Entry, 0, len(rows))
|
||||
for _, r := range rows {
|
||||
entries = append(entries, scopestore.Entry{
|
||||
Name: r.Name,
|
||||
TransportKey: r.TransportKey,
|
||||
KeyFingerprint: r.KeyFingerprint,
|
||||
})
|
||||
}
|
||||
return entries, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetTransportScopeByName(ctx context.Context, name string) (int32, error) {
|
||||
return s.q.GetTransportScopeByName(ctx, name)
|
||||
}
|
||||
+173
@@ -0,0 +1,173 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
sqlc "github.com/MeshCore-Tower/tower-server/db/sqlc"
|
||||
"github.com/MeshCore-Tower/tower-server/internal/api"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
)
|
||||
|
||||
func (s *Store) GetStatsOverview(ctx context.Context, iata string) (*api.StatsOverview, error) {
|
||||
row, err := s.q.GetStatsOverview(ctx, iata)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &api.StatsOverview{
|
||||
TotalPackets: row.TotalPackets,
|
||||
TotalObservations: row.TotalObservations,
|
||||
ActiveObservers: row.ActiveObservers,
|
||||
ActiveIATAs: row.ActiveIatas,
|
||||
WindowHours: 24,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetStatsObservations(ctx context.Context, iata string, since time.Time) ([]api.ObservationPoint, error) {
|
||||
if since.IsZero() {
|
||||
since = time.Now().Add(-7 * 24 * time.Hour)
|
||||
}
|
||||
interval := time.Since(since)
|
||||
rows, err := s.q.GetHourlyStats(ctx, sqlc.GetHourlyStatsParams{
|
||||
Column1: iata,
|
||||
Column2: pgtype.Interval{Microseconds: int64(interval.Hours()) * 3600 * 1e6, Valid: true},
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
points := make([]api.ObservationPoint, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
points = append(points, api.ObservationPoint{
|
||||
Hour: v.Hour.Time.UnixMilli(),
|
||||
IATA: v.Iata,
|
||||
ObservationCount: v.ObservationCount,
|
||||
UniquePackets: v.UniquePackets,
|
||||
ActiveObservers: v.ActiveObservers,
|
||||
})
|
||||
}
|
||||
return points, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetStatsPayloadBreakdown(ctx context.Context, iata string, since time.Time) ([]api.PayloadBreakdownItem, error) {
|
||||
if since.IsZero() {
|
||||
since = time.Now().Add(-24 * time.Hour)
|
||||
}
|
||||
rows, err := s.q.GetStatsPayloadBreakdown(ctx, sqlc.GetStatsPayloadBreakdownParams{
|
||||
HeardAt: pgtype.Timestamptz{Time: since, Valid: true},
|
||||
Column2: iata,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
items := make([]api.PayloadBreakdownItem, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
items = append(items, api.PayloadBreakdownItem{
|
||||
PayloadType: v.PayloadType,
|
||||
PayloadTypeName: api.PayloadTypeName(v.PayloadType),
|
||||
Count: v.Count,
|
||||
})
|
||||
}
|
||||
return items, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetStatsTopNodes(ctx context.Context, iata string, limit int32) ([]api.TopNode, error) {
|
||||
rows, err := s.q.GetTopNodes(ctx, sqlc.GetTopNodesParams{
|
||||
Column1: iata,
|
||||
Limit: limit,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
items := make([]api.TopNode, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
var count int64
|
||||
if v.ObservationCount != nil {
|
||||
count = *v.ObservationCount
|
||||
}
|
||||
items = append(items, api.TopNode{
|
||||
NodeID: v.NodeID,
|
||||
NodeName: v.Name,
|
||||
NodeType: v.NodeType,
|
||||
NodeTypeName: api.NodeTypeName(v.NodeType),
|
||||
IATA: v.Iata,
|
||||
ObservationCount: count,
|
||||
LastHeard: v.LastHeard.Time.UnixMilli(),
|
||||
})
|
||||
}
|
||||
return items, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetStatsTopObservers(ctx context.Context, iata string, since time.Time, limit int32) ([]api.TopObserver, error) {
|
||||
if since.IsZero() {
|
||||
since = time.Now().Add(-24 * time.Hour)
|
||||
}
|
||||
rows, err := s.q.GetStatsTopObservers(ctx, sqlc.GetStatsTopObserversParams{
|
||||
HeardAt: pgtype.Timestamptz{Time: since, Valid: true},
|
||||
Column2: iata,
|
||||
Limit: limit,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
items := make([]api.TopObserver, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
iata, _ := v.Iata.(string)
|
||||
items = append(items, api.TopObserver{
|
||||
ObserverID: v.ID,
|
||||
DisplayName: v.DisplayName,
|
||||
ObserverType: v.ObserverType,
|
||||
IATA: iata,
|
||||
ObservationCount: v.ObservationCount,
|
||||
})
|
||||
}
|
||||
return items, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetRadioPresets(ctx context.Context, preset, iata string) ([]api.RadioPreset, error) {
|
||||
rows, err := s.q.GetRadioPresets(ctx, sqlc.GetRadioPresetsParams{
|
||||
Column1: preset,
|
||||
Column2: iata,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
items := make([]api.RadioPreset, 0, len(rows))
|
||||
for _, v := range rows {
|
||||
items = append(items, api.RadioPreset{
|
||||
Preset: v.Preset,
|
||||
IATA: v.Iata,
|
||||
SourceType: v.SourceType,
|
||||
Count: v.Count,
|
||||
})
|
||||
}
|
||||
return items, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetScopeStats(ctx context.Context) ([]api.ScopeStats, error) {
|
||||
rows, err := s.q.GetScopeStats(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
items := make([]api.ScopeStats, 0, len(rows))
|
||||
for _, r := range rows {
|
||||
items = append(items, api.ScopeStats{
|
||||
Name: r.Name,
|
||||
PacketCount: r.PacketCount,
|
||||
ObserverCount: r.ObserverCount,
|
||||
NodeCount: r.NodeCount,
|
||||
})
|
||||
}
|
||||
return items, nil
|
||||
}
|
||||
|
||||
func (s *Store) RefreshHourlyStats(ctx context.Context) error {
|
||||
return s.q.RefreshHourlyStats(ctx)
|
||||
}
|
||||
|
||||
func (s *Store) RefreshTopNodes(ctx context.Context) error {
|
||||
return s.q.RefreshTopNodes(ctx)
|
||||
}
|
||||
|
||||
func (s *Store) RefreshRadioPresets(ctx context.Context) error {
|
||||
return s.q.RefreshRadioPresets(ctx)
|
||||
}
|
||||
+11
-1472
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user