fix: improve local Pulsar connectivity checks
This commit is contained in:
+20
-7
@@ -16,11 +16,14 @@ limitations under the License.
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"net"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
"go.uber.org/zap"
|
||||
"it2000.com.cn/tele-recv/serial"
|
||||
"it2000.com.cn/tele-recv/storage"
|
||||
"it2000.com.cn/tele-recv/transport"
|
||||
)
|
||||
|
||||
// testCmd represents the test command
|
||||
@@ -56,16 +59,26 @@ func test() {
|
||||
store.Close()
|
||||
}
|
||||
|
||||
// Test Pulsar
|
||||
sender, err := transport.NewPulsarSender(pulsarUrl, topic, name)
|
||||
// Test the Pulsar endpoint without constructing a producer. The legacy Pulsar
|
||||
// client used by the service is not compatible with newer Go runtimes on macOS.
|
||||
endpoint, err := url.Parse(pulsarUrl)
|
||||
if err != nil {
|
||||
logger.Warn("pulsar connection failed (expected if no broker)", zap.Error(err))
|
||||
logger.Warn("pulsar URL is invalid", zap.Error(err))
|
||||
} else {
|
||||
logger.Info("pulsar ok")
|
||||
sender.Close()
|
||||
address := endpoint.Host
|
||||
if endpoint.Port() == "" {
|
||||
address = net.JoinHostPort(endpoint.Hostname(), "6650")
|
||||
}
|
||||
conn, err := net.DialTimeout("tcp", address, 3*time.Second)
|
||||
if err != nil {
|
||||
logger.Warn("pulsar connection failed (expected if no broker)", zap.Error(err))
|
||||
} else {
|
||||
conn.Close()
|
||||
logger.Info("pulsar endpoint ok")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func init() {
|
||||
rootCmd.AddCommand(testCmd)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user