#include "conn_session.h" #ifdef VOICECAT_HAS_NET #include #include #include #include #include #include "db.h" #include "session_registry.h" #include "core/worker_pool.h" #include "protocol/envelope.h" namespace voicecat::server { static voicecat::v1::Envelope make_env(uint64_t req_id = 0) { voicecat::v1::Envelope e; e.set_request_id(req_id); return e; } static voicecat::v1::Permissions all_permissions() { voicecat::v1::Permissions p; p.set_can_create_temp_channel(true); p.set_can_kick(true); p.set_can_ban(true); p.set_can_move_users(true); p.set_can_admin_accounts(true); p.set_is_admin(true); return p; } static voicecat::v1::Permissions no_permissions() { voicecat::v1::Permissions p; return p; } ConnSession::ConnSession(std::shared_ptr db, std::shared_ptr registry, std::shared_ptr workers, const std::array& server_fp, bool allow_guests, uint16_t udp_media_port) : db_(std::move(db)), registry_(std::move(registry)), workers_(std::move(workers)), server_fp_(server_fp), allow_guests_(allow_guests), udp_media_port_(udp_media_port) { randombytes_buf(udp_token_.data(), udp_token_.size()); } void ConnSession::set_io(SendFn send_fn, CloseFn close_fn) { send_fn_ = std::move(send_fn); close_fn_ = std::move(close_fn); } void ConnSession::begin() { // Nothing to do at TCP level — wait for ClientHello } void ConnSession::on_frame(std::vector frame) { voicecat::v1::Envelope env; if (!protocol::decode_envelope(frame, env)) return; auto st = state_.load(std::memory_order_acquire); switch (env.body_case()) { case voicecat::v1::Envelope::kClientHello: if (st == State::WaitingHello) handle_client_hello(env.request_id(), env.client_hello()); break; case voicecat::v1::Envelope::kAuthRequest: if (st == State::WaitingAuth) handle_auth_request(env.request_id(), env.auth_request()); break; case voicecat::v1::Envelope::kJoinChannel: if (st == State::Authenticated) handle_join_channel(env.request_id(), env.join_channel()); break; case voicecat::v1::Envelope::kTextMessage: if (st == State::Authenticated) handle_text_message(env.text_message()); break; case voicecat::v1::Envelope::kPing: handle_ping(env.ping()); break; case voicecat::v1::Envelope::kLeaveChannel: if (st == State::Authenticated) handle_leave_channel(); break; case voicecat::v1::Envelope::kUdpBinding: if (st == State::Authenticated) handle_udp_binding(env.request_id(), env.udp_binding()); break; case voicecat::v1::Envelope::kStreamAnnounce: if (st == State::Authenticated) handle_stream_announce(env.request_id(), env.stream_announce()); break; case voicecat::v1::Envelope::kStreamStop: if (st == State::Authenticated) handle_stream_stop(env.stream_stop()); break; // ── M5 moderation / admin ───────────────────────────────────────────── case voicecat::v1::Envelope::kKick: if (st == State::Authenticated) handle_kick_request(env.request_id(), env.kick()); break; case voicecat::v1::Envelope::kBan: if (st == State::Authenticated) handle_ban_request(env.request_id(), env.ban()); break; case voicecat::v1::Envelope::kSetPermission: if (st == State::Authenticated) handle_set_permission(env.request_id(), env.set_permission()); break; case voicecat::v1::Envelope::kServerMute: if (st == State::Authenticated) handle_server_mute_request(env.request_id(), env.server_mute()); break; case voicecat::v1::Envelope::kMoveUser: if (st == State::Authenticated) handle_move_user(env.request_id(), env.move_user()); break; case voicecat::v1::Envelope::kCreateChannel: if (st == State::Authenticated) handle_create_channel(env.request_id(), env.create_channel()); break; case voicecat::v1::Envelope::kEditChannel: if (st == State::Authenticated) handle_edit_channel(env.request_id(), env.edit_channel()); break; case voicecat::v1::Envelope::kDeleteChannel: if (st == State::Authenticated) handle_delete_channel(env.request_id(), env.delete_channel()); break; case voicecat::v1::Envelope::kCreateAccount: if (st == State::Authenticated) handle_create_account(env.request_id(), env.create_account()); break; case voicecat::v1::Envelope::kResetPassword: if (st == State::Authenticated) handle_reset_password(env.request_id(), env.reset_password()); break; case voicecat::v1::Envelope::kDeleteAccount: if (st == State::Authenticated) handle_delete_account(env.request_id(), env.delete_account()); break; case voicecat::v1::Envelope::kListAccounts: if (st == State::Authenticated) handle_list_accounts(env.request_id(), env.list_accounts()); break; default: break; } } void ConnSession::on_disconnect() { close(); } void ConnSession::send_envelope(const voicecat::v1::Envelope& env) { if (!send_fn_ || closed_.load()) return; // Serialize to raw protobuf bytes; send_fn_ (→ TcpServerConn::send_frame) // adds the [4-byte len] framing, so we must NOT pre-frame here. std::string bytes; if (!env.SerializeToString(&bytes)) return; std::vector raw(bytes.begin(), bytes.end()); send_fn_(std::move(raw)); } void ConnSession::close() { if (closed_.exchange(true)) return; state_.store(State::Disconnecting, std::memory_order_release); uint32_t uid = user_id_.load(); if (uid) registry_->remove_user(uid); if (session_id_) registry_->unregister_session(session_id_); if (close_fn_) close_fn_(); } // ── M2: media crypto ───────────────────────────────────────────────────────── void ConnSession::set_media_crypto( std::unique_ptr send, std::unique_ptr recv) { std::lock_guard lk(crypto_mu_); send_crypto_ = std::move(send); recv_crypto_ = std::move(recv); } voicecat::crypto::SodiumMediaCrypto* ConnSession::send_crypto() { std::lock_guard lk(crypto_mu_); return send_crypto_.get(); } voicecat::crypto::SodiumMediaCrypto* ConnSession::recv_crypto() { std::lock_guard lk(crypto_mu_); return recv_crypto_.get(); } // ── M2: UDP endpoint ───────────────────────────────────────────────────────── void ConnSession::set_udp_endpoint(asio::ip::udp::endpoint ep) { { std::lock_guard lk(udp_ep_mu_); udp_ep_ = ep; } has_udp_ep_.store(true, std::memory_order_release); registry_->register_udp_endpoint(ep, session_id_); } asio::ip::udp::endpoint ConnSession::udp_endpoint() const { std::lock_guard lk(udp_ep_mu_); return udp_ep_; } // ── Permission helpers ─────────────────────────────────────────────────────── bool ConnSession::has_permission(bool (voicecat::v1::Permissions::* getter)() const) const { return (permissions_.*getter)(); } void ConnSession::set_permissions(const voicecat::v1::Permissions& perms) { permissions_ = perms; registry_->set_session_permissions(session_id_, perms); } void ConnSession::send_generic_result(uint64_t req_id, bool ok, uint32_t code, const std::string& message) { auto env = make_env(req_id); auto* gr = env.mutable_generic_result(); gr->set_ok(ok); gr->set_code(code); gr->set_message(message); send_envelope(env); } // ── Handlers ───────────────────────────────────────────────────────────────── void ConnSession::handle_client_hello(uint64_t req_id, const voicecat::v1::ClientHello& msg) { if (msg.proto_version() != 1) { send_disconnect_and_close(1, "unsupported protocol version"); return; } auto env = make_env(req_id); auto* hello = env.mutable_server_hello(); hello->set_proto_version(1); hello->set_server_name("VoiceCat Server"); hello->set_server_version("0.1.0"); if (allow_guests_) hello->add_auth_methods("guest"); hello->add_auth_methods("password"); hello->set_server_identity_fingerprint(server_fp_.data(), server_fp_.size()); if (udp_media_port_) hello->set_udp_port(udp_media_port_); send_envelope(env); state_.store(State::WaitingAuth, std::memory_order_release); } void ConnSession::handle_auth_request(uint64_t req_id, const voicecat::v1::AuthRequest& msg) { if (msg.has_guest()) { finish_guest_auth(msg.guest(), req_id); } else if (msg.has_password()) { finish_password_auth(msg.password().username(), msg.password().password(), req_id); } else { auto env = make_env(req_id); env.mutable_auth_result()->set_ok(false); env.mutable_auth_result()->set_error("unknown auth method"); send_envelope(env); } } void ConnSession::finish_guest_auth(const voicecat::v1::GuestAuth& guest, uint64_t req_id) { if (!allow_guests_) { auto env = make_env(req_id); env.mutable_auth_result()->set_ok(false); env.mutable_auth_result()->set_error("guest login not permitted"); send_envelope(env); return; } voicecat::v1::User user; user.set_nickname(guest.nickname().empty() ? "Guest" : guest.nickname()); user.set_is_guest(true); user.set_channel_id(1); uint32_t uid = registry_->add_user(session_id_, user); user.set_id(uid); user_id_.store(uid, std::memory_order_relaxed); state_.store(State::Authenticated, std::memory_order_release); permissions_ = no_permissions(); registry_->set_session_permissions(session_id_, permissions_); registry_->register_udp_token(udp_token_, session_id_); { auto env = make_env(req_id); auto* res = env.mutable_auth_result(); res->set_ok(true); res->set_session_id(session_id_); *res->mutable_self() = user; *res->mutable_permissions() = permissions_; res->set_udp_token(udp_token_.data(), udp_token_.size()); send_envelope(env); } broadcast_user_joined(user); send_state_snapshot(); } void ConnSession::finish_password_auth(const std::string& username, const std::string& password, uint64_t req_id) { // Argon2id runs on the worker pool (deliberately slow). auto self = shared_from_this(); workers_->post([self, username, password, req_id] { // M5: check username bans before verifying password. if (self->db_->ban_check("username", username)) { auto env = make_env(req_id); env.mutable_auth_result()->set_ok(false); env.mutable_auth_result()->set_error("account banned"); self->send_envelope(env); return; } auto acc = self->db_->authenticate(username, password); if (!acc) { auto env = make_env(req_id); env.mutable_auth_result()->set_ok(false); env.mutable_auth_result()->set_error("invalid credentials"); self->send_envelope(env); return; } voicecat::v1::User user; user.set_nickname(acc->username); user.set_is_guest(false); user.set_channel_id(1); uint32_t uid = self->registry_->add_user(self->session_id_, user); user.set_id(uid); self->user_id_.store(uid, std::memory_order_relaxed); self->state_.store(State::Authenticated, std::memory_order_release); voicecat::v1::Permissions perms = acc->is_admin ? all_permissions() : no_permissions(); self->permissions_ = perms; self->registry_->set_session_permissions(self->session_id_, perms); self->registry_->register_udp_token(self->udp_token_, self->session_id_); { auto env = make_env(req_id); auto* res = env.mutable_auth_result(); res->set_ok(true); res->set_session_id(self->session_id_); *res->mutable_self() = user; *res->mutable_permissions() = perms; res->set_udp_token(self->udp_token_.data(), self->udp_token_.size()); self->send_envelope(env); } self->broadcast_user_joined(user); self->send_state_snapshot(); }); } void ConnSession::send_state_snapshot() { auto env = make_env(); auto* snap = env.mutable_server_state(); for (auto& ch : registry_->channel_snapshot()) *snap->add_channels() = ch; for (auto& u : registry_->user_snapshot()) *snap->add_users() = u; send_envelope(env); } void ConnSession::broadcast_user_joined(const voicecat::v1::User& user) { auto bcast = make_env(); auto* ue = bcast.mutable_user_event(); ue->set_kind(voicecat::v1::UserEvent::JOINED); *ue->mutable_user() = user; registry_->broadcast(bcast, session_id_); } void ConnSession::handle_join_channel(uint64_t req_id, const voicecat::v1::JoinChannelRequest& msg) { uint32_t uid = user_id_.load(); auto ch = registry_->get_channel(msg.channel_id()); if (!ch) { auto env = make_env(req_id); auto* res = env.mutable_join_channel_result(); res->set_ok(false); res->set_error("channel not found"); send_envelope(env); return; } if (ch->password_protected()) { if (!registry_->check_channel_password(msg.channel_id(), msg.password())) { auto env = make_env(req_id); auto* res = env.mutable_join_channel_result(); res->set_ok(false); res->set_error("invalid channel password"); send_envelope(env); return; } } if (ch->max_users() > 0) { auto members = registry_->find_channel_sessions(msg.channel_id(), session_id_); if (members.size() >= ch->max_users()) { auto env = make_env(req_id); auto* res = env.mutable_join_channel_result(); res->set_ok(false); res->set_error("channel is full"); send_envelope(env); return; } } bool ok = registry_->set_user_channel(uid, msg.channel_id()); auto env = make_env(req_id); auto* res = env.mutable_join_channel_result(); res->set_ok(ok); if (!ok) { res->set_error("channel not found"); } else { res->set_channel_id(msg.channel_id()); *res->mutable_audio() = ch->audio(); // Broadcast that this user changed channel. if (auto updated_user = registry_->user_snapshot_user(uid)) { auto bcast = make_env(); auto* ue = bcast.mutable_user_event(); ue->set_kind(voicecat::v1::UserEvent::UPDATED); *ue->mutable_user() = *updated_user; registry_->broadcast(bcast, session_id_); } } send_envelope(env); } void ConnSession::handle_leave_channel() { uint32_t uid = user_id_.load(); if (!uid) return; if (!registry_->set_user_channel(uid, 1)) return; if (auto updated_user = registry_->user_snapshot_user(uid)) { auto bcast = make_env(); auto* ue = bcast.mutable_user_event(); ue->set_kind(voicecat::v1::UserEvent::UPDATED); *ue->mutable_user() = *updated_user; registry_->broadcast(bcast, session_id_); } } void ConnSession::handle_text_message(const voicecat::v1::TextMessage& msg) { using namespace std::chrono; int64_t now_ms = duration_cast( system_clock::now().time_since_epoch()).count(); voicecat::v1::TextMessage relay = msg; relay.set_sender_id(user_id_.load(std::memory_order_relaxed)); relay.set_sent_at_unix_ms(now_ms); voicecat::v1::Envelope fwd; *fwd.mutable_text_message() = relay; auto targets = registry_->resolve_text_targets(session_id_, msg.scope(), msg.target_id()); for (auto& t : targets) t->send_envelope(fwd); // Ack auto env = make_env(); auto* ack = env.mutable_text_message_ack(); ack->set_client_msg_id(msg.client_msg_id()); ack->set_ok(true); send_envelope(env); } void ConnSession::handle_ping(const voicecat::v1::Ping& msg) { auto env = make_env(); env.mutable_pong()->set_nonce(msg.nonce()); send_envelope(env); } void ConnSession::handle_udp_binding(uint64_t req_id, const voicecat::v1::UdpBinding& msg) { if (msg.ack()) return; // server→client direction; ignore if echoed back const std::string& tok = msg.udp_token(); if (tok.size() != 16 || std::memcmp(tok.data(), udp_token_.data(), 16) != 0) { // Bad token — silently ignore (don't leak timing information) return; } // Ack over TCP; MediaRelay will set the UDP endpoint when the UDP binding packet arrives. auto env = make_env(req_id); env.mutable_udp_binding()->set_ack(true); send_envelope(env); } void ConnSession::handle_stream_announce(uint64_t req_id, const voicecat::v1::StreamAnnounce& msg) { uint32_t ssrc = registry_->assign_ssrc(session_id_); uint32_t stream_id = next_stream_id_++; auto env = make_env(req_id); auto* res = env.mutable_stream_announce_result(); res->set_ok(true); res->set_stream_id(stream_id); res->set_ssrc(ssrc); // Per-channel AudioConfig is authoritative (docs/voice.md §3): the channel's mode/ // frame_ms/application/fec/dtx/complexity/expected_packet_loss apply to every stream // announced into it, regardless of kind. bitrate_bps is clamped (not overridden) to the // channel's ceiling so a client may still request less. sample_rate stays // client-requested-or-48000 — everything runs at 48kHz internally per voice.md §3. auto* eff = res->mutable_effective_audio(); auto chan_cfg = registry_->channel_audio_config(registry_->user_channel(user_id_.load())); uint32_t requested_bps = msg.has_requested_audio() ? msg.requested_audio().bitrate_bps() : 0; uint32_t requested_rate = msg.has_requested_audio() ? msg.requested_audio().sample_rate() : 0; if (chan_cfg) { *eff = *chan_cfg; eff->set_bitrate_bps(requested_bps > 0 ? std::min(requested_bps, chan_cfg->bitrate_bps()) : chan_cfg->bitrate_bps()); eff->set_sample_rate(requested_rate > 0 ? requested_rate : 48000); } else if (msg.has_requested_audio()) { *eff = msg.requested_audio(); } else { eff->set_codec(0); // OPUS eff->set_sample_rate(48000); eff->set_bitrate_bps(24000); eff->set_frame_ms(20); eff->set_fec(true); } if (eff->sample_rate() == 0) eff->set_sample_rate(48000); if (eff->bitrate_bps() == 0) eff->set_bitrate_bps(24000); if (eff->frame_ms() == 0) eff->set_frame_ms(20); announced_stream_ids_.push_back(stream_id); voicecat::v1::StreamInfo info; info.set_stream_id(stream_id); info.set_ssrc(ssrc); info.set_kind(msg.kind()); *info.mutable_audio() = *eff; info.set_label(msg.label()); send_envelope(env); auto updated = registry_->set_user_stream(user_id_.load(), info); if (updated) { auto bcast = make_env(); auto* ue = bcast.mutable_user_event(); ue->set_kind(voicecat::v1::UserEvent::UPDATED); *ue->mutable_user() = *updated; registry_->broadcast(bcast, session_id_); } } void ConnSession::handle_stream_stop(const voicecat::v1::StreamStop& msg) { auto it = std::find(announced_stream_ids_.begin(), announced_stream_ids_.end(), msg.stream_id()); if (it == announced_stream_ids_.end()) return; // not ours — ignore (no spoofed stops) announced_stream_ids_.erase(it); auto updated = registry_->clear_user_stream(user_id_.load(), msg.stream_id()); if (updated) { auto bcast = make_env(); auto* ue = bcast.mutable_user_event(); ue->set_kind(voicecat::v1::UserEvent::UPDATED); *ue->mutable_user() = *updated; registry_->broadcast(bcast, session_id_); } } // ── M5 handlers ────────────────────────────────────────────────────────────── void ConnSession::handle_kick_request(uint64_t req_id, const voicecat::v1::KickRequest& msg) { if (!is_admin() && !has_permission(&voicecat::v1::Permissions::can_kick)) { send_generic_result(req_id, false, 6, "permission denied"); return; } bool ok = registry_->kick_user(msg.user_id(), msg.reason()); send_generic_result(req_id, ok, ok ? 0 : 3, ok ? "" : "user not found"); } void ConnSession::handle_ban_request(uint64_t req_id, const voicecat::v1::BanRequest& msg) { if (!is_admin() && !has_permission(&voicecat::v1::Permissions::can_ban)) { send_generic_result(req_id, false, 6, "permission denied"); return; } // Ban by username for persistence (nickname == username for password users). // Also ban by runtime user_id for immediate effect. if (auto target = registry_->find_session_by_user_id(msg.user_id())) { if (!target->user_id()) { send_generic_result(req_id, false, 3, "user not found"); return; } if (auto nick = registry_->user_nickname(target->user_id())) { std::string err; db_->ban_create("username", *nick, msg.reason(), static_cast(msg.expires_unix_ms()), err); } } bool ok = registry_->ban_user(msg.user_id(), msg.reason(), static_cast(msg.expires_unix_ms())); send_generic_result(req_id, ok, ok ? 0 : 3, ok ? "" : "user not found"); } void ConnSession::handle_set_permission(uint64_t req_id, const voicecat::v1::SetPermissionRequest& msg) { if (!is_admin() && !has_permission(&voicecat::v1::Permissions::can_admin_accounts)) { send_generic_result(req_id, false, 6, "permission denied"); return; } auto target = registry_->find_session_by_user_id(msg.user_id()); if (!target) { send_generic_result(req_id, false, 3, "user not online"); return; } target->set_permissions(msg.permissions()); send_generic_result(req_id, true, 0, ""); } void ConnSession::handle_server_mute_request(uint64_t req_id, const voicecat::v1::ServerMuteRequest& msg) { if (!is_admin() && !has_permission(&voicecat::v1::Permissions::can_kick)) { send_generic_result(req_id, false, 6, "permission denied"); return; } bool ok = registry_->set_server_mute(msg.user_id(), msg.muted(), msg.deafened()); send_generic_result(req_id, ok, ok ? 0 : 3, ok ? "" : "user not found"); } void ConnSession::handle_move_user(uint64_t req_id, const voicecat::v1::MoveUserRequest& msg) { if (!is_admin() && !has_permission(&voicecat::v1::Permissions::can_move_users)) { send_generic_result(req_id, false, 6, "permission denied"); return; } bool ok = registry_->move_user(msg.user_id(), msg.channel_id()); send_generic_result(req_id, ok, ok ? 0 : 3, ok ? "" : "user or channel not found"); } void ConnSession::handle_create_channel(uint64_t req_id, const voicecat::v1::CreateChannelRequest& msg) { if (!is_admin() && !has_permission(&voicecat::v1::Permissions::can_create_temp_channel)) { send_generic_result(req_id, false, 6, "permission denied"); return; } std::string error; uint32_t id = registry_->create_channel(msg.channel(), msg.password(), error); send_generic_result(req_id, id != 0, id != 0 ? 0 : 3, id != 0 ? "" : (error.empty() ? "create failed" : error)); } void ConnSession::handle_edit_channel(uint64_t req_id, const voicecat::v1::EditChannelRequest& msg) { if (!is_admin()) { send_generic_result(req_id, false, 6, "permission denied"); return; } std::string error; bool ok = registry_->update_channel(msg.channel(), msg.password(), error); send_generic_result(req_id, ok, ok ? 0 : 3, ok ? "" : (error.empty() ? "update failed" : error)); } void ConnSession::handle_delete_channel(uint64_t req_id, const voicecat::v1::DeleteChannelRequest& msg) { if (!is_admin()) { send_generic_result(req_id, false, 6, "permission denied"); return; } std::string error; bool ok = registry_->delete_channel(msg.channel_id(), error); send_generic_result(req_id, ok, ok ? 0 : 3, ok ? "" : (error.empty() ? "delete failed" : error)); } void ConnSession::handle_create_account(uint64_t req_id, const voicecat::v1::CreateAccountRequest& msg) { if (!is_admin() && !has_permission(&voicecat::v1::Permissions::can_admin_accounts)) { send_generic_result(req_id, false, 6, "permission denied"); return; } auto self = shared_from_this(); workers_->post([self, req_id, msg]() mutable { std::string error; auto acc = self->db_->create_account(msg.username(), msg.password(), false, error); self->send_generic_result(req_id, acc.has_value(), acc.has_value() ? 0 : 3, acc.has_value() ? "" : error); }); } void ConnSession::handle_reset_password(uint64_t req_id, const voicecat::v1::ResetPasswordRequest& msg) { if (!is_admin() && !has_permission(&voicecat::v1::Permissions::can_admin_accounts)) { send_generic_result(req_id, false, 6, "permission denied"); return; } auto self = shared_from_this(); workers_->post([self, req_id, msg]() mutable { std::string error; bool ok = self->db_->reset_password(msg.username(), msg.new_password(), error); self->send_generic_result(req_id, ok, ok ? 0 : 3, ok ? "" : error); }); } void ConnSession::handle_delete_account(uint64_t req_id, const voicecat::v1::DeleteAccountRequest& msg) { if (!is_admin() && !has_permission(&voicecat::v1::Permissions::can_admin_accounts)) { send_generic_result(req_id, false, 6, "permission denied"); return; } std::string error; bool ok = db_->delete_account(msg.username(), error); send_generic_result(req_id, ok, ok ? 0 : 3, ok ? "" : error); } void ConnSession::handle_list_accounts(uint64_t req_id, const voicecat::v1::ListAccountsRequest& /*msg*/) { if (!is_admin() && !has_permission(&voicecat::v1::Permissions::can_admin_accounts)) { send_generic_result(req_id, false, 6, "permission denied"); return; } auto env = make_env(req_id); auto* lr = env.mutable_list_accounts_result(); for (const auto& acc : db_->list_accounts()) { auto* e = lr->add_accounts(); e->set_username(acc.username); e->set_is_admin(acc.is_admin); e->set_created_at_unix_ms(static_cast(acc.created_at) * 1000); e->set_last_login_unix_ms(static_cast(acc.last_login) * 1000); } send_envelope(env); } void ConnSession::send_disconnect_and_close(uint32_t code, const std::string& reason) { auto env = make_env(); auto* d = env.mutable_disconnect(); d->set_code(code); d->set_reason(reason); send_envelope(env); close(); } } // namespace voicecat::server #endif // VOICECAT_HAS_NET