package main import ( "fmt" "log" "os" "os/signal" "syscall" "time" rabbitmq "github.com/wagslane/go-rabbitmq" ) func main() { conn, err := rabbitmq.NewConn( "amqp://guest:guest@localhost", rabbitmq.WithConnectionOptionsLogging, ) if err != nil { log.Fatal(err) } defer conn.Close() publisher, err := rabbitmq.NewPublisher( conn, rabbitmq.WithPublisherOptionsLogging, rabbitmq.WithPublisherOptionsExchangeName("events"), rabbitmq.WithPublisherOptionsExchangeDeclare, ) if err != nil { log.Fatal(err) } defer publisher.Close() publisher.NotifyReturn(func(r rabbitmq.Return) { log.Printf("message returned from server: %s", string(r.Body)) }) publisher.NotifyPublish(func(c rabbitmq.Confirmation) { log.Printf("message confirmed from server. tag: %v, ack: %v", c.DeliveryTag, c.Ack) }) publisher2, err := rabbitmq.NewPublisher( conn, rabbitmq.WithPublisherOptionsLogging, rabbitmq.WithPublisherOptionsExchangeName("events"), rabbitmq.WithPublisherOptionsExchangeDeclare, ) if err != nil { log.Fatal(err) } defer publisher2.Close() publisher2.NotifyReturn(func(r rabbitmq.Return) { log.Printf("message returned from server: %s", string(r.Body)) }) publisher2.NotifyPublish(func(c rabbitmq.Confirmation) { log.Printf("message confirmed from server. tag: %v, ack: %v", c.DeliveryTag, c.Ack) }) // block main thread - wait for shutdown signal sigs := make(chan os.Signal, 1) done := make(chan bool, 1) signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM) go func() { sig := <-sigs fmt.Println() fmt.Println(sig) done <- true }() fmt.Println("awaiting signal") ticker := time.NewTicker(time.Second) for { select { case <-ticker.C: err = publisher.Publish( []byte("hello, world"), []string{"my_routing_key"}, rabbitmq.WithPublishOptionsContentType("application/json"), rabbitmq.WithPublishOptionsMandatory, rabbitmq.WithPublishOptionsPersistentDelivery, rabbitmq.WithPublishOptionsExchange("events"), ) if err != nil { log.Println(err) } err = publisher2.Publish( []byte("hello, world 2"), []string{"my_routing_key_2"}, rabbitmq.WithPublishOptionsContentType("application/json"), rabbitmq.WithPublishOptionsMandatory, rabbitmq.WithPublishOptionsPersistentDelivery, rabbitmq.WithPublishOptionsExchange("events"), ) if err != nil { log.Println(err) } case <-done: fmt.Println("stopping publisher") return } } }