syncthing/discover/cmd/discosrv/main.go

289 lines
5.5 KiB
Go
Raw Normal View History

2014-06-01 20:50:14 +00:00
// Copyright (C) 2014 Jakob Borg and other contributors. All rights reserved.
// Use of this source code is governed by an MIT-style license that can be
// found in the LICENSE file.
2013-12-23 02:35:05 +00:00
package main
import (
2014-02-20 16:40:15 +00:00
"encoding/binary"
"encoding/hex"
"flag"
2014-04-19 21:14:56 +00:00
"fmt"
"io"
2013-12-23 02:35:05 +00:00
"log"
"net"
2014-02-20 16:40:15 +00:00
"os"
2013-12-23 02:35:05 +00:00
"sync"
"time"
2013-12-23 02:35:05 +00:00
"github.com/calmh/syncthing/discover"
"github.com/calmh/syncthing/protocol"
2014-04-03 21:38:32 +00:00
"github.com/golang/groupcache/lru"
"github.com/juju/ratelimit"
2013-12-23 02:35:05 +00:00
)
type node struct {
addresses []address
updated time.Time
2014-02-20 16:40:15 +00:00
}
type address struct {
ip []byte
port uint16
2013-12-23 02:35:05 +00:00
}
var (
nodes = make(map[protocol.NodeID]node)
2014-06-27 20:39:03 +00:00
lock sync.Mutex
queries = 0
announces = 0
answered = 0
limited = 0
unknowns = 0
debug = false
lruSize = 1024
limitAvg = 1
limitBurst = 10
limiter *lru.Cache
2013-12-23 02:35:05 +00:00
)
func main() {
2014-02-20 16:40:15 +00:00
var listen string
var timestamp bool
2014-04-19 21:14:56 +00:00
var statsIntv int
var statsFile string
2014-02-20 16:40:15 +00:00
flag.StringVar(&listen, "listen", ":22025", "Listen address")
flag.BoolVar(&debug, "debug", false, "Enable debug output")
flag.BoolVar(&timestamp, "timestamp", true, "Timestamp the log output")
2014-04-19 21:14:56 +00:00
flag.IntVar(&statsIntv, "stats-intv", 0, "Statistics output interval (s)")
flag.StringVar(&statsFile, "stats-file", "/var/log/discosrv.stats", "Statistics file name")
2014-06-27 20:39:03 +00:00
flag.IntVar(&lruSize, "limit-cache", lruSize, "Limiter cache entries")
flag.IntVar(&limitAvg, "limit-avg", limitAvg, "Allowed average package rate, per 10 s")
flag.IntVar(&limitBurst, "limit-burst", limitBurst, "Allowed burst size, packets")
2014-02-20 16:40:15 +00:00
flag.Parse()
2014-06-27 20:39:03 +00:00
limiter = lru.New(lruSize)
2014-02-20 16:40:15 +00:00
log.SetOutput(os.Stdout)
if !timestamp {
log.SetFlags(0)
}
addr, _ := net.ResolveUDPAddr("udp", listen)
2013-12-23 02:35:05 +00:00
conn, err := net.ListenUDP("udp", addr)
if err != nil {
2014-04-03 20:44:40 +00:00
log.Fatal(err)
2013-12-23 02:35:05 +00:00
}
2014-04-19 21:14:56 +00:00
if statsIntv > 0 {
go logStats(statsFile, statsIntv)
}
2013-12-23 02:35:05 +00:00
var buf = make([]byte, 1024)
for {
2014-02-20 16:40:15 +00:00
buf = buf[:cap(buf)]
2013-12-23 02:35:05 +00:00
n, addr, err := conn.ReadFromUDP(buf)
2014-04-03 21:38:32 +00:00
if limit(addr) {
// Rate limit in effect for source
continue
}
2013-12-23 02:35:05 +00:00
if err != nil {
2014-04-03 20:44:40 +00:00
log.Fatal(err)
2013-12-23 02:35:05 +00:00
}
2014-04-03 21:38:32 +00:00
2014-02-20 16:40:15 +00:00
if n < 4 {
log.Printf("Received short packet (%d bytes)", n)
2013-12-23 02:35:05 +00:00
continue
}
2014-02-20 16:40:15 +00:00
buf = buf[:n]
magic := binary.BigEndian.Uint32(buf)
switch magic {
2014-04-03 20:44:40 +00:00
case discover.AnnouncementMagicV2:
handleAnnounceV2(addr, buf)
2014-02-20 16:40:15 +00:00
2014-04-03 20:44:40 +00:00
case discover.QueryMagicV2:
handleQueryV2(conn, addr, buf)
2014-04-19 21:14:56 +00:00
default:
lock.Lock()
unknowns++
lock.Unlock()
2014-04-03 20:44:40 +00:00
}
}
}
2014-02-20 16:40:15 +00:00
2014-04-03 21:38:32 +00:00
func limit(addr *net.UDPAddr) bool {
key := addr.IP.String()
lock.Lock()
defer lock.Unlock()
bkt, ok := limiter.Get(key)
if ok {
bkt := bkt.(*ratelimit.Bucket)
if bkt.TakeAvailable(1) != 1 {
// Rate limit exceeded; ignore packet
if debug {
2014-04-16 13:06:54 +00:00
log.Println("Rate limit exceeded for", key)
2014-04-03 21:38:32 +00:00
}
limited++
return true
}
} else {
if debug {
2014-04-16 13:06:54 +00:00
log.Println("New limiter for", key)
2014-04-03 21:38:32 +00:00
}
// One packet per ten seconds average rate, burst ten packets
2014-06-27 20:39:03 +00:00
limiter.Add(key, ratelimit.NewBucket(10*time.Second/time.Duration(limitAvg), int64(limitBurst)))
2014-04-03 21:38:32 +00:00
}
return false
}
2014-04-03 20:44:40 +00:00
func handleAnnounceV2(addr *net.UDPAddr, buf []byte) {
var pkt discover.AnnounceV2
err := pkt.UnmarshalXDR(buf)
if err != nil && err != io.EOF {
2014-04-03 20:44:40 +00:00
log.Println("AnnounceV2 Unmarshal:", err)
log.Println(hex.Dump(buf))
return
}
if debug {
log.Printf("<- %v %#v", addr, pkt)
}
2014-04-19 21:14:56 +00:00
lock.Lock()
announces++
lock.Unlock()
2014-04-03 20:44:40 +00:00
ip := addr.IP.To4()
if ip == nil {
ip = addr.IP.To16()
}
var addrs []address
for _, addr := range pkt.This.Addresses {
2014-04-03 20:44:40 +00:00
tip := addr.IP
if len(tip) == 0 {
tip = ip
}
addrs = append(addrs, address{
ip: tip,
port: addr.Port,
2014-04-03 20:44:40 +00:00
})
}
node := node{
addresses: addrs,
updated: time.Now(),
2014-04-03 20:44:40 +00:00
}
var id protocol.NodeID
if len(pkt.This.ID) == 32 {
// Raw node ID
copy(id[:], pkt.This.ID)
} else {
id.UnmarshalText(pkt.This.ID)
}
2014-04-03 20:44:40 +00:00
lock.Lock()
nodes[id] = node
2014-04-03 20:44:40 +00:00
lock.Unlock()
}
func handleQueryV2(conn *net.UDPConn, addr *net.UDPAddr, buf []byte) {
var pkt discover.QueryV2
err := pkt.UnmarshalXDR(buf)
if err != nil {
log.Println("QueryV2 Unmarshal:", err)
log.Println(hex.Dump(buf))
return
}
if debug {
log.Printf("<- %v %#v", addr, pkt)
}
var id protocol.NodeID
if len(pkt.NodeID) == 32 {
// Raw node ID
copy(id[:], pkt.NodeID)
} else {
id.UnmarshalText(pkt.NodeID)
}
2014-04-03 20:44:40 +00:00
lock.Lock()
node, ok := nodes[id]
2014-04-03 20:44:40 +00:00
queries++
lock.Unlock()
if ok && len(node.addresses) > 0 {
ann := discover.AnnounceV2{
Magic: discover.AnnouncementMagicV2,
This: discover.Node{
ID: pkt.NodeID,
},
2014-04-03 20:44:40 +00:00
}
for _, addr := range node.addresses {
ann.This.Addresses = append(ann.This.Addresses, discover.Address{IP: addr.ip, Port: addr.port})
2014-04-03 20:44:40 +00:00
}
if debug {
log.Printf("-> %v %#v", addr, pkt)
}
tb := ann.MarshalXDR()
2014-04-03 20:44:40 +00:00
_, _, err = conn.WriteMsgUDP(tb, nil, addr)
if err != nil {
log.Println("QueryV2 response write:", err)
}
lock.Lock()
answered++
lock.Unlock()
}
}
2014-04-19 21:14:56 +00:00
func next(intv int) time.Time {
d := time.Duration(intv) * time.Second
t0 := time.Now()
t1 := t0.Add(d).Truncate(d)
time.Sleep(t1.Sub(t0))
return t1
}
func logStats(file string, intv int) {
f, err := os.OpenFile(file, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644)
if err != nil {
log.Fatal(err)
}
2014-04-03 20:44:40 +00:00
for {
2014-04-19 21:14:56 +00:00
t := next(intv)
2014-04-03 20:44:40 +00:00
lock.Lock()
var deleted = 0
for id, node := range nodes {
if time.Since(node.updated) > 60*time.Minute {
2014-04-03 20:44:40 +00:00
delete(nodes, id)
deleted++
2013-12-23 02:35:05 +00:00
}
}
2014-04-19 21:14:56 +00:00
fmt.Fprintf(f, "%d Nr:%d Ne:%d Qt:%d Qa:%d A:%d U:%d Lq:%d Lc:%d\n",
t.Unix(), len(nodes), deleted, queries, answered, announces, unknowns, limited, limiter.Len())
f.Sync()
2014-04-03 20:44:40 +00:00
queries = 0
2014-04-19 21:14:56 +00:00
announces = 0
2014-04-03 20:44:40 +00:00
answered = 0
2014-04-03 21:38:32 +00:00
limited = 0
2014-04-19 21:14:56 +00:00
unknowns = 0
2014-04-03 20:44:40 +00:00
lock.Unlock()
2013-12-23 02:35:05 +00:00
}
}