bisheng LDAP 组织定时校对与 SSO 换型 Relink:基于 Celery 的 6h 自动同步与冲突治理实战
【免费下载链接】bishengBISHENG is an open LLM devops platform for next generation Enterprise AI applications. Powerful and comprehensive features include: GenAI workflow, RAG, Agent, Unified model management, Evaluation, SFT, Dataset Management, Enterprise-level System Management, Observability and more.项目地址: https://gitcode.com/GitHub_Trending/bi/bisheng
导读
本文围绕 bisheng 开源 LLM DevOps 平台在 v2.5.1 中落地的F015 特性(LDAP 组织定时校对 + SSO 换型 relink),系统讲解如何用 Celery Beat 每 6 小时强制校对 SSO/LDAP 部门树与 bisheng 内部组织结构,自动修复漂移,并在 SSO 更换 HR 系统后通过 relink 接口重建 external_id 映射。读完本文你将掌握:org_sync_log事件化改造与复合索引设计、基于 Redis SETNX 的并发幂等锁、last_sync_ts时间戳仲裁的冲突决策表(INV-T12)、ts 冲突周报与每日升级告警链路,以及 external_id_map / path_plus_name 两种 relink 匹配策略的 HTTP 调用方式。文中所有结论均有仓库源码与测试用例佐证,可直接对照实现进行二次开发或运维排障。
1. 背景:为什么需要“定时校对 + 换型 relink”
1.1 集团 IT 的用户故事
F015 对应的用户故事来自多租户需求文档(spec.md §1):
作为集团 IT,我希望系统每 6 小时自动校对 SSO 端部门树与 bisheng 内部结构,发现漂移时自动修复;SSO 换 HR 系统时提供 relink 接口重建映射,以便部门归属长期保持一致,换型不丢失用户归属。
在真实的企业环境里,SSO 侧的部门树会不断变化(新增、改名、删除、跨部门移动),而 Gateway 实时同步只能覆盖“在线发生”的变更;一旦出现断网、Provider 短暂不可用或批量历史数据导入,bisheng 内部结构与 SSO 就会产生漂移。F015 的价值在于:
- 兜底:用 Celery Beat 每 6 小时做一次全量拉取式校对(spec.md AD-01 决策为 6h,非 1h/24h);
- 自动修复:新增/命名变更自动 upsert,删除进入 orphaned(孤儿子部门并标记),跨租户移动自动触发用户租户重算;
- 换型逃生门:SSO 从旧 HR 系统迁到新系统时 external_id 全部变化,通过 relink 接口用路径+名称或显式映射重建关联,避免用户归属丢失。
1.2 与既有特性的关系
F015 不是从零开发,而是站在 v2.5.1 一系列前置特性之上(tasks.md「前置」节):
| 前置能力 | 来源 | F015 中的用途 |
|---|---|---|
OrgSyncTsGuard | F014 T04 | 每 op 的 ts 决策(INV-T12) |
DeptUpsertService | F014 T07 | upsert 落库 |
_OrgSyncLogBuffer+flush_log | F014 T12 | 批摘要日志 |
verify_hmac | F014 T05 | relink HTTP 签名校验 |
DepartmentDeletionHandler.on_deleted(dept_id, deletion_source) | F011 | 归档后孤儿租户处理 |
UserTenantSyncService.sync_user(user_id, trigger) | F012 | 跨租户移动后的用户租户重算 |
DepartmentDao.aget_by_source_external_id / aupsert_by_external_id / aarchive_by_external_id | F014 T03 | 部门 DAO 基础操作 |
AuditLogDao.ainsert_v2 | F011 | 冲突/relink 审计留痕 |
DeletionSource.CELERY_RECONCILE/UserTenantSyncTrigger.CELERY_RECONCILE | F011 constants | 标注操作来源 |
这些前置实现的当前仓库位置分别为 org_sync/domain/services/ts_guard.py、sso_sync/domain/services/dept_upsert_service.py、sso_sync/domain/services/org_sync_log_writer.py 等。
2. 总体架构与任务拆解
2.1 决策锁定(D1–D8)
tasks.md「决策锁定」明确了 8 项架构决策,是理解代码布局的钥匙:
- D1:Service 放
bisheng/org_sync/domain/services/;API endpoint 放bisheng/org_sync/api/endpoints/relink.py(F009 扩展); - D2:6h Beat 用独立
reconcile_all_organizations,不复用 F009 的check_org_sync_schedules(后者响应用户 cron,F015 是系统强制); - D3:
org_sync_log表加 4 列 + 复合索引,不拆事件表——摘要行event_type='',事件行非空; - D4:relink 多候选冲突存 Redis Hash
relink_conflict:{dept_id},TTL 7 天; - D5:前端 UI 不纳入,归 F019 admin-scope console;
- D6:细粒度 15 个任务(对齐 F014 的 Test-Alongside 模式);
- D7:独立 worktree 开发;
- D8:错误码 19314–19318 写入
errcode/sso_sync.py(MMM=193 与 F014 共享)。
从当前仓库的实际文件布局看,这些决策已全部落地:org_sync/domain/services/下存在reconcile_service.py、remote_dept_differ.py、relink_service.py、relink_conflict_store.py、ts_guard.py、reconciler.py;org_sync/api/endpoints/relink.py存在;错误码定义于 common/errcode/sso_sync.py。
2.2 依赖图与并行建议
tasks.md 给出了完整依赖图,关键路径如下:
T01 (errcode + ReconcileConf) ├─→ T02 (Alembic 迁移) → T03 (OrgSyncLog ORM+DAO) → T04 (aget_all_active) → T06 (主编排) │ └─→ T10 (TsConflictReporter) → T13 (weekly/daily beats) ├─→ T05 (RemoteDeptDiffer 纯函数) → T06 └─→ T07 (Relink schemas+service) → T08 (Relink API) ─┐ └─→ T12 (RelinkConflictStore) T09 (Celery 6h fan-out) / T11 (event 持久化) / T14 (SETNX 并发幂等) T01-T14 ──→ T15 (AC 对照 + e2e + 压测占位)并行建议:T01 后三路并行 — (T02→T03→T04→T06) / (T05→T07→T08→T12) / (T10);T09/T11/T13/T14 收尾;T15 统一验收。这种拆法让纯函数(T05 diff)、无 IO 的决策(T01 常量)可以最早启动,数据层与编排层则串行推进。
2.3 开发模式约定(Test-Alongside)
- 单任务 = 实现代码 + 单测/集成测共 2–4 文件;
- Celery 任务用直接调
.apply()同步触发或 mockapply_async; - Redis SETNX/ZSET/Hash 用
mock_redisfixture; - Provider
fetch_departments用MagicMock()返回 DTO 列表; - 迁移通过 MySQL 手工
alembic upgrade head/downgrade -1往返 + SQLitetable_definitions.py同步; - 冲突告警复用 F011
send_inbox_notice+list_global_super_admin_ids。
3. 配置与错误码(T01)
3.1 ReconcileConf 配置模型
ReconcileConf定义在 core/config/reconcile.py(tasks.md 中字段更全,仓库当前实现保留了锁与冲突 TTL 两项,其余字段按设计在 settings 注册):
class ReconcileConf(BaseModel): beat_cron_reconcile: str = '0 */6 * * *' # 6h 校对 beat_cron_weekly_report: str = '0 9 * * MON' # 周一 09:00 周报 beat_cron_daily_escalation: str = '0 9 * * *' # 每日 09:00 升级 redis_lock_ttl_seconds: int = 1800 # 校对锁 TTL(30min) relink_conflict_ttl_seconds: int = 604800 # relink 冲突 7d weekly_conflict_threshold: int = 3 # 周报冲突阈值 daily_escalation_days: int = 5 # 5 天未解决升级 task_time_limit: int = 1800 # Celery 硬超时 task_soft_time_limit: int = 1500 # Celery 软超时仓库实现中redis_lock_ttl_seconds的注释解释了设计意图:30 分钟与 Celery 软超时对齐,保证卡死的 worker 不会永久持有锁;relink_conflict_ttl_seconds要求管理员必须在 7 天内通过resolve-conflict接口解决冲突。这些配置经Settings注册为reconcile子配置,运行时通过settings.reconcile.xxx访问(见 reconcile_service.py 中settings.reconcile.redis_lock_ttl_seconds的用法)。
3.2 错误码 19314–19318
F015 的错误码全部追加在 common/errcode/sso_sync.py(MMM=193 与 F014 共享模块,19310–19313 已被 F014 占用):
| Code | 错误类 | 含义 |
|---|---|---|
| 19314 | SsoReconcileLockBusyError | 同一 config 的校对正在执行,Redis SETNX 锁忙,跳过本次触发 |
| 19315 | SsoRelinkStrategyUnsupportedError | relinkmatching_strategy不是external_id_map/path_plus_name |
| 19316 | SsoRelinkConflictUnresolvedError | 冲突候选列表为空,或所选new_external_id不在已存候选内 |
| 19317 | SsoSameTsRemoveAppliedWarnError | 同 ts 的 upsert/remove 冲突已按 remove 为准应用(告警类,不抛出) |
| 19318 | SsoReconcileReservedError | 预留 |
仓库源码与 tasks.md 的编号完全一致。特别值得注意的是 19317 在 reconcile_service.py 中只记日志不抛出——管理员通过org_sync_log事件行和audit_log观察,而不是走 HTTP 错误面。
4. 数据层改造(T02/T03/T04)
4.1 Alembic 迁移:org_sync_log 加 4 列 + 复合索引
T02 是数据层的地基,迁移文件v2_5_1_f015_reconcile_log_fields.py复用 F014 的_column_exists/_index_exists幂等 helper,为org_sync_log表追加:
| 列 | 类型 | 默认值 | 用途 |
|---|---|---|---|
event_type | String(32) | '' | 事件类型:空=批摘要行 /ts_conflict/stale_ts/conflict_weekly_sent/conflict_daily_escalation_sent |
level | String(16) | 'info' | 日志级别:info / warn / error |
external_id | String(128) | NULL | 事件行关联的部门 external_id |
source_ts | BigInteger | NULL | INV-T12 捕获的 incoming ts(审计用) |
同时创建复合索引:
CREATE INDEX idx_conflict_lookup ON org_sync_log (level, event_type, external_id, create_time);这个索引直接服务于 5.5.3 节的冲突计数查询(level='warn' AND event_type='ts_conflict' AND external_id=? AND create_time > now-7d)。spec.md §5.5.3 还给出了兼容补丁 SQL(ALTER TABLE ... ADD INDEX / ADD COLUMN),用于 v2.5.0/F009 已建表但缺字段/索引的升级场景。SQLite 侧的test/fixtures/table_definitions.py同步更新,T03 的 SQLite 测试绿灯即为同步佐证;MySQL 侧做alembic upgrade head && downgrade -1 && upgrade head往返验证。
4.2 OrgSyncLog ORM 与 DAO 扩展
OrgSyncLogORM 追加上述 4 列后,OrgSyncLogDao新增 3 个 classmethod(tasks.md T03):
acreate_event(event_type, level, external_id, source_ts, config_id, error_details=None, tenant_id=1)— 事件行快捷构造,摘要计数器全 0;acount_recent_conflicts(external_id, days=7) -> int— 利用idx_conflict_lookup统计冲突次数;aget_conflicts_since(since, event_type='ts_conflict', level='warn') -> list[OrgSyncLog]— 聚合输入,Service 层按 external_id group。
T03 的 6 条测试覆盖了事件行持久化、摘要行计数器保持为 0、按 external_id 与时间窗过滤、窗口外返回 0、按序返回,以及摘要行(event_type='')与事件行(event_type='ts_conflict')共存可分别查出。
4.3 OrgSyncConfigDao.aget_all_active
T04 为 6h fan-out 提供入口:aget_all_active扫描所有租户的 active 配置。与 F009aget_active_cron_configs的关键区别是不按 schedule_type 过滤——F015 是系统强制校对,与用户配置的 cron 无关。调用方负责过滤provider == 'sso_realtime'(F014 seed id=9999,仅用于 HMAC 实时日志,不经过fetch_departments)。
5. Diff 引擎:RemoteDeptDiffer(T05)
5.1 输出 DTO
RemoteDeptDiffer定义于 org_sync/domain/services/remote_dept_differ.py,输出四类 dataclass:
@dataclass class UpsertOp: # CREATE / UPDATE 语义 external_id: str; name: str; parent_external_id: Optional[str] sort_order: int; incoming_ts: int; is_new: bool existing_dept_id: Optional[int] = None # UPDATE 时带上已有 id @dataclass class ArchiveOp: # 远端不再列出 → 软归档 external_id: str; dept_id: int; mounted_tenant_id: Optional[int]; incoming_ts: int @dataclass class MoveOp: # 父节点变更 + 叶子租户影响标志 external_id: str; dept_id: int; new_parent_external_id: Optional[str] crosses_tenant: bool; incoming_ts: int @dataclass class ReconcileDiff: upserts: list[UpsertOp]; archives: list[ArchiveOp]; moves: list[MoveOp]5.2 纯函数设计
diff(remote_depts, local_depts, source, ts)是无 IO 的纯函数:它包装 F009 的reconcile_departments(见 org_sync/domain/services/reconciler.py),并把拓扑排序保证(upsert 父先于子、archive 子先于父)继承下来,同时给每个 op 注入incoming_ts与crosses_tenant。
crosses_tenant的计算是叶子租户推导:沿 local 父链向上找最近的is_tenant_root=1节点,取其mounted_tenant_id作为叶子租户;移动前后叶子租户不同则crosses_tenant=True。仓库实现中_derive_leaf_tenant_id还带环路防御(visited 集合)与“找不到挂载点视为 Root 租户”的兜底,属于对 pathological 数据的防御性处理。
6. 主编排服务:OrgReconcileService(T06)
6.1 11 步主流程
org_sync/domain/services/reconcile_service.py 的reconcile_config(config_id)是 F015 的心脏,完整流程如下:
1) 加载 config;provider=='sso_realtime' 或 status!='active' → skipped 2) 获取 Redis SETNX 锁 org_reconcile:{config_id}(TTL=redis_lock_ttl_seconds) 3) Provider.authenticate() + fetch_departments(sync_scope.root_dept_ids) 4) local = DepartmentDao.aget_active_by_tenant(config.tenant_id) 5) diff = RemoteDeptDiffer.diff(remote, local, source, ts=now) 6) Upsert 循环:逐 op 走 OrgSyncTsGuard;SKIP_TS → stale_ts 事件行 + warn;APPLY → DeptUpsertService 7) Archive 循环:逐 op 走 Guard;同 ts 冲突(remove wins)→ ts_conflict 事件行 + audit;archive + DeletionHandler 8) 跨租户移动 → 主部门成员逐个 UserTenantSyncService.sync_user(uid, CELERY_RECONCILE) 9) flush_log 批摘要 10) 逐条持久化 event_rows(acreate_event,单条失败不影响其他) 11) 返回 ReconcileResult 供 Celery 记录仓库实现比 tasks.md 骨架更完善的地方包括:租户上下文显式管理——Celery 任务入口没有 HTTP 中间件,服务在bypass_tenant_filter()下设置current_tenant_id,避免 SQLAlchemy tenant_filter 事件报 20004 缺租户上下文;逐 op try/except——单个部门 upsert/archive 失败不中断整轮校对,失败进入result.errors;逐事件行容错——event row 持久化单条失败只记日志。
6.2 三个关键子流程
Upsert(AC-02):对每个UpsertOp查DepartmentDao.aget_by_source_external_id,OrgSyncTsGuard.check_and_update(existing, incoming_ts, 'upsert')决策。APPLY 时通过DeptUpsertService.upsert_from_sync_payload落库。is_new决定 buffer 计数是dept_created还是dept_updated,与 F009 批摘要语义对齐,保证管理端历史面板数字符合预期。mount 标记(is_tenant_root)在 upsert 中不被修改——这正是 AC-02 “新部门 upsert 不动挂载标记”的落点。
Archive(AC-03 / AC-11):对每个ArchiveOp同样走 Guard。同 ts 冲突的检测条件是last_sync_ts == incoming_ts且is_deleted == 0——说明本批 upsert 已应用但随后又收到同 ts 的 remove,此时:
- 写
ts_conflict事件行(error_details.resolution='remove_wins'); - 写
audit_log.action='dept.sync_conflict'(审计失败不中断操作); - 记 19317 告警日志;
DepartmentDao.aarchive_by_external_id落删除;DepartmentArchiveCleanupService.arun_for_archived_department清理;DepartmentDeletionHandler.on_deleted(dept_id, DeletionSource.CELERY_RECONCILE)统一处理孤儿租户(INV-T8:mount 部门归档后Tenant.status='orphaned'并告警)。
跨租户移动(AC-04 / INV-T2):crosses_tenant=True的移动,对部门所有主部门成员(is_primary=True)逐个调用UserTenantSyncService.sync_user(uid, CELERY_RECONCILE),触发token_version +1使 JWT 的tenant_id失效刷新。
6.3 Redis SETNX 锁(AC-13)
redis = await get_redis_client() key = f'org_reconcile:{config_id}' ok = await redis.async_connection.set(key, b'1', nx=True, ex=ex) if not ok: raise SsoReconcileLockBusyError.http_exception()锁在finally中释放(删除失败靠 TTL 兜底)。tasks.md 的 T14 专项测试覆盖:同 config 并发一成一 19314、不同 config 锁相互独立、异常时锁仍释放、SETNX 调用 ex=1800。
7. Relink 子系统(T07/T08/T12)
7.1 使用场景
SSO 换型(更换 HR 系统)时,新系统为部门生成了全新的external_id,旧 id 全部失效。如果不处理,下次实时同步会把所有旧部门判定为“已删除”并归档,用户归属随之丢失。relink 的作用就是在迁移窗口内把旧 external_id 重新指向新 id。
7.2 两种匹配策略
org_sync/domain/services/relink_service.py 实现两种策略:
external_id_map:运维提供old_ext -> new_ext显式映射字典,逐条查 dept(by source+old_ext)后重写external_id;path_plus_name:对每个 old_ext,在同 source 且未被占用的 active 部门中按(path, name)精确匹配找候选;单候选自动 apply,多候选存入RelinkConflictStore并返回 conflicts 列表,由管理员人工确认(spec.md AD-02 决策:人工确认避免误操作)。
策略不识别时抛SsoRelinkStrategyUnsupportedError(19315)。每次实际重写都写审计action='dept.relink_applied'(自动路径)或'dept.relink_resolved'(人工解决路径),metadata 含{old_ext, new_ext, strategy},满足 INV-T7 的“挂载/解绑强制 audit”。
7.3 dry_run 预演
dry_run=True时只收集would_apply清单,不写 DB、不写冲突存储,用于迁移前验证匹配结果。这是 AC-07 的核心,也是运维上线前必做的安全检查。
7.4 多候选冲突存储
org_sync/domain/services/relink_conflict_store.py 用 Redis Hash 存储:relink_conflict:{dept_id}的每个 field 是候选new_external_id,value 是 JSON(含 path/name/score)。TTL 7 天,超时自动过期,忘掉的冲突重跑 relink 即可重建。仓库实现特意在save前先delete保证干净替换,并在get中对 bytes/str 双客户端形态做兼容、对损坏 JSON 容错。模块 docstring 解释了“为什么用 Redis 不用表”:冲突是短生命周期的临时数据、丢失可接受、避免引入迁移/模型/清理 cron。
7.5 resolve-conflict
管理员从候选列表中选择一个new_external_id提交:若所选不在候选内(或候选已过期),抛SsoRelinkConflictUnresolvedError(19316)且保留存储;合法则重写 external_id + 写dept.relink_resolved审计 + 删除存储条目。
8. HTTP API 层(T08)
org_sync/api/endpoints/relink.py 提供两个 HMAC 签名的内部端点,挂在router.py,并在utils/http_middleware.py的TENANT_CHECK_EXEMPT_PATHS中登记以绕过租户上下文中间件:
POST /api/v1/internal/departments/relink Body: { "old_external_ids": ["abc", ...], "matching_strategy": "external_id_map" | "path_plus_name", "external_id_map": {"abc": "new_abc", ...}, // external_id_map 策略必填 "source": "sso", "dry_run": false } Returns: {"applied": [...], "would_apply": [...], "conflicts": [...]} POST /api/v1/internal/departments/relink/resolve-conflict Body: {"dept_id": 123, "chosen_new_external_id": "new_abc"} Returns: {"dept_id": 123, "old_external_id": "abc", "new_external_id": "new_abc"}两个端点都依赖 F014 的verify_hmac做请求签名校验(无签名 401),Service 层在 ROOT_TENANT_ID +bypass_tenant_filter下运行,与 F014 LoginSyncService 的模式一致。
9. Celery 任务:6h 校对 + 冲突告警(T09/T13)
9.1 6h 校对 fan-out
worker/org_sync/reconcile_tasks.py 定义两个任务:
@bisheng_celery.task(acks_late=True) def reconcile_all_organizations(): """6h Beat entry: fan out reconcile per active OrgSyncConfig.""" loop = asyncio.new_event_loop() try: loop.run_until_complete(_fan_out_all()) finally: loop.close() async def _fan_out_all() -> None: configs = await OrgSyncConfigDao.aget_all_active() for c in configs: if c.provider == 'sso_realtime': continue reconcile_single_config.apply_async(args=[c.id], queue='knowledge_celery') @bisheng_celery.task(acks_late=True, time_limit=1800, soft_time_limit=1500) def reconcile_single_config(config_id: int): """Execute one reconcile run; swallow lock-busy without retry.""" loop = asyncio.new_event_loop() try: loop.run_until_complete(OrgReconcileService.reconcile_config(config_id)) except SsoReconcileLockBusyError: logger.warning(f'reconcile_single_config {config_id} skipped: lock busy') except Exception: logger.exception(f'reconcile_single_config {config_id} failed') finally: loop.close()两个关键设计:锁忙不重试(SsoReconcileLockBusyError只记 warning,避免重试风暴叠加);硬超时 1800s 与 Redis 锁 TTL 一致(卡死 worker 也不会永久持锁)。子任务进knowledge_celery队列。Beat 注册在CeleryConf.validate中追加:
if 'reconcile_all_organizations' not in self.beat_schedule: self.beat_schedule['reconcile_all_organizations'] = { 'task': 'bisheng.worker.org_sync.reconcile_tasks.reconcile_all_organizations', 'schedule': crontab.from_string('0 */6 * * *'), # every 6h }9.2 冲突周报与每日升级
TsConflictReporter(org_sync/domain/services/ts_conflict_reporter.py)提供两个告警方法,由两条 Beat 独立触发:
| 任务 | cron | 逻辑 |
|---|---|---|
report_ts_conflicts_weekly | 0 9 * * MON(周一 09:00) | 聚合过去 7 天ts_conflict事件按 external_id 计数,count >= weekly_conflict_threshold(3)的条目发全局超管站内消息(payload 含冲突详情 + “是否需要 relink”建议 + 本周总冲突数),并写conflict_weekly_sent标记行 |
report_ts_conflicts_daily_escalation | 0 9 * * *(每日 09:00) | 若最近一次周报标记已超过daily_escalation_days(5)天且冲突仍存在,升级为每日告警;全部解决则reason='resolved'不升级 |
升级状态机的三种退出路径值得注意:无周报标记(no_weekly_marker)、仍在宽限期内(within_grace)、冲突已解决(resolved)。告警复用 F011 的send_inbox_notice+list_global_super_admin_ids,对应的conflict_weekly_sent/conflict_daily_escalation_sent标记行也写入org_sync_log,使告警行为本身可审计。
10. 冲突决策核心:OrgSyncTsGuard 与 INV-T12
10.1 决策表
org_sync/domain/services/ts_guard.py 是 F014(Gateway 实时)与 F015(Celery 校对)共享的纯决策函数,把 spec.md §5.5.2 的决策表落地:
| incoming ts vs last_sync_ts | 动作 |
|---|---|
incoming_ts > last_sync_ts | 应用变更+ 更新last_sync_ts(AC-10) |
incoming_ts == last_sync_ts,同 source 同方向 | 幂等跳过(已应用) |
incoming_ts == last_sync_ts,upsert + remove 双向 | 以 remove 为准(AC-11,从严避免幽灵部门),remove 已应用时后续 upsert 观察is_deleted=1而 SKIP_TS |
incoming_ts < last_sync_ts | 跳过+ 写org_sync_loglevel=warn(陈旧消息丢弃,AC-09) |
10.2 纯函数实现
Guard 不执行任何 I/O,只有两个输出:APPLY/SKIP_TS。对从未见过的 external_id:upsert 放行、remove 静默丢弃(没有历史可保护)。决策与写入分离的设计让 8 种组合无需 DB 往返即可单测。“同 ts remove wins”不变量的成立依赖一个关键事实:第一个写入方把is_deleted=1落库,后续同 ts 的 upsert 在 Guard 中读到该标记即判 SKIP_TS——这与 F015 同批内“先 upsert 后 archive”的冲突(由reconcile_service的is_same_ts_conflict检测并审计)形成互补,两层共同保证 INV-T12。
10.3 不变量映射小结
- INV-T12(ts 最大为准 + 同 ts remove 优先):T05 diff 注入
incoming_ts→ T06 逐 op 走 Guard → 同 ts 冲突走 AC-11(remove wins + auditdept.sync_conflict+ 19317)→ T11 持久化 → T14 并发兜底; - INV-T8(孤儿 Tenant):T06 归档 mount 部门后调
DepartmentDeletionHandler.on_deleted(dept_id, CELERY_RECONCILE)→ F011 将Tenant.status='orphaned'并告警; - INV-T7(挂载/解绑强制 audit):T07
resolve_conflict写dept.relink_resolved;T06 同 ts 冲突写dept.sync_conflict; - INV-T2(用户唯一叶子):T06
crosses_tenant=True主动sync_user(uid, CELERY_RECONCILE)→token_version +1。
11. 验收标准矩阵(AC-01 ~ AC-13)
| AC | 验收点 | 关键测试(仓库中已存在) |
|---|---|---|
| AC-01 | 每 6h 执行 SSO 全量校对 | test_beat_schedule_registers_reconcile_all_every_6h、test_reconcile_all_dispatches_single_config_per_active_config(test/celery/test_reconcile_celery_tasks.py) |
| AC-02 | 新部门自动 upsert 且不动挂载标记 | test_reconcile_new_dept_upserts_preserves_mount |
| AC-03 | 删除部门标记 is_deleted、挂载点 orphaned | test_reconcile_removed_dept_triggers_department_deletion_handler |
| AC-04 | 主部门跨 Tenant 变更触发 UserTenantSyncService | test_reconcile_primary_dept_change_triggers_user_tenant_sync |
| AC-05 | external_id_map 策略应用成功 | test_relink_external_id_map_strategy_applies(test/department/test_department_relink_service.py) |
| AC-06 | path_plus_name 单候选自动 apply、多候选 conflicts | test_relink_path_plus_name_single_candidate_auto_apply、test_relink_path_plus_name_multi_candidate_returns_conflicts |
| AC-07 | dry_run 返回 would_apply 不写入 | test_relink_dry_run_returns_would_apply_no_db_write+ HTTP dry_run(test/department/test_relink_api_integration.py) |
| AC-08 | 10 万部门校对 < 30 min | locust 压测占位(scripts/performance/locust_ldap_reconcile_100k.py,发版前专项,不在 CI) |
| AC-09 | 陈旧 ts 跳过 + warn 事件行 | test_reconcile_stale_ts_skipped_writes_warn_event |
| AC-10 | 更新 ts 应用并更新 last_sync_ts | test_reconcile_newer_ts_applies_and_updates_last_sync_ts |
| AC-11 | 同 ts upsert/remove 以 remove 为准 + audit + 19317 | test_reconcile_same_ts_upsert_then_remove_prefers_remove |
| AC-12 | 周报 ≥3 次冲突告警 + 5 天升级每日 | test_weekly_report_aggregates_conflicts_above_threshold、test_daily_escalation_triggers_after_5_days_unresolved |
| AC-13 | 并发同 external_id 同 ts 去重幂等 | test_concurrent_same_ts_same_config_deduped_by_setnx(test/org_sync/test_org_reconcile_service.py) |
上表的关键测试均已实际存在于当前仓库的测试目录中(如 test/org_sync/test_org_reconcile_service.py 包含全部 16 条 reconcile 相关用例,涵盖 AC-02/03/04/09/10/11/13 与 event 持久化、锁独立性、异常释放、TTL 断言等),可用pytest test/ -k "org_sync or tenant or sso or reconcile or relink"一键回归。
12. 开发与运维命令速查
# 单任务跑测 .venv/bin/pytest test/test_org_reconcile_service.py -v .venv/bin/pytest test/test_reconcile_celery_tasks.py -v .venv/bin/pytest test/test_ts_conflict_reporter.py -v .venv/bin/pytest test/test_department_relink_service.py -v .venv/bin/pytest test/test_relink_conflict_store.py -v .venv/bin/pytest test/test_remote_dept_differ.py -v .venv/bin/pytest test/test_org_sync_log_dao_f015.py -v # 迁移往返验证(T02) .venv/bin/alembic upgrade head && .venv/bin/alembic downgrade -1 && .venv/bin/alembic upgrade head # 全 feature 回归 .venv/bin/pytest test/ -k "org_sync or tenant or sso or reconcile or relink" -v # Worker + Beat 手工 QA .venv/bin/celery -A bisheng.worker.main:bisheng_celery worker -Q knowledge_celery -l info & .venv/bin/celery -A bisheng.worker.main:bisheng_celery beat -l info & # 手动触发 6h 校对(不等待 Beat) .venv/bin/python -c "from bisheng.worker.org_sync.reconcile_tasks import reconcile_all_organizations; reconcile_all_organizations.delay()" # relink HTTP 调用(HMAC 签名,X-Signature 用共享 secret 对 # 'POST\n/api/v1/internal/departments/relink\n<body>' 做 SHA256 得到) curl -X POST http://localhost:7860/api/v1/internal/departments/relink \ -H "X-Signature: <sha256>" -H "Content-Type: application/json" \ -d '{"old_external_ids":["abc"], "matching_strategy":"path_plus_name", "dry_run":true}'注意:以上命令中的相对路径均相对src/backend目录执行;reconcile_all_organizations.delay()是运维验证 6h 校对链路的快捷入口,正常生产环境依赖 Celery Beat 的0 */6 * * *调度。
13. 边界情况与运维要点
- 校对期间部门变更:Redis SETNX 锁避免并发同步(AC-13)。若锁忙,Celery 任务只记 warning 不重试,下一轮 6h 自动补偿。
- relink 冲突长期未解决:候选在 Redis 中 7 天过期;若 5 天内未通过
resolve-conflict解决,周报升级为每日告警,直到冲突消失或人工处理。 - SSO 端暂时不可用:
authenticate或fetch_departments失败时整轮跳过(result.errors记录),不写任何变更,下轮重试;审计失败、事件行持久化失败均为单点容错,不中断整轮校对。 - 陈旧消息:任何
incoming_ts < last_sync_ts的消息(无论来自 Gateway 还是 Celery)一律丢弃并写 warn 事件行,不覆盖 bisheng 当前状态(AC-09)。 - 性能专项:AC-08 的 10 万部门压测(wall time < 30 min、MySQL/Redis 压力基线)是发版前 2 周的专项工作,不在 CI 范围,基线数据记录于
features/v2.5.1/015-ldap-reconcile-celery/ac-verification.md。
14. 扩展阅读
- 特性规格与验收标准:features/v2.5.1/015-ldap-reconcile-celery/spec.md
- 任务拆解与依赖图(本文骨架来源):features/v2.5.1/015-ldap-reconcile-celery/tasks.md
- 主编排服务:src/backend/bisheng/org_sync/domain/services/reconcile_service.py
- ts 决策守卫:src/backend/bisheng/org_sync/domain/services/ts_guard.py
- Diff 引擎:src/backend/bisheng/org_sync/domain/services/remote_dept_differ.py
- Relink 服务与冲突存储:src/backend/bisheng/org_sync/domain/services/relink_service.py、src/backend/bisheng/org_sync/domain/services/relink_conflict_store.py
- relink HTTP 端点:src/backend/bisheng/org_sync/api/endpoints/relink.py
- 错误码定义:src/backend/bisheng/common/errcode/sso_sync.py
- 集成测试:src/backend/test/org_sync/test_org_reconcile_service.py、src/backend/test/celery/test_reconcile_celery_tasks.py
【免费下载链接】bishengBISHENG is an open LLM devops platform for next generation Enterprise AI applications. Powerful and comprehensive features include: GenAI workflow, RAG, Agent, Unified model management, Evaluation, SFT, Dataset Management, Enterprise-level System Management, Observability and more.项目地址: https://gitcode.com/GitHub_Trending/bi/bisheng
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考