ircclient.py 8.8 KB

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