Lines 94.72% 305 / 322
Methods 85.36% 35 / 41
Classes 0.00% 0 / 1
Name Lines Methods CRAP
 __construct 100.00% 1 / 1 100.00% 1 / 1 1
 table 100.00% 1 / 1 100.00% 1 / 1 1
 select 100.00% 17 / 17 100.00% 1 / 1 2
 selectColumn 100.00% 31 / 31 100.00% 1 / 1 9
 aggregateRows 100.00% 14 / 14 100.00% 1 / 1 7
 cursor 100.00% 1 / 1 100.00% 1 / 1 1
 insert 100.00% 11 / 11 100.00% 1 / 1 2
 insertGetId 100.00% 2 / 2 100.00% 1 / 1 1
 update 100.00% 14 / 14 100.00% 1 / 1 4
 delete 70.00% 7 / 10 0.00% 0 / 1 2.11
 applyWheres 100.00% 7 / 7 100.00% 1 / 1 2
 matchesWheres 100.00% 8 / 8 100.00% 1 / 1 6
 matchesWhere 100.00% 10 / 10 100.00% 1 / 1 10
 matchesNull 100.00% 2 / 2 100.00% 1 / 1 3
 matchesBetween 100.00% 2 / 2 100.00% 1 / 1 3
 matchesBasic 100.00% 13 / 13 100.00% 1 / 1 14
 valuesEqual 100.00% 7 / 7 100.00% 1 / 1 9
 normalizeCell 100.00% 3 / 3 100.00% 1 / 1 3
 valuesIn 100.00% 4 / 4 100.00% 1 / 1 3
 like 100.00% 7 / 7 100.00% 1 / 1 5
 applyOrders 90.90% 10 / 11 0.00% 0 / 1 4.01
 compareCells 100.00% 3 / 3 100.00% 1 / 1 3
 applyLimit 100.00% 3 / 3 100.00% 1 / 1 2
 splitColumns 88.23% 15 / 17 0.00% 0 / 1 6.06
 project 100.00% 8 / 8 100.00% 1 / 1 4
 computeAggregates 100.00% 15 / 15 100.00% 1 / 1 13
 assertWritable 100.00% 2 / 2 100.00% 1 / 1 2
 isStale 100.00% 1 / 1 100.00% 1 / 1 1
 markStale 100.00% 1 / 1 100.00% 1 / 1 1
 matchesAll 100.00% 1 / 1 100.00% 1 / 1 1
 readRows 100.00% 3 / 3 100.00% 1 / 1 1
 readRowsUnlocked 100.00% 12 / 12 100.00% 1 / 1 4
 releaseIfHeld 100.00% 2 / 2 100.00% 1 / 1 2
 acquireLock 62.50% 5 / 8 0.00% 0 / 1 4.84
 lockPath 100.00% 1 / 1 100.00% 1 / 1 1
 releaseLock 100.00% 2 / 2 100.00% 1 / 1 1
 writeRows 82.50% 33 / 40 0.00% 0 / 1 17.37
 canonicalColumns 100.00% 6 / 6 100.00% 1 / 1 4
 neutralizeFormula 85.71% 6 / 7 0.00% 0 / 1 5.07
 [BlueprintAU\Radiant\Database\Concerns\NormalizesInsertRows] normalizeInsertRows 100.00% 3 / 3 100.00% 1 / 1 4
 [BlueprintAU\Radiant\Database\Concerns\NormalizesInsertRows] assertUniformInsertRows 100.00% 11 / 11 100.00% 1 / 1 4
