ircclient.py 8.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269
  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. log.info(self.jobs)
  127. def handle_connect(self):
  128. self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
  129. self.socket.connect((self.HOST, self.PORT))
  130. self.socket.settimeout(2)
  131. log.info("connected")
  132. self.connected = True
  133. def handle_close(self):
  134. log.info("closed")
  135. self.connected = False
  136. self.socket.close()
  137. def get_data_from_irc(self):
  138. data = self.socket.recv(4096)
  139. if len(data) < 2:
  140. return ""
  141. # newline at the end?
  142. if not (data[-1] == 10 and data[-2] == 13):
  143. self.readbuffer += data
  144. return ""
  145. data = self.readbuffer + data
  146. self.readbuffer = b''
  147. # purge crap
  148. data = (data.replace(b'\x95', b'').replace(b'0xc2', b'').decode('utf-8', 'ignore'))
  149. return data
  150. def handle_read(self):
  151. lines = self.get_data_from_irc().splitlines()
  152. for line in lines:
  153. words = line.split(' ')
  154. if len(words) < 2:
  155. continue
  156. msg_from = words[0]
  157. comm = words[1]
  158. if self.joined_channel:
  159. # log.debug(line)
  160. pass
  161. if comm in self.IGNORE:
  162. continue
  163. if comm == "JOIN":
  164. if self.name not in msg_from: # msg "NICK joined the channel" not about me
  165. continue
  166. log.info("Joined channel %s", self.CHANNEL)
  167. self.time_joined = time.time()
  168. self.joined_channel = True
  169. continue
  170. if comm == "PRIVMSG":
  171. # private message not addressed to me
  172. if words[2] != self.name:
  173. continue
  174. log.info("privmsg: %s", line)
  175. dcc_args = get_dcc_args(line)
  176. if dcc_args is None:
  177. continue
  178. ip, port, size, filename = dcc_args
  179. job = filename_to_job(filename)
  180. self.set_job_state(job, 'downloading')
  181. self.busy = True
  182. downloaded_filename = self.netcat(ip, port, size, filename)
  183. self.set_job_state(job, 'unarchiving')
  184. # TODO save state in redis on ip port size filename + output of handle files
  185. files = unar(downloaded_filename, self.PATH)
  186. self.busy = False
  187. self.set_job_state(job, 'done')
  188. self.results_queue.put({'type': 'files', 'files': files})
  189. if comm == "PING" or msg_from == "PING": # respond ping to avoid getting kicked
  190. self.pong(line)
  191. continue
  192. if comm == "376": # END MOTD
  193. # MOTD complete, lets join the channel
  194. self.join_channel(self.CHANNEL)
  195. continue
  196. def netcat(self, ip, port, size, filename, job_key):
  197. filename = os.path.basename(filename).replace(" ", "_")
  198. log.info('netcat: %s %d %d %s', ip, port, size, filename)
  199. s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
  200. log.info("Receiving file")
  201. s.connect((ip, port))
  202. fname = os.path.join(self.PATH, filename)
  203. f = open(fname, 'wb')
  204. count = 0
  205. last_perc = 0
  206. while True:
  207. data = s.recv(16384)
  208. if len(data) == 0:
  209. log.info("No data received - finished")
  210. break
  211. count += len(data)
  212. f.write(data)
  213. perc = int(100 * count / size)
  214. if perc % 10 == 0 and perc != last_perc:
  215. log.info("Download percentage: %d", perc)
  216. self.set_job_state(job_key, "%d%%" % perc)
  217. last_perc = perc
  218. if count >= size:
  219. break
  220. s.close()
  221. f.close()
  222. return fname
  223. def pong(self, data):
  224. msg = data.replace("PING ", "")
  225. self.send_queue.put("PONG %s" % msg)
  226. def join_channel(self, channel):
  227. self.send_queue.put("JOIN %s" % channel)
  228. def process_send_queue(self):
  229. if self.send_queue.empty():
  230. return
  231. data = self.send_queue.get()
  232. if not data.startswith("PONG"):
  233. log.info("Sending %s", data)
  234. add = bytes(str(data), "utf-8") + bytes([13, 10])
  235. self.socket.send(add)