Mainflux.mainflux/readers/mocks/messages.go

164 lines
3.3 KiB
Go
Raw Normal View History

// Copyright (c) Mainflux
// SPDX-License-Identifier: Apache-2.0
package mocks
import (
"encoding/json"
"sync"
"github.com/mainflux/mainflux/pkg/transformers/senml"
"github.com/mainflux/mainflux/readers"
)
var _ readers.MessageRepository = (*messageRepositoryMock)(nil)
type messageRepositoryMock struct {
mutex sync.Mutex
MF-1264 - Add support for JSON readers (#1295) * MF-1254 - Create universal JSON writer (#1260) Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Add JSON support to Readers Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Fix Influx Reader tests Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Fix messages format query Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Fix Postgres reader Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Fix Cassandra Readers and writers Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Fix Mongo reader Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Extract utility method to the JSON transformer Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Fix Influx and Postgres count Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Update JSON transformer Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Fix Influxdb Reader total count Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Refactor init.go for Cassandra writer Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Create a Payload type Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Add comments for defaults Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Fix variable declarations Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Replace interface{} with a new type Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Don't set channel just to overwrite it later Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Fix range search Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> * Rename Messages field Signed-off-by: dusanb94 <dusan.borovcanin@mainflux.com> Co-authored-by: Manuel Imperiale <manuel.imperiale@gmail.com>
2020-12-30 22:43:04 +08:00
messages map[string][]readers.Message
}
// NewMessageRepository returns mock implementation of message repository.
func NewMessageRepository(chanID string, messages []readers.Message) readers.MessageRepository {
repo := map[string][]readers.Message{
chanID: messages,
}
return &messageRepositoryMock{
mutex: sync.Mutex{},
messages: repo,
}
}
func (repo *messageRepositoryMock) ReadAll(chanID string, rpm readers.PageMetadata) (readers.MessagesPage, error) {
repo.mutex.Lock()
defer repo.mutex.Unlock()
if rpm.Format != "" && rpm.Format != "messages" {
return readers.MessagesPage{}, nil
}
var query map[string]interface{}
meta, _ := json.Marshal(rpm)
json.Unmarshal(meta, &query)
var msgs []readers.Message
for _, m := range repo.messages[chanID] {
senml := m.(senml.Message)
ok := true
for name := range query {
switch name {
case "subtopic":
if rpm.Subtopic != senml.Subtopic {
ok = false
}
case "publisher":
if rpm.Publisher != senml.Publisher {
ok = false
}
case "name":
if rpm.Name != senml.Name {
ok = false
}
case "protocol":
if rpm.Protocol != senml.Protocol {
ok = false
}
case "v":
if senml.Value == nil {
ok = false
}
val, okQuery := query["comparator"]
if okQuery {
switch val.(string) {
case readers.LowerThanKey:
if senml.Value != nil &&
*senml.Value >= rpm.Value {
ok = false
}
case readers.LowerThanEqualKey:
if senml.Value != nil &&
*senml.Value > rpm.Value {
ok = false
}
case readers.GreaterThanKey:
if senml.Value != nil &&
*senml.Value <= rpm.Value {
ok = false
}
case readers.GreaterThanEqualKey:
if senml.Value != nil &&
*senml.Value < rpm.Value {
ok = false
}
case readers.EqualKey:
default:
if senml.Value != nil &&
*senml.Value != rpm.Value {
ok = false
}
}
}
case "vb":
if senml.BoolValue == nil ||
(senml.BoolValue != nil &&
*senml.BoolValue != rpm.BoolValue) {
ok = false
}
case "vs":
if senml.StringValue == nil ||
(senml.StringValue != nil &&
*senml.StringValue != rpm.StringValue) {
ok = false
}
case "vd":
if senml.DataValue == nil ||
(senml.DataValue != nil &&
*senml.DataValue != rpm.DataValue) {
ok = false
}
case "from":
if senml.Time < rpm.From {
ok = false
}
case "to":
if senml.Time >= rpm.To {
ok = false
}
}
if !ok {
break
}
}
if ok {
msgs = append(msgs, m)
}
}
numOfMessages := uint64(len(msgs))
if rpm.Offset >= numOfMessages {
return readers.MessagesPage{}, nil
}
if rpm.Limit < 1 {
return readers.MessagesPage{}, nil
}
end := rpm.Offset + rpm.Limit
if rpm.Offset+rpm.Limit > numOfMessages {
end = numOfMessages
}
return readers.MessagesPage{
PageMetadata: rpm,
Total: uint64(len(msgs)),
Messages: msgs[rpm.Offset:end],
}, nil
}