// Cyphesis Online RPG Server and AI Engine // Copyright (C) 2000-2007 Alistair Riddoch // Copyright (C) 2018 Erik Ogenvik // // This program is free software; you can redistribute it and/or modify // it under the terms of the GNU General Public License as published by // the Free Software Foundation; either version 2 of the License, or // (at your option) any later version. // // This program is distributed in the hope that it will be useful, // but WITHOUT ANY WARRANTY; without even the implied warranty of // MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the // GNU General Public License for more details. // // You should have received a copy of the GNU General Public License // along with this program; if not, write to the Free Software Foundation, // Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA #include "DatabasePostgres.h" #include "id.h" #include "log.h" #include "debug.h" #include "globals.h" #include "compose.hpp" #include "const.h" #include #include #include using Atlas::Message::Element; using Atlas::Message::MapType; using Atlas::Objects::Root; using String::compose; typedef Atlas::Codecs::XML Serialiser; static const bool debug_flag = false; static void databaseNotice(void*, const char* message) { log(NOTICE, "Notice from database:"); log_formatted(NOTICE, message); } DatabasePostgres::DatabasePostgres() : Database(), m_connection(nullptr) { } DatabasePostgres::~DatabasePostgres() { if (!pendingQueries.empty()) { log(ERROR, compose("Database delete with %1 queries pending", pendingQueries.size())); } shutdownConnection(); } bool DatabasePostgres::tuplesOk() { assert(m_connection != nullptr); bool status = false; PGresult* res; while ((res = PQgetResult(m_connection)) != nullptr) { if (PQresultStatus(res) == PGRES_TUPLES_OK) { status = true; } PQclear(res); }; return status; } int DatabasePostgres::commandOk() { assert(m_connection != nullptr); int status = -1; PGresult* res; while ((res = PQgetResult(m_connection)) != nullptr) { if (PQresultStatus(res) == PGRES_COMMAND_OK) { status = 0; } else { reportError(); } PQclear(res); }; return status; } int DatabasePostgres::connect(const std::string& context, std::string& error_msg) { std::stringstream conninfos; std::string db_server; if (readConfigItem(context, "dbserver", db_server) == 0) { if (db_server.empty()) { log(WARNING, "Empty database hostname specified in config file. " "Using none."); } else { conninfos << "host=" << db_server << " "; } } std::string dbname; if (context == CYPHESIS) { dbname = CYPHESIS; } else { dbname = compose("%1_%2", CYPHESIS, ::instance); } readConfigItem(context, "dbname", dbname); conninfos << "dbname=" << dbname << " "; std::string db_user; if (readConfigItem(context, "dbuser", db_user) == 0) { if (db_user.empty()) { log(WARNING, "Empty username specified in config file. " "Using current user."); } else { conninfos << "user=" << db_user << " "; } } std::string db_passwd; if (readConfigItem(context, "dbpasswd", db_passwd) == 0) { conninfos << "password=" << db_passwd << " "; } const std::string cinfo = conninfos.str(); m_connection = PQconnectdb(cinfo.c_str()); if (m_connection == nullptr) { error_msg = "Unknown error"; return -1; } if (PQstatus(m_connection) != CONNECTION_OK) { error_msg = PQerrorMessage(m_connection); PQfinish(m_connection); m_connection = nullptr; return -1; } return 0; } int DatabasePostgres::initConnection() { std::string error_message; if (connect(::instance, error_message) != 0) { log(ERROR, "Connection to database failed:"); log_formatted(ERROR, error_message); return -1; } PQsetNoticeProcessor(m_connection, databaseNotice, nullptr); return 0; } void DatabasePostgres::shutdownConnection() { if (m_connection != nullptr) { PQfinish(m_connection); m_connection = nullptr; } } int DatabasePostgres::encodeObject(const MapType& o, std::string& data) { std::stringstream str; Serialiser codec(str, str, m_d); Atlas::Message::Encoder enc(codec); codec.streamBegin(); enc.streamMessageElement(o); codec.streamEnd(); std::string raw = str.str(); char safe[raw.size() * 2 + 1]; int errcode; PQescapeStringConn(m_connection, safe, raw.c_str(), raw.size(), &errcode); if (errcode != 0) { std::cerr << "ERROR: " << errcode << std::endl << std::flush; } data = safe; return 0; } int DatabasePostgres::getObject(const std::string& table, const std::string& key, MapType& o) { assert(m_connection != nullptr); debug(std::cout << "Database::getObject() " << table << "." << key << std::endl << std::flush;); std::string query = std::string("SELECT * FROM ") + table + " WHERE id = '" + key + "'"; clearPendingQuery(); int status = PQsendQuery(m_connection, query.c_str()); if (!status) { reportError(); return -1; } PGresult* res; if ((res = PQgetResult(m_connection)) == nullptr) { debug(std::cout << "Error accessing " << key << " in " << table << " table" << std::endl << std::flush;); return -1; } if (PQntuples(res) < 1 || PQnfields(res) < 2) { debug(std::cout << "No entry for " << key << " in " << table << " table" << std::endl << std::flush;); PQclear(res); while ((res = PQgetResult(m_connection)) != nullptr) { PQclear(res); } return -1; } const char* data = PQgetvalue(res, 0, 1); debug(std::cout << "Got record " << key << " from database, value " << data << std::endl << std::flush;); int ret = decodeMessage(data, o); PQclear(res); while ((res = PQgetResult(m_connection)) != nullptr) { PQclear(res); log(ERROR, "Extra database result to simple query."); }; return ret; } void DatabasePostgres::reportError() { assert(m_connection != nullptr); char* message = PQerrorMessage(m_connection); assert(message != nullptr); if (strlen(message) < 2) { log(WARNING, "Zero length database error message"); } std::string msg = std::string("DATABASE: ") + message; msg = msg.substr(0, msg.size() - 1); log(ERROR, msg); } DatabaseResult DatabasePostgres::runSimpleSelectQuery(const std::string& query) { assert(m_connection != nullptr); debug(std::cout << "QUERY: " << query << std::endl << std::flush;); clearPendingQuery(); int status = PQsendQuery(m_connection, query.c_str()); if (!status) { log(ERROR, "runSimpleSelectQuery(): Database query error."); reportError(); return DatabaseResult(nullptr); } debug(std::cout << "done" << std::endl << std::flush;); PGresult* res; if ((res = PQgetResult(m_connection)) == nullptr) { log(ERROR, "Error selecting."); reportError(); debug(std::cout << "Row query didn't work" << std::endl << std::flush;); return DatabaseResult(nullptr); } if (PQresultStatus(res) != PGRES_TUPLES_OK) { log(ERROR, "Error selecting row."); debug(std::cout << "QUERY: " << query << std::endl << std::flush;); reportError(); PQclear(res); res = nullptr; } PGresult* nres; while ((nres = PQgetResult(m_connection)) != nullptr) { PQclear(nres); log(ERROR, "Extra database result to simple query."); }; return DatabaseResult(std::unique_ptr(new DatabaseResultWorkerPostgres(res))); } int DatabasePostgres::runCommandQuery(const std::string& query) { assert(m_connection != nullptr); clearPendingQuery(); int status = PQsendQuery(m_connection, query.c_str()); if (!status) { log(ERROR, "runCommandQuery(): Database query error."); reportError(); return -1; } if (commandOk() != 0) { log(ERROR, "Error running command query row."); log(NOTICE, query); reportError(); debug(std::cout << "Row query didn't work" << std::endl << std::flush;); } else { debug(std::cout << "Query worked" << std::endl << std::flush;); return 0; } return -1; } int DatabasePostgres::registerRelation(std::string& tablename, const std::string& sourcetable, const std::string& targettable, RelationType kind) { assert(m_connection != nullptr); tablename = sourcetable + "_" + targettable; std::string query = compose("SELECT * FROM %1 WHERE source = 0 " "AND target = 0", tablename); debug(std::cout << "QUERY: " << query << std::endl << std::flush;); clearPendingQuery(); int status = PQsendQuery(m_connection, query.c_str()); if (!status) { log(ERROR, "registerRelation(): Database query error."); reportError(); return -1; } if (!tuplesOk()) { debug(reportError();); debug(std::cout << "Table does not yet exist" << std::endl << std::flush;); } else { debug(std::cout << "Table exists" << std::endl << std::flush;); allTables.insert(tablename); return 0; } query = "CREATE TABLE "; query += tablename; if (kind == OneToOne || kind == ManyToOne) { query += " (source integer UNIQUE REFERENCES "; } else { query += " (source integer REFERENCES "; } query += sourcetable; if (kind == OneToOne || kind == OneToMany) { query += " (id), target integer UNIQUE REFERENCES "; } else { query += " (id), target integer REFERENCES "; } query += targettable; query += " (id) ON DELETE CASCADE ) WITHOUT OIDS"; debug(std::cout << "CREATE QUERY: " << query << std::endl << std::flush;); if (runCommandQuery(query) != 0) { return -1; } allTables.insert(tablename); #if 0 if (kind == ManyToOne || kind == OneToOne) { return true; } else { std::string indexQuery = "CREATE INDEX "; indexQuery += tablename; indexQuery += "_source_idx ON "; indexQuery += tablename; indexQuery += " (source)"; return runCommandQuery(indexQuery) == 0; } #else return 0; #endif } int DatabasePostgres::registerSimpleTable(const std::string& name, const MapType& row) { assert(m_connection != nullptr); if (row.empty()) { log(ERROR, "Attempt to create empty database table"); } // Check whether the table exists std::string query = "SELECT * FROM "; std::string createquery = "CREATE TABLE "; query += name; createquery += name; query += " WHERE id = 0"; createquery += " (id integer UNIQUE PRIMARY KEY"; auto Iend = row.end(); for (auto I = row.begin(); I != Iend; ++I) { query += " AND "; createquery += ", "; const std::string& column = I->first; query += column; createquery += column; const Element& type = I->second; if (type.isString()) { query += " LIKE 'foo'"; std::size_t size = type.String().size(); if (size == 0) { createquery += " text"; } else { char buf[32]; snprintf(buf, 32, "%zd", size); createquery += " varchar("; createquery += buf; createquery += ")"; } } else if (type.isInt()) { query += " = 1"; createquery += " integer"; } else if (type.isFloat()) { query += " = 1.0"; createquery += " float"; } else { log(ERROR, "Illegal column type in database simple row"); } } debug(std::cout << "QUERY: " << query << std::endl << std::flush;); clearPendingQuery(); int status = PQsendQuery(m_connection, query.c_str()); if (!status) { log(ERROR, "registerSimpleTable(): Database query error."); reportError(); return -1; } if (!tuplesOk()) { debug(reportError();); debug(std::cout << "Table does not yet exist" << std::endl << std::flush;); } else { debug(std::cout << "Table exists" << std::endl << std::flush;); allTables.insert(name); return 0; } createquery += ") WITHOUT OIDS"; debug(std::cout << "CREATE QUERY: " << createquery << std::endl << std::flush;); int ret = runCommandQuery(createquery); if (ret == 0) { allTables.insert(name); } return ret; } int DatabasePostgres::registerEntityIdGenerator() { assert(m_connection != nullptr); clearPendingQuery(); int status = PQsendQuery(m_connection, "SELECT * FROM entity_ent_id_seq"); if (!status) { log(ERROR, "registerEntityIdGenerator(): Database query error."); reportError(); return -1; } if (!tuplesOk()) { debug(reportError();); debug(std::cout << "Sequence does not yet exist" << std::endl << std::flush;); } else { debug(std::cout << "Sequence exists" << std::endl << std::flush;); return 0; } return runCommandQuery("CREATE SEQUENCE entity_ent_id_seq"); } long DatabasePostgres::newId(std::string& id) { assert(m_connection != nullptr); clearPendingQuery(); int status = PQsendQuery(m_connection, "SELECT nextval('entity_ent_id_seq')"); if (!status) { log(ERROR, "newId(): Database query error."); reportError(); return -1; } PGresult* res; if ((res = PQgetResult(m_connection)) == nullptr) { log(ERROR, "Error getting new ID."); reportError(); return -1; } const char* cid = PQgetvalue(res, 0, 0); id = cid; PQclear(res); while ((res = PQgetResult(m_connection)) != nullptr) { PQclear(res); log(ERROR, "Extra database result to simple query."); }; if (id.empty()) { log(ERROR, "Unknown error getting ID from database."); return -1; } return forceIntegerId(id); } int DatabasePostgres::registerEntityTable(const std::map& chunks) { assert(m_connection != nullptr); clearPendingQuery(); int status = PQsendQuery(m_connection, "SELECT * FROM entities"); if (!status) { log(ERROR, "registerEntityIdGenerator(): Database query error."); reportError(); return -1; } if (!tuplesOk()) { debug(reportError();); debug(std::cout << "Table does not yet exist" << std::endl << std::flush;); } else { allTables.insert("entities"); // FIXME Flush out the whole state of the databases, to ensure they // don't clog up while we are testing. // runCommandQuery("DELETE FROM properties"); // runCommandQuery(compose("DELETE FROM entities WHERE id!=%1", // consts::rootWorldIntId)); debug(std::cout << "Table exists" << std::endl << std::flush;); return 0; } std::string query = compose("CREATE TABLE entities (" "id integer UNIQUE PRIMARY KEY, " "loc integer, " "type varchar(%1), " "seq integer", consts::id_len); auto I = chunks.begin(); auto Iend = chunks.end(); for (; I != Iend; ++I) { query += compose(", %1 varchar(1024)", I->first); } query += ")"; if (runCommandQuery(query) != 0) { return -1; } allTables.insert("entities"); query = compose("INSERT INTO entities VALUES (%1, null, 'world')", consts::rootWorldIntId); if (runCommandQuery(query) != 0) { return -1; } return 0; } int DatabasePostgres::registerPropertyTable() { assert(m_connection != nullptr); clearPendingQuery(); int status = PQsendQuery(m_connection, "SELECT * FROM properties"); if (!status) { log(ERROR, "registerPropertyIdGenerator(): Database query error."); reportError(); return -1; } if (!tuplesOk()) { debug(reportError();); debug(std::cout << "Table does not yet exist" << std::endl << std::flush;); } else { allTables.insert("properties"); debug(std::cout << "Table exists" << std::endl << std::flush;); return 0; } allTables.insert("properties"); std::string query = compose("CREATE TABLE properties (" "id integer REFERENCES entities " "ON DELETE CASCADE, " "name varchar(%1), " "value text)", consts::id_len); if (runCommandQuery(query) != 0) { reportError(); return -1; } query = "CREATE INDEX property_names on properties (name)"; if (runCommandQuery(query) != 0) { reportError(); return -1; } return 0; } int DatabasePostgres::registerThoughtsTable() { assert(m_connection != nullptr); clearPendingQuery(); int status = PQsendQuery(m_connection, "SELECT * FROM thoughts"); if (!status) { log(ERROR, "registerThoughtsTable(): Database query error."); reportError(); return -1; } if (!tuplesOk()) { debug(reportError();); debug(std::cout << "Table does not yet exist" << std::endl << std::flush;); } else { allTables.insert("thoughts"); debug(std::cout << "Table exists" << std::endl << std::flush;); return 0; } allTables.insert("properties"); std::string query = "CREATE TABLE thoughts (" "id integer REFERENCES entities " "ON DELETE CASCADE, " "thought text)"; if (runCommandQuery(query) != 0) { reportError(); return -1; } return 0; } // General functions for handling queries at the low level. void DatabasePostgres::queryResult(ExecStatusType status) { if (!m_queryInProgress || pendingQueries.empty()) { log(ERROR, "Got database result when no query was pending."); return; } DatabaseQuery& q = pendingQueries.front(); if (q.second == PGRES_EMPTY_QUERY) { log(ERROR, "Got database result which is already done."); return; } if (q.second == status) { debug(std::cout << "Query status ok" << std::endl << std::flush;); // Mark this query as done q.second = PGRES_EMPTY_QUERY; } else { log(ERROR, "Database error from async query"); std::cerr << "Query error in : " << q.first << std::endl << std::flush; reportError(); q.second = PGRES_EMPTY_QUERY; } } void DatabasePostgres::queryComplete() { if (!m_queryInProgress || pendingQueries.empty()) { log(ERROR, "Got database query complete when no query was pending"); return; } DatabaseQuery& q = pendingQueries.front(); if (q.second != PGRES_EMPTY_QUERY) { log(ERROR, "Got database query complete when query was not done"); return; } debug(std::cout << "Query complete" << std::endl << std::flush;); pendingQueries.pop_front(); m_queryInProgress = false; } int DatabasePostgres::launchNewQuery() { if (m_connection == nullptr) { log(ERROR, "Can't launch new query while database is offline."); return -1; } if (m_queryInProgress) { log(ERROR, "Launching new query when query is in progress"); return -1; } if (pendingQueries.empty()) { debug(std::cout << "No queries to launch" << std::endl << std::flush;); return -1; } debug(std::cout << pendingQueries.size() << " queries pending" << std::endl << std::flush;); DatabaseQuery& q = pendingQueries.front(); debug(std::cout << "Launching async query: " << q.first << std::endl << std::flush;); int status = PQsendQuery(m_connection, q.first.c_str()); if (!status) { log(ERROR, "Database query error when launching."); reportError(); return -1; } else { m_queryInProgress = true; PQflush(m_connection); return 0; } } int DatabasePostgres::scheduleCommand(const std::string& query) { pendingQueries.push_back(std::make_pair(query, PGRES_COMMAND_OK)); if (!m_queryInProgress) { debug(std::cout << "Query: " << query << " launched" << std::endl << std::flush;); return launchNewQuery(); } else { debug(std::cout << "Query: " << query << " scheduled" << std::endl << std::flush;); return 0; } } int DatabasePostgres::clearPendingQuery() { if (!m_queryInProgress) { return 0; } assert(!pendingQueries.empty()); debug(std::cout << "Clearing a pending query" << std::endl << std::flush;); DatabaseQuery& q = pendingQueries.front(); if (q.second == PGRES_COMMAND_OK) { m_queryInProgress = false; pendingQueries.pop_front(); return commandOk(); } else { log(ERROR, "Pending query wants unknown status"); return -1; } } int DatabasePostgres::runMaintainance(unsigned int command) { // VACUUM and REINDEX tables from a common store if ((command & MAINTAIN_REINDEX) == MAINTAIN_REINDEX) { std::string query("REINDEX TABLE "); auto Iend = allTables.end(); for (auto I = allTables.begin(); I != Iend; ++I) { scheduleCommand(query + *I); } } if ((command & MAINTAIN_VACUUM) == MAINTAIN_VACUUM) { std::string query("VACUUM "); if ((command & MAINTAIN_VACUUM_ANALYZE) == MAINTAIN_VACUUM_ANALYZE) { query += "ANALYZE "; } if ((command & MAINTAIN_VACUUM_FULL) == MAINTAIN_VACUUM_FULL) { query += "FULL "; } auto Iend = allTables.end(); for (auto I = allTables.begin(); I != Iend; ++I) { scheduleCommand(query + *I); } } return 0; } // //const char* DatabaseResultPostgres::field(const char* column, int row) const //{ // int col_num = PQfnumber(m_res.get(), column); // if (col_num == -1) { // return ""; // } // return PQgetvalue(m_res.get(), row, col_num); //} // //const char* DatabaseResultPostgres::const_iterator_postgres::column(const char* column) const //{ // int col_num = PQfnumber(m_dr.m_res.get(), column); // if (col_num == -1) { // return ""; // } // return PQgetvalue(m_dr.m_res.get(), m_row, col_num); //} // //void DatabaseResultPostgres::const_iterator_postgres::readColumn(const char* column, // int& val) const //{ // int col_num = PQfnumber(m_dr.m_res.get(), column); // if (col_num == -1) { // return; // } // const char* v = PQgetvalue(m_dr.m_res.get(), m_row, col_num); // val = static_cast(strtol(v, 0, 10)); //} // //void DatabaseResultPostgres::const_iterator_postgres::readColumn(const char* column, // float& val) const //{ // int col_num = PQfnumber(m_dr.m_res.get(), column); // if (col_num == -1) { // return; // } // const char* v = PQgetvalue(m_dr.m_res.get(), m_row, col_num); // val = strtof(v, 0); //} // //void DatabaseResultPostgres::const_iterator_postgres::readColumn(const char* column, // double& val) const //{ // int col_num = PQfnumber(m_dr.m_res.get(), column); // if (col_num == -1) { // return; // } // const char* v = PQgetvalue(m_dr.m_res.get(), m_row, col_num); // val = strtod(v, 0); //} // //void DatabaseResultPostgres::const_iterator_postgres::readColumn(const char* column, // std::string& val) const //{ // int col_num = PQfnumber(m_dr.m_res.get(), column); // if (col_num == -1) { // return; // } // const char* v = PQgetvalue(m_dr.m_res.get(), m_row, col_num); // val = v; //} // //void DatabaseResultPostgres::const_iterator_postgres::readColumn(const char* column, // MapType& val) const //{ // int col_num = PQfnumber(m_dr.m_res.get(), column); // if (col_num == -1) { // return; // } // const char* v = PQgetvalue(m_dr.m_res.get(), m_row, col_num); // Database::instance().decodeMessage(v, val); //}