lightning-terminal/session/server.go

167 lines
3.7 KiB
Go
Raw Permalink Normal View History

2021-08-10 19:06:49 +02:00
package session
import (
"context"
2021-08-10 19:06:49 +02:00
"crypto/tls"
"fmt"
"sync"
"time"
2021-08-10 19:06:49 +02:00
"github.com/btcsuite/btcd/btcec/v2"
2021-08-10 19:06:49 +02:00
"github.com/lightninglabs/lightning-node-connect/mailbox"
"github.com/lightningnetwork/lnd/fn"
2021-08-10 19:06:49 +02:00
"github.com/lightningnetwork/lnd/keychain"
"google.golang.org/grpc"
"google.golang.org/grpc/keepalive"
2021-08-10 19:06:49 +02:00
)
type sessionID [33]byte
type GRPCServerCreator func(sessionID ID,
opts ...grpc.ServerOption) *grpc.Server
2021-08-10 19:06:49 +02:00
type mailboxSession struct {
server *grpc.Server
cancel fn.Option[context.CancelFunc]
wg sync.WaitGroup
quit chan struct{}
}
func newMailboxSession() *mailboxSession {
return &mailboxSession{
quit: make(chan struct{}),
}
2021-08-10 19:06:49 +02:00
}
func (m *mailboxSession) start(session *Session,
serverCreator GRPCServerCreator, authData []byte,
onUpdate func(ctx context.Context, id ID,
remote *btcec.PublicKey) error,
onNewStatus func(s mailbox.ServerStatus)) error {
2021-08-10 19:06:49 +02:00
tlsConfig := &tls.Config{}
if session.DevServer {
tlsConfig = &tls.Config{InsecureSkipVerify: true}
}
ecdh := &keychain.PrivKeyECDH{PrivKey: session.LocalPrivateKey}
ctx, cancel := context.WithCancel(context.Background())
m.cancel = fn.Some(cancel)
keys := mailbox.NewConnData(
ecdh, session.RemotePublicKey, session.PairingSecret[:],
authData, func(key *btcec.PublicKey) error {
return onUpdate(ctx, session.ID, key)
}, nil,
)
2021-08-10 19:06:49 +02:00
// Start the mailbox gRPC server.
mailboxServer, err := mailbox.NewServer(
session.ServerAddr, keys, onNewStatus,
grpc.WithTransportCredentials(
NewMailboxTLSCredentials(tlsConfig),
),
grpc.WithKeepaliveParams(keepalive.ClientParameters{
Time: 2 * time.Minute,
}),
2021-08-10 19:06:49 +02:00
)
if err != nil {
return err
}
noiseConn := mailbox.NewNoiseGrpcConn(keys)
m.server = serverCreator(session.ID, grpc.Creds(noiseConn))
2021-08-10 19:06:49 +02:00
m.wg.Add(1)
go m.run(mailboxServer)
return nil
}
func (m *mailboxSession) run(mailboxServer *mailbox.Server) {
defer m.wg.Done()
log.Infof("Mailbox RPC server listening on %s", mailboxServer.Addr())
if err := m.server.Serve(mailboxServer); err != nil {
log.Errorf("Unable to serve mailbox gRPC: %v", err)
}
}
func (m *mailboxSession) stop() {
m.cancel.WhenSome(func(fn context.CancelFunc) { fn() })
2021-08-10 19:06:49 +02:00
m.server.Stop()
close(m.quit)
2021-08-10 19:06:49 +02:00
m.wg.Wait()
}
type Server struct {
serverCreator GRPCServerCreator
activeSessions map[sessionID]*mailboxSession
activeSessionsMtx sync.Mutex
quit chan struct{}
}
func NewServer(serverCreator GRPCServerCreator) *Server {
return &Server{
serverCreator: serverCreator,
activeSessions: make(map[sessionID]*mailboxSession),
quit: make(chan struct{}),
}
}
func (s *Server) StartSession(session *Session, authData []byte,
onUpdate func(ctx context.Context, id ID,
remote *btcec.PublicKey) error,
onNewStatus func(s mailbox.ServerStatus)) (chan struct{}, error) {
2021-08-10 19:06:49 +02:00
s.activeSessionsMtx.Lock()
defer s.activeSessionsMtx.Unlock()
var id sessionID
copy(id[:], session.LocalPublicKey.SerializeCompressed())
_, ok := s.activeSessions[id]
if ok {
return nil, fmt.Errorf("session %x is already active", id[:])
2021-08-10 19:06:49 +02:00
}
sess := newMailboxSession()
s.activeSessions[id] = sess
return sess.quit, sess.start(
session, s.serverCreator, authData, onUpdate, onNewStatus,
)
2021-08-10 19:06:49 +02:00
}
func (s *Server) StopSession(localPublicKey *btcec.PublicKey) error {
s.activeSessionsMtx.Lock()
defer s.activeSessionsMtx.Unlock()
var id sessionID
copy(id[:], localPublicKey.SerializeCompressed())
_, ok := s.activeSessions[id]
if !ok {
return fmt.Errorf("session %x is not active", id[:])
}
s.activeSessions[id].stop()
delete(s.activeSessions, id)
return nil
}
func (s *Server) Stop() {
s.activeSessionsMtx.Lock()
defer s.activeSessionsMtx.Unlock()
for id, session := range s.activeSessions {
session.stop()
delete(s.activeSessions, id)
}
}