Add support for query parameters, JOINs, SELECT * and three-part table names
This commit is contained in:
+48
-18
@@ -22,33 +22,63 @@ class QueryExecutor:
|
||||
self._stats = stats
|
||||
|
||||
def execute(self, parsed: ParsedQuery) -> list[dict]:
|
||||
table = parsed.table
|
||||
columns = parsed.columns
|
||||
for table in parsed.tables:
|
||||
self._ensure_table(table, parsed)
|
||||
return self._run_in_memory(parsed)
|
||||
|
||||
def _ensure_table(self, table: str, parsed: ParsedQuery) -> None:
|
||||
if table in parsed.wildcard_tables:
|
||||
self._ensure_full(table)
|
||||
else:
|
||||
self._ensure_columns(table, parsed.columns_by_table[table])
|
||||
|
||||
def _ensure_full(self, table: str) -> None:
|
||||
"""Load every column of *table* (SELECT * / t.*), refetching unless already full."""
|
||||
if self._cache.is_table_cached(table) and self._cache.is_table_full(table):
|
||||
logger.debug(f"Cache hit (full): {table!r}")
|
||||
self._stats.record_hit()
|
||||
return
|
||||
|
||||
if self._cache.is_table_cached(table):
|
||||
logger.warning(f"Re-fetching {table!r} in full — SELECT * requested.")
|
||||
self._stats.record_refetch()
|
||||
else:
|
||||
self._stats.record_miss()
|
||||
|
||||
columns = self._cache.discover_columns(table, self._source_conn)
|
||||
self._cache.load_table(table, columns, self._source_conn, full=True)
|
||||
self._registry.update(table, columns)
|
||||
|
||||
def _ensure_columns(self, table: str, columns: list[str]) -> None:
|
||||
"""Load *table* with at least *columns*, refetching only when columns are missing."""
|
||||
missing = self._registry.needs_refetch(table, columns)
|
||||
table_cached = self._cache.is_table_cached(table)
|
||||
|
||||
if missing or not table_cached:
|
||||
if table_cached and missing:
|
||||
logger.warning(
|
||||
f"Re-fetching {table!r} — new columns requested: {missing}. "
|
||||
f"Expanding cache from {self._registry.get_columns(table)} + {missing}"
|
||||
)
|
||||
self._stats.record_refetch()
|
||||
else:
|
||||
self._stats.record_miss()
|
||||
all_columns = list(self._registry.get_columns(table)) + missing
|
||||
self._cache.load_table(table, all_columns, self._source_conn)
|
||||
self._registry.update(table, all_columns)
|
||||
else:
|
||||
if not missing and table_cached:
|
||||
logger.debug(f"Cache hit: {table!r} columns={columns}")
|
||||
self._stats.record_hit()
|
||||
return
|
||||
|
||||
return self._run_in_memory(parsed)
|
||||
if table_cached and missing:
|
||||
logger.warning(
|
||||
f"Re-fetching {table!r} — new columns requested: {missing}. "
|
||||
f"Expanding cache from {self._registry.get_columns(table)} + {missing}"
|
||||
)
|
||||
self._stats.record_refetch()
|
||||
else:
|
||||
self._stats.record_miss()
|
||||
|
||||
all_columns = list(self._registry.get_columns(table)) + missing
|
||||
self._cache.load_table(table, all_columns, self._source_conn)
|
||||
self._registry.update(table, all_columns)
|
||||
|
||||
def _run_in_memory(self, parsed: ParsedQuery) -> list[dict]:
|
||||
logger.debug(f"Executing in SQLite RAM: {parsed.original_sql!r}")
|
||||
cursor = self._cache.connection.execute(parsed.original_sql)
|
||||
logger.debug(f"Executing in SQLite RAM: {parsed.sqlite_sql!r} params={parsed.params!r}")
|
||||
conn = self._cache.connection
|
||||
if parsed.params is None:
|
||||
cursor = conn.execute(parsed.sqlite_sql)
|
||||
else:
|
||||
cursor = conn.execute(parsed.sqlite_sql, parsed.params)
|
||||
col_names = [desc[0] for desc in cursor.description]
|
||||
rows = cursor.fetchall()
|
||||
return [dict(zip(col_names, row)) for row in rows]
|
||||
|
||||
Reference in New Issue
Block a user