吾爱破解 - 52pojie.cn

 找回密码
 注册[Register]

QQ登录

只需一步,快速开始

查看: 399|回复: 1
上一主题 下一主题
收起左侧

[Python 转载] 蓝牙,传讯 (Win)

[复制链接]
跳转到指定楼层
楼主
pyjiujiu 发表于 2026-10-6 21:35 回帖奖励



一款 用来 借助蓝牙 分享信息的软件

初心:
有台设备 网络模块不能用,U盘复制太麻烦(主要小文件 文字),于是想到蓝牙,
最初想着 是奔着 剪贴板同步应用, 开发中途发现 agent 开发聊天应用 太顺手,遂完全转向 聊天应用的模样

---

# 使用相关
限于 windows (android 开发没完)
蓝牙 BLE  速度很慢(省电),就好象回到  2008年时代的 qq,发文字,.txt 是完全可以的(剪贴板 分享文章)
目前 电脑睡眠唤醒 还支持恢复重连,需要重启软件
还有一些bug,偶然会卡(需要重启软件)

---
AI 背景:
- mimo 2.6 flash ,负责原型搭建 (opencode 免费版)
- ds 4.1 flash 负责最初计划(web端构思+测试),还有后期迭代(DSH)

使用来说,模型好不好用,个人也只片面,,就本项目来说,mimo 出原型还可以,修深入的bug 还是不如 ds 4.1
winRT 个人完全没法写,没 agent的话, 这个项目就不会出现

遇到个新的东西,有种跟不上 大模型的节奏
`§` 第一次看到 大模型(mimo) 频繁使用这个术语符号(大概意思是 标记章节的)


---
软件截图:


---

项目 维护在 github:      github.com/fun-tailor/ble_chat

---

源代码
(摘录 ‘blechat\ble\server.py’  和 winRT 相关 都在这个)


[Python] 纯文本查看 复制代码
"""Host 侧传输:WinRT GattServiceProvider(广播 + GATT server)。"""

from __future__ import annotations

import asyncio
import logging
import threading
import time
from typing import Callable

from winrt.windows.devices.bluetooth.genericattributeprofile import (
    GattCharacteristicProperties as Props,
    GattLocalCharacteristic,
    GattLocalCharacteristicParameters,
    GattProtectionLevel,
    GattServiceProvider,
    GattServiceProviderAdvertisementStatus,
    GattServiceProviderAdvertisingParameters,
)
from winrt.windows.storage.streams import DataReader, DataWriter

from ..protocol import DATA_HEADER
from .uuid_defs import (
    ADVERTISE_RETRY,
    ADVERTISE_WAIT,
    CHAR_CTRL,
    CHAR_TX,
    UUID_CTRL,
    UUID_RX,
    UUID_SERVICE,
    UUID_TX,
)

log = logging.getLogger("blechat.ble.server")

# watchdog 最短两次 revive 之间的间隔:反复 stop/start_advertising 会把在线对端踢掉
REVIVE_COOLDOWN = 30.0
# 会话多久没有任何入站数据就算"僵尸"。
# 两端各有 20s 一跳的 PING/PONG,所以**真活着**的对端不可能静默这么久;
# 而对端机器睡眠/崩溃时 WinRT 可能既不回调 `session_status_changed` 也不改
# `subscribed_clients`,`_sessions` 里就留着一条"看着还活着"的僵尸。
# 这条僵尸会让 `advertising_ok` 恒为真 ⇒ watchdog 永远不去恢复广播 ⇒
# 对端唤醒后一直连不上,而 host 这边看着一切正常。
SESSION_IDLE_TRUST = 60.0


# debug phone
def _dump_recv(channel: str, peer_id: str, data: bytes, *, limit: int = 256) -> None:
    """把手机写过来的字节打出来(hex + 可打印文本),过长的只打印前 limit 字节。"""
    head = data[:limit]
    tail = "" if len(data) <= limit else f" ...(+{len(data) - limit}B)"
    hexs = head.hex(" ")
    txt = "".join(chr(b) if 32 <= b < 127 else "." for b in head)
    log.info(
        "recv %s from %s: %d bytes%s\n  HEX: %s\n  TXT: %s",
        channel, peer_id[-12:], len(data), tail, hexs, txt,
    )

