PHP code example of hiblaphp / stream

1. Go to this page and download the library: Download hiblaphp/stream library. Choose the download type require.

2. Extract the ZIP file and open the index.php.

3. Add this code to the index.php.
    
        
<?php
require_once('vendor/autoload.php');

/* Start to develop here. Best regards https://php-download.com/ */

    

hiblaphp / stream example snippets


use Hibla\Stream\ReadableResourceStream;
use Hibla\Stream\WritableResourceStream;
use Hibla\Stream\DuplexResourceStream;
use Hibla\Stream\PromiseReadableStream;
use Hibla\Stream\PromiseWritableStream;
use Hibla\Stream\CompositeStream;
use Hibla\Stream\ThroughStream;

$readable  = new ReadableResourceStream(fopen('/path/to/input.log', 'rb'));
$writable  = new WritableResourceStream(fopen('/path/to/output.log', 'wb'));
$duplex    = new DuplexResourceStream(fopen('/path/to/data.bin', 'r+b'));
$composite = new CompositeStream($readable, $writable);
$through   = new ThroughStream(fn(string $data) => strtoupper($data));

// Promise-based variants
$readable = new PromiseReadableStream(fopen('/path/to/input.log', 'rb'));
$writable = new PromiseWritableStream(fopen('/path/to/output.log', 'wb'));

// Named constructor alternative for promise streams
$readable = PromiseReadableStream::fromResource(fopen('/path/to/input.log', 'rb'));
$writable = PromiseWritableStream::fromResource(fopen('/path/to/output.log', 'wb'));

use Hibla\Stream\Stream;

$readable = Stream::readableFile('/path/to/input.log');
$writable = Stream::writableFile('/path/to/output.log');
$stdin    = Stream::stdin();
$stdout   = Stream::stdout();

use Hibla\Stream\ReadableResourceStream;

$stream = new ReadableResourceStream(fopen('/var/log/app.log', 'rb'));

// Attach all listeners first — no data flows yet
$stream->on('data', function (string $chunk) {
    echo $chunk;
});

$stream->on('end', function () {
    echo "Stream fully consumed\n";
});

$stream->on('error', function (\Throwable $e) {
    echo "Error: " . $e->getMessage() . "\n";
});

// Start the flow only after all listeners are in place
$stream->resume();

$stream->on('data', function (string $chunk) use ($stream) {
    $stream->pause();
    processChunk($chunk);
    $stream->resume();
});

$stream = new ReadableResourceStream(fopen('/tmp/data.bin', 'rb'));

$stream->on('data', function (string $chunk) use ($stream) {
    echo "Read: " . strlen($chunk) . " bytes\n";
    $stream->pause();
});

$stream->resume();

// Rewind to the beginning and read again
$stream->seek(0);
$stream->resume();

// WRONG — seek() on a pipe silently returns false
$stream = new ReadableResourceStream(STDIN);
$stream->seek(0); // false — silently ignored

// CORRECT
if ($stream->seek(0) === false) {
    // resource is non-seekable — handle accordingly
}

$position = $stream->tell();
$stream->seek(512, SEEK_SET);   // seek to byte 512
$stream->seek(0, SEEK_END);     // seek to end of file
$stream->seek(-128, SEEK_CUR);  // seek relative to current position

use Hibla\Stream\WritableResourceStream;

$stream = new WritableResourceStream(fopen('/tmp/output.txt', 'wb'));

$stream->on('finish', fn() => echo "All data written\n");
$stream->on('error', fn(\Throwable $e) => echo "Write error: " . $e->getMessage() . "\n");

$stream->write("First line\n");
$stream->write("Second line\n");
$stream->end("Final line\n");

use Hibla\Stream\WritableResourceStream;

$socket   = stream_socket_client('tcp://example.com:9000');
$writable = new WritableResourceStream($socket, softLimit: 65536);

function pump(string $data, WritableResourceStream $writable): void
{
    $canContinue = $writable->write($data);

    if ($canContinue === false) {
        $writable->once('drain', function () use ($writable) {
            pump(getNextChunk(), $writable);
        });
    }
}

use Hibla\Stream\ReadableResourceStream;
use Hibla\Stream\WritableResourceStream;

$source      = new ReadableResourceStream(fopen('/tmp/input.bin', 'rb'));
$destination = new WritableResourceStream(fopen('/tmp/output.bin', 'wb'));

$destination->on('finish', fn() => echo "Transfer complete\n");

// pipe() calls resume() internally — no need to call it yourself
$source->pipe($destination);

use Hibla\Stream\ThroughStream;

