diff --git a/src/db-copy.cpp b/src/db-copy.cpp index 1a7dedec5..08b5a00fd 100644 --- a/src/db-copy.cpp +++ b/src/db-copy.cpp @@ -20,13 +20,18 @@ void db_deleter_by_id_t::delete_rows(std::string const &table, std::string const &column, + std::string const &history, pg_conn_t const &db_connection) { fmt::memory_buffer sql; // Each deletable contributes an OSM ID and a comma. The highest node ID // currently has 10 digits, so 15 characters should do for a couple of years. // Add 50 characters for the SQL statement itself. - sql.reserve(m_deletables.size() * 15 + 50); + sql.reserve(m_deletables.size() * 15 + 50 + history.size() * 2); + + if (!history.empty()) { + fmt::format_to(std::back_inserter(sql), FMT_STRING("WITH d AS (")); + } fmt::format_to(std::back_inserter(sql), FMT_STRING("DELETE FROM {} WHERE {} IN ("), table, column); @@ -36,12 +41,20 @@ void db_deleter_by_id_t::delete_rows(std::string const &table, } sql[sql.size() - 1] = ')'; + if (!history.empty()) { + fmt::format_to( + std::back_inserter(sql), + FMT_STRING(" RETURNING *) INSERT INTO {} SELECT d.*, now() FROM d"), + history); + } + sql.push_back('\0'); db_connection.exec(sql.data()); } void db_deleter_by_type_and_id_t::delete_rows(std::string const &table, std::string const &column, + std::string const &history, pg_conn_t const &db_connection) { assert(!m_deletables.empty()); @@ -50,7 +63,11 @@ void db_deleter_by_type_and_id_t::delete_rows(std::string const &table, // Need a VALUES line for each deletable: type (3 bytes), id (15 bytes), // braces etc. (4 bytes). And additional space for the remainder of the // SQL command. - sql.reserve(m_deletables.size() * 22 + 200); + sql.reserve(m_deletables.size() * 22 + 200 + history.size() * 2); + + if (!history.empty()) { + fmt::format_to(std::back_inserter(sql), FMT_STRING("WITH d AS (")); + } if (m_has_type) { fmt::format_to(std::back_inserter(sql), @@ -71,6 +88,13 @@ void db_deleter_by_type_and_id_t::delete_rows(std::string const &table, ") AS t (osm_type, osm_id) WHERE" " p.{} = t.osm_type::char(1) AND p.{} = t.osm_id", type, column.c_str() + pos + 1); + + if (!history.empty()) { + fmt::format_to(std::back_inserter(sql), + FMT_STRING(" RETURNING p.*) INSERT INTO {}" + " SELECT d.*, now() FROM d"), + history); + } } else { fmt::format_to(std::back_inserter(sql), FMT_STRING("DELETE FROM {} WHERE {} IN ("), table, @@ -80,6 +104,13 @@ void db_deleter_by_type_and_id_t::delete_rows(std::string const &table, format_to(std::back_inserter(sql), FMT_STRING("{},"), item.osm_id); } sql[sql.size() - 1] = ')'; + + if (!history.empty()) { + fmt::format_to(std::back_inserter(sql), + FMT_STRING(" RETURNING *) INSERT INTO {}" + " SELECT d.*, now() FROM d"), + history); + } } sql.push_back('\0'); diff --git a/src/db-copy.hpp b/src/db-copy.hpp index 3a9d3f05e..0f375d7cf 100644 --- a/src/db-copy.hpp +++ b/src/db-copy.hpp @@ -46,8 +46,10 @@ class db_target_descr_t std::string const &name() const noexcept { return m_name; } std::string const &id() const noexcept { return m_id; } std::string const &rows() const noexcept { return m_rows; } + std::string const &history() const noexcept { return m_history; } void set_rows(std::string rows) { m_rows = std::move(rows); } + void set_history(std::string history) { m_history = std::move(history); } /** * Check if the buffer would use exactly the same copy operation. @@ -68,6 +70,8 @@ class db_target_descr_t std::string m_id; /// Comma-separated list of rows for copy operation (when empty: all rows) std::string m_rows; + /// Qualified name of the history table (when empty: history disabled) + std::string m_history; }; /** @@ -87,6 +91,7 @@ class db_deleter_by_id_t void add(osmid_t osm_id) { m_deletables.push_back(osm_id); } void delete_rows(std::string const &table, std::string const &column, + std::string const &history, pg_conn_t const &db_connection); bool is_full() const noexcept { return m_deletables.size() > MAX_ENTRIES; } @@ -127,6 +132,7 @@ class db_deleter_by_type_and_id_t } void delete_rows(std::string const &table, std::string const &column, + std::string const &history, pg_conn_t const &db_connection); bool is_full() const noexcept { return m_deletables.size() > MAX_ENTRIES; } @@ -192,7 +198,7 @@ class db_cmd_copy_delete_t : public db_cmd_copy_t if (m_deleter.has_data()) { m_deleter.delete_rows( qualified_name(target->schema(), target->name()), target->id(), - db_connection); + target->history(), db_connection); } } diff --git a/src/flex-lua-table.cpp b/src/flex-lua-table.cpp index b88c58261..d99ca4561 100644 --- a/src/flex-lua-table.cpp +++ b/src/flex-lua-table.cpp @@ -446,6 +446,37 @@ TRAMPOLINE_WRAPPED_OBJECT(table, schema) } // anonymous namespace +void setup_flex_table_history(lua_State *lua_state, flex_table_t *table) +{ + assert(lua_state); + assert(table); + + lua_getfield(lua_state, -1, "history"); + if (lua_isboolean(lua_state, -1)) { + table->set_has_history(lua_toboolean(lua_state, -1)); + } else if (!lua_isnil(lua_state, -1)) { + throw fmt_error("The 'history' field in table '{}' must be a boolean.", + table->name()); + } + lua_pop(lua_state, 1); + + if (!table->has_history()) { + return; + } + + for (char const *name : {"valid_from", "valid_to"}) { + if (util::find_by_name(table->columns(), name)) { + throw fmt_error( + "Table '{}' has history enabled, so column '{}' is added by" + " osm2pgsql and must not be defined in the Lua config.", + table->name(), name); + } + } + + table->add_column("valid_from", "timestamptz", "timestamptz DEFAULT now()") + .set_create_only(); +} + int setup_flex_table(lua_State *lua_state, std::vector *tables, std::vector *expire_outputs, std::string const &default_schema, bool updatable, @@ -461,6 +492,7 @@ int setup_flex_table(lua_State *lua_state, std::vector *tables, setup_flex_table_columns(lua_state, &new_table, expire_outputs, append_mode); setup_flex_table_indexes(lua_state, &new_table, updatable); + setup_flex_table_history(lua_state, &new_table); void *ptr = lua_newuserdata(lua_state, sizeof(std::size_t)); auto *num = new (ptr) std::size_t{}; diff --git a/src/flex-table.cpp b/src/flex-table.cpp index 263a9cedc..a4871eb40 100644 --- a/src/flex-table.cpp +++ b/src/flex-table.cpp @@ -60,6 +60,11 @@ std::string flex_table_t::full_name() const return qualified_name(schema(), name()); } +std::string flex_table_t::full_history_name() const +{ + return qualified_name(schema(), history_name()); +} + std::string flex_table_t::full_tmp_name() const { return qualified_name(schema(), name() + "_tmp"); @@ -207,6 +212,63 @@ flex_table_t::build_sql_create_table(table_type ttype, return sql; } +std::string flex_table_t::build_sql_create_history_table() const +{ + assert(!m_columns.empty()); + assert(m_has_history); + + std::string sql = + fmt::format("CREATE TABLE IF NOT EXISTS {} (", full_history_name()); + + util::string_joiner_t joiner{','}; + for (auto const &column : m_columns) { + joiner.add(column.sql_create()); + } + joiner.add(R"("valid_to" timestamptz)"); + + sql += joiner(); + sql += ')'; + sql += tablespace_clause(m_data_tablespace); + + return sql; +} + +std::string flex_table_t::build_sql_dedup_history() const +{ + assert(m_has_history); + + std::string match; + util::string_joiner_t hist{','}; + util::string_joiner_t live{','}; + + for (auto const &column : m_columns) { + if (column.type() == table_column_type::id_type || + column.type() == table_column_type::id_num) { + if (!match.empty()) { + match += " AND "; + } + match += fmt::format(R"(h."{0}" = b."{0}")", column.name()); + continue; + } + if (column.create_only()) { + continue; + } + hist.add(fmt::format(R"(h."{}")", column.name())); + live.add(fmt::format(R"(b."{}")", column.name())); + } + + if (match.empty() || hist.empty()) { + return {}; + } + + return fmt::format( + R"(DELETE FROM {} h USING {} b WHERE {})" + R"( AND h."valid_to" = (SELECT max("valid_to") FROM {}))" + R"( AND ROW({}) IS NOT DISTINCT FROM ROW({}))", + full_history_name(), full_name(), match, full_history_name(), hist(), + live()); +} + std::string flex_table_t::build_sql_column_list() const { assert(!m_columns.empty()); @@ -313,6 +375,10 @@ void table_connection_t::start(pg_conn_t const &db_connection, table().full_name())); enable_check_trigger(db_connection, table()); + + if (table().has_history()) { + db_connection.exec(table().build_sql_create_history_table()); + } } table().prepare(db_connection); @@ -324,6 +390,12 @@ void table_connection_t::stop(pg_conn_t const &db_connection, bool updateable, m_copy_mgr.sync(); if (append) { + if (table().has_history() && table().has_id_column()) { + auto const sql = table().build_sql_dedup_history(); + if (!sql.empty()) { + db_connection.exec(sql); + } + } return; } diff --git a/src/flex-table.hpp b/src/flex-table.hpp index 663375327..d72446962 100644 --- a/src/flex-table.hpp +++ b/src/flex-table.hpp @@ -174,6 +174,17 @@ class flex_table_t std::string full_name() const; std::string full_tmp_name() const; + bool has_history() const noexcept { return m_has_history; } + + void set_has_history(bool value) noexcept { m_has_history = value; } + + std::string history_name() const { return m_name + "_history"; } + std::string full_history_name() const; + + std::string build_sql_create_history_table() const; + + std::string build_sql_dedup_history() const; + bool has_multiple_geom_columns() const noexcept { return m_has_multiple_geom_columns; @@ -266,6 +277,8 @@ class flex_table_t /// Does this table have more than one geometry column? bool m_has_multiple_geom_columns = false; + bool m_has_history = false; + /// Always build the id index, not only when it is needed for updates? bool m_always_build_id_index = false; @@ -291,6 +304,9 @@ class table_connection_t table->build_sql_column_list())), m_copy_mgr(copy_thread) { + if (table->has_history()) { + m_target->set_history(table->full_history_name()); + } } void start(pg_conn_t const &db_connection, bool append) const;