Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
xindoo
redis
提交
abb81c63
R
redis
项目概览
xindoo
/
redis
通知
2
Star
2
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
R
redis
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
体验新版 GitCode,发现更多精彩内容 >>
提交
abb81c63
编写于
2月 11, 2020
作者:
A
antirez
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Tracking: BCAST: registration in the prefix table.
上级
77da9608
变更
3
隐藏空白更改
内联
并排
Showing
3 changed file
with
67 addition
and
20 deletion
+67
-20
src/networking.c
src/networking.c
+10
-11
src/server.h
src/server.h
+3
-6
src/tracking.c
src/tracking.c
+54
-3
未找到文件。
src/networking.c
浏览文件 @
abb81c63
...
...
@@ -154,7 +154,7 @@ client *createClient(connection *conn) {
c
->
peerid
=
NULL
;
c
->
client_list_node
=
NULL
;
c
->
client_tracking_redirection
=
0
;
c
->
client_tracking_prefix
_nod
es
=
NULL
;
c
->
client_tracking_prefixes
=
NULL
;
c
->
auth_callback
=
NULL
;
c
->
auth_callback_privdata
=
NULL
;
c
->
auth_module
=
NULL
;
...
...
@@ -2028,7 +2028,6 @@ int clientSetNameOrReply(client *c, robj *name) {
void
clientCommand
(
client
*
c
)
{
listNode
*
ln
;
listIter
li
;
client
*
client
;
if
(
c
->
argc
==
2
&&
!
strcasecmp
(
c
->
argv
[
1
]
->
ptr
,
"help"
))
{
const
char
*
help
[]
=
{
...
...
@@ -2142,7 +2141,7 @@ NULL
/* Iterate clients killing all the matching clients. */
listRewind
(
server
.
clients
,
&
li
);
while
((
ln
=
listNext
(
&
li
))
!=
NULL
)
{
client
=
listNodeValue
(
ln
);
client
*
client
=
listNodeValue
(
ln
);
if
(
addr
&&
strcmp
(
getClientPeerId
(
client
),
addr
)
!=
0
)
continue
;
if
(
type
!=
-
1
&&
getClientType
(
client
)
!=
type
)
continue
;
if
(
id
!=
0
&&
client
->
id
!=
id
)
continue
;
...
...
@@ -2229,7 +2228,7 @@ NULL
size_t
numprefix
=
0
;
/* Parse the options. */
if
(
for
int
j
=
3
;
j
<
argc
;
j
++
)
{
for
(
int
j
=
3
;
j
<
c
->
argc
;
j
++
)
{
int
moreargs
=
(
c
->
argc
-
1
)
-
j
;
if
(
!
strcasecmp
(
c
->
argv
[
j
]
->
ptr
,
"redirect"
)
&&
moreargs
)
{
...
...
@@ -2246,10 +2245,10 @@ NULL
}
}
else
if
(
!
strcasecmp
(
c
->
argv
[
j
]
->
ptr
,
"bcast"
))
{
bcast
++
;
}
else
if
(
!
strcasecmp
(
c
->
argv
[
j
]
->
ptr
,
"prefix"
)
&&
morearg
)
{
}
else
if
(
!
strcasecmp
(
c
->
argv
[
j
]
->
ptr
,
"prefix"
)
&&
morearg
s
)
{
j
++
;
prefix
=
zrealloc
(
sizeof
(
robj
*
)
*
(
numprefix
+
1
));
prefix
[
numprefix
++
]
=
argv
[
j
];
prefix
=
zrealloc
(
prefix
,
sizeof
(
robj
*
)
*
(
numprefix
+
1
));
prefix
[
numprefix
++
]
=
c
->
argv
[
j
];
}
else
{
addReply
(
c
,
shared
.
syntaxerr
);
return
;
...
...
@@ -2259,16 +2258,16 @@ NULL
/* Make sure options are compatible among each other and with the
* current state of the client. */
if
(
!
bcast
&&
numprefix
)
{
addReplyError
(
"PREFIX option requires BCAST mode to be enabled"
);
addReplyError
(
c
,
"PREFIX option requires BCAST mode to be enabled"
);
zfree
(
prefix
);
return
;
}
if
(
c
lient
->
flags
&
CLIENT_TRACKING
)
{
int
oldbcast
=
!!
c
lient
->
flags
&
CLIENT_TRACKING_BCAST
;
if
(
c
->
flags
&
CLIENT_TRACKING
)
{
int
oldbcast
=
!!
c
->
flags
&
CLIENT_TRACKING_BCAST
;
if
(
oldbcast
!=
bcast
)
{
}
addReplyError
(
addReplyError
(
c
,
"You can't switch BCAST mode on/off before disabling "
"tracking for this client, and then re-enabling it with "
"a different mode."
);
...
...
src/server.h
浏览文件 @
abb81c63
...
...
@@ -823,12 +823,9 @@ typedef struct client {
* invalidation messages for keys fetched by this client will be send to
* the specified client ID. */
uint64_t
client_tracking_redirection
;
list
*
client_tracking_prefix_nodes
;
/* This list contains listNode pointers
to the nodes we have in every list
of clients in the tracking bcast
table. This way we can remove our
client in O(1) for each list. */
rax
*
client_tracking_prefixes
;
/* A dictionary of prefixes we are already
subscribed to in BCAST mode, in the
context of client side caching. */
/* Response buffer */
int
bufpos
;
char
buf
[
PROTO_REPLY_CHUNK_BYTES
];
...
...
src/tracking.c
浏览文件 @
abb81c63
...
...
@@ -49,6 +49,15 @@ uint64_t TrackingTableTotalItems = 0; /* Total number of IDs stored across
are using server side for CSC. */
robj
*
TrackingChannelName
;
/* This is the structure that we have as value of the PrefixTable, and
* represents the list of keys modified, and the list of clients that need
* to be notified, for a given prefix. */
typedef
struct
bcastState
{
rax
*
keys
;
/* Keys modified in the current event loop cycle. */
rax
*
clients
;
/* Clients subscribed to the notification events for this
prefix. */
}
bcastState
;
/* Remove the tracking state from the client 'c'. Note that there is not much
* to do for us here, if not to decrement the counter of the clients in
* tracking mode, because we just store the ID of the client in the tracking
...
...
@@ -56,9 +65,51 @@ robj *TrackingChannelName;
* client with many entries in the table is removed, it would cost a lot of
* time to do the cleanup. */
void
disableTracking
(
client
*
c
)
{
/* If this client is in broadcasting mode, we need to unsubscribe it
* from all the prefixes it is registered to. */
if
(
c
->
flags
&
CLIENT_TRACKING_BCAST
)
{
raxIterator
ri
;
raxStart
(
&
ri
,
c
->
client_tracking_prefixes
);
raxSeek
(
&
ri
,
"^"
,
NULL
,
0
);
while
(
raxNext
(
&
ri
))
{
bcastState
*
bs
=
raxFind
(
PrefixTable
,
ri
.
key
,
ri
.
key_len
);
serverAssert
(
bs
!=
raxNotFound
);
raxRemove
(
bs
->
clients
,(
unsigned
char
*
)
&
c
,
sizeof
(
c
),
NULL
);
/* Was it the last client? Remove the prefix from the
* table. */
if
(
raxSize
(
bs
->
clients
)
==
0
)
{
raxFree
(
bs
->
clients
);
raxFree
(
bs
->
keys
);
zfree
(
bs
);
raxRemove
(
PrefixTable
,
ri
.
key
,
ri
.
key_len
,
NULL
);
}
}
raxStop
(
&
ri
);
}
/* Clear flags and adjust the count. */
if
(
c
->
flags
&
CLIENT_TRACKING
)
{
server
.
tracking_clients
--
;
c
->
flags
&=
~
(
CLIENT_TRACKING
|
CLIENT_TRACKING_BROKEN_REDIR
);
c
->
flags
&=
~
(
CLIENT_TRACKING
|
CLIENT_TRACKING_BROKEN_REDIR
|
CLIENT_TRACKING_BCAST
);
}
}
/* Set the client 'c' to track the prefix 'prefix'. If the client 'c' is
* already registered for the specified prefix, no operation is performed. */
void
enableBcastTrackingForPrefix
(
client
*
c
,
char
*
prefix
,
size_t
plen
)
{
bcastState
*
bs
=
raxFind
(
PrefixTable
,(
unsigned
char
*
)
prefix
,
sdslen
(
prefix
));
/* If this is the first client subscribing to such prefix, create
* the prefix in the table. */
if
(
bs
==
raxNotFound
)
{
bs
=
zmalloc
(
sizeof
(
*
bs
));
bs
->
keys
=
raxNew
();
bs
->
clients
=
raxNew
();
raxInsert
(
PrefixTable
,(
unsigned
char
*
)
prefix
,
plen
,
bs
,
NULL
);
}
if
(
raxTryInsert
(
bs
->
clients
,(
unsigned
char
*
)
&
c
,
sizeof
(
c
),
NULL
,
NULL
))
{
raxInsert
(
c
->
client_tracking_prefixes
,
(
unsigned
char
*
)
prefix
,
plen
,
NULL
,
NULL
);
}
}
...
...
@@ -83,9 +134,9 @@ void enableTracking(client *c, uint64_t redirect_to, int bcast, robj **prefix, s
if
(
bcast
)
{
c
->
flags
|=
CLIENT_TRACKING_BCAST
;
if
(
numprefix
==
0
)
enableBcastTrackingForPrefix
(
c
,
""
,
0
);
for
(
in
t
j
=
0
;
j
<
numprefix
;
j
++
)
{
for
(
size_
t
j
=
0
;
j
<
numprefix
;
j
++
)
{
sds
sdsprefix
=
prefix
[
j
]
->
ptr
;
enableBcastTrackingForPrefix
(
c
,
sdsprefix
,
sdslen
(
prefix
));
enableBcastTrackingForPrefix
(
c
,
sdsprefix
,
sdslen
(
sds
prefix
));
}
}
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录