feat: add dead letter queue
All checks were successful
continuous-integration/drone/push Build is passing
All checks were successful
continuous-integration/drone/push Build is passing
This commit is contained in:
parent
1cf9d23491
commit
0b042780c3
@ -3,6 +3,7 @@ package app
|
|||||||
import (
|
import (
|
||||||
"log/slog"
|
"log/slog"
|
||||||
|
|
||||||
|
"git.front.kjuulh.io/kjuulh/orbis/internal/deadletter"
|
||||||
"git.front.kjuulh.io/kjuulh/orbis/internal/executor"
|
"git.front.kjuulh.io/kjuulh/orbis/internal/executor"
|
||||||
"git.front.kjuulh.io/kjuulh/orbis/internal/modelschedule"
|
"git.front.kjuulh.io/kjuulh/orbis/internal/modelschedule"
|
||||||
"git.front.kjuulh.io/kjuulh/orbis/internal/scheduler"
|
"git.front.kjuulh.io/kjuulh/orbis/internal/scheduler"
|
||||||
@ -48,7 +49,11 @@ func (a *App) WorkScheduler() *workscheduler.WorkScheduler {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (a *App) WorkProcessor() *workprocessor.WorkProcessor {
|
func (a *App) WorkProcessor() *workprocessor.WorkProcessor {
|
||||||
return workprocessor.NewWorkProcessor(a.WorkScheduler(), a.logger)
|
return workprocessor.NewWorkProcessor(a.WorkScheduler(), a.logger, a.DeadLetter())
|
||||||
|
}
|
||||||
|
|
||||||
|
func (a *App) DeadLetter() *deadletter.DeadLetter {
|
||||||
|
return deadletter.NewDeadLetter(Postgres(), a.logger)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (a *App) ModelSchedule() *modelschedule.ModelSchedule {
|
func (a *App) ModelSchedule() *modelschedule.ModelSchedule {
|
||||||
|
35
internal/deadletter/deadletter.go
Normal file
35
internal/deadletter/deadletter.go
Normal file
@ -0,0 +1,35 @@
|
|||||||
|
package deadletter
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"log/slog"
|
||||||
|
|
||||||
|
"git.front.kjuulh.io/kjuulh/orbis/internal/deadletter/repositories"
|
||||||
|
"github.com/google/uuid"
|
||||||
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
|
)
|
||||||
|
|
||||||
|
type DeadLetter struct {
|
||||||
|
db *pgxpool.Pool
|
||||||
|
logger *slog.Logger
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewDeadLetter(db *pgxpool.Pool, logger *slog.Logger) *DeadLetter {
|
||||||
|
return &DeadLetter{
|
||||||
|
db: db,
|
||||||
|
logger: logger,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (d *DeadLetter) InsertDeadLetter(ctx context.Context, schedule uuid.UUID) error {
|
||||||
|
repo := repositories.New(d.db)
|
||||||
|
|
||||||
|
d.logger.WarnContext(ctx, "deadlettering schedule", "schedule", schedule)
|
||||||
|
if err := repo.InsertDeadLetter(ctx, schedule); err != nil {
|
||||||
|
return fmt.Errorf("failed to insert item into dead letter: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
|
||||||
|
}
|
12
internal/deadletter/queries.sql
Normal file
12
internal/deadletter/queries.sql
Normal file
@ -0,0 +1,12 @@
|
|||||||
|
-- name: Ping :one
|
||||||
|
SELECT 1;
|
||||||
|
|
||||||
|
-- name: InsertDeadLetter :exec
|
||||||
|
INSERT INTO dead_letter
|
||||||
|
(
|
||||||
|
schedule_id
|
||||||
|
)
|
||||||
|
VALUES
|
||||||
|
(
|
||||||
|
$1
|
||||||
|
);
|
32
internal/deadletter/repositories/db.go
Normal file
32
internal/deadletter/repositories/db.go
Normal file
@ -0,0 +1,32 @@
|
|||||||
|
// Code generated by sqlc. DO NOT EDIT.
|
||||||
|
// versions:
|
||||||
|
// sqlc v1.23.0
|
||||||
|
|
||||||
|
package repositories
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5"
|
||||||
|
"github.com/jackc/pgx/v5/pgconn"
|
||||||
|
)
|
||||||
|
|
||||||
|
type DBTX interface {
|
||||||
|
Exec(context.Context, string, ...interface{}) (pgconn.CommandTag, error)
|
||||||
|
Query(context.Context, string, ...interface{}) (pgx.Rows, error)
|
||||||
|
QueryRow(context.Context, string, ...interface{}) pgx.Row
|
||||||
|
}
|
||||||
|
|
||||||
|
func New(db DBTX) *Queries {
|
||||||
|
return &Queries{db: db}
|
||||||
|
}
|
||||||
|
|
||||||
|
type Queries struct {
|
||||||
|
db DBTX
|
||||||
|
}
|
||||||
|
|
||||||
|
func (q *Queries) WithTx(tx pgx.Tx) *Queries {
|
||||||
|
return &Queries{
|
||||||
|
db: tx,
|
||||||
|
}
|
||||||
|
}
|
35
internal/deadletter/repositories/models.go
Normal file
35
internal/deadletter/repositories/models.go
Normal file
@ -0,0 +1,35 @@
|
|||||||
|
// Code generated by sqlc. DO NOT EDIT.
|
||||||
|
// versions:
|
||||||
|
// sqlc v1.23.0
|
||||||
|
|
||||||
|
package repositories
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/google/uuid"
|
||||||
|
"github.com/jackc/pgx/v5/pgtype"
|
||||||
|
)
|
||||||
|
|
||||||
|
type DeadLetter struct {
|
||||||
|
ScheduleID uuid.UUID `json:"schedule_id"`
|
||||||
|
UpdatedAt pgtype.Timestamptz `json:"updated_at"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type ModelSchedule struct {
|
||||||
|
ModelName string `json:"model_name"`
|
||||||
|
LastRun pgtype.Timestamptz `json:"last_run"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type WorkSchedule struct {
|
||||||
|
ScheduleID uuid.UUID `json:"schedule_id"`
|
||||||
|
WorkerID uuid.UUID `json:"worker_id"`
|
||||||
|
StartRun pgtype.Timestamptz `json:"start_run"`
|
||||||
|
EndRun pgtype.Timestamptz `json:"end_run"`
|
||||||
|
UpdatedAt pgtype.Timestamptz `json:"updated_at"`
|
||||||
|
State string `json:"state"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type WorkerRegister struct {
|
||||||
|
WorkerID uuid.UUID `json:"worker_id"`
|
||||||
|
Capacity int32 `json:"capacity"`
|
||||||
|
HeartBeat pgtype.Timestamptz `json:"heart_beat"`
|
||||||
|
}
|
18
internal/deadletter/repositories/querier.go
Normal file
18
internal/deadletter/repositories/querier.go
Normal file
@ -0,0 +1,18 @@
|
|||||||
|
// Code generated by sqlc. DO NOT EDIT.
|
||||||
|
// versions:
|
||||||
|
// sqlc v1.23.0
|
||||||
|
|
||||||
|
package repositories
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
|
||||||
|
"github.com/google/uuid"
|
||||||
|
)
|
||||||
|
|
||||||
|
type Querier interface {
|
||||||
|
InsertDeadLetter(ctx context.Context, scheduleID uuid.UUID) error
|
||||||
|
Ping(ctx context.Context) (int32, error)
|
||||||
|
}
|
||||||
|
|
||||||
|
var _ Querier = (*Queries)(nil)
|
39
internal/deadletter/repositories/queries.sql.go
Normal file
39
internal/deadletter/repositories/queries.sql.go
Normal file
@ -0,0 +1,39 @@
|
|||||||
|
// Code generated by sqlc. DO NOT EDIT.
|
||||||
|
// versions:
|
||||||
|
// sqlc v1.23.0
|
||||||
|
// source: queries.sql
|
||||||
|
|
||||||
|
package repositories
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
|
||||||
|
"github.com/google/uuid"
|
||||||
|
)
|
||||||
|
|
||||||
|
const insertDeadLetter = `-- name: InsertDeadLetter :exec
|
||||||
|
INSERT INTO dead_letter
|
||||||
|
(
|
||||||
|
schedule_id
|
||||||
|
)
|
||||||
|
VALUES
|
||||||
|
(
|
||||||
|
$1
|
||||||
|
)
|
||||||
|
`
|
||||||
|
|
||||||
|
func (q *Queries) InsertDeadLetter(ctx context.Context, scheduleID uuid.UUID) error {
|
||||||
|
_, err := q.db.Exec(ctx, insertDeadLetter, scheduleID)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
const ping = `-- name: Ping :one
|
||||||
|
SELECT 1
|
||||||
|
`
|
||||||
|
|
||||||
|
func (q *Queries) Ping(ctx context.Context) (int32, error) {
|
||||||
|
row := q.db.QueryRow(ctx, ping)
|
||||||
|
var column_1 int32
|
||||||
|
err := row.Scan(&column_1)
|
||||||
|
return column_1, err
|
||||||
|
}
|
21
internal/deadletter/sqlc.yaml
Normal file
21
internal/deadletter/sqlc.yaml
Normal file
@ -0,0 +1,21 @@
|
|||||||
|
version: "2"
|
||||||
|
sql:
|
||||||
|
- queries: queries.sql
|
||||||
|
schema: ../persistence/migrations/
|
||||||
|
engine: "postgresql"
|
||||||
|
gen:
|
||||||
|
go:
|
||||||
|
out: "repositories"
|
||||||
|
package: "repositories"
|
||||||
|
sql_package: "pgx/v5"
|
||||||
|
emit_json_tags: true
|
||||||
|
emit_prepared_queries: true
|
||||||
|
emit_interface: true
|
||||||
|
emit_empty_slices: true
|
||||||
|
emit_result_struct_pointers: true
|
||||||
|
emit_params_struct_pointers: true
|
||||||
|
overrides:
|
||||||
|
- db_type: "uuid"
|
||||||
|
go_type:
|
||||||
|
import: "github.com/google/uuid"
|
||||||
|
type: "UUID"
|
@ -9,6 +9,11 @@ import (
|
|||||||
"github.com/jackc/pgx/v5/pgtype"
|
"github.com/jackc/pgx/v5/pgtype"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
type DeadLetter struct {
|
||||||
|
ScheduleID uuid.UUID `json:"schedule_id"`
|
||||||
|
UpdatedAt pgtype.Timestamptz `json:"updated_at"`
|
||||||
|
}
|
||||||
|
|
||||||
type ModelSchedule struct {
|
type ModelSchedule struct {
|
||||||
ModelName string `json:"model_name"`
|
ModelName string `json:"model_name"`
|
||||||
LastRun pgtype.Timestamptz `json:"last_run"`
|
LastRun pgtype.Timestamptz `json:"last_run"`
|
||||||
|
@ -20,3 +20,8 @@ CREATE TABLE work_schedule (
|
|||||||
|
|
||||||
CREATE INDEX idx_work_schedule_worker ON work_schedule (worker_id);
|
CREATE INDEX idx_work_schedule_worker ON work_schedule (worker_id);
|
||||||
CREATE INDEX idx_work_schedule_worker_updated ON work_schedule (worker_id, updated_at DESC);
|
CREATE INDEX idx_work_schedule_worker_updated ON work_schedule (worker_id, updated_at DESC);
|
||||||
|
|
||||||
|
CREATE TABLE dead_letter (
|
||||||
|
schedule_id UUID PRIMARY KEY NOT NULL
|
||||||
|
, updated_at TIMESTAMPTZ NOT NULL default now()
|
||||||
|
);
|
||||||
|
@ -9,6 +9,11 @@ import (
|
|||||||
"github.com/jackc/pgx/v5/pgtype"
|
"github.com/jackc/pgx/v5/pgtype"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
type DeadLetter struct {
|
||||||
|
ScheduleID uuid.UUID `json:"schedule_id"`
|
||||||
|
UpdatedAt pgtype.Timestamptz `json:"updated_at"`
|
||||||
|
}
|
||||||
|
|
||||||
type ModelSchedule struct {
|
type ModelSchedule struct {
|
||||||
ModelName string `json:"model_name"`
|
ModelName string `json:"model_name"`
|
||||||
LastRun pgtype.Timestamptz `json:"last_run"`
|
LastRun pgtype.Timestamptz `json:"last_run"`
|
||||||
|
@ -4,8 +4,10 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
|
"math/rand/v2"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"git.front.kjuulh.io/kjuulh/orbis/internal/deadletter"
|
||||||
"git.front.kjuulh.io/kjuulh/orbis/internal/workscheduler"
|
"git.front.kjuulh.io/kjuulh/orbis/internal/workscheduler"
|
||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
)
|
)
|
||||||
@ -13,12 +15,18 @@ import (
|
|||||||
type WorkProcessor struct {
|
type WorkProcessor struct {
|
||||||
workscheduler *workscheduler.WorkScheduler
|
workscheduler *workscheduler.WorkScheduler
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
|
deadletter *deadletter.DeadLetter
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewWorkProcessor(workscheduler *workscheduler.WorkScheduler, logger *slog.Logger) *WorkProcessor {
|
func NewWorkProcessor(
|
||||||
|
workscheduler *workscheduler.WorkScheduler,
|
||||||
|
logger *slog.Logger,
|
||||||
|
deadletter *deadletter.DeadLetter,
|
||||||
|
) *WorkProcessor {
|
||||||
return &WorkProcessor{
|
return &WorkProcessor{
|
||||||
workscheduler: workscheduler,
|
workscheduler: workscheduler,
|
||||||
logger: logger,
|
logger: logger,
|
||||||
|
deadletter: deadletter,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -41,10 +49,17 @@ func (w *WorkProcessor) ProcessNext(ctx context.Context, workerID uuid.UUID) err
|
|||||||
}
|
}
|
||||||
|
|
||||||
time.Sleep(time.Millisecond * 10)
|
time.Sleep(time.Millisecond * 10)
|
||||||
|
// Process or achive
|
||||||
|
|
||||||
if err := w.workscheduler.Archive(ctx, *schedule); err != nil {
|
if err := w.workscheduler.Archive(ctx, *schedule); err != nil {
|
||||||
return fmt.Errorf("failed to archive item: %w", err)
|
return fmt.Errorf("failed to archive item: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if rand.Float64() > 0.9 {
|
||||||
|
if err := w.deadletter.InsertDeadLetter(ctx, *schedule); err != nil {
|
||||||
|
return fmt.Errorf("failed to deadletter message: %w", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
@ -9,6 +9,11 @@ import (
|
|||||||
"github.com/jackc/pgx/v5/pgtype"
|
"github.com/jackc/pgx/v5/pgtype"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
type DeadLetter struct {
|
||||||
|
ScheduleID uuid.UUID `json:"schedule_id"`
|
||||||
|
UpdatedAt pgtype.Timestamptz `json:"updated_at"`
|
||||||
|
}
|
||||||
|
|
||||||
type ModelSchedule struct {
|
type ModelSchedule struct {
|
||||||
ModelName string `json:"model_name"`
|
ModelName string `json:"model_name"`
|
||||||
LastRun pgtype.Timestamptz `json:"last_run"`
|
LastRun pgtype.Timestamptz `json:"last_run"`
|
||||||
|
Loading…
Reference in New Issue
Block a user