要看看实际发生了什么,让我们重写程序以(某种程度上)涵盖react 所做的事情。我将忽略许多与手头的问题无关紧要的细节。
为了使事情更紧凑,我将重写所提供程序的这一部分:
react {
whenever signal(SIGINT, |%scheduler) {
say "Got signal";
exit;
}
whenever Supply.from-list($*IN.lines, |%scheduler) {
say "Got line";
exit if $++ == 1 ;
}
}
首先,react { ... } 与await supply { ... } 非常相似——也就是说,它点击supply { } 块,awaits 结束。
await supply {
whenever signal(SIGINT, |%scheduler) {
say "Got signal";
exit;
}
whenever Supply.from-list($*IN.lines, |%scheduler) {
say "Got line";
exit if $++ == 1 ;
}
}
但是supply 块是什么?在其核心,supply(以及react)提供:
- 并发控制,因此它将一次处理一条消息(因此我们需要某种锁;它为此使用
Lock::Async)
- 完成和错误传播(在公然的作弊中,我们将使用
Promise 来实现这一点,因为我们只需要react 部分;真实的东西会产生Supply,我们可以emit 值入)
- 当没有未完成的订阅时自动完成(我们将在
SetHash 中跟踪这些订阅)
因此,我们可以将程序改写成这样:
await do {
# Concurency control
my $lock = Lock::Async.new;
# Completion/error conveyance
my $done = Promise.new;
# What's active?
my %active is SetHash;
# An implementation a bit like that behind the `whenever` keyword, but with
# plenty of things that don't matter for this question missing.
sub whenever-impl(Supply $s, &block) {
# Tap the Supply
my $tap;
$s.tap:
# When it gets tapped, add the tap to our active set.
tap => {
$tap = $_;
%active{$_} = True;
},
# Run the handler for any events
{ $lock.protect: { block($_) } },
# When this one is done, remove it from the %active list; if it's
# the last thing, we're done overall.
done => {
$lock.protect: {
%active{$tap}:delete;
$done.keep() unless %active;
}
},
# If there's an async error, close all taps and pass it along.
quit => {
$lock.protect: -> $err {
.close for %active.keys;
$done.quit($err);
}
}
}
# We hold the lock while doing initial setup, so you can rely on having
# done all initialization before processing a first message.
$lock.protect: {
whenever-impl signal(SIGINT, |%scheduler), {
say "Got signal";
exit;
}
whenever-impl Supply.from-list($*IN.lines, |%scheduler), {
say "Got line";
exit if $++ == 1 ;
}
}
$done
}
请注意,这里没有任何关于调度程序或事件循环的内容; supply 或 react 不关心消息来自谁,它只关心自己的完整性,通过 Lock::Async 强制执行。另请注意,它也没有引入任何并发:它实际上只是一个并发控制构造。
通常情况下,将supply 和react 与数据源一起使用,您可以在其中tap 它们并立即获得控制权。然后我们继续设置进一步的whenever 块,退出设置阶段,并且锁定可用于我们收到的任何消息。这种行为是您通常遇到的几乎所有用品都会出现的情况。 signal(...) 就是这种情况。当你给Supply.from-list(...) 一个显式调度器,传入$*SCHEDULER 时也是这种情况;在这种情况下,它会在池中安排从$*IN 读取的循环并立即交还控制权。
当我们遇到不是那样的行为时,问题就来了。如果我们点击Supply.from-list($*IN.lines),它默认在当前线程上从$*IN 读取以生成emit 的值,因为Supply.from-list 使用CurrentThreadScheduler 作为其默认值。那有什么作用呢?只需运行它要求立即运行的代码!
然而,这给我们留下了另一个谜团。鉴于Lock::Async 不可重入,那么如果我们:
- 获取锁以进行设置
- 在
Supply.from-list(...) 上调用tap,它同步运行并尝试为emit 一个值
- 尝试获取锁,以便我们可以处理值
然后我们就会遇到死锁,因为我们正试图获取一个已经被我们持有的不可重入锁。事实上,如果你在这里运行我对程序的 desugar,这正是发生的事情:它挂起。但是,原代码并没有挂起;它只是表现得有点尴尬。什么给了?
真正实现的其中一件事是在设置阶段检测此类情况;然后它会继续,并在设置阶段完成后恢复它。这意味着我们可以执行以下操作:
my $primes = supply {
.emit for ^Inf .grep(*.is-prime);
}
react {
whenever $primes { .say }
whenever Promise.in(3) { done }
}
然后让它解决。我不会在这里重现那种的乐趣,但如果足够巧妙地使用gather/take,它应该是可能的。