Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
爱吃血肠
spring-framework
提交
b7c924ca
S
spring-framework
项目概览
爱吃血肠
/
spring-framework
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
S
spring-framework
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
前往新版Gitcode,体验更适合开发者的 AI 搜索 >>
提交
b7c924ca
编写于
11月 21, 2017
作者:
R
Rossen Stoyanchev
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Polish
上级
d8099adc
变更
5
隐藏空白更改
内联
并排
Showing
5 changed file
with
27 addition
and
16 deletion
+27
-16
spring-web/src/main/java/org/springframework/http/server/reactive/AbstractListenerReadPublisher.java
...k/http/server/reactive/AbstractListenerReadPublisher.java
+2
-3
spring-web/src/main/java/org/springframework/http/server/reactive/ServletServerHttpRequest.java
...mework/http/server/reactive/ServletServerHttpRequest.java
+5
-0
spring-web/src/test/java/org/springframework/http/server/reactive/ListenerReadPublisherTests.java
...work/http/server/reactive/ListenerReadPublisherTests.java
+5
-0
spring-webflux/src/main/java/org/springframework/web/reactive/socket/adapter/AbstractListenerWebSocketSession.java
...tive/socket/adapter/AbstractListenerWebSocketSession.java
+11
-10
spring-webflux/src/main/java/org/springframework/web/reactive/socket/adapter/TomcatWebSocketSession.java
...k/web/reactive/socket/adapter/TomcatWebSocketSession.java
+4
-3
未找到文件。
spring-web/src/main/java/org/springframework/http/server/reactive/AbstractListenerReadPublisher.java
浏览文件 @
b7c924ca
...
...
@@ -116,10 +116,9 @@ public abstract class AbstractListenerReadPublisher<T> implements Publisher<T> {
protected
abstract
T
read
()
throws
IOException
;
/**
* Suspend reading
. Defaults to no-op
.
* Suspend reading
, if the underlying API provides such a mechanism
.
*/
protected
void
suspendReading
()
{
}
protected
abstract
void
suspendReading
();
// Private methods for use in State...
...
...
spring-web/src/main/java/org/springframework/http/server/reactive/ServletServerHttpRequest.java
浏览文件 @
b7c924ca
...
...
@@ -266,6 +266,11 @@ class ServletServerHttpRequest extends AbstractServerHttpRequest {
return
null
;
}
@Override
protected
void
suspendReading
()
{
// no-op
}
private
class
RequestBodyPublisherReadListener
implements
ReadListener
{
...
...
spring-web/src/test/java/org/springframework/http/server/reactive/ListenerReadPublisherTests.java
浏览文件 @
b7c924ca
...
...
@@ -67,6 +67,11 @@ public class ListenerReadPublisherTests {
return
mock
(
DataBuffer
.
class
);
}
@Override
protected
void
suspendReading
()
{
// No-op
}
public
int
getReadCalls
()
{
return
this
.
readCalls
;
}
...
...
spring-webflux/src/main/java/org/springframework/web/reactive/socket/adapter/AbstractListenerWebSocketSession.java
浏览文件 @
b7c924ca
...
...
@@ -150,13 +150,13 @@ public abstract class AbstractListenerWebSocketSession<T> extends AbstractWebSoc
protected
abstract
void
resumeReceiving
();
/**
* Returns {@code true} if receiving new message(s) is suspended otherwise
* {@code false}.
* Whether receiving new message(s) is suspended.
* <p><strong>Note:</strong> if the underlying WebSocket API does not provide
* flow control for receiving messages, and this method should return
* {@code false} and {@link #canSuspendReceiving()} should return {@code false}.
* @return returns {@code true} if receiving new message(s) is suspended
* otherwise {@code false}.
* flow control for receiving messages, then this method as well as
* {@link #canSuspendReceiving()} should both return {@code false}.
* @return returns {@code true} if receiving new message(s) is suspended,
* or otherwise {@code false}.
* @since 5.0.2
*/
protected
abstract
boolean
isSuspended
();
...
...
@@ -226,14 +226,15 @@ public abstract class AbstractListenerWebSocketSession<T> extends AbstractWebSoc
private
final
class
WebSocketReceivePublisher
extends
AbstractListenerReadPublisher
<
WebSocketMessage
>
{
private
volatile
Queue
<
Object
>
pendingWebSocketMessages
=
Queues
.
unbounded
().
get
();
private
volatile
Queue
<
Object
>
pendingMessages
=
Queues
.
unbounded
(
Queues
.
SMALL_BUFFER_SIZE
).
get
();
@Override
protected
void
checkOnDataAvailable
()
{
if
(
isSuspended
())
{
resumeReceiving
();
}
if
(!
pendingWebSocket
Messages
.
isEmpty
())
{
if
(!
this
.
pending
Messages
.
isEmpty
())
{
onDataAvailable
();
}
}
...
...
@@ -246,7 +247,7 @@ public abstract class AbstractListenerWebSocketSession<T> extends AbstractWebSoc
@Override
@Nullable
protected
WebSocketMessage
read
()
throws
IOException
{
return
(
WebSocketMessage
)
pendingWebSocket
Messages
.
poll
();
return
(
WebSocketMessage
)
this
.
pending
Messages
.
poll
();
}
@Override
...
...
@@ -258,7 +259,7 @@ public abstract class AbstractListenerWebSocketSession<T> extends AbstractWebSoc
}
void
handleMessage
(
WebSocketMessage
webSocketMessage
)
{
this
.
pending
WebSocket
Messages
.
offer
(
webSocketMessage
);
this
.
pendingMessages
.
offer
(
webSocketMessage
);
onDataAvailable
();
}
}
...
...
spring-webflux/src/main/java/org/springframework/web/reactive/socket/adapter/TomcatWebSocketSession.java
浏览文件 @
b7c924ca
...
...
@@ -17,16 +17,15 @@
package
org.springframework.web.reactive.socket.adapter
;
import
java.util.concurrent.atomic.AtomicIntegerFieldUpdater
;
import
javax.websocket.Session
;
import
org.apache.tomcat.websocket.WsSession
;
import
reactor.core.publisher.MonoProcessor
;
import
org.springframework.core.io.buffer.DataBufferFactory
;
import
org.springframework.web.reactive.socket.HandshakeInfo
;
import
org.springframework.web.reactive.socket.WebSocketSession
;
import
reactor.core.publisher.MonoProcessor
;
/**
* Spring {@link WebSocketSession} adapter for Tomcat's
* {@link javax.websocket.Session}.
...
...
@@ -35,8 +34,10 @@ import reactor.core.publisher.MonoProcessor;
* @since 5.0
*/
public
class
TomcatWebSocketSession
extends
StandardWebSocketSession
{
private
static
final
AtomicIntegerFieldUpdater
<
TomcatWebSocketSession
>
SUSPENDED
=
AtomicIntegerFieldUpdater
.
newUpdater
(
TomcatWebSocketSession
.
class
,
"suspended"
);
private
volatile
int
suspended
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录