# -*- coding: utf-8 -*- """Обнаружение сеансов пользователей по домену. Логика повторяет то, как это делает Dameware NT Utilities: 1. список компьютеров берём из AD (ADSI/LDAP, RSAT не нужен); 2. быстро отсеиваем выключенные машины проверкой порта 445 (асинхронно); 3. живые опрашиваем через WTS API — тот же интерфейс, что у диспетчера задач, без разбора текстового вывода quser (он зависит от локали Windows). Модуль не содержит GUI и может использоваться отдельно. """ import ctypes import os import socket import threading import time from collections import deque from ctypes import wintypes import win32api import win32con import win32security # WTS дёргаем через ctypes, а не через pywin32, сознательно: pywin32 не отпускает # GIL на время нативного вызова, из-за чего опрос машин выполняется строго # последовательно и никакие потоки не помогают. ctypes.WinDLL GIL освобождает, # и параллельный опрос домена начинает работать по-настоящему. _wts = ctypes.WinDLL("wtsapi32.dll", use_last_error=True) WTS_CURRENT_SERVER_HANDLE = 0 # Классы информации о сеансе (WTS_INFO_CLASS) _WTS_USER_NAME = 5 _WTS_DOMAIN_NAME = 7 _WTS_CLIENT_NAME = 10 class _WTS_SESSION_INFOW(ctypes.Structure): _fields_ = [ ("SessionId", wintypes.DWORD), ("pWinStationName", wintypes.LPWSTR), ("State", ctypes.c_int), ] _wts.WTSOpenServerW.argtypes = [wintypes.LPWSTR] _wts.WTSOpenServerW.restype = wintypes.HANDLE _wts.WTSCloseServer.argtypes = [wintypes.HANDLE] _wts.WTSCloseServer.restype = None _wts.WTSEnumerateSessionsW.argtypes = [ wintypes.HANDLE, wintypes.DWORD, wintypes.DWORD, ctypes.POINTER(ctypes.POINTER(_WTS_SESSION_INFOW)), ctypes.POINTER(wintypes.DWORD), ] _wts.WTSEnumerateSessionsW.restype = wintypes.BOOL _wts.WTSQuerySessionInformationW.argtypes = [ wintypes.HANDLE, wintypes.DWORD, ctypes.c_int, ctypes.POINTER(ctypes.c_void_p), ctypes.POINTER(wintypes.DWORD), ] _wts.WTSQuerySessionInformationW.restype = wintypes.BOOL _wts.WTSFreeMemory.argtypes = [ctypes.c_void_p] _wts.WTSFreeMemory.restype = None # Значения по умолчанию для сканирования домена. # Опрос упирается не в процессор, а в ожидание сети, поэтому потоков берём много. DEFAULT_WORKERS = 256 # Жёсткий предел на весь опрос. Машины, у которых открыт порт 445, но закрыт RPC, # вешают WTSOpenServer на 20-45 секунд — без общего дедлайна они растягивают скан. DEFAULT_DEADLINE = 30.0 DEFAULT_PORT_TIMEOUT = 0.4 # Сколько ждать отставшие машины после того, как очередь разобрана. # Зависшие на RPC не дождутся никогда, а живые обычно отвечают за доли секунды. DEFAULT_GRACE = 4.0 # Состояния сеанса WTS WTS_ACTIVE = 0 WTS_DISCONNECTED = 4 STATE_NAMES = { 0: "Активен", 1: "Подключается", 2: "Запрос подключения", 3: "Теневой", 4: "Отключён", 5: "Простой", 6: "Ожидание", 7: "Сброс", 8: "Отключение", 9: "Инициализация", } # Сеансы, которые держат профиль пользователя загруженным. Отключённый (Disconnected) # сеанс профиль НЕ выгружает, поэтому при перемещаемых профилях писать в него нужно # так же, как в активный — иначе правку затрёт при выходе пользователя. PROFILE_HELD_STATES = (WTS_ACTIVE, WTS_DISCONNECTED) PROFILE_LIST_KEY = r"SOFTWARE\Microsoft\Windows NT\CurrentVersion\ProfileList" # ------------------------------------------------------------------ AD def list_domain_computers(): """Список компьютеров домена (активные учётки) через ADSI. Возвращает список имён (FQDN, если заполнен dNSHostName). Используется только pywin32 — дополнительных зависимостей и RSAT не требуется. """ import pythoncom import win32com.client # COM инициализируется отдельно в КАЖДОМ потоке. Ленивый импорт win32com делает # это неявно только для того потока, где импорт случился первым, поэтому повторный # поиск (кнопка «Пересканировать» создаёт новый поток) падал с MK_E_SYNTAX # «Синтаксическая ошибка». Инициализируем явно и симметрично освобождаем. try: pythoncom.CoInitialize() com_ready = True except Exception: com_ready = False # поток уже в другой модели — работаем как есть root = conn = cmd = rs = None names = [] try: root = win32com.client.GetObject("LDAP://RootDSE") base_dn = root.Get("defaultNamingContext") conn = win32com.client.Dispatch("ADODB.Connection") conn.Provider = "ADsDSOObject" conn.Open("Active Directory Provider") cmd = win32com.client.Dispatch("ADODB.Command") cmd.ActiveConnection = conn # Без Page Size AD вернёт максимум 1000 записей и молча обрежет остальные cmd.Properties("Page Size").Value = 1000 cmd.Properties("Timeout").Value = 30 cmd.CommandText = ( f";" # (!userAccountControl:...:=2) — отбрасываем отключённые учётки компьютеров "(&(objectCategory=computer)(!userAccountControl:1.2.840.113556.1.4.803:=2));" "dNSHostName,name;subtree" ) rs = cmd.Execute() if isinstance(rs, tuple): # позднее связывание отдаёт (recordset, records_affected) rs = rs[0] while not rs.EOF: dns = rs.Fields.Item("dNSHostName").Value name = rs.Fields.Item("name").Value host = dns or name if host: names.append(str(host)) rs.MoveNext() rs.Close() finally: if conn is not None: try: conn.Close() except Exception: pass # Отпускаем COM-объекты ДО CoUninitialize: иначе их деструкторы сработают, # когда COM в потоке уже деинициализирован, и посыплется # "Win32 exception occurred releasing IUnknown". rs = cmd = conn = root = None if com_ready: try: pythoncom.CoUninitialize() except Exception: pass return names def current_domain(): """NetBIOS-имя домена текущего пользователя (пустая строка, если не в домене).""" return os.environ.get("USERDOMAIN", "") # ------------------------------------------------------- Проверка доступности def is_alive(host, port=445, timeout=0.4): """Быстрая проверка, что машина включена и отвечает по SMB.""" try: with socket.create_connection((host, port), timeout=timeout): return True except OSError: return False def _parallel(func, items, workers, deadline, grace=DEFAULT_GRACE, progress=None, stage=""): """Выполняем func по всем items в несколько потоков, не ожидая зависших. Своя реализация вместо ThreadPoolExecutor по трём причинам: * потоки демонические — зависший в нативном вызове WTSOpenServer поток не мешает закрыть программу (ThreadPoolExecutor ждёт свои потоки на выходе); * прервать зависший вызов нельзя, но можно перестать его ждать: как только очередь машин разобрана, даём отставшим ровно `grace` секунд и уходим. Именно это отличает 5 секунд от полутора минут — машины с закрытым RPC висят по 20-45 с каждая, и ждать их бессмысленно; * `deadline` остаётся страховкой на случай очень большого парка. Возвращает (results, unfinished), где results — список (item, value, error). """ items = list(items) total = len(items) if not total: return [], [] pending = deque(items) pending_lock = threading.Lock() results = [] results_lock = threading.Lock() counter = [0] in_flight = [0] end_at = time.monotonic() + deadline def worker(): while True: if time.monotonic() >= end_at: return with pending_lock: if not pending: return item = pending.popleft() in_flight[0] += 1 try: value, error = func(item), None except Exception as exc: value, error = None, f"{type(exc).__name__}: {exc}" with results_lock: results.append((item, value, error)) counter[0] += 1 done = counter[0] with pending_lock: in_flight[0] -= 1 # Прогресс обновляем пачками, чтобы не забивать очередь событий Tk if progress and (done % 5 == 0 or done == total): progress(stage, done, total) threads = [threading.Thread(target=worker, daemon=True) for _ in range(min(workers, total))] for t in threads: t.start() drain_started = None while True: now = time.monotonic() if now >= end_at: break with pending_lock: queue_empty = not pending running = in_flight[0] if queue_empty: if running == 0: break # все машины честно опрошены if drain_started is None: drain_started = now elif now - drain_started >= grace: break # остальные зависли — дальше не ждём time.sleep(0.05) with results_lock: handled = {item for item, _, _ in results} snapshot = list(results) unfinished = [i for i in items if i not in handled] return snapshot, unfinished def filter_alive(hosts, timeout=DEFAULT_PORT_TIMEOUT, workers=DEFAULT_WORKERS): """Оставляем только машины, ответившие на порт 445.""" results, _ = _parallel(lambda h: is_alive(h, timeout=timeout), hosts, workers, DEFAULT_DEADLINE) return [host for host, alive, err in results if alive and not err] # ------------------------------------------------------------------ WTS def _query_session(handle, session_id, info_class): """Строковое свойство сеанса; пустая строка, если недоступно.""" buffer = ctypes.c_void_p() returned = wintypes.DWORD() ok = _wts.WTSQuerySessionInformationW( handle, session_id, info_class, ctypes.byref(buffer), ctypes.byref(returned)) if not ok or not buffer: return "" try: return ctypes.cast(buffer, ctypes.c_wchar_p).value or "" finally: _wts.WTSFreeMemory(buffer) def enum_sessions(machine=None): """Сеансы одной машины. machine=None — локальная. Возвращает список словарей. Сеансы без пользователя (службы, listener) пропускаем. """ remote = bool(machine) if remote: handle = _wts.WTSOpenServerW(machine) if not handle: raise OSError(ctypes.WinError(ctypes.get_last_error())) else: handle = WTS_CURRENT_SERVER_HANDLE info_ptr = ctypes.POINTER(_WTS_SESSION_INFOW)() count = wintypes.DWORD() result = [] try: ok = _wts.WTSEnumerateSessionsW( handle, 0, 1, ctypes.byref(info_ptr), ctypes.byref(count)) if not ok: raise OSError(ctypes.WinError(ctypes.get_last_error())) try: for i in range(count.value): entry = info_ptr[i] session_id = entry.SessionId user = _query_session(handle, session_id, _WTS_USER_NAME) if not user: continue result.append({ "machine": machine or os.environ.get("COMPUTERNAME", ""), "session_id": session_id, "station": entry.pWinStationName or "", "state": entry.State, "state_name": STATE_NAMES.get(entry.State, str(entry.State)), "user": user, "domain": _query_session(handle, session_id, _WTS_DOMAIN_NAME), "client": _query_session(handle, session_id, _WTS_CLIENT_NAME), }) finally: _wts.WTSFreeMemory(ctypes.cast(info_ptr, ctypes.c_void_p)) finally: if remote: _wts.WTSCloseServer(handle) return result _OFFLINE = object() # маркер: машина не ответила на порт 445 def scan_domain(computers=None, workers=DEFAULT_WORKERS, port_timeout=DEFAULT_PORT_TIMEOUT, deadline=DEFAULT_DEADLINE, grace=DEFAULT_GRACE, progress=None): """Снимаем сеансы со всех машин домена одним проходом. Проверка порта и опрос WTS выполняются в одной задаче: поток, освободившись, сразу берёт следующую машину. Раньше это были две последовательные фазы, и вторая простаивала, пока первая доделывала самую медленную машину. Возвращает (sessions, errors, stats) — сеансы всех пользователей, без фильтра. """ t_start = time.monotonic() if progress: progress("Получаем список компьютеров из AD...", 0, 0) t0 = time.monotonic() hosts = list(computers) if computers else list_domain_computers() ad_seconds = time.monotonic() - t0 stats = {"total": len(hosts), "alive": 0, "offline": 0, "failed": 0, "unfinished": 0, "sessions_total": 0, "ad_seconds": round(ad_seconds, 2), "scan_seconds": 0.0, "total_seconds": 0.0} if not hosts: return [], [], stats def probe(host): # Быстрый отсев выключенных машин: без него WTSOpenServer будет ждать # RPC-таймаут в десятки секунд на каждой недоступной машине. if not is_alive(host, timeout=port_timeout): return _OFFLINE return enum_sessions(host) t0 = time.monotonic() results, unfinished = _parallel(probe, hosts, workers, deadline, grace, progress, "Опрашиваем машины...") scan_seconds = time.monotonic() - t0 sessions, errors = [], [] for host, value, error in results: if error: errors.append((host, error)) elif value is _OFFLINE: stats["offline"] += 1 else: stats["alive"] += 1 sessions.extend(value) stats["failed"] = len(errors) stats["unfinished"] = len(unfinished) stats["sessions_total"] = len(sessions) stats["scan_seconds"] = round(scan_seconds, 2) stats["total_seconds"] = round(time.monotonic() - t_start, 2) return sessions, errors, stats def filter_sessions(sessions, user=None, domain=None, include_disconnected=True): """Отбираем из готового списка сеансы нужного пользователя. Вынесено отдельно, чтобы поиск другого пользователя по уже собранным данным происходил мгновенно, без повторного опроса домена. """ allowed = PROFILE_HELD_STATES if include_disconnected else (WTS_ACTIVE,) out = [s for s in sessions if s["state"] in allowed] if user: target = user.strip().lower() out = [s for s in out if s["user"].lower() == target] if domain: dom = domain.strip().lower() # У локальных учёток в domain стоит имя машины — такие отсеиваем out = [s for s in out if s["domain"].lower() == dom] out.sort(key=lambda s: (s["machine"].lower(), s["session_id"])) return out def find_sessions(user=None, computers=None, domain=None, include_disconnected=True, port_timeout=DEFAULT_PORT_TIMEOUT, workers=DEFAULT_WORKERS, deadline=DEFAULT_DEADLINE, progress=None): """Опрашиваем домен и сразу отбираем сеансы одного пользователя.""" sessions, errors, stats = scan_domain( computers=computers, workers=workers, port_timeout=port_timeout, deadline=deadline, progress=progress) matched = filter_sessions(sessions, user=user, domain=domain, include_disconnected=include_disconnected) stats["matched"] = len(matched) return matched, errors, stats def verify_sessions(sessions): """Перепроверяем перед записью, что сеансы ещё живы. Данные скана могут быть слегка устаревшими — пользователь мог выйти. Машин здесь единицы, так что проверка почти мгновенная. """ still_there, gone = [], [] checked = {} for sess in sessions: machine = sess["machine"] if machine not in checked: try: checked[machine] = enum_sessions(machine) except Exception: checked[machine] = None # не смогли проверить — не мешаем записи current = checked[machine] if current is None: still_there.append(sess) continue match = any(c["user"].lower() == sess["user"].lower() and c["domain"].lower() == sess["domain"].lower() and c["state"] in PROFILE_HELD_STATES for c in current) (still_there if match else gone).append(sess) return still_there, gone # -------------------------------------------------------------- Путь профиля def _to_unc(machine, local_path): """C:\\Users\\Ivanov на машине PC1 -> \\\\PC1\\C$\\Users\\Ivanov""" drive, rest = os.path.splitdrive(local_path) if not drive: return local_path return f"\\\\{machine}\\{drive[0]}$" + rest def resolve_profile_dir(machine, domain, user): """Каталог профиля пользователя на машине, в виде UNC-пути. Основной способ — SID + реестр ProfileList: корректно разрешает случаи вида Ivanov.CORP или Ivanov.000, которые не угадать по имени. Если удалённый реестр недоступен (служба RemoteRegistry остановлена), перебираем типовые варианты имени папки. """ account = f"{domain}\\{user}" if domain else user sid_str = None for lookup_host in (machine, None): try: sid_obj, _, _ = win32security.LookupAccountName(lookup_host, account) sid_str = win32security.ConvertSidToStringSid(sid_obj) break except Exception: continue if sid_str: try: root = win32api.RegConnectRegistry(f"\\\\{machine}", win32con.HKEY_LOCAL_MACHINE) key = win32api.RegOpenKeyEx(root, PROFILE_LIST_KEY + "\\" + sid_str, 0, win32con.KEY_READ) path, _ = win32api.RegQueryValueEx(key, "ProfileImagePath") win32api.RegCloseKey(key) path = win32api.ExpandEnvironmentStrings(path) if path: return _to_unc(machine, path) except Exception: pass # Фоллбэк: перебираем типовые имена папок профиля candidates = [user] if domain: candidates.append(f"{user}.{domain}") candidates.append(f"{user}.000") for name in candidates: unc = f"\\\\{machine}\\C$\\Users\\{name}" try: if os.path.isdir(unc): return unc except Exception: continue raise RuntimeError( "не удалось определить папку профиля " "(нет доступа к удалённому реестру и папка не найдена перебором)" ) def ibases_path_for_session(session): """Путь к ibases.v8i внутри профиля пользователя из сеанса.""" profile = resolve_profile_dir(session["machine"], session["domain"], session["user"]) return os.path.join(profile, "AppData", "Roaming", "1C", "1CEStart", "ibases.v8i") def short_name(machine): """SERVERTS.CORP.local -> SERVERTS (для компактного отображения в списке).""" return machine.split(".")[0] if machine else machine