Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 33 additions & 2 deletions src/db-copy.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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());
Expand All @@ -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),
Expand All @@ -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,
Expand All @@ -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');
Expand Down
8 changes: 7 additions & 1 deletion src/db-copy.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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;
};

/**
Expand All @@ -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; }
Expand Down Expand Up @@ -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; }
Expand Down Expand Up @@ -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);
}
}

Expand Down
32 changes: 32 additions & 0 deletions src/flex-lua-table.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<flex_table_t> *tables,
std::vector<expire_output_t> *expire_outputs,
std::string const &default_schema, bool updatable,
Expand All @@ -461,6 +492,7 @@ int setup_flex_table(lua_State *lua_state, std::vector<flex_table_t> *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{};
Expand Down
72 changes: 72 additions & 0 deletions src/flex-table.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down Expand Up @@ -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());
Expand Down Expand Up @@ -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);
Expand All @@ -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;
}

Expand Down
16 changes: 16 additions & 0 deletions src/flex-table.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;

Expand All @@ -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;
Expand Down
Loading