ircclient.py 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270
  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. busy = False
  76. def __init__(self, command_queue, results_queue):
  77. super(IRCClient, self).__init__(daemon=True)
  78. self.socket = None
  79. self.command_queue = command_queue
  80. self.results_queue = results_queue
  81. self.send_queue = queue.Queue()
  82. nickstr = "NICK %s" % self.name
  83. userstr = "USER %s %s bla :%s" % (self.name, self.HOST, self.name)
  84. self.send_queue.put(nickstr)
  85. self.send_queue.put(userstr)
  86. def run(self):
  87. self.handle_connect()
  88. while True:
  89. if self.connected and self.joined_channel:
  90. self.handle_commands()
  91. self.process_send_queue()
  92. try:
  93. self.handle_read()
  94. except socket.timeout:
  95. continue
  96. except socket.error as e:
  97. log.error('socket error')
  98. log.exception(e)
  99. self.handle_connect()
  100. except Exception as e:
  101. log.exception(e)
  102. break
  103. self.handle_close()
  104. log.info('Exiting RUN')
  105. def handle_commands(self):
  106. if self.command_queue.empty():
  107. time.sleep(0.2)
  108. return
  109. elapsed = time.time() - self.time_joined
  110. if elapsed < self.TIME_TO_FIRST_COMMAND:
  111. log.info("commands to process, but we have to wait %d", self.TIME_TO_FIRST_COMMAND - elapsed)
  112. time.sleep(1)
  113. return
  114. command = self.command_queue.get()
  115. log.info("command %s", command)
  116. job = query_to_job_key(command['query'])
  117. self.set_job_state(job, 'pending')
  118. if command['mode'] == MODE_SEARCH:
  119. self.send_queue.put("PRIVMSG %s :@%s %s " % (self.CHANNEL, self.SEARCH_BOT, command['query']))
  120. return
  121. self.send_queue.put("PRIVMSG %s :%s " % (self.CHANNEL, command['query']))
  122. def set_job_state(self, job, state):
  123. log.info("Setting job [%s] to %s", job, state)
  124. self.jobs[job].update({'state': state})
  125. self.results_queue.put({'type': 'status', 'status': state, 'key': job})
  126. if state == 'done':
  127. del self.jobs[job]
  128. def handle_connect(self):
  129. self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
  130. self.socket.connect((self.HOST, self.PORT))
  131. self.socket.settimeout(2)
  132. log.info("connected")
  133. self.connected = True
  134. def handle_close(self):
  135. log.info("closed")
  136. self.connected = False
  137. self.socket.close()
  138. def get_data_from_irc(self):
  139. data = self.socket.recv(4096)
  140. if len(data) < 2:
  141. return ""
  142. # newline at the end?
  143. if not (data[-1] == 10 and data[-2] == 13):
  144. self.readbuffer += data
  145. return ""
  146. data = self.readbuffer + data
  147. self.readbuffer = b''
  148. # purge crap
  149. data = (data.replace(b'\x95', b'').replace(b'0xc2', b'').decode('utf-8', 'ignore'))
  150. return data
  151. def handle_read(self):
  152. lines = self.get_data_from_irc().splitlines()
  153. for line in lines:
  154. words = line.split(' ')
  155. if len(words) < 2:
  156. continue
  157. msg_from = words[0]
  158. comm = words[1]
  159. if self.joined_channel:
  160. # log.debug(line)
  161. pass
  162. if comm in self.IGNORE:
  163. continue
  164. if comm == "JOIN":
  165. if self.name not in msg_from: # msg "NICK joined the channel" not about me
  166. continue
  167. log.info("Joined channel %s", self.CHANNEL)
  168. self.time_joined = time.time()
  169. self.joined_channel = True
  170. continue
  171. if comm == "PRIVMSG":
  172. # private message not addressed to me
  173. if words[2] != self.name:
  174. continue
  175. log.info("privmsg: %s", line)
  176. dcc_args = get_dcc_args(line)
  177. if dcc_args is None:
  178. continue
  179. ip, port, size, filename = dcc_args
  180. job = filename_to_job(filename)
  181. self.set_job_state(job, 'downloading')
  182. self.busy = True
  183. downloaded_filename = self.netcat(ip, port, size, filename, job)
  184. self.set_job_state(job, 'unarchiving')
  185. # TODO save state in redis on ip port size filename + output of handle files
  186. files = unar(downloaded_filename, self.PATH)
  187. self.busy = False
  188. self.results_queue.put({'type': 'files', 'files': files})
  189. self.set_job_state(job, 'done')
  190. if comm == "PING" or msg_from == "PING": # respond ping to avoid getting kicked
  191. self.pong(line)
  192. continue
  193. if comm == "376": # END MOTD
  194. # MOTD complete, lets join the channel
  195. self.join_channel(self.CHANNEL)
  196. continue
  197. def netcat(self, ip, port, size, filename, job_key):
  198. filename = os.path.basename(filename).replace(" ", "_")
  199. log.info('netcat: %s %d %d %s', ip, port, size, filename)
  200. s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
  201. log.info("Receiving file")
  202. s.connect((ip, port))
  203. fname = os.path.join(self.PATH, filename)
  204. f = open(fname, 'wb')
  205. count = 0
  206. last_perc = 0
  207. while True:
  208. data = s.recv(16384)
  209. if len(data) == 0:
  210. log.info("No data received - finished")
  211. break
  212. count += len(data)
  213. f.write(data)
  214. perc = int(100 * count / size)
  215. if perc % 10 == 0 and perc != last_perc:
  216. log.info("Download percentage: %d", perc)
  217. self.set_job_state(job_key, "%d%%" % perc)
  218. last_perc = perc
  219. if count >= size:
  220. break
  221. s.close()
  222. f.close()
  223. return fname
  224. def pong(self, data):
  225. msg = data.replace("PING ", "")
  226. self.send_queue.put("PONG %s" % msg)
  227. def join_channel(self, channel):
  228. self.send_queue.put("JOIN %s" % channel)
  229. def process_send_queue(self):
  230. if self.send_queue.empty():
  231. return
  232. data = self.send_queue.get()
  233. if not data.startswith("PONG"):
  234. log.info("Sending %s", data)
  235. add = bytes(str(data), "utf-8") + bytes([13, 10])
  236. self.socket.send(add)