@@ -4,10 +4,20 @@ use std::sync::Arc;
44use std:: time:: { Duration , Instant } ;
55
66use aruna_core:: NodeId ;
7+ use aruna_core:: effects:: { IterStart , StorageEffect } ;
78use aruna_core:: errors:: AuthorizationError ;
9+ use aruna_core:: events:: { Event , StorageEvent } ;
810use aruna_core:: id:: short_display_id;
11+ use aruna_core:: keyspaces:: {
12+ METADATA_EVENT_LOG_KEYSPACE , METADATA_GRAPH_LIFECYCLE_KEYSPACE ,
13+ METADATA_PENDING_PROJECTION_KEYSPACE ,
14+ } ;
915use aruna_core:: metadata:: {
10- MetadataError , MetadataQueryResults , MetadataRoCratePage , MetadataSearchHit ,
16+ MetadataCreateEventRecord , MetadataError , MetadataGraphLifecycleRecord , MetadataQueryResults ,
17+ MetadataRoCratePage , MetadataSearchHit ,
18+ } ;
19+ use aruna_core:: storage_entries:: {
20+ metadata_event_log_key, metadata_graph_lifecycle_key, metadata_pending_projection_target,
1121} ;
1222use aruna_core:: structs:: { AuthContext , MetadataRegistryRecord , Permission , RealmId } ;
1323use aruna_core:: telemetry:: record_elapsed_ms;
@@ -28,7 +38,7 @@ use crate::get_metadata_document::{
2838use crate :: get_realm_nodes:: GetRealmNodesOperation ;
2939use crate :: list_groups:: ListGroupOperation ;
3040use crate :: list_metadata_documents:: ListMetadataDocumentsOperation ;
31- use crate :: metadata:: repository:: StorageReadError ;
41+ use crate :: metadata:: repository:: { LIST_METADATA_PAGE_SIZE , StorageReadError } ;
3242
3343const DEFAULT_LIST_METADATA_LIMIT : usize = 50 ;
3444const MAX_LIST_METADATA_LIMIT : usize = 1_000 ;
@@ -198,14 +208,20 @@ pub async fn list_visible_metadata_documents(
198208 let offset = request. offset . unwrap_or ( 0 ) ;
199209
200210 let mut records = Vec :: new ( ) ;
211+ let include_pending_projection = request. include_summary ;
201212 if let Some ( group_id) = request. group_id {
202- records. extend ( load_group_metadata_records ( context, group_id) . await ?) ;
213+ records. extend (
214+ load_group_metadata_records ( context, group_id, include_pending_projection) . await ?,
215+ ) ;
203216 } else {
204217 let groups = drive ( ListGroupOperation :: new ( ) , context)
205218 . await
206219 . map_err ( |error| MetadataApiError :: Internal ( error. to_string ( ) ) ) ?;
207220 for group in groups {
208- records. extend ( load_group_metadata_records ( context, group. group_id ) . await ?) ;
221+ records. extend (
222+ load_group_metadata_records ( context, group. group_id , include_pending_projection)
223+ . await ?,
224+ ) ;
209225 }
210226 }
211227
@@ -237,9 +253,14 @@ pub async fn list_visible_metadata_documents(
237253 )
238254 . await ;
239255 for ( record, summary) in selected. into_iter ( ) . zip ( summaries) {
256+ let rocrate_summary_jsonld = match summary {
257+ Ok ( summary) => Some ( summary) ,
258+ Err ( MetadataApiError :: ServiceUnavailable ) => None ,
259+ Err ( error) => return Err ( error) ,
260+ } ;
240261 documents. push ( ListedMetadataDocument {
241262 record,
242- rocrate_summary_jsonld : Some ( summary? ) ,
263+ rocrate_summary_jsonld,
243264 } ) ;
244265 }
245266 } else {
@@ -403,10 +424,11 @@ pub async fn load_metadata_realm_nodes(
403424async fn load_group_metadata_records (
404425 context : & DriverContext ,
405426 group_id : GroupId ,
427+ include_pending_projection : bool ,
406428) -> Result < Vec < MetadataRegistryRecord > , MetadataApiError > {
407429 // Listing remains eventually consistent: the handle-owned visibility cache
408430 // serves stale snapshots while one refill updates the operation-owned read path.
409- if let Some ( metadata_handle) = context. metadata_handle . as_ref ( ) {
431+ if !include_pending_projection && let Some ( metadata_handle) = context. metadata_handle . as_ref ( ) {
410432 match metadata_handle
411433 . list_cached_registry_records_for_group ( group_id)
412434 . await
@@ -421,9 +443,164 @@ async fn load_group_metadata_records(
421443 }
422444 }
423445
424- drive ( ListMetadataDocumentsOperation :: new ( group_id) , context)
446+ let mut records = drive ( ListMetadataDocumentsOperation :: new ( group_id) , context)
447+ . await
448+ . map_err ( |err| MetadataApiError :: Internal ( err. to_string ( ) ) ) ?;
449+
450+ if include_pending_projection {
451+ merge_pending_metadata_records (
452+ & mut records,
453+ load_pending_group_metadata_records ( context, group_id) . await ?,
454+ ) ;
455+ records. sort_by_key ( |record| record. document_id ) ;
456+ }
457+
458+ Ok ( records)
459+ }
460+
461+ async fn load_pending_group_metadata_records (
462+ context : & DriverContext ,
463+ group_id : GroupId ,
464+ ) -> Result < Vec < MetadataRegistryRecord > , MetadataApiError > {
465+ let mut records = Vec :: new ( ) ;
466+ let mut start_after = None ;
467+
468+ loop {
469+ let page = context
470+ . storage_handle
471+ . send_storage_effect ( StorageEffect :: Iter {
472+ key_space : METADATA_PENDING_PROJECTION_KEYSPACE . to_string ( ) ,
473+ prefix : None ,
474+ start : start_after. take ( ) . map ( IterStart :: After ) ,
475+ limit : LIST_METADATA_PAGE_SIZE ,
476+ txn_id : None ,
477+ } )
478+ . await ;
479+ let ( values, next_start_after) = match page {
480+ Event :: Storage ( StorageEvent :: IterResult {
481+ values,
482+ next_start_after,
483+ } ) => ( values, next_start_after) ,
484+ Event :: Storage ( StorageEvent :: Error { error } ) => {
485+ return Err ( MetadataApiError :: Internal ( error. to_string ( ) ) ) ;
486+ }
487+ other => return Err ( MetadataApiError :: Internal ( format ! ( "{other:?}" ) ) ) ,
488+ } ;
489+
490+ for ( key, _) in values {
491+ let Some ( ( document_id, event_id) ) = metadata_pending_projection_target ( key. as_ref ( ) )
492+ else {
493+ continue ;
494+ } ;
495+ let Some ( event) = read_metadata_create_event ( context, document_id, event_id) . await ?
496+ else {
497+ continue ;
498+ } ;
499+ let record = event. record ;
500+ if record. group_id != group_id {
501+ continue ;
502+ }
503+ if metadata_graph_is_deleted ( context, & record. graph_iri ) . await ? {
504+ continue ;
505+ }
506+ records. push ( record) ;
507+ }
508+
509+ if next_start_after. is_none ( ) {
510+ break ;
511+ }
512+ start_after = next_start_after;
513+ }
514+
515+ Ok ( records)
516+ }
517+
518+ async fn read_metadata_create_event (
519+ context : & DriverContext ,
520+ document_id : Ulid ,
521+ event_id : Ulid ,
522+ ) -> Result < Option < MetadataCreateEventRecord > , MetadataApiError > {
523+ let value = match context
524+ . storage_handle
525+ . send_storage_effect ( StorageEffect :: Read {
526+ key_space : METADATA_EVENT_LOG_KEYSPACE . to_string ( ) ,
527+ key : metadata_event_log_key ( document_id, event_id) ,
528+ txn_id : None ,
529+ } )
530+ . await
531+ {
532+ Event :: Storage ( StorageEvent :: ReadResult { value, .. } ) => value,
533+ Event :: Storage ( StorageEvent :: Error { error } ) => {
534+ return Err ( MetadataApiError :: Internal ( error. to_string ( ) ) ) ;
535+ }
536+ other => return Err ( MetadataApiError :: Internal ( format ! ( "{other:?}" ) ) ) ,
537+ } ;
538+ let Some ( value) = value else {
539+ return Ok ( None ) ;
540+ } ;
541+
542+ let event: MetadataCreateEventRecord = postcard:: from_bytes ( & value)
543+ . map_err ( |error| MetadataApiError :: Internal ( error. to_string ( ) ) ) ?;
544+ if event. record . document_id != document_id || event. event_id != event_id {
545+ return Err ( MetadataApiError :: Internal ( format ! (
546+ "metadata create event log target {document_id}/{event_id} did not match payload {}/{}" ,
547+ event. record. document_id, event. event_id
548+ ) ) ) ;
549+ }
550+ Ok ( Some ( event) )
551+ }
552+
553+ async fn metadata_graph_is_deleted (
554+ context : & DriverContext ,
555+ graph_iri : & str ,
556+ ) -> Result < bool , MetadataApiError > {
557+ match context
558+ . storage_handle
559+ . send_storage_effect ( StorageEffect :: Read {
560+ key_space : METADATA_GRAPH_LIFECYCLE_KEYSPACE . to_string ( ) ,
561+ key : metadata_graph_lifecycle_key ( graph_iri) ,
562+ txn_id : None ,
563+ } )
425564 . await
426- . map_err ( |err| MetadataApiError :: Internal ( err. to_string ( ) ) )
565+ {
566+ Event :: Storage ( StorageEvent :: ReadResult {
567+ value : Some ( value) , ..
568+ } ) => {
569+ let record: MetadataGraphLifecycleRecord = postcard:: from_bytes ( & value)
570+ . map_err ( |error| MetadataApiError :: Internal ( error. to_string ( ) ) ) ?;
571+ Ok ( record. is_deleted ( ) )
572+ }
573+ Event :: Storage ( StorageEvent :: ReadResult { value : None , .. } ) => Ok ( false ) ,
574+ Event :: Storage ( StorageEvent :: Error { error } ) => {
575+ Err ( MetadataApiError :: Internal ( error. to_string ( ) ) )
576+ }
577+ other => Err ( MetadataApiError :: Internal ( format ! ( "{other:?}" ) ) ) ,
578+ }
579+ }
580+
581+ fn merge_pending_metadata_records (
582+ records : & mut Vec < MetadataRegistryRecord > ,
583+ pending_records : Vec < MetadataRegistryRecord > ,
584+ ) {
585+ let mut positions = records
586+ . iter ( )
587+ . enumerate ( )
588+ . map ( |( index, record) | ( record. document_id , index) )
589+ . collect :: < HashMap < _ , _ > > ( ) ;
590+
591+ for pending_record in pending_records {
592+ if let Some ( & index) = positions. get ( & pending_record. document_id ) {
593+ let existing_record = & records[ index] ;
594+ if ( pending_record. updated_at_ms , pending_record. last_event_id )
595+ > ( existing_record. updated_at_ms , existing_record. last_event_id )
596+ {
597+ records[ index] = pending_record;
598+ }
599+ } else {
600+ positions. insert ( pending_record. document_id , records. len ( ) ) ;
601+ records. push ( pending_record) ;
602+ }
603+ }
427604}
428605
429606async fn load_record_by_document (
0 commit comments