NOISSUE - Revert alarms Atom migration

This reverts commits afd2665c9c and 45c426c8bf.
This commit is contained in:
dusan
2026-06-29 17:49:57 +02:00
parent 3cf6b2f0f4
commit cdc065a7b2
15 changed files with 1383 additions and 907 deletions
+58 -21
View File
@@ -1,6 +1,6 @@
# Alarms
The Alarms service stores, manages and exposes alarms raised by rules and device activity. It consumes alarm events from the message broker, stores alarms as Atom resources, and provides an HTTP API for listing, viewing, updating, and deleting alarms with full authn/authz, metrics, and tracing support.
The Alarms service stores, manages and exposes alarms raised by rules and device activity. It consumes alarm events from the message broker, persists them to PostgreSQL, and provides an HTTP API for listing, viewing, updating, and deleting alarms with full authn/authz, metrics, and tracing support.
## Configuration
@@ -13,19 +13,29 @@ The service is configured using the following environment variables (values show
| `MG_ALARMS_HTTP_PORT` | HTTP port to bind | `8050` |
| `MG_ALARMS_HTTP_SERVER_CERT` | Path to PEM-encoded HTTPS server certificate | "" |
| `MG_ALARMS_HTTP_SERVER_KEY` | Path to PEM-encoded HTTPS server key | "" |
| `MG_ALARMS_DB_HOST` | PostgreSQL host | `alarms-db` |
| `MG_ALARMS_DB_PORT` | PostgreSQL port | `5432` |
| `MG_ALARMS_DB_USER` | PostgreSQL user | `magistrala` |
| `MG_ALARMS_DB_PASS` | PostgreSQL password | `magistrala` |
| `MG_ALARMS_DB_NAME` | PostgreSQL database name | `alarms` |
| `MG_ALARMS_DB_SSL_MODE` | PostgreSQL SSL mode | `disable` |
| `MG_ALARMS_DB_SSL_CERT` | PostgreSQL SSL client cert | "" |
| `MG_ALARMS_DB_SSL_KEY` | PostgreSQL SSL client key | "" |
| `MG_ALARMS_DB_SSL_ROOT_CERT` | PostgreSQL SSL root cert | "" |
| `MG_ALARMS_INSTANCE_ID` | Instance ID for tracing/health | "" |
| `MG_MESSAGE_BROKER_URL` | Message broker URL for alarm ingestion | `nats://nats:4222` |
| `MG_JAEGER_URL` | Jaeger collector endpoint | `http://jaeger:4318/v1/traces` |
| `MG_JAEGER_TRACE_RATIO` | Trace sampling ratio | `1.0` |
| `ATOM_URL` | Atom HTTP endpoint | `http://atom:8080` |
| `ATOM_SERVICE_TOKEN` | Atom service token used by the alarms service. In Docker Compose this is populated from `MG_ATOM_TOKEN_ALARMS`. | "" |
| `ATOM_JWKS_URL` | Atom JWKS endpoint for JWT verification | `http://atom:8080/.well-known/jwks.json` |
| `ATOM_ADMIN_USERNAME` | Atom admin login for service projections | `atom-admin` |
| `ATOM_ADMIN_SECRET` | Atom admin secret for service projections | `change-me` |
| `ATOM_TIMEOUT` | Atom request timeout | `5s` |
| `MG_ALLOW_UNVERIFIED_USER` | Allow unverified users to access | `true` |
## Features
- **Alarm ingestion**: Consumes alarms from the message broker and persists them as Atom resources with `kind=alarm`.
- **Alarm ingestion**: Consumes alarms from the message broker and persists them to PostgreSQL.
- **Stateful updates**: Updates assignee, acknowledgment, resolution, and metadata fields.
- **Filtering and paging**: Lists alarms by domain, rule, channel, client, subtopic, status, severity, and time range.
- **Observability**: `/metrics` Prometheus endpoint and Jaeger tracing support.
@@ -37,30 +47,50 @@ The service is configured using the following environment variables (values show
1. The message broker publishes alarm events under the `alarms.>` subject.
2. The Alarms consumer decodes the event payload, enriches it with message metadata, validates it, and calls `CreateAlarm`.
3. The repository writes to Atom resources while deduplicating repeated active alarms with the same severity.
3. The repository writes to PostgreSQL while deduplicating repeated active alarms with the same severity.
4. The HTTP API exposes list/view/update/delete operations with authn/authz, metrics, and tracing middleware.
### Components
- **HTTP API**: `alarms/api` exposes REST endpoints and health/metrics handlers.
- **Service layer**: `alarms/service.go` validates requests and coordinates repository operations.
- **Repository**: `alarms/atom_repository.go` stores alarms as Atom resources and applies alarm filtering.
- **Repository**: `alarms/postgres/alarms.go` implements persistence and filtering.
- **Consumer**: `alarms/consumer` processes broker messages and creates alarms.
- **Message broker**: `alarms/brokers` uses NATS JetStream with stream `alarms` and subject `alarms.>`.
- **Atom**: stores alarms as `resources` rows with `kind=alarm`.
- **Migrations**: `alarms/postgres/init.go` defines the alarms schema and indexes.
### Atom Resource Mapping
### Alarms table
The service stores each alarm as an Atom resource:
Defined in `alarms/postgres/init.go`:
| Alarm field | Atom resource |
| --- | --- |
| `id` | `resources.id` and `resources.name` |
| `domain_id` | `resources.tenant_id` |
| `assignee_id` | `resources.owner_id` on create and `attributes.assignee_id` |
| `status` | `attributes.status` and `attributes.alarm_status` |
| `metadata` | `attributes.metadata` |
| alarm-specific fields | `attributes.rule_id`, `attributes.channel_id`, `attributes.client_id`, `attributes.subtopic`, `attributes.measurement`, `attributes.value`, `attributes.unit`, `attributes.threshold`, `attributes.cause`, `attributes.severity`, and lifecycle fields |
| Column | Type | Description |
| --- | --- | --- |
| `id` | `VARCHAR(36)` | Alarm UUID (primary key) |
| `rule_id` | `VARCHAR(36)` | Rule ID that triggered the alarm |
| `domain_id` | `VARCHAR(36)` | Domain ID |
| `channel_id` | `VARCHAR(36)` | Channel ID |
| `subtopic` | `TEXT` | Subtopic associated with the alarm |
| `client_id` | `VARCHAR(36)` | Client ID |
| `measurement` | `TEXT` | Measurement name |
| `value` | `TEXT` | Measured value |
| `unit` | `TEXT` | Measurement unit |
| `threshold` | `TEXT` | Threshold value |
| `cause` | `TEXT` | Cause/description |
| `status` | `SMALLINT` | 0 = active, 1 = cleared |
| `severity` | `SMALLINT` | Severity (0-100) |
| `assignee_id` | `VARCHAR(36)` | Assignee ID |
| `created_at` | `TIMESTAMPTZ` | Creation timestamp |
| `updated_at` | `TIMESTAMPTZ` | Last update timestamp |
| `updated_by` | `VARCHAR(36)` | User who updated |
| `assigned_at` | `TIMESTAMPTZ` | When assigned |
| `assigned_by` | `VARCHAR(36)` | Who assigned |
| `acknowledged_at` | `TIMESTAMPTZ` | When acknowledged |
| `acknowledged_by` | `VARCHAR(36)` | Who acknowledged |
| `resolved_at` | `TIMESTAMPTZ` | When resolved |
| `resolved_by` | `VARCHAR(36)` | Who resolved |
| `metadata` | `JSONB` | Custom metadata |
Index: `idx_alarms_state (domain_id, rule_id, channel_id, subtopic, client_id, measurement, created_at DESC)`
## Deployment
@@ -71,18 +101,25 @@ make alarms
MG_ALARMS_LOG_LEVEL=debug \
MG_ALARMS_HTTP_PORT=8050 \
MG_ALARMS_DB_HOST=localhost \
MG_ALARMS_DB_PORT=5432 \
MG_ALARMS_DB_USER=magistrala \
MG_ALARMS_DB_PASS=magistrala \
MG_ALARMS_DB_NAME=alarms \
MG_MESSAGE_BROKER_URL=nats://localhost:4222 \
ATOM_URL=http://localhost:8080 \
ATOM_SERVICE_TOKEN=<alarms-service-token> \
MG_AUTH_GRPC_URL=localhost:7001 \
MG_AUTH_GRPC_TIMEOUT=300s \
MG_DOMAINS_GRPC_URL=localhost:7003 \
MG_DOMAINS_GRPC_TIMEOUT=300s \
./build/alarms
```
### Docker Compose
The service is available as a Docker container. Refer to [docker/docker-compose.yaml](https://github.com/absmach/magistrala/blob/main/docker/docker-compose.yaml) for the `alarms` service and its environment variables. For a full local stack, make sure Atom, `atom-bootstrap`, NATS, and nginx are also running.
The service is available as a Docker container. Refer to [docker/docker-compose.yaml](https://github.com/absmach/magistrala/blob/main/docker/docker-compose.yaml) for the `alarms` and `alarms-db` services and their environment variables. For a full local stack, make sure the auth, domains, and message broker services are also running.
```bash
docker compose -f docker/docker-compose.yaml up alarms
docker compose -f docker/docker-compose.yaml up alarms alarms-db
```
### Health check
@@ -95,7 +132,7 @@ curl -X GET http://localhost:8050/health \
## Testing
```bash
go test ./alarms ./cmd/alarms
go test ./alarms/...
```
## Usage
+48 -16
View File
@@ -4,17 +4,63 @@
package alarms
import (
"fmt"
"context"
"time"
"github.com/absmach/magistrala/internal/atom"
"github.com/absmach/magistrala/pkg/authn"
)
type atomService struct {
Service
projector atom.Projector
}
func WithAtom(svc Service, projector atom.Projector) Service {
if projector == nil {
return svc
}
return atomService{Service: svc, projector: projector}
}
func (svc atomService) CreateAlarm(ctx context.Context, alarm Alarm) (Alarm, error) {
created, err := svc.Service.CreateAlarm(ctx, alarm)
if err != nil {
return created, err
}
if created.ID == "" {
return created, nil
}
if err := svc.projector.UpsertResource(ctx, alarmProjection(created)); err != nil {
return created, nil
}
return created, nil
}
func (svc atomService) UpdateAlarm(ctx context.Context, session authn.Session, alarm Alarm) (Alarm, error) {
updated, err := svc.Service.UpdateAlarm(ctx, session, alarm)
if err != nil {
return updated, err
}
if err := svc.projector.UpsertResource(ctx, alarmProjection(updated)); err != nil {
return updated, nil
}
return updated, nil
}
func (svc atomService) DeleteAlarm(ctx context.Context, session authn.Session, id string) error {
if err := svc.Service.DeleteAlarm(ctx, session, id); err != nil {
return err
}
_ = svc.projector.DeleteResource(ctx, id)
return nil
}
func alarmProjection(a Alarm) atom.Resource {
res := atom.ResourceFromFields(atom.ObjectFields{
ID: a.ID,
Kind: atom.KindAlarm,
Name: alarmName(a),
Name: a.Cause,
TenantID: a.DomainID,
OwnerID: a.AssigneeID,
Status: a.Status.String(),
@@ -33,7 +79,6 @@ func alarmProjection(a Alarm) atom.Resource {
res.Attributes["unit"] = a.Unit
res.Attributes["threshold"] = a.Threshold
res.Attributes["cause"] = a.Cause
res.Attributes["alarm_status"] = uint8(a.Status)
res.Attributes["assignee_id"] = a.AssigneeID
res.Attributes["assigned_at"] = alarmTimeString(a.AssignedAt)
res.Attributes["assigned_by"] = a.AssignedBy
@@ -44,19 +89,6 @@ func alarmProjection(a Alarm) atom.Resource {
return res
}
func alarmName(a Alarm) string {
if a.Cause != "" && a.Measurement != "" {
return fmt.Sprintf("%s: %s", a.Measurement, a.Cause)
}
if a.Cause != "" {
return a.Cause
}
if a.Measurement != "" {
return fmt.Sprintf("%s alarm", a.Measurement)
}
return a.ID
}
func alarmTimeString(ts time.Time) string {
if ts.IsZero() {
return ""
-531
View File
@@ -1,531 +0,0 @@
// Copyright (c) Abstract Machines
// SPDX-License-Identifier: Apache-2.0
package alarms
import (
"context"
"encoding/json"
"math"
"sort"
"strconv"
"time"
api "github.com/absmach/magistrala/api/http"
"github.com/absmach/magistrala/internal/atom"
"github.com/absmach/magistrala/pkg/errors"
repoerr "github.com/absmach/magistrala/pkg/errors/repository"
)
const atomAlarmListLimit uint64 = 1000
type atomResourceStore interface {
CreateResource(ctx context.Context, resource atom.Resource) (atom.Resource, error)
UpdateResource(ctx context.Context, id string, resource atom.Resource) (atom.Resource, error)
GetResource(ctx context.Context, id string) (atom.Resource, error)
ListResources(ctx context.Context, q atom.Query) (atom.ResourceList, error)
DeleteResource(ctx context.Context, id string) error
}
type atomRepository struct {
store atomResourceStore
}
var _ Repository = (*atomRepository)(nil)
func NewAtomRepository(store atomResourceStore) Repository {
return &atomRepository{store: store}
}
func (repo *atomRepository) CreateAlarm(ctx context.Context, alarm Alarm) (Alarm, error) {
ok, err := repo.shouldCreateAlarm(ctx, alarm)
if err != nil {
return Alarm{}, errors.Wrap(repoerr.ErrCreateEntity, err)
}
if !ok {
return Alarm{}, repoerr.ErrNotFound
}
res, err := repo.store.CreateResource(ctx, alarmProjection(alarm))
if err != nil {
return Alarm{}, atomRepositoryError(repoerr.ErrCreateEntity, err)
}
created, err := alarmFromResource(res)
if err != nil {
return Alarm{}, errors.Wrap(repoerr.ErrCreateEntity, err)
}
return created, nil
}
func (repo *atomRepository) UpdateAlarm(ctx context.Context, alarm Alarm) (Alarm, error) {
current, err := repo.loadAlarm(ctx, alarm.ID, "", repoerr.ErrUpdateEntity)
if err != nil {
return Alarm{}, err
}
updated := mergeAlarmUpdate(current, alarm)
res, err := repo.store.UpdateResource(ctx, updated.ID, alarmProjection(updated))
if err != nil {
return Alarm{}, atomRepositoryError(repoerr.ErrUpdateEntity, err)
}
out, err := alarmFromResource(res)
if err != nil {
return Alarm{}, errors.Wrap(repoerr.ErrUpdateEntity, err)
}
return out, nil
}
func (repo *atomRepository) ViewAlarm(ctx context.Context, alarmID, domainID string) (Alarm, error) {
return repo.loadAlarm(ctx, alarmID, domainID, repoerr.ErrViewEntity)
}
func (repo *atomRepository) loadAlarm(ctx context.Context, alarmID, domainID string, wrapper error) (Alarm, error) {
res, err := repo.store.GetResource(ctx, alarmID)
if err != nil {
return Alarm{}, atomRepositoryError(wrapper, err)
}
if res.Kind != atom.KindAlarm {
return Alarm{}, repoerr.ErrNotFound
}
if domainID != "" && res.TenantID != domainID {
return Alarm{}, repoerr.ErrNotFound
}
alarm, err := alarmFromResource(res)
if err != nil {
return Alarm{}, errors.Wrap(wrapper, err)
}
return alarm, nil
}
func (repo *atomRepository) ListAllAlarms(ctx context.Context, pm PageMetadata) (AlarmsPage, error) {
items, err := repo.listAlarms(ctx, pm.DomainID, alarmPageAttributesContains(pm))
if err != nil {
return AlarmsPage{}, errors.Wrap(repoerr.ErrViewEntity, err)
}
filtered := make([]Alarm, 0, len(items))
for _, alarm := range items {
if matchesAlarmPage(alarm, pm) {
filtered = append(filtered, alarm)
}
}
sortAlarms(filtered, pm)
total := uint64(len(filtered))
start := pm.Offset
if start > total {
start = total
}
end := total
if pm.Limit > 0 && start+pm.Limit < end {
end = start + pm.Limit
}
return AlarmsPage{
Offset: pm.Offset,
Limit: pm.Limit,
Total: total,
Alarms: filtered[start:end],
}, nil
}
func (repo *atomRepository) DeleteAlarm(ctx context.Context, id string) error {
if _, err := repo.loadAlarm(ctx, id, "", repoerr.ErrRemoveEntity); err != nil {
return err
}
if err := repo.store.DeleteResource(ctx, id); err != nil {
return atomRepositoryError(repoerr.ErrRemoveEntity, err)
}
return nil
}
func (repo *atomRepository) shouldCreateAlarm(ctx context.Context, alarm Alarm) (bool, error) {
items, err := repo.listAlarms(ctx, alarm.DomainID, alarmIdentityAttributes(alarm))
if err != nil {
return false, err
}
var latest *Alarm
for i := range items {
item := items[i]
if !sameAlarmState(item, alarm) {
continue
}
if !item.CreatedAt.IsZero() && !alarm.CreatedAt.IsZero() && item.CreatedAt.After(alarm.CreatedAt) {
continue
}
if latest == nil || item.CreatedAt.After(latest.CreatedAt) {
latest = &item
}
}
if latest == nil {
return alarm.Status == ActiveStatus, nil
}
if latest.Status != alarm.Status {
return true, nil
}
return alarm.Status == ActiveStatus && latest.Severity != alarm.Severity, nil
}
func (repo *atomRepository) listAlarms(ctx context.Context, domainID string, attributes atom.Attributes) ([]Alarm, error) {
var out []Alarm
var offset uint64
for {
page, err := repo.store.ListResources(ctx, atom.Query{
Kind: atom.KindAlarm,
TenantID: domainID,
AttributesContains: attributes,
Limit: atomAlarmListLimit,
Offset: offset,
})
if err != nil {
return nil, err
}
for _, res := range page.Items {
if res.Kind != atom.KindAlarm {
continue
}
alarm, err := alarmFromResource(res)
if err != nil {
return nil, err
}
out = append(out, alarm)
}
offset += uint64(len(page.Items))
if len(page.Items) == 0 || offset >= page.Total {
break
}
}
return out, nil
}
func alarmIdentityAttributes(alarm Alarm) atom.Attributes {
attrs := atom.Attributes{}
setAttrIfNotEmpty(attrs, "rule_id", alarm.RuleID)
setAttrIfNotEmpty(attrs, "channel_id", alarm.ChannelID)
setAttrIfNotEmpty(attrs, "client_id", alarm.ClientID)
setAttrIfNotEmpty(attrs, "subtopic", alarm.Subtopic)
setAttrIfNotEmpty(attrs, "measurement", alarm.Measurement)
return attrs
}
func alarmPageAttributesContains(pm PageMetadata) atom.Attributes {
attrs := atom.Attributes{}
setAttrIfNotEmpty(attrs, "rule_id", pm.RuleID)
setAttrIfNotEmpty(attrs, "channel_id", pm.ChannelID)
setAttrIfNotEmpty(attrs, "client_id", pm.ClientID)
setAttrIfNotEmpty(attrs, "subtopic", pm.Subtopic)
setAttrIfNotEmpty(attrs, "measurement", pm.Measurement)
if pm.Status != AllStatus {
attrs["status"] = pm.Status.String()
}
if pm.Severity != math.MaxUint8 {
attrs["severity"] = pm.Severity
}
setAttrIfNotEmpty(attrs, "assignee_id", pm.AssigneeID)
setAttrIfNotEmpty(attrs, "updated_by", pm.UpdatedBy)
setAttrIfNotEmpty(attrs, "assigned_by", pm.AssignedBy)
setAttrIfNotEmpty(attrs, "acknowledged_by", pm.AcknowledgedBy)
setAttrIfNotEmpty(attrs, "resolved_by", pm.ResolvedBy)
return attrs
}
func setAttrIfNotEmpty(attrs atom.Attributes, key, value string) {
if value != "" {
attrs[key] = value
}
}
func mergeAlarmUpdate(current, update Alarm) Alarm {
if update.Status != 0 {
current.Status = update.Status
}
if update.AssigneeID != "" {
current.AssigneeID = update.AssigneeID
}
if !update.AssignedAt.IsZero() {
current.AssignedAt = update.AssignedAt
}
if update.AssignedBy != "" {
current.AssignedBy = update.AssignedBy
}
if update.AcknowledgedBy != "" {
current.AcknowledgedBy = update.AcknowledgedBy
}
if !update.AcknowledgedAt.IsZero() {
current.AcknowledgedAt = update.AcknowledgedAt
}
if update.ResolvedBy != "" {
current.ResolvedBy = update.ResolvedBy
}
if !update.ResolvedAt.IsZero() {
current.ResolvedAt = update.ResolvedAt
}
if update.Metadata != nil {
current.Metadata = update.Metadata
}
current.UpdatedAt = update.UpdatedAt
current.UpdatedBy = update.UpdatedBy
return current
}
func alarmFromResource(res atom.Resource) (Alarm, error) {
attrs := res.Attributes
alarm := Alarm{
ID: res.ID,
DomainID: res.TenantID,
CreatedAt: res.CreatedAt,
UpdatedAt: res.UpdatedAt,
}
var err error
alarm.RuleID = stringAttr(attrs, "rule_id")
alarm.ChannelID = stringAttr(attrs, "channel_id")
alarm.ClientID = stringAttr(attrs, "client_id")
alarm.Subtopic = stringAttr(attrs, "subtopic")
alarm.Measurement = stringAttr(attrs, "measurement")
alarm.Value = stringAttr(attrs, "value")
alarm.Unit = stringAttr(attrs, "unit")
alarm.Threshold = stringAttr(attrs, "threshold")
alarm.Cause = stringAttr(attrs, "cause")
alarm.AssigneeID = stringAttr(attrs, "assignee_id")
alarm.UpdatedBy = stringAttr(attrs, "updated_by")
alarm.AssignedBy = stringAttr(attrs, "assigned_by")
alarm.AcknowledgedBy = stringAttr(attrs, "acknowledged_by")
alarm.ResolvedBy = stringAttr(attrs, "resolved_by")
if alarm.Status, err = statusAttr(attrs); err != nil {
return Alarm{}, err
}
if alarm.Severity, err = uint8Attr(attrs, "severity"); err != nil {
return Alarm{}, err
}
if alarm.CreatedAt, err = timeAttr(attrs, "created_at", alarm.CreatedAt); err != nil {
return Alarm{}, err
}
if alarm.UpdatedAt, err = timeAttr(attrs, "updated_at", alarm.UpdatedAt); err != nil {
return Alarm{}, err
}
if alarm.AssignedAt, err = timeAttr(attrs, "assigned_at", time.Time{}); err != nil {
return Alarm{}, err
}
if alarm.AcknowledgedAt, err = timeAttr(attrs, "acknowledged_at", time.Time{}); err != nil {
return Alarm{}, err
}
if alarm.ResolvedAt, err = timeAttr(attrs, "resolved_at", time.Time{}); err != nil {
return Alarm{}, err
}
alarm.Metadata, err = metadataAttr(attrs)
if err != nil {
return Alarm{}, err
}
return alarm, nil
}
func matchesAlarmPage(a Alarm, pm PageMetadata) bool {
if pm.DomainID != "" && a.DomainID != pm.DomainID {
return false
}
if pm.RuleID != "" && a.RuleID != pm.RuleID {
return false
}
if pm.ChannelID != "" && a.ChannelID != pm.ChannelID {
return false
}
if pm.Subtopic != "" && a.Subtopic != pm.Subtopic {
return false
}
if pm.ClientID != "" && a.ClientID != pm.ClientID {
return false
}
if pm.Measurement != "" && a.Measurement != pm.Measurement {
return false
}
if pm.Status != AllStatus && a.Status != pm.Status {
return false
}
if pm.Severity != math.MaxUint8 && a.Severity != pm.Severity {
return false
}
if pm.AssigneeID != "" && a.AssigneeID != pm.AssigneeID {
return false
}
if pm.UpdatedBy != "" && a.UpdatedBy != pm.UpdatedBy {
return false
}
if pm.ResolvedBy != "" && a.ResolvedBy != pm.ResolvedBy {
return false
}
if pm.AcknowledgedBy != "" && a.AcknowledgedBy != pm.AcknowledgedBy {
return false
}
if pm.AssignedBy != "" && a.AssignedBy != pm.AssignedBy {
return false
}
if !pm.CreatedFrom.IsZero() && a.CreatedAt.Before(pm.CreatedFrom) {
return false
}
if !pm.CreatedTo.IsZero() && a.CreatedAt.After(pm.CreatedTo) {
return false
}
return true
}
func sameAlarmState(a, b Alarm) bool {
return a.DomainID == b.DomainID &&
a.RuleID == b.RuleID &&
a.ChannelID == b.ChannelID &&
a.ClientID == b.ClientID &&
a.Subtopic == b.Subtopic &&
a.Measurement == b.Measurement
}
func sortAlarms(items []Alarm, pm PageMetadata) {
desc := pm.Dir != api.AscDir
sort.SliceStable(items, func(i, j int) bool {
it, jt := alarmOrderTime(items[i], pm), alarmOrderTime(items[j], pm)
if !it.Equal(jt) {
if desc {
return it.After(jt)
}
return it.Before(jt)
}
if desc {
return items[i].ID > items[j].ID
}
return items[i].ID < items[j].ID
})
}
func alarmOrderTime(a Alarm, pm PageMetadata) time.Time {
if pm.Order == api.CreatedAtOrder {
return a.CreatedAt
}
if !a.UpdatedAt.IsZero() {
return a.UpdatedAt
}
return a.CreatedAt
}
func stringAttr(attrs atom.Attributes, key string) string {
if attrs == nil {
return ""
}
if v, ok := attrs[key].(string); ok {
return v
}
return ""
}
func statusAttr(attrs atom.Attributes) (Status, error) {
if v := stringAttr(attrs, "status"); v != "" {
return ToStatus(v)
}
n, ok, err := numericAttr(attrs, "alarm_status")
if err != nil || !ok {
return ActiveStatus, err
}
return Status(n), nil
}
func uint8Attr(attrs atom.Attributes, key string) (uint8, error) {
n, _, err := numericAttr(attrs, key)
return n, err
}
func numericAttr(attrs atom.Attributes, key string) (uint8, bool, error) {
if attrs == nil {
return 0, false, nil
}
switch v := attrs[key].(type) {
case nil:
return 0, false, nil
case uint8:
return v, true, nil
case uint64:
if v > math.MaxUint8 {
return 0, true, repoerr.ErrMalformedEntity
}
return uint8(v), true, nil
case int:
if v < 0 || v > math.MaxUint8 {
return 0, true, repoerr.ErrMalformedEntity
}
return uint8(v), true, nil
case float64:
if v < 0 || v > math.MaxUint8 || v != math.Trunc(v) {
return 0, true, repoerr.ErrMalformedEntity
}
return uint8(v), true, nil
case json.Number:
i, err := strconv.ParseUint(v.String(), 10, 8)
if err != nil {
return 0, true, errors.Wrap(repoerr.ErrMalformedEntity, err)
}
return uint8(i), true, nil
default:
return 0, true, repoerr.ErrMalformedEntity
}
}
func timeAttr(attrs atom.Attributes, key string, fallback time.Time) (time.Time, error) {
if attrs == nil {
return fallback, nil
}
switch v := attrs[key].(type) {
case nil:
return fallback, nil
case string:
if v == "" {
return time.Time{}, nil
}
t, err := time.Parse(time.RFC3339Nano, v)
if err != nil {
return time.Time{}, errors.Wrap(repoerr.ErrMalformedEntity, err)
}
return t, nil
case time.Time:
return v, nil
default:
return time.Time{}, repoerr.ErrMalformedEntity
}
}
func metadataAttr(attrs atom.Attributes) (Metadata, error) {
if attrs == nil || attrs["metadata"] == nil {
return nil, nil
}
switch v := attrs["metadata"].(type) {
case map[string]any:
return Metadata(v), nil
case atom.Attributes:
return Metadata(v), nil
default:
return nil, repoerr.ErrMalformedEntity
}
}
func atomRepositoryError(wrapper, err error) error {
switch {
case atom.IsNotFound(err) || errors.Contains(err, repoerr.ErrNotFound):
return repoerr.ErrNotFound
case atom.IsConflict(err) || errors.Contains(err, repoerr.ErrConflict):
return errors.Wrap(repoerr.ErrConflict, err)
default:
return errors.Wrap(wrapper, err)
}
}
-257
View File
@@ -1,257 +0,0 @@
// Copyright (c) Abstract Machines
// SPDX-License-Identifier: Apache-2.0
package alarms
import (
"context"
"math"
"reflect"
"testing"
"time"
api "github.com/absmach/magistrala/api/http"
"github.com/absmach/magistrala/internal/atom"
"github.com/absmach/magistrala/pkg/errors"
repoerr "github.com/absmach/magistrala/pkg/errors/repository"
)
func TestAtomRepositoryCreateAlarmSuppressesDuplicateActiveSeverity(t *testing.T) {
ts := time.Date(2026, 6, 29, 10, 0, 0, 0, time.UTC)
existing := testAlarm("alarm-1", ts)
store := &alarmAtomStore{
resources: []atom.Resource{alarmProjection(existing)},
}
repo := NewAtomRepository(store)
_, err := repo.CreateAlarm(context.Background(), testAlarm("alarm-2", ts.Add(time.Minute)))
if !errors.Contains(err, repoerr.ErrNotFound) {
t.Fatalf("expected duplicate suppression, got %v", err)
}
if store.created.ID != "" {
t.Fatalf("duplicate alarm should not be created: %#v", store.created)
}
if len(store.queries) != 1 {
t.Fatalf("expected one duplicate lookup query, got %d", len(store.queries))
}
assertAttributesContain(t, store.queries[0].AttributesContains, atom.Attributes{
"rule_id": existing.RuleID,
"channel_id": existing.ChannelID,
"client_id": existing.ClientID,
"subtopic": existing.Subtopic,
"measurement": existing.Measurement,
})
}
func TestAtomRepositoryCreateAlarmCreatesOnSeverityChange(t *testing.T) {
ts := time.Date(2026, 6, 29, 10, 0, 0, 0, time.UTC)
existing := testAlarm("alarm-1", ts)
next := testAlarm("alarm-2", ts.Add(time.Minute))
next.Severity = existing.Severity + 1
store := &alarmAtomStore{
resources: []atom.Resource{alarmProjection(existing)},
}
repo := NewAtomRepository(store)
created, err := repo.CreateAlarm(context.Background(), next)
if err != nil {
t.Fatalf("create alarm: %v", err)
}
if created.ID != next.ID || created.Severity != next.Severity {
t.Fatalf("unexpected created alarm: %#v", created)
}
if store.created.ID != next.ID || store.created.Kind != atom.KindAlarm || store.created.Name != alarmName(next) {
t.Fatalf("unexpected created resource: %#v", store.created)
}
}
func TestAtomRepositoryUpdateAlarmMergesMutableFields(t *testing.T) {
ts := time.Date(2026, 6, 29, 10, 0, 0, 0, time.UTC)
current := testAlarm("alarm-1", ts)
store := &alarmAtomStore{
resources: []atom.Resource{alarmProjection(current)},
}
repo := NewAtomRepository(store)
updatedAt := ts.Add(time.Hour)
got, err := repo.UpdateAlarm(context.Background(), Alarm{
ID: current.ID,
Status: ClearedStatus,
UpdatedAt: updatedAt,
UpdatedBy: "user-1",
ResolvedBy: "user-1",
ResolvedAt: updatedAt,
AcknowledgedBy: "user-2",
Metadata: Metadata{"note": "resolved"},
})
if err != nil {
t.Fatalf("update alarm: %v", err)
}
if got.Status != ClearedStatus || got.RuleID != current.RuleID || got.ResolvedBy != "user-1" || got.AcknowledgedBy != "user-2" {
t.Fatalf("unexpected updated alarm: %#v", got)
}
if got.Metadata["note"] != "resolved" {
t.Fatalf("metadata not updated: %#v", got.Metadata)
}
if store.updated.ID != current.ID || store.updated.Attributes["status"] != Cleared {
t.Fatalf("unexpected updated resource: %#v", store.updated)
}
}
func TestAtomRepositoryListAlarmsFiltersAndSorts(t *testing.T) {
ts := time.Date(2026, 6, 29, 10, 0, 0, 0, time.UTC)
first := testAlarm("alarm-1", ts)
second := testAlarm("alarm-2", ts.Add(time.Hour))
second.Severity = 20
otherChannel := testAlarm("alarm-3", ts.Add(2*time.Hour))
otherChannel.ChannelID = "other-channel"
store := &alarmAtomStore{
resources: []atom.Resource{
alarmProjection(first),
alarmProjection(second),
alarmProjection(otherChannel),
},
}
repo := NewAtomRepository(store)
page, err := repo.ListAllAlarms(context.Background(), PageMetadata{
DomainID: "domain-1",
ChannelID: "channel-1",
Status: AllStatus,
Severity: math.MaxUint8,
Offset: 0,
Limit: 10,
Order: api.CreatedAtOrder,
Dir: api.DescDir,
})
if err != nil {
t.Fatalf("list alarms: %v", err)
}
if page.Total != 2 || len(page.Alarms) != 2 {
t.Fatalf("unexpected page: %#v", page)
}
if page.Alarms[0].ID != second.ID || page.Alarms[1].ID != first.ID {
t.Fatalf("unexpected order: %#v", page.Alarms)
}
if len(store.queries) != 1 {
t.Fatalf("expected one list query, got %d", len(store.queries))
}
assertAttributesContain(t, store.queries[0].AttributesContains, atom.Attributes{
"channel_id": "channel-1",
})
if _, ok := store.queries[0].AttributesContains["status"]; ok {
t.Fatalf("status should not be pushed for all-status query: %#v", store.queries[0].AttributesContains)
}
if _, ok := store.queries[0].AttributesContains["severity"]; ok {
t.Fatalf("severity should not be pushed for sentinel query: %#v", store.queries[0].AttributesContains)
}
}
func testAlarm(id string, createdAt time.Time) Alarm {
return Alarm{
ID: id,
RuleID: "rule-1",
DomainID: "domain-1",
ChannelID: "channel-1",
ClientID: "client-1",
Subtopic: "temperature",
Status: ActiveStatus,
Measurement: "temperature",
Value: "91.2",
Unit: "C",
Threshold: "80",
Cause: "high temperature",
Severity: 90,
CreatedAt: createdAt,
}
}
type alarmAtomStore struct {
resources []atom.Resource
created atom.Resource
updated atom.Resource
deleted string
queries []atom.Query
}
func (s *alarmAtomStore) CreateResource(_ context.Context, resource atom.Resource) (atom.Resource, error) {
s.created = resource
s.resources = append(s.resources, resource)
return resource, nil
}
func (s *alarmAtomStore) UpdateResource(_ context.Context, id string, resource atom.Resource) (atom.Resource, error) {
s.updated = resource
for i, res := range s.resources {
if res.ID == id {
s.resources[i] = resource
return resource, nil
}
}
return atom.Resource{}, repoerr.ErrNotFound
}
func (s *alarmAtomStore) GetResource(_ context.Context, id string) (atom.Resource, error) {
for _, res := range s.resources {
if res.ID == id {
return res, nil
}
}
return atom.Resource{}, repoerr.ErrNotFound
}
func (s *alarmAtomStore) ListResources(_ context.Context, q atom.Query) (atom.ResourceList, error) {
s.queries = append(s.queries, q)
var items []atom.Resource
for _, res := range s.resources {
if q.Kind != "" && res.Kind != q.Kind {
continue
}
if q.TenantID != "" && res.TenantID != q.TenantID {
continue
}
if !resourceAttributesContain(res.Attributes, q.AttributesContains) {
continue
}
items = append(items, res)
}
total := uint64(len(items))
start := q.Offset
if start > total {
start = total
}
end := total
if q.Limit > 0 && start+q.Limit < end {
end = start + q.Limit
}
return atom.ResourceList{Items: items[start:end], Total: total}, nil
}
func resourceAttributesContain(attrs, contains atom.Attributes) bool {
for key, want := range contains {
if !reflect.DeepEqual(attrs[key], want) {
return false
}
}
return true
}
func assertAttributesContain(t *testing.T, got, want atom.Attributes) {
t.Helper()
for key, wantValue := range want {
if !reflect.DeepEqual(got[key], wantValue) {
t.Fatalf("attribute %q mismatch: got %#v want %#v in %#v", key, got[key], wantValue, got)
}
}
}
func (s *alarmAtomStore) DeleteResource(_ context.Context, id string) error {
s.deleted = id
for i, res := range s.resources {
if res.ID == id {
s.resources = append(s.resources[:i], s.resources[i+1:]...)
return nil
}
}
return repoerr.ErrNotFound
}
+65 -51
View File
@@ -4,66 +4,80 @@
package alarms
import (
"context"
"testing"
"github.com/absmach/magistrala/internal/atom"
"github.com/absmach/magistrala/pkg/authn"
)
func TestAlarmProjectionBuildsAtomResource(t *testing.T) {
resource := alarmProjection(Alarm{
ID: "alarm-1",
RuleID: "rule-1",
DomainID: "domain-1",
ChannelID: "channel-1",
ClientID: "client-1",
Cause: "high temperature",
Measurement: "temperature",
Value: "92.4",
Unit: "C",
Threshold: "80",
Severity: 90,
Status: ActiveStatus,
})
func TestAtomServiceCreateAlarmProjectsCreatedAlarm(t *testing.T) {
projector := &alarmProjector{}
svc := WithAtom(alarmService{
create: Alarm{
ID: "alarm-1",
RuleID: "rule-1",
DomainID: "domain-1",
ChannelID: "channel-1",
ClientID: "client-1",
Cause: "high temperature",
Measurement: "temperature",
Value: "92.4",
Unit: "C",
Threshold: "80",
Severity: 90,
Status: ActiveStatus,
},
}, projector)
if resource.ID != "alarm-1" || resource.Kind != atom.KindAlarm || resource.Name != "temperature: high temperature" {
t.Fatalf("unexpected projection: %#v", resource)
created, err := svc.CreateAlarm(context.Background(), Alarm{RuleID: "rule-1"})
if err != nil {
t.Fatalf("create alarm: %v", err)
}
if resource.Attributes["rule_id"] != "rule-1" {
t.Fatalf("missing rule projection: %#v", resource.Attributes)
if created.ID != "alarm-1" {
t.Fatalf("unexpected created alarm: %#v", created)
}
if resource.Attributes["value"] != "92.4" || resource.Attributes["threshold"] != "80" {
t.Fatalf("missing alarm value projection: %#v", resource.Attributes)
if projector.resource.ID != "alarm-1" || projector.resource.Kind != atom.KindAlarm {
t.Fatalf("unexpected projection: %#v", projector.resource)
}
if projector.resource.Attributes["rule_id"] != "rule-1" {
t.Fatalf("missing rule projection: %#v", projector.resource.Attributes)
}
if projector.resource.Attributes["value"] != "92.4" || projector.resource.Attributes["threshold"] != "80" {
t.Fatalf("missing alarm value projection: %#v", projector.resource.Attributes)
}
}
func TestAlarmNameFallbacks(t *testing.T) {
cases := []struct {
name string
alarm Alarm
want string
}{
{
name: "cause only",
alarm: Alarm{ID: "alarm-1", Cause: "high temperature"},
want: "high temperature",
},
{
name: "measurement only",
alarm: Alarm{ID: "alarm-1", Measurement: "temperature"},
want: "temperature alarm",
},
{
name: "id fallback",
alarm: Alarm{ID: "alarm-1"},
want: "alarm-1",
},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
if got := alarmName(tc.alarm); got != tc.want {
t.Fatalf("unexpected alarm name: got %q want %q", got, tc.want)
}
})
}
type alarmService struct {
create Alarm
}
func (svc alarmService) CreateAlarm(context.Context, Alarm) (Alarm, error) {
return svc.create, nil
}
func (svc alarmService) UpdateAlarm(context.Context, authn.Session, Alarm) (Alarm, error) {
return Alarm{}, nil
}
func (svc alarmService) ViewAlarm(context.Context, authn.Session, string) (Alarm, error) {
return Alarm{}, nil
}
func (svc alarmService) ListAlarms(context.Context, authn.Session, PageMetadata) (AlarmsPage, error) {
return AlarmsPage{}, nil
}
func (svc alarmService) DeleteAlarm(context.Context, authn.Session, string) error {
return nil
}
type alarmProjector struct {
atom.Projector
resource atom.Resource
}
func (p *alarmProjector) UpsertResource(_ context.Context, resource atom.Resource) error {
p.resource = resource
return nil
}
+507
View File
@@ -0,0 +1,507 @@
// Copyright (c) Abstract Machines
// SPDX-License-Identifier: Apache-2.0
package postgres
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"math"
"strings"
"time"
"github.com/absmach/magistrala/alarms"
api "github.com/absmach/magistrala/api/http"
"github.com/absmach/magistrala/pkg/errors"
repoerr "github.com/absmach/magistrala/pkg/errors/repository"
"github.com/absmach/magistrala/pkg/postgres"
"github.com/jmoiron/sqlx"
)
const alarmColumns = `alarms.id, alarms.rule_id, alarms.domain_id, alarms.channel_id, alarms.client_id, alarms.subtopic, alarms.measurement, alarms.value, alarms.unit,
alarms.threshold, alarms.cause, alarms.status, alarms.severity, alarms.assignee_id, alarms.created_at, alarms.updated_at, alarms.updated_by, alarms.assigned_at,
alarms.assigned_by, alarms.acknowledged_at, alarms.acknowledged_by, alarms.resolved_at, alarms.resolved_by, alarms.metadata`
type repository struct {
db *sqlx.DB
}
var _ alarms.Repository = (*repository)(nil)
func NewAlarmsRepo(db *sqlx.DB) alarms.Repository {
return &repository{db: db}
}
func (r *repository) CreateAlarm(ctx context.Context, alarm alarms.Alarm) (alarms.Alarm, error) {
query := `
WITH existing AS (
SELECT status, severity
FROM alarms
WHERE domain_id = :domain_id
AND rule_id = :rule_id
AND channel_id = :channel_id
AND client_id = :client_id
AND subtopic = :subtopic
AND measurement = :measurement
AND created_at <= :created_at
ORDER BY created_at DESC
LIMIT 1
)
INSERT INTO alarms (
id, rule_id, domain_id, channel_id, client_id, subtopic, measurement,
value, unit, threshold, cause, status, severity, assignee_id,
created_at, updated_at, updated_by, assigned_at, assigned_by,
acknowledged_at, acknowledged_by, resolved_at, resolved_by, metadata
)
SELECT
:id, :rule_id, :domain_id, :channel_id, :client_id, :subtopic, :measurement,
:value, :unit, :threshold, :cause, :status, :severity, :assignee_id,
:created_at, :updated_at, :updated_by, :assigned_at, :assigned_by,
:acknowledged_at, :acknowledged_by, :resolved_at, :resolved_by, :metadata
WHERE (
EXISTS (
SELECT 1 FROM existing
WHERE existing.status IS DISTINCT FROM :status
OR (:status = 0 AND existing.status = 0 AND existing.severity IS DISTINCT FROM :severity)
)
OR (
NOT EXISTS (SELECT 1 FROM existing) AND :status = 0
)
)
RETURNING
id, rule_id, domain_id, channel_id, client_id, subtopic, measurement,
value, unit, threshold, cause, status, severity, created_at,
assignee_id, updated_at, updated_by, assigned_at, assigned_by,
acknowledged_at, acknowledged_by, resolved_at, resolved_by, metadata
;
`
dba, err := toDBAlarm(alarm)
if err != nil {
return alarms.Alarm{}, errors.Wrap(repoerr.ErrCreateEntity, err)
}
row, err := r.db.NamedQueryContext(ctx, query, dba)
if err != nil {
return alarms.Alarm{}, postgres.HandleError(repoerr.ErrCreateEntity, err)
}
defer row.Close()
if !row.Next() {
return alarms.Alarm{}, repoerr.ErrNotFound
}
dba = dbAlarm{}
if err := row.StructScan(&dba); err != nil {
return alarms.Alarm{}, errors.Wrap(repoerr.ErrCreateEntity, err)
}
return toAlarm(dba)
}
func (r *repository) UpdateAlarm(ctx context.Context, alarm alarms.Alarm) (alarms.Alarm, error) {
var query []string
var upq string
if alarm.Status != 0 {
query = append(query, "status = :status,")
}
if alarm.AssigneeID != "" {
query = append(query, "assignee_id = :assignee_id,")
}
if !alarm.AssignedAt.IsZero() {
query = append(query, "assigned_at = :assigned_at,")
}
if alarm.AssignedBy != "" {
query = append(query, "assigned_by = :assigned_by,")
}
if alarm.AcknowledgedBy != "" {
query = append(query, "acknowledged_by = :acknowledged_by,")
}
if !alarm.AcknowledgedAt.IsZero() {
query = append(query, "acknowledged_at = :acknowledged_at,")
}
if alarm.ResolvedBy != "" {
query = append(query, "resolved_by = :resolved_by,")
}
if !alarm.ResolvedAt.IsZero() {
query = append(query, "resolved_at = :resolved_at,")
}
if alarm.Metadata != nil {
query = append(query, "metadata = :metadata,")
}
if len(query) > 0 {
upq = strings.Join(query, " ")
}
q := fmt.Sprintf(`UPDATE alarms SET %s updated_by = :updated_by, updated_at = :updated_at WHERE id = :id
RETURNING id, rule_id, domain_id, channel_id, client_id, subtopic, measurement, value, unit, threshold,
cause, status, severity, assignee_id, assigned_at, assigned_by, acknowledged_at, acknowledged_by,
resolved_by, resolved_at, metadata, created_at, updated_by, updated_at;`, upq)
dba, err := toDBAlarm(alarm)
if err != nil {
return alarms.Alarm{}, errors.Wrap(repoerr.ErrUpdateEntity, err)
}
row, err := r.db.NamedQueryContext(ctx, q, dba)
if err != nil {
return alarms.Alarm{}, postgres.HandleError(repoerr.ErrUpdateEntity, err)
}
defer row.Close()
if !row.Next() {
return alarms.Alarm{}, repoerr.ErrNotFound
}
dba = dbAlarm{}
if err := row.StructScan(&dba); err != nil {
return alarms.Alarm{}, errors.Wrap(repoerr.ErrUpdateEntity, err)
}
return toAlarm(dba)
}
func (r *repository) ViewAlarm(ctx context.Context, alarmID, domainID string) (alarms.Alarm, error) {
query := `SELECT * FROM alarms WHERE id = :id AND domain_id = :domain_id;`
row, err := r.db.NamedQueryContext(ctx, query, map[string]any{
"id": alarmID, "domain_id": domainID,
})
if err != nil {
return alarms.Alarm{}, postgres.HandleError(repoerr.ErrViewEntity, err)
}
defer row.Close()
if !row.Next() {
return alarms.Alarm{}, repoerr.ErrNotFound
}
dba := dbAlarm{}
if err := row.StructScan(&dba); err != nil {
return alarms.Alarm{}, errors.Wrap(repoerr.ErrViewEntity, err)
}
alarm, err := toAlarm(dba)
if err != nil {
return alarms.Alarm{}, errors.Wrap(repoerr.ErrViewEntity, err)
}
return alarm, nil
}
func (r *repository) ListAllAlarms(ctx context.Context, pm alarms.PageMetadata) (alarms.AlarmsPage, error) {
query, err := pageQuery(pm)
if err != nil {
return alarms.AlarmsPage{}, errors.Wrap(repoerr.ErrViewEntity, err)
}
comQuery := fmt.Sprintf(`SELECT %s FROM alarms %s`, alarmColumns, query)
return r.alarmsPage(ctx, comQuery, pm)
}
func (r *repository) alarmsPage(ctx context.Context, comQuery string, pm alarms.PageMetadata) (alarms.AlarmsPage, error) {
dir := api.DescDir
if pm.Dir == api.AscDir {
dir = api.AscDir
}
var orderClause string
switch pm.Order {
case api.CreatedAtOrder:
orderClause = fmt.Sprintf("ORDER BY created_at %s, id %s", dir, dir)
default:
orderClause = fmt.Sprintf("ORDER BY COALESCE(updated_at, created_at) %s, id %s", dir, dir)
}
q := fmt.Sprintf(`SELECT * FROM (%s) AS sub_query %s LIMIT :limit OFFSET :offset;`, comQuery, orderClause)
cq := fmt.Sprintf(`SELECT COUNT(*) AS total_count FROM (%s) AS sub_query;`, comQuery)
rows, err := r.db.NamedQueryContext(ctx, q, pm)
if err != nil {
return alarms.AlarmsPage{}, errors.Wrap(repoerr.ErrViewEntity, err)
}
defer rows.Close()
var items []alarms.Alarm
for rows.Next() {
dba := dbAlarm{}
if err := rows.StructScan(&dba); err != nil {
return alarms.AlarmsPage{}, errors.Wrap(repoerr.ErrViewEntity, err)
}
a, err := toAlarm(dba)
if err != nil {
return alarms.AlarmsPage{}, err
}
items = append(items, a)
}
total, err := postgres.Total(ctx, r.db, cq, pm)
if err != nil {
return alarms.AlarmsPage{}, errors.Wrap(repoerr.ErrViewEntity, err)
}
return alarms.AlarmsPage{
Total: total,
Offset: pm.Offset,
Limit: pm.Limit,
Alarms: items,
}, nil
}
func (r *repository) DeleteAlarm(ctx context.Context, id string) error {
query := `DELETE FROM alarms WHERE id = :id;`
result, err := r.db.NamedExecContext(ctx, query, map[string]any{"id": id})
if err != nil {
return errors.Wrap(repoerr.ErrRemoveEntity, err)
}
rowsAffected, err := result.RowsAffected()
if err != nil {
return errors.Wrap(repoerr.ErrRemoveEntity, err)
}
if rowsAffected == 0 {
return repoerr.ErrNotFound
}
return nil
}
type dbAlarm struct {
ID string `db:"id"`
RuleID string `db:"rule_id"`
DomainID string `db:"domain_id"`
ChannelID string `db:"channel_id"`
ClientID string `db:"client_id"`
Subtopic string `db:"subtopic"`
Measurement string `db:"measurement"`
Value string `db:"value"`
Unit string `db:"unit"`
Cause string `db:"cause"`
Threshold string `db:"threshold"`
Status alarms.Status `db:"status"`
Severity uint8 `db:"severity"`
AssigneeID string `db:"assignee_id"`
CreatedAt time.Time `db:"created_at"`
UpdatedAt sql.NullTime `db:"updated_at,omitempty"`
UpdatedBy *string `db:"updated_by,omitempty"`
AssignedAt sql.NullTime `db:"assigned_at,omitempty"`
AssignedBy *string `db:"assigned_by,omitempty"`
AcknowledgedAt sql.NullTime `db:"acknowledged_at,omitempty"`
AcknowledgedBy *string `db:"acknowledged_by,omitempty"`
ResolvedAt sql.NullTime `db:"resolved_at,omitempty"`
ResolvedBy *string `db:"resolved_by,omitempty"`
Metadata []byte `db:"metadata,omitempty"`
}
func toDBAlarm(a alarms.Alarm) (dbAlarm, error) {
if a.CreatedAt.IsZero() {
a.CreatedAt = time.Now()
}
var updatedBy *string
if a.UpdatedBy != "" {
updatedBy = &a.UpdatedBy
}
var updatedAt sql.NullTime
if a.UpdatedAt != (time.Time{}) {
updatedAt = sql.NullTime{Time: a.UpdatedAt, Valid: true}
}
var acknowledgedBy *string
if a.AcknowledgedBy != "" {
acknowledgedBy = &a.AcknowledgedBy
}
var acknowledgedAt sql.NullTime
if a.AcknowledgedAt != (time.Time{}) {
acknowledgedAt = sql.NullTime{Time: a.AcknowledgedAt, Valid: true}
}
var resolvedBy *string
if a.ResolvedBy != "" {
resolvedBy = &a.ResolvedBy
}
var resolvedAt sql.NullTime
if a.ResolvedAt != (time.Time{}) {
resolvedAt = sql.NullTime{Time: a.ResolvedAt, Valid: true}
}
var assignedBy *string
if a.AssignedBy != "" {
assignedBy = &a.AssignedBy
}
var assignedAt sql.NullTime
if a.AssignedAt != (time.Time{}) {
assignedAt = sql.NullTime{Time: a.AssignedAt, Valid: true}
}
metadata := []byte("{}")
if len(a.Metadata) > 0 {
b, err := json.Marshal(a.Metadata)
if err != nil {
return dbAlarm{}, errors.Wrap(repoerr.ErrMalformedEntity, err)
}
metadata = b
}
return dbAlarm{
ID: a.ID,
RuleID: a.RuleID,
DomainID: a.DomainID,
ChannelID: a.ChannelID,
ClientID: a.ClientID,
Subtopic: a.Subtopic,
Measurement: a.Measurement,
Value: a.Value,
Unit: a.Unit,
Cause: a.Cause,
Threshold: a.Threshold,
Status: a.Status,
Severity: a.Severity,
AssigneeID: a.AssigneeID,
CreatedAt: a.CreatedAt,
UpdatedAt: updatedAt,
UpdatedBy: updatedBy,
AssignedAt: assignedAt,
AssignedBy: assignedBy,
AcknowledgedAt: acknowledgedAt,
AcknowledgedBy: acknowledgedBy,
ResolvedAt: resolvedAt,
ResolvedBy: resolvedBy,
Metadata: metadata,
}, nil
}
func toAlarm(dbr dbAlarm) (alarms.Alarm, error) {
var updatedBy string
if dbr.UpdatedBy != nil {
updatedBy = *dbr.UpdatedBy
}
var updatedAt time.Time
if dbr.UpdatedAt.Valid {
updatedAt = dbr.UpdatedAt.Time
}
var assignedBy string
if dbr.AssignedBy != nil {
assignedBy = *dbr.AssignedBy
}
var assignedAt time.Time
if dbr.AssignedAt.Valid {
assignedAt = dbr.AssignedAt.Time
}
var acknowledgedBy string
if dbr.AcknowledgedBy != nil {
acknowledgedBy = *dbr.AcknowledgedBy
}
var acknowledgedAt time.Time
if dbr.AcknowledgedAt.Valid {
acknowledgedAt = dbr.AcknowledgedAt.Time
}
var resolvedBy string
if dbr.ResolvedBy != nil {
resolvedBy = *dbr.ResolvedBy
}
var resolvedAt time.Time
if dbr.ResolvedAt.Valid {
resolvedAt = dbr.ResolvedAt.Time
}
var metadata map[string]any
if len(dbr.Metadata) > 0 {
err := json.Unmarshal(dbr.Metadata, &metadata)
if err != nil {
return alarms.Alarm{}, errors.Wrap(repoerr.ErrMalformedEntity, err)
}
}
return alarms.Alarm{
ID: dbr.ID,
RuleID: dbr.RuleID,
DomainID: dbr.DomainID,
ChannelID: dbr.ChannelID,
ClientID: dbr.ClientID,
Subtopic: dbr.Subtopic,
Measurement: dbr.Measurement,
Value: dbr.Value,
Unit: dbr.Unit,
Threshold: dbr.Threshold,
Cause: dbr.Cause,
Status: dbr.Status,
Severity: dbr.Severity,
AssigneeID: dbr.AssigneeID,
CreatedAt: dbr.CreatedAt,
UpdatedAt: updatedAt,
UpdatedBy: updatedBy,
AssignedAt: assignedAt,
AssignedBy: assignedBy,
AcknowledgedAt: acknowledgedAt,
AcknowledgedBy: acknowledgedBy,
ResolvedAt: resolvedAt,
ResolvedBy: resolvedBy,
Metadata: metadata,
}, nil
}
func pageQuery(pm alarms.PageMetadata) (string, error) {
query := pageQueryConditions(pm)
var emq string
if len(query) > 0 {
emq = fmt.Sprintf("WHERE %s", strings.Join(query, " AND "))
}
return emq, nil
}
func pageQueryConditions(pm alarms.PageMetadata) []string {
var query []string
if pm.DomainID != "" {
query = append(query, "alarms.domain_id = :domain_id")
}
if pm.RuleID != "" {
query = append(query, "alarms.rule_id = :rule_id")
}
if pm.ChannelID != "" {
query = append(query, "alarms.channel_id = :channel_id")
}
if pm.Subtopic != "" {
query = append(query, "alarms.subtopic = :subtopic")
}
if pm.ClientID != "" {
query = append(query, "alarms.client_id = :client_id")
}
if pm.Measurement != "" {
query = append(query, "alarms.measurement = :measurement")
}
if pm.Status != alarms.AllStatus {
query = append(query, "alarms.status = :status")
}
if pm.Severity != math.MaxUint8 {
query = append(query, "alarms.severity = :severity")
}
if pm.AssigneeID != "" {
query = append(query, "alarms.assignee_id = :assignee_id")
}
if pm.UpdatedBy != "" {
query = append(query, "alarms.updated_by = :updated_by")
}
if pm.ResolvedBy != "" {
query = append(query, "alarms.resolved_by = :resolved_by")
}
if pm.AcknowledgedBy != "" {
query = append(query, "alarms.acknowledged_by = :acknowledged_by")
}
if pm.AssignedBy != "" {
query = append(query, "alarms.assigned_by = :assigned_by")
}
if !pm.CreatedFrom.IsZero() {
query = append(query, "alarms.created_at >= :created_from")
}
if !pm.CreatedTo.IsZero() {
query = append(query, "alarms.created_at <= :created_to")
}
return query
}
+481
View File
@@ -0,0 +1,481 @@
// Copyright (c) Abstract Machines
// SPDX-License-Identifier: Apache-2.0
package postgres_test
import (
"context"
"fmt"
"strings"
"testing"
"time"
"github.com/0x6flab/namegenerator"
"github.com/absmach/magistrala/alarms"
"github.com/absmach/magistrala/alarms/postgres"
"github.com/absmach/magistrala/pkg/errors"
repoerr "github.com/absmach/magistrala/pkg/errors/repository"
"github.com/absmach/magistrala/pkg/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
var (
namegen = namegenerator.NewGenerator()
idProvider = uuid.New()
)
func TestCreateAlarm(t *testing.T) {
t.Cleanup(func() {
_, err := db.Exec("DELETE FROM alarms")
require.Nil(t, err, fmt.Sprintf("clean alarms unexpected error: %s", err))
})
repo := postgres.NewAlarmsRepo(db)
alarm := alarms.Alarm{
ID: generateUUID(t),
RuleID: generateUUID(t),
DomainID: generateUUID(t),
ChannelID: generateUUID(t),
ClientID: generateUUID(t),
Subtopic: namegen.Generate(),
Measurement: namegen.Generate(),
Value: namegen.Generate(),
Unit: namegen.Generate(),
Threshold: namegen.Generate(),
Cause: namegen.Generate(),
Status: 0,
AssigneeID: generateUUID(t),
CreatedAt: time.Now().UTC(),
Metadata: map[string]any{
"key": "value",
},
}
cases := []struct {
desc string
alarm alarms.Alarm
err error
}{
{
desc: "valid alarm",
alarm: alarm,
err: nil,
},
{
desc: "duplicate alarm",
alarm: alarm,
err: repoerr.ErrNotFound,
},
{
desc: "missing rule id",
alarm: alarms.Alarm{
ID: generateUUID(t),
DomainID: generateUUID(t),
ChannelID: generateUUID(t),
ClientID: generateUUID(t),
Subtopic: namegen.Generate(),
Measurement: namegen.Generate(),
Value: namegen.Generate(),
Unit: namegen.Generate(),
Threshold: namegen.Generate(),
Cause: namegen.Generate(),
Status: 0,
AssigneeID: generateUUID(t),
CreatedAt: time.Now().UTC(),
Metadata: map[string]any{
"key": "value",
},
},
err: repoerr.ErrCreateEntity,
},
{
desc: "invalid alarm",
alarm: alarms.Alarm{
ID: generateUUID(t),
DomainID: generateUUID(t),
ChannelID: generateUUID(t),
ClientID: generateUUID(t),
Subtopic: namegen.Generate(),
Measurement: namegen.Generate(),
Value: namegen.Generate(),
Unit: namegen.Generate(),
Threshold: namegen.Generate(),
Cause: namegen.Generate(),
Status: 0,
AssigneeID: generateUUID(t),
CreatedAt: time.Now().UTC(),
Metadata: map[string]any{
"key": make(chan int),
},
},
err: repoerr.ErrCreateEntity,
},
{
desc: "empty alarm",
alarm: alarms.Alarm{},
err: repoerr.ErrCreateEntity,
},
}
for _, tc := range cases {
t.Run(tc.desc, func(t *testing.T) {
alarm, err := repo.CreateAlarm(context.Background(), tc.alarm)
if tc.err != nil {
assert.True(t, errors.Contains(err, tc.err), fmt.Sprintf("%s: expected %s got %s\n", tc.desc, tc.err, err))
return
}
assert.Nil(t, err, fmt.Sprintf("unexpected error: %s", err))
assert.NotEmpty(t, alarm.ID)
assert.Equal(t, tc.alarm.RuleID, alarm.RuleID)
assert.Equal(t, tc.alarm.Measurement, alarm.Measurement)
assert.Equal(t, tc.alarm.Value, alarm.Value)
assert.Equal(t, tc.alarm.Unit, alarm.Unit)
assert.Equal(t, tc.alarm.Cause, alarm.Cause)
assert.Equal(t, tc.alarm.Status, alarm.Status)
assert.Equal(t, tc.alarm.DomainID, alarm.DomainID)
assert.Equal(t, tc.alarm.AssigneeID, alarm.AssigneeID)
assert.Equal(t, tc.alarm.Metadata, alarm.Metadata)
})
}
}
func TestUpdateAlarm(t *testing.T) {
t.Cleanup(func() {
_, err := db.Exec("DELETE FROM alarms")
require.Nil(t, err, fmt.Sprintf("clean alarms unexpected error: %s", err))
})
repo := postgres.NewAlarmsRepo(db)
alarm := alarms.Alarm{
ID: generateUUID(t),
RuleID: generateUUID(t),
DomainID: generateUUID(t),
ChannelID: generateUUID(t),
ClientID: generateUUID(t),
Measurement: namegen.Generate(),
Value: namegen.Generate(),
Unit: namegen.Generate(),
Threshold: namegen.Generate(),
Cause: namegen.Generate(),
Status: 0,
AssigneeID: generateUUID(t),
CreatedAt: time.Now().UTC(),
Metadata: map[string]any{
"key": "value",
},
}
alarm, err := repo.CreateAlarm(context.Background(), alarm)
require.Nil(t, err, fmt.Sprintf("unexpected error: %s", err))
cases := []struct {
desc string
alarm alarms.Alarm
err error
}{
{
desc: "valid alarm",
alarm: alarms.Alarm{
ID: alarm.ID,
Status: alarms.ClearedStatus,
DomainID: alarm.DomainID,
AssigneeID: generateUUID(t),
AssignedBy: generateUUID(t),
AssignedAt: time.Now().UTC(),
AcknowledgedBy: generateUUID(t),
AcknowledgedAt: time.Now().UTC(),
CreatedAt: alarm.CreatedAt,
UpdatedAt: time.Now().UTC(),
UpdatedBy: generateUUID(t),
ResolvedAt: time.Now().UTC(),
ResolvedBy: generateUUID(t),
Metadata: map[string]any{
"key": "value",
},
},
err: nil,
},
{
desc: "non existing alarm",
alarm: alarms.Alarm{
ID: generateUUID(t),
},
err: repoerr.ErrNotFound,
},
{
desc: "invalid alarm",
alarm: alarms.Alarm{
ID: alarm.ID,
RuleID: generateUUID(t),
Status: 0,
DomainID: generateUUID(t),
AssigneeID: strings.Repeat("a", 40),
CreatedAt: time.Now().UTC(),
Metadata: map[string]any{
"key": "value",
},
},
err: repoerr.ErrMalformedEntity,
},
{
desc: "empty alarm",
alarm: alarms.Alarm{},
err: repoerr.ErrNotFound,
},
}
for _, tc := range cases {
t.Run(tc.desc, func(t *testing.T) {
alarm, err := repo.UpdateAlarm(context.Background(), tc.alarm)
if tc.err != nil {
assert.True(t, errors.Contains(err, tc.err), fmt.Sprintf("%s: expected %s got %s\n", tc.desc, tc.err, err))
return
}
assert.Nil(t, err, fmt.Sprintf("unexpected error: %s", err))
assert.NotEmpty(t, alarm.ID)
assert.Equal(t, tc.alarm.Status, alarm.Status)
assert.Equal(t, tc.alarm.DomainID, alarm.DomainID)
assert.Equal(t, tc.alarm.AssigneeID, alarm.AssigneeID)
assert.Equal(t, tc.alarm.UpdatedBy, alarm.UpdatedBy)
assert.Equal(t, tc.alarm.ResolvedBy, alarm.ResolvedBy)
assert.Equal(t, tc.alarm.AcknowledgedBy, alarm.AcknowledgedBy)
assert.Equal(t, tc.alarm.Metadata, alarm.Metadata)
})
}
}
func TestViewAlarm(t *testing.T) {
t.Cleanup(func() {
_, err := db.Exec("DELETE FROM alarms")
require.Nil(t, err, fmt.Sprintf("clean alarms unexpected error: %s", err))
})
repo := postgres.NewAlarmsRepo(db)
alarm := alarms.Alarm{
ID: generateUUID(t),
RuleID: generateUUID(t),
DomainID: generateUUID(t),
ChannelID: generateUUID(t),
ClientID: generateUUID(t),
Measurement: namegen.Generate(),
Value: namegen.Generate(),
Unit: namegen.Generate(),
Threshold: namegen.Generate(),
Cause: namegen.Generate(),
Status: 0,
AssigneeID: generateUUID(t),
CreatedAt: time.Now().UTC(),
Metadata: map[string]any{
"key": "value",
},
}
alarm, err := repo.CreateAlarm(context.Background(), alarm)
require.Nil(t, err, fmt.Sprintf("unexpected error: %s", err))
cases := []struct {
desc string
id string
domainID string
err error
}{
{
desc: "valid alarm",
id: alarm.ID,
domainID: alarm.DomainID,
err: nil,
},
{
desc: "non existing alarm id",
id: generateUUID(t),
domainID: alarm.DomainID,
err: repoerr.ErrNotFound,
},
{
desc: "non existing domain id",
id: alarm.ID,
domainID: generateUUID(t),
err: repoerr.ErrNotFound,
},
}
for _, tc := range cases {
t.Run(tc.desc, func(t *testing.T) {
alarm, err := repo.ViewAlarm(context.Background(), tc.id, tc.domainID)
if tc.err != nil {
assert.True(t, errors.Contains(err, tc.err), fmt.Sprintf("%s: expected %s got %s\n", tc.desc, tc.err, err))
return
}
assert.Nil(t, err, fmt.Sprintf("unexpected error: %s", err))
assert.NotEmpty(t, alarm.ID)
assert.Equal(t, tc.id, alarm.ID)
})
}
}
func TestListAlarms(t *testing.T) {
t.Cleanup(func() {
_, err := db.Exec("DELETE FROM alarms")
require.Nil(t, err, fmt.Sprintf("clean alarms unexpected error: %s", err))
})
repo := postgres.NewAlarmsRepo(db)
items := make([]alarms.Alarm, 1000)
for i := range 1000 {
items[i] = alarms.Alarm{
ID: generateUUID(t),
RuleID: generateUUID(t),
DomainID: generateUUID(t),
ChannelID: generateUUID(t),
ClientID: generateUUID(t),
Measurement: namegen.Generate(),
Value: namegen.Generate(),
Unit: namegen.Generate(),
Threshold: namegen.Generate(),
Cause: namegen.Generate(),
Status: 0,
AssigneeID: generateUUID(t),
CreatedAt: time.Now().UTC(),
Metadata: map[string]any{
"key": "value",
},
}
alarm, err := repo.CreateAlarm(context.Background(), items[i])
require.Nil(t, err, fmt.Sprintf("unexpected error: %s", err))
items[i].ID = alarm.ID
}
cases := []struct {
desc string
pm alarms.PageMetadata
response []alarms.Alarm
err error
}{
{
desc: "valid page",
pm: alarms.PageMetadata{
Offset: 0,
Limit: 10,
},
response: items[:10],
err: nil,
},
{
desc: "offset and limit",
pm: alarms.PageMetadata{
Offset: 10,
Limit: 50,
},
response: items[10:60],
err: nil,
},
{
desc: "empty page",
pm: alarms.PageMetadata{},
response: []alarms.Alarm{},
err: nil,
},
{
desc: "invalid page",
pm: alarms.PageMetadata{
Offset: 1000,
Limit: 10,
},
response: []alarms.Alarm{},
err: nil,
},
{
desc: "invalid assignee id",
pm: alarms.PageMetadata{
Offset: 0,
Limit: 10,
AssigneeID: generateUUID(t),
},
response: []alarms.Alarm{},
err: nil,
},
}
for _, tc := range cases {
t.Run(tc.desc, func(t *testing.T) {
alarms, err := repo.ListAllAlarms(context.Background(), tc.pm)
if tc.err != nil {
assert.True(t, errors.Contains(err, tc.err), fmt.Sprintf("%s: expected %s got %s\n", tc.desc, tc.err, err))
return
}
assert.Nil(t, err, fmt.Sprintf("unexpected error: %s", err))
assert.Equal(t, len(tc.response), len(alarms.Alarms))
})
}
}
func TestDeleteAlarm(t *testing.T) {
t.Cleanup(func() {
_, err := db.Exec("DELETE FROM alarms")
require.Nil(t, err, fmt.Sprintf("clean alarms unexpected error: %s", err))
})
repo := postgres.NewAlarmsRepo(db)
alarm := alarms.Alarm{
ID: generateUUID(t),
RuleID: generateUUID(t),
DomainID: generateUUID(t),
ChannelID: generateUUID(t),
ClientID: generateUUID(t),
Measurement: namegen.Generate(),
Value: namegen.Generate(),
Unit: namegen.Generate(),
Threshold: namegen.Generate(),
Cause: namegen.Generate(),
Status: 0,
AssigneeID: generateUUID(t),
CreatedAt: time.Now().UTC(),
Metadata: map[string]any{
"key": "value",
},
}
alarm, err := repo.CreateAlarm(context.Background(), alarm)
require.Nil(t, err, fmt.Sprintf("unexpected error: %s", err))
cases := []struct {
desc string
id string
err error
}{
{
desc: "valid alarm",
id: alarm.ID,
err: nil,
},
{
desc: "non existing alarm",
id: generateUUID(t),
err: repoerr.ErrNotFound,
},
}
for _, tc := range cases {
t.Run(tc.desc, func(t *testing.T) {
err := repo.DeleteAlarm(context.Background(), tc.id)
if tc.err != nil {
assert.True(t, errors.Contains(err, tc.err), fmt.Sprintf("%s: expected %s got %s\n", tc.desc, tc.err, err))
return
}
assert.Nil(t, err, fmt.Sprintf("unexpected error: %s", err))
})
}
}
func generateUUID(t *testing.T) string {
ulid, err := idProvider.ID()
require.Nil(t, err, fmt.Sprintf("unexpected error: %s", err))
return ulid
}
+55
View File
@@ -0,0 +1,55 @@
// Copyright (c) Abstract Machines
// SPDX-License-Identifier: Apache-2.0
package postgres
import (
_ "github.com/jackc/pgx/v5/stdlib" // required for SQL access
migrate "github.com/rubenv/sql-migrate"
)
// Migration of Alarms service.
func Migration() (*migrate.MemoryMigrationSource, error) {
alarmsMigration := &migrate.MemoryMigrationSource{
Migrations: []*migrate.Migration{
{
Id: "alarms_01",
// VARCHAR(36) for columns with IDs as UUIDS have a maximum of 36 characters
Up: []string{
`CREATE TABLE IF NOT EXISTS alarms (
id VARCHAR(36) PRIMARY KEY,
rule_id VARCHAR(36) NOT NULL CHECK (length(rule_id) > 0),
domain_id VARCHAR(36) NOT NULL,
channel_id VARCHAR(36) NOT NULL,
subtopic TEXT NOT NULL,
client_id VARCHAR(36) NOT NULL,
measurement TEXT NOT NULL,
value TEXT NOT NULL,
unit TEXT NOT NULL,
threshold TEXT NOT NULL,
cause TEXT NOT NULL,
status SMALLINT NOT NULL DEFAULT 0 CHECK (status >= 0),
severity SMALLINT NOT NULL DEFAULT 0 CHECK (severity >= 0),
assignee_id VARCHAR(36),
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMPTZ NULL,
updated_by VARCHAR(36) NULL,
assigned_at TIMESTAMPTZ NULL,
assigned_by VARCHAR(36) NULL,
acknowledged_at TIMESTAMPTZ NULL,
acknowledged_by VARCHAR(36) NULL,
resolved_at TIMESTAMPTZ NULL,
resolved_by VARCHAR(36) NULL,
metadata JSONB
);`,
"CREATE INDEX IF NOT EXISTS idx_alarms_state ON alarms (domain_id, rule_id, channel_id, subtopic, client_id, measurement, created_at DESC);",
},
Down: []string{
`DROP TABLE IF EXISTS alarms`,
},
},
},
}
return alarmsMigration, nil
}
+97
View File
@@ -0,0 +1,97 @@
// Copyright (c) Abstract Machines
// SPDX-License-Identifier: Apache-2.0
package postgres_test
import (
"database/sql"
"fmt"
"log"
"os"
"testing"
"time"
apostgres "github.com/absmach/magistrala/alarms/postgres"
"github.com/absmach/magistrala/pkg/postgres"
"github.com/jmoiron/sqlx"
dockertest "github.com/ory/dockertest/v3"
"github.com/ory/dockertest/v3/docker"
"go.opentelemetry.io/otel"
)
var (
db *sqlx.DB
database postgres.Database
tracer = otel.Tracer("repo_tests")
)
func TestMain(m *testing.M) {
pool, err := dockertest.NewPool("")
if err != nil {
log.Fatalf("Could not connect to docker: %s", err)
}
container, err := pool.RunWithOptions(&dockertest.RunOptions{
Repository: "postgres",
Tag: "16.2-alpine",
Env: []string{
"POSTGRES_USER=test",
"POSTGRES_PASSWORD=test",
"POSTGRES_DB=test",
"listen_addresses = '*'",
},
}, func(config *docker.HostConfig) {
config.AutoRemove = true
config.RestartPolicy = docker.RestartPolicy{Name: "no"}
})
if err != nil {
log.Fatalf("Could not start container: %s", err)
}
port := container.GetPort("5432/tcp")
// exponential backoff-retry, because the application in the container might not be ready to accept connections yet
pool.MaxWait = 120 * time.Second
if err := pool.Retry(func() error {
url := fmt.Sprintf("host=localhost port=%s user=test dbname=test password=test sslmode=disable", port)
db, err := sql.Open("pgx", url)
if err != nil {
return err
}
return db.Ping()
}); err != nil {
log.Fatalf("Could not connect to docker: %s", err)
}
dbConfig := postgres.Config{
Host: "localhost",
Port: port,
User: "test",
Pass: "test",
Name: "test",
SSLMode: "disable",
SSLCert: "",
SSLKey: "",
SSLRootCert: "",
}
migration, err := apostgres.Migration()
if err != nil {
log.Fatalf("Could not get migration: %s", err)
}
if db, err = postgres.Setup(dbConfig, *migration); err != nil {
log.Fatalf("Could not setup test DB connection: %s", err)
}
database = postgres.NewDatabase(db, dbConfig, tracer)
code := m.Run()
// Defers will not be run when using os.Exit
db.Close()
if err := pool.Purge(container); err != nil {
log.Fatalf("Could not purge container: %s", err)
}
os.Exit(code)
}
+28 -3
View File
@@ -16,6 +16,7 @@ import (
"github.com/absmach/magistrala/alarms/consumer"
"github.com/absmach/magistrala/alarms/middleware"
"github.com/absmach/magistrala/alarms/operations"
alarmsRepo "github.com/absmach/magistrala/alarms/postgres"
"github.com/absmach/magistrala/internal/atom"
mglog "github.com/absmach/magistrala/logger"
smqauthn "github.com/absmach/magistrala/pkg/authn"
@@ -24,6 +25,7 @@ import (
"github.com/absmach/magistrala/pkg/messaging"
brokerstracing "github.com/absmach/magistrala/pkg/messaging/brokers/tracing"
"github.com/absmach/magistrala/pkg/permissions"
"github.com/absmach/magistrala/pkg/postgres"
"github.com/absmach/magistrala/pkg/prometheus"
"github.com/absmach/magistrala/pkg/server"
httpserver "github.com/absmach/magistrala/pkg/server/http"
@@ -34,7 +36,9 @@ import (
const (
svcName = "alarms"
envPrefixDB = "MG_ALARMS_DB_"
envPrefixHTTP = "MG_ALARMS_HTTP_"
defDB = "alarms"
defSvcHTTPPort = "8050"
alarmEntity = "alarm"
)
@@ -78,21 +82,42 @@ func main() {
}()
tracer := tp.Tracer(svcName)
dbConfig := postgres.Config{Name: defDB}
if err := env.ParseWithOptions(&dbConfig, env.Options{Prefix: envPrefixDB}); err != nil {
logger.Error(err.Error())
}
migrations, err := alarmsRepo.Migration()
if err != nil {
logger.Error(fmt.Sprintf("failed to load migrations: %s", err))
exitCode = 1
return
}
db, err := postgres.Setup(dbConfig, *migrations)
if err != nil {
logger.Error(err.Error())
exitCode = 1
return
}
defer db.Close()
repo := alarmsRepo.NewAlarmsRepo(db)
atomCfg := atom.LoadConfig()
if atomCfg.URL == "" {
logger.Error("ATOM_URL is required")
exitCode = 1
return
}
atomClient := atom.NewClient(atomCfg)
logger.Info("AuthN configured to use Atom bearer tokens")
logger.Info("AuthZ configured to use Atom PDP")
am := smqauthn.NewAuthNMiddleware(atomauthn.NewAuthentication())
idp := uuid.New()
repo := alarms.NewAtomRepository(atomClient)
svc := alarms.NewService(idp, repo)
svc = alarms.WithAtom(svc, atom.NewClient(atomCfg))
permConfig, err := permissions.ParsePermissionsFile(cfg.PermissionsFile)
if err != nil {
@@ -122,7 +147,7 @@ func main() {
return
}
svc, err = middleware.NewAtomAuthorizationMiddleware(svc, atomClient, entitiesOps)
svc, err = middleware.NewAtomAuthorizationMiddleware(svc, atom.NewClient(atomCfg), entitiesOps)
if err != nil {
logger.Error(fmt.Sprintf("failed to create authorization middleware: %s", err))
exitCode = 1
+4
View File
@@ -739,3 +739,7 @@ MG_UI_CLI_WS_URL=ws://localhost:80/mqtt
MG_UI_CLI_COAP_HOST=0.0.0.0
MG_UI_CLI_COAP_PORT=5684
MG_UI_CLI_HTTP_URL=http://localhost:80/http
# Atom
ATOM_LOG_LEVEL=info
ATOM_LOG_FORMAT=text
+28
View File
@@ -16,6 +16,7 @@ volumes:
magistrala-ui-backend-db-volume:
magistrala-journal-volume:
magistrala-re-db-volume:
magistrala-alarms-db-volume:
magistrala-reports-db-volume:
magistrala-timescale-writer-volume:
magistrala-fluxmq-node1-volume:
@@ -878,10 +879,28 @@ services:
- ./permission.yaml:${MG_PERMISSIONS_FILE}
- ./templates/${MG_RE_EMAIL_TEMPLATE}:/email.tmpl
alarms-db:
image: docker.io/postgres:18.0-alpine3.22
container_name: magistrala-alarms-db
restart: on-failure
command: postgres -c "max_connections=${MG_POSTGRES_MAX_CONNECTIONS}"
environment:
POSTGRES_USER: ${MG_ALARMS_DB_USER}
POSTGRES_PASSWORD: ${MG_ALARMS_DB_PASS}
POSTGRES_DB: ${MG_ALARMS_DB_NAME}
ports:
- 6019:5432
networks:
- magistrala-base-net
volumes:
- magistrala-alarms-db-volume:/var/lib/postgresql/data
alarms:
image: ghcr.io/absmach/magistrala/alarms:${MG_RELEASE_TAG}
container_name: magistrala-alarms
depends_on:
alarms-db:
condition: service_started
atom-bootstrap:
condition: service_completed_successfully
nginx:
@@ -899,6 +918,15 @@ services:
MG_ALARMS_HTTP_HOST: ${MG_ALARMS_HTTP_HOST}
MG_ALARMS_HTTP_SERVER_CERT: ${MG_ALARMS_HTTP_SERVER_CERT}
MG_ALARMS_HTTP_SERVER_KEY: ${MG_ALARMS_HTTP_SERVER_KEY}
MG_ALARMS_DB_HOST: ${MG_ALARMS_DB_HOST}
MG_ALARMS_DB_PORT: ${MG_ALARMS_DB_PORT}
MG_ALARMS_DB_USER: ${MG_ALARMS_DB_USER}
MG_ALARMS_DB_PASS: ${MG_ALARMS_DB_PASS}
MG_ALARMS_DB_NAME: ${MG_ALARMS_DB_NAME}
MG_ALARMS_DB_SSL_MODE: ${MG_ALARMS_DB_SSL_MODE}
MG_ALARMS_DB_SSL_CERT: ${MG_ALARMS_DB_SSL_CERT}
MG_ALARMS_DB_SSL_KEY: ${MG_ALARMS_DB_SSL_KEY}
MG_ALARMS_DB_SSL_ROOT_CERT: ${MG_ALARMS_DB_SSL_ROOT_CERT}
MG_MESSAGE_BROKER_URL: ${MG_MESSAGE_BROKER_URL}
MG_ES_URL: ${MG_ES_URL}
MG_JAEGER_URL: ${MG_JAEGER_URL}
+2 -5
View File
@@ -554,11 +554,8 @@ func (c *Client) ListResources(ctx context.Context, q Query) (ResourceList, erro
if q.Name != "" && q.Q == "" {
vars["q"] = q.Name
}
if len(q.AttributesContains) > 0 {
vars["attributesContains"] = q.AttributesContains
}
err := c.graphQL(ctx, `query Resources($q: String, $kind: String, $tenantId: ID, $attributesContains: JSON, $limit: Int, $offset: Int) {
resources(q: $q, kind: $kind, tenantId: $tenantId, attributesContains: $attributesContains, limit: $limit, offset: $offset) {
err := c.graphQL(ctx, `query Resources($q: String, $kind: String, $tenantId: ID, $limit: Int, $offset: Int) {
resources(q: $q, kind: $kind, tenantId: $tenantId, limit: $limit, offset: $offset) {
total
items { id kind name tenant_id: tenantId owner_id: ownerId attributes created_at: createdAt updated_at: updatedAt }
}
+1 -13
View File
@@ -65,22 +65,14 @@ func TestListResources(t *testing.T) {
t.Fatalf("unexpected request: %s %s", r.Method, r.URL.Path)
}
var payload struct {
Query string `json:"query"`
Variables map[string]any `json:"variables"`
}
if err := json.NewDecoder(r.Body).Decode(&payload); err != nil {
t.Fatalf("decode request: %v", err)
}
if !strings.Contains(payload.Query, "$attributesContains: JSON") || !strings.Contains(payload.Query, "attributesContains: $attributesContains") {
t.Fatalf("query does not include attributesContains: %s", payload.Query)
}
if payload.Variables["kind"] != KindRule || payload.Variables["tenantId"] != testDomainID {
t.Fatalf("unexpected variables: %+v", payload.Variables)
}
attrs, ok := payload.Variables["attributesContains"].(map[string]any)
if !ok || attrs["status"] != "active" {
t.Fatalf("unexpected attributesContains: %+v", payload.Variables)
}
_ = json.NewEncoder(w).Encode(map[string]any{
"data": map[string]any{
"resources": map[string]any{
@@ -93,11 +85,7 @@ func TestListResources(t *testing.T) {
defer srv.Close()
client := NewClient(Config{URL: srv.URL, Timeout: time.Second})
got, err := client.ListResources(context.Background(), Query{
Kind: KindRule,
TenantID: testDomainID,
AttributesContains: Attributes{"status": "active"},
})
got, err := client.ListResources(context.Background(), Query{Kind: KindRule, TenantID: testDomainID})
if err != nil {
t.Fatalf("list failed: %v", err)
}
+9 -10
View File
@@ -55,16 +55,15 @@ type Resource struct {
}
type Query struct {
IDs []string
Q string
Kind string
TenantID string
Name string
Route string
Status string
AttributesContains Attributes
Limit uint64
Offset uint64
IDs []string
Q string
Kind string
TenantID string
Name string
Route string
Status string
Limit uint64
Offset uint64
}
type AuthzRequest struct {