ircclient.py 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263
  1. #!/usr/bin/env python3
  2. import logging
  3. import os
  4. import queue
  5. import socket
  6. import time
  7. import utils
  8. import re
  9. from unzipper import unar
  10. from threading import Thread
  11. from collections import defaultdict
  12. logging.basicConfig(level=logging.DEBUG)
  13. log = logging.getLogger(__name__)
  14. log.setLevel(logging.DEBUG)
  15. MODE_SEARCH = 'search'
  16. MODE_BOOK = 'book'
  17. results_key = re.compile(r'_results_for[_ ]+(?P<key>.*?)\.', re.I)
  18. def query_to_job_key(query):
  19. job = query
  20. if job.startswith('!'):
  21. # !Ook Brandon Sanderson - [Skyward 01] - Skyward (retail) (epub).rar ::INFO:: 3.5MB
  22. # !Horla-new Brandon Sanderson - Skyward (US) (epub).epub
  23. job = ' '.join(job.split(' ')[1:]).strip()
  24. # Brandon Sanderson - [Skyward 01] - Skyward (retail) (epub).rar ::INFO:: 3.5MB
  25. # Brandon Sanderson - Skyward (US) (epub).epub
  26. job = re.sub(r'::.*$', '', job).strip()
  27. return job.lower()
  28. def filename_to_job(fname):
  29. match = results_key.search(fname)
  30. if match: # list of results
  31. return match.group('key').replace('_', ' ').lower()
  32. # brandon_sanderson_-_skyward_(uk)_(epub).rar
  33. return fname.replace('_', ' ').lower()
  34. def get_dcc_args(msg):
  35. msg = msg.split(':')[2]
  36. msg = msg.replace("\x01", "")
  37. if not msg.startswith("DCC"):
  38. return None
  39. args = msg.replace("DCC SEND ", "").split(" ")
  40. size = int(args.pop())
  41. port = int(args.pop())
  42. ip = utils.ip_from_decimal(int(args.pop()))
  43. filename = "_".join(args).replace('"', '')
  44. return ip, port, size, filename
  45. class IRCClient(Thread):
  46. TIME_TO_FIRST_COMMAND = 30
  47. HOST = "irc.irchighway.net"
  48. PORT = 6667
  49. CHANNEL = "#ebooks"
  50. PATH = "/tmp/"
  51. SEARCH_BOT = "searchook"
  52. # ^ config
  53. IGNORE = [
  54. "NOTICE",
  55. "PART",
  56. "QUIT",
  57. "332",
  58. "333",
  59. "372",
  60. "353",
  61. "366",
  62. "251",
  63. "252",
  64. "254",
  65. "255",
  66. "265",
  67. "266",
  68. "396"]
  69. joined_channel = False
  70. connected = False
  71. name = "bookbot" + utils.random_hash()
  72. readbuffer = b''
  73. time_joined = None
  74. jobs = defaultdict(dict)
  75. def __init__(self, command_queue, results_queue):
  76. super(IRCClient, self).__init__(daemon=True)
  77. self.socket = None
  78. self.command_queue = command_queue
  79. self.results_queue = results_queue
  80. self.send_queue = queue.Queue()
  81. nickstr = "NICK %s" % self.name
  82. userstr = "USER %s %s bla :%s" % (self.name, self.HOST, self.name)
  83. self.send_queue.put(nickstr)
  84. self.send_queue.put(userstr)
  85. def run(self):
  86. self.handle_connect()
  87. while True:
  88. if self.connected and self.joined_channel:
  89. self.handle_commands()
  90. self.process_send_queue()
  91. try:
  92. self.handle_read()
  93. except socket.timeout:
  94. continue
  95. except socket.error as e:
  96. log.error('socket error')
  97. log.exception(e)
  98. self.handle_connect()
  99. except Exception as e:
  100. log.exception(e)
  101. break
  102. self.handle_close()
  103. log.info('Exiting RUN')
  104. def handle_commands(self):
  105. if self.command_queue.empty():
  106. time.sleep(0.2)
  107. return
  108. elapsed = time.time() - self.time_joined
  109. if elapsed < self.TIME_TO_FIRST_COMMAND:
  110. log.info("commands to process, but we have to wait %d", self.TIME_TO_FIRST_COMMAND - elapsed)
  111. time.sleep(1)
  112. return
  113. command = self.command_queue.get()
  114. log.info("command %s", command)
  115. job = query_to_job_key(command['query'])
  116. self.set_job_state(job, 'pending')
  117. if command['mode'] == MODE_SEARCH:
  118. self.send_queue.put("PRIVMSG %s :@%s %s " % (self.CHANNEL, self.SEARCH_BOT, command['query']))
  119. return
  120. self.send_queue.put("PRIVMSG %s :%s " % (self.CHANNEL, command['query']))
  121. def set_job_state(self, job, state):
  122. log.info("Setting job [%s] to %s", job, state)
  123. self.jobs[job].update({'state': state})
  124. log.info(self.jobs)
  125. def handle_connect(self):
  126. self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
  127. self.socket.connect((self.HOST, self.PORT))
  128. self.socket.settimeout(2)
  129. log.info("connected")
  130. self.connected = True
  131. def handle_close(self):
  132. log.info("closed")
  133. self.connected = False
  134. self.socket.close()
  135. def get_data_from_irc(self):
  136. data = self.socket.recv(4096)
  137. if len(data) < 2:
  138. return ""
  139. # newline at the end?
  140. if not (data[-1] == 10 and data[-2] == 13):
  141. self.readbuffer += data
  142. return ""
  143. data = self.readbuffer + data
  144. self.readbuffer = b''
  145. # purge crap
  146. data = (data.replace(b'\x95', b'').replace(b'0xc2', b'').decode('utf-8', 'ignore'))
  147. return data
  148. def handle_read(self):
  149. lines = self.get_data_from_irc().splitlines()
  150. for line in lines:
  151. words = line.split(' ')
  152. if len(words) < 2:
  153. continue
  154. msg_from = words[0]
  155. comm = words[1]
  156. if self.joined_channel:
  157. # log.debug(line)
  158. pass
  159. if comm in self.IGNORE:
  160. continue
  161. if comm == "JOIN":
  162. if self.name not in msg_from: # msg "NICK joined the channel" not about me
  163. continue
  164. log.info("Joined channel %s", self.CHANNEL)
  165. self.time_joined = time.time()
  166. self.joined_channel = True
  167. continue
  168. if comm == "PRIVMSG":
  169. # private message not addressed to me
  170. if words[2] != self.name:
  171. continue
  172. log.info("privmsg: %s", line)
  173. dcc_args = get_dcc_args(line)
  174. if dcc_args is None:
  175. continue
  176. ip, port, size, filename = dcc_args
  177. job = filename_to_job(filename)
  178. self.set_job_state(job, 'downloading')
  179. downloaded_filename = self.netcat(ip, port, size, filename)
  180. self.set_job_state(job, 'unarchiving')
  181. # TODO save state in redis on ip port size filename + output of handle files
  182. files = unar(downloaded_filename, self.PATH)
  183. self.set_job_state(job, 'done')
  184. self.results_queue.put({'type': 'files', 'files': files})
  185. if comm == "PING" or msg_from == "PING": # respond ping to avoid getting kicked
  186. self.pong(line)
  187. continue
  188. if comm == "376": # END MOTD
  189. # MOTD complete, lets join the channel
  190. self.join_channel(self.CHANNEL)
  191. continue
  192. def netcat(self, ip, port, size, filename):
  193. filename = os.path.basename(filename).replace(" ", "_")
  194. log.info('netcat: %s %d %d %s', ip, port, size, filename)
  195. s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
  196. log.info("Receiving file")
  197. s.connect((ip, port))
  198. fname = os.path.join(self.PATH, filename)
  199. f = open(fname, 'wb')
  200. count = 0
  201. last_perc = 0
  202. while True:
  203. data = s.recv(16384)
  204. if len(data) == 0:
  205. log.info("No data received - finished")
  206. break
  207. count += len(data)
  208. f.write(data)
  209. perc = int(100 * count / size)
  210. if perc % 10 == 0 and perc != last_perc:
  211. log.info("Download percentage: %d", perc)
  212. last_perc = perc
  213. if count >= size:
  214. break
  215. s.close()
  216. f.close()
  217. return fname
  218. def pong(self, data):
  219. msg = data.replace("PING ", "")
  220. self.send_queue.put("PONG %s" % msg)
  221. def join_channel(self, channel):
  222. self.send_queue.put("JOIN %s" % channel)
  223. def process_send_queue(self):
  224. if self.send_queue.empty():
  225. return
  226. data = self.send_queue.get()
  227. log.info("Sending %s", data)
  228. add = bytes(str(data), "utf-8") + bytes([13, 10])
  229. self.socket.send(add)