Skip to content

Commit e86fc61

Browse files
committed
update to panic instead of fatal, add more tests, fix job tests
1 parent 6a7acfe commit e86fc61

11 files changed

Lines changed: 435 additions & 99 deletions

database/dbJob.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -97,7 +97,7 @@ func (r JobDBHandler) CreateTable() error {
9797
SELECT create_hypertable('job_archive', by_range('updated_at'), if_not_exists => TRUE);`,
9898
)
9999
if err != nil {
100-
log.Fatalf("error creating job table: %#v", err)
100+
log.Panicf("error creating job table: %#v", err)
101101
}
102102

103103
_, err = r.db.Instance.ExecContext(
@@ -107,7 +107,7 @@ func (r JobDBHandler) CreateTable() error {
107107
FOR EACH ROW EXECUTE PROCEDURE notify_event();`,
108108
)
109109
if err != nil {
110-
log.Fatalf("error creating notify trigger on job table: %#v", err)
110+
log.Panicf("error creating notify trigger on job table: %#v", err)
111111
}
112112

113113
err = r.db.CreateIndexes("job", "worker_id", "worker_rid", "status", "created_at", "updated_at") // Indexes on common search/filter fields

database/dbWorker.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,7 @@ func (r WorkerDBHandler) CreateTable() error {
7878
)`,
7979
)
8080
if err != nil {
81-
log.Fatalf("error creating worker table: %#v", err)
81+
log.Panicf("error creating worker table: %#v", err)
8282
}
8383

8484
err = r.db.CreateIndexes("worker", "rid", "name", "status")

helper/database.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@ func NewDatabase(name string, dbConfig *DatabaseConfiguration) *Database {
3434

3535
err := db.AddNotifyFunction()
3636
if err != nil {
37-
logger.Fatalf("failed to add notify function: %v", err)
37+
logger.Panicf("failed to add notify function: %v", err)
3838
}
3939

4040
return db
@@ -87,7 +87,7 @@ func (d *Database) ConnectToDatabase(dbConfig *DatabaseConfiguration, logger *lo
8787
connectOnce.Do(func() {
8888
db, err = sql.Open("postgres", dbConfig.DatabaseConnectionString())
8989
if err != nil {
90-
logger.Fatalf("error establishing connection to db: %v. Trying again.", err.Error())
90+
logger.Panicf("error establishing connection to db: %v. Trying again.", err.Error())
9191
}
9292

9393
db.SetMaxOpenConns(10)
@@ -307,7 +307,7 @@ func (d *Database) Health() map[string]string {
307307
if err != nil {
308308
stats["status"] = "down"
309309
stats["error"] = fmt.Sprintf("db down: %v", err)
310-
log.Fatalf("db down: %v", err) // Log the error and terminate the program
310+
log.Panicf("db down: %v", err) // Log the error and terminate the program
311311
return stats
312312
}
313313

main_test.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,8 @@ import (
1313
"github.com/testcontainers/testcontainers-go"
1414
)
1515

16+
const maxDeviation = 50 * time.Millisecond
17+
1618
var dbPort string
1719

1820
func TestMain(m *testing.M) {

queuer.go

Lines changed: 21 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -35,14 +35,19 @@ type Queuer struct {
3535
log *log.Logger
3636
}
3737

38+
// NewQueuer creates a new Queuer instance with the given name and max concurrency.
39+
// It initializes the database connection, job listeners, and worker.
40+
// If options are provided, it creates a worker with those options.
41+
// If any error occurs during initialization, it logs a panic error and exits the program.
42+
// It returns a pointer to the newly created Queuer instance.
3843
func NewQueuer(name string, maxConcurrency int, options ...*model.OnError) *Queuer {
3944
// Logger
4045
logger := log.New(os.Stdout, "Queuer: ", log.Ltime)
4146

4247
// Database
4348
dbConfig, err := helper.NewDatabaseConfiguration()
4449
if err != nil {
45-
logger.Fatalf("failed to create database configuration: %v", err)
50+
logger.Panicf("failed to create database configuration: %v", err)
4651
}
4752
dbConnection := helper.NewDatabase(
4853
"queuer",
@@ -54,44 +59,44 @@ func NewQueuer(name string, maxConcurrency int, options ...*model.OnError) *Queu
5459
var dbWorker database.WorkerDBHandlerFunctions
5560
dbJob, err = database.NewJobDBHandler(dbConnection)
5661
if err != nil {
57-
logger.Fatalf("failed to create job db handler: %v", err)
62+
logger.Panicf("failed to create job db handler: %v", err)
5863
}
5964
dbWorker, err = database.NewWorkerDBHandler(dbConnection)
6065
if err != nil {
61-
logger.Fatalf("failed to create worker db handler: %v", err)
66+
logger.Panicf("failed to create worker db handler: %v", err)
6267
}
6368

6469
// Job listeners
6570
jobInsertListener, err := database.NewQueuerListener(dbConfig, "job.INSERT")
6671
if err != nil {
67-
logger.Fatalf("failed to create job insert listener: %v", err)
72+
logger.Panicf("failed to create job insert listener: %v", err)
6873
}
6974
jobUpdateListener, err := database.NewQueuerListener(dbConfig, "job.UPDATE")
7075
if err != nil {
71-
logger.Fatalf("failed to create job update listener: %v", err)
76+
logger.Panicf("failed to create job update listener: %v", err)
7277
}
7378
jobDeleteListener, err := database.NewQueuerListener(dbConfig, "job.DELETE")
7479
if err != nil {
75-
logger.Fatalf("failed to create job delete listener: %v", err)
80+
logger.Panicf("failed to create job delete listener: %v", err)
7681
}
7782

7883
// Inserting worker
7984
var newWorker *model.Worker
8085
if len(options) > 0 {
8186
newWorker, err = model.NewWorkerWithOptions(name, maxConcurrency, options[0])
8287
if err != nil {
83-
logger.Fatalf("error creating new worker with options: %v", err)
88+
logger.Panicf("error creating new worker with options: %v", err)
8489
}
8590
} else {
8691
newWorker, err = model.NewWorker(name, maxConcurrency)
8792
if err != nil {
88-
logger.Fatalf("error creating new worker: %v", err)
93+
logger.Panicf("error creating new worker: %v", err)
8994
}
9095
}
9196

9297
worker, err := dbWorker.InsertWorker(newWorker)
9398
if err != nil {
94-
logger.Fatalf("error inserting worker: %v", err)
99+
logger.Panicf("error inserting worker: %v", err)
95100
}
96101
logger.Printf("Worker %s created with RID %s", worker.Name, worker.RID.String())
97102

@@ -108,9 +113,12 @@ func NewQueuer(name string, maxConcurrency int, options ...*model.OnError) *Queu
108113
}
109114
}
110115

116+
// Start starts the queuer by initializing the job listeners and starting the job poll ticker.
117+
// It checks if the queuer is initialized properly, and if not, it logs a panic error and exits the program.
118+
// It runs the job processing in a separate goroutine and listens for job events.
111119
func (q *Queuer) Start(ctx context.Context, cancel context.CancelFunc) {
112120
if q.dbJob == nil || q.dbWorker == nil || q.jobInsertListener == nil || q.jobUpdateListener == nil || q.jobDeleteListener == nil {
113-
q.log.Fatalln("worker is not initialized properly")
121+
q.log.Panicln("worker is not initialized properly")
114122
}
115123

116124
q.ctx = ctx
@@ -135,6 +143,8 @@ func (q *Queuer) Start(ctx context.Context, cancel context.CancelFunc) {
135143
}()
136144
}
137145

146+
// Stop stops the queuer by closing the job listeners, cancelling all queued and running jobs,
147+
// and cancelling the context to stop the queuer.
138148
func (q *Queuer) Stop() error {
139149
if q.jobInsertListener != nil {
140150
err := q.jobInsertListener.Listener.Close()
@@ -156,7 +166,7 @@ func (q *Queuer) Stop() error {
156166
}
157167

158168
// Cancel all queued and running jobs
159-
err := q.CancelAllJobsByWorker(q.worker.RID)
169+
err := q.CancelAllJobsByWorker(q.worker.RID, 100)
160170
if err != nil {
161171
return fmt.Errorf("error cancelling all jobs by worker: %v", err)
162172
}

queuerJob.go

Lines changed: 12 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -80,8 +80,8 @@ func (q *Queuer) CancelJob(jobRid uuid.UUID) (*model.Job, error) {
8080
return job, nil
8181
}
8282

83-
func (q *Queuer) CancelAllJobsByWorker(workerRid uuid.UUID) error {
84-
jobs, err := q.dbJob.SelectAllJobsByWorkerRID(workerRid, 0, 0)
83+
func (q *Queuer) CancelAllJobsByWorker(workerRid uuid.UUID, entries int) error {
84+
jobs, err := q.dbJob.SelectAllJobsByWorkerRID(workerRid, 0, entries)
8585
if err != nil {
8686
return fmt.Errorf("error selecting jobs by worker RID %v: %v", workerRid, err)
8787
}
@@ -97,21 +97,21 @@ func (q *Queuer) CancelAllJobsByWorker(workerRid uuid.UUID) error {
9797
}
9898

9999
// ReaddJobFromArchive readds a job from the archive back to the queue.
100-
func (q *Queuer) ReaddJobFromArchive(jobRid uuid.UUID) error {
100+
func (q *Queuer) ReaddJobFromArchive(jobRid uuid.UUID) (*model.Job, error) {
101101
job, err := q.dbJob.SelectJobFromArchive(jobRid)
102102
if err != nil {
103-
return fmt.Errorf("error selecting job from archive with rid %v: %v", jobRid, err)
103+
return nil, fmt.Errorf("error selecting job from archive with rid %v: %v", jobRid, err)
104104
}
105105

106106
// Readd the job to the queue
107107
newJob, err := q.AddJobWithOptions(job.Options, job.TaskName, job.Parameters...)
108108
if err != nil {
109-
return fmt.Errorf("error readding job: %v", err)
109+
return nil, fmt.Errorf("error readding job: %v", err)
110110
}
111111

112112
q.log.Printf("Job readded with RID %v", newJob.RID)
113113

114-
return nil
114+
return newJob, nil
115115
}
116116

117117
// GetJob retrieves a job by its RID.
@@ -285,8 +285,8 @@ func (q *Queuer) scheduleJob(job *model.Job) {
285285
}
286286

287287
func (q *Queuer) cancelJob(job *model.Job) error {
288-
var err error
289-
if job.Status == model.JobStatusRunning {
288+
switch job.Status {
289+
case model.JobStatusRunning:
290290
jobRunner, found := q.activeRunners.Load(job.RID)
291291
if !found {
292292
return fmt.Errorf("job with rid %v not found or not running", job.RID)
@@ -295,16 +295,16 @@ func (q *Queuer) cancelJob(job *model.Job) error {
295295
runner := jobRunner.(*core.Runner)
296296
runner.Cancel(func() {
297297
job.Status = model.JobStatusCancelled
298-
job, err = q.dbJob.UpdateJobFinal(job)
298+
_, err := q.dbJob.UpdateJobFinal(job)
299299
if err != nil {
300300
q.log.Printf("error updating job status to cancelled: %v", err)
301301
}
302302
})
303-
} else if job.Status == model.JobStatusScheduled || job.Status == model.JobStatusQueued {
303+
case model.JobStatusScheduled, model.JobStatusQueued:
304304
job.Status = model.JobStatusCancelled
305-
job, err = q.dbJob.UpdateJobFinal(job)
305+
_, err := q.dbJob.UpdateJobFinal(job)
306306
if err != nil {
307-
return fmt.Errorf("error updating job status to cancelled: %v", err)
307+
q.log.Printf("error updating job status to cancelled: %v", err)
308308
}
309309
}
310310
return nil

0 commit comments

Comments
 (0)