PHP code example of olexin-pro / data-processing-pipeline

1. Go to this page and download the library: Download olexin-pro/data-processing-pipeline 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/ */

    

olexin-pro / data-processing-pipeline example snippets


namespace App\Pipelines\Steps;

use DataProcessingPipeline\Pipelines\Contracts\PipelineStepInterface;
use DataProcessingPipeline\Pipelines\Contracts\PipelineContextInterface;
use DataProcessingPipeline\Pipelines\Results\GenericPipelineResult;
use DataProcessingPipeline\Pipelines\Enums\ConflictPolicy;

class EmailFormatterStep implements PipelineStepInterface
{
    public function handle(PipelineContextInterface $context): GenericPipelineResult
    {
        $email = $context->getContent('user.email', '');
        
        return new GenericPipelineResult(
            key: 'email',
            data: ['value' => strtolower(trim($email))],
            policy: ConflictPolicy::MERGE,
            priority: 10,
            provenance: self::class
        );
    }
}



namespace App\Pipelines\Steps;

use DataProcessingPipeline\Pipelines\Contracts\PipelineStepInterface;
use DataProcessingPipeline\Pipelines\Contracts\PipelineContextInterface;
use DataProcessingPipeline\Pipelines\Results\GenericPipelineResult;
use DataProcessingPipeline\Pipelines\Enums\ConflictPolicy;

class EmailValidatorStep implements PipelineStepInterface
{
    public function handle(PipelineContextInterface $context): GenericPipelineResult
    {
        $emailData = $context->getResult('email')->getData();
        $email = $emailData['value'] ?? null;

        $isValid = filter_var($email, FILTER_VALIDATE_EMAIL) !== false;

        return new GenericPipelineResult(
            key: 'email',
            data: [
                'value' => $email
                'status' => $isValid ? 'verified' : 'invalid',
            ],
            policy: ConflictPolicy::MERGE,
            priority: 20,
            provenance: self::class
        );
    }
}


use DataProcessingPipeline\Pipelines\Context\PipelineContext;
use DataProcessingPipeline\Pipelines\Contracts\PipelineRunnerInterface;
use DataProcessingPipeline\Pipelines\History\PipelineHistoryRecorder;
use App\Pipelines\Steps\EmailFormatterStep;
use App\Pipelines\Steps\EmailValidatorStep;

$context = PipelineContext::make(['user' => ['email' => '[email protected]']]);
$recorder = new PipelineHistoryRecorder('user-processing');

$runner = app(PipelineRunnerInterface::class);
$runner->setRecorder($recorder)
    ->addStep(new EmailFormatterStep())
    ->addStep(new EmailValidatorStep());

$result = $runner->run($context);

// Access results
$emailData = $result->getResult('email')->getData();
// ['value' => '[email protected]', 'status' => 'verified']

// Build data
$built = $result->build();
// ['email' => ['value' => '[email protected]', 'status' => 'verified']]

use App\Pipelines\Steps\EmailFormatterStep;
use App\Pipelines\Steps\EmailValidatorStep;
use DataProcessingPipeline\Pipelines\Context\PipelineContext;
use DataProcessingPipeline\Jobs\ProcessPipelineJob;
use DataProcessingPipeline\Services\Notifiers\LogNotifier;

$payload = ['user' => ['email' => '[email protected]']];

$steps = [
    EmailFormatterStep::class,
    EmailValidatorStep::class,
];


$context = PipelineContext::make($payload);

ProcessPipelineJob::dispatch(
    contextData: $context->toArray(),
    stepClasses: $steps,
    pipelineName: 'email-processing'
    recordHistory: false
    notifierClass: LogNotifier::class
);


use DataProcessingPipeline\Pipelines\Context\PipelineContext;
use DataProcessingPipeline\Pipelines\Contracts\ConflictResolverInterface;

$context = new PipelineContext::make(
    payload: ['user' => ['id' => 1, 'email' => '[email protected]']],
    meta: ['request_id' => 'abc-123'],
    conflictResolver: app()->make(ConflictResolverInterface::class)
);

interface PipelineResultInterface
{
    public function getKey(): string;
    public function getData(): int|float|array|bool|string|null;
    public function getPolicy(): ConflictPolicy;
    public function getPriority(): int;
    public function getProvenance(): string;
    public function getStatus(): ResultStatus;
    public function getMeta(): array;
}

// Step 1
['name' => 'John', 'age' => 30]

// Step 2
['age' => 31, 'city' => 'NYC']

// Result
['name' => 'John', 'age' => 31, 'city' => 'NYC']

return new GenericPipelineResult(
    key: 'user',
    data: ['id' => 2, 'name' => 'Jane'],
    policy: ConflictPolicy::OVERWRITE
);

return new GenericPipelineResult(
    key: 'config',
    data: ['version' => 2],
    policy: ConflictPolicy::SKIP
);

