import { BadRequestException, ConflictException, Injectable, NotFoundException, } from '@nestjs/common'; import { DataSource, EntityManager } from 'typeorm'; import { AuditService } from '../audit/audit.service'; import { administrationAuditContext, isUniqueViolation } from '../administration/common/administration-audit'; import type { AuthPrincipal, RequestWithContext } from '../common/http/request-context'; import { AreaLegalRight, AreaLegalRightOrganization, Asset, AssetExternalIdentifier, AssetSourceDocument, AssetTypeOperationalRole, AssetVersionChangeType, AuditAction, OrganizationKind, OrganizationMembership, OrganizationProfile, SourceDocument, } from '../database/entities'; import { AssetHistoryService } from './asset-history.service'; import type { AddAreaLegalRightOrganizationDto } from './dto/add-area-legal-right-organization.dto'; import type { AddOrganizationMembershipDto } from './dto/add-organization-membership.dto'; import type { CreateAreaLegalRightDto } from './dto/create-area-legal-right.dto'; import type { CreateExternalIdentifierDto } from './dto/create-external-identifier.dto'; import type { CreateSourceDocumentDto } from './dto/create-source-document.dto'; import type { EndAreaLegalRightOrganizationDto } from './dto/end-area-legal-right-organization.dto'; import type { EndExternalIdentifierDto } from './dto/end-external-identifier.dto'; import type { EndOrganizationMembershipDto } from './dto/end-organization-membership.dto'; import type { LinkAssetSourceDocumentDto } from './dto/link-asset-source-document.dto'; import type { UpdateAreaLegalRightDto } from './dto/update-area-legal-right.dto'; import type { UpsertOrganizationProfileDto } from './dto/upsert-organization-profile.dto'; function registryNotFound(entity: string): NotFoundException { return new NotFoundException({ code: 'ASSET_REGISTRY_NOT_FOUND', message: `${entity} no encontrado` }); } @Injectable() export class AssetRegistryService { constructor( private readonly dataSource: DataSource, private readonly audit: AuditService, private readonly history: AssetHistoryService, ) {} async get(assetId: string) { await this.requireAsset(this.dataSource.manager, assetId); const [organizationProfile] = await this.dataSource.query(` SELECT asset_id AS "assetId", organization_kind AS "organizationKind", legal_name AS "legalName", tax_id AS "taxId", notification_email AS "notificationEmail", notes, created_at AS "createdAt", updated_at AS "updatedAt", updated_by AS "updatedBy" FROM organization_profiles WHERE asset_id=$1 `, [assetId]); const organizationMemberships = await this.dataSource.query(` SELECT m.id, JSONB_BUILD_OBJECT('id',p.id,'code',p.code,'name',p.name) AS parent, JSONB_BUILD_OBJECT('id',member.id,'code',member.code,'name',member.name) AS member, m.role, m.participation_percent::double precision AS "participationPercent", m.valid_from AS "validFrom", m.valid_until AS "validUntil", m.source_document_id AS "sourceDocumentId", m.notes, m.end_reason AS "endReason", m.created_at AS "createdAt", m.updated_at AS "updatedAt" FROM organization_memberships m JOIN assets p ON p.id=m.parent_organization_id JOIN assets member ON member.id=m.member_organization_id WHERE m.parent_organization_id=$1 OR m.member_organization_id=$1 ORDER BY m.valid_until NULLS FIRST, m.valid_from DESC `, [assetId]); const externalIdentifiers = await this.dataSource.query(` SELECT id, namespace, value, valid_from AS "validFrom", valid_until AS "validUntil", source_document_id AS "sourceDocumentId", notes, end_reason AS "endReason", created_at AS "createdAt" FROM asset_external_identifiers WHERE asset_id=$1 ORDER BY valid_until NULLS FIRST, namespace, valid_from DESC `, [assetId]); const sourceDocuments = await this.dataSource.query(` SELECT link.id AS "linkId", link.relation_type AS "relationType", link.notes AS "linkNotes", doc.id, doc.document_type AS "documentType", doc.document_number AS "documentNumber", doc.title, doc.issuer, doc.document_date AS "documentDate", doc.external_reference AS "externalReference", doc.notes FROM asset_source_documents link JOIN source_documents doc ON doc.id=link.document_id WHERE link.asset_id=$1 ORDER BY doc.document_date DESC NULLS LAST, doc.created_at DESC `, [assetId]); const legalRights = await this.dataSource.query(` SELECT r.id, r.right_type AS "rightType", r.name, r.instrument_number AS "instrumentNumber", r.valid_from AS "validFrom", r.valid_until AS "validUntil", r.status, r.source_document_id AS "sourceDocumentId", r.notes, COALESCE((SELECT JSONB_AGG(JSONB_BUILD_OBJECT( 'id',o.id,'organizationId',o.organization_id,'organizationName',a.name,'role',o.role, 'participationPercent',o.participation_percent::double precision,'validFrom',o.valid_from, 'validUntil',o.valid_until,'notes',o.notes,'endReason',o.end_reason ) ORDER BY o.valid_until NULLS FIRST,o.valid_from DESC) FROM area_legal_right_organizations o JOIN assets a ON a.id=o.organization_id WHERE o.right_id=r.id),'[]'::jsonb) AS organizations FROM area_legal_rights r WHERE r.area_id=$1 ORDER BY r.valid_until DESC NULLS FIRST, r.valid_from DESC NULLS LAST `, [assetId]); return { organizationProfile: organizationProfile ?? null, organizationMemberships, externalIdentifiers, sourceDocuments, legalRights }; } async listSourceDocuments(search?: string) { const q = search?.trim(); const params: unknown[] = []; const where = q ? `WHERE title ILIKE $1 OR document_number ILIKE $1 OR issuer ILIKE $1` : ''; if (q) params.push(`%${q}%`); const data = await this.dataSource.query(` SELECT id, document_type AS "documentType", document_number AS "documentNumber", title, issuer, document_date AS "documentDate", external_reference AS "externalReference", notes, created_at AS "createdAt", updated_at AS "updatedAt" FROM source_documents ${where} ORDER BY document_date DESC NULLS LAST, created_at DESC LIMIT 100 `, params); return { data }; } async createSourceDocument(dto: CreateSourceDocumentDto, principal: AuthPrincipal, request: RequestWithContext) { try { return await this.dataSource.transaction(async manager => { const doc = manager.getRepository(SourceDocument).create({ documentType: dto.documentType, documentNumber: dto.documentNumber ?? null, title: dto.title, issuer: dto.issuer ?? null, documentDate: dto.documentDate ?? null, externalReference: dto.externalReference ?? null, notes: dto.notes ?? null, createdBy: principal.userId, updatedBy: principal.userId, }); await manager.getRepository(SourceDocument).save(doc); await this.audit.record({ ...administrationAuditContext(principal, request), action: AuditAction.SOURCE_DOCUMENT_CREATED, entityType: 'source_document', entityId: doc.id, afterData: { ...doc }, }, manager); return doc; }); } catch (error) { if (isUniqueViolation(error)) throw new ConflictException({ code:'SOURCE_DOCUMENT_EXISTS', message:'Ya existe un documento con ese número y emisor' }); throw error; } } async upsertOrganizationProfile(assetId: string, dto: UpsertOrganizationProfileDto, principal: AuthPrincipal, request: RequestWithContext) { try { return await this.dataSource.transaction(async manager => { const asset = await this.requireRole(manager, assetId, AssetTypeOperationalRole.COMPANY); const repo = manager.getRepository(OrganizationProfile); const before = await repo.findOne({ where: { assetId } }); if (before && before.organizationKind !== dto.organizationKind) { const [usage] = await manager.query(`SELECT 1 FROM organization_memberships WHERE (parent_organization_id=$1 OR member_organization_id=$1) AND valid_until IS NULL LIMIT 1`, [assetId]); if (usage) throw new ConflictException({ code:'ORGANIZATION_KIND_IN_USE', message:'No se puede cambiar el tipo de organización mientras tenga una composición UTE activa' }); } const profile = repo.create({ ...(before ?? {}), assetId, organizationKind:dto.organizationKind, legalName:dto.legalName ?? asset.name, taxId:dto.taxId ?? null, notificationEmail:dto.notificationEmail ?? null, notes:dto.notes ?? null, updatedBy:principal.userId }); await repo.save(profile); const versionNumber = await this.history.capture(manager, assetId, AssetVersionChangeType.REGISTRY_UPDATED, principal, request); await this.audit.record({ ...administrationAuditContext(principal,request), action:AuditAction.ASSET_REGISTRY_UPDATED, entityType:'organization_profile', entityId:assetId, beforeData:before ? { ...before }:null, afterData:{ ...profile }, metadata:{versionNumber} }, manager); return profile; }); } catch (error) { if (isUniqueViolation(error)) throw new ConflictException({ code:'ORGANIZATION_TAX_ID_EXISTS', message:'El identificador fiscal ya pertenece a otra organización' }); throw error; } } async addOrganizationMembership(parentId:string, dto:AddOrganizationMembershipDto, principal:AuthPrincipal, request:RequestWithContext) { try { return await this.dataSource.transaction(async manager => { await this.requireRole(manager,parentId,AssetTypeOperationalRole.COMPANY); await manager.query('SELECT id FROM assets WHERE id=$1 FOR UPDATE',[parentId]); await this.requireRole(manager,dto.memberOrganizationId,AssetTypeOperationalRole.COMPANY); const parentProfile = await manager.getRepository(OrganizationProfile).findOne({where:{assetId:parentId}}); const memberProfile = await manager.getRepository(OrganizationProfile).findOne({where:{assetId:dto.memberOrganizationId}}); if (parentProfile?.organizationKind !== OrganizationKind.UTE) throw new BadRequestException({code:'UTE_REQUIRED',message:'La organización contenedora debe estar configurada como UTE'}); if (memberProfile?.organizationKind === OrganizationKind.UTE) throw new BadRequestException({code:'UTE_MEMBER_INVALID',message:'Una UTE no puede ser miembro directo de otra UTE'}); if (dto.sourceDocumentId) await this.requireDocument(manager,dto.sourceDocumentId); const validFrom = dto.validFrom ?? this.today(); const [overlap] = await manager.query(` SELECT 1 FROM organization_memberships WHERE parent_organization_id=$1 AND member_organization_id=$2 AND role=$3 AND daterange(valid_from,COALESCE(valid_until,'infinity'::date),'[]') && daterange($4::date,'infinity'::date,'[]') LIMIT 1 `,[parentId,dto.memberOrganizationId,dto.role,validFrom]); if (overlap) throw new ConflictException({code:'UTE_MEMBERSHIP_OVERLAP',message:'La participación se superpone con una vigencia histórica existente'}); if (dto.participationPercent != null) { const [sum] = await manager.query(`SELECT COALESCE(SUM(participation_percent),0)::double precision AS total FROM organization_memberships WHERE parent_organization_id=$1 AND valid_until IS NULL`,[parentId]); if (Number(sum?.total ?? 0) + dto.participationPercent > 100.0001) throw new BadRequestException({code:'UTE_PARTICIPATION_EXCEEDS_100',message:'La participación activa de la UTE no puede superar el 100%'}); } const repo=manager.getRepository(OrganizationMembership); const membership=repo.create({parentOrganizationId:parentId,memberOrganizationId:dto.memberOrganizationId,role:dto.role,participationPercent:dto.participationPercent == null ? null:String(dto.participationPercent),validFrom,validUntil:null,sourceDocumentId:dto.sourceDocumentId ?? null,notes:dto.notes ?? null,endReason:null,createdBy:principal.userId,endedBy:null}); await repo.save(membership); await this.captureAssets(manager,[parentId,dto.memberOrganizationId],principal,request); await this.audit.record({ ...administrationAuditContext(principal,request),action:AuditAction.ASSET_REGISTRY_UPDATED,entityType:'organization_membership',entityId:membership.id,afterData:{...membership}},manager); return membership; }); } catch(error){ if(isUniqueViolation(error)) throw new ConflictException({code:'UTE_MEMBERSHIP_EXISTS',message:'La participación ya se encuentra activa'}); throw error; } } async endOrganizationMembership(id:string,dto:EndOrganizationMembershipDto,principal:AuthPrincipal,request:RequestWithContext){ return this.dataSource.transaction(async manager=>{ const repo=manager.getRepository(OrganizationMembership); const membership=await repo.createQueryBuilder('m').where('m.id=:id',{id}).setLock('pessimistic_write').getOne(); if(!membership) throw registryNotFound('Participación de organización'); if(membership.validUntil) throw new ConflictException({code:'UTE_MEMBERSHIP_ENDED',message:'La participación ya está finalizada'}); const end=dto.validUntil ?? this.today(); this.validateEndDate(end,membership.validFrom); membership.validUntil=end; membership.endReason=dto.reason; membership.endedBy=principal.userId; await repo.save(membership); await this.captureAssets(manager,[membership.parentOrganizationId,membership.memberOrganizationId],principal,request); await this.audit.record({...administrationAuditContext(principal,request),action:AuditAction.ASSET_REGISTRY_UPDATED,entityType:'organization_membership',entityId:id,afterData:{...membership}},manager); return membership; }); } async linkDocument(assetId:string,documentId:string,dto:LinkAssetSourceDocumentDto,principal:AuthPrincipal,request:RequestWithContext){ try { return await this.dataSource.transaction(async manager=>{ await this.requireAsset(manager,assetId); await this.requireDocument(manager,documentId); const repo=manager.getRepository(AssetSourceDocument); const link=repo.create({assetId,documentId,relationType:dto.relationType,notes:dto.notes??null,createdBy:principal.userId}); await repo.save(link); const versionNumber=await this.history.capture(manager,assetId,AssetVersionChangeType.REGISTRY_UPDATED,principal,request); await this.audit.record({...administrationAuditContext(principal,request),action:AuditAction.ASSET_REGISTRY_UPDATED,entityType:'asset_source_document',entityId:link.id,afterData:{...link},metadata:{versionNumber}},manager); return link; }); } catch(error){ if(isUniqueViolation(error)) throw new ConflictException({code:'ASSET_DOCUMENT_LINK_EXISTS',message:'El documento ya está vinculado al activo con ese tipo de relación'}); throw error; } } async addExternalIdentifier(assetId:string,dto:CreateExternalIdentifierDto,principal:AuthPrincipal,request:RequestWithContext){ try { return await this.dataSource.transaction(async manager=>{ await this.requireAsset(manager,assetId); if(dto.sourceDocumentId) await this.requireDocument(manager,dto.sourceDocumentId); const repo=manager.getRepository(AssetExternalIdentifier); const identifier=repo.create({assetId,namespace:dto.namespace,value:dto.value,validFrom:dto.validFrom ? new Date(dto.validFrom):new Date(),validUntil:null,sourceDocumentId:dto.sourceDocumentId??null,notes:dto.notes??null,endReason:null,createdBy:principal.userId,endedBy:null}); await repo.save(identifier); const versionNumber=await this.history.capture(manager,assetId,AssetVersionChangeType.REGISTRY_UPDATED,principal,request); await this.audit.record({...administrationAuditContext(principal,request),action:AuditAction.ASSET_REGISTRY_UPDATED,entityType:'asset_external_identifier',entityId:identifier.id,afterData:{...identifier},metadata:{versionNumber}},manager); return identifier; }); } catch(error){ if(isUniqueViolation(error)) throw new ConflictException({code:'EXTERNAL_IDENTIFIER_EXISTS',message:'Ese identificador externo ya está activo'}); throw error; } } async endExternalIdentifier(id:string,dto:EndExternalIdentifierDto,principal:AuthPrincipal,request:RequestWithContext){ return this.dataSource.transaction(async manager=>{ const repo=manager.getRepository(AssetExternalIdentifier); const item=await repo.createQueryBuilder('i').where('i.id=:id',{id}).setLock('pessimistic_write').getOne(); if(!item) throw registryNotFound('Identificador'); if(item.validUntil) throw new ConflictException({code:'EXTERNAL_IDENTIFIER_ENDED',message:'El identificador ya está finalizado'}); item.validUntil=new Date(); item.endReason=dto.reason; item.endedBy=principal.userId; await repo.save(item); const versionNumber=await this.history.capture(manager,item.assetId,AssetVersionChangeType.REGISTRY_UPDATED,principal,request); await this.audit.record({...administrationAuditContext(principal,request),action:AuditAction.ASSET_REGISTRY_UPDATED,entityType:'asset_external_identifier',entityId:id,afterData:{...item},metadata:{versionNumber}},manager); return item; }); } async createLegalRight(areaId:string,dto:CreateAreaLegalRightDto,principal:AuthPrincipal,request:RequestWithContext){ return this.dataSource.transaction(async manager=>{ await this.requireRole(manager,areaId,AssetTypeOperationalRole.AREA); if(dto.sourceDocumentId) await this.requireDocument(manager,dto.sourceDocumentId); this.validateDateRange(dto.validFrom??null,dto.validUntil??null); const repo=manager.getRepository(AreaLegalRight); const right=repo.create({areaId,rightType:dto.rightType,name:dto.name,instrumentNumber:dto.instrumentNumber??null,validFrom:dto.validFrom??null,validUntil:dto.validUntil??null,status:dto.status,sourceDocumentId:dto.sourceDocumentId??null,notes:dto.notes??null,createdBy:principal.userId,updatedBy:principal.userId}); await repo.save(right); const versionNumber=await this.history.capture(manager,areaId,AssetVersionChangeType.REGISTRY_UPDATED,principal,request); await this.audit.record({...administrationAuditContext(principal,request),action:AuditAction.AREA_LEGAL_RIGHT_CREATED,entityType:'area_legal_right',entityId:right.id,afterData:{...right},metadata:{versionNumber}},manager); return right; }); } async updateLegalRight(id:string,dto:UpdateAreaLegalRightDto,principal:AuthPrincipal,request:RequestWithContext){ return this.dataSource.transaction(async manager=>{ const repo=manager.getRepository(AreaLegalRight); const right=await repo.createQueryBuilder('r').where('r.id=:id',{id}).setLock('pessimistic_write').getOne(); if(!right) throw registryNotFound('Derecho hidrocarburífero'); const before={...right}; if(dto.sourceDocumentId) await this.requireDocument(manager,dto.sourceDocumentId); const validFrom=dto.validFrom===undefined?right.validFrom:dto.validFrom; const validUntil=dto.validUntil===undefined?right.validUntil:dto.validUntil; this.validateDateRange(validFrom,validUntil); if(dto.name!==undefined)right.name=dto.name;if(dto.instrumentNumber!==undefined)right.instrumentNumber=dto.instrumentNumber;if(dto.validFrom!==undefined)right.validFrom=dto.validFrom;if(dto.validUntil!==undefined)right.validUntil=dto.validUntil;if(dto.status!==undefined)right.status=dto.status;if(dto.sourceDocumentId!==undefined)right.sourceDocumentId=dto.sourceDocumentId;if(dto.notes!==undefined)right.notes=dto.notes;right.updatedBy=principal.userId;await repo.save(right);const versionNumber=await this.history.capture(manager,right.areaId,AssetVersionChangeType.REGISTRY_UPDATED,principal,request);await this.audit.record({...administrationAuditContext(principal,request),action:AuditAction.AREA_LEGAL_RIGHT_UPDATED,entityType:'area_legal_right',entityId:id,beforeData:before,afterData:{...right},metadata:{reason:dto.reason,versionNumber}},manager);return right; }); } async addLegalRightOrganization(rightId:string,dto:AddAreaLegalRightOrganizationDto,principal:AuthPrincipal,request:RequestWithContext){ try { return await this.dataSource.transaction(async manager=>{ const right=await manager.getRepository(AreaLegalRight).createQueryBuilder('r').where('r.id=:rightId',{rightId}).setLock('pessimistic_write').getOne(); if(!right) throw registryNotFound('Derecho hidrocarburífero'); await this.requireRole(manager,dto.organizationId,AssetTypeOperationalRole.COMPANY); const validFrom=dto.validFrom??this.today(); const [overlap]=await manager.query(`SELECT 1 FROM area_legal_right_organizations WHERE right_id=$1 AND organization_id=$2 AND role=$3 AND daterange(valid_from,COALESCE(valid_until,'infinity'::date),'[]') && daterange($4::date,'infinity'::date,'[]') LIMIT 1`,[rightId,dto.organizationId,dto.role,validFrom]); if(overlap) throw new ConflictException({code:'LEGAL_RIGHT_ORGANIZATION_OVERLAP',message:'La participación se superpone con una vigencia histórica existente'}); const repo=manager.getRepository(AreaLegalRightOrganization); const item=repo.create({rightId,organizationId:dto.organizationId,role:dto.role,participationPercent:dto.participationPercent==null?null:String(dto.participationPercent),validFrom,validUntil:null,notes:dto.notes??null,endReason:null,createdBy:principal.userId,endedBy:null}); await repo.save(item); await this.captureAssets(manager,[right.areaId,dto.organizationId],principal,request); await this.audit.record({...administrationAuditContext(principal,request),action:AuditAction.ASSET_REGISTRY_UPDATED,entityType:'area_legal_right_organization',entityId:item.id,afterData:{...item}},manager); return item; }); } catch(error){ if(isUniqueViolation(error)) throw new ConflictException({code:'LEGAL_RIGHT_ORGANIZATION_EXISTS',message:'La organización ya tiene ese rol activo en el derecho'}); throw error; } } async endLegalRightOrganization(id:string,dto:EndAreaLegalRightOrganizationDto,principal:AuthPrincipal,request:RequestWithContext){ return this.dataSource.transaction(async manager=>{ const repo=manager.getRepository(AreaLegalRightOrganization); const item=await repo.createQueryBuilder('o').where('o.id=:id',{id}).setLock('pessimistic_write').getOne(); if(!item) throw registryNotFound('Participación legal'); if(item.validUntil) throw new ConflictException({code:'LEGAL_RIGHT_ORGANIZATION_ENDED',message:'La participación ya está finalizada'}); const end=dto.validUntil??this.today();this.validateEndDate(end,item.validFrom);item.validUntil=end;item.endReason=dto.reason;item.endedBy=principal.userId;await repo.save(item);const right=await manager.getRepository(AreaLegalRight).findOne({where:{id:item.rightId}});if(right)await this.captureAssets(manager,[right.areaId,item.organizationId],principal,request);await this.audit.record({...administrationAuditContext(principal,request),action:AuditAction.AREA_LEGAL_RIGHT_ORGANIZATION_ENDED,entityType:'area_legal_right_organization',entityId:id,afterData:{...item}},manager);return item; }); } private async captureAssets(manager:EntityManager,ids:string[],principal:AuthPrincipal,request:RequestWithContext){ for(const id of [...new Set(ids)]) await this.history.capture(manager,id,AssetVersionChangeType.REGISTRY_UPDATED,principal,request); } private async requireAsset(manager:EntityManager,id:string):Promise{ const asset=await manager.getRepository(Asset).findOne({where:{id}}); if(!asset) throw registryNotFound('Activo'); return asset; } private async requireRole(manager:EntityManager,id:string,role:AssetTypeOperationalRole):Promise{ const [row]=await manager.query(`SELECT a.id FROM assets a JOIN asset_types t ON t.id=a.asset_type_id WHERE a.id=$1 AND t.operational_role=$2`,[id,role]); if(!row) throw new BadRequestException({code:'ASSET_ROLE_INVALID',message:role===AssetTypeOperationalRole.AREA?'El activo debe ser un Área':'El activo debe ser una Organización'}); return this.requireAsset(manager,id); } private async requireDocument(manager:EntityManager,id:string){ const doc=await manager.getRepository(SourceDocument).findOne({where:{id}}); if(!doc) throw registryNotFound('Documento fuente'); return doc; } private today(){ return new Date().toISOString().slice(0,10); } private validateEndDate(end:string,start:string){ const today=this.today(); if(endtoday) throw new BadRequestException({code:'FUTURE_END_DATE',message:'La fecha de finalización no puede estar en el futuro'}); } private validateDateRange(start:string|null,end:string|null){ if(start&&end&&end