2022-06-21 13:29:57 +02:00
|
|
|
package rpcmiddleware
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"context"
|
|
|
|
|
"sync"
|
|
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
"github.com/lightninglabs/lndclient"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
// Manager is the main middleware manager service.
|
|
|
|
|
type Manager struct {
|
|
|
|
|
interceptTimeout time.Duration
|
|
|
|
|
lndClient lndclient.LightningClient
|
|
|
|
|
interceptors []RequestInterceptor
|
|
|
|
|
|
2022-10-15 09:04:01 +02:00
|
|
|
mainErrChan chan<- error
|
|
|
|
|
wg sync.WaitGroup
|
|
|
|
|
cancel context.CancelFunc
|
|
|
|
|
quit chan struct{}
|
|
|
|
|
stopOnce sync.Once
|
2022-06-21 13:29:57 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// NewManager returns a new middleware manager.
|
|
|
|
|
func NewManager(interceptTimeout time.Duration,
|
2022-10-15 09:04:01 +02:00
|
|
|
lndClient lndclient.LightningClient, errChan chan<- error,
|
2022-06-21 13:29:57 +02:00
|
|
|
interceptors ...RequestInterceptor) *Manager {
|
|
|
|
|
|
|
|
|
|
return &Manager{
|
|
|
|
|
interceptTimeout: interceptTimeout,
|
|
|
|
|
lndClient: lndClient,
|
|
|
|
|
interceptors: interceptors,
|
2022-10-15 09:04:01 +02:00
|
|
|
mainErrChan: errChan,
|
2022-06-21 13:29:57 +02:00
|
|
|
quit: make(chan struct{}),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Start starts the firewall by registering the interceptors with lnd.
|
2025-01-13 07:04:36 +02:00
|
|
|
func (f *Manager) Start(ctx context.Context) error {
|
|
|
|
|
ctxc, cancel := context.WithCancel(ctx)
|
2022-06-21 13:29:57 +02:00
|
|
|
f.cancel = cancel
|
|
|
|
|
|
|
|
|
|
for _, i := range f.interceptors {
|
|
|
|
|
errChan, err := f.lndClient.RegisterRPCMiddleware(
|
|
|
|
|
ctxc, i.Name(), i.CustomCaveatName(), i.ReadOnly(),
|
|
|
|
|
f.interceptTimeout, i.Intercept,
|
|
|
|
|
)
|
|
|
|
|
if err != nil {
|
|
|
|
|
cancel()
|
|
|
|
|
f.wg.Wait()
|
|
|
|
|
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
f.wg.Add(1)
|
|
|
|
|
go func(i RequestInterceptor, errChan chan error) {
|
|
|
|
|
defer f.wg.Done()
|
|
|
|
|
|
|
|
|
|
for {
|
|
|
|
|
select {
|
|
|
|
|
case <-f.quit:
|
|
|
|
|
log.Debugf("Quitting interceptor %v, "+
|
|
|
|
|
"shutting down", i.Name())
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
case <-ctxc.Done():
|
|
|
|
|
log.Debugf("Quitting interceptor %v, "+
|
|
|
|
|
"context canceled", i.Name())
|
|
|
|
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
case err := <-errChan:
|
|
|
|
|
log.Errorf("Error in interceptor: %v",
|
|
|
|
|
err)
|
|
|
|
|
|
|
|
|
|
select {
|
2022-10-15 09:04:01 +02:00
|
|
|
case f.mainErrChan <- err:
|
2022-06-21 13:29:57 +02:00
|
|
|
case <-f.quit:
|
|
|
|
|
case <-ctxc.Done():
|
|
|
|
|
}
|
2022-10-15 09:03:59 +02:00
|
|
|
|
|
|
|
|
return
|
2022-06-21 13:29:57 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}(i, errChan)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Stop shuts down the middleware manager.
|
|
|
|
|
func (f *Manager) Stop() {
|
|
|
|
|
f.stopOnce.Do(func() {
|
|
|
|
|
close(f.quit)
|
|
|
|
|
f.cancel()
|
|
|
|
|
|
|
|
|
|
f.wg.Wait()
|
|
|
|
|
})
|
|
|
|
|
}
|