@@ -1794,17 +1794,26 @@ test('owner release revalidates its authoritative winner inside the sibling comp
17941794 const result = await broker . releaseOwner ( socket , ownerId , [ ] , Date . now ( ) + 2_000 ) ; assert . deepEqual ( result . releasedSessionIds , [ ] ) ; assert . deepEqual ( result . failedSessionIds . sort ( ) , [ sessionId , siblingId ] . sort ( ) ) ; assert . deepEqual ( JSON . parse ( await readFile ( ownershipPath , 'utf8' ) ) . sessions , before ) ; assert . equal ( broker . uncertainOwnerReleases . size , 0 ) ; await rm ( directory , { recursive : true , force : true } ) ;
17951795} ) ;
17961796
1797- test ( 'owner release aborts its unlocked winner read after a reset compensation misses the deadline' , { timeout : scaleTestTimeout ( 3_000 ) } , async ( t ) => {
1798- const directory = await mkdtemp ( join ( tmpdir ( ) , 'zcode-broker-release-compensation-deadline-' ) ) ; const endpoint = join ( directory , 'broker.sock' ) ; const ownershipPath = `${ endpoint } .owners.json` ; const ownerId = 'release-compensation-deadline-owner' ; const sessionId = 'release-compensation-deadline-session' ; const siblingId = 'release-compensation-deadline-sibling' ; const socket = { writable : true , destroyed : false , zcodeWriter : { write ( ) { } } , destroy ( ) { } } ; const broker = newTestBroker ( { endpoint, brokerToken : '2' . repeat ( 64 ) , workspace : directory , launch : { command : process . execPath , args : [ fixture ] , target : fixture } } ) ; broker . sessionOwners . set ( sessionId , { ownerId, socket, claimToken : null } ) ; broker . sessionOwners . set ( siblingId , { ownerId, socket, claimToken : null } ) ; broker . activeSessionSockets . set ( sessionId , { socket, token : 'release-compensation-deadline-turn' , baseline : 1 , inputId : 'release-compensation-deadline-input' } ) ; broker . activeSessions . add ( sessionId ) ; await writeFile ( ownershipPath , JSON . stringify ( { version : 1 , sessions : { [ sessionId ] : ownerId , [ siblingId ] : ownerId } } ) ) ; broker . ownershipStoreEstablished = true ; const protocol = { request : async ( ) => ( { } ) , cancelTurn ( ) { } } ; broker . protocol = protocol ; let writes = 0 ; let observedSignal ; let compensationSignal ; let releaseOutcome ; broker . writeOwnerStore = async ( sessions , options ) => { writes += 1 ; if ( writes === 1 ) { await atomicWriteJson ( ownershipPath , { version : 1 , sessions } ) ; retireTestSessionLease ( broker , siblingId ) ; broker . clearProtocolGeneration ( protocol ) ; return ; } compensationSignal = options . signal ; await new Promise ( ( resolvePromise , rejectPromise ) => { if ( options . signal . aborted ) { rejectPromise ( options . signal . reason ) ; return ; } options . signal . addEventListener ( 'abort' , ( ) => rejectPromise ( options . signal . reason ) , { once : true } ) ; } ) ; } ;
1797+ test ( 'owner release aborts its unlocked winner read after a reset compensation misses the deadline' , { timeout : scaleTestTimeout ( 5_000 ) } , async ( t ) => {
1798+ const directory = await mkdtemp ( join ( tmpdir ( ) , 'zcode-broker-release-compensation-deadline-' ) ) ; const endpoint = join ( directory , 'broker.sock' ) ; const ownershipPath = `${ endpoint } .owners.json` ; const ownerId = 'release-compensation-deadline-owner' ; const sessionId = 'release-compensation-deadline-session' ; const siblingId = 'release-compensation-deadline-sibling' ; const socket = { writable : true , destroyed : false , zcodeWriter : { write ( ) { } } , destroy ( ) { } } ; const broker = newTestBroker ( { endpoint, brokerToken : '2' . repeat ( 64 ) , workspace : directory , launch : { command : process . execPath , args : [ fixture ] , target : fixture } } ) ; broker . sessionOwners . set ( sessionId , { ownerId, socket, claimToken : null } ) ; broker . sessionOwners . set ( siblingId , { ownerId, socket, claimToken : null } ) ; broker . activeSessionSockets . set ( sessionId , { socket, token : 'release-compensation-deadline-turn' , baseline : 1 , inputId : 'release-compensation-deadline-input' } ) ; broker . activeSessions . add ( sessionId ) ; await writeFile ( ownershipPath , JSON . stringify ( { version : 1 , sessions : { [ sessionId ] : ownerId , [ siblingId ] : ownerId } } ) ) ; broker . ownershipStoreEstablished = true ; const protocol = { request : async ( ) => ( { } ) , cancelTurn ( ) { } } ; broker . protocol = protocol ; let writes = 0 ; let observedSignal ; let compensationSignal ; let releaseOutcome ; let releaseOutcomeSettled = false ; broker . writeOwnerStore = async ( sessions , options ) => { writes += 1 ; if ( writes === 1 ) { await atomicWriteJson ( ownershipPath , { version : 1 , sessions } ) ; retireTestSessionLease ( broker , siblingId ) ; broker . clearProtocolGeneration ( protocol ) ; return ; } compensationSignal = options . signal ; await new Promise ( ( resolvePromise , rejectPromise ) => { if ( options . signal . aborted ) { rejectPromise ( options . signal . reason ) ; return ; } options . signal . addEventListener ( 'abort' , ( ) => rejectPromise ( options . signal . reason ) , { once : true } ) ; } ) ; } ;
17991799 let enterSecondWrite = ( ) => { } ; const secondWriteEntered = new Promise ( ( resolvePromise ) => { enterSecondWrite = resolvePromise ; } ) ; const writeOwnerStore = broker . writeOwnerStore ; broker . writeOwnerStore = async ( ...args ) => { const operation = writeOwnerStore ( ...args ) ; if ( writes === 2 ) enterSecondWrite ( ) ; return operation ; } ;
18001800 broker . readOwnerStoreUnlocked = async ( _allowMissing , options = { } ) => { observedSignal = options . signal ; assert . ok ( observedSignal , 'unlocked winner read must receive the compensation signal' ) ; assert . equal ( observedSignal , compensationSignal ) ; observedSignal . throwIfAborted ( ) ; return { exists : true , sessions : Object . create ( null ) } ; } ;
18011801 t . after ( async ( ) => {
1802- if ( releaseOutcome ) await releaseOutcome ;
1803- for ( let turn = 0 ; turn < 100 && broker . releaseTasks . size ; turn += 1 ) await new Promise ( ( resolvePromise ) => setImmediate ( resolvePromise ) ) ;
1804- if ( broker . releaseTasks . size ) await Promise . allSettled ( [ ...broker . releaseTasks ] ) ;
1805- await rm ( directory , { recursive : true , force : true } ) ;
1802+ const cleanupErrors = [ ] ;
1803+ const waitForCleanup = async ( operation , describeResidual ) => {
1804+ let timer ;
1805+ try {
1806+ await Promise . race ( [ operation , new Promise ( ( _ , rejectPromise ) => { timer = setTimeout ( ( ) => rejectPromise ( new Error ( `owner release cleanup timed out: ${ describeResidual ( ) } ` ) ) , scaleTestTimeout ( 500 ) ) ; } ) ] ) ;
1807+ } finally { clearTimeout ( timer ) ; }
1808+ } ;
1809+ try {
1810+ try { await waitForCleanup ( releaseOutcome ?? Promise . resolve ( ) , ( ) => `caller residual=${ releaseOutcome ? ( releaseOutcomeSettled ? 'settled' : 'pending' ) : 'not-started' } ` ) ; } catch ( error ) { cleanupErrors . push ( error ) ; }
1811+ for ( let turn = 0 ; turn < 100 && broker . releaseTasks . size ; turn += 1 ) await new Promise ( ( resolvePromise ) => setImmediate ( resolvePromise ) ) ;
1812+ try { await waitForCleanup ( Promise . allSettled ( [ ...broker . releaseTasks ] ) , ( ) => `task residual=${ broker . releaseTasks . size } ` ) ; } catch ( error ) { cleanupErrors . push ( error ) ; }
1813+ if ( cleanupErrors . length ) throw new AggregateError ( cleanupErrors , `owner release cleanup failed: caller residual=${ releaseOutcome ? ( releaseOutcomeSettled ? 'settled' : 'pending' ) : 'not-started' } ; task residual=${ broker . releaseTasks . size } ` ) ;
1814+ } finally { await rm ( directory , { recursive : true , force : true } ) ; }
18061815 } ) ;
1807- const releasing = broker . releaseOwner ( socket , ownerId , [ ] , Date . now ( ) + 1_000 ) ; releaseOutcome = releasing . then ( ( value ) => ( { kind : 'fulfilled' , value } ) , ( error ) => ( { kind : 'rejected' , error } ) ) ;
1816+ const releasing = broker . releaseOwner ( socket , ownerId , [ ] , Date . now ( ) + scaleTestTimeout ( 1_000 ) ) ; releaseOutcome = releasing . then ( ( value ) => ( { kind : 'fulfilled' , value } ) , ( error ) => ( { kind : 'rejected' , error } ) ) ; void releaseOutcome . then ( ( ) => { releaseOutcomeSettled = true ; } ) ;
18081817 const boundary = await Promise . race ( [ secondWriteEntered . then ( ( ) => 'second-write' ) , releaseOutcome . then ( ( ) => 'release-settled' ) ] ) ; assert . equal ( boundary , 'second-write' ) ; assert . equal ( writes , 2 ) ;
18091818 await assert . rejects ( releasing , { code : 'ZCODE_OWNER_RELEASE_TIMEOUT' } ) ;
18101819 for ( let turn = 0 ; turn < 100 && broker . releaseTasks . size ; turn += 1 ) await new Promise ( ( resolvePromise ) => setImmediate ( resolvePromise ) ) ;
0 commit comments