Compare commits
11 Commits
prod
...
71db409b15
| Author | SHA1 | Date | |
|---|---|---|---|
| 71db409b15 | |||
| 153943b3ed | |||
| b17ed7e7a9 | |||
| a564621d80 | |||
| 59869821f3 | |||
| d9dcfc1a3c | |||
| 932c9b609e | |||
| 71e40b932c | |||
| 39829a7348 | |||
| 9d9c11cf11 | |||
| 75a59c3e5c |
@@ -6,11 +6,30 @@ services:
|
||||
image: mathwave/sprint-repo:queues
|
||||
networks:
|
||||
- queues-development
|
||||
- monitoring
|
||||
- mongo-development
|
||||
- queues-mongo-development
|
||||
environment:
|
||||
MONGO_HOST: "mongo.develop.sprinthub.ru"
|
||||
MONGO_PASSWORD: $MONGO_PASSWORD_DEV
|
||||
REDIS_HOST: "redis.develop.sprinthub.ru"
|
||||
REDIS_PASSWORD: $REDIS_PASSWORD_DEV
|
||||
STAGE: "development"
|
||||
deploy:
|
||||
mode: replicated
|
||||
restart_policy:
|
||||
condition: any
|
||||
update_config:
|
||||
parallelism: 1
|
||||
order: start-first
|
||||
|
||||
storage:
|
||||
image: mongo:6.0.2
|
||||
networks:
|
||||
- queues-mongo-development
|
||||
volumes:
|
||||
- /sprint-data/queues-mongo:/data/db
|
||||
environment:
|
||||
MONGO_INITDB_ROOT_USERNAME: mongo
|
||||
MONGO_INITDB_ROOT_PASSWORD: password
|
||||
deploy:
|
||||
mode: replicated
|
||||
restart_policy:
|
||||
@@ -24,7 +43,5 @@ services:
|
||||
networks:
|
||||
queues-development:
|
||||
external: true
|
||||
monitoring:
|
||||
external: true
|
||||
mongo-development:
|
||||
external: true
|
||||
queues-mongo-development:
|
||||
driver: overlay
|
||||
|
||||
@@ -6,17 +6,36 @@ services:
|
||||
image: mathwave/sprint-repo:queues
|
||||
networks:
|
||||
- queues
|
||||
- monitoring
|
||||
- mongo
|
||||
- queues-mongo-production
|
||||
environment:
|
||||
MONGO_HOST: "mongo.sprinthub.ru"
|
||||
MONGO_PASSWORD: $MONGO_PASSWORD_PROD
|
||||
REDIS_HOST: "redis.sprinthub.ru"
|
||||
REDIS_PASSWORD: $REDIS_PASSWORD_PROD
|
||||
STAGE: "production"
|
||||
deploy:
|
||||
mode: replicated
|
||||
restart_policy:
|
||||
condition: any
|
||||
update_config:
|
||||
parallelism: 1
|
||||
order: start-first
|
||||
|
||||
storage:
|
||||
image: mongo:6.0.2
|
||||
networks:
|
||||
- queues-mongo-production
|
||||
volumes:
|
||||
- /sprint-data/queues-mongo:/data/db
|
||||
environment:
|
||||
MONGO_INITDB_ROOT_USERNAME: mongo
|
||||
MONGO_INITDB_ROOT_PASSWORD: password
|
||||
deploy:
|
||||
mode: replicated
|
||||
restart_policy:
|
||||
condition: any
|
||||
placement:
|
||||
constraints: [node.labels.stage == production]
|
||||
constraints: [node.labels.stage == development]
|
||||
update_config:
|
||||
parallelism: 1
|
||||
order: start-first
|
||||
@@ -24,7 +43,5 @@ services:
|
||||
networks:
|
||||
queues:
|
||||
external: true
|
||||
monitoring:
|
||||
external: true
|
||||
mongo:
|
||||
external: true
|
||||
queues-mongo-production:
|
||||
driver: overlay
|
||||
|
||||
@@ -26,6 +26,13 @@ jobs:
|
||||
steps:
|
||||
- name: push
|
||||
run: docker push mathwave/sprint-repo:queues
|
||||
create_dir:
|
||||
name: Create dir
|
||||
runs-on: [ dev ]
|
||||
needs: build
|
||||
steps:
|
||||
- name: create_dir
|
||||
run: mkdir /sprint-data/queues-mongo || true
|
||||
deploy-dev:
|
||||
name: Deploy dev
|
||||
runs-on: [prod]
|
||||
|
||||
@@ -26,6 +26,13 @@ jobs:
|
||||
steps:
|
||||
- name: push
|
||||
run: docker push mathwave/sprint-repo:queues
|
||||
create_dir:
|
||||
name: Create dir
|
||||
runs-on: [ prod ]
|
||||
needs: build
|
||||
steps:
|
||||
- name: create_dir
|
||||
run: mkdir /sprint-data/queues-mongo || true
|
||||
deploy-prod:
|
||||
name: Deploy prod
|
||||
runs-on: [prod]
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -9,7 +9,6 @@ import (
|
||||
"go.mongodb.org/mongo-driver/bson"
|
||||
"go.mongodb.org/mongo-driver/bson/primitive"
|
||||
"go.mongodb.org/mongo-driver/mongo"
|
||||
"go.mongodb.org/mongo-driver/mongo/options"
|
||||
)
|
||||
|
||||
type Task struct {
|
||||
@@ -33,14 +32,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 {
|
||||
@@ -71,19 +62,7 @@ func Take(queue string) (*Task, error) {
|
||||
if task == nil {
|
||||
return nil, nil
|
||||
}
|
||||
_, err = collection().UpdateByID(
|
||||
context.TODO(),
|
||||
task.Id,
|
||||
bson.M{
|
||||
"$set": bson.M{
|
||||
"taken_at": now,
|
||||
"attempts": task.Attempts + 1,
|
||||
"available_from": now.Add(
|
||||
time.Duration(task.SecondsToExecute+task.Attempts) * time.Second,
|
||||
),
|
||||
},
|
||||
},
|
||||
)
|
||||
_, err = collection().UpdateByID(context.TODO(), task.Id, bson.M{"$set": bson.M{"taken_at": now, "attempts": task.Attempts + 1}})
|
||||
if err != nil {
|
||||
println("Error updaing")
|
||||
println(err.Error())
|
||||
@@ -99,7 +78,6 @@ 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),
|
||||
)
|
||||
if err != nil {
|
||||
println("Error find")
|
||||
@@ -116,8 +94,15 @@ func findTask(queue string, now time.Time) (*Task, error) {
|
||||
}
|
||||
|
||||
for _, task := range results {
|
||||
if task.TakenAt == nil {
|
||||
return &task, nil
|
||||
}
|
||||
takenAt := *task.TakenAt
|
||||
|
||||
if takenAt.Add(time.Second * time.Duration(task.SecondsToExecute)).Before(now) {
|
||||
return &task, nil
|
||||
}
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
|
||||
56
main.go
56
main.go
@@ -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)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user