From 42238034490fb5479d787bd1695750387d508200 Mon Sep 17 00:00:00 2001 From: Adam Date: Mon, 24 Nov 2014 14:27:23 -0500 Subject: Rewrite serializable to have field level granularity Represent serializable objects in a digraph, and as a result made most object relationships implicitly defined, and use the graph to trace references between objects to determine relationships. Edges may also be marked as having a dependency of the object they point to, which allows for automatic cleanup and deletion of most objects when no longer needed. Additionally, this allows not having to require in-memory copies of everything when using external databases. db_sql has been rewritten for this and now always requires a database to function. db_sql with MySQL now requires InnoDB to make use of transactions and foreign key constraints. --- modules/database/db_flatfile.cpp | 368 +++++--------------- modules/database/db_old.cpp | 468 ++++++++++++++------------ modules/database/db_redis.cpp | 703 ++++++++++++++------------------------- modules/database/db_sql.cpp | 488 ++++++++++++++++----------- modules/database/db_sql_live.cpp | 268 --------------- 5 files changed, 860 insertions(+), 1435 deletions(-) delete mode 100644 modules/database/db_sql_live.cpp (limited to 'modules/database') diff --git a/modules/database/db_flatfile.cpp b/modules/database/db_flatfile.cpp index a2226a0ef..91e43e14e 100644 --- a/modules/database/db_flatfile.cpp +++ b/modules/database/db_flatfile.cpp @@ -10,104 +10,9 @@ #include "module.h" -#ifndef _WIN32 -#include -#endif - -class SaveData : public Serialize::Data -{ - public: - Anope::string last; - std::fstream *fs; - - SaveData() : fs(NULL) { } - - std::iostream& operator[](const Anope::string &key) override - { - if (key != last) - { - *fs << "\nDATA " << key << " "; - last = key; - } - - return *fs; - } -}; - -class LoadData : public Serialize::Data -{ - public: - std::fstream *fs; - unsigned int id; - std::map data; - std::stringstream ss; - bool read; - - LoadData() : fs(NULL), id(0), read(false) { } - - std::iostream& operator[](const Anope::string &key) override - { - if (!read) - { - for (Anope::string token; std::getline(*this->fs, token.str());) - { - if (token.find("ID ") == 0) - { - try - { - this->id = convertTo(token.substr(3)); - } - catch (const ConvertException &) { } - - continue; - } - else if (token.find("DATA ") != 0) - break; - - size_t sp = token.find(' ', 5); // Skip DATA - if (sp != Anope::string::npos) - data[token.substr(5, sp - 5)] = token.substr(sp + 1); - } - - read = true; - } - - ss.clear(); - this->ss << this->data[key]; - return this->ss; - } - - std::set KeySet() const override - { - std::set keys; - for (std::map::const_iterator it = this->data.begin(), it_end = this->data.end(); it != it_end; ++it) - keys.insert(it->first); - return keys; - } - - size_t Hash() const override - { - size_t hash = 0; - for (std::map::const_iterator it = this->data.begin(), it_end = this->data.end(); it != it_end; ++it) - if (!it->second.empty()) - hash ^= Anope::hash_cs()(it->second); - return hash; - } - - void Reset() - { - id = 0; - read = false; - data.clear(); - } -}; - -class DBFlatFile : public Module, public Pipe - , public EventHook - , public EventHook +class DBFlatFile : public Module , public EventHook , public EventHook - , public EventHook { /* Day the last backup was on */ int last_day; @@ -115,8 +20,6 @@ class DBFlatFile : public Module, public Pipe std::map > backups; bool loaded; - int child_pid; - void BackupDatabase() { tm *tm = localtime(&Anope::CurTime); @@ -125,16 +28,16 @@ class DBFlatFile : public Module, public Pipe { last_day = tm->tm_mday; - const std::vector &type_order = Serialize::Type::GetTypeOrder(); + const std::map &types = Serialize::TypeBase::GetTypes(); std::set dbs; dbs.insert(Config->GetModule(this)->Get("database", "anope.db")); - for (unsigned i = 0; i < type_order.size(); ++i) + for (const std::pair &p : types) { - Serialize::Type *stype = Serialize::Type::Find(type_order[i]); + Serialize::TypeBase *stype = p.second; - if (stype && stype->GetOwner()) + if (stype->GetOwner()) dbs.insert("module_" + stype->GetOwner()->name + ".db"); } @@ -177,66 +80,17 @@ class DBFlatFile : public Module, public Pipe public: DBFlatFile(const Anope::string &modname, const Anope::string &creator) : Module(modname, creator, DATABASE | VENDOR) - , EventHook("OnRestart") - , EventHook("OnShutdown") , EventHook("OnLoadDatabase") , EventHook("OnSaveDatabase") - , EventHook("OnSerializeTypeCreate") , last_day(0) , loaded(false) - , child_pid(-1) { } -#ifndef _WIN32 - void OnRestart() override - { - OnShutdown(); - } - - void OnShutdown() override - { - if (child_pid > -1) - { - Log(this) << "Waiting for child to exit..."; - - int status; - waitpid(child_pid, &status, 0); - - Log(this) << "Done"; - } - } -#endif - - void OnNotify() override - { - char buf[512]; - int i = this->Read(buf, sizeof(buf) - 1); - if (i <= 0) - return; - buf[i] = 0; - - child_pid = -1; - - if (!*buf) - { - Log(this) << "Finished saving databases"; - return; - } - - Log(this) << "Error saving databases: " << buf; - - if (!Config->GetModule(this)->Get("nobackupok")) - Anope::Quitting = true; - } - EventReturn OnLoadDatabase() override { - const std::vector &type_order = Serialize::Type::GetTypeOrder(); - std::set tried_dbs; - - const Anope::string &db_name = Anope::DataDir + "/" + Config->GetModule(this)->Get("database", "anope.db"); + const Anope::string &db_name = Anope::DataDir + "/" + Config->GetModule(this)->Get("database", "anope.db"); std::fstream fd(db_name.c_str(), std::ios_base::in | std::ios_base::binary); if (!fd.is_open()) @@ -245,32 +99,58 @@ class DBFlatFile : public Module, public Pipe return EVENT_STOP; } - std::map > positions; - + Serialize::TypeBase *type = nullptr; + Serialize::Object *obj = nullptr; for (Anope::string buf; std::getline(fd, buf.str());) + { + Log() << buf; if (buf.find("OBJECT ") == 0) - positions[buf.substr(7)].push_back(fd.tellg()); + { + Anope::string t = buf.substr(7); + if (obj) + Log(LOG_DEBUG) << "obj != null but got OBJECT"; + if (type) + Log(LOG_DEBUG) << "type != null but got OBJECT"; + type = Serialize::TypeBase::Find(t); + obj = nullptr; + } + else if (buf.find("ID ") == 0) + { + if (!type || obj) + continue; - LoadData ld; - ld.fs = &fd; + try + { + Serialize::ID id = convertTo(buf.substr(3)); + obj = type->Require(id); + } + catch (const ConvertException &) + { + Log(LOG_DEBUG) << "Unable to parse object id " << buf.substr(3); + } + } + else if (buf.find("DATA ") == 0) + { + if (!type) + continue; - for (unsigned i = 0; i < type_order.size(); ++i) - { - Serialize::Type *stype = Serialize::Type::Find(type_order[i]); - if (!stype || stype->GetOwner()) - continue; + if (!obj) + obj = type->Create(); - std::vector &pos = positions[stype->GetName()]; + size_t sp = buf.find(' ', 5); // Skip DATA + if (sp == Anope::string::npos) + continue; - for (unsigned j = 0; j < pos.size(); ++j) - { - fd.clear(); - fd.seekg(pos[j]); + Anope::string key = buf.substr(5, sp - 5), value = buf.substr(sp + 1); - Serializable *obj = stype->Unserialize(NULL, ld); - if (obj != NULL) - obj->id = ld.id; - ld.Reset(); + Serialize::FieldBase *field = type->GetField(key); + if (field) + field->UnserializeFromString(obj, value); + } + else if (buf.find("END") == 0) + { + type = nullptr; + obj = nullptr; } } @@ -283,142 +163,44 @@ class DBFlatFile : public Module, public Pipe void OnSaveDatabase() override { - if (child_pid > -1) - { - Log(this) << "Database save is already in progress!"; - return; - } - BackupDatabase(); - int i = -1; -#ifndef _WIN32 - if (!Anope::Quitting && Config->GetModule(this)->Get("fork")) - { - i = fork(); - if (i > 0) - { - child_pid = i; - return; - } - else if (i < 0) - Log(this) << "Unable to fork for database save"; - } -#endif + Anope::string db_name = Anope::DataDir + "/" + Config->GetModule(this)->Get("database", "anope.db"); - try - { - std::map databases; - - /* First open the databases of all of the registered types. This way, if we have a type with 0 objects, that database will be properly cleared */ - for (std::map::const_iterator it = Serialize::Type::GetTypes().begin(), it_end = Serialize::Type::GetTypes().end(); it != it_end; ++it) - { - Serialize::Type *s_type = it->second; - - if (databases[s_type->GetOwner()]) - continue; + if (Anope::IsFile(db_name)) + rename(db_name.c_str(), (db_name + ".tmp").c_str()); - Anope::string db_name; - if (s_type->GetOwner()) - db_name = Anope::DataDir + "/module_" + s_type->GetOwner()->name + ".db"; - else - db_name = Anope::DataDir + "/" + Config->GetModule(this)->Get("database", "anope.db"); - - if (Anope::IsFile(db_name)) - rename(db_name.c_str(), (db_name + ".tmp").c_str()); - - std::fstream *fs = databases[s_type->GetOwner()] = new std::fstream(db_name.c_str(), std::ios_base::out | std::ios_base::trunc | std::ios_base::binary); - - if (!fs->is_open()) - Log(this) << "Unable to open " << db_name << " for writing"; - } - - SaveData data; - const std::list &items = Serializable::GetItems(); - for (std::list::const_iterator it = items.begin(), it_end = items.end(); it != it_end; ++it) - { - Serializable *base = *it; - Serialize::Type *s_type = base->GetSerializableType(); - - data.fs = databases[s_type->GetOwner()]; - if (!data.fs || !data.fs->is_open()) - continue; + std::fstream f(db_name.c_str(), std::ios_base::out | std::ios_base::trunc | std::ios_base::binary); - *data.fs << "OBJECT " << s_type->GetName(); - if (base->id) - *data.fs << "\nID " << base->id; - base->Serialize(data); - *data.fs << "\nEND\n"; - } - - for (std::map::iterator it = databases.begin(), it_end = databases.end(); it != it_end; ++it) - { - std::fstream *f = it->second; - const Anope::string &db_name = Anope::DataDir + "/" + (it->first ? (it->first->name + ".db") : Config->GetModule(this)->Get("database", "anope.db")); - - if (!f->is_open() || !f->good()) - { - this->Write("Unable to write database " + db_name); - - f->close(); - - if (Anope::IsFile((db_name + ".tmp").c_str())) - rename((db_name + ".tmp").c_str(), db_name.c_str()); - } - else - { - f->close(); - unlink((db_name + ".tmp").c_str()); - } - - delete f; - } - } - catch (...) + if (!f.is_open()) { - if (i) - throw; + Log(this) << "Unable to open " << db_name << " for writing"; } - - if (!i) + else { - this->Notify(); - exit(0); + for (std::pair p : Serialize::objects) + { + Serialize::Object *object = p.second; + Serialize::TypeBase *s_type = object->GetSerializableType(); + + f << "OBJECT " << s_type->GetName() << "\n"; + f << "ID " << object->id << "\n"; + for (Serialize::FieldBase *field : s_type->fields) + f << "DATA " << field->GetName() << " " << field->SerializeToString(object) << "\n"; + f << "END\n"; + } } - } - - /* Load just one type. Done if a module is reloaded during runtime */ - void OnSerializeTypeCreate(Serialize::Type *stype) override - { - if (!loaded) - return; - - Anope::string db_name; - if (stype->GetOwner()) - db_name = Anope::DataDir + "/module_" + stype->GetOwner()->name + ".db"; - else - db_name = Anope::DataDir + "/" + Config->GetModule(this)->Get("database", "anope.db"); - std::fstream fd(db_name.c_str(), std::ios_base::in | std::ios_base::binary); - if (!fd.is_open()) + if (!f.is_open() || !f.good()) { - Log(this) << "Unable to open " << db_name << " for reading!"; - return; + f.close(); + rename((db_name + ".tmp").c_str(), db_name.c_str()); } - - LoadData ld; - ld.fs = &fd; - - for (Anope::string buf; std::getline(fd, buf.str());) + else { - if (buf == "OBJECT " + stype->GetName()) - { - stype->Unserialize(NULL, ld); - ld.Reset(); - } + f.close(); + unlink((db_name + ".tmp").c_str()); } - - fd.close(); } }; diff --git a/modules/database/db_old.cpp b/modules/database/db_old.cpp index 587a8f227..fe7a47dc7 100644 --- a/modules/database/db_old.cpp +++ b/modules/database/db_old.cpp @@ -14,9 +14,12 @@ #include "modules/cs_mode.h" #include "modules/bs_badwords.h" #include "modules/os_news.h" -#include "modules/suspend.h" #include "modules/os_forbid.h" #include "modules/cs_entrymsg.h" +#include "modules/ns_suspend.h" +#include "modules/cs_suspend.h" +#include "modules/cs_access.h" +#include "modules/ns_access.h" #define READ(x) \ if (true) \ @@ -93,6 +96,18 @@ else \ #define OLD_NEWS_OPER 1 #define OLD_NEWS_RANDOM 2 +enum +{ + TTB_BOLDS, + TTB_COLORS, + TTB_REVERSES, + TTB_UNDERLINES, + TTB_BADWORDS, + TTB_CAPS, + TTB_FLOOD, + TTB_REPEAT, +}; + static struct mlock_info { char c; @@ -144,19 +159,21 @@ enum static void process_mlock(ChanServ::Channel *ci, uint32_t lock, bool status, uint32_t *limit, Anope::string *key) { - ModeLocks *ml = ci->Require("modelocks"); + if (!mlocks) + return; + for (unsigned i = 0; i < (sizeof(mlock_infos) / sizeof(mlock_info)); ++i) if (lock & mlock_infos[i].m) { ChannelMode *cm = ModeManager::FindChannelModeByChar(mlock_infos[i].c); - if (cm && ml) + if (cm) { if (limit && mlock_infos[i].c == 'l') - ml->SetMLock(cm, status, stringify(*limit)); + mlocks->SetMLock(ci, cm, status, stringify(*limit)); else if (key && mlock_infos[i].c == 'k') - ml->SetMLock(cm, status, *key); + mlocks->SetMLock(ci, cm, status, *key); else - ml->SetMLock(cm, status); + mlocks->SetMLock(ci, cm, status); } } } @@ -436,7 +453,6 @@ static void LoadNicks() { if (!NickServ::service) return; - ServiceReference forbid("ForbidService", "forbid"); dbFILE *f = open_db_read("NickServ", "nick.db", 14); if (f == NULL) return; @@ -446,27 +462,33 @@ static void LoadNicks() Anope::string buffer; READ(read_string(buffer, f)); - NickServ::Account *nc = NickServ::service->CreateAccount(buffer); + + NickServ::Account *nc = NickServ::account.Create(); + nc->SetDisplay(buffer); const Anope::string settings[] = { "killprotect", "kill_quick", "ns_secure", "ns_private", "hide_email", "hide_mask", "hide_quit", "memo_signon", "memo_receive", "autoop", "msg", "ns_keepmodes" }; for (unsigned j = 0; j < sizeof(settings) / sizeof(Anope::string); ++j) - nc->Shrink(settings[j].upper()); + nc->UnsetS(settings[j].upper()); char pwbuf[32]; READ(read_buffer(pwbuf, f)); if (hashm == "plain") - my_b64_encode(pwbuf, nc->pass); + { + Anope::string p; + my_b64_encode(pwbuf, p); + nc->SetPassword(p); + } else if (hashm == "md5" || hashm == "oldmd5") - nc->pass = Hex(pwbuf, 16); + nc->SetPassword(Hex(pwbuf, 16)); else if (hashm == "sha1") - nc->pass = Hex(pwbuf, 20); + nc->SetPassword(Hex(pwbuf, 20)); else - nc->pass = Hex(pwbuf, strlen(pwbuf)); - nc->pass = hashm + ":" + nc->pass; + nc->SetPassword(Hex(pwbuf, strlen(pwbuf))); + nc->SetPassword(hashm + ":" + nc->GetPassword()); READ(read_string(buffer, f)); - nc->email = buffer; + nc->SetEmail(buffer); READ(read_string(buffer, f)); if (!buffer.empty()) @@ -481,103 +503,111 @@ static void LoadNicks() READ(read_uint32(&u32, f)); if (u32 & OLD_NI_KILLPROTECT) - nc->Extend("KILLPROTECT"); + nc->SetS("KILLPROTECT", true); if (u32 & OLD_NI_SECURE) - nc->Extend("NS_SECURE"); + nc->SetS("NS_SECURE", true); if (u32 & OLD_NI_MSG) - nc->Extend("MSG"); + nc->SetS("MSG", true); if (u32 & OLD_NI_MEMO_HARDMAX) - nc->Extend("MEMO_HARDMAX"); + nc->SetS("MEMO_HARDMAX", true); if (u32 & OLD_NI_MEMO_SIGNON) - nc->Extend("MEMO_SIGNON"); + nc->SetS("MEMO_SIGNON", true); if (u32 & OLD_NI_MEMO_RECEIVE) - nc->Extend("MEMO_RECEIVE"); + nc->SetS("MEMO_RECEIVE", true); if (u32 & OLD_NI_PRIVATE) - nc->Extend("NS_PRIVATE"); + nc->SetS("NS_PRIVATE", true); if (u32 & OLD_NI_HIDE_EMAIL) - nc->Extend("HIDE_EMAIL"); + nc->SetS("HIDE_EMAIL", true); if (u32 & OLD_NI_HIDE_MASK) - nc->Extend("HIDE_MASK"); + nc->SetS("HIDE_MASK", true); if (u32 & OLD_NI_HIDE_QUIT) - nc->Extend("HIDE_QUIT"); + nc->SetS("HIDE_QUIT", true); if (u32 & OLD_NI_KILL_QUICK) - nc->Extend("KILL_QUICK"); + nc->SetS("KILL_QUICK", true); if (u32 & OLD_NI_KILL_IMMED) - nc->Extend("KILL_IMMED"); + nc->SetS("KILL_IMMED", true); if (u32 & OLD_NI_MEMO_MAIL) - nc->Extend("MEMO_MAIL"); + nc->SetS("MEMO_MAIL", true); if (u32 & OLD_NI_HIDE_STATUS) - nc->Extend("HIDE_STATUS"); + nc->SetS("HIDE_STATUS", true); if (u32 & OLD_NI_SUSPENDED) { - SuspendInfo si; - si.what = nc->display; - si.when = si.expires = 0; - nc->Extend("NS_SUSPENDED", si); + if (nssuspendinfo) + { + NSSuspendInfo *si = nssuspendinfo.Create(); + si->SetAccount(nc); + } } if (!(u32 & OLD_NI_AUTOOP)) - nc->Extend("AUTOOP"); + nc->SetS("AUTOOP", true); uint16_t u16; READ(read_uint16(&u16, f)); switch (u16) { case LANG_ES: - nc->language = "es_ES"; + nc->SetLanguage("es_ES"); break; case LANG_PT: - nc->language = "pt_PT"; + nc->SetLanguage("pt_PT"); break; case LANG_FR: - nc->language = "fr_FR"; + nc->SetLanguage("fr_FR"); break; case LANG_TR: - nc->language = "tr_TR"; + nc->SetLanguage("tr_TR"); break; case LANG_IT: - nc->language = "it_IT"; + nc->SetLanguage("it_IT"); break; case LANG_DE: - nc->language = "de_DE"; + nc->SetLanguage("de_DE"); break; case LANG_CAT: - nc->language = "ca_ES"; // yes, iso639 defines catalan as CA + nc->SetLanguage("ca_ES"); // yes, iso639 defines catalan as CA break; case LANG_GR: - nc->language = "el_GR"; + nc->SetLanguage("el_GR"); break; case LANG_NL: - nc->language = "nl_NL"; + nc->SetLanguage("nl_NL"); break; case LANG_RU: - nc->language = "ru_RU"; + nc->SetLanguage("ru_RU"); break; case LANG_HUN: - nc->language = "hu_HU"; + nc->SetLanguage("hu_HU"); break; case LANG_PL: - nc->language = "pl_PL"; + nc->SetLanguage("pl_PL"); break; case LANG_EN_US: case LANG_JA_JIS: case LANG_JA_EUC: case LANG_JA_SJIS: // these seem to be unused default: - nc->language = "en"; + nc->SetLanguage("en"); } READ(read_uint16(&u16, f)); for (uint16_t j = 0; j < u16; ++j) { READ(read_string(buffer, f)); - nc->access.push_back(buffer); + + if (nsaccess) + { + NickAccess *a = nsaccess.Create(); + a->SetAccount(nc); + a->SetMask(buffer); + } } int16_t i16; READ(read_int16(&i16, f)); READ(read_int16(&i16, f)); - if (nc->memos) - nc->memos->memomax = i16; + MemoServ::MemoInfo *mi = nc->GetMemos(); + if (mi) + mi->SetMemoMax(i16); for (int16_t j = 0; j < i16; ++j) { MemoServ::Memo *m = MemoServ::service ? MemoServ::service->CreateMemo() : nullptr; @@ -587,26 +617,20 @@ static void LoadNicks() int32_t tmp32; READ(read_int32(&tmp32, f)); if (m) - m->time = tmp32; + m->SetTime(tmp32); char sbuf[32]; READ(read_buffer(sbuf, f)); if (m) - m->sender = sbuf; + m->SetSender(sbuf); Anope::string text; READ(read_string(text, f)); if (m) - m->text = text; - if (m) - m->owner = nc->display; - if (nc->memos && m) - nc->memos->memos->push_back(m); - else - delete m; + m->SetText(text); } READ(read_uint16(&u16, f)); READ(read_int16(&i16, f)); - Log(LOG_DEBUG) << "Loaded NickServ::Account " << nc->display; + Log(LOG_DEBUG) << "Loaded NickServ::Account " << nc->GetDisplay(); } for (int i = 0; i < 1024; ++i) @@ -640,41 +664,40 @@ static void LoadNicks() if (tmpu16 & OLD_NS_VERBOTEN) { - if (!forbid) + if (!forbiddata) { delete nc; continue; } - if (nc->display.find_first_of("?*") != Anope::string::npos) + if (nc->GetDisplay().find_first_of("?*") != Anope::string::npos) { delete nc; continue; } - ForbidData *d = forbid->CreateForbid(); - d->mask = nc->display; - d->creator = last_usermask; - d->reason = last_realname; - d->expires = 0; - d->created = 0; - d->type = FT_NICK; + ForbidData *d = forbiddata.Create(); + d->SetMask(nc->GetDisplay()); + d->SetCreator(last_usermask); + d->SetReason(last_realname); + d->SetType(FT_NICK); delete nc; - forbid->AddForbid(d); continue; } - NickServ::Nick *na = NickServ::service->CreateNick(nick, nc); - na->last_usermask = last_usermask; - na->last_realname = last_realname; - na->last_quit = last_quit; - na->time_registered = time_registered; - na->last_seen = last_seen; + NickServ::Nick *na = NickServ::nick.Create(); + na->SetNick(nick); + na->SetAccount(nc); + na->SetLastUsermask(last_usermask); + na->SetLastRealname(last_realname); + na->SetLastQuit(last_quit); + na->SetTimeRegistered(time_registered); + na->SetLastSeen(last_seen); if (tmpu16 & OLD_NS_NO_EXPIRE) - na->Extend("NS_NO_EXPIRE"); + na->SetS("NS_NO_EXPIRE", true); - Log(LOG_DEBUG) << "Loaded NickServ::Nick " << na->nick; + Log(LOG_DEBUG) << "Loaded NickServ::Nick " << na->GetNick(); } close_db(f); /* End of section Ia */ @@ -706,7 +729,7 @@ static void LoadVHosts() na->SetVhost(ident, host, creator, vtime); - Log() << "Loaded vhost for " << na->nick; + Log() << "Loaded vhost for " << na->GetNick(); } close_db(f); @@ -732,13 +755,14 @@ static void LoadBots() READ(read_int32(&created, f)); READ(read_int16(&chancount, f)); - BotInfo *bi = BotInfo::Find(nick, true); - if (!bi) - bi = new BotInfo(nick, user, host, real); - bi->created = created; + ServiceBot *bi = ServiceBot::Find(nick, true); + //XXX + // if (!bi) + // bi = new ServiceBot(nick, user, host, real); + bi->bi->SetCreated(created); if (flags & OLD_BI_PRIVATE) - bi->oper_only = true; + bi->bi->SetOperOnly(true); Log(LOG_DEBUG) << "Loaded bot " << bi->nick; } @@ -762,12 +786,13 @@ static void LoadChannels() Anope::string buffer; char namebuf[64]; READ(read_buffer(namebuf, f)); - ChanServ::Channel *ci = ChanServ::service->Create(namebuf); + ChanServ::Channel *ci = ChanServ::channel.Create(); + ci->SetName(namebuf); const Anope::string settings[] = { "keeptopic", "peace", "cs_private", "restricted", "cs_secure", "secureops", "securefounder", "signkick", "signkick_level", "topiclock", "persist", "noautoop", "cs_keepmodes" }; for (unsigned j = 0; j < sizeof(settings) / sizeof(Anope::string); ++j) - ci->Shrink(settings[j].upper()); + ci->UnsetS(settings[j].upper()); READ(read_string(buffer, f)); ci->SetFounder(NickServ::FindAccount(buffer)); @@ -778,71 +803,75 @@ static void LoadChannels() char pwbuf[32]; READ(read_buffer(pwbuf, f)); - READ(read_string(ci->desc, f)); + Anope::string desc; + READ(read_string(desc, f)); + ci->SetDesc(desc); READ(read_string(buffer, f)); READ(read_string(buffer, f)); int32_t tmp32; READ(read_int32(&tmp32, f)); - ci->time_registered = tmp32; + ci->SetTimeRegistered(tmp32); READ(read_int32(&tmp32, f)); - ci->last_used = tmp32; + ci->SetLastUsed(tmp32); - READ(read_string(ci->last_topic, f)); + Anope::string last_topic; + READ(read_string(last_topic, f)); + ci->SetLastTopic(last_topic); READ(read_buffer(pwbuf, f)); - ci->last_topic_setter = pwbuf; + ci->SetLastTopicSetter(pwbuf); READ(read_int32(&tmp32, f)); - ci->last_topic_time = tmp32; + ci->SetLastTopicTime(tmp32); uint32_t tmpu32; READ(read_uint32(&tmpu32, f)); // Temporary flags cleanup tmpu32 &= ~0x80000000; if (tmpu32 & OLD_CI_KEEPTOPIC) - ci->Extend("KEEPTOPIC"); + ci->SetS("KEEPTOPIC", true); if (tmpu32 & OLD_CI_SECUREOPS) - ci->Extend("SECUREOPS"); + ci->SetS("SECUREOPS", true); if (tmpu32 & OLD_CI_PRIVATE) - ci->Extend("CS_PRIVATE"); + ci->SetS("CS_PRIVATE", true); if (tmpu32 & OLD_CI_TOPICLOCK) - ci->Extend("TOPICLOCK"); + ci->SetS("TOPICLOCK", true); if (tmpu32 & OLD_CI_RESTRICTED) - ci->Extend("RESTRICTED"); + ci->SetS("RESTRICTED", true); if (tmpu32 & OLD_CI_PEACE) - ci->Extend("PEACE"); + ci->SetS("PEACE", true); if (tmpu32 & OLD_CI_SECURE) - ci->Extend("CS_SECURE"); + ci->SetS("CS_SECURE", true); if (tmpu32 & OLD_CI_NO_EXPIRE) - ci->Extend("CS_NO_EXPIRE"); + ci->SetS("CS_NO_EXPIRE", true); if (tmpu32 & OLD_CI_MEMO_HARDMAX) - ci->Extend("MEMO_HARDMAX"); + ci->SetS("MEMO_HARDMAX", true); if (tmpu32 & OLD_CI_SECUREFOUNDER) - ci->Extend("SECUREFOUNDER"); + ci->SetS("SECUREFOUNDER", true); if (tmpu32 & OLD_CI_SIGNKICK) - ci->Extend("SIGNKICK"); + ci->SetS("SIGNKICK", true); if (tmpu32 & OLD_CI_SIGNKICK_LEVEL) - ci->Extend("SIGNKICK_LEVEL"); + ci->SetS("SIGNKICK_LEVEL", true); Anope::string forbidby, forbidreason; READ(read_string(forbidby, f)); READ(read_string(forbidreason, f)); if (tmpu32 & OLD_CI_SUSPENDED) { - SuspendInfo si; - si.what = ci->name; - si.by = forbidby; - si.reason = forbidreason; - si.when = si.expires = 0; - ci->Extend("CS_SUSPENDED", si); + if (cssuspendinfo) + { + CSSuspendInfo *si = cssuspendinfo.Create(); + si->SetChannel(ci); + si->SetBy(forbidby); + } } bool forbid_chan = tmpu32 & OLD_CI_VERBOTEN; int16_t tmp16; READ(read_int16(&tmp16, f)); - ci->bantype = tmp16; + ci->SetBanType(tmp16); READ(read_int16(&tmp16, f)); if (tmp16 > 36) @@ -856,13 +885,12 @@ static void LoadChannels() level = ChanServ::ACCESS_FOUNDER; if (j == 10 && level < 0) // NOJOIN - ci->Shrink("RESTRICTED"); // If CSDefRestricted was enabled this can happen + ci->UnsetS("RESTRICTED"); // If CSDefRestricted was enabled this can happen ci->SetLevel(GetLevelName(j), level); } bool xop = tmpu32 & OLD_CI_XOP; - ServiceReference provider_access("AccessProvider", "access/access"), provider_xop("AccessProvider", "access/xop"); uint16_t tmpu16; READ(read_uint16(&tmpu16, f)); for (uint16_t j = 0; j < tmpu16; ++j) @@ -875,15 +903,17 @@ static void LoadChannels() if (xop) { - if (provider_xop) - access = provider_xop->Create(); + if (xopchanaccess) + access = xopchanaccess.Create(); } else - if (provider_access) - access = provider_access->Create(); + { + if (accesschanaccess) + access = accesschanaccess.Create(); + } if (access) - access->ci = ci; + access->SetChannel(ci); int16_t level; READ(read_int16(&level, f)); @@ -914,16 +944,19 @@ static void LoadChannels() Anope::string mask; READ(read_string(mask, f)); if (access) - access->SetMask(mask, ci); + { + access->SetMask(mask); + NickServ::Nick *na = NickServ::FindNick(mask); + if (na) + na->SetAccount(na->GetAccount()); + } READ(read_int32(&tmp32, f)); if (access) { - access->last_seen = tmp32; - access->creator = "Unknown"; - access->created = Anope::CurTime; - - ci->AddAccess(access); + access->SetLastSeen(tmp32); + access->SetCreator("Unknown"); + access->SetCreated(Anope::CurTime); } } } @@ -958,8 +991,9 @@ static void LoadChannels() READ(read_int16(&tmp16, f)); READ(read_int16(&tmp16, f)); - if (ci->memos) - ci->memos->memomax = tmp16; + MemoServ::MemoInfo *mi = ci->GetMemos(); + if (mi) + mi->SetMemoMax(tmp16); for (int16_t j = 0; j < tmp16; ++j) { READ(read_uint32(&tmpu32, f)); @@ -967,74 +1001,56 @@ static void LoadChannels() MemoServ::Memo *m = MemoServ::service ? MemoServ::service->CreateMemo() : nullptr; READ(read_int32(&tmp32, f)); if (m) - m->time = tmp32; + m->SetTime(tmp32); char sbuf[32]; READ(read_buffer(sbuf, f)); if (m) - m->sender = sbuf; + m->SetSender(sbuf); Anope::string text; READ(read_string(text, f)); if (m) - m->text = text; - if (m) - m->owner = ci->name; - if (ci->memos && m) - ci->memos->memos->push_back(m); - else - delete m; + m->SetText(text); } READ(read_string(buffer, f)); if (!buffer.empty()) { - EntryMessageList *eml = ci->Require("entrymsg"); - if (eml) + if (entrymsg) { - EntryMsg *e = eml->Create(); - - e->chan = ci->name; - e->creator = "Unknown"; - e->message = buffer; - e->when = Anope::CurTime; - - (*eml)->push_back(e); + EntryMsg *e = entrymsg.Create(); + e->SetChannel(ci); + e->SetCreator("Unknown"); + e->SetMessage(buffer); + e->SetWhen(Anope::CurTime); } } READ(read_string(buffer, f)); - ci->bi = BotInfo::Find(buffer, true); + ci->SetBot(ServiceBot::Find(buffer, true)); READ(read_int32(&tmp32, f)); if (tmp32 & OLD_BS_DONTKICKOPS) - ci->Extend("BS_DONTKICKOPS"); + ci->SetS("BS_DONTKICKOPS", true); if (tmp32 & OLD_BS_DONTKICKVOICES) - ci->Extend("BS_DONTKICKVOICES"); + ci->SetS("BS_DONTKICKVOICES", true); if (tmp32 & OLD_BS_FANTASY) - ci->Extend("BS_FANTASY"); + ci->SetS("BS_FANTASY", true); if (tmp32 & OLD_BS_GREET) - ci->Extend("BS_GREET"); + ci->SetS("BS_GREET", true); if (tmp32 & OLD_BS_NOBOT) - ci->Extend("BS_NOBOT"); + ci->SetS("BS_NOBOT", true); - KickerData *kd = ci->Require("kickerdata"); + KickerData *kd = GetKickerData(ci); if (kd) { - if (tmp32 & OLD_BS_KICK_BOLDS) - kd->bolds = true; - if (tmp32 & OLD_BS_KICK_COLORS) - kd->colors = true; - if (tmp32 & OLD_BS_KICK_REVERSES) - kd->reverses = true; - if (tmp32 & OLD_BS_KICK_UNDERLINES) - kd->underlines = true; - if (tmp32 & OLD_BS_KICK_BADWORDS) - kd->badwords = true; - if (tmp32 & OLD_BS_KICK_CAPS) - kd->caps = true; - if (tmp32 & OLD_BS_KICK_FLOOD) - kd->flood = true; - if (tmp32 & OLD_BS_KICK_REPEAT) - kd->repeat = true; + kd->SetBolds(tmp32 & OLD_BS_KICK_BOLDS); + kd->SetColors(tmp32 & OLD_BS_KICK_COLORS); + kd->SetReverses(tmp32 & OLD_BS_KICK_REVERSES); + kd->SetUnderlines(tmp32 & OLD_BS_KICK_UNDERLINES); + kd->SetBadwords(tmp32 & OLD_BS_KICK_BADWORDS); + kd->SetCaps(tmp32 & OLD_BS_KICK_CAPS); + kd->SetFlood(tmp32 & OLD_BS_KICK_FLOOD); + kd->SetRepeat(tmp32 & OLD_BS_KICK_REPEAT); } READ(read_int16(&tmp16, f)); @@ -1042,27 +1058,51 @@ static void LoadChannels() { int16_t ttb; READ(read_int16(&ttb, f)); - if (j < TTB_SIZE && kd) - kd->ttb[j] = ttb; + switch (j) + { + case TTB_BOLDS: + kd->SetTTBBolds(ttb); + break; + case TTB_COLORS: + kd->SetTTBColors(ttb); + break; + case TTB_REVERSES: + kd->SetTTBReverses(ttb); + break; + case TTB_UNDERLINES: + kd->SetTTBUnderlines(ttb); + break; + case TTB_BADWORDS: + kd->SetTTBBadwords(ttb); + break; + case TTB_CAPS: + kd->SetTTBCaps(ttb); + break; + case TTB_FLOOD: + kd->SetTTBFlood(ttb); + break; + case TTB_REPEAT: + kd->SetTTBRepeat(ttb); + break; + } } READ(read_int16(&tmp16, f)); if (kd) - kd->capsmin = tmp16; + kd->SetCapsMin(tmp16); READ(read_int16(&tmp16, f)); if (kd) - kd->capspercent = tmp16; + kd->SetCapsPercent(tmp16); READ(read_int16(&tmp16, f)); if (kd) - kd->floodlines = tmp16; + kd->SetFloodLines(tmp16); READ(read_int16(&tmp16, f)); if (kd) - kd->floodsecs = tmp16; + kd->SetFloodSecs(tmp16); READ(read_int16(&tmp16, f)); if (kd) - kd->repeattimes = tmp16; + kd->SetRepeatTimes(tmp16); - BadWords *bw = ci->Require("badwords"); READ(read_uint16(&tmpu16, f)); for (uint16_t j = 0; j < tmpu16; ++j) { @@ -1082,38 +1122,35 @@ static void LoadChannels() else if (type == 3) bwtype = BW_END; - if (bw) - bw->AddBadWord(buffer, bwtype); + if (badwords) + badwords->AddBadWord(ci, buffer, bwtype); } } if (forbid_chan) { - if (!forbid) + if (!forbiddata) { delete ci; continue; } - if (ci->name.find_first_of("?*") != Anope::string::npos) + if (ci->GetName().find_first_of("?*") != Anope::string::npos) { delete ci; continue; } - ForbidData *d = forbid->CreateForbid(); - d->mask = ci->name; - d->creator = forbidby; - d->reason = forbidreason; - d->expires = 0; - d->created = 0; - d->type = FT_CHAN; + ForbidData *d = forbiddata.Create(); + d->SetMask(ci->GetName()); + d->SetCreator(forbidby); + d->SetReason(forbidreason); + d->SetType(FT_CHAN); delete ci; - forbid->AddForbid(d); continue; } - Log(LOG_DEBUG) << "Loaded channel " << ci->name; + Log(LOG_DEBUG) << "Loaded channel " << ci->GetName(); } close_db(f); @@ -1128,9 +1165,8 @@ static void LoadOper() XLineManager *akill, *sqline, *snline, *szline; akill = sqline = snline = szline = NULL; - for (std::list::iterator it = XLineManager::XLineManagers.begin(), it_end = XLineManager::XLineManagers.end(); it != it_end; ++it) + for (XLineManager *xl : XLineManager::XLineManagers) { - XLineManager *xl = *it; if (xl->Type() == 'G') akill = xl; else if (xl->Type() == 'Q') @@ -1163,7 +1199,7 @@ static void LoadOper() continue; XLine *x = new XLine(user + "@" + host, by, expires, reason, XLineManager::GenerateUID()); - x->created = seton; + x->SetCreated(seton); akill->AddXLine(x); } @@ -1183,7 +1219,7 @@ static void LoadOper() continue; XLine *x = new XLine(mask, by, expires, reason, XLineManager::GenerateUID()); - x->created = seton; + x->SetCreated(seton); snline->AddXLine(x); } @@ -1203,7 +1239,7 @@ static void LoadOper() continue; XLine *x = new XLine(mask, by, expires, reason, XLineManager::GenerateUID()); - x->created = seton; + x->SetCreated(seton); sqline->AddXLine(x); } @@ -1223,7 +1259,7 @@ static void LoadOper() continue; XLine *x = new XLine(mask, by, expires, reason, XLineManager::GenerateUID()); - x->created = seton; + x->SetCreated(seton); szline->AddXLine(x); } @@ -1255,14 +1291,16 @@ static void LoadExceptions() READ(read_int32(&time, f)); READ(read_int32(&expires, f)); - Exception *exception = session_service->CreateException(); - exception->mask = mask; - exception->limit = limit; - exception->who = who; - exception->time = time; - exception->expires = expires; - exception->reason = reason; - session_service->AddException(exception); + if (exception && session_service) + { + Exception *e = exception.Create(); + e->SetMask(mask); + e->SetLimit(limit); + e->SetWho(who); + e->SetTime(time); + e->SetExpires(expires); + e->SetReason(reason); + } } close_db(f); @@ -1270,7 +1308,7 @@ static void LoadExceptions() static void LoadNews() { - if (!news_service) + if (!newsitem) return; dbFILE *f = open_db_read("OperServ", "news.db", 9); @@ -1284,37 +1322,37 @@ static void LoadNews() for (int16_t i = 0; i < n; i++) { int16_t type; - NewsItem *ni = news_service->CreateNewsItem(); + NewsItem *ni = newsitem.Create(); READ(read_int16(&type, f)); switch (type) { case OLD_NEWS_LOGON: - ni->type = NEWS_LOGON; + ni->SetNewsType(NEWS_LOGON); break; case OLD_NEWS_OPER: - ni->type = NEWS_OPER; + ni->SetNewsType(NEWS_OPER); break; case OLD_NEWS_RANDOM: - ni->type = NEWS_RANDOM; + ni->SetNewsType(NEWS_RANDOM); break; } int32_t unused; READ(read_int32(&unused, f)); - READ(read_string(ni->text, f)); + Anope::string text; + READ(read_string(text, f)); + ni->SetText(text); char who[32]; READ(read_buffer(who, f)); - ni->who = who; + ni->SetWho(who); int32_t tmp; READ(read_int32(&tmp, f)); - ni->time = tmp; - - news_service->AddNewsItem(ni); + ni->SetTime(tmp); } close_db(f); @@ -1324,8 +1362,8 @@ class DBOld : public Module , public EventHook , public EventHook { - PrimitiveExtensibleItem mlock_on, mlock_off, mlock_limit; - PrimitiveExtensibleItem mlock_key; + ExtensibleItem mlock_on, mlock_off, mlock_limit; + ExtensibleItem mlock_key; public: DBOld(const Anope::string &modname, const Anope::string &creator) : Module(modname, creator, DATABASE | VENDOR) diff --git a/modules/database/db_redis.cpp b/modules/database/db_redis.cpp index dd5ff0301..6105bfa65 100644 --- a/modules/database/db_redis.cpp +++ b/modules/database/db_redis.cpp @@ -16,100 +16,35 @@ using namespace Redis; class DatabaseRedis; static DatabaseRedis *me; -class Data : public Serialize::Data -{ - public: - std::map data; - - ~Data() - { - for (std::map::iterator it = data.begin(), it_end = data.end(); it != it_end; ++it) - delete it->second; - } - - std::iostream& operator[](const Anope::string &key) override - { - std::stringstream* &stream = data[key]; - if (!stream) - stream = new std::stringstream(); - return *stream; - } - - std::set KeySet() const override - { - std::set keys; - for (std::map::const_iterator it = this->data.begin(), it_end = this->data.end(); it != it_end; ++it) - keys.insert(it->first); - return keys; - } - - size_t Hash() const override - { - size_t hash = 0; - for (std::map::const_iterator it = this->data.begin(), it_end = this->data.end(); it != it_end; ++it) - if (!it->second->str().empty()) - hash ^= Anope::hash_cs()(it->second->str()); - return hash; - } -}; - class TypeLoader : public Interface { - Anope::string type; + Serialize::TypeBase *type; + public: - TypeLoader(Module *creator, const Anope::string &t) : Interface(creator), type(t) { } + TypeLoader(Module *creator, Serialize::TypeBase *t) : Interface(creator), type(t) { } void OnResult(const Reply &r) override; }; class ObjectLoader : public Interface { - Anope::string type; - int64_t id; + Serialize::Object *obj; public: - ObjectLoader(Module *creator, const Anope::string &t, int64_t i) : Interface(creator), type(t), id(i) { } + ObjectLoader(Module *creator, Serialize::Object *s) : Interface(creator), obj(s) { } void OnResult(const Reply &r) override; }; -class IDInterface : public Interface +class FieldLoader : public Interface { - Reference o; - public: - IDInterface(Module *creator, Serializable *obj) : Interface(creator), o(obj) { } - - void OnResult(const Reply &r) override; -}; + Serialize::Object *obj; + Serialize::FieldBase *field; -class Deleter : public Interface -{ - Anope::string type; - int64_t id; public: - Deleter(Module *creator, const Anope::string &t, int64_t i) : Interface(creator), type(t), id(i) { } + FieldLoader(Module *creator, Serialize::Object *o, Serialize::FieldBase *f) : Interface(creator), obj(o), field(f) { } - void OnResult(const Reply &r) override; -}; - -class Updater : public Interface -{ - Anope::string type; - int64_t id; - public: - Updater(Module *creator, const Anope::string &t, int64_t i) : Interface(creator), type(t), id(i) { } - - void OnResult(const Reply &r) override; -}; - -class ModifiedObject : public Interface -{ - Anope::string type; - int64_t id; - public: - ModifiedObject(Module *creator, const Anope::string &t, int64_t i) : Interface(creator), type(t), id(i) { } - - void OnResult(const Reply &r) override; + void OnResult(const Reply &) override; }; class SubscriptionListener : public Interface @@ -120,341 +55,263 @@ class SubscriptionListener : public Interface void OnResult(const Reply &r) override; }; -class DatabaseRedis : public Module, public Pipe +class DatabaseRedis : public Module , public EventHook - , public EventHook - , public EventHook - , public EventHook - , public EventHook + , public EventHook { SubscriptionListener sl; - std::set updated_items; public: ServiceReference redis; DatabaseRedis(const Anope::string &modname, const Anope::string &creator) : Module(modname, creator, DATABASE | VENDOR) , EventHook("OnLoadDatabase") - , EventHook("OnSerializeTypeCreate") - , EventHook("OnSerializableConstruct") - , EventHook("OnSerializableDestruct") - , EventHook("OnSerializableUpdate") + , EventHook("OnSerialize") , sl(this) { me = this; - - } - - /* Insert or update an object */ - void InsertObject(Serializable *obj) - { - Serialize::Type *t = obj->GetSerializableType(); - - /* If there is no id yet for ths object, get one */ - if (!obj->id) - redis->SendCommand(new IDInterface(this, obj), "INCR id:" + t->GetName()); - else - { - Data data; - obj->Serialize(data); - - if (obj->IsCached(data)) - return; - - obj->UpdateCache(data); - - std::vector args; - args.push_back("HGETALL"); - args.push_back("hash:" + t->GetName() + ":" + stringify(obj->id)); - - /* Get object attrs to clear before updating */ - redis->SendCommand(new Updater(this, t->GetName(), obj->id), args); - } - } - - void OnNotify() override - { - for (std::set::iterator it = this->updated_items.begin(), it_end = this->updated_items.end(); it != it_end; ++it) - { - Serializable *s = *it; - - this->InsertObject(s); - } - - this->updated_items.clear(); } void OnReload(Configuration::Conf *conf) override { Configuration::Block *block = conf->GetModule(this); - this->redis = ServiceReference("Redis::Provider", block->Get("engine", "redis/main")); + this->redis = ServiceReference("Redis::Provider", block->Get("engine", "redis/main")); } EventReturn OnLoadDatabase() override { - const std::vector type_order = Serialize::Type::GetTypeOrder(); - for (unsigned i = 0; i < type_order.size(); ++i) - { - Serialize::Type *sb = Serialize::Type::Find(type_order[i]); - this->OnSerializeTypeCreate(sb); - } + if (!redis) + return EVENT_STOP; + + const std::map &types = Serialize::TypeBase::GetTypes(); + for (const std::pair &p : types) + this->OnSerializeTypeCreate(p.second); while (redis->BlockAndProcess()); - redis->Subscribe(&this->sl, "__keyspace@*__:hash:*"); + redis->Subscribe(&this->sl, "anope"); return EVENT_STOP; } - void OnSerializeTypeCreate(Serialize::Type *sb) override + void OnSerializeTypeCreate(Serialize::TypeBase *sb) { - if (!redis) - return; - - std::vector args; - args.push_back("SMEMBERS"); - args.push_back("ids:" + sb->GetName()); + std::vector args = { "SMEMBERS", "ids:" + sb->GetName() }; - redis->SendCommand(new TypeLoader(this, sb->GetName()), args); + redis->SendCommand(new TypeLoader(this, sb), args); } - void OnSerializableConstruct(Serializable *obj) override + EventReturn OnSerializeList(Serialize::TypeBase *type, std::vector &ids) override { - this->updated_items.insert(obj); - this->Notify(); + return EVENT_CONTINUE; } - void OnSerializableDestruct(Serializable *obj) override + EventReturn OnSerializeFind(Serialize::TypeBase *type, Serialize::FieldBase *field, const Anope::string &value, Serialize::ID &id) override { - Serialize::Type *t = obj->GetSerializableType(); - - std::vector args; - args.push_back("HGETALL"); - args.push_back("hash:" + t->GetName() + ":" + stringify(obj->id)); + return EVENT_CONTINUE; + } - /* Get all of the attributes for this object */ - redis->SendCommand(new Deleter(this, t->GetName(), obj->id), args); + EventReturn OnSerializeGet(Serialize::Object *object, Serialize::FieldBase *field, Anope::string &value) override + { + return EVENT_CONTINUE; + } - this->updated_items.erase(obj); - t->objects.erase(obj->id); - this->Notify(); + EventReturn OnSerializeGetRefs(Serialize::Object *object, Serialize::TypeBase *type, std::vector &) override + { + return EVENT_CONTINUE; } - void OnSerializableUpdate(Serializable *obj) override + EventReturn OnSerializeDeref(Serialize::ID id, Serialize::TypeBase *type) override { - this->updated_items.insert(obj); - this->Notify(); + return EVENT_CONTINUE; } -}; -void TypeLoader::OnResult(const Reply &r) -{ - if (r.type != Reply::MULTI_BULK || !me->redis) + EventReturn OnSerializeGetSerializable(Serialize::Object *object, Serialize::FieldBase *field, Anope::string &type, Serialize::ID &value) override { - delete this; - return; + return EVENT_CONTINUE; } - for (unsigned i = 0; i < r.multi_bulk.size(); ++i) + EventReturn OnSerializeSet(Serialize::Object *object, Serialize::FieldBase *field, const Anope::string &value) override { - const Reply *reply = r.multi_bulk[i]; + std::vector args; - if (reply->type != Reply::BULK) - continue; + redis->StartTransaction(); - int64_t id; - try - { - id = convertTo(reply->bulk); - } - catch (const ConvertException &) - { - continue; - } + const Anope::string &old = field->SerializeToString(object); + args = { "SREM", "lookup:" + object->GetSerializableType()->GetName() + ":" + field->GetName() + ":" + old, stringify(object->id) }; + redis->SendCommand(nullptr, args); - std::vector args; - args.push_back("HGETALL"); - args.push_back("hash:" + this->type + ":" + stringify(id)); + // add object to type set + args = { "SADD", "ids:" + object->GetSerializableType()->GetName(), stringify(object->id) }; + redis->SendCommand(nullptr, args); - me->redis->SendCommand(new ObjectLoader(me, this->type, id), args); - } + // add key to key set + args = { "SADD", "keys:" + stringify(object->id), field->GetName() }; + redis->SendCommand(nullptr, args); - delete this; -} + // set value + args = { "SET", "values:" + stringify(object->id) + ":" + field->GetName(), value }; + redis->SendCommand(nullptr, args); -void ObjectLoader::OnResult(const Reply &r) -{ - Serialize::Type *st = Serialize::Type::Find(this->type); + // lookup + args = { "SADD", "lookup:" + object->GetSerializableType()->GetName() + ":" + field->GetName() + ":" + value, stringify(object->id) }; + redis->SendCommand(nullptr, args); - if (r.type != Reply::MULTI_BULK || r.multi_bulk.empty() || !me->redis || !st) - { - delete this; - return; + redis->CommitTransaction(); + + return EVENT_CONTINUE; } - Data data; + EventReturn OnSerializeSetSerializable(Serialize::Object *object, Serialize::FieldBase *field, Serialize::Object *value) override + { + return OnSerializeSet(object, field, stringify(value->id)); + } - for (unsigned i = 0; i + 1 < r.multi_bulk.size(); i += 2) + EventReturn OnSerializeUnset(Serialize::Object *object, Serialize::FieldBase *field) override { - const Reply *key = r.multi_bulk[i], - *value = r.multi_bulk[i + 1]; + std::vector args; - data[key->bulk] << value->bulk; + redis->StartTransaction(); + + const Anope::string &old = field->SerializeToString(object); + args = { "SREM", "lookup:" + object->GetSerializableType()->GetName() + ":" + field->GetName() + ":" + old, stringify(object->id) }; + redis->SendCommand(nullptr, args); + + // remove field from set + args = { "SREM", "keys:" + stringify(object->id), field->GetName() }; + redis->SendCommand(nullptr, args); + + redis->CommitTransaction(); + + return EVENT_CONTINUE; } - Serializable* &obj = st->objects[this->id]; - obj = st->Unserialize(obj, data); - if (obj) + EventReturn OnSerializeUnsetSerializable(Serialize::Object *object, Serialize::FieldBase *field) override { - obj->id = this->id; - obj->UpdateCache(data); + return OnSerializeUnset(object, field); } - delete this; -} - -void IDInterface::OnResult(const Reply &r) -{ - if (!o || r.type != Reply::INT || !r.i) + EventReturn OnSerializeHasField(Serialize::Object *object, Serialize::FieldBase *field) override { - delete this; - return; + return EVENT_CONTINUE; } - Serializable* &obj = o->GetSerializableType()->objects[r.i]; - if (obj) - /* This shouldn't be possible */ - obj->id = 0; - - o->id = r.i; - obj = o; + EventReturn OnSerializableGetId(Serialize::ID &id) override + { + std::vector args = { "INCR", "id" }; - /* Now that we have the id, insert this object for real */ - anope_dynamic_static_cast(this->owner)->InsertObject(o); + auto f = [&](const Reply &r) + { + id = r.i; + }; - delete this; -} + FInterface inter(this, f); + redis->SendCommand(&inter, args); + while (redis->BlockAndProcess()); + return EVENT_ALLOW; + } -void Deleter::OnResult(const Reply &r) -{ - if (r.type != Reply::MULTI_BULK || !me->redis || r.multi_bulk.empty()) + void OnSerializableCreate(Serialize::Object *) override { - delete this; - return; } - /* Transaction start */ - me->redis->StartTransaction(); + void OnSerializableDelete(Serialize::Object *obj) override + { + std::vector args; - std::vector args; - args.push_back("DEL"); - args.push_back("hash:" + this->type + ":" + stringify(this->id)); + redis->StartTransaction(); - /* Delete hash object */ - me->redis->SendCommand(NULL, args); + for (Serialize::FieldBase *field : obj->GetSerializableType()->fields) + { + Anope::string value = field->SerializeToString(obj); - args.clear(); - args.push_back("SREM"); - args.push_back("ids:" + this->type); - args.push_back(stringify(this->id)); + args = { "SREM", "lookup:" + obj->GetSerializableType()->GetName() + ":" + field->GetName() + ":" + value, stringify(obj->id) }; + redis->SendCommand(nullptr, args); - /* Delete id from ids set */ - me->redis->SendCommand(NULL, args); + args = { "DEL", "values:" + stringify(obj->id) + ":" + field->GetName() }; + redis->SendCommand(nullptr, args); - for (unsigned i = 0; i + 1 < r.multi_bulk.size(); i += 2) - { - const Reply *key = r.multi_bulk[i], - *value = r.multi_bulk[i + 1]; + args = { "SREM", "keys:" + stringify(obj->id), field->GetName() }; + redis->SendCommand(nullptr, args); + } - args.clear(); - args.push_back("SREM"); - args.push_back("value:" + this->type + ":" + key->bulk + ":" + value->bulk); - args.push_back(stringify(this->id)); + args = { "SREM", "ids:" + obj->GetSerializableType()->GetName(), stringify(obj->id) }; + redis->SendCommand(nullptr, args); - /* Delete value -> object id */ - me->redis->SendCommand(NULL, args); + redis->CommitTransaction(); } +}; - /* Transaction end */ - me->redis->CommitTransaction(); - - delete this; -} - -void Updater::OnResult(const Reply &r) +void TypeLoader::OnResult(const Reply &r) { - Serialize::Type *st = Serialize::Type::Find(this->type); - - if (!st) + if (r.type != Reply::MULTI_BULK || !me->redis) { delete this; return; } - Serializable *obj = st->objects[this->id]; - if (!obj) + for (unsigned i = 0; i < r.multi_bulk.size(); ++i) { - delete this; - return; - } + const Reply *reply = r.multi_bulk[i]; - Data data; - obj->Serialize(data); + if (reply->type != Reply::BULK) + continue; - /* Transaction start */ - me->redis->StartTransaction(); + int64_t id; + try + { + id = convertTo(reply->bulk); + } + catch (const ConvertException &) + { + continue; + } - for (unsigned i = 0; i + 1 < r.multi_bulk.size(); i += 2) - { - const Reply *key = r.multi_bulk[i], - *value = r.multi_bulk[i + 1]; + Serialize::Object *obj = type->Require(id); + if (obj == nullptr) + { + Log(LOG_DEBUG) << "redis: Unable to require object #" << id << " of type " << type->GetName(); + continue; + } - std::vector args; - args.push_back("SREM"); - args.push_back("value:" + this->type + ":" + key->bulk + ":" + value->bulk); - args.push_back(stringify(this->id)); + std::vector args = { "SMEMBERS", "keys:" + stringify(id) }; - /* Delete value -> object id */ - me->redis->SendCommand(NULL, args); + me->redis->SendCommand(new ObjectLoader(me, obj), args); } - /* Add object id to id set for this type */ - std::vector args; - args.push_back("SADD"); - args.push_back("ids:" + this->type); - args.push_back(stringify(obj->id)); - me->redis->SendCommand(NULL, args); - - args.clear(); - args.push_back("HMSET"); - args.push_back("hash:" + this->type + ":" + stringify(obj->id)); + delete this; +} - typedef std::map items; - for (items::iterator it = data.data.begin(), it_end = data.data.end(); it != it_end; ++it) +void ObjectLoader::OnResult(const Reply &r) +{ + if (r.type != Reply::MULTI_BULK || r.multi_bulk.empty() || !me->redis) { - const Anope::string &key = it->first; - std::stringstream *value = it->second; + delete this; + return; + } - args.push_back(key); - args.push_back(value->str()); + Serialize::TypeBase *type = obj->GetSerializableType(); - std::vector args2; + for (Reply *reply : r.multi_bulk) + { + const Anope::string &key = reply->bulk; + Serialize::FieldBase *field = type->GetField(key); - args2.push_back("SADD"); - args2.push_back("value:" + this->type + ":" + key + ":" + value->str()); - args2.push_back(stringify(obj->id)); + if (field == nullptr) + continue; - /* Add to value -> object id set */ - me->redis->SendCommand(NULL, args2); - } + std::vector args = { "GET", "values:" + stringify(obj->id) + ":" + key }; - ++obj->redis_ignore; + me->redis->SendCommand(new FieldLoader(me, obj, field), args); + } - /* Add object */ - me->redis->SendCommand(NULL, args); + delete this; +} - /* Transaction end */ - me->redis->CommitTransaction(); +void FieldLoader::OnResult(const Reply &r) +{ + Log(LOG_DEBUG_2) << "redis: Setting field " << field->GetName() << " of object #" << obj->id << " of type " << obj->GetSerializableType()->GetName() << " to " << r.bulk; + field->UnserializeFromString(obj, r.bulk); delete this; } @@ -462,196 +319,118 @@ void Updater::OnResult(const Reply &r) void SubscriptionListener::OnResult(const Reply &r) { /* - * [May 15 13:59:35.645839 2013] Debug: pmessage - * [May 15 13:59:35.645866 2013] Debug: __keyspace@*__:anope:hash:* - * [May 15 13:59:35.645880 2013] Debug: __keyspace@0__:anope:hash:type:id - * [May 15 13:59:35.645893 2013] Debug: hset + * message + * anope + * message + * + * set 4 email adam@anope.org + * unset 4 email + * create 4 NickCore + * delete 4 */ - if (r.multi_bulk.size() != 4) - return; - - size_t sz = r.multi_bulk[2]->bulk.find(':'); - if (sz == Anope::string::npos) - return; - const Anope::string &key = r.multi_bulk[2]->bulk.substr(sz + 1), - &op = r.multi_bulk[3]->bulk; + const Anope::string &message = r.multi_bulk[2]->bulk; + Anope::string command; + spacesepstream sep(message); - sz = key.rfind(':'); - if (sz == Anope::string::npos) - return; - - const Anope::string &id = key.substr(sz + 1); - - size_t sz2 = key.rfind(':', sz - 1); - if (sz2 == Anope::string::npos) - return; - const Anope::string &type = key.substr(sz2 + 1, sz - sz2 - 1); + sep.GetToken(command); - Serialize::Type *s_type = Serialize::Type::Find(type); - - if (s_type == NULL) - return; - - uint64_t obj_id; - try + if (command == "set" || command == "unset") { - obj_id = convertTo(id); - } - catch (const ConvertException &) - { - return; - } + Anope::string sid, key, value; - if (op == "hset" || op == "hdel") - { - Serializable *s = s_type->objects[obj_id]; + sep.GetToken(sid); + sep.GetToken(key); + value = sep.GetRemaining(); - if (s && s->redis_ignore) + Serialize::ID id; + try { - --s->redis_ignore; - Log(LOG_DEBUG) << "redis: notify: got modify for object id " << obj_id << " of type " << type << ", but I am ignoring it"; + id = convertTo(sid); } - else + catch (const ConvertException &ex) { - Log(LOG_DEBUG) << "redis: notify: got modify for object id " << obj_id << " of type " << type; - - std::vector args; - args.push_back("HGETALL"); - args.push_back("hash:" + type + ":" + id); - - me->redis->SendCommand(new ModifiedObject(me, type, obj_id), args); - } - } - else if (op == "del") - { - Serializable* &s = s_type->objects[obj_id]; - if (s == NULL) + Log(LOG_DEBUG) << "redis: unable to get id for SL update key " << sid; return; + } - Log(LOG_DEBUG) << "redis: notify: deleting object id " << obj_id << " of type " << type; - - Data data; - - s->Serialize(data); - - /* Transaction start */ - me->redis->StartTransaction(); - - typedef std::map items; - for (items::iterator it = data.data.begin(), it_end = data.data.end(); it != it_end; ++it) + Serialize::Object *obj = Serialize::GetID(id); + if (obj == nullptr) { - const Anope::string &k = it->first; - std::stringstream *value = it->second; - - std::vector args; - args.push_back("SREM"); - args.push_back("value:" + type + ":" + k + ":" + value->str()); - args.push_back(id); - - /* Delete value -> object id */ - me->redis->SendCommand(NULL, args); + Log(LOG_DEBUG) << "redis: pmessage for unknown object #" << id; + return; } - std::vector args; - args.push_back("SREM"); - args.push_back("ids:" + type); - args.push_back(stringify(s->id)); - - /* Delete object from id set */ - me->redis->SendCommand(NULL, args); - - /* Transaction end */ - me->redis->CommitTransaction(); - - delete s; - s = NULL; - } -} - -void ModifiedObject::OnResult(const Reply &r) -{ - Serialize::Type *st = Serialize::Type::Find(this->type); + Serialize::FieldBase *field = obj->GetSerializableType()->GetField(key); + if (field == nullptr) + { + Log(LOG_DEBUG) << "redis: pmessage for unknown field of object #" << id << ": " << key; + return; + } - if (!st) - { - delete this; - return; + Log(LOG_DEBUG_2) << "redis: Setting field " << field->GetName() << " of object #" << obj->id << " of type " << obj->GetSerializableType()->GetName() << " to " << value; + field->UnserializeFromString(obj, value); } - - Serializable* &obj = st->objects[this->id]; - - /* Transaction start */ - me->redis->StartTransaction(); - - /* Erase old object values */ - if (obj) + else if (command == "create") { - Data data; + Anope::string sid, stype; - obj->Serialize(data); + sep.GetToken(sid); + sep.GetToken(stype); - typedef std::map items; - for (items::iterator it = data.data.begin(), it_end = data.data.end(); it != it_end; ++it) + Serialize::ID id; + try + { + id = convertTo(sid); + } + catch (const ConvertException &ex) { - const Anope::string &key = it->first; - std::stringstream *value = it->second; + Log(LOG_DEBUG) << "redis: unable to get id for SL update key " << sid; + return; + } - std::vector args; - args.push_back("SREM"); - args.push_back("value:" + st->GetName() + ":" + key + ":" + value->str()); - args.push_back(stringify(this->id)); + Serialize::TypeBase *type = Serialize::TypeBase::Find(stype); + if (type == nullptr) + { + Log(LOG_DEBUG) << "redis: pmessage create for nonexistant type " << stype; + return; + } - /* Delete value -> object id */ - me->redis->SendCommand(NULL, args); + Serialize::Object *obj = type->Require(id); + if (obj == nullptr) + { + Log(LOG_DEBUG) << "redis: require for pmessage create type " << type->GetName() << " id #" << id << " returned nullptr"; + return; } } - - Data data; - - for (unsigned i = 0; i + 1 < r.multi_bulk.size(); i += 2) + else if (command == "delete") { - const Reply *key = r.multi_bulk[i], - *value = r.multi_bulk[i + 1]; - - data[key->bulk] << value->bulk; - } + Anope::string sid; - obj = st->Unserialize(obj, data); - if (obj) - { - obj->id = this->id; - obj->UpdateCache(data); + sep.GetToken(sid); - /* Insert new object values */ - typedef std::map items; - for (items::iterator it = data.data.begin(), it_end = data.data.end(); it != it_end; ++it) + Serialize::ID id; + try { - const Anope::string &key = it->first; - std::stringstream *value = it->second; - - std::vector args; - args.push_back("SADD"); - args.push_back("value:" + st->GetName() + ":" + key + ":" + value->str()); - args.push_back(stringify(obj->id)); - - /* Add to value -> object id set */ - me->redis->SendCommand(NULL, args); + id = convertTo(sid); + } + catch (const ConvertException &ex) + { + Log(LOG_DEBUG) << "redis: unable to get id for SL update key " << sid; + return; } - std::vector args; - args.push_back("SADD"); - args.push_back("ids:" + st->GetName()); - args.push_back(stringify(obj->id)); + Serialize::Object *obj = Serialize::GetID(id); + if (obj == nullptr) + { + Log(LOG_DEBUG) << "redis: message for unknown object #" << id; + return; + } - /* Add to type -> id set */ - me->redis->SendCommand(NULL, args); + obj->Delete(); } - - /* Transaction end */ - me->redis->CommitTransaction(); - - delete this; + else + Log(LOG_DEBUG) << "redis: unknown message: " << message; } MODULE_INIT(DatabaseRedis) diff --git a/modules/database/db_sql.cpp b/modules/database/db_sql.cpp index 5d213cb9a..3481d3a89 100644 --- a/modules/database/db_sql.cpp +++ b/modules/database/db_sql.cpp @@ -1,279 +1,373 @@ -/* - * (C) 2003-2014 Anope Team - * Contact us at team@anope.org - * - * Please read COPYING and README for further details. - * - * Based on the original code of Epona by Lara. - * Based on the original code of Services by Andy Church. - */ - #include "module.h" #include "modules/sql.h" using namespace SQL; -class SQLSQLInterface : public Interface +class DBMySQL : public Module, public Pipe + , public EventHook { - public: - SQLSQLInterface(Module *o) : Interface(o) { } + private: + bool transaction = false; + bool inited = false; + Anope::string prefix; + ServiceReference SQL; - void OnResult(const Result &r) override + Result Run(const Query &query) { - Log(LOG_DEBUG) << "SQL successfully executed query: " << r.finished_query; + if (!SQL) +