Compare commits

...

2 Commits

Author SHA1 Message Date
cuddle-please
7f37bda8b8 chore(release): 0.1.0
All checks were successful
continuous-integration/drone/push Build is passing
continuous-integration/drone/pr Build is passing
2025-01-19 10:34:00 +00:00
0b042780c3
feat: add dead letter queue
All checks were successful
continuous-integration/drone/push Build is passing
2025-01-19 11:33:39 +01:00
14 changed files with 276 additions and 2 deletions

42
CHANGELOG.md Normal file
View File

@ -0,0 +1,42 @@
# Changelog
All notable changes to this project will be documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
## [Unreleased]
## [0.1.0] - 2025-01-19
### Added
- add dead letter queue
- move schedules to registered workers
- prune old workers
- deregister worker on close
- add worker distributor and model registry
- enable worker process
- add migration
- add executor (#3)
Adds an executor which can process and dispatch events to a set of workers.
Co-authored-by: kjuulh <contact@kjuulh.io>
Co-committed-by: kjuulh <contact@kjuulh.io>
- add basic main.go
- add default
### Fixed
- use dedicated connection for scheduler process
- *(deps)* update module gitlab.com/greyxor/slogor to v1.6.1
- orbis demo
### Other
- add more specific log for when leader is acquired
- add basic leader election on top of postgres
- add orbis demo
- add basic scheduler
- add utility scripts
- add utility scripts
- add logger
- add cmd

View File

@ -3,6 +3,7 @@ package app
import (
"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/modelschedule"
"git.front.kjuulh.io/kjuulh/orbis/internal/scheduler"
@ -48,7 +49,11 @@ func (a *App) WorkScheduler() *workscheduler.WorkScheduler {
}
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 {

View 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
}

View File

@ -0,0 +1,12 @@
-- name: Ping :one
SELECT 1;
-- name: InsertDeadLetter :exec
INSERT INTO dead_letter
(
schedule_id
)
VALUES
(
$1
);

View 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,
}
}

View 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"`
}

View 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)

View 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
}

View 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"

View File

@ -9,6 +9,11 @@ import (
"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"`

View File

@ -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_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()
);

View File

@ -9,6 +9,11 @@ import (
"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"`

View File

@ -4,8 +4,10 @@ import (
"context"
"fmt"
"log/slog"
"math/rand/v2"
"time"
"git.front.kjuulh.io/kjuulh/orbis/internal/deadletter"
"git.front.kjuulh.io/kjuulh/orbis/internal/workscheduler"
"github.com/google/uuid"
)
@ -13,12 +15,18 @@ import (
type WorkProcessor struct {
workscheduler *workscheduler.WorkScheduler
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{
workscheduler: workscheduler,
logger: logger,
deadletter: deadletter,
}
}
@ -41,10 +49,17 @@ func (w *WorkProcessor) ProcessNext(ctx context.Context, workerID uuid.UUID) err
}
time.Sleep(time.Millisecond * 10)
// Process or achive
if err := w.workscheduler.Archive(ctx, *schedule); err != nil {
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
}

View File

@ -9,6 +9,11 @@ import (
"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"`