Compare commits
No commits in common. "0fbafe4355d60591bf8db73e3c9cdf07cd1f9c41" and "748a4c2d9a6faf7a6d1554766e88af2bee94b8c5" have entirely different histories.
0fbafe4355
...
748a4c2d9a
20
priv/2024-01-18-log.txt
Normal file
20
priv/2024-01-18-log.txt
Normal file
@ -0,0 +1,20 @@
|
||||
{"device_uuid":"20118448129232486417014256831677", "key":"epi", "label":"kWh", "name":"正向总有功电能", "timestamp":1705576081, "type":"AI", "unit":5, "value":314498}
|
||||
{"device_uuid":"20118448129232486417014256831677", "key":"a_voltage", "label":"V", "name":"A相电压", "timestamp":1705576081, "type":"AI", "unit":0, "value":224.1}
|
||||
{"device_uuid":"20118448129232486417014256831677", "key":"a_current", "label":"A", "name":"A相电流", "timestamp":1705576081, "type":"AI", "unit":3, "value":0.394}
|
||||
{"device_uuid":"20118448129232486417014256831677", "key":"b_voltage", "label":"V", "name":"B相电压", "timestamp":1705576081, "type":"AI", "unit":0, "value":225.3}
|
||||
{"device_uuid":"20118448129232486417014256831677", "key":"b_current", "label":"A", "name":"B相电流", "timestamp":1705576081, "type":"AI", "unit":3, "value":0.323}
|
||||
{"device_uuid":"20118448129232486417014256831677", "key":"c_voltage", "label":"V", "name":"C相电压", "timestamp":1705576081, "type":"AI", "unit":0, "value":225}
|
||||
{"device_uuid":"20118448129232486417014256831677", "key":"c_current", "label":"A", "name":"C相电流", "timestamp":1705576081, "type":"AI", "unit":3, "value":0.269}
|
||||
{"device_uuid":"20118448129232486417014256831677", "key":"active_power", "label":"kW", "name":"瞬时总有功功率", "timestamp":1705576081, "type":"AI", "unit":23, "value":20.91}
|
||||
{"device_uuid":"20118448129232486417014256831677", "key":"power_factor", "label":"", "name":"总功率因数", "timestamp":1705576081, "type":"AI", "unit":16, "value":0.994}
|
||||
{"device_uuid":"20118448549921177617014256841701", "key":"epi", "label":"kWh", "name":"正向总有功电能", "timestamp":1705576038, "type":"AI", "unit":5, "value":299493.6}
|
||||
{"device_uuid":"20118448549921177617014256841701", "key":"a_voltage", "label":"V", "name":"A相电压", "timestamp":1705576038, "type":"AI", "unit":0, "value":225.9}
|
||||
{"device_uuid":"20118448549921177617014256841701", "key":"a_current", "label":"A", "name":"A相电流", "timestamp":1705576038, "type":"AI", "unit":3, "value":0.318}
|
||||
{"device_uuid":"20118448549921177617014256841701", "key":"b_voltage", "label":"V", "name":"B相电压", "timestamp":1705576038, "type":"AI", "unit":0, "value":225.7}
|
||||
{"device_uuid":"20118448549921177617014256841701", "key":"b_current", "label":"A", "name":"B相电流", "timestamp":1705576038, "type":"AI", "unit":3, "value":0.305}
|
||||
{"device_uuid":"20118448549921177617014256841701", "key":"c_voltage", "label":"V", "name":"C相电压", "timestamp":1705576038, "type":"AI", "unit":0, "value":225}
|
||||
{"device_uuid":"20118448549921177617014256841701", "key":"c_current", "label":"A", "name":"C相电流", "timestamp":1705576038, "type":"AI", "unit":3, "value":0.3}
|
||||
{"device_uuid":"20118448549921177617014256841701", "key":"active_power", "label":"kW", "name":"瞬时总有功功率", "timestamp":1705576038, "type":"AI", "unit":23, "value":15.056}
|
||||
{"device_uuid":"20118448549921177617014256841701", "key":"power_factor", "label":"", "name":"总功率因数", "timestamp":1705576038, "type":"AI", "unit":16, "value":1}
|
||||
{"device_uuid":"20118448970190438417014256851724", "key":"epi", "label":"kWh", "name":"正向总有功电能", "timestamp":1705576117, "type":"AI", "unit":5, "value":285011}
|
||||
{"device_uuid":"20118448970190438417014256851724", "key":"a_voltage", "label":"V", "name":"A相电压", "timestamp":1705576117, "type":"AI", "unit":0, "value":225.1}
|
||||
@ -1,23 +0,0 @@
|
||||
-----BEGIN CERTIFICATE-----
|
||||
MIIDwzCCAqugAwIBAgIUMhNyuVo6MieWs2I9LsyTTiDKtYgwDQYJKoZIhvcNAQEL
|
||||
BQAwgYkxCzAJBgNVBAYTAkNOMRAwDgYDVQQIDAdCZWlqaW5nMRAwDgYDVQQHDAdC
|
||||
ZWlqaW5nMSEwHwYDVQQKDBhJbnRlcm5ldCBXaWRnaXRzIFB0eSBMdGQxEjAQBgNV
|
||||
BAMMCWFubGljaGVuZzEfMB0GCSqGSIb3DQEJARYQMjQ0MTA4NzE1QHFxLmNvbTAe
|
||||
Fw0yNTA0MjEwMzUyNTlaFw0yNjA0MjEwMzUyNTlaMIGJMQswCQYDVQQGEwJDTjEQ
|
||||
MA4GA1UECAwHQmVpamluZzEQMA4GA1UEBwwHQmVpamluZzEhMB8GA1UECgwYSW50
|
||||
ZXJuZXQgV2lkZ2l0cyBQdHkgTHRkMRIwEAYDVQQDDAlhbmxpY2hlbmcxHzAdBgkq
|
||||
hkiG9w0BCQEWEDI0NDEwODcxNUBxcS5jb20wggEiMA0GCSqGSIb3DQEBAQUAA4IB
|
||||
DwAwggEKAoIBAQDdxxyYG0zV2KzeiVH5AGj1X1h7vLAHVS8FGkPL2QBBqW5/PrJ/
|
||||
z/sv3QJz6dB3ElOnk14GsY8lWk8uX/fjgNIPawN1G2/bCdMrlPOCEJEIOqZvzQWg
|
||||
PLRPGjoZx+wjCM/H7h47KVr4GNbo8MakLJeg6QB3rEEIiPrhQIoq9N7AhcjbEvWR
|
||||
NIxIQrpSqMCE2A5RoKUxAMum9rEoe/6PBw2GgEsQl5E6suRZw3wavg3aUU+6MMyx
|
||||
iiKEA6fCD5gDqgSo/xjgDUWJIOTTrmW6RBTcWP9iBUjWnCovQQ9zc6CzDXRgfsz9
|
||||
7uXQt9fuQEi++lcGnrsePjC99PnirCvmj/C/AgMBAAGjITAfMB0GA1UdDgQWBBT6
|
||||
n4jUTWQNcastC5jPHuU7CDzLFjANBgkqhkiG9w0BAQsFAAOCAQEAZJ9fY2Z+4vKr
|
||||
bBqwHmfBjEnGgS7L2mC7uS2x/x2meBRKlAlw5+nKdaUBccyI2baI3P1mh/iV72Wr
|
||||
OTcwUwSS6gIOB7JeWSB0UT7KbEOJ1tM/1HYs5F9tOT94P4Adm2gcKY81UlJMfSNQ
|
||||
WyWFLjWOk/5fUP42BmgbUIafTT9p+LeP6NOyEs6b4hGpF3q5L1QDwMUfASpOWtHn
|
||||
O7VyBvFyCGkVchnorWJ3ZXPaa7hy+2ULOK/d9zH3xxq4LKRclAS5XAMWHuw6/+tu
|
||||
CoxA+RPnaqHexPSzYlEElOT286FDyZHadjUDD0q/0Um92NUM3r+UrNEs/OHp7/nU
|
||||
M3BVM7H6LQ==
|
||||
-----END CERTIFICATE-----
|
||||
@ -1,18 +0,0 @@
|
||||
-----BEGIN CERTIFICATE REQUEST-----
|
||||
MIICzzCCAbcCAQAwgYkxCzAJBgNVBAYTAkNOMRAwDgYDVQQIDAdCZWlqaW5nMRAw
|
||||
DgYDVQQHDAdCZWlqaW5nMSEwHwYDVQQKDBhJbnRlcm5ldCBXaWRnaXRzIFB0eSBM
|
||||
dGQxEjAQBgNVBAMMCWFubGljaGVuZzEfMB0GCSqGSIb3DQEJARYQMjQ0MTA4NzE1
|
||||
QHFxLmNvbTCCASIwDQYJKoZIhvcNAQEBBQADggEPADCCAQoCggEBAN3HHJgbTNXY
|
||||
rN6JUfkAaPVfWHu8sAdVLwUaQ8vZAEGpbn8+sn/P+y/dAnPp0HcSU6eTXgaxjyVa
|
||||
Ty5f9+OA0g9rA3Ubb9sJ0yuU84IQkQg6pm/NBaA8tE8aOhnH7CMIz8fuHjspWvgY
|
||||
1ujwxqQsl6DpAHesQQiI+uFAiir03sCFyNsS9ZE0jEhCulKowITYDlGgpTEAy6b2
|
||||
sSh7/o8HDYaASxCXkTqy5FnDfBq+DdpRT7owzLGKIoQDp8IPmAOqBKj/GOANRYkg
|
||||
5NOuZbpEFNxY/2IFSNacKi9BD3NzoLMNdGB+zP3u5dC31+5ASL76Vwaeux4+ML30
|
||||
+eKsK+aP8L8CAwEAAaAAMA0GCSqGSIb3DQEBCwUAA4IBAQBHdTrwcS0Ip5XO/zJ/
|
||||
wFZvxe8wcdMBaHjSRx09Lgy3V1L9d0DPOfuUP4LP8UCxniucA2JLugJO1wr/5nZA
|
||||
FS9GcyCgbMKX/EskFkpuz72EA11WhyqFp/9nWZsxZB0t1cs3bUFeoFd2SE6QIo4N
|
||||
UY+guZCcF8hVwppceOxyaUQ9cudAClH9JSMR1XoIwvv3X7FJkPhM12DvcsLr062S
|
||||
fwG7/1h+VoSpAM/UYtBHBkU0MIVn1Gw6lJvkabUV0oYYBzl9ejsq1Qo2nUMdOEQY
|
||||
SS/VhVLONMyMCIERt3/GJTXdrJ12VqawJACFwv8g2i6NxbQvbZpF7TpS8amduO8+
|
||||
aRE6
|
||||
-----END CERTIFICATE REQUEST-----
|
||||
@ -1,28 +0,0 @@
|
||||
-----BEGIN PRIVATE KEY-----
|
||||
MIIEvgIBADANBgkqhkiG9w0BAQEFAASCBKgwggSkAgEAAoIBAQDdxxyYG0zV2Kze
|
||||
iVH5AGj1X1h7vLAHVS8FGkPL2QBBqW5/PrJ/z/sv3QJz6dB3ElOnk14GsY8lWk8u
|
||||
X/fjgNIPawN1G2/bCdMrlPOCEJEIOqZvzQWgPLRPGjoZx+wjCM/H7h47KVr4GNbo
|
||||
8MakLJeg6QB3rEEIiPrhQIoq9N7AhcjbEvWRNIxIQrpSqMCE2A5RoKUxAMum9rEo
|
||||
e/6PBw2GgEsQl5E6suRZw3wavg3aUU+6MMyxiiKEA6fCD5gDqgSo/xjgDUWJIOTT
|
||||
rmW6RBTcWP9iBUjWnCovQQ9zc6CzDXRgfsz97uXQt9fuQEi++lcGnrsePjC99Pni
|
||||
rCvmj/C/AgMBAAECggEAM4j6RwpI/4RbH1cvmjoTKbmfORmumfWceIYS7QKfAaMa
|
||||
jy0Fk5fD2ep0kHTrwU+b6tvexJVsGxTyQ2d/lfkwVu7aHdNjWbXdwUnakAXDffML
|
||||
C/3LaeHRUHRavfTsFXQNvHrDwaGphu9WuUiCEFJgZb7fIfAAKLiT+9XghXzjaj79
|
||||
nGANBB++kfBGacfj2NMedCxsbZ++pD7DbTG4utnHmzgwB3WEEm2kIbt0XpuDSZqe
|
||||
/tCA1gCOADjV17C31booT6xdoj+s/Se0YxY91JaZFkI5wTSDIhB9scZzh8piy8YQ
|
||||
QaEaRUUYNbzdI1MyvgewT6wQLONNlytOOc9hYnPgAQKBgQDyzaublY2q+yzqCQ68
|
||||
JSKNTCSFRwumJ1kmdfz/0QwzS363l7LFnfIb55h27q3vuiagPfRbi0GK/tVXhOe+
|
||||
8mntNMoJJQJzowpy4ANqKfGndojXejE1s1ExpuGyqWV4MyFGiwAnnzyWRENw4laS
|
||||
roY0OcynNuq++isgrVMPlDLAvwKBgQDp1OTh2wcSfhsaVbMjNVCOGD82OFK057JB
|
||||
3C62Li2lCybFwcooypCyZthUdskY2Gf/mDdnpLSLvXr/7bO3IyLopA3WpnYqMBt+
|
||||
7sP/VQJMoyobLFMnJh+N+CxwkhRNAib64054WKwiJcccJ5fplUvNjOLBgGx3a/3N
|
||||
OIxx+L7QAQKBgQC+EK3zPuEFJVYFZk24jkE75oz4H6NIz6iD6PzBrH0mckShpwh0
|
||||
la1+lo7NGw3hiRDPg3ATcTE/gziyKAHZgZ3V5+r3uZbvuoNlZWKG6oqWkr2QH8EB
|
||||
znsSqRYoa15Y931m4Uyft5EP+CPq6+LlM+UuYMiJZP3hvaehDszkdxg7tQKBgGjv
|
||||
vMPLCpJ2+2zHFvxu+ba7FOsdPain7ix2RpRFhwBxT7Yh8Lp7pZIaa20EXd0DiTCA
|
||||
PCUMGmY345IlN6ixYQIsVXWGALOQIVbGijj1CnIgK05EhxCjoDeTL0ZZmDizBZFE
|
||||
HzwM9zrf30o4TolqEbmuRzj1jDfPw/esMAMZ0XABAoGBANpjvUN0Qvy9WmLyPe6U
|
||||
DBzIAUyyIL3mnCuHLCQ+dE6nBf5l20sVjTdIl/eLlvmD7dwT/D9Zsg32/c2NBCdY
|
||||
DltKKZl2u9nuSJ1g8DWTZUQHkYSTeKPOoAl/5mbCgyTxzmmmtZF6MySKuif0rzwS
|
||||
X2Rwou9zTpwgsEEmQ/RH4G94
|
||||
-----END PRIVATE KEY-----
|
||||
@ -2,52 +2,7 @@
|
||||
%%% @author anlicheng
|
||||
%%% @copyright (C) 2025, <COMPANY>
|
||||
%%% @doc
|
||||
%%% Endpoint 数据分发订阅索引。
|
||||
%%%
|
||||
%%% 当前实现面向较大订阅量场景,publish 热路径不再全量遍历订阅表,而是把
|
||||
%%% 订阅拆成 exact 和 wildcard 两类:
|
||||
%%%
|
||||
%%% <ul>
|
||||
%%% <li>精确订阅写入 endpoint_subscription_exact,publish 时按 RouteKey
|
||||
%%% 直接 ets:lookup/2。</li>
|
||||
%%% <li>包含通配符的订阅写入 ETS trie。trie 边存放在
|
||||
%%% endpoint_subscription_trie_edge,节点订阅存放在
|
||||
%%% endpoint_subscription_trie_sub。</li>
|
||||
%%% <li>endpoint_subscription_reverse 保存 SubscriberPid 到订阅对象的反向索引,
|
||||
%%% 用于 unsubscribe 和进程 DOWN 清理。</li>
|
||||
%%% <li>endpoint_subscription_pid 记录每个 SubscriberPid 的 monitor 和订阅计数,
|
||||
%%% 同一个 pid 只 monitor 一次。</li>
|
||||
%%% </ul>
|
||||
%%%
|
||||
%%% 所有 matcher 共用同一个 root node,root node id 固定为 0。新 trie 节点使用
|
||||
%%% 单调递增的 integer node id。subscribe/unsubscribe 通过本 gen_server 串行写 ETS;
|
||||
%%% publish/2 直接并发读取 ETS,避免中心 gen_server 成为数据分发瓶颈。
|
||||
%%%
|
||||
%%% 匹配规则:
|
||||
%%%
|
||||
%%% <ul>
|
||||
%%% <li>普通 segment 必须完全相等,例如 <<"device/a/temp">> 只精确匹配
|
||||
%%% 同名订阅。</li>
|
||||
%%% <li><<"*">> 是单级通配,只匹配一个 segment。例如
|
||||
%%% <<"device/*/temp">> 匹配 <<"device/a/temp">>,不匹配
|
||||
%%% <<"device/a/b/temp">>。</li>
|
||||
%%% <li><<"+">> 是末尾多级通配,只允许出现在 matcher 最后一段,并且至少
|
||||
%%% 匹配一个剩余 segment。例如 <<"device/+">> 匹配 <<"device/a">>
|
||||
%%% 和 <<"device/a/temp">>,不匹配 <<"device">>。</li>
|
||||
%%% </ul>
|
||||
%%%
|
||||
%%% publish/2 的匹配流程:
|
||||
%%%
|
||||
%%% <ol>
|
||||
%%% <li>从 exact 表直接查找 RouteKey 对应的订阅。</li>
|
||||
%%% <li>把 RouteKey 按 "/" 拆分后,从 trie root 开始逐层匹配。每层最多查找
|
||||
%%% 当前 segment 分支和 <<"*">> 分支。</li>
|
||||
%%% <li>每消费一层前,收集当前候选节点上的 plus 订阅;全部 segment 消费完后,
|
||||
%%% 收集候选节点上的 exact 订阅。</li>
|
||||
%%% <li>合并 exact 和 trie 命中结果,按原有优先级排序,再按 SubscriberPid 去重。
|
||||
%%% 同一个 pid 同时被多个 matcher 命中时,保留优先级最高的订阅。</li>
|
||||
%%% <li>对最终订阅者执行 endpoint:forward/2。</li>
|
||||
%%% </ol>
|
||||
%%% @end
|
||||
%%% Created : 07. 11月 2025 16:27
|
||||
%%%-------------------------------------------------------------------
|
||||
@ -59,18 +14,13 @@
|
||||
%% API
|
||||
-export([start_link/0]).
|
||||
-export([subscribe/2, unsubscribe/2, publish/2, get_subscribers/0]).
|
||||
-export([is_valid_components/1, of_components/1]).
|
||||
-export([match_components/2, is_valid_components/1, of_components/1]).
|
||||
|
||||
%% gen_server callbacks
|
||||
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]).
|
||||
|
||||
-define(SERVER, ?MODULE).
|
||||
-define(EXACT_TAB, endpoint_subscription_exact).
|
||||
-define(TRIE_EDGE_TAB, endpoint_subscription_trie_edge).
|
||||
-define(TRIE_SUB_TAB, endpoint_subscription_trie_sub).
|
||||
-define(REVERSE_TAB, endpoint_subscription_reverse).
|
||||
-define(PID_TAB, endpoint_subscription_pid).
|
||||
-define(ROOT_NODE, 0).
|
||||
-define(SUBSCRIBER_TAB, endpoint_subscription_subscribers).
|
||||
|
||||
%% 定义订阅者
|
||||
-record(subscriber, {
|
||||
@ -86,12 +36,7 @@
|
||||
}).
|
||||
|
||||
-record(state, {
|
||||
exact_tid :: ets:tid(),
|
||||
edge_tid :: ets:tid(),
|
||||
trie_sub_tid :: ets:tid(),
|
||||
reverse_tid :: ets:tid(),
|
||||
pid_tid :: ets:tid(),
|
||||
next_node_id = 1 :: pos_integer()
|
||||
tid :: ets:tid()
|
||||
}).
|
||||
|
||||
%%%===================================================================
|
||||
@ -106,22 +51,23 @@ subscribe(Topic, SubscriberPid) when is_binary(Topic), is_pid(SubscriberPid) ->
|
||||
unsubscribe(Topic, SubscriberPid) when is_binary(Topic), is_pid(SubscriberPid) ->
|
||||
gen_server:call(?SERVER, {unsubscribe, Topic, SubscriberPid}).
|
||||
|
||||
-spec get_subscribers() -> {ok, Subscribers :: list()}.
|
||||
-spec get_subscribers() -> {ok, Subscribers :: map()}.
|
||||
get_subscribers() ->
|
||||
gen_server:call(?SERVER, get_subscribers).
|
||||
|
||||
-spec publish(RouteKey :: binary(), Content :: binary()) -> ok.
|
||||
publish(RouteKey, Content) when is_binary(RouteKey), is_binary(Content) ->
|
||||
case ets:info(?EXACT_TAB) of
|
||||
case ets:info(?SUBSCRIBER_TAB) of
|
||||
undefined ->
|
||||
ok;
|
||||
_ ->
|
||||
MatchedSubscribers = match_route_key(RouteKey),
|
||||
Subscribers = ets:tab2list(?SUBSCRIBER_TAB),
|
||||
MatchedSubscribers = match_subscribers(Subscribers, RouteKey),
|
||||
lists:foreach(fun(#subscriber{subscriber_pid = SubscriberPid}) ->
|
||||
endpoint:forward(SubscriberPid, Content)
|
||||
end, MatchedSubscribers),
|
||||
maybe_log_unmatched_publish(RouteKey, Content, MatchedSubscribers),
|
||||
logger:debug("[endpoint_subscription] route_key: ~p, match_count: ~p", [RouteKey, length(MatchedSubscribers)]),
|
||||
logger:debug("[efka_subscription] route_key: ~p, metric: ~p, match subscribers: ~p", [RouteKey, Content, MatchedSubscribers]),
|
||||
ok
|
||||
end.
|
||||
|
||||
@ -142,18 +88,8 @@ start_link() ->
|
||||
{stop, Reason :: term()} | ignore).
|
||||
init([]) ->
|
||||
ok = iot_log:set_metadata(),
|
||||
ExactTid = ets:new(?EXACT_TAB, [named_table, protected, bag, {read_concurrency, true}]),
|
||||
EdgeTid = ets:new(?TRIE_EDGE_TAB, [named_table, protected, set, {read_concurrency, true}]),
|
||||
TrieSubTid = ets:new(?TRIE_SUB_TAB, [named_table, protected, bag, {read_concurrency, true}]),
|
||||
ReverseTid = ets:new(?REVERSE_TAB, [named_table, protected, bag]),
|
||||
PidTid = ets:new(?PID_TAB, [named_table, protected, set]),
|
||||
{ok, #state{
|
||||
exact_tid = ExactTid,
|
||||
edge_tid = EdgeTid,
|
||||
trie_sub_tid = TrieSubTid,
|
||||
reverse_tid = ReverseTid,
|
||||
pid_tid = PidTid
|
||||
}}.
|
||||
Tid = ets:new(?SUBSCRIBER_TAB, [named_table, protected, bag, {keypos, 2}]),
|
||||
{ok, #state{tid = Tid}}.
|
||||
|
||||
%% @private
|
||||
%% @doc Handling call messages
|
||||
@ -166,19 +102,20 @@ init([]) ->
|
||||
{stop, Reason :: term(), Reply :: term(), NewState :: #state{}} |
|
||||
{stop, Reason :: term(), NewState :: #state{}}).
|
||||
%% 同一个SubscriberPid只能订阅同一个topic一次
|
||||
handle_call(get_subscribers, _From, State = #state{exact_tid = ExactTid, trie_sub_tid = TrieSubTid}) ->
|
||||
Subscribers = exact_subscribers(ExactTid) ++ trie_subscribers(TrieSubTid),
|
||||
handle_call(get_subscribers, _From, State = #state{tid = Tid}) ->
|
||||
Subscribers = ets:tab2list(Tid),
|
||||
{reply, {ok, Subscribers}, State};
|
||||
|
||||
handle_call({subscribe, Topic, SubscriberPid}, _From, State = #state{reverse_tid = ReverseTid}) ->
|
||||
handle_call({subscribe, Topic, SubscriberPid}, _From, State = #state{tid = Tid}) ->
|
||||
Components = of_components(Topic),
|
||||
case is_valid_components(Components) of
|
||||
true ->
|
||||
case has_subscription(ReverseTid, Topic, SubscriberPid) of
|
||||
case has_subscription(Tid, Topic, SubscriberPid) of
|
||||
true ->
|
||||
{reply, ok, State};
|
||||
false ->
|
||||
{MonitorRef, State1} = ensure_pid_monitor(SubscriberPid, State),
|
||||
%% 建立到SubscriberPid的monitor,进程退出需要清理订阅
|
||||
MonitorRef = erlang:monitor(process, SubscriberPid),
|
||||
Sub = #subscriber{
|
||||
topic = Topic,
|
||||
subscriber_pid = SubscriberPid,
|
||||
@ -186,19 +123,23 @@ handle_call({subscribe, Topic, SubscriberPid}, _From, State = #state{reverse_tid
|
||||
monitor_ref = MonitorRef,
|
||||
order = order_num(Components)
|
||||
},
|
||||
State2 = insert_subscription(Topic, Components, SubscriberPid, Sub, State1),
|
||||
{reply, ok, State2}
|
||||
true = ets:insert(Tid, Sub),
|
||||
{reply, ok, State}
|
||||
end;
|
||||
false ->
|
||||
{reply, {error, <<"invalid topic name">>}, State}
|
||||
end;
|
||||
|
||||
handle_call({unsubscribe, Topic, SubscriberPid}, _From, State = #state{reverse_tid = ReverseTid}) ->
|
||||
Removed = [Reverse || Reverse = {SubscriberPid0, Topic0, _Kind, _NodeId, _Sub} <- ets:lookup(ReverseTid, SubscriberPid),
|
||||
SubscriberPid =:= SubscriberPid0, Topic =:= Topic0],
|
||||
lists:foreach(fun(Reverse) -> delete_subscription(Reverse, State) end, Removed),
|
||||
State1 = release_pid_monitor(SubscriberPid, length(Removed), State),
|
||||
{reply, ok, State1}.
|
||||
handle_call({unsubscribe, Topic, SubscriberPid}, _From, State = #state{tid = Tid}) ->
|
||||
Removed = [Sub || Sub = #subscriber{subscriber_pid = SubscriberPid0} <- ets:lookup(Tid, Topic),
|
||||
SubscriberPid =:= SubscriberPid0],
|
||||
lists:foreach(fun(#subscriber{monitor_ref = MonitorRef}) when is_reference(MonitorRef) ->
|
||||
erlang:demonitor(MonitorRef, [flush]);
|
||||
(_) ->
|
||||
ok
|
||||
end, Removed),
|
||||
lists:foreach(fun(Sub) -> true = ets:delete_object(Tid, Sub) end, Removed),
|
||||
{reply, ok, State}.
|
||||
|
||||
%% @private
|
||||
%% @doc Handling cast messages
|
||||
@ -215,20 +156,18 @@ handle_cast(_Request, State = #state{}) ->
|
||||
{noreply, NewState :: #state{}} |
|
||||
{noreply, NewState :: #state{}, timeout() | hibernate} |
|
||||
{stop, Reason :: term(), NewState :: #state{}}).
|
||||
handle_info({'DOWN', MonitorRef, process, SubscriberPid, Reason}, State = #state{reverse_tid = ReverseTid, pid_tid = PidTid}) ->
|
||||
logger:debug("[endpoint_subscription] subscriber: ~p, down with reason: ~p", [SubscriberPid, Reason]),
|
||||
Removed = ets:lookup(ReverseTid, SubscriberPid),
|
||||
lists:foreach(fun(Reverse) -> delete_subscription(Reverse, State) end, Removed),
|
||||
case ets:lookup(PidTid, SubscriberPid) of
|
||||
[{SubscriberPid, MonitorRef, _Count}] ->
|
||||
ets:delete(PidTid, SubscriberPid);
|
||||
_ ->
|
||||
ok
|
||||
end,
|
||||
handle_info({'DOWN', _Ref, process, SubscriberPid, Reason}, State = #state{tid = Tid}) ->
|
||||
logger:debug("[efka_subscription] subscriber: ~p, down with reason: ~p", [SubscriberPid, Reason]),
|
||||
Subscribers = ets:tab2list(Tid),
|
||||
lists:foreach(fun(Sub = #subscriber{subscriber_pid = Pid0}) when Pid0 =:= SubscriberPid ->
|
||||
true = ets:delete_object(Tid, Sub);
|
||||
(_) ->
|
||||
ok
|
||||
end, Subscribers),
|
||||
{noreply, State};
|
||||
|
||||
handle_info(Info, State = #state{}) ->
|
||||
logger:debug("[endpoint_subscription] get unknown info: ~p", [Info]),
|
||||
logger:debug("[efka_subscription] get unknown info: ~p", [Info]),
|
||||
{noreply, State}.
|
||||
|
||||
%% @private
|
||||
@ -253,11 +192,14 @@ code_change(_OldVsn, State = #state{}, _Extra) ->
|
||||
%%% Internal functions
|
||||
%%%===================================================================
|
||||
|
||||
-spec match_route_key(binary()) -> [#subscriber{}].
|
||||
match_route_key(RouteKey) when is_binary(RouteKey) ->
|
||||
ExactSubs = [Sub || {_Topic, Sub} <- ets:lookup(?EXACT_TAB, RouteKey)],
|
||||
TrieSubs = match_trie(of_components(RouteKey)),
|
||||
Sorted = lists:sort(fun compare_subscriber/2, ExactSubs ++ TrieSubs),
|
||||
%% 查找满足条件订阅者
|
||||
-spec match_subscribers(Subscribers :: [#subscriber{}], Topic :: binary()) -> [#subscriber{}].
|
||||
match_subscribers(Subscribers, Topic) when is_list(Subscribers), is_binary(Topic) ->
|
||||
Components = of_components(Topic),
|
||||
Matched = lists:filter(fun(#subscriber{components = Components0}) ->
|
||||
match_components(Components0, Components)
|
||||
end, Subscribers),
|
||||
Sorted = lists:sort(fun compare_subscriber/2, Matched),
|
||||
dedupe_subscribers(Sorted).
|
||||
|
||||
-spec maybe_log_unmatched_publish(binary(), binary(), [#subscriber{}]) -> ok.
|
||||
@ -266,6 +208,22 @@ maybe_log_unmatched_publish(RouteKey, Content, []) ->
|
||||
maybe_log_unmatched_publish(_RouteKey, _Content, _MatchedSubscribers) ->
|
||||
ok.
|
||||
|
||||
%% 开始对比订阅的topic和发布的topic的Components信息
|
||||
%% *表示单级匹配,+表示多级匹配;+只能出现一次,并且只能在末尾
|
||||
-spec match_components(list(), list()) -> boolean().
|
||||
match_components(A, B) when is_list(A), is_list(B) ->
|
||||
match_components(A, B, false).
|
||||
match_components([<<"+">>], [_|_], _) ->
|
||||
true;
|
||||
match_components([], [], _) ->
|
||||
true;
|
||||
match_components([<<"*">>|T0], [_|T1], _) ->
|
||||
match_components(T0, T1, false);
|
||||
match_components([C0|T0], [C0|T1], _) ->
|
||||
match_components(T0, T1, false);
|
||||
match_components(_, _, _) ->
|
||||
false.
|
||||
|
||||
-spec of_components(Topic :: binary()) -> [binary()].
|
||||
of_components(Topic) when is_binary(Topic) ->
|
||||
binary:split(Topic, <<$/>>, [global]).
|
||||
@ -290,11 +248,11 @@ order_num([_|Tail]) ->
|
||||
order_num(Tail).
|
||||
|
||||
-spec has_subscription(ets:tid(), binary(), pid()) -> boolean().
|
||||
has_subscription(ReverseTid, Topic, SubscriberPid) ->
|
||||
Subscriptions = ets:lookup(ReverseTid, SubscriberPid),
|
||||
lists:any(fun({SubscriberPid0, Topic0, _Kind, _NodeId, _Sub}) ->
|
||||
has_subscription(Tid, Topic, SubscriberPid) ->
|
||||
Subscribers = ets:lookup(Tid, Topic),
|
||||
lists:any(fun(#subscriber{topic = Topic0, subscriber_pid = SubscriberPid0}) ->
|
||||
Topic =:= Topic0 andalso SubscriberPid =:= SubscriberPid0
|
||||
end, Subscriptions).
|
||||
end, Subscribers).
|
||||
|
||||
-spec compare_subscriber(#subscriber{}, #subscriber{}) -> boolean().
|
||||
compare_subscriber(#subscriber{order = Order0, topic = Topic0}, #subscriber{order = Order1, topic = Topic1}) ->
|
||||
@ -316,142 +274,3 @@ dedupe_subscribers(Subscribers) ->
|
||||
end
|
||||
end, {sets:new(), []}, Subscribers),
|
||||
lists:reverse(Result).
|
||||
|
||||
-spec insert_subscription(binary(), [binary()], pid(), #subscriber{}, #state{}) -> #state{}.
|
||||
insert_subscription(Topic, Components, SubscriberPid, Sub, State = #state{exact_tid = ExactTid, trie_sub_tid = TrieSubTid, reverse_tid = ReverseTid}) ->
|
||||
case has_wildcard(Components) of
|
||||
false ->
|
||||
true = ets:insert(ExactTid, {Topic, Sub}),
|
||||
true = ets:insert(ReverseTid, {SubscriberPid, Topic, exact, undefined, Sub}),
|
||||
State;
|
||||
true ->
|
||||
{Kind, NodeId, State1} = insert_trie_subscription(Components, State),
|
||||
true = ets:insert(TrieSubTid, {{NodeId, subscription_match_type(Kind)}, Sub}),
|
||||
true = ets:insert(ReverseTid, {SubscriberPid, Topic, Kind, NodeId, Sub}),
|
||||
State1
|
||||
end.
|
||||
|
||||
-spec insert_trie_subscription([binary()], #state{}) -> {trie_exact | trie_plus, non_neg_integer(), #state{}}.
|
||||
insert_trie_subscription(Components, State) ->
|
||||
case Components of
|
||||
[] ->
|
||||
{trie_exact, ?ROOT_NODE, State};
|
||||
_ ->
|
||||
case lists:last(Components) of
|
||||
<<"+">> ->
|
||||
Prefix = lists:sublist(Components, length(Components) - 1),
|
||||
{NodeId, State1} = ensure_path(Prefix, State),
|
||||
{trie_plus, NodeId, State1};
|
||||
_ ->
|
||||
{NodeId, State1} = ensure_path(Components, State),
|
||||
{trie_exact, NodeId, State1}
|
||||
end
|
||||
end.
|
||||
|
||||
-spec subscription_match_type(trie_exact | trie_plus) -> exact | plus.
|
||||
subscription_match_type(trie_exact) ->
|
||||
exact;
|
||||
subscription_match_type(trie_plus) ->
|
||||
plus.
|
||||
|
||||
-spec ensure_path([binary()], #state{}) -> {non_neg_integer(), #state{}}.
|
||||
ensure_path(Components, State) ->
|
||||
ensure_path(?ROOT_NODE, Components, State).
|
||||
|
||||
-spec ensure_path(non_neg_integer(), [binary()], #state{}) -> {non_neg_integer(), #state{}}.
|
||||
ensure_path(NodeId, [], State) ->
|
||||
{NodeId, State};
|
||||
ensure_path(NodeId, [Segment | Rest], State = #state{edge_tid = EdgeTid, next_node_id = NextNodeId}) ->
|
||||
EdgeKey = {NodeId, Segment},
|
||||
case ets:lookup(EdgeTid, EdgeKey) of
|
||||
[{EdgeKey, ChildNodeId}] ->
|
||||
ensure_path(ChildNodeId, Rest, State);
|
||||
[] ->
|
||||
true = ets:insert(EdgeTid, {EdgeKey, NextNodeId}),
|
||||
ensure_path(NextNodeId, Rest, State#state{next_node_id = NextNodeId + 1})
|
||||
end.
|
||||
|
||||
-spec delete_subscription(tuple(), #state{}) -> ok.
|
||||
delete_subscription(Reverse = {_SubscriberPid, Topic, exact, undefined, Sub}, #state{exact_tid = ExactTid, reverse_tid = ReverseTid}) ->
|
||||
true = ets:delete_object(ExactTid, {Topic, Sub}),
|
||||
true = ets:delete_object(ReverseTid, Reverse),
|
||||
ok;
|
||||
delete_subscription(Reverse = {_SubscriberPid, _Topic, Kind, NodeId, Sub}, #state{trie_sub_tid = TrieSubTid, reverse_tid = ReverseTid}) ->
|
||||
true = ets:delete_object(TrieSubTid, {{NodeId, subscription_match_type(Kind)}, Sub}),
|
||||
true = ets:delete_object(ReverseTid, Reverse),
|
||||
ok.
|
||||
|
||||
-spec ensure_pid_monitor(pid(), #state{}) -> {reference(), #state{}}.
|
||||
ensure_pid_monitor(SubscriberPid, State = #state{pid_tid = PidTid}) ->
|
||||
case ets:lookup(PidTid, SubscriberPid) of
|
||||
[{SubscriberPid, MonitorRef, Count}] ->
|
||||
true = ets:insert(PidTid, {SubscriberPid, MonitorRef, Count + 1}),
|
||||
{MonitorRef, State};
|
||||
[] ->
|
||||
MonitorRef = erlang:monitor(process, SubscriberPid),
|
||||
true = ets:insert(PidTid, {SubscriberPid, MonitorRef, 1}),
|
||||
{MonitorRef, State}
|
||||
end.
|
||||
|
||||
-spec release_pid_monitor(pid(), non_neg_integer(), #state{}) -> #state{}.
|
||||
release_pid_monitor(_SubscriberPid, 0, State) ->
|
||||
State;
|
||||
release_pid_monitor(SubscriberPid, RemovedCount, State = #state{pid_tid = PidTid}) ->
|
||||
case ets:lookup(PidTid, SubscriberPid) of
|
||||
[{SubscriberPid, MonitorRef, Count}] when Count =< RemovedCount ->
|
||||
erlang:demonitor(MonitorRef, [flush]),
|
||||
ets:delete(PidTid, SubscriberPid),
|
||||
State;
|
||||
[{SubscriberPid, MonitorRef, Count}] ->
|
||||
true = ets:insert(PidTid, {SubscriberPid, MonitorRef, Count - RemovedCount}),
|
||||
State;
|
||||
[] ->
|
||||
State
|
||||
end.
|
||||
|
||||
-spec match_trie([binary()]) -> [#subscriber{}].
|
||||
match_trie(Components) ->
|
||||
match_trie([?ROOT_NODE], Components, []).
|
||||
|
||||
-spec match_trie([non_neg_integer()], [binary()], [#subscriber{}]) -> [#subscriber{}].
|
||||
match_trie(NodeIds, [], Acc) ->
|
||||
collect_trie_subscribers(NodeIds, exact) ++ Acc;
|
||||
match_trie(NodeIds, [Segment | Rest], Acc) ->
|
||||
PlusSubs = collect_trie_subscribers(NodeIds, plus),
|
||||
NextNodeIds = next_trie_nodes(NodeIds, Segment),
|
||||
match_trie(NextNodeIds, Rest, PlusSubs ++ Acc).
|
||||
|
||||
-spec collect_trie_subscribers([non_neg_integer()], exact | plus) -> [#subscriber{}].
|
||||
collect_trie_subscribers(NodeIds, MatchType) ->
|
||||
lists:flatmap(fun(NodeId) ->
|
||||
[Sub || {{NodeId0, MatchType0}, Sub} <- ets:lookup(?TRIE_SUB_TAB, {NodeId, MatchType}),
|
||||
NodeId0 =:= NodeId, MatchType0 =:= MatchType]
|
||||
end, NodeIds).
|
||||
|
||||
-spec next_trie_nodes([non_neg_integer()], binary()) -> [non_neg_integer()].
|
||||
next_trie_nodes(NodeIds, Segment) ->
|
||||
lists:usort(lists:flatmap(fun(NodeId) ->
|
||||
lookup_child(NodeId, Segment) ++ lookup_child(NodeId, <<"*">>)
|
||||
end, NodeIds)).
|
||||
|
||||
-spec lookup_child(non_neg_integer(), binary()) -> [non_neg_integer()].
|
||||
lookup_child(NodeId, Segment) ->
|
||||
EdgeKey = {NodeId, Segment},
|
||||
case ets:lookup(?TRIE_EDGE_TAB, EdgeKey) of
|
||||
[{EdgeKey, ChildNodeId}] ->
|
||||
[ChildNodeId];
|
||||
[] ->
|
||||
[]
|
||||
end.
|
||||
|
||||
-spec has_wildcard([binary()]) -> boolean().
|
||||
has_wildcard(Components) ->
|
||||
lists:any(fun(Component) -> Component =:= <<"*">> orelse Component =:= <<"+">> end, Components).
|
||||
|
||||
-spec exact_subscribers(ets:tid()) -> [#subscriber{}].
|
||||
exact_subscribers(ExactTid) ->
|
||||
[Sub || {_Topic, Sub} <- ets:tab2list(ExactTid)].
|
||||
|
||||
-spec trie_subscribers(ets:tid()) -> [#subscriber{}].
|
||||
trie_subscribers(TrieSubTid) ->
|
||||
[Sub || {_Key, Sub} <- ets:tab2list(TrieSubTid)].
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user