store

package
v1.0.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 1, 2026 License: Apache-2.0 Imports: 22 Imported by: 0

Documentation

Overview

Package store keeps everything the central knows in PostgreSQL with TimescaleDB. Migrations are embedded and applied on start.

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrNotFound     = errors.New("not found")
	ErrTokenInvalid = errors.New("enrollment token is invalid, expired or already used")
	ErrNameTaken    = errors.New("name already in use")
)
View Source
var (
	ErrUsernameTaken = errors.New("username already in use")
	// ErrLastAdmin protects the last enabled admin from removal.
	ErrLastAdmin = errors.New("at least one enabled admin must remain")
)

Functions

func ASPath

func ASPath(hops []PathHop) []int64

ASPath lists the distinct AS numbers along hops, unknown hops skipped.

Types

type Address

type Address struct {
	Source     Source
	Family     string // ipv4 or ipv6
	IP         netip.Addr
	Net        NetInfo
	ValidFrom  time.Time
	ValidTo    *time.Time
	LastSeenAt time.Time
}

type Alert

type Alert struct {
	ID         int64
	RuleID     int64
	RuleName   string
	Kind       RuleKind
	EdgeID     uuid.UUID
	EdgeName   string
	CheckID    *int64
	TargetID   *uuid.UUID
	TargetName string
	Summary    string
	Value      float64
	StartedAt  time.Time
	ResolvedAt *time.Time
	WebhookIDs []int64
}

type AuditEntry

type AuditEntry struct {
	ID       int64
	At       time.Time
	Actor    string // username, "cli" or "anonymous"
	Action   string // e.g. "SaveTarget", "Login"
	Outcome  string // "ok" or an error code
	ClientIP string
	Detail   json.RawMessage // request, secrets redacted
}

AuditEntry is one recorded action.

type Certificate

type Certificate struct {
	Serial    string
	NotBefore time.Time
	NotAfter  time.Time
}

type Check

type Check struct {
	ID        int64
	VersionID int64
	TargetID  uuid.UUID
	Host      string // host of the current version
	Spec      check.Spec
	Selector  map[string]string
}

type CheckInput

type CheckInput struct {
	ID       int64
	Spec     check.Spec
	Selector map[string]string
}

CheckInput is one check of a SaveTarget call. ID 0 creates a check.

type CheckRow

type CheckRow struct {
	Time       time.Time
	Edge       string
	Target     string
	Host       string
	Spec       check.Spec
	ResolvedIP netip.Addr
	DNSSeconds *float64
	Sent       int
	Received   int
	RTTMin     *float64
	RTTAvg     *float64
	RTTMax     *float64
	Jitter     *float64
	Error      string
	Metrics    map[string]float64
}

CheckRow is one probe cycle, with names, for exports. Seconds.

type Delivery

type Delivery struct {
	ID          int64
	WebhookID   int64
	Event       string
	Summary     string
	Payload     []byte
	Attempts    int
	NextAttempt *time.Time
	DeliveredAt *time.Time
	StatusCode  int
	Error       string
	CreatedAt   time.Time
}

func (Delivery) State

func (d Delivery) State() string

State is delivered, pending or failed.

type Edge

type Edge struct {
	ID           uuid.UUID
	Name         string
	Labels       map[string]string
	Status       Status
	EnrolledAt   time.Time
	RevokedAt    *time.Time
	LastSeenAt   *time.Time
	AgentVersion string
	OS           string
	Arch         string
	Hostname     string
	ClockSkew    *float64 // seconds
	CertNotAfter *time.Time
	Addresses    []Address // current ones, filled by ListEdges and Edge
}

type EdgeInfo

type EdgeInfo struct {
	AgentVersion, OS, Arch, Hostname string
}

EdgeInfo is sent by the edge when its session opens.

type ExportFilter

type ExportFilter struct {
	From, To time.Time
	EdgeID   *uuid.UUID
	TargetID *uuid.UUID
}

ExportFilter narrows an export; nil fields are ignored.