use DataProcessingPipeline\Pipelines\Contracts\ConflictResolverInterface;
use DataProcessingPipeline\Pipelines\Contracts\PipelineResultInterface;
use DataProcessingPipeline\Pipelines\Contracts\PipelineContextInterface;

class PriorityConflictResolver implements ConflictResolverInterface
{
    public function resolve(
        PipelineResultInterface $existing,
        PipelineResultInterface $incoming,
        PipelineContextInterface $context
    ): PipelineResultInterface {
        return $incoming->getPriority() > $existing->getPriority()
            ? $incoming
            : $existing;
    }
}

return new GenericPipelineResult(
    key: 'user',
    data: ['role' => 'admin'],
    policy: ConflictPolicy::CUSTOM,
    meta: ['resolver' => PriorityConflictResolver::class]
);

// Low priority
new GenericPipelineResult(
    key: 'settings',
    data: ['theme' => 'light'],
    priority: 5
);

// High priority
new GenericPipelineResult(
    key: 'settings',
    data: ['theme' => 'dark'],
    priority: 20
);

// Result: theme = 'dark'

use DataProcessingPipeline\Pipelines\History\PipelineHistoryRecorder;

$recorder = new PipelineHistoryRecorder('user-processing');

$runner = new PipelineRunner(
    steps: [
        new EmailFormatterStep(),
        new EmailValidatorStep(),
    ],
    recorder: $recorder
);

// or
$recorder = app()->makeWith(
    PipelineHistoryRecorderInterface::class, 
    [
        'pipelineName' => $pipelineName, 
        'enabled' => true
    ]
);

$runner = app(PipelineRunnerInterface::class)
    ->setRecorder($recorder)
    ->addStep(new EmailFormatterStep())
    ->addStep(new EmailValidatorStep());

$runner->run($context);

class RiskyStep implements PipelineStepInterface
{
    public function handle(PipelineContext $context): PipelineResultInterface
    {
        if ($someCondition) {
            throw new \RuntimeException('Processing failed');
        }
        
        return new GenericPipelineResult(
            key: 'result',
            data: ['success' => true]
        );
    }
}

$result = $runner->run($context);

if (!empty($result->meta['errors'])) {
    foreach ($result->meta['errors'] as $error) {
        Log::error('Pipeline step failed', [
            'step' => $error['step'],
            'message' => $error['message']
        ]);
    }
}

$recorder = new PipelineHistoryRecorder('user-processing');
$runner = app(PipelineRunnerInterface::class);

$runner->setRecorder($recorder)
       ->addStep(new EmailFormatterStep())
       ->addStep(new EmailValidatorStep())
       ->addStep(new EmailDomainCheckerStep());


$runner->run($context);

class ConditionalStep implements PipelineStepInterface
{
    public function handle(PipelineContextInterface $context): PipelineResultInterface
    {
        if (!$context->getContent('process_email')) {
            return null;
        }

        // normal logic...
    }
}

// Domain: E-commerce Order Processing


use DataProcessingPipeline\Pipelines\Contracts\{
    PipelineContextInterface,
    PipelineResultInterface,
    PipelineStepInterface
};
use DataProcessingPipeline\Pipelines\Results\GenericPipelineResult;
use DataProcessingPipeline\Pipelines\Enums\ConflictPolicy;
use DataProcessingPipeline\Pipelines\Context\PipelineContext;
use DataProcessingPipeline\Pipelines\Runner\PipelineRunner;

/**
 * Validate order
 */
final class ValidateOrderStep implements PipelineStepInterface
{
    public function handle(PipelineContextInterface $context): PipelineResultInterface
    {
        $order = $context->getContent('order', []);
        $errors = [];

        if (empty($order['items'])) {
            $errors[] = 'Order must contain at least one item.';
        }

        if (!isset($order['id'])) {
            $errors[] = 'Order ID is missing.';
        }

        return new GenericPipelineResult(
            key: 'validation',
            data: [
                'valid'  => empty($errors),
                'errors' => $errors,
            ],
            priority: 100
        );
    }
}

/**
 * Calculate per-product totals
 */
final class CalculateProductsStep implements PipelineStepInterface
{
    public function handle(PipelineContextInterface $context): PipelineResultInterface
    {
        $items = $context->getContent('order.items', []);

        $products = collect($items)->map(fn ($p) => [
            'name'  => $p['name'],
            'price' => (float) $p['price'],
            'qty'   => (int) $p['qty'],
            'total' => (float) $p['price'] * (int) $p['qty'],
        ])->toArray();

        return new GenericPipelineResult(
            key: 'products',
            data: $products,
            policy: ConflictPolicy::MERGE
        );
    }
}

/**
 * Calculate totals (subtotal, tax, total)
 */
