Search by

quazardous / gramphp

quazardous

A small workflow graph for work queues that already live in a storage: nodes, joins, retries, leases — over pluggable drivers (memory, MariaDB). A PHP port of grampy.

Package info

github.com/quazardous/gramphp

pkg:composer/quazardous/gramphp

Statistics

Installs: 11

Dependents: 1

Suggesters: 0

Stars: 1

Open Issues: 1

v0.1.0 2026-10-03 15:50 UTC

This package is auto-updated.

Last update: 2026-10-03 16:45:05 UTC


README

CI Packagist License: MIT

A small workflow graph for work queues that already live in a storage — a database, or plain memory: the storage is a driver.

A PHP port of grampy.

composer require quazardous/gramphp

PHP 8.2 or later. The MariaDB driver needs ext-pdo_mysql, or doctrine/dbal 4 to share a Doctrine connection.

You declare a DAG of nodes. Each subject (a job, an order, a feed…) goes through the nodes, and a node journal keeps track of how far each one has gone. Workers claim a node for eligible subjects, conclude it, and the graph decides what becomes claimable next — forks run in parallel, joins wait for their parents. Your subjects keep their own states: the journal sits beside them.

use Quazardous\GramPHP\Dag;
use Quazardous\GramPHP\Driver\Memory\MemoryDriver;
use Quazardous\GramPHP\Node;
use Quazardous\GramPHP\NodeJournal;
use Quazardous\GramPHP\Retry;

$journal = new NodeJournal(new MemoryDriver(), new Dag(
    new Node('pay'),
    new Node('ship', parents: ['pay'], retry: new Retry(limit: 3, delay: '30s')),
));

$lease = $journal->claim('pay', 10, candidates: [1, 2, 3]);   // your subjects, in priority order
foreach ($lease as $order) {
    // … pay the order …
}
$journal->conclude('pay', $lease, $lease->token);

$journal->claim('ship', 10, candidates: [1, 2, 3]);           // pay is done: ship is next

Why not a status column?

A status column and a few UPDATEs hold one line of steps, one worker at a time. What gramphp adds, each proven by the shared driver contract:

  • Two workers, one subject. A claim is atomic and returns a token; a slow worker whose lease went to another writes nothing (a fencing token).
  • A graph, not a line. Forks run in parallel, joins wait for their parents, optional steps and exclusive choices do not block what follows, a failure edge can run a compensation.
  • Never backwards. A node is never taken again once a descendant has started; going back is forget, and nothing is lost: every row taken away is archived in the history.
  • Time is declared. Retries with backoff and jitter, leases released by a janitor (expire), waits on signals received even before the wait began, optional nodes skipped after a grace delay.
  • One rule, one place. What is claimable is decided by the journal, in PHP, on rows the driver read — never re-expressed in SQL. A driver stores rows and offers a few atomic operations.

Drivers

driver storage concurrency
Driver\Memory\MemoryDriver arrays, one process the reference; for tests and single-process use
Driver\Mariadb\MariadbDriver InnoDB tables over PDO or Doctrine DBAL (MariaDB, MySQL 8) READ COMMITTED, in the caller's transaction
use Quazardous\GramPHP\Driver\Mariadb\MariadbDriver;
use Quazardous\GramPHP\Driver\Mariadb\Query;

foreach (MariadbDriver::schema(subjectType: 'int') as $ddl) {
    $pdo->exec($ddl);                             // the tables are yours
}
$pdo->exec('SET SESSION TRANSACTION ISOLATION LEVEL READ COMMITTED');
$journal = new NodeJournal(new MariadbDriver($pdo, subjectType: 'int'), $dag);

$pdo->beginTransaction();
$lease = $journal->claim('ship', 50, new Query(
    'SELECT id FROM orders WHERE paid = 1 ORDER BY priority DESC, id',
));
$pdo->commit();

The journal never commits: the application opens and ends the transaction around each call, so a node row and the application's own writes can go in one transaction.

With Doctrine, share the application's DBAL connection:

use Quazardous\GramPHP\Driver\Mariadb\DbalSql;

