Skip to content

Commit 526067f

Browse files
authored
add keyed parameters functionality (#17)
1 parent f3aa68f commit 526067f

13 files changed

Lines changed: 193 additions & 139 deletions

File tree

database/dbJob.go

Lines changed: 25 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -143,10 +143,11 @@ func (r JobDBHandler) DropTables() error {
143143
func (r JobDBHandler) InsertJob(job *model.Job) (*model.Job, error) {
144144
newJob := &model.Job{}
145145
row := r.db.Instance.QueryRow(
146-
`SELECT * FROM insert_job($1, $2, $3, $4, $5, $6)`,
146+
`SELECT * FROM insert_job($1, $2, $3, $4, $5, $6, $7)`,
147147
job.Options,
148148
job.TaskName,
149149
job.Parameters,
150+
job.ParametersKeyed,
150151
job.Status,
151152
job.ScheduledAt,
152153
job.ScheduleCount,
@@ -160,6 +161,7 @@ func (r JobDBHandler) InsertJob(job *model.Job) (*model.Job, error) {
160161
&newJob.Options,
161162
&newJob.TaskName,
162163
&newJob.Parameters,
164+
&newJob.ParametersKeyed,
163165
&newJob.Status,
164166
&newJob.ScheduledAt,
165167
&newJob.ScheduleCount,
@@ -178,10 +180,11 @@ func (r JobDBHandler) InsertJob(job *model.Job) (*model.Job, error) {
178180
func (r JobDBHandler) InsertJobTx(tx *sql.Tx, job *model.Job) (*model.Job, error) {
179181
newJob := &model.Job{}
180182
row := tx.QueryRow(
181-
`SELECT * FROM insert_job($1, $2, $3, $4, $5, $6)`,
183+
`SELECT * FROM insert_job($1, $2, $3, $4, $5, $6, $7)`,
182184
job.Options,
183185
job.TaskName,
184186
job.Parameters,
187+
job.ParametersKeyed,
185188
job.Status,
186189
job.ScheduledAt,
187190
job.ScheduleCount,
@@ -195,6 +198,7 @@ func (r JobDBHandler) InsertJobTx(tx *sql.Tx, job *model.Job) (*model.Job, error
195198
&newJob.Options,
196199
&newJob.TaskName,
197200
&newJob.Parameters,
201+
&newJob.ParametersKeyed,
198202
&newJob.Status,
199203
&newJob.ScheduledAt,
200204
&newJob.ScheduleCount,
@@ -220,7 +224,7 @@ func (r JobDBHandler) BatchInsertJobs(jobs []*model.Job) error {
220224
return helper.NewError("transaction start", err)
221225
}
222226

223-
stmt, err := tx.Prepare(pq.CopyIn("job", "options", "task_name", "parameters", "scheduled_at"))
227+
stmt, err := tx.Prepare(pq.CopyIn("job", "options", "task_name", "parameters", "parameters_keyed", "scheduled_at"))
224228
if err != nil {
225229
return helper.NewError("statement preparation", err)
226230
}
@@ -229,6 +233,7 @@ func (r JobDBHandler) BatchInsertJobs(jobs []*model.Job) error {
229233
var err error
230234
optionsJSON := []byte("{}")
231235
parametersJSON := []byte("[]")
236+
parametersKeyedJSON := []byte("{}")
232237

233238
if job.Options != nil {
234239
optionsJSON, err = job.Options.Marshal()
@@ -244,10 +249,18 @@ func (r JobDBHandler) BatchInsertJobs(jobs []*model.Job) error {
244249
}
245250
}
246251

252+
if job.ParametersKeyed != nil {
253+
parametersKeyedJSON, err = job.ParametersKeyed.Marshal()
254+
if err != nil {
255+
return helper.NewError("marshaling job parameters_keyed", err)
256+
}
257+
}
258+
247259
_, err = stmt.Exec(
248260
string(optionsJSON),
249261
job.TaskName,
250262
string(parametersJSON),
263+
string(parametersKeyedJSON),
251264
job.ScheduledAt,
252265
)
253266
if err != nil {
@@ -304,6 +317,7 @@ func (r JobDBHandler) UpdateJobsInitial(worker *model.Worker) ([]*model.Job, err
304317
&job.Options,
305318
&job.TaskName,
306319
&job.Parameters,
320+
&job.ParametersKeyed,
307321
&job.Status,
308322
&job.ScheduledAt,
309323
&job.StartedAt,
@@ -364,6 +378,7 @@ func (r JobDBHandler) UpdateJobFinal(job *model.Job) (*model.Job, error) {
364378
&archivedJob.Options,
365379
&archivedJob.TaskName,
366380
&archivedJob.Parameters,
381+
&archivedJob.ParametersKeyed,
367382
&archivedJob.Status,
368383
&archivedJob.ScheduledAt,
369384
&archivedJob.StartedAt,
@@ -418,6 +433,7 @@ func (r JobDBHandler) SelectJob(rid uuid.UUID) (*model.Job, error) {
418433
&job.Options,
419434
&job.TaskName,
420435
&job.Parameters,
436+
&job.ParametersKeyed,
421437
&job.Status,
422438
&job.ScheduledAt,
423439
&job.StartedAt,
@@ -461,6 +477,7 @@ func (r JobDBHandler) SelectAllJobs(lastID int, entries int) ([]*model.Job, erro
461477
&job.Options,
462478
&job.TaskName,
463479
&job.Parameters,
480+
&job.ParametersKeyed,
464481
&job.Status,
465482
&job.ScheduledAt,
466483
&job.StartedAt,
@@ -513,6 +530,7 @@ func (r JobDBHandler) SelectAllJobsByWorkerRID(workerRid uuid.UUID, lastID int,
513530
&job.Options,
514531
&job.TaskName,
515532
&job.Parameters,
533+
&job.ParametersKeyed,
516534
&job.Status,
517535
&job.ScheduledAt,
518536
&job.StartedAt,
@@ -569,6 +587,7 @@ func (r JobDBHandler) SelectAllJobsBySearch(search string, lastID int, entries i
569587
&job.Options,
570588
&job.TaskName,
571589
&job.Parameters,
590+
&job.ParametersKeyed,
572591
&job.Status,
573592
&job.ScheduledAt,
574593
&job.StartedAt,
@@ -653,6 +672,7 @@ func (r JobDBHandler) SelectJobFromArchive(rid uuid.UUID) (*model.Job, error) {
653672
&job.Options,
654673
&job.TaskName,
655674
&job.Parameters,
675+
&job.ParametersKeyed,
656676
&job.Status,
657677
&job.ScheduledAt,
658678
&job.StartedAt,
@@ -696,6 +716,7 @@ func (r JobDBHandler) SelectAllJobsFromArchive(lastID int, entries int) ([]*mode
696716
&job.Options,
697717
&job.TaskName,
698718
&job.Parameters,
719+
&job.ParametersKeyed,
699720
&job.Status,
700721
&job.ScheduledAt,
701722
&job.StartedAt,
@@ -749,6 +770,7 @@ func (r JobDBHandler) SelectAllJobsFromArchiveBySearch(search string, lastID int
749770
&job.Options,
750771
&job.TaskName,
751772
&job.Parameters,
773+
&job.ParametersKeyed,
752774
&job.Status,
753775
&job.ScheduledAt,
754776
&job.StartedAt,

database/dbJob_test.go

Lines changed: 23 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -100,7 +100,7 @@ func TestJobInsertJob(t *testing.T) {
100100
jobDbHandler, err := NewJobDBHandler(database, true)
101101
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
102102

103-
job, err := model.NewJob("TestTask", nil)
103+
job, err := model.NewJob("TestTask", nil, nil)
104104
require.NoError(t, err, "Expected NewJob to not return an error")
105105

106106
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -124,7 +124,7 @@ func TestJobInsertJobTx(t *testing.T) {
124124
jobDbHandler, err := NewJobDBHandler(database, true)
125125
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
126126

127-
job, err := model.NewJob("TestTask", nil)
127+
job, err := model.NewJob("TestTask", nil, nil)
128128
require.NoError(t, err, "Expected NewJob to not return an error")
129129

130130
tx, err := database.Instance.Begin()
@@ -158,7 +158,7 @@ func TestJobBatchInsertJobs(t *testing.T) {
158158
jobCount := 5
159159
jobs := []*model.Job{}
160160
for i := 0; i < jobCount; i++ {
161-
job, err := model.NewJob("TestTask", nil)
161+
job, err := model.NewJob("TestTask", nil, nil)
162162
require.NoError(t, err, "Expected NewJob to not return an error")
163163
require.NotNil(t, job, "Expected NewJob to return a non-nil job")
164164
jobs = append(jobs, job)
@@ -179,7 +179,7 @@ func TestJobBatchInsertJobs(t *testing.T) {
179179
RetryBackoff: model.RETRY_BACKOFF_EXPONENTIAL,
180180
MaxRetries: 3,
181181
},
182-
})
182+
}, nil)
183183
require.NoError(t, err, "Expected NewJob to not return an error")
184184
require.NotNil(t, job, "Expected NewJob to return a non-nil job")
185185
jobs = append(jobs, job)
@@ -199,7 +199,7 @@ func TestJobBatchInsertJobs(t *testing.T) {
199199
Interval: time.Second * 5,
200200
MaxCount: 10,
201201
},
202-
})
202+
}, nil)
203203
require.NoError(t, err, "Expected NewJob to not return an error")
204204
require.NotNil(t, job, "Expected NewJob to return a non-nil job")
205205
jobs = append(jobs, job)
@@ -213,7 +213,7 @@ func TestJobBatchInsertJobs(t *testing.T) {
213213
jobCount := 5
214214
jobs := []*model.Job{}
215215
for i := 0; i < jobCount; i++ {
216-
job, err := model.NewJob("TestTask", nil)
216+
job, err := model.NewJob("TestTask", nil, nil)
217217
job.Parameters = []interface{}{i, fmt.Sprintf("param-%d", i)}
218218
require.NoError(t, err, "Expected NewJob to not return an error")
219219
require.NotNil(t, job, "Expected NewJob to return a non-nil job")
@@ -252,7 +252,7 @@ func TestJobUpdateJobsInitial(t *testing.T) {
252252
jobDbHandler, err := NewJobDBHandler(database, true)
253253
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
254254

255-
job, err := model.NewJob("TestTask", nil)
255+
job, err := model.NewJob("TestTask", nil, nil)
256256
require.NoError(t, err, "Expected NewJob to not return an error")
257257

258258
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -279,7 +279,7 @@ func TestJobUpdateJobFinal(t *testing.T) {
279279
jobDbHandler, err := NewJobDBHandler(database, true)
280280
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
281281

282-
job, err := model.NewJob("TestTask", nil)
282+
job, err := model.NewJob("TestTask", nil, nil)
283283
require.NoError(t, err, "Expected NewJob to not return an error")
284284

285285
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -305,7 +305,7 @@ func TestJobUpdateJobFinalEncrypted(t *testing.T) {
305305
jobDbHandler, err := NewJobDBHandler(database, true, "test-encryption-key")
306306
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
307307

308-
job, err := model.NewJob("TestTask", nil)
308+
job, err := model.NewJob("TestTask", nil, nil)
309309
require.NoError(t, err, "Expected NewJob to not return an error")
310310

311311
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -375,7 +375,7 @@ func TestUpdateStaleJobs(t *testing.T) {
375375

376376
jobs := make([]*model.Job, len(testCases))
377377
for i, tc := range testCases {
378-
job, err := model.NewJob("test-task", nil)
378+
job, err := model.NewJob("test-task", nil, nil)
379379
require.NoError(t, err)
380380

381381
if tc.workerRID == "stopped" {
@@ -426,7 +426,7 @@ func TestUpdateStaleJobs(t *testing.T) {
426426
insertedWorker, err := workerDbHandler.InsertWorker(worker)
427427
require.NoError(t, err)
428428

429-
job, err := model.NewJob("test-task", nil)
429+
job, err := model.NewJob("test-task", nil, nil)
430430
require.NoError(t, err)
431431
job.WorkerRID = insertedWorker.RID
432432
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -453,7 +453,7 @@ func TestJobDeleteJob(t *testing.T) {
453453
jobDbHandler, err := NewJobDBHandler(database, true)
454454
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
455455

456-
job, err := model.NewJob("TestTask", nil)
456+
job, err := model.NewJob("TestTask", nil, nil)
457457
require.NoError(t, err, "Expected NewJob to not return an error")
458458

459459
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -487,7 +487,7 @@ func TestJobSelectJob(t *testing.T) {
487487
jobDbHandler, err := NewJobDBHandler(database, true)
488488
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
489489

490-
job, err := model.NewJob("TestTask", nil)
490+
job, err := model.NewJob("TestTask", nil, nil)
491491
require.NoError(t, err, "Expected NewJob to not return an error")
492492

493493
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -510,7 +510,7 @@ func TestJobSelectJobEncrypted(t *testing.T) {
510510
jobDbHandler, err := NewJobDBHandler(database, true, "test-encryption-key")
511511
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
512512

513-
job, err := model.NewJob("TestTask", nil)
513+
job, err := model.NewJob("TestTask", nil, nil)
514514
require.NoError(t, err, "Expected NewJob to not return an error")
515515

516516
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -547,7 +547,7 @@ func TestJobSelectAllJobs(t *testing.T) {
547547
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
548548

549549
for i := 0; i < newJobCount; i++ {
550-
job, err := model.NewJob(fmt.Sprintf("TestJob%v", i), nil)
550+
job, err := model.NewJob(fmt.Sprintf("TestJob%v", i), nil, nil)
551551
require.NoError(t, err, "Expected NewJob to not return an error")
552552

553553
_, err = jobDbHandler.InsertJob(job)
@@ -579,7 +579,7 @@ func TestJobSelectAllJobsEncrypted(t *testing.T) {
579579
expectedResultsByTaskName := make(map[string]model.Parameters)
580580
for i := 0; i < newJobCount; i++ {
581581
taskName := fmt.Sprintf("TestJob%v", i)
582-
job, err := model.NewJob(taskName, nil)
582+
job, err := model.NewJob(taskName, nil, nil)
583583
require.NoError(t, err, "Expected NewJob to not return an error")
584584

585585
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -645,7 +645,7 @@ func TestJobSelectAllJobsByWorkerRID(t *testing.T) {
645645
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
646646

647647
for i := 0; i < newJobCount; i++ {
648-
job, err := model.NewJob("TestTask", nil)
648+
job, err := model.NewJob("TestTask", nil, nil)
649649
require.NoError(t, err, "Expected NewJob to not return an error")
650650

651651
_, err = jobDbHandler.InsertJob(job)
@@ -678,15 +678,15 @@ func TestJobSelectAllJobsBySearch(t *testing.T) {
678678

679679
// Insert multiple jobs with different names
680680
for i := 0; i < newJobCountSearch; i++ {
681-
job, err := model.NewJob(searchTerm, nil)
681+
job, err := model.NewJob(searchTerm, nil, nil)
682682
require.NoError(t, err, "Expected NewJob to not return an error")
683683

684684
_, err = jobDbHandler.InsertJob(job)
685685
require.NoError(t, err, "Expected InsertJob to not return an error")
686686
}
687687

688688
for i := 0; i < newJobCountOther; i++ {
689-
job, err := model.NewJob("TestTask", nil)
689+
job, err := model.NewJob("TestTask", nil, nil)
690690
require.NoError(t, err, "Expected NewJob to not return an error")
691691

692692
_, err = jobDbHandler.InsertJob(job)
@@ -714,7 +714,7 @@ func TestJobSelectJobFromArchive(t *testing.T) {
714714
jobDbHandler, err := NewJobDBHandler(database, true)
715715
require.NoError(t, err, "Expected NewJobDBHandler to not return an error")
716716

717-
job, err := model.NewJob("TestTask", nil)
717+
job, err := model.NewJob("TestTask", nil, nil)
718718
require.NoError(t, err, "Expected NewJob to not return an error")
719719

720720
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -746,7 +746,7 @@ func TestJobSelectAllJobsFromArchive(t *testing.T) {
746746

747747
newJobCount := 5
748748
for i := 0; i < newJobCount; i++ {
749-
job, err := model.NewJob(fmt.Sprintf("TestJob%v", i), nil)
749+
job, err := model.NewJob(fmt.Sprintf("TestJob%v", i), nil, nil)
750750
require.NoError(t, err, "Expected NewJob to not return an error")
751751

752752
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -785,7 +785,7 @@ func TestJobSelectAllJobsFromArchiveBySearch(t *testing.T) {
785785

786786
// Insert multiple jobs with different names
787787
for i := 0; i < newJobCountSearch; i++ {
788-
job, err := model.NewJob(searchTerm, nil)
788+
job, err := model.NewJob(searchTerm, nil, nil)
789789
require.NoError(t, err, "Expected NewJob to not return an error")
790790

791791
insertedJob, err := jobDbHandler.InsertJob(job)
@@ -798,7 +798,7 @@ func TestJobSelectAllJobsFromArchiveBySearch(t *testing.T) {
798798
}
799799

800800
for i := 0; i < newJobCountOther; i++ {
801-
job, err := model.NewJob("TestTask", nil)
801+
job, err := model.NewJob("TestTask", nil, nil)
802802
require.NoError(t, err, "Expected NewJob to not return an error")
803803

804804
insertedJob, err := jobDbHandler.InsertJob(job)

example/exampleEasy.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@ func ExampleEasy() {
2121
q.Start(ctx, cancel)
2222

2323
// Add a job to the queue
24-
job, err := q.AddJob(ShortTask, 5, "12")
24+
job, err := q.AddJob(ShortTask, nil, 5, "12")
2525
if err != nil {
2626
log.Fatalf("Error adding job: %v", err)
2727
}

example/exampleFull.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,7 @@ func ExampleFull() {
3737
})
3838

3939
// Example adding a single job to the queue
40-
_, err := q.AddJob(ShortTask, 5, "12")
40+
_, err := q.AddJob(ShortTask, nil, 5, "12")
4141
if err != nil {
4242
log.Fatalf("Error adding job: %v", err)
4343
}
@@ -69,7 +69,7 @@ func ExampleFull() {
6969
Interval: time.Second * 5,
7070
},
7171
}
72-
_, err = q.AddJobWithOptions(options, ShortTask, 5, "12")
72+
_, err = q.AddJobWithOptions(options, ShortTask, nil, 5, "12")
7373
if err != nil {
7474
log.Fatalf("Error adding job with options: %v", err)
7575
}
@@ -82,7 +82,7 @@ func ExampleFull() {
8282
MaxCount: 10,
8383
},
8484
}
85-
job, err := q.AddJobWithOptions(options, LongTask, time.Now().Second(), "1")
85+
job, err := q.AddJobWithOptions(options, LongTask, nil, time.Now().Second(), "1")
8686
if err != nil {
8787
log.Fatalf("Error adding job with schedule options: %v", err)
8888
}

0 commit comments

Comments
 (0)