$source
    ->pipe(new ThroughStream(fn($d) => gzencode($d)))
    ->pipe(new WritableResourceStream(fopen('/tmp/compressed.gz', 'wb')));

// Source A's 'end' will NOT call dest->end()
$sourceA->pipe($dest, ['end' => false]);
$sourceA->on('end', function () use ($sourceB, $dest) {
    $sourceB->pipe($dest); // this one WILL close dest when finished
});

use Hibla\Stream\ThroughStream;

// Transform: compress mid-pipe
$source
    ->pipe(new ThroughStream(fn(string $data) => gzencode($data)))
    ->pipe($destination);

// Spy: inspect data mid-pipe without modifying it
$spy = new ThroughStream(function (string $data) {
    fwrite(STDERR, sprintf("[spy] %d bytes\n", strlen($data)));
    return $data; // must return data to pass it through
});

$source->pipe($spy)->pipe($destination);

use Hibla\Stream\DuplexResourceStream;

$socket = stream_socket_client('tcp://api.example.com:80');
$duplex = new DuplexResourceStream($socket);

$duplex->on('data', function (string $response) use ($duplex) {
    echo $response;
    $duplex->close();
});

$duplex->on('error', fn(\Throwable $e) => echo "Error: " . $e->getMessage() . "\n");

// Write is always available immediately
$duplex->write("GET / HTTP/1.0\r\nHost: api.example.com\r\n\r\n");

// Start receiving after listeners are in place
$duplex->resume();

use Hibla\Stream\CompositeStream;
use Hibla\Stream\ReadableResourceStream;
use Hibla\Stream\WritableResourceStream;

$process = proc_open('ffmpeg -i pipe:0 -f mp3 pipe:1', [
    0 => ['pipe', 'r'],
    1 => ['pipe', 'w'],
], $pipes);

$composite = new CompositeStream(
    new ReadableResourceStream($pipes[1]), // process stdout → our readable
    new WritableResourceStream($pipes[0])  // our writable → process stdin
);

// Attach listeners before resuming
$composite->on('data', fn(string $chunk) => saveChunk($chunk));
$composite->on('end', fn() => proc_close($process));

$composite->resume();
$inputStream->pipe($composite);

use Hibla\Stream\PromiseReadableStream;
use function Hibla\await;

$stream = new PromiseReadableStream(fopen('/tmp/data.txt', 'rb'));

// Read the next chunk (up to $length bytes)
$chunk = await($stream->readAsync(1024));

// Read a full line including the trailing newline character
$line = await($stream->readLineAsync());

// Read the entire stream into a single string
$contents = await($stream->readAllAsync());

// CORRECT — stops only on null (EOF)
while (($line = await($stream->readLineAsync())) !== null) {
    processLine(rtrim($line));
}

// WRONG — stops on any falsy chunk, including valid data like "0" or "\n"
while ($line = await($stream->readLineAsync())) {
    processLine($line);
}

$contents = await($stream->readAllAsync());

$contents = await($stream->readAllAsync(maxLength: 524288));  // 512 KiB limit
$line     = await($stream->readLineAsync(maxLength: 4096));   // 4 KiB per line

/**
 * Read exactly $length bytes from a stream.
 * Returns null if EOF is reached before $length bytes are available.
 *
 * @return string|null
 */
function readExact(PromiseReadableStream $stream, int $length): ?string
{
    $buffer    = '';
    $remaining = $length;

    while ($remaining > 0) {
        $chunk = await($stream->readAsync($remaining));

        if ($chunk === null) {
            return null; // EOF before enough bytes arrived
        }

        $buffer    .= $chunk;
        $remaining -= strlen($chunk);
    }

    return $buffer;
}

// Read a length-prefixed binary message:
// [ 4-byte uint32 length ][ N bytes payload ]

$header = readExact($stream, 4);
if ($header === null) {
    return; // clean EOF — no more messages
}

$payloadLength = unpack('N', $header)[1];
$payload       = readExact($stream, $payloadLength);

if ($payload === null) {
    throw new \RuntimeException("Truncated message: stream ended early");
}

use Hibla\Stream\PromiseWritableStream;
use function Hibla\await;

$stream = new PromiseWritableStream(fopen('/tmp/out.txt', 'wb'));

// Resolves with the number of bytes buffered
// Backpressure is handled internally — no drain listener needed
$bytes = await($stream->writeAsync("Hello, world\n"));

// Write a line (appends "\n" automatically)
await($stream->writeLineAsync("Another line"));

