ipc.go 7.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274
  1. package main
  2. import (
  3. "bytes"
  4. "container/heap"
  5. "fmt"
  6. "log"
  7. "net"
  8. "time"
  9. "git.torproject.org/pluggable-transports/snowflake.git/common/messages"
  10. // "github.com/prometheus/client_golang/prometheus"
  11. )
  12. const (
  13. ClientTimeout = 10
  14. ProxyTimeout = 10
  15. NATUnknown = "unknown"
  16. NATRestricted = "restricted"
  17. NATUnrestricted = "unrestricted"
  18. )
  19. type clientVersion int
  20. const (
  21. v1 clientVersion = iota
  22. )
  23. type IPC struct {
  24. ctx *BrokerContext
  25. }
  26. func (i *IPC) Debug(_ interface{}, response *string) error {
  27. var webexts, browsers, standalones, unknowns int
  28. var natRestricted, natUnrestricted, natUnknown int
  29. i.ctx.snowflakeLock.Lock()
  30. s := fmt.Sprintf("current snowflakes available: %d\n", len(i.ctx.idToSnowflake))
  31. for _, snowflake := range i.ctx.idToSnowflake {
  32. if snowflake.proxyType == "badge" {
  33. browsers++
  34. } else if snowflake.proxyType == "webext" {
  35. webexts++
  36. } else if snowflake.proxyType == "standalone" {
  37. standalones++
  38. } else {
  39. unknowns++
  40. }
  41. switch snowflake.natType {
  42. case NATRestricted:
  43. natRestricted++
  44. case NATUnrestricted:
  45. natUnrestricted++
  46. default:
  47. natUnknown++
  48. }
  49. }
  50. i.ctx.snowflakeLock.Unlock()
  51. s += fmt.Sprintf("\tstandalone proxies: %d", standalones)
  52. s += fmt.Sprintf("\n\tbrowser proxies: %d", browsers)
  53. s += fmt.Sprintf("\n\twebext proxies: %d", webexts)
  54. s += fmt.Sprintf("\n\tunknown proxies: %d", unknowns)
  55. s += fmt.Sprintf("\nNAT Types available:")
  56. s += fmt.Sprintf("\n\trestricted: %d", natRestricted)
  57. s += fmt.Sprintf("\n\tunrestricted: %d", natUnrestricted)
  58. s += fmt.Sprintf("\n\tunknown: %d", natUnknown)
  59. *response = s
  60. return nil
  61. }
  62. func (i *IPC) ProxyPolls(arg messages.Arg, response *[]byte) error {
  63. sid, proxyType, natType, err := messages.DecodePollRequest(arg.Body)
  64. if err != nil {
  65. return messages.ErrBadRequest
  66. }
  67. // Log geoip stats
  68. remoteIP, _, err := net.SplitHostPort(arg.RemoteAddr)
  69. if err != nil {
  70. log.Println("Error processing proxy IP: ", err.Error())
  71. } else {
  72. i.ctx.metrics.lock.Lock()
  73. i.ctx.metrics.UpdateCountryStats(remoteIP, proxyType, natType)
  74. i.ctx.metrics.lock.Unlock()
  75. }
  76. var b []byte
  77. // Wait for a client to avail an offer to the snowflake, or timeout if nil.
  78. offer := i.ctx.RequestOffer(sid, proxyType, natType)
  79. if offer == nil {
  80. i.ctx.metrics.lock.Lock()
  81. i.ctx.metrics.proxyIdleCount++
  82. // i.ctx.metrics.promMetrics.ProxyPollTotal.With(prometheus.Labels{"nat": natType, "status": "idle"}).Inc()
  83. i.ctx.metrics.lock.Unlock()
  84. b, err = messages.EncodePollResponse("", false, "")
  85. if err != nil {
  86. return messages.ErrInternal
  87. }
  88. *response = b
  89. return nil
  90. }
  91. // i.ctx.metrics.promMetrics.ProxyPollTotal.With(prometheus.Labels{"nat": natType, "status": "matched"}).Inc()
  92. b, err = messages.EncodePollResponse(string(offer.sdp), true, offer.natType)
  93. if err != nil {
  94. return messages.ErrInternal
  95. }
  96. *response = b
  97. return nil
  98. }
  99. func sendClientResponse(resp *messages.ClientPollResponse, response *[]byte) error {
  100. data, err := resp.EncodePollResponse()
  101. if err != nil {
  102. log.Printf("error encoding answer")
  103. return messages.ErrInternal
  104. } else {
  105. *response = []byte(data)
  106. return nil
  107. }
  108. }
  109. func (i *IPC) ClientOffers(arg messages.Arg, response *[]byte) error {
  110. var version clientVersion
  111. startTime := time.Now()
  112. body := arg.Body
  113. parts := bytes.SplitN(body, []byte("\n"), 2)
  114. if len(parts) < 2 {
  115. // no version number found
  116. err := fmt.Errorf("unsupported message version")
  117. return sendClientResponse(&messages.ClientPollResponse{Error: err.Error()}, response)
  118. }
  119. body = parts[1]
  120. if string(parts[0]) == "1.0" {
  121. version = v1
  122. } else {
  123. err := fmt.Errorf("unsupported message version")
  124. return sendClientResponse(&messages.ClientPollResponse{Error: err.Error()}, response)
  125. }
  126. var offer *ClientOffer
  127. switch version {
  128. case v1:
  129. req, err := messages.DecodeClientPollRequest(body)
  130. if err != nil {
  131. return sendClientResponse(&messages.ClientPollResponse{Error: err.Error()}, response)
  132. }
  133. offer = &ClientOffer{
  134. natType: req.NAT,
  135. sdp: []byte(req.Offer),
  136. }
  137. default:
  138. panic("unknown version")
  139. }
  140. // Only hand out known restricted snowflakes to unrestricted clients
  141. var snowflakeHeap *SnowflakeHeap
  142. if offer.natType == NATUnrestricted {
  143. snowflakeHeap = i.ctx.restrictedSnowflakes
  144. } else {
  145. snowflakeHeap = i.ctx.snowflakes
  146. }
  147. // Immediately fail if there are no snowflakes available.
  148. i.ctx.snowflakeLock.Lock()
  149. numSnowflakes := snowflakeHeap.Len()
  150. i.ctx.snowflakeLock.Unlock()
  151. if numSnowflakes <= 0 {
  152. i.ctx.metrics.lock.Lock()
  153. i.ctx.metrics.clientDeniedCount++
  154. // i.ctx.metrics.promMetrics.ClientPollTotal.With(prometheus.Labels{"nat": offer.natType, "status": "denied"}).Inc()
  155. if offer.natType == NATUnrestricted {
  156. i.ctx.metrics.clientUnrestrictedDeniedCount++
  157. } else {
  158. i.ctx.metrics.clientRestrictedDeniedCount++
  159. }
  160. i.ctx.metrics.lock.Unlock()
  161. switch version {
  162. case v1:
  163. resp := &messages.ClientPollResponse{Error: "no snowflake proxies currently available"}
  164. return sendClientResponse(resp, response)
  165. default:
  166. panic("unknown version")
  167. }
  168. }
  169. // Otherwise, find the most available snowflake proxy, and pass the offer to it.
  170. // Delete must be deferred in order to correctly process answer request later.
  171. i.ctx.snowflakeLock.Lock()
  172. snowflake := heap.Pop(snowflakeHeap).(*Snowflake)
  173. i.ctx.snowflakeLock.Unlock()
  174. snowflake.offerChannel <- offer
  175. var err error
  176. // Wait for the answer to be returned on the channel or timeout.
  177. select {
  178. case answer := <-snowflake.answerChannel:
  179. i.ctx.metrics.lock.Lock()
  180. i.ctx.metrics.clientProxyMatchCount++
  181. // i.ctx.metrics.promMetrics.ClientPollTotal.With(prometheus.Labels{"nat": offer.natType, "status": "matched"}).Inc()
  182. i.ctx.metrics.lock.Unlock()
  183. switch version {
  184. case v1:
  185. resp := &messages.ClientPollResponse{Answer: answer}
  186. err = sendClientResponse(resp, response)
  187. default:
  188. panic("unknown version")
  189. }
  190. // Initial tracking of elapsed time.
  191. i.ctx.metrics.clientRoundtripEstimate = time.Since(startTime) / time.Millisecond
  192. case <-time.After(time.Second * ClientTimeout):
  193. log.Println("Client: Timed out.")
  194. switch version {
  195. case v1:
  196. resp := &messages.ClientPollResponse{
  197. Error: "timed out waiting for answer!"}
  198. err = sendClientResponse(resp, response)
  199. default:
  200. panic("unknown version")
  201. }
  202. }
  203. i.ctx.snowflakeLock.Lock()
  204. // i.ctx.metrics.promMetrics.AvailableProxies.With(prometheus.Labels{"nat": snowflake.natType, "type": snowflake.proxyType}).Dec()
  205. delete(i.ctx.idToSnowflake, snowflake.id)
  206. i.ctx.snowflakeLock.Unlock()
  207. return err
  208. }
  209. func (i *IPC) ProxyAnswers(arg messages.Arg, response *[]byte) error {
  210. answer, id, err := messages.DecodeAnswerRequest(arg.Body)
  211. if err != nil || answer == "" {
  212. return messages.ErrBadRequest
  213. }
  214. var success = true
  215. i.ctx.snowflakeLock.Lock()
  216. snowflake, ok := i.ctx.idToSnowflake[id]
  217. i.ctx.snowflakeLock.Unlock()
  218. if !ok || snowflake == nil {
  219. // The snowflake took too long to respond with an answer, so its client
  220. // disappeared / the snowflake is no longer recognized by the Broker.
  221. success = false
  222. }
  223. b, err := messages.EncodeAnswerResponse(success)
  224. if err != nil {
  225. log.Printf("Error encoding answer: %s", err.Error())
  226. return messages.ErrInternal
  227. }
  228. *response = b
  229. if success {
  230. snowflake.answerChannel <- answer
  231. }
  232. return nil
  233. }