mirror of
https://github.com/sot-tech/mochi.git
synced 2026-07-27 17:48:11 -07:00
initial
This commit is contained in:
@@ -0,0 +1,81 @@
|
||||
package server
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"errors"
|
||||
"log"
|
||||
|
||||
"github.com/jzelinskie/chihaya/config"
|
||||
"github.com/jzelinskie/chihaya/storage"
|
||||
)
|
||||
|
||||
func (h *handler) serveAnnounce(w *http.ResponseWriter, r *http.Request) {
|
||||
buf := h.bufferpool.Take()
|
||||
defer h.bufferpool.Give(buf)
|
||||
defer h.writeResponse(&w, r, buf)
|
||||
|
||||
user, err := validatePasskey(dir, h.storage)
|
||||
if err != nil {
|
||||
fail(err, buf)
|
||||
return
|
||||
}
|
||||
|
||||
pq, err := parseQuery(r.URL.RawQuery)
|
||||
if err != nil {
|
||||
fail(errors.New("Error parsing query"), buf)
|
||||
return
|
||||
}
|
||||
|
||||
ip, err := determineIP(r, pq)
|
||||
if err != nil {
|
||||
fail(err, buf)
|
||||
return
|
||||
}
|
||||
|
||||
err := validateParsedQuery(pq)
|
||||
if err != nil {
|
||||
fail(errors.New("Malformed request"), buf)
|
||||
return
|
||||
}
|
||||
|
||||
if !whitelisted(peerId, h.conf) {
|
||||
fail(errors.New("Your client is not approved"), buf)
|
||||
return
|
||||
}
|
||||
|
||||
torrent, exists, err := h.storage.FindTorrent(infohash)
|
||||
if err != nil {
|
||||
panic("server: failed to find torrent")
|
||||
}
|
||||
if !exists {
|
||||
fail(errors.New("This torrent does not exist"), buf)
|
||||
return
|
||||
}
|
||||
|
||||
if torrent.Status == 1 && left == 0 {
|
||||
err := h.storage.UnpruneTorrent(torrent)
|
||||
if err != nil {
|
||||
panic("server: failed to unprune torrent")
|
||||
}
|
||||
torrent.Status = 0
|
||||
} else if torrent.Status != 0 {
|
||||
fail(
|
||||
fmt.Errorf(
|
||||
"This torrent does not exist (status: %d, left: %d)",
|
||||
torrent.Status,
|
||||
left,
|
||||
),
|
||||
buf,
|
||||
)
|
||||
return
|
||||
}
|
||||
|
||||
//go
|
||||
}
|
||||
|
||||
func whitelisted(peerId string, conf config.Config) bool {
|
||||
// TODO Decide if whitelist should be in storage or config
|
||||
}
|
||||
|
||||
func newPeer() {
|
||||
}
|
||||
+119
@@ -0,0 +1,119 @@
|
||||
package server
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"strconv"
|
||||
)
|
||||
|
||||
type parsedQuery struct {
|
||||
infohashes []string
|
||||
params map[string]string
|
||||
}
|
||||
|
||||
func (pq *parsedQuery) getUint64(key string) (uint64, bool) {
|
||||
str, exists := pq[key]
|
||||
if !exists {
|
||||
return 0, false
|
||||
}
|
||||
val, err := strconv.Uint64(str, 10, 64)
|
||||
if err != nil {
|
||||
return 0, false
|
||||
}
|
||||
return val, true
|
||||
}
|
||||
|
||||
func parseQuery(query string) (*parsedQuery, error) {
|
||||
var (
|
||||
keyStart, keyEnd int
|
||||
valStart, valEnd int
|
||||
firstInfohash string
|
||||
|
||||
onKey = true
|
||||
hasInfohash = false
|
||||
|
||||
pq = &parsedQuery{
|
||||
infohashes: nil,
|
||||
params: make(map[string]string),
|
||||
}
|
||||
)
|
||||
|
||||
for i, length := 0, len(query); i < length; i++ {
|
||||
separator := query[i] == '&' || query[i] == ';' || query[i] == '?'
|
||||
if separator || i == length-1 {
|
||||
if onKey {
|
||||
keyStart = i + 1
|
||||
continue
|
||||
}
|
||||
|
||||
if i == length-1 && !separator {
|
||||
if query[i] == '=' {
|
||||
continue
|
||||
}
|
||||
valEnd = i
|
||||
}
|
||||
|
||||
keyStr, err := url.QueryUnescape(query[keyStart : keyEnd+1])
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
valStr, err := url.QueryUnescape(query[valStart : valEnd+1])
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
pq.params[keyStr] = valStr
|
||||
|
||||
if keyStr == "info_hash" {
|
||||
if hasInfohash {
|
||||
// Multiple infohashes
|
||||
if pq.infohashes == nil {
|
||||
pq.infohashes = []string{firstInfoHash}
|
||||
}
|
||||
pq.infohashes = append(pq.infohashes, valStr)
|
||||
} else {
|
||||
firstInfohash = valStr
|
||||
hasInfohash = true
|
||||
}
|
||||
}
|
||||
|
||||
onKey = true
|
||||
keyStart = i + 1
|
||||
} else if query[i] == '=' {
|
||||
onKey = false
|
||||
valStart = i + 1
|
||||
} else if onKey {
|
||||
keyEnd = i
|
||||
} else {
|
||||
valEnd = i
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func validateParsedQuery(pq *parsedQuery) error {
|
||||
infohash, ok := pq["info_hash"]
|
||||
if infohash == "" {
|
||||
return errors.New("infohash does not exist")
|
||||
}
|
||||
peerId, ok := pq["peer_id"]
|
||||
if peerId == "" {
|
||||
return errors.New("peerId does not exist")
|
||||
}
|
||||
port, ok := pq.getUint64("port")
|
||||
if ok == false {
|
||||
return errors.New("port does not exist")
|
||||
}
|
||||
uploaded, ok := pq.getUint64("uploaded")
|
||||
if ok == false {
|
||||
return errors.New("uploaded does not exist")
|
||||
}
|
||||
downloaded, ok := pq.getUint64("downloaded")
|
||||
if ok == false {
|
||||
return errors.New("downloaded does not exist")
|
||||
}
|
||||
left, ok := pq.getUint64("left")
|
||||
if ok == false {
|
||||
return errors.New("left does not exist")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,155 @@
|
||||
package server
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"errors"
|
||||
"net"
|
||||
"net/http"
|
||||
"path"
|
||||
"strconv"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
|
||||
"github.com/jzelinskie/bufferpool"
|
||||
|
||||
"github.com/jzelinskie/chihaya/config"
|
||||
"github.com/jzelinskie/chihaya/storage"
|
||||
)
|
||||
|
||||
type Server struct {
|
||||
http.Server
|
||||
listener *net.Listener
|
||||
}
|
||||
|
||||
func New(conf *config.Config) {
|
||||
return &Server{
|
||||
Addr: conf.Addr,
|
||||
Handler: newHandler(conf),
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Server) Start() error {
|
||||
s.listener, err = net.Listen("tcp", config.Addr)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
s.Handler.terminated = false
|
||||
s.Serve(s.listener)
|
||||
s.Handler.waitgroup.Wait()
|
||||
s.Handler.storage.Shutdown()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Server) Stop() error {
|
||||
s.Handler.waitgroup.Wait()
|
||||
s.Handler.terminated = true
|
||||
return s.Handler.listener.Close()
|
||||
}
|
||||
|
||||
type handler struct {
|
||||
bufferpool *bufferpool.BufferPool
|
||||
conf *config.Config
|
||||
deltaRequests int64
|
||||
storage *storage.Storage
|
||||
terminated bool
|
||||
waitgroup sync.WaitGroup
|
||||
}
|
||||
|
||||
func newHandler(conf *config.Config) {
|
||||
return &Handler{
|
||||
bufferpool: bufferpool.New(conf.BufferPoolSize, 500),
|
||||
conf: conf,
|
||||
storage: storage.New(&conf.Storage),
|
||||
}
|
||||
}
|
||||
|
||||
func (h *handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
if h.terminated {
|
||||
return
|
||||
}
|
||||
|
||||
h.waitgroup.Add(1)
|
||||
defer h.waitgroup.Done()
|
||||
|
||||
if r.URL.Path == "/stats" {
|
||||
h.serveStats(&w, r)
|
||||
return
|
||||
}
|
||||
|
||||
dir, action := path.Split(requestPath)
|
||||
switch action {
|
||||
case "announce":
|
||||
h.serveAnnounce(&w, r)
|
||||
return
|
||||
case "scrape":
|
||||
// TODO
|
||||
h.serveScrape(&w, r)
|
||||
return
|
||||
default:
|
||||
buf := h.bufferpool.Take()
|
||||
fail(errors.New("Unknown action"), buf)
|
||||
h.writeResponse(&w, r, buf)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
func writeResponse(w *http.ResponseWriter, r *http.Request, buf *bytes.Buffer) {
|
||||
r.Close = true
|
||||
w.Header().Add("Content-Type", "text/plain")
|
||||
w.Header().Add("Connection", "close")
|
||||
w.Header().Add("Content-Length", strconv.Itoa(buf.Len()))
|
||||
w.Write(buf.Bytes())
|
||||
w.(http.Flusher).Flush()
|
||||
atomic.AddInt64(h.deltaRequests, 1)
|
||||
}
|
||||
|
||||
func fail(err error, buf *bytes.Buffer) {
|
||||
buf.WriteString("d14:failure reason")
|
||||
buf.WriteString(strconv.Itoa(len(err)))
|
||||
buf.WriteRune(':')
|
||||
buf.WriteString(err)
|
||||
buf.WriteRune('e')
|
||||
}
|
||||
|
||||
func validatePasskey(dir string, s *storage.Storage) (storage.User, error) {
|
||||
if len(dir) != 34 {
|
||||
return nil, errors.New("Your passkey is invalid")
|
||||
}
|
||||
passkey := dir[1:33]
|
||||
|
||||
user, exists, err := s.FindUser(passkey)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !exists {
|
||||
return nil, errors.New("Passkey not found")
|
||||
}
|
||||
|
||||
return user, nil
|
||||
}
|
||||
|
||||
func determineIP(r *http.Request, pq *parsedQuery) (string, error) {
|
||||
ip, ok := pq.params["ip"]
|
||||
if !ok {
|
||||
ip, ok = pq.params["ipv4"]
|
||||
if !ok {
|
||||
ips, ok := r.Header["X-Real-Ip"]
|
||||
if ok && len(ips) > 0 {
|
||||
ip = ips[0]
|
||||
} else {
|
||||
portIndex := len(r.RemoteAddr) - 1
|
||||
for ; portIndex >= 0; portIndex-- {
|
||||
if r.RemoteAddr[portIndex] == ':' {
|
||||
break
|
||||
}
|
||||
}
|
||||
if portIndex != -1 {
|
||||
ip = r.RemoteAddr[0:portIndex]
|
||||
} else {
|
||||
return "", errors.New("Failed to parse IP address")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return &ip, nil
|
||||
}
|
||||
Reference in New Issue
Block a user