$journal = new NodeJournal(new MariadbDriver(new DbalSql($connection), subjectType: 'int'), $dag);
$lease = $connection->transactional(fn() => $journal->claim('ship', 50, $candidates));

A transaction mixing several claims and forgets on the same subjects can meet a deadlock: InnoDB reports it (error 1213, SQLSTATE 40001), rolls the transaction back, and the caller retries it — nothing is ever half-written. Transaction does exactly that, over PDO or DBAL:

use Quazardous\GramPHP\Driver\Mariadb\Transaction;

$lease = (new Transaction($pdo))->run(fn() => $journal->claim('ship', 50, $candidates));

It begins, runs the unit, commits; on a deadlock it runs the unit again from a fresh transaction, after a short random pause, at most attempts times. Anything else is rolled back and thrown as it came. The unit must be safe to run again — keep mails and calls after run returns.

Policies: one workflow, several ways of pushing it

Subjects differ in how hard they may be pushed — a slow partner, a bulk customer. A policy changes a node's settings, never its structure:

use Quazardous\GramPHP\Document;
use Quazardous\GramPHP\Graph;

$graph = new Graph(new Document('offers', version: '1'), $dag, policies: [
    'slow-partner' => ['call' => ['retry' => new Retry(5, '1m'), 'lease' => '2h']],
    'bulk' => ['call' => ['rate' => [new Rate(1000, '1h')]]],   // needs per: Per::Policy
]);
$journal = new NodeJournal($driver, $graph);
$journal->enroll($subjects, 'slow-partner');

A policy may change retry, lease, timeout, grace, rate, concurrency and lane (tune it, never add or remove one); the graph is checked under every policy. Each subject is read through its own: its retries, its lease (expire), its timeout and grace (settle), its budget and its lane.

Versions: a graph that changes under subjects in flight

A journal built on a Graph pins each subject, on its first write, to the graph it started on — the whole Document::identity(), namespace/name@version. It never touches a subject pinned elsewhere: version 1 finishes its subjects while version 2 takes the new ones. To move subjects across, on purpose:

$v2 = new NodeJournal($driver, $graphV2);
$v2->migrate($subjects, $graphV1, ['crop' => 'trim', 'legacy' => null]);   // renamed, dropped

