| File: | blib/lib/Database/Join.pm |
| Coverage: | 95.5% |
| line | stmt | bran | cond | sub | time | code |
|---|---|---|---|---|---|---|
| 1 | package Database::Join; | |||||
| 2 | ||||||
| 3 | # ABSTRACT: Combined view across two or more Database::Abstraction objects | |||||
| 4 | ||||||
| 5 | 25 25 | 3848949 45 | use 5.010001; | |||
| 6 | 25 25 25 | 43 18 183 | use strict; | |||
| 7 | 25 25 25 | 34 25 496 | use warnings; | |||
| 8 | 25 25 25 | 983 38154 61 | use autodie qw(:all); | |||
| 9 | ||||||
| 10 | 25 25 25 | 81530 22 523 | use Carp qw(croak carp); | |||
| 11 | 25 25 25 | 36 19 233 | use File::Spec; | |||
| 12 | 25 25 25 | 34 23 415 | use List::Util qw(max); | |||
| 13 | 25 25 25 | 2404 361540 232 | use Log::Abstraction; | |||
| 14 | 25 25 25 | 55 20 402 | use Readonly; | |||
| 15 | 25 25 25 | 30 16 309 | use Scalar::Util qw(blessed); | |||
| 16 | 25 25 25 | 1984 35486 246 | use Object::Configure; | |||
| 17 | 25 25 25 | 36 19 362 | use Params::Get qw(get_params); | |||
| 18 | 25 25 25 | 34 18 286 | use Params::Validate::Strict qw(validate_strict); | |||
| 19 | 25 25 25 | 1322 26392 78 | use Sub::Protected; | |||
| 20 | ||||||
| 21 | # Named-pair keys accepted by add_database; kept here so the guard and the | |||||
| 22 | # validate_strict schema cannot silently diverge. | |||||
| 23 | Readonly::Array my @_ADD_DB_KEYS => qw(database join_column filter remove_columns); | |||||
| 24 | ||||||
| 25 | # SQL comparison operators that are safe to interpolate into WHERE clauses. | |||||
| 26 | # Scalar-binding operators: the value is always passed as a bind param (?), so | |||||
| 27 | # no quoting or escaping is required on the value side. LIKE and NOT LIKE are | |||||
| 28 | # safe for the same reason â the pattern is bound, not interpolated. | |||||
| 29 | Readonly::Hash my %SAFE_SQL_OPS => map { $_ => 1 } ('>', '<', '>=', '<=', '!=', '=', 'LIKE', 'NOT LIKE'); | |||||
| 30 | ||||||
| 31 | # List-membership operators whose criterion value is an arrayref. | |||||
| 32 | # Each element is bound as a separate ?, so injection is impossible. | |||||
| 33 | # IN with an empty list â WHERE 1=0 (no rows match, correct SQL semantics). | |||||
| 34 | # NOT IN with an empty list â no WHERE term (all rows match, correct semantics). | |||||
| 35 | Readonly::Hash my %SAFE_LIST_OPS => map { $_ => 1 } ('IN', 'NOT IN'); | |||||
| 36 | ||||||
| 37 | # Nullability operators: no bind parameter. The hashref value is ignored â | |||||
| 38 | # any value (undef, 1, ...) signals intent; only the key selects the operator. | |||||
| 39 | # A bare undef criterion value (col => undef) also generates IS NULL. | |||||
| 40 | Readonly::Hash my %SAFE_NOARG_OPS => map { $_ => 1 } ('IS NULL', 'IS NOT NULL'); | |||||
| 41 | ||||||
| 42 | our $VERSION = '0.008.2'; | |||||
| 43 | ||||||
| 44 | # Package-level cache for threads availability. undef = not yet checked; | |||||
| 45 | # 1 = available; 0 = not available. Checked lazily on the first parallel | |||||
| 46 | # query and never re-evaluated (require is cached in %INC after success). | |||||
| 47 | my $HAS_THREADS; | |||||
| 48 | ||||||
| 49 | # --------------------------------------------------------------------------- | |||||
| 50 | # All user-facing strings route through this dictionary. Supply an i18n | |||||
| 51 | # object with a translate($key, @sprintf_args) method to localise them. | |||||
| 52 | # --------------------------------------------------------------------------- | |||||
| 53 | Readonly::Hash my %MESSAGES => ( | |||||
| 54 | error_no_databases => 'At least one Database::Abstraction object is required', | |||||
| 55 | error_invalid_db => 'databases[%d] does not support the selectall_arrayref/columns interface', | |||||
| 56 | error_join_col_missing => 'join_column "%s" is absent from databases[%d] (%s)', | |||||
| 57 | error_col_conflict => 'Column "%s" exists in multiple databases; use the owning database directly or rename the column', | |||||
| 58 | error_remove_join_col => 'Cannot remove join_column "%s"; it is required for the join', | |||||
| 59 | warn_unknown_column => 'Column "%s" is not present in any configured database; criterion ignored', | |||||
| 60 | error_join_criterion => 'join => criteria are not supported on Database::Join (SQL JOINs target a single DA; use multiple DAs in the databases => [] constructor instead)', | |||||
| 61 | error_query_unsupported => 'query() chained builder is not supported on Database::Join (Database::Abstraction::Query targets a single DA, not the merged view); use selectall_arrayref, selectall_array, fetchrow_hashref, count, or each_row instead', | |||||
| 62 | error_execute_unsupported => 'execute() raw SQL is not supported on Database::Join', | |||||
| 63 | error_unknown_message => 'Unknown message key "%s"', | |||||
| 64 | error_invalid_prefix => 'collision_prefix[%d] must be a plain string, not a reference; passing a reference would leak a heap address into column names', | |||||
| 65 | error_invalid_backend => 'backend must be "array", "sqlite", or "auto"; got "%s"', | |||||
| 66 | error_sqlite_connect => 'Failed to open temporary SQLite database for join backend: %s', | |||||
| 67 | warn_schema_type_mismatch => 'column "%s" has type "%s" in database[%d] but type "%s" in database[%d]; use collision_prefix to preserve both values without silent type coercion', | |||||
| 68 | error_invalid_callback => 'each_row: first argument must be a code reference', | |||||
| 69 | ); | |||||
| 70 | ||||||
| 71 - 732 | =head1 NAME
Database::Join - Read-only combined view across two or more Database::Abstraction objects
=head1 VERSION
Version 0.008.2
=head1 SYNOPSIS
B<Basic two-database join>
use Database::Join;
# Step 1: create each component database the normal way
my $customers = Database::Customers->new(directory => '/data');
my $loyalty = Database::Loyalty->new(directory => '/data');
# Step 2: combine them on the shared key column 'entry'
my $join = Database::Join->new(
databases => [ $customers, $loyalty ],
join_column => 'entry',
);
# Step 3: query exactly as you would a single Database::Abstraction object
my $all_rows = $join->selectall_arrayref();
my $vip_rows = $join->selectall_arrayref(tier => 'gold');
my $one_row = $join->fetchrow_hashref(entry => 'C001');
my $total = $join->count();
my $col_names = $join->columns();
B<Hiding internal columns>
my $join = Database::Join->new(
databases => [ $customers, $loyalty ],
join_column => 'entry',
remove_columns => [ 'internal_id', 'audit_ts' ],
);
# 'internal_id' and 'audit_ts' never appear in results or columns()
B<join_map: when the key column has different names in each database>
# $cities (index 0) has a column called 'statecode' -- matches join_column
# $stnames (index 1) has a column called 'entry' -- different name
my $join = Database::Join->new(
databases => [ $cities, $stnames ],
# index 0 index 1
join_column => 'statecode',
join_map => { 1 => 'entry' }, # index 1 calls its join key 'entry'
);
# All returned rows use 'statecode'; 'entry' is never exposed
my $rows = $join->selectall_arrayref();
B<filters: permanently restrict a database's visible rows>
# Only show orders placed more than 60 days ago, without repeating
# the criterion on every query call.
my $join = Database::Join->new(
databases => [ $customers, $orders ],
join_column => 'entry',
filters => { 1 => { age_days => { '>' => 60 } } },
);
my $rows = $join->selectall_arrayref(); # all old orders
my $vip = $join->selectall_arrayref(tier => 'gold'); # old + gold tier
B<Inner and outer join types>
my $inner = Database::Join->new(
databases => [ $customers, $loyalty ],
join_column => 'entry',
join_type => 'inner', # only keys present in BOTH databases
);
my $outer = Database::Join->new(
databases => [ $customers, $loyalty ],
join_column => 'entry',
join_type => 'outer', # all keys from EITHER database
);
B<Building the view incrementally with add_database>
my $join = Database::Join->new(
databases => [ $customers ],
join_column => 'entry',
);
$join->add_database($loyalty)
->add_database($scores, remove_columns => ['raw_score']);
B<AUTOLOAD column shortcut>
# Returns the 'name' value for entry 'C001' (scalar context)
my $name = $join->name(entry => 'C001');
# Returns all 'tier' values (list context)
my @tiers = $join->tier();
B<SQLite join backend for large datasets>
# 'auto' (default): switches to SQLite automatically above the threshold
my $join = Database::Join->new(
databases => [ $customers, $loyalty ],
join_column => 'entry',
backend => 'auto', # default
max_array_rows => 50_000, # use SQLite when combined rows > 50,000
);
# Always use SQLite -- useful when you know the data is large
my $join = Database::Join->new(
databases => [ $customers, $loyalty ],
join_column => 'entry',
backend => 'sqlite',
tmpdir => '/fast/nvme/tmp', # optional: faster temp disk
);
# Always use the original in-memory path
my $join = Database::Join->new(
databases => [ $customers, $loyalty ],
join_column => 'entry',
backend => 'array',
);
=head1 DESCRIPTION
C<Database::Join> merges two or more L<Database::Abstraction> objects into a
single logical, read-only view. Each component database is queried
independently through its own C<Database::Abstraction> interface. The results
are combined using a shared key column (C<join_column>).
In effect, this means that you can view data from more than one database using
an intuitive, non-SQL interface.
Every storage format that C<Database::Abstraction> supports works as a
component database: CSV, PSV, TSV, SQLite, JSON, XML, XLSX, BerkeleyDB, HTML
URL, JSON URL, or any custom subclass. Component databases may mix formats
within the same join.
The module exposes the same read-only API as C<Database::Abstraction>:
C<selectall_arrayref>, C<selectall_array>, C<fetchrow_hashref>, C<count>,
C<columns>, C<schema>, C<updated>, C<set_logger>, and the AUTOLOAD column
shortcut. Callers do not need to know how many underlying databases are
involved.
Think of it as a virtual database table that is assembled on demand from
several real tables, one per component database.
B<Join backends>
By default (C<backend =E<gt> 'auto'>), C<Database::Join> first checks the
combined source row count. For small datasets (up to C<max_array_rows>,
default 10,000 rows) it merges entirely in Perl memory. For larger datasets it
automatically spills source rows into a temporary SQLite database and executes
a single SQL JOIN there, keeping peak RAM to roughly one times the source data
size instead of three. You can also force either path unconditionally with
C<backend =E<gt> 'sqlite'> or C<backend =E<gt> 'array'>.
=head2 Join semantics
The C<join_type> parameter controls what happens when a particular key value
exists in some component databases but not all:
=over 4
=item C<left> (the default)
All rows from the I<primary> (first) database are returned. Columns from
subsequent databases are included where a matching row is found, and simply
absent from the hashref where there is no match. If you are familiar with
SQL, this is a LEFT OUTER JOIN on the first table.
=item C<inner>
Only rows whose join-column value is present in I<every> component database
are returned. This is equivalent to a SQL INNER JOIN.
=item C<outer>
Every join-column value found in I<any> component database is returned.
Columns from databases that do not have that key value are absent from the
merged row. This is a FULL OUTER JOIN.
=back
B<Important override rule:> whenever you pass a query criterion for a column
that belongs to a secondary database, that database automatically acts as an
inner-join partner for that query only -- regardless of C<join_type>. This
gives WHERE-clause semantics. For example, if you have a LEFT join but query
C<< tier => 'gold' >> on a secondary database, only rows whose secondary entry
has tier = 'gold' are returned (rows with no secondary entry are excluded, just
as a WHERE clause would exclude them).
=head2 Column ownership and routing
At construction time, C<Database::Join> calls C<columns()> on each component
database and builds an internal index that maps every column name to the
database that owns it.
When you pass criteria to a query method, each key-value pair is automatically
routed to the right database. You never need to say which database a column
belongs to.
The C<join_column> is special: criteria on it are broadcast to I<all>
databases so that each database fetches only the relevant rows before the
in-memory merge.
When the same non-join column name exists in more than one database, the
I<last> database in the C<databases> array wins by default: its value
overwrites earlier ones in merged rows. Use C<collision_prefix> to
preserve both values under distinct names instead.
=head1 LIMITATIONS
=over 4
=item Memory usage (array backend)
When C<backend> is C<'array'> (or C<'auto'> and the dataset is small), all
matching rows are fetched into Perl memory. Peak RAM is roughly three times
the source data size. For large datasets use C<backend =E<gt> 'sqlite'>, or
leave C<backend =E<gt> 'auto'> and set C<max_array_rows> appropriately.
=item SQLite backend uses a persistent cache file
When the SQLite path is active, a single C<.db> file with a randomly generated
name (chosen by C<File::Temp> to avoid collisions) is created in C<tmpdir> the
first time a query runs on a given C<Database::Join> object. Source data is
spilled into that file once; subsequent queries against the same object reuse
the file without re-fetching the source data.
The cache is automatically invalidated and rebuilt whenever any source
database's C<updated()> timestamp changes (indicating new data), or when
C<add_database()> is called.
The file is deleted when the C<Database::Join> object is destroyed (typically
when it goes out of scope). At any given moment no more than one such file
exists per object. The directory must be writable and have enough free space
for the full source data (once, not per-query).
=item No chained builder or raw SQL
C<query()> and C<execute()> are not implemented. Use C<selectall_arrayref>
or C<fetchrow_hashref> instead.
=item Single-column equi-join only
Joining on more than one column simultaneously, or on expressions, is not
supported. When the join key has different names in different databases,
use C<join_map> to declare each database's local column name.
=item Sort order
Results are sorted ascending by C<join_column> by default. Pass
C<< sort_by => 'colname' >> (or C<< sort_by => ['colname', 'DESC'] >>)
to any query method to override this. The array path uses string comparison
(C<cmp>); for accurate numeric ordering on large datasets use the SQLite
backend, which sorts natively by type.
=item count() on the array backend fetches all rows
On the array backend, C<count()> executes the full in-memory join and counts
the resulting rows in Perl; no C<COUNT(*)> is pushed to the component
databases. On the SQLite backend, C<count()> executes a C<SELECT COUNT(*)>
SQL query against the cached join tables, avoiding a full row transfer.
=back
=head1 COMMON PITFALLS
=over 4
=item The join_column must exist in every component database
If even one database is missing the join key column, C<new()> (or
C<add_database()>) will C<croak> immediately. Use C<join_map> when the
column has a different local name in some databases.
=item Criteria on a removed column are silently dropped
If you call C<remove_column('tier')> and later query
C<< selectall_arrayref(tier => 'gold') >>, the criterion is ignored (with a
C<carp> warning) and all rows are returned. Always pass criteria before
removing columns, or restructure your code to avoid this.
=item You cannot remove the join column
C<< $join->remove_column($join->join_column) >> will C<croak>. The join key
is required for the merge to work.
=item Left join does not guarantee all columns are populated
Under a LEFT join, rows from the primary database that have no matching row
in a secondary database will be returned with I<no keys> from that secondary
database. Accessing C<< $row->{score} >> on such a row returns C<undef> --
not zero, not an empty string. Always test C<defined $row->{score}> rather
than just C<$row->{score}> when the secondary match is optional.
=item Filters act as inner-join partners
Any database that has a C<filters> entry is promoted to an inner-join partner,
regardless of C<join_type>. A row whose join-key value does not appear in
the filtered database's result is removed from the merged output entirely, not
merely missing its secondary columns. This is intentional but can be
surprising if you expected LEFT join semantics.
=item Criteria-merging replaces scalar filters
When both a base filter and a query criterion target the same column, and both
are operator hashrefs (e.g. C<< { '>' => 60 } >>), the operators are combined
(AND semantics). But if the query criterion is a plain scalar (e.g.
C<< score => 75 >>), it I<replaces> the base filter for that column entirely --
the base filter is ignored for that query.
=item AUTOLOAD sees the full merged join when filters or join_map are active
When either C<filters> or C<join_map> is in effect, the AUTOLOAD shortcut
(C<< $join->columnname(...) >>) runs the full join query rather than
delegating directly to the owning database. This is necessary for correctness
but means the result respects all active filters and join-key translations,
which may differ from what the owning database would return on its own.
=item Duplicate column names: last database wins (unless collision_prefix is set)
When two component databases each have a column called C<notes>, the second
database's value silently overwrites the first in every merged row. Use
C<collision_prefix => { 1 => 'right' }> to publish the second database's
C<notes> as C<right.notes> so both values survive, or use C<remove_columns>
(or C<remove_column>) to drop the unwanted duplicate entirely.
=item Mutating the filters hashref after construction has no effect
C<Database::Join> deep-copies the C<filters> hashref (and any C<filter>
passed to C<add_database>) at the moment of construction. The original hashref
you passed in is never stored. If you later modify it -- for example, to
tighten or loosen a filter criterion -- the joined view is I<not> affected.
Construct a new C<Database::Join> object, or use a component
C<Database::Abstraction> that supports dynamic filter modification.
=item auto mode may always use the array path for some DAs
C<backend =E<gt> 'auto'> counts rows cheaply only when each component database
either implements C<dbi_source()> (SQLite-backed) or directly defines a
C<count()> method in its own package. A DA that merely I<inherits> C<count()>
from C<Database::Abstraction> is treated as uncountable, because the parent
class C<count()> expects a key argument and behaves differently from a
"return total row count" function. In that case C<Database::Join>
conservatively uses the array path for the whole query, even if the dataset is
large. To opt in to the SQLite path for such a DA, either add your own
C<count()> override that returns the total row count, implement C<dbi_source()>,
or use C<backend =E<gt> 'sqlite'> unconditionally.
=item dbi_source() ATTACH is unconditional - query-time criteria go into WHERE
When a component database implements C<dbi_source()>, C<Database::Join> always
uses the zero-copy ATTACH path, I<even when the current query includes criteria
for columns in that database>. The criteria are translated into parameterised
SQL C<WHERE> clauses applied against the ATTACHed table; no row-level copy is
performed. (Prior to 0.005.0 the presence of any query-time criteria would
force a spill; that restriction has been removed.)
=item Broadcast join-column criterion does not force secondaries into inner-join
When a caller passes a join-column criterion (e.g. C<entry =E<gt> 'k1'>),
C<Database::Join> broadcasts it to all component databases so each DA can
filter its fetch to the requested key. Prior to 0.006.0 this broadcast was
incorrectly counted as "having criteria" for secondary databases, causing
C<left> and C<outer> joins to silently behave as C<inner> joins when a
join-column criterion was present. The fix: only non-join-column criteria
(e.g. column filters from the caller or base C<filters =E<gt> {...}>) promote
a secondary to inner-join status. The broadcast itself is now a transparent
key-range selector that does not affect join semantics.
=item LIKE and NOT LIKE work on the SQLite path; other pattern operators do not yet
C<LIKE>, C<NOT LIKE>, C<IN>, and C<NOT IN> are fully supported on the SQLite
backend. C<LIKE>/C<NOT LIKE> take a scalar pattern; C<IN>/C<NOT IN> take an
arrayref of values. All are injection-safe because values are passed as bind
parameters, never interpolated.
# LIKE
my $rows = $join->selectall_arrayref(name => { LIKE => 'A%' });
# IN
my $rows = $join->selectall_arrayref(tier => { IN => ['gold', 'silver'] });
# NOT IN
my $rows = $join->selectall_arrayref(tier => { 'NOT IN' => ['bronze'] });
C<IN> with an empty arrayref matches no rows (SQL semantics: C<IN ()> is
always false). C<NOT IN> with an empty arrayref matches all rows (no
constraint added).
C<IS NULL> and C<IS NOT NULL> are supported on the SQLite backend using an
explicit operator hashref:
my $rows = $join->selectall_arrayref(score => { 'IS NULL' => undef });
my $rows = $join->selectall_arrayref(score => { 'IS NOT NULL' => 1 });
The hashref value is ignored; only the key selects the operator. A bare
C<undef> criterion value (C<< score => undef >>) also generates C<IS NULL>
on the SQLite path. Note: on the in-memory array path, the component DA
may treat an C<undef> criterion value as C<< no filter >> rather than
C<IS NULL>, so use the explicit hashref form for consistent behaviour
across backends.
Any operator not in the supported set is skipped on the SQLite path with a
C<carp> warning. The array path forwards all operators to the component DA
unchanged, which may or may not honour them.
=item Temp file directory must be writable and have free space
The SQLite path creates one temporary C<.db> file per C<Database::Join> object
in C<tmpdir> (default: C<File::Spec-E<gt>tmpdir()>, usually C</tmp> on Unix).
The filename is randomly generated by C<File::Temp> -- you cannot predict it,
only the directory is under your control. The file is created on the first
query and kept alive until the object is destroyed; it is not re-created on
every query call. If the directory is not writable, or the filesystem is full,
the call will C<croak> with C<error_sqlite_connect>. Check permissions and
free space if you see that error.
=item C<join =E<gt>> criteria are not supported
C<Database::Abstraction> accepts a C<join =E<gt> { table =E<gt> ..., on =E<gt> ... }>
key in its criteria hashrefs to express an SQL JOIN within a single table.
C<Database::Join> cannot route this to a component DA meaningfully: the merged
view has no concept of a single underlying table. Passing C<join =E<gt>> to any
query method will C<croak> with C<error_join_criterion>.
B<Fix:> Model the joined table as a separate C<Database::Abstraction> object and
add it to the C<databases =E<gt> []> list. Use C<join_column> or C<join_map> to
specify the shared key.
=back
=head1 METHODS
=head2 new
=head3 SYNOPSIS
my $join = Database::Join->new(
databases => [ $db1, $db2 ],
join_column => 'entry',
join_type => 'left',
join_map => { 1 => 'local_col' },
filters => { 1 => { score => { '>' => 60 } } },
collision_prefix => { 1 => 'right' },
remove_columns => [ 'email', 'internal_id' ],
backend => 'auto', # 'auto' | 'sqlite' | 'array'
max_array_rows => 10_000, # threshold for 'auto' mode
tmpdir => '/tmp', # directory for temp SQLite file
logger => $log,
i18n => $locale,
);
=head3 DESCRIPTION
Constructs and returns a new C<Database::Join> object.
Each element of C<databases> must be an already-instantiated subclass of
C<Database::Abstraction>. The constructor calls C<columns()> on every
database to build an internal column-routing table and verifies that
C<join_column> (or its local alias from C<join_map>) is present in each one.
Columns listed in C<remove_columns> are hidden immediately: they do not appear
in C<columns()>, C<schema()>, or any returned row hashref. This is equivalent
to calling C<remove_column> once per name after construction.
=head3 API SPECIFICATION
=head4 Input
databases => { type => 'arrayref', required => 1 }
# One or more Database::Abstraction subclass objects.
#
# DOMAIN -- EP valid: non-empty arrayref of blessed DA subclasses.
# DOMAIN -- EP invalid: scalar, hashref, or absent => croak.
# DOMAIN -- BVA size: minimum 1 element; no documented upper bound.
# DOMAIN -- BVA elem: each element must pass isa('Database::Abstraction').
join_column => { type => 'string', optional => 1, default => 'entry' }
# The column name shared by all databases (the join key).
#
# DOMAIN -- EP valid: any non-empty string present in every component DA.
# DOMAIN -- EP invalid: column absent from any DA => croak join_col_missing.
# DOMAIN -- BVA: empty string '' is treated as a column name and
# will croak if (as expected) it is absent from every DA.
# DOMAIN -- NOTE: matching is case-sensitive and exact.
join_type => { type => 'string', optional => 1, default => 'left',
enum => ['inner', 'left', 'outer'] }
# Controls which keys appear in the result when not all
# databases share the same key values.
#
# DOMAIN -- EP valid: exactly 'inner', 'left', or 'outer'.
# DOMAIN -- EP invalid: any other string including 'INNER', 'LEFT',
# 'OUTER' (enum check is case-sensitive), 'cross',
# or '' => croak from validate_strict.
join_map => { type => 'hashref', optional => 1 }
# Zero-based database index => local column name.
# See the join_map section for full details.
#
# DOMAIN -- EP valid: hashref values must be plain strings.
# DOMAIN -- EP invalid: reference value (hashref, arrayref, coderef, etc.)
# => croak; the guard prevents heap-address leakage.
# DOMAIN -- BVA: out-of-range keys (beyond the databases array) are
# silently ignored.
filters => { type => 'hashref', optional => 1 }
# Zero-based database index => criteria hashref.
# Permanent row restrictions on individual databases.
# See the filters section for full details.
base_criteria => { type => 'hashref', optional => 1 }
# Column name => value hashref.
# Permanent criteria applied to every query, specified by
# column name rather than database index. Each key must
# be a column that appears in the merged view (or the
# join_column); each value is a plain scalar or an
# operator hashref in the same format as
# selectall_arrayref accepts. Criteria are automatically
# routed to the database that owns each column (same
# routing used by selectall_arrayref).
# Equivalent to filters but more convenient when you know
# the column name but not which database index owns it.
# When a column appears in both base_criteria and filters,
# the filters entry takes precedence.
# See the base_criteria section for full details.
#
# DOMAIN -- EP valid: hashref of column-name => scalar or
# operator-hashref pairs. Unknown column
# names emit warn_unknown_column (carp)
# and are silently dropped.
# DOMAIN -- EP invalid: non-hashref value => croak from
# validate_strict.
# DOMAIN -- BVA: {} empty hashref is a safe no-op.
collision_prefix => { type => 'hashref', optional => 1 }
# Zero-based database index (>0) => prefix string.
# When a secondary database has a column that collides
# with a column already present in the merged view, the
# secondary column is published as "$prefix.$col" instead
# of silently overwriting the earlier value.
# Index 0 entries are silently ignored.
# Omitting this parameter preserves the original
# last-database-wins behaviour.
# See the collision_prefix section for full details.
#
# DOMAIN -- EP valid: absent or {} => last-database-wins (no change).
# DOMAIN -- EP valid: { N => 'prefix' } where N > 0 => colliding
# columns from DB[N] published as "$prefix.$col";
# non-colliding columns from the same DB added plain.
# DOMAIN -- EP note: index-0 entries are silently ignored.
# DOMAIN -- Invariant: join_column is never prefixed regardless of
# collision_prefix configuration.
remove_columns => { type => 'arrayref', optional => 1 }
# Column names to hide from the merged view.
#
# DOMAIN -- EP valid: arrayref of any strings; non-existent columns
# are silently ignored (idempotent).
# DOMAIN -- EP invalid: join_column itself => croak remove_join_col.
# DOMAIN -- BVA: [] empty arrayref is a safe no-op.
backend => { type => 'string', optional => 1, default => 'auto',
enum => ['array', 'sqlite', 'auto'] }
# Controls which join strategy is used.
# 'auto' -- (default) use 'array' when combined source row count
# <= max_array_rows, 'sqlite' otherwise.
# 'sqlite' -- always spill to a temporary SQLite database.
# 'array' -- always use the in-memory merge path.
#
# DOMAIN -- EP valid: 'array', 'sqlite', or 'auto' (case-sensitive).
# DOMAIN -- EP invalid: any other string => croak error_invalid_backend.
# DOMAIN -- Default: 'auto'.
max_array_rows => { type => 'integer', optional => 1, default => 10_000 }
# Row-count threshold for 'auto' mode. When the combined
# source row count exceeds this value, the SQLite path is used.
# Ignored when backend is 'array' or 'sqlite'.
#
# DOMAIN -- EP valid: any non-negative integer.
# DOMAIN -- BVA: 0 means always use SQLite (all counts exceed 0).
# DOMAIN -- Default: 10,000.
tmpdir => { type => 'string', optional => 1 }
# Directory for the per-call temporary SQLite database file.
# The file is created securely by File::Temp and removed when
# the query completes. Ignored when backend is 'array'.
#
# DOMAIN -- EP valid: any writable directory path string.
# DOMAIN -- EP absent: uses File::Spec->tmpdir() (system temp dir).
parallel => { type => 'integer', optional => 1, default => 0 }
# When set to 1 and the join has more than 2 databases (primary +
# 2 or more secondaries), secondary DA fetches are issued in
# parallel Perl threads. Requires the 'threads' module; falls
# back to sequential with a carp warning when unavailable.
# Has no effect on the SQLite backend (which uses a single SQL
# JOIN). DBI-backed DAs are not thread-safe by default; only
# enable this for in-memory or otherwise thread-safe DA backends.
#
# DOMAIN -- EP valid: 0 (sequential, default) or 1 (parallel).
# DOMAIN -- EP invalid: any other integer is treated as truthy/falsy.
# DOMAIN -- Default: 0.
logger => { type => 'object', optional => 1 }
# Logger object propagated to all component databases.
i18n => { type => 'object', optional => 1 }
# Localisation object with a translate($key, @args) method.
=head4 Output
A blessed Database::Join object.
=head3 EXAMPLE
# Customers database: entry | name | email
# Loyalty database: entry | tier | points
my $join = Database::Join->new(
databases => [ $customers, $loyalty ],
join_column => 'entry',
join_type => 'inner', # only customers who also have loyalty records
remove_columns => [ 'email' ], # hide PII from query results
filters => { 1 => { points => { '>' => 0 } } }, # ignore zero-point records
);
my $rows = $join->selectall_arrayref();
# Each row: { entry => ..., name => ..., tier => ..., points => ... }
# 'email' is absent. Zero-point loyalty records are excluded.
=head3 PSEUDOCODE
validate all parameters with validate_strict
croak if databases is empty
croak if any element of databases is not a Database::Abstraction subclass
bless the object with all fields initialised
call _build_col_index to map every column to its owning database
and verify join_column presence in each database
if base_criteria given:
partition base_criteria by column ownership into per-db slices
for each database with a non-empty slice:
merge the slice into _filters[i]
(explicit filters win on plain-scalar conflicts)
for each column in remove_columns: call remove_column
return the new object
=head3 MESSAGES
error_no_databases -- databases arrayref was empty
error_invalid_db -- an element of databases is not a D::A subclass
error_join_col_missing -- join_column (or its join_map alias) not found in a database
error_invalid_backend -- backend value is not 'array', 'sqlite', or 'auto'
error_sqlite_connect -- temporary SQLite database could not be created (backend='sqlite'/'auto')
warn_schema_type_mismatch -- (carp) a shared column has different types across databases;
use collision_prefix to preserve both values
=cut | |||||
| 733 | ||||||
| 734 | sub new { | |||||
| 735 | 1249 | 2083697 | my ($class, @args) = @_; | |||
| 736 | ||||||
| 737 | 1249 | 10293 | my $p = validate_strict( | |||
| 738 | schema => { | |||||
| 739 | # databases => { type => 'arrayref', element_type => 'object' }, | |||||
| 740 | databases => { type => 'arrayref' }, | |||||
| 741 | join_column => { type => 'string', optional => 1, default => 'entry' }, | |||||
| 742 | join_type => { | |||||
| 743 | type => 'string', | |||||
| 744 | optional => 1, | |||||
| 745 | default => 'left', | |||||
| 746 | enum => ['inner', 'left', 'outer'] | |||||
| 747 | }, | |||||
| 748 | join_map => { type => 'hashref', optional => 1 }, | |||||
| 749 | filters => { type => 'hashref', optional => 1 }, | |||||
| 750 | base_criteria => { type => 'hashref', optional => 1 }, | |||||
| 751 | collision_prefix => { type => 'hashref', optional => 1 }, | |||||
| 752 | remove_columns => { type => 'arrayref', optional => 1 }, | |||||
| 753 | backend => { | |||||
| 754 | type => 'string', | |||||
| 755 | optional => 1, | |||||
| 756 | default => 'auto', | |||||
| 757 | enum => ['array', 'sqlite', 'auto'] | |||||
| 758 | }, | |||||
| 759 | max_array_rows => { type => 'integer', optional => 1, default => 10_000 }, | |||||
| 760 | tmpdir => { type => 'string', optional => 1 }, | |||||
| 761 | parallel => { type => 'integer', optional => 1, default => 0 }, | |||||
| 762 | logger => { type => 'object', optional => 1 }, | |||||
| 763 | i18n => { type => 'object', optional => 1 }, | |||||
| 764 | }, | |||||
| 765 | input => get_params(undef, \@args) // {}, | |||||
| 766 | ); | |||||
| 767 | ||||||
| 768 | croak _msg($p->{i18n}, 'error_no_databases') | |||||
| 769 | 1229 1229 | 317604 1731 | unless @{ $p->{databases} }; | |||
| 770 | ||||||
| 771 | # Capture caller-supplied logger and i18n object BEFORE Object::Configure::configure | |||||
| 772 | # overwrites them with class-level defaults. configure() reads default values | |||||
| 773 | # from Database::Abstraction's class config and silently replaces any caller- | |||||
| 774 | # supplied value if a class default exists for that key. | |||||
| 775 | 1214 | 1015 | my $caller_logger = $p->{logger}; | |||
| 776 | 1214 | 741 | my $caller_i18n = $p->{i18n}; | |||
| 777 | ||||||
| 778 | 1214 | 1598 | $p = Object::Configure::configure($class, $p); | |||
| 779 | ||||||
| 780 | 1214 1214 | 1971091 2033 | for my $i (0 .. $#{ $p->{databases} }) { | |||
| 781 | croak _msg($p->{i18n}, 'error_invalid_db', $i) | |||||
| 782 | unless blessed($p->{databases}[$i]) | |||||
| 783 | && $p->{databases}[$i]->can('selectall_arrayref') | |||||
| 784 | 2201 | 9509 | && $p->{databases}[$i]->can('columns'); | |||
| 785 | } | |||||
| 786 | ||||||
| 787 | # Cache the primary database's internal primary-key column name once at | |||||
| 788 | # construction. AUTOLOAD uses this to map a bare positional argument to the | |||||
| 789 | # correct column when join_map is active and the primary DB's primary key | |||||
| 790 | # differs from join_column (e.g. join_col='statecode' but primary key='entry'). | |||||
| 791 | # Accessing {id} here â at construction time, before the object is shared â | |||||
| 792 | # is the single permitted point of coupling to DA's internal field; caching | |||||
| 793 | # avoids repeating the hash intrusion on every AUTOLOAD call. | |||||
| 794 | 1199 | 2086 | my $primary_pk = $p->{databases}[0]{id} // $p->{join_column}; | |||
| 795 | ||||||
| 796 | my $self = bless { | |||||
| 797 | _dbs => $p->{databases}, | |||||
| 798 | _join_col => $p->{join_column}, | |||||
| 799 | _join_type => $p->{join_type}, | |||||
| 800 | _join_map => $p->{join_map} // {}, # db_index => local join col name | |||||
| 801 | # Security: deep-copy filters so post-construction mutation of the caller's | |||||
| 802 | # hashref cannot silently bypass the inner-join row-security guarantee. | |||||
| 803 | # Two-level copy mirrors the broadcast-copy idiom in _partition_criteria: | |||||
| 804 | # outer keys are db indices (integers); inner values are criteria hashrefs | |||||
| 805 | # whose operator sub-hashrefs are also shallow-copied one level deeper. | |||||
| 806 | _filters => _copy_filters($p->{filters}), # db_index => criteria hashref | |||||
| 807 | _collision_prefix => $p->{collision_prefix} // {}, # db_index => prefix string | |||||
| 808 | _logger => $caller_logger, | |||||
| 809 | _i18n => $caller_i18n, | |||||
| 810 | _col_db => {}, # published_col_name => db_index | |||||
| 811 | _db_cols => [], # per-db column-presence hashref | |||||
| 812 | _removed_cols => {}, # published_col_name => 1 (hidden from view) | |||||
| 813 | _col_cache => undef, # memoised columns() result | |||||
| 814 | _schema_cache => undef, # memoised schema() result | |||||
| 815 | _col_rename => [], # per-db: { orig_col => published_col } for collisions | |||||
| 816 | _col_unrename => [], # per-db: { published_col => orig_col } reverse map | |||||
| 817 | _autoload_pk => $primary_pk, # primary DB's key col; positional arg for AUTOLOAD | |||||
| 818 | _backend => $p->{backend}, | |||||
| 819 | _max_array_rows => $p->{max_array_rows}, | |||||
| 820 | _tmpdir => $p->{tmpdir} // File::Spec->tmpdir, | |||||
| 821 | 1199 | 3072 | _parallel => $p->{parallel} // 0, | |||
| 822 | }, $class; | |||||
| 823 | ||||||
| 824 | 1199 | 2412 | $self->_build_col_index(); | |||
| 825 | 1170 | 1455 | $self->_validate_schema_types(); | |||
| 826 | ||||||
| 827 | # base_criteria: partition by column ownership into _filters. | |||||
| 828 | # Security: _copy_criteria() deep-copies before partitioning so | |||||
| 829 | # post-construction mutation of the caller's hashref cannot bypass filters. | |||||
| 830 | # Explicit filters (db-indexed) take precedence: they are treated as | |||||
| 831 | # "extra" in _merge_criteria so their values win on plain-scalar conflicts. | |||||
| 832 | 1170 | 1211 | if (my $bc = $p->{base_criteria}) { | |||
| 833 | 9 | 11 | my $partitioned = $self->_partition_criteria(_copy_criteria($bc)); | |||
| 834 | 9 9 | 10 10 | for my $i (0 .. $#{ $self->{_dbs} }) { | |||
| 835 | 18 18 | 6 21 | next unless %{ $partitioned->[$i] }; | |||
| 836 | $self->{_filters}{$i} = _merge_criteria( | |||||
| 837 | $partitioned->[$i], | |||||
| 838 | 8 | 17 | $self->{_filters}{$i} // {}, | |||
| 839 | ); | |||||
| 840 | } | |||||
| 841 | } | |||||
| 842 | ||||||
| 843 | # Propagate the logger to every component database if one was supplied. | |||||
| 844 | # set_logger() is used here (rather than a direct hash write) to honour each | |||||
| 845 | # DA's own logging setup hook and remain decoupled from DA internals. | |||||
| 846 | 1170 | 1061 | if (my $log = $self->{_logger}) { | |||
| 847 | 3 3 | 3 10 | $_->set_logger($log) for @{ $self->{_dbs} }; | |||
| 848 | } | |||||
| 849 | ||||||
| 850 | # Apply column removals requested in the constructor | |||||
| 851 | 1170 | 1067 | if (my $rc = $p->{remove_columns}) { | |||
| 852 | 12 12 | 8 35 | $self->remove_column($_) for @{$rc}; | |||
| 853 | } | |||||
| 854 | ||||||
| 855 | 1168 | 2205 | return $self; | |||
| 856 | } | |||||
| 857 | ||||||
| 858 | # --------------------------------------------------------------------------- | |||||
| 859 | # Public API (mirrors Database::Abstraction) | |||||
| 860 | # --------------------------------------------------------------------------- | |||||
| 861 | ||||||
| 862 - 1456 | =head2 join_map - joining on differently-named columns
By default every component database must have a column whose name matches
C<join_column>. If a database uses a different local name for the join key,
declare the mapping with C<join_map>.
C<join_map> is a hashref. Each B<key> is the B<zero-based position> of a
database in the C<databases> array (0 = first, 1 = second, and so on). Each
B<value> is the name that B<that particular database> uses for the join key.
Databases not listed in C<join_map> are assumed to already have a column
named C<join_column> and need no entry.
Throughout the merged view the join key is I<always> referred to by the name
given in C<join_column>. The local alias is never exposed in returned rows,
in C<columns()>, or in C<schema()>.
B<When do you need join_map?>
You need C<join_map> when you have two tables like:
cities table : entry (the city name) | statecode
stnames table : entry (the state code) | state
Here you want to join cities.statecode to stnames.entry. You choose
C<< join_column => 'statecode' >> as the canonical name, but stnames calls
that same concept C<entry>, so you declare:
join_map => { 1 => 'entry' } # stnames (index 1) calls it 'entry'
B<Example>
# index 0 index 1
my @databases = ( $cities, $stnames );
# join key column: 'statecode' 'entry'
# join_column: 'statecode' (chosen canonical name)
# stnames differs, so declare the alias:
my $join = Database::Join->new(
databases => \@databases,
join_column => 'statecode',
join_map => { 1 => 'entry' },
);
my $rows = $join->selectall_arrayref();
# Each $row has keys: entry (city), statecode, state
# 'entry' from stnames is never exposed directly.
my $row = $join->fetchrow_hashref(statecode => 'CA');
B<Using add_database instead>
If you build the join incrementally with C<add_database>, pass
C<join_column> directly to that call instead of using C<join_map>:
my $join = Database::Join->new(
databases => [ $cities ],
join_column => 'statecode',
);
$join->add_database($stnames, join_column => 'entry');
This is exactly equivalent to the C<join_map> form above.
=head2 filters - permanent per-database row filters
C<filters> lets you restrict a component database to a subset of its rows
permanently, without repeating the criterion on every query call.
Think of it as telling the join: "whenever you query this database, always
add these extra conditions". Callers never need to specify the restriction
themselves and can never accidentally omit it.
C<filters> is a hashref. Each B<key> is the B<zero-based position> of a
database in the C<databases> array (same numbering as C<join_map>). Each
B<value> is a criteria hashref in the same format as C<selectall_arrayref>
accepts.
B<Key-set semantics>
A filtered database always acts as an inner-join partner, regardless of the
C<join_type> setting. Any join-key value that does not pass the filter is
excluded from the merged output entirely -- not just missing its secondary
columns. This ensures the filter genuinely restricts the view rather than
simply hiding a few fields.
B<Criteria merging>
When a query call also passes a criterion for a column that already has a base
filter, the two constraints are combined:
=over 4
=item *
When both the base filter value and the query criterion are operator hashrefs
(e.g. C<< { '>' => 60 } >> and C<< { '<' => 365 } >>), their operators are
merged: I<both> constraints apply simultaneously (AND semantics).
=item *
When either value is a plain scalar, or the operators conflict, the
query-time criterion wins and the base filter for that column is ignored for
that one call.
=back
B<Example -- only show orders placed more than 60 days ago>
my $join = Database::Join->new(
databases => [ $customers, $orders ],
join_column => 'entry',
filters => { 1 => { age_days => { '>' => 60 } } },
);
# Every query automatically sees only old orders
my $rows = $join->selectall_arrayref();
# Additional criteria layer on top -- gold tier AND old order
my $vip = $join->selectall_arrayref(tier => 'gold');
# Range intersection: age_days > 60 AND age_days < 365
my $mid = $join->selectall_arrayref(age_days => { '<' => 365 });
When using C<add_database>, pass C<filter> (singular) to set the base
criteria for the new database:
$join->add_database($orders, filter => { age_days => { '>' => 60 } });
=head2 base_criteria - permanent view-level row filters by column name
C<base_criteria> is a convenience alternative to C<filters> for callers who
know the column names they want to restrict but prefer not to track database
indices.
my $join = Database::Join->new(
databases => [ $customers, $orders ],
join_column => 'entry',
base_criteria => { active => 1, deleted_at => undef },
);
# Every query automatically sees only active, non-deleted rows.
my $rows = $join->selectall_arrayref();
At construction time C<base_criteria> is partitioned by column ownership
using the same routing logic as C<selectall_arrayref>. Each criterion is
sent to the database that owns that column and stored as a permanent base
filter (equivalent to the corresponding C<filters> entry).
B<Unknown columns> emit a C<warn_unknown_column> carp and are silently
dropped, just as they would be in a query call.
B<Key-set semantics> are identical to C<filters>: any database that receives
a C<base_criteria> slice acts as an inner-join partner regardless of
C<join_type>.
B<Precedence>: when a column appears in both C<base_criteria> and C<filters>,
the C<filters> entry wins on plain-scalar conflicts; operator-hashref values
are combined with AND semantics.
B<Use cases>
=over 4
=item *
Row-level security: C<< base_criteria => { tenant_id => $tid } >>
=item *
Soft-delete filtering: C<< base_criteria => { deleted_at => undef } >>
=item *
Status gates: C<< base_criteria => { active => 1 } >>
=back
=head2 collision_prefix - preserve colliding columns from secondary databases
By default, when a column name appears in more than one database the I<last>
database wins: its value silently overwrites earlier ones in merged rows.
This loses data and makes the origin invisible.
C<collision_prefix> changes this for secondary databases you designate.
When a secondary database at index N has a column that already exists in the
merged view, and C<collision_prefix-E<gt>{N}> is set, the colliding column is
published as C<"$prefix.$col"> instead of overwriting. Both values are then
visible: the original column keeps its name (from the earlier database), and
the collision gets the prefixed name.
C<collision_prefix> is a hashref. Each B<key> is the B<zero-based index> of
a secondary database in the C<databases> array (same numbering as C<join_map>).
Each B<value> is the prefix string to prepend. An index-0 entry is
meaningless and silently ignored. Omitting C<collision_prefix> entirely
preserves the previous last-wins behaviour and changes nothing.
Non-colliding columns from a secondary database are always added as-is with
no prefix, whether or not C<collision_prefix> is configured.
B<Example -- sales table and products table, both with a "product" column>
# $sales columns: id, product, amount, date
# $products columns: sku, product, price, category
# join on 'product' (left key) matched against 'sku' (right key via join_map)
my $join = Database::Join->new(
databases => [$sales, $products],
join_column => 'product',
join_map => { 1 => 'sku' },
collision_prefix => { 1 => 'products' },
);
$join->columns;
# => ['amount', 'category', 'date', 'id', 'price', 'product', 'products.product']
# ^-- prefixed collision
my $row = $join->fetchrow_hashref(product => 'widget');
# $row->{product} -- value from $sales
# $row->{'products.product'} -- value from $products (different row, same column name)
# $row->{price} -- from $products, no collision, kept as-is
B<Querying on a prefixed column>
Use the full published name as the criterion key:
my $rows = $join->selectall_arrayref('products.product' => 'widget');
# Internally routes as: product => 'widget' to $products
B<Interaction with remove_column>
C<remove_column> operates on published names. To suppress a prefixed
collision column entirely, pass the prefixed name:
$join->remove_column('products.product');
=head2 backend - SQLite join backend for large datasets
C<Database::Join> can merge component databases in two different ways,
controlled by the C<backend> constructor parameter.
=over 4
=item C<backend =E<gt> 'array'> -- in-memory merge (original behaviour)
All matching rows are fetched from every component database into Perl hashes
and merged there. Simple and fast for small and medium datasets. Peak RAM
is roughly three times the combined source data size (one copy per database
plus one merged copy).
=item C<backend =E<gt> 'sqlite'> -- SQL JOIN via a cached temporary file
C<Database::Join> creates a temporary SQLite database file, spills source
rows into it (one table per component database), then executes a single SQL
C<JOIN> statement per query call. Peak RAM drops to roughly one times the
source data size.
The temporary file is created once and reused across multiple query calls on
the same object (the cache). Only query-time criteria vary per call; they
are applied as SQL C<WHERE> clauses against the cached data. The cache is
automatically rebuilt when any source database's C<updated()> timestamp
changes. The file is deleted when the object is destroyed (goes out of
scope). See I<Temporary file: name, location, and lifetime> below for
details.
Requires C<DBD::SQLite E<gt>= 1.70> (C<FULL OUTER JOIN> support was added in
SQLite 3.39.0; DBD::SQLite 1.70 ships SQLite 3.39.2).
=item C<backend =E<gt> 'auto'> (default)
C<Database::Join> counts the total rows from all component databases cheaply
-- without fetching them -- and then decides:
=over 4
=item *
If the combined count is less than or equal to C<max_array_rows> (default
10,000), use the array path.
=item *
If the combined count exceeds C<max_array_rows>, use the SQLite path.
=back
For counting to work without fetching, each component database must either
implement the C<dbi_source()> interface (for SQLite-backed sources, where a
C<COUNT(*)> SQL query is issued directly), or directly define a C<count()>
method in its own package -- not just inherit one from a parent class. If
neither is available for a particular database, C<Database::Join> plays it
safe and uses the array path for the whole query without fetching any rows.
=back
B<Choosing max_array_rows>
The default of 10,000 is a reasonable starting point. Adjust it to match
your hardware and typical row width. For wide rows (many columns or long
strings) you may want a lower threshold; for narrow rows you can raise it.
B<Temporary file: name, location, and lifetime>
When the SQLite path is active, a single temporary SQLite database file acts
as the join cache for the life of the C<Database::Join> object.
B<Name>: the filename is randomly generated by C<File::Temp>, for example:
/tmp/Cj8xK7mP2Q.db
The random portion (ten characters) is chosen automatically to avoid
collisions. Only the directory is under your control; you cannot specify
the filename itself.
B<Location>: controlled by the C<tmpdir> constructor parameter.
If C<tmpdir> is not specified, C<File::Spec-E<gt>tmpdir()> is used (usually
C</tmp> on Unix, or the value of the C<TEMP> or C<TMP> environment variable
on Windows).
B<Lifetime: one file per object, deleted when the object is destroyed>: the
file is created on the first query call that uses the SQLite path and kept
alive until the C<Database::Join> object is destroyed (i.e. when it goes out
of scope or is explicitly C<undef>-d). At most one file exists per object at
any given moment. Calling C<selectall_arrayref()> ten times on the same
object creates and uses I<one> file, not ten.
B<Cache invalidation>: the cache is automatically rebuilt (the old file is
replaced with a new one) when any source database's C<updated()> return value
changes, or when C<add_database()> is called. Base-filter criteria
(C<filters> constructor parameter) are applied once at build time for spilled
sources; query-time criteria are applied per-call as SQL C<WHERE> clauses.
To use a different directory -- for example a RAM-backed filesystem or a
faster local disk:
my $join = Database::Join->new(
databases => [ $db1, $db2 ],
join_column => 'entry',
backend => 'sqlite',
tmpdir => '/dev/shm', # Linux RAM disk
);
B<Parallel secondary fetches (C<parallel> constructor parameter)>
By default, component databases are queried sequentially - the primary first,
then each secondary in order. When the component databases are network- or
disk-backed and have non-trivial per-query latency, the sequential fetch means
total latency is the I<sum> of all per-DA latencies.
Setting C<< parallel => 1 >> in the constructor enables concurrent fetching
of secondary databases using Perl C<threads>. With C<parallel => 1> and
two or more secondary databases (C<n E<gt> 2> total), secondary fetches run in
parallel after the primary fetch completes; total latency drops to
I<max(secondary latencies)> instead of I<sum(secondary latencies)>.
my $join = Database::Join->new(
databases => [ $customers, $loyalty, $scores ], # 3 DAs - 2 secondaries
join_column => 'entry',
parallel => 1, # loyalty and scores fetched concurrently
);
Requirements and caveats:
=over 4
=item *
The C<threads> module must be available. Most distributions ship it, but it
requires a Perl binary compiled with C<-Dusethreads>. When threads are
unavailable, a C<carp> warning is emitted and fetching falls back to
sequential; the result is identical, only slower.
=item *
Parallel fetching is only active when the join has three or more total
databases (C<n E<gt> 2>). With two databases (one secondary), the thread
creation overhead exceeds the benefit of concurrency; sequential is used
regardless of C<parallel>.
=item *
Component databases must be safe to call from Perl threads. In-memory
databases (CSV, JSON, TSV after slurp) are safe. DBI-backed databases
whose handles were created in the same thread may not be safe - consult your
DBD driver's thread documentation. The array backend is recommended for
DBI-backed sources; the SQLite backend performs its join in a single SQL
statement and does not use parallel fetching.
=item *
C<parallel => 1> applies to the array backend only. The SQLite backend
performs a single SQL JOIN after spilling source data, so per-DA parallelism
is irrelevant.
=back
B<Zero-copy ATTACH (C<dbi_source()> interface)>
Normally, when the SQLite path is active, rows from each component database
are fetched one by one and inserted into the temporary SQLite file. This is
efficient but does involve INSERT overhead.
If a component database is itself SQLite-backed and implements a
C<dbi_source()> method, C<Database::Join> can skip the row-by-row copy
entirely and instead use C<ATTACH DATABASE> to link the source file directly
to the temporary join connection. This is the zero-copy path and is
significantly faster for large SQLite sources.
The C<dbi_source()> method must return a hashref with two keys:
=over 4
=item C<dbh>
A connected C<DBD::SQLite> database handle (C<DBI> connection object).
=item C<table>
The name of the table in that database that holds the source rows.
=back
Example implementation:
package My::SQLiteDatabase;
use parent 'Database::Abstraction';
sub dbi_source {
my ($self) = @_;
return {
dbh => $self->{_dbh}, # connected DBD::SQLite handle
table => $self->{_table}, # table name in that database
};
}
1;
The zero-copy ATTACH path is always used when a component database implements
C<dbi_source()>. Query-time criteria are applied as parameterised SQL
C<WHERE> clauses against the ATTACHed source table, so no row-level copy is
needed even when the current call includes column filters. (Prior to 0.005.0
any query-time criteria forced a spill; that restriction was removed in
0.005.0.)
B<Result identity>
Both the array path and the SQLite path produce identical results for any
given query. You can switch between them freely without changing callers.
The C<collision_prefix> column renaming, C<join_map> key translation, and all
three join types (left, inner, outer) work identically on both paths.
=head2 selectall_arrayref
=head3 SYNOPSIS
my $rows = $join->selectall_arrayref();
my $rows = $join->selectall_arrayref(tier => 'gold');
my $rows = $join->selectall_arrayref(score => { '>' => 80 });
my $rows = $join->selectall_arrayref('C001'); # positional: entry => 'C001'
my $rows = $join->selectall_arrayref(sort_by => 'name');
my $rows = $join->selectall_arrayref(tier => 'gold', sort_by => ['score', 'DESC']);
my $rows = $join->selectall_arrayref(limit => 10);
my $rows = $join->selectall_arrayref(limit => 10, offset => 20);
my $rows = $join->selectall_arrayref(tier => 'gold', sort_by => 'name', limit => 5);
=head3 DESCRIPTION
Returns an arrayref of hashrefs representing the merged view of all component
databases, optionally filtered by the given criteria.
Criteria for columns that live in different databases are routed
automatically: each database is queried with only the criteria that apply to
its own columns. The results are combined in memory using C<join_column>.
Accepts the same criteria syntax as C<Database::Abstraction::selectall_arrayref>.
A single plain scalar argument is interpreted as the C<join_column> value
(equivalent to C<< entry => 'C001' >> when C<join_column> is C<'entry'>).
=head3 API SPECIFICATION
=head4 Input
Calling conventions (in order of precedence):
1. No arguments -- returns all rows
2. One plain scalar -- shorthand for join_column => $scalar
3. Key-value pairs or
a criteria hashref -- routed per-database
Values may be:
Plain scalar -- exact match
Hashref of operators -- e.g. { '>' => 80 }
Optional parameters (mixed in with any of the above):
sort_by => 'colname' -- sort ascending by that column
sort_by => ['colname', 'DESC'] -- sort descending
sort_by => ['colname', 'ASC'] -- sort ascending (explicit)
limit => N -- return at most N rows (positive integer)
offset => M -- skip the first M rows (non-negative integer)
The column named in sort_by must be present in the merged view (i.e. it
must appear in columns()). An unknown column or an invalid direction emits
a carp warning and falls back to the default join_column ascending sort.
limit and offset are applied after ordering. offset without limit skips
rows but returns all remaining rows. limit without offset starts from
the first qualifying row. An invalid limit or offset emits a carp warning
and the parameter is ignored (treated as absent).
DOMAIN -- sort_by:
EP absent: result sorted by join_column ASC (default).
EP string: any column name in columns(); sorts ASC by that column.
EP ['col','ASC']: explicit ascending; equivalent to the string form.
EP ['col','DESC']: descending sort by the named column.
EP ['col']: single-element arrayref; direction defaults to ASC.
EP []: empty arrayref; column is undef -> carp + join_col ASC fallback.
EP invalid column: column not in columns() -> carp + join_col ASC fallback.
EP invalid dir: direction not 'ASC' or 'DESC' -> carp + ASC used.
Sort is lexicographic (cmp); use backend=>'sqlite' for numeric ORDER BY.
DOMAIN -- limit:
EP absent: no truncation; all qualifying rows are returned.
EP 0: not a positive integer -> carp + ignored (all rows returned).
BVA min valid = 1: exactly 1 row returned.
BVA at count: limit == total rows -> all rows returned (no truncation).
BVA above count: limit > total rows -> all rows returned.
EP invalid: negative integer, float string, or non-numeric string
=> carp + ignored (all rows returned).
Valid domain: integers in [1, INF); matched by /^\d+\z/a with value >= 1.
DOMAIN -- offset:
EP absent: no rows skipped; result starts from row 0.
BVA min valid = 0: no rows skipped (zero is a valid non-negative integer).
BVA offset=1: first row skipped; result starts from row 1.
BVA offset=N-1: N-1 rows skipped; only the last row returned.
BVA offset=N: all N rows skipped; empty result returned.
BVA offset>N: all rows skipped; empty result returned.
EP invalid: negative integer, float string, or non-numeric string
=> carp + ignored (no rows skipped).
Valid domain: integers in [0, INF); matched by /^\d+\z/a.
=head4 Output
Arrayref of hashrefs; one hashref per qualifying merged row.
Sorted ascending by join_column by default; caller-controlled via sort_by.
At most C<limit> rows when limit is given; the first C<offset> rows are
skipped when offset is given.
Returns a reference to an empty array when no rows match.
=head3 EXAMPLE
# All rows from both databases
my $all = $join->selectall_arrayref();
# Only rows where the 'tier' column (from the loyalty database)
# equals 'gold' -- the criterion is routed to the right database
my $vip = $join->selectall_arrayref(tier => 'gold');
# Operator hashref: score > 80
my $high = $join->selectall_arrayref(score => { '>' => 80 });
# Access each merged row
for my $row (@{$vip}) {
printf "%-10s tier=%-8s score=%d\n",
$row->{entry}, $row->{tier}, $row->{score} // 0;
}
=head3 MESSAGES
warn_unknown_column (carp)
-- A criterion key names a column not present in any component database;
the criterion is silently dropped and all rows are returned.
error_join_criterion (croak)
-- A criterion hashref contains a join => key (DA SQL-JOIN syntax).
Use a separate DA object for the joined table instead.
sort_by column unknown (carp)
-- The column given in sort_by is not in the merged view; the result
is returned in the default join_column ascending order instead.
sort_by direction invalid (carp)
-- The direction given in sort_by is not 'ASC' or 'DESC'; ASC is used.
limit invalid (carp)
-- The value given for limit is not a positive integer; it is ignored.
offset invalid (carp)
-- The value given for offset is not a non-negative integer; it is ignored.
operator-unsupported (carp, SQLite/auto path only)
-- An operator hashref key is not in the supported set (>, <, >=, <=,
!=, =, LIKE, NOT LIKE, IS NULL, IS NOT NULL, IN, NOT IN); the
individual operator term is dropped from the WHERE clause (other
operators in the same hashref still apply).
IN/NOT IN wrong value type (carp, SQLite/auto path only)
-- An IN or NOT IN criterion was given a non-arrayref value; the operator
term is skipped.
error_sqlite_connect (croak, SQLite/auto path only)
-- The temporary SQLite join file could not be created; check tmpdir
permissions and available disk space.
=cut | |||||
| 1457 | ||||||
| 1458 | sub selectall_arrayref { | |||||
| 1459 | 725 | 43544 | my ($self, @args) = @_; | |||
| 1460 | 725 | 953 | my $params = $self->_parse_query_args(undef, @args); | |||
| 1461 | 725 | 4056 | my $sort_by = delete $params->{sort_by}; | |||
| 1462 | 725 | 502 | my $limit = delete $params->{limit}; | |||
| 1463 | 725 | 457 | my $offset = delete $params->{offset}; | |||
| 1464 | 725 | 961 | return $self->_joined_query($params, sort_by => $sort_by, limit => $limit, offset => $offset); | |||
| 1465 | } | |||||
| 1466 | ||||||
| 1467 - 1509 | =head2 selectall_array
=head3 SYNOPSIS
my @rows = $join->selectall_array(tier => 'gold');
# Scalar context: only the first matching row
my $first = $join->selectall_array(entry => 'C001');
=head3 DESCRIPTION
In list context returns a list of merged hashrefs -- the same rows that
C<selectall_arrayref> would return, just as a flat list rather than an
arrayref.
In scalar context returns only the first matching hashref (or C<undef> if
nothing matches).
=head3 API SPECIFICATION
=head4 Input
Same as selectall_arrayref.
=head4 Output
List context: list of hashrefs (may be empty).
Scalar context: single hashref or undef.
=head3 EXAMPLE
my @all = $join->selectall_array();
print scalar @all, " rows\n";
# First gold-tier customer only
my $first_vip = $join->selectall_array(tier => 'gold');
print $first_vip->{name}, "\n" if defined $first_vip;
=head3 MESSAGES
Same messages as C<selectall_arrayref>.
=cut | |||||
| 1510 | ||||||
| 1511 | sub selectall_array { | |||||
| 1512 | 22 | 1749 | my ($self, @args) = @_; | |||
| 1513 | 22 | 30 | my $params = $self->_parse_query_args(undef, @args); | |||
| 1514 | 22 | 78 | my $sort_by = delete $params->{sort_by}; | |||
| 1515 | 22 | 16 | my $limit = delete $params->{limit}; | |||
| 1516 | 22 | 17 | my $offset = delete $params->{offset}; | |||
| 1517 | 22 | 30 | my $rows = $self->_joined_query($params, sort_by => $sort_by, limit => $limit, offset => $offset); | |||
| 1518 | 22 13 | 41 26 | return wantarray ? @{$rows} : $rows->[0]; | |||
| 1519 | } | |||||
| 1520 | ||||||
| 1521 - 1527 | =head2 selectall_hashref Deprecated alias for L</selectall_arrayref>, present for compatibility with callers written against C<Database::Abstraction>'s deprecated API. Use C<selectall_arrayref> in new code. =cut | |||||
| 1528 | ||||||
| 1529 | sub selectall_hashref { | |||||
| 1530 | 5 | 96 | my $self = shift; | |||
| 1531 | 5 | 57 | carp 'Database::Join::selectall_hashref is deprecated; use selectall_arrayref'; | |||
| 1532 | 5 | 656 | return $self->selectall_arrayref(@_); | |||
| 1533 | } | |||||
| 1534 | ||||||
| 1535 - 1541 | =head2 selectall_hash Deprecated alias for L</selectall_array>, present for compatibility with callers written against C<Database::Abstraction>'s deprecated API. Use C<selectall_array> in new code. =cut | |||||
| 1542 | ||||||
| 1543 | sub selectall_hash { | |||||
| 1544 | 5 | 97 | my $self = shift; | |||
| 1545 | 5 | 62 | carp 'Database::Join::selectall_hash is deprecated; use selectall_array'; | |||
| 1546 | 5 | 664 | return $self->selectall_array(@_); | |||
| 1547 | } | |||||
| 1548 | ||||||
| 1549 - 1590 | =head2 fetchrow_hashref
=head3 SYNOPSIS
my $row = $join->fetchrow_hashref(entry => 'C001');
my $row = $join->fetchrow_hashref('C001'); # positional shorthand
=head3 DESCRIPTION
Returns a single merged hashref for the first row matching the given
criteria, or C<undef> when nothing matches.
Equivalent to calling C<selectall_arrayref> and taking only the first element.
All the same criteria conventions apply.
=head3 API SPECIFICATION
=head4 Input
Same as selectall_arrayref.
=head4 Output
Hashref, or undef when no row matches.
=head3 EXAMPLE
my $row = $join->fetchrow_hashref(entry => 'C001');
if (defined $row) {
print "Name: $row->{name}, Tier: $row->{tier}\n";
} else {
print "No record for C001\n";
}
# Positional: works when join_column is 'entry'
my $row2 = $join->fetchrow_hashref('C001');
=head3 MESSAGES
Same messages as C<selectall_arrayref>.
=cut | |||||
| 1591 | ||||||
| 1592 | sub fetchrow_hashref { | |||||
| 1593 | 67 | 9728 | my ($self, @args) = @_; | |||
| 1594 | 67 | 99 | my $params = $self->_parse_query_args(undef, @args); | |||
| 1595 | 67 | 551 | my $sort_by = delete $params->{sort_by}; | |||
| 1596 | 67 | 56 | delete $params->{limit}; # fetchrow_hashref always returns one row; limit is meaningless | |||
| 1597 | 67 | 50 | delete $params->{offset}; # offset would change which row is "first"; not supported here | |||
| 1598 | 67 | 95 | my $rows = $self->_joined_query($params, sort_by => $sort_by); | |||
| 1599 | 67 | 359 | return $rows->[0]; | |||
| 1600 | } | |||||
| 1601 | ||||||
| 1602 - 1640 | =head2 count
=head3 SYNOPSIS
my $total = $join->count();
my $active = $join->count(tier => 'gold');
=head3 DESCRIPTION
Returns the number of merged rows that satisfy the given criteria.
On the array backend, the full in-memory join is performed and the resulting
rows are counted in Perl. On the SQLite backend, a C<SELECT COUNT(*)> SQL
query is executed against the cached join tables, avoiding a full row fetch.
=head3 API SPECIFICATION
=head4 Input
Same criteria syntax as selectall_arrayref.
=head4 Output
Non-negative integer.
=head3 EXAMPLE
my $total = $join->count();
my $gold = $join->count(tier => 'gold');
my $high = $join->count(score => { '>' => 90 });
printf "%d total, %d gold-tier, %d high-scorers\n",
$total, $gold, $high;
=head3 MESSAGES
Same messages as C<selectall_arrayref>.
=cut | |||||
| 1641 | ||||||
| 1642 | sub count { | |||||
| 1643 | 77 | 12545 | my ($self, @args) = @_; | |||
| 1644 | 77 | 122 | my $params = $self->_parse_query_args(undef, @args); | |||
| 1645 | 77 | 298 | delete $params->{sort_by}; # row ordering is irrelevant for a count | |||
| 1646 | 77 | 58 | delete $params->{limit}; # count returns total matching rows, not a page | |||
| 1647 | 77 | 47 | delete $params->{offset}; | |||
| 1648 | # On the SQLite path, push COUNT(*) into SQL to avoid fetching all rows. | |||||
| 1649 | return $self->_sqlite_join($params, count_only => 1) | |||||
| 1650 | 77 | 168 | unless $self->{_backend} eq 'array'; | |||
| 1651 | 5 5 | 5 10 | return scalar @{ $self->_joined_query_array($params) }; | |||
| 1652 | } | |||||
| 1653 | ||||||
| 1654 - 1716 | =head2 each_row
=head3 SYNOPSIS
my $count = $join->each_row(sub { my ($row) = @_; ... });
my $count = $join->each_row(sub { ... }, tier => 'gold');
my $count = $join->each_row(sub { ... }, sort_by => 'name', limit => 100);
=head3 DESCRIPTION
Calls C<\&callback> once for every row in the merged view that matches the
given criteria, then returns the total count of rows visited.
Accepts the same criteria, C<sort_by>, C<limit>, and C<offset> parameters as
C<selectall_arrayref>.
B<Note on memory usage>: C<Database::Join> always materialises the complete
merged result set before invoking the callback, because merging rows from
multiple independent sources requires that every source be queried and
cross-referenced first. True constant-memory streaming (as C<Database::Abstraction>
provides on single-source SQL queries) is not possible at the join layer.
For genuinely large result sets use the SQLite backend (C<backend =E<gt>
'sqlite'>), which keeps peak RAM to approximately one times the source data
size rather than three.
Any exception thrown inside the callback propagates to the caller after the
current row; subsequent rows are not visited.
=head3 API SPECIFICATION
=head4 Input
\&callback Positional coderef (required).
Called as $callback->($row_hashref) for each merged row.
Additional arguments follow the same calling conventions as
selectall_arrayref: no args, a single scalar join-column value,
or key-value criteria pairs (including sort_by, limit, offset).
=head4 Output
Non-negative integer: the count of rows for which the callback was invoked.
=head3 EXAMPLE
# Print every gold-tier customer's name, sorted by name
my $count = $join->each_row(
sub { my ($row) = @_; print "$row->{name}\n" },
tier => 'gold',
sort_by => 'name',
);
print "$count gold-tier customers\n";
# Accumulate without holding the full result
my $total_score = 0;
$join->each_row(sub { $total_score += $_[0]->{score} // 0 });
=head3 MESSAGES
error_invalid_callback (croak)
-- First argument is not a code reference.
=cut | |||||
| 1717 | ||||||
| 1718 | sub each_row { | |||||
| 1719 | 16 | 573 | my ($self, $cb, @args) = @_; | |||
| 1720 | 16 | 31 | croak $self->_err('error_invalid_callback') | |||
| 1721 | unless ref($cb) eq 'CODE'; | |||||
| 1722 | # Materialise the full merged result first; multi-source joins cannot | |||||
| 1723 | # stream rows one at a time because the merge key set is not known until | |||||
| 1724 | # all sources have been queried and cross-referenced. | |||||
| 1725 | 11 | 12 | my $rows = $self->selectall_arrayref(@args); | |||
| 1726 | 11 | 8 | my $count = 0; | |||
| 1727 | 11 11 | 5 10 | for my $row (@{$rows}) { | |||
| 1728 | 17 | 17 | $cb->($row); | |||
| 1729 | 16 | 22 | $count++; | |||
| 1730 | } | |||||
| 1731 | 10 | 15 | return $count; | |||
| 1732 | } | |||||
| 1733 | ||||||
| 1734 - 1795 | =head2 dbi_source
=head3 SYNOPSIS
# Use a Database::Join object as a zero-copy SQLite source inside a
# parent Database::Join, giving the parent ATTACHed-speed access to
# the child's joined data without iterating through Perl.
my $child = Database::Join->new(databases => [$da, $db], join_column => 'id', backend => 'sqlite');
my $parent = Database::Join->new(databases => [$child, $dc], join_column => 'id', backend => 'sqlite');
=head3 DESCRIPTION
Returns a hashref C<< { dbh => $sqlite_dbh, table => '_dj_result' } >> that
allows a parent C<Database::Join> (or any other caller that understands the
C<dbi_source()> interface) to ATTACH the child's temporary SQLite database
file and query the materialised join result directly via SQL, without routing
rows through Perl.
The first call builds the SQLite cache (if not already current) and
materialises the full join result - with C<filters> applied but no query-time
criteria - into a real table named C<_dj_result> inside the cache file.
Subsequent calls within the same cache cycle reuse the existing table.
Returns C<undef> when the backend is C<'array'> (no SQLite file exists).
=head3 API SPECIFICATION
=head4 Input
None.
=head4 Output
On the SQLite/auto backend:
Hashref with keys:
dbh => DBI handle to the child's temporary SQLite database file.
table => '_dj_result' (the materialized join table inside that file).
On the array backend:
undef
=head3 EXAMPLE
# The parent automatically ATTACHes the child's SQLite file and queries
# _dj_result for zero-copy composable nested joins.
my $inner = Database::Join->new(
databases => [$customers, $loyalty],
join_column => 'entry',
backend => 'sqlite',
);
my $outer = Database::Join->new(
databases => [$inner, $scores],
join_column => 'entry',
backend => 'sqlite',
);
my $rows = $outer->selectall_arrayref(tier => 'gold');
=head3 MESSAGES
error_sqlite_connect (croak)
-- The temporary SQLite join file could not be created.
=cut | |||||
| 1796 | ||||||
| 1797 | sub dbi_source { | |||||
| 1798 | 26 | 422 | my ($self) = @_; | |||
| 1799 | ||||||
| 1800 | # The array backend has no SQLite handle to expose. | |||||
| 1801 | 26 | 45 | return undef if $self->{_backend} eq 'array'; | |||
| 1802 | ||||||
| 1803 | # Build or refresh the per-source SQLite cache. On 'auto' backend this | |||||
| 1804 | # forces the SQLite path regardless of the auto-threshold decision â a | |||||
| 1805 | # parent that wants to ATTACH us requires a real on-disk file. | |||||
| 1806 | 20 | 29 | $self->_build_sqlite_cache() unless $self->_cache_fresh(); | |||
| 1807 | ||||||
| 1808 | 18 | 20 | my $cache = $self->{_sqlite_cache}; | |||
| 1809 | ||||||
| 1810 | # If the materialised result table is already current, reuse it. | |||||
| 1811 | # _dj_built is reset implicitly when _build_sqlite_cache creates a fresh | |||||
| 1812 | # cache hashref (the old hashref is replaced, so its _dj_built is gone). | |||||
| 1813 | 18 | 27 | unless ($cache->{_dj_built}) { | |||
| 1814 | # Materialise all joined rows (filter criteria only, no query-time | |||||
| 1815 | # criteria) into _dj_result so a parent connection can ATTACH and | |||||
| 1816 | # query it as a plain table. | |||||
| 1817 | 16 | 46 | $self->_sqlite_join({}, create_table => '_dj_result'); | |||
| 1818 | 16 | 36 | $cache->{_dj_built} = 1; | |||
| 1819 | } | |||||
| 1820 | ||||||
| 1821 | 18 | 75 | return { dbh => $cache->{dbh}, table => '_dj_result' }; | |||
| 1822 | } | |||||
| 1823 | ||||||
| 1824 - 1862 | =head2 columns
=head3 SYNOPSIS
my $cols = $join->columns();
=head3 DESCRIPTION
Returns an arrayref of all column names visible in the merged view,
deduplicated and sorted alphabetically.
The C<join_column> appears exactly once, even if it exists under different
local names in some databases (see C<join_map>). Columns that have been
hidden with C<remove_column> or C<remove_columns> do not appear.
The result is memoised: repeated calls are cheap.
=head3 API SPECIFICATION
=head4 Input
None.
=head4 Output
Arrayref of column name strings, sorted alphabetically.
=head3 EXAMPLE
my $cols = $join->columns();
print join(', ', @{$cols}), "\n";
# e.g. "entry, name, score, tier"
=head3 MESSAGES
C<columns()> does not itself emit any warnings or errors. Any exception thrown
by a component database's C<columns()> method propagates uncaught.
=cut | |||||
| 1863 | ||||||
| 1864 | sub columns { | |||||
| 1865 | 185 | 11292 | my ($self) = @_; | |||
| 1866 | ||||||
| 1867 | 185 | 231 | return $self->{_col_cache} if $self->{_col_cache}; | |||
| 1868 | ||||||
| 1869 | 165 | 118 | my %seen; | |||
| 1870 | my @cols; | |||||
| 1871 | 165 | 131 | my $join_col = $self->{_join_col}; | |||
| 1872 | ||||||
| 1873 | 165 165 | 131 166 | for my $i (0 .. $#{ $self->{_dbs} }) { | |||
| 1874 | 341 | 254 | my $local_jc = $self->{_join_map}{$i}; | |||
| 1875 | 341 | 303 | my $renames = $self->{_col_rename}[$i] // {}; | |||
| 1876 | 341 341 | 179 312 | for my $col (@{ $self->{_dbs}[$i]->columns() }) { | |||
| 1877 | # The local join-key alias is not a data column; the canonical name | |||||
| 1878 | # is already contributed by the database that owns it under that name. | |||||
| 1879 | 828 | 862 | next if $local_jc && $col eq $local_jc && $col ne $join_col; | |||
| 1880 | # Use the published name (prefixed if this column was a collision rename) | |||||
| 1881 | 820 | 739 | my $pub = $renames->{$col} // $col; | |||
| 1882 | 820 | 672 | next if $seen{$pub}++; | |||
| 1883 | 648 | 475 | next if $self->{_removed_cols}{$pub}; | |||
| 1884 | 601 | 457 | push @cols, $pub; | |||
| 1885 | } | |||||
| 1886 | } | |||||
| 1887 | ||||||
| 1888 | 165 | 367 | $self->{_col_cache} = [ sort @cols ]; | |||
| 1889 | 165 | 294 | return $self->{_col_cache}; | |||
| 1890 | } | |||||
| 1891 | ||||||
| 1892 - 1935 | =head2 schema
=head3 SYNOPSIS
my $schema = $join->schema();
=head3 DESCRIPTION
Returns a merged schema hashref for all visible columns across all component
databases. Each key is a column name; each value is the schema metadata
hashref returned by C<Database::Abstraction::schema()> for that column
(typically C<{ type, nullable, default, pk }>).
When the same column name appears in more than one database the I<last>
database's metadata is used. Columns hidden with C<remove_column> are not
included.
The result is memoised.
=head3 API SPECIFICATION
=head4 Input
None.
=head4 Output
Hashref: column_name => { type => ..., nullable => ..., default => ..., pk => ... }.
=head3 EXAMPLE
my $schema = $join->schema();
for my $col (sort keys %{$schema}) {
my $info = $schema->{$col};
printf "%-15s type=%-10s nullable=%s\n",
$col, $info->{type}, $info->{nullable} ? 'yes' : 'no';
}
=head3 MESSAGES
C<schema()> does not itself emit any warnings or errors. Any exception thrown
by a component database's C<schema()> method propagates uncaught.
=cut | |||||
| 1936 | ||||||
| 1937 | sub schema { | |||||
| 1938 | 56 | 5129 | my ($self) = @_; | |||
| 1939 | ||||||
| 1940 | 56 | 62 | return $self->{_schema_cache} if $self->{_schema_cache}; | |||
| 1941 | ||||||
| 1942 | # Hash slice assignment (@merged{keys} = values) is O(M) per database; | |||||
| 1943 | # the previous (%merged = (%merged, %s)) pattern was O(NÃM) per iteration, | |||||
| 1944 | # totalling O(N²ÃM) over all databases for the same result. | |||||
| 1945 | 49 | 37 | my %merged; | |||
| 1946 | 49 49 | 44 53 | for my $i (0 .. $#{ $self->{_dbs} }) { | |||
| 1947 | 97 | 87 | my $s = $self->{_dbs}[$i]->schema() // {}; | |||
| 1948 | 97 | 165 | my $local_jc = $self->{_join_map}{$i}; | |||
| 1949 | 97 | 110 | my $renames = $self->{_col_rename}[$i] // {}; | |||
| 1950 | 97 | 90 | my $skip_jc = $local_jc && $local_jc ne $self->{_join_col}; | |||
| 1951 | 97 | 79 | if (!$skip_jc && !%{$renames}) { | |||
| 1952 | # Fast path: no join-key alias to exclude and no collision renames | |||||
| 1953 | 88 88 88 | 45 127 67 | @merged{keys %{$s}} = values %{$s}; | |||
| 1954 | } else { | |||||
| 1955 | 9 9 | 5 12 | for my $col (keys %{$s}) { | |||
| 1956 | 24 | 30 | next if $skip_jc && $col eq $local_jc; | |||
| 1957 | 20 | 35 | $merged{ $renames->{$col} // $col } = $s->{$col}; | |||
| 1958 | } | |||||
| 1959 | } | |||||
| 1960 | } | |||||
| 1961 | ||||||
| 1962 | 10 | 17 | delete @merged{keys %{ $self->{_removed_cols} }} | |||
| 1963 | 49 49 | 39 58 | if %{ $self->{_removed_cols} }; | |||
| 1964 | ||||||
| 1965 | 49 | 52 | $self->{_schema_cache} = \%merged; | |||
| 1966 | 49 | 76 | return $self->{_schema_cache}; | |||
| 1967 | } | |||||
| 1968 | ||||||
| 1969 - 2010 | =head2 updated
=head3 SYNOPSIS
my $ts = $join->updated();
=head3 DESCRIPTION
Returns the Unix timestamp of the most recent modification across all
component databases. This is the maximum of all individual C<updated()>
return values.
Use this to implement simple cache-invalidation logic: if C<updated()>
has advanced since your last snapshot, re-query.
=head3 API SPECIFICATION
=head4 Input
None.
=head4 Output
Unix timestamp (positive integer).
=head3 EXAMPLE
my $last_modified = $join->updated();
if ($last_modified > $my_cache_timestamp) {
$my_cache = $join->selectall_arrayref();
$my_cache_timestamp = $last_modified;
}
=head3 MESSAGES
C<updated()> does not emit any warnings or errors. Component databases that do
not implement C<updated()>, or whose C<updated()> throws, are silently skipped;
only defined return values contribute to the maximum. If no component database
implements C<updated()>, C<undef> is returned (same as C<List::Util::max> on an
empty list).
=cut | |||||
| 2011 | ||||||
| 2012 | sub updated { | |||||
| 2013 | 35 | 1956 | my ($self) = @_; | |||
| 2014 | 35 | 26 | my @timestamps; | |||
| 2015 | 35 35 | 26 46 | for my $db (@{ $self->{_dbs} }) { | |||
| 2016 | 71 | 41 | my $ts; | |||
| 2017 | 71 71 71 71 | 45 37 47 143 | do { local $@; $ts = eval { $db->updated() } }; | |||
| 2018 | 71 | 1111 | push @timestamps, $ts if defined $ts; | |||
| 2019 | } | |||||
| 2020 | 35 | 112 | return max(@timestamps); | |||
| 2021 | } | |||||
| 2022 | ||||||
| 2023 - 2070 | =head2 set_logger =head3 SYNOPSIS $join->set_logger($log); $join->set_logger(logger => $log); # named-pair form also accepted =head3 DESCRIPTION Attaches a new logger object to the join and propagates it to every component database. The logger is used for diagnostic output by all component databases. Non-blessed values (a log-level string such as C<"debug">, a filename, or a code reference) are wrapped in C<Log::Abstraction->new(...)> automatically, matching the behaviour of C<Database::Abstraction::set_logger>. =head3 API SPECIFICATION =head4 Input logger Positional or named: a logger object, log-level string, filename, or code reference (required). Non-blessed values are wrapped in Log::Abstraction automatically. =head4 Output Returns C<$self> for method chaining. =head3 EXAMPLE # Log::Any is used here as an example; any object that implements # debug() and info() (or whichever methods your component databases # call internally) works equally well. use Log::Any qw($log); my $join = Database::Join->new(databases => [$db1, $db2], join_column => 'entry'); $join->set_logger($log); # $log is now used by $join and by $db1 and $db2 # Named-pair form (mirrors Database::Abstraction API): $join->set_logger(logger => $log); =head3 MESSAGES (croak) Usage: set_logger(logger => $logger) -- Called with an undefined argument. Pass a valid logger object. =cut | |||||
| 2071 | ||||||
| 2072 | sub set_logger { | |||||
| 2073 | 17 | 167 | my $self = shift; | |||
| 2074 | 17 | 22 | my $p = Params::Get::get_params('logger', @_); | |||
| 2075 | 17 | 203 | my $logger = $p->{'logger'}; | |||
| 2076 | ||||||
| 2077 | 17 | 56 | croak 'Usage: set_logger(logger => $logger)' unless defined $logger; | |||
| 2078 | ||||||
| 2079 | # Wrap non-blessed values (log-level strings, filenames, coderefs) exactly | |||||
| 2080 | # as Database::Abstraction does, so callers get identical behaviour from DJ. | |||||
| 2081 | 14 | 21 | $logger = Log::Abstraction->new($logger) | |||
| 2082 | unless Scalar::Util::blessed($logger); | |||||
| 2083 | ||||||
| 2084 | 14 | 14 | $self->{_logger} = $logger; | |||
| 2085 | 14 14 | 7 34 | $_->set_logger($logger) for @{ $self->{_dbs} }; | |||
| 2086 | ||||||
| 2087 | 14 | 124 | return $self; | |||
| 2088 | } | |||||
| 2089 | ||||||
| 2090 - 2220 | =head2 add_database
=head3 SYNOPSIS
# Positional: database object as first argument
$join->add_database($db);
# Named: equivalent to the above
$join->add_database(database => $db);
# With options (mixed positional + named)
$join->add_database($db, remove_columns => ['internal_id']);
$join->add_database($db, join_column => 'local_key_name');
$join->add_database($db, filter => { score => { '>' => 60 } });
# Chainable
$join->add_database($db1)->add_database($db2, remove_columns => ['notes']);
=head3 DESCRIPTION
Adds one more C<Database::Abstraction> subclass object to the logical view
and immediately updates the column-ownership index.
After the call, all query methods return rows that include columns from the
newly added database, and criteria on those new columns are routed to it
automatically.
When a column name in the new database already exists in an earlier database,
the new database becomes the authoritative source for that column
(last-database-wins, the same rule that applies at construction time).
The join-column must be present in the new database (or declared via
C<join_column>). The logger is propagated to the new database if one is set.
C<add_database> is the runtime equivalent of listing the database in the
C<databases> array to C<new>. The optional C<join_column> parameter is
equivalent to a C<join_map> entry; the optional C<filter> parameter is
equivalent to a C<filters> entry.
=head3 API SPECIFICATION
=head4 Input
database => { type => 'object', required => 1 }
# A Database::Abstraction subclass instance.
#
# DOMAIN -- EP valid: blessed object that passes
# isa('Database::Abstraction').
# DOMAIN -- EP invalid: non-reference, unblessed ref, wrong class,
# or non-reference non-key scalar (the guard at
# the top of add_database rejects it with
# error_invalid_db before validate_strict runs).
join_column => { type => 'string', optional => 1 }
# The name of the join key in THIS new database,
# when it differs from the canonical join_column.
#
# DOMAIN -- EP valid: any string that exists as a column in the
# new database.
# DOMAIN -- EP invalid: string absent from the new database's columns()
# => croak error_join_col_missing.
filter => { type => 'hashref', optional => 1 }
# Permanent criteria for this database only.
# Same format as selectall_arrayref.
#
# DOMAIN -- EP valid: hashref of criteria (may be {} for no-op).
# DOMAIN -- EP absent: no permanent filter applied; all rows visible.
# DOMAIN -- Key-set: a non-empty filter makes this DB an inner-join
# partner regardless of the outer join_type.
remove_columns => { type => 'arrayref', optional => 1 }
# Column names from this database to hide.
#
# DOMAIN -- EP valid: arrayref of strings; non-existent columns silently
# ignored; empty [] is a safe no-op.
# DOMAIN -- EP invalid: join_column itself => croak error_remove_join_col.
=head4 Output
Returns C<$self> to support method chaining.
=head3 EXAMPLE
my $join = Database::Join->new(
databases => [ $customers ],
join_column => 'entry',
);
# Add loyalty data; hide internal columns from it
$join->add_database($loyalty, remove_columns => ['audit_ts']);
# Add score data; only include rows with score > 60
$join->add_database($scores, filter => { score => { '>' => 60 } });
# Add a database whose join key has a different local name
$join->add_database($stnames, join_column => 'state_code');
# All three options combined, and chained
$join->add_database($db4,
join_column => 'ref_id',
filter => { active => 1 },
remove_columns => ['legacy_col'],
);
=head3 PSEUDOCODE
determine the new database's index (length of current _dbs array)
extract the database object from positional or named argument
croak if it is not a Database::Abstraction subclass
register join_column alias in _join_map if different from canonical
register filter in _filters if provided
fetch column list from the new database
croak if the join key is missing from the new database
append the new database to _dbs and _db_cols
update _col_db: for each new column, point it at the new index
(last-database-wins; skip removed columns and the local join alias)
invalidate _col_cache and _schema_cache
propagate logger if set
apply remove_columns if provided
return $self
=head3 MESSAGES
error_invalid_db -- argument is not a Database::Abstraction subclass
error_join_col_missing -- join_column not found in the new database
warn_schema_type_mismatch -- (carp) the new database has a shared column whose type
differs from the type already in the view; use
collision_prefix to preserve both values
=cut | |||||
| 2221 | ||||||
| 2222 | sub add_database { | |||||
| 2223 | 113 | 6470 | my ($self, @args) = @_; | |||
| 2224 | ||||||
| 2225 | 113 113 | 71 132 | my $idx = scalar @{ $self->{_dbs} }; | |||
| 2226 | 113 | 97 | my $db; | |||
| 2227 | ||||||
| 2228 | # Fail-fast guard: a non-reference first arg must be a recognised named-pair key. | |||||
| 2229 | # Modus Ponens: !ref(x) â§ x â @_ADD_DB_KEYS â cannot be a database object â croak. | |||||
| 2230 | # De Morgan reduction: the elsif below is logically equivalent to (@args && ref(args[0])) | |||||
| 2231 | # because the !ref branch was already handled; exhaustion makes the ref() check redundant. | |||||
| 2232 | 113 | 284 | if (@args && !ref($args[0])) { | |||
| 2233 | croak $self->_err('error_invalid_db', $idx) | |||||
| 2234 | 14 48 | 59 163 | unless defined($args[0]) && grep { $args[0] eq $_ } @_ADD_DB_KEYS; | |||
| 2235 | } elsif (@args) { | |||||
| 2236 | # Positional form: first arg is a reference â extract it before get_params | |||||
| 2237 | # to avoid the mixed positional+named-pairs confusion. | |||||
| 2238 | 99 | 80 | $db = shift @args; | |||
| 2239 | } | |||||
| 2240 | ||||||
| 2241 | 104 | 531 | my $p = validate_strict( | |||
| 2242 | schema => { | |||||
| 2243 | database => { | |||||
| 2244 | type => 'object', | |||||
| 2245 | optional => 1, | |||||
| 2246 | can => ['selectall_arrayref', 'columns'] | |||||
| 2247 | }, | |||||
| 2248 | join_column => { type => 'string', optional => 1 }, | |||||
| 2249 | filter => { type => 'hashref', optional => 1 }, | |||||
| 2250 | remove_columns => { type => 'arrayref', optional => 1 }, | |||||
| 2251 | }, | |||||
| 2252 | input => (@args ? get_params(undef, @args) : {}) // {}, | |||||
| 2253 | ); | |||||
| 2254 | ||||||
| 2255 | 104 | 8358 | $db //= $p->{database}; | |||
| 2256 | ||||||
| 2257 | 104 | 436 | croak $self->_err('error_invalid_db', $idx) | |||
| 2258 | unless blessed($db) | |||||
| 2259 | && $db->can('selectall_arrayref') | |||||
| 2260 | && $db->can('columns'); | |||||
| 2261 | ||||||
| 2262 | # Determine and register the local join column name for this database | |||||
| 2263 | 97 | 174 | my $local_jc = $p->{join_column} // $self->{_join_col}; | |||
| 2264 | 97 | 100 | $self->{_join_map}{$idx} = $local_jc if $p->{join_column}; | |||
| 2265 | # Security: deep-copy the filter; same rationale as the constructor's _copy_filters call. | |||||
| 2266 | 97 | 106 | $self->{_filters}{$idx} = _copy_criteria($p->{filter}) if $p->{filter}; | |||
| 2267 | ||||||
| 2268 | 97 | 96 | my $cols = $db->columns(); | |||
| 2269 | 97 211 97 | 1185 202 79 | my %col_presence = map { $_ => 1 } @{$cols}; | |||
| 2270 | ||||||
| 2271 | croak $self->_err('error_join_col_missing', $local_jc, $idx, ref($db)) | |||||
| 2272 | 97 | 138 | unless $col_presence{$local_jc}; | |||
| 2273 | ||||||
| 2274 | # Register the new database | |||||
| 2275 | 93 93 | 63 90 | push @{ $self->{_dbs} }, $db; | |||
| 2276 | 93 93 | 120 85 | push @{ $self->{_db_cols} }, \%col_presence; | |||
| 2277 | ||||||
| 2278 | # Update column routing: last-database-wins for duplicates, unless a | |||||
| 2279 | # collision_prefix is configured for this index (in which case the | |||||
| 2280 | # duplicate is published under "$prefix.$col" instead of overwriting). | |||||
| 2281 | 93 | 90 | my $prefix = $self->{_collision_prefix}{$idx}; | |||
| 2282 | 93 | 180 | $self->{_col_rename}[$idx] //= {}; | |||
| 2283 | 93 | 150 | $self->{_col_unrename}[$idx] //= {}; | |||
| 2284 | ||||||
| 2285 | 93 93 | 53 80 | for my $col (@{$cols}) { | |||
| 2286 | 207 | 202 | next if $local_jc ne $self->{_join_col} && $col eq $local_jc; | |||
| 2287 | ||||||
| 2288 | 199 | 104 | my $pub; | |||
| 2289 | 199 | 198 | if (defined $prefix && exists $self->{_col_db}{$col} && $col ne $self->{_join_col}) { | |||
| 2290 | # Same guard as _build_col_index: never prefix the join_column itself. | |||||
| 2291 | 4 | 5 | $pub = "$prefix.$col"; | |||
| 2292 | 4 | 4 | $self->{_col_rename}[$idx]{$col} = $pub; | |||
| 2293 | 4 | 5 | $self->{_col_unrename}[$idx]{$pub} = $col; | |||
| 2294 | } else { | |||||
| 2295 | 195 | 111 | $pub = $col; | |||
| 2296 | } | |||||
| 2297 | ||||||
| 2298 | 199 | 182 | next if $self->{_removed_cols}{$pub}; | |||
| 2299 | 197 | 183 | $self->{_col_db}{$pub} = $idx; | |||
| 2300 | } | |||||
| 2301 | ||||||
| 2302 | # Invalidate memoisation caches | |||||
| 2303 | 93 | 83 | $self->{_col_cache} = undef; | |||
| 2304 | 93 | 60 | $self->{_schema_cache} = undef; | |||
| 2305 | ||||||
| 2306 | # Warn about schema type mismatches introduced by the new database. | |||||
| 2307 | 93 | 105 | $self->_validate_schema_types(); | |||
| 2308 | ||||||
| 2309 | # Invalidate the SQLite join cache: a new source requires a full rebuild. | |||||
| 2310 | 93 | 110 | if (my $old = delete $self->{_sqlite_cache}) { | |||
| 2311 | 5 | 5 | local $@; | |||
| 2312 | 5 5 | 34 106 | eval { $old->{dbh}->disconnect } if $old->{dbh}; | |||
| 2313 | } | |||||
| 2314 | ||||||
| 2315 | # Propagate logger if one is configured | |||||
| 2316 | 93 | 109 | if (my $log = $self->{_logger}) { | |||
| 2317 | 4 | 4 | $db->set_logger($log); | |||
| 2318 | } | |||||
| 2319 | ||||||
| 2320 | # Apply any column removals requested for this database | |||||
| 2321 | 93 | 120 | if (my $rc = $p->{remove_columns}) { | |||
| 2322 | 6 6 | 4 14 | $self->remove_column($_) for @{$rc}; | |||
| 2323 | } | |||||
| 2324 | ||||||
| 2325 | 93 | 192 | return $self; | |||
| 2326 | } | |||||
| 2327 | ||||||
| 2328 - 2397 | =head2 remove_column
=head3 SYNOPSIS
$join->remove_column('email');
# Chainable
$join->remove_column('internal_id')->remove_column('audit_ts');
=head3 DESCRIPTION
Permanently hides a column from the merged view. After this call:
=over 4
=item *
The column does not appear in C<columns()> or C<schema()>.
=item *
Returned row hashrefs do not contain the column key.
=item *
Any query criterion that references the removed column is silently dropped
(with a C<carp> warning).
=back
The C<join_column> cannot be removed; attempting to do so will C<croak>.
Removing a column that does not exist in any database is silently ignored
(the call is idempotent and safe). The C<columns()> and C<schema()>
memoisation caches are cleared automatically.
=head3 API SPECIFICATION
=head4 Input
$col Positional string: the column name to remove.
DOMAIN -- EP valid: any string; non-existent columns are silently
ignored (idempotent call, returns $self).
DOMAIN -- EP invalid: join_column value => croak error_remove_join_col.
DOMAIN -- BVA: undef and '' are explicit no-ops (returns $self).
These are below the minimum meaningful string
length and are handled without any warning.
=head4 Output
Returns C<$self> to support method chaining.
=head3 EXAMPLE
# Hide private fields immediately after construction
my $join = Database::Join->new(
databases => [ $customers, $loyalty ],
join_column => 'entry',
)->remove_column('email')
->remove_column('internal_notes');
# Verify they are gone
my $cols = $join->columns();
# 'email' and 'internal_notes' are absent
=head3 MESSAGES
error_remove_join_col -- attempt to remove the join_column itself
=cut | |||||
| 2398 | ||||||
| 2399 | sub remove_column { | |||||
| 2400 | 127 | 5062 | my ($self, $col) = @_; | |||
| 2401 | ||||||
| 2402 | # Premise: undef and '' are provably no-ops (nothing to remove). | |||||
| 2403 | # Conclusion: guard at the top eliminates two separate defined() checks below. | |||||
| 2404 | 127 | 231 | return $self unless defined $col && length $col; | |||
| 2405 | ||||||
| 2406 | # Premise: $col is defined (proven above) and join_col is always a non-empty string. | |||||
| 2407 | # Conclusion: direct string comparison is safe without a redundant defined() check. | |||||
| 2408 | croak $self->_err('error_remove_join_col', $col) | |||||
| 2409 | 113 | 139 | if $col eq $self->{_join_col}; | |||
| 2410 | ||||||
| 2411 | 99 | 108 | $self->{_removed_cols}{$col} = 1; | |||
| 2412 | 99 | 122 | delete $self->{_col_db}{$col}; | |||
| 2413 | 99 | 84 | $self->{_col_cache} = undef; | |||
| 2414 | 99 | 61 | $self->{_schema_cache} = undef; | |||
| 2415 | 99 | 78 | $self->{_removed_list} = undef; # invalidate the cached removed-column list | |||
| 2416 | ||||||
| 2417 | 99 | 99 | return $self; | |||
| 2418 | } | |||||
| 2419 | ||||||
| 2420 - 2430 | =head2 query Not supported. C<Database::Join> does not implement the C<Database::Abstraction::Query> chained builder because the builder's C<.all()> / C<.first()> methods would be targeting a single component DA rather than the merged view. Calling this method will always C<croak>. Use C<selectall_arrayref>, C<selectall_array>, C<fetchrow_hashref>, C<count>, or C<each_row> against the C<Database::Join> object instead. =cut | |||||
| 2431 | ||||||
| 2432 | sub query { | |||||
| 2433 | 8 | 409 | my $self = $_[0]; | |||
| 2434 | 8 | 14 | croak $self->_err('error_query_unsupported'); | |||
| 2435 | } | |||||
| 2436 | ||||||
| 2437 - 2444 | =head2 execute Not supported. Raw SQL cannot span heterogeneous backends that may use different database engines. Calling this method will always C<croak>. Use C<selectall_arrayref> or C<fetchrow_hashref> to query the joined view. =cut | |||||
| 2445 | ||||||
| 2446 | sub execute { | |||||
| 2447 | 8 | 690 | my $self = $_[0]; | |||
| 2448 | 8 | 15 | croak $self->_err('error_execute_unsupported'); | |||
| 2449 | } | |||||
| 2450 | ||||||
| 2451 - 2536 | =head2 AUTOLOAD - column shortcut
Calling an unknown method whose name matches a visible column name performs
a column lookup across the merged view.
=head3 SYNOPSIS
# Scalar context: value from the first matching row
my $name = $join->name(entry => 'C001');
# List context: values from every matching row
my @tiers = $join->tier();
# With a positional join-key argument (when join_column is 'entry')
my $score = $join->score('C001');
=head3 DESCRIPTION
AUTOLOAD routes the call to the appropriate component database by looking up
the column name in the internal column-ownership index.
When either C<join_map> or C<filters> is active, AUTOLOAD performs a full
join query instead of delegating directly to the owning database. This is
necessary because:
=over 4
=item *
With C<join_map>, the owning database's primary key may differ from the
canonical join key used in the call arguments.
=item *
With C<filters>, bypassing the join would return rows that the filter is
meant to exclude.
=back
In list context, every matching merged row contributes one value to the
returned list. In scalar context, only the first row's value is returned.
Calling a method whose name begins with C<_> (a private method) via AUTOLOAD
will C<croak> with a clear error message rather than being silently ignored.
=head3 EXAMPLE
# Lookup a single customer's name (scalar context)
my $name = $join->name('C001'); # 'C001' maps to entry => 'C001'
print "Name: $name\n";
# Get every tier value in the view (list context)
my @all_tiers = $join->tier();
my %freq;
$freq{$_}++ for @all_tiers;
# join_map active: AUTOLOAD runs a full join so the criteria are
# translated correctly between the canonical and local key names.
my @leesburg_states = sort $join->state('Leesburg');
# ['Florida', 'Virginia'] if Leesburg appears in two states
=head3 PSEUDOCODE
extract column name from $AUTOLOAD
return if DESTROY
croak if column name starts with '_' (private method guard)
croak if column name is not in _col_db (unknown column)
if join_map or filters are active:
parse calling arguments using _parse_query_args
call _joined_query to get all merged rows
return map { $_->{col} } @rows in list context
return $rows[0]{col} in scalar context
else:
delegate directly to the owning database
=head3 MESSAGES
(croak) Database::Join: cannot call private method '_NAME' via AUTOLOAD
-- Method name begins with '_'. Private methods must be called directly,
not via AUTOLOAD. This is a programming error.
(croak) Database::Join: unknown column 'NAME'
-- Method name does not match any visible column in the merged view.
Check spelling, or whether the column was removed with remove_column().
=cut | |||||
| 2537 | ||||||
| 2538 | our $AUTOLOAD; | |||||
| 2539 | ||||||
| 2540 | sub AUTOLOAD { | |||||
| 2541 | 56 | 2488 | my $self = shift; | |||
| 2542 | ||||||
| 2543 | 56 | 165 | my ($col) = $AUTOLOAD =~ / | |||
| 2544 | :: # package separator â skip the fully-qualified prefix | |||||
| 2545 | (\w++) # method name: possessive quantifier commits immediately; | |||||
| 2546 | # no backtrack possible because \w chars cannot match \z | |||||
| 2547 | \z # strict end-of-string (\z never matches a trailing newline, | |||||
| 2548 | # unlike $ which can â important if $AUTOLOAD ever embeds \n) | |||||
| 2549 | /x; | |||||
| 2550 | ||||||
| 2551 | # Private methods must not be reached via AUTOLOAD â croak immediately so | |||||
| 2552 | # typos like $join->_join_col are not silently swallowed. | |||||
| 2553 | # substr() avoids regex-engine overhead for this single-character prefix check. | |||||
| 2554 | 56 | 192 | croak ref($self), ": cannot call private method '$col' via AUTOLOAD" | |||
| 2555 | if substr($col, 0, 1) eq '_'; | |||||
| 2556 | ||||||
| 2557 | 46 | 57 | my $db_idx = $self->{_col_db}{$col}; | |||
| 2558 | 46 | 196 | croak ref($self), ": unknown column '$col'" unless defined $db_idx; | |||
| 2559 | ||||||
| 2560 | # Use a full join query when join_map OR filters are active. Direct | |||||
| 2561 | # delegation to the owning database would bypass the join key translation | |||||
| 2562 | # (join_map) and skip any permanent per-database row filters (filters). | |||||
| 2563 | 35 35 24 | 26 55 38 | if (%{ $self->{_join_map} } || %{ $self->{_filters} }) { | |||
| 2564 | # _autoload_pk was captured once at construction from the primary DA's | |||||
| 2565 | # {id} field; using the cached value avoids re-introspecting the blessed | |||||
| 2566 | # hash on every call and isolates the coupling to a single known site. | |||||
| 2567 | 27 | 41 | my $params = $self->_parse_query_args($self->{_autoload_pk}, @_); | |||
| 2568 | 27 | 94 | my $sort_by = delete $params->{sort_by}; | |||
| 2569 | 27 | 21 | my $limit = delete $params->{limit}; | |||
| 2570 | 27 | 21 | my $offset = delete $params->{offset}; | |||
| 2571 | 27 | 49 | my $rows = $self->_joined_query($params, sort_by => $sort_by, limit => $limit, offset => $offset); | |||
| 2572 | 27 10 9 | 35 20 12 | return map { $_->{$col} } @{$rows} if wantarray; | |||
| 2573 | 18 18 | 13 62 | return @{$rows} ? $rows->[0]{$col} : undef; | |||
| 2574 | } | |||||
| 2575 | ||||||
| 2576 | # $db is resolved here (not earlier) to avoid a dead store on the join path above. | |||||
| 2577 | 8 | 12 | my $db = $self->{_dbs}[$db_idx]; | |||
| 2578 | 8 | 18 | return $db->$col(@_); | |||
| 2579 | } | |||||
| 2580 | ||||||
| 2581 | sub DESTROY { | |||||
| 2582 | 1236 | 476608 | my ($self) = @_; | |||
| 2583 | # Disconnect and release the cached SQLite handle (if any) so File::Temp | |||||
| 2584 | # can unlink the temp file before the object is freed. | |||||
| 2585 | 1236 | 6930 | if (my $cache = delete $self->{_sqlite_cache}) { | |||
| 2586 | 181 | 146 | local $@; | |||
| 2587 | 181 181 | 302 6469 | eval { $cache->{dbh}->disconnect } if $cache->{dbh}; | |||
| 2588 | # $cache->{tmpfile} (File::Temp, UNLINK => 1) is released here. | |||||
| 2589 | } | |||||
| 2590 | } | |||||
| 2591 | ||||||
| 2592 | # --------------------------------------------------------------------------- | |||||
| 2593 | # Private helpers | |||||
| 2594 | # --------------------------------------------------------------------------- | |||||
| 2595 | ||||||
| 2596 | # _parse_query_args( $self, $positional_key, @caller_args ) -> \%params | |||||
| 2597 | # Purpose: Normalise the three calling conventions used by every public query | |||||
| 2598 | # method and AUTOLOAD into a single criteria hashref. | |||||
| 2599 | # Entry: $positional_key -- the column name mapped to a bare scalar argument; | |||||
| 2600 | # pass undef to use the join_column (the default for public methods). | |||||
| 2601 | # Exit: Always returns a hashref; never undef. | |||||
| 2602 | sub _parse_query_args :Protected { | |||||
| 2603 | 508 | 523 | my ($self, $key, @args) = @_; | |||
| 2604 | # D~ elimination (Modus Tollens): when @args is empty the early return fires | |||||
| 2605 | # before $key is ever read, making the $key //= assignment a dead store. | |||||
| 2606 | # Moving the guard above the assignment removes the wasted hash dereference. | |||||
| 2607 | 573 | 666 | return {} unless @args; | |||
| 2608 | 327 | 658 | $key //= $self->{_join_col}; | |||
| 2609 | 377 | 438 | return { $key => $args[0] } if @args == 1 && !ref($args[0]); | |||
| 2610 | 371 | 1040 | return get_params(undef, @args) // {}; | |||
| 2611 | 25 25 25 | 39840 40 390 | } | |||
| 2612 | ||||||
| 2613 | # _err( $self, $msg_key, @sprintf_args ) -> $string | |||||
| 2614 | # Convenience wrapper around _msg for use after construction, so callers do | |||||
| 2615 | # not have to extract $self->{_i18n} at every error site. | |||||
| 2616 | sub _err :Protected { | |||||
| 2617 | 176 | 186 | my ($self, $key, @args) = @_; | |||
| 2618 | 176 | 436 | return _msg($self->{_i18n}, $key, @args); | |||
| 2619 | 25 25 25 | 3820 28 147 | } | |||
| 2620 | ||||||
| 2621 | # _build_col_index() | |||||
| 2622 | # Purpose: Populate _col_db (column_name => db_index) and _db_cols | |||||
| 2623 | # (per-db column-presence hashrefs) by calling columns() on each | |||||
| 2624 | # component database at construction time. | |||||
| 2625 | # Entry: _dbs, _join_col, _join_map must already be set. | |||||
| 2626 | # Exit: _col_db and _db_cols are set; join_column verified in every db. | |||||
| 2627 | # Effects: Croaks if any database is missing its join key column. | |||||
| 2628 | sub _build_col_index :Protected { | |||||
| 2629 | 576 | 395 | my ($self) = @_; | |||
| 2630 | ||||||
| 2631 | 576 | 720 | my $join_col = $self->{_join_col}; | |||
| 2632 | 576 | 512 | my $cp = $self->{_collision_prefix} // {}; | |||
| 2633 | 576 | 688 | my %col_db; | |||
| 2634 | my @db_cols; | |||||
| 2635 | 576 | 146 | my @col_rename; # per-db: { orig_col => published_col } for renamed collisions | |||
| 2636 | 576 | 539 | my @col_unrename; # per-db: { published_col => orig_col } reverse map | |||
| 2637 | ||||||
| 2638 | 481 481 | 324 544 | for my $i (0 .. $#{ $self->{_dbs} }) { | |||
| 2639 | 851 | 938 | my $db = $self->{_dbs}[$i]; | |||
| 2640 | 851 | 1065 | my $local_jc = $self->{_join_map}{$i} // $join_col; | |||
| 2641 | 851 | 1332 | my $cols = $db->columns(); | |||
| 2642 | 851 1805 851 | 1006 1661 552 | $db_cols[$i] = { map { $_ => 1 } @{$cols} }; | |||
| 2643 | 977 | 920 | $col_rename[$i] = {}; | |||
| 2644 | 977 | 730 | $col_unrename[$i] = {}; | |||
| 2645 | ||||||
| 2646 | # Guard: a join_map value that is a reference (e.g. a hashref) would | |||||
| 2647 | # stringify to "HASH(0x...)" when interpolated into an error message, | |||||
| 2648 | # leaking a heap address to callers. Reject early with a clear message. | |||||
| 2649 | 977 | 2585 | croak $self->_err('error_join_col_missing', "(join_map[$i] must be a string)", $i, ref($db)) | |||
| 2650 | if ref $local_jc; | |||||
| 2651 | ||||||
| 2652 | croak $self->_err('error_join_col_missing', $local_jc, $i, ref($db)) | |||||
| 2653 | 972 | 1016 | unless $db_cols[$i]{$local_jc}; | |||
| 2654 | ||||||
| 2655 | # collision_prefix only applies to secondary databases (index > 0); | |||||
| 2656 | # an index-0 entry is meaningless and silently ignored. | |||||
| 2657 | 967 | 922 | my $prefix = ($i > 0) ? $cp->{$i} : undef; | |||
| 2658 | ||||||
| 2659 | # Guard: a collision_prefix value that is a reference would stringify | |||||
| 2660 | # to "HASH(0x...)" or "ARRAY(0x...)" when interpolated into "$prefix.$col", | |||||
| 2661 | # leaking a heap address into every column name, columns(), schema(), and | |||||
| 2662 | # merged row hashref. This is the same class of leak as the join_map guard | |||||
| 2663 | # above; reject early with a clear message before any column name is built. | |||||
| 2664 | 967 | 841 | croak $self->_err('error_invalid_prefix', $i) | |||
| 2665 | if defined $prefix && ref $prefix; | |||||
| 2666 | ||||||
| 2667 | 1143 966 | 668 645 | for my $col (@{$cols}) { | |||
| 2668 | # Skip the local alias for the join key â it is not a data column | |||||
| 2669 | 1909 | 1343 | next if $local_jc ne $join_col && $col eq $local_jc; | |||
| 2670 | ||||||
| 2671 | 1899 | 191454 | if (defined $prefix && exists $col_db{$col} && $col ne $join_col) { | |||
| 2672 | # Column already claimed by an earlier database AND a prefix is | |||||
| 2673 | # configured: publish the collision as "$prefix.$col" so both | |||||
| 2674 | # values survive in the merged row rather than one silently winning. | |||||
| 2675 | # The join_column itself is never prefixed â it is the shared merge | |||||
| 2676 | # key and is always broadcast by name; renaming it would break routing. | |||||
| 2677 | 155 | 176826 | my $pub = "$prefix.$col"; | |||
| 2678 | 155 | 295 | $col_rename[$i]{$col} = $pub; | |||
| 2679 | 155 | 118 | $col_unrename[$i]{$pub} = $col; | |||
| 2680 | 332 | 272 | $col_db{$pub} = $i; | |||
| 2681 | } else { | |||||
| 2682 | # No collision, or no prefix configured: last database wins | |||||
| 2683 | # (preserved backward-compatible behaviour). | |||||
| 2684 | 1876 | 1823 | $col_db{$col} = $i; | |||
| 2685 | } | |||||
| 2686 | } | |||||
| 2687 | } | |||||
| 2688 | ||||||
| 2689 | 596 | 545 | $self->{_col_db} = \%col_db; | |||
| 2690 | 596 | 589 | $self->{_db_cols} = \@db_cols; | |||
| 2691 | 596 | 450 | $self->{_col_rename} = \@col_rename; | |||
| 2692 | 596 | 4809 | $self->{_col_unrename} = \@col_unrename; | |||
| 2693 | ||||||
| 2694 | 596 | 664 | return; | |||
| 2695 | 25 25 25 | 7072 19 138 | } | |||
| 2696 | ||||||
| 2697 | # _validate_schema_types() | |||||
| 2698 | # Purpose: Carp when two databases share a column name (without collision_prefix) | |||||
| 2699 | # but disagree on its type, which would cause silent type coercion. | |||||
| 2700 | # Entry: _dbs, _col_rename, _join_map, _join_col are all populated. | |||||
| 2701 | # Exit: Emits one carp per mismatched column; no other side effects. | |||||
| 2702 | sub _validate_schema_types :Protected { | |||||
| 2703 | 636 | 1261 | my ($self) = @_; | |||
| 2704 | ||||||
| 2705 | 636 | 499 | my $join_col = $self->{_join_col}; | |||
| 2706 | 947 | 520 | my %seen; # col_name => { idx => $i, type => $type_str } | |||
| 2707 | ||||||
| 2708 | 1599 947 | 7471 747 | for my $i (0 .. $#{ $self->{_dbs} }) { | |||
| 2709 | 1375 | 1177 | my $db = $self->{_dbs}[$i]; | |||
| 2710 | 932 932 932 1064 | 424 501 543 158090 | my $s = do { local $@; eval { $db->schema() } } // {}; | |||
| 2711 | 1064 | 2361 | my $local_jc = $self->{_join_map}{$i} // $join_col; | |||
| 2712 | 1064 | 1180 | my $renames = $self->{_col_rename}[$i] // {}; | |||
| 2713 | ||||||
| 2714 | 1005 1005 | 489 1048 | for my $orig_col (keys %{$s}) { | |||
| 2715 | # Skip the join-key alias (a different name for the same join key) | |||||
| 2716 | 583 | 417 | next if $local_jc ne $join_col && $orig_col eq $local_jc; | |||
| 2717 | ||||||
| 2718 | # Skip the canonical join column itself â types may legitimately | |||||
| 2719 | # differ between DAs (e.g. INTEGER PK vs TEXT) without causing issues | |||||
| 2720 | # because the join column is not a data column. | |||||
| 2721 | 581 | 437 | next if $orig_col eq $join_col; | |||
| 2722 | ||||||
| 2723 | # Skip columns that have a collision prefix configured for this DB: | |||||
| 2724 | # they are published under distinct names, so no silent merge occurs. | |||||
| 2725 | 411 | 971 | next if exists $renames->{$orig_col}; | |||
| 2726 | ||||||
| 2727 | # Normalise to a plain type string; skip if the DA returns no type. | |||||
| 2728 | 392 | 414 | my $entry = $s->{$orig_col}; | |||
| 2729 | my $type_str = ref($entry) eq 'HASH' | |||||
| 2730 | 327 | 608 | ? uc($entry->{type} // '') | |||
| 2731 | : uc($entry // ''); | |||||
| 2732 | 327 | 294 | next unless length $type_str; | |||
| 2733 | ||||||
| 2734 | 525 | 456 | if (exists $seen{$orig_col}) { | |||
| 2735 | 273 | 388 | my $prev = $seen{$orig_col}; | |||
| 2736 | 273 | 217 | if ($prev->{type} ne $type_str) { | |||
| 2737 | carp $self->_err( | |||||
| 2738 | 'warn_schema_type_mismatch', | |||||
| 2739 | $orig_col, | |||||
| 2740 | $prev->{type}, $prev->{idx}, | |||||
| 2741 | 272 | 169 | $type_str, $i, | |||
| 2742 | ); | |||||
| 2743 | } | |||||
| 2744 | } else { | |||||
| 2745 | 523 | 516 | $seen{$orig_col} = { idx => $i, type => $type_str }; | |||
| 2746 | } | |||||
| 2747 | } | |||||
| 2748 | } | |||||
| 2749 | ||||||
| 2750 | 775 | 722 | return; | |||
| 2751 | 25 25 25 | 5418 20 560 | } | |||
| 2752 | ||||||
| 2753 | # _partition_criteria( \%params ) -> \@per_db | |||||
| 2754 | # Purpose: Split a flat criteria hashref into one slice per component database. | |||||
| 2755 | # Entry: $params is a criteria hashref; all keys must be column names or | |||||
| 2756 | # join_column. | |||||
| 2757 | # Exit: Returns an arrayref of per-database criteria hashrefs. The | |||||
| 2758 | # join_column criterion is broadcast to every database using each | |||||
| 2759 | # database's local join-key name. Unknown columns trigger a carp. | |||||
| 2760 | # Effects: Carps for each unrecognised column name. | |||||
| 2761 | sub _partition_criteria :Protected { | |||||
| 2762 | 837 | 656 | my ($self, $params) = @_; | |||
| 2763 | ||||||
| 2764 | 837 | 520 | my $join_col = $self->{_join_col}; | |||
| 2765 | 837 837 | 465 530 | my $n = scalar @{ $self->{_dbs} }; | |||
| 2766 | 837 1288 | 737 998 | my @per_db = map { {} } 1 .. $n; | |||
| 2767 | ||||||
| 2768 | 837 837 | 646 634 | for my $col (keys %{$params}) { | |||
| 2769 | 417 | 473 | if ($col eq $join_col) { | |||
| 2770 | # Broadcast to every database using each one's local key column name. | |||||
| 2771 | # Shallow-copy operator hashrefs so a malicious component DA that | |||||
| 2772 | # mutates its criteria hashref contents cannot corrupt siblings. | |||||
| 2773 | 554 | 613 | my $val = $params->{$col}; | |||
| 2774 | 554 | 302 | for my $i (0 .. $n - 1) { | |||
| 2775 | 607 | 738 | my $local = $self->{_join_map}{$i} // $join_col; | |||
| 2776 | 377 281 | 276 269 | $per_db[$i]{$local} = ref($val) eq 'HASH' ? { %{$val} } : $val; | |||
| 2777 | } | |||||
| 2778 | } elsif ($col eq 'join') { | |||||
| 2779 | # DA accepts a 'join =>' criteria key for SQL JOINs within one table. | |||||
| 2780 | # DJ cannot route such a criterion to a component DA meaningfully, so | |||||
| 2781 | # croak early with a diagnostic rather than silently dropping it. | |||||
| 2782 | 230 | 307 | croak $self->_err('error_join_criterion'); | |||
| 2783 | } elsif (defined(my $idx = $self->{_col_db}{$col})) { | |||||
| 2784 | # Translate the published column name back to the database's own name | |||||
| 2785 | # when the column was renamed for a collision (e.g. "pfx.col" -> "col"). | |||||
| 2786 | # Invariant: _col_unrename[$idx] is always initialised to {} by _build_col_index | |||||
| 2787 | # and add_database, so the // {} fallback can never trigger (transitive reduction). | |||||
| 2788 | 314 | 236 | my $db_col = $self->{_col_unrename}[$idx]{$col} // $col; | |||
| 2789 | 314 | 407 | $per_db[$idx]{$db_col} = $params->{$col}; | |||
| 2790 | } else { | |||||
| 2791 | 239 | 156 | carp $self->_err('warn_unknown_column', $col); | |||
| 2792 | } | |||||
| 2793 | } | |||||
| 2794 | ||||||
| 2795 | 796 | 1847 | return \@per_db; | |||
| 2796 | 25 25 25 | 5143 17 125 | } | |||
| 2797 | ||||||
| 2798 | # _fetch_indexed( $db_idx, \%criteria ) -> \%join_val_to_\@rows | |||||
| 2799 | # Purpose: Query one component database and index its rows by join-key value. | |||||
| 2800 | # Entry: $db_idx is the zero-based database index; $criteria is the | |||||
| 2801 | # pre-partitioned criteria hashref for this database. | |||||
| 2802 | # Exit: Returns a hashref: join-key value => arrayref of row hashrefs. | |||||
| 2803 | # Multiple rows sharing the same join-key value are all preserved | |||||
| 2804 | # (important for the primary database when one key maps to many rows). | |||||
| 2805 | # Effects: Calls selectall_arrayref on the component database. | |||||
| 2806 | sub _fetch_indexed :Protected { | |||||
| 2807 | 740 | 553 | my ($self, $db_idx, $criteria) = @_; | |||
| 2808 | ||||||
| 2809 | 740 | 672 | my $db = $self->{_dbs}[$db_idx]; | |||
| 2810 | 781 | 939 | my $local_jc = $self->{_join_map}{$db_idx} // $self->{_join_col}; | |||
| 2811 | ||||||
| 2812 | 694 | 533 | my $rows = $db->selectall_arrayref($criteria); | |||
| 2813 | 690 | 5134 | $rows //= []; | |||
| 2814 | ||||||
| 2815 | 690 | 380 | my %indexed; | |||
| 2816 | 699 699 | 425 463 | for my $row (@{$rows}) { | |||
| 2817 | 1265 | 808 | my $key = $row->{$local_jc}; | |||
| 2818 | 1265 | 1108 | next unless defined $key; | |||
| 2819 | 1264 1081 | 926 1007 | push @{ $indexed{$key} }, $row; | |||
| 2820 | } | |||||
| 2821 | ||||||
| 2822 | 516 | 585 | return \%indexed; | |||
| 2823 | 25 25 25 | 3749 19 111 | } | |||
| 2824 | ||||||
| 2825 | # _joined_query( \%params, %opts ) -> \@merged_rows | |||||
| 2826 | # | |||||
| 2827 | # Purpose: Dispatcher â routes to the array (in-memory) or SQLite join backend | |||||
| 2828 | # based on $self->{_backend}. %opts are passed through to the backend | |||||
| 2829 | # (currently: sort_by, limit, offset). | |||||
| 2830 | sub _joined_query :Protected { | |||||
| 2831 | 379 | 480 | my ($self, $params, %opts) = @_; | |||
| 2832 | 379 | 295 | my $backend = $self->{_backend}; | |||
| 2833 | 379 | 485 | return $self->_joined_query_array($params, %opts) if $backend eq 'array'; | |||
| 2834 | 261 | 400 | return $self->_sqlite_join($params, %opts); | |||
| 2835 | 25 25 25 | 2742 15 93 | } | |||
| 2836 | ||||||
| 2837 | # _joined_query_array( \%params ) -> \@merged_rows | |||||
| 2838 | # | |||||
| 2839 | # Purpose: Core in-memory join algorithm. Partitions criteria, fetches per-database | |||||
| 2840 | # results, resolves the key set, and merges rows. | |||||
| 2841 | # | |||||
| 2842 | # Key-set resolution (applied for each secondary database after the primary): | |||||
| 2843 | # | |||||
| 2844 | # If the database had criteria in this query call (after base filter overlay), | |||||
| 2845 | # it acts as an INNER-JOIN partner: only keys present in its filtered result | |||||
| 2846 | # survive. This gives WHERE-clause semantics even under a LEFT join. | |||||
| 2847 | # | |||||
| 2848 | # If the database had NO effective criteria: | |||||
| 2849 | # inner -> intersect (standard inner join) | |||||
| 2850 | # left -> no change (primary defines the key set) | |||||
| 2851 | # outer -> union (all keys from any database) | |||||
| 2852 | # | |||||
| 2853 | # Row merge: for each qualifying primary row, secondary rows are overlaid in | |||||
| 2854 | # index order. For duplicate columns, later databases win. Local join-key | |||||
| 2855 | # aliases are renamed to the canonical join_column before merging. | |||||
| 2856 | # Removed columns are deleted from every merged row. | |||||
| 2857 | # _validate_pagination( $limit, $offset ) -> ($validated_limit, $validated_offset) | |||||
| 2858 | # | |||||
| 2859 | # Purpose: Single source of truth for pagination parameter validation. | |||||
| 2860 | # Shared by _joined_query_array (array path) and _sqlite_join (SQLite path). | |||||
| 2861 | # Entry: $limit and $offset are raw caller values -- may be undef, negative, | |||||
| 2862 | # non-integer, etc. | |||||
| 2863 | # Exit: Returns the original values when valid; returns undef (with carp) when | |||||
| 2864 | # invalid. Premise: after this call, limit â â¤+ ⪠{undef} and | |||||
| 2865 | # offset â â¤â¥0 ⪠{undef}. All downstream guards may safely rely on this. | |||||
| 2866 | sub _validate_pagination :Protected { | |||||
| 2867 | 316 | 239 | my ($self, $limit, $offset) = @_; | |||
| 2868 | # Syllogism: limit must be a positive integer â§ matches /^\d+\z/a â§ >= 1. | |||||
| 2869 | # Conclusion: any value that fails either check is treated as absent. | |||||
| 2870 | # \z (not $): rejects strings ending with \n that $ would silently accept. | |||||
| 2871 | # /a flag: restricts \d to ASCII [0-9]; rejects Unicode decimal digits | |||||
| 2872 | # (e.g. Arabic-Indic Ù£) that \d matches but Perl's numeric | |||||
| 2873 | # coercion would silently treat as 0, bypassing the >= 1 guard. | |||||
| 2874 | 316 | 278 | if (defined $limit) { | |||
| 2875 | 47 | 114 | if ($limit !~ /^\d+\z/a || $limit < 1) { | |||
| 2876 | 10 | 122 | carp "Database::Join: limit must be a positive integer; ignored"; | |||
| 2877 | 203 | 1564 | undef $limit; | |||
| 2878 | } | |||||
| 2879 | } | |||||
| 2880 | # Syllogism: offset must be a non-negative integer â§ matches /^\d+\z/a. | |||||
| 2881 | # Conclusion: any value that fails the check is treated as absent. | |||||
| 2882 | # Same \z / /a rationale as limit above. | |||||
| 2883 | 509 | 436 | if (defined $offset) { | |||
| 2884 | 233 | 153 | if ($offset !~ /^\d+\z/a) { | |||
| 2885 | 198 | 531 | carp "Database::Join: offset must be a non-negative integer; ignored"; | |||
| 2886 | 27 | 781 | undef $offset; | |||
| 2887 | } | |||||
| 2888 | } | |||||
| 2889 | 338 | 351 | return ($limit, $offset); | |||
| 2890 | 25 25 25 | 3859 17 111 | } | |||
| 2891 | ||||||
| 2892 | sub _joined_query_array :Protected { | |||||
| 2893 | 317 | 331 | my ($self, $params, %opts) = @_; | |||
| 2894 | 317 | 228 | my $sort_by = $opts{sort_by}; | |||
| 2895 | 317 | 207 | my $limit = $opts{limit}; | |||
| 2896 | 317 | 211 | my $offset = $opts{offset}; | |||
| 2897 | ||||||
| 2898 | # Fast path: skip Sub::Protected dispatch entirely when no pagination params are | |||||
| 2899 | # supplied (the common case). Saves one method-lookup overhead per query. | |||||
| 2900 | 317 | 417 | ($limit, $offset) = $self->_validate_pagination($limit, $offset) | |||
| 2901 | if defined $limit || defined $offset; | |||||
| 2902 | ||||||
| 2903 | 295 | 243 | my $join_col = $self->{_join_col}; | |||
| 2904 | 295 | 193 | my $join_type = $self->{_join_type}; | |||
| 2905 | 317 317 | 174 248 | my $n = scalar @{ $self->{_dbs} }; | |||
| 2906 | ||||||
| 2907 | 466 | 398 | my $per_db = $self->_partition_criteria($params); | |||
| 2908 | ||||||
| 2909 | # Overlay any per-database base filters onto the partitioned criteria. | |||||
| 2910 | # A filtered database always has effective criteria, so $had_criteria will | |||||
| 2911 | # be true for it â giving inner-join key-set semantics regardless of join_type. | |||||
| 2912 | 466 | 378 | for my $i (0 .. $n - 1) { | |||
| 2913 | 700 | 743 | my $base = $self->{_filters}{$i} // {}; | |||
| 2914 | 693 690 | 460 1045 | next unless %{$base}; | |||
| 2915 | 36 | 55 | $per_db->[$i] = _merge_criteria($base, $per_db->[$i]); | |||
| 2916 | } | |||||
| 2917 | ||||||
| 2918 | # Fetch and index each database with its own criteria slice. | |||||
| 2919 | 389 | 301 | my @indexed; | |||
| 2920 | 383 | 2011 | $indexed[0] = $self->_fetch_indexed(0, $per_db->[0]); | |||
| 2921 | ||||||
| 2922 | # Early exit: for inner and left joins, an empty primary result means the | |||||
| 2923 | # key set is provably empty (left: primary defines it; inner: â© â = â ). | |||||
| 2924 | # Skipping secondary fetches avoids up to N-1 unnecessary DA round-trips. | |||||
| 2925 | # | |||||
| 2926 | # Two guards prevent premature exit: | |||||
| 2927 | # (a) join-column broadcast: when a join-col criterion is present it must | |||||
| 2928 | # be physically delivered to each secondary DA (the call itself is what | |||||
| 2929 | # forwards it; the partition only prepared the per-db slice). | |||||
| 2930 | # (b) secondary-owned criteria: a secondary with its own criteria (e.g. | |||||
| 2931 | # score => $val) must still be queried so those criteria are delivered. | |||||
| 2932 | # Without the call the DA never receives them â breaking the partition- | |||||
| 2933 | # isolation invariant the security tests verify. | |||||
| 2934 | 382 | 437 | my $local_jc_0_early = $self->{_join_map}{0} // $join_col; | |||
| 2935 | 382 308 308 | 348 239 292 | my $sec_has_criteria = grep { %{ $per_db->[$_] } } 1 .. $n - 1; | |||
| 2936 | 382 | 390 | return [] if !%{ $indexed[0] } | |||
| 2937 | && $join_type ne 'outer' | |||||
| 2938 | 382 | 237 | && !exists $per_db->[0]{$local_jc_0_early} | |||
| 2939 | && !$sec_has_criteria; | |||||
| 2940 | ||||||
| 2941 | # Fetch and index secondary databases. | |||||
| 2942 | # With parallel => 1 and 2+ secondaries: spawn one Perl thread per secondary | |||||
| 2943 | # so all secondaries are queried concurrently. Total latency becomes | |||||
| 2944 | # max(DA latencies) instead of sum(DA latencies). The primary was already | |||||
| 2945 | # fetched sequentially above (needed for the early-exit short-circuit). | |||||
| 2946 | # The thread closure captures only $db, $crit, and $local_jc (plain scalars | |||||
| 2947 | # or in-memory data) â $self is intentionally not captured to avoid copying | |||||
| 2948 | # the full blessed hashref (including all DA refs) into each thread. | |||||
| 2949 | # Falls back to sequential when: threads module is unavailable, n <= 2, or | |||||
| 2950 | # parallel => 0 (the default). | |||||
| 2951 | # | |||||
| 2952 | # Windows / ithreads safety: cloned DBI handles inside DA objects are not | |||||
| 2953 | # thread-safe. Pre-initialise every secondary slot to {} so that if a thread | |||||
| 2954 | # fails the slot holds a valid (empty) hashref rather than undef, which would | |||||
| 2955 | # corrupt the inner-join key-set resolution. Any secondary whose thread does | |||||
| 2956 | # not deliver a valid hashref is re-fetched sequentially after all joins | |||||
| 2957 | # complete, preserving correctness on every platform. | |||||
| 2958 | 377 | 461 | $indexed[$_] = {} for 1 .. $n - 1; | |||
| 2959 | ||||||
| 2960 | 377 | 405 | if ($self->{_parallel} && $n > 2) { | |||
| 2961 | 96 168 168 168 22 | 108 163 85 1519 19 | $HAS_THREADS //= do { local $@; eval { require threads; 1 } ? 1 : 0 }; | |||
| 2962 | 30 | 2239 | if ($HAS_THREADS) { | |||
| 2963 | 22 | 12 | my @thr; | |||
| 2964 | 22 | 39 | for my $i (1 .. $n - 1) { | |||
| 2965 | my ($db, $crit, $local_jc) = ( | |||||
| 2966 | $self->{_dbs}[$i], $per_db->[$i], | |||||
| 2967 | 22 | 37 | $self->{_join_map}{$i} // $join_col, | |||
| 2968 | ); | |||||
| 2969 | # Wrap thread creation: it can fail if the DA cannot be cloned | |||||
| 2970 | # (e.g. DBI handles on Windows). On failure we skip the push so | |||||
| 2971 | # the sequential fallback below handles that secondary. | |||||
| 2972 | 22 | 17 | local $@; | |||
| 2973 | 22 | 38 | my $t = eval { | |||
| 2974 | threads->create(sub { | |||||
| 2975 | # eval inside the thread: prevents a DA exception from | |||||
| 2976 | # killing the thread and returning an empty list to join(). | |||||
| 2977 | 19 | 22 | local $@; | |||
| 2978 | 19 9 | 63 33 | my $rows = eval { $db->selectall_arrayref($crit) } // []; | |||
| 2979 | 9 | 12 | my %idx; | |||
| 2980 | 1 1 | 23 246 | for my $row (@{$rows}) { | |||
| 2981 | 8 | 5 | my $key = $row->{$local_jc}; | |||
| 2982 | 8 8 | 9 9 | push @{$idx{$key}}, $row if defined $key; | |||
| 2983 | } | |||||
| 2984 | 4 | 7 | return ($i, \%idx); | |||
| 2985 | 19 | 11 | }); | |||
| 2986 | }; | |||||
| 2987 | 4 | 47 | push @thr, $t if $t; | |||
| 2988 | } | |||||
| 2989 | ||||||
| 2990 | 4 | 5 | my %done; | |||
| 2991 | 4 | 4 | for my $t (@thr) { | |||
| 2992 | # eval on join: a thread that died (e.g. uncaught exception) | |||||
| 2993 | # causes join() to rethrow on some platforms. | |||||
| 2994 | 4 | 8 | local $@; | |||
| 2995 | 10 | 47 | my ($i, $idx); | |||
| 2996 | 3 3 | 11 6 | eval { ($i, $idx) = $t->join() }; | |||
| 2997 | 7 | 24 | if (defined $i && ref($idx) eq 'HASH') { | |||
| 2998 | 0 | 0 | $indexed[$i] = $idx; | |||
| 2999 | 0 | 0 | $done{$i} = 1; | |||
| 3000 | } | |||||
| 3001 | } | |||||
| 3002 | ||||||
| 3003 | # Re-fetch sequentially any secondary not delivered by a thread. | |||||
| 3004 | # This covers: thread creation failure, thread death, or an empty | |||||
| 3005 | # result that may indicate a non-thread-safe DA (e.g. DBI on Windows). | |||||
| 3006 | 7 | 27 | for my $i (1 .. $n - 1) { | |||
| 3007 | 7 | 12 | next if $done{$i}; | |||
| 3008 | 1 | 2 | $indexed[$i] = $self->_fetch_indexed($i, $per_db->[$i]); | |||
| 3009 | } | |||||
| 3010 | } else { | |||||
| 3011 | 10 | 81 | carp 'Database::Join: parallel => 1 requires the threads module; falling back to sequential'; | |||
| 3012 | 10 | 889 | $indexed[$_] = $self->_fetch_indexed($_, $per_db->[$_]) for 1 .. $n - 1; | |||
| 3013 | } | |||||
| 3014 | } else { | |||||
| 3015 | 369 | 426 | $indexed[$_] = $self->_fetch_indexed($_, $per_db->[$_]) for 1 .. $n - 1; | |||
| 3016 | } | |||||
| 3017 | ||||||
| 3018 | # Premise: the key-set resolution loop starts at i=1 (primary seeds %key_set). | |||||
| 3019 | # Conclusion: $had_criteria[0] is a dead store (D~); compute only for i >= 1. | |||||
| 3020 | # | |||||
| 3021 | # The broadcast join-column criterion (entry=>'A3' delivered to ALL databases) | |||||
| 3022 | # must NOT count as "had criteria" for secondaries. It is a key-range selector | |||||
| 3023 | # on the merged view, not a predicate that bounds what the secondary contributes. | |||||
| 3024 | # Only base filters and non-join-column query-time criteria trigger inner-join. | |||||
| 3025 | 374 | 348 | my @had_criteria; | |||
| 3026 | 374 | 299 | for my $i (1 .. $n - 1) { | |||
| 3027 | 300 | 300 | my $local_jc = $self->{_join_map}{$i} // $join_col; | |||
| 3028 | 300 287 | 190 377 | my $has_filter = !!%{ $self->{_filters}{$i} // {} }; | |||
| 3029 | 287 | 337 | my $crit = $per_db->[$i]; | |||
| 3030 | # Avoid copying the criteria hash just to delete the join-col key. | |||||
| 3031 | # Arithmetic is O(1) allocations: count total keys, subtract 1 when the | |||||
| 3032 | # local join-col key is present. Any result > 0 is truthy (has own criteria). | |||||
| 3033 | 287 300 | 183 302 | my $own = (keys %{$crit}) - (exists $crit->{$local_jc} ? 1 : 0); | |||
| 3034 | 217 | 347 | $had_criteria[$i] = $has_filter || $own; | |||
| 3035 | } | |||||
| 3036 | ||||||
| 3037 | # Seed the key set from the primary database. | |||||
| 3038 | # Hash-slice assignment avoids the intermediate 2K-element flat list that | |||||
| 3039 | # map { $_ => 1 } would allocate before assigning to %key_set. Values are | |||||
| 3040 | # undef; only exists() is used for lookups, so the sentinel value is irrelevant. | |||||
| 3041 | 291 | 422 | my %key_set; | |||
| 3042 | 369 369 | 235 465 | @key_set{ keys %{ $indexed[0] } } = (); | |||
| 3043 | ||||||
| 3044 | # Merge in each secondary database. | |||||
| 3045 | # Premise 1: indexed[$i] is a valid hashref (returned by _fetch_indexed). | |||||
| 3046 | # Premise 2: join_type â {left, inner, outer} (enforced by validate_strict). | |||||
| 3047 | # Conclusion: the three branches below are exhaustive and mutually exclusive. | |||||
| 3048 | 369 | 292 | for my $i (1 .. $n - 1) { | |||
| 3049 | 295 | 442 | if ($had_criteria[$i] || $join_type eq 'inner') { | |||
| 3050 | # Intersect: single-pass delete for keys absent from this secondary. | |||||
| 3051 | # A single loop avoids the intermediate list that grep would allocate | |||||
| 3052 | # before the delete loop could iterate it (saves O(K) allocations). | |||||
| 3053 | 77 | 68 | for my $k (keys %key_set) { | |||
| 3054 | 170 | 189 | delete $key_set{$k} unless exists $indexed[$i]{$k}; | |||
| 3055 | } | |||||
| 3056 | } elsif ($join_type eq 'outer') { | |||||
| 3057 | # Union: hash slice assignment is a single Perl op, not a per-key loop. | |||||
| 3058 | 14 94 | 12 72 | @key_set{ keys %{ $indexed[$i] } } = (); | |||
| 3059 | } | |||||
| 3060 | # left + no criteria: key_set unchanged (primary defines the set). | |||||
| 3061 | } | |||||
| 3062 | ||||||
| 3063 | # Pre-hoist per-secondary constants outside the key loop. | |||||
| 3064 | # $sec_local_jc[$i], the rename flag, and $sec_renames[$i] are all invariant | |||||
| 3065 | # across every key and every primary row. Computing them inside the key loop | |||||
| 3066 | # wastes K dereferences per secondary database (K = number of qualifying keys). | |||||
| 3067 | # Splitting the inner column loop on the rename flag eliminates the flag check | |||||
| 3068 | # from inside the per-column loop, saving R-1 branch evaluations per secondary | |||||
| 3069 | # per row (R = columns in the secondary row). | |||||
| 3070 | 368 | 295 | my (@sec_local_jc, @sec_rename, @sec_renames); | |||
| 3071 | 369 | 310 | for my $i (1 .. $n - 1) { | |||
| 3072 | 295 | 235 | $sec_local_jc[$i] = $self->{_join_map}{$i}; | |||
| 3073 | 295 | 296 | $sec_rename[$i] = ($sec_local_jc[$i] && $sec_local_jc[$i] ne $join_col) ? 1 : 0; | |||
| 3074 | # _col_rename[$i] is always initialised to {} by _build_col_index / add_database | |||||
| 3075 | # (transitive reduction: the // {} fallback can never trigger). | |||||
| 3076 | 295 | 257 | $sec_renames[$i] = $self->{_col_rename}[$i]; | |||
| 3077 | } | |||||
| 3078 | ||||||
| 3079 | # Cache the removed-column list across calls; avoids extracting keys %hash every | |||||
| 3080 | # query. Lazily built here and invalidated to undef by remove_column(). | |||||
| 3081 | 498 370 | 507 429 | my $removed = ($self->{_removed_list} //= [keys %{ $self->{_removed_cols} }]); | |||
| 3082 | ||||||
| 3083 | # Parse sort_by once before the merge loop so we can decide whether the | |||||
| 3084 | # initial O(K log K) sort of %key_set is necessary. | |||||
| 3085 | # When sort_by targets a column other than the join_col, that initial sort | |||||
| 3086 | # is overridden by the final Schwarzian pass -- skip it to save a full sort. | |||||
| 3087 | # When sort_by is absent or targets the join_col, the initial sort IS the | |||||
| 3088 | # final order and must be kept. | |||||
| 3089 | 415 | 342 | my ($ob_col, $ob_dir) = ($join_col, 'ASC'); | |||
| 3090 | 415 | 330 | if (defined $sort_by) { | |||
| 3091 | 117 88 | 138 107 | my ($req_col, $req_dir) = ref($sort_by) eq 'ARRAY' ? @{$sort_by} : ($sort_by, 'ASC'); | |||
| 3092 | 104 | 102 | $req_dir = uc($req_dir // 'ASC'); | |||
| 3093 | 104 | 120 | unless ($req_dir eq 'ASC' || $req_dir eq 'DESC') { | |||
| 3094 | 73 | 143 | carp "Database::Join: sort_by direction '$req_dir' is not supported; using ASC"; | |||
| 3095 | 144 | 445 | $req_dir = 'ASC'; | |||
| 3096 | } | |||||
| 3097 | 105 | 100 | if ($req_col ne $join_col && !exists $self->{_col_db}{$req_col}) { | |||
| 3098 | 76 | 127 | carp "Database::Join: sort_by column '$req_col' is not in the merged view; result sorted by join_column"; | |||
| 3099 | } else { | |||||
| 3100 | 100 | 88 | ($ob_col, $ob_dir) = ($req_col, $req_dir); | |||
| 3101 | } | |||||
| 3102 | } | |||||
| 3103 | # True when sort_by targets a non-join column: the initial sort is redundant. | |||||
| 3104 | 357 | 920 | my $ob_override = ($ob_col ne $join_col); | |||
| 3105 | ||||||
| 3106 | # Build one merged result row for every primary-database row that qualifies. | |||||
| 3107 | # Secondary databases act as lookup tables: when a key maps to multiple | |||||
| 3108 | # secondary rows, the last one wins (consistent with construction-time | |||||
| 3109 | # last-database-wins column routing). | |||||
| 3110 | 357 | 213 | my @result; | |||
| 3111 | # When sort_by will override the join_col order, iterate keys unsorted (O(K)) | |||||
| 3112 | # instead of sorted (O(K log K)); the Schwarzian pass at the end reorders. | |||||
| 3113 | 357 | 526 | for my $key ($ob_override ? keys %key_set : sort keys %key_set) { | |||
| 3114 | # Iterate directly over the arrayref: avoids copying primary rows into a | |||||
| 3115 | # new @base_rows array (saves P element copies per key, P = rows per key). | |||||
| 3116 | # [{}] ensures outer-join keys absent from the primary produce one merged row. | |||||
| 3117 | 719 731 | 398 612 | for my $prow (@{ $indexed[0]{$key} // [{}] }) { | |||
| 3118 | 655 655 | 311 722 | my %merged = %{$prow}; | |||
| 3119 | ||||||
| 3120 | 655 | 1320 | for my $i (1 .. $n - 1) { | |||
| 3121 | 389 | 8417 | my $sec_arr = $indexed[$i]{$key}; | |||
| 3122 | 460 412 | 374 381 | next unless $sec_arr && @{$sec_arr}; | |||
| 3123 | ||||||
| 3124 | # Write secondary columns directly into %merged without copying | |||||
| 3125 | # the source row into a temporary hash first. | |||||
| 3126 | # Before: %row_copy = %{$src} then %merged = (%merged,%row_copy) | |||||
| 3127 | # â 2 full hash copies per secondary per row: O(C) + O(|merged|+C) | |||||
| 3128 | # After: per-key loop writes straight into %merged | |||||
| 3129 | # â O(C) key assignments only; no intermediate allocation | |||||
| 3130 | 342 | 202 | my $src = $sec_arr->[-1]; | |||
| 3131 | ||||||
| 3132 | 339 | 244 | if ($sec_rename[$i]) { | |||
| 3133 | 14 | 20 | my $local_jc = $sec_local_jc[$i]; | |||
| 3134 | 14 | 22 | my $renames = $sec_renames[$i]; | |||
| 3135 | 8 8 | 26 237 | for my $k (keys %{$src}) { | |||
| 3136 | 21 | 38 | if ($k eq $local_jc) { | |||
| 3137 | # Translate local join-key alias to the canonical join_column name | |||||
| 3138 | 8 | 30 | $merged{$join_col} = $src->{$k}; | |||
| 3139 | } elsif (my $pub = $renames->{$k}) { | |||||
| 3140 | 6 | 8 | $merged{$pub} = $src->{$k}; | |||
| 3141 | } else { | |||||
| 3142 | 84 | 378 | $merged{$k} = $src->{$k}; | |||
| 3143 | } | |||||
| 3144 | } | |||||
| 3145 | } else { | |||||
| 3146 | 405 | 312 | my $renames = $sec_renames[$i]; | |||
| 3147 | 405 405 | 211 352 | for my $k (keys %{$src}) { | |||
| 3148 | 701 | 613 | if (my $pub = $renames->{$k}) { | |||
| 3149 | # Collision-renamed column: write under the published prefixed name | |||||
| 3150 | 26 | 26 | $merged{$pub} = $src->{$k}; | |||
| 3151 | } else { | |||||
| 3152 | 680 | 602 | $merged{$k} = $src->{$k}; | |||
| 3153 | } | |||||
| 3154 | } | |||||
| 3155 | } | |||||
| 3156 | } | |||||
| 3157 | ||||||
| 3158 | 652 25 650 | 366 21 472 | delete @merged{@{$removed}} if @{$removed}; | |||
| 3159 | 650 | 533 | push @result, \%merged; | |||
| 3160 | } | |||||
| 3161 | } | |||||
| 3162 | ||||||
| 3163 | # Caller-specified ORDER BY. | |||||
| 3164 | # $ob_col / $ob_dir / $ob_override were parsed BEFORE the merge loop. | |||||
| 3165 | # Three cases: | |||||
| 3166 | # 1. Non-join-col override ($ob_override true): Schwarzian transform | |||||
| 3167 | # O(R) key extractions + O(R log R) scalar comparisons -- cheaper than | |||||
| 3168 | # O(2R log R) hash dereferences that a naive sort block would make. | |||||
| 3169 | # 2. join_col DESC: O(R) reverse -- already sorted ASC by the loop. | |||||
| 3170 | # 3. join_col ASC (default): already in order, nothing to do. | |||||
| 3171 | 363 | 515 | if ($ob_override) { | |||
| 3172 | # Schwarzian: decorate, sort, undecorate. | |||||
| 3173 | # String comparison (cmp). For accurate numeric ordering on large | |||||
| 3174 | # numeric columns use backend => 'sqlite', which sorts by SQL type. | |||||
| 3175 | 101 149 | 5265 343 | my @tagged = map { [$_, $_->{$ob_col} // ''] } @result; | |||
| 3176 | 72 | 65 | if ($ob_dir eq 'DESC') { | |||
| 3177 | 52 59 57 | 32 61 41 | @result = map { $_->[0] } sort { $b->[1] cmp $a->[1] } @tagged; | |||
| 3178 | } else { | |||||
| 3179 | 68 84 69 | 95 103 62 | @result = map { $_->[0] } sort { $a->[1] cmp $b->[1] } @tagged; | |||
| 3180 | } | |||||
| 3181 | } elsif ($ob_dir eq 'DESC') { | |||||
| 3182 | # join_col DESC: O(R) reverse is cheaper than O(R log R) re-sort | |||||
| 3183 | # because the merge loop already produced join_col ASC order. | |||||
| 3184 | 17 | 15 | @result = reverse @result; | |||
| 3185 | } | |||||
| 3186 | # join_col ASC: @result is already in join_column ascending order. | |||||
| 3187 | ||||||
| 3188 | # LIMIT / OFFSET pagination â applied after ordering. | |||||
| 3189 | # splice removes elements from the front (offset) then truncates to limit. | |||||
| 3190 | 299 | 277 | if (defined $offset && $offset > 0) { | |||
| 3191 | 27 | 538 | splice(@result, 0, $offset); | |||
| 3192 | } | |||||
| 3193 | 334 | 325 | if (defined $limit && $limit < scalar @result) { | |||
| 3194 | 18 | 19 | splice(@result, $limit); | |||
| 3195 | } | |||||
| 3196 | ||||||
| 3197 | 286 | 1054 | return \@result; | |||
| 3198 | 25 25 25 | 20106 18 194 | } | |||
| 3199 | ||||||
| 3200 | # _cache_fresh() -> bool | |||||
| 3201 | # Purpose: Check whether the SQLite join cache is still valid. | |||||
| 3202 | # Entry: $self->{_sqlite_cache} may or may not be set. | |||||
| 3203 | # Exit: Returns 1 if the cache exists, the DBI handle is active, the source | |||||
| 3204 | # count matches, and all source updated() timestamps match. Returns 0 | |||||
| 3205 | # if any of these conditions fail (caller must rebuild the cache). | |||||
| 3206 | sub _cache_fresh :Protected { | |||||
| 3207 | 112 | 429 | my ($self) = @_; | |||
| 3208 | ||||||
| 3209 | 112 | 234 | my $cache = $self->{_sqlite_cache} // return 0; | |||
| 3210 | 34 34 | 18 37 | my $n = scalar @{ $self->{_dbs} }; | |||
| 3211 | ||||||
| 3212 | 34 | 66 | return 0 if ($cache->{n} // 0) != $n; | |||
| 3213 | ||||||
| 3214 | # Verify the DBI handle is still usable. | |||||
| 3215 | 33 33 33 33 | 21 39 25 152 | return 0 unless do { local $@; eval { $cache->{dbh}{Active} } }; | |||
| 3216 | ||||||
| 3217 | # Verify that no source has been updated since the cache was built. | |||||
| 3218 | # If a source does not implement updated(), skip the timestamp check for | |||||
| 3219 | # it (the data is assumed stable; the cache stays valid indefinitely for | |||||
| 3220 | # that source unless add_database() is called or the object is destroyed). | |||||
| 3221 | 32 | 40 | for my $i (0 .. $n - 1) { | |||
| 3222 | 55 | 69 | my $cached_ts = $cache->{updated}{$i} // next; # not captured â skip | |||
| 3223 | 29 | 13 | my $current_ts; | |||
| 3224 | 29 29 29 29 | 42 15 19 28 | do { local $@; $current_ts = eval { $self->{_dbs}[$i]->updated() } }; | |||
| 3225 | 29 | 50 | next unless defined $current_ts; # no updated() â skip | |||
| 3226 | 29 | 41 | return 0 if $current_ts != $cached_ts; | |||
| 3227 | } | |||||
| 3228 | ||||||
| 3229 | 27 | 39 | return 1; | |||
| 3230 | 25 25 25 | 4566 19 142 | } | |||
| 3231 | ||||||
| 3232 | # _build_sqlite_cache() | |||||
| 3233 | # Purpose: Create (or rebuild) the persistent SQLite join cache. Spills each | |||||
| 3234 | # source database into a temp SQLite file using filter-only criteria; | |||||
| 3235 | # SQLite-backed sources are zero-copy ATTACHed instead of spilled. | |||||
| 3236 | # Query-time criteria are NOT applied here â they become WHERE clauses | |||||
| 3237 | # in the per-call SQL generated by _sqlite_join. | |||||
| 3238 | # _sql_quote_identifier( $name ) -> $quoted | |||||
| 3239 | # Purpose: Produce a properly double-quoted SQL identifier, escaping any | |||||
| 3240 | # embedded double-quote characters by doubling them (SQL standard). | |||||
| 3241 | # Defence-in-depth: prevents SQL identifier injection when DA-supplied | |||||
| 3242 | # column names or table names contain literal double-quote characters. | |||||
| 3243 | # SQLite, like all ANSI SQL databases, represents a literal " inside a | |||||
| 3244 | # double-quoted identifier as ""; this routine applies that transform. | |||||
| 3245 | # Entry: $name â raw identifier string (column name, table name, or alias). | |||||
| 3246 | # Exit: Returns the double-quoted, injection-safe SQL identifier string. | |||||
| 3247 | sub _sql_quote_identifier { | |||||
| 3248 | 4822 | 3481 | my ($name) = @_; | |||
| 3249 | 4822 | 3849 | (my $safe = $name) =~ s/"/""/g; | |||
| 3250 | 4822 | 8725 | return "\"$safe\""; | |||
| 3251 | } | |||||
| 3252 | ||||||
| 3253 | # Entry: _dbs, _join_map, _filters, _tmpdir must be set. | |||||
| 3254 | # Exit: $self->{_sqlite_cache} holds {dbh, tmpfile, table_refs, source_cols, | |||||
| 3255 | # is_attached, updated, n}. Any previous cache is disconnected first. | |||||
| 3256 | # Each spilled table has a B-tree index on its join column. | |||||
| 3257 | # Effects: Creates a File::Temp file (SUFFIX='.db', DIR=_tmpdir, UNLINK=1). | |||||
| 3258 | # Croaks with error_sqlite_connect if DBI::connect fails. | |||||
| 3259 | sub _build_sqlite_cache :Protected { | |||||
| 3260 | my ($self) = @_; | |||||
| 3261 | ||||||
| 3262 | # Disconnect any previous cache to release the old temp file. | |||||
| 3263 | if (my $old = delete $self->{_sqlite_cache}) { | |||||
| 3264 | local $@; | |||||
| 3265 | eval { $old->{dbh}->disconnect } if $old->{dbh}; | |||||
| 3266 | } | |||||
| 3267 | ||||||
| 3268 | require DBI; | |||||
| 3269 | require File::Temp; | |||||
| 3270 | ||||||
| 3271 | my $join_col = $self->{_join_col}; | |||||
| 3272 | my $n = scalar @{ $self->{_dbs} }; | |||||
| 3273 | ||||||
| 3274 | my $tmpfile = File::Temp->new( | |||||
| 3275 | SUFFIX => '.db', | |||||
| 3276 | DIR => $self->{_tmpdir}, | |||||
| 3277 | UNLINK => 1, | |||||
| 3278 | ); | |||||
| 3279 | ||||||
| 3280 | my $tmpdbh = DBI->connect( | |||||
| 3281 | 'dbi:SQLite:dbname=' . $tmpfile->filename, '', '', | |||||
| 3282 | { RaiseError => 1, PrintError => 0, AutoCommit => 1 }, | |||||
| 3283 | ) or croak $self->_err('error_sqlite_connect', DBI->errstr // 'unknown error'); | |||||
| 3284 | ||||||
| 3285 | my (@table_refs, @source_cols, @is_attached); | |||||
| 3286 | ||||||
| 3287 | for my $i (0 .. $n - 1) { | |||||
| 3288 | my $db = $self->{_dbs}[$i]; | |||||
| 3289 | my $local_jc = $self->{_join_map}{$i} // $join_col; | |||||
| 3290 | ||||||
| 3291 | # Zero-copy ATTACH path: unconditionally available when the source | |||||
| 3292 | # implements dbi_source() returning a live SQLite handle. Query-time | |||||
| 3293 | # criteria for this source will go into the SQL WHERE clause. | |||||
| 3294 | # eval wraps can() to suppress ISA warnings from stub packages in tests. | |||||
| 3295 | if (do { local $@; eval { $db->can('dbi_source') } }) { | |||||
| 3296 | # Probe dbi_source() in an isolated scope: local $@ prevents leaking | |||||
| 3297 | # eval failure into the caller's $@; local $SIG{__WARN__} suppresses | |||||
| 3298 | # spurious warnings from inherited dbi_source() probing non-SQLite DAs | |||||
| 3299 | # (e.g. Database::Abstraction 0.46 base-class dbi_source() calling | |||||
| 3300 | # _open() on a DA with no backing file). | |||||
| 3301 | my $src = do { | |||||
| 3302 | local $SIG{__WARN__} = sub {}; | |||||
| 3303 | local $@; | |||||
| 3304 | eval { $db->dbi_source() } | |||||
| 3305 | }; | |||||
| 3306 | if ($src && ref($src) eq 'HASH' && $src->{dbh} && $src->{table} | |||||
| 3307 | && eval { $src->{dbh}{Driver}{Name} } eq 'SQLite') { | |||||
| 3308 | my ($db_file) = $src->{dbh}->selectrow_array( | |||||
| 3309 | "SELECT file FROM pragma_database_list WHERE name='main'" | |||||
| 3310 | ); | |||||
| 3311 | my $alias = "ext$i"; | |||||
| 3312 | $tmpdbh->do(sprintf("ATTACH DATABASE %s AS %s", | |||||
| 3313 | $tmpdbh->quote($db_file), $alias)); | |||||
| 3314 | $table_refs[$i] = $alias . '.' . _sql_quote_identifier($src->{table}); | |||||
| 3315 | $is_attached[$i] = 1; | |||||
| 3316 | # Transitive reduction: new() and add_database() both validate | |||||
| 3317 | # can('columns') before registering any DA (P1 invariant). | |||||
| 3318 | # The else branch is dead code; the guard is vacuous. | |||||
| 3319 | $source_cols[$i] = $db->columns(); | |||||
| 3320 | next; | |||||
| 3321 | } | |||||
| 3322 | } | |||||
| 3323 | ||||||
| 3324 | # Spill path: fetch rows using filter-only criteria. Query-time | |||||
| 3325 | # criteria are NOT applied here â they become WHERE clauses per call. | |||||
| 3326 | my $filter_crit = $self->{_filters}{$i} // {}; | |||||
| 3327 | my $rows = $db->selectall_arrayref($filter_crit) // []; | |||||
| 3328 | ||||||
| 3329 | # No need to check if $rows exists or not | |||||
| 3330 | # Transitive reduction (P1 invariant): can('columns') is guaranteed for | |||||
| 3331 | # all _dbs elements | |||||
| 3332 | $source_cols[$i] = $db->columns(); | |||||
| 3333 | ||||||
| 3334 | my $tbl = "t$i"; | |||||
| 3335 | my $cols = $source_cols[$i] // []; | |||||
| 3336 | ||||||
| 3337 | my $col_defs = join(', ', map { _sql_quote_identifier($_) . ' TEXT' } @{$cols}); | |||||
| 3338 | $tmpdbh->do('CREATE TABLE ' . _sql_quote_identifier($tbl) . " ($col_defs)"); | |||||
| 3339 | # Index on the join column: upgrades ON-clause equality lookups from | |||||
| 3340 | # an O(N²) full-table nested-loop scan to O(N log N) b-tree seek. | |||||
| 3341 | # SQLite query planner uses it for INNER JOIN / LEFT JOIN ON expressions. | |||||
| 3342 | $tmpdbh->do('CREATE INDEX ' . _sql_quote_identifier("${tbl}_jc") | |||||
| 3343 | . ' ON ' . _sql_quote_identifier($tbl) | |||||
| 3344 | . ' (' . _sql_quote_identifier($local_jc) . ')'); | |||||
| 3345 | ||||||
| 3346 | if (@{$rows}) { | |||||
| 3347 | my $col_list = join(', ', map { _sql_quote_identifier($_) } @{$cols}); | |||||
| 3348 | my $placeholders = join(', ', ('?') x scalar @{$cols}); | |||||
| 3349 | my $sth = $tmpdbh->prepare( | |||||
| 3350 | 'INSERT INTO ' . _sql_quote_identifier($tbl) . " ($col_list) VALUES ($placeholders)" | |||||
| 3351 | ); | |||||
| 3352 | my $batch = 0; | |||||
| 3353 | $tmpdbh->begin_work; | |||||
| 3354 | for my $row (@{$rows}) { | |||||
| 3355 | $sth->execute(map { $row->{$_} } @{$cols}); | |||||
| 3356 | if (++$batch >= 1_000) { | |||||
| 3357 | $tmpdbh->commit; | |||||
| 3358 | $tmpdbh->begin_work; | |||||
| 3359 | $batch = 0; | |||||
| 3360 | } | |||||
| 3361 | } | |||||
| 3362 | $tmpdbh->commit; | |||||
| 3363 | } | |||||
| 3364 | $table_refs[$i] = _sql_quote_identifier($tbl); | |||||
| 3365 | $is_attached[$i] = 0; | |||||
| 3366 | } | |||||
| 3367 | ||||||
| 3368 | # Snapshot updated() timestamps for cache-validity checks. | |||||
| 3369 | my %updated; | |||||
| 3370 | for my $i (0 .. $n - 1) { | |||||
| 3371 | local $@; | |||||
| 3372 | my $ts = eval { $self->{_dbs}[$i]->updated() }; | |||||
| 3373 | $updated{$i} = $ts unless $@; | |||||
| 3374 | } | |||||
| 3375 | ||||||
| 3376 | $self->{_sqlite_cache} = { | |||||
| 3377 | dbh => $tmpdbh, | |||||
| 3378 | tmpfile => $tmpfile, | |||||
| 3379 | table_refs => \@table_refs, | |||||
| 3380 | source_cols => \@source_cols, | |||||
| 3381 | is_attached => \@is_attached, | |||||
| 3382 | updated => \%updated, | |||||
| 3383 | n => $n, | |||||
| 3384 | }; | |||||
| 3385 | ||||||
| 3386 | return; | |||||
| 3387 | 25 25 25 | 12373 19 167 | } | |||
| 3388 | ||||||
| 3389 | # _sqlite_join( \%params ) -> \@merged_rows | |||||
| 3390 | # | |||||
| 3391 | # Purpose: Join via a persistent SQLite database cache. The first call (or | |||||
| 3392 | # any call after a source updated() changes) spills source data into | |||||
| 3393 | # a File::Temp SQLite file via _build_sqlite_cache; subsequent calls | |||||
| 3394 | # reuse the same file and handle. A single SQL JOIN with a per-call | |||||
| 3395 | # WHERE clause (built from query-time criteria) produces the result. | |||||
| 3396 | # For 'auto' mode, uses count() or dbi_source() COUNT(*) to check | |||||
| 3397 | # the threshold without fetching rows; falls back to | |||||
| 3398 | # _joined_query_array when count <= $self->{_max_array_rows} or when | |||||
| 3399 | # no count method is available. | |||||
| 3400 | # Entry: $params is the query criteria hashref. | |||||
| 3401 | # Exit: Returns arrayref of merged hashrefs sorted by join_column. | |||||
| 3402 | # Effects: On the first call (or after cache invalidation), creates a | |||||
| 3403 | # File::Temp SQLite file in _tmpdir; the file persists until the | |||||
| 3404 | # Database::Join object is destroyed or the source data changes. | |||||
| 3405 | sub _sqlite_join :Protected { | |||||
| 3406 | my ($self, $params, %opts) = @_; | |||||
| 3407 | my $count_only = $opts{count_only} // 0; | |||||
| 3408 | my $sort_by = $opts{sort_by}; | |||||
| 3409 | my $limit = $opts{limit}; | |||||
| 3410 | my $offset = $opts{offset}; | |||||
| 3411 | my $create_table = $opts{create_table}; # when set, materialize into a real table | |||||
| 3412 | ||||||
| 3413 | # Transitive reduction: _validate_pagination is the single validation site. | |||||
| 3414 | # count() and dbi_source() never supply limit/offset (both methods delete them | |||||
| 3415 | # before calling _sqlite_join), so for those callers this is a cheap undef-check. | |||||
| 3416 | ($limit, $offset) = $self->_validate_pagination($limit, $offset); | |||||
| 3417 | ||||||
| 3418 | my $backend = $self->{_backend}; | |||||
| 3419 | my $join_col = $self->{_join_col}; | |||||
| 3420 | my $join_type = $self->{_join_type}; | |||||
| 3421 | my $n = scalar @{ $self->{_dbs} }; | |||||
| 3422 | ||||||
| 3423 | # Partition query-time criteria only (no filter overlay). | |||||
| 3424 | # Filters are applied at cache-build time for spilled sources, and via the | |||||
| 3425 | # SQL WHERE clause for ATTACHed sources. Keeping them separate means the | |||||
| 3426 | # cached tables can serve any query without rebuilding. | |||||
| 3427 | my $per_db_query = $self->_partition_criteria($params); | |||||
| 3428 | ||||||
| 3429 | # Compute the full merged criteria (filter + query) for each source. | |||||
| 3430 | # Used for had_criteria (join-type semantics) and the WHERE clause for | |||||
| 3431 | # ATTACHed sources (which were not filtered at spill time). | |||||
| 3432 | my @per_db_full; | |||||
| 3433 | for my $i (0 .. $n - 1) { | |||||
| 3434 | my $base = $self->{_filters}{$i} // {}; | |||||
| 3435 | $per_db_full[$i] = %{$base} | |||||
| 3436 | ? _merge_criteria($base, $per_db_query->[$i]) | |||||
| 3437 | : $per_db_query->[$i]; | |||||
| 3438 | } | |||||
| 3439 | ||||||
| 3440 | # Determine which secondary sources had effective criteria (inner-join semantics). | |||||
| 3441 | # The broadcast join-column criterion must NOT count â it is a key-range selector | |||||
| 3442 | # on the merged view, not a predicate that restricts the secondary's contribution. | |||||
| 3443 | # Base filters always count (documented: a filtered db is always inner-join). | |||||
| 3444 | my @had_criteria; | |||||
| 3445 | for my $i (1 .. $n - 1) { | |||||
| 3446 | my $local_jc = $self->{_join_map}{$i} // $join_col; | |||||
| 3447 | my $has_filter = !!%{ $self->{_filters}{$i} // {} }; | |||||
| 3448 | my %q = %{ $per_db_query->[$i] }; | |||||
| 3449 | delete $q{$local_jc}; | |||||
| 3450 | $had_criteria[$i] = $has_filter || !!%q; | |||||
| 3451 | } | |||||
| 3452 | ||||||
| 3453 | # For 'auto' mode: check total row count without fetching rows. | |||||
| 3454 | # Count(*) is used for dbi_source() sources; count() for others. | |||||
| 3455 | # If any source supports neither, fall back to the array path. | |||||
| 3456 | # Skipped when create_table is set â dbi_source() has already committed to | |||||
| 3457 | # the SQLite path and we must materialise regardless of row count. | |||||
| 3458 | if ($backend eq 'auto' && !defined $create_table) { | |||||
| 3459 | my $total = 0; | |||||
| 3460 | my $can_count = 1; | |||||
| 3461 | for my $i (0 .. $n - 1) { | |||||
| 3462 | my $db = $self->{_dbs}[$i]; | |||||
| 3463 | # dbi_source() path: COUNT(*) against the entire source table | |||||
| 3464 | # (no WHERE) gives a conservative upper bound on the spilled size. | |||||
| 3465 | if (do { local $@; eval { $db->can('dbi_source') } }) { | |||||
| 3466 | my $src = do { | |||||
| 3467 | local $SIG{__WARN__} = sub {}; | |||||
| 3468 | local $@; | |||||
| 3469 | eval { $db->dbi_source() } | |||||
| 3470 | }; | |||||
| 3471 | if ($src && ref($src) eq 'HASH' && $src->{dbh} && $src->{table} | |||||
| 3472 | && eval { $src->{dbh}{Driver}{Name} } eq 'SQLite') { | |||||
| 3473 | my ($cnt) = $src->{dbh}->selectrow_array( | |||||
| 3474 | 'SELECT COUNT(*) FROM ' | |||||
| 3475 | . _sql_quote_identifier($src->{table}) | |||||
| 3476 | ); | |||||
| 3477 | $total += $cnt // 0; | |||||
| 3478 | next; | |||||
| 3479 | } | |||||
| 3480 | } | |||||
| 3481 | # count() path: only use it when the DA's own class directly defines | |||||
| 3482 | # count() (not inherited). Database::Abstraction's inherited count($entry) | |||||
| 3483 | # takes a key argument and emits uninitialized-value warnings when called | |||||
| 3484 | # with no args, so we must not invoke it for the threshold probe. | |||||
| 3485 | my $pkg = ref($db) // ''; | |||||
| 3486 | 25 25 25 | 7375 22 18189 | if ($pkg && do { no strict 'refs'; defined &{"${pkg}::count"} }) { | |||
| 3487 | my ($cnt, $failed); | |||||
| 3488 | do { local $@; $cnt = eval { $db->count() }; $failed = $@ }; | |||||
| 3489 | if ($failed) { | |||||
| 3490 | $can_count = 0; | |||||
| 3491 | last; | |||||
| 3492 | } | |||||
| 3493 | $total += $cnt // 0; | |||||
| 3494 | next; | |||||
| 3495 | } | |||||
| 3496 | $can_count = 0; | |||||
| 3497 | last; | |||||
| 3498 | } | |||||
| 3499 | if (!$can_count || $total <= $self->{_max_array_rows}) { | |||||
| 3500 | my $rows = $self->_joined_query_array($params, | |||||
| 3501 | sort_by => $sort_by, limit => $limit, offset => $offset); | |||||
| 3502 | return $count_only ? scalar @{$rows} : $rows; | |||||
| 3503 | } | |||||
| 3504 | } | |||||
| 3505 | ||||||
| 3506 | # Ensure the SQLite cache is valid; rebuild if stale or absent. | |||||
| 3507 | $self->_build_sqlite_cache() unless $self->_cache_fresh(); | |||||
| 3508 | ||||||
| 3509 | my $cache = $self->{_sqlite_cache}; | |||||
| 3510 | my $tmpdbh = $cache->{dbh}; | |||||
| 3511 | my @table_refs = @{ $cache->{table_refs} }; | |||||
| 3512 | my @source_cols = @{ $cache->{source_cols} }; | |||||
| 3513 | my @is_attached = @{ $cache->{is_attached} }; | |||||
| 3514 | ||||||
| 3515 | # Build the WHERE clause from per-call criteria. | |||||
| 3516 | # Spilled sources: query-only criteria (filter already applied to spilled data). | |||||
| 3517 | # ATTACHed sources: full criteria (filter + query), since source was not filtered. | |||||
| 3518 | # Secondary tables (i>0): the broadcast join-column criterion is omitted because | |||||
| 3519 | # it is already enforced by the ON clause; adding it to WHERE nullifies LEFT JOIN. | |||||
| 3520 | my (@where_parts, @bind_vals); | |||||
| 3521 | for my $i (0 .. $n - 1) { | |||||
| 3522 | my $crit = $is_attached[$i] ? $per_db_full[$i] : $per_db_query->[$i]; | |||||
| 3523 | next unless %{$crit}; | |||||
| 3524 | my $tref = $table_refs[$i]; | |||||
| 3525 | my $local_jc_i = $i > 0 ? ($self->{_join_map}{$i} // $join_col) : undef; | |||||
| 3526 | for my $col (sort keys %{$crit}) { | |||||
| 3527 | next if defined $local_jc_i && $col eq $local_jc_i; | |||||
| 3528 | my $val = $crit->{$col}; | |||||
| 3529 | if (ref($val) eq 'HASH') { | |||||
| 3530 | for my $op (sort keys %{$val}) { | |||||
| 3531 | if ($SAFE_LIST_OPS{$op}) { | |||||
| 3532 | # IN / NOT IN: value must be an arrayref; each element is bound. | |||||
| 3533 | my $arr = $val->{$op}; | |||||
| 3534 | unless (ref($arr) eq 'ARRAY') { | |||||
| 3535 | carp "Database::Join: operator '$op' requires an arrayref value; criterion skipped"; | |||||
| 3536 | next; | |||||
| 3537 | } | |||||
| 3538 | my @items = @{$arr}; | |||||
| 3539 | if (!@items) { | |||||
| 3540 | # IN () â always false: add a tautologically-false term. | |||||
| 3541 | # NOT IN () â always true: omit the term (all rows match). | |||||
| 3542 | push @where_parts, '1 = 0' if $op eq 'IN'; | |||||
| 3543 | next; | |||||
| 3544 | } | |||||
| 3545 | push @where_parts, | |||||
| 3546 | $tref . '.' . _sql_quote_identifier($col) | |||||
| 3547 | . " $op (" . join(', ', ('?') x scalar @items) . ')'; | |||||
| 3548 | push @bind_vals, @items; | |||||
| 3549 | next; | |||||
| 3550 | } | |||||
| 3551 | if ($SAFE_NOARG_OPS{$op}) { | |||||
| 3552 | # IS NULL / IS NOT NULL: no bind parameter at all. | |||||
| 3553 | push @where_parts, $tref . '.' . _sql_quote_identifier($col) . " $op"; | |||||
| 3554 | next; | |||||
| 3555 | } | |||||
| 3556 | unless ($SAFE_SQL_OPS{$op}) { | |||||
| 3557 | carp "Database::Join: operator '$op' is not supported on the SQLite backend; criterion skipped (use the array backend or a supported operator)"; | |||||
| 3558 | next; | |||||
| 3559 | } | |||||
| 3560 | push @where_parts, $tref . '.' . _sql_quote_identifier($col) . " $op ?"; | |||||
| 3561 | push @bind_vals, $val->{$op}; | |||||
| 3562 | } | |||||
| 3563 | } elsif (!defined $val) { | |||||
| 3564 | # A bare undef value means "WHERE col IS NULL". | |||||
| 3565 | # (col = NULL is always UNKNOWN in SQL and would match nothing.) | |||||
| 3566 | push @where_parts, $tref . '.' . _sql_quote_identifier($col) . ' IS NULL'; | |||||
| 3567 | } else { | |||||
| 3568 | push @where_parts, $tref . '.' . _sql_quote_identifier($col) . ' = ?'; | |||||
| 3569 | push @bind_vals, $val; | |||||
| 3570 | } | |||||
| 3571 | } | |||||
| 3572 | } | |||||
| 3573 | my $where_sql = @where_parts ? ' WHERE ' . join(' AND ', @where_parts) : ''; | |||||
| 3574 | ||||||
| 3575 | my $local_jc_0 = $self->{_join_map}{0} // $join_col; | |||||
| 3576 | ||||||
| 3577 | # Build JOIN clauses. Hoisted before the SELECT list so the count_only | |||||
| 3578 | # path can return early without building the (unused) column expressions. | |||||
| 3579 | # Join type mirrors the key-set semantics of _joined_query_array: | |||||
| 3580 | # had_criteria[i] OR inner => INNER JOIN | |||||
| 3581 | # outer (no criteria) => FULL OUTER JOIN | |||||
| 3582 | # left (no criteria) => LEFT JOIN | |||||
| 3583 | my $from = $table_refs[0]; | |||||
| 3584 | my $join_sql = ''; | |||||
| 3585 | for my $i (1 .. $n - 1) { | |||||
| 3586 | my $local_jc = $self->{_join_map}{$i} // $join_col; | |||||
| 3587 | my $join_kw = ($had_criteria[$i] || $join_type eq 'inner') ? 'JOIN' | |||||
| 3588 | : ($join_type eq 'outer') ? 'FULL OUTER JOIN' | |||||
| 3589 | : 'LEFT JOIN'; | |||||
| 3590 | $join_sql .= " $join_kw $table_refs[$i]" | |||||
| 3591 | . ' ON ' . $table_refs[0] . '.' . _sql_quote_identifier($local_jc_0) | |||||
| 3592 | . ' = ' . $table_refs[$i] . '.' . _sql_quote_identifier($local_jc); | |||||
| 3593 | } | |||||
| 3594 | ||||||
| 3595 | # COUNT(*) short-circuit: WHERE and JOIN are built; SELECT list and ORDER BY | |||||
| 3596 | # are not needed. selectrow_array returns a single integer without fetching rows. | |||||
| 3597 | if ($count_only) { | |||||
| 3598 | my ($cnt) = $tmpdbh->selectrow_array( | |||||
| 3599 | 'SELECT COUNT(*) FROM ' . $from . $join_sql . $where_sql, | |||||
| 3600 | undef, | |||||
| 3601 | @bind_vals, | |||||
| 3602 | ); | |||||
| 3603 | return $cnt // 0; | |||||
| 3604 | } | |||||
| 3605 | ||||||
| 3606 | # Build SELECT clause. | |||||
| 3607 | # Walk sources in order, applying collision_prefix renaming exactly as | |||||
| 3608 | # _build_col_index does: first occurrence of a column name wins; subsequent | |||||
| 3609 | # occurrences in a source that has a collision_prefix are published as | |||||
| 3610 | # "$prefix.$col"; removed columns are omitted. | |||||
| 3611 | my %pub_seen; | |||||
| 3612 | my @selects; | |||||
| 3613 | ||||||
| 3614 | # Join-column expression: for outer joins, use COALESCE across all sources | |||||
| 3615 | # so that B-only (primary-absent) rows carry their join key rather than NULL. | |||||
| 3616 | unless ($self->{_removed_cols}{$join_col}) { | |||||
| 3617 | my $jc_expr; | |||||
| 3618 | if ($join_type eq 'outer' && $n > 1) { | |||||
| 3619 | $jc_expr = 'COALESCE(' | |||||
| 3620 | . join(', ', map { | |||||
| 3621 | my $lc = $self->{_join_map}{$_} // $join_col; | |||||
| 3622 | $table_refs[$_] . '.' . _sql_quote_identifier($lc) | |||||
| 3623 | } 0 .. $n - 1) | |||||
| 3624 | . ') AS ' . _sql_quote_identifier($join_col); | |||||
| 3625 | } else { | |||||
| 3626 | $jc_expr = $table_refs[0] . '.' . _sql_quote_identifier($local_jc_0); | |||||
| 3627 | $jc_expr .= ' AS ' . _sql_quote_identifier($join_col) if $local_jc_0 ne $join_col; | |||||
| 3628 | } | |||||
| 3629 | push @selects, $jc_expr; | |||||
| 3630 | $pub_seen{$join_col} = 1; | |||||
| 3631 | } | |||||
| 3632 | ||||||
| 3633 | # Non-join columns from the primary table. | |||||
| 3634 | for my $col (@{$source_cols[0] // []}) { | |||||
| 3635 | next if $col eq $local_jc_0; # join column already handled above | |||||
| 3636 | next if $self->{_removed_cols}{$col}; | |||||
| 3637 | $pub_seen{$col} = 1; | |||||
| 3638 | push @selects, $table_refs[0] . '.' . _sql_quote_identifier($col); | |||||
| 3639 | } | |||||
| 3640 | ||||||
| 3641 | # Non-join columns from secondary tables, with collision_prefix renaming. | |||||
| 3642 | for my $i (1 .. $n - 1) { | |||||
| 3643 | my $local_jc = $self->{_join_map}{$i} // $join_col; | |||||
| 3644 | my $prefix = $self->{_collision_prefix}{$i}; | |||||
| 3645 | for my $col (@{$source_cols[$i] // []}) { | |||||
| 3646 | next if $col eq $local_jc; # join key already contributed above | |||||
| 3647 | my $pub = $col; | |||||
| 3648 | $pub = "$prefix.$col" | |||||
| 3649 | if defined $prefix && exists $pub_seen{$col} && $col ne $join_col; | |||||
| 3650 | next if $self->{_removed_cols}{$pub}; | |||||
| 3651 | $pub_seen{$pub} = 1; | |||||
| 3652 | my $expr = $table_refs[$i] . '.' . _sql_quote_identifier($col); | |||||
| 3653 | $expr .= ' AS ' . _sql_quote_identifier($pub) if $pub ne $col; | |||||
| 3654 | push @selects, $expr; | |||||
| 3655 | } | |||||
| 3656 | } | |||||
| 3657 | ||||||
| 3658 | # dbi_source() materialisation path: CREATE TABLE name AS SELECT ... | |||||
| 3659 | # Executed BEFORE ORDER BY / LIMIT / OFFSET because they are irrelevant here â | |||||
| 3660 | # the parent join will impose its own ordering and pagination per-call. | |||||
| 3661 | # _sql_quote_identifier guards against any injection via the table name. | |||||
| 3662 | # D~: $sort_by, $limit, $offset are dead stores on this path. They were | |||||
| 3663 | # parsed and validated above (shared with the normal query path) but the | |||||
| 3664 | # caller of create_table always passes undef for these, so _validate_pagination | |||||
| 3665 | # is a no-op and the variables are harmlessly abandoned at this return. | |||||
| 3666 | if (defined $create_table) { | |||||
| 3667 | my $mat_sql = 'SELECT ' . join(', ', @selects) | |||||
| 3668 | . ' FROM ' . $from . $join_sql . $where_sql; | |||||
| 3669 | $tmpdbh->do('DROP TABLE IF EXISTS ' . _sql_quote_identifier($create_table)); | |||||
| 3670 | $tmpdbh->do('CREATE TABLE ' . _sql_quote_identifier($create_table) | |||||
| 3671 | . ' AS ' . $mat_sql, | |||||
| 3672 | undef, @bind_vals); | |||||
| 3673 | return; | |||||
| 3674 | } | |||||
| 3675 | ||||||
| 3676 | # ORDER BY: default is join_column ascending. Caller may override via | |||||
| 3677 | # sort_by => 'col' or sort_by => ['col', 'DESC']. | |||||
| 3678 | # For the join column on an outer join, reference the COALESCE alias. | |||||
| 3679 | # For any other column, reference the published SELECT-list alias â | |||||
| 3680 | # SQLite resolves ORDER BY aliases from the SELECT clause. | |||||
| 3681 | my ($ob_col, $ob_dir) = ($join_col, 'ASC'); | |||||
| 3682 | if (defined $sort_by) { | |||||
| 3683 | my ($req_col, $req_dir) = ref($sort_by) eq 'ARRAY' ? @{$sort_by} : ($sort_by, 'ASC'); | |||||
| 3684 | $req_dir = uc($req_dir // 'ASC'); | |||||
| 3685 | unless ($req_dir eq 'ASC' || $req_dir eq 'DESC') { | |||||
| 3686 | carp "Database::Join: sort_by direction '$req_dir' is not supported; using ASC"; | |||||
| 3687 | $req_dir = 'ASC'; | |||||
| 3688 | } | |||||
| 3689 | if ($req_col ne $join_col && !exists $self->{_col_db}{$req_col}) { | |||||
| 3690 | carp "Database::Join: sort_by column '$req_col' is not in the merged view; result sorted by join_column"; | |||||
| 3691 | } else { | |||||
| 3692 | ($ob_col, $ob_dir) = ($req_col, $req_dir); | |||||
| 3693 | } | |||||
| 3694 | } | |||||
| 3695 | my $order_expr = ($ob_col eq $join_col) | |||||
| 3696 | ? (($join_type eq 'outer' && $n > 1) | |||||
| 3697 | ? _sql_quote_identifier($join_col) | |||||
| 3698 | : $table_refs[0] . '.' . _sql_quote_identifier($local_jc_0)) | |||||
| 3699 | : _sql_quote_identifier($ob_col); | |||||
| 3700 | my $sql = 'SELECT ' | |||||
| 3701 | . join(', ', @selects) | |||||
| 3702 | . ' FROM ' . $from . $join_sql | |||||
| 3703 | . $where_sql | |||||
| 3704 | . ' ORDER BY ' . $order_expr . ($ob_dir eq 'DESC' ? ' DESC' : ''); | |||||
| 3705 | ||||||
| 3706 | # LIMIT / OFFSET pagination â appended after ORDER BY as bind parameters | |||||
| 3707 | # (never interpolated) to prevent any SQL injection from caller values. | |||||
| 3708 | # OFFSET without LIMIT uses LIMIT -1 (SQLite extension: "all rows from offset"). | |||||
| 3709 | my @page_bind; | |||||
| 3710 | if (defined $limit) { | |||||
| 3711 | $sql .= ' LIMIT ?'; | |||||
| 3712 | push @page_bind, $limit; | |||||
| 3713 | if (defined $offset) { | |||||
| 3714 | $sql .= ' OFFSET ?'; | |||||
| 3715 | push @page_bind, $offset; | |||||
| 3716 | } | |||||
| 3717 | } elsif (defined $offset) { | |||||
| 3718 | $sql .= ' LIMIT -1 OFFSET ?'; | |||||
| 3719 | push @page_bind, $offset; | |||||
| 3720 | } | |||||
| 3721 | ||||||
| 3722 | # prepare_cached reuses the parsed statement when the same SQL is executed | |||||
| 3723 | # again (e.g. identical criteria pattern in a pagination or batch loop), | |||||
| 3724 | # avoiding repeated statement compilation overhead. fetchall_arrayref | |||||
| 3725 | # always exhausts the result set, so the handle is never left active. | |||||
| 3726 | my $sth = $tmpdbh->prepare_cached($sql); | |||||
| 3727 | $sth->execute(@bind_vals, @page_bind); | |||||
| 3728 | return $sth->fetchall_arrayref({}); | |||||
| 3729 | 25 25 25 | 56 18 180 | } | |||
| 3730 | ||||||
| 3731 | # _merge_criteria( \%base, \%extra ) -> \%merged | |||||
| 3732 | # Purpose: Merge two criteria hashrefs for the same database column set. | |||||
| 3733 | # Entry: %base is the permanent filter; %extra is the query-time criteria. | |||||
| 3734 | # Exit: Returns a new hashref with both applied. | |||||
| 3735 | # Merging rule: when both values for the same column are operator hashrefs | |||||
| 3736 | # (e.g. { '>' => 60 } and { '<' => 365 }), the operators are combined | |||||
| 3737 | # so both constraints apply simultaneously (AND semantics). | |||||
| 3738 | # Otherwise the extra (query-time) value overwrites the base value. | |||||
| 3739 | sub _merge_criteria :Protected { | |||||
| 3740 | my ($base, $extra) = @_; | |||||
| 3741 | my %merged = %{$base}; | |||||
| 3742 | for my $col (keys %{$extra}) { | |||||
| 3743 | if (exists $merged{$col} | |||||
| 3744 | && ref($merged{$col}) eq 'HASH' | |||||
| 3745 | && ref($extra->{$col}) eq 'HASH') { | |||||
| 3746 | $merged{$col} = { %{ $merged{$col} }, %{ $extra->{$col} } }; | |||||
| 3747 | } else { | |||||
| 3748 | $merged{$col} = $extra->{$col}; | |||||
| 3749 | } | |||||
| 3750 | } | |||||
| 3751 | return \%merged; | |||||
| 3752 | 25 25 25 | 4521 33 145 | } | |||
| 3753 | ||||||
| 3754 | # _copy_criteria( \%criteria ) -> \%copy | |||||
| 3755 | # Purpose: Return a two-level deep copy of a single criteria hashref so that | |||||
| 3756 | # post-construction mutation of the caller's hash cannot change the | |||||
| 3757 | # stored filter. Operator sub-hashrefs (e.g. { '>' => 80 }) are | |||||
| 3758 | # shallow-copied one additional level, matching the broadcast-copy | |||||
| 3759 | # idiom used in _partition_criteria for join-column criteria. | |||||
| 3760 | # Entry: $criteria is a hashref (may be undef). | |||||
| 3761 | # Exit: Returns a new hashref; never returns the input reference itself. | |||||
| 3762 | sub _copy_criteria { | |||||
| 3763 | 177 | 156 | my ($criteria) = @_; | |||
| 3764 | 177 100 | 216 173 | return {} unless $criteria && %{$criteria}; | |||
| 3765 | return { | |||||
| 3766 | map { | |||||
| 3767 | $_ => ref($criteria->{$_}) eq 'HASH' | |||||
| 3768 | 137 | 5255 | ? { %{ $criteria->{$_} } } # shallow-copy operator sub-hashref | |||
| 3769 | 95 | 580 | : $criteria->{$_} | |||
| 3770 | 94 171 | 68 41291 | } keys %{$criteria} | |||
| 3771 | }; | |||||
| 3772 | } | |||||
| 3773 | ||||||
| 3774 | # _copy_filters( \%filters ) -> \%copy | |||||
| 3775 | # Purpose: Deep-copy the filters hashref (db_index => criteria_hashref) so | |||||
| 3776 | # that post-construction mutation of the caller's hash cannot silently | |||||
| 3777 | # bypass the inner-join row-security guarantee. | |||||
| 3778 | # Entry: $filters may be undef. | |||||
| 3779 | # Exit: Returns a new hashref; never returns the input reference itself. | |||||
| 3780 | sub _copy_filters { | |||||
| 3781 | 1279 | 1399 | my ($filters) = @_; | |||
| 3782 | 1279 158 | 13139 200 | return {} unless $filters && %{$filters}; | |||
| 3783 | 157 151 150 | 283 11922 42743 | return { map { $_ => _copy_criteria($filters->{$_}) } keys %{$filters} }; | |||
| 3784 | } | |||||
| 3785 | ||||||
| 3786 | # _msg( $i18n, $key, @sprintf_args ) -> $string | |||||
| 3787 | # Purpose: Format a user-facing message, routing through the i18n object when | |||||
| 3788 | # one is provided. Falls back to the built-in %MESSAGES dictionary. | |||||
| 3789 | # Entry: $i18n may be undef. $key must be a key in %MESSAGES. | |||||
| 3790 | # Exit: Returns the formatted string. | |||||
| 3791 | sub _msg :Protected { | |||||
| 3792 | my ($i18n, $key, @args) = @_; | |||||
| 3793 | ||||||
| 3794 | if ($i18n && $i18n->can('translate')) { | |||||
| 3795 | return $i18n->translate($key, @args); | |||||
| 3796 | } | |||||
| 3797 | ||||||
| 3798 | my $fmt = $MESSAGES{$key} | |||||
| 3799 | // sprintf($MESSAGES{error_unknown_message}, $key); | |||||
| 3800 | ||||||
| 3801 | return @args ? sprintf($fmt, @args) : $fmt; | |||||
| 3802 | 25 25 25 | 4873 72 151 | } | |||
| 3803 | ||||||
| 3804 | 1; | |||||
| 3805 | ||||||