NAME
Acme::Parataxis::Channel - A simple message queue for inter-fiber communication
SYNOPSIS
use Acme::Parataxis;
use Acme::Parataxis::Channel;
my $q = Acme::Parataxis::Channel->new( 4 );
async {
fiber { $q->put( $_ ) for 1 .. 8 }; # producers
say $q->get for 1 .. 8; # consumer
};
DESCRIPTION
A simple message queue that allows you to send and receive data between fibers. If the channel is full, writers block; if it is empty, readers block. Both ends can be used by as many fibers as you want concurrently.
A channel of size 1 is a rendezvous point (no buffering: put waits for a matching get); to buffer one element use size 2, and so on.
Channels are internally implemented using two Acme::Parataxis::Semaphore instances to coordinate producers and consumers.
CONSTRUCTOR
new( [...] )
my $ch = Acme::Parataxis::Channel->new;
my $ch = Acme::Parataxis::Channel->new(capacity => 10);
Creates a new channel. The optional capacity parameter sets the maximum number of items the channel can hold and defaults to 2_000_000_000.
METHODS
put( $value )
$ch->put($value);
Append a value to the channel. Blocks the current fiber if the channel is at capacity, waiting until space becomes available.
get( )
my $value = $ch->get;
Remove and return the next value from the channel. Blocks the current fiber if the channel is empty, waiting until a value is available.
size( )
my $n = $ch->size;
Returns the number of items currently in the channel.
adjust( $diff )
$ch->adjust( $diff );
Adjust the channel capacity by $diff. This can be used to grow or shrink the effective capacity at runtime.
shutdown( )
$ch->shutdown;
Shut down the channel by releasing a large number of permits on the internal get semaphore. This unblocks any fibers waiting on get( ), allowing them to drain remaining items.
AUTHOR
Sanko Robinson https://github.com/sanko
LICENSE
Copyright (C) Sanko Robinson.
This library is free software; you can redistribute it and/or modify it under the terms found in the Artistic License 2.