type IssueFunc

type IssueFunc func(id uuid.UUID, name string) (Certificate, error)

IssueFunc signs the certificate inside the enrollment transaction.

type Latest

type Latest struct {
	EdgeID     uuid.UUID
	EdgeName   string
	TargetID   uuid.UUID
	TargetName string
	CheckID    int64
	Spec       check.Spec
	Time       time.Time
	ResolvedIP netip.Addr
	Sent       int
	Received   int
	RTTMin     *float64
	RTTAvg     *float64
	RTTMax     *float64
	Jitter     *float64
	Error      string
	Metrics    map[string]float64 // per protocol, nil if none
	Trace      *TraceSummary      // traceroute checks only
}

Latest is the last result of an edge for a check, with its context.

type LatestPathChange

type LatestPathChange struct {
	EdgeID    uuid.UUID
	CheckID   int64
	TargetID  uuid.UUID
	ASPath    []int64
	ChangedAt time.Time
}

PathChange of a traceroute, as seen in trace_latest.

type NetInfo

type NetInfo struct {
	ASN   uint32
	ASOrg string
	DB    string // identifies the database build that answered
}

NetInfo is the AS of an address.

type PathChange

type PathChange struct {
	Time   time.Time
	ASPath []int64
}

PathChange is the first traceroute with a new AS path.

type PathHop

type PathHop struct {
	IP  netip.Addr `json:"ip"`
	Net NetInfo    `json:"net"`
}

PathHop is one TTL of a path. IP is zero for a silent hop.

type Point

type Point struct {
	Time     time.Time
	Sent     int64
	Received int64
	RTTAvg   *float64
	RTTMin   *float64
	RTTMax   *float64
	Jitter   *float64
}

Point is one bucket of a series. Seconds; RTT fields are nil without replies.

type RangeCount

type RangeCount struct {
	Table string
	Rows  int64
}

RangeCount is the number of rows of a table in a time range.

type Reportable

type Reportable struct {
	Protocol check.Protocol
	Since    time.Time // the version exists from then: no result may be older
}

Reportable is a check version an edge may send results for.

type Result

type Result struct {
	CheckVersionID int64
	Time           time.Time
	ResolvedIP     netip.Addr
	DNSSeconds     float64 // 0 when no lookup was needed
	Sent           int
	Received       int
	Samples        []float32
	RTTMin         float64
	RTTAvg         float64
	RTTMax         float64
	RTTStdDev      float64
	Jitter         float64
	Error          string
	Metrics        map[string]float64 // per protocol, nil if none
}

Result is one cycle, with statistics already computed. Seconds.

type Role

type Role string
const (
	RoleAdmin  Role = "admin"
	RoleViewer Role = "viewer"
)

type Rule

type Rule struct {
	ID         int64
	Name       string
	Kind       RuleKind
	Threshold  float64 // percent for loss, ms for rtt
	Window     time.Duration
	TargetID   *uuid.UUID
	Selector   map[string]string
	WebhookIDs []int64
	Enabled    bool
}

type RuleKind

type RuleKind string
const (
	RuleEdgeOffline RuleKind = "edge_offline"
	RuleLoss        RuleKind = "loss"
	RuleRTT         RuleKind = "rtt"
	RulePathChange  RuleKind = "path_change"
)

type Session

type Session struct {
	TokenHash  []byte
	UserID     uuid.UUID
	CreatedAt  time.Time
	LastSeenAt time.Time
	ExpiresAt  time.Time
	ClientIP   string
	UserAgent  string
}

Session is a login. The token itself is never stored.

type Source

type Source string

Source says who observed an address.

const (
	SourceSession Source = "session" // source address of the edge's session
	SourceSTUN    Source = "stun"    // mapped address seen by the STUN server
	SourceLocal   Source = "local"   // the edge's own socket address
)

type Status

type Status string
const (
	StatusActive  Status = "active"
	StatusRevoked Status = "revoked"
)

type Store

type Store struct {
	// contains filtered or unexported fields
}

func Open

func Open(ctx context.Context, url string) (*Store, error)

Open connects and brings the schema up to date.

func (*Store) ActiveChecks

func (s *Store) ActiveChecks(ctx context.Context) ([]Check, error)

ActiveChecks returns the current version of every live check.

func (*Store) AddressHistory

func (s *Store) AddressHistory(ctx context.Context, id uuid.UUID) ([]Address, error)

AddressHistory returns every address period of an edge, newest first.

func (*Store) Audit

func (s *Store) Audit(ctx context.Context, e AuditEntry) error

func (*Store) AuditLog

func (s *Store) AuditLog(ctx context.Context, before int64, limit int) ([]AuditEntry, error)

AuditLog returns entries newest first, before an id (0: from the newest).

func (*Store) Authenticate

func (s *Store) Authenticate(ctx context.Context, edgeID uuid.UUID, serial string) (string, error)

Authenticate checks a certificate serial and returns the edge name.

func (*Store) ClaimDeliveries

func (s *Store) ClaimDeliveries(ctx context.Context, now, lease time.Time, limit int) ([]Delivery, error)

ClaimDeliveries takes up to limit due deliveries and pushes their next attempt to lease, so a slow send is not picked twice.

func (*Store) Close

func (s *Store) Close()

func (*Store) CountRange

func (s *Store) CountRange(ctx context.Context, from, to time.Time) ([]RangeCount, error)

CountRange counts the measurements in [from, to).

func (*Store) CreateSession

func (s *Store) CreateSession(ctx context.Context, sess Session) error

CreateSession stores a session and records the login.

func (*Store) CreateToken

func (s *Store) CreateToken(ctx context.Context, t Token) error

func (*Store) CreateUser

func (s *Store) CreateUser(ctx context.Context, u User) error

CreateUser stores a new user; the username must be lowercase.

func (*Store) DataStats

func (s *Store) DataStats(ctx context.Context) ([]TableStats, int64, error)

DataStats returns the size of each measurement table and of the database.

func (*Store) DeleteRange

func (s *Store) DeleteRange(ctx context.Context, from, to time.Time) ([]RangeCount, error)

DeleteRange deletes the measurements in [from, to), then recomputes the hourly rollup over the range. Latest results and alerts are kept.

A plain DELETE, not drop_chunks: dropped chunks do not invalidate the rollup, which would keep the deleted data.

func (*Store) DeleteRule

func (s *Store) DeleteRule(ctx context.Context, id int64) error

func (*Store) DeleteSession

func (s *Store) DeleteSession(ctx context.Context, tokenHash []byte) error

func (*Store) DeleteTarget

func (s *Store) DeleteTarget(ctx context.Context, id uuid.UUID, now time.Time) error

DeleteTarget soft deletes a target and its checks. Results are kept.

func (*Store) DeleteUser

func (s *Store) DeleteUser(ctx context.Context, id uuid.UUID) error

DeleteUser removes a user and its sessions. Refuses to leave no enabled admin.

func (*Store) DeleteWebhook

func (s *Store) DeleteWebhook(ctx context.Context, id int64) error

DeleteWebhook removes a webhook, its deliveries, and it from the rules.

func (*Store) Edge

func (s *Store) Edge(ctx context.Context, id uuid.UUID) (Edge, error)

func (*Store) EnqueueDelivery

func (s *Store) EnqueueDelivery(ctx context.Context, webhookIDs []int64, event, summary string, payload any, now time.Time) error

EnqueueDelivery queues a payload for each enabled webhook of ids.

func (*Store) Enroll

func (s *Store) Enroll(ctx context.Context, tokenID string, secretHash []byte, now time.Time, issue IssueFunc) (Edge, Certificate, error)

Enroll consumes a token and creates (or takes over) the edge, all or nothing.

func (*Store) EnsureGrafanaLogin

func (s *Store) EnsureGrafanaLogin(ctx context.Context, name, password string) error

EnsureGrafanaLogin creates or updates a login role that may only read the grafana schema (role netprobe_grafana), with short queries and few connections. Needs CREATEROLE.

func (*Store) ExportChecks

func (s *Store) ExportChecks(ctx context.Context, f ExportFilter, fn func(CheckRow) error) error

ExportChecks streams probe results in time order to fn.

func (*Store) ExportTraces

func (s *Store) ExportTraces(ctx context.Context, f ExportFilter, fn func(TraceRow) error) error

ExportTraces streams traceroutes in time order to fn.

func (*Store) FiringAlerts

func (s *Store) FiringAlerts(ctx context.Context) ([]Alert, error)

FiringAlerts returns the alerts not resolved yet, newest first.

func (*Store) Heartbeat

func (s *Store) Heartbeat(ctx context.Context, id uuid.UUID, now time.Time, clockSkew float64) error

func (*Store) InsertResults

func (s *Store) InsertResults(ctx context.Context, edgeID uuid.UUID, results []Result, now time.Time) (int, error)

InsertResults stores results of an edge, ignoring duplicates and unknown versions. Returns how many rows were new.

func (*Store) InsertTraces

func (s *Store) InsertTraces(ctx context.Context, edgeID uuid.UUID, results []TraceResult, now time.Time) (int, error)

InsertTraces stores traceroutes of an edge, ignoring duplicates and unknown versions. Returns how many rows were new.

func (*Store) LatestResults

func (s *Store) LatestResults(ctx context.Context, edgeID, targetID *uuid.UUID) ([]Latest, error)

LatestResults returns the last result per (edge, check), traceroutes included; nil filters are ignored.

func (*Store) ListDeliveries

func (s *Store) ListDeliveries(ctx context.Context, webhookID int64, limit int) ([]Delivery, error)

func (*Store) ListEdges

func (s *Store) ListEdges(ctx context.Context) ([]Edge, error)

func (*Store) ListRules

func (s *Store) ListRules(ctx context.Context) ([]Rule, error)

func (*Store) ListTargets

func (s *Store) ListTargets(ctx context.Context) ([]Target, error)

ListTargets returns live targets with their live checks.

func (*Store) ListUsers

func (s *Store) ListUsers(ctx context.Context) ([]User, error)

func (*Store) ListWebhooks

func (s *Store) ListWebhooks(ctx context.Context) ([]Webhook, error)

func (*Store) LogDelivery

func (s *Store) LogDelivery(ctx context.Context, d Delivery) error

LogDelivery stores a finished delivery (a test), for the history.

func (*Store) ObserveAddress

func (s *Store) ObserveAddress(ctx context.Context, edgeID uuid.UUID, src Source, family string, ip netip.Addr, info NetInfo, now time.Time) (bool, error)

ObserveAddress extends the current period or opens a new one. Returns true on change.

func (*Store) OpenAlert

func (s *Store) OpenAlert(ctx context.Context, a Alert) (int64, error)

OpenAlert stores a new alert; resolved may be set for one-shot events. Returns 0 when the same alert is already firing.

func (*Store) PathChanges

func (s *Store) PathChanges(ctx context.Context, edgeID uuid.UUID, checkID int64, from, to time.Time, limit int) ([]PathChange, error)

PathChanges lists AS path changes over a range, newest first. The first entry of the range is included, as the starting path.

func (*Store) PathChangesSince

func (s *Store) PathChangesSince(ctx context.Context, since time.Time) ([]LatestPathChange, error)

PathChangesSince lists traceroutes whose AS path changed after since.

func (*Store) Ping

func (s *Store) Ping(ctx context.Context) error

Ping checks that the database answers.

func (*Store) PurgeDeliveries

func (s *Store) PurgeDeliveries(ctx context.Context, before time.Time) (int64, error)

PurgeDeliveries deletes finished deliveries created before a time.

func (*Store) PurgeSessions

func (s *Store) PurgeSessions(ctx context.Context, now, idleSince time.Time) (int64, error)