class PeerGone(Exception):
    pass


# RO_E_CLOSED = 0x80000013 "The object has been closed."(WinRT 对象已被关闭)
RO_E_CLOSED = 0x80000013
RO_E_CLOSED_SIGNED = -2147483629


def is_closed_error(exc: BaseException) -> bool:
    """判断异常是否是 WinRT `RO_E_CLOSED`(对端/会话对象已销毁)。

    pywinrt 抛出的 OSError 里 HRESULT 落在 `winerror`,且**可能是有符号的**
    (真机实测 `OSError(22, '该对象已关闭。', None, -2147483629)`,`errno=22` 与
    HRESULT 无关,`str()` 只有本地化文本、没有 hex)。因此:
    1. `winerror` 先做 `& 0xFFFFFFFF` 归一化再比;
    2. `errno` 分支保留(两参数构造 `OSError(0x80000013, ...)` 会把 HRESULT 放这里);
    3. 文本分支补十进制与本地化关键词 —— 中文系统上消息是「该对象已关闭。」。
    """
    if isinstance(exc, OSError):
        winerror = getattr(exc, "winerror", None)
        if winerror is not None and (winerror & 0xFFFFFFFF) == RO_E_CLOSED:
            return True
        if exc.errno in (RO_E_CLOSED, RO_E_CLOSED_SIGNED):
            return True
    text = str(exc)
    if "0x80000013" in text or "80000013" in text:
        return True
    if str(RO_E_CLOSED_SIGNED) in text:  # '[WinError -2147483629] ...'
        return True
    low = text.lower()
    return "has been closed" in low or "已关闭" in text


def is_link_dead(exc: BaseException) -> bool:
    """判断异常是否意味着**链路已经死了**(必须关会话触发重连)。

    早期只认 `RO_E_CLOSED`,导致 `GATT Protocol Error: Unlikely Error`
    这类"对端已不在"的错误只被打 WARNING、会话保持 READY → 永远重连不上。
    """
    if isinstance(exc, PeerGone):
        return True
    if is_closed_error(exc):
        return True
    text = str(exc).lower()
    if "gatt protocol error" in text:
        return True
    if "not connected" in text or "not connected to" in text:
        return True
    if "unexpected" in text and "error" in text:
        return True
    if "device" in text and "connect" in text and "fail" in text:
        return True
    if "unreachable" in text or "host is down" in text:
        return True
    # `GattServer._subscribed` 抛的 `PeerGone(peer … not subscribed to 0000a003)`。
    # 包一层 PeerGone 时已经 True,但同文案也可能从 WinRT 回调原样冒出来,
    # dev 的日志里就见过 —— 漏判会话会停在 READY,直到 20s PING 才发现。
    if "not subscribed" in text:
        return True
    return False


def to_buffer(data: bytes):
    w = DataWriter()
    w.write_bytes(data)
    return w.detach_buffer()


def from_buffer(buf) -> bytes:
    r = DataReader.from_buffer(buf)
    out = bytearray(buf.length)
    r.read_bytes(out)
    return bytes(out)


def session_peer_id(session) -> str:
    """GattSession → 稳定 peer_id。

    `GattSession.device_id` 返回的是 `BluetoothDeviceId` **对象**,每次访问都会
    生成新的包装对象,`str()` 得到的 repr(`<BluetoothDeviceId object at …>`)
    每次都不同 —— 直接当 peer_id 会导致同一次连接的每次写都新建会话、且
    `subscribed_clients` 匹配永远失败(服务端发不出任何通知)。
    这里取其 `.id` 字符串(设备接口路径,按设备稳定)。
    """
    dev = getattr(session, "device_id", "")
    ident = getattr(dev, "id", None)
    if ident:
        return str(ident)
    return str(dev)


