diff --git a/registry/data-product-consumer-policies.json b/registry/data-product-consumer-policies.json index 7fcb36a..db74c4f 100644 --- a/registry/data-product-consumer-policies.json +++ b/registry/data-product-consumer-policies.json @@ -68,6 +68,87 @@ }, "removeMode": "canonical-tombstone-or-snapshot-rebase" }, + { + "id": "map-moving-object-unit-profile-current-v1", + "version": "1.0.0", + "dataProductId": "fleet.units.profile.current.v1", + "productVersion": "1.0.0", + "freshness": "none", + "staleAfterMs": null, + "terminalStatuses": [], + "consumerContract": { + "ontologyRevision": "ontology.map.moving_object.v3", + "deliveryMode": "snapshot+patch", + "semanticTypes": [ + "map.moving_object" + ], + "fieldProjection": [ + "corrected_engine_hours_factor", + "corrected_mileage_factor", + "display_name", + "filter_by_satellite_count_enabled", + "filter_by_satellite_count_value", + "filter_emissions", + "hardware_manufacturer_name", + "hardware_port", + "hardware_type_class", + "hardware_type_name", + "limit_acceleration", + "lost_connection_enabled", + "lost_connection_time_value", + "maximum_permissible_speed", + "maximum_valid_height", + "maximum_valid_speed", + "mileage_by_ignition", + "minimum_movement_speed", + "minimum_movement_time", + "minimum_parking_time", + "minimum_stop_time", + "minimum_trip_distance", + "minimum_valid_height", + "speed_parameter", + "trip_detection_type", + "unit_type_class", + "unit_type_name", + "use_odometer" + ], + "fieldTypes": { + "corrected_engine_hours_factor": "number", + "corrected_mileage_factor": "number", + "display_name": "string", + "filter_by_satellite_count_enabled": "boolean", + "filter_by_satellite_count_value": "number", + "filter_emissions": "number", + "hardware_manufacturer_name": "string", + "hardware_port": "number", + "hardware_type_class": "string", + "hardware_type_name": "string", + "limit_acceleration": "number", + "lost_connection_enabled": "boolean", + "lost_connection_time_value": "number", + "maximum_permissible_speed": "number", + "maximum_valid_height": "number", + "maximum_valid_speed": "number", + "mileage_by_ignition": "boolean", + "minimum_movement_speed": "number", + "minimum_movement_time": "number", + "minimum_parking_time": "number", + "minimum_stop_time": "number", + "minimum_trip_distance": "number", + "minimum_valid_height": "number", + "speed_parameter": "string", + "trip_detection_type": "string", + "unit_type_class": "string", + "unit_type_name": "string", + "use_odometer": "boolean" + }, + "geometryTypes": [], + "subjectIdentity": "semantic-type+source-id", + "snapshotMode": "atomic-replace", + "patchMode": "atomic" + }, + "removeMode": "canonical-tombstone-or-snapshot-rebase" + }, { "id": "map-zone-current-v1", "version": "1.0.0", diff --git a/scripts/validate-registry.mjs b/scripts/validate-registry.mjs index b5f83f1..7097f7c 100644 --- a/scripts/validate-registry.mjs +++ b/scripts/validate-registry.mjs @@ -70,6 +70,19 @@ requireValue(new Set(consumerPolicyProducts).size === consumerPolicyProducts.len for (const policy of consumerPolicies) { const statusContract = policy.statusContract; const freshness = policy.freshness ?? statusContract?.freshness ?? "observed-at"; + const allowedPolicyKeys = new Set([ + "id", + "version", + "dataProductId", + "productVersion", + "freshness", + "staleAfterMs", + "terminalStatuses", + "statusContract", + "consumerContract", + "removeMode", + ]); + requireValue(Object.keys(policy).every((key) => allowedPolicyKeys.has(key)), `unsupported data product consumer policy key: ${policy.id ?? "unknown"}`); requireValue(typeof policy.id === "string" && /^[a-z0-9._-]{1,160}$/.test(policy.id), `invalid data product consumer policy id: ${policy.id ?? "unknown"}`); requireValue(typeof policy.version === "string" && /^\d+\.\d+\.\d+$/.test(policy.version), `invalid data product consumer policy version: ${policy.id ?? "unknown"}`); requireValue(typeof policy.dataProductId === "string" && /^[A-Za-z0-9._:-]{1,160}$/.test(policy.dataProductId), `invalid data product consumer policy product: ${policy.id ?? "unknown"}`); @@ -90,7 +103,14 @@ for (const policy of consumerPolicies) { requireValue(contract?.fieldTypes && typeof contract.fieldTypes === "object" && !Array.isArray(contract.fieldTypes), `invalid consumer field types: ${policy.id ?? "unknown"}`); requireValue(JSON.stringify(Object.keys(contract?.fieldTypes ?? {})) === JSON.stringify(contract?.fieldProjection ?? []), `consumer field type order/scope mismatch: ${policy.id ?? "unknown"}`); requireValue(Object.values(contract?.fieldTypes ?? {}).every((value) => ["string", "number", "boolean"].includes(value)), `invalid consumer field type: ${policy.id ?? "unknown"}`); - requireValue(JSON.stringify(contract?.geometryTypes) === JSON.stringify(["Polygon", "MultiPolygon"]), `invalid consumer geometry types: ${policy.id ?? "unknown"}`); + const supportedGeometryTypes = ["Point", "LineString", "Polygon", "MultiPolygon"]; + requireValue( + Array.isArray(contract?.geometryTypes) + && contract.geometryTypes.length <= supportedGeometryTypes.length + && new Set(contract.geometryTypes).size === contract.geometryTypes.length + && contract.geometryTypes.every((value) => supportedGeometryTypes.includes(value)), + `invalid consumer geometry types: ${policy.id ?? "unknown"}`, + ); requireValue(contract?.subjectIdentity === "semantic-type+source-id", `invalid consumer subject identity: ${policy.id ?? "unknown"}`); requireValue(contract?.snapshotMode === "atomic-replace" && contract?.patchMode === "atomic", `invalid consumer commit modes: ${policy.id ?? "unknown"}`); } @@ -147,6 +167,44 @@ requireValue(JSON.stringify(zoneV2Policy?.consumerContract?.fieldProjection) === requireValue(JSON.stringify(zoneV2Policy?.consumerContract?.geometryTypes) === JSON.stringify(["Polygon", "MultiPolygon"]), "map.zones.current.v2 geometry contract must be Polygon/MultiPolygon"); requireValue(!/(provider|tenant|connection|endpoint|credential|token|secret|authorization)/i.test(JSON.stringify(zoneV2Policy)), "map.zones.current.v2 consumer policy crosses the transport/secret boundary"); +const unitProfileV1Policy = consumerPolicies.find((policy) => policy.dataProductId === "fleet.units.profile.current.v1" && policy.productVersion === "1.0.0"); +const unitProfileV1Projection = [ + "corrected_engine_hours_factor", + "corrected_mileage_factor", + "display_name", + "filter_by_satellite_count_enabled", + "filter_by_satellite_count_value", + "filter_emissions", + "hardware_manufacturer_name", + "hardware_port", + "hardware_type_class", + "hardware_type_name", + "limit_acceleration", + "lost_connection_enabled", + "lost_connection_time_value", + "maximum_permissible_speed", + "maximum_valid_height", + "maximum_valid_speed", + "mileage_by_ignition", + "minimum_movement_speed", + "minimum_movement_time", + "minimum_parking_time", + "minimum_stop_time", + "minimum_trip_distance", + "minimum_valid_height", + "speed_parameter", + "trip_detection_type", + "unit_type_class", + "unit_type_name", + "use_odometer", +]; +requireValue(unitProfileV1Policy?.id === "map-moving-object-unit-profile-current-v1" && unitProfileV1Policy?.version === "1.0.0", "fleet.units.profile.current.v1 consumer policy is missing"); +requireValue(unitProfileV1Policy?.freshness === "none" && unitProfileV1Policy?.staleAfterMs === null, "fleet.units.profile.current.v1 must not invent freshness states"); +requireValue(JSON.stringify(unitProfileV1Policy?.consumerContract?.semanticTypes) === JSON.stringify(["map.moving_object"]), "fleet.units.profile.current.v1 semantic scope must be exact"); +requireValue(JSON.stringify(unitProfileV1Policy?.consumerContract?.fieldProjection) === JSON.stringify(unitProfileV1Projection), "fleet.units.profile.current.v1 field projection must be exact"); +requireValue(JSON.stringify(unitProfileV1Policy?.consumerContract?.geometryTypes) === JSON.stringify([]), "fleet.units.profile.current.v1 must remain a non-geometric aspect"); +requireValue(!/(provider|tenant|endpoint|credential|token|secret|authorization)/i.test(JSON.stringify(unitProfileV1Policy)), "fleet.units.profile.current.v1 consumer policy crosses the transport/secret boundary"); + const ids = components.components.map((component) => component.id); requireValue(new Set(ids).size === ids.length, "component ids must be unique"); const candidateIds = candidates.candidates.map((component) => component.id); diff --git a/server/foundry-data-product-consumer.mjs b/server/foundry-data-product-consumer.mjs index a688b35..b1fed4a 100644 --- a/server/foundry-data-product-consumer.mjs +++ b/server/foundry-data-product-consumer.mjs @@ -157,6 +157,7 @@ function validateConsumerContract(value) { const semanticTypes = Array.isArray(value.semanticTypes) ? value.semanticTypes : []; const fieldProjection = Array.isArray(value.fieldProjection) ? value.fieldProjection : []; const geometryTypes = Array.isArray(value.geometryTypes) ? value.geometryTypes : []; + const supportedGeometryTypes = new Set(["Point", "LineString", "Polygon", "MultiPolygon"]); const fieldTypes = value.fieldTypes && typeof value.fieldTypes === "object" && !Array.isArray(value.fieldTypes) ? value.fieldTypes : null; @@ -174,7 +175,9 @@ function validateConsumerContract(value) { || !fieldTypes || JSON.stringify(Object.keys(fieldTypes)) !== JSON.stringify(fieldProjection) || Object.values(fieldTypes).some((type) => !["string", "number", "boolean"].includes(type)) - || JSON.stringify(geometryTypes) !== JSON.stringify(["Polygon", "MultiPolygon"]) + || geometryTypes.length > supportedGeometryTypes.size + || new Set(geometryTypes).size !== geometryTypes.length + || geometryTypes.some((type) => !supportedGeometryTypes.has(type)) || value.subjectIdentity !== "semantic-type+source-id" || value.snapshotMode !== "atomic-replace" || value.patchMode !== "atomic" @@ -198,6 +201,8 @@ function assertConsumerContract(contract, product, binding) { field, product?.fieldContracts?.[field]?.type, ])); + const productHasGeometry = Array.isArray(product.fields) && product.fields.includes("geometry"); + const contractHasGeometry = contract.geometryTypes.length > 0; if ( product.ontologyRevision !== contract.ontologyRevision || product.deliveryMode !== contract.deliveryMode @@ -207,6 +212,8 @@ function assertConsumerContract(contract, product, binding) { || !Array.isArray(product.fields) || contract.fieldProjection.some((field) => !product.fields.includes(field)) || JSON.stringify(productFieldTypes) !== JSON.stringify(contract.fieldTypes) + || productHasGeometry !== contractHasGeometry + || (contractHasGeometry && product?.fieldContracts?.geometry?.type !== "geometry") ) throw consumerError("data_product_consumer_contract_mismatch", 409); } @@ -219,11 +226,17 @@ function assertEnvelopeProduct(envelope, product, errorCode) { function assertFactContract(fact, policy) { const contract = policy.consumerContract; if (!contract) return; - if ( - !contract.semanticTypes.includes(fact.semanticType) - || !fact.geometry - || !contract.geometryTypes.includes(fact.geometry.type) - ) throw consumerError("data_product_consumer_fact_contract_invalid", 502); + if (!contract.semanticTypes.includes(fact.semanticType)) { + throw consumerError("data_product_consumer_fact_contract_invalid", 502); + } + const hasGeometry = fact.geometry !== undefined && fact.geometry !== null; + if (contract.geometryTypes.length === 0) { + if (hasGeometry) throw consumerError("data_product_consumer_fact_contract_invalid", 502); + return; + } + if (!hasGeometry || !contract.geometryTypes.includes(fact.geometry.type)) { + throw consumerError("data_product_consumer_fact_contract_invalid", 502); + } } function validatePolicy(policy, product) { diff --git a/server/foundry-data-product-consumer.test.mjs b/server/foundry-data-product-consumer.test.mjs index 4b5653c..df3eda2 100644 --- a/server/foundry-data-product-consumer.test.mjs +++ b/server/foundry-data-product-consumer.test.mjs @@ -667,6 +667,7 @@ const zoneV2Product = { schedule_timezone: { type: "string", required: false }, applies_to_couriers: { type: "boolean", required: false }, applies_to_kicksharing: { type: "boolean", required: false }, + geometry: { type: "geometry", required: true }, }, active: true, }; @@ -862,3 +863,219 @@ test("zone v2 consumer fails closed on projection, product, envelope and geometr } } }); + +const unitProfileV1Projection = [ + "corrected_engine_hours_factor", + "corrected_mileage_factor", + "display_name", + "filter_by_satellite_count_enabled", + "filter_by_satellite_count_value", + "filter_emissions", + "hardware_manufacturer_name", + "hardware_port", + "hardware_type_class", + "hardware_type_name", + "limit_acceleration", + "lost_connection_enabled", + "lost_connection_time_value", + "maximum_permissible_speed", + "maximum_valid_height", + "maximum_valid_speed", + "mileage_by_ignition", + "minimum_movement_speed", + "minimum_movement_time", + "minimum_parking_time", + "minimum_stop_time", + "minimum_trip_distance", + "minimum_valid_height", + "speed_parameter", + "trip_detection_type", + "unit_type_class", + "unit_type_name", + "use_odometer", +]; + +const unitProfileV1FieldTypes = { + corrected_engine_hours_factor: "number", + corrected_mileage_factor: "number", + display_name: "string", + filter_by_satellite_count_enabled: "boolean", + filter_by_satellite_count_value: "number", + filter_emissions: "number", + hardware_manufacturer_name: "string", + hardware_port: "number", + hardware_type_class: "string", + hardware_type_name: "string", + limit_acceleration: "number", + lost_connection_enabled: "boolean", + lost_connection_time_value: "number", + maximum_permissible_speed: "number", + maximum_valid_height: "number", + maximum_valid_speed: "number", + mileage_by_ignition: "boolean", + minimum_movement_speed: "number", + minimum_movement_time: "number", + minimum_parking_time: "number", + minimum_stop_time: "number", + minimum_trip_distance: "number", + minimum_valid_height: "number", + speed_parameter: "string", + trip_detection_type: "string", + unit_type_class: "string", + unit_type_name: "string", + use_odometer: "boolean", +}; + +const unitProfileV1Policy = { + id: "map-moving-object-unit-profile-current-v1", + version: "1.0.0", + dataProductId: "fleet.units.profile.current.v1", + productVersion: "1.0.0", + freshness: "none", + staleAfterMs: null, + terminalStatuses: [], + consumerContract: { + ontologyRevision: "ontology.map.moving_object.v3", + deliveryMode: "snapshot+patch", + semanticTypes: ["map.moving_object"], + fieldProjection: unitProfileV1Projection, + fieldTypes: unitProfileV1FieldTypes, + geometryTypes: [], + subjectIdentity: "semantic-type+source-id", + snapshotMode: "atomic-replace", + patchMode: "atomic", + }, + removeMode: "canonical-tombstone-or-snapshot-rebase", +}; + +const unitProfileV1Product = { + id: "fleet.units.profile.current.v1", + version: "1.0.0", + ontologyRevision: "ontology.map.moving_object.v3", + deliveryMode: "snapshot+patch", + semanticTypes: ["map.moving_object"], + fields: [...unitProfileV1Projection], + fieldContracts: Object.fromEntries(unitProfileV1Projection.map((field) => [ + field, + { type: unitProfileV1FieldTypes[field], required: field === "display_name" }, + ])), + active: true, +}; + +function unitProfileV1Target(fieldProjection = unitProfileV1Projection) { + return { + application: { id: "66666666-6666-4666-8666-666666666666" }, + page: { id: "map" }, + binding: { + id: "trike-unit-profile", + dataProductId: unitProfileV1Product.id, + slotId: "points", + delivery: "snapshot+patch", + semanticTypes: ["map.moving_object"], + fieldProjection: [...fieldProjection], + }, + }; +} + +function unitProfileV1Fact({ geometry } = {}) { + return { + sourceId: "gelios-unit-501", + semanticType: "map.moving_object", + observedAt: "2026-07-23T06:28:52.000Z", + receivedAt: "2026-07-23T06:28:53.000Z", + attributes: { + display_name: "501 B2", + hardware_manufacturer_name: "Gelios", + hardware_type_name: "Trike tracker", + maximum_permissible_speed: 25, + minimum_movement_speed: 3, + use_odometer: true, + }, + ...(geometry === undefined ? {} : { geometry }), + }; +} + +function unitProfileV1Snapshot(facts = [unitProfileV1Fact()]) { + return { + schemaVersion: "nodedc.data-product.snapshot/v1", + dataProduct: { id: unitProfileV1Product.id, version: unitProfileV1Product.version }, + generatedAt: "2026-07-23T06:28:54.000Z", + cursor: "107", + facts, + }; +} + +test("non-geometric unit profile aspect is strictly typed and joins by stable subject identity", async () => { + const stateDir = await mkdtemp(join(tmpdir(), "foundry-consumer-unit-profile-v1-")); + const currentTarget = unitProfileV1Target(); + const manager = createFoundryDataProductConsumerManager({ + stateDir, + dataPlaneUrl: "http://edp.test", + resolveTarget: async () => structuredClone(currentTarget), + readReaderToken: async () => "ndc_edprb_profile-reader-capability", + inspectReaderGrant: async () => ({ product: unitProfileV1Product, readerGrantAction: "reuse", readerGrantGeneration: 1 }), + resolvePolicy: () => unitProfileV1Policy, + sanitizeSnapshot: (value) => value, + sanitizePatch: (value) => value, + fetchImpl: async () => Response.json(unitProfileV1Snapshot()), + }); + const input = { applicationId: currentTarget.application.id, pageId: currentTarget.page.id, bindingId: currentTarget.binding.id }; + try { + const plan = await manager.plan(input); + assert.deepEqual(plan.configuration.policy.consumerContract.geometryTypes, []); + assert.deepEqual(plan.configuration.policy.consumerContract.fieldProjection, unitProfileV1Projection); + const applied = await manager.apply({ ...input, planId: plan.planId }); + assert.equal(applied.consumer.subjectCount, 1); + assert.equal(applied.consumer.subjects[0].subjectId, "gelios-unit-501"); + assert.equal(applied.consumer.subjects[0].status, "active"); + const persisted = await manager.snapshot(currentTarget); + assert.equal(persisted.facts[0].geometry, undefined); + assert.equal(persisted.facts[0].attributes.hardware_type_name, "Trike tracker"); + } finally { + await manager.shutdown(); + await rm(stateDir, { recursive: true, force: true }); + } +}); + +test("non-geometric unit profile aspect fails closed on projection drift and injected geometry", async () => { + const variants = [ + { + name: "projection", + target: unitProfileV1Target(unitProfileV1Projection.slice(1)), + snapshot: unitProfileV1Snapshot(), + planError: /data_product_consumer_contract_mismatch/, + }, + { + name: "geometry", + target: unitProfileV1Target(), + snapshot: unitProfileV1Snapshot([unitProfileV1Fact({ geometry: { type: "Point", coordinates: [37.61, 55.75] } })]), + applyError: /data_product_consumer_fact_contract_invalid/, + }, + ]; + for (const variant of variants) { + const stateDir = await mkdtemp(join(tmpdir(), `foundry-consumer-unit-profile-v1-${variant.name}-`)); + const manager = createFoundryDataProductConsumerManager({ + stateDir, + dataPlaneUrl: "http://edp.test", + resolveTarget: async () => structuredClone(variant.target), + readReaderToken: async () => "ndc_edprb_profile-reader-capability", + inspectReaderGrant: async () => ({ product: unitProfileV1Product, readerGrantAction: "reuse", readerGrantGeneration: 1 }), + resolvePolicy: () => unitProfileV1Policy, + sanitizeSnapshot: (value) => value, + sanitizePatch: (value) => value, + fetchImpl: async () => Response.json(variant.snapshot), + }); + const input = { applicationId: variant.target.application.id, pageId: variant.target.page.id, bindingId: variant.target.binding.id }; + try { + if (variant.planError) { + await assert.rejects(manager.plan(input), variant.planError); + } else { + const plan = await manager.plan(input); + await assert.rejects(manager.apply({ ...input, planId: plan.planId }), variant.applyError); + } + } finally { + await manager.shutdown(); + await rm(stateDir, { recursive: true, force: true }); + } + } +});