from weixin import WebWeixin
import sys
import time
import queue
import logging
import os
import requests
class Subscription:
def __init__(self, queue):
self.lastRead = time.time()
self.queue = queue
# 扩展一下将数据存到数据库
class WXBot(WebWeixin):
def __init__(self, db_conn):
super(WXBot, self).__init__()
self.db_conn = db_conn
self.subscriptions = [] # type: List[Subscription]
def getUserNickName(self, id):
if id == self.User["UserName"]:
return self.User["NickName"]
for member in self.ContactList:
if member["UserName"] == id:
return member["NickName"]
return None
# 获取用户备注(没有则返回昵称)或群名称,若本地无相关数据返回_unknown
def getUserOrGroupName(self, id):
name = "_unknown"
if id == self.User["UserName"]:
return self.User["NickName"] # 自己
if id[:2] == "@@":
# 群
name = self.getGroupName(id)
else:
# 特殊账号
for member in self.SpecialUsersList:
if member["UserName"] == id:
name = (
member["RemarkName"]
if member["RemarkName"]
else member["NickName"]
)
# 公众号或服务号
for member in self.PublicUsersList:
if member["UserName"] == id:
name = (
member["RemarkName"]
if member["RemarkName"]
else member["NickName"]
)
# 直接联系人
for member in self.ContactList:
if member["UserName"] == id:
name = (
member["RemarkName"]
if member["RemarkName"]
else member["NickName"]
)
# 群友
for member in self.GroupMemeberList:
if member["UserName"] == id:
name = (
member["DisplayName"]
if member["DisplayName"]
else member["NickName"]
)
return name
# 订阅新消息
def subscribe(self):
print("new subscription")
sub = Subscription(queue.Queue(1024))
self.subscriptions.append(sub)
return sub
def unsubscribe(self, subscription):
self.subscriptions.remove(subscription)
def composeMessageID(self, msgID, timestamp):
return time.strftime("%Y%m%d%H%M%S", time.localtime(timestamp)) + msgID
def saveMessageImage(self, msgid, composedID):
url = self.base_uri + "/webwxgetmsgimg?MsgID=%s&skey=%s" % (msgid, self.skey)
data = self._get(url, "webwxgetmsgimg")
return self.saveFile(composedID[:4], composedID + ".jpg", data)
def saveVoice(self, msgid, composedID):
url = self.base_uri + "/webwxgetvoice?msgid=%s&skey=%s" % (msgid, self.skey)
data = self._get(url, api="webwxgetvoice")
return self.saveFile(composedID[:4], composedID + ".mp3", data)
def saveVideo(self, msgid, composedID):
url = self.base_uri + "/webwxgetvideo?msgid=%s&skey=%s" % (msgid, self.skey)
data = self._get(url, api="webwxgetvideo")
return self.saveFile(composedID[:4], composedID + ".mp4", data)
def saveFile(self, subdir, filename, data):
if data == "":
return False
dirName = os.path.join(self.saveFolder, subdir)
if not os.path.exists(dirName):
os.makedirs(dirName)
fn = os.path.join(dirName, filename)
with open(fn, "wb") as f:
f.write(data)
f.close()
return True
def handleMsg(self, r):
if self.DEBUG:
print("处理消息, NewMessageCount =", r["AddMsgCount"])
for msg in r["AddMsgList"]:
msgID = msg["MsgId"]
composedID = self.composeMessageID(msgID, msg["CreateTime"])
msgType = msg["MsgType"]
if (
msg["FromUserName"] in self.SpecialUsers
or msg["ToUserName"] in self.SpecialUsers
):
# 过滤特殊帐号消息
continue
elif msgType in (50, 51, 52, 53, 9999):
# 无用消息
continue
sql = "INSERT INTO msg (\
msg_id,\
create_time,\
msg_type,\
content,\
from_name,\
from_nickname,\
to_name,\
to_nickname,\
group_name\
) VALUES (\
:msg_id,\
:create_time,\
:msg_type,\
:content,\
:from_name,\
:from_nickname,\
:to_name,\
:to_nickname,\
:group_name\
)"
params = {
"msg_id": composedID,
"create_time": msg["CreateTime"],
"msg_type": msgType,
"content": None,
"from_id": msg["ToUserName"],
"from_name": self.getUserOrGroupName(msg["FromUserName"]),
"from_nickname": self.getUserNickName(msg["FromUserName"]),
"to_id": msg["ToUserName"],
"to_name": self.getUserOrGroupName(msg["ToUserName"]),
"to_nickname": self.getUserNickName(msg["ToUserName"]),
"group_id": None,
"group_name": None,
}
if msg["FromUserName"] == self.User["UserName"]:
params["from_name"] = "_self"
elif msg["ToUserName"] == self.User["UserName"]:
params["to_name"] = "_self"
if msg["FromUserName"][:2] == "@@":
# 群消息
params["from_id"] = params["from_name"] = params["from_nickname"] = None
params["to_id"] = params["to_name"] = params["to_nickname"] = None
params["group_id"] = msg["FromUserName"]
params["group_name"] = self.getUserOrGroupName(msg["FromUserName"])
if ":
" in msg["Content"]:
[people, content] = msg["Content"].split(":
", 1)
params["from_id"] = people
params["from_name"] = self.getUserOrGroupName(people)
elif msg["ToUserName"][:2] == "@@":
# 发送的群消息
params["group_id"] = msg["ToUserName"]
params["group_name"] = self.getUserOrGroupName(msg["ToUserName"])
params["from_id"] = self.User["UserName"]
params["from_name"] = self.User["NickName"]
params["to_id"] = params["to_name"] = params["to_nickname"] = None
if msgType == 1:
# 文字消息
if params["from_name"] == "_unknown" or params["to_name"] == "_unknown":
# 有陌生人的文字消息,可能是新添加的好友,重新获取联系人列表
self.webwxgetcontact()
# 重新获取用户名称
params["from_name"] = self.getUserOrGroupName(msg["FromUserName"])
params["to_name"] = self.getUserOrGroupName(msg["ToUserName"])
params["content"] = (
msg["Content"].replace("<", "<").replace(">", ">")
)
elif msgType == 3:
params["content"] = "[图片](%s)" % composedID
self.saveMessageImage(msgID, composedID)
elif msgType == 34:
params["content"] = "[语音](%s)" % composedID
self.saveVoice(msgID, composedID)
elif msgType == 43:
params["content"] = "[视频](%s)" % composedID
self.saveVideo(msgID, composedID)
elif msgType == 62:
params["msg_type"] = 43
params["content"] = "[小视频](%s)" % composedID
self.saveVideo(msgID, composedID)
elif msgType == 42:
params["msg_type"] = 0
params["content"] = "[名片](%s)" % msg["RecommendInfo"]["NickName"]
elif msgType == 47:
if msg["HasProductId"] == 1:
# 有版权的表情包,无法在Web端查看
params["msg_type"] = 0
params["content"] = "[表情包图片]"
else:
params["msg_type"] = 3
params["content"] = "[动画表情](%s)" % composedID
self.saveMessageImage(msgID, composedID)
elif msgType == 49:
params["msg_type"] = 0
params["content"] = "[链接:%s](%s)" % (msg["FileName"], msg["Url"])
elif msgType == 10000:
params["msg_type"] = 0
params["content"] = "[系统消息](%s)" % (msg["Content"])
elif msgType == 10002:
params["msg_type"] = 0
params["content"] = "[系统消息](%s撤回了一条消息)" % params["from_name"]
else:
params["msg_type"] = 0
params["content"] = "[未识别消息]" + msg["Content"]
self.db_conn.execute(sql, params)
self.db_conn.commit()
# 工作时间发送新消息提醒到飞书
curr_time_struct = time.localtime(time.time())
if params['from_name'] <> '_self' and curr_time_struct.tm_hour >= 10 and curr_time_struct.tm_hour < 19:
requests.post(
"https://open.feishu.cn/open-apis/bot/v2/hook/30bda541-46de-40b7-8881-31b77deadc35",
headers={"Content-Type": "application/json"},
json={
"msg_type": "post",
"content": {
"post": {
"zh_cn": {
"title": "【" + params["from_name"] + "】发来一条新消息",
"content": [
[
{
"tag": "text",
"text": "【微信】收到一条来自【"
+ params["from_name"]
+ "】的消息",
}
]
],
}
}
},
},
)
for subscription in self.subscriptions:
if (
not subscription.queue.empty()
and time.time() - subscription.lastRead > 120
):
# 超时关闭
subscription.queue.put_nowait(None)
self.subscriptions.remove(subscription)
# 将消息放入队列,推送到客户端
sent_by_self = msg["FromUserName"] == self.User["UserName"]
message = {
"uid": params["to_id"] if sent_by_self else params["from_id"],
"content": params["content"],
"name": params["to_name"] if sent_by_self else params["from_name"],
"self": sent_by_self,
"type": params["msg_type"],
"gid": params["group_id"],
"group": params["group_name"],
}
if subscription.queue.full():
subscription.queue.get_nowait()
subscription.queue.put_nowait(message)