class GattServer:
    def __init__(self) -> None:
        self._provider: GattServiceProvider | None = None
        self._chars: dict[str, GattLocalCharacteristic] = {}
        self._sessions: dict[str, object] = {}
        self._last_activity: dict[str, float] = {}
        self._started_at = time.monotonic()
        self._loop: asyncio.AbstractEventLoop | None = None
        self._loop_thread: int | None = None
        self._running = False
        self._last_revive: float = -1e9

        self.on_peer_ctrl: Callable[[str, bytes], None] | None = None
        self.on_peer_data: Callable[[str, bytes], None] | None = None
        self.on_peer_connected: Callable[[str], None] | None = None
        self.on_peer_disconnected: Callable[[str, str], None] | None = None

    @property
    def running(self) -> bool:
        return self._running

    @property
    def advertising_status(self) -> str:
        if not self._provider:
            return "NONE"
        try:
            return GattServiceProviderAdvertisementStatus(self._provider.advertisement_status).name
        except Exception:
            return "UNKNOWN"

    @property
    def advertising_ok(self) -> bool:
        """广播是否真的在跑(睡眠唤醒后 WinRT 会悄悄停掉)。

        注意:已经有对端连上来时,Windows 可能会把 `advertisement_status`
        变成 `ABORTED` —— 这是**正常**的。若仍按状态判断,watchdog 会每 5s
        `stop_advertising()` 一次把在线客户端踢掉(表现为对端 Unlikely Error
        / RO_E_CLOSED、server 端 UI 反复重建目标菜单)。

        但"有会话就算健康"也有反例(第八轮真机):对端机器睡眠/崩溃后
        WinRT 可能留下一条**僵尸会话**,于是这里恒为真、watchdog 永远不去
        恢复广播 —— 对端唤醒后一直连不上,host 这边却一切正常。
        所以对端必须**最近说过话**(`SESSION_IDLE_TRUST` 内)才算数;
        太久没声音的会话当作僵尸,让 watchdog 去 revive。
        """
        if not self._running or self._provider is None:
            return False
        if self.advertising_status in (
            "STARTED",
            "STARTED_WITHOUT_ALL_ADVERTISEMENT_DATA",
        ):
            return True
        live = [pid for pid in self._sessions if self.idle_seconds(pid) < SESSION_IDLE_TRUST]
        if live:
            return True
        if self._sessions:
            log.info(
                "advertising status=%s and every session idle > %.0fs -> needs revive (%s)",
                self.advertising_status,
                SESSION_IDLE_TRUST,
                {pid[:24]: round(self.idle_seconds(pid)) for pid in self._sessions},
            )
        return False

    def touch(self, peer_id: str) -> None:
        """记一次"这个对端刚说过话"(用于僵尸会话判定,见 `advertising_ok`)。"""
        if peer_id:
            self._last_activity[peer_id] = time.monotonic()

    def idle_seconds(self, peer_id: str) -> float:
        """这个对端静默了多久(没有记录的按"从 server 启动算起")。"""
        return time.monotonic() - self._last_activity.get(peer_id, self._started_at)

    def stale_sessions(self, limit: float = SESSION_IDLE_TRUST) -> list[str]:
        return [pid for pid in self._sessions if self.idle_seconds(pid) >= limit]

    @property
    def peer_count(self) -> int:
        return len(self._sessions)

    @property
    def session_count(self) -> int:
        """当前记录的底层会话数(含尚未完成握手的),供诊断日志使用。"""
        return len(self._sessions)

    def has_peers(self) -> bool:
        return bool(self._sessions)

    async def revive(self, *, hard: bool = False) -> bool:
        """重新拉起广播。返回是否恢复成功。带 30s 冷却,避免反复踢对端。

        `hard=True`(本机睡眠唤醒后)走 `restart()`:把 provider 整段重建。
        光 `stop_advertising()` + `start_advertising()` 拉不回来 —— 唤醒后
        Windows 可能已经把服务的注册丢了,`advertisement_status` 会一直停在
        `ABORTED`,只有重建 provider 才能恢复。
        """
        if not self._running or self._provider is None:
            return False
        if self._sessions and not self.stale_sessions():
            # 真有对端在线:不要动广播,动了就是把人家踢下线
            return True
        if hard:
            return await self.restart(reason="hard revive")
        now = time.monotonic()
        if now - self._last_revive < REVIVE_COOLDOWN:
            return self.advertising_ok
        self._last_revive = now
        self.sweep_stale(reason="revive")
        try:
            self._provider.stop_advertising()
        except Exception:
            pass
        await asyncio.sleep(0.3)
        await self._advertise()
        return self.advertising_ok

    async def restart(self, *, reason: str) -> bool:
        """整段重建 GATT server(provider + 特征 + 广播)。

        用于两种"WinRT 对象已经不可信"的场合:本机睡眠唤醒、以及
        `advertisement_status` 卡在非 STARTED 又拉不回来。
        **只在没有活跃对端时调用**(调用方负责判断,见 `revive`):
        重建意味着旧的特征对象全部作废,在线对端会立刻掉线。
        """
        log.warning("restarting gatt server (%s)", reason)
        self._running = False
        provider, self._provider = self._provider, None
        if provider is not None:
            try:
                provider.stop_advertising()
            except Exception:
                pass
        for peer_id in list(self._sessions):
            self.drop_peer(peer_id)
        self._sessions.clear()
        self._last_activity.clear()
        self._chars.clear()
        # 旧 provider 要能被回收(COM 引用释放 + WinRT 侧注销)才能建新的,
        # 立刻重建常常撞 "服务已被注册"。
        await asyncio.sleep(0.5)
        try:
            await self.start()
        except Exception as exc:
            log.error("gatt server restart failed (%s): %s", reason, exc)
            return False
        log.warning(
            "gatt server restarted (%s): advertising=%s", reason, self.advertising_status
        )
        return True

    def sweep_stale(self, *, reason: str) -> int:
        """清掉**没有活跃会话**的残留订阅与僵尸会话,返回清掉的条目数。

        为什么需要:对端软件直接下线/崩溃时,WinRT 可能既不回调
        `session_status_changed` 也不回调 `subscribed_clients_changed`,
        `_sessions`(以及 GATT server 的订阅表)里就会留着上一次的幽灵条目。
        之后新连接上来时,`subscribed_clients` 里混着幽灵 + 新条目,
        `_subscribed()` 可能挑中幽灵 &#8658; 通知发不出去 &#8658; dev 看到的
        「双断开后点加入网络,server 单方显示已连接、client 按钮灰」。

        只在**没有任何底层会话还活着**时动手(启动 / 恢复广播),
        因此不会误伤在线对端。

        注意:`GattLocalCharacteristic.subscribed_clients` 是 WinRT 内部维护的
        只读集合,**没有 API 可以主动移除订阅者** —— 这里能做的是把 Python 侧
        的记录清干净,并让 `_subscribed()` 优先挑活着的条目;真正的清理依赖
        WinRT 在链路断开后自己回收,以及两端的 `PING&#8596;PONG` 探活兜底。
        """
        if any(self._session_alive(s) for s in self._sessions.values()):
            return 0
        removed = len(self._sessions)
        self._sessions.clear()
        self._last_activity.clear()
        for key in (CHAR_TX, CHAR_CTRL):
            char = self._chars.get(key)
            if char is None:
                continue
            try:
                subs = list(char.subscribed_clients or ())
            except Exception:
                continue
            if subs:
                removed += 1
                log.info(
                    "%s: %s still holds %s subscriber(s) with no live session",
                    reason,
                    key[:8],
                    len(subs),
                )
        if removed:
            log.info("%s: swept %s stale peer record(s)", reason, removed)
        return removed

    # ------------------------------------------------------------- lifecycle

    async def start(self) -> None:
        self._loop = asyncio.get_running_loop()
        self._loop_thread = threading.get_ident()
        self._started_at = time.monotonic()

        result = await GattServiceProvider.create_async(UUID_SERVICE)
        if result.error != 0:
            raise RuntimeError(f"GattServiceProvider 创建失败: {result.error}")
        provider = result.service_provider
        self._provider = provider
        service = provider.service

        for uuid_, props in (
            (UUID_RX, Props.WRITE | Props.WRITE_WITHOUT_RESPONSE),
            (UUID_TX, Props.INDICATE),
            (UUID_CTRL, Props.WRITE | Props.WRITE_WITHOUT_RESPONSE | Props.NOTIFY),
        ):
            params = GattLocalCharacteristicParameters()
            params.characteristic_properties = props
            params.read_protection_level = GattProtectionLevel.PLAIN
            params.write_protection_level = GattProtectionLevel.PLAIN
            cres = await service.create_characteristic_async(uuid_, params)
            if cres.error != 0:
                raise RuntimeError(f"特征 {uuid_} 创建失败: {cres.error}")
            self._chars[str(uuid_).lower()] = cres.characteristic

        self._chars[str(UUID_RX).lower()].add_write_requested(self._make_write_handler("rx"))
        self._chars[str(UUID_CTRL).lower()].add_write_requested(self._make_write_handler("ctrl"))
        for char in self._chars.values():
            char.add_subscribed_clients_changed(self._on_subscribed_changed)

        provider.add_advertisement_status_changed(self._on_adv_status)
        # 广播之前先把上一轮留下的幽灵会话/订阅清掉:本进程刚起来,不可能有
        # 在线对端,留着只会让第一条新连接撞上 `subscribed_clients` 里的死条目。
        self.sweep_stale(reason="server start")
        await self._advertise()
        self._running = True
        log.info("GATT server started, advertising=%s", self.advertising_status)

    async def _advertise(self) -> None:
        assert self._provider is not None
        adv = GattServiceProviderAdvertisingParameters()
        adv.is_connectable = True
        adv.is_discoverable = True
        for attempt in range(ADVERTISE_RETRY):
            try:
                self._provider.start_advertising_with_parameters(adv)
            except Exception as exc:
                log.warning("start_advertising attempt %s failed: %s", attempt, exc)
            for _ in range(int(ADVERTISE_WAIT / 0.05)):
                await asyncio.sleep(0.05)
                status = self._provider.advertisement_status
                if status == GattServiceProviderAdvertisementStatus.STARTED:
                    return
                if status == GattServiceProviderAdvertisementStatus.STARTED_WITHOUT_ALL_ADVERTISEMENT_DATA:
                    return
            log.warning("advertising status=%s, retrying", self.advertising_status)
        log.error("advertising failed, status=%s", self.advertising_status)

    async def stop(self) -> None:
        self._running = False
        provider, self._provider = self._provider, None
        if provider:
            try:
                provider.stop_advertising()
            except Exception:
                pass
        for peer_id in list(self._sessions):
            self.drop_peer(peer_id)
        self._chars.clear()
        log.info("GATT server stopped")

    # ------------------------------------------------------------- events

    def _schedule(self, fn: Callable) -> None:
        loop = self._loop
        if loop is None:
            return
        if threading.get_ident() == self._loop_thread:
            fn()
        else:
            loop.call_soon_threadsafe(fn)

    def _make_write_handler(self, channel: str):
        def handler(sender, args) -> None:
            # deferral / session / request 必须在 WinRT 事件回调内同步取得:
            # 事件返回后再取会抛 E_ILLEGAL_METHOD_CALL (0x8000000E)。
            deferral = None
            try:
                deferral = args.get_deferral()
                session = args.session
                peer_id = session_peer_id(session)
                request_op = args.get_request_async()
            except Exception as exc:
                log.warning("write prepare failed (%s): %s", channel, exc)
                self._release_deferral(deferral)
                return

            async def handle() -> None:
                # 这次写请求是否已经被"应答"过。真机日志里的
                # `respond failed: [WinError -2147483618] The object has been committed.`
                # 就是**重复应答**:`request.value` 这个 getter 会隐式结束
                # deferral(WinRT 的 DataReader 语义),系统随即自动应答一次;
                # 之后我们再去 `respond()` 自然就撞 "already committed"。
                # 那不是错误,只是我们在多此一举 —— 记录状态、别再无脑重试。
                try:
                    request = await request_op
                    if request is None:
                        log.warning("write request empty (%s)", channel)
                        return
                    try:
                        request.respond()  # 先应答,避免解析失败时客户端一直等
                    except Exception as exc:
                        log.debug(
                            "respond skipped (%s, 系统应已自动应答): %s", channel, exc
                        )
                    data = from_buffer(request.value)
                    self._track_session(session, peer_id)
                    self.touch(peer_id)

                    # === 调试:打印手机写过来的原始数据 ===
                    _dump_recv(channel, peer_id, data)

                    cb = self.on_peer_data if channel == "rx" else self.on_peer_ctrl
                    if cb:
                        cb(peer_id, data)
                    # log.debug(
                    #     "recv %s from %s: %d bytes (responded=%s)",
                    #     channel,
                    #     peer_id,
                    #     len(data),
                    #     responded,
                    # )
                except asyncio.CancelledError:
                    raise
                except Exception as exc:
                    log.warning("write handling failed (%s): %s", channel, exc)
                finally:
                    self._release_deferral(deferral)

            try:
                self._schedule(lambda: asyncio.ensure_future(handle()))
            except Exception as exc:
                log.warning("schedule write failed (%s): %s", channel, exc)
                self._release_deferral(deferral)

        return handler

    @staticmethod
    def _release_deferral(deferral) -> None:
        if deferral is None:
            return
        try:
            deferral.complete()
        except Exception as exc:
            log.debug("deferral complete failed: %s", exc)

    @staticmethod
    def _session_alive(session) -> bool:
        """底层会话是否还活着(对象已被关掉时读属性会抛 `RO_E_CLOSED`)。"""
        try:
            return int(session.session_status) != 0  # GattSessionStatus.CLOSED == 0
        except Exception:
            return False

    def _track_session(self, session, peer_id: str) -> None:
        self.touch(peer_id)
        existing = self._sessions.get(peer_id)
        if existing is session:
            return
        if existing is not None and self._session_alive(existing):
            # WinRT 每次读 `args.session` / `sub.session` 都可能返回**新的包装对象**,
            # 底层却是同一条会话。按对象身份判断会让**每一帧**都走到「换绑」分支,
            # 于是每帧 `log.info` 同步写盘 + 每帧 `add_session_status_changed`
            # (token 丢弃、从不反注册)。一次 4MB 传输 ≈ 8500 帧 &#8658; 8500 条日志 +
            # 8500 个泄漏回调,把 GUI 线程拖死 —— dev 反馈的「server 端 UI 卡死、
            # 无报错、对端 write response 超时」就是这么来的。
            # 底层会话还活着 &#8658; 只是换了包装,换引用即可,别动事件、别打日志。
            self._sessions[peer_id] = session
            return
        if existing is not None:
            # 真的是换了一条 GattSession(对端重连):换绑并重新挂钩。
            log.info("peer %s rebound to new gatt session", peer_id)
            self._sessions.pop(peer_id, None)
        self._sessions[peer_id] = session

        def on_status(s, a) -> None:
            def fire() -> None:
                # 事件参数是 GattSessionStatusChangedEventArgs,只有 `status`/`error`
                # 两个属性(`session_status` 是 GattSession 的属性,事件上没有)。
                # 且 CLOSED=0 / ACTIVE=1,所以「关闭」是 == 0。
                try:
                    closed = int(a.status) == 0
                except Exception:
                    closed = True
                if not closed:
                    if self.on_peer_connected:
                        self.on_peer_connected(peer_id)
                    return
                # 身份判断不能比对象(包装每次都新),改成「现在存的这条是否还活着」:
                # 还活着 &#8658; 对端早已重连、这条是旧会话迟到的事件,忽略。
                cur = self._sessions.get(peer_id)
                if cur is not None and self._session_alive(cur):
                    return
                self._sessions.pop(peer_id, None)
                log.info("peer %s disconnected", peer_id)
                if self.on_peer_disconnected:
                    self.on_peer_disconnected(peer_id, "link lost")

            self._schedule(fire)

        try:
            session.add_session_status_changed(on_status)
        except Exception as exc:
            log.debug("session status hook failed: %s", exc)
        if self.on_peer_connected:
            self.on_peer_connected(peer_id)

    def _on_subscribed_changed(self, sender, args) -> None:
        def fire() -> None:
            uuid_str = str(sender.uuid).lower()
            log.debug("subscribed %s -> %d", uuid_str, len(sender.subscribed_clients))
            # 订阅集合按 **TX + CTRL 两条 notify 通道取并集**(握手走 CTRL、数据走 TX)。
            # 以前这里只"加"不"减":对端软件下线/重启后 `_sessions` 里留下一条僵尸会话,
            # keepalive PING 要等 20s 才发现 `not subscribed` —— 这 20s 里双方状态不一致,
            # 界面上就是 dev 报的"掉线了却还报一堆 `UNKNOWN:peer … not subscribed …`"。
            sessions: dict[str, object] = {}
            duplicates: set[str] = set()
            for key in (CHAR_TX, CHAR_CTRL):
                char = self._chars.get(key)
                if char is None:
                    continue
                try:
                    subs = list(char.subscribed_clients or ())
                except Exception as exc:
                    log.debug("subscribed_clients unavailable on %s: %s", key[:8], exc)
                    continue
                for sub in subs:
                    try:
                        peer = session_peer_id(sub.session)
                    except Exception as exc:
                        log.debug("skip subscription: %s", exc)
                        continue
                    if peer in sessions:
                        # 同一个 peer 出现两条订阅:WinRT 的订阅表在"对端下线又
                        # 重新连上"之后确实会留下幽灵条目。这是排查"单方已连接"
                        # 的第一手证据,所以留一条 DEBUG。
                        duplicates.add(peer)
                    else:
                        sessions[peer] = sub.session
            if duplicates:
                log.debug(
                    "duplicate subscriptions for %s on %s+%s",
                    sorted(duplicates),
                    CHAR_TX[:8],
                    CHAR_CTRL[:8],
                )
            for peer_id, session in sessions.items():
                self._track_session(session, peer_id)
            if uuid_str not in (CHAR_TX, CHAR_CTRL):
                return  # RX 是写通道,没有 notify 订阅,不参与掉线判定
            for peer_id in list(self._sessions):
                if peer_id in sessions:
                    continue
                if not self._session_alive(self._sessions[peer_id]):
                    self._sessions.pop(peer_id, None)  # 底层已关,`on_status` 可能报过
                    continue
                self._sessions.pop(peer_id, None)
                log.info("peer %s unsubscribed from all notify channels", peer_id)
                if self.on_peer_disconnected:
                    self.on_peer_disconnected(peer_id, "unsubscribed")

        self._schedule(fire)

    def _on_adv_status(self, sender, args) -> None:
        def fire() -> None:
            try:
                status = GattServiceProviderAdvertisementStatus(args.status).name
            except Exception:
                status = str(args.status)
            log.info("advertising status -> %s", status)

        self._schedule(fire)

    # ------------------------------------------------------------- outbound

    def _subscribed(self, char_uuid: str, peer_id: str):
        """找出该对端在这条通道上的订阅,**优先挑底层会话还活着的**。

        对端软件直接下线时,`subscribed_clients` 里常常还留着幽灵条目,
        同时新连接又登记了一条 —— 老代码取"第一个匹配",可能正好挑中幽灵,
        于是 `notify_value_for_subscribed_client_async` 静默失败/抛 PeerGone,
        上层看到的就是「单方已连接」。
        """
        char = self._chars.get(char_uuid)
        if not char:
            raise PeerGone("no characteristic")
        fallback = None
        for sub in char.subscribed_clients:
            if session_peer_id(sub.session) != peer_id:
                continue
            if self._session_alive(sub.session):
                return char, sub
            if fallback is None:
                fallback = sub
        if fallback is not None:
            log.debug("subscribed(%s): only a dead session for %s", char_uuid[:8], peer_id)
            return char, fallback
        raise PeerGone(f"peer {peer_id} not subscribed to {char_uuid}")

    async def send_ctrl(self, peer_id: str, data: bytes) -> None:
        char, sub = self._subscribed(CHAR_CTRL, peer_id)
        self._check_notify_size(sub, data, peer_id, "ctrl")
        try:
            await char.notify_value_for_subscribed_client_async(to_buffer(data), sub)
        except OSError as exc:
            if is_link_dead(exc):
                raise PeerGone(f"peer {peer_id} session closed") from exc
            raise

    async def send_data(self, peer_id: str, data: bytes) -> None:
        char, sub = self._subscribed(CHAR_TX, peer_id)
        self._check_notify_size(sub, data, peer_id, "data")
        try:
            await char.notify_value_for_subscribed_client_async(to_buffer(data), sub)
        except OSError as exc:
            if is_link_dead(exc):
                raise PeerGone(f"peer {peer_id} session closed") from exc
            raise

    @staticmethod
    def _check_notify_size(sub, data: bytes, peer_id: str, channel: str) -> None:
        try:
            limit = int(sub.max_notification_size)
        except Exception:
            return
        if limit and len(data) > limit:
            log.warning(
                "%s frame %d bytes exceeds max_notification_size %d (peer=%s)",
                channel,
                len(data),
                limit,
                peer_id,
            )

    def peer_mtu(self, peer_id: str) -> int:
        for char in self._chars.values():
            for sub in char.subscribed_clients:
                if session_peer_id(sub.session) == peer_id:
                    try:
                        return max(23, int(sub.session.max_pdu_size))
                    except Exception:
                        return 23
        return 23

    def peer_chunk(self, peer_id: str) -> int:
        mtu = self.peer_mtu(peer_id)
        for char in self._chars.values():
            for sub in char.subscribed_clients:
                if session_peer_id(sub.session) == peer_id:
                    try:
                        size = int(sub.max_notification_size)
                    except Exception:
                        size = 0
                    if size <= 0:
                        size = max(23, mtu - 3)
                    return max(8, min(size, mtu) - 3 - DATA_HEADER)
        return 8

    def list_peers(self) -> list[str]:
        return list(self._sessions)

    def drop_peer(self, peer_id: str) -> None:
        session = self._sessions.pop(peer_id, None)
        if session is None:
            return
        try:
            session.close()
        except Exception as exc:
            log.debug("close session %s failed: %s", peer_id, exc)