PurgeSessions deletes expired and idle sessions.

func (*Store) RecordAttempt

func (s *Store) RecordAttempt(ctx context.Context, id int64, ok bool, status int, errText string, next *time.Time, now time.Time) error

RecordAttempt stores the outcome of a try. next nil ends the delivery: delivered if ok, else failed.

func (*Store) RecordCertificate

func (s *Store) RecordCertificate(ctx context.Context, edgeID uuid.UUID, c Certificate, now time.Time) error

func (*Store) RefreshRollup

func (s *Store) RefreshRollup(ctx context.Context, from, to time.Time) error

RefreshRollup recomputes the hourly rollup over a range (late results). TimescaleDB refuses two refreshes at once, and its own policy runs in the background: a refused refresh waits for its turn, up to rollupWait.

func (*Store) ReportableVersions

func (s *Store) ReportableVersions(ctx context.Context, edgeID uuid.UUID, versions []int64) (map[int64]Reportable, error)

ReportableVersions keeps, among versions, those of live checks assigned to the edge (its labels match the selector of the check). Others are unknown to it and their results must be dropped.

func (*Store) ResolveAlert

func (s *Store) ResolveAlert(ctx context.Context, id int64, at time.Time) error

func (*Store) ResolvedAlerts

func (s *Store) ResolvedAlerts(ctx context.Context, limit int) ([]Alert, error)

ResolvedAlerts returns the last resolved alerts, newest first.

func (*Store) RevokeEdge

func (s *Store) RevokeEdge(ctx context.Context, id uuid.UUID, now time.Time) error

func (*Store) SaveRule

func (s *Store) SaveRule(ctx context.Context, r Rule, now time.Time) (int64, error)

SaveRule creates (ID 0) or updates a rule. Returns its id.

func (*Store) SaveTarget

func (s *Store) SaveTarget(ctx context.Context, in TargetInput, now time.Time) (uuid.UUID, error)

SaveTarget creates or updates a target with its full list of checks. Checks missing from the input are deleted (soft). Returns the target id.

func (*Store) SaveWebhook

func (s *Store) SaveWebhook(ctx context.Context, w Webhook, now time.Time) (int64, error)

func (*Store) Series

func (s *Store) Series(ctx context.Context, checkID int64, from, to time.Time, bucket time.Duration, rollup bool) (map[uuid.UUID][]Point, error)

Series returns bucketed points per edge for all versions of a check. With rollup set, the hourly aggregate is used (bucket must be >= 1h).

func (*Store) SessionUser

func (s *Store) SessionUser(ctx context.Context, tokenHash []byte, now, idleSince time.Time) (Session, User, error)

SessionUser returns a live session and its enabled user. A session idle since before idleSince, expired, or of a disabled user is not found.

func (*Store) SetPassword

func (s *Store) SetPassword(ctx context.Context, id uuid.UUID, hash string, keep []byte, now time.Time) error

SetPassword changes the password of a user and ends its other sessions (all but the one whose hash is keep).

func (*Store) SetWebhookSecret

func (s *Store) SetWebhookSecret(ctx context.Context, id int64, secret string) error

SaveWebhook creates (ID 0) or updates a webhook. An empty secret keeps the current one on update. SetWebhookSecret replaces the stored secret of a webhook, e.g. to seal it.

func (*Store) TouchSession

func (s *Store) TouchSession(ctx context.Context, tokenHash []byte, now time.Time) error

func (*Store) Trace

func (s *Store) Trace(ctx context.Context, edgeID uuid.UUID, checkID int64, at *time.Time) (*TraceRun, error)

Trace returns the last traceroute of an edge for a check, at or before at when set. Nil when there is none.

func (*Store) UpdateEdgeInfo

func (s *Store) UpdateEdgeInfo(ctx context.Context, id uuid.UUID, info EdgeInfo, now time.Time) error

func (*Store) UpdateUser

func (s *Store) UpdateUser(ctx context.Context, id uuid.UUID, u UserUpdate, now time.Time) error

