Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Support
Keyboard shortcuts
?
Submit feedback
Contribute to GitLab
Sign in / Register
Toggle navigation
A
Actor Framework
Project overview
Project overview
Details
Activity
Releases
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Issues
0
Issues
0
List
Boards
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Operations
Operations
Metrics
Environments
Analytics
Analytics
CI / CD
Repository
Value Stream
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
cpp-libs
Actor Framework
Commits
8562c1af
Commit
8562c1af
authored
Apr 04, 2020
by
Dominik Charousset
Browse files
Options
Browse Files
Download
Plain Diff
Merge branch 'issue/1076'
parents
08f7a298
538ee173
Changes
2
Show whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
10 additions
and
3 deletions
+10
-3
libcaf_core/src/scheduled_actor.cpp
libcaf_core/src/scheduled_actor.cpp
+8
-2
libcaf_core/src/stream_manager.cpp
libcaf_core/src/stream_manager.cpp
+2
-1
No files found.
libcaf_core/src/scheduled_actor.cpp
View file @
8562c1af
...
...
@@ -1170,9 +1170,15 @@ scheduled_actor::advance_streams(actor_clock::time_point now) {
auto
bitmask
=
stream_ticks_
.
timeouts
(
now
,
{
max_batch_delay_ticks_
,
credit_round_ticks_
});
// Force batches on all output paths.
if
((
bitmask
&
0x01
)
!=
0
)
{
if
((
bitmask
&
0x01
)
!=
0
&&
!
stream_managers_
.
empty
())
{
std
::
vector
<
stream_manager
*>
managers
;
managers
.
reserve
(
stream_managers_
.
size
());
for
(
auto
&
kvp
:
stream_managers_
)
kvp
.
second
->
out
().
force_emit_batches
();
managers
.
emplace_back
(
kvp
.
second
.
get
());
std
::
sort
(
managers
.
begin
(),
managers
.
end
());
auto
e
=
std
::
unique
(
managers
.
begin
(),
managers
.
end
());
for
(
auto
i
=
managers
.
begin
();
i
!=
e
;
++
i
)
(
*
i
)
->
out
().
force_emit_batches
();
}
// Fill up credit on each input path.
if
((
bitmask
&
0x02
)
!=
0
)
{
...
...
libcaf_core/src/stream_manager.cpp
View file @
8562c1af
...
...
@@ -114,7 +114,8 @@ void stream_manager::handle(stream_slots slots, upstream_msg::ack_batch& x) {
void
stream_manager
::
handle
(
stream_slots
slots
,
upstream_msg
::
drop
&
)
{
CAF_LOG_TRACE
(
CAF_ARG
(
slots
));
out
().
close
(
slots
.
receiver
);
error
reason
;
out
().
remove_path
(
slots
.
receiver
,
reason
,
false
);
}
void
stream_manager
::
handle
(
stream_slots
slots
,
upstream_msg
::
forced_drop
&
x
)
{
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment