#include "store/queues.hpp" #include #include #include #include #include #include namespace velox::daemon::store { namespace proto = velox::proto; namespace { // One queue row (columns queue_id, name, state, max_concurrent, schedule, on_complete, in // that order) plus its member taskIds, read off the row a caller has already step()'d to. DbResult project_row(Db& db, Stmt& st) { proto::Queue q; q.queueId = st.column_text(0); q.name = st.column_text(1); if (auto s = proto::parse_QueueState(st.column_text(2))) q.state = *s; q.maxConcurrent = st.column_int(3); if (!st.column_is_null(4)) { auto j = nlohmann::json::parse(st.column_text(4), nullptr, false); if (auto sched = proto::parse(j, "schedule")) q.schedule = *sched; } if (auto oc = proto::parse_QueueOnComplete(st.column_text(5))) q.onComplete = *oc; auto ts = db.prepare("SELECT task_id FROM tasks WHERE queue_id = ?1 ORDER BY queue_position"); if (!ts) return std::unexpected(ts.error()); if (auto b = ts->bind(1, std::string_view(q.queueId)); !b) return std::unexpected(b.error()); std::vector ids; for (;;) { auto r = ts->step(); if (!r) return std::unexpected(r.error()); if (!*r) break; ids.push_back(ts->column_text(0)); } q.taskIds = std::move(ids); return q; } } // namespace DbResult> Queues::list() { auto st = db_.prepare( "SELECT queue_id, name, state, max_concurrent, schedule, on_complete FROM queues " "ORDER BY name"); if (!st) return std::unexpected(st.error()); std::vector out; for (;;) { auto row = st->step(); if (!row) return std::unexpected(row.error()); if (!*row) break; auto q = project_row(db_, *st); if (!q) return std::unexpected(q.error()); out.push_back(std::move(*q)); } return out; } DbResult> Queues::get(std::string_view queue_id) { auto st = db_.prepare( "SELECT queue_id, name, state, max_concurrent, schedule, on_complete FROM queues " "WHERE queue_id = ?1"); if (!st) return std::unexpected(st.error()); if (auto b = st->bind(1, queue_id); !b) return std::unexpected(b.error()); auto row = st->step(); if (!row) return std::unexpected(row.error()); if (!*row) return std::optional{}; auto q = project_row(db_, *st); if (!q) return std::unexpected(q.error()); return std::optional{std::move(*q)}; } DbResult Queues::set_state(std::string_view queue_id, std::string_view state) { auto st = db_.prepare("UPDATE queues SET state = ?2 WHERE queue_id = ?1"); if (!st) return std::unexpected(st.error()); if (auto b = st->bind(1, queue_id); !b) return std::unexpected(b.error()); if (auto b = st->bind(2, state); !b) return std::unexpected(b.error()); if (auto r = st->step(); !r) return std::unexpected(r.error()); return sqlite3_changes(db_.raw()) > 0; } DbResult Queues::set_schedule(std::string_view queue_id, const std::optional& schedule) { auto st = db_.prepare("UPDATE queues SET schedule = ?2 WHERE queue_id = ?1"); if (!st) return std::unexpected(st.error()); if (auto b = st->bind(1, queue_id); !b) return std::unexpected(b.error()); if (auto r = schedule ? st->bind(2, std::string_view(nlohmann::json(*schedule).dump())) : st->bind_null(2); !r) return std::unexpected(r.error()); if (auto r = st->step(); !r) return std::unexpected(r.error()); return sqlite3_changes(db_.raw()) > 0; } DbResult Queues::reorder(std::string_view queue_id, const std::vector& task_ids) { std::vector current; { auto st = db_.prepare( "SELECT task_id FROM tasks WHERE queue_id = ?1 ORDER BY queue_position"); if (!st) return std::unexpected(st.error()); if (auto b = st->bind(1, queue_id); !b) return std::unexpected(b.error()); for (;;) { auto row = st->step(); if (!row) return std::unexpected(row.error()); if (!*row) break; current.push_back(st->column_text(0)); } } // Exact permutation: same size, same members, order aside. std::vector a = current, b = task_ids; std::sort(a.begin(), a.end()); std::sort(b.begin(), b.end()); if (a != b) return false; auto txn = db_.transaction([&]() -> DbResult { for (std::size_t i = 0; i < task_ids.size(); ++i) { auto st = db_.prepare("UPDATE tasks SET queue_position = ?2 WHERE task_id = ?1"); if (!st) return std::unexpected(st.error()); if (auto bd = st->bind(1, std::string_view(task_ids[i])); !bd) return std::unexpected(bd.error()); if (auto bd = st->bind(2, static_cast(i)); !bd) return std::unexpected(bd.error()); if (auto r = st->step(); !r) return std::unexpected(r.error()); } return {}; }); if (!txn) return std::unexpected(txn.error()); return true; } DbResult Queues::upsert(proto::Queue queue) { if (queue.queueId.empty()) { std::random_device rd; std::uniform_int_distribution d; char buf[17]; std::snprintf(buf, sizeof(buf), "%016llx", static_cast(d(rd))); queue.queueId = std::string(buf); } // A create defaults to 'stopped' (never auto-runs a brand-new queue); a replace keeps // whatever run state the queue is already in — queue.upsert edits the config, not the // run state (that's queue.start/stop). std::string state = "stopped"; if (auto existing = get(queue.queueId); existing && existing->has_value()) state = std::string(proto::to_string((*existing)->state)); const std::string schedule_json = queue.schedule ? nlohmann::json(*queue.schedule).dump() : std::string(); const std::string on_complete = std::string(proto::to_string(queue.onComplete.value_or(proto::QueueOnComplete::Nothing))); auto st = db_.prepare( "INSERT INTO queues(queue_id, name, state, max_concurrent, schedule, on_complete) " "VALUES(?1,?2,?3,?4,?5,?6) " "ON CONFLICT(queue_id) DO UPDATE SET " "name=excluded.name, max_concurrent=excluded.max_concurrent, " "schedule=excluded.schedule, on_complete=excluded.on_complete"); if (!st) return std::unexpected(st.error()); if (auto b = st->bind(1, std::string_view(queue.queueId)); !b) return std::unexpected(b.error()); if (auto b = st->bind(2, std::string_view(queue.name)); !b) return std::unexpected(b.error()); if (auto b = st->bind(3, std::string_view(state)); !b) return std::unexpected(b.error()); if (auto b = st->bind(4, queue.maxConcurrent); !b) return std::unexpected(b.error()); if (auto r = queue.schedule ? st->bind(5, std::string_view(schedule_json)) : st->bind_null(5); !r) return std::unexpected(r.error()); if (auto b = st->bind(6, std::string_view(on_complete)); !b) return std::unexpected(b.error()); if (auto r = st->step(); !r) return std::unexpected(r.error()); auto stored = get(queue.queueId); if (!stored) return std::unexpected(stored.error()); if (!stored->has_value()) return std::unexpected(DbError{0, "queue.upsert: row vanished after insert"}); return **stored; } } // namespace velox::daemon::store