Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
张重言
rails
提交
16a66039
R
rails
项目概览
张重言
/
rails
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
R
rails
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
体验新版 GitCode,发现更多精彩内容 >>
提交
16a66039
编写于
1月 25, 2016
作者:
M
Matthew Draper
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Synchronize the lazy setters in Server
They're all at risk of races on the first requests.
上级
a928aa3d
变更
5
隐藏空白更改
内联
并排
Showing
5 changed file
with
47 addition
and
15 deletion
+47
-15
actioncable/lib/action_cable/server/base.rb
actioncable/lib/action_cable/server/base.rb
+16
-7
actioncable/lib/action_cable/subscription_adapter/async.rb
actioncable/lib/action_cable/subscription_adapter/async.rb
+2
-2
actioncable/lib/action_cable/subscription_adapter/inline.rb
actioncable/lib/action_cable/subscription_adapter/inline.rb
+10
-1
actioncable/lib/action_cable/subscription_adapter/postgresql.rb
...cable/lib/action_cable/subscription_adapter/postgresql.rb
+6
-1
actioncable/lib/action_cable/subscription_adapter/redis.rb
actioncable/lib/action_cable/subscription_adapter/redis.rb
+13
-4
未找到文件。
actioncable/lib/action_cable/server/base.rb
浏览文件 @
16a66039
require
'thread'
module
ActionCable
module
Server
# A singleton ActionCable::Server instance is available via ActionCable.server. It's used by the rack process that starts the cable server, but
...
...
@@ -13,7 +15,12 @@ class Base
def
self
.
logger
;
config
.
logger
;
end
delegate
:logger
,
to: :config
attr_reader
:mutex
def
initialize
@mutex
=
Mutex
.
new
@remote_connections
=
@stream_event_loop
=
@worker_pool
=
@channel_classes
=
@pubsub
=
nil
end
# Called by rack to setup the server.
...
...
@@ -29,29 +36,31 @@ def disconnect(identifiers)
# Gateway to RemoteConnections. See that class for details.
def
remote_connections
@remote_connections
||
=
RemoteConnections
.
new
(
self
)
@remote_connections
||
@mutex
.
synchronize
{
@remote_connections
||=
RemoteConnections
.
new
(
self
)
}
end
def
stream_event_loop
@stream_event_loop
||
=
ActionCable
::
Connection
::
StreamEventLoop
.
new
@stream_event_loop
||
@mutex
.
synchronize
{
@stream_event_loop
||=
ActionCable
::
Connection
::
StreamEventLoop
.
new
}
end
# The thread worker pool for handling all the connection work on this server. Default size is set by config.worker_pool_size.
def
worker_pool
@worker_pool
||
=
ActionCable
::
Server
::
Worker
.
new
(
max_size:
config
.
worker_pool_size
)
@worker_pool
||
@mutex
.
synchronize
{
@worker_pool
||=
ActionCable
::
Server
::
Worker
.
new
(
max_size:
config
.
worker_pool_size
)
}
end
# Requires and returns a hash of all the channel class constants keyed by name.
def
channel_classes
@channel_classes
||=
begin
config
.
channel_paths
.
each
{
|
channel_path
|
require
channel_path
}
config
.
channel_class_names
.
each_with_object
({})
{
|
name
,
hash
|
hash
[
name
]
=
name
.
constantize
}
@channel_classes
||
@mutex
.
synchronize
do
@channel_classes
||=
begin
config
.
channel_paths
.
each
{
|
channel_path
|
require
channel_path
}
config
.
channel_class_names
.
each_with_object
({})
{
|
name
,
hash
|
hash
[
name
]
=
name
.
constantize
}
end
end
end
# Adapter used for all streams/broadcasting.
def
pubsub
@pubsub
||
=
config
.
pubsub_adapter
.
new
(
self
)
@pubsub
||
@mutex
.
synchronize
{
@pubsub
||=
config
.
pubsub_adapter
.
new
(
self
)
}
end
# All the identifiers applied to the connection class associated with this server.
...
...
actioncable/lib/action_cable/subscription_adapter/async.rb
浏览文件 @
16a66039
...
...
@@ -4,8 +4,8 @@ module ActionCable
module
SubscriptionAdapter
class
Async
<
Inline
# :nodoc:
private
def
subscriber_map
@subscriber_map
||=
AsyncSubscriberMap
.
new
def
new_
subscriber_map
AsyncSubscriberMap
.
new
end
class
AsyncSubscriberMap
<
SubscriberMap
...
...
actioncable/lib/action_cable/subscription_adapter/inline.rb
浏览文件 @
16a66039
module
ActionCable
module
SubscriptionAdapter
class
Inline
<
Base
# :nodoc:
def
initialize
(
*
)
super
@subscriber_map
=
nil
end
def
broadcast
(
channel
,
payload
)
subscriber_map
.
broadcast
(
channel
,
payload
)
end
...
...
@@ -19,7 +24,11 @@ def shutdown
private
def
subscriber_map
@subscriber_map
||=
SubscriberMap
.
new
@subscriber_map
||
@server
.
mutex
.
synchronize
{
@subscriber_map
||=
new_subscriber_map
}
end
def
new_subscriber_map
SubscriberMap
.
new
end
end
end
...
...
actioncable/lib/action_cable/subscription_adapter/postgresql.rb
浏览文件 @
16a66039
...
...
@@ -5,6 +5,11 @@
module
ActionCable
module
SubscriptionAdapter
class
PostgreSQL
<
Base
# :nodoc:
def
initialize
(
*
)
super
@listener
=
nil
end
def
broadcast
(
channel
,
payload
)
with_connection
do
|
pg_conn
|
pg_conn
.
exec
(
"NOTIFY
#{
pg_conn
.
escape_identifier
(
channel
)
}
, '
#{
pg_conn
.
escape_string
(
payload
)
}
'"
)
...
...
@@ -37,7 +42,7 @@ def with_connection(&block) # :nodoc:
private
def
listener
@listener
||
=
Listener
.
new
(
self
)
@listener
||
@server
.
mutex
.
synchronize
{
@listener
||=
Listener
.
new
(
self
)
}
end
class
Listener
<
SubscriberMap
...
...
actioncable/lib/action_cable/subscription_adapter/redis.rb
浏览文件 @
16a66039
...
...
@@ -13,6 +13,11 @@ module SubscriptionAdapter
class
Redis
<
Base
# :nodoc:
@@mutex
=
Mutex
.
new
def
initialize
(
*
)
super
@redis_connection_for_broadcasts
=
@redis_connection_for_subscriptions
=
nil
end
def
broadcast
(
channel
,
payload
)
redis_connection_for_broadcasts
.
publish
(
channel
,
payload
)
end
...
...
@@ -35,15 +40,19 @@ def shutdown
private
def
redis_connection_for_subscriptions
ensure_reactor_running
@redis_connection_for_subscriptions
||=
EM
::
Hiredis
.
connect
(
@server
.
config
.
cable
[
:url
]).
tap
do
|
redis
|
redis
.
on
(
:reconnect_failed
)
do
@logger
.
info
"[ActionCable] Redis reconnect failed."
@redis_connection_for_subscriptions
||
@server
.
mutex
.
synchronize
do
@redis_connection_for_subscriptions
||=
EM
::
Hiredis
.
connect
(
@server
.
config
.
cable
[
:url
]).
tap
do
|
redis
|
redis
.
on
(
:reconnect_failed
)
do
@logger
.
info
"[ActionCable] Redis reconnect failed."
end
end
end
end
def
redis_connection_for_broadcasts
@redis_connection_for_broadcasts
||=
::
Redis
.
new
(
@server
.
config
.
cable
)
@redis_connection_for_broadcasts
||
@server
.
mutex
.
synchronize
do
@redis_connection_for_broadcasts
||=
::
Redis
.
new
(
@server
.
config
.
cable
)
end
end
def
ensure_reactor_running
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录