PHP code example of saravanasai / stream-pulse

1. Go to this page and download the library: Download saravanasai/stream-pulse 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/ */

    

saravanasai / stream-pulse example snippets


// In your producer Laravel application
// Publish an event
StreamPulse::publish('orders', ['id' => 1234, 'amount' => 99.99]);

// In your consumer Laravel application
// Register an event handler - app service provider
StreamPulse::on('orders', function ($payload, $messageId) {
    OrderProcessor::process($payload);
});

// Run the consumer to process events (in the consumer app)
// php artisan streampulse:consume orders

use StreamPulse\StreamPulse\Facades\StreamPulse;

// Publish an event to a topic
StreamPulse::publish('orders', [
    'id' => 1234,
    'customer' => 'John Doe',
    'total' => 99.99,
    'items' => [
        ['product_id' => 101, 'quantity' => 2, 'price' => 49.99]
    ]
]);

// Publish after DB transaction commits
StreamPulse::publishAfterCommit('orders', $orderData);

use Illuminate\Support\Facades\DB;
use StreamPulse\StreamPulse\Facades\StreamPulse;

DB::transaction(function () {
    // Create an order in the database
    $order = Order::create([
        'customer_id' => 123,
        'amount' => 99.99,
    ]);

    // This event will only be published if the transaction commits successfully
    StreamPulse::publishAfterCommit('orders', [
        'id' => $order->id,
        'status' => 'created',
        'customer_id' => $order->customer_id,
        'amount' => $order->amount,
    ]);

    // If the transaction fails or is rolled back, no event will be published
});

// In a service provider
StreamPulse::on('orders', function ($payload, $messageId) {
    // Process the order
    OrderProcessor::process($payload);
});

// Then run the consumer command
// php artisan streampulse:consume orders

// Define global defaults in config/streampulse.php
'defaults' => [
    'retention' => 1000, // Keep 1000 events per stream by default
],

// Override for specific topics
'topics' => [
    'orders' => [
        'retention' => 5000, // Keep more events for important topics
    ],
    'logs' => [
        'retention' => 500, // Keep fewer events for high-volume topics
    ],
],

return [
    /*
    | Default Driver
    | Available: "redis", "nats" (NATS support coming in future releases)
    */
    'driver' => env('STREAMPULSE_DRIVER', 'redis'),

    /*
    | Strict Mode - When enabled, only explicitly defined topics can be used
    */
    'strict_mode' => env('STREAMPULSE_STRICT_MODE', true),

    /*
    | Auto Processing - Automatically process pending messages that exceed retry limits
    */
    'auto_process_pending' => env('STREAMPULSE_AUTO_PROCESS', true),

    /*
    | Global Defaults - Applied to all topics unless overridden
    */
    'defaults' => [
        'max_retries' => 3,
        'dlq' => 'dead_letter',
        'retention' => 1000,
        'min_idle_time' => 30000,  // Minimum time (ms) before re-processing pending messages
        'preserve_order' => false, // Whether to enforce strict message ordering
    ],

    /*
    | Topics - Per-topic configuration overrides
    */
    'topics' => [
        'orders' => [
            'max_retries' => 5,
            'dlq' => 'orders_dlq',
            'retention' => 5000,
            'min_idle_time' => 60000,
            'preserve_order' => true,
        ],
        // Other topics...
    ],

    /*
    | Drivers - Backend-specific configuration
    */
    'drivers' => [
        'redis' => [
            'connection' => env('REDIS_CONNECTION', 'default'),
            'stream_prefix' => 'streampulse:',
        ],
    ],

    /*
    | UI Settings - Dashboard configuration
    */
    'ui' => [
        'enabled' => env('STREAMPULSE_UI_ENABLED', true),
        'page_size' => env('STREAMPULSE_UI_PAGE_SIZE', 50),
        'route_prefix' => 'stream-pulse',
    ],
];

use StreamPulse\StreamPulse\Facades\StreamPulse;

// Consume events from a topic with a consumer group
StreamPulse::consume('orders', 'order-processors', function ($payload, $messageId) {
    try {
        // Process the event
        OrderProcessor::process($payload);

        // Acknowledge the message as processed
        StreamPulse::ack('orders', $messageId, 'order-processors');
    } catch (\Exception $e) {
        // Handle error
        // The message will remain in pending state and can be retried
        // After max retries, it will be moved to the DLQ automatically
    }
});
bash
php artisan vendor:publish --tag="stream-pulse-config"
bash
php artisan vendor:publish --tag="stream-pulse-views"