final class CalculateTotalsStep implements PipelineStepInterface
{
    public function handle(PipelineContextInterface $context): PipelineResultInterface
    {
        $products = $context->getResult('products')?->getData() ?? [];

        $subtotal = array_sum(array_column($products, 'total'));
        $tax = round($subtotal * 0.1, 2);
        $total = $subtotal + $tax;

        return new GenericPipelineResult(
            key: 'totals',
            data: compact('subtotal', 'tax', 'total')
        );
    }
}

/**
 * Apply coupon discounts
 */
final class ApplyDiscountStep implements PipelineStepInterface
{
    public function handle(PipelineContextInterface $context): PipelineResultInterface
    {
        $totals = $context->getResult('totals')?->getData() ?? [];
        $coupon = $context->getContent('coupon_code');

        $discountRate = match ($coupon) {
            'SAVE10' => 0.10,
            'SAVE20' => 0.20,
            default  => 0.0,
        };

        $discount = round(($totals['total'] ?? 0) * $discountRate, 2);

        return new GenericPipelineResult(
            key: 'totals',
            data: [
                'discount' => $discount,
                'total'    => ($totals['total'] ?? 0) - $discount,
            ],
            policy: ConflictPolicy::MERGE
        );
    }
}

// ----------------------------------------------------
// ▶ Example usage
// ----------------------------------------------------

$context = new PipelineContext([
    'order' => [
        'id' => 123,
        'items' => [
            ['name' => 'Product A', 'price' => 100, 'qty' => 1],
            ['name' => 'Product B', 'price' => 50,  'qty' => 8],
        ],
    ],
    'coupon_code' => 'SAVE10',
]);

$runner = new PipelineRunner([
    new CalculateProductsStep(),
    new CalculateTotalsStep(),
    new ValidateOrderStep(),
    new ApplyDiscountStep(),
]);

$result = $runner->run($context);

// ----------------------------------------------------
// ▶ Results
// ----------------------------------------------------

$totals = $result->getResult('totals')->getData();
/*
[
    'subtotal' => 500,        // 100*1 + 50*8
    'tax' => 50,
    'discount' => 55,
    'total' => 495
]
*/

$data = $result->build();
/*
[
    'products' => [
        ['name' => 'Product A', 'price' => 100.0, 'qty' => 1, 'total' => 100.0],
        ['name' => 'Product B', 'price' => 50.0,  'qty' => 8, 'total' => 400.0],,
    ],
    'totals' => [
        'subtotal' => 500,
        'tax' => 50,
        'discount' => 55,
        'total' => 495,
    ],
    'validation' => [
        'valid' => true,
        'errors' => [],
    ],
]
*/


use DataProcessingPipeline\Pipelines\Context\PipelineContext;
use DataProcessingPipeline\Pipelines\Contracts\ConflictResolverInterface;

new PipelineContext(
    array $payload,
    array $results = [],
    array $meta = [],
    ?ConflictResolverInterface $conflictResolver = null
);
// or
PipelineContext::make(
    array $payload,
    array $results = [],
    array $meta = [],
    ?ConflictResolverInterface $conflictResolver = null
);

$context->addResult(PipelineResultInterface $result): void
$context->getResult(string $key): ?PipelineResultInterface
$context->getContent(string $key, mixed $default = null): mixed
$context->hasResult(string $key): bool
$context->toArray(): array
$context->build(): array

use DataProcessingPipeline\Pipelines\Enums\ConflictPolicy;
use DataProcessingPipeline\Pipelines\Enums\ResultStatus;

new GenericPipelineResult(
    string $key,
    int|float|array|bool|string|null $data,
    ConflictPolicy $policy = ConflictPolicy::MERGE,
    int $priority = 10,
    string $provenance = '',
    ResultStatus $status = ResultStatus::OK,
    array $meta = []
);

GenericPipelineResult::make(
        string $key,
        int|float|array|bool|string|null $data,
        ConflictPolicy $policy = ConflictPolicy::MERGE,
        int $priority = 10,
        string $provenance = '',
        ResultStatus $status = ResultStatus::OK,
        array $meta =[],
    ): self
bash
php artisan vendor:publish --provider="DataProcessingPipeline\PipelineServiceProvider" --tag=pipeline-migrations
php artisan migrate
bash
php artisan make:pipeline-step EmailFormatterStep
# => app/Pipeline/Steps/EmailFormatterStep.php
# => namespace App\Pipeline\Steps;
bash
php artisan make:step Order/TotalCalculationStep
# => app/Pipeline/Steps/Order/TotalCalculationStep.php
# => namespace App\Pipeline\Steps\Order;
bash
php artisan make:pipeline-step EmailValidatorStep --policy=MERGE --priority=20
php artisan make:pipeline-step "Order/Discount/ApplyCouponStep" --key=coupon --policy=OVERWRITE