diff --git a/apps/api/src/test/java/com/orgmemory/api/assetregistry/AssetRegistryIntegrationTests.java b/apps/api/src/test/java/com/orgmemory/api/assetregistry/AssetRegistryIntegrationTests.java index 90a87c971..cb10af9c4 100644 --- a/apps/api/src/test/java/com/orgmemory/api/assetregistry/AssetRegistryIntegrationTests.java +++ b/apps/api/src/test/java/com/orgmemory/api/assetregistry/AssetRegistryIntegrationTests.java @@ -21,7 +21,7 @@ import com.orgmemory.core.assetregistry.api.AssetRole; import com.orgmemory.core.assetregistry.api.AssetType; import com.orgmemory.core.assetregistry.api.AssetUnavailableException; -import com.orgmemory.core.assetregistry.AssetAuthorizationConvergenceService; +import com.orgmemory.core.assetregistry.authorization.AssetAuthorizationConvergenceService; import com.orgmemory.core.assetregistry.AssetAvailability; import com.orgmemory.core.assetregistry.AssetCatalogSort; import com.orgmemory.core.assetregistry.AssetDeliveryService; diff --git a/apps/worker/src/main/java/com/orgmemory/worker/authorization/AssetAuthorizationConvergenceScheduler.java b/apps/worker/src/main/java/com/orgmemory/worker/authorization/AssetAuthorizationConvergenceScheduler.java index 75d876b8f..864843ac8 100644 --- a/apps/worker/src/main/java/com/orgmemory/worker/authorization/AssetAuthorizationConvergenceScheduler.java +++ b/apps/worker/src/main/java/com/orgmemory/worker/authorization/AssetAuthorizationConvergenceScheduler.java @@ -1,6 +1,6 @@ package com.orgmemory.worker.authorization; -import com.orgmemory.core.assetregistry.AssetAuthorizationConvergenceService; +import com.orgmemory.core.assetregistry.authorization.AssetAuthorizationConvergenceService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationBatch.java b/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationBatch.java deleted file mode 100644 index e96fa379a..000000000 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationBatch.java +++ /dev/null @@ -1,25 +0,0 @@ -package com.orgmemory.core.assetregistry; - -import com.orgmemory.core.authorization.RelationshipTuple; -import java.util.List; -import java.util.Objects; -import java.util.UUID; - -record AssetAuthorizationBatch( - UUID organizationId, - UUID assetId, - UUID claimToken, - List outboxIds, - List tuples) { - - AssetAuthorizationBatch { - organizationId = Objects.requireNonNull(organizationId, "organizationId"); - assetId = Objects.requireNonNull(assetId, "assetId"); - claimToken = Objects.requireNonNull(claimToken, "claimToken"); - outboxIds = List.copyOf(outboxIds); - tuples = List.copyOf(tuples); - if (outboxIds.size() != tuples.size()) { - throw new IllegalArgumentException("Outbox ids and tuples must have equal size"); - } - } -} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationTarget.java b/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationTarget.java deleted file mode 100644 index 638f71ef1..000000000 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationTarget.java +++ /dev/null @@ -1,12 +0,0 @@ -package com.orgmemory.core.assetregistry; - -import com.orgmemory.core.assetregistry.api.AssetType; -import java.util.UUID; - -record AssetAuthorizationTarget( - UUID organizationId, - UUID assetId, - UUID knowledgeSpaceId, - AssetType type, - boolean authorizationReady) { -} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetRepository.java b/core/src/main/java/com/orgmemory/core/assetregistry/AssetCatalogReadModelRepository.java similarity index 90% rename from core/src/main/java/com/orgmemory/core/assetregistry/AssetRepository.java rename to core/src/main/java/com/orgmemory/core/assetregistry/AssetCatalogReadModelRepository.java index ed9b6ccc1..2d40f3276 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetRepository.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/AssetCatalogReadModelRepository.java @@ -1,19 +1,16 @@ package com.orgmemory.core.assetregistry; import com.orgmemory.core.assetregistry.api.AssetType; -import jakarta.persistence.LockModeType; import java.util.Collection; import java.util.List; -import java.util.Optional; import java.util.UUID; import org.springframework.data.domain.Page; import org.springframework.data.domain.Pageable; -import org.springframework.data.jpa.repository.JpaRepository; -import org.springframework.data.jpa.repository.Lock; import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.Repository; import org.springframework.data.repository.query.Param; -interface AssetRepository extends JpaRepository { +interface AssetCatalogReadModelRepository extends Repository { String CATALOG_FROM_AND_PREDICATES = """ from Asset asset @@ -109,22 +106,6 @@ like concat('%', :query, '%') asset.id asc """; - Optional findByIdAndOrganizationId(UUID id, UUID organizationId); - - @Lock(LockModeType.PESSIMISTIC_WRITE) - @Query(""" - select asset - from Asset asset - where asset.id = :id - and asset.organizationId = :organizationId - """) - Optional findForUpdate( - @Param("id") UUID id, - @Param("organizationId") UUID organizationId); - - Optional findByOrganizationIdAndNamespaceAndSlug( - UUID organizationId, String namespace, String slug); - @Query(""" select new com.orgmemory.core.assetregistry.AssetSummary( asset.id, diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetDraftRepository.java b/core/src/main/java/com/orgmemory/core/assetregistry/AssetDraftRepository.java index cc09b7ab7..67b072adb 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetDraftRepository.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/AssetDraftRepository.java @@ -2,9 +2,23 @@ import java.util.Optional; import java.util.UUID; +import jakarta.persistence.LockModeType; import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Lock; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; interface AssetDraftRepository extends JpaRepository { Optional findByAssetIdAndOrganizationId(UUID assetId, UUID organizationId); + + @Lock(LockModeType.PESSIMISTIC_WRITE) + @Query(""" + select draft from AssetDraft draft + where draft.assetId = :assetId + and draft.organizationId = :organizationId + """) + Optional findForUpdate( + @Param("assetId") UUID assetId, + @Param("organizationId") UUID organizationId); } diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetRegistryCoordinator.java b/core/src/main/java/com/orgmemory/core/assetregistry/AssetRegistryCoordinator.java index b5d8f6102..405dcca49 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetRegistryCoordinator.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/AssetRegistryCoordinator.java @@ -1,12 +1,19 @@ package com.orgmemory.core.assetregistry; import com.orgmemory.core.assetregistry.api.AssetConflictException; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationTarget; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationTargetQuery; +import com.orgmemory.core.assetregistry.api.AssetIdentity; +import com.orgmemory.core.assetregistry.api.AssetIdentityQuery; import com.orgmemory.core.assetregistry.api.AssetNotFoundException; +import com.orgmemory.core.assetregistry.api.AssetPortfolioCommand; +import com.orgmemory.core.assetregistry.api.AssetRegistrationCommand; import com.orgmemory.core.assetregistry.api.AssetRole; +import com.orgmemory.core.assetregistry.api.AssetRoleCommand; +import com.orgmemory.core.assetregistry.api.AssetRoleQuery; import com.orgmemory.core.assetregistry.api.AssetType; import com.orgmemory.core.assetregistry.api.AssetUnavailableException; import com.orgmemory.core.authorization.PrincipalRef; -import com.orgmemory.core.authorization.RelationshipTuple; import com.orgmemory.core.organization.CurrentActor; import com.orgmemory.core.shared.error.BusinessValidationException; import java.time.Instant; @@ -27,15 +34,19 @@ class AssetRegistryCoordinator { static final String REVIEW_POLICY_VERSION = "asset-review-v1"; - private final AssetRepository assets; + private final AssetRegistrationCommand registrations; + private final AssetIdentityQuery identities; + private final AssetAuthorizationTargetQuery targets; + private final AssetRoleCommand roleCommands; + private final AssetRoleQuery roleQueries; + private final AssetPortfolioCommand portfolios; + private final AssetCatalogReadModelRepository catalog; private final AssetDraftRepository drafts; private final AssetRevisionRepository revisions; private final AssetReviewCaseRepository reviews; private final AssetReviewDecisionRepository decisions; private final AssetReleaseRepository releases; private final AssetReleaseAvailabilityRepository availability; - private final AssetRoleAssignmentRepository roles; - private final AssetAuthorizationOutboxRepository outbox; private final AssetAuditEventRepository audit; private final AssetPayloadReferenceRepository payloadReferences; private final SkillPackageSupersessionRepository packageSupersessions; @@ -44,30 +55,38 @@ class AssetRegistryCoordinator { private final AssetPayloadDigester digester; AssetRegistryCoordinator( - AssetRepository assets, + AssetRegistrationCommand registrations, + AssetIdentityQuery identities, + AssetAuthorizationTargetQuery targets, + AssetRoleCommand roleCommands, + AssetRoleQuery roleQueries, + AssetPortfolioCommand portfolios, + AssetCatalogReadModelRepository catalog, AssetDraftRepository drafts, AssetRevisionRepository revisions, AssetReviewCaseRepository reviews, AssetReviewDecisionRepository decisions, AssetReleaseRepository releases, AssetReleaseAvailabilityRepository availability, - AssetRoleAssignmentRepository roles, - AssetAuthorizationOutboxRepository outbox, AssetAuditEventRepository audit, AssetPayloadReferenceRepository payloadReferences, SkillPackageSupersessionRepository packageSupersessions, AssetTypeProfileRegistry profiles, SkillPackageSpecReader skillPackages, AssetPayloadDigester digester) { - this.assets = assets; + this.registrations = registrations; + this.identities = identities; + this.targets = targets; + this.roleCommands = roleCommands; + this.roleQueries = roleQueries; + this.portfolios = portfolios; + this.catalog = catalog; this.drafts = drafts; this.revisions = revisions; this.reviews = reviews; this.decisions = decisions; this.releases = releases; this.availability = availability; - this.roles = roles; - this.outbox = outbox; this.audit = audit; this.payloadReferences = payloadReferences; this.packageSupersessions = packageSupersessions; @@ -125,30 +144,25 @@ private UUID create( } AssetPayloadDigester.CanonicalAssetPayload canonical = validateDraft(type, input); - Asset asset; + UUID assetId; try { - asset = new Asset( + assetId = registrations.register(new AssetRegistrationCommand.NewAsset( actor.organizationId(), type, namespace, slug, - knowledgeSpaceId); + knowledgeSpaceId, + actor.principal(), + actor.userId())); } catch (IllegalArgumentException invalidIdentity) { throw new BusinessValidationException( "asset.identity-invalid", "The Asset namespace, slug, type, or Knowledge Space is invalid", invalidIdentity); } - try { - assets.saveAndFlush(asset); - } catch (DataIntegrityViolationException duplicate) { - throw new AssetConflictException( - "An Asset already uses this namespace and slug, or the Space is unavailable", - duplicate); - } AssetDraft draft = drafts.saveAndFlush(new AssetDraft( actor.organizationId(), - asset.getId(), + assetId, canonical.title(), canonical.summary(), canonical.classification(), @@ -165,59 +179,19 @@ private UUID create( payloadReferences.saveAndFlush( AssetPayloadReference.forDraft(draft, storedPackage)); } - Instant now = Instant.now(); - AssetRoleAssignment owner = roles.saveAndFlush(new AssetRoleAssignment( - actor.organizationId(), - asset.getId(), - actor.principal(), - AssetRole.OWNER, - actor.userId(), - now)); - String object = "asset:" + asset.getId(); - outbox.saveAllAndFlush(List.of( - new AssetAuthorizationOutbox( - actor.organizationId(), - asset.getId(), - null, - RelationshipTuple.of( - "organization:" + actor.organizationId(), - "organization", - object)), - new AssetAuthorizationOutbox( - actor.organizationId(), - asset.getId(), - null, - RelationshipTuple.of( - "knowledge_space:" + knowledgeSpaceId, - "space", - object)), - new AssetAuthorizationOutbox( - actor.organizationId(), - asset.getId(), - owner.getId(), - RelationshipTuple.of( - actor.principal().openFgaUser(), - AssetRole.OWNER.relation(), - object)))); recordAudit( actor, - asset.getId(), + assetId, "ASSET_CREATED", "DRAFT", draft.getId(), "{\"permission\":\"can_create_asset\"}"); - return asset.getId(); + return assetId; } @Transactional(readOnly = true, propagation = Propagation.REQUIRES_NEW) Optional target(UUID organizationId, UUID assetId) { - return assets.findByIdAndOrganizationId(assetId, organizationId) - .map(asset -> new AssetAuthorizationTarget( - asset.getOrganizationId(), - asset.getId(), - asset.getKnowledgeSpaceId(), - asset.getType(), - asset.isAuthorizationReady())); + return targets.find(organizationId, assetId); } @Transactional(readOnly = true, propagation = Propagation.REQUIRES_NEW) @@ -230,7 +204,7 @@ List summaries( return List.of(); } String normalizedQuery = query == null ? "" : query.trim().toLowerCase(java.util.Locale.ROOT); - return assets.searchAuthorized( + return catalog.searchAuthorized( organizationId, ids, normalizedQuery, @@ -266,7 +240,7 @@ AssetSummaryPage ownedSummaryPage( if (visibleIds.isEmpty()) { return AssetSummaryPage.empty(page, pageSize, sort); } - var ids = new java.util.LinkedHashSet<>(roles.findActiveAssetIdsForUserRole( + var ids = new java.util.LinkedHashSet<>(roleQueries.activeAssetIdsForUserRole( organizationId, userId.toString(), AssetRole.OWNER, @@ -277,7 +251,7 @@ AssetSummaryPage ownedSummaryPage( } String normalizedQuery = query == null ? "" : query.trim().toLowerCase(java.util.Locale.ROOT); - Page result = assets.searchOwnedSummaries( + Page result = catalog.searchOwnedSummaries( organizationId, ids, normalizedQuery, @@ -307,7 +281,7 @@ AssetRecommendationPage recommendationPage( } String normalizedQuery = query == null ? "" : query.trim().toLowerCase(java.util.Locale.ROOT); - Page result = assets.searchAuthorizedRecommendations( + Page result = catalog.searchAuthorizedRecommendations( organizationId, ids, normalizedQuery, @@ -369,7 +343,7 @@ private static boolean matchesRecommendation( @Transactional(readOnly = true, propagation = Propagation.REQUIRES_NEW) AssetView view(UUID organizationId, UUID assetId) { - Asset asset = requiredAsset(organizationId, assetId); + AssetIdentity asset = requiredAsset(organizationId, assetId); return view(asset); } @@ -378,7 +352,7 @@ AssetConsumptionRelease consumptionRelease( UUID organizationId, UUID assetId, UUID releaseId) { - Asset asset = requiredAsset(organizationId, assetId); + AssetIdentity asset = requiredAsset(organizationId, assetId); AssetRelease release = releases .findByIdAndAssetIdAndOrganizationId( releaseId, assetId, organizationId) @@ -389,12 +363,12 @@ AssetConsumptionRelease consumptionRelease( "The requested Asset release is not available for new use"); } return new AssetConsumptionRelease( - asset.getId(), + asset.id(), release.getId(), release.getRevisionId(), - asset.getType(), - asset.getNamespace(), - asset.getSlug(), + asset.type(), + asset.namespace(), + asset.slug(), release.getVersionLabel(), release.getPublicationMode(), release.getTitle(), @@ -410,7 +384,7 @@ AssetConsumptionRelease consumptionRelease( @Transactional(readOnly = true, propagation = Propagation.REQUIRES_NEW) AssetConsumptionRelease latestConsumptionRelease( UUID organizationId, UUID assetId) { - Asset asset = requiredAsset(organizationId, assetId); + AssetIdentity asset = requiredAsset(organizationId, assetId); AssetRelease release = releases .findLatestUsable( assetId, @@ -421,7 +395,7 @@ AssetConsumptionRelease latestConsumptionRelease( .findFirst() .orElseThrow(AssetNotFoundException::new); return consumptionRelease( - organizationId, asset.getId(), release.getId()); + organizationId, asset.id(), release.getId()); } @Transactional(propagation = Propagation.REQUIRES_NEW) @@ -430,8 +404,8 @@ AssetView updateDraft( UUID assetId, long expectedLockVersion, AssetDraftInput input) { - Asset asset = requiredAsset(actor.organizationId(), assetId); - if (asset.getType() == AssetType.SKILL) { + AssetIdentity asset = requiredAsset(actor.organizationId(), assetId); + if (asset.type() == AssetType.SKILL) { throw new BusinessValidationException( "skill.generic-edit-unsupported", "Skill drafts cannot be changed through the generic payload endpoint"); @@ -442,7 +416,7 @@ AssetView updateDraft( "The Asset draft changed; reload it before saving"); } AssetPayloadDigester.CanonicalAssetPayload canonical = - validateDraft(asset.getType(), input); + validateDraft(asset.type(), input); draft.update( canonical.title(), canonical.summary(), @@ -468,13 +442,13 @@ SkillDraftReplacement replaceSkillDraft( long expectedLockVersion, AssetDraftInput input, SkillPackageStoragePort.StoredSkillPackage storedPackage) { - Asset asset = requiredAssetForUpdate(actor.organizationId(), assetId); - if (asset.getType() != AssetType.SKILL) { + AssetDraft draft = requiredDraftForUpdate(actor.organizationId(), assetId); + AssetIdentity asset = requiredAsset(actor.organizationId(), assetId); + if (asset.type() != AssetType.SKILL) { throw new BusinessValidationException( "skill.package-replacement-unsupported", "Only Skill drafts can replace a package"); } - AssetDraft draft = requiredDraft(actor.organizationId(), assetId); if (draft.getVersion() != expectedLockVersion) { throw new AssetConflictException( "The Skill draft changed; reload it before replacing the package"); @@ -524,13 +498,13 @@ SkillDraftReplacement replaceSkillDraft( @Transactional(propagation = Propagation.REQUIRES_NEW) AssetView submit(CurrentActor actor, UUID assetId, String changeNote) { String validatedChangeNote = requireChangeNote(changeNote); - Asset asset = requiredAssetForUpdate(actor.organizationId(), assetId); + AssetDraft draft = requiredDraftForUpdate(actor.organizationId(), assetId); + AssetIdentity asset = requiredAsset(actor.organizationId(), assetId); if (reviews.existsByAssetIdAndOrganizationIdAndState( assetId, actor.organizationId(), AssetReviewState.IN_REVIEW)) { throw new AssetConflictException("This Asset already has a revision in review"); } - AssetDraft draft = requiredDraft(actor.organizationId(), assetId); - profiles.require(asset.getType()).requireSupported(draft.getSchemaVersion()); + profiles.require(asset.type()).requireSupported(draft.getSchemaVersion()); AssetPayloadDigester.CanonicalAssetPayload canonical = digester.canonicalize( draft.getTitle(), draft.getSummary(), @@ -555,7 +529,7 @@ AssetView submit(CurrentActor actor, UUID assetId, String changeNote) { } AssetReviewCase review = reviews.saveAndFlush(new AssetReviewCase( revision, REVIEW_POLICY_VERSION, actor.userId())); - if (asset.getType() == AssetType.SKILL) { + if (asset.type() == AssetType.SKILL) { SkillPackageSpec spec = skillSpec(draft.getPayload()); AssetPayloadReference draftReference = payloadReferences .findByDraftIdAndOrganizationId( @@ -573,7 +547,7 @@ AssetView submit(CurrentActor actor, UUID assetId, String changeNote) { "REVIEW_CASE", review.getId(), "{\"digest\":\"" + revision.getDigest() + "\"}"); - return view(asset); + return view(requiredAsset(actor.organizationId(), assetId)); } @Transactional(propagation = Propagation.REQUIRES_NEW) @@ -583,7 +557,7 @@ AssetView decide( UUID reviewCaseId, AssetReviewDecisionType decision, String comment) { - Asset asset = requiredAsset(actor.organizationId(), assetId); + AssetIdentity asset = requiredAsset(actor.organizationId(), assetId); AssetReviewCase review = reviews.findByIdAndAssetIdAndOrganizationId( reviewCaseId, assetId, actor.organizationId()) .orElseThrow(AssetNotFoundException::new); @@ -707,7 +681,7 @@ AssetView publish( UUID revisionId, String versionLabel) { String validatedVersionLabel = validatedVersionLabel(versionLabel); - Asset asset = requiredAsset(actor.organizationId(), assetId); + AssetIdentity asset = requiredAsset(actor.organizationId(), assetId); AssetRevision revision = revisions.findByIdAndAssetIdAndOrganizationId( revisionId, assetId, actor.organizationId()) .orElseThrow(AssetNotFoundException::new); @@ -748,7 +722,7 @@ AssetView publish( "RELEASE", release.getId(), "{\"digest\":\"" + release.getDigest() + "\"}"); - return view(asset); + return view(requiredAsset(actor.organizationId(), assetId)); } @Transactional(propagation = Propagation.REQUIRES_NEW) @@ -757,8 +731,9 @@ AssetView publishSkillDraft( UUID assetId, String versionLabel) { String validatedVersionLabel = validatedVersionLabel(versionLabel); - Asset asset = requiredAssetForUpdate(actor.organizationId(), assetId); - if (asset.getType() != AssetType.SKILL) { + AssetDraft draft = requiredDraftForUpdate(actor.organizationId(), assetId); + AssetIdentity asset = requiredAsset(actor.organizationId(), assetId); + if (asset.type() != AssetType.SKILL) { throw new BusinessValidationException( "skill.direct-publish-unsupported", "Direct publication is available only for Skill Assets"); @@ -768,7 +743,6 @@ AssetView publishSkillDraft( throw new AssetConflictException( "This Skill already has a revision in review; finish or cancel it first"); } - AssetDraft draft = requiredDraft(actor.organizationId(), assetId); profiles.require(AssetType.SKILL).requireSupported(draft.getSchemaVersion()); AssetPayloadDigester.CanonicalAssetPayload canonical = digester.canonicalize( draft.getTitle(), @@ -819,12 +793,12 @@ AssetView publishSkillDraft( "{\"policy\":\"skill-direct-v1\"," + "\"permission\":\"can_publish_skill\"," + "\"digest\":\"" + release.getDigest() + "\"}"); - return view(asset); + return view(requiredAsset(actor.organizationId(), assetId)); } private AssetRelease createRelease( CurrentActor actor, - Asset asset, + AssetIdentity asset, AssetRevision revision, String versionLabel, AssetPublicationMode publicationMode, @@ -834,7 +808,7 @@ private AssetRelease createRelease( try { release = releases.saveAndFlush(new AssetRelease( revision, - releases.maxSequence(asset.getId(), actor.organizationId()) + 1, + releases.maxSequence(asset.id(), actor.organizationId()) + 1, versionLabel, publicationMode, actor.userId(), @@ -857,7 +831,7 @@ private AssetRelease createRelease( availabilityReason, actor.userId(), now)); - if (asset.getType() == AssetType.SKILL) { + if (asset.type() == AssetType.SKILL) { AssetPayloadReference revisionReference = payloadReferences .findByRevisionIdAndOrganizationId( revision.getId(), actor.organizationId()) @@ -868,8 +842,7 @@ private AssetRelease createRelease( payloadReferences.saveAndFlush( AssetPayloadReference.forRelease(release, revisionReference)); } - asset.activate(); - assets.save(asset); + portfolios.activateAfterRelease(actor.organizationId(), asset.id()); return release; } @@ -883,7 +856,7 @@ AssetView changeAvailability( if (next == AssetAvailability.AVAILABLE) { throw new IllegalArgumentException("A release cannot be made available again"); } - Asset asset = requiredAsset(actor.organizationId(), assetId); + AssetIdentity asset = requiredAsset(actor.organizationId(), assetId); AssetRelease release = releases.findByIdAndAssetIdAndOrganizationId( releaseId, assetId, actor.organizationId()) .orElseThrow(AssetNotFoundException::new); @@ -896,12 +869,13 @@ AssetView changeAvailability( Instant now = Instant.now(); availability.saveAndFlush(new AssetReleaseAvailabilityEvent( release, next, requireReason(reason), actor.userId(), now)); + AssetIdentity updated; if (next == AssetAvailability.WITHDRAWN && allReleasesWithdrawn(asset, releaseId)) { - asset.retire(); + updated = portfolios.retireAfterFinalWithdrawal(actor.organizationId(), assetId); } else { - asset.startSunsetting(); + updated = portfolios.startSunsettingAfterReleaseChange( + actor.organizationId(), assetId); } - assets.save(asset); recordAudit( actor, assetId, @@ -909,7 +883,7 @@ AssetView changeAvailability( "RELEASE", releaseId, "{\"previous\":\"" + current + "\"}"); - return view(asset); + return view(updated); } @Transactional(propagation = Propagation.REQUIRES_NEW) @@ -918,45 +892,29 @@ UUID assignRole( UUID assetId, PrincipalRef principal, AssetRole role) { - Asset asset = requiredAssetForUpdate(actor.organizationId(), assetId); - if (roles.findByAssetIdAndPrincipalTypeAndPrincipalIdAndRoleAndValidUntilIsNull( - assetId, principal.type(), principal.id(), role) - .isPresent()) { - throw new AssetConflictException("This active Asset role assignment already exists"); - } - Instant now = Instant.now(); - AssetRoleAssignment assignment = roles.saveAndFlush(new AssetRoleAssignment( + UUID assignmentId = roleCommands.assign(new AssetRoleCommand.Assignment( actor.organizationId(), assetId, principal, role, - actor.userId(), - now)); - outbox.saveAndFlush(new AssetAuthorizationOutbox( - actor.organizationId(), - assetId, - assignment.getId(), - RelationshipTuple.of( - principal.openFgaUser(), - role.relation(), - "asset:" + assetId))); + actor.userId())); recordAudit( actor, assetId, "ROLE_ASSIGNED", "ROLE_ASSIGNMENT", - assignment.getId(), + assignmentId, "{\"role\":\"" + role + "\"}"); - return asset.getId(); + return assetId; } - private Asset requiredAsset(UUID organizationId, UUID assetId) { - return assets.findByIdAndOrganizationId(assetId, organizationId) + private AssetIdentity requiredAsset(UUID organizationId, UUID assetId) { + return identities.findById(organizationId, assetId) .orElseThrow(AssetNotFoundException::new); } - private Asset requiredAssetForUpdate(UUID organizationId, UUID assetId) { - return assets.findForUpdate(assetId, organizationId) + private AssetDraft requiredDraftForUpdate(UUID organizationId, UUID assetId) { + return drafts.findForUpdate(assetId, organizationId) .orElseThrow(AssetNotFoundException::new); } @@ -965,9 +923,9 @@ private AssetDraft requiredDraft(UUID organizationId, UUID assetId) { .orElseThrow(AssetNotFoundException::new); } - private boolean allReleasesWithdrawn(Asset asset, UUID releaseBeingWithdrawn) { + private boolean allReleasesWithdrawn(AssetIdentity asset, UUID releaseBeingWithdrawn) { return releases.findByAssetIdAndOrganizationIdOrderBySequenceDesc( - asset.getId(), asset.getOrganizationId()) + asset.id(), asset.organizationId()) .stream() .allMatch(release -> release.getId().equals(releaseBeingWithdrawn) || currentAvailability(release.getId()) == AssetAvailability.WITHDRAWN); @@ -981,28 +939,28 @@ private AssetAvailability currentAvailability(UUID releaseId) { "Asset release is missing availability history")); } - private AssetView view(Asset asset) { - AssetDraft draft = requiredDraft(asset.getOrganizationId(), asset.getId()); + private AssetView view(AssetIdentity asset) { + AssetDraft draft = requiredDraft(asset.organizationId(), asset.id()); List assetRevisions = revisions.findByAssetIdAndOrganizationIdOrderBySequenceDesc( - asset.getId(), asset.getOrganizationId()); + asset.id(), asset.organizationId()); List assetReviews = reviews.findByAssetIdAndOrganizationIdOrderByCreatedAtDesc( - asset.getId(), asset.getOrganizationId()); + asset.id(), asset.organizationId()); List assetReleases = releases.findByAssetIdAndOrganizationIdOrderBySequenceDesc( - asset.getId(), asset.getOrganizationId()); - List assignments = - roles.findByAssetIdOrderByValidFromAsc(asset.getId()); + asset.id(), asset.organizationId()); Instant viewedAt = Instant.now(); + AssetRoleQuery.RoleHistory roleHistory = + roleQueries.history(asset.organizationId(), asset.id(), viewedAt); return new AssetView( - asset.getId(), - asset.getType(), - asset.getNamespace(), - asset.getSlug(), - asset.getKnowledgeSpaceId(), - asset.getPortfolioState(), - asset.isAuthorizationReady(), + asset.id(), + asset.type(), + asset.namespace(), + asset.slug(), + asset.knowledgeSpaceId(), + asset.portfolioState(), + asset.authorizationReady(), new AssetView.Draft( draft.getId(), draft.getVersion(), @@ -1016,31 +974,19 @@ private AssetView view(Asset asset) { assetRevisions.stream().map(AssetRegistryCoordinator::revisionView).toList(), assetReviews.stream().map(this::reviewView).toList(), assetReleases.stream().map(this::releaseView).toList(), - ownershipHealth(assignments, viewedAt), - assignments.stream().map(AssetRegistryCoordinator::roleView).toList()); + ownershipHealth(roleHistory.ownershipHealth()), + roleHistory.assignments().stream() + .map(AssetRegistryCoordinator::roleView) + .toList()); } private static AssetView.OwnershipHealth ownershipHealth( - List assignments, Instant viewedAt) { - boolean ownerPresent = hasActiveRole( - assignments, AssetRole.OWNER, viewedAt); - boolean backupOwnerPresent = hasActiveRole( - assignments, AssetRole.BACKUP_OWNER, viewedAt); + AssetRoleQuery.OwnershipHealth ownershipHealth) { return new AssetView.OwnershipHealth( - ownerPresent, - backupOwnerPresent, - !ownerPresent && !backupOwnerPresent, - !ownerPresent || !backupOwnerPresent); - } - - private static boolean hasActiveRole( - List assignments, - AssetRole role, - Instant viewedAt) { - return assignments.stream().anyMatch(assignment -> - assignment.getRole() == role - && (assignment.getValidUntil() == null - || assignment.getValidUntil().isAfter(viewedAt))); + ownershipHealth.ownerPresent(), + ownershipHealth.backupOwnerPresent(), + ownershipHealth.orphaned(), + ownershipHealth.continuityAtRisk()); } private AssetView.Review reviewView(AssetReviewCase review) { @@ -1125,16 +1071,16 @@ private static AssetView.Revision revisionView(AssetRevision revision) { revision.getCreatedAt()); } - private static AssetView.RoleAssignment roleView(AssetRoleAssignment assignment) { + private static AssetView.RoleAssignment roleView(AssetRoleQuery.RoleAssignment assignment) { return new AssetView.RoleAssignment( - assignment.getId(), - assignment.getPrincipalType(), - assignment.getPrincipalId(), - assignment.getRole(), - assignment.getValidFrom(), - assignment.getValidUntil(), - assignment.getAssignedByUserId(), - assignment.getProjectedAt()); + assignment.id(), + assignment.principalType(), + assignment.principalId(), + assignment.role(), + assignment.validFrom(), + assignment.validUntil(), + assignment.assignedByUserId(), + assignment.projectedAt()); } private static String requireReason(String reason) { diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetRegistryService.java b/core/src/main/java/com/orgmemory/core/assetregistry/AssetRegistryService.java index 1849c5c41..dbdb76095 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetRegistryService.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/AssetRegistryService.java @@ -1,5 +1,7 @@ package com.orgmemory.core.assetregistry; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationTarget; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationProjectionCommand; import com.orgmemory.core.assetregistry.api.AssetConflictException; import com.orgmemory.core.assetregistry.api.AssetNotFoundException; import com.orgmemory.core.assetregistry.api.AssetRole; @@ -42,13 +44,13 @@ public class AssetRegistryService { PermissionKey.of("can_manage_roles"); private final AssetRegistryCoordinator coordinator; - private final AssetAuthorizationProjectionService projection; + private final AssetAuthorizationProjectionCommand projection; private final RelationshipAuthorizationPort authorization; private final RelationshipAuthorizationSetPort authorizationSets; AssetRegistryService( AssetRegistryCoordinator coordinator, - AssetAuthorizationProjectionService projection, + AssetAuthorizationProjectionCommand projection, RelationshipAuthorizationPort authorization, RelationshipAuthorizationSetPort authorizationSets) { this.coordinator = coordinator; diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/SkillDistributionService.java b/core/src/main/java/com/orgmemory/core/assetregistry/SkillDistributionService.java index 0ceda4f64..f63d7f248 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/SkillDistributionService.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/SkillDistributionService.java @@ -1,5 +1,7 @@ package com.orgmemory.core.assetregistry; +import com.orgmemory.core.assetregistry.api.AssetIdentity; +import com.orgmemory.core.assetregistry.api.AssetIdentityQuery; import com.orgmemory.core.assetregistry.api.AssetNotFoundException; import com.orgmemory.core.assetregistry.api.AssetType; import com.orgmemory.core.assetregistry.api.AssetUnavailableException; @@ -23,7 +25,7 @@ public class SkillDistributionService { Pattern.compile("[a-z0-9]+(?:[._-][a-z0-9]+)*"); private final AssetRegistryService assets; - private final AssetRepository assetRepository; + private final AssetIdentityQuery identities; private final AssetReleaseRepository releaseRepository; private final AssetPayloadReferenceRepository references; private final SkillPackageSpecReader specs; @@ -31,13 +33,13 @@ public class SkillDistributionService { SkillDistributionService( AssetRegistryService assets, - AssetRepository assetRepository, + AssetIdentityQuery identities, AssetReleaseRepository releaseRepository, AssetPayloadReferenceRepository references, SkillPackageSpecReader specs, SkillPackageStoragePort storage) { this.assets = assets; - this.assetRepository = assetRepository; + this.identities = identities; this.releaseRepository = releaseRepository; this.references = references; this.specs = specs; @@ -59,12 +61,12 @@ public SkillInstallManifest manifest( String slug, String version) { Objects.requireNonNull(actor, "actor"); - Asset asset = assetRepository - .findByOrganizationIdAndNamespaceAndSlug( + AssetIdentity asset = identities + .findByCoordinate( actor.organizationId(), normalizeCoordinate(namespace, "namespace"), normalizeCoordinate(slug, "slug")) - .filter(value -> value.getType() == AssetType.SKILL) + .filter(value -> value.type() == AssetType.SKILL) .orElseThrow(AssetNotFoundException::new); String versionLabel; try { @@ -74,11 +76,11 @@ public SkillInstallManifest manifest( } AssetRelease release = releaseRepository .findByAssetIdAndOrganizationIdAndVersionLabel( - asset.getId(), + asset.id(), actor.organizationId(), versionLabel) .orElseThrow(AssetNotFoundException::new); - return manifest(actor, asset.getId(), release.getId()); + return manifest(actor, asset.id(), release.getId()); } public SkillPackageContent open( diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetAuthorizationProjectionCommand.java b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetAuthorizationProjectionCommand.java new file mode 100644 index 000000000..d44422f62 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetAuthorizationProjectionCommand.java @@ -0,0 +1,9 @@ +package com.orgmemory.core.assetregistry.api; + +import java.util.UUID; + +/** Projects the canonical authorization relationships for one Asset. */ +public interface AssetAuthorizationProjectionCommand { + + void project(UUID organizationId, UUID assetId); +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetAuthorizationTarget.java b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetAuthorizationTarget.java new file mode 100644 index 000000000..0e9cec1f9 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetAuthorizationTarget.java @@ -0,0 +1,19 @@ +package com.orgmemory.core.assetregistry.api; + +import java.util.Objects; +import java.util.UUID; + +public record AssetAuthorizationTarget( + UUID organizationId, + UUID assetId, + UUID knowledgeSpaceId, + AssetType type, + boolean authorizationReady) { + + public AssetAuthorizationTarget { + Objects.requireNonNull(organizationId, "organizationId"); + Objects.requireNonNull(assetId, "assetId"); + Objects.requireNonNull(knowledgeSpaceId, "knowledgeSpaceId"); + Objects.requireNonNull(type, "type"); + } +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetAuthorizationTargetQuery.java b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetAuthorizationTargetQuery.java new file mode 100644 index 000000000..96616c067 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetAuthorizationTargetQuery.java @@ -0,0 +1,10 @@ +package com.orgmemory.core.assetregistry.api; + +import java.util.Optional; +import java.util.UUID; + +/** Fail-closed authorization target lookup backed by canonical Asset identity. */ +public interface AssetAuthorizationTargetQuery { + + Optional find(UUID organizationId, UUID assetId); +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetIdentity.java b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetIdentity.java new file mode 100644 index 000000000..c13b87061 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetIdentity.java @@ -0,0 +1,26 @@ +package com.orgmemory.core.assetregistry.api; + +import java.util.Objects; +import java.util.UUID; + +/** Immutable parent-facing view of canonical Asset identity and portfolio state. */ +public record AssetIdentity( + UUID organizationId, + UUID id, + AssetType type, + String namespace, + String slug, + UUID knowledgeSpaceId, + AssetPortfolioState portfolioState, + boolean authorizationReady) { + + public AssetIdentity { + Objects.requireNonNull(organizationId, "organizationId"); + Objects.requireNonNull(id, "id"); + Objects.requireNonNull(type, "type"); + Objects.requireNonNull(namespace, "namespace"); + Objects.requireNonNull(slug, "slug"); + Objects.requireNonNull(knowledgeSpaceId, "knowledgeSpaceId"); + Objects.requireNonNull(portfolioState, "portfolioState"); + } +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetIdentityQuery.java b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetIdentityQuery.java new file mode 100644 index 000000000..7640acfbf --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetIdentityQuery.java @@ -0,0 +1,13 @@ +package com.orgmemory.core.assetregistry.api; + +import java.util.Optional; +import java.util.UUID; + +/** Read access to canonical Asset identity without exposing Kernel persistence. */ +public interface AssetIdentityQuery { + + Optional findById(UUID organizationId, UUID assetId); + + Optional findByCoordinate( + UUID organizationId, String namespace, String slug); +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetPortfolioCommand.java b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetPortfolioCommand.java new file mode 100644 index 000000000..fe965674e --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetPortfolioCommand.java @@ -0,0 +1,13 @@ +package com.orgmemory.core.assetregistry.api; + +import java.util.UUID; + +/** Atomic Asset portfolio transitions derived from parent-owned release outcomes. */ +public interface AssetPortfolioCommand { + + AssetIdentity activateAfterRelease(UUID organizationId, UUID assetId); + + AssetIdentity startSunsettingAfterReleaseChange(UUID organizationId, UUID assetId); + + AssetIdentity retireAfterFinalWithdrawal(UUID organizationId, UUID assetId); +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetRegistrationCommand.java b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetRegistrationCommand.java new file mode 100644 index 000000000..9ccb7a3f8 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetRegistrationCommand.java @@ -0,0 +1,31 @@ +package com.orgmemory.core.assetregistry.api; + +import com.orgmemory.core.authorization.PrincipalRef; +import java.util.Objects; +import java.util.UUID; + +/** Atomic registration of Asset identity, initial ownership, and authorization intent. */ +public interface AssetRegistrationCommand { + + UUID register(NewAsset command); + + record NewAsset( + UUID organizationId, + AssetType type, + String namespace, + String slug, + UUID knowledgeSpaceId, + PrincipalRef owner, + UUID assignedByUserId) { + + public NewAsset { + Objects.requireNonNull(organizationId, "organizationId"); + Objects.requireNonNull(type, "type"); + Objects.requireNonNull(namespace, "namespace"); + Objects.requireNonNull(slug, "slug"); + Objects.requireNonNull(knowledgeSpaceId, "knowledgeSpaceId"); + Objects.requireNonNull(owner, "owner"); + Objects.requireNonNull(assignedByUserId, "assignedByUserId"); + } + } +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetRoleCommand.java b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetRoleCommand.java new file mode 100644 index 000000000..ac461f770 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetRoleCommand.java @@ -0,0 +1,27 @@ +package com.orgmemory.core.assetregistry.api; + +import com.orgmemory.core.authorization.PrincipalRef; +import java.util.Objects; +import java.util.UUID; + +/** Atomic accountable-role assignment plus its authorization intent. */ +public interface AssetRoleCommand { + + UUID assign(Assignment command); + + record Assignment( + UUID organizationId, + UUID assetId, + PrincipalRef principal, + AssetRole role, + UUID assignedByUserId) { + + public Assignment { + Objects.requireNonNull(organizationId, "organizationId"); + Objects.requireNonNull(assetId, "assetId"); + Objects.requireNonNull(principal, "principal"); + Objects.requireNonNull(role, "role"); + Objects.requireNonNull(assignedByUserId, "assignedByUserId"); + } + } +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetRoleQuery.java b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetRoleQuery.java new file mode 100644 index 000000000..b2de39ef6 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/api/AssetRoleQuery.java @@ -0,0 +1,39 @@ +package com.orgmemory.core.assetregistry.api; + +import java.time.Instant; +import java.util.List; +import java.util.UUID; + +/** Immutable accountable-role history and ownership health. */ +public interface AssetRoleQuery { + + RoleHistory history(UUID organizationId, UUID assetId, Instant viewedAt); + + List activeAssetIdsForUserRole( + UUID organizationId, String userId, AssetRole role, Instant viewedAt); + + record RoleHistory(OwnershipHealth ownershipHealth, List assignments) { + + public RoleHistory { + assignments = List.copyOf(assignments); + } + } + + record RoleAssignment( + UUID id, + String principalType, + String principalId, + AssetRole role, + Instant validFrom, + Instant validUntil, + UUID assignedByUserId, + Instant projectedAt) { + } + + record OwnershipHealth( + boolean ownerPresent, + boolean backupOwnerPresent, + boolean orphaned, + boolean continuityAtRisk) { + } +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationConvergenceReport.java b/core/src/main/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationConvergenceReport.java similarity index 87% rename from core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationConvergenceReport.java rename to core/src/main/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationConvergenceReport.java index f581ec2f8..bebc69db9 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationConvergenceReport.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationConvergenceReport.java @@ -1,4 +1,4 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.authorization; public record AssetAuthorizationConvergenceReport( int candidates, diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationConvergenceService.java b/core/src/main/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationConvergenceService.java similarity index 60% rename from core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationConvergenceService.java rename to core/src/main/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationConvergenceService.java index 39c34da7d..ff3eedef1 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationConvergenceService.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationConvergenceService.java @@ -1,24 +1,29 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.authorization; +import com.orgmemory.core.assetregistry.kernel.AssetAuthorizationBatch; +import com.orgmemory.core.assetregistry.kernel.AssetAuthorizationProjectionQueue; import java.util.List; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; @Service public class AssetAuthorizationConvergenceService { - private final AssetAuthorizationCoordinator coordinator; + private final AssetAuthorizationProjectionQueue queue; private final AssetAuthorizationProjectionService projection; AssetAuthorizationConvergenceService( - AssetAuthorizationCoordinator coordinator, + AssetAuthorizationProjectionQueue queue, AssetAuthorizationProjectionService projection) { - this.coordinator = coordinator; + this.queue = queue; this.projection = projection; } + @Transactional(propagation = Propagation.NEVER) public AssetAuthorizationConvergenceReport reconcile(int limit) { List candidates = - coordinator.claimPendingBatches(limit); + queue.claimPending(limit); int applied = 0; int failed = 0; for (AssetAuthorizationBatch candidate : candidates) { diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationProjectionService.java b/core/src/main/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationProjectionService.java similarity index 50% rename from core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationProjectionService.java rename to core/src/main/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationProjectionService.java index 2d8f665fd..d010a8506 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationProjectionService.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationProjectionService.java @@ -1,35 +1,46 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.authorization; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationProjectionCommand; import com.orgmemory.core.assetregistry.api.AssetUnavailableException; +import com.orgmemory.core.assetregistry.kernel.AssetAuthorizationBatch; +import com.orgmemory.core.assetregistry.kernel.AssetAuthorizationProjectionQueue; import com.orgmemory.core.authorization.RelationshipTupleWritePort; import com.orgmemory.core.authorization.RelationshipTupleWriteRequest; import com.orgmemory.core.authorization.RelationshipTupleWriteResult; import java.util.Objects; import java.util.UUID; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; @Service -class AssetAuthorizationProjectionService { +class AssetAuthorizationProjectionService implements AssetAuthorizationProjectionCommand { - private final AssetAuthorizationCoordinator coordinator; + private static final Logger log = + LoggerFactory.getLogger(AssetAuthorizationProjectionService.class); + + private final AssetAuthorizationProjectionQueue queue; private final RelationshipTupleWritePort relationshipTuples; AssetAuthorizationProjectionService( - AssetAuthorizationCoordinator coordinator, + AssetAuthorizationProjectionQueue queue, RelationshipTupleWritePort relationshipTuples) { - this.coordinator = coordinator; + this.queue = queue; this.relationshipTuples = relationshipTuples; } - void project(UUID organizationId, UUID assetId) { - AssetAuthorizationBatch batch = coordinator.startAttempt(organizationId, assetId); - if (batch.tuples().isEmpty()) { - throw new AssetUnavailableException( - "Asset authorization is already being projected"); - } + @Override + @Transactional(propagation = Propagation.NEVER) + public void project(UUID organizationId, UUID assetId) { + AssetAuthorizationBatch batch = queue.claimForAsset(organizationId, assetId) + .orElseThrow(() -> new AssetUnavailableException( + "Asset authorization is already being projected")); project(batch); } + @Transactional(propagation = Propagation.NEVER) void project(AssetAuthorizationBatch batch) { RelationshipTupleWriteResult result; try { @@ -37,7 +48,12 @@ void project(AssetAuthorizationBatch batch) { relationshipTuples.write(new RelationshipTupleWriteRequest(batch.tuples())), "relationship tuple write result"); } catch (RuntimeException exception) { - coordinator.recordFailure( + log.warn( + "Asset authorization projection failed for organization {} and asset {}", + batch.organizationId(), + batch.assetId(), + exception); + queue.fail( batch, "OPENFGA_WRITE_FAILED", "The Asset authorization relationship could not be applied"); @@ -45,13 +61,13 @@ void project(AssetAuthorizationBatch batch) { "Asset authorization is waiting for projection", exception); } if (!result.applied()) { - coordinator.recordFailure( + queue.fail( batch, result.reasonCode(), "The Asset authorization relationship could not be confirmed"); throw new AssetUnavailableException( "Asset authorization is waiting for projection"); } - coordinator.complete(batch, result.policyVersion()); + queue.complete(batch, result.policyVersion()); } } diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/authorization/package-info.java b/core/src/main/java/com/orgmemory/core/assetregistry/authorization/package-info.java new file mode 100644 index 000000000..071302b3a --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/authorization/package-info.java @@ -0,0 +1,4 @@ +@org.springframework.modulith.ApplicationModule( + type = org.springframework.modulith.ApplicationModule.Type.CLOSED, + allowedDependencies = {"assetregistry.kernel", "assetregistry::api", "authorization"}) +package com.orgmemory.core.assetregistry.authorization; diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/Asset.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/Asset.java similarity index 98% rename from core/src/main/java/com/orgmemory/core/assetregistry/Asset.java rename to core/src/main/java/com/orgmemory/core/assetregistry/kernel/Asset.java index 28f8be72b..c9055deb5 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/Asset.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/Asset.java @@ -1,4 +1,4 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.kernel; import com.orgmemory.core.assetregistry.api.AssetConflictException; import com.orgmemory.core.assetregistry.api.AssetPortfolioState; diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationBatch.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationBatch.java new file mode 100644 index 000000000..3941604ce --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationBatch.java @@ -0,0 +1,52 @@ +package com.orgmemory.core.assetregistry.kernel; + +import com.orgmemory.core.authorization.RelationshipTuple; +import java.util.List; +import java.util.Objects; +import java.util.UUID; + +/** Opaque persisted-lease capability presented to the external projection flow. */ +public final class AssetAuthorizationBatch { + + private final UUID organizationId; + private final UUID assetId; + private final UUID claimToken; + private final List outboxIds; + private final List tuples; + + AssetAuthorizationBatch( + UUID organizationId, + UUID assetId, + UUID claimToken, + List outboxIds, + List tuples) { + this.organizationId = Objects.requireNonNull(organizationId, "organizationId"); + this.assetId = Objects.requireNonNull(assetId, "assetId"); + this.claimToken = Objects.requireNonNull(claimToken, "claimToken"); + this.outboxIds = List.copyOf(outboxIds); + this.tuples = List.copyOf(tuples); + if (this.outboxIds.size() != this.tuples.size()) { + throw new IllegalArgumentException("Outbox ids and tuples must have equal size"); + } + } + + public UUID organizationId() { + return organizationId; + } + + public UUID assetId() { + return assetId; + } + + public List tuples() { + return tuples; + } + + UUID claimToken() { + return claimToken; + } + + List outboxIds() { + return outboxIds; + } +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationCoordinator.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationCoordinator.java similarity index 76% rename from core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationCoordinator.java rename to core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationCoordinator.java index 5ce368b5b..2a7f87ec9 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationCoordinator.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationCoordinator.java @@ -1,4 +1,4 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.kernel; import com.orgmemory.core.assetregistry.api.AssetConflictException; import com.orgmemory.core.assetregistry.api.AssetNotFoundException; @@ -7,6 +7,7 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.Set; import java.util.stream.Collectors; import java.util.UUID; @@ -16,7 +17,7 @@ import org.springframework.transaction.annotation.Transactional; @Service -class AssetAuthorizationCoordinator { +class AssetAuthorizationCoordinator implements AssetAuthorizationProjectionQueue { private static final Duration CLAIM_LEASE = Duration.ofMinutes(2); @@ -34,20 +35,25 @@ class AssetAuthorizationCoordinator { } @Transactional(propagation = Propagation.REQUIRES_NEW) - AssetAuthorizationBatch startAttempt(UUID organizationId, UUID assetId) { + @Override + public Optional claimForAsset( + UUID organizationId, UUID assetId) { Asset asset = assets.findForUpdate(assetId, organizationId) .orElseThrow(AssetNotFoundException::new); Instant now = Instant.now(); List claimable = outbox.findClaimableForAsset(asset.getId(), now); - AssetAuthorizationBatch batch = - claimBatch(organizationId, assetId, claimable, now); + if (claimable.isEmpty()) { + return Optional.empty(); + } + AssetAuthorizationBatch batch = claimBatch(organizationId, assetId, claimable, now); outbox.saveAllAndFlush(claimable); - return batch; + return Optional.of(batch); } @Transactional(propagation = Propagation.REQUIRES_NEW) - List claimPendingBatches(int limit) { + @Override + public List claimPending(int limit) { if (limit <= 0) { throw new IllegalArgumentException("Convergence limit must be positive"); } @@ -71,7 +77,8 @@ List claimPendingBatches(int limit) { } @Transactional(propagation = Propagation.REQUIRES_NEW) - void complete(AssetAuthorizationBatch batch, String modelId) { + @Override + public void complete(AssetAuthorizationBatch batch, String modelId) { Asset asset = assets.findForUpdate(batch.assetId(), batch.organizationId()) .orElseThrow(AssetNotFoundException::new); Instant appliedAt = Instant.now(); @@ -79,17 +86,18 @@ void complete(AssetAuthorizationBatch batch, String modelId) { if (records.size() != batch.outboxIds().size()) { throw new IllegalStateException("Asset authorization outbox batch is incomplete"); } + requireBatchOwnership(batch, records); Set roleIds = records.stream() .map(AssetAuthorizationOutbox::getRoleAssignmentId) .filter(java.util.Objects::nonNull) .collect(Collectors.toSet()); records.forEach( record -> record.markApplied(batch.claimToken(), modelId, appliedAt)); - outbox.saveAll(records); + outbox.saveAllAndFlush(records); if (!roleIds.isEmpty()) { List assignments = roles.findAllById(roleIds); assignments.forEach(role -> role.markProjected(appliedAt)); - roles.saveAll(assignments); + roles.saveAllAndFlush(assignments); } if (outbox.countUnresolved(batch.assetId()) == 0) { asset.markAuthorizationReady(); @@ -98,16 +106,29 @@ record -> record.markApplied(batch.claimToken(), modelId, appliedAt)); } @Transactional(propagation = Propagation.REQUIRES_NEW) - void recordFailure(AssetAuthorizationBatch batch, String code, String message) { + @Override + public void fail(AssetAuthorizationBatch batch, String code, String message) { List records = outbox.findByIdIn(batch.outboxIds()); if (records.size() != batch.outboxIds().size()) { throw new AssetConflictException( "Asset authorization outbox batch is incomplete"); } + requireBatchOwnership(batch, records); Instant failedAt = Instant.now(); records.forEach(record -> record.recordFailure( batch.claimToken(), code, message, failedAt)); - outbox.saveAll(records); + outbox.saveAllAndFlush(records); + } + + private static void requireBatchOwnership( + AssetAuthorizationBatch batch, List records) { + boolean mismatched = records.stream().anyMatch(record -> + !record.getOrganizationId().equals(batch.organizationId()) + || !record.getAssetId().equals(batch.assetId())); + if (mismatched) { + throw new AssetConflictException( + "Asset authorization outbox batch crosses an Asset boundary"); + } } private static AssetAuthorizationBatch claimBatch( diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationOutbox.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationOutbox.java similarity index 99% rename from core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationOutbox.java rename to core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationOutbox.java index 39fff6231..e6a4f9076 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationOutbox.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationOutbox.java @@ -1,4 +1,4 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.kernel; import com.orgmemory.core.assetregistry.api.AssetConflictException; import com.orgmemory.core.authorization.RelationshipTuple; diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationOutboxRepository.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationOutboxRepository.java similarity index 88% rename from core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationOutboxRepository.java rename to core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationOutboxRepository.java index 1079a29a2..a77a175b6 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationOutboxRepository.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationOutboxRepository.java @@ -1,4 +1,4 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.kernel; import jakarta.persistence.LockModeType; import java.util.Collection; @@ -19,10 +19,10 @@ interface AssetAuthorizationOutboxRepository from AssetAuthorizationOutbox outbox where outbox.assetId = :assetId and ( - (outbox.status = com.orgmemory.core.assetregistry.AssetAuthorizationStatus.PENDING + (outbox.status = com.orgmemory.core.assetregistry.kernel.AssetAuthorizationStatus.PENDING and outbox.nextAttemptAt <= :now) or - (outbox.status = com.orgmemory.core.assetregistry.AssetAuthorizationStatus.IN_FLIGHT + (outbox.status = com.orgmemory.core.assetregistry.kernel.AssetAuthorizationStatus.IN_FLIGHT and outbox.leaseUntil <= :now) ) order by outbox.createdAt @@ -36,11 +36,11 @@ List findClaimableForAsset( select outbox from AssetAuthorizationOutbox outbox where ( - outbox.status = com.orgmemory.core.assetregistry.AssetAuthorizationStatus.PENDING + outbox.status = com.orgmemory.core.assetregistry.kernel.AssetAuthorizationStatus.PENDING and outbox.nextAttemptAt <= :now ) or ( - outbox.status = com.orgmemory.core.assetregistry.AssetAuthorizationStatus.IN_FLIGHT + outbox.status = com.orgmemory.core.assetregistry.kernel.AssetAuthorizationStatus.IN_FLIGHT and outbox.leaseUntil <= :now ) order by outbox.createdAt @@ -53,7 +53,7 @@ List findClaimable( select count(outbox) from AssetAuthorizationOutbox outbox where outbox.assetId = :assetId - and outbox.status <> com.orgmemory.core.assetregistry.AssetAuthorizationStatus.APPLIED + and outbox.status <> com.orgmemory.core.assetregistry.kernel.AssetAuthorizationStatus.APPLIED """) long countUnresolved(@Param("assetId") UUID assetId); diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationProjectionQueue.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationProjectionQueue.java new file mode 100644 index 000000000..857d38ea1 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationProjectionQueue.java @@ -0,0 +1,17 @@ +package com.orgmemory.core.assetregistry.kernel; + +import java.util.List; +import java.util.Optional; +import java.util.UUID; + +/** Transaction-owning canonical queue around external Asset authorization projection. */ +public interface AssetAuthorizationProjectionQueue { + + Optional claimForAsset(UUID organizationId, UUID assetId); + + List claimPending(int limit); + + void complete(AssetAuthorizationBatch batch, String authorizationModelId); + + void fail(AssetAuthorizationBatch batch, String code, String message); +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationStatus.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationStatus.java similarity index 65% rename from core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationStatus.java rename to core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationStatus.java index a7abb90c9..edc423ea1 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetAuthorizationStatus.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationStatus.java @@ -1,4 +1,4 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.kernel; enum AssetAuthorizationStatus { PENDING, diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetKernelService.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetKernelService.java new file mode 100644 index 000000000..80d1efb64 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetKernelService.java @@ -0,0 +1,242 @@ +package com.orgmemory.core.assetregistry.kernel; + +import com.orgmemory.core.assetregistry.api.AssetAuthorizationTarget; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationTargetQuery; +import com.orgmemory.core.assetregistry.api.AssetConflictException; +import com.orgmemory.core.assetregistry.api.AssetIdentity; +import com.orgmemory.core.assetregistry.api.AssetIdentityQuery; +import com.orgmemory.core.assetregistry.api.AssetNotFoundException; +import com.orgmemory.core.assetregistry.api.AssetPortfolioCommand; +import com.orgmemory.core.assetregistry.api.AssetRegistrationCommand; +import com.orgmemory.core.assetregistry.api.AssetRole; +import com.orgmemory.core.assetregistry.api.AssetRoleCommand; +import com.orgmemory.core.assetregistry.api.AssetRoleQuery; +import com.orgmemory.core.authorization.RelationshipTuple; +import java.time.Instant; +import java.util.List; +import java.util.Optional; +import java.util.UUID; +import org.springframework.dao.DataIntegrityViolationException; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +@Service +class AssetKernelService implements + AssetRegistrationCommand, + AssetRoleCommand, + AssetPortfolioCommand, + AssetIdentityQuery, + AssetAuthorizationTargetQuery, + AssetRoleQuery { + + private final AssetRepository assets; + private final AssetRoleAssignmentRepository roles; + private final AssetAuthorizationOutboxRepository outbox; + + AssetKernelService( + AssetRepository assets, + AssetRoleAssignmentRepository roles, + AssetAuthorizationOutboxRepository outbox) { + this.assets = assets; + this.roles = roles; + this.outbox = outbox; + } + + @Override + @Transactional(propagation = Propagation.MANDATORY) + public UUID register(NewAsset command) { + Asset asset = new Asset( + command.organizationId(), + command.type(), + command.namespace(), + command.slug(), + command.knowledgeSpaceId()); + try { + assets.saveAndFlush(asset); + } catch (DataIntegrityViolationException duplicate) { + throw new AssetConflictException( + "An Asset already uses this namespace and slug, or the Space is unavailable", + duplicate); + } + Instant now = Instant.now(); + AssetRoleAssignment owner = roles.saveAndFlush(new AssetRoleAssignment( + command.organizationId(), + asset.getId(), + command.owner(), + AssetRole.OWNER, + command.assignedByUserId(), + now)); + String object = "asset:" + asset.getId(); + outbox.saveAllAndFlush(List.of( + new AssetAuthorizationOutbox( + command.organizationId(), + asset.getId(), + null, + RelationshipTuple.of( + "organization:" + command.organizationId(), + "organization", + object)), + new AssetAuthorizationOutbox( + command.organizationId(), + asset.getId(), + null, + RelationshipTuple.of( + "knowledge_space:" + command.knowledgeSpaceId(), + "space", + object)), + new AssetAuthorizationOutbox( + command.organizationId(), + asset.getId(), + owner.getId(), + RelationshipTuple.of( + command.owner().openFgaUser(), + AssetRole.OWNER.relation(), + object)))); + return asset.getId(); + } + + @Override + @Transactional(propagation = Propagation.MANDATORY) + public UUID assign(Assignment command) { + Asset asset = requiredForUpdate(command.organizationId(), command.assetId()); + if (roles.findByAssetIdAndPrincipalTypeAndPrincipalIdAndRoleAndValidUntilIsNull( + asset.getId(), + command.principal().type(), + command.principal().id(), + command.role()) + .isPresent()) { + throw new AssetConflictException("This active Asset role assignment already exists"); + } + Instant now = Instant.now(); + AssetRoleAssignment assignment = roles.saveAndFlush(new AssetRoleAssignment( + command.organizationId(), + asset.getId(), + command.principal(), + command.role(), + command.assignedByUserId(), + now)); + outbox.saveAndFlush(new AssetAuthorizationOutbox( + command.organizationId(), + asset.getId(), + assignment.getId(), + RelationshipTuple.of( + command.principal().openFgaUser(), + command.role().relation(), + "asset:" + asset.getId()))); + return assignment.getId(); + } + + @Override + @Transactional(propagation = Propagation.MANDATORY) + public AssetIdentity activateAfterRelease(UUID organizationId, UUID assetId) { + Asset asset = requiredForUpdate(organizationId, assetId); + asset.activate(); + return identity(assets.save(asset)); + } + + @Override + @Transactional(propagation = Propagation.MANDATORY) + public AssetIdentity startSunsettingAfterReleaseChange( + UUID organizationId, UUID assetId) { + Asset asset = requiredForUpdate(organizationId, assetId); + asset.startSunsetting(); + return identity(assets.save(asset)); + } + + @Override + @Transactional(propagation = Propagation.MANDATORY) + public AssetIdentity retireAfterFinalWithdrawal(UUID organizationId, UUID assetId) { + Asset asset = requiredForUpdate(organizationId, assetId); + asset.retire(); + return identity(assets.save(asset)); + } + + @Override + @Transactional(readOnly = true) + public Optional findById(UUID organizationId, UUID assetId) { + return assets.findByIdAndOrganizationId(assetId, organizationId) + .map(AssetKernelService::identity); + } + + @Override + @Transactional(readOnly = true) + public Optional findByCoordinate( + UUID organizationId, String namespace, String slug) { + return assets.findByOrganizationIdAndNamespaceAndSlug(organizationId, namespace, slug) + .map(AssetKernelService::identity); + } + + @Override + @Transactional(readOnly = true) + public Optional find(UUID organizationId, UUID assetId) { + return findById(organizationId, assetId) + .map(asset -> new AssetAuthorizationTarget( + asset.organizationId(), + asset.id(), + asset.knowledgeSpaceId(), + asset.type(), + asset.authorizationReady())); + } + + @Override + @Transactional(readOnly = true) + public RoleHistory history(UUID organizationId, UUID assetId, Instant viewedAt) { + findById(organizationId, assetId).orElseThrow(AssetNotFoundException::new); + List assignments = roles.findByAssetIdOrderByValidFromAsc(assetId); + boolean ownerPresent = hasActiveRole(assignments, AssetRole.OWNER, viewedAt); + boolean backupOwnerPresent = hasActiveRole(assignments, AssetRole.BACKUP_OWNER, viewedAt); + return new RoleHistory( + new OwnershipHealth( + ownerPresent, + backupOwnerPresent, + !ownerPresent && !backupOwnerPresent, + !ownerPresent || !backupOwnerPresent), + assignments.stream().map(AssetKernelService::roleView).toList()); + } + + @Override + @Transactional(readOnly = true) + public List activeAssetIdsForUserRole( + UUID organizationId, String userId, AssetRole role, Instant viewedAt) { + return roles.findActiveAssetIdsForUserRole( + organizationId, userId, role, viewedAt); + } + + private Asset requiredForUpdate(UUID organizationId, UUID assetId) { + return assets.findForUpdate(assetId, organizationId) + .orElseThrow(AssetNotFoundException::new); + } + + private static AssetIdentity identity(Asset asset) { + return new AssetIdentity( + asset.getOrganizationId(), + asset.getId(), + asset.getType(), + asset.getNamespace(), + asset.getSlug(), + asset.getKnowledgeSpaceId(), + asset.getPortfolioState(), + asset.isAuthorizationReady()); + } + + private static boolean hasActiveRole( + List assignments, AssetRole role, Instant viewedAt) { + return assignments.stream().anyMatch(assignment -> + assignment.getRole() == role + && (assignment.getValidUntil() == null + || assignment.getValidUntil().isAfter(viewedAt))); + } + + private static RoleAssignment roleView(AssetRoleAssignment assignment) { + return new RoleAssignment( + assignment.getId(), + assignment.getPrincipalType(), + assignment.getPrincipalId(), + assignment.getRole(), + assignment.getValidFrom(), + assignment.getValidUntil(), + assignment.getAssignedByUserId(), + assignment.getProjectedAt()); + } +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetRepository.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetRepository.java new file mode 100644 index 000000000..cd5143132 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetRepository.java @@ -0,0 +1,28 @@ +package com.orgmemory.core.assetregistry.kernel; + +import jakarta.persistence.LockModeType; +import java.util.Optional; +import java.util.UUID; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Lock; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; + +interface AssetRepository extends JpaRepository { + + Optional findByIdAndOrganizationId(UUID id, UUID organizationId); + + @Lock(LockModeType.PESSIMISTIC_WRITE) + @Query(""" + select asset + from Asset asset + where asset.id = :id + and asset.organizationId = :organizationId + """) + Optional findForUpdate( + @Param("id") UUID id, + @Param("organizationId") UUID organizationId); + + Optional findByOrganizationIdAndNamespaceAndSlug( + UUID organizationId, String namespace, String slug); +} diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetRoleAssignment.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetRoleAssignment.java similarity index 98% rename from core/src/main/java/com/orgmemory/core/assetregistry/AssetRoleAssignment.java rename to core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetRoleAssignment.java index 96a64f364..b42767a0c 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetRoleAssignment.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetRoleAssignment.java @@ -1,4 +1,4 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.kernel; import com.orgmemory.core.assetregistry.api.AssetRole; import com.orgmemory.core.authorization.PrincipalRef; diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/AssetRoleAssignmentRepository.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetRoleAssignmentRepository.java similarity index 96% rename from core/src/main/java/com/orgmemory/core/assetregistry/AssetRoleAssignmentRepository.java rename to core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetRoleAssignmentRepository.java index 417b8e13e..f464f7a4a 100644 --- a/core/src/main/java/com/orgmemory/core/assetregistry/AssetRoleAssignmentRepository.java +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/AssetRoleAssignmentRepository.java @@ -1,4 +1,4 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.kernel; import com.orgmemory.core.assetregistry.api.AssetRole; import java.time.Instant; diff --git a/core/src/main/java/com/orgmemory/core/assetregistry/kernel/package-info.java b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/package-info.java new file mode 100644 index 000000000..77b85c8e0 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assetregistry/kernel/package-info.java @@ -0,0 +1,10 @@ +/** + * Canonical Asset identity, accountable-role, authorization-intent, and readiness ledger. + * + *

Parent-facing commands and immutable values belong to {@code assetregistry::api}. The + * Kernel exposes only its opaque projection queue to the sibling authorization module. + */ +@org.springframework.modulith.ApplicationModule( + type = org.springframework.modulith.ApplicationModule.Type.CLOSED, + allowedDependencies = {"assetregistry::api", "authorization", "shared"}) +package com.orgmemory.core.assetregistry.kernel; diff --git a/core/src/test/java/com/orgmemory/core/ModulithVerificationTests.java b/core/src/test/java/com/orgmemory/core/ModulithVerificationTests.java index d43864fd5..f1ac31c94 100644 --- a/core/src/test/java/com/orgmemory/core/ModulithVerificationTests.java +++ b/core/src/test/java/com/orgmemory/core/ModulithVerificationTests.java @@ -5,10 +5,19 @@ import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationProjectionCommand; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationTarget; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationTargetQuery; import com.orgmemory.core.assetregistry.api.AssetConflictException; +import com.orgmemory.core.assetregistry.api.AssetIdentity; +import com.orgmemory.core.assetregistry.api.AssetIdentityQuery; import com.orgmemory.core.assetregistry.api.AssetNotFoundException; +import com.orgmemory.core.assetregistry.api.AssetPortfolioCommand; import com.orgmemory.core.assetregistry.api.AssetPortfolioState; +import com.orgmemory.core.assetregistry.api.AssetRegistrationCommand; import com.orgmemory.core.assetregistry.api.AssetRole; +import com.orgmemory.core.assetregistry.api.AssetRoleCommand; +import com.orgmemory.core.assetregistry.api.AssetRoleQuery; import com.orgmemory.core.assetregistry.api.AssetType; import com.orgmemory.core.assetregistry.api.AssetUnavailableException; import com.orgmemory.core.knowledge.catalog.KnowledgeCatalogEntry; @@ -1164,7 +1173,84 @@ void assetRegistryApiIsAnExactExplicitNamedInterface() { AssetPortfolioState.class.getName(), AssetRole.class.getName(), AssetType.class.getName(), - AssetUnavailableException.class.getName()), + AssetUnavailableException.class.getName(), + AssetAuthorizationProjectionCommand.class.getName(), + AssetAuthorizationTarget.class.getName(), + AssetAuthorizationTargetQuery.class.getName(), + AssetIdentity.class.getName(), + AssetIdentityQuery.class.getName(), + AssetPortfolioCommand.class.getName(), + AssetRegistrationCommand.class.getName(), + AssetRegistrationCommand.NewAsset.class.getName(), + AssetRoleCommand.class.getName(), + AssetRoleCommand.Assignment.class.getName(), + AssetRoleQuery.class.getName(), + AssetRoleQuery.OwnershipHealth.class.getName(), + AssetRoleQuery.RoleAssignment.class.getName(), + AssetRoleQuery.RoleHistory.class.getName()), exposedTypes); } + + @Test + void assetRegistryKernelIsAClosedNestedModule() { + var kernel = modules.getModuleByName("assetregistry.kernel").orElseThrow(); + var allowedDependencies = kernel.getAllowedDependencies(modules).stream() + .map(Object::toString) + .map(dependency -> dependency.replace(" :: ", "::")) + .collect(TreeSet::new, Set::add, Set::addAll); + + assertFalse(kernel.isOpen()); + assertEquals( + Set.of("assetregistry::api", "authorization", "shared"), + allowedDependencies); + } + + @Test + void assetRegistryKernelExposesOnlyProjectionQueueContracts() { + var publicRootTypes = new ClassFileImporter() + .withImportOption(ImportOption.Predefined.DO_NOT_INCLUDE_TESTS) + .importPackages("com.orgmemory.core.assetregistry.kernel") + .stream() + .filter(type -> type.getPackageName().equals( + "com.orgmemory.core.assetregistry.kernel")) + .filter(type -> type.getModifiers().contains(JavaModifier.PUBLIC)) + .map(type -> type.getName()) + .collect(TreeSet::new, Set::add, Set::addAll); + + assertEquals( + Set.of( + "com.orgmemory.core.assetregistry.kernel.AssetAuthorizationBatch", + "com.orgmemory.core.assetregistry.kernel.AssetAuthorizationProjectionQueue"), + publicRootTypes); + } + + @Test + void assetRegistryKernelDoesNotDependOnParentPersistenceOrProjection() { + noClasses() + .that() + .resideInAPackage("com.orgmemory.core.assetregistry.kernel..") + .should() + .dependOnClassesThat() + // Exact parent-package match keeps assetregistry.api available to the kernel. + .resideInAnyPackage( + "com.orgmemory.core.assetregistry", + "com.orgmemory.core.assetregistry.authorization..") + .check(new ClassFileImporter() + .withImportOption(ImportOption.Predefined.DO_NOT_INCLUDE_TESTS) + .importPackages("com.orgmemory.core.assetregistry.kernel")); + } + + @Test + void assetRegistryAuthorizationIsAClosedProjectionModule() { + var authorization = modules.getModuleByName("assetregistry.authorization").orElseThrow(); + var allowedDependencies = authorization.getAllowedDependencies(modules).stream() + .map(Object::toString) + .map(dependency -> dependency.replace(" :: ", "::")) + .collect(TreeSet::new, Set::add, Set::addAll); + + assertFalse(authorization.isOpen()); + assertEquals( + Set.of("assetregistry.kernel", "assetregistry::api", "authorization"), + allowedDependencies); + } } diff --git a/core/src/test/java/com/orgmemory/core/assetregistry/AssetRegistryServiceTests.java b/core/src/test/java/com/orgmemory/core/assetregistry/AssetRegistryServiceTests.java index 1cb6bba9d..25cf4a378 100644 --- a/core/src/test/java/com/orgmemory/core/assetregistry/AssetRegistryServiceTests.java +++ b/core/src/test/java/com/orgmemory/core/assetregistry/AssetRegistryServiceTests.java @@ -8,6 +8,8 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationProjectionCommand; +import com.orgmemory.core.assetregistry.api.AssetAuthorizationTarget; import com.orgmemory.core.assetregistry.api.AssetRole; import com.orgmemory.core.assetregistry.api.AssetType; import com.orgmemory.core.authorization.AuthorizationDecision; @@ -28,6 +30,30 @@ class AssetRegistryServiceTests { + @Test + void authorizationTargetRequiresCanonicalIdentityFields() { + UUID organizationId = UUID.randomUUID(); + UUID assetId = UUID.randomUUID(); + UUID knowledgeSpaceId = UUID.randomUUID(); + + assertThrows( + NullPointerException.class, + () -> new AssetAuthorizationTarget( + null, assetId, knowledgeSpaceId, AssetType.SKILL, true)); + assertThrows( + NullPointerException.class, + () -> new AssetAuthorizationTarget( + organizationId, null, knowledgeSpaceId, AssetType.SKILL, true)); + assertThrows( + NullPointerException.class, + () -> new AssetAuthorizationTarget( + organizationId, assetId, null, AssetType.SKILL, true)); + assertThrows( + NullPointerException.class, + () -> new AssetAuthorizationTarget( + organizationId, assetId, knowledgeSpaceId, null, true)); + } + @Test void ownedWorkspaceResolvesVisibilityAndCanonicalOwnerAssignments() { UUID organizationId = UUID.randomUUID(); @@ -59,7 +85,7 @@ void ownedWorkspaceResolvesVisibilityAndCanonicalOwnerAssignments() { .thenReturn(expected); AssetRegistryService service = new AssetRegistryService( coordinator, - mock(AssetAuthorizationProjectionService.class), + mock(AssetAuthorizationProjectionCommand.class), mock(RelationshipAuthorizationPort.class), authorizationSets); @@ -111,7 +137,7 @@ void governanceActionsUseLivePermissionsAfterRequiringAssetView() { }); AssetRegistryService service = new AssetRegistryService( coordinator, - mock(AssetAuthorizationProjectionService.class), + mock(AssetAuthorizationProjectionCommand.class), authorization, mock(RelationshipAuthorizationSetPort.class)); @@ -156,7 +182,7 @@ void rejectsMissingRolePrincipalFieldsAsBusinessValidation() { AuthorizationDecision.allow("model-v1")); AssetRegistryService service = new AssetRegistryService( coordinator, - mock(AssetAuthorizationProjectionService.class), + mock(AssetAuthorizationProjectionCommand.class), authorization, mock(RelationshipAuthorizationSetPort.class)); diff --git a/core/src/test/java/com/orgmemory/core/assetregistry/AssetValidationTests.java b/core/src/test/java/com/orgmemory/core/assetregistry/AssetValidationTests.java index 7566f100d..91add90cd 100644 --- a/core/src/test/java/com/orgmemory/core/assetregistry/AssetValidationTests.java +++ b/core/src/test/java/com/orgmemory/core/assetregistry/AssetValidationTests.java @@ -2,33 +2,12 @@ import static org.junit.jupiter.api.Assertions.assertThrows; -import com.orgmemory.core.assetregistry.api.AssetType; import java.time.Instant; import java.util.UUID; import org.junit.jupiter.api.Test; class AssetValidationTests { - @Test - void coordinatesRejectValuesLongerThanTheirDatabaseColumns() { - assertThrows( - IllegalArgumentException.class, - () -> new Asset( - UUID.randomUUID(), - AssetType.PROMPT_TEMPLATE, - "a".repeat(129), - "valid-slug", - UUID.randomUUID())); - assertThrows( - IllegalArgumentException.class, - () -> new Asset( - UUID.randomUUID(), - AssetType.PROMPT_TEMPLATE, - "valid.namespace", - "a".repeat(129), - UUID.randomUUID())); - } - @Test void releaseLabelRejectsValuesLongerThanItsDatabaseColumn() { UUID actorId = UUID.randomUUID(); diff --git a/core/src/test/java/com/orgmemory/core/assetregistry/SkillDistributionServiceTests.java b/core/src/test/java/com/orgmemory/core/assetregistry/SkillDistributionServiceTests.java index ae56925b6..0d55892a1 100644 --- a/core/src/test/java/com/orgmemory/core/assetregistry/SkillDistributionServiceTests.java +++ b/core/src/test/java/com/orgmemory/core/assetregistry/SkillDistributionServiceTests.java @@ -7,7 +7,10 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import com.orgmemory.core.assetregistry.api.AssetIdentity; +import com.orgmemory.core.assetregistry.api.AssetIdentityQuery; import com.orgmemory.core.assetregistry.api.AssetNotFoundException; +import com.orgmemory.core.assetregistry.api.AssetPortfolioState; import com.orgmemory.core.assetregistry.api.AssetType; import com.orgmemory.core.assetregistry.api.AssetUnavailableException; import com.orgmemory.core.organization.CurrentActor; @@ -76,13 +79,11 @@ void closesAndRejectsStoredBytesWhoseMetadataNoLongerMatchesTheRelease() { @Test void resolvesCoordinateAndVersionBeforeApplyingTheSameLiveUseCheck() { Fixture fixture = fixture(); - Asset asset = mock(Asset.class); + AssetIdentity asset = assetIdentity(); AssetRelease release = mock(AssetRelease.class); - when(asset.getId()).thenReturn(ASSET_ID); - when(asset.getType()).thenReturn(AssetType.SKILL); when(release.getId()).thenReturn(RELEASE_ID); - when(fixture.assetRepository - .findByOrganizationIdAndNamespaceAndSlug( + when(fixture.identities + .findByCoordinate( ORGANIZATION_ID, "support", "triage")) @@ -106,10 +107,9 @@ void resolvesCoordinateAndVersionBeforeApplyingTheSameLiveUseCheck() { @Test void keepsTheInvalidVersionCauseBehindTheOpaqueNotFoundError() { Fixture fixture = fixture(); - Asset asset = mock(Asset.class); - when(asset.getType()).thenReturn(AssetType.SKILL); - when(fixture.assetRepository - .findByOrganizationIdAndNamespaceAndSlug( + AssetIdentity asset = assetIdentity(); + when(fixture.identities + .findByCoordinate( ORGANIZATION_ID, "support", "triage")) @@ -123,9 +123,25 @@ void keepsTheInvalidVersionCauseBehindTheOpaqueNotFoundError() { assertTrue(failure.getCause() instanceof IllegalArgumentException); } + @Test + void rejectsACoordinateThatResolvesToANonSkillAsset() { + Fixture fixture = fixture(); + when(fixture.identities + .findByCoordinate( + ORGANIZATION_ID, + "support", + "triage")) + .thenReturn(Optional.of(assetIdentity(AssetType.PROMPT_TEMPLATE))); + + assertThrows( + AssetNotFoundException.class, + () -> fixture.service.manifest( + ACTOR, "support", "triage", "1.2.0")); + } + private static Fixture fixture() { AssetRegistryService assets = mock(AssetRegistryService.class); - AssetRepository assetRepository = mock(AssetRepository.class); + AssetIdentityQuery identities = mock(AssetIdentityQuery.class); AssetReleaseRepository releaseRepository = mock(AssetReleaseRepository.class); AssetPayloadReferenceRepository references = @@ -151,13 +167,13 @@ private static Fixture fixture() { return new Fixture( new SkillDistributionService( assets, - assetRepository, + identities, releaseRepository, references, specs, storage), assets, - assetRepository, + identities, releaseRepository, storage); } @@ -182,6 +198,22 @@ private static AssetConsumptionRelease release() { java.time.Instant.parse("2026-07-27T10:00:00Z")); } + private static AssetIdentity assetIdentity() { + return assetIdentity(AssetType.SKILL); + } + + private static AssetIdentity assetIdentity(AssetType type) { + return new AssetIdentity( + ORGANIZATION_ID, + ASSET_ID, + type, + "support", + "triage", + UUID.randomUUID(), + AssetPortfolioState.DRAFT_ONLY, + true); + } + private static SkillPackageSpec spec() { return new SkillPackageSpec( "triage", @@ -204,7 +236,7 @@ private static SkillPackageSpec spec() { private record Fixture( SkillDistributionService service, AssetRegistryService assets, - AssetRepository assetRepository, + AssetIdentityQuery identities, AssetReleaseRepository releaseRepository, SkillPackageStoragePort storage) { } diff --git a/core/src/test/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationProjectionServiceTests.java b/core/src/test/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationProjectionServiceTests.java new file mode 100644 index 000000000..a8d92fbbf --- /dev/null +++ b/core/src/test/java/com/orgmemory/core/assetregistry/authorization/AssetAuthorizationProjectionServiceTests.java @@ -0,0 +1,159 @@ +package com.orgmemory.core.assetregistry.authorization; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.orgmemory.core.assetregistry.api.AssetUnavailableException; +import com.orgmemory.core.assetregistry.kernel.AssetAuthorizationBatch; +import com.orgmemory.core.assetregistry.kernel.AssetAuthorizationProjectionQueue; +import com.orgmemory.core.authorization.RelationshipTuple; +import com.orgmemory.core.authorization.RelationshipTupleWritePort; +import com.orgmemory.core.authorization.RelationshipTupleWriteRequest; +import com.orgmemory.core.authorization.RelationshipTupleWriteResult; +import java.util.List; +import java.util.Optional; +import java.util.UUID; +import org.junit.jupiter.api.Test; +import org.springframework.aop.framework.ProxyFactory; +import org.springframework.transaction.IllegalTransactionStateException; +import org.springframework.transaction.TransactionDefinition; +import org.springframework.transaction.annotation.AnnotationTransactionAttributeSource; +import org.springframework.transaction.interceptor.TransactionInterceptor; +import org.springframework.transaction.support.AbstractPlatformTransactionManager; +import org.springframework.transaction.support.DefaultTransactionStatus; +import org.springframework.transaction.support.TransactionTemplate; + +class AssetAuthorizationProjectionServiceTests { + + @Test + void completesTheClaimOnlyAfterTheRelationshipWriteIsConfirmed() { + AssetAuthorizationProjectionQueue queue = mock(AssetAuthorizationProjectionQueue.class); + RelationshipTupleWritePort writer = mock(RelationshipTupleWritePort.class); + AssetAuthorizationBatch batch = mock(AssetAuthorizationBatch.class); + UUID organizationId = UUID.randomUUID(); + UUID assetId = UUID.randomUUID(); + when(batch.tuples()).thenReturn(List.of(RelationshipTuple.of( + "user:" + UUID.randomUUID(), "owner", "asset:" + assetId))); + when(queue.claimForAsset(organizationId, assetId)).thenReturn(Optional.of(batch)); + when(writer.write(any(RelationshipTupleWriteRequest.class))) + .thenReturn(RelationshipTupleWriteResult.applied("model-1")); + AssetAuthorizationProjectionService service = + new AssetAuthorizationProjectionService(queue, writer); + + service.project(organizationId, assetId); + + verify(queue).complete(batch, "model-1"); + } + + @Test + void recordsAnOpaqueFailureWhenTheExternalWriterThrows() { + AssetAuthorizationProjectionQueue queue = mock(AssetAuthorizationProjectionQueue.class); + RelationshipTupleWritePort writer = mock(RelationshipTupleWritePort.class); + AssetAuthorizationBatch batch = mock(AssetAuthorizationBatch.class); + UUID organizationId = UUID.randomUUID(); + UUID assetId = UUID.randomUUID(); + when(batch.tuples()).thenReturn(List.of(RelationshipTuple.of( + "user:" + UUID.randomUUID(), "owner", "asset:" + assetId))); + when(queue.claimForAsset(organizationId, assetId)).thenReturn(Optional.of(batch)); + when(writer.write(any(RelationshipTupleWriteRequest.class))) + .thenThrow(new IllegalStateException("provider detail")); + AssetAuthorizationProjectionService service = + new AssetAuthorizationProjectionService(queue, writer); + + assertThrows( + AssetUnavailableException.class, + () -> service.project(organizationId, assetId)); + + verify(queue).fail( + batch, + "OPENFGA_WRITE_FAILED", + "The Asset authorization relationship could not be applied"); + } + + @Test + void rejectsAClaimThatIsAlreadyBeingProjected() { + AssetAuthorizationProjectionQueue queue = mock(AssetAuthorizationProjectionQueue.class); + UUID organizationId = UUID.randomUUID(); + UUID assetId = UUID.randomUUID(); + when(queue.claimForAsset(organizationId, assetId)).thenReturn(Optional.empty()); + AssetAuthorizationProjectionService service = new AssetAuthorizationProjectionService( + queue, mock(RelationshipTupleWritePort.class)); + + assertThrows( + AssetUnavailableException.class, + () -> service.project(organizationId, assetId)); + } + + @Test + void externalProjectionRejectsAnAmbientDatabaseTransactionAtRuntime() { + TestTransactionManager transactions = new TestTransactionManager(); + AssetAuthorizationProjectionQueue queue = mock(AssetAuthorizationProjectionQueue.class); + AssetAuthorizationProjectionService projection = transactionalProxy( + new AssetAuthorizationProjectionService( + queue, mock(RelationshipTupleWritePort.class)), + AssetAuthorizationProjectionService.class, + transactions); + AssetAuthorizationConvergenceService convergence = transactionalProxy( + new AssetAuthorizationConvergenceService(queue, projection), + AssetAuthorizationConvergenceService.class, + transactions); + AssetAuthorizationBatch batch = mock(AssetAuthorizationBatch.class); + TransactionTemplate outer = new TransactionTemplate(transactions); + + assertThrows( + IllegalTransactionStateException.class, + () -> outer.executeWithoutResult(status -> projection.project(batch))); + assertThrows( + IllegalTransactionStateException.class, + () -> outer.executeWithoutResult(status -> convergence.reconcile(10))); + } + + private static T transactionalProxy( + T target, Class targetType, TestTransactionManager transactions) { + TransactionInterceptor interceptor = new TransactionInterceptor(); + interceptor.setTransactionManager(transactions); + interceptor.setTransactionAttributeSource(new AnnotationTransactionAttributeSource(false)); + ProxyFactory factory = new ProxyFactory(target); + factory.setProxyTargetClass(true); + factory.addAdvice(interceptor); + return targetType.cast(factory.getProxy()); + } + + private static final class TestTransactionManager extends AbstractPlatformTransactionManager { + + private final ThreadLocal active = ThreadLocal.withInitial(() -> false); + + @Override + protected Object doGetTransaction() { + return new Object(); + } + + @Override + protected boolean isExistingTransaction(Object transaction) { + return active.get(); + } + + @Override + protected void doBegin(Object transaction, TransactionDefinition definition) { + active.set(true); + } + + @Override + protected void doCommit(DefaultTransactionStatus status) { + // No resource commit is needed for transaction-boundary verification. + } + + @Override + protected void doRollback(DefaultTransactionStatus status) { + // No resource rollback is needed for transaction-boundary verification. + } + + @Override + protected void doCleanupAfterCompletion(Object transaction) { + active.remove(); + } + } +} diff --git a/core/src/test/java/com/orgmemory/core/assetregistry/AssetAuthorizationOutboxTests.java b/core/src/test/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationOutboxTests.java similarity index 98% rename from core/src/test/java/com/orgmemory/core/assetregistry/AssetAuthorizationOutboxTests.java rename to core/src/test/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationOutboxTests.java index e565f973f..74d5399af 100644 --- a/core/src/test/java/com/orgmemory/core/assetregistry/AssetAuthorizationOutboxTests.java +++ b/core/src/test/java/com/orgmemory/core/assetregistry/kernel/AssetAuthorizationOutboxTests.java @@ -1,4 +1,4 @@ -package com.orgmemory.core.assetregistry; +package com.orgmemory.core.assetregistry.kernel; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; diff --git a/core/src/test/java/com/orgmemory/core/assetregistry/kernel/AssetKernelServiceTests.java b/core/src/test/java/com/orgmemory/core/assetregistry/kernel/AssetKernelServiceTests.java new file mode 100644 index 000000000..3e208cb96 --- /dev/null +++ b/core/src/test/java/com/orgmemory/core/assetregistry/kernel/AssetKernelServiceTests.java @@ -0,0 +1,178 @@ +package com.orgmemory.core.assetregistry.kernel; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.orgmemory.core.assetregistry.api.AssetConflictException; +import com.orgmemory.core.assetregistry.api.AssetPortfolioState; +import com.orgmemory.core.assetregistry.api.AssetRegistrationCommand; +import com.orgmemory.core.assetregistry.api.AssetRole; +import com.orgmemory.core.assetregistry.api.AssetRoleCommand; +import com.orgmemory.core.assetregistry.api.AssetType; +import com.orgmemory.core.authorization.PrincipalRef; +import java.lang.reflect.Method; +import java.util.List; +import java.util.Optional; +import java.util.UUID; +import org.junit.jupiter.api.Test; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +class AssetKernelServiceTests { + + private static final UUID ORGANIZATION_ID = UUID.randomUUID(); + private static final UUID SPACE_ID = UUID.randomUUID(); + private static final UUID USER_ID = UUID.randomUUID(); + private static final PrincipalRef OWNER = PrincipalRef.user(USER_ID); + + @Test + void registrationWritesIdentityOwnerAndThreeAuthorizationIntentsTogether() { + Fixture fixture = fixture(); + when(fixture.assets.saveAndFlush(any(Asset.class))) + .thenAnswer(invocation -> invocation.getArgument(0)); + when(fixture.roles.saveAndFlush(any(AssetRoleAssignment.class))) + .thenAnswer(invocation -> invocation.getArgument(0)); + + UUID assetId = fixture.service.register(new AssetRegistrationCommand.NewAsset( + ORGANIZATION_ID, + AssetType.PROMPT_TEMPLATE, + "Support", + "triage", + SPACE_ID, + OWNER, + USER_ID)); + + verify(fixture.assets).saveAndFlush(any(Asset.class)); + verify(fixture.roles).saveAndFlush(any(AssetRoleAssignment.class)); + verify(fixture.outbox).saveAllAndFlush(org.mockito.ArgumentMatchers.argThat(records -> { + List intents = java.util.stream.StreamSupport + .stream(records.spliterator(), false) + .toList(); + return intents.size() == 3 + && intents.stream().allMatch(intent -> intent.getAssetId().equals(assetId)) + && intents.stream().map(AssetAuthorizationOutbox::tuple).anyMatch(tuple -> + tuple.user().equals("organization:" + ORGANIZATION_ID) + && tuple.relation().equals("organization")) + && intents.stream().map(AssetAuthorizationOutbox::tuple).anyMatch(tuple -> + tuple.user().equals("knowledge_space:" + SPACE_ID) + && tuple.relation().equals("space")) + && intents.stream().map(AssetAuthorizationOutbox::tuple).anyMatch(tuple -> + tuple.user().equals(OWNER.openFgaUser()) + && tuple.relation().equals("owner")); + })); + } + + @Test + void duplicateActiveRoleDoesNotAppendAnotherAuthorizationIntent() { + Fixture fixture = fixture(); + Asset asset = asset(); + when(fixture.assets.findForUpdate(asset.getId(), ORGANIZATION_ID)) + .thenReturn(Optional.of(asset)); + when(fixture.roles + .findByAssetIdAndPrincipalTypeAndPrincipalIdAndRoleAndValidUntilIsNull( + asset.getId(), "user", USER_ID.toString(), AssetRole.OWNER)) + .thenReturn(Optional.of(mock(AssetRoleAssignment.class))); + + assertThrows( + AssetConflictException.class, + () -> fixture.service.assign(new AssetRoleCommand.Assignment( + ORGANIZATION_ID, + asset.getId(), + OWNER, + AssetRole.OWNER, + USER_ID))); + + verify(fixture.outbox, never()).saveAndFlush(any()); + } + + @Test + void portfolioTransitionsRemainOwnedByTheLockedCanonicalAsset() { + Fixture fixture = fixture(); + Asset asset = asset(); + when(fixture.assets.findForUpdate(asset.getId(), ORGANIZATION_ID)) + .thenReturn(Optional.of(asset)); + when(fixture.assets.save(asset)).thenReturn(asset); + + assertEquals( + AssetPortfolioState.ACTIVE, + fixture.service.activateAfterRelease(ORGANIZATION_ID, asset.getId()) + .portfolioState()); + assertEquals( + AssetPortfolioState.SUNSETTING, + fixture.service.startSunsettingAfterReleaseChange( + ORGANIZATION_ID, asset.getId()) + .portfolioState()); + assertEquals( + AssetPortfolioState.RETIRED, + fixture.service.retireAfterFinalWithdrawal(ORGANIZATION_ID, asset.getId()) + .portfolioState()); + } + + @Test + void kernelCommandsJoinTheParentTransactionAndQueueOperationsOwnShortTransactions() + throws NoSuchMethodException { + assertPropagation( + AssetKernelService.class.getMethod( + "register", AssetRegistrationCommand.NewAsset.class), + Propagation.MANDATORY); + assertPropagation( + AssetKernelService.class.getMethod( + "assign", AssetRoleCommand.Assignment.class), + Propagation.MANDATORY); + assertPropagation( + AssetKernelService.class.getMethod( + "activateAfterRelease", UUID.class, UUID.class), + Propagation.MANDATORY); + assertPropagation( + AssetAuthorizationCoordinator.class.getMethod( + "claimForAsset", UUID.class, UUID.class), + Propagation.REQUIRES_NEW); + assertPropagation( + AssetAuthorizationCoordinator.class.getMethod("claimPending", int.class), + Propagation.REQUIRES_NEW); + assertPropagation( + AssetAuthorizationCoordinator.class.getMethod( + "complete", AssetAuthorizationBatch.class, String.class), + Propagation.REQUIRES_NEW); + assertPropagation( + AssetAuthorizationCoordinator.class.getMethod( + "fail", AssetAuthorizationBatch.class, String.class, String.class), + Propagation.REQUIRES_NEW); + } + + private static void assertPropagation(Method method, Propagation expected) { + Transactional transactional = method.getAnnotation(Transactional.class); + assertTrue(transactional != null, () -> method + " must declare @Transactional"); + assertEquals(expected, transactional.propagation()); + } + + private static Asset asset() { + return new Asset( + ORGANIZATION_ID, + AssetType.PROMPT_TEMPLATE, + "support", + "triage", + SPACE_ID); + } + + private static Fixture fixture() { + AssetRepository assets = mock(AssetRepository.class); + AssetRoleAssignmentRepository roles = mock(AssetRoleAssignmentRepository.class); + AssetAuthorizationOutboxRepository outbox = + mock(AssetAuthorizationOutboxRepository.class); + return new Fixture(new AssetKernelService(assets, roles, outbox), assets, roles, outbox); + } + + private record Fixture( + AssetKernelService service, + AssetRepository assets, + AssetRoleAssignmentRepository roles, + AssetAuthorizationOutboxRepository outbox) { + } +} diff --git a/core/src/test/java/com/orgmemory/core/assetregistry/kernel/AssetValidationTests.java b/core/src/test/java/com/orgmemory/core/assetregistry/kernel/AssetValidationTests.java new file mode 100644 index 000000000..f561f842f --- /dev/null +++ b/core/src/test/java/com/orgmemory/core/assetregistry/kernel/AssetValidationTests.java @@ -0,0 +1,30 @@ +package com.orgmemory.core.assetregistry.kernel; + +import static org.junit.jupiter.api.Assertions.assertThrows; + +import com.orgmemory.core.assetregistry.api.AssetType; +import java.util.UUID; +import org.junit.jupiter.api.Test; + +class AssetValidationTests { + + @Test + void coordinatesRejectValuesLongerThanTheirDatabaseColumns() { + assertThrows( + IllegalArgumentException.class, + () -> new Asset( + UUID.randomUUID(), + AssetType.PROMPT_TEMPLATE, + "a".repeat(129), + "valid-slug", + UUID.randomUUID())); + assertThrows( + IllegalArgumentException.class, + () -> new Asset( + UUID.randomUUID(), + AssetType.PROMPT_TEMPLATE, + "valid.namespace", + "a".repeat(129), + UUID.randomUUID())); + } +} diff --git a/docs/increments/active/2026-07-31-spring-modulith-package-refactor/assetregistry-authorization-challenge-verdict.md b/docs/increments/active/2026-07-31-spring-modulith-package-refactor/assetregistry-authorization-challenge-verdict.md index d5634aab6..5fc0c5cc8 100644 --- a/docs/increments/active/2026-07-31-spring-modulith-package-refactor/assetregistry-authorization-challenge-verdict.md +++ b/docs/increments/active/2026-07-31-spring-modulith-package-refactor/assetregistry-authorization-challenge-verdict.md @@ -106,6 +106,25 @@ canonical writes, queue transactions, and OpenFGA calls are forbidden. Every PR contains code, remains below both its slice cap and the repository 100-file ceiling, and may not weaken closure or expose persistence types. +## Executable Topology Correction — 2026-08-02 + +The full `modules.verify()` gate rejected the planned intermediate parent to +Kernel projection dependency as a real cycle: Kernel consumes the parent-owned +`assetregistry::api`, so a parent import of Kernel or Authorization cannot be +retained. The strongest alternative was to keep PR 2 and PR 3 separate by +leaving the outbox coordinator in the parent temporarily. That would contradict +the selected atomic ownership boundary and make Kernel incomplete on arrival. + +The binding implementation therefore combines delivery steps 2 and 3 while +retaining their total 60-path ceiling. Parent orchestration depends only on the +new `AssetAuthorizationProjectionCommand` in `assetregistry::api`; +package-private Authorization projection implements it. Authorization publicly +exposes only convergence service/report for Worker, while its exact dependency +allowlist remains Kernel, `assetregistry::api`, and Authorization. This +supersedes the parent-to-Authorization arrow and the public projection-service +statement above; all aggregate ownership, transaction rules, and rejected +alternatives remain unchanged. + ## Required Evidence - Characterization tests fail first for module existence/closure, exact public diff --git a/docs/increments/active/2026-07-31-spring-modulith-package-refactor/design.md b/docs/increments/active/2026-07-31-spring-modulith-package-refactor/design.md index fcbac1321..aa60f24d8 100644 --- a/docs/increments/active/2026-07-31-spring-modulith-package-refactor/design.md +++ b/docs/increments/active/2026-07-31-spring-modulith-package-refactor/design.md @@ -407,6 +407,14 @@ and parent orchestration remain in the parent Asset Registry module. In particular, `AssetAvailability` and the broad catalog queries cannot move into kernel because they depend on parent-owned release persistence. +Executable Modulith verification corrected one dependency detail in the +reviewed topology. Parent orchestration cannot import the sibling Authorization +implementation while Kernel consumes the parent-owned API, because that forms a +module cycle. The parent therefore invokes projection through a narrow +`assetregistry::api` command implemented by Authorization; only Worker imports +the public convergence entry point. This preserves the selected ownership and +transaction boundaries without opening either nested module. + Executable Modulith verification corrected the first vocabulary placement: unrelated top-level modules cannot import a nested Kernel directly. The six cross-module Asset enums and business exceptions therefore move first to the diff --git a/docs/increments/active/2026-07-31-spring-modulith-package-refactor/plan.md b/docs/increments/active/2026-07-31-spring-modulith-package-refactor/plan.md index 62092ddc7..e81e8426d 100644 --- a/docs/increments/active/2026-07-31-spring-modulith-package-refactor/plan.md +++ b/docs/increments/active/2026-07-31-spring-modulith-package-refactor/plan.md @@ -1399,19 +1399,19 @@ authorization persistence module. Role, outbox, lease completion, and readiness belong with Asset in one transactional kernel; the authorization module is the external OpenFGA projection edge only. -- [ ] PR 1: move the six cross-module Asset vocabulary/error types to the exact +- [x] PR 1: move the six cross-module Asset vocabulary/error types to the exact parent-owned `assetregistry::api` named interface; keep Kernel absent and stay at or below 70 changed paths. -- [ ] PR 2: introduce and immediately close Kernel; move the canonical +- [x] PR 2: introduce and immediately close Kernel; move the canonical Asset/role/outbox ledger, narrow the identity repository, extract the parent catalog read model, add parent draft locking, and implement parent-facing `assetregistry::api` command/query contracts; stay at or below 60 changed paths. -- [ ] PR 3: move projection and convergence entry points, enforce transaction +- [x] PR 3: move projection and convergence entry points, enforce transaction propagation `NEVER` around OpenFGA calls, close `assetregistry.authorization`, and update Worker wiring; stay at or below 20 changed paths. -- [ ] Run the focused Core/API/Worker gates, docs and release-policy checks, +- [x] Run the focused Core/API/Worker gates, docs and release-policy checks, static analysis, and a terminating clean repository test for every slice. - [ ] Merge each code-bearing PR through CI and CodeRabbit before starting the next branch. Release only after all remaining Asset Registry slices and the @@ -1428,3 +1428,41 @@ passed on exact Node 24.15.0, and the terminating repository-wide `clean test` passed 99 tasks in 2m07s. JetBrains semantic inspection was unavailable, so the documented Gradle/mechanical fallback was used. The complete PR remains at 68 changed paths. + +PR 2's executable Modulith verifier exposed a cycle in the planned intermediate +state: parent projection code importing Kernel while Kernel consumed the parent +`assetregistry::api` interface produced `assetregistry -> kernel -> +assetregistry`. The already-selected PR 3 topology was therefore delivered in +the same code commit rather than weakening verification. Closed Kernel now owns +the canonical Asset, role, outbox/lease, and readiness ledger; closed +Authorization owns only external projection/convergence; the parent invokes +projection through `assetregistry::api`. Registration, role, and portfolio +commands join the parent transaction with `MANDATORY`, queue operations own +short `REQUIRES_NEW` transactions, OpenFGA calls reject ambient transactions +with `NEVER`, and the three mutable Skill/review flows serialize on the Draft +row. Code commit `573c1d1f` contains 41 changed paths. Full Core passed 453 +tests, the Asset Registry API integration suite passed 22, and Worker passed 65, +all with zero failures. Documentation hygiene passed for 509 Markdown files and +8 mirrored domain pairs; all 23 release-policy tests passed on exact Node +24.15.0. JetBrains semantic inspection was unavailable, so the documented +Gradle, `modules.verify()`, import-boundary, and `git diff --check` fallback was +used. The first clean run reached 94 tasks before an orphaned Gradle daemon +caused native-memory exhaustion; after removing that process, the terminating +single-worker `clean test` passed all 99 tasks in 4m18s with a 2 GB daemon heap. + +CodeRabbit's full review of PR #270 produced ten inline findings and one +outside-diff finding. Commit `01e26c23` accepted eight valid improvements: +provider-failure logging without tuple disclosure, required-field validation +for the public authorization target, the unclaimed projection branch, runtime +Spring-proxy enforcement of both `NEVER` boundaries including the package- +private batch overload, both completion transaction boundaries, non-Skill +coordinate rejection, class-literal API pinning, and the intentional exact- +package ArchUnit comment. The Draft lock-timeout suggestion was rejected +because the repository has no established finite value and its PostgreSQL +baseline explicitly leaves `lock_timeout` at zero. The terminal portfolio +finding was rejected as unreachable after the existing withdrawn-release +guard, and the optional role partial index remains outside this no-schema +refactor because the locked canonical Asset already serializes all Kernel role +writes. All review threads were answered and resolved. Full Core passed 456 +tests with zero failures; the terminating sequential repository-wide +`clean test` then passed all 99 tasks in 8m52s on the reviewed code head. diff --git a/docs/specs/domains/asset-registry.md b/docs/specs/domains/asset-registry.md index 45df84715..c4aa7134d 100644 --- a/docs/specs/domains/asset-registry.md +++ b/docs/specs/domains/asset-registry.md @@ -11,7 +11,7 @@ Source: `core/src/main/java/com/orgmemory/core/assetregistry`, `apps/cli/package.json`, `.github/workflows/publish-cli.yml`, and `apps/web/src/features/assets`. -Reconciled: `2026-08-02-spring-modulith-package-refactor (fdf3cca4)`. +Reconciled: `2026-08-02-spring-modulith-package-refactor (573c1d1f)`. ## Current Behavior @@ -26,6 +26,22 @@ exact parent-owned `assetregistry::api` named interface. Nested implementation modules implement parent-facing contracts instead of being imported directly by unrelated top-level modules. +The closed `assetregistry.kernel` module owns the package-private canonical +Asset identity/portfolio entity, accountable roles, authorization outbox +leases, readiness state, and their repositories. Registration and role writes +join parent transactions and persist their authorization intent atomically; +portfolio transitions remain behind immutable parent-facing commands. The +parent retains Draft, revision, review, release, availability, audit, and +catalog read-model persistence. Skill package replacement, submission, and +direct publication serialize on the Draft row rather than exposing a Kernel +lock. + +The closed `assetregistry.authorization` module owns only the external OpenFGA +projection and convergence edge. Parent orchestration invokes its projection +through `assetregistry::api`; Worker uses the public convergence entry point. +Queue claim/completion/failure use short independent transactions, while the +external OpenFGA write explicitly rejects an ambient database transaction. + Consumers always address an exact authorized release. A withdrawn release cannot start new consumption. Forking creates a new Asset draft from an exact release payload and does not copy reviews or approvals. diff --git a/docs/tests/domains/asset-registry.md b/docs/tests/domains/asset-registry.md index 5e2fd6271..78f3288db 100644 --- a/docs/tests/domains/asset-registry.md +++ b/docs/tests/domains/asset-registry.md @@ -10,11 +10,13 @@ Source: `core/src/test/java/com/orgmemory/core/assetregistry`, `scripts/npm-publish-workflow-policy.test.mjs`, and `apps/web/src/features/assets/**/*.test.ts`. -Reconciled: `2026-08-02-spring-modulith-package-refactor (fdf3cca4)`. +Reconciled: `2026-08-02-spring-modulith-package-refactor (573c1d1f)`. | Behavior | Evidence | Status | | --- | --- | --- | -| The parent-owned `assetregistry::api` named interface exposes exactly the six cross-module vocabulary/error types, and the full module graph rejects invalid nested-module references | `ModulithVerificationTests#assetRegistryApiIsAnExactExplicitNamedInterface`, `ModulithVerificationTests#modulesAreWellFormed` | covered | +| The parent-owned `assetregistry::api` named interface exposes exactly the shared vocabulary/errors plus immutable parent-facing contracts, and the full module graph rejects invalid nested-module references | `ModulithVerificationTests#assetRegistryApiIsAnExactExplicitNamedInterface`, `ModulithVerificationTests#modulesAreWellFormed` | covered | +| Closed Kernel owns canonical Asset/role/outbox/readiness persistence with an exact allowlist and exposes only its opaque projection queue; closed Authorization owns the external edge with its exact allowlist | `ModulithVerificationTests#assetRegistryKernelIsAClosedNestedModule`, `#assetRegistryKernelExposesOnlyProjectionQueueContracts`, `#assetRegistryKernelDoesNotDependOnParentPersistenceOrProjection`, `#assetRegistryAuthorizationIsAClosedProjectionModule` | covered | +| Registration writes Asset, OWNER, and three authorization intents atomically; duplicate roles emit no intent; lifecycle transitions stay on the locked canonical Asset; commands join the parent transaction, queue operations own short transactions, and OpenFGA projection rejects ambient transactions | `AssetKernelServiceTests`, `AssetAuthorizationOutboxTests`, `AssetAuthorizationProjectionServiceTests` | covered | | Prompt, Work Instruction, Pack, and Skill schemas reject invalid payloads | `AssetProfileValidationTests` | covered | | Skill ZIP inspection rejects traversal, case collisions, symlinks, invalid frontmatter, invalid UTF-8, and bounded-size violations without extraction | `SkillPackageInspectorTests` | covered | | Unauthorized Skill import is rejected before object storage and pre-identity failures clean up staged objects | `SkillRegistryServiceTests` | covered |