Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ endef

all: edgex-app

.PHONY: all
.PHONY: all docker


edgex-app:
Expand All @@ -20,8 +20,13 @@ clean:
install:
cp ${BUILD_DIR}/* $(GOBIN)

docker:
docker build --no-cache --build-arg edgex-app --tag=edgex-app -f docker/Dockerfile .


test:
GOCACHE=off go test -v -race -tags test $(shell go list ./... | grep -v 'vendor\|cmd')




43 changes: 34 additions & 9 deletions cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,11 @@ const (
envDBSSLCert = "MF_EDGEX_DB_SSL_CERT"
envDBSSLKey = "MF_EDGEX_DB_SSL_KEY"
envDBSSLRootCert = "MF_EDGEX_DB_SSL_ROOT_CERT"
envMqttHost = "MF_EDGEX_MQTT_HOST"
envMqttPassword = "MF_EDGEX_MQTT_PASS"
envMqttUsername = "MF_EDGEX_MQTT_USERNAME"
envMqttSndTopic = "MF_EDGEX_MQTT_TOPIC"
envMqttClientID = "MF_EDGEX_MQTT_CLIENT_ID"
topicUnknown = "out.unknown"
)

Expand All @@ -64,26 +69,33 @@ func main() {
if err != nil {
log.Fatalf(err.Error())
}
// Connect to a NATS server

// Connect to a DB server
db := connectToDB(cfg.dbConfig, logger)
if db == nil {
log.Fatalf("cannot connect to db")
}
defer db.Close()

logger.Info(fmt.Sprintf("connecting %s", defNatsURL))
nc, err := nats.Connect(defNatsURL)
if err != nil {
logger.Error("Failed to connect to nats")
// Make a MQTT publish connection
mfMqttHost := Env(envMqttHost, "")
mfMqttPassword := Env(envMqttPassword, "")
mfMqttUsername := Env(envMqttUsername, "")
mfMqttClientID := Env(envMqttClientID, "")
mfMqttSndTopic := Env(envMqttSndTopic, "")
if mfMqttPassword != "" && mfMqttUsername != "" && mfMqttSndTopic != "" {
_, err := api.NewMQTTPublisher(&logger, mfMqttHost, mfMqttUsername, mfMqttPassword, mfMqttClientID, mfMqttSndTopic)
if err != nil {
logger.Error(fmt.Sprintf("Failed to create mqtt publisher: %s", err))
}
}
defer closeConn(nc, logger)

// create a service
svc := newService(db, logger)

logger.Info(fmt.Sprintf("pid: %d connecting to nats\n", os.Getpid()))
nc.Subscribe(topicUnknown, exapp.NatsMSGHandler(svc))
// connect to NATS
connectToNats(svc, logger)

// configure endpoints and start http server
err = http.ListenAndServe(fmt.Sprintf(":%s", cfg.Port), httpapi.MakeHandler(svc, logger))
if err != nil {
logger.Error(fmt.Sprintf("Failed to init http server on port %s: ", cfg.Port))
Expand Down Expand Up @@ -137,7 +149,19 @@ func loadConfig() config {
dbConfig: dbConfig,
}
}
func connectToNats(svc exapp.Service, logger logger.Logger) {

logger.Info(fmt.Sprintf("connecting %s", defNatsURL))

nc, err := nats.Connect(defNatsURL)
defer closeConn(nc, logger)

if err != nil {
logger.Error("Failed to connect to nats")
} else {
nc.Subscribe(topicUnknown, exapp.NatsMSGHandler(svc))
}
}
func connectToDB(dbConfig postgres.Config, logger logger.Logger) *sql.DB {
db, err := postgres.Connect(dbConfig, logger)

Expand All @@ -163,6 +187,7 @@ func newService(db *sql.DB, logger logger.Logger) exapp.Service {
eventsRepository := postgres.New(db, logger)
svc := exapp.New(eventsRepository, logger)
svc = api.LoggingMiddleware(svc)

svc = api.MetricsMiddleware(
svc,
kitprometheus.NewCounterFrom(stdprometheus.CounterOpts{
Expand Down
13 changes: 13 additions & 0 deletions docker/Dockerfile
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
FROM golang:1.11.2-alpine AS builder


WORKDIR /go/src/github.com/mteodor/edgex-app
COPY . .

RUN apk update \
&& apk add make bash \
&& mv build/edgex-app /exe

FROM scratch
COPY --from=builder /exe /
ENTRYPOINT ["/exe"]
44 changes: 44 additions & 0 deletions docker/addons/edgex-app/docker-compose.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
###
# This docker-compose file contains optional InfluxDB-reader service for the Mainflux
# platform. Since this service is optional, this file is dependent on the docker-compose.yml
# file from <project_root>/docker/. In order to run InfluxDB-reader service, core services,
# as well as the network from the core composition, should be already running.
###

version: "3"

networks:
mainflux-base-net:
driver: bridge

services:
events-db:
image: postgres:10.2-alpine
container_name: edgexapp-events-db
restart: on-failure
environment:
POSTGRES_USER: mainflux
POSTGRES_PASSWORD: mainflux
POSTGRES_DB: events
networks:
- mainflux-base-net

edgex-app:
image: mainflux/users:latest
container_name: edgex-app
depends_on:
- events-db
expose:
- 8000
restart: on-failure
environment:
MF_EDGEX_APP_LOG_LEVEL: debug
MF_EDGEX_DB_HOST: events-db
MF_EDGEX_DB_PORT: 5432
MF_EDGEX_DB_USER: mainflux
MF_EDGEX_DB_PASS: mainflux
MF_EDGEX_DB: events
ports:
- 8000:8000
networks:
- mainflux-base-net
58 changes: 58 additions & 0 deletions docker/docker-compose.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
###
# Copyright (c) 2015-2017 Mainflux
#
# Mainflux is licensed under an Apache license, version 2.0 license.
# All rights not explicitly granted in the Apache license, version 2.0 are reserved.
# See the included LICENSE file for more details.
###

version: "3"

networks:
mainflux-base-net:
driver: bridge

services:
events-db:
image: postgres:10.2-alpine
container_name: edgexapp-events-db
restart: on-failure
environment:
POSTGRES_USER: mainflux
POSTGRES_PASSWORD: mainflux
POSTGRES_DB: events
networks:
- mainflux-base-net

edgex-app:
image: edgex-app:latest
container_name: edgex-app
depends_on:
- events-db
expose:
- 8000
restart: on-failure
environment:
MF_EDGEX_APP_LOG_LEVEL: debug
MF_EDGEX_DB_HOST: events-db
MF_EDGEX_DB_PORT: 5432
MF_EDGEX_DB_USER: mainflux
MF_EDGEX_DB_PASS: mainflux
MF_EDGEX_DB: events
ports:
- 8000:8000
networks:
- mainflux-base-net

nats:
image: nats:1.3.0
container_name: edgexapp-nats
expose:
- 4222
restart: on-failure
networks:
- mainflux-base-net
ports:
- 4222:4222


59 changes: 59 additions & 0 deletions exapp/api/mqttpublisher.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
package api

import (
"fmt"

mqtt "github.com/eclipse/paho.mqtt.golang"
log "github.com/mainflux/mainflux/logger"
)

type MQTTPublisher interface {
connect() error
Publish([]byte)
}

type mqttPublisher struct {
logger *log.Logger
mqttClient mqtt.Client
host string
username string
password string
clientID string
topic string
}

var _ MQTTPublisher = (*mqttPublisher)(nil)

// NewMQTTPublisher creates a new MQTT publisher, configures it's topic and connects it to the MQTT server.
func NewMQTTPublisher(logger *log.Logger, host, username, password, clientID, topic string) (*mqttPublisher, error) {
p := mqttPublisher{logger: logger, host: host, username: username, password: password, clientID: clientID, topic: topic}
if err := p.connect(); err != nil {
return &p, err
}
return &p, nil
}

func (m *mqttPublisher) connect() error {
opts := mqtt.NewClientOptions().AddBroker(m.host)
opts.SetUsername(m.username)
opts.SetPassword(m.password)
opts.SetClientID(m.clientID)

mc := mqtt.NewClient(opts)

if token := mc.Connect(); token.Wait() && token.Error() != nil {
return token.Error()
}
m.mqttClient = mc
return nil
}

// Publish publishes a message to the configured topic.
func (m *mqttPublisher) Publish(payload []byte) {

(*m.logger).Debug(fmt.Sprintf("Publishing message %s", payload))
if token := m.mqttClient.Publish(m.topic, 0, false, payload); token.Wait() && token.Error() != nil {
(*m.logger).Error(fmt.Sprintf("Failed to publish message on topic %s : %s", m.topic, token.Error()))
}

}
64 changes: 64 additions & 0 deletions exapp/swagger.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
swagger: "2.0"
info:
title: Edgex app service
description: Service for recieving and manipulating edgex specific format
version: "1.0.0"
consumes:
- "application/json"
produces:
- "application/json"
paths:
/status:
post:
summary: Gives a status and greeting to user
description: |
Gives a status and greeting to user.
tags:
- status
parameters:
- name: status
description: JSON-formatted message providing username.
in: body
schema:
$ref: "#/definitions/StatReq"
required: true
responses:
201:
description: Greeting to user with status of service.
headers:
Location:
type: string
400:
description: Failed due to malformed JSON.
404:
description: Page does not exist..
/version:
get:
summary: Retrieves version info
tags:
- version
responses:
200:
description: Service version.
schema:
$ref: "#/definitions/VerRes"
404:
description: Page does not exist.
responses:
ServiceError:
description: Unexpected server-side error occured.
definitions:
VerRes:
type: object
properties:
Version:
type: string
description: Version of the service.
StatReq:
type: object
properties:
Name:
type: string
description: Name of user.
required:
- type
Loading