From cdcc6fb3ab6530821ffe4e557f78df3209f9c39c Mon Sep 17 00:00:00 2001 From: Ta-Ching Chen Date: Tue, 19 Mar 2019 21:13:49 +0800 Subject: [PATCH] Add connection lost handler for NATS-streaming (#1125) --- mqtrigger/messageQueue/nats.go | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/mqtrigger/messageQueue/nats.go b/mqtrigger/messageQueue/nats.go index 8e4eb129..4b8a6ec7 100644 --- a/mqtrigger/messageQueue/nats.go +++ b/mqtrigger/messageQueue/nats.go @@ -47,7 +47,16 @@ type ( ) func makeNatsMessageQueue(logger *zap.Logger, routerUrl string, mqCfg MessageQueueConfig) (MessageQueue, error) { - conn, err := ns.Connect(natsClusterID, natsClientID, ns.NatsURL(mqCfg.Url)) + conn, err := ns.Connect(natsClusterID, natsClientID, ns.NatsURL(mqCfg.Url), + ns.SetConnectionLostHandler(func(conn ns.Conn, reason error) { + // TODO: Better way to handle connection lost problem. + // Currently, MessageQueue has no such interface to expose the status of underlying + // messaging service, hence MessageQueueTriggerManager has no way to detect and handle + // such situation properly. It takes some time to redesign interface of MessageQueue. + // For now, we simply fatal here. + logger.Fatal("Connection lost", zap.Error(reason)) + }), + ) if err != nil { return nil, err }