def _fmt_addr(value: int) -> str:
    if not value:
        return ""
    return ":".join(f"{(value >> (8 * i)) & 0xFF:02X}" for i in reversed(range(6)))


async def adapter_address() -> str:
    """本机蓝牙 MAC(用于直连/展示)。"""
    try:
        from winrt.windows.devices.bluetooth import BluetoothAdapter

        adapter = await BluetoothAdapter.get_default_async()
        if adapter is None:
            return ""
        return _fmt_addr(int(adapter.bluetooth_address))
    except Exception as exc:
        log.debug("adapter address failed: %s", exc)
        return ""


__all__ = [
    "GattServer",
    "PeerGone",
    "is_closed_error",
    "is_link_dead",
    "to_buffer",
    "from_buffer",
    "adapter_address",
    "session_peer_id",
]




发帖前要善用【论坛搜索】功能,那里可能会有你要找的答案或者已经有人发布过相同内容了,请勿重复发帖。

沙发
追逐飞翔 发表于 2026-10-6 21:47
我也需要,可以分享下成品吗?
您需要登录后才可以回帖 登录 | 注册[Register]

本版积分规则

返回列表

RSS订阅|小黑屋|处罚记录|联系我们|吾爱破解 - 52pojie.cn ( 京ICP备16042023号 | 京公网安备 11010502030087号 )

GMT+8, 2026-10-7 02:58

Powered by Discuz!

Copyright © 2001-2020, Tencent Cloud.

快速回复 返回顶部 返回列表