A subject moves only if its journal could have been written on the new graph (every row's parents joined there), and never while a dropped node is held. All or nothing: one non-compliant subject raises MigrationError, naming each one and why, and nothing moves. Dropped rows go to the history.

Limits and groups

What a node uses is protected where the claim is decided:

use Quazardous\GramPHP\Group;
use Quazardous\GramPHP\Per;
use Quazardous\GramPHP\Rate;

new Node('call', rate: [new Rate(100, '1m'), new Rate(1000, '1h')]);   // bands, all apply
new Node('gpu', concurrency: 3);                                         // at most 3 running
new Node('api', concurrency: 1, per: Per::Policy);                       // one budget per policy
new Node('pack', parents: ['sort'], group: new Group(5, maxWait: '1h')); // five of a key, one lease

A rate band is a generic cell rate algorithm: limit per period, spread evenly, up to burst at once. A claim takes no more than the bands and the cap let through; under contention it takes fewer, never too many — the driver's guard holds the read and the write together. A grouping claim hands out a whole group of subjects sharing a key (grampy_key in SQL, Keyed otherwise) or nothing; past maxWait an incomplete group goes as it is. A lane's door honours its node's rate too.

Lanes: subjects that come back

A subject often comes back — a listing updated again, a file re-uploaded. A lane is the node it comes back through: it waits there, merges with the version already waiting, and settle lets it in once due, archiving the previous pass in the same write.

use Quazardous\GramPHP\Lane;

$journal = new NodeJournal($driver, new Dag(
    new Node('arrive', lane: Lane::throttle(cooldown: '1h')),
    new Node('scrape', parents: ['arrive']),
    new Node('publish', parents: ['scrape']),
));

$journal->arrive('arrive', [42], ref: 'v7');     // ['queued' => 1, 'merged' => 0, 'skipped' => 0]
$journal->settle($candidates);                   // ['arrive' => ['entered' => 1]] once due
preset keeps place under the names other tools gave it
Lane::throttle($cooldown) the last ref the first's throttle
Lane::debounce($delay) the last ref to the back debounce
Lane::dedupe() the first ref the first's dedupe
Lane::batch(maxSize: 100) every ref, in order the first's aggregator

maxWait bounds the wait whatever the rest, whileRunning: WhileRunning::Skip drops an arrival while a pass runs, urgent: true skips the cooldown — never a running pass. A lane may also merge with a function of yours, named in the graph (Merge::fn('ends')) and given to the journal (mergers: ['ends' => $fn]). Every merge, drop and entry is noted in the history.

Items: speak objects, not ids

The core works on ids and knows nothing about your data. The optional items layer holds the handlers that read your own objects and turns their answers into calls the journal already understands:

use Quazardous\GramPHP\Items\Adapter;
use Quazardous\GramPHP\Items\Items;

final class Orders extends Adapter
{
    public function idOf(mixed $candidate): int|string
    {
        return $candidate instanceof Order ? $candidate->id : $candidate;
    }

    public function inflate(array $candidates): iterable   // ids in, orders out, in one call
    {
        $ids = array_filter($candidates, is_int(...));
        return [...array_diff_key($candidates, $ids), ...$this->repository->findByIds($ids)];
    }

    public function policyOf(mixed $order): ?string { return $order->plan; }
    public function branch(mixed $order, string $node): ?string { return $order->digital ? 'mail' : 'ship'; }
    public function applies(mixed $order, string $node): bool { return 'gift-wrap' !== $node || $order->gift; }
}

$items = new Items($journal, new Orders($repository));
$lease = $items->claim('pay', 10, $orders);      // your objects in — or ids, or a Query
foreach ($lease as $order) {                     // your objects out
    // …
}
$items->conclude('pay', $lease);                 // the token travels with the lease

Objects handed in travel with the claim and are never loaded twice; ids — and a driver's Query — are loaded with inflate, once per call, and an id nothing loads lands in $lease->missing. applies() gives an optional node up for the items it is not for; branch() names the way out of a choice. The layer translates and stops there: everything it does is a call the id-based API could have made by hand.

The graph as data, and drawn

A graph has one canonical form — grampy's format (grampy/1), so the same file is read and written by the Python and the PHP implementations alike:

file_put_contents('offers.graph.json', $graph->toJson());   // only what differs from a default
$graph = Graph::fromJson(file_get_contents('offers.graph.json'));   // strict: a lie is refused with its path

And drawn, every mechanism with a shape of its own, counts overlaid if given:

use Quazardous\GramPHP\Diagram;

echo Diagram::mermaid($graph, Diagram::overlay($journal));   // flowchart, live counts per node
echo Diagram::stateDiagram($graph);                          // the statechart reading
echo Diagram::dot($graph);                                   // Graphviz

Monitoring

$journal->snapshot($candidates) gives, per node, plain numbers to sample and alert on: how many subjects stand where, oldest_running (a stuck worker past its lease), next_due (negative: retries overdue), ready (how many a claim could take now) and oldest_ready (starvation).

How it works

docs/concepts.md: the journal's own vocabulary, the claim rule, the history, and how a node maps onto your subject's states.

Status of the port

Ported: the graph and its claim rule, joins (on, need), choices, loops, retries, leases, waits and signals, grace, skip, adopt, forget, release, the history, counts, stages and snapshots, rate limits, concurrency caps and groups, policies, versions and migration, lanes (throttle, debounce, dedupe, batch, merge functions), the items layer, the graph's JSON form and the diagrams — with the memory and MariaDB drivers, both certified by the shared contract, concurrency included.

Every mechanism grampy's graph declares runs here. Not ported: grampy's SQLite and PostgreSQL drivers, and the optional one-statement fast paths a driver may offer (skip_where, ready_count).

Development

Everything runs in Docker:

make build install
make test      # unit, memory and MariaDB contract (concurrency included)
make lint      # PHPStan (max level), PHP-CS-Fixer, composer validate
make check     # what CI runs

License

MIT.