32final class CsvConnection implements ConnectionInterface
33{
34    use NormalizesInsertRows;
35
36    /** Lock mode for {@see acquireLock()}: shared (reads). */
37    private const LOCK_SHARED = false;
38
39    /**
40     * @param  string  $filePath  The CSV file to read from and write to.
41     * @param  bool  $readOnly  When true, write operations throw instead of modifying the file.
42     */
43    public function __construct(
44        protected string $filePath,
45        protected bool $readOnly = false,
46    ) {}
47
48    /**
49     * Start a fluent query against a table, bound to this connection.
50     *
51     * @param  string  $identifier
52     * @return QueryBuilder
53     */
54    #[Override]
55    public function table(string $identifier): QueryBuilder
56    {
57        return new QueryBuilder($this, $identifier);
58    }
59
60    /**
61     * Run the query and return the matching rows.
62     *
63     * @param  QueryBuilder  $query
64     * @return Collection<int,\stdClass>
65     * @throws UnsupportedFeatureException
66     */
67    #[Override]
68    public function select(QueryBuilder $query): Collection
69    {
70        $query->assertSupports(
71            SqlFeature::Aggregates,
72        );
73
74        $rows = $this->applyWheres($query, $this->readRows());
75        $rows = $this->applyOrders($query, $rows);
76
77        // Split the requested columns into plain fields and typed
78        // Aggregate declarations.
79        [$fields, $aggregates] = $this->splitColumns($query->getColumns());
80
81        // No aggregates â†’ apply limit/offset to the raw rows, then project.
82        if ($aggregates === []) {
83            $rows = $this->applyLimit($query, $rows);
84            return Collection::make(array_map(
85                fn (array $row) => (object) $this->project($row, $fields),
86                $rows,
87            ));
88        }
89
90        // Aggregates present â†’ group by the group columns (if any) and
91        // compute one row per group, carrying the group field values
92        // alongside the aggregates â€” the same shape the SQL backend
93        // produces. LIMIT/OFFSET apply *after* aggregation, matching SQL.
94        $out = $this->aggregateRows($query, $rows, $aggregates);
95
96        return Collection::make(array_map(
97            fn (array $row) => (object) $row,
98            $this->applyLimit($query, $out),
99        ));
100    }
101
102    /**
103     * Run the query and return the first selected column's values.
104     *
105     * @param  QueryBuilder  $query
106     * @return Collection<int, mixed>
107     * @throws UnsupportedFeatureException
108     */
109    #[Override]
110    public function selectColumn(QueryBuilder $query): Collection
111    {
112        $query->assertSupports(
113            SqlFeature::Aggregates,
114        );
115
116        $rows = $this->applyWheres($query, $this->readRows());
117        $rows = $this->applyOrders($query, $rows);
118
119        [$fields, $aggregates] = $this->splitColumns($query->getColumns());
120
121        if ($aggregates === []) {
122            $rows = $this->applyLimit($query, $rows);
123            $first = $fields[0] ?? null;
124
125            if ($first === null || $first === '*') {
126                throw new \InvalidArgumentException(
127                    'selectColumn() requires a single named column; got '
128                    . ($first === null ? 'an empty select list.' : "a wildcard select [{$first}]."),
129                );
130            }
131
132            $out = [];
133            foreach ($rows as $row) {
134                if (!array_key_exists($first, $row)) {
135                    throw new \InvalidArgumentException("Unknown column [{$first}] on CSV connection.");
136                }
137                $out[] = $row[$first];
138            }
139            return Collection::make($out);
140        }
141
142        // Aggregates: the first selected column IS the aggregate (the
143        // builders' scalar reads select exactly one). Run the same
144        // group â†’ compute â†’ order pipeline as select(), then read the
145        // aggregate's alias positionally.
146        $out = $this->aggregateRows($query, $rows, $aggregates);
147        $out = $this->applyLimit($query, $out);
148
149        $first = array_key_first($out[0] ?? []);
150        if ($first === null) {
151            throw new \InvalidArgumentException(
152                'selectColumn() requires a single named column; got an empty select list.',
153            );
154        }
155
156        $values = [];
157        foreach ($out as $row) {
158            $values[] = $row[$first];
159        }
160        return Collection::make($values);
161    }
162
163    /**
164     * Group the filtered rows and compute the aggregates â€” the shared
165     * aggregate pipeline behind {@see select()} and {@see selectColumn()}.
166     *
167     * @param  QueryBuilder  $query
168     * @param  list<array<string,mixed>>  $rows
169     * @param  array<string, array{0: string, string|Expression}>  $aggregates  Alias â†’ [function, column].
170     * @return list<array<string,mixed>>
171     */
172    private function aggregateRows(QueryBuilder $query, array $rows, array $aggregates): array
173    {
174        $groups = $query->getGroups();
175        $buckets = [];
176        foreach ($rows as $row) {
177            $key = $groups === [] ? '' : implode("\0", array_map(fn (string $c) => $row[$c] ?? null, $groups));
178            $buckets[$key][] = $row;
179        }
180
181        $out = [];
182        foreach ($buckets as $bucket) {
183            $computed = $this->computeAggregates($bucket, $aggregates);
184            foreach ($groups as $c) {
185                $computed[$c] = $bucket[0][$c] ?? null;
186            }
187            $out[] = $computed;
188        }
189
190        if ($out === [] && $groups === []) {
191            $out[] = $this->computeAggregates([], $aggregates);
192        }
193
194        return $this->applyOrders($query, $out);
195    }
196
197    /**
198     * Run the query and yield each matching row as it arrives.
199     *
200     * @param  QueryBuilder  $query
201     * @return \Generator<int,\stdClass>
202     * @throws UnsupportedFeatureException
203     */
204    #[Override]
205    public function cursor(QueryBuilder $query): \Generator
206    {
207        yield from $this->select($query);
208    }
209
210    /**
211     * Insert one or more rows into the file.
212     *
213     * @param  QueryBuilder  $query
214     * @param  array<string,mixed>|list<array<string,mixed>>  $values
215     * @return int
216     */
217    #[Override]
218    public function insert(QueryBuilder $query, array $values): int
219    {
220        $this->assertWritable();
221        // One exclusive lock spans the whole read-modify-write: no other
222        // process can read between our read and our write, so no lost
223        // updates. writeRows() reuses (and releases) the lock.
224        $lock = $this->acquireLock();
225        try {
226            $rows = $this->readRowsUnlocked();
227            $normalized = $this->normalizeInsertRows($values);
228
229            // Fail fast on ragged rows BEFORE the read-modify-write â€” a
230            // throw here releases the lock without touching the file.
231            $this->assertUniformInsertRows($normalized);
232
233            array_push($rows, ...$normalized);
234            $this->writeRows($rows, $lock);
235        } catch (\Throwable $e) {
236            $this->releaseIfHeld($lock);
237            throw $e;
238        }
239        return count($normalized);
240    }
241
242    /**
243     * Insert a single row and return its generated id.
244     *
245     * CSV has no auto-increment id, so always returns null.
246     *
247     * @param  QueryBuilder  $query
248     * @param  array<string,mixed>  $values
249     * @return string|int|null
250     */
251    #[Override]
252    public function insertGetId(QueryBuilder $query, array $values): string|int|null
253    {
254        $this->insert($query, $values);
255        return null;
256    }
257
258    /**
259     * Update the rows matching the query's conditions.
260     *
261     * @param  QueryBuilder  $query
262     * @param  array<string,mixed>  $values
263     * @return int
264     */
265    #[Override]
266    public function update(QueryBuilder $query, array $values): int
267    {
268        $this->assertWritable();
269        // Same read-modify-write lock discipline as insert().
270        $lock = $this->acquireLock();
271        try {
272            $rows = $this->readRowsUnlocked();
273            $affected = 0;
274            foreach ($rows as &$row) {
275                if ($this->matchesAll($query, $row)) {
276                    $row = array_merge($row, $values);
277                    $affected++;
278                }
279            }
280            unset($row);
281            $this->writeRows($rows, $lock);
282        } catch (\Throwable $e) {
283            $this->releaseIfHeld($lock);
284            throw $e;
285        }
286        return $affected;
287    }
288
289    /**
290     * Delete the rows matching the query's conditions.
291     *
292     * @param  QueryBuilder  $query
293     * @return int
294     */
295    #[Override]
296    public function delete(QueryBuilder $query): int
297    {
298        $this->assertWritable();
299        // Same read-modify-write lock discipline as insert().
300        $lock = $this->acquireLock();
301        try {
302            $rows = $this->readRowsUnlocked();
303            $kept = array_filter($rows, fn (array $row) => !$this->matchesAll($query, $row));
304            $affected = count($rows) - count($kept);
305            $this->writeRows(array_values($kept), $lock);
306        } catch (\Throwable $e) {
307            $this->releaseIfHeld($lock);
308            throw $e;
309        }
310        return $affected;
311    }
312
313    // ---- Pipeline helpers ----
314
315    /**
316     * Apply the query's where clauses to the rows, honoring boolean
317     * connectors and nested groups.
318     *
319     * @param  QueryBuilder  $query
320     * @param  list<array<string,mixed>>  $rows
321     * @return list<array<string,mixed>>
322     * @throws UnsupportedFeatureException
323     */
324    private function applyWheres(QueryBuilder $query, array $rows): array
325    {
326        $wheres = $query->getWheres();
327        if ($wheres === []) {
328            return $rows;
329        }
330
331        return array_values(array_filter(
332            $rows,
333            fn (array $row) => $this->matchesWheres($wheres, $row),
334        ));
335    }
336
337    /**
338     * Evaluate a list of where clauses against a row, honoring the boolean
339     * connectors between them.
340     *
341     * @param  list<array<string,mixed>>  $wheres
342     * @param  array<string,mixed>  $row
343     * @return bool
344     * @throws UnsupportedFeatureException
345     */
346    private function matchesWheres(array $wheres, array $row): bool
347    {
348        // An empty constraint list matches everything (SQL semantics: an
349        // UPDATE/DELETE with no WHERE affects every row). The builder
350        // rejects empty *nested* groups at declaration time, so this guard
351        // only ever fires for the top-level no-clause case.
352        if ($wheres === []) {
353            return true;
354        }
355        $result = $this->matchesWhere($wheres[0], $row);
356        for ($i = 1, $count = count($wheres); $i < $count; $i++) {
357            $matches = $this->matchesWhere($wheres[$i], $row);
358            $boolean = $wheres[$i]['boolean'] ?? WhereBoolean::And;
359            $result = $boolean === WhereBoolean::Or ? $result || $matches : $result && $matches;
360        }
361        return $result;
362    }
363
364    /**
365     * Evaluate a single where clause against a row.
366     *
367     * @param  array<string,mixed>  $where
368     * @param  array<string,mixed>  $row
369     * @return bool
370     * @throws UnsupportedFeatureException
371     */
372    private function matchesWhere(array $where, array $row): bool
373    {
374        return match ($where['type']) {
375            WhereType::Nested => $this->matchesWheres($where['group']->wheres, $row),
376            WhereType::Raw => throw new UnsupportedFeatureException('This connection does not support raw where clauses.'),
377            WhereType::Column => throw new UnsupportedFeatureException('This connection does not support column-to-column where clauses.'),
378            WhereType::Exists => throw new UnsupportedFeatureException('This connection does not support exists where clauses.'),
379            WhereType::InSub => throw new UnsupportedFeatureException('This connection does not support subquery IN where clauses.'),
380            WhereType::Null => $this->matchesNull($row[$where['column']] ?? null, $where['operator']),
381            WhereType::Between => $this->matchesBetween($row[$where['column']] ?? null, $where['operator'], $where['value']),
382            WhereType::Basic => $this->matchesBasic($row[$where['column']] ?? null, $where['operator'], $where['value']),
383            default => throw new \LogicException('Unknown where type on a CSV connection: ' . get_debug_type($where['type'])),
384        };
385    }
386
387    /**
388     * Evaluate an IS NULL / IS NOT NULL clause.
389     *
390     * An empty cell counts as null â€” CSV has no null representation, so
391     * treating `''` as null is the only way IS NULL semantics survive a
392     * write/read round trip.
393     *
394     * @param  mixed  $value
395     * @param  WhereOperator  $operator
396     * @return bool
397     */
398    private function matchesNull(mixed $value, WhereOperator $operator): bool
399    {
400        $isNull = $value === null || $value === '';
401        return $operator === WhereOperator::NotNull ? !$isNull : $isNull;
402    }
403
404    /**
405     * Evaluate a BETWEEN / NOT BETWEEN clause.
406     *
407     * @param  mixed  $value
408     * @param  WhereOperator  $operator
409     * @param  array{0: mixed, 1: mixed}  $range
410     * @return bool
411     */
412    private function matchesBetween(mixed $value, WhereOperator $operator, array $range): bool
413    {
414        $between = $value >= $range[0] && $value <= $range[1];
415        return $operator === WhereOperator::NotBetween ? !$between : $between;
416    }
417
418    /**
419     * Evaluate a basic comparison (also handles IN / NOT IN, LIKE / NOT LIKE).
420     *
421     * @param  mixed  $value
422     * @param  WhereOperator  $operator
423     * @param  mixed  $operand
424     * @return bool
425     * @throws UnsupportedFeatureException
426     */
427    private function matchesBasic(mixed $value, WhereOperator $operator, mixed $operand): bool
428    {
429        return match ($operator) {
430            WhereOperator::Eq => $this->valuesEqual($value, $operand),
431            WhereOperator::NotEq => !$this->valuesEqual($value, $operand),
432            WhereOperator::Lt => $value < $operand,
433            WhereOperator::LtEq => $value <= $operand,
434            WhereOperator::Gt => $value > $operand,
435            WhereOperator::GtEq => $value >= $operand,
436            WhereOperator::Like => is_string($value) && $this->like($value, (string) $operand),
437            WhereOperator::NotLike => !(is_string($value) && $this->like($value, (string) $operand)),
438            WhereOperator::In => $this->valuesIn($value, (array) $operand),
439            WhereOperator::NotIn => !$this->valuesIn($value, (array) $operand),
440            default => throw new UnsupportedFeatureException(
441                'This connection does not support the ' . $operator->value . ' operator.',
442            ),
443        };
444    }
445
446    /**
447     * The canonical CSV comparator.
448     *
449     * Numeric integer strings collapse to int on both sides, everything
450     * else compares strictly â€” `'5'` matches `5`, `'0e1'` does not match
451     * `0`.
452     *
453     * @param  mixed  $value
454     * @param  mixed  $operand
455     * @return bool
456     */
457    private function valuesEqual(mixed $value, mixed $operand): bool
458    {
459        if ($value === null || $operand === null) {
460            return $value === $operand;
461        }
462
463        if (is_array($value) || is_array($operand)
464            || is_object($value) || is_object($operand)
465            || is_bool($value) || is_bool($operand)) {
466            // Non-scalar or boolean operands have no CSV-cell meaning; a
467            // strict identity check is the honest answer (and never the
468            // type-juggling match `==` would produce).
469            return $value === $operand;
470        }
471
472        return $this->normalizeCell($value) === $this->normalizeCell($operand);
473    }
474
475    /**
476     * Normalize a scalar for strict comparison â€” integer numeric strings
477     * collapse to int, everything else passes through.
478     *
479     * @param  mixed  $value
480     * @return mixed
481     */
482    private function normalizeCell(mixed $value): mixed
483    {
484        if (is_string($value) && preg_match('/^-?\d+$/', $value) === 1) {
485            return (int) $value;
486        }
487
488        return $value;
489    }
490
491    /**
492     * Set membership through the same canonical comparator as equality â€”
493     * `IN` must never be stricter than `=` on the same backend.
494     *
495     * @param  mixed  $value
496     * @param  array<mixed>  $operands
497     * @return bool
498     */
499    private function valuesIn(mixed $value, array $operands): bool
500    {
501        foreach ($operands as $operand) {
502            if ($this->valuesEqual($value, $operand)) {
503                return true;
504            }
505        }
506        return false;
507    }
508
509    /**
510     * SQL LIKE semantics for the few basic operators that need it.
511     *
512     * Matching is case-insensitive (a documented divergence from the
513     * case-sensitive LIKE of Postgres/SQLite).
514     *
515     * @param  string  $value
516     * @param  string  $pattern
517     * @return bool
518     */
519    private function like(string $value, string $pattern): bool
520    {
521        $regex = '';
522        foreach (mb_str_split($pattern) as $char) {
523            $regex .= match ($char) {
524                '%' => '.*',
525                '_' => '.',
526                default => preg_quote($char, '~'),
527            };
528        }
529        return preg_match("~^{$regex}$~is", $value) === 1;
530    }
531
532    /**
533     * Apply the query's order-by clauses to the rows.
534     *
535     * @param  QueryBuilder  $query
536     * @param  list<array<string,mixed>>  $rows
537     * @return list<array<string,mixed>>
538     */
539    private function applyOrders(QueryBuilder $query, array $rows): array
540    {
541        $orders = $query->getOrders();
542        // Reverse so the last-listed order wins as the primary sort key.
543        foreach (array_reverse($orders) as $order) {
544            $column = $order['column'];
545            if ($column instanceof Expression) {
546                throw new UnsupportedFeatureException('This connection does not support raw order-by expressions.');
547            }
548            $direction = $order['direction'] === SortDirection::Desc ? -1 : 1;
549            usort(
550                $rows,
551                fn (array $a, array $b) => $direction * $this->compareCells($a[$column] ?? null, $b[$column] ?? null),
552            );
553        }
554        return $rows;
555    }
556
557    /**
558     * Compare two CSV cells for ordering.
559     *
560     * When both cells are numeric, compare as numbers; otherwise compare
561     * as strings.
562     *
563     * @param  mixed  $a
564     * @param  mixed  $b
565     * @return int
566     */
567    private function compareCells(mixed $a, mixed $b): int
568    {
569        if (is_numeric($a) && is_numeric($b)) {
570            return (+$a) <=> (+$b);
571        }
572        return ($a ?? '') <=> ($b ?? '');
573    }
574
575    /**
576     * Apply the limit and offset.
577     *
578     * @param  QueryBuilder  $query
579     * @param  list<array<string,mixed>>  $rows
580     * @return list<array<string,mixed>>
581     */
582    private function applyLimit(QueryBuilder $query, array $rows): array
583    {
584        $limit = $query->getLimit();
585        $offset = $query->getOffset() ?? 0;
586        return $limit === null ? $rows : array_slice($rows, $offset, $limit);
587    }
588
589    /**
590     * Split the requested columns into plain fields and aggregates.
591     *
592     * @param  list<string|Expression|Aggregate|SubquerySelect>  $columns
593     * @return array{0: list<string>, 1: array<string, array{0: string, string|Expression}>}
594     * @throws UnsupportedFeatureException
595     */
596    private function splitColumns(array $columns): array
597    {
598        $fields = [];
599        $aggregates = [];
600        foreach ($columns as $column) {
601            if ($column instanceof Expression) {
602                throw new UnsupportedFeatureException('This connection does not support raw select expressions.');
603            }
604            if ($column instanceof SubquerySelect) {
605                throw new UnsupportedFeatureException('This connection does not support subquery select columns.');
606            }
607            if ($column instanceof Aggregate) {
608                // The read-back key: the alias when given, else the derived
609                // call text (`count(*)`, `sum(age)`) â€” the same value SQL
610                // returns for an aliased aggregate and the same shape the
611                // ordering path matches against.
612                $columnKey = $column->column instanceof Expression
613                    ? $column->column->value
614                    : $column->column;
615                $aggregates[$column->alias ?? "{$column->function}({$columnKey})"] = [
616                    $column->function,
617                    $column->column,
618                ];
619            } else {
620                $fields[] = $column;
621            }
622        }
623        return [$fields, $aggregates];
624    }
625
626    /**
627     * Project a row to only the requested plain fields.
628     *
629     * @param  array<string,mixed>  $row
630     * @param  list<string>  $fields
631     * @return array<string,mixed>
632     * @throws \InvalidArgumentException
633     */
634    private function project(array $row, array $fields): array
635    {
636        if (in_array('*', $fields, true)) {
637            return $row;
638        }
639        $out = [];
640        foreach ($fields as $field) {
641            if (!array_key_exists($field, $row)) {
642                throw new \InvalidArgumentException("Unknown column [{$field}] on CSV connection.");
643            }
644            $out[$field] = $row[$field];
645        }
646        return $out;
647    }
648
649    /**
650     * Compute aggregate functions over a group of rows.
651     *
652     * @param  list<array<string,mixed>>  $rows
653     * @param  array<string, array{0: string, string|Expression}>  $aggregates  Alias â†’ [function, column].
654     * @return array<string, mixed>
655     * @throws \InvalidArgumentException
656     * @throws UnsupportedFeatureException
657     */
658    private function computeAggregates(array $rows, array $aggregates): array
659    {
660        $result = [];
661        foreach ($aggregates as $alias => [$function, $column]) {
662            if ($column instanceof Expression) {
663                throw new UnsupportedFeatureException(
664                    'This connection cannot compute an aggregate over a raw Expression argument.',
665                );
666            }
667            $values = $column === '*' ? $rows : array_column($rows, $column);
668            $result[$alias] = match ($function) {
669                'count' => count($values),
670                'max' => $values === [] ? null : max($values),
671                'min' => $values === [] ? null : min($values),
672                'sum' => array_sum($values),
673                'avg' => count($values) ? array_sum($values) / count($values) : null,
674                default => throw new \InvalidArgumentException("Unsupported aggregate [{$function}] on a CSV connection."),
675            };
676        }
677        return $result;
678    }
679
680    // ---- Row access ----
681
682    /**
683     * Fail fast when the connection is read-only.
684     *
685     * @throws UnsupportedFeatureException
686     */
687    private function assertWritable(): void
688    {
689        if ($this->readOnly) {
690            throw new UnsupportedFeatureException('This CSV connection is read-only.');
691        }
692    }
693
694    /**
695     * A file-backed connection has no transport to lose â€” it is never stale.
696     *
697     * @return bool
698     */
699    #[Override]
700    public function isStale(): bool
701    {
702        return false;
703    }
704
705    /**
706     * A no-op for the CSV backend â€” see {@see isStale()}.
707     */
708    #[Override]
709    public function markStale(): void
710    {
711        // Nothing to lose: the file is opened per operation.
712    }
713
714    /**
715     * Whether a row matches every where clause of the query.
716     *
717     * @param  QueryBuilder  $query
718     * @param  array<string,mixed>  $row
719     * @return bool
720     */
721    private function matchesAll(QueryBuilder $query, array $row): bool
722    {
723        return $this->matchesWheres($query->getWheres(), $row);
724    }
725
726    /**
727     * Read the CSV file into an array of associative rows, under a shared
728     * lock that is released before returning.
729     *
730     * @return list<array<string,mixed>>
731     * @throws \RuntimeException
732     */
733    private function readRows(): array
734    {
735        $lock = $this->acquireLock(self::LOCK_SHARED);
736        try {
737            return $this->readRowsUnlocked();
738        } finally {
739            $this->releaseLock($lock);
740        }
741    }
742
743    /**
744     * Read rows from the CSV file without taking any lock.
745     *
746     * The data handle is opened and closed here â€” it must never be confused
747     * with the lock handle, because the data file's inode is replaced by
748     * rename() on every write while the lock file's is stable.
749     *
750     * @return list<array<string,mixed>>
751     */
752    private function readRowsUnlocked(): array
753    {
754        $handle = fopen($this->filePath, 'r');
755        if ($handle === false) {
756            throw new \RuntimeException("Could not open CSV file [{$this->filePath}].");
757        }
758        try {
759            $rows = [];
760            $header = fgetcsv($handle, escape: '');
761            if ($header === false) {
762                return [];
763            }
764            $header = array_map(strval(...), $header);
765            while (($line = fgetcsv($handle, escape: '')) !== false) {
766                $rows[] = array_combine($header, array_map(strval(...), $line));
767            }
768            return $rows;
769        } finally {
770            fclose($handle);
771        }
772    }
773
774    /**
775     * Best-effort lock release on a failure path.
776     *
777     * @param  resource|null  $lock
778     */
779    private function releaseIfHeld($lock): void
780    {
781        if (is_resource($lock)) {
782            $this->releaseLock($lock);
783        }
784    }
785
786    /**
787     * Acquire an advisory lock on the CSV file's sidecar lock file.
788     *
789     * The lock must not be taken on the CSV file itself: writes go through
790     * temp-file + rename(), which replaces the CSV's inode. The sidecar's
791     * inode never changes, so every cooperating process serializes on the
792     * same object for the file's whole lifetime.
793     *
794     * @param  bool  $exclusive  True for LOCK_EX (writes), false for LOCK_SH (reads).
795     * @return resource
796     * @throws \RuntimeException
797     */
798    private function acquireLock(bool $exclusive = true)
799    {
800        $lockPath = $this->lockPath();
801        $lock = fopen($lockPath, 'c');
802        if ($lock === false) {
803            throw new \RuntimeException("Could not open CSV lock file [{$lockPath}].");
804        }
805        if (!flock($lock, $exclusive ? LOCK_EX : LOCK_SH)) {
806            fclose($lock);
807            throw new \RuntimeException("Could not lock CSV file [{$this->filePath}].");
808        }
809        return $lock;
810    }
811
812    /**
813     * The sidecar lock file path for the CSV file.
814     *
815     * @return string
816     */
817    private function lockPath(): string
818    {
819        return $this->filePath . '.lock';
820    }
821
822    /**
823     * Release the sidecar lock and close its handle.
824     *
825     * @param  resource  $lock
826     */
827    private function releaseLock($lock): void
828    {
829        flock($lock, LOCK_UN);
830        fclose($lock);
831    }
832
833    /**
834     * Write rows back to the file atomically via temp-file + rename().
835     *
836     * Every value is passed through {@see neutralizeFormula()} so a value
837     * that begins with `=`, `+`, `-`, `@`, tab, or CR cannot execute as a
838     * spreadsheet formula when the file is opened in Excel/Sheets. Rows are
839     * aligned to the canonical column order â€” the header row.
840     *
841     * @param  list<array<string,mixed>>  $rows
842     * @param  resource|null  $lock  An existing locked sidecar handle to reuse, or null to acquire the lock for this write.
843     * @throws \RuntimeException
844     */
845    private function writeRows(array $rows, $lock = null): void
846    {
847        $ownsLock = $lock === null;
848        if ($ownsLock) {
849            $lock = $this->acquireLock();
850        }
851
852        // A UNIQUE temp path per write: a shared fixed temp name lets a
853        // second writer's fopen('w') truncate the first writer's in-flight
854        // temp (silent lost updates across processes). Process id + random
855        // suffix; same directory so rename() stays same-filesystem atomic.
856        $tempPath = sprintf(
857            '%s.radiant-%s-%s.tmp',
858            $this->filePath,
859            (string) (getmypid() ?: 'unknown'),
860            bin2hex(random_bytes(6)),
861        );
862        $temp = fopen($tempPath, 'w');
863        if ($temp === false) {
864            if ($ownsLock) {
865                $this->releaseLock($lock);
866            }
867            throw new \RuntimeException("Could not write CSV file [{$tempPath}].");
868        }
869
870        // try/finally guarantees the temp file cannot outlive this call â€”
871        // a TypeError from fputcsv (or any other unwinding failure) would
872        // otherwise orphan a partial, data-bearing temp file per failure.
873        try {
874            $columns = $this->canonicalColumns($rows);
875            // The header row is neutralized too â€” a hostile column name is
876            // just as able to execute as a spreadsheet formula as a cell.
877            $written = fputcsv($temp, array_map(
878                fn (string $column) => $this->neutralizeFormula($column),
879                $columns,
880            ), escape: '') !== false;
881            foreach ($rows as $row) {
882                $aligned = [];
883                foreach ($columns as $column) {
884                    $value = $row[$column] ?? null;
885                    $aligned[] = is_string($value) ? $this->neutralizeFormula($value) : $value;
886                }
887                if (fputcsv($temp, $aligned, escape: '') === false) {
888                    $written = false;
889                    break;
890                }
891            }
892
893            if (!fclose($temp) || !$written) {
894                throw new \RuntimeException("Could not write CSV file [{$tempPath}].");
895            }
896
897            // rename() replaces the original â€” carry its permissions over so
898            // a 0600 file is not demoted to umask defaults on every write.
899            $originalPerms = @fileperms($this->filePath);
900            if ($originalPerms !== false) {
901                @chmod($tempPath, $originalPerms & 0o777);
902            }
903
904            if (!rename($tempPath, $this->filePath)) {
905                throw new \RuntimeException("Could not replace CSV file [{$this->filePath}].");
906            }
907        } finally {
908            // After a successful rename the temp no longer exists; after any
909            // failure it does â€” unlink it best-effort so no partial copy of
910            // the data is ever left behind.
911            if (is_resource($temp)) {
912                fclose($temp);
913            }
914            if (file_exists($tempPath)) {
915                @unlink($tempPath);
916            }
917            if ($ownsLock) {
918                $this->releaseLock($lock);
919            }
920        }
921    }
922
923    /**
924     * The canonical column order for a write: the union of the header row's
925     * keys and every row's keys, in first-seen order.
926     *
927     * @param  list<array<string,mixed>>  $rows
928     * @return list<string>
929     */
930    private function canonicalColumns(array $rows): array
931    {
932        $columns = [];
933        foreach ($rows as $row) {
934            foreach (array_keys($row) as $column) {
935                if (!in_array($column, $columns, true)) {
936                    $columns[] = $column;
937                }
938            }
939        }
940        return $columns;
941    }
942
943    /**
944     * Neutralize a value that a spreadsheet would evaluate as a formula.
945     *
946     * @param  string  $value
947     * @return string
948     */
949    private function neutralizeFormula(string $value): string
950    {
951        $trimmed = ltrim($value, " \t\r\n\0\v\f\xC2\xA0\xE2\x80\x8B\xEF\xBB\xBF");
952        $first = $trimmed === '' ? '' : $trimmed[0];
953        if (in_array($first, ['=', '@', '|', "\t", "\r"], true)) {
954            return "'" . $trimmed;
955        }
956        if (in_array($first, ['+', '-'], true) && !is_numeric($trimmed)) {
957            return "'" . $trimmed;
958        }
959        return $value;
960    }
961}

