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.