From 77f7e80d9a7285502d578ef129012460e3f4da49 Mon Sep 17 00:00:00 2001 From: Nemoo <18210033860@163.com> Date: Sat, 30 Sep 2023 14:21:36 +0800 Subject: [PATCH] Update threadpool_sessctl.cpp --- .../process/threadpool/threadpool_sessctl.cpp | 723 ++++++++++-------- 1 file changed, 415 insertions(+), 308 deletions(-) diff --git a/src/gausskernel/process/threadpool/threadpool_sessctl.cpp b/src/gausskernel/process/threadpool/threadpool_sessctl.cpp index 1267f0ad1..cb972b2d7 100755 --- a/src/gausskernel/process/threadpool/threadpool_sessctl.cpp +++ b/src/gausskernel/process/threadpool/threadpool_sessctl.cpp @@ -59,51 +59,63 @@ #endif +// 线程池会话控制类的构造函数(参数:内存上下文) ThreadPoolSessControl::ThreadPoolSessControl(MemoryContext context) { + // 自动上下文切换到指定的内存上下文 AutoContextSwitch acontext(context); m_context = context; - pthread_mutex_init(&m_sessCtrlock, NULL); - m_sessionId = 0; - m_activeSessionCount = 0; - m_maxActiveSessionCount = GLOBAL_MAX_SESSION_NUM; - m_maxReserveSessionCount = GLOBAL_RESERVE_SESSION_NUM; - DLInitList(&m_activelist); - DLInitList(&m_freelist); + pthread_mutex_init(&m_sessCtrlock, NULL); // 初始化线程互斥锁 + m_sessionId = 0; // 会话 ID 初始化为 0 + m_activeSessionCount = 0; // 活跃会话计数器初始化为 0 + m_maxActiveSessionCount = GLOBAL_MAX_SESSION_NUM; // 最大活跃会话数初始化为全局最大会话数 + m_maxReserveSessionCount = GLOBAL_RESERVE_SESSION_NUM; // 最大保留会话数初始化为全局保留会话数 + DLInitList(&m_activelist); // 初始化活跃会话列表 + DLInitList(&m_freelist); // 初始化空闲会话列表 + // 分配存储空间并初始化基础会话控制结构 m_base = (knl_sess_control*)palloc0(sizeof(knl_sess_control) * m_maxActiveSessionCount); for (int i = 0; i < m_maxActiveSessionCount; i++) { - m_base[i].idx = (m_maxReserveSessionCount + i); - m_base[i].lock = 0; - DLInitElem(&m_base[i].elem, &m_base[i]); - DLAddHead(&m_freelist, &m_base[i].elem); + m_base[i].idx = (m_maxReserveSessionCount + i); // 设置会话索引 + m_base[i].lock = 0; // 会话锁定标志初始化为 0 + DLInitElem(&m_base[i].elem, &m_base[i]); // 初始化双向链表元素 + DLAddHead(&m_freelist, &m_base[i].elem); // 将当前会话添加到空闲会话列表头部 } } +// 线程池会话控制类的析构函数 ThreadPoolSessControl::~ThreadPoolSessControl() { - pfree_ext(m_base); - m_context = NULL; + pfree_ext(m_base); // 释放基础会话控制结构的内存 + m_context = NULL; // 将内存上下文指针置为空 } +// 创建新的线程池会话(参数:端口指针) knl_session_context* ThreadPoolSessControl::CreateSession(Port* port) { knl_session_context* sc = NULL; - /* We use u_sess->session_id to mark memory context. */ + /* 使用 u_sess->session_id 标记内存上下文。 */ m_sessionId++; + // 创建会话上下文 sc = create_session_context(m_context, m_sessionId); if (sc == NULL) { - ereport(WARNING, (errmsg("can't allocate memory for session"))); + ereport(WARNING, (errmsg("can't allocate memory for session"))); // 内存分配失败报警 return NULL; } + // 切换内存上下文到默认内存上下文 MemoryContext old_cxt = MemoryContextSwitchTo(sc->mcxt_group->GetMemCxtGroup(MEMORY_CONTEXT_DEFAULT)); + // 分配并复制端口信息到新的会话上下文 sc->proc_cxt.MyProcPort = (Port*)palloc0(sizeof(Port)); MemoryContextSwitchTo(old_cxt); + + // 复制端口信息 int rc = memcpy_s(sc->proc_cxt.MyProcPort, sizeof(Port), port, sizeof(Port)); securec_check(rc, "\0", "\0"); - ereport(DEBUG4, (errmsg("CreateSession fd:[%d] to dispatch.", port->sock))); + ereport(DEBUG4, (errmsg("CreateSession fd:[%d] to dispatch.", port->sock))); // 创建会话并分派报告日志 + + // 分配会话槽位,如果成功则返回会话,否则清理资源返回 NULL if (AllocateSlot(sc)) { sc->stat_cxt.trackedBytes += u_sess->stat_cxt.trackedBytes; sc->stat_cxt.trackedMemChunks += u_sess->stat_cxt.trackedMemChunks; @@ -113,7 +125,7 @@ knl_session_context* ThreadPoolSessControl::CreateSession(Port* port) } else { u_sess->stat_cxt.trackedBytes += sc->stat_cxt.trackedBytes; u_sess->stat_cxt.trackedMemChunks += sc->stat_cxt.trackedMemChunks; - // sc->top_mem_cxt has been sealed in create_session_context + // sc->top_mem_cxt 在 create_session_context 中已经被封存 MemoryContextUnSeal(sc->top_mem_cxt); MemoryContextDeleteChildren(sc->top_mem_cxt); MemoryContextDelete(sc->top_mem_cxt); @@ -122,117 +134,121 @@ knl_session_context* ThreadPoolSessControl::CreateSession(Port* port) } } +// 分配一个会话槽位并返回会话控制结构指针(参数:会话上下文) knl_sess_control* ThreadPoolSessControl::AllocateSlot(knl_session_context* sc) { - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); + AutoMutexLock alock(&m_sessCtrlock); // 自动锁定互斥锁 + alock.lock(); // 锁定互斥锁 + + // 检查是否可以接受新连接 if (canAcceptConnections(true) != CAC_OK) { - alock.unLock(); - /* CN is in the process of starting recovery, redo is not completed, so the connection error is reasonable */ - ereport(WARNING, - (errmsg("ThreadPool cannot start new session due to PM state, PM state is %s", GetPMState(pmState)))); + alock.unLock(); // 解锁互斥锁 + /* 数据库正在执行恢复过程,此时连接错误是合理的 */ + ereport(WARNING, (errmsg("ThreadPool cannot start new session due to PM state, PM state is %s", GetPMState(pmState)))); return NULL; } + // 检查活跃会话数是否已达上限 if (m_activeSessionCount == m_maxActiveSessionCount) { - alock.unLock(); - ereport(WARNING, - (errmsg("ThreadPool cannot start new session due to to many sessions, current upper bound is %d", - m_maxActiveSessionCount))); + alock.unLock(); // 解锁互斥锁 + ereport(WARNING, (errmsg("ThreadPool cannot start new session due to to many sessions, current upper bound is %d", + m_maxActiveSessionCount))); // 活跃会话数已达上限报警 return NULL; } - Assert(DLListLength(&m_freelist) != 0); + Assert(DLListLength(&m_freelist) != 0); // 确保空闲会话列表不为空 - /* remove from free list */ + // 从空闲会话列表中移除一个会话控制结构 knl_sess_control* ctrl = (knl_sess_control*)DLE_VAL(DLRemHead(&m_freelist)); ctrl->sess = sc; - sc->session_ctr_index = ctrl->idx; - DLAddHead(&m_activelist, &ctrl->elem); - m_activeSessionCount++; - alock.unLock(); + sc->session_ctr_index = ctrl->idx; // 设置会话槽位索引 + DLAddHead(&m_activelist, &ctrl->elem); // 将会话控制结构添加到活跃会话列表头部 + m_activeSessionCount++; // 活跃会话计数器加一 + alock.unLock(); // 解锁互斥锁 return ctrl; } +// 释放一个会话槽位(参数:会话槽位索引) void ThreadPoolSessControl::FreeSlot(int ctrl_index) { if (!IsValidCtrlIndex(ctrl_index)) { return; } - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); + AutoMutexLock alock(&m_sessCtrlock); // 自动锁定互斥锁 + alock.lock(); // 锁定互斥锁 - knl_sess_control *ctrl = &m_base[ctrl_index - m_maxReserveSessionCount]; - Assert(ctrl->elem.dle_list == &m_activelist); + knl_sess_control *ctrl = &m_base[ctrl_index - m_maxReserveSessionCount]; // 获取会话控制结构 + Assert(ctrl->elem.dle_list == &m_activelist); // 确保会话控制结构在活跃会话列表中 - m_activeSessionCount--; + m_activeSessionCount--; // 活跃会话计数器减一 volatile sig_atomic_t* plock = &ctrl->lock; sig_atomic_t val; do { if (*plock == 0) { - /* perform an atomic compare and swap. */ + /* 进行原子比较和交换操作 */ val = __sync_val_compare_and_swap(plock, 0, 1); if (val == 0) { ctrl->sess = NULL; - /* restore the value. */ + /* 恢复值 */ ctrl->lock = 0; - DLRemove(&ctrl->elem); - DLAddHead(&m_freelist, &ctrl->elem); + DLRemove(&ctrl->elem); // 从活跃会话列表中移除 + DLAddHead(&m_freelist, &ctrl->elem); // 将会话控制结构添加到空闲会话列表头部 break; } } - pg_usleep(100); + pg_usleep(100); // 休眠 100 微秒 } while (true); - alock.unLock(); + alock.unLock(); // 解锁互斥锁 } +// 标记所有会话为关闭状态 void ThreadPoolSessControl::MarkAllSessionClose() { - /* Mark all session to be closed. */ - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); + /* 标记所有会话为关闭状态。*/ + AutoMutexLock alock(&m_sessCtrlock); // 自动锁定互斥锁 + alock.lock(); // 锁定互斥锁 knl_sess_control* ctrl = NULL; - Dlelem* elem = DLGetHead(&m_activelist); + Dlelem* elem = DLGetHead(&m_activelist); // 获取活跃会话列表的头部 while (elem != NULL) { - ctrl= (knl_sess_control*)DLE_VAL(elem); - ctrl->sess->status = KNL_SESS_CLOSE; - CloseClientSocket(ctrl->sess, false); - elem = DLGetSucc(elem); + ctrl = (knl_sess_control*)DLE_VAL(elem); + ctrl->sess->status = KNL_SESS_CLOSE; // 设置会话状态为关闭 + CloseClientSocket(ctrl->sess, false); // 关闭客户端套接字 + elem = DLGetSucc(elem); // 获取下一个元素 } - alock.unLock(); + alock.unLock(); // 解锁互斥锁 ereport(LOG, (errmodule(MOD_THREAD_POOL), - errmsg("pmState:%d, mark all threadpool sessions closed.", pmState))); + errmsg("pmState:%d, mark all threadpool sessions closed.", pmState))); // 记录日志 } +// 检查发送信号的权限(参数:会话上下文,锁指针) void ThreadPoolSessControl::CheckPermissionForSendSignal(knl_session_context* sess, sig_atomic_t* lock) { - /* User id is invalid only when sometimes dealing with cancel signal. Because that permission is ensured - by random cancel key, so we don't have to check the permission again. */ + /* 当处理取消信号时,用户 ID 无效。因为取消信号的权限是通过随机取消密钥保证的,所以不需要再次检查权限。*/ if (!OidIsValid(u_sess->misc_cxt.CurrentUserId)) { return; } - /* Only users with sysadmin privilege or the member of gs_role_signal_backend role - * or the owner of the database or user himself have the permission to send singal. */ + /* 只有具有 sysadmin 特权、gs_role_signal_backend 角色的成员、数据库的所有者或用户本身才具有发送信号的权限。 */ bool role_signal_backend_permission = is_member_of_role(GetUserId(), DEFAULT_ROLE_SIGNAL_BACKENDID) && (sess->proc_cxt.MyRoleId != BOOTSTRAP_SUPERUSERID && !is_role_persistence(sess->proc_cxt.MyRoleId)); if (!superuser() && !pg_database_ownercheck(sess->proc_cxt.MyDatabaseId, u_sess->misc_cxt.CurrentUserId) && !role_signal_backend_permission) { if (sess->proc_cxt.MyRoleId != GetUserId()) { - *lock = 0; + *lock = 0; // 解锁 ereport(ERROR, (errcode(ERRCODE_INSUFFICIENT_PRIVILEGE), (errmsg("must have sysadmin privilege or a member of the gs_role_signal_backend role or the " - "owner of the database or the same user to terminate other backend")))); + "owner of the database or the same user to terminate other backend")))); // 发送错误报告 } } } +// 如果必要,释放锁(通常用于处理异常情况) void ThreadPoolSessControl::releaseLockIfNecessary() { if (unlikely(t_thrd.sig_cxt.cur_ctrl_index != 0)) { - knl_sess_control* ctrl = &m_base[t_thrd.sig_cxt.cur_ctrl_index - m_maxReserveSessionCount]; + knl_sess_control* ctrl = &m_base[t_thrd.sig_cxt.cur_ctrl_index - m_maxReserveSessionCount]; // 获取会话控制结构 volatile sig_atomic_t plock = ctrl->lock; if (plock != 0) { plock = 0; @@ -241,468 +257,549 @@ void ThreadPoolSessControl::releaseLockIfNecessary() } } +// 向指定会话发送信号(参数:会话槽位索引,信号编号) int ThreadPoolSessControl::SendSignal(int ctrl_index, int signal) { - Assert(signal != SIGHUP); - int status = ESRCH; + Assert(signal != SIGHUP); // 断言:信号不能为 SIGHUP + int status = ESRCH; // 初始状态为 ESRCH(没有这样的进程) + + // 检查会话槽位索引是否有效 if (!IsValidCtrlIndex(ctrl_index)) { - return ESRCH; + return ESRCH; // 返回 ESRCH(没有这样的进程) } + // 获取指定会话的会话控制结构 knl_sess_control* ctrl = &m_base[ctrl_index - m_maxReserveSessionCount]; - t_thrd.sig_cxt.cur_ctrl_index = ctrl_index; - volatile sig_atomic_t* plock = &ctrl->lock; + t_thrd.sig_cxt.cur_ctrl_index = ctrl_index; // 设置当前控制结构的索引 + volatile sig_atomic_t* plock = &ctrl->lock; // 获取控制结构中的锁 sig_atomic_t val; + do { if (*plock == 0) { - /* perform an atomic compare and swap. */ + /* 执行原子比较和交换操作 */ val = __sync_val_compare_and_swap(plock, 0, 1); if (val == 0) { - knl_session_context* sess = ctrl->sess; - /* Session may be NULL when the session exits during the clean connection process. - We do nothing if the session is NULL */ + knl_session_context* sess = ctrl->sess; // 获取会话上下文 + /* 当会话在清理连接过程中退出时,会话可能为 NULL。此时我们不执行任何操作。 */ if (sess == NULL) { - /* restore the value. */ + /* 恢复值。 */ ctrl->lock = 0; - status = ESRCH; + status = ESRCH; // 设置状态为 ESRCH(没有这样的进程) break; } - /* Check user permission, and we dont have user id for cancel request. */ + + /* 检查用户权限,对于取消请求,我们没有用户 ID。 */ CheckPermissionForSendSignal(sess, (sig_atomic_t*)plock); + + // 如果会话状态为已连接(KNL_SESS_ATTACH) if (sess->status == KNL_SESS_ATTACH) { - t_thrd.sig_cxt.gs_sigale_check_type = SIGNAL_CHECK_SESS_KEY; - t_thrd.sig_cxt.session_id = sess->session_id; - status = gs_signal_send(sess->attachPid, signal); - t_thrd.sig_cxt.gs_sigale_check_type = SIGNAL_CHECK_NONE; - t_thrd.sig_cxt.session_id = 0; - } else if (sess->status == KNL_SESS_DETACH) { + t_thrd.sig_cxt.gs_sigale_check_type = SIGNAL_CHECK_SESS_KEY; // 设置信号检查类型 + t_thrd.sig_cxt.session_id = sess->session_id; // 设置会话 ID + status = gs_signal_send(sess->attachPid, signal); // 向进程发送信号 + t_thrd.sig_cxt.gs_sigale_check_type = SIGNAL_CHECK_NONE; // 恢复信号检查类型 + t_thrd.sig_cxt.session_id = 0; // 恢复会话 ID + } + // 如果会话状态为已分离(KNL_SESS_DETACH) + else if (sess->status == KNL_SESS_DETACH) { switch (signal) { - case SIGTERM: - sess->status = KNL_SESS_CLOSE; - CloseClientSocket(sess, false); - status = 0; + case SIGTERM: // 如果信号是 SIGTERM(终止信号) + sess->status = KNL_SESS_CLOSE; // 设置会话状态为关闭 + CloseClientSocket(sess, false); // 关闭客户端套接字 + status = 0; // 设置状态为 0(成功) break; default: break; } } else { - status = ESRCH; + status = ESRCH; // 设置状态为 ESRCH(没有这样的进程) } - /* restore the value. */ + + /* 恢复值。 */ ctrl->lock = 0; break; } } - pg_usleep(100); + pg_usleep(100); // 休眠 100 微秒 } while (true); - t_thrd.sig_cxt.cur_ctrl_index = 0; - return status; + t_thrd.sig_cxt.cur_ctrl_index = 0; // 恢复当前控制结构的索引为 0 + + return status; // 返回状态 } -void ThreadPoolSessControl::SendProcSignal(int ctrl_index, ProcSignalReason reason, uint64 query_id) +// 向指定的会话发送信号(参数:会话槽位索引,信号编号) +int ThreadPoolSessControl::SendSignal(int ctrl_index, int signal) { + Assert(signal != SIGHUP); // 断言:信号不能为 SIGHUP + + int status = ESRCH; // 初始状态为 ESRCH(没有这样的进程) + + // 检查会话槽位索引是否有效 if (!IsValidCtrlIndex(ctrl_index)) { - return; + return ESRCH; // 返回 ESRCH(没有这样的进程) } + // 获取指定会话的会话控制结构 knl_sess_control* ctrl = &m_base[ctrl_index - m_maxReserveSessionCount]; - - volatile sig_atomic_t* plock = &ctrl->lock; + t_thrd.sig_cxt.cur_ctrl_index = ctrl_index; // 设置当前控制结构的索引 + volatile sig_atomic_t* plock = &ctrl->lock; // 获取控制结构中的锁 sig_atomic_t val; + do { if (*plock == 0) { - /* perform an atomic compare and swap. */ + /* 执行原子比较和交换操作 */ val = __sync_val_compare_and_swap(plock, 0, 1); if (val == 0) { - if (ctrl->sess != NULL) { - switch (reason) { - case PROCSIG_EXECUTOR_FLAG: { - if (IS_PGXC_DATANODE && ctrl->sess->debug_query_id == query_id) { - ctrl->sess->exec_cxt.executorStopFlag = true; - } - break; - } - default: { - Assert(0); - ctrl->lock = 0; - ereport(ERROR, - (errcode(ERRCODE_CONNECTION_EXCEPTION), errmsg("Unexpected receive proc signal."))); - } - } + knl_session_context* sess = ctrl->sess; // 获取会话上下文 + + /* 当会话在清理连接过程中退出时,会话可能为 NULL。此时我们不执行任何操作。 */ + if (sess == NULL) { + /* 恢复锁的值。 */ + ctrl->lock = 0; + status = ESRCH; // 设置状态为 ESRCH(没有这样的进程) + break; } - /* restore the value. */ + + /* 检查用户权限。对于取消请求,我们没有用户 ID。 */ + CheckPermissionForSendSignal(sess, (sig_atomic_t*)plock); + + // 如果会话状态为已连接(KNL_SESS_ATTACH) + if (sess->status == KNL_SESS_ATTACH) { + t_thrd.sig_cxt.gs_sigale_check_type = SIGNAL_CHECK_SESS_KEY; // 设置信号检查类型 + t_thrd.sig_cxt.session_id = sess->session_id; // 设置会话 ID + status = gs_signal_send(sess->attachPid, signal); // 向进程发送信号 + t_thrd.sig_cxt.gs_sigale_check_type = SIGNAL_CHECK_NONE; // 恢复信号检查类型 + t_thrd.sig_cxt.session_id = 0; // 恢复会话 ID + } + // 如果会话状态为已分离(KNL_SESS_DETACH) + else if (sess->status == KNL_SESS_DETACH) { + switch (signal) { + case SIGTERM: // 如果信号是 SIGTERM(终止信号) + sess->status = KNL_SESS_CLOSE; // 设置会话状态为关闭 + CloseClientSocket(sess, false); // 关闭客户端套接字 + status = 0; // 设置状态为 0(成功) + break; + default: + break; + } + } else { + status = ESRCH; // 设置状态为 ESRCH(没有这样的进程) + } + + /* 恢复锁的值。 */ ctrl->lock = 0; break; } } - pg_usleep(100); + pg_usleep(100); // 休眠 100 微秒 } while (true); + + t_thrd.sig_cxt.cur_ctrl_index = 0; // 恢复当前控制结构的索引为 0 + + return status; // 返回状态 } +// 计算指定数据库ID的会话数量 int ThreadPoolSessControl::CountDBSessions(Oid dbId) { - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + alock.lock(); // 获取锁 - knl_sess_control* ctrl = NULL; - Dlelem* elem = DLGetHead(&m_activelist); - int count = 0; + knl_sess_control* ctrl = NULL; // 会话控制结构指针 + Dlelem* elem = DLGetHead(&m_activelist); // 获取活动会话链表的头指针 + int count = 0; // 会话数量计数器 + // 遍历活动会话链表 while (elem != NULL) { - ctrl= (knl_sess_control*)DLE_VAL(elem); + ctrl = (knl_sess_control*)DLE_VAL(elem); // 获取会话控制结构指针 + // 如果会话的数据库ID匹配指定的数据库ID if (ctrl->sess->proc_cxt.MyDatabaseId == dbId) { - count++; + count++; // 增加会话数量计数器 + // 如果会话的应用程序名称是 "WDRXdb" if (strcmp(ctrl->sess->attr.attr_common.application_name, "WDRXdb") == 0) { - ThreadPoolSessControl::SendSignal(ctrl->idx, SIGTERM); - ThreadPoolSessControl::SendSignal(ctrl->idx, SIGUSR2); + ThreadPoolSessControl::SendSignal(ctrl->idx, SIGTERM); // 发送SIGTERM信号给会话 + ThreadPoolSessControl::SendSignal(ctrl->idx, SIGUSR2); // 发送SIGUSR2信号给会话 } } - - elem = DLGetSucc(elem); + elem = DLGetSucc(elem); // 获取下一个会话元素 } - alock.unLock(); + alock.unLock(); // 释放锁 - return count; + return count; // 返回会话数量 } +// 验证指定的数据库OID和用户OID是否有效 bool ThreadPoolSessControl::ValidDBoidAndUseroid(Oid dbOid, Oid userOid, knl_sess_control* ctrl) { /* - * Thread are 3 situation in CLEAN CONNECTION: - * 1. Only database, for example: CLEAN CONNECTION TO ALL FORCE FOR DATABASE xxx; - * 2. Only user, for example: CLEAN CONNECTION TO ALL FORCE TO USER xxx; - * 3. Both database and user, for example: CLEAN CONNECTION TO ALL FORCE FOR DATABASE xxx TO USER xxx; + * 在清理连接过程中,线程有3种情况: + * 1. 只有数据库,例如:CLEAN CONNECTION TO ALL FORCE FOR DATABASE xxx; + * 2. 只有用户,例如:CLEAN CONNECTION TO ALL FORCE TO USER xxx; + * 3. 既有数据库又有用户,例如:CLEAN CONNECTION TO ALL FORCE FOR DATABASE xxx TO USER xxx; */ if (((dbOid != InvalidOid) && (userOid == InvalidOid) && (ctrl->sess->proc_cxt.MyDatabaseId == dbOid)) || ((dbOid == InvalidOid) && (userOid != InvalidOid) && (ctrl->sess->proc_cxt.MyRoleId == userOid)) || ((ctrl->sess->proc_cxt.MyDatabaseId == dbOid) && (ctrl->sess->proc_cxt.MyRoleId == userOid))) { - return true; + return true; // 验证通过 } - return false; + return false; // 验证未通过 } +// 计算未清理的指定数据库OID和用户OID的会话数量 int ThreadPoolSessControl::CountDBSessionsNotCleaned(Oid dbOid, Oid userOid) { + // 如果数据库OID和用户OID都是无效的(可能为NULL) if ((dbOid == InvalidOid) && (userOid == InvalidOid)) { ereport(WARNING, (errmsg("DB oid and user oid are all Invalid (may be NULL). Shut down clean activite sessions."))); - return 0; + return 0; // 返回0,表示没有会话被清理 } - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + alock.lock(); // 获取锁 - int count = 0; - knl_sess_control* ctrl = NULL; - Dlelem* elem = DLGetHead(&m_activelist); + int count = 0; // 会话数量计数器 + knl_sess_control* ctrl = NULL; // 会话控制结构指针 + Dlelem* elem = DLGetHead(&m_activelist); // 获取活动会话链表的头指针 + // 遍历活动会话链表 while (elem != NULL) { - ctrl= (knl_sess_control*)DLE_VAL(elem); + ctrl = (knl_sess_control*)DLE_VAL(elem); // 获取会话控制结构指针 + // 如果数据库OID和用户OID与会话的数据库ID和用户ID匹配 if (ValidDBoidAndUseroid(dbOid, userOid, ctrl)) { - int status = ThreadPoolSessControl::SendSignal(ctrl->idx, 0); + int status = ThreadPoolSessControl::SendSignal(ctrl->idx, 0); // 向会话发送信号(0表示查询是否已经终止) + // 如果终止操作尚未完成 if (status == 0) { - /* Termination not done yet */ - count++; + count++; // 增加会话数量计数器 } } - elem = DLGetSucc(elem); + elem = DLGetSucc(elem); // 获取下一个会话元素 } - alock.unLock(); - return count; + alock.unLock(); // 释放锁 + + return count; // 返回会话数量 } +// 清理指定数据库OID和用户OID的会话 int ThreadPoolSessControl::CleanDBSessions(Oid dbOid, Oid userOid) { + // 如果数据库OID和用户OID都是无效的(可能为NULL) if ((dbOid == InvalidOid) && (userOid == InvalidOid)) { ereport(WARNING, (errmsg("DB oid and user oid are all Invalid (may be NULL). Shut down clean activite sessions."))); - return 0; + return 0; // 返回0,表示没有会话被清理 } - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + alock.lock(); // 获取锁 - int count = 0; - knl_sess_control* ctrl = NULL; - Dlelem* elem = DLGetHead(&m_activelist); + int count = 0; // 会话数量计数器 + knl_sess_control* ctrl = NULL; // 会话控制结构指针 + Dlelem* elem = DLGetHead(&m_activelist); // 获取活动会话链表的头指针 + // 遍历活动会话链表 while (elem != NULL) { - ctrl= (knl_sess_control*)DLE_VAL(elem); + ctrl = (knl_sess_control*)DLE_VAL(elem); // 获取会话控制结构指针 + // 如果数据库OID和用户OID与会话的数据库ID和用户ID匹配 if (ValidDBoidAndUseroid(dbOid, userOid, ctrl)) { - count++; - ThreadPoolSessControl::SendSignal(ctrl->idx, SIGTERM); - ThreadPoolSessControl::SendSignal(ctrl->idx, SIGUSR2); + count++; // 增加会话数量计数器 + ThreadPoolSessControl::SendSignal(ctrl->idx, SIGTERM); // 向会话发送SIGTERM信号 + ThreadPoolSessControl::SendSignal(ctrl->idx, SIGUSR2); // 向会话发送SIGUSR2信号 } - elem = DLGetSucc(elem); + elem = DLGetSucc(elem); // 获取下一个会话元素 } - alock.unLock(); - return count; + alock.unLock(); // 释放锁 + + return count; // 返回会话数量 } +// 处理SIGHUP信号 void ThreadPoolSessControl::SigHupHandler() { - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + alock.lock(); // 获取锁 - knl_sess_control* ctrl = NULL; - Dlelem* elem = DLGetHead(&m_activelist); + knl_sess_control* ctrl = NULL; // 会话控制结构指针 + Dlelem* elem = DLGetHead(&m_activelist); // 获取活动会话链表的头指针 + + // 遍历活动会话链表 while (elem != NULL) { - ctrl= (knl_sess_control*)DLE_VAL(elem); - ctrl->sess->sig_cxt.got_SIGHUP = true; - elem = DLGetSucc(elem); + ctrl = (knl_sess_control*)DLE_VAL(elem); // 获取会话控制结构指针 + ctrl->sess->sig_cxt.got_SIGHUP = true; // 标记接收到SIGHUP信号 + elem = DLGetSucc(elem); // 获取下一个会话元素 } - alock.unLock(); + alock.unLock(); // 释放锁 } +// 处理连接池重新加载事件 void ThreadPoolSessControl::HandlePoolerReload() { + // 如果是PGXC数据节点,直接返回,不处理连接池重新加载事件 if (IS_PGXC_DATANODE) { return; } - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + alock.lock(); // 获取锁 - knl_sess_control* ctrl = NULL; - Dlelem* elem = DLGetHead(&m_activelist); + knl_sess_control* ctrl = NULL; // 会话控制结构指针 + Dlelem* elem = DLGetHead(&m_activelist); // 获取活动会话链表的头指针 + + // 遍历活动会话链表 while (elem != NULL) { - ctrl= (knl_sess_control*)DLE_VAL(elem); - /* we have already send got_pool_reload to threads */ - ctrl->sess->sig_cxt.got_pool_reload = true; - ctrl->sess->sig_cxt.cp_PoolReload = true; - elem = DLGetSucc(elem); + ctrl = (knl_sess_control*)DLE_VAL(elem); // 获取会话控制结构指针 + /* 已经向线程发送了got_pool_reload信号 */ + ctrl->sess->sig_cxt.got_pool_reload = true; // 标记接收到连接池重新加载事件 + ctrl->sess->sig_cxt.cp_PoolReload = true; // 标记连接池重新加载事件 + elem = DLGetSucc(elem); // 获取下一个会话元素 } - alock.unLock(); + alock.unLock(); // 释放锁 } +// 计算会话的内存上下文统计信息 void ThreadPoolSessControl::calculateSessMemCxtStats( knl_session_context* sess, const MemoryContext context, Tuplestorestate* tupStore, TupleDesc tupDesc) { - AllocSetContext* set = (AllocSetContext*)context; + AllocSetContext* set = (AllocSetContext*)context; // 获取内存上下文的分配集上下文 - char sessId[SESSION_ID_LEN] = {0}; - ThreadId threadId = 0; - pg_time_t sessStartTime = 0; - uint64 sessionId = 0; + char sessId[SESSION_ID_LEN] = {0}; // 会话ID + ThreadId threadId = 0; // 线程ID + pg_time_t sessStartTime = 0; // 会话开始时间 + uint64 sessionId = 0; // 会话ID if (sess != NULL) { - sessionId = sess->session_id; - threadId = sess->attachPid; - sessStartTime = timestamptz_to_time_t(sess->proc_cxt.MyProcPort->SessionStartTime); + sessionId = sess->session_id; // 获取会话ID + threadId = sess->attachPid; // 获取线程ID + sessStartTime = timestamptz_to_time_t(sess->proc_cxt.MyProcPort->SessionStartTime); // 获取会话开始时间 } - getSessionID(sessId, sessStartTime, sessionId); + getSessionID(sessId, sessStartTime, sessionId); // 生成会话ID - /* build one tuple and save it in tuplestore. */ - Datum values[NUM_SESSION_MEMORY_DETAIL_ELEM] = {0}; - bool nulls[NUM_SESSION_MEMORY_DETAIL_ELEM] = {false}; + /* 构建一个元组并将其保存在元组存储器中。*/ + Datum values[NUM_SESSION_MEMORY_DETAIL_ELEM] = {0}; // 元组数据 + bool nulls[NUM_SESSION_MEMORY_DETAIL_ELEM] = {false}; // 元组的空标志位 - values[0] = CStringGetTextDatum(sessId); - values[1] = Int64GetDatum(threadId); + values[0] = CStringGetTextDatum(sessId); // 会话ID + values[1] = Int64GetDatum(threadId); // 线程ID - values[2] = CStringGetTextDatum(pstrdup(context->name)); - values[3] = Int16GetDatum(context->level); + values[2] = CStringGetTextDatum(pstrdup(context->name)); // 内存上下文名称 + values[3] = Int16GetDatum(context->level); // 内存上下文层级 if (context->level > 0 && context->parent != NULL) - values[4] = CStringGetTextDatum(pstrdup(context->parent->name)); + values[4] = CStringGetTextDatum(pstrdup(context->parent->name)); // 父级内存上下文名称 else - nulls[4] = true; - values[5] = Int64GetDatum(set->totalSpace); - values[6] = Int64GetDatum(set->freeSpace); - values[7] = Int64GetDatum(set->totalSpace - set->freeSpace); + nulls[4] = true; // 如果没有父级内存上下文,将空标志位设置为true + values[5] = Int64GetDatum(set->totalSpace); // 总共分配的内存空间 + values[6] = Int64GetDatum(set->freeSpace); // 剩余的内存空间 + values[7] = Int64GetDatum(set->totalSpace - set->freeSpace); // 已使用的内存空间 - tuplestore_putvalues(tupStore, tupDesc, values, nulls); + tuplestore_putvalues(tupStore, tupDesc, values, nulls); // 将元组数据插入到元组存储器中 } +// 递归计算会话的内存上下文统计信息 void ThreadPoolSessControl::recursiveSessMemCxt( knl_session_context* sess, const MemoryContext context, Tuplestorestate* tupStore, TupleDesc tupDesc) { - /* calculate MemoryContext Stats */ + /* 计算内存上下文统计信息 */ calculateSessMemCxtStats(sess, context, tupStore, tupDesc); - /* recursive MemoryContext's child */ + /* 递归处理内存上下文的子级 */ for (MemoryContext child = context->firstchild; child != NULL; child = child->nextchild) { recursiveSessMemCxt(sess, child, tupStore, tupDesc); } } +// 获取会话的内存详细信息 void ThreadPoolSessControl::getSessionMemoryDetail(Tuplestorestate* tupStore, TupleDesc tupDesc, knl_sess_control** sess) { - AutoMutexLock alock(&m_sessCtrlock); - knl_sess_control* ctrl = NULL; - Dlelem* elem = NULL; + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + knl_sess_control* ctrl = NULL; // 会话控制结构指针 + Dlelem* elem = NULL; // 双向链表元素指针 PG_TRY(); { - HOLD_INTERRUPTS(); - alock.lock(); + HOLD_INTERRUPTS(); // 暂时屏蔽中断 + alock.lock(); // 获取锁 - /* collect all the Memory Context status, put in data */ - elem = DLGetHead(&m_activelist); + /* 收集所有内存上下文的状态信息,将其放入数据中 */ + elem = DLGetHead(&m_activelist); // 获取活动会话链表的头指针 while (elem != NULL) { - ctrl = (knl_sess_control*)DLE_VAL(elem); - *sess = ctrl; + ctrl = (knl_sess_control*)DLE_VAL(elem); // 获取会话控制结构指针 + *sess = ctrl; // 更新会话控制结构指针 if (ctrl->sess) { - (void)syscalllockAcquire(&ctrl->sess->utils_cxt.deleMemContextMutex); - recursiveSessMemCxt(ctrl->sess, ctrl->sess->top_mem_cxt, tupStore, tupDesc); - (void)syscalllockRelease(&ctrl->sess->utils_cxt.deleMemContextMutex); + (void)syscalllockAcquire(&ctrl->sess->utils_cxt.deleMemContextMutex); // 获取内存上下文互斥锁 + recursiveSessMemCxt(ctrl->sess, ctrl->sess->top_mem_cxt, tupStore, tupDesc); // 递归计算内存上下文统计信息 + (void)syscalllockRelease(&ctrl->sess->utils_cxt.deleMemContextMutex); // 释放内存上下文互斥锁 } - elem = DLGetSucc(elem); + elem = DLGetSucc(elem); // 获取下一个会话元素 } - alock.unLock(); + alock.unLock(); // 释放锁 - RESUME_INTERRUPTS(); + RESUME_INTERRUPTS(); // 恢复中断 } PG_CATCH(); { if (*sess != NULL) { ctrl = *sess; - (void)syscalllockRelease(&ctrl->sess->utils_cxt.deleMemContextMutex); + (void)syscalllockRelease(&ctrl->sess->utils_cxt.deleMemContextMutex); // 释放内存上下文互斥锁 } - alock.unLock(); - PG_RE_THROW(); + alock.unLock(); // 释放锁 + PG_RE_THROW(); // 重新抛出异常 } - PG_END_TRY(); + PG_END_TRY(); // 结束异常处理 + } +// 计算客户端信息 void ThreadPoolSessControl::calculateClientInfo( knl_session_context* sess, Tuplestorestate* tupStore, TupleDesc tupDesc) { - /* build one tuple and save it in tuplestore. */ - const int COLUMN_NUM = 2; - Datum values[COLUMN_NUM] = {0}; - bool nulls[COLUMN_NUM] = {false}; - - values[0] = Int64GetDatum(sess->session_id); + /* 构建一个元组并将其保存在元组存储器中。*/ + const int COLUMN_NUM = 2; // 元组的列数 + Datum values[COLUMN_NUM] = {0}; // 元组数据 + bool nulls[COLUMN_NUM] = {false}; // 元组的空标志位 + + values[0] = Int64GetDatum(sess->session_id); // 会话ID if (sess->plsql_cxt.client_info != NULL) { - values[1] = CStringGetTextDatum(sess->plsql_cxt.client_info); + values[1] = CStringGetTextDatum(sess->plsql_cxt.client_info); // 客户端信息 } else { - nulls[1] = true; + nulls[1] = true; // 如果客户端信息为空,将空标志位设置为true } - tuplestore_putvalues(tupStore, tupDesc, values, nulls); + tuplestore_putvalues(tupStore, tupDesc, values, nulls); // 将元组数据插入到元组存储器中 } +// 获取会话的客户端信息 void ThreadPoolSessControl::getSessionClientInfo(Tuplestorestate* tupStore, TupleDesc tupDesc) { - AutoMutexLock alock(&m_sessCtrlock); - knl_sess_control* ctrl = NULL; - Dlelem* elem = NULL; - knl_sess_control* sess = NULL; - + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + knl_sess_control* ctrl = NULL; // 会话控制结构指针 + Dlelem* elem = NULL; // 双向链表元素指针 + knl_sess_control* sess = NULL; // 会话控制结构指针 + PG_TRY(); { - HOLD_INTERRUPTS(); - alock.lock(); - - /* collect all the Memory Context status, put in data */ - elem = DLGetHead(&m_activelist); - + HOLD_INTERRUPTS(); // 暂时屏蔽中断 + alock.lock(); // 获取锁 + + /* 收集所有内存上下文的状态信息,将其放入数据中 */ + elem = DLGetHead(&m_activelist); // 获取活动会话链表的头指针 + while (elem != NULL) { - ctrl = (knl_sess_control*)DLE_VAL(elem); - sess = ctrl; + ctrl = (knl_sess_control*)DLE_VAL(elem); // 获取会话控制结构指针 + sess = ctrl; // 更新会话控制结构指针 if (ctrl->sess) { - (void)syscalllockAcquire(&ctrl->sess->plsql_cxt.client_info_lock); - calculateClientInfo(ctrl->sess, tupStore, tupDesc); - (void)syscalllockRelease(&ctrl->sess->plsql_cxt.client_info_lock); + (void)syscalllockAcquire(&ctrl->sess->plsql_cxt.client_info_lock); // 获取客户端信息互斥锁 + calculateClientInfo(ctrl->sess, tupStore, tupDesc); // 计算客户端信息并将其保存在元组存储器中 + (void)syscalllockRelease(&ctrl->sess->plsql_cxt.client_info_lock); // 释放客户端信息互斥锁 } - elem = DLGetSucc(elem); + elem = DLGetSucc(elem); // 获取下一个会话元素 } - alock.unLock(); - sess = NULL; - - RESUME_INTERRUPTS(); + alock.unLock(); // 释放锁 + sess = NULL; // 重置会话控制结构指针 + + RESUME_INTERRUPTS(); // 恢复中断 } PG_CATCH(); { if (sess != NULL) { ctrl = sess; - (void)syscalllockRelease(&ctrl->sess->plsql_cxt.client_info_lock); + (void)syscalllockRelease(&ctrl->sess->plsql_cxt.client_info_lock); // 释放客户端信息互斥锁 } - alock.unLock(); - PG_RE_THROW(); + alock.unLock(); // 释放锁 + PG_RE_THROW(); // 重新抛出异常 } - PG_END_TRY(); + PG_END_TRY(); // 结束异常处理 + } -void ThreadPoolSessControl::getSessionMemoryContextInfo(const char* ctx_name, - StringInfoData* buf, knl_sess_control** sess) +// 获取指定内存上下文的详细信息,并将结果保存在给定的缓冲区中 +void ThreadPoolSessControl::getSessionMemoryContextInfo(const char* ctx_name, StringInfoData* buf, knl_sess_control** sess) { -#ifdef MEMORY_CONTEXT_TRACK - AutoMutexLock alock(&m_sessCtrlock); - knl_sess_control* ctrl = NULL; - Dlelem* elem = NULL; +#ifdef MEMORY_CONTEXT_TRACK // 如果启用内存上下文追踪 + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + knl_sess_control* ctrl = NULL; // 会话控制结构指针 + Dlelem* elem = NULL; // 双向链表元素指针 PG_TRY(); { - HOLD_INTERRUPTS(); - alock.lock(); + HOLD_INTERRUPTS(); // 暂时屏蔽中断 + alock.lock(); // 获取锁 - /* collect all the Memory Context status, put in data */ - elem = DLGetHead(&m_activelist); + /* 收集所有内存上下文的状态信息,将其放入数据中 */ + elem = DLGetHead(&m_activelist); // 获取活动会话链表的头指针 while (elem != NULL) { - ctrl = (knl_sess_control*)DLE_VAL(elem); - *sess = ctrl; + ctrl = (knl_sess_control*)DLE_VAL(elem); // 获取会话控制结构指针 + *sess = ctrl; // 更新会话控制结构指针 if (ctrl->sess) { - (void)syscalllockAcquire(&ctrl->sess->utils_cxt.deleMemContextMutex); - gs_recursive_unshared_memory_context(ctrl->sess->top_mem_cxt, ctx_name, buf); - (void)syscalllockRelease(&ctrl->sess->utils_cxt.deleMemContextMutex); + (void)syscalllockAcquire(&ctrl->sess->utils_cxt.deleMemContextMutex); // 获取内存上下文互斥锁 + gs_recursive_unshared_memory_context(ctrl->sess->top_mem_cxt, ctx_name, buf); // 递归获取内存上下文信息 + (void)syscalllockRelease(&ctrl->sess->utils_cxt.deleMemContextMutex); // 释放内存上下文互斥锁 } - elem = DLGetSucc(elem); + elem = DLGetSucc(elem); // 获取下一个会话元素 } - alock.unLock(); + alock.unLock(); // 释放锁 - RESUME_INTERRUPTS(); + RESUME_INTERRUPTS(); // 恢复中断 } PG_CATCH(); { if (*sess != NULL) { ctrl = *sess; - (void)syscalllockRelease(&ctrl->sess->utils_cxt.deleMemContextMutex); + (void)syscalllockRelease(&ctrl->sess->utils_cxt.deleMemContextMutex); // 释放内存上下文互斥锁 } - alock.unLock(); - PG_RE_THROW(); + alock.unLock(); // 释放锁 + PG_RE_THROW(); // 重新抛出异常 } - PG_END_TRY(); + PG_END_TRY(); // 结束异常处理 #endif } +// 根据会话槽位索引获取会话上下文 knl_session_context* ThreadPoolSessControl::GetSessionByIdx(int idx) { if (IsValidCtrlIndex(idx)) { - return m_base[idx - m_maxReserveSessionCount].sess; + return m_base[idx - m_maxReserveSessionCount].sess; // 返回指定槽位索引的会话上下文 } else { - return NULL; + return NULL; // 如果索引无效,则返回 NULL } } +// 根据会话ID查找会话槽位索引 int ThreadPoolSessControl::FindCtrlIdxBySessId(uint64 id) { - int cidx = 0; + int cidx = 0; // 控制结构索引 for (cidx = 0; cidx < m_maxActiveSessionCount; cidx++) { if (m_base[cidx].sess != NULL && m_base[cidx].sess->session_id == id) { - return cidx + m_maxReserveSessionCount; + return cidx + m_maxReserveSessionCount; // 返回找到的会话槽位索引 } } - ereport(LOG, (errcode(ERRCODE_INVALID_PARAMETER_VALUE), errmsg("Invalid session id"))); + ereport(LOG, (errcode(ERRCODE_INVALID_PARAMETER_VALUE), errmsg("Invalid session id"))); // 记录日志,表示会话ID无效 - return -1; + return -1; // 返回-1表示未找到 } + +// 检查会话超时 void ThreadPoolSessControl::CheckSessionTimeout() { - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); - int cidx; - TimestampTz now = GetCurrentTimestamp(); + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + alock.lock(); // 获取锁 + int cidx; // 控制结构索引 + TimestampTz now = GetCurrentTimestamp(); // 获取当前时间戳 + + // 遍历所有活动会话的控制结构 for (cidx = 0; cidx < m_maxActiveSessionCount; cidx++) { - knl_sess_control* ctrl = &m_base[cidx]; - knl_session_context* sess = ctrl->sess; + knl_sess_control* ctrl = &m_base[cidx]; // 获取会话控制结构指针 + knl_session_context* sess = ctrl->sess; // 获取会话上下文指针 + + // 如果会话不为NULL且设置了会话超时 if (sess != NULL && sess->attr.attr_common.SessionTimeout != 0) { + // 如果会话的存储上下文中的会话超时活动标志为真 if (sess->storage_cxt.session_timeout_active) { + // 如果当前时间超过了会话的结束时间,并且会话的超时计数小于10 if (now >= sess->storage_cxt.session_fin_time && sess->attr.attr_common.SessionTimeoutCount < 10) { #ifdef HAVE_INT64_TIMESTAMP elog(LOG, "close session : %lu for %d times due to session timeout : %d, max finish time is %ld. But now is:%ld", @@ -713,25 +810,28 @@ void ThreadPoolSessControl::CheckSessionTimeout() sess->session_id, sess->attr.attr_common.SessionTimeoutCount + 1, sess->attr.attr_common.SessionTimeout, sess->storage_cxt.session_fin_time, now); #endif - sess->attr.attr_common.SessionTimeoutCount++; - CloseClientSocket(sess, false); + sess->attr.attr_common.SessionTimeoutCount++; // 增加会话的超时计数 + CloseClientSocket(sess, false); // 关闭客户端套接字 } } } } - alock.unLock(); + alock.unLock(); // 释放锁 } +// 列出所有会话的全局临时表(GTT)的冻结事务ID TransactionId ThreadPoolSessControl::ListAllSessionGttFrozenxids(int maxSize, ThreadId *pids, TransactionId *xids, int *n) { - TransactionId result = InvalidTransactionId; - int i = 0; + TransactionId result = InvalidTransactionId; // 初始化冻结事务ID为无效事务ID + int i = 0; // 计数器 + // 如果未启用全局临时表,直接返回 if (u_sess->attr.attr_storage.max_active_gtt <= 0) { return 0; } + // 如果传入的参数合法,初始化相关变量 if (maxSize > 0) { Assert(pids); Assert(xids); @@ -739,51 +839,58 @@ TransactionId ThreadPoolSessControl::ListAllSessionGttFrozenxids(int maxSize, *n = 0; } + // 如果数据库中未启用全局临时表,返回无效事务ID if (u_sess->attr.attr_storage.max_active_gtt <= 0) { return InvalidTransactionId; } + // 如果处于恢复状态,返回无效事务ID if (RecoveryInProgress()) { return InvalidTransactionId; } - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); - knl_sess_control *ctl = nullptr; - knl_session_context *session = nullptr; - Dlelem* elem = DLGetHead(&m_activelist); + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + alock.lock(); // 获取锁 + knl_sess_control *ctl = nullptr; // 会话控制结构指针 + knl_session_context *session = nullptr; // 会话上下文指针 + Dlelem* elem = DLGetHead(&m_activelist); // 获取活动会话链表的头指针 while (elem != nullptr) { - ctl = (knl_sess_control*)DLE_VAL(elem); - session = ctl->sess; + ctl = (knl_sess_control*)DLE_VAL(elem); // 获取会话控制结构指针 + session = ctl->sess; // 获取会话上下文指针 + // 如果会话的数据库ID与当前会话的数据库ID相同,并且会话的GTT冻结事务ID有效 if (session->proc_cxt.MyDatabaseId == u_sess->proc_cxt.MyDatabaseId && TransactionIdIsNormal(session->gtt_ctx.gtt_session_frozenxid)) { + // 如果result为无效事务ID,或者当前会话的GTT冻结事务ID较小,则更新result if (result == InvalidTransactionId) { result = session->gtt_ctx.gtt_session_frozenxid; } else if (TransactionIdPrecedes(session->gtt_ctx.gtt_session_frozenxid, result)) { result = session->gtt_ctx.gtt_session_frozenxid; } + // 如果传入的参数允许保存信息 if (maxSize > 0) { - pids[i] = session->attachPid; - xids[i] = session->gtt_ctx.gtt_session_frozenxid; - i++; + pids[i] = session->attachPid; // 保存会话的进程ID + xids[i] = session->gtt_ctx.gtt_session_frozenxid; // 保存会话的GTT冻结事务ID + i++; // 计数器加1 } } - elem = DLGetSucc(elem); + elem = DLGetSucc(elem); // 获取下一个会话 } - alock.unLock(); + alock.unLock(); // 释放锁 + // 如果传入的参数允许保存信息,更新保存信息的数量 if (maxSize > 0) { *n = i; } - return result; + return result; // 返回最大的GTT冻结事务ID } +// 检查活动会话链表是否为空 bool ThreadPoolSessControl::IsActiveListEmpty() { - AutoMutexLock alock(&m_sessCtrlock); - alock.lock(); - bool res = (m_activelist.dll_len == 0); - alock.unLock(); - return res; -} + AutoMutexLock alock(&m_sessCtrlock); // 自动互斥锁,确保线程安全 + alock.lock(); // 获取锁 + bool res = (m_activelist.dll_len == 0); // 检查活动会话链表的长度是否为0 + alock.unLock(); // 释放锁 + return res; // 返回结果,true表示链表为空,false表示链表不为空 +} \ No newline at end of file