mirror of
https://github.com/absmach/magistrala.git
synced 2026-08-07 07:14:46 +00:00
14d6db9684
Property Based Tests / api-test (push) Has been cancelled
Continuous Delivery / lint-and-build (push) Has been cancelled
Deploy GitHub Pages / swagger-ui (push) Has been cancelled
CI Pipeline / Lint Proto (push) Has been cancelled
CI Pipeline / Detect Changes (push) Has been cancelled
Continuous Delivery / Build and Push Docker Images (push) Has been cancelled
CI Pipeline / lint-and-build (push) Has been cancelled
CI Pipeline / Test ${{ matrix.module }} (push) Has been cancelled
CI Pipeline / Upload Coverage (push) Has been cancelled
Signed-off-by: fbugarski <filipbugarski@gmail.com> Co-authored-by: Claude Sonnet 5 <noreply@anthropic.com>
174 lines
4.9 KiB
Go
174 lines
4.9 KiB
Go
// Copyright (c) Abstract Machines
|
|
// SPDX-License-Identifier: Apache-2.0
|
|
|
|
package grpc
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
|
|
grpcReadersV1 "github.com/absmach/magistrala/api/grpc/readers/v1"
|
|
grpcapi "github.com/absmach/magistrala/auth/api/grpc"
|
|
"github.com/absmach/magistrala/pkg/transformers/senml"
|
|
"github.com/absmach/magistrala/readers"
|
|
kitgrpc "github.com/go-kit/kit/transport/grpc"
|
|
)
|
|
|
|
var _ grpcReadersV1.ReadersServiceServer = (*readersGrpcServer)(nil)
|
|
|
|
type readersGrpcServer struct {
|
|
grpcReadersV1.UnimplementedReadersServiceServer
|
|
readMessages kitgrpc.Handler
|
|
}
|
|
|
|
func NewReadersServer(svc readers.MessageRepository) grpcReadersV1.ReadersServiceServer {
|
|
return &readersGrpcServer{
|
|
readMessages: kitgrpc.NewServer(
|
|
(readMessagesEndpoint(svc)),
|
|
decodeReadMessagesRequest,
|
|
encodeReadMessagesResponse,
|
|
),
|
|
}
|
|
}
|
|
|
|
func decodeReadMessagesRequest(_ context.Context, grpcReq any) (any, error) {
|
|
req := grpcReq.(*grpcReadersV1.ReadMessagesReq)
|
|
return readMessagesReq{
|
|
chanID: req.GetChannelId(),
|
|
domain: req.GetDomainId(),
|
|
pageMeta: readers.PageMetadata{
|
|
Offset: req.GetPageMetadata().GetOffset(),
|
|
Limit: req.GetPageMetadata().GetLimit(),
|
|
Comparator: req.GetPageMetadata().GetComparator(),
|
|
Aggregation: stringifyAggregation(req.GetPageMetadata().GetAggregation()),
|
|
From: req.GetPageMetadata().GetFrom(),
|
|
To: req.GetPageMetadata().GetTo(),
|
|
Interval: req.GetPageMetadata().GetInterval(),
|
|
Subtopic: req.GetPageMetadata().GetSubtopic(),
|
|
Publisher: req.GetPageMetadata().GetPublisher(),
|
|
Publishers: req.GetPageMetadata().GetPublishers(),
|
|
Protocol: req.GetPageMetadata().GetProtocol(),
|
|
Name: req.GetPageMetadata().GetName(),
|
|
Value: req.GetPageMetadata().GetValue(),
|
|
BoolValue: req.GetPageMetadata().GetBoolValue(),
|
|
StringValue: req.GetPageMetadata().GetStringValue(),
|
|
DataValue: req.GetPageMetadata().GetDataValue(),
|
|
Format: req.GetPageMetadata().GetFormat(),
|
|
Order: req.GetPageMetadata().GetOrder(),
|
|
Dir: req.GetPageMetadata().GetDir(),
|
|
},
|
|
}, nil
|
|
}
|
|
|
|
func encodeReadMessagesResponse(_ context.Context, grpcRes any) (any, error) {
|
|
res := grpcRes.(readMessagesRes)
|
|
|
|
resp := &grpcReadersV1.ReadMessagesRes{
|
|
Total: res.Total,
|
|
Messages: toResponseMessages(res.Messages),
|
|
PageMetadata: &grpcReadersV1.PageMetadata{
|
|
Offset: res.PageMetadata.Offset,
|
|
Limit: res.PageMetadata.Limit,
|
|
Order: res.PageMetadata.Order,
|
|
Dir: res.PageMetadata.Dir,
|
|
},
|
|
}
|
|
return resp, nil
|
|
}
|
|
|
|
func (s *readersGrpcServer) ReadMessages(ctx context.Context, req *grpcReadersV1.ReadMessagesReq) (*grpcReadersV1.ReadMessagesRes, error) {
|
|
_, res, err := s.readMessages.ServeGRPC(ctx, req)
|
|
if err != nil {
|
|
return nil, grpcapi.EncodeError(err)
|
|
}
|
|
return res.(*grpcReadersV1.ReadMessagesRes), nil
|
|
}
|
|
|
|
func toResponseMessages(messages []readers.Message) []*grpcReadersV1.Message {
|
|
var res []*grpcReadersV1.Message
|
|
for _, m := range messages {
|
|
switch typed := m.(type) {
|
|
case senml.Message:
|
|
res = append(res, &grpcReadersV1.Message{
|
|
Payload: &grpcReadersV1.Message_Senml{
|
|
Senml: &grpcReadersV1.SenMLMessage{
|
|
Base: &grpcReadersV1.BaseMessage{
|
|
Channel: typed.Channel,
|
|
Subtopic: typed.Subtopic,
|
|
Publisher: typed.Publisher,
|
|
Protocol: typed.Protocol,
|
|
},
|
|
Name: typed.Name,
|
|
Unit: typed.Unit,
|
|
Time: typed.Time,
|
|
UpdateTime: typed.UpdateTime,
|
|
Value: typed.Value,
|
|
StringValue: typed.StringValue,
|
|
DataValue: typed.DataValue,
|
|
BoolValue: typed.BoolValue,
|
|
Sum: typed.Sum,
|
|
},
|
|
},
|
|
})
|
|
case map[string]any:
|
|
payload := typed["payload"]
|
|
data, err := json.Marshal(payload)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
res = append(res, &grpcReadersV1.Message{
|
|
Payload: &grpcReadersV1.Message_Json{
|
|
Json: &grpcReadersV1.JsonMessage{
|
|
Base: &grpcReadersV1.BaseMessage{
|
|
Channel: safeString(typed["channel"]),
|
|
Subtopic: safeString(typed["subtopic"]),
|
|
Publisher: safeString(typed["publisher"]),
|
|
Protocol: safeString(typed["protocol"]),
|
|
},
|
|
Created: safeInt64(typed["created"]),
|
|
Payload: data,
|
|
},
|
|
},
|
|
})
|
|
}
|
|
}
|
|
return res
|
|
}
|
|
|
|
func stringifyAggregation(agg grpcReadersV1.Aggregation) string {
|
|
switch agg {
|
|
case grpcReadersV1.Aggregation_AGGREGATION_UNSPECIFIED:
|
|
return ""
|
|
case grpcReadersV1.Aggregation_AGGREGATION_MAX:
|
|
return aggregationMax
|
|
case grpcReadersV1.Aggregation_AGGREGATION_MIN:
|
|
return aggregationMin
|
|
case grpcReadersV1.Aggregation_AGGREGATION_AVG:
|
|
return aggregationAvg
|
|
case grpcReadersV1.Aggregation_AGGREGATION_SUM:
|
|
return aggregationSum
|
|
case grpcReadersV1.Aggregation_AGGREGATION_COUNT:
|
|
return aggregationCount
|
|
default:
|
|
return ""
|
|
}
|
|
}
|
|
|
|
func safeString(v any) string {
|
|
if s, ok := v.(string); ok {
|
|
return s
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func safeInt64(v any) int64 {
|
|
switch v := v.(type) {
|
|
case float64:
|
|
return int64(v)
|
|
case int64:
|
|
return v
|
|
default:
|
|
return 0
|
|
}
|
|
}
|