Skip to content

Commit 279d5a8

Browse files
add producer wrapper with timestamp behavior test
1 parent 3599c3b commit 279d5a8

2 files changed

Lines changed: 67 additions & 0 deletions

File tree

pkg/producer/producer.go

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,31 @@
1+
package producer
2+
3+
import (
4+
"context"
5+
"time"
6+
7+
"distributed_job_queue/pkg/metrics"
8+
"distributed_job_queue/pkg/queue"
9+
)
10+
11+
type Publisher struct {
12+
broker queue.Broker
13+
metrics *metrics.Collector
14+
}
15+
16+
func New(broker queue.Broker, m *metrics.Collector) *Publisher {
17+
return &Publisher{broker: broker, metrics: m}
18+
}
19+
20+
func (p *Publisher) Enqueue(ctx context.Context, job queue.Job) error {
21+
if job.EnqueuedAt.IsZero() {
22+
job.EnqueuedAt = time.Now().UTC()
23+
}
24+
if err := p.broker.Enqueue(ctx, job); err != nil {
25+
return err
26+
}
27+
if p.metrics != nil {
28+
p.metrics.Enqueued.Inc()
29+
}
30+
return nil
31+
}

pkg/producer/producer_test.go

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,36 @@
1+
package producer
2+
3+
import (
4+
"context"
5+
"testing"
6+
7+
"distributed_job_queue/pkg/queue"
8+
)
9+
10+
type fakeBroker struct{ job queue.Job }
11+
12+
func (f *fakeBroker) Enqueue(_ context.Context, job queue.Job) error { f.job = job; return nil }
13+
func (f *fakeBroker) Reserve(context.Context, string) (queue.Job, queue.Lease, error) {
14+
return queue.Job{}, queue.Lease{}, nil
15+
}
16+
func (f *fakeBroker) Ack(context.Context, queue.Lease) error { return nil }
17+
func (f *fakeBroker) Nack(context.Context, queue.Lease, queue.RetryDecision) error { return nil }
18+
func (f *fakeBroker) RequeueExpired(context.Context) (int64, error) { return 0, nil }
19+
20+
func TestEnqueueSetsTimestamp(t *testing.T) {
21+
fb := &fakeBroker{}
22+
p := New(fb, nil)
23+
err := p.Enqueue(context.Background(), queue.Job{
24+
ID: "job-1",
25+
Type: "t",
26+
IdempotencyKey: "idem-1",
27+
Priority: queue.PriorityHigh,
28+
MaxAttempts: 3,
29+
})
30+
if err != nil {
31+
t.Fatalf("enqueue failed: %v", err)
32+
}
33+
if fb.job.EnqueuedAt.IsZero() {
34+
t.Fatal("expected enqueued timestamp")
35+
}
36+
}

0 commit comments

Comments
 (0)