Files
clip-sync/internal/peer/peer.go
T

236 lines
5.4 KiB
Go

// Package peer sends clipboard content to remote clip-sync peers via HTTP(S) POST.
package peer
import (
"bytes"
"crypto/tls"
"crypto/x509"
"encoding/json"
"fmt"
"log"
"net"
"net/http"
"os"
"sync"
"time"
"git.dracodev.net/Projets/clip-sync/internal/config"
"git.dracodev.net/Projets/clip-sync/internal/dedup"
"git.dracodev.net/Projets/clip-sync/internal/transfer"
)
const (
clipPath = "/clip"
filePath = "/file"
)
// Options configures peer transport security.
type Options struct {
SharedKey string
TLS bool // global TLS default
CertFile string // CA (self-signed cert) to trust
InsecureSkipVerify bool
}
// Broadcaster sends clip payloads to all configured peers concurrently.
type Broadcaster struct {
mu sync.RWMutex
peers []config.PeerConfig
origin string
client *http.Client
opts Options
}
// NewBroadcaster creates a Broadcaster for the configured peer list.
func NewBroadcaster(peers []config.PeerConfig, origin string, opts Options) *Broadcaster {
client := &http.Client{
Timeout: 5 * time.Second,
Transport: &http.Transport{
DialContext: (&net.Dialer{
Timeout: 3 * time.Second,
}).DialContext,
MaxIdleConnsPerHost: 2,
TLSClientConfig: buildTLSConfig(opts),
},
}
return &Broadcaster{
peers: peers,
origin: origin,
client: client,
opts: opts,
}
}
func buildTLSConfig(opts Options) *tls.Config {
if !opts.TLS {
return nil
}
cfg := &tls.Config{
MinVersion: tls.VersionTLS12,
InsecureSkipVerify: opts.InsecureSkipVerify, // #nosec G402 — opt-in LAN trust
}
if !opts.InsecureSkipVerify && opts.CertFile != "" {
if pem, err := os.ReadFile(opts.CertFile); err == nil {
pool := x509.NewCertPool()
if pool.AppendCertsFromPEM(pem) {
cfg.RootCAs = pool
}
}
}
return cfg
}
// Broadcast sends the given text to all configured peers in parallel.
// Errors are logged but not returned — a single unreachable peer should not
// block other peers or the daemon loop.
func (b *Broadcaster) Broadcast(text string) {
body, ok := b.marshalPayload(text)
if !ok {
return
}
snap := b.snapshot()
for i := range snap {
go b.sendToPeer(snap[i], body)
}
}
// BroadcastSync sends the given text to all peers and blocks until every send
// has completed (or failed). Used by the one-shot mode so the process does not
// exit before the payload is on the wire.
func (b *Broadcaster) BroadcastSync(text string) {
body, ok := b.marshalPayload(text)
if !ok {
return
}
snap := b.snapshot()
var wg sync.WaitGroup
for i := range snap {
wg.Add(1)
go func(p config.PeerConfig) {
defer wg.Done()
b.sendToPeer(p, body)
}(snap[i])
}
wg.Wait()
}
func (b *Broadcaster) marshalPayload(text string) ([]byte, bool) {
payload := dedup.Payload{
Text: text,
Ts: time.Now().UnixNano(),
Origin: b.origin,
}
body, err := json.Marshal(payload)
if err != nil {
log.Printf("peer: marshal payload: %v", err)
return nil, false
}
return body, true
}
// snapshot returns a copy of the current peer list under read lock.
func (b *Broadcaster) snapshot() []config.PeerConfig {
b.mu.RLock()
defer b.mu.RUnlock()
return append([]config.PeerConfig(nil), b.peers...)
}
// PeerCount returns the number of currently configured peers.
func (b *Broadcaster) PeerCount() int {
b.mu.RLock()
defer b.mu.RUnlock()
return len(b.peers)
}
// AddPeer adds a peer if its address is not already present.
func (b *Broadcaster) AddPeer(p config.PeerConfig) {
b.mu.Lock()
defer b.mu.Unlock()
for _, existing := range b.peers {
if existing.Addr == p.Addr {
return
}
}
b.peers = append(b.peers, p)
}
// RemovePeer removes a peer by address.
func (b *Broadcaster) RemovePeer(addr string) {
b.mu.Lock()
defer b.mu.Unlock()
out := b.peers[:0]
for _, p := range b.peers {
if p.Addr != addr {
out = append(out, p)
}
}
b.peers = out
}
// SendFile sends a binary payload to all configured peers in parallel.
func (b *Broadcaster) SendFile(f transfer.File) {
body, err := json.Marshal(f)
if err != nil {
log.Printf("peer: marshal file: %v", err)
return
}
snap := b.snapshot()
for i := range snap {
go b.sendToPath(snap[i], filePath, body)
}
}
// SendFileSync sends a binary payload to all peers and blocks until done.
func (b *Broadcaster) SendFileSync(f transfer.File) {
body, err := json.Marshal(f)
if err != nil {
log.Printf("peer: marshal file: %v", err)
return
}
snap := b.snapshot()
var wg sync.WaitGroup
for i := range snap {
wg.Add(1)
go func(p config.PeerConfig) {
defer wg.Done()
b.sendToPath(p, filePath, body)
}(snap[i])
}
wg.Wait()
}
func (b *Broadcaster) sendToPeer(p config.PeerConfig, body []byte) {
b.sendToPath(p, clipPath, body)
}
func (b *Broadcaster) sendToPath(p config.PeerConfig, path string, body []byte) {
scheme := "http"
if p.TLS || b.opts.TLS {
scheme = "https"
}
url := fmt.Sprintf("%s://%s%s", scheme, p.Addr, path)
req, err := http.NewRequest(http.MethodPost, url, bytes.NewReader(body))
if err != nil {
log.Printf("peer [%s]: create request: %v", p.Name, err)
return
}
req.Header.Set("Content-Type", "application/json")
if b.opts.SharedKey != "" {
req.Header.Set("Authorization", "Bearer "+b.opts.SharedKey)
}
resp, err := b.client.Do(req)
if err != nil {
log.Printf("peer [%s] %s: %v", p.Name, p.Addr, err)
return
}
resp.Body.Close()
if resp.StatusCode != http.StatusNoContent {
log.Printf("peer [%s] %s: unexpected status %d", p.Name, p.Addr, resp.StatusCode)
}
}