Compare commits

..

4 Commits

Author SHA1 Message Date
ecc9bef172 Merge pull request 'fix' (#21) from master into prod
Reviewed-on: #21
2025-01-13 13:47:26 +03:00
a9f14411ce Merge pull request 'master' (#19) from master into prod
Reviewed-on: #19
2025-01-07 02:30:43 +03:00
a4a3176cb0 Merge pull request 'master' (#8) from master into prod
Reviewed-on: #8
2025-01-01 14:26:00 +03:00
bf152e6d13 Merge pull request 'master' (#6) from master into prod
Reviewed-on: #6
2024-12-31 03:04:59 +03:00
6 changed files with 5 additions and 107 deletions

View File

@@ -6,11 +6,9 @@ services:
image: mathwave/sprint-repo:queues
networks:
- queues-development
- monitoring
- mongo-development
environment:
MONGO_HOST: "mongo.develop.sprinthub.ru"
MONGO_PASSWORD: $MONGO_PASSWORD_DEV
STAGE: "development"
deploy:
mode: replicated
restart_policy:
@@ -24,7 +22,3 @@ services:
networks:
queues-development:
external: true
monitoring:
external: true
mongo-development:
external: true

View File

@@ -6,11 +6,9 @@ services:
image: mathwave/sprint-repo:queues
networks:
- queues
- monitoring
- mongo
environment:
MONGO_HOST: "mongo.sprinthub.ru"
MONGO_PASSWORD: $MONGO_PASSWORD_PROD
STAGE: "production"
deploy:
mode: replicated
restart_policy:
@@ -24,7 +22,3 @@ services:
networks:
queues:
external: true
monitoring:
external: true
mongo:
external: true

View File

@@ -1,15 +1,9 @@
package routers
import (
"bytes"
"fmt"
"log"
"net/http"
"os"
tasks "queues-go/app/storage/mongo/collections"
"strconv"
"sync"
"time"
)
type TaskResponse struct {
@@ -24,24 +18,6 @@ type TakeResponse struct {
var MutexMap map[string]*sync.Mutex
func sendLatency(timestamp time.Time, latency int) error {
loc, _ := time.LoadLocation("Europe/Moscow")
s := fmt.Sprintf(
`{"timestamp":"%s","service":"queues","environment":"%s","name":"latency","count":%s}`,
timestamp.In(loc).Format("2006-01-02T15:04:05Z"),
os.Getenv("STAGE"),
strconv.Itoa(latency),
)
data := []byte(s)
r := bytes.NewReader(data)
_, err := http.Post("http://monitoring:1237/api/v1/metrics/increment", "application/json", r)
if err != nil {
log.Printf("ERROR %s", err.Error())
}
return err
}
func Take(r *http.Request) (interface{}, int) {
queue := r.Header.Get("queue")
mutex, ok := MutexMap[queue]
@@ -65,8 +41,6 @@ func Take(r *http.Request) (interface{}, int) {
Attempt: task.Attempts,
Payload: task.Payload,
}
now := time.Now()
go sendLatency(now, int(now.Sub(task.PutAt).Milliseconds()))
}
return response, http.StatusOK

View File

@@ -12,7 +12,7 @@ import (
var Database mongo.Database
func Connect() {
connectionString := fmt.Sprintf("mongodb://mongo:%s@mongo:27017/", utils.GetEnv("MONGO_PASSWORD", "password"))
connectionString := fmt.Sprintf("mongodb://mongo:%s@%s:27017/", utils.GetEnv("MONGO_PASSWORD", "password"), utils.GetEnv("MONGO_HOST", "localhost"))
client, err := mongo.Connect(context.TODO(), options.Client().ApplyURI(connectionString))
if err != nil {
panic(err)

View File

@@ -33,14 +33,6 @@ type InsertedTask struct {
Attempts int `bson:"attempts"`
}
func Count() (*int64, error) {
count, err := collection().CountDocuments(context.TODO(), bson.M{})
if err != nil {
return nil, err
}
return &count, nil
}
func Add(task InsertedTask) error {
_, err := collection().InsertOne(context.TODO(), task)
if err != nil {
@@ -79,7 +71,7 @@ func Take(queue string) (*Task, error) {
"taken_at": now,
"attempts": task.Attempts + 1,
"available_from": now.Add(
time.Duration(task.SecondsToExecute+task.Attempts) * time.Second,
time.Duration(task.SecondsToExecute) * time.Second,
),
},
},
@@ -99,7 +91,7 @@ func findTask(queue string, now time.Time) (*Task, error) {
"queue": queue,
"available_from": bson.M{"$lte": now},
},
options.Find().SetSort(bson.D{{Key: "available_from", Value: 1}}).SetLimit(1),
options.Find().SetLimit(1),
)
if err != nil {
println("Error find")

56
main.go
View File

@@ -1,39 +1,15 @@
package main
import (
"bytes"
"encoding/json"
"fmt"
"log"
"net/http"
"os"
"queues-go/app/routers"
client "queues-go/app/storage/mongo"
tasks "queues-go/app/storage/mongo/collections"
"strconv"
"sync"
"time"
)
func writeMetric(timestamp time.Time, endpoint string, statusCode int, responseTime int, method string) {
loc, _ := time.LoadLocation("Europe/Moscow")
s := fmt.Sprintf(
`{"timestamp":"%s","service":"queues","environment":"%s","endpoint":"%s","status_code":%s,"response_time":%s,"method":"%s"}`,
timestamp.In(loc).Format("2006-01-02T15:04:05Z"),
os.Getenv("STAGE"),
endpoint,
strconv.Itoa(statusCode),
strconv.Itoa(responseTime),
method,
)
data := []byte(s)
r := bytes.NewReader(data)
_, err := http.Post("http://monitoring:1237/api/v1/metrics/endpoint", "application/json", r)
if err != nil {
log.Printf("Error sending metrics %s", err.Error())
}
}
func handlerWrapper(f func(*http.Request) (interface{}, int)) func(http.ResponseWriter, *http.Request) {
return func(w http.ResponseWriter, r *http.Request) {
start := time.Now()
@@ -49,48 +25,16 @@ func handlerWrapper(f func(*http.Request) (interface{}, int)) func(http.Response
} else {
w.WriteHeader(status)
}
go writeMetric(start, r.URL.Path, status, int(time.Since(start).Milliseconds()), r.Method)
log.Printf("%s %d %s", r.URL, status, time.Since(start))
}
}
func metricProxy(w http.ResponseWriter, r *http.Request) {
http.Post("http://monitoring:1237/api/v1/metrics/task", "application/json", r.Body)
w.WriteHeader(202)
}
func writeCount() {
for {
count, err := tasks.Count()
if err != nil {
log.Printf("Failed getting docs count: %s", err.Error())
} else {
loc, _ := time.LoadLocation("Europe/Moscow")
s := fmt.Sprintf(
`{"timestamp":"%s","service":"queues","environment":"%s","name":"tasks","count":%s}`,
time.Now().In(loc).Format("2006-01-02T15:04:05Z"),
os.Getenv("STAGE"),
strconv.Itoa(int(*count)),
)
data := []byte(s)
r := bytes.NewReader(data)
_, err := http.Post("http://monitoring:1237/api/v1/metrics/increment", "application/json", r)
if err != nil {
log.Printf("ERROR %s", err.Error())
}
}
time.Sleep(time.Second)
}
}
func main() {
client.Connect()
routers.MutexMap = make(map[string]*sync.Mutex)
http.HandleFunc("/api/v1/take", handlerWrapper(routers.Take))
http.HandleFunc("/api/v1/finish", handlerWrapper(routers.Finish))
http.HandleFunc("/api/v1/put", handlerWrapper(routers.Put))
http.HandleFunc("/api/v1/metric", metricProxy)
log.Printf("Server started")
go writeCount()
http.ListenAndServe(":1239", nil)
}