2021-08-10 19:06:49 +02:00
|
|
|
package session
|
|
|
|
|
|
|
|
|
|
import (
|
2025-02-25 16:13:16 +02:00
|
|
|
"context"
|
2021-08-10 19:06:49 +02:00
|
|
|
"crypto/tls"
|
|
|
|
|
"fmt"
|
|
|
|
|
"sync"
|
2022-02-08 12:14:54 +02:00
|
|
|
"time"
|
2021-08-10 19:06:49 +02:00
|
|
|
|
2022-03-28 22:03:17 +02:00
|
|
|
"github.com/btcsuite/btcd/btcec/v2"
|
2021-08-10 19:06:49 +02:00
|
|
|
"github.com/lightninglabs/lightning-node-connect/mailbox"
|
2025-02-25 16:13:16 +02:00
|
|
|
"github.com/lightningnetwork/lnd/fn"
|
2021-08-10 19:06:49 +02:00
|
|
|
"github.com/lightningnetwork/lnd/keychain"
|
|
|
|
|
"google.golang.org/grpc"
|
2022-02-08 12:14:54 +02:00
|
|
|
"google.golang.org/grpc/keepalive"
|
2021-08-10 19:06:49 +02:00
|
|
|
)
|
|
|
|
|
|
|
|
|
|
type sessionID [33]byte
|
|
|
|
|
|
2025-04-20 10:55:34 +02:00
|
|
|
type GRPCServerCreator func(sessionID ID,
|
|
|
|
|
opts ...grpc.ServerOption) *grpc.Server
|
2021-08-10 19:06:49 +02:00
|
|
|
|
|
|
|
|
type mailboxSession struct {
|
|
|
|
|
server *grpc.Server
|
|
|
|
|
|
2025-02-25 16:13:16 +02:00
|
|
|
cancel fn.Option[context.CancelFunc]
|
|
|
|
|
wg sync.WaitGroup
|
|
|
|
|
quit chan struct{}
|
2022-01-17 17:13:58 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func newMailboxSession() *mailboxSession {
|
|
|
|
|
return &mailboxSession{
|
|
|
|
|
quit: make(chan struct{}),
|
|
|
|
|
}
|
2021-08-10 19:06:49 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (m *mailboxSession) start(session *Session,
|
2022-04-29 14:15:35 +02:00
|
|
|
serverCreator GRPCServerCreator, authData []byte,
|
2025-03-04 17:39:03 +02:00
|
|
|
onUpdate func(ctx context.Context, id ID,
|
2025-02-25 16:13:16 +02:00
|
|
|
remote *btcec.PublicKey) error,
|
2022-08-18 18:10:44 +02:00
|
|
|
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}
|
|
|
|
|
}
|
|
|
|
|
|
2022-04-29 14:15:35 +02:00
|
|
|
ecdh := &keychain.PrivKeyECDH{PrivKey: session.LocalPrivateKey}
|
|
|
|
|
|
2025-02-25 16:13:16 +02:00
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
|
m.cancel = fn.Some(cancel)
|
|
|
|
|
|
2022-04-29 14:15:35 +02:00
|
|
|
keys := mailbox.NewConnData(
|
|
|
|
|
ecdh, session.RemotePublicKey, session.PairingSecret[:],
|
|
|
|
|
authData, func(key *btcec.PublicKey) error {
|
2025-03-04 17:39:03 +02:00
|
|
|
return onUpdate(ctx, session.ID, key)
|
2022-04-29 14:15:35 +02:00
|
|
|
}, nil,
|
|
|
|
|
)
|
|
|
|
|
|
2021-08-10 19:06:49 +02:00
|
|
|
// Start the mailbox gRPC server.
|
|
|
|
|
mailboxServer, err := mailbox.NewServer(
|
2022-08-18 18:10:44 +02:00
|
|
|
session.ServerAddr, keys, onNewStatus,
|
2026-03-17 16:52:26 -05:00
|
|
|
grpc.WithTransportCredentials(
|
|
|
|
|
NewMailboxTLSCredentials(tlsConfig),
|
|
|
|
|
),
|
2022-02-08 12:14:54 +02:00
|
|
|
grpc.WithKeepaliveParams(keepalive.ClientParameters{
|
|
|
|
|
Time: 2 * time.Minute,
|
|
|
|
|
}),
|
2021-08-10 19:06:49 +02:00
|
|
|
)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
2022-04-29 14:15:35 +02:00
|
|
|
noiseConn := mailbox.NewNoiseGrpcConn(keys)
|
2025-04-20 10:55:34 +02:00
|
|
|
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() {
|
2025-02-25 16:13:16 +02:00
|
|
|
m.cancel.WhenSome(func(fn context.CancelFunc) { fn() })
|
2021-08-10 19:06:49 +02:00
|
|
|
m.server.Stop()
|
2022-01-17 17:13:58 +02:00
|
|
|
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{}),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2022-04-29 14:15:35 +02:00
|
|
|
func (s *Server) StartSession(session *Session, authData []byte,
|
2025-03-04 17:39:03 +02:00
|
|
|
onUpdate func(ctx context.Context, id ID,
|
2025-02-25 16:13:16 +02:00
|
|
|
remote *btcec.PublicKey) error,
|
2022-08-18 18:10:44 +02:00
|
|
|
onNewStatus func(s mailbox.ServerStatus)) (chan struct{}, error) {
|
2022-01-17 17:13:58 +02:00
|
|
|
|
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 {
|
2022-01-17 17:13:58 +02:00
|
|
|
return nil, fmt.Errorf("session %x is already active", id[:])
|
2021-08-10 19:06:49 +02:00
|
|
|
}
|
|
|
|
|
|
2022-01-17 17:13:58 +02:00
|
|
|
sess := newMailboxSession()
|
|
|
|
|
s.activeSessions[id] = sess
|
|
|
|
|
|
2022-04-29 14:15:35 +02:00
|
|
|
return sess.quit, sess.start(
|
2022-08-18 18:10:44 +02:00
|
|
|
session, s.serverCreator, authData, onUpdate, onNewStatus,
|
2022-04-29 14:15:35 +02:00
|
|
|
)
|
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)
|
|
|
|
|
}
|
|
|
|
|
}
|