-
Notifications
You must be signed in to change notification settings - Fork 1.7k
/
redis.py
275 lines (256 loc) · 9.63 KB
/
redis.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
# -*- coding: UTF-8 -*-
"""
@author: hhyo、yyukai
@license: Apache Licence
@file: redis.py
@time: 2019/03/26
"""
import json
import re
import shlex
import redis
import logging
import traceback
from common.utils.timer import FuncTimer
from . import EngineBase
from .models import ResultSet, ReviewSet, ReviewResult
__author__ = "hhyo"
logger = logging.getLogger("default")
class RedisEngine(EngineBase):
def get_connection(self, db_name=None):
db_name = db_name or self.db_name
if self.mode == "cluster":
return redis.cluster.RedisCluster(
host=self.host,
port=self.port,
username=self.user,
password=self.password or None,
encoding_errors="ignore",
decode_responses=True,
socket_connect_timeout=10,
ssl=self.instance.is_ssl,
)
else:
return redis.Redis(
host=self.host,
port=self.port,
db=db_name,
username=self.user,
password=self.password or None,
encoding_errors="ignore",
decode_responses=True,
socket_connect_timeout=10,
ssl=self.instance.is_ssl,
)
name = "Redis"
info = "Redis engine"
def test_connection(self):
return self.get_all_databases()
def get_all_databases(self, **kwargs):
"""
获取数据库列表
:return:
"""
result = ResultSet(full_sql="CONFIG GET databases")
conn = self.get_connection()
try:
rows = conn.config_get("databases")["databases"]
except Exception as e:
"""
由于尝试获取databases配置失败,下面的代码块将通过解析info命令的输出来确定数据库的数量。
失败场景1:AWS-ElastiCache(Redis)服务不支持部分命令行。比如: config get xx, acl 部分命令
失败场景2:使用了没有管理员权限(-@admin)的Redis用户。 (异常信息:this user has no permissions to run the 'config' command or its subcommand)
步骤:
- 通过info("Keyspace")获取所有的数据库键空间信息。
- 从键空间信息中提取数据库编号(如db0, db1等)。
- 计算数据库数量,至少会返回0到15共16个数据库。
"""
logger.warning(f"Redis CONFIG GET databases 执行报错,异常信息:{e}")
dbs = [
int(i.split("db")[1])
for i in conn.info("Keyspace").keys()
if len(i.split("db")) == 2
]
rows = max(dbs + [15]) + 1
db_list = [str(x) for x in range(int(rows))]
result.rows = db_list
return result
def query_check(self, db_name=None, sql="", limit_num=0):
"""提交查询前的检查"""
result = {"msg": "", "bad_query": True, "filtered_sql": sql, "has_star": False}
safe_cmd = [
"scan",
"exists",
"ttl",
"pttl",
"type",
"get",
"mget",
"strlen",
"hgetall",
"hexists",
"hget",
"hmget",
"hkeys",
"hvals",
"smembers",
"scard",
"sdiff",
"sunion",
"sismember",
"llen",
"lrange",
"lindex",
"zrange",
"zrangebyscore",
"zscore",
"zcard",
"zcount",
"zrank",
"info",
]
# 命令校验,仅可以执行safe_cmd内的命令
for cmd in safe_cmd:
if re.match(rf"^{cmd}", sql.strip(), re.I):
result["bad_query"] = False
break
if result["bad_query"]:
result["msg"] = "禁止执行该命令!"
return result
def processlist(self, command_type, **kwargs):
"""获取连接信息"""
sql = "client list"
result_set = ResultSet(full_sql=sql)
conn = self.get_connection(db_name=0)
clients = conn.client_list()
# 根据空闲时间排序
sort_by = "idle"
reverse = False
clients = sorted(
clients, key=lambda client: client.get(sort_by), reverse=reverse
)
result_set.rows = clients
return result_set
def query(self, db_name=None, sql="", limit_num=0, close_conn=True, **kwargs):
"""返回 ResultSet"""
result_set = ResultSet(full_sql=sql)
try:
conn = self.get_connection(db_name=db_name)
rows = conn.execute_command(*shlex.split(sql))
result_set.column_list = ["Result"]
if isinstance(rows, list) or isinstance(rows, tuple):
if re.match(rf"^scan", sql.strip(), re.I):
keys = [[row] for row in rows[1]]
keys.insert(0, [rows[0]])
result_set.rows = tuple(keys)
result_set.affected_rows = len(rows[1])
else:
result_set.rows = tuple([row] for row in rows)
result_set.affected_rows = len(rows)
elif isinstance(rows, dict):
result_set.column_list = ["field", "value"]
# 对Redis的返回结果进行类型判断,如果是dict,list转为json字符串。
pairs_list = []
for k, v in rows.items():
if isinstance(v, dict):
processed_value = json.dumps(v)
elif isinstance(v, list):
processed_value = json.dumps(v)
else:
processed_value = v
# 添加处理后的键值对到列表
pairs_list.append([k, processed_value])
# 将列表转换为元组并赋值给 result_set.rows
result_set.rows = tuple(pairs_list)
result_set.affected_rows = len(result_set.rows)
else:
result_set.rows = tuple([[rows]])
result_set.affected_rows = 1 if rows else 0
if limit_num > 0:
result_set.rows = result_set.rows[0:limit_num]
except Exception as e:
logger.warning(
f"Redis命令执行报错,语句:{sql}, 错误信息:{traceback.format_exc()}"
)
result_set.error = str(e)
return result_set
def filter_sql(self, sql="", limit_num=0):
return sql.strip()
def query_masking(self, db_name=None, sql="", resultset=None):
"""不做脱敏"""
return resultset
def execute_check(self, db_name=None, sql=""):
"""上线单执行前的检查, 返回Review set"""
check_result = ReviewSet(full_sql=sql)
split_sql = [cmd.strip() for cmd in sql.split("\n") if cmd.strip()]
line = 1
for cmd in split_sql:
result = ReviewResult(
id=line,
errlevel=0,
stagestatus="Audit completed",
errormessage="暂不支持显示影响行数",
sql=cmd,
affected_rows=0,
execute_time=0,
)
check_result.rows += [result]
line += 1
return check_result
def execute_workflow(self, workflow):
"""执行上线单,返回Review set"""
sql = workflow.sqlworkflowcontent.sql_content
split_sql = [cmd.strip() for cmd in sql.split("\n") if cmd.strip()]
execute_result = ReviewSet(full_sql=sql)
line = 1
cmd = None
try:
conn = self.get_connection(db_name=workflow.db_name)
for cmd in split_sql:
with FuncTimer() as t:
conn.execute_command(*shlex.split(cmd))
execute_result.rows.append(
ReviewResult(
id=line,
errlevel=0,
stagestatus="Execute Successfully",
errormessage="暂不支持显示影响行数",
sql=cmd,
affected_rows=0,
execute_time=t.cost,
)
)
line += 1
except Exception as e:
logger.warning(
f"Redis命令执行报错,语句:{cmd or sql}, 错误信息:{traceback.format_exc()}"
)
# 追加当前报错语句信息到执行结果中
execute_result.error = str(e)
execute_result.rows.append(
ReviewResult(
id=line,
errlevel=2,
stagestatus="Execute Failed",
errormessage=f"异常信息:{e}",
sql=cmd,
affected_rows=0,
execute_time=0,
)
)
line += 1
# 报错语句后面的语句标记为审核通过、未执行,追加到执行结果中
for statement in split_sql[line - 1 :]:
execute_result.rows.append(
ReviewResult(
id=line,
errlevel=0,
stagestatus="Audit completed",
errormessage=f"前序语句失败, 未执行",
sql=statement,
affected_rows=0,
execute_time=0,
)
)
line += 1
return execute_result