// End the stream and wait for all data to flush
// Resolves only after 'finish' fires — all data is durably written
await($stream->endAsync());

use Hibla\Stream\PromiseReadableStream;
use Hibla\Stream\WritableResourceStream;
use function Hibla\await;

$source = new PromiseReadableStream(fopen('/tmp/large.bin', 'rb'));
$dest   = new WritableResourceStream(fopen('/tmp/copy.bin', 'wb'));

$totalBytes = await($source->pipeAsync($dest));
echo "Transferred: $totalBytes bytes\n";

use Hibla\EventLoop\Loop;
use function Hibla\await;

$readable    = new PromiseReadableStream(fopen('/tmp/large.log', 'rb'));
$readPromise = $readable->readLineAsync();

$timerId = Loop::addTimer(2.0, function () use ($readPromise) {
    $readPromise->cancel();
});

try {
    $line = await($readPromise);
    Loop::cancelTimer($timerId);
} catch (\Hibla\Promise\Exceptions\CancelledException $e) {
    echo "Read cancelled — no data within 2 seconds\n";
}

$transferPromise = $source->pipeAsync($dest);

Loop::addTimer(5.0, fn() => $transferPromise->cancel());

try {
    $totalBytes = await($transferPromise);
} catch (\Hibla\Promise\Exceptions\CancelledException $e) {
    echo "Transfer cancelled\n";
    // $source and $dest are still open — you decide what to do with them
}

use Hibla\Cancellation\CancellationTokenSource;
use function Hibla\await;

$cts = new CancellationTokenSource(30.0); // 30 second hard limit

try {
    while (($line = await($stream->readLineAsync(), $cts->token)) !== null) {
        processLine(rtrim($line));
    }
} catch (\Hibla\Promise\Exceptions\CancelledException $e) {
    echo "Stream read timed out after 30 seconds\n";
}

use Hibla\Stream\PromiseReadableStream;
use Hibla\Stream\WritableResourceStream;
use Hibla\Stream\CompositeStream;
use Hibla\Stream\ReadableResourceStream;
use function Hibla\await;

// Read from STDIN line by line
$stdin = new PromiseReadableStream(STDIN);

while (($line = await($stdin->readLineAsync())) !== null) {
    processLine(rtrim($line));
}

// Write to STDOUT respecting backpressure
$stdout = new WritableResourceStream(STDOUT);
$stdout->write("Hello from async PHP\n");

// Write errors to STDERR
$stderr = new WritableResourceStream(STDERR);
$stderr->write("Something went wrong\n");

// Combined interactive console — STDIN readable, STDOUT writable
$stdio = new CompositeStream(
    new ReadableResourceStream(STDIN),
    new WritableResourceStream(STDOUT)
);

$stdio->on('data', fn(string $input) => $stdio->write("Echo: $input"));
$stdio->resume();

$stream = new ReadableResourceStream(fopen('/var/log/app.log', 'rb'));

$stream->on('data',  fn(string $chunk)  => echo $chunk);
$stream->on('end',   fn()              => echo "Done reading\n");
$stream->on('close', fn()              => echo "Stream closed\n");
$stream->on('error', fn(\Throwable $e) => echo "Error: " . $e->getMessage() . "\n");

$stream->resume();

$stream = new WritableResourceStream(fopen('/tmp/output.log', 'wb'));

$stream->on('drain',  fn()             => echo "Drained, can write again\n");
$stream->on('finish', fn()             => echo "All data written\n");
$stream->on('close',  fn()             => echo "Stream closed\n");
$stream->on('error',  fn(\Throwable $e) => echo "Write error: " . $e->getMessage() . "\n");

$stream->write("Hello\n");
$stream->end("Goodbye\n");

// CORRECT — just handle the error; close fires on its own
$stream->on('error', function (\Throwable $e) {
    echo "Error: " . $e->getMessage() . "\n";
});

$stream->on('close', function () {
    cleanupResources();
});

// Wrong — buffer may be discarded if $stream goes out of scope
$stream = new WritableResourceStream(fopen('/tmp/output.txt', 'wb'));
$stream->write("Important data\n");
// $stream goes out of scope — destructor calls close(), buffer is discarded

// Correct — drain the buffer before releasing the stream
$stream = new WritableResourceStream(fopen('/tmp/output.txt', 'wb'));
$stream->on('finish', fn() => echo "All data flushed\n");
$stream->write("Important data\n");
$stream->end();

$stream = new PromiseWritableStream(fopen('/tmp/output.txt', 'wb'));
await($stream->writeAsync("Important data\n"));
await($stream->endAsync()); // all data flushed before this resolves