UpdateUser applies u and ends the sessions of the user when its access changes. Refuses to leave no enabled admin.

func (*Store) User

func (s *Store) User(ctx context.Context, id uuid.UUID) (User, error)

func (*Store) UserByName

func (s *Store) UserByName(ctx context.Context, username string) (User, error)

func (*Store) Webhook

func (s *Store) Webhook(ctx context.Context, id int64) (Webhook, error)

func (*Store) WindowStats

func (s *Store) WindowStats(ctx context.Context, since time.Time) ([]WindowStat, error)

WindowStats sums the results of live checks since a time.

type TableStats

type TableStats struct {
	Table  string
	Rows   int64 // approximate
	Bytes  int64 // with indexes, compressed chunks included
	Chunks int
	Oldest *time.Time
	Newest *time.Time
}

TableStats describes one measurement table.

type Target

type Target struct {
	ID        uuid.UUID
	Name      string
	Host      string
	Labels    map[string]string
	CreatedAt time.Time
	Checks    []Check
}

type TargetInput

type TargetInput struct {
	ID     uuid.UUID
	Name   string
	Host   string
	Labels map[string]string
	Checks []CheckInput
}

TargetInput is a SaveTarget call. uuid.Nil creates a target.

type Token

type Token struct {
	ID         string
	SecretHash []byte
	EdgeName   string
	Labels     map[string]string
	Replace    bool
	CreatedAt  time.Time
	ExpiresAt  time.Time
}

type TraceResult

type TraceResult struct {
	CheckVersionID int64
	Time           time.Time
	ResolvedIP     netip.Addr
	Reached        bool
	Hops           []PathHop // metadata is only used for a path never seen before
	Sent           []int16
	Received       []int16
	RTT            []*float32 // average, nil for a silent hop
	Error          string
}

TraceResult is one traceroute. Hop slices are indexed by TTL - 1.

type TraceRow

type TraceRow struct {
	Time       time.Time
	Edge       string
	Target     string
	Host       string
	Spec       check.Spec
	ResolvedIP netip.Addr
	Reached    bool
	Hops       []PathHop
	ASPath     []int64
	Received   []int16
	RTT        []*float32
	Error      string
}

TraceRow is one traceroute, with its path, for exports.

type TraceRun

type TraceRun struct {
	Time       time.Time
	ResolvedIP netip.Addr
	Reached    bool
	Hops       []PathHop
	ASPath     []int64
	Sent       []int16
	Received   []int16
	RTT        []*float32
	Error      string
}

TraceRun is a stored traceroute with its path.

type TraceSummary

type TraceSummary struct {
	Reached   bool
	Hops      int
	ASPath    []int64
	ChangedAt *time.Time
}

TraceSummary is the traceroute part of a Latest.

type User

type User struct {
	ID                uuid.UUID
	Username          string
	PasswordHash      string
	Role              Role
	CreatedAt         time.Time
	PasswordChangedAt time.Time
	LastLoginAt       *time.Time
	DisabledAt        *time.Time
}

type UserUpdate

type UserUpdate struct {
	Role         *Role
	Disabled     *bool
	PasswordHash *string
}

UserUpdate changes a user; nil fields are kept.

type Webhook

type Webhook struct {
	ID             int64
	Name           string
	URL            string
	Secret         string
	Enabled        bool
	CreatedAt      time.Time
	LastStatus     string
	LastDeliveryAt *time.Time
}

type WindowStat

type WindowStat struct {
	EdgeID   uuid.UUID
	CheckID  int64
	TargetID uuid.UUID
	Sent     int64
	Received int64
	RTTAvg   *float64 // seconds, nil without replies
}

WindowStat sums the results of one edge and check over a window.

Directories

Path Synopsis
Package storetest gives tests a fresh database, created from NETPROBE_TEST_DATABASE_URL and dropped at the end.
Package storetest gives tests a fresh database, created from NETPROBE_TEST_DATABASE_URL and dropped at the end.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL