wulf-pulse/lib/services/sync-service.ts
lorentz ff9e34cafe feat: add Autotask tags sync (tag groups, tags, ticket tag associations)
- Migration 057: autotask_tag_groups, autotask_tags, and junction tables
  (ticket_tags, company_tags, configuration_item_tags, contact_tags)
- Add TAG_GROUPS and TAGS to EntityType enum and dependency map
- Add mapTagGroup() and mapTag() entity mapper functions
- Add syncTagGroups(), syncTags(), syncTicketTagAssociations() methods
- Wire TicketTagAssociations bulk sync into full/incremental sync flow
- Add 'exist' operator to QueryFilter type
- No FK on ticket_id (tagged tickets may be outside 2yr sync window)

Synced: 26 tag groups, 7299 tags, 10141 ticket-tag associations
2026-03-20 09:22:40 -04:00

460 lines
13 KiB
TypeScript

/**
* Sync Service
* Main orchestration service for syncing Autotask data to PostgreSQL
*/
import postgresClient from './postgres-client';
import autotaskRateLimiter from './rate-limiter';
import { AutotaskClient } from './autotask-client';
import { createEntitySyncService, EntitySyncService } from './entity-sync';
import {
EntityType,
SyncType,
SyncStatus,
SyncConfig,
SyncResult,
SyncProgress,
EntitySyncResult,
SyncHistoryRecord
} from '../types/sync';
import {
getEntitySyncOrder,
getAllEntitiesInOrder,
generateSyncId,
getEntityDisplayName
} from '../utils/sync-helpers';
import { getLastSyncTime } from '../utils/db-helpers';
import { createSyncLogger, SyncPhase, categorizeError } from '../utils/sync-logger';
/**
* Main Sync Service Class
*/
export class SyncService {
private currentSyncId: string | null = null;
private isSyncing = false;
private autotaskClient: AutotaskClient;
private entitySyncService: EntitySyncService;
private logger = createSyncLogger({ component: 'SyncService' });
constructor(autotaskClient: AutotaskClient) {
this.autotaskClient = autotaskClient;
this.entitySyncService = createEntitySyncService(autotaskClient);
}
/**
* Start a full sync of all entities
* @param triggeredBy User or system identifier
* @param yearsBack Number of years to look back for time-based entities
* @returns Sync result
*/
async fullSync(triggeredBy?: string, yearsBack?: number): Promise<SyncResult> {
const config: SyncConfig = {
syncType: SyncType.FULL,
entities: getAllEntitiesInOrder(),
triggeredBy,
yearsBack,
};
return await this.executeSync(config);
}
/**
* Start an incremental sync of all entities
* @param triggeredBy User or system identifier
* @param yearsBack Number of years to look back for time-based entities
* @returns Sync result
*/
async incrementalSync(triggeredBy?: string, yearsBack?: number): Promise<SyncResult> {
const config: SyncConfig = {
syncType: SyncType.INCREMENTAL,
entities: getAllEntitiesInOrder(),
triggeredBy,
yearsBack,
};
return await this.executeSync(config);
}
/**
* Sync specific entities
* @param entities Array of entity types to sync
* @param syncType Type of sync (full or incremental)
* @param triggeredBy User or system identifier
* @param yearsBack Number of years to look back for time-based entities
* @returns Sync result
*/
async syncEntities(
entities: EntityType[],
syncType: SyncType = SyncType.ENTITY_SPECIFIC,
triggeredBy?: string,
yearsBack?: number
): Promise<SyncResult> {
const config: SyncConfig = {
syncType,
entities: getEntitySyncOrder(entities),
triggeredBy,
yearsBack,
};
return await this.executeSync(config);
}
/**
* Execute sync operation
* @param config Sync configuration
* @returns Sync result
*/
private async executeSync(config: SyncConfig): Promise<SyncResult> {
if (this.isSyncing) {
throw new Error('A sync operation is already in progress');
}
this.isSyncing = true;
const syncId = generateSyncId();
this.currentSyncId = syncId;
const startTime = new Date();
const entityResults: EntitySyncResult[] = [];
const errors: string[] = [];
try {
const syncLogger = this.logger.child({ syncId, syncType: config.syncType });
syncLogger.info(`Starting ${config.syncType} sync`, {
entities: config.entities.map(e => getEntityDisplayName(e)).join(', '),
triggeredBy: config.triggeredBy,
yearsBack: config.yearsBack,
});
// Sync each entity in order
for (const entity of config.entities) {
try {
const entityLogger = syncLogger.child({ entityType: entity });
const entityStartTime = entityLogger.start(`Syncing ${getEntityDisplayName(entity)}`);
// Create sync history record
const historyId = await this.createSyncHistory(
entity,
config.syncType,
config.triggeredBy
);
let recordsAdded = 0;
let recordsUpdated = 0;
let recordsDeleted = 0;
// Determine if incremental sync
const isIncremental = config.syncType === SyncType.INCREMENTAL;
const yearsBack = config.yearsBack || 2; // Default to 2 years
// Execute entity sync
const syncStats = await this.entitySyncService.syncEntity(entity, isIncremental, yearsBack);
recordsAdded = syncStats.recordsAdded;
recordsUpdated = syncStats.recordsUpdated;
recordsDeleted = syncStats.recordsDeleted;
const duration = Date.now() - entityStartTime;
// Update sync history
await this.updateSyncHistory(
historyId,
SyncStatus.COMPLETED,
recordsAdded,
recordsUpdated,
recordsDeleted
);
entityResults.push({
entityType: entity,
success: true,
recordsAdded,
recordsUpdated,
recordsDeleted,
duration,
});
entityLogger.complete(`Syncing ${getEntityDisplayName(entity)}`, entityStartTime, {
recordsAdded,
recordsUpdated,
recordsDeleted,
});
} catch (error) {
const err = error instanceof Error ? error : new Error(String(error));
const entityLogger = syncLogger.child({ entityType: entity });
const entityName = getEntityDisplayName(entity);
const errorCategory = categorizeError(err);
// Log detailed error information
entityLogger.error(`Failed to sync ${entityName}`, {
errorCategory,
syncType: config.syncType,
}, err);
const fullErrorMessage = `[${errorCategory}] ${err.message}`;
errors.push(`${entityName}: ${fullErrorMessage}`);
// Try to update sync history with error
try {
const historyId = await this.createSyncHistory(
entity,
config.syncType,
config.triggeredBy
);
await this.updateSyncHistory(
historyId,
SyncStatus.FAILED,
0,
0,
0,
fullErrorMessage
);
} catch (historyError) {
const histErr = historyError instanceof Error ? historyError : new Error(String(historyError));
entityLogger.error('Failed to update sync history with error', {}, histErr);
}
entityResults.push({
entityType: entity,
success: false,
recordsAdded: 0,
recordsUpdated: 0,
recordsDeleted: 0,
duration: 0,
error: fullErrorMessage,
});
}
}
// After all entities are synced, sync TicketTagAssociations (junction table)
try {
const tagAssocLogger = syncLogger.child({ entityType: 'ticket_tag_associations' });
tagAssocLogger.info('Syncing TicketTagAssociations');
const tagStats = await this.entitySyncService.syncTicketTagAssociations();
tagAssocLogger.info('TicketTagAssociations sync complete', { recordsAdded: tagStats.recordsAdded });
} catch (error) {
const err = error instanceof Error ? error : new Error(String(error));
syncLogger.warn('TicketTagAssociations sync failed, continuing', {}, err);
errors.push(`TicketTagAssociations: ${err.message}`);
}
const endTime = new Date();
const totalDuration = endTime.getTime() - startTime.getTime();
const result: SyncResult = {
syncId,
syncType: config.syncType,
status: errors.length === 0 ? SyncStatus.COMPLETED : SyncStatus.FAILED,
entities: entityResults,
totalRecordsAdded: entityResults.reduce((sum, r) => sum + r.recordsAdded, 0),
totalRecordsUpdated: entityResults.reduce((sum, r) => sum + r.recordsUpdated, 0),
totalRecordsDeleted: entityResults.reduce((sum, r) => sum + r.recordsDeleted, 0),
startedAt: startTime,
completedAt: endTime,
duration: totalDuration,
errors,
};
syncLogger.info(`Sync completed`, {
duration: totalDuration,
status: result.status,
recordsAdded: result.totalRecordsAdded,
recordsUpdated: result.totalRecordsUpdated,
recordsDeleted: result.totalRecordsDeleted,
errorCount: errors.length,
});
return result;
} catch (error) {
const err = error instanceof Error ? error : new Error(String(error));
const syncLogger = this.logger.child({ syncId, syncType: config.syncType });
syncLogger.error('Sync operation failed', {
entitiesAttempted: config.entities.join(', '),
successfulEntities: entityResults.filter(r => r.success).length,
failedEntities: entityResults.filter(r => !r.success).length,
}, err);
throw error;
} finally {
this.isSyncing = false;
this.currentSyncId = null;
}
}
/**
* Create sync history record
* @param entity Entity type
* @param syncType Sync type
* @param triggeredBy User identifier
* @returns Sync history ID
*/
async createSyncHistory(
entity: EntityType,
syncType: SyncType,
triggeredBy?: string
): Promise<number> {
const query = `
INSERT INTO sync_history (
entity_type,
sync_type,
status,
started_at,
records_added,
records_updated,
records_deleted,
triggered_by
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
RETURNING id
`;
const result = await postgresClient.query<{ id: number }>(query, [
entity,
syncType,
SyncStatus.STARTED,
new Date(),
0,
0,
0,
triggeredBy || 'system',
]);
return result.rows[0].id;
}
/**
* Update sync history record
* @param id Sync history ID
* @param status Sync status
* @param recordsAdded Number of records added
* @param recordsUpdated Number of records updated
* @param recordsDeleted Number of records deleted
* @param errorMessage Optional error message
*/
async updateSyncHistory(
id: number,
status: SyncStatus,
recordsAdded: number,
recordsUpdated: number,
recordsDeleted: number,
errorMessage?: string
): Promise<void> {
const query = `
UPDATE sync_history
SET status = $1,
completed_at = $2,
records_added = $3,
records_updated = $4,
records_deleted = $5,
error_message = $6
WHERE id = $7
`;
await postgresClient.query(query, [
status,
new Date(),
recordsAdded,
recordsUpdated,
recordsDeleted,
errorMessage || null,
id,
]);
}
/**
* Get sync history
* @param limit Number of records to return
* @param entityType Optional entity type filter
* @param offset Offset for pagination
* @returns Array of sync history records
*/
async getSyncHistory(
limit: number = 50,
entityType?: EntityType,
offset: number = 0
): Promise<SyncHistoryRecord[]> {
let query = `
SELECT *
FROM sync_history
`;
const params: any[] = [];
if (entityType) {
query += ` WHERE entity_type = $1`;
params.push(entityType);
}
query += ` ORDER BY started_at DESC LIMIT $${params.length + 1} OFFSET $${params.length + 2}`;
params.push(limit, offset);
const result = await postgresClient.query<SyncHistoryRecord>(query, params);
return result.rows;
}
/**
* Get last sync info for all entities
* @returns Map of entity type to last sync record
*/
async getLastSyncInfo(): Promise<Map<EntityType, SyncHistoryRecord>> {
const query = `
SELECT DISTINCT ON (entity_type) *
FROM sync_history
WHERE status = 'completed'
ORDER BY entity_type, completed_at DESC
`;
const result = await postgresClient.query<SyncHistoryRecord>(query);
const map = new Map<EntityType, SyncHistoryRecord>();
for (const row of result.rows) {
// Check if the entity_type value exists in the EntityType enum values
const entityTypeValues = Object.values(EntityType);
if (entityTypeValues.includes(row.entity_type as EntityType)) {
map.set(row.entity_type as EntityType, row);
}
}
return map;
}
/**
* Check if sync is currently running
* @returns True if sync is in progress
*/
isSyncInProgress(): boolean {
return this.isSyncing;
}
/**
* Get current sync ID
* @returns Current sync ID or null
*/
getCurrentSyncId(): string | null {
return this.currentSyncId;
}
/**
* Cancel current sync operation
*/
async cancelSync(): Promise<void> {
if (!this.isSyncing) {
throw new Error('No sync operation in progress');
}
// TODO: Implement graceful cancellation
this.isSyncing = false;
const cancelledSyncId = this.currentSyncId || 'unknown';
this.currentSyncId = null;
this.logger.warn('Sync operation cancelled', { syncId: cancelledSyncId });
}
}
/**
* Create sync service instance
* @param autotaskClient Autotask client instance
* @returns SyncService instance
*/
export function createSyncService(autotaskClient: AutotaskClient): SyncService {
return new SyncService(autotaskClient);
}