This is an automated email from the ASF dual-hosted git repository.
hansva pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git
The following commit(s) were added to refs/heads/main by this push:
new 1edaaa18b4 clear cache button code hardening, fixes #3312 (#8541)
1edaaa18b4 is described below
commit 1edaaa18b475a71a45024c067a98de2029a54431
Author: Hans Van Akelyen <[email protected]>
AuthorDate: Thu Sep 24 09:50:14 2026 +0200
clear cache button code hardening, fixes #3312 (#8541)
* clear cache button code hardening, fixes #3312
* code hardening
* add extra qualifiers
---
.../src/main/java/org/apache/hop/core/DbCache.java | 31 +++--
.../org/apache/hop/core/database/Database.java | 10 +-
.../hop/core/database/SqlQueryClassifier.java | 128 +++++++++++++++++++++
.../test/java/org/apache/hop/core/DbCacheTest.java | 109 ++++++++++++++++++
.../hop/core/database/SqlQueryClassifierTest.java | 97 ++++++++++++++++
.../java/org/apache/hop/pipeline/PipelineMeta.java | 26 +++++
.../PipelineMetaDbCacheInvalidationTest.java | 121 +++++++++++++++++++
.../context/metadata/MetadataContextHandler.java | 11 +-
8 files changed, 516 insertions(+), 17 deletions(-)
diff --git a/core/src/main/java/org/apache/hop/core/DbCache.java
b/core/src/main/java/org/apache/hop/core/DbCache.java
index 396a7ce87c..3153667ce9 100644
--- a/core/src/main/java/org/apache/hop/core/DbCache.java
+++ b/core/src/main/java/org/apache/hop/core/DbCache.java
@@ -19,6 +19,7 @@ package org.apache.hop.core;
import java.util.Enumeration;
import java.util.Hashtable;
+import java.util.concurrent.atomic.AtomicInteger;
import lombok.Getter;
import lombok.Setter;
import org.apache.hop.core.row.IRowMeta;
@@ -28,9 +29,16 @@ import org.apache.hop.core.row.IRowMeta;
* often launched to the databases to get information on tables etc.
*/
public class DbCache {
- private static DbCache dbCache;
+ private static final DbCache dbCache = new DbCache();
- private Hashtable<DbCacheEntry, IRowMeta> cache;
+ private volatile Hashtable<DbCacheEntry, IRowMeta> cache;
+
+ /**
+ * Bumped every time entries are removed from this cache. Anything which
derives row metadata from
+ * this cache and keeps the result around can compare the generation it last
saw with {@link
+ * #getGeneration()} to find out whether its own copy went stale.
+ */
+ private final AtomicInteger generation = new AtomicInteger();
@Getter @Setter private boolean active;
@@ -85,6 +93,20 @@ public class DbCache {
}
}
}
+ // Only bump the generation once the entries are gone. A reader which sees
the new generation
+ // and re-derives its row metadata must not be able to pick up the entries
we are removing.
+ generation.incrementAndGet();
+ }
+
+ /**
+ * The number of times this cache was cleared. Callers which cache anything
derived from the
+ * database cache can store this value alongside their own copy and drop
that copy as soon as the
+ * generation changes.
+ *
+ * @return the current generation of this cache
+ */
+ public int getGeneration() {
+ return generation.get();
}
private DbCache() {
@@ -93,14 +115,9 @@ public class DbCache {
}
/**
- * Create the database cache instance by loading it from disk
- *
* @return the database cache instance.
*/
public static DbCache getInstance() {
- if (dbCache == null) {
- dbCache = new DbCache();
- }
return dbCache;
}
diff --git a/core/src/main/java/org/apache/hop/core/database/Database.java
b/core/src/main/java/org/apache/hop/core/database/Database.java
index 21ff73bbc5..63d5fda9b2 100644
--- a/core/src/main/java/org/apache/hop/core/database/Database.java
+++ b/core/src/main/java/org/apache/hop/core/database/Database.java
@@ -1548,11 +1548,11 @@ public class Database implements IVariables,
ILoggingObject, AutoCloseable {
countAffectedRows(result, sql, count);
}
- // See if a cache needs to be cleared...
- String upperSql = sql.toUpperCase();
- if (upperSql.startsWith("ALTER TABLE")
- || upperSql.startsWith("DROP TABLE")
- || upperSql.startsWith("CREATE TABLE")) {
+ // A statement which changes the layout of a table or a view invalidates
anything we cached
+ // for this connection. The classifier also recognises the modifier
forms: CREATE OR REPLACE
+ // VIEW, DROP TABLE IF EXISTS, CREATE MATERIALIZED VIEW, ...
+ //
+ if (SqlQueryClassifier.isSchemaChange(sql)) {
DbCache.getInstance().clear(databaseMeta.getName());
}
} catch (SQLException ex) {
diff --git
a/core/src/main/java/org/apache/hop/core/database/SqlQueryClassifier.java
b/core/src/main/java/org/apache/hop/core/database/SqlQueryClassifier.java
index 56d9c59d25..8dc9ffcb31 100644
--- a/core/src/main/java/org/apache/hop/core/database/SqlQueryClassifier.java
+++ b/core/src/main/java/org/apache/hop/core/database/SqlQueryClassifier.java
@@ -35,6 +35,68 @@ public final class SqlQueryClassifier {
private static final Set<String> QUERY_STARTERS =
Set.of("SELECT", "SHOW", "EXPLAIN", "DESCRIBE", "DESC", "VALUES",
"TABLE");
+ /**
+ * Verbs which can change the layout of a table or a view. {@code TRUNCATE}
is deliberately
+ * absent: it removes rows, not columns.
+ */
+ private static final Set<String> SCHEMA_CHANGE_VERBS =
+ Set.of("CREATE", "ALTER", "DROP", "RENAME");
+
+ /**
+ * Keywords which are allowed between the verb and the object type, so that
{@code CREATE OR
+ * REPLACE VIEW}, {@code CREATE OR ALTER VIEW} and {@code DROP TABLE IF
EXISTS} are recognised as
+ * well as the plain forms. A modifier may carry a value ({@code
ALGORITHM=MERGE}, {@code
+ * DEFINER=`root`@`localhost`}). Treating a statement as a schema change
when it is not only costs
+ * a cache clear, so this list errs on the generous side.
+ */
+ private static final Set<String> SCHEMA_CHANGE_MODIFIERS =
+ Set.of(
+ "OR",
+ "REPLACE",
+ "ALTER",
+ "TEMP",
+ "TEMPORARY",
+ "GLOBAL",
+ "LOCAL",
+ "UNLOGGED",
+ "MATERIALIZED",
+ "EXTERNAL",
+ "VIRTUAL",
+ "FOREIGN",
+ "IF",
+ "NOT",
+ "EXISTS",
+ // Oracle: CREATE OR REPLACE FORCE EDITIONABLE VIEW, CREATE PUBLIC
SYNONYM
+ "FORCE",
+ "NOFORCE",
+ "EDITIONABLE",
+ "NONEDITIONABLE",
+ "EDITIONING",
+ "PUBLIC",
+ // MySQL / MariaDB, as SHOW CREATE VIEW and mysqldump write it:
+ // CREATE ALGORITHM=UNDEFINED DEFINER=`root`@`localhost` SQL
SECURITY DEFINER VIEW
+ "ALGORITHM",
+ "UNDEFINED",
+ "MERGE",
+ "TEMPTABLE",
+ "DEFINER",
+ "SQL",
+ "SECURITY",
+ "INVOKER",
+ // PostgreSQL: CREATE RECURSIVE VIEW
+ "RECURSIVE",
+ // Snowflake: CREATE TRANSIENT TABLE, CREATE SECURE VIEW, CREATE
DYNAMIC TABLE, ...
+ "TRANSIENT",
+ "SECURE",
+ "DYNAMIC",
+ "HYBRID",
+ // Teradata: CREATE MULTISET TABLE, CREATE VOLATILE TABLE
+ "MULTISET",
+ "VOLATILE");
+
+ /** Object types whose layout is reflected in cached row metadata. */
+ private static final Set<String> SCHEMA_CHANGE_OBJECTS = Set.of("TABLE",
"VIEW", "SYNONYM");
+
/**
* First keywords of a complete statement. Leftover clauses after a
semicolon ({@code WHERE},
* {@code AND}, {@code ORDER}, …) are not in this set.
@@ -153,6 +215,72 @@ public final class SqlQueryClassifier {
return first;
}
+ /**
+ * Whether a statement changes the layout of a table or a view, which makes
any row metadata
+ * cached for that connection unreliable.
+ *
+ * @param sql one statement, comments allowed
+ * @return {@code true} for statements such as {@code ALTER TABLE ...},
{@code DROP TABLE IF
+ * EXISTS ...} or {@code CREATE OR REPLACE VIEW ...}
+ */
+ public static boolean isSchemaChange(String sql) {
+ if (Utils.isEmpty(sql)) {
+ return false;
+ }
+ int i = skipTrivia(sql, 0);
+ String keyword = keywordAt(sql, i);
+ if (keyword == null || !SCHEMA_CHANGE_VERBS.contains(keyword)) {
+ return false;
+ }
+ i = skipKeyword(sql, i);
+ while ((keyword = keywordAt(sql, i)) != null) {
+ if (SCHEMA_CHANGE_OBJECTS.contains(keyword)) {
+ return true;
+ }
+ if (!SCHEMA_CHANGE_MODIFIERS.contains(keyword)) {
+ return false;
+ }
+ i = skipOptionValue(sql, skipKeyword(sql, i));
+ }
+ return false;
+ }
+
+ /**
+ * Skips the {@code = value} of a modifier such as {@code ALGORITHM=MERGE}
or {@code
+ * DEFINER=`root`@`localhost`}, if there is one. The value is a keyword, a
number or a quoted
+ * identifier, optionally followed by {@code @host} (a MySQL account) or
{@code ()} ({@code
+ * CURRENT_USER()}).
+ */
+ private static int skipOptionValue(String sql, int i) {
+ int j = skipTrivia(sql, i);
+ if (j >= sql.length() || sql.charAt(j) != '=') {
+ return i;
+ }
+ j = skipValue(sql, skipTrivia(sql, j + 1));
+ int k = skipTrivia(sql, j);
+ if (k < sql.length() && sql.charAt(k) == '@') {
+ j = skipValue(sql, skipTrivia(sql, k + 1));
+ k = skipTrivia(sql, j);
+ }
+ if (k + 1 < sql.length() && sql.charAt(k) == '(' && sql.charAt(k + 1) ==
')') {
+ j = k + 2;
+ }
+ return j;
+ }
+
+ private static int skipValue(String sql, int i) {
+ if (i >= sql.length()) {
+ return i;
+ }
+ if (isQuote(sql.charAt(i))) {
+ return skipQuoted(sql, i);
+ }
+ while (i < sql.length() && (isIdentPart(sql.charAt(i)) || sql.charAt(i) ==
'.')) {
+ i++;
+ }
+ return i;
+ }
+
/**
* @param sql one statement, comments allowed
* @return {@code true} when {@code sql} starts with a SQL verb (query, DML
or DDL), {@code false}
diff --git a/core/src/test/java/org/apache/hop/core/DbCacheTest.java
b/core/src/test/java/org/apache/hop/core/DbCacheTest.java
new file mode 100644
index 0000000000..f19a0aea24
--- /dev/null
+++ b/core/src/test/java/org/apache/hop/core/DbCacheTest.java
@@ -0,0 +1,109 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.hop.core;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaString;
+import org.junit.jupiter.api.Test;
+
+class DbCacheTest {
+
+ @Test
+ void entriesAreRemovedPerConnection() {
+ DbCache cache = DbCache.getInstance();
+ cache.clear(null);
+
+ cache.put(new DbCacheEntry("one", "select * from t"), rowMeta());
+ cache.put(new DbCacheEntry("two", "select * from t"), rowMeta());
+
+ cache.clear("ONE");
+
+ assertNull(cache.get(new DbCacheEntry("one", "select * from t")));
+ assertNotNull(cache.get(new DbCacheEntry("two", "select * from t")));
+ }
+
+ @Test
+ void clearingBumpsTheGeneration() {
+ DbCache cache = DbCache.getInstance();
+
+ int before = cache.getGeneration();
+ cache.put(new DbCacheEntry("one", "select * from t"), rowMeta());
+ assertEquals(before, cache.getGeneration(), "storing an entry is not a
change of generation");
+
+ cache.clear(null);
+ int afterClearAll = cache.getGeneration();
+ assertNotEquals(before, afterClearAll);
+
+ cache.clear("one");
+ assertNotEquals(afterClearAll, cache.getGeneration(), "clearing one
connection counts as well");
+ }
+
+ @Test
+ void concurrentClearsAreAllCounted() throws Exception {
+ DbCache cache = DbCache.getInstance();
+ int threads = 8;
+ int clearsPerThread = 1000;
+ int before = cache.getGeneration();
+
+ ExecutorService executor = Executors.newFixedThreadPool(threads);
+ CountDownLatch start = new CountDownLatch(1);
+ List<Future<Void>> futures = new ArrayList<>();
+ try {
+ for (int t = 0; t < threads; t++) {
+ futures.add(
+ executor.submit(
+ () -> {
+ start.await();
+ assertSame(cache, DbCache.getInstance());
+ for (int i = 0; i < clearsPerThread; i++) {
+ cache.clear("one");
+ }
+ return null;
+ }));
+ }
+ start.countDown();
+ } finally {
+ executor.shutdown();
+ }
+ assertTrue(executor.awaitTermination(30, TimeUnit.SECONDS));
+ for (Future<Void> future : futures) {
+ future.get(); // rethrows any assertion failure from the worker threads
+ }
+
+ assertEquals(before + threads * clearsPerThread, cache.getGeneration());
+ }
+
+ private RowMeta rowMeta() {
+ RowMeta rowMeta = new RowMeta();
+ rowMeta.addValueMeta(new ValueMetaString("a"));
+ return rowMeta;
+ }
+}
diff --git
a/core/src/test/java/org/apache/hop/core/database/SqlQueryClassifierTest.java
b/core/src/test/java/org/apache/hop/core/database/SqlQueryClassifierTest.java
index 88f7655c7d..900882ca75 100644
---
a/core/src/test/java/org/apache/hop/core/database/SqlQueryClassifierTest.java
+++
b/core/src/test/java/org/apache/hop/core/database/SqlQueryClassifierTest.java
@@ -164,4 +164,101 @@ class SqlQueryClassifierTest {
assertNull(SqlQueryClassifier.statementVerb(" "));
assertNull(SqlQueryClassifier.statementVerb("(SELECT 1)"));
}
+
+ @Test
+ void schemaChangesOnTablesAndViewsAreDetected() {
+ assertTrue(SqlQueryClassifier.isSchemaChange("ALTER TABLE t ADD COLUMN c
INT"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("create table t (id int)"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("DROP TABLE t"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("RENAME TABLE a TO b"));
+ }
+
+ @Test
+ void schemaChangesAreDetectedPastModifiersAndTrivia() {
+ // The old check was a startsWith() on the upper-cased statement, so all
of these were missed.
+ assertTrue(SqlQueryClassifier.isSchemaChange("\n\t ALTER TABLE t ADD
COLUMN c INT"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("-- fix the layout\nALTER
TABLE t DROP COLUMN c"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("/* ticket 42 */ DROP TABLE
IF EXISTS t"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE TABLE IF NOT EXISTS t
(id int)"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE OR REPLACE VIEW v AS
SELECT * FROM t"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE OR ALTER VIEW v AS
SELECT * FROM t"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("ALTER VIEW v AS SELECT *
FROM t"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("DROP VIEW v"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE MATERIALIZED VIEW v
AS SELECT 1"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE GLOBAL TEMPORARY
TABLE t (id int)"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("DROP SYNONYM s"));
+ }
+
+ @Test
+ void dialectSpecificModifiersAreSchemaChanges() {
+ // Oracle, DBMS_METADATA emits the FORCE EDITIONABLE form verbatim
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE OR REPLACE FORCE VIEW
v AS SELECT 1"));
+ assertTrue(
+ SqlQueryClassifier.isSchemaChange(
+ "CREATE OR REPLACE FORCE EDITIONABLE VIEW v AS SELECT 1 FROM
dual"));
+ assertTrue(
+ SqlQueryClassifier.isSchemaChange("CREATE OR REPLACE EDITIONABLE VIEW
v AS SELECT 1"));
+ assertTrue(
+ SqlQueryClassifier.isSchemaChange("CREATE OR REPLACE NONEDITIONABLE
VIEW v AS SELECT 1"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE PUBLIC SYNONYM s FOR
t"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("DROP PUBLIC SYNONYM s"));
+ // PostgreSQL
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE RECURSIVE VIEW v (a)
AS SELECT 1"));
+ // Snowflake
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE OR REPLACE TRANSIENT
TABLE t (id int)"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE SECURE VIEW v AS
SELECT 1"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE OR REPLACE DYNAMIC
TABLE t AS SELECT 1"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE HYBRID TABLE t (id
int PRIMARY KEY)"));
+ // Teradata
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE MULTISET TABLE t (id
int)"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("CREATE VOLATILE TABLE t (id
int)"));
+ // Oracle forms outside the DBMS_METADATA default
+ assertTrue(
+ SqlQueryClassifier.isSchemaChange(
+ "CREATE OR REPLACE NOFORCE EDITIONABLE VIEW v AS SELECT 1 FROM
dual"));
+ assertTrue(
+ SqlQueryClassifier.isSchemaChange(
+ "CREATE OR REPLACE EDITIONING VIEW v AS SELECT a FROM t"));
+ }
+
+ @Test
+ void mysqlViewDefinitionsAreSchemaChanges() {
+ // The text SHOW CREATE VIEW returns and mysqldump writes, verbatim
+ assertTrue(
+ SqlQueryClassifier.isSchemaChange(
+ "CREATE ALGORITHM=UNDEFINED DEFINER=`root`@`localhost` SQL
SECURITY DEFINER VIEW `v`"
+ + " AS select `t`.`id` AS `id` from `t`"));
+ assertTrue(
+ SqlQueryClassifier.isSchemaChange(
+ "CREATE OR REPLACE ALGORITHM = MERGE DEFINER = 'app'@'%' SQL
SECURITY INVOKER VIEW v"
+ + " AS SELECT 1"));
+ assertTrue(
+ SqlQueryClassifier.isSchemaChange(
+ "CREATE DEFINER=CURRENT_USER() SQL SECURITY INVOKER VIEW v AS
SELECT 1"));
+ assertTrue(SqlQueryClassifier.isSchemaChange("ALTER ALGORITHM=TEMPTABLE
VIEW v AS SELECT 1"));
+ // A modifier with a value still has to be followed by a table or a view
+ assertFalse(
+ SqlQueryClassifier.isSchemaChange(
+ "CREATE DEFINER=`root`@`localhost` TRIGGER tr BEFORE INSERT ON t
FOR EACH ROW SET"
+ + " @x = 1"));
+ assertFalse(
+ SqlQueryClassifier.isSchemaChange(
+ "CREATE DEFINER=`root`@`localhost` PROCEDURE p() SELECT 1"));
+ }
+
+ @Test
+ void otherStatementsAreNotSchemaChanges() {
+ assertFalse(SqlQueryClassifier.isSchemaChange(null));
+ assertFalse(SqlQueryClassifier.isSchemaChange(" "));
+ assertFalse(SqlQueryClassifier.isSchemaChange("SELECT * FROM t"));
+ assertFalse(SqlQueryClassifier.isSchemaChange("INSERT INTO t VALUES (1)"));
+ assertFalse(SqlQueryClassifier.isSchemaChange("DELETE FROM t"));
+ // TRUNCATE removes rows, it does not change the layout
+ assertFalse(SqlQueryClassifier.isSchemaChange("TRUNCATE t"));
+ assertFalse(SqlQueryClassifier.isSchemaChange("TRUNCATE TABLE t"));
+ assertFalse(SqlQueryClassifier.isSchemaChange("CREATE INDEX i ON t (a)"));
+ assertFalse(SqlQueryClassifier.isSchemaChange("CREATE SEQUENCE s"));
+ assertFalse(SqlQueryClassifier.isSchemaChange("ALTER SESSION SET x = 1"));
+ assertFalse(SqlQueryClassifier.isSchemaChange("DROP INDEX i"));
+ }
}
diff --git a/engine/src/main/java/org/apache/hop/pipeline/PipelineMeta.java
b/engine/src/main/java/org/apache/hop/pipeline/PipelineMeta.java
index 78c2f90cfa..1d57779ed4 100644
--- a/engine/src/main/java/org/apache/hop/pipeline/PipelineMeta.java
+++ b/engine/src/main/java/org/apache/hop/pipeline/PipelineMeta.java
@@ -38,6 +38,7 @@ import org.apache.commons.vfs2.FileSystemException;
import org.apache.hop.base.AbstractMeta;
import org.apache.hop.core.CheckResult;
import org.apache.hop.core.Const;
+import org.apache.hop.core.DbCache;
import org.apache.hop.core.HopVersionProvider;
import org.apache.hop.core.ICheckResult;
import org.apache.hop.core.IProgressMonitor;
@@ -151,6 +152,14 @@ public class PipelineMeta extends AbstractMeta
/** The transforms fields cache. */
protected Map<String, IRowMeta> transformFieldsCache;
+ /**
+ * The {@link DbCache} generation the transform fields cache was filled
against. Transforms like
+ * Table Input derive their output fields from the database cache, so
clearing that cache has to
+ * invalidate the fields we cached here as well. Without this, clearing the
database cache only
+ * takes effect after the pipeline is reloaded.
+ */
+ protected int transformFieldsCacheDbGeneration;
+
/** The loop cache. */
protected Map<String, Boolean> loopCache;
@@ -286,6 +295,7 @@ public class PipelineMeta extends AbstractMeta
maxUndo = Const.MAX_UNDO;
undoPosition = -1;
transformFieldsCache = new HashMap<>();
+ transformFieldsCacheDbGeneration = DbCache.getInstance().getGeneration();
loopCache = new HashMap<>();
previousTransformCache = new HashMap<>();
super.clear();
@@ -1210,6 +1220,8 @@ public class PipelineMeta extends AbstractMeta
return row;
}
+ discardTransformFieldsCacheIfDatabaseCacheCleared();
+
String fromToCacheEntry = calculateFieldsCacheEntryKey(transformMeta,
targetTransform);
IRowMeta rowMeta = transformFieldsCache.get(fromToCacheEntry);
if (rowMeta != null) {
@@ -3449,6 +3461,20 @@ public class PipelineMeta extends AbstractMeta
/** Clears the transform fields cache. */
private void clearTransformFieldsCache() {
transformFieldsCache.clear();
+ transformFieldsCacheDbGeneration = DbCache.getInstance().getGeneration();
+ }
+
+ /**
+ * Drop the cached transform fields when the database cache was cleared
since we filled them.
+ * Transforms which read their layout from the database (Table Input, Table
Output, Database
+ * Lookup, ...) go through {@link DbCache}, so a stale entry here survives
clearing that cache and
+ * keeps showing the old columns until the pipeline is reloaded.
+ */
+ private void discardTransformFieldsCacheIfDatabaseCacheCleared() {
+ int currentGeneration = DbCache.getInstance().getGeneration();
+ if (currentGeneration != transformFieldsCacheDbGeneration) {
+ clearTransformFieldsCache();
+ }
}
/** Clears the loop cache. */
diff --git
a/engine/src/test/java/org/apache/hop/pipeline/PipelineMetaDbCacheInvalidationTest.java
b/engine/src/test/java/org/apache/hop/pipeline/PipelineMetaDbCacheInvalidationTest.java
new file mode 100644
index 0000000000..c286339024
--- /dev/null
+++
b/engine/src/test/java/org/apache/hop/pipeline/PipelineMetaDbCacheInvalidationTest.java
@@ -0,0 +1,121 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.hop.pipeline;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.when;
+
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.hop.core.DbCache;
+import org.apache.hop.core.IProgressMonitor;
+import org.apache.hop.core.exception.HopTransformException;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.value.ValueMetaString;
+import org.apache.hop.core.variables.IVariables;
+import org.apache.hop.core.variables.Variables;
+import org.apache.hop.pipeline.transform.ITransformMeta;
+import org.apache.hop.pipeline.transform.TransformIOMeta;
+import org.apache.hop.pipeline.transform.TransformMeta;
+import org.apache.hop.pipeline.transforms.dummy.DummyMeta;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.stubbing.Answer;
+
+/**
+ * Clearing the database cache has to invalidate the row metadata a pipeline
cached for transforms
+ * which read their layout from a database. See <a
+ * href="https://github.com/apache/hop/issues/3312">issue 3312</a>: without
this, the "clear cache"
+ * buttons only took effect after the pipeline was reloaded.
+ */
+class PipelineMetaDbCacheInvalidationTest {
+
+ /** Stands in for the columns the database currently reports for the
transform's query. */
+ private final List<String> databaseColumns = new ArrayList<>(List.of("a",
"b"));
+
+ private IVariables variables;
+ private PipelineMeta pipelineMeta;
+ private TransformMeta databaseTransform;
+
+ @BeforeEach
+ void setUp() throws Exception {
+ DbCache.getInstance().clear(null);
+ variables = new Variables();
+
+ databaseTransform = new TransformMeta("database input",
databaseBackedTransformMeta());
+ TransformMeta after = new TransformMeta("after", new DummyMeta());
+
+ pipelineMeta = new PipelineMeta();
+ pipelineMeta.addTransform(databaseTransform);
+ pipelineMeta.addTransform(after);
+ pipelineMeta.addPipelineHop(new PipelineHopMeta(databaseTransform, after));
+ }
+
+ @Test
+ void fieldsAreCachedUntilTheDatabaseCacheIsCleared() throws Exception {
+ assertArrayEquals(new String[] {"a", "b"}, transformFields());
+
+ // The table gains a column outside of Hop, so what we cached is now stale.
+ databaseColumns.add("c");
+ assertArrayEquals(new String[] {"a", "b"}, transformFields());
+
+ // Which is exactly what the "clear cache" buttons in the GUI are for.
+ DbCache.getInstance().clear(null);
+ assertArrayEquals(new String[] {"a", "b", "c"}, transformFields());
+ }
+
+ @Test
+ void clearingASingleConnectionAlsoInvalidatesCachedFields() throws Exception
{
+ assertArrayEquals(new String[] {"a", "b"}, transformFields());
+
+ databaseColumns.add("c");
+ DbCache.getInstance().clear("some connection");
+
+ assertArrayEquals(new String[] {"a", "b", "c"}, transformFields());
+ }
+
+ private String[] transformFields() throws HopTransformException {
+ return pipelineMeta
+ .getTransformFields(variables, databaseTransform, null,
mock(IProgressMonitor.class))
+ .getFieldNames();
+ }
+
+ /** A transform which reports whatever columns the database currently has,
like Table Input. */
+ private ITransformMeta databaseBackedTransformMeta() throws
HopTransformException {
+ TransformIOMeta transformIOMeta = mock(TransformIOMeta.class);
+ when(transformIOMeta.getInfoTransformNames()).thenReturn(new String[0]);
+
+ ITransformMeta transformMeta = spy(new DummyMeta());
+ when(transformMeta.getTransformIOMeta()).thenReturn(transformIOMeta);
+ doAnswer(
+ (Answer<Void>)
+ invocation -> {
+ IRowMeta rowMeta = (IRowMeta) invocation.getArguments()[0];
+ databaseColumns.forEach(
+ column -> rowMeta.addValueMeta(new
ValueMetaString(column)));
+ return null;
+ })
+ .when(transformMeta)
+ .getFields(any(), any(), any(), any(), any(), any());
+
+ return transformMeta;
+ }
+}
diff --git
a/ui/src/main/java/org/apache/hop/ui/hopgui/context/metadata/MetadataContextHandler.java
b/ui/src/main/java/org/apache/hop/ui/hopgui/context/metadata/MetadataContextHandler.java
index d1342a4dfa..58f5e029ae 100644
---
a/ui/src/main/java/org/apache/hop/ui/hopgui/context/metadata/MetadataContextHandler.java
+++
b/ui/src/main/java/org/apache/hop/ui/hopgui/context/metadata/MetadataContextHandler.java
@@ -130,11 +130,12 @@ public class MetadataContextHandler implements
IGuiContextHandler {
BaseMessages.getString(
PKG,
"HopGui.Context.Database.Menu.ClearDatabaseCache.Tooltip"),
null,
- (shiftClicked, controlClicked, parameters) ->
- DbCache.getInstance().clear((String) parameters[0]));
- newAction.setClassLoader(metadataObjectClass.getClassLoader());
- newAction.setCategory(CONST_METADATA);
- newAction.setCategoryOrder("3");
+ // No connection is selected in this context, and the action is
executed without
+ // parameters: clear the cache of every connection.
+ (shiftClicked, controlClicked, parameters) ->
DbCache.getInstance().clear(null));
+
databaseClearCacheAction.setClassLoader(metadataObjectClass.getClassLoader());
+ databaseClearCacheAction.setCategory(CONST_METADATA);
+ databaseClearCacheAction.setCategoryOrder("3");
actions.add(databaseClearCacheAction);
}