fix(uno): close tcp connexions to the LibreOffice listener, better health check that take the restart into account

This commit is contained in:
Julien Neuhart
2022-02-08 16:12:41 +01:00
parent 5cad5ca11a
commit b7a8b95820
5 changed files with 229 additions and 66 deletions

View File

@@ -15,6 +15,7 @@ import (
type listener interface {
start(logger *zap.Logger) error
stop(logger *zap.Logger) error
restart(logger *zap.Logger) error
lock(ctx context.Context, logger *zap.Logger) error
unlock(logger *zap.Logger) error
port() int
@@ -33,6 +34,8 @@ type libreOfficeListener struct {
cfgMu sync.RWMutex
usage int
restarting bool
restartingMu sync.RWMutex
queueLength int
queueLengthMu sync.RWMutex
lockChan chan struct{}
@@ -106,8 +109,12 @@ func (listener *libreOfficeListener) start(logger *zap.Logger) error {
return fmt.Errorf("waiting for the LibreOffice listener socket to be available: %w", ctx.Err())
}
_, err = net.DialTimeout("tcp", fmt.Sprintf("127.0.0.1:%d", port), time.Duration(1)*time.Second)
conn, err := net.DialTimeout("tcp", fmt.Sprintf("127.0.0.1:%d", port), time.Duration(1)*time.Second)
if err == nil {
err = conn.Close()
if err != nil {
logger.Debug(fmt.Sprintf("close connection after health checking the LibreOffice listener: %v", err))
}
break
}
}
@@ -149,6 +156,32 @@ func (listener *libreOfficeListener) stop(logger *zap.Logger) error {
return nil
}
func (listener *libreOfficeListener) restart(logger *zap.Logger) error {
listener.restartingMu.Lock()
listener.restarting = true
listener.restartingMu.Unlock()
defer func() {
listener.restartingMu.Lock()
listener.restarting = false
listener.restartingMu.Unlock()
}()
err := listener.stop(logger)
if err != nil {
return fmt.Errorf("stop LibreOffice listener: %w", err)
}
err = listener.start(logger)
if err != nil {
return fmt.Errorf("start LibreOffice listener: %w", err)
}
listener.usage = 0
return nil
}
func (listener *libreOfficeListener) lock(ctx context.Context, logger *zap.Logger) error {
listener.queueLengthMu.Lock()
listener.queueLength += 1
@@ -162,6 +195,17 @@ func (listener *libreOfficeListener) lock(ctx context.Context, logger *zap.Logge
listener.queueLength -= 1
listener.queueLengthMu.Unlock()
if !listener.healthy() {
logger.Debug("LibreOffice listener is unhealthy, restarting it...")
err := listener.restart(logger)
if err == nil {
return nil
}
return fmt.Errorf("restart LibreOffice listener: %w", err)
}
return nil
case <-ctx.Done():
logger.Debug("failed to acquire LibreOffice listener lock before deadline")
@@ -175,22 +219,6 @@ func (listener *libreOfficeListener) lock(ctx context.Context, logger *zap.Logge
}
func (listener *libreOfficeListener) unlock(logger *zap.Logger) error {
restart := func() error {
err := listener.stop(logger)
if err != nil {
return fmt.Errorf("stop LibreOffice listener: %w", err)
}
err = listener.start(logger)
if err != nil {
return fmt.Errorf("start LibreOffice listener: %w", err)
}
listener.usage = 0
return nil
}
defer func() {
<-listener.lockChan
logger.Debug("LibreOffice listener lock released")
@@ -199,7 +227,7 @@ func (listener *libreOfficeListener) unlock(logger *zap.Logger) error {
if !listener.healthy() {
logger.Debug("LibreOffice listener is unhealthy, restarting it...")
err := restart()
err := listener.restart(logger)
if err == nil {
return nil
}
@@ -214,7 +242,7 @@ func (listener *libreOfficeListener) unlock(logger *zap.Logger) error {
logger.Debug("LibreOffice listener threshold reached, restarting it...")
err := restart()
err := listener.restart(logger)
if err == nil {
return nil
}
@@ -237,9 +265,24 @@ func (listener *libreOfficeListener) queue() int {
}
func (listener *libreOfficeListener) healthy() bool {
_, err := net.DialTimeout("tcp", fmt.Sprintf("127.0.0.1:%d", listener.port()), time.Duration(1)*time.Second)
listener.restartingMu.RLock()
defer listener.restartingMu.RUnlock()
return err == nil
if listener.restarting {
return true
}
conn, err := net.DialTimeout("tcp", fmt.Sprintf("127.0.0.1:%d", listener.port()), time.Duration(1)*time.Second)
if err == nil {
err := conn.Close()
if err != nil {
listener.logger.Debug(fmt.Sprintf("close connection after health checking the LibreOffice listener: %v", err))
}
return true
}
return false
}
// Interface guards.

View File

@@ -2,7 +2,6 @@ package uno
import (
"context"
"errors"
"os"
"testing"
"time"
@@ -74,7 +73,7 @@ func TestListener_stop(t *testing.T) {
}
}
func TestListener_lock(t *testing.T) {
func TestListener_restart(t *testing.T) {
listener := newLibreOfficeListener(
zap.NewNop(),
os.Getenv("LIBREOFFICE_BIN_PATH"),
@@ -82,18 +81,112 @@ func TestListener_lock(t *testing.T) {
10,
)
ctx, cancel := context.WithTimeout(context.Background(), time.Duration(10)*time.Second)
err := listener.lock(ctx, zap.NewNop())
err := listener.start(zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.lock(), but got: %v", err)
t.Fatalf("expected no error from listener.start(), but got: %v", err)
}
cancel()
err = listener.restart(zap.NewNop())
if err != nil {
t.Errorf("expected no error from listener.stop(), but got: %v", err)
}
err = listener.lock(ctx, zap.NewNop())
if !errors.Is(err, context.Canceled) {
t.Errorf("expected %v error, but got: %v", context.Canceled, err)
if !listener.healthy() {
t.Error("expected an healthy LibreOffice listener")
}
}
func TestListener_lock(t *testing.T) {
tests := []struct {
name string
listener listener
ctx context.Context
teardown func(listener listener) error
expectLockErr bool
}{
{
name: "nominal behavior",
listener: func() listener {
listener := newLibreOfficeListener(zap.NewNop(), os.Getenv("LIBREOFFICE_BIN_PATH"), time.Duration(10)*time.Second, 10)
err := listener.start(zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.start(), but got: %v", err)
}
return listener
}(),
ctx: context.Background(),
teardown: func(listener listener) error {
return listener.stop(zap.NewNop())
},
},
{
name: "unhealthy listener",
listener: func() listener {
listener := newLibreOfficeListener(zap.NewNop(), os.Getenv("LIBREOFFICE_BIN_PATH"), time.Duration(10)*time.Second, 10)
err := listener.start(zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.start(), but got: %v", err)
}
err = listener.stop(zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.stop(), but got: %v", err)
}
return listener
}(),
ctx: context.Background(),
teardown: func(listener listener) error {
return listener.stop(zap.NewNop())
},
},
{
name: "context done",
listener: func() listener {
listener := newLibreOfficeListener(zap.NewNop(), os.Getenv("LIBREOFFICE_BIN_PATH"), time.Duration(10)*time.Second, 10)
err := listener.start(zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.start(), but got: %v", err)
}
return listener
}(),
ctx: func() context.Context {
ctx, cancel := context.WithCancel(context.Background())
cancel()
return ctx
}(),
expectLockErr: true,
teardown: func(listener listener) error {
return listener.stop(zap.NewNop())
},
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
defer func() {
err := tc.teardown(tc.listener)
if err != nil {
t.Errorf("expected no error from tc.teardown(), but got: %v", err)
}
}()
err := tc.listener.lock(tc.ctx, zap.NewNop())
if tc.expectLockErr && err == nil {
t.Fatalf("expected listener.lock() error, but got none")
}
if !tc.expectLockErr && err != nil {
t.Fatalf("expected no error from listener.lock(), but got: %v", err)
}
})
}
}
@@ -112,6 +205,12 @@ func TestListener_unlock(t *testing.T) {
if err != nil {
t.Fatalf("expected no error from listener.start(), but got: %v", err)
}
err = listener.lock(context.Background(), zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.lock(), but got: %v", err)
}
return listener
}(),
teardown: func(listener listener) error {
@@ -128,6 +227,11 @@ func TestListener_unlock(t *testing.T) {
t.Fatalf("expected no error from listener.start(), but got: %v", err)
}
err = listener.lock(context.Background(), zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.lock(), but got: %v", err)
}
err = listener.stop(zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.stop(), but got: %v", err)
@@ -148,6 +252,12 @@ func TestListener_unlock(t *testing.T) {
if err != nil {
t.Fatalf("expected no error from listener.start(), but got: %v", err)
}
err = listener.lock(context.Background(), zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.lock(), but got: %v", err)
}
return listener
}(),
teardown: func(listener listener) error {
@@ -158,23 +268,17 @@ func TestListener_unlock(t *testing.T) {
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), time.Duration(10)*time.Second)
defer cancel()
defer func() {
err := tc.teardown(tc.listener)
if err != nil {
t.Errorf("expected no error from tc.teardown(), but got: %v", err)
}
}()
err := tc.listener.lock(ctx, zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.lock(), but got: %v", err)
}
err = tc.listener.unlock(zap.NewNop())
err := tc.listener.unlock(zap.NewNop())
if err != nil {
t.Errorf("expected no error from listener.unlock(), but got: %v", err)
}
err = tc.teardown(tc.listener)
if err != nil {
t.Errorf("expected no error from tc.teardown(), but got: %v", err)
}
})
}
}
@@ -211,6 +315,18 @@ func TestListener_queue(t *testing.T) {
10,
)
err := listener.start(zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.start(), but got: %v", err)
}
defer func() {
err := listener.stop(zap.NewNop())
if err != nil {
t.Errorf("expected no error from listener.stop(), but got: %v", err)
}
}()
queueLength := listener.queue()
if queueLength != 0 {
t.Fatalf("expected a zero value from listener.queue(), but got %d", queueLength)
@@ -218,7 +334,7 @@ func TestListener_queue(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), time.Duration(10)*time.Second)
err := listener.lock(ctx, zap.NewNop())
err = listener.lock(ctx, zap.NewNop())
if err != nil {
t.Fatalf("expected no error from listener.lock(), but got: %v", err)
}
@@ -261,12 +377,13 @@ func TestListener_queue(t *testing.T) {
}
func TestListener_healthy(t *testing.T) {
listener := newLibreOfficeListener(
zap.NewNop(),
os.Getenv("LIBREOFFICE_BIN_PATH"),
time.Duration(10)*time.Second,
10,
)
listener := &libreOfficeListener{
binPath: os.Getenv("LIBREOFFICE_BIN_PATH"),
startTimeout: time.Duration(10) * time.Second,
threshold: 10,
lockChan: make(chan struct{}, 1),
logger: zap.NewNop(),
}
err := listener.start(zap.NewNop())
if err != nil {
@@ -285,4 +402,10 @@ func TestListener_healthy(t *testing.T) {
if listener.healthy() {
t.Errorf("expected a non-healthy LibreOffice listener")
}
listener.restarting = true
if !listener.healthy() {
t.Error("expected an healthy LibreOffice listener")
}
}

View File

@@ -179,10 +179,10 @@ func (mod UNO) Start() error {
// StartupMessage returns a custom startup message.
func (mod UNO) StartupMessage() string {
if mod.libreOfficeRestartThreshold == 0 {
return "Long-running LibreOffice listener disabled"
return "long-running LibreOffice listener disabled"
}
return "Long-running LibreOffice listener started"
return "long-running LibreOffice listener started"
}
// Stop stops the long-running LibreOffice Listener if it exists.
@@ -288,10 +288,6 @@ func (mod UNO) Checks() ([]health.CheckerOption, error) {
return errors.New("long-running LibreOffice listener unhealthy")
},
// The long-running LibreOffice listener may be restarting, so we
// wait a given amount of time until we consider the module
// unavailable.
MaxTimeInError: mod.libreOfficeStartTimeout,
}),
}, nil
}
@@ -339,6 +335,8 @@ func (mod UNO) PDF(ctx context.Context, logger *zap.Logger, inputPath, outputPat
err := mod.listener.unlock(logger)
if err != nil {
mod.logger.Error(fmt.Sprintf("unlock long-running LibreOffice listener: %v", err))
return
}
}()
}()

View File

@@ -258,14 +258,14 @@ func TestUNO_StartupMessage(t *testing.T) {
mod: UNO{
libreOfficeRestartThreshold: 10,
},
expectMessage: "Long-running LibreOffice listener started",
expectMessage: "long-running LibreOffice listener started",
},
{
name: "long-running LibreOffice listener disabled",
mod: UNO{
libreOfficeRestartThreshold: 0,
},
expectMessage: "Long-running LibreOffice listener disabled",
expectMessage: "long-running LibreOffice listener disabled",
},
}
@@ -755,7 +755,7 @@ func TestUNO_PDF(t *testing.T) {
},
logger: zap.NewNop(),
},
ctx: nil,
ctx: context.Background(),
logger: zap.NewNop(),
inputPath: "/tests/test/testdata/libreoffice/sample1.docx",
expectPDFErr: true,
@@ -824,6 +824,7 @@ func TestUNO_UNO(t *testing.T) {
type listenerMock struct {
startMock func(logger *zap.Logger) error
stopMock func(logger *zap.Logger) error
restartMock func(logger *zap.Logger) error
lockMock func(ctx context.Context, logger *zap.Logger) error
unlockMock func(logger *zap.Logger) error
portMock func() int
@@ -839,6 +840,10 @@ func (listener listenerMock) stop(logger *zap.Logger) error {
return listener.stopMock(logger)
}
func (listener listenerMock) restart(logger *zap.Logger) error {
return listener.restartMock(logger)
}
func (listener listenerMock) lock(ctx context.Context, logger *zap.Logger) error {
return listener.lockMock(ctx, logger)
}