From BlueprintAU\Radiant\Database\Concerns\NormalizesInsertRows

15trait NormalizesInsertRows
16{
17    /**
18     * Normalize a single row or a list of rows into a list of rows.
19     *
20     * @param  array<string,mixed>|list<array<string,mixed>>  $values
21     * @return list<array<string,mixed>>
22     */
23    protected function normalizeInsertRows(array $values): array
24    {
25        // A list of rows: [[...], [...]] â€” each element is an associative row.
26        if (array_is_list($values) && isset($values[0]) && is_array($values[0])) {
27            /** @var list<array<string, mixed>> $values */
28            return $values;
29        }
30
31        /** @var array<string, mixed> $values */
32        return [$values];
33    }
34
35    /**
36     * Assert every row in a bulk insert carries the SAME column set.
37     *
38     * A multi-row INSERT compiles ONE column list and ONE placeholder
39     * group per row â€” rows of differing arity produce a malformed
40     * statement or a placeholder/binding mismatch. Padding a missing
41     * column with NULL would silently write NULL into a nullable column
42     * the caller never named, so ragged rows fail fast instead.
43     *
44     * @param  list<array<string,mixed>>  $rows
45     * @return void
46     *
47     * @throws \InvalidArgumentException
48     */
49    protected function assertUniformInsertRows(array $rows): void
50    {
51        if (count($rows) < 2) {
52            return;
53        }
54
55        $expected = array_keys($rows[0]);
56
57        foreach (array_slice($rows, 1) as $i => $row) {
58            if (array_keys($row) !== $expected) {
59                throw new \InvalidArgumentException(
60                    'A bulk insert requires every row to carry the same columns; row '
61                        . ($i + 1) . ' differs from row 0. Split the call or give every '
62                        . 'row the same column set (an absent column would otherwise '
63                        . 'silently write NULL).'
64                );
65            }
66        }
67    }
68}