Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
97647ad
chore: add elasticsearch client
nhassl3 Jul 27, 2026
cdbe6b0
refactor: do embed sqlc structs and use standard map for Moderation a…
nhassl3 Jul 28, 2026
bf881df
refactor: move notification_consumer to transport layer
nhassl3 Jul 28, 2026
0dfd9c7
refactor: remove '&' from product object because link already provide…
nhassl3 Jul 28, 2026
49987e5
feat: add new commands for elasticsearch udpate in Makefile
nhassl3 Jul 28, 2026
fb35edc
feat: new configuration variables
nhassl3 Jul 28, 2026
cea27d3
refactor: global changes with kafka. Now how mush topic created in co…
nhassl3 Jul 28, 2026
914dbe9
refactor: change db to json structs tags because it's domain structs
nhassl3 Jul 28, 2026
0b3bb17
feat: add mocks for elasticsearch update
nhassl3 Jul 28, 2026
31c36b3
feat: add elasticsearch searching and cruds operation in elasticsearch
nhassl3 Jul 28, 2026
18db11a
feat: add interface for general transport layer handlers
nhassl3 Jul 28, 2026
2a78ca2
feat: add es-reindex bin file for re-index all products in database. …
nhassl3 Jul 28, 2026
fa41ed1
feat: add Elasticsearch client and repository with json managing
nhassl3 Jul 28, 2026
1f51b16
refactor: remove Dockerfile for consumer and start.sh for start entry…
nhassl3 Jul 28, 2026
2fe4eb7
feat: add BuildKit with cache for build files, mark .env like depreca…
nhassl3 Jul 28, 2026
621073b
feat: add .dockerignore for that the docker doesn't create a large co…
nhassl3 Jul 28, 2026
c316bae
feat: now for list and search response with []*domain.Product
nhassl3 Jul 28, 2026
7333576
feat: add fallback on PG search FTS when elasticsearch response with …
nhassl3 Jul 28, 2026
91bfabd
refactor: remove giving link for mapping proto responses for product …
nhassl3 Jul 28, 2026
7ba3baf
refactor: move log
nhassl3 Jul 28, 2026
2340eeb
refactor: remove ptr and put product from list
nhassl3 Jul 28, 2026
62c9635
fix: repair bug with date and bulk
nhassl3 Jul 28, 2026
cf425e9
fix: repair local Elasticsearch docker
nhassl3 Jul 28, 2026
ffd4ed3
refactor: change format for date with milliseconds
nhassl3 Jul 28, 2026
8c4381c
refactor: remove logs in search service layer
nhassl3 Jul 29, 2026
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
32 changes: 32 additions & 0 deletions .dockerignore
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
.git
.gitignore
.claude
.github


.idea
.vscode

tmp
bin
dist
build
doc

coverage
*.out
*.log

node_modules

Dockerfile*
docker-compose*.yml

README.md
AGENTS.md
CONTRIBUTING.md
LICENSE
Makefile
servicehub_backend_sshkey
servicehub_backend_sshkey.pub
sqlc.yaml
35 changes: 26 additions & 9 deletions Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,24 @@ RUN apk add --no-cache git gcc musl-dev

WORKDIR /app

COPY go.mod go.sum ./

RUN --mount=type=cache,target=/go/pkg/mod \
go mod download

COPY . .

RUN CGO_ENABLED=0 GOOS=linux go build -ldflags="-w -s" -o /bin/servicehub ./cmd/servicehub
RUN --mount=type=cache,target=/go/pkg/mod \
--mount=type=cache,target=/root/.cache/go-build \
CGO_ENABLED=0 GOOS=linux go build -ldflags="-w -s" -o /bin/servicehub ./cmd/servicehub

RUN --mount=type=cache,target=/go/pkg/mod \
--mount=type=cache,target=/root/.cache/go-build \
CGO_ENABLED=0 GOOS=linux go build -ldflags="-w -s" -o /bin/consumer ./cmd/consumer

RUN --mount=type=cache,target=/go/pkg/mod \
--mount=type=cache,target=/root/.cache/go-build \
CGO_ENABLED=0 GOOS=linux go build -ldflags="-w -s" -o /bin/es-reindex ./cmd/es-reindex

# Install migrate via go install (uses already-downloaded Go modules cache)
RUN go install -tags 'postgres' github.com/golang-migrate/migrate/v4/cmd/migrate@v4.19.1
Expand All @@ -20,19 +35,21 @@ RUN apk add --no-cache ca-certificates tzdata

WORKDIR /app

# Copy binaries
COPY --from=builder /bin/servicehub .
COPY --from=builder /app/migrations ./migrations
COPY --from=builder /bin/es-reindex .
COPY --from=builder /bin/consumer .
COPY --from=builder /go/bin/migrate ./migrate

# Copy configuration files (DO NOT include .env - it should be provided at runtime)
COPY --from=builder /app/migrations ./migrations
COPY --from=builder /app/config/prod.yaml config/prod.yaml
COPY --from=builder /app/config/local.yaml config/local.yaml
COPY --from=builder /app/config/dev.yaml config/dev.yaml
COPY --from=builder /app/.env .
COPY --from=builder /app/pkg/mailer/templates/*.html pkg/mailer/templates/
COPY --from=builder /app/start.sh .

RUN chmod +x /app/start.sh /app/migrate

EXPOSE 8082 50051
# Copy templates if they exist
COPY --from=builder /app/pkg/mailer/templates/*.html pkg/mailer/templates/
RUN chmod +x /app/migrate /app/servicehub /app/consumer /app/es-reindex

ENTRYPOINT ["/app/start.sh"]
CMD ["/app/servicehub"]
EXPOSE 8082 50051
20 changes: 17 additions & 3 deletions Makefile
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
.PHONY: build run runb test lint mock sqlc migrate-up migrate-down migrate-force clean docker-build postgres opendb dropdb createdb generate-data redis cli-redis minio minio-stop build-consumer run-consumer runb-consumer kafka-docker \
.els-docker
.els-docker els-docker-stop els-reindex-build runb-els-reindex

.DEFAULT_GOAL := build

Expand All @@ -19,6 +19,7 @@ BINARY_NAME=servicehub
BUILD_DIR=./bin
CMD_PATH=./cmd/servicehub
CONSUMER_PATH=./cmd/consumer
ES_REINDEX_PATH=./cmd/es-reindex
ENVIRONMENT=local

# Migrations
Expand Down Expand Up @@ -177,8 +178,21 @@ get-contracts:

## ─── Elasticsearch (FTS) ──────────────────────────────────────────────────────
els-docker:
@docker run -d --name servicehub-elasticsearch-local-9.3.8 \
@docker run -d --name servicehub-elasticsearch-local \
-p 9200:9200 \
-e "discovery.type=single-node" \
-e "xpack.security.enabled=false" \
elasticsearch:9.3.8
elasticsearch:9.3.8
@echo "Elasticsearch started on http://localhost:9200"

els-docker-stop:
@docker stop servicehub-elasticsearch-local && docker rm servicehub-elasticsearch-local
@echo "Elasticsearch stopped"

els-reindex-build:
go build -o $(BUILD_DIR)/$(BINARY_NAME)-els-reindex-$(GOOS)-$(GOARCH) $(ES_REINDEX_PATH)
@chmod +x $(BUILD_DIR)/$(BINARY_NAME)-els-reindex-$(GOOS)-$(GOARCH)
@echo "Successfully built es-reindex"

runb-els-reindex:
@ENVIRONMENT=$(ENVIRONMENT) ./$(BUILD_DIR)/$(BINARY_NAME)-els-reindex-$(GOOS)-$(GOARCH)
29 changes: 0 additions & 29 deletions cmd/consumer/Dockerfile

This file was deleted.

91 changes: 63 additions & 28 deletions cmd/consumer/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,15 +2,23 @@ package main

import (
"context"
"errors"
"fmt"
"os/signal"
"slices"
"sync"
"syscall"
"time"

"github.com/elastic/go-elasticsearch/v9"
"github.com/nhassl3/servicehub-backend/internal/domain"
elsRepo "github.com/nhassl3/servicehub-backend/internal/repository/elasticsearch"
elsPkg "github.com/nhassl3/servicehub-backend/pkg/elasticsearch"
"github.com/nhassl3/servicehub-backend/pkg/mailer"
"go.uber.org/zap"

"github.com/nhassl3/servicehub-backend/cmd"
transportKafka "github.com/nhassl3/servicehub-backend/internal/transport/kafka"
"github.com/nhassl3/servicehub-backend/pkg/kafka"
pkgkafka "github.com/nhassl3/servicehub-backend/pkg/kafka"
)
Expand All @@ -36,14 +44,15 @@ func main() {
// Топики создаются явно ДО подписки консьюмеров — иначе при auto.create.topics.enable=true
// consumer group может попытаться присоединиться к топику, которого физически ещё нет
// (он создастся позже, лениво, при первой публикации из backend), и не восстановится сама.
topics := []pkgkafka.TopicSpec{
{Name: cfg.Kafka.Topics.OrderEvents, NumPartitions: 3, ReplicationFactor: 1},
{Name: cfg.Kafka.Topics.TransactionEvents, NumPartitions: 3, ReplicationFactor: 1},
topics := make([]kafka.TopicSpec, 0, len(cfg.Kafka.Topics.Events))
for _, e := range cfg.Kafka.Topics.Events {
topics = append(topics, kafka.TopicSpec{Name: e, NumPartitions: 3, ReplicationFactor: 1})
}
if err := pkgkafka.EnsureTopics(ctx, cfg.Kafka.Brokers, topics, log); err != nil {
if err := kafka.EnsureTopics(ctx, cfg.Kafka.Brokers, topics, log); err != nil {
log.Fatal("kafka: failed to ensure topics exist, exiting for restart", zap.Error(err))
}

// SMTP connect
var (
smtpClient mailer.Notifier
err error
Expand All @@ -59,39 +68,65 @@ func main() {
}
}

// Elasticsearch connect
elsClient, err := elsPkg.New(ctx, cfg.ELS.Hosts, cfg.ELS.Username, cfg.ELS.Password)
if err != nil {
log.Fatal(err.Error())
}
defer func(elsClient *elasticsearch.Client) {
_ = elsClient.Close(ctx)
}(elsClient)

elsRepository := elsRepo.NewProductESRepo(elsClient, log)

dlqProducer := kafka.NewProducer(cfg.Kafka.Brokers, cfg.Kafka.Topics.NotificationsDLQ, log)
defer dlqProducer.Close()

orderConsumer := pkgkafka.NewConsumer(
cfg.Kafka.Brokers, cfg.Kafka.Topics.OrderEvents, cfg.Kafka.GroupID+"-order-events", log,
).WithDLQ(dlqProducer)
txConsumer := pkgkafka.NewConsumer(
cfg.Kafka.Brokers, cfg.Kafka.Topics.TransactionEvents, cfg.Kafka.GroupID+"-transaction-events", log,
)
defer orderConsumer.Close()
defer txConsumer.Close()
consumers := make(map[domain.Topic]*pkgkafka.Consumer, len(cfg.Kafka.Topics.Events))
for _, c := range cfg.Kafka.Topics.Events {
var consumer *pkgkafka.Consumer
if slices.Contains(domain.DLQTopics, domain.Topic(c)) {
consumer = pkgkafka.NewConsumer(
cfg.Kafka.Brokers, c, fmt.Sprintf("%s-%s", cfg.Kafka.GroupID, c), log,
).WithDLQ(dlqProducer)
} else {
consumer = pkgkafka.NewConsumer(
cfg.Kafka.Brokers, c, fmt.Sprintf("%s-%s", cfg.Kafka.GroupID, c), log,
)
}
consumers[domain.Topic(c)] = consumer
}

orderNotifConsumer := kafka.NewNotificationConsumer(orderConsumer, smtpClient, log)
txNotifConsumer := kafka.NewNotificationConsumer(txConsumer, smtpClient, log)
handlers := make([]transportKafka.ConsumerHandler, 0, len(consumers))
for t, consumer := range consumers {
switch t {
case domain.TopicOrderEvent, domain.TopicTransactionEvent:
handlers = append(handlers, transportKafka.NewNotificationConsumer(consumer, smtpClient, log))
case domain.TopicProductEvent:
handlers = append(handlers, transportKafka.NewProductConsumer(consumer, elsRepository, log))
}
}

var wg sync.WaitGroup
wg.Add(2) // len(cfg.Kafka.Topic)

go func() {
defer wg.Done()
if err := orderNotifConsumer.Run(ctx); err != nil {
log.Error("order notification consumer stopped", zap.Error(err))
}
}()
wg.Add(len(handlers))

go func() {
defer wg.Done()
if err := txNotifConsumer.Run(ctx); err != nil {
log.Error("transaction notification consumer stopped", zap.Error(err))
}
}()
for _, h := range handlers {
go func() {
defer wg.Done()
if err = h.Run(ctx); err != nil {
log.Error("kafka: consumer stopped", zap.Error(err))
}
}()
}

log.Info("kafka consumer service started", zap.String("env", cfg.Environment), zap.String("Mode", cfg.Log.Level))
wg.Wait()
var consumerErrors error
for _, consumer := range consumers {
consumerErrors = errors.Join(consumer.Close())
}
if consumerErrors != nil {
log.Fatal(consumerErrors.Error())
}
log.Info("kafka consumer service stopped gracefully")
}
69 changes: 69 additions & 0 deletions cmd/es-reindex/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
package main

import (
"context"
"log"

"github.com/nhassl3/servicehub-backend/cmd"
"github.com/nhassl3/servicehub-backend/internal/db"
"github.com/nhassl3/servicehub-backend/internal/domain"
repoES "github.com/nhassl3/servicehub-backend/internal/repository/elasticsearch"
repoPostgres "github.com/nhassl3/servicehub-backend/internal/repository/postgres"
pkgES "github.com/nhassl3/servicehub-backend/pkg/elasticsearch"
"github.com/nhassl3/servicehub-backend/pkg/postgres"
)

func main() {
ctx := context.Background()
cfg := cmd.MustLoadConfig()
logger := cmd.MustLoadLogger(cfg.Log.Level)
defer func() { _ = logger.Sync() }()

dsn := postgres.DSN(cfg.DB.Host, cfg.DB.Port, cfg.DB.User, cfg.DB.Password, cfg.DB.Name, cfg.DB.SSLMode)
pool, err := postgres.New(ctx, dsn)
if err != nil {
log.Fatalf("postgres: %s", err)
}
defer pool.Close()

store := db.NewStore(pool)
productRepo := repoPostgres.NewProductRepo(store)

esClient, err := pkgES.New(ctx, cfg.ELS.Hosts, cfg.ELS.Username, cfg.ELS.Password)
if err != nil {
log.Fatalf("elasticsearch: %s", err)
}
defer func() { _ = esClient.Close(context.Background()) }()

esProductRepo := repoES.NewProductESRepo(esClient, logger)

if err := esProductRepo.EnsureIndex(ctx); err != nil {
log.Fatalf("ensure index: %s", err)
}

var offset int32
const batchSize = 100

for {
products, _, err := productRepo.List(ctx, domain.ListProductsParams{
Limit: batchSize,
Offset: offset,
})
if err != nil {
log.Fatalf("list products: %s", err)
}

if len(products) == 0 {
break
}

if err := esProductRepo.BulkIndexProducts(ctx, products); err != nil {
log.Fatalf("bulk index: %s", err)
}

logger.Sugar().Infof("reindexed %d products (offset %d)", len(products), offset)
offset += batchSize
}

logger.Info("reindex complete")
}
10 changes: 8 additions & 2 deletions config/dev.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -22,14 +22,20 @@ redis:

kafka:
topics:
order_events: "order_events"
transaction_events: "transaction_events"
events:
- "order_events"
- "transaction_events"
- "product_events"
notifications: "notifications"
notifications_dlq: "notification_dlq"
brokers:
- "kafka:9092"
group_id: "servicehub-backend"

els:
hosts:
- "http://elasticsearch:9200"

auth:
access_token_ttl: 15m
refresh_token_ttl: 168h
Expand Down
Loading
Loading