@@ -9,7 +9,8 @@ import {performance} from "node:perf_hooks";
99import { parseArgs } from "node:util" ;
1010import { fileURLToPath } from "node:url" ;
1111import type { JsonObject } from "../../loopx/control_plane/effect_program.ts" ;
12- import type { AuthorityStore , AuthorityStoreCommit } from "../../loopx/control_plane/coordination/authority_store.ts" ;
12+ import type { AuthorityStore , AuthorityStoreCommit , AuthorityStoreCommitResult } from "../../loopx/control_plane/coordination/authority_store.ts" ;
13+ import { canonicalAuthorityBytes } from "../../loopx/control_plane/coordination/authority_store_codec.ts" ;
1314import { FileAuthorityStore } from "../../loopx/control_plane/coordination/file_authority_store.ts" ;
1415import { SqliteAuthorityStore } from "../../loopx/control_plane/coordination/sqlite_authority_store.ts" ;
1516import { sqliteRuntimeIdentity } from "../../loopx/control_plane/coordination/sqlite_runtime.ts" ;
@@ -84,6 +85,9 @@ async function measure(root: string) {
8485 const finalProjection = projectionAt ( count - 1 ) ;
8586 let revision : string | null = null ;
8687 let first : AuthorityStoreCommit | undefined ;
88+ let firstResult : Extract < AuthorityStoreCommitResult , { status : "applied" } > | undefined ;
89+ const receiptsAt = ( index : number ) => [ { operation_id : `op-${ index } ` , index,
90+ metadata : { checked : true , labels : [ "synthetic" , "保留" ] } } ] ;
8791 let filePublicationBytes = 0 ;
8892 const fileBytes = ( ) : number => readdirSync ( root ) . reduce ( ( sum , name ) => sum + statSync ( join ( root , name ) ) . size , 0 ) ;
8993 const timed = async < T > ( action : ( ) => Promise < T > , into : number [ ] ) : Promise < T > => {
@@ -93,11 +97,11 @@ async function measure(root: string) {
9397 for ( let index = 0 ; index < count ; index ++ ) {
9498 const input : AuthorityStoreCommit = { expected_provider_revision : revision , operation_id : `op-${ index } ` ,
9599 next_projection : projectionAt ( index ) , events : [ { kind : "observation" , index} ] ,
96- receipts : [ { operation_id : `op- ${ index } ` , index , metadata : { checked : true , labels : [ "synthetic" , "保留" ] } } ] } ;
100+ receipts : receiptsAt ( index ) } ;
97101 const result = await timed ( ( ) => store . commitAuthority ( input ) , commits ) ;
98102 assert . equal ( result . status , "applied" ) ; if ( result . status !== "applied" ) throw new Error ( "commit rejected" ) ;
99103 revision = result . provider_revision ;
100- if ( index === 0 ) first = input ;
104+ if ( index === 0 ) { first = input ; firstResult = result ; }
101105 if ( values . provider === "file" ) filePublicationBytes += statSync ( ( store as FileAuthorityStore ) . path ) . size ;
102106 if ( ( index + 1 ) % 128 === 0 ) process . stderr . write ( `${ values . provider } ${ values . workload } : ${ index + 1 } /${ count } \n` ) ;
103107 }
@@ -109,6 +113,10 @@ async function measure(root: string) {
109113 const receiptIndex = Math . floor ( index * ( count - 1 ) / ( samples - 1 ) ) ;
110114 const receipt = await timed ( ( ) => store . readReceipt ( `op-${ receiptIndex } ` ) , reads ) ;
111115 assert . equal ( receipt . status , "found" ) ;
116+ if ( receipt . status === "found" ) {
117+ assert . equal ( receipt . cursor , String ( receiptIndex + 1 ) ) ;
118+ assert . deepEqual ( receipt . receipts , receiptsAt ( receiptIndex ) ) ;
119+ }
112120 const page = await timed ( ( ) => store . scanCommitted ( String ( count - 100 ) , 100 ) , scans ) ;
113121 assert . equal ( page . status , "page" ) ;
114122 if ( page . status === "page" ) {
@@ -117,8 +125,8 @@ async function measure(root: string) {
117125 const ordinal = count - 100 + offset ;
118126 assert . equal ( row . operation_id , `op-${ ordinal } ` ) ;
119127 assert . deepEqual ( row . projection , projectionAt ( ordinal ) ) ;
120- assert . deepEqual ( row . receipts , [ { operation_id : `op- ${ ordinal } ` , index : ordinal ,
121- metadata : { checked : true , labels : [ "synthetic" , "保留" ] } } ] ) ;
128+ assert . deepEqual ( row . events , [ { kind : "observation" , index : ordinal } ] ) ;
129+ assert . deepEqual ( row . receipts , receiptsAt ( ordinal ) ) ;
122130 }
123131 }
124132 }
@@ -132,13 +140,62 @@ async function measure(root: string) {
132140 cold . push ( performance . now ( ) - started ) ; assert . equal ( child . status , 0 , child . stderr ) ;
133141 assert . deepEqual ( JSON . parse ( child . stdout ) . head , finalProjection ) ;
134142 }
135- const reopened = openStore ( root ) , replay = await reopened . commitAuthority ( first ! ) ;
136- assert . equal ( replay . status , "conflict" ) ; // Current store contract reconciles via readReceipt.
137- const original = await reopened . readReceipt ( first ! . operation_id ) ;
143+ // Qualification is outside the timings. Walk bounded pages without retaining
144+ // N full projections, checking independently generated input and exact history.
145+ const reopened = openStore ( root ) ;
146+ const historyDigest = async ( ) => {
147+ const digest = createHash ( "sha256" ) ;
148+ let cursor : string | null = null , checked = 0 ;
149+ for ( ; ; ) {
150+ const page = await reopened . scanCommitted ( cursor , 100 ) ;
151+ assert . equal ( page . status , "page" ) ; if ( page . status !== "page" ) throw new Error ( "history read rejected" ) ;
152+ assert ( page . transactions . length > 0 ) ;
153+ for ( const row of page . transactions ) {
154+ assert ( checked < count , "history includes an unexpected transaction" ) ;
155+ assert . equal ( row . operation_id , `op-${ checked } ` ) ;
156+ assert . equal ( row . cursor , String ( checked + 1 ) ) ;
157+ assert . deepEqual ( row . projection , projectionAt ( checked ) ) ;
158+ assert . deepEqual ( row . events , [ { kind : "observation" , index : checked } ] ) ;
159+ assert . deepEqual ( row . receipts , receiptsAt ( checked ) ) ;
160+ digest . update ( canonicalAuthorityBytes ( row ) ) ; digest . update ( "\n" ) ; checked ++ ;
161+ }
162+ if ( ! page . has_more ) break ;
163+ assert . equal ( page . next_cursor , String ( checked ) ) ;
164+ cursor = page . next_cursor ;
165+ }
166+ assert . equal ( checked , count ) ;
167+ return digest . digest ( "hex" ) ;
168+ } ;
169+ assert ( first && firstResult ) ;
170+ const before = await reopened . loadAuthority ( ) , original = await reopened . readReceipt ( first . operation_id ) ;
171+ assert . equal ( before . status , "loaded" ) ;
172+ if ( before . status === "loaded" ) {
173+ assert . equal ( before . provider_revision , revision ) ; assert . equal ( before . cursor , String ( count ) ) ;
174+ assert . deepEqual ( before . head , finalProjection ) ;
175+ }
138176 assert . equal ( original . status , "found" ) ;
139- if ( original . status === "found" ) assert . deepEqual ( original . receipts , first ! . receipts ) ;
140- const after = await reopened . loadAuthority ( ) ;
141- assert . equal ( after . status , "loaded" ) ; if ( after . status === "loaded" ) assert . equal ( after . provider_revision , revision ) ;
177+ if ( original . status === "found" ) {
178+ assert . equal ( original . provider_revision , firstResult . provider_revision ) ;
179+ assert . equal ( original . cursor , firstResult . cursor ) ; assert . deepEqual ( original . receipts , first . receipts ) ;
180+ }
181+ const retainedHistory = await historyDigest ( ) ;
182+ // An identical historical intent returns its original result despite a stale
183+ // basis. Projection-, event- and receipt-only drift must conflict, not append.
184+ for ( const change of [ "none" , "projection" , "events" , "receipts" ] as const ) {
185+ const input = structuredClone ( first ) ;
186+ if ( change === "projection" ) input . next_projection . replay_marker = "different" ;
187+ if ( change === "events" ) input . events = [ { kind : "observation" , index : - 1 } ] ;
188+ if ( change === "receipts" ) input . receipts = [ { ...receiptsAt ( 0 ) [ 0 ] , replay_marker : "different" } ] ;
189+ const replay = await reopened . commitAuthority ( input ) ;
190+ if ( change === "none" ) assert . deepEqual ( replay , firstResult ) ;
191+ else {
192+ assert . equal ( replay . status , "conflict" , `${ change } -only drift was accepted` ) ;
193+ if ( replay . status === "conflict" ) assert . equal ( replay . conflict_kind , "operation_id_exists" ) ;
194+ }
195+ assert . deepEqual ( await reopened . loadAuthority ( ) , before , `${ change } changed the later head` ) ;
196+ assert . deepEqual ( await reopened . readReceipt ( first . operation_id ) , original , `${ change } changed the original receipt` ) ;
197+ }
198+ assert . equal ( await historyDigest ( ) , retainedHistory , "replay attempts changed retained history" ) ;
142199 assert . deepEqual ( sourceIdentity ( ) , source , "measurement source changed while running" ) ;
143200 return { schema_version : "loopx_local_provider_comparison_v0" , provider : values . provider , workload : values . workload ,
144201 source, node : process . version , sqlite : sqliteRuntimeIdentity ( ) , platform : process . platform , arch : process . arch ,
@@ -148,5 +205,6 @@ async function measure(root: string) {
148205 post_fill_rss_bytes : postFillRss ,
149206 final_store_bytes : fileBytes ( ) , file_document_publication_bytes : values . provider === "file" ? filePublicationBytes : null ,
150207 complete_record_and_receipt_checks : "passed" , original_receipt_recovery_after_reopen : "passed" ,
208+ historical_replay_and_conflict_checks : "passed" ,
151209 limits : "bounded sequential store experiment; cold process includes module loading, not cold OS cache; RSS includes fixture and verification allocations, not a steady-state qualification; no CLI, concurrent writers, crash, soak or formal D2 qualification; File publication bytes are application bytes, not physical writes; SQLite WAL traffic is measured by the separate capacity runner" } ;
152210}
0 commit comments