File tree Expand file tree Collapse file tree
Expand file tree Collapse file tree Original file line number Diff line number Diff line change 55
66package main
77
8- import "github.com/goforj/queue"
8+ import (
9+ "context"
10+
11+ "github.com/goforj/queue"
12+ )
913
1014func main () {
1115 // Observe forwards an event to the configured channel.
1216
1317 // Example: channel observer
1418 ch := make (chan queue.Event , 1 )
1519 observer := queue.ChannelObserver {Events : ch }
16- observer .Observe (queue.Event {Kind : queue .EventProcessStarted , Queue : "default" })
20+ observer .Observe (context . Background (), queue.Event {Kind : queue .EventProcessStarted , Queue : "default" })
1721 event := <- ch
1822 _ = event
1923}
Original file line number Diff line number Diff line change 66package main
77
88import (
9+ "context"
910 "fmt"
1011 "github.com/goforj/queue"
1112)
@@ -17,9 +18,9 @@ func main() {
1718 events := make (chan queue.Event , 2 )
1819 observer := queue .MultiObserver (
1920 queue.ChannelObserver {Events : events },
20- queue .ObserverFunc (func (queue.Event ) {}),
21+ queue .ObserverFunc (func (context. Context , queue.Event ) {}),
2122 )
22- observer .Observe (queue.Event {Kind : queue .EventEnqueueAccepted })
23+ observer .Observe (context . Background (), queue.Event {Kind : queue .EventEnqueueAccepted })
2324 fmt .Println (len (events ))
2425 // Output: 1
2526}
Original file line number Diff line number Diff line change @@ -22,7 +22,7 @@ func main() {
2222 var flakyAttempts atomic.Int32
2323 ctx := context .Background ()
2424
25- runtimeObserver := queue .ObserverFunc (func (event queue.Event ) {
25+ runtimeObserver := queue .ObserverFunc (func (ctx context. Context , event queue.Event ) {
2626 logger .Info ("runtime event" ,
2727 "kind" , event .Kind ,
2828 "driver" , event .Driver ,
@@ -35,7 +35,7 @@ func main() {
3535 )
3636 })
3737
38- workflowObserver := queue .WorkflowObserverFunc (func (event queue.WorkflowEvent ) {
38+ workflowObserver := queue .WorkflowObserverFunc (func (ctx context. Context , event queue.WorkflowEvent ) {
3939 logger .Info ("workflow event" ,
4040 "kind" , event .Kind ,
4141 "dispatch_id" , event .DispatchID ,
Original file line number Diff line number Diff line change 55
66package main
77
8- import "github.com/goforj/queue"
8+ import (
9+ "context"
10+
11+ "github.com/goforj/queue"
12+ )
913
1014func main () {
1115 // Observe handles a queue runtime event.
1216
1317 // Example: observe runtime event
1418 var observer queue.Observer
15- observer .Observe (queue.Event {
19+ observer .Observe (context . Background (), queue.Event {
1620 Kind : queue .EventEnqueueAccepted ,
1721 Driver : queue .DriverSync ,
1822 Queue : "default" ,
Original file line number Diff line number Diff line change 66package main
77
88import (
9+ "context"
910 "github.com/goforj/queue"
1011 "log/slog"
1112 "os"
@@ -16,7 +17,7 @@ func main() {
1617
1718 // Example: observer func logging hook
1819 logger := slog .New (slog .NewTextHandler (os .Stdout , nil ))
19- observer := queue .ObserverFunc (func (event queue.Event ) {
20+ observer := queue .ObserverFunc (func (ctx context. Context , event queue.Event ) {
2021 logger .Info ("queue event" ,
2122 "kind" , event .Kind ,
2223 "driver" , event .Driver ,
@@ -28,10 +29,10 @@ func main() {
2829 "err" , event .Err ,
2930 )
3031 })
31- observer .Observe (queue.Event {
32- Kind : queue .EventProcessSucceeded ,
33- Driver : queue .DriverSync ,
34- Queue : "default" ,
32+ observer .Observe (context . Background (), queue.Event {
33+ Kind : queue .EventProcessSucceeded ,
34+ Driver : queue .DriverSync ,
35+ Queue : "default" ,
3536 JobType : "emails:send" ,
3637 })
3738}
Original file line number Diff line number Diff line change 66package main
77
88import (
9+ "context"
910 "github.com/goforj/queue"
1011 "time"
1112)
@@ -15,7 +16,7 @@ func main() {
1516
1617 // Example: observe event
1718 collector := queue .NewStatsCollector ()
18- collector .Observe (queue.Event {
19+ collector .Observe (context . Background (), queue.Event {
1920 Kind : queue .EventEnqueueAccepted ,
2021 Driver : queue .DriverSync ,
2122 Queue : "default" ,
Original file line number Diff line number Diff line change 66package main
77
88import (
9+ "context"
910 "fmt"
1011 "github.com/goforj/queue"
1112 "time"
@@ -16,24 +17,24 @@ func main() {
1617
1718 // Example: snapshot print
1819 collector := queue .NewStatsCollector ()
19- collector .Observe (queue.Event {
20+ collector .Observe (context . Background (), queue.Event {
2021 Kind : queue .EventEnqueueAccepted ,
2122 Driver : queue .DriverSync ,
2223 Queue : "default" ,
2324 Time : time .Now (),
2425 })
25- collector .Observe (queue.Event {
26+ collector .Observe (context . Background (), queue.Event {
2627 Kind : queue .EventProcessStarted ,
2728 Driver : queue .DriverSync ,
2829 Queue : "default" ,
2930 JobKey : "job-1" ,
3031 Time : time .Now (),
3132 })
32- collector .Observe (queue.Event {
33+ collector .Observe (context . Background (), queue.Event {
3334 Kind : queue .EventProcessSucceeded ,
3435 Driver : queue .DriverSync ,
3536 Queue : "default" ,
36- JobKey : "job-1" ,
37+ JobKey : "job-1" ,
3738 Duration : 12 * time .Millisecond ,
3839 Time : time .Now (),
3940 })
Original file line number Diff line number Diff line change 66package main
77
88import (
9+ "context"
910 "fmt"
1011 "github.com/goforj/queue"
1112 "time"
@@ -16,7 +17,7 @@ func main() {
1617
1718 // Example: paused count getter
1819 collector := queue .NewStatsCollector ()
19- collector .Observe (queue.Event {
20+ collector .Observe (context . Background (), queue.Event {
2021 Kind : queue .EventQueuePaused ,
2122 Driver : queue .DriverSync ,
2223 Queue : "default" ,
Original file line number Diff line number Diff line change 66package main
77
88import (
9+ "context"
910 "fmt"
1011 "github.com/goforj/queue"
1112 "time"
@@ -16,7 +17,7 @@ func main() {
1617
1718 // Example: queue counters getter
1819 collector := queue .NewStatsCollector ()
19- collector .Observe (queue.Event {
20+ collector .Observe (context . Background (), queue.Event {
2021 Kind : queue .EventEnqueueAccepted ,
2122 Driver : queue .DriverSync ,
2223 Queue : "default" ,
Original file line number Diff line number Diff line change 66package main
77
88import (
9+ "context"
910 "fmt"
1011 "github.com/goforj/queue"
1112 "time"
@@ -16,7 +17,7 @@ func main() {
1617
1718 // Example: list queues
1819 collector := queue .NewStatsCollector ()
19- collector .Observe (queue.Event {
20+ collector .Observe (context . Background (), queue.Event {
2021 Kind : queue .EventEnqueueAccepted ,
2122 Driver : queue .DriverSync ,
2223 Queue : "critical" ,
You can’t perform that action at this time.
0 commit comments