package main import ( "context" "flag" "fmt" "os" "os/signal" "time" "github.com/nats-io/nats.go" "github.com/nats-io/nats.go/jetstream" ) func main() { var testJetStream bool var password, host string flag.StringVar(&password, "p", "", "password") flag.StringVar(&host, "h", "nats", "hostname") flag.BoolVar(&testJetStream, "js", false, "jetstream") flag.Parse() ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt) defer stop() nc, err := nats.Connect(fmt.Sprintf("nats://zerops:%s@%s:4222", password, host)) if err != nil { panic(err) } defer nc.Close() js, err := jetstream.New(nc) if err != nil { panic(err) } stream, err := js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{Name: "stream1"}) if err != nil { panic(err) } consumer, err := stream.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{Name: "consumer1"}) if err != nil { panic(err) } _, err = consumer.Consume(func(msg jetstream.Msg) { fmt.Printf("==> %s\n", string(msg.Data())) _ = msg.Ack() }) if err != nil { panic(err) } for { if ctx.Err() != nil { break } message := time.Now().Format(time.TimeOnly) fmt.Printf("%s ==>\n", message) _, err = js.Publish(ctx, "stream1", []byte(message)) if err != nil { fmt.Println(err) break } time.Sleep(time.Second) } <-ctx.Done() }