19 KiB
钉钉/企微同步对接 Implementation Plan
For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) to implement this plan task-by-task. Steps use checkbox (
- [ ]) syntax for tracking.
Goal: Replace sync.service.ts stubs with real DingTalk/WeCom API calls — fetch departments and users, persist to DB, record sync logs.
Architecture: Add source/sourceId/parentSourceId to Department entity for idempotent sync matching. Create two integration services (DingTalkService, WeComService) under src/integration/, one shared IntegrationModule, wire into SyncModule. Each service self-checks env vars and degrades gracefully when not configured.
Tech Stack: NestJS 11 + TypeORM + @nestjs/schedule + native fetch
Global Constraints
+- MUST check process.env.DINGTALK_APP_KEY / WECOM_CORP_ID before attempting API calls
+- When env vars missing: log warning, return 0 records, set sync log to success (not failed — "not configured" is not an error)
+- Sync count returned from perform*Sync is the number of new/updated records persisted
+- Follow existing NestJS module structure: one directory per concern
+- Each service is independently injectable; SyncModule imports IntegrationModule
+- Use native fetch (Node 18+) — no extra HTTP client dependency
+- Department entity gains nullable source / sourceId / parentSourceId — existing records unaffected
+- Sync preserves existing tree structure: departments matched by sourceId, users by username
Task 1: Add Source Tracking to Department Entity
Files:
+- Modify: apps/server/src/entities/department.entity.ts
Interfaces:
+- Produces: Department entity with nullable source, sourceId, parentSourceId columns
+- [ ] Step 1: Add columns to Department entity
// apps/server/src/entities/department.entity.ts — add after 'status' field (line ~42):
@Column({ length: 20, nullable: true })
source: string; // 'dingtalk' | 'wecom' | null (null = manual)
@Column({ name: 'source_id', length: 50, nullable: true })
sourceId: string; // external dept ID for idempotent sync matching
@Column({ name: 'parent_source_id', length: 50, nullable: true })
parentSourceId: string; // external parent dept ID (resolved in post-processing)
+- [ ] Step 2: Verify compilation
Run: cd apps/server && npx tsc --noEmit 2>&1 | head -20
Expected: No new errors from department.entity.ts
+- [ ] Step 3: Commit
git add apps/server/src/entities/department.entity.ts
git commit -m "feat: add source tracking fields to Department for sync idempotency"
Task 2: DingTalk Integration Service
Files:
+- Create: apps/server/src/integration/dingtalk.service.ts
Interfaces:
+- Produces: DingTalkService with syncAll(): Promise<{ deptCount: number; userCount: number }>
+- Consumes: Department repo, User repo
+- [ ] Step 1: Create the full DingTalkService
// apps/server/src/integration/dingtalk.service.ts
import { Injectable, Logger } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { Department } from '../entities/department.entity';
import { User } from '../entities/user.entity';
interface DingTalkTokenResponse {
errcode: number;
errmsg: string;
access_token: string;
expires_in: number;
}
interface DingTalkDeptListResponse {
errcode: number;
errmsg: string;
result: Array<{ dept_id: number; name: string; parent_id: number }>;
}
interface DingTalkUserListResponse {
errcode: number;
errmsg: string;
result: {
has_more: boolean;
list: Array<{
userid: string;
name: string;
mobile: string;
dept_id_list: number[];
}>;
};
}
@Injectable()
export class DingTalkService {
private readonly logger = new Logger(DingTalkService.name);
private accessToken: string | null = null;
private tokenExpiresAt = 0;
constructor(
@InjectRepository(Department)
private readonly deptRepo: Repository<Department>,
@InjectRepository(User)
private readonly userRepo: Repository<User>,
) {}
private get configured(): boolean {
return !!(process.env.DINGTALK_APP_KEY && process.env.DINGTALK_APP_SECRET);
}
private async getAccessToken(): Promise<string> {
if (this.accessToken && Date.now() < this.tokenExpiresAt - 60_000) {
return this.accessToken;
}
const appKey = process.env.DINGTALK_APP_KEY!;
const appSecret = process.env.DINGTALK_APP_SECRET!;
const url = `https://oapi.dingtalk.com/gettoken?appkey=${appKey}&appsecret=${appSecret}`;
const res = await fetch(url);
const body: DingTalkTokenResponse = await res.json();
if (body.errcode !== 0) {
throw new Error(`DingTalk gettoken failed: ${body.errmsg} (${body.errcode})`);
}
this.accessToken = body.access_token;
this.tokenExpiresAt = Date.now() + body.expires_in * 1000;
return this.accessToken;
}
private async fetchDepartments(token: string): Promise<Array<{ dept_id: number; name: string; parent_id: number }>> {
const url = `https://oapi.dingtalk.com/topapi/v2/department/listsub?access_token=${token}`;
const res = await fetch(url, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ dept_id: 1 }),
});
const body: DingTalkDeptListResponse = await res.json();
if (body.errcode !== 0) {
throw new Error(`DingTalk department list failed: ${body.errmsg} (${body.errcode})`);
}
return body.result;
}
private async fetchUsers(
token: string,
deptId: number,
): Promise<Array<{ userid: string; name: string; mobile: string; dept_id_list: number[] }>> {
const allUsers: DingTalkUserListResponse['result']['list'] = [];
let cursor = 0;
while (true) {
const url = `https://oapi.dingtalk.com/topapi/v2/user/list?access_token=${token}`;
const res = await fetch(url, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ dept_id: deptId, cursor, size: 100 }),
});
const body: DingTalkUserListResponse = await res.json();
if (body.errcode !== 0) {
throw new Error(`DingTalk user list failed: ${body.errmsg} (${body.errcode})`);
}
allUsers.push(...body.result.list);
if (!body.result.has_more) break;
cursor = allUsers.length;
}
return allUsers;
}
async syncAll(): Promise<{ deptCount: number; userCount: number }> {
if (!this.configured) {
this.logger.warn('DingTalk not configured (DINGTALK_APP_KEY / DINGTALK_APP_SECRET missing), skipping sync');
return { deptCount: 0, userCount: 0 };
}
const token = await this.getAccessToken();
const dingDepts = await this.fetchDepartments(token);
// Upsert departments
let deptCount = 0;
for (const dd of dingDepts) {
const sourceId = String(dd.dept_id);
let dept = await this.deptRepo.findOne({ where: { source: 'dingtalk', sourceId } });
if (dept) {
dept.name = dd.name;
dept.parentSourceId = dd.parent_id ? String(dd.parent_id) : null;
} else {
dept = this.deptRepo.create({
name: dd.name,
source: 'dingtalk',
sourceId,
parentSourceId: dd.parent_id ? String(dd.parent_id) : null,
type: 'department',
});
deptCount++;
}
await this.deptRepo.save(dept);
}
// Resolve parentSourceId → parentId for tree linking
const syncedDepts = await this.deptRepo.find({ where: { source: 'dingtalk' } });
const idMap = new Map(syncedDepts.map((d) => [d.sourceId, d.id]));
for (const dept of syncedDepts) {
if (dept.parentSourceId && idMap.has(dept.parentSourceId)) {
dept.parentId = idMap.get(dept.parentSourceId)!;
} else if (dept.parentSourceId === '1' || dept.parentSourceId === '0') {
dept.parentId = null; // root
}
}
await this.deptRepo.save(syncedDepts);
// Upsert users across all departments
let userCount = 0;
const seenUserIds = new Set<string>();
for (const dd of dingDepts) {
const dingUsers = await this.fetchUsers(token, dd.dept_id);
for (const du of dingUsers) {
if (seenUserIds.has(du.userid)) continue;
seenUserIds.add(du.userid);
let user = await this.userRepo.findOne({ where: { username: du.userid } });
if (user) {
user.name = du.name;
} else {
user = this.userRepo.create({
username: du.userid,
name: du.name,
passwordHash: '',
isActive: true,
});
userCount++;
}
await this.userRepo.save(user);
}
}
this.logger.log(`DingTalk sync done: ${deptCount} new depts, ${userCount} new users`);
return { deptCount, userCount };
}
}
+- [ ] Step 2: Verify file exists
Run: wc -l apps/server/src/integration/dingtalk.service.ts
Expected: ~160 lines
Task 3: WeCom Integration Service
Files:
+- Create: apps/server/src/integration/wecom.service.ts
Interfaces:
+- Produces: WeComService with syncAll(): Promise<{ deptCount: number; userCount: number }>
+- [ ] Step 1: Create the full WeComService
// apps/server/src/integration/wecom.service.ts
import { Injectable, Logger } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { Department } from '../entities/department.entity';
import { User } from '../entities/user.entity';
interface WeComTokenResponse {
errcode: number;
errmsg: string;
access_token: string;
expires_in: number;
}
interface WeComDeptListResponse {
errcode: number;
errmsg: string;
department: Array<{ id: number; name: string; parentid: number }>;
}
interface WeComUserListResponse {
errcode: number;
errmsg: string;
userlist: Array<{
userid: string;
name: string;
mobile: string;
department: number[];
}>;
}
@Injectable()
export class WeComService {
private readonly logger = new Logger(WeComService.name);
private accessToken: string | null = null;
private tokenExpiresAt = 0;
constructor(
@InjectRepository(Department)
private readonly deptRepo: Repository<Department>,
@InjectRepository(User)
private readonly userRepo: Repository<User>,
) {}
private get configured(): boolean {
return !!(process.env.WECOM_CORP_ID && process.env.WECOM_CORP_SECRET);
}
private async getAccessToken(): Promise<string> {
if (this.accessToken && Date.now() < this.tokenExpiresAt - 60_000) {
return this.accessToken;
}
const corpId = process.env.WECOM_CORP_ID!;
const corpSecret = process.env.WECOM_CORP_SECRET!;
const url = `https://qyapi.weixin.qq.com/cgi-bin/gettoken?corpid=${corpId}&corpsecret=${corpSecret}`;
const res = await fetch(url);
const body: WeComTokenResponse = await res.json();
if (body.errcode !== 0) {
throw new Error(`WeCom gettoken failed: ${body.errmsg} (${body.errcode})`);
}
this.accessToken = body.access_token;
this.tokenExpiresAt = Date.now() + body.expires_in * 1000;
return this.accessToken;
}
private async fetchDepartments(
token: string,
parentId = 1,
): Promise<Array<{ id: number; name: string; parentid: number }>> {
const all: WeComDeptListResponse['department'] = [];
const url = `https://qyapi.weixin.qq.com/cgi-bin/department/list?access_token=${token}&id=${parentId}`;
const res = await fetch(url);
const body: WeComDeptListResponse = await res.json();
if (body.errcode !== 0) {
// Empty department list is OK for leaf departments
if (body.errcode === 60003) return all;
throw new Error(`WeCom department list failed: ${body.errmsg} (${body.errcode})`);
}
for (const dept of body.department) {
all.push(dept);
if (dept.id !== parentId) {
const children = await this.fetchDepartments(token, dept.id);
all.push(...children);
}
}
return all;
}
private async fetchUsers(
token: string,
deptId: number,
): Promise<Array<{ userid: string; name: string; mobile: string; department: number[] }>> {
const url = `https://qyapi.weixin.qq.com/cgi-bin/user/simplelist?access_token=${token}&department_id=${deptId}&fetch_child=1`;
const res = await fetch(url);
const body: WeComUserListResponse = await res.json();
if (body.errcode !== 0) {
throw new Error(`WeCom user list failed: ${body.errmsg} (${body.errcode})`);
}
return body.userlist;
}
async syncAll(): Promise<{ deptCount: number; userCount: number }> {
if (!this.configured) {
this.logger.warn('WeCom not configured (WECOM_CORP_ID / WECOM_CORP_SECRET missing), skipping sync');
return { deptCount: 0, userCount: 0 };
}
const token = await this.getAccessToken();
const wxDepts = await this.fetchDepartments(token);
// Upsert departments
let deptCount = 0;
for (const wd of wxDepts) {
const sourceId = String(wd.id);
let dept = await this.deptRepo.findOne({ where: { source: 'wecom', sourceId } });
if (dept) {
dept.name = wd.name;
dept.parentSourceId = wd.parentid ? String(wd.parentid) : null;
} else {
dept = this.deptRepo.create({
name: wd.name,
source: 'wecom',
sourceId,
parentSourceId: wd.parentid ? String(wd.parentid) : null,
type: 'department',
});
deptCount++;
}
await this.deptRepo.save(dept);
}
// Resolve parentSourceId → parentId
const syncedDepts = await this.deptRepo.find({ where: { source: 'wecom' } });
const idMap = new Map(syncedDepts.map((d) => [d.sourceId, d.id]));
for (const dept of syncedDepts) {
if (dept.parentSourceId && idMap.has(dept.parentSourceId)) {
dept.parentId = idMap.get(dept.parentSourceId)!;
} else if (dept.parentSourceId === '0' || dept.parentSourceId === '1') {
dept.parentId = null;
}
}
await this.deptRepo.save(syncedDepts);
// Upsert users
let userCount = 0;
const seenUserIds = new Set<string>();
for (const wd of wxDepts) {
const wxUsers = await this.fetchUsers(token, wd.id);
for (const wu of wxUsers) {
if (seenUserIds.has(wu.userid)) continue;
seenUserIds.add(wu.userid);
let user = await this.userRepo.findOne({ where: { username: wu.userid } });
if (user) {
user.name = wu.name;
} else {
user = this.userRepo.create({
username: wu.userid,
name: wu.name,
passwordHash: '',
isActive: true,
});
userCount++;
}
await this.userRepo.save(user);
}
}
this.logger.log(`WeCom sync done: ${deptCount} new depts, ${userCount} new users`);
return { deptCount, userCount };
}
}
+- [ ] Step 2: Verify file exists
Run: wc -l apps/server/src/integration/wecom.service.ts
Expected: ~170 lines
Task 4: Integration Module + Wiring
Files:
+- Create: apps/server/src/integration/integration.module.ts
+- Modify: apps/server/src/sync/sync.module.ts
+- Modify: apps/server/src/sync/sync.service.ts
Interfaces:
+- Consumes: DingTalkService.syncAll(), WeComService.syncAll()
+- Produces: SyncService with real sync calls replacing stubs
+- [ ] Step 1: Create IntegrationModule
// apps/server/src/integration/integration.module.ts
import { Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { Department, User } from '../entities';
import { DingTalkService } from './dingtalk.service';
import { WeComService } from './wecom.service';
@Module({
imports: [TypeOrmModule.forFeature([Department, User])],
providers: [DingTalkService, WeComService],
exports: [DingTalkService, WeComService],
})
export class IntegrationModule {}
+- [ ] Step 2: Wire IntegrationModule into SyncModule
In apps/server/src/sync/sync.module.ts: add IntegrationModule to the imports array and the import statement:
// Add at top:
import { IntegrationModule } from '../integration/integration.module';
// In @Module decorator, add IntegrationModule to imports:
@Module({
imports: [
ScheduleModule.forRoot(),
TypeOrmModule.forFeature([SyncLog, SyncState]),
IntegrationModule,
],
// ... rest unchanged
+- [ ] Step 3: Inject services and replace stubs in SyncService
In apps/server/src/sync/sync.service.ts:
Add imports:
import { DingTalkService } from '../integration/dingtalk.service';
import { WeComService } from '../integration/wecom.service';
Add to constructor parameters:
constructor(
@InjectRepository(SyncLog) private readonly syncLogRepo: Repository<SyncLog>,
@InjectRepository(SyncState) private readonly syncStateRepo: Repository<SyncState>,
private readonly dingTalkService: DingTalkService,
private readonly weComService: WeComService,
) {}
Replace performDingTalkSync (lines 154-166):
private async performDingTalkSync(_lastSyncAt: Date | null): Promise<number> {
const result = await this.dingTalkService.syncAll();
return result.deptCount + result.userCount;
}
Replace performWeComSync (lines 168-179):
private async performWeComSync(_lastSyncAt: Date | null): Promise<number> {
const result = await this.weComService.syncAll();
return result.deptCount + result.userCount;
}
Also remove the stale JSDoc comments above the old stubs.
+- [ ] Step 4: Verify compilation
Run: cd apps/server && npx tsc --noEmit 2>&1 | head -30
Expected: No type errors from integration/ or sync/ modules
+- [ ] Step 5: Commit
git add apps/server/src/integration/
git add apps/server/src/sync/sync.module.ts apps/server/src/sync/sync.service.ts
git commit -m "feat: wire DingTalk/WeCom integration services into sync pipeline"
Task 5: Verification
Files: +- (none modified — verification only)
+- [ ] Step 1: Start dev server with missing env vars
cd apps/server && npm run start:dev &
sleep 5
Check logs: should show DingTalk not configured... skipping sync and WeCom not configured... skipping sync at startup (or wait for the 2 AM cron, or trigger manually).
+- [ ] Step 2: Test manual trigger endpoint
curl -s http://localhost:3002/api/sync/trigger | python3 -m json.tool 2>/dev/null || curl -s http://localhost:3002/api/sync/trigger
Expected: JSON array of sync log objects with status: "success" and recordsCount: 0
+- [ ] Step 3: Check sync logs endpoint
curl -s http://localhost:3002/api/sync/logs | python3 -m json.tool 2>/dev/null | head -30
Expected: Array of sync log entries with fields platform, status, recordsCount, startedAt, finishedAt
+- [ ] Step 4: Verify server still serves other endpoints
curl -s http://localhost:3002/api/students?pageSize=1 | python3 -m json.tool 2>/dev/null | head -10
Expected: Normal student list response (no regression)
+- [ ] Step 5: Stop server and commit verification
kill %1 2>/dev/null
# If all checks passed, no additional commits needed