1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707270827092710271127122713271427152716271727182719272027212722272327242725272627272728272927302731273227332734273527362737273827392740274127422743274427452746274727482749275027512752275327542755275627572758275927602761276227632764276527662767276827692770277127722773277427752776277727782779278027812782278327842785278627872788278927902791279227932794279527962797279827992800280128022803280428052806280728082809281028112812281328142815281628172818281928202821282228232824282528262827282828292830283128322833283428352836283728382839284028412842284328442845284628472848284928502851285228532854285528562857285828592860286128622863286428652866286728682869287028712872287328742875287628772878287928802881288228832884288528862887288828892890289128922893289428952896289728982899290029012902290329042905290629072908290929102911291229132914291529162917291829192920292129222923292429252926292729282929293029312932293329342935293629372938293929402941294229432944294529462947294829492950295129522953295429552956295729582959296029612962296329642965296629672968296929702971297229732974297529762977297829792980298129822983298429852986298729882989299029912992299329942995299629972998299930003001300230033004300530063007300830093010301130123013301430153016301730183019302030213022302330243025302630273028302930303031303230333034303530363037303830393040304130423043304430453046304730483049305030513052305330543055305630573058305930603061306230633064306530663067306830693070307130723073307430753076307730783079308030813082308330843085308630873088308930903091309230933094309530963097309830993100310131023103310431053106310731083109311031113112311331143115311631173118311931203121312231233124312531263127312831293130313131323133313431353136313731383139314031413142314331443145314631473148314931503151315231533154315531563157315831593160316131623163316431653166316731683169317031713172317331743175317631773178317931803181318231833184318531863187318831893190319131923193319431953196319731983199320032013202320332043205320632073208320932103211321232133214321532163217321832193220322132223223322432253226322732283229323032313232323332343235323632373238323932403241324232433244324532463247324832493250325132523253325432553256325732583259326032613262326332643265326632673268326932703271327232733274327532763277327832793280328132823283328432853286328732883289329032913292329332943295329632973298329933003301330233033304330533063307330833093310331133123313331433153316331733183319332033213322332333243325332633273328332933303331333233333334333533363337333833393340334133423343334433453346334733483349335033513352335333543355335633573358335933603361336233633364336533663367336833693370337133723373337433753376337733783379338033813382338333843385338633873388338933903391339233933394339533963397339833993400340134023403340434053406340734083409341034113412341334143415341634173418341934203421342234233424342534263427342834293430343134323433343434353436343734383439344034413442344334443445344634473448344934503451345234533454345534563457345834593460346134623463346434653466346734683469347034713472347334743475347634773478347934803481348234833484348534863487348834893490349134923493349434953496349734983499350035013502350335043505350635073508350935103511351235133514351535163517351835193520352135223523352435253526352735283529353035313532353335343535353635373538353935403541354235433544354535463547354835493550355135523553355435553556355735583559356035613562356335643565356635673568356935703571357235733574357535763577357835793580(* Phase 1 store. Two backends share the same interface:
- [Mem] — pure in-memory [Bytes_map]-per-tree (the Phase 0 backend).
Used by [create ()]. No I/O, no size limits, no errors.
- [Btree] — CoW B+-tree over a Pager over a BLOCK device, given as
I/O callbacks via [open_block]/[open_block_wal]. Persists across
reopen. Inherits the B+-tree leaf-cell size limits (512-byte keys,
1024-byte values). Unix-file convenience constructors live in the
[granary.unix] driver library, not here, so the core stays
platform-agnostic (#170).
The two are wrapped in a sum type so callers see one [Store.t]. *)openLwt.SyntaxmoduleBtree=Granary_storage.BtreemodulePager=Granary_storage.PagermoduleHeader=Granary_storage.HeadermoduleFreelist=Granary_storage.FreelistmodulePager_event=Granary_storage.Pager_eventmodulePage=Granary_storage.PagemoduleGeometry=Granary_storage.GeometrymoduleCrypto=Granary_storage.CryptomoduleVarint=Granary_encoding.VarintmoduleBytes_map=Map.Make(Bytes)typerotyperwtypetree_id=inttypeerror=|Block_errorofstring|Corruptionofstring|Key_too_largeofint|Value_too_largeofint|Header_errorofstring|Encryption_key_required(** DB is encrypted but no key was supplied *)|Encryption_key_mismatch(** supplied key fails the header canary *)|Not_encrypted(** a key was supplied for a plaintext DB *)|Encryption_rng_unseeded(** a key was supplied but {!Mirage_crypto_rng} is not seeded, so no per-page
nonce can be generated — the application must seed the RNG at boot *)|History_unavailable(** as-of API used on a store opened without the feature *)|History_pruned(** as-of target is older than the retained floor *)|History_misconfigured(** [as_of_history:true] but no history sink supplied *)letpp_errorfmt=function|Block_errors->Format.fprintffmt"Block_error(%s)"s|Corruptions->Format.fprintffmt"Corruption(%s)"s|Key_too_largen->Format.fprintffmt"Key_too_large(%d)"n|Value_too_largen->Format.fprintffmt"Value_too_large(%d)"n|Header_errors->Format.fprintffmt"Header_error(%s)"s|Encryption_key_required->Format.pp_print_stringfmt"Encryption_key_required"|Encryption_key_mismatch->Format.pp_print_stringfmt"Encryption_key_mismatch"|Not_encrypted->Format.pp_print_stringfmt"Not_encrypted"|Encryption_rng_unseeded->Format.pp_print_stringfmt"Encryption_rng_unseeded"|History_unavailable->Format.fprintffmt"as-of time travel is not enabled on this database"|History_pruned->Format.fprintffmt"as-of target is older than the retained history horizon"|History_misconfigured->Format.fprintffmt"as_of_history was requested but no history log was supplied";;(* ------------------------------------------------------------------ *)(* Btree-backend internal state *)(* ------------------------------------------------------------------ *)(* The meta-tree is stored separately from user trees (it doesn't live
in the [trees] hashtable). It tracks the root_page of every tree_id
created via [get/put/del]; its OWN root_page is what we commit into
the header. *)typebt_savepoint={sp_name:string;sp_meta_root:int64;sp_tree_roots:(tree_id*int64)list;sp_freelist:Freelist.t;sp_n_pages:int64;sp_dirty:Pager.dirty_snapshot;sp_txn_pool:int64list(** #297: txn-owned page pool snapshot, restored on savepoint rollback. *)}(* Per-store commit queue for WAL-mode group commit (#77, #151). After
a writer has staged its WAL frames it releases [lock] and joins
this queue to await one shared fsync. [drainer] is set to [true] by
the first arriving writer (acting as coordinator); subsequent writers
register a resolver on [waiters] and block until the drainer wakes
them with the sync result. [pending] tracks the number of waiters so
the drainer can yield additional ticks while new arrivals keep
registering, widening the batch. Cooperative Lwt scheduling makes
the [drainer]/[pending]/[waiters] transitions atomic (no implicit
yield between read and write).
The per-batch resolvers carry [(unit, exn) result] so an fsync
failure in the drainer propagates to every joiner in the same batch
instead of being silently dropped by a unit-broadcast (#151). *)typecommit_queue={mutabledrainer:bool;mutablepending:int;mutablewaiters:(unit,exn)resultLwt.ulist}letcreate_commit_queue()={drainer=false;pending=0;waiters=[]}typebt_state={close_fn:unit->unitLwt.t;pager:Pager.t;cipher:Crypto.toption(** #84: the page cipher when the DB is encrypted, else [None]. Mirrors
the cipher captured by the read/write callback closures; retained here
so [copy_to]/[rekey_to] can re-encrypt the snapshot page image. *);mutablemeta:Btree.t;trees:(tree_id,Btree.t)Hashtbl.t;tree_tags:(tree_id,int32)Hashtbl.t(** #174: per-tree page-header stamp (low 32 bits of the schema
fingerprint), set by the catalog via {!set_tree_tag}. Pages written
for a tree carry its tag in the reserved header bytes; untagged trees
(default) carry 0. *);mutablecurrent_header:Header.t;schema_version:int64;mutabletxn_freelist_snapshot:Freelist.toption;(* Snapshot of freelist taken at rw_begin; restored on rollback. None when no RW txn is active. *)active_readers:(int64,int)Hashtbl.t;(* Maps snap_txn_id -> reference count of active RO txns at that snapshot *)mutablebt_savepoints:bt_savepointlist;(* Stack of named savepoints; newest at front. Cleared on commit/rollback. *)bt_append:(tree_id,Btree.append_cursor)Hashtbl.t;(* #356: per-tree append cursor for O(1) bulk sequential inserts. Set by
[put_x] after an append, consumed by the next append. Invalidated on
commit/rollback/savepoint-rollback and on any non-append mutation of the
tree. Re-validated against the live page on every use, so a stale entry
can only force the slow path, never corrupt the tree. *)wal:Granary_storage.Wal.toption;(* When set, commits append to this WAL instead of writing to the main
DB; reads route through it via the Pager hook. *)wal_close:(unit->unitLwt.t)option;mutablewal_autocheckpoint_threshold:int;(* When > 0 and committed WAL frames reach this number, the next
commit triggers an inline checkpoint (still under [lock]) so
the WAL stays bounded. 0 disables auto-checkpoint. Per-connection,
not persisted. *)commit_queue:commit_queue;(* WAL-mode group commit (#77). Used only when [wal] is [Some];
allocated unconditionally to keep [bt_state] uniform. *)active_reader_frames:(int,int)Hashtbl.t;(* WAL committed_frames snapshot value -> refcount of RO snapshots
captured at that value. Lets [min_active_ro_reader_frames] compute
the lowest snapshot bound currently in flight in O(distinct
snapshots) which is bounded by the number of concurrent readers. *)reader_done_cond:unitLwt_condition.t;(* Broadcast on every [ro_end] so a waiting checkpoint can re-check
[min_active_ro_reader_frames] without busy-waiting. *)mutableautockpt_in_flight:bool(* True iff a background autocheckpoint fiber is currently running.
Used to coalesce: if a commit crosses the threshold while a
checkpoint is already running, we skip rescheduling. *);mutablereplication_shipped_frames:int(* WAL frame index up to which the replication consumer (if any) has
acknowledged shipment. Initialized to [max_int] so that when no
consumer is active it does not gate checkpoint truncation. When a
consumer registers it sets this to its shipped position; [checkpoint]
then waits (subject to [replication_gate_max_yields]) for this to
reach [committed_frames] via [wait_for_readers_past]. *);mutablereplication_gate_max_yields:int(* Bounded-yield "timeout" for the checkpoint gate's wait on the
replication floor (#207). When the floor (a standby's acked
position, plumbed in by the app via [update_replication_position])
is below the checkpoint target, the gate yields up to this many
times before proceeding anyway — a dead or slow standby must not
wedge the master's WAL forever. Pure-Mirage has no ambient clock,
so the "timeout" is a bounded count of cooperative [Lwt.pause]
yields (the project's [wait_for] idiom), not wall-clock time.
[max_int] (the default) means unbounded: wait indefinitely on the
broadcast condition, exactly as before this knob existed. Local
RO readers are NEVER abandoned by this budget — only the
replication floor. On timeout the standby falls outside the live
un-checkpointed window and must re-base (see #208). *);mutablebackup_shipped_frames:int(* WAL frame index up to which the backup consumer (if any) has
captured frames. Analogous to [replication_shipped_frames] but for
incremental backup (#265). Initialized to [max_int] so that when no
backup consumer is active it does not gate checkpoint truncation.
The backup consumer calls {!update_backup_position} to advance this
as frames are captured and stored. *);mutablebackup_gate_max_yields:int(* Bounded-yield "timeout" for the checkpoint gate's wait on the
backup floor (#265). Same semantics as
[replication_gate_max_yields]: when the backup consumer has not yet
captured frames up to the checkpoint target, the gate yields up to
this many times before proceeding anyway. [max_int] (the default)
means unbounded — wait indefinitely. A finite budget bounds the
wait: once spent, the checkpoint proceeds and un-captured frames are
recycled (the backup must re-base). Negative inputs clamp to [0]. *);mutableon_committed_frames:(epoch:int64->base_idx:int->count:int->unitLwt.t)option(* Optional callback invoked asynchronously after each WAL commit batch.
Receives ~epoch, ~base_idx (starting WAL frame index of the batch),
~count (number of frames in the batch). The application reads the
individual frames via [Wal.read_frame] and ships them to the object
store. Fired via [Lwt.async] so it never blocks the commit path.
[None] when no sink is registered. *);mutableon_event:(Store_event.t->unit)option(* #382: optional, synchronous, fire-and-forget observer for internal
events (the internals monitor). [None] = zero overhead. Invoked via
[emit_event], which swallows any exception so a faulty observer can
never break a transaction. Btree backend only — Mem has no bt_state. *);history:History.sinkoption(* #266: append-only as-of commit log. [Some] iff the store was opened
with [~as_of_history:true] AND a sink was supplied; [None] disables
the whole feature (zero commit-path overhead). *);history_now:unit->int64(* #266: wall-clock (ms since epoch) stamped onto each commit-log record.
Injected at open; defaults to a constant 0 when history is disabled. *);mutablehistory_floor:int64option(* #266: retention floor. When [Some t], [min_safe] is capped at [t+1]
so pages reachable from roots >= t are never reused. *);mutablecurrent_tree:tree_idoption(* #385: the tree id of the in-flight read/write/cursor op, set at the
bt_get_tree(_ro) chokepoint and stamped onto page events by
translate_pager_event. Best-effort (single mutable shared across fibers),
same spirit as the pager's txn_id. [None] => stamp tree = -1. *);mutablefollower:bool(* When true, [rw_begin] rejects with an error. Set by the standby
consumer while following the master's WAL stream; cleared on
promotion or when the follower loop exits. The in-memory backend
ignores this flag (Mem stores have no standby semantics). *);mutablefollower_ack_position:intoption(* [Wal.committed_frames] at the time the last apply batch completed on
this follower, or [None] when no position has been recorded yet.
When [Some n], [ro_begin] caps the RO snapshot's visible WAL frames
to [min committed_frames n] so readers never observe frames past the
follower's last-applied commit (#263). Stored in local committed-frame
count space so it compares correctly against [Wal.committed_frames]
and survives local epoch resets. *);mutablesync_mode:[`Full|`Batched|`Off](* #298: durability mode. [`Full] = fsync every group-commit (default).
[`Batched] = defer fsync until [batch_commits] or [batch_interval_ms].
[`Off] = never fsync on commit. Only consulted in WAL mode. *);mutablebatch_commits:int(* #298: batched N threshold (default 256) *);mutablebatch_interval_ms:int(* #298: batched T threshold ms (default 100) *);mutableunsynced_commits:int(* #298: committed-but-unsynced batches since last fsync. *);mutablelast_sync_time:float(* #298: clock () at last commit fsync; for the T trigger. *);mutableclock:unit->float(* #298: wall-clock source; default returns 0. *);mutablesink_shipped_frames:int(* #298/#1: per-epoch count of WAL frames already shipped to the
replication sink ([on_committed_frames]). The sink is fired ONLY for
frames that have been fsynced, so a standby can never lead a
crash-recovered master. Reset to 0 on checkpoint (new epoch). *);mutableclosing:bool(* #338: set by [close] to signal teardown. [maybe_autockpt_after_commit]
then dispatches no fresh checkpoint, and an in-flight/parked one unwinds
without touching the pager/WAL fds ([checkpoint_unlocked] and
[wait_for_readers_past] bail on it). Lets [close] drain checkpoints
without acquiring [t.lock] — which an abandoned write txn holds until
commit/rollback, so taking it would hang close. *);mutableckpt_io_in_flight:int(* #338 (review r2): count of checkpoints that have passed the gate and are
actively performing pager/WAL fd I/O ([checkpoint_unlocked], both the auto
and manual paths). Incremented AFTER lock acquisition + the [closing]
check, so a checkpoint merely parked on [acquire_write] (e.g. behind an
abandoned write txn) is NOT counted — [close] therefore never waits on it
(it aborts on [closing] if it ever acquires the lock). [close] drains
this to 0 (together with [sink_ships_in_flight]) before fd teardown.
Distinct from [autockpt_in_flight], which is dispatch-intent (coalescing)
only. *);mutablesink_ships_in_flight:int(* #337: count of async sink ships dispatched but not yet completed. The
ship callback reads WAL frame payloads LAZILY ([Wal.read_frame]); a
checkpoint's [Wal.reset] would recycle/zero those frames and bump the
epoch out from under an in-flight reader ([Corrupt_frame] / stale
epoch). Incremented synchronously at dispatch (before the [Lwt.async]);
decremented in the callback's finalize. [checkpoint_unlocked] waits for
this to reach 0 before [Wal.reset]. *)}letdefault_wal_autocheckpoint_threshold=1000letdefault_batch_commits=256letdefault_batch_interval_ms=100(** #298: per-deployment durability mode. *)typedurability=|Full|Batchedof{commits:int;interval_ms:int}|Offtypebackend=|Memof(tree_id,Bytes.tBytes_map.tref)Hashtbl.t|Btreeofbt_statetypet={backend:backend;lock:Rwlock.t;(* Shadow copies of Mem backend tree contents for the active RW txn.
Writes during the txn go to the shadow — the live tree is NEVER
modified until commit. This prevents readers (both RO snapshots
and subsequent RW txns) from ever seeing uncommitted state.
None when no RW transaction is active.
#178: without shadow writes, concurrent RO reads could observe
uncommitted mutations because Rwlock's acquire_read never blocks. *)mutablemem_rw_shadow:(tree_id*Bytes.tBytes_map.t)listoption;(* Savepoint stack for the Mem backend; newest entry at front.
Each entry is (savepoint_name, snapshot_of_shadow). *)mutablemem_savepoints:(string*(tree_id*Bytes.tBytes_map.t)list)list}letppfmtt=Format.fprintffmt"Store.t { backend = %s }"(matcht.backendwith|Mem_->"Mem"|Btree_->"Btree");;(* #95/#176: the page geometry this store is backed by. The in-memory backend
has no on-disk geometry, so it reports {!Geometry.default}; VACUUM uses this
to rebuild the temp file at the source's page_size/reserved. *)letgeometryt=matcht.backendwith|Mem_->Geometry.default|Btreest->Pager.geomst.pager;;typero_snapshot={rs_store:t;rs_snap_txn_id:int64;rs_snap_meta_root:int64;rs_snap_trees:(tree_id,Btree.t)Hashtbl.t;rs_snap_frames:int;(* WAL committed_frames at ro_begin; 0 when no WAL is in effect. *)rs_pinned:(int64,unit)Hashtbl.t(* Page ids this snapshot has pinned in the Pager cache (#159). Every
snapshot read records the pages it materialises here; [ro_end]
releases them via [Pager.unpin_all]. Unused for the Mem backend. *);rs_mem_snap:(tree_id*Bytes.tBytes_map.t)listoption(** #178: for the in-memory backend, a deep copy of every tree's
contents taken at [ro_begin] so RO reads never see uncommitted
writes from a concurrent (but rollback-destined) writer.
[None] for Btree backend. *)}type'atxn=|Ro:ro_snapshot->rotxn|Rw:t->rwtxntypeseek_result=|Foundofbytes|Not_foundof[`Greaterofbytes|`End](* Cursor over either backend.
For the in-memory backend, the cursor holds an immutable snapshot of
the bindings as a list (matches Phase 0 semantics).
For the B+-tree backend we similarly materialise a snapshot (list of
(k,v) pairs) at cursor_open time. This is acceptable for Phase 1 and
makes seek/next semantics identical to the in-memory implementation;
true streaming cursors come later.
In both cases [ready] and [remaining] together implement the
"pre-positioned" semantics from [store.mli]: the first [cursor_next]
after positioning returns the positioned entry without advancing. *)typecursor={all:(bytes*bytes)list;mutableremaining:(bytes*bytes)list;mutableready:bool}(* ------------------------------------------------------------------ *)(* Backend helpers — Mem *)(* ------------------------------------------------------------------ *)letmem_treetreestid=matchHashtbl.find_opttreestidwith|Somer->r|None->letr=refBytes_map.emptyinHashtbl.addtreestidr;r;;(* #178: look up a tree in the snapshot taken at [ro_begin] for the
in-memory backend. Returns [Bytes_map.empty] when the tree didn't
exist at snapshot time — an RO reader should see an empty tree, not
the live (possibly uncommitted) contents. *)letmem_tree_snap(snap:(tree_id*Bytes.tBytes_map.t)list)(tid:tree_id)=matchList.assoc_opttidsnapwith|Somemap->map|None->Bytes_map.empty;;(* Shadow helpers for the in-memory backend (#178).
During a RW transaction, all writes go to a per-txn shadow.
The live tree is never mutated until commit, so RO txn
snapshots always capture committed-only state. *)(* Get a tree's content from the shadow, falling back to the live
tree when the tree hasn't been touched by this txn yet. *)letshadow_get(shadow:(tree_id*Bytes.tBytes_map.t)list)(trees:(tree_id,Bytes.tBytes_map.tref)Hashtbl.t)(tid:tree_id)=matchList.assoc_opttidshadowwith|Somemap->map|None->!(mem_treetreestid);;(* Update a tree in the shadow. The tree is lazy-copied from the live
tree on first access (via [shadow_get]). *)letshadow_update(shadow:(tree_id*Bytes.tBytes_map.t)list)(trees:(tree_id,Bytes.tBytes_map.tref)Hashtbl.t)(tid:tree_id)(f:Bytes.tBytes_map.t->Bytes.tBytes_map.t)=letmap=shadow_getshadowtreestidin(tid,fmap)::List.remove_assoctidshadow;;(* ------------------------------------------------------------------ *)(* Backend helpers — Btree *)(* ------------------------------------------------------------------ *)(* tree_id <-> bytes encoding via varint (zigzag, since negative ids are
reserved for internal use; we don't actually persist negative ids but
using signed encoding lets us round-trip safely). *)letencode_tree_id(tid:tree_id):bytes=letbuf=Buffer.create8inVarint.encode_int64buf(Int64.of_inttid);Buffer.to_bytesbuf;;letencode_root_page(pid:int64):bytes=letbuf=Buffer.create8inVarint.encode_uint64bufpid;Buffer.to_bytesbuf;;letdecode_root_page(b:bytes):int64=letv,_=Varint.decode_uint64b0inv;;letmap_btree_err:Btree.error->error=function|Btree.Pager_error(Pager.Block_errors)->Block_errors|Btree.Pager_error(Pager.Corruptions)->Corruptions|Btree.Key_too_largen->Key_too_largen|Btree.Value_too_largen->Value_too_largen|Btree.Tree_corrupts->Corruptions;;(* The B+-tree treats Bytes by [Bytes.compare]; cursor_seek consumes the
raw bytes; everything is byte-clean. *)(* Lookup-or-build the Btree handle for a tree_id. Looks up the
tree_id's root page in the meta-tree; if absent (new tree), creates a
fresh empty Btree (root_page = 0L). *)letbt_get_treest(tid:tree_id):(Btree.t,error)resultLwt.t=(* #385/#174: meta-tree pages read while resolving the root are system pages —
keep [current_tree] clear (stamps tree = -1) during the lookup, then stamp
[tid] so the caller's subsequent data-page reads are attributed to it. *)st.current_tree<-None;let*r=matchHashtbl.find_optst.treestidwith|Somebt->Lwt.return_okbt|None->letkey=encode_tree_idtidinlet*r=Btree.getst.metakeyin(matchrwith|Errore->Lwt.return_error(map_btree_erre)|OkNone->letbt=Btree.createst.pager~root_page:0LinHashtbl.replacest.treestidbt;Lwt.return_okbt|Ok(Somev)->letroot_page=decode_root_pagevinletbt=Btree.createst.pager~root_pageinHashtbl.replacest.treestidbt;Lwt.return_okbt)inst.current_tree<-Sometid;Lwt.returnr;;(* #174: the page-header stamp for [tid] (0 when untagged). *)lettree_tagst(tid:tree_id):int32=Option.value~default:0l(Hashtbl.find_optst.tree_tagstid);;(* Convert a result with [error] payload to an Lwt-failing version. The
public [get/put/del/cursor_open] signatures don't return [result], so
B+-tree errors are surfaced as Lwt exceptions. *)letunwrap_errorr=matchrwith|Okv->Lwt.returnv|Errore->Lwt.fail_with(Format.asprintf"Store: %a"pp_errore);;letmin_active_reader_txnst=Hashtbl.fold(funtxn_id_acc->matchaccwith|None->Sometxn_id|Somem->Some(Int64.minmtxn_id))st.active_readersNone;;(* Lowest WAL frame index pinned by an in-flight RO snapshot, ignoring the
replication floor. A checkpoint must NEVER recycle past this (the
snapshot would observe a broken WAL), so this gate is honored
unconditionally — unlike the replication floor, which the #207 timeout
may abandon. *)letmin_active_ro_reader_frames(st:bt_state):intoption=Hashtbl.fold(funk_acc->matchaccwith|None->Somek|Somem->Some(minmk))st.active_reader_framesNone;;(* True iff an in-flight RO snapshot still needs WAL frames below [target].
This gate is honored unconditionally by the checkpoint wait. *)letro_readers_below(st:bt_state)~target=matchmin_active_ro_reader_framesstwith|Somem->m<target|None->false;;(* True iff a replication consumer is active and its acked floor is below
[target]. This gate is subject to the #207 bounded-yield timeout. *)letreplication_floor_below(st:bt_state)~target=st.replication_shipped_frames<>max_int&&st.replication_shipped_frames<target;;(* True iff a backup consumer is active and its captured floor is below
[target]. Analogous to [replication_floor_below] but for the
incremental backup watermark (#265). Subject to a bounded-yield
timeout like the replication floor. *)letbackup_floor_below(st:bt_state)~target=st.backup_shipped_frames<>max_int&&st.backup_shipped_frames<target;;(* Lookup-or-build the Btree handle for a tree_id using a snapshot's
pinned meta root page rather than the live meta tree. *)letbt_get_tree_ro(snap:ro_snapshot)(st:bt_state)(tid:tree_id):(Btree.t,error)resultLwt.t=(* #385/#174: see bt_get_tree — meta reads stay unattributed (tree = -1);
[tid] is stamped only after the handle is resolved. *)st.current_tree<-None;let*r=matchHashtbl.find_optsnap.rs_snap_treestidwith|Somebt->Lwt.return_okbt|None->letsnap_frames=ifsnap.rs_snap_frames=0thenNoneelseSomesnap.rs_snap_framesinletsnap_meta=Btree.create?snapshot_frames:snap_frames~pin_set:snap.rs_pinnedst.pager~root_page:snap.rs_snap_meta_rootinletkey=encode_tree_idtidinlet*r=Btree.getsnap_metakeyin(matchrwith|Errore->Lwt.return_error(map_btree_erre)|OkNone->letbt=Btree.create?snapshot_frames:snap_frames~pin_set:snap.rs_pinnedst.pager~root_page:0LinHashtbl.replacesnap.rs_snap_treestidbt;Lwt.return_okbt|Ok(Somev)->letroot_page=decode_root_pagevinletbt=Btree.create?snapshot_frames:snap_frames~pin_set:snap.rs_pinnedst.pager~root_pageinHashtbl.replacesnap.rs_snap_treestidbt;Lwt.return_okbt)inst.current_tree<-Sometid;Lwt.returnr;;(* ------------------------------------------------------------------ *)(* Freelist page I/O helpers (forward-declared here; used by open_block *)(* and commit below) *)(* ------------------------------------------------------------------ *)(* Walk the freelist page chain starting at [first_page], collect all
entries, and return a reconstructed [Freelist.t]. *)letread_freelist_pagespager~first_page:Freelist.tLwt.t=ifInt64.equalfirst_page0LthenLwt.returnFreelist.emptyelse(letreclooppidacc=ifInt64.equalpid0LthenLwt.return(Freelist.of_list(List.revacc))elselet*r=Pager.readpagerpidinmatchrwith|Error_->Lwt.return(Freelist.of_list(List.revacc))|Okbuf->letcommon=Page.read_commonbufinletn=mincommon.Page.n_keys(Pager.max_freelist_entries_per_pagepager)inletnext_pid=Int64.logand0xFFFFFFFFL(Int64.of_int32common.Page.right_page)inletentries=List.initn(funi->lete=Page.freelist_entry_atbuf~index:iine.Page.page_id,e.Page.freed_at_txn_id)inloopnext_pid(List.rev_appendentriesacc)inloopfirst_page[]);;(* ------------------------------------------------------------------ *)(* create / open_block / close *)(* ------------------------------------------------------------------ *)letcreate():t={backend=Mem(Hashtbl.create16);lock=Rwlock.create();mem_rw_shadow=None;mem_savepoints=[]};;letmap_header_err(e:Header.error):error=matchewith|Header.Ios->Header_errors|Header.Both_headers_corrupt->Header_error"both header pages corrupt"|Header.Unsupported_formatv->Header_error(Printf.sprintf"unsupported on-disk format_version %ld"v);;(* Build a fully-initialised [t] wrapping a B-tree-backed [bt_state] from the
given pager/meta/header. [wal]/[wal_close] default to None (plain opens);
WAL opens pass [Some _]. *)letmake_btree_store?(wal=None)?(wal_close=None)?(cipher=None)?(history=None)?(history_now=fun()->0L)~close_fn~pager~meta~(h:Header.t)()=letst={close_fn;pager;cipher;meta;trees=Hashtbl.create16;current_tree=None;tree_tags=Hashtbl.create16;current_header=h;schema_version=h.schema_version;txn_freelist_snapshot=None;active_readers=Hashtbl.create4;bt_savepoints=[];bt_append=Hashtbl.create8;wal;wal_close;wal_autocheckpoint_threshold=default_wal_autocheckpoint_threshold;commit_queue=create_commit_queue();active_reader_frames=Hashtbl.create4;reader_done_cond=Lwt_condition.create();autockpt_in_flight=false;replication_shipped_frames=max_int;replication_gate_max_yields=max_int;backup_shipped_frames=max_int;backup_gate_max_yields=max_int;on_committed_frames=None;on_event=None;history;history_now;history_floor=None;follower=false;follower_ack_position=None;sync_mode=`Full;batch_commits=default_batch_commits;batch_interval_ms=default_batch_interval_ms;unsynced_commits=0;last_sync_time=0.;clock=(fun()->0.);sink_shipped_frames=0;sink_ships_in_flight=0;closing=false;ckpt_io_in_flight=0}in{backend=Btreest;lock=Rwlock.create();mem_rw_shadow=None;mem_savepoints=[]};;(* #338 (review r2): event-driven wait until [pred] holds, parking on
[reader_done_cond] (broadcast whenever an in-flight counter changes). Shared
by [close] (drain checkpoint fd-I/O + sink ships) and [checkpoint_unlocked]
(drain sink ships before [Wal.reset]). Cooperative Lwt: the pred check and
the [Lwt_condition.wait] register with no yield between, so no wakeup is
lost. *)letrecwait_until(st:bt_state)(pred:unit->bool):unitLwt.t=ifpred()thenLwt.return_unitelselet*()=Lwt_condition.waitst.reader_done_condinwait_untilstpred;;letclose(t:t):unitLwt.t=matcht.backendwith|Mem_->Lwt.return_unit|Btreest->(* #338: an async checkpoint or sink ship touches the pager/WAL fds; tearing
them down underneath one corrupts it (a checkpoint's error is swallowed
and relies on WAL-replay self-healing; a ship loses tail frames the
standby then misses). Rather than take [t.lock] for teardown — which an
abandoned write txn holds until commit/rollback, so [close] would hang on
it (review r2 #3) — signal teardown via [st.closing]:
- [maybe_autockpt_after_commit] dispatches no fresh checkpoint once set;
- a checkpoint parked on the gate or [acquire_write] unwinds without fd
I/O ([wait_for_readers_past]/[checkpoint_unlocked] bail on [closing]);
- the broadcast wakes a checkpoint parked on the replication floor.
We then drain — event-driven — the work that is ACTUALLY mid-fd-I/O:
checkpoints past the gate ([ckpt_io_in_flight], covers the auto AND manual
paths — review r2 #3) and async sink ships ([sink_ships_in_flight], whose
lazy [Wal.read_frame] would hit a closed fd — review r2 #2). A checkpoint
merely parked on [acquire_write] is invisible here (it never incremented),
so an abandoned txn cannot wedge close (review r2 #1). Callers must still
quiesce their own writers before [close] (see store.mli). *)st.closing<-true;Lwt_condition.broadcastst.reader_done_cond();let*()=wait_untilst(fun()->st.ckpt_io_in_flight=0&&st.sink_ships_in_flight=0)in(* #298: in batched/off mode the last acked commits may never have been
fsynced. Decide on the WAL's actual committed-frame state rather than the
in-memory unsynced counter. A redundant fsync here (frames already
durable) is cheap and safe; skipping a needed one is not. Full mode syncs
every commit. *)letneeds_final_sync=st.sync_mode<>`Full&&matchst.walwith|Somew->Granary_storage.Wal.committed_framesw>0|None->falsein(* #298/#2: a failed final fsync still releases the fds (wal_close/close_fn)
but THEN raises — close is a durability anchor, so silently reporting
success on EIO/ENOSPC is wrong (matches the commit/checkpoint convention
of surfacing sync errors). *)let*sync_err=ifneeds_final_syncthenlet*r=Pager.wal_syncst.pagerinmatchrwith|Ok()->st.unsynced_commits<-0;Lwt.return_none|Errore->Lwt.return_someeelseLwt.return_noneinlet*()=matchst.wal_closewith|None->Lwt.return_unit|Somef->f()inlet*()=st.close_fn()in(matchsync_errwith|None->Lwt.return_unit|Somee->Lwt.fail_with(Format.asprintf"Store.close: final wal_sync: %a"Pager.pp_errore));;(* #95: discover the file's geometry by reading page 0's leading bytes through
the raw block callback (page 0 is always at offset 0, so this works whatever
the backend's addressing page size, and it bypasses the pager cache). Falls
back to [fallback] for a fresh/empty/zeroed device, which a subsequent
[Header.init] then stamps. *)letpeek_geometry~read_page~fallback=letbuf=Cstruct.createGeometry.default.page_sizeinlet%lwtr=read_page~page_id:0Lbufinmatchrwith|Ok()->Lwt.return(Option.value(Header.peek_geometrybuf)~default:fallback)|Error_->Lwt.returnfallback;;(* ------------------------------------------------------------------ *)(* Opt-in page encryption (#84). *)(* *)(* The pager and B+-tree only ever see PLAINTEXT. A supplied key *)(* builds a cipher; we then (a) force the fresh-creation geometry to *)(* carve [Crypto.overhead] reserved bytes off each page's tail, (b) *)(* wrap the raw read/write callbacks so pages >= 2 are *)(* decrypted/encrypted (pages 0,1 are the headers and pass through *)(* plaintext), and (c) stamp / verify a key-check canary in the header. *)(* ------------------------------------------------------------------ *)letbuild_cipher=function|None->OkNone|Somek->(matchCrypto.create~key:kwith|Okc->Ok(Somec)|Error`Bad_key_length->Error(Block_error"encryption key must be 32 bytes"));;(* When a key is in play we draw a fresh nonce on every encrypted write (and one
for the header canary at creation). [Mirage_crypto_rng.generate] raises if
the application never seeded the RNG ([lib/] is Mirage-clean and never seeds):
[No_default_generator] when no generator was installed at all (the common
forgot-to-seed case) and [Unseeded_generator] when one was installed but not
seeded. Either would otherwise surface as a raw exception on the first write
rather than a [Store.error]. Probe once at open time so the foot-gun is
caught at the entry point the caller controls; the RNG is process-global, so
a seed present here is present for later writes. *)letensure_rng_seeded=function|None->Ok()|Some_->(tryignore(Mirage_crypto_rng.generate1:string);Ok()with|Mirage_crypto_rng.Unseeded_generator|Mirage_crypto_rng.No_default_generator->ErrorEncryption_rng_unseeded);;(* Force a fresh-creation geometry to carry the crypto overhead in its
reserved tail. If bumping [reserved] to [Crypto.overhead] is rejected we
surface a clean error rather than silently proceeding with a geometry whose
reserved tail is too small — encryption would then write nonce+tag into bytes
the B+-tree believes are usable, corrupting the page. (Unreachable with the
4096-multiple page sizes [Geometry.create] permits, since they always leave
>= 480 payload after reserving 32 bytes; kept as defense-in-depth.) *)letgeom_for_ciphercipher(g:Geometry.t)=matchcipherwith|None->Okg|Some_->ifg.reserved_bytes_per_page>=Crypto.overheadthenOkgelse(matchGeometry.create~page_size:g.page_size~reserved_bytes_per_page:(maxg.reserved_bytes_per_pageCrypto.overhead)with|Okg'->Okg'|Errore->Error(Block_error(Format.asprintf"encryption needs %d reserved bytes/page, but the geometry rejects it: %a"Crypto.overheadGeometry.pp_errore)));;letwrap_callbackscipher~read_page~write_page=matchcipherwith|None->read_page,write_page|Somec->letrd~page_idbuf=let*r=read_page~page_idbufinmatchrwith|Error_ase->Lwt.returne|Ok()->ifInt64.comparepage_id2L<0thenLwt.return_ok()else(matchCrypto.decrypt_pagec~page_idbufwith|Ok()->Lwt.return_ok()|Error`Tag_mismatch->Lwt.return_error"decrypt: tag mismatch")inletwr~page_idbuf=ifInt64.comparepage_id2L<0thenwrite_page~page_idbufelse(lettmp=Cstruct.create(Cstruct.lengthbuf)inCstruct.blitbuf0tmp0(Cstruct.lengthbuf);Crypto.encrypt_pagec~page_idtmp;write_page~page_idtmp)inrd,wr;;letmake_enc_info=function|None->None|Somec->letnonce=Mirage_crypto_rng.generateCrypto.nonce_leninlettag=Crypto.make_canaryc~nonceinSome{Header.canary_nonce=nonce;canary_tag=tag};;letcheck_key(h:Header.t)cipher=matchh.Header.enc,cipherwith|None,None->Ok()|Some_,None->ErrorEncryption_key_required|None,Some_->ErrorNot_encrypted|Somee,Somec->ifCrypto.check_canaryc~nonce:e.Header.canary_nonce~tag:e.Header.canary_tagthenOk()elseErrorEncryption_key_mismatch;;letopen_block?(as_of_history=false)?(history:History.sinkoption)?(now:(unit->int64)option)?(key:stringoption)?(geom=Geometry.default)~(init_if_corrupt:bool)~(read_page:page_id:int64->Cstruct.t->(unit,string)resultLwt.t)~(write_page:page_id:int64->Cstruct.t->(unit,string)resultLwt.t)~(sync:unit->(unit,string)resultLwt.t)~(resize:n_pages:int64->(unit,string)resultLwt.t)~(n_pages:int64)~(close:unit->unitLwt.t)():(t,error)resultLwt.t=(* #266: guard the misconfig FIRST, before opening any device/fds, so a bad
request can never leak resources. *)ifas_of_history&&Option.is_nonehistorythenLwt.return_errorHistory_misconfiguredelse(lethistory=ifas_of_historythenhistoryelseNoneinlethistory_now=matchnowwith|Somef->f|None->fun()->0Linmatchlet(let*)=Result.bindinlet*cipher=build_cipherkeyinlet*()=ensure_rng_seededcipherinlet*geom=geom_for_cipherciphergeominOk(cipher,geom)with|Errore->Lwt.return_errore|Ok(cipher,geom)->letread_page,write_page=wrap_callbackscipher~read_page~write_pageinletpager=Pager.create~read_page~write_page~sync~resize~n_pages~freelist:Freelist.emptyin(* Adopt the file's real geometry (peeked for an existing file, [geom] for a
fresh one) before any header read so buffers are sized correctly (#95). *)let%lwteff_geom=peek_geometry~read_page~fallback:geominPager.set_geompagereff_geom;let%lwthr=Header.read_livepagerin(matchhrwith|ErrorHeader.Both_headers_corruptwheninit_if_corrupt->(* Fresh device — initialise headers. Disabled via [~init_if_corrupt:false]
so an existing-but-corrupt device surfaces [Header_error] instead of
being silently re-initialised (a Unix-file open must not clobber). *)let%lwtir=Header.init~enc:(make_enc_infocipher)pagerin(matchirwith|Errore->Lwt.return_error(map_header_erre)|Ok()->Pager.set_n_pagespager2L;let%lwthr2=Header.read_livepagerin(matchhr2with|Errore->Lwt.return_error(map_header_erre)|Okh->letmeta=Btree.createpager~root_page:0LinLwt.return_ok(make_btree_store~cipher~history~history_now~close_fn:close~pager~meta~h())))|Errore->Lwt.return_error(map_header_erre)|Okh->(matchcheck_keyhcipherwith|Errore->Lwt.return_errore|Ok()->Pager.set_n_pagespagerh.n_pages_total;let%lwtfl=read_freelist_pagespager~first_page:h.freelist_pageinPager.set_freelistpagerfl;letmeta=Btree.createpager~root_page:h.root_pageinLwt.return_ok(make_btree_store~cipher~history~history_now~close_fn:close~pager~meta~h()))));;(* ------------------------------------------------------------------ *)(* WAL-mode opens *)(* ------------------------------------------------------------------ *)moduleWal=Granary_storage.Walletinstall_wal_hook(pager:Pager.t)(wal:Wal.t)=letcb:Pager.wal_callbacks={wal_find_page=(funpid->Wal.find_pagewalpid);wal_find_page_at=(funpid~max_frame->Wal.find_page_atwalpid~max_frame);wal_read_frame=(funidx->let*r=Wal.read_framewalidxinmatchrwith|Okpage->Lwt.return_okpage|Errore->Lwt.return_error(Format.asprintf"%a"Wal.pp_errore));wal_append_commit=(funpages->let*r=Wal.append_commitwalpagesinmatchrwith|Ok()->Lwt.return_ok()|Errore->Lwt.return_error(Format.asprintf"%a"Wal.pp_errore));wal_append_commit_no_sync=(funpages->let*r=Wal.append_commit_no_syncwalpagesinmatchrwith|Ok()->Lwt.return_ok()|Errore->Lwt.return_error(Format.asprintf"%a"Wal.pp_errore));wal_sync=(fun()->let*r=Wal.flush_syncwalinmatchrwith|Ok()->Lwt.return_ok()|Errore->Lwt.return_error(Format.asprintf"%a"Wal.pp_errore))}inPager.set_walpager(Somecb);;(* After the WAL hook is installed, re-read the (now WAL-aware) header,
reconcile [n_pages] for a freshly-initialised DB, load the freelist, and
build the WAL-backed store. *)letfinish_wal_open~cipher~history~history_now~close~wal_close~pager~wal~was_fresh=let%lwthr2=Header.read_livepagerinmatchhr2with|Errore->Lwt.return_error(map_header_erre)|Okh->(matchcheck_keyhcipherwith|Errore->Lwt.return_errore|Ok()->(* If the header n_pages_total is below the pager's current allocation,
prefer the pager's value (freshly-init'd headers carry
n_pages_total = 0). *)letchosen_n_pages=ifwas_freshthenInt64.maxh.n_pages_total(Pager.n_pagespager)elseh.n_pages_totalinPager.set_n_pagespagerchosen_n_pages;let%lwtfl=read_freelist_pagespager~first_page:h.freelist_pageinPager.set_freelistpagerfl;letmeta=Btree.createpager~root_page:h.root_pageinLwt.return_ok(make_btree_store~cipher~history~history_now~wal:(Somewal)~wal_close:(Somewal_close)~close_fn:close~pager~meta~h()));;letopen_block_wal?(as_of_history=false)?(history:History.sinkoption)?(now:(unit->int64)option)?(key:stringoption)?(geom=Geometry.default)~(read_page:page_id:int64->Cstruct.t->(unit,string)resultLwt.t)~(write_page:page_id:int64->Cstruct.t->(unit,string)resultLwt.t)~(sync:unit->(unit,string)resultLwt.t)~(resize:n_pages:int64->(unit,string)resultLwt.t)~(n_pages:int64)~(wal_read_at:offset:int64->Cstruct.t->(unit,string)resultLwt.t)~(wal_write_at:offset:int64->Cstruct.t->(unit,string)resultLwt.t)~(wal_sync:unit->(unit,string)resultLwt.t)~(wal_size_bytes:int64)~(close:unit->unitLwt.t)~(wal_close:unit->unitLwt.t)():(t,error)resultLwt.t=(* #266: guard the misconfig FIRST, before opening any device/fds, so a bad
request can never leak resources. *)ifas_of_history&&Option.is_nonehistorythenLwt.return_errorHistory_misconfiguredelse(lethistory=ifas_of_historythenhistoryelseNoneinlethistory_now=matchnowwith|Somef->f|None->fun()->0Linmatchlet(let*)=Result.bindinlet*cipher=build_cipherkeyinlet*()=ensure_rng_seededcipherinlet*geom=geom_for_cipherciphergeominOk(cipher,geom)with|Errore->Lwt.return_errore|Ok(cipher,geom)->letread_page,write_page=wrap_callbackscipher~read_page~write_pageinletpager=Pager.create~read_page~write_page~sync~resize~n_pages~freelist:Freelist.emptyin(* Adopt the file's real geometry before any header read or WAL open so the
main-DB buffers and the WAL frame size both match it (#95). *)let%lwteff_geom=peek_geometry~read_page~fallback:geominPager.set_geompagereff_geom;(* Step 1: read the main-DB header (or initialise if fresh). The WAL
hook is NOT installed yet, so writes go directly to the main DB.
[was_fresh] flag preserves the post-init n_pages override below. *)let%lwthr=Header.read_livepagerinlet%lwtinit_result=matchhrwith|ErrorHeader.Both_headers_corrupt->let%lwtir=Header.init~enc:(make_enc_infocipher)pagerin(matchirwith|Errore->Lwt.return_error(map_header_erre)|Ok()->Pager.set_n_pagespager2L;Lwt.return_oktrue)|Errore->Lwt.return_error(map_header_erre)|Ok_->Lwt.return_okfalsein(matchinit_resultwith|Errore->Lwt.return_errore|Okwas_fresh->(* Step 2: open the WAL and recover its index. *)let%lwtwr=Wal.open_~cipher~page_size:(Pager.page_sizepager)~read_at:wal_read_at~write_at:wal_write_at~sync:wal_sync~size_bytes:wal_size_bytes()in(matchwrwith|Errore->Lwt.return_error(Block_error(Format.asprintf"wal open: %a"Wal.pp_errore))|Okwal->(* Step 3: install the hook so subsequent reads consult the WAL. *)install_wal_hookpagerwal;(* Step 4: re-read the header (now WAL-aware) and build the store. *)finish_wal_open~cipher~history~history_now~close~wal_close~pager~wal~was_fresh)));;(* ------------------------------------------------------------------ *)(* Transactions *)(* ------------------------------------------------------------------ *)(* Build an RO snapshot against an explicit (txn_id, meta_root). [ro_begin]
passes the live header; [ro_begin_as_of] passes a retained historical root.
Registers the snapshot in [active_readers]/[active_reader_frames] so
reclamation respects it, exactly as the live path does. Uses the CURRENT
committed-frames horizon: CoW never rewrites a retained page-id, so each
retained page has exactly one WAL frame and "latest up to head" == the
historical content (the floor/reader pin prevents reuse). *)letro_begin_attst~snap_txn_id~snap_meta_root=letcommitted_frames=matchst.walwith|None->0|Somew->Wal.committed_frameswinletsnap_frames=ifst.followerthen(matchst.follower_ack_positionwith|Somen->mincommitted_framesn|None->committed_frames)elsecommitted_framesinletcount=Option.value~default:0(Hashtbl.find_optst.active_readerssnap_txn_id)inHashtbl.replacest.active_readerssnap_txn_id(count+1);letframe_count=Option.value~default:0(Hashtbl.find_optst.active_reader_framessnap_frames)inHashtbl.replacest.active_reader_framessnap_frames(frame_count+1);Ro{rs_store=t;rs_snap_txn_id=snap_txn_id;rs_snap_meta_root=snap_meta_root;rs_snap_trees=Hashtbl.create4;rs_snap_frames=snap_frames;rs_pinned=Hashtbl.create64;rs_mem_snap=None};;letro_begint=(* [Rwlock.acquire_read] is a counter bump, not an exclusion: under
snapshot isolation readers and writers don't conflict, so the call
never blocks regardless of writer state. It only matters for the
checkpoint coordinator that wants to know "are any RO snapshots
still in flight?" *)let*()=Rwlock.acquire_readt.lockinletis_closing=matcht.backendwith|Btreest->st.closing|Mem_->falseinifis_closingthen((* #338 (review r3): fail fast on a snapshot begun after [close] signalled
teardown — matches [rw_begin], avoiding an obscure pager EBADF later. *)Rwlock.release_readt.lock;Lwt.fail_with"Store.ro_begin: store is closing — read transactions are rejected")else(matcht.backendwith|Memtrees->(* #178: snapshot every tree so RO reads never observe uncommitted
writes from a concurrent writer that later rolls back. The
Btree backend gets snapshot isolation from the pager/WAL layer;
the mem backend must provide it here. *)letsnap=Hashtbl.fold(funtidracc->(tid,!r)::acc)trees[]inLwt.return(Ro{rs_store=t;rs_snap_txn_id=0L;rs_snap_meta_root=0L;rs_snap_trees=Hashtbl.create1;rs_snap_frames=0;rs_pinned=Hashtbl.create1;rs_mem_snap=Somesnap})|Btreest->Lwt.return(ro_begin_attst~snap_txn_id:st.current_header.txn_id~snap_meta_root:st.current_header.root_page));;(* #266: as-of retention API + time-travel read path. [History_error] carries
an {!error} out to the caller (the SQL layer maps it to a friendly message). *)exceptionHistory_erroroferrorletbt_oft=matcht.backendwith|Btrees->Somes|Mem_->None;;lethistory_pint~txn_id=matchbt_oftwith|Somest->st.history_floor<-Sometxn_id|None->();;lethistory_floort=matchbt_oftwith|Somest->st.history_floor|None->None;;lethistory_releaset=matchbt_oftwith|Somest->st.history_floor<-None|None->();;lethistory_logt=matchbt_oftwith|Some{history=Somesink;_}->sink.History.load()|_->Lwt.return[];;lethistory_enabledt=matchbt_oftwith|Some{history=Some_;_}->true|_->false;;letro_begin_as_oft(target:History.target)=(* #266 (review): resolve the historical target — which requires loading the
history log from the sink (real I/O for the Unix file sink: openfile/fstat/
read, any of which may reject with EIO/EMFILE/…) — BEFORE acquiring the read
lock. Holding the read lock across a rejecting [load] would leak it (the
coordinator would then never see the reader drain, stalling checkpoint and
[close]). Only the snapshot registration in [ro_begin_at] and the [closing]
recheck need the lock, exactly as [ro_begin] holds it during registration. *)matchbt_oftwith|None->Lwt.fail(History_errorHistory_unavailable)|Some{history=None;_}->Lwt.fail(History_errorHistory_unavailable)|Some({history=Somesink;_}asst)->let*records=sink.History.load()in(matchHistory.resolverecordstargetwith|None->Lwt.fail(History_errorHistory_pruned)|Somer->letpruned=matchst.history_floorwith|Somef->Int64.comparer.History.txn_idf<0|None->(* #266 (review): no floor ⇒ nothing is retained; refuse rather than
serve recycled pages. Without a pin, [rw_begin] leaves [min_safe]
uncapped, so the freelist recycles the resolved root's pages and
the snapshot would read garbage. *)trueinifprunedthenLwt.fail(History_errorHistory_pruned)elselet*()=Rwlock.acquire_readt.lockinifst.closingthen(Rwlock.release_readt.lock;Lwt.fail_with"Store.ro_begin_as_of: store is closing")elseLwt.return(ro_begin_attst~snap_txn_id:r.History.txn_id~snap_meta_root:r.History.root_page));;letemit_event(st:bt_state)(ev:Store_event.t)=matchst.on_eventwith|None->()|Somef->(tryfevwith|_->());;(* The id the currently-active rw txn will commit as. The header is not bumped
until commit, so every event of one txn shares this id (one writer at a time
under the write lock). *)letactive_txn_id(st:bt_state)=Int64.addst.current_header.txn_id1Lletrw_begint=let*()=Rwlock.acquire_writet.lockinletis_follower=matcht.backendwith|Btreest->st.follower|Mem_->falseinletis_closing=matcht.backendwith|Btreest->st.closing|Mem_->falseinifis_closingthen((* #338 (review r2 #4): fail fast on a write begun after [close] signalled
teardown, rather than letting the commit surface an obscure EBADF from a
torn-down fd. [close] does not take [t.lock], so a write can still race
in here; this is best-effort, paired with the quiesce-before-close
contract documented on [close]. *)Rwlock.release_writet.lock;Lwt.fail_with"Store.rw_begin: store is closing — write transactions are rejected")elseifis_followerthen(Rwlock.release_writet.lock;Lwt.fail_with"Store.rw_begin: store is in follower mode — write transactions are rejected while \
following")else((matcht.backendwith|Memtrees->letsnap=Hashtbl.fold(funtidracc->(tid,!r)::acc)trees[]int.mem_rw_shadow<-Somesnap;t.mem_savepoints<-[]|Btreest->letcurrent_rw_txn_id=Int64.addst.current_header.txn_id1Linemit_eventst(Store_event.Txn_begin{txn_id=current_rw_txn_id});Pager.set_txn_idst.pagercurrent_rw_txn_id;letmin_safe=matchmin_active_reader_txnstwith|None->current_rw_txn_id|Somem->Int64.mincurrent_rw_txn_idmin(* #266: cap [min_safe] at the retention floor so pages reachable from
roots >= the floor are never reused by the allocator. *)letmin_safe=matchst.history_floorwith|None->min_safe|Somef->Int64.minmin_safe(Int64.addf1L)inPager.set_alloc_min_safest.pagermin_safe;(* #297: same-txn page reuse happens via the txn_owned_pool (pages
allocated above n_pages_at_rw_begin), NOT the main freelist, so
alloc_min_safe is unchanged from the pre-#297 baseline. The
None branch (no readers) and Some m branch (reader exists) both
keep the original guard — committed-tree pages freed at
current_rw_txn_id are never eligible for same-txn reuse via the
main freelist regardless of reader state. The txn_owned_pool,
checked before the main freelist by Pager.alloc, provides
same-txn reuse independently of the freelist guard. *)Pager.set_n_pages_at_rw_beginst.pager(Pager.n_pagesst.pager);(* #297: defensive reset — any leftover from the previous txn is stale. *)Pager.txn_owned_pool_setst.pager[];st.txn_freelist_snapshot<-Some(Pager.freelistst.pager));Lwt.return(Rwt));;letro_end(Rosnap:rotxn)=(matchsnap.rs_store.backendwith|Mem_->()|Btreest->lettid=snap.rs_snap_txn_idin(matchHashtbl.find_optst.active_readerstidwith|None|Some1->Hashtbl.removest.active_readerstid|Somen->Hashtbl.replacest.active_readerstid(n-1));(matchHashtbl.find_optst.active_reader_framessnap.rs_snap_frameswith|None|Some1->Hashtbl.removest.active_reader_framessnap.rs_snap_frames|Somen->Hashtbl.replacest.active_reader_framessnap.rs_snap_frames(n-1));(* Release the pages this snapshot pinned (#159) so they become
evictable again. *)Pager.unpin_allst.pagersnap.rs_pinned;Lwt_condition.broadcastst.reader_done_cond());Rwlock.release_readsnap.rs_store.lock;Lwt.return_unit;;letwith_rotf=let*tx=ro_begintinLwt.finalize(fun()->ftx)(fun()->ro_endtx);;(* Free the previous freelist page chain back into the pager's in-memory
freelist (stamped with the current txn_id). *)letfree_old_freelist_pagespager~first_page=letreclooppid=ifInt64.equalpid0LthenLwt.return_unitelselet*r=Pager.readpagerpidinletnext_pid=matchrwith|Error_->0L|Okbuf->letc=Page.read_commonbufinInt64.logand0xFFFFFFFFL(Int64.of_int32c.Page.right_page)in(* Note: if read fails mid-chain, remaining pages beyond this point are
orphaned (leaked). This is acceptable only because a corrupt freelist
page implies a deeper storage invariant violation. *)Pager.freepager~page_id:pid~freed_at_txn_id:(Pager.get_txn_idpager);loopnext_pidinloopfirst_page;;(* Serialize the current pager freelist to a new page chain.
Returns the first page id (0L if the freelist is empty). *)(* Build and write a single freelist page holding [chunk] (possibly empty),
chaining to [next]. *)letwrite_one_freelist_pagepager~pid~next~chunk=letbuf=Cstruct.create(Pager.page_sizepager)inCstruct.memsetbuf0;Page.write_commonbuf{Page.kind=Page.Freelist;flags=0;n_keys=List.lengthchunk;right_page=Int64.to_int32next;crc32=0l};List.iteri(funj(page_id,freed_at_txn_id)->Page.freelist_set_entrybuf~index:j~page_id~freed_at_txn_id)chunk;Pager.writepagerpidbuf;;letwrite_freelist_pagespager:int64Lwt.t=letentries_before=Freelist.to_list(Pager.freelistpager)inletn_entries=List.lengthentries_beforeinletmax_per=Pager.max_freelist_entries_per_pagepagerinletn_fl_pages=(n_entries+max_per-1)/max_perinifn_fl_pages=0thenLwt.return0Lelse(* Allocate all needed pages *)let*page_ids=Lwt_list.map_s(fun()->let*r=Pager.allocpagerinmatchrwith|Okpid->Lwt.returnpid|Errore->Lwt.fail_with(Format.asprintf"write_freelist_pages: %a"Pager.pp_errore))(List.initn_fl_pages(fun_->()))in(* Get FINAL freelist state after allocations *)letfinal_entries=Freelist.to_list(Pager.freelistpager)in(* Split into chunks of max_per *)letrecchunkify=function|[]->[]|lst->letchunk=List.filteri(funi_->i<max_per)lstinletrest=List.filteri(funi_->i>=max_per)lstinchunk::chunkifyrestinletchunks=chunkifyfinal_entriesinletn_chunks=List.lengthchunksinletpid_arr=Array.of_listpage_idsinletnext_ofi=ifi+1<Array.lengthpid_arrthenpid_arr.(i+1)else0Lin(* Write each chunk to a freelist page *)List.iteri(funichunk->write_one_freelist_pagepager~pid:pid_arr.(i)~next:(next_ofi)~chunk)chunks;(* Any extra allocated pages (n_fl_pages > n_chunks) get empty freelist pages *)fori=n_chunkston_fl_pages-1dowrite_one_freelist_pagepager~pid:pid_arr.(i)~next:(next_ofi)~chunk:[]done;Lwt.returnpid_arr.(0);;(* Block until it is safe to recycle WAL frames below [target]. Used by
[checkpoint_unlocked] before [Wal.reset] truncates the index — otherwise
an in-flight reader's [find_page_at] would resolve to a recycled frame
index after the next writer's append.
Two distinct gates, with different urgency (#207):
- RO snapshots ([ro_readers_below]): a local reader still needs frames
below [target]. Recycling past it corrupts its snapshot, so this gate
is honored UNCONDITIONALLY — we wait on the broadcast, which a local
reader always eventually fires via [ro_end].
- Replication floor ([replication_floor_below]): a standby's acked
position, plumbed in by the app. A dead or slow standby must not wedge
the master's WAL forever, so when this is the SOLE remaining blocker we
honor a bounded-yield budget [max_floor_yields] and then proceed anyway
(the standby falls outside the live window and must re-base — #208).
[max_floor_yields = max_int] means unbounded: we wait on the broadcast
([update_replication_position] fires it when the floor advances), exactly
as before this knob existed — no busy-poll. A finite budget polls via
[Lwt.pause] (the project's [wait_for] idiom) because a dead standby
produces no broadcast to wake on.
No [~mutex] is passed to [Lwt_condition.wait]: under cooperative Lwt the
gate check + wait register atomically (no yield between them), so the
standard POSIX condvar mutex pairing isn't needed. Would need revisiting
under a preemptive or effect-based multicore runtime. *)(** Body of [checkpoint] without mutex management. Caller MUST already
hold [t.lock] (e.g. during [commit]). Defined here so [commit]
can invoke it via [maybe_autocheckpoint] below. *)letrecwait_for_readers_past(st:bt_state)~target~replication_max_yields~backup_max_yields=ifst.closingthen(* #338: close is tearing down — stop gating so a parked checkpoint unwinds;
[checkpoint_unlocked] then aborts before any fd I/O. *)Lwt.return_unitelseifro_readers_belowst~targetthen(* A local reader blocks: wait unconditionally on the broadcast. *)let*()=Lwt_condition.waitst.reader_done_condinwait_for_readers_pastst~target~replication_max_yields~backup_max_yieldselseifreplication_floor_belowst~target&&replication_max_yields>0then(* Replication floor is behind and budget remains. Guard with > 0 so
an exhausted budget falls through to the backup floor check below
(review #2); the <= 0 sub-branch is therefore never entered. *)ifreplication_max_yields=max_intthen(* Unbounded: efficient event-driven wait, no busy-poll. *)let*()=Lwt_condition.waitst.reader_done_condinwait_for_readers_pastst~target~replication_max_yields~backup_max_yieldselselet*()=Lwt.pause()inwait_for_readers_pastst~target~replication_max_yields:(replication_max_yields-1)~backup_max_yieldselseifbackup_floor_belowst~targetthenifbackup_max_yields=max_intthenlet*()=Lwt_condition.waitst.reader_done_condinwait_for_readers_pastst~target~replication_max_yields~backup_max_yieldselseifbackup_max_yields<=0thenLwt.return_unitelselet*()=Lwt.pause()inwait_for_readers_pastst~target~replication_max_yields~backup_max_yields:(backup_max_yields-1)elseLwt.return_unit;;letcheckpoint_unlocked(st:bt_state)(wal:Wal.t):unitLwt.t=lettarget=Wal.committed_frameswalin(* #382: [Checkpoint_begin] is a best-effort signal — if the checkpoint
aborts early (store closing, or an I/O error before [Wal.reset]), no
matching [Checkpoint_end] is emitted. Monitor consumers must tolerate an
unbalanced begin. *)emit_eventst(Store_event.Checkpoint_begin{target_frames=target});let*()=wait_for_readers_pastst~target~replication_max_yields:st.replication_gate_max_yields~backup_max_yields:st.backup_gate_max_yieldsinifst.closingthen(* #338: [close] signalled teardown while we were gated — abort before any
pager/WAL fd I/O. [close] fsyncs the WAL itself; the un-migrated frames
replay on next open. No data loss. *)Lwt.return_unitelse((* #338 (review r2): past the gate and about to touch fds — register as
in-flight so [close] drains us before teardown (covers BOTH the auto path
and the manual [checkpoint] path, which share this function). A
checkpoint still parked above on [acquire_write]/the gate is NOT yet
counted, so it cannot wedge close. *)st.ckpt_io_in_flight<-st.ckpt_io_in_flight+1;Lwt.finalize(fun()->letpairs=ref[]inWal.iter_indexwal(funpididx->pairs:=(pid,idx)::!pairs);(* #382: count the pages actually migrated to the main file so the
[Checkpoint_end] event reports the real work done. *)letmigrated=ref0inletrecwrite_each=function|[]->Lwt.return_unit|(pid,idx)::rest->let*r=Wal.read_framewalidxin(matchrwith|Errore->Lwt.fail_with(Format.asprintf"checkpoint read: %a"Wal.pp_errore)|Okpage->let*wr=Pager.flush_one_to_mainst.pager~page_id:pid~buf:pagein(matchwrwith|Errore->Lwt.fail_with(Format.asprintf"checkpoint write: %a"Pager.pp_errore)|Ok()->incrmigrated;write_eachrest))inlet*()=write_each!pairsinlet*sr=Pager.flush_sync_mainst.pagerinmatchsrwith|Errore->Lwt.fail_with(Format.asprintf"checkpoint sync: %a"Pager.pp_errore)|Ok()->(* #337: an async sink ship dispatched from [commit_wal] reads its
frame payloads LAZILY. [Wal.reset] below recycles/zeroes those
frames and bumps the epoch, so wait for any in-flight ship to finish
reading first. The ship runs without [t.lock] (it only reads
frames), so it makes progress while this fiber holds the lock and
parks here. The check-then-reset is yield-free, so no ship
dispatched after the count reaches 0 can slip in before reset. *)let*()=wait_untilst(fun()->st.sink_ships_in_flight=0)inWal.resetwal;(* #382: WAL is reset to a fresh epoch and the migration is done.
Read [Wal.epoch] AFTER reset for the new epoch. *)emit_eventst(Store_event.Wal_reset{epoch=Wal.epochwal});emit_eventst(Store_event.Checkpoint_end{pages_migrated=!migrated});(* #298/#1: checkpoint is a full-sync durability anchor — everything
is now durable and the WAL starts a fresh epoch at frame 0. Reset
the sink ship counter (new epoch) and the batched durability counters
so a long unsynced window doesn't carry stale state across the
anchor. *)st.sink_shipped_frames<-0;st.unsynced_commits<-0;st.last_sync_time<-st.clock();(* Re-pin the replication floor for the new epoch. [Wal.reset] zeroes
committed_frames, but [replication_shipped_frames] still refers to
the old epoch's absolute count. Without re-pinning, the next
checkpoint would see a stale floor that appears to be past the new
target, silently allowing frame recycling before the sink ships
them. *)ifst.on_committed_frames<>Nonethenst.replication_shipped_frames<-Wal.committed_frameswal;(* Re-pin the backup floor for the new epoch (#265). Same
reasoning: without re-pinning, the next checkpoint would see a
stale backup floor from the old epoch and recycle frames before
the backup consumer has captured them. *)ifst.backup_shipped_frames<>max_intthenst.backup_shipped_frames<-Wal.committed_frameswal;Lwt.return_unit)(fun()->st.ckpt_io_in_flight<-st.ckpt_io_in_flight-1;Lwt_condition.broadcastst.reader_done_cond();Lwt.return_unit));;(** Called from [commit] while [lock] is still held (exclusive). If the WAL has
grown past the per-connection threshold, migrate it inline so
subsequent commits start fresh. Best-effort: a checkpoint failure
is swallowed (the commit itself already succeeded). *)letmaybe_autocheckpoint(st:bt_state):unitLwt.t=matchst.walwith|None->Lwt.return_unit|Somewal->letthr=st.wal_autocheckpoint_thresholdinifthr<=0thenLwt.return_unitelseifWal.committed_frameswal<thrthenLwt.return_unitelseLwt.catch(fun()->checkpoint_unlockedstwal)(fun_->Lwt.return_unit);;(* Group-commit coordinator (#77, #151). One fiber per [commit_queue]
runs the actual fsync via [sync_fn]; concurrent writers register a
per-batch resolver and block until the drainer wakes them with the
sync result. Returns the role this fiber played so the caller can
attach drainer-only side work (e.g. autocheckpoint).
Each joiner pushes a [(unit, exn) result Lwt.u] resolver onto
[waiters] and increments [pending] so the drainer can detect
concurrent arrivals. The drainer yields via [Lwt.pause] once to let
any ready-to-write fibers reach the queue, then loops pausing while
[pending] keeps growing. This widens the batch from "2 commits per
fsync" (one Lwt.pause yields one continuation) to "N concurrent
writers per fsync" while adding only one tick of latency to a lone
writer.
Failure semantics (#151): if [sync_fn] raises, the drainer wakes
every joiner's resolver with [Error exn] (so each joiner re-raises
the same exception via [Lwt.fail]) and then re-raises to its own
caller. All N writers in the current batch observe the failure;
none see a spurious [Ok].
Late joiners that arrive while [sync_fn] is in flight register on
the same [waiters] list (since [q.drainer] is still [true]) and so
ride along with the current sync's result — matching the pre-#151
broadcast behaviour. Whether those late frames are physically
flushed by the in-flight fsync is timing-dependent at the kernel
level (POSIX only guarantees flushing of writes queued before the
fsync syscall); this is a pre-existing concern, not introduced by
the error-channel rework. *)letgroup_commit_sync(q:commit_queue)(sync_fn:unit->unitLwt.t):[`Drainer|`Joiner]Lwt.t=ifq.drainerthen(letp,u=Lwt.wait()inq.waiters<-u::q.waiters;q.pending<-q.pending+1;let*r=pinmatchrwith|Ok()->Lwt.return`Joiner|Errorexn->Lwt.failexn)else(q.drainer<-true;(* Gather: one initial pause to let the next-in-line writer reach
the queue; then keep pausing while [pending] keeps growing.
Stops as soon as a pause completes without seeing any new
arrival — keeping per-commit overhead bounded for solo
writers. *)let*()=Lwt.pause()inletrecgatherlast_seen=letnow_seen=q.pendinginifnow_seen>last_seenthenlet*()=Lwt.pause()ingathernow_seenelseLwt.return_unitinlet*()=gather0in(* Run sync first, capturing success or failure; then atomically
snapshot the (possibly grown) waiter list, clear queue state for
the next batch, and wake each joiner with the same result.
Cooperative scheduling guarantees no yield between try_bind's
handler and the iter, so no joiner can register after the
snapshot. *)Lwt.try_bindsync_fn(fun()->letwaiters=q.waitersinq.waiters<-[];q.pending<-0;q.drainer<-false;List.iter(funu->Lwt.wakeup_lateru(Ok()))waiters;Lwt.return`Drainer)(funexn->letwaiters=q.waitersinq.waiters<-[];q.pending<-0;q.drainer<-false;List.iter(funu->Lwt.wakeup_lateru(Errorexn))waiters;Lwt.failexn));;(* Prepare phase of [commit] for the Btree backend. Pushes every tree's
latest root_page through the meta-tree, writes the freelist pages,
and invokes [~header_commit] (either {!Header.commit} for the inline-
sync path or {!Header.commit_no_sync} for group commit). On success
advances [st.current_header] and resets per-txn state. On failure
raises via [Lwt.fail_with] without touching the mutex. *)letcommit_prepare_btree~(header_commit:Pager.t->prev_header:Header.t->new_state:Header.t->(unit,Header.error)resultLwt.t)(st:bt_state):unitLwt.t=let*()=free_old_freelist_pagesst.pager~first_page:st.current_header.freelist_pageinletbindings=Hashtbl.fold(funtidbtacc->(tid,bt)::acc)st.trees[]in(* #174: meta-tree pages are system pages — never stamped with a tree tag. *)Pager.set_write_tagst.pager0l;let*()=Lwt_list.iter_s(fun(tid,bt)->letkey=encode_tree_idtidinletv=encode_root_page(Btree.root_pagebt)inlet*r=Btree.putst.metakeyvinmatchrwith|Okmeta'->st.meta<-meta';Lwt.return_unit|Errore->Lwt.fail_with(Format.asprintf"Store.commit: %a"pp_error(map_btree_erre)))bindingsinlet*freelist_first_page=write_freelist_pagesst.pagerinletnew_state:Header.t={txn_id=0L;(* overwritten by header_commit *)root_page=Btree.root_pagest.meta;freelist_page=freelist_first_page;n_pages_total=Pager.n_pagesst.pager;schema_version=st.schema_version;(* Preserve the on-disk format version this db was opened with (#174);
never silently upgrade or downgrade it here. *)format_version=st.current_header.format_version;(* Preserve the file's page geometry (#95); fixed at creation. *)geom=st.current_header.geom;(* Preserve the encryption marker/canary across commits (#84). *)enc=st.current_header.enc}inlet*r=header_commitst.pager~prev_header:st.current_header~new_stateinmatchrwith|Errore->Lwt.fail_with(Format.asprintf"Store.commit: %a"pp_error(map_header_erre))|Ok()->st.current_header<-{new_statewithtxn_id=Int64.addst.current_header.txn_id1L};st.txn_freelist_snapshot<-None;st.bt_savepoints<-[];(* #297: discard the txn-owned pool — these pages are now part
of the committed tree (freed file-extension pages reuse
within the txn; any not reused by commit are orphans). *)Pager.txn_owned_pool_setst.pager[];(* #266: record this committed root in the as-of commit log. Shared by both
the inline-sync and WAL group-commit paths (both call this function). *)(matchst.historywith|None->Lwt.return_unit|Somesink->letr={History.txn_id=st.current_header.txn_id;timestamp=st.history_now();root_page=st.current_header.root_page}in(* Best-effort: a failed append must never fail a durable commit. *)Lwt.catch(fun()->sink.History.appendr)(fun_->Lwt.return_unit));;(* commit:
- Mem backend: no I/O, just release the writer lock.
- Btree backend without WAL: flush all currently-open trees'
root_pages into the meta-tree, then write a new header pointing at
the new meta root, sync inline, autocheckpoint, release.
- Btree backend with WAL: same prepare phase but using
[Header.commit_no_sync] so the writer can release [lock]
before the fsync. Writers then converge on a per-store
[commit_queue]; one drainer fsyncs and resolves all waiters. Only
the drainer attempts the autocheckpoint (single check per fsync
covers the whole batch).
Note: the Btree.create/put/del API returns a NEW Btree.t after every
mutation (root_page may have changed). We update [st.trees] each
time; here we additionally persist the latest root_page for each
touched tree into the meta-tree (whose own root we then commit via
the header alternating-pages protocol). *)(* After a WAL group-commit, the drainer kicks off an async autocheckpoint if
the WAL has grown past the threshold and none is already in flight. *)letmaybe_autockpt_after_committst=ifst.closing(* #338: no fresh checkpoint once close has signalled teardown. *)||st.wal_autocheckpoint_threshold<=0||Wal.committed_frames(matchst.walwith|Somew->w|None->assertfalse)<st.wal_autocheckpoint_threshold||st.autockpt_in_flightthenLwt.return_unitelse(st.autockpt_in_flight<-true;Lwt.async(fun()->Lwt.finalize(fun()->Lwt.catch(fun()->let*()=Rwlock.acquire_writet.lockinLwt.finalize(fun()->matchst.walwith|None->Lwt.return_unit|Somewal->checkpoint_unlockedstwal)(fun()->Rwlock.release_writet.lock;Lwt.return_unit))(fun_->Lwt.return_unit))(fun()->st.autockpt_in_flight<-false;(* #338: wake a [close] awaiting the in-flight checkpoint to drain. *)Lwt_condition.broadcastst.reader_done_cond();Lwt.return_unit));Lwt.return_unit);;(* WAL-mode commit: prepare the btree (no sync), release the write lock early,
then group-commit-sync the WAL. The elected drainer may autocheckpoint. *)letcommit_waltst=letunlocked=reffalseinletunlock_once()=ifnot!unlockedthen(unlocked:=true;Rwlock.release_writet.lock)inletwal=matchst.walwith|Somew->w|None->assertfalsein(* #382: capture the cumulative committed frame count and this txn's id BEFORE
the append/commit work, so the [Wal_append] emitted at the end reports the
authoritative batch this commit produced (relative to nothing — an absolute
base/count pair), independent of the replication cursor. *)letframes_before=Wal.committed_frameswalinletappend_txn_id=active_txn_idstinLwt.catch(fun()->let*()=commit_prepare_btree~header_commit:Header.commit_no_syncstin(* #382: capture the appended frame count under the write lock, before
[unlock_once] and the group-commit yields — a concurrent
auto-checkpoint can [Wal.reset] (zeroing committed_frames) during
those yields, which would make the count negative if read later. *)letframes_after=Wal.committed_frameswalin(* #298: decide the sync policy under the write lock so the counter is
race-free across concurrent writers, then release the lock. *)letdo_sync=matchst.sync_modewith|`Full->true(* always sync; unsynced_commits is not tracked in Full mode *)|`Off->st.unsynced_commits<-st.unsynced_commits+1;false|`Batched->st.unsynced_commits<-st.unsynced_commits+1;(* #298/#8: a non-positive threshold DISABLES that trigger (the
codebase's [0 = disabled] convention), rather than firing every
commit. If BOTH are 0, batched never syncs on commit. *)letn_trig=st.batch_commits>0&&st.unsynced_commits>=st.batch_commitsinlett_trig=st.batch_interval_ms>0&&(st.clock()-.st.last_sync_time)*.1000.>=float_of_intst.batch_interval_msinn_trig||t_trigin(* #338/#2: reset the batched loss-window counters HERE, under the write
lock, at the moment we decide to sync — not after [unlock_once] post
fsync. The old unlocked post-fsync write raced both the locked
checkpoint reset and a concurrent committer's increment, drifting the
batched-N trigger by one. The frames become durable when the fsync
below completes; a failed fsync raises (the counter is then moot). *)ifdo_syncthen(st.unsynced_commits<-0;st.last_sync_time<-st.clock());unlock_once();let*()=ifdo_syncthen(let*role=group_commit_syncst.commit_queue(fun()->let*r=Pager.wal_syncst.pagerinmatchrwith|Ok()->Lwt.return_unit|Errore->Lwt.fail_with(Format.asprintf"Store.commit: wal_sync: %a"Pager.pp_errore))in(* #298/#1: ship synced frames to the replication sink. The sink
must NEVER see a frame that has not been fsynced, so this fires
ONLY here (in the sync success branch), shipping the whole synced
range since the last ship. In Full mode this fires every commit
(one batch each); in Batched it fires at each sync (the whole
accumulated batch). In Off it never fires on commit — only
checkpoint/close make frames durable. *)(matchst.on_committed_frameswith|None->()|Somecb->letsynced=Wal.committed_frameswalin(* #338 (review r3): do NOT dispatch a fresh ship once [close] has
signalled teardown — its lazy [Wal.read_frame] would race
[wal_close] (a commit mid-fsync when close starts can reach here
AFTER close's drain saw [sink_ships_in_flight = 0]). The frames
are already fsynced; the standby re-syncs from the WAL on
reconnect, same recovery story as an aborted checkpoint. The
check is yield-free up to [Lwt.async], so close cannot set
[closing] between this check and the dispatch. *)if(notst.closing)&&synced>st.sink_shipped_framesthen(letbase=st.sink_shipped_framesinletcount=synced-baseinst.sink_shipped_frames<-synced;letepoch=Wal.epochwalin(* #337: register the ship as in-flight SYNCHRONOUSLY (before the
[Lwt.async] yields) so a checkpoint dispatched right after
sees the count and waits in [checkpoint_unlocked] before
[Wal.reset]; decrement + wake the gate when the lazy reader
completes. *)st.sink_ships_in_flight<-st.sink_ships_in_flight+1;Lwt.async(fun()->Lwt.finalize(fun()->cb~epoch~base_idx:base~count)(fun()->st.sink_ships_in_flight<-st.sink_ships_in_flight-1;Lwt_condition.broadcastst.reader_done_cond();Lwt.return_unit))));matchrolewith|`Joiner->Lwt.return_unit|`Drainer->maybe_autockpt_after_committst)else(* No fsync this commit: still bound the WAL via autocheckpoint
(checkpoint is a full-sync durability anchor). *)maybe_autockpt_after_committstin(* #382: emit the authoritative frame batch for this commit. Uses
[frames_after] captured under the write lock (above) rather than
re-reading [Wal.committed_frames] here — a concurrent auto-checkpoint
could have [Wal.reset] during the group-commit yields, which would
make a re-read count negative. Synchronous and lock-free, consistent
with the [Txn_commit] emit that runs after lock release. *)letappended=frames_after-frames_beforeinemit_eventst(Store_event.Wal_append{txn_id=append_txn_id;base_idx=frames_before;count=appended});Lwt.returnappended)(funexn->unlock_once();Lwt.failexn);;letcommit(Rwt:rwtxn):unitLwt.t=matcht.backendwith|Memtrees->(* #178: merge the shadow back into the live tree. The live tree was
never mutated during the txn — only the shadow was touched — so
commit is the first and only time the live tree sees the txn's
writes. *)(matcht.mem_rw_shadowwith|None->()|Someshadow->List.iter(fun(tid,map)->letr=mem_treetreestidinr:=map)shadow);t.mem_rw_shadow<-None;t.mem_savepoints<-[];Rwlock.release_writet.lock;Lwt.return_unit|Btreest->(* #356: the append cursor is only valid within a txn (its leaf is dirty);
commit flushes dirty pages, so drop it. *)Hashtbl.clearst.bt_append;(* #382/#386: capture the id this txn commits as BEFORE the header is bumped
(commit_prepare_btree advances [st.current_header.txn_id]). [frames] is
the authoritative WAL-appended count returned by [commit_wal] (0 for the
non-WAL path, which appends no WAL frames). *)letcommitted_id=active_txn_idstinlet*frames=matchst.walwith|None->let*()=Lwt.finalize(fun()->let*()=commit_prepare_btree~header_commit:Header.commitstinmaybe_autocheckpointst)(fun()->Rwlock.release_writet.lock;Lwt.return_unit)inLwt.return0(* non-WAL: no WAL frames appended *)|Some_->commit_waltstinemit_eventst(Store_event.Txn_commit{txn_id=committed_id;frames});Lwt.return_unit;;(* rollback:
- Mem: discard the per-txn shadow. With shadow writes (#178) the live
tree is never mutated during a RW txn, so rollback does not need to
restore anything — it just drops the uncommitted shadow.
- Btree: drop cached tree handles so subsequent reads pick up
last-committed roots from the meta-tree, then restore the freelist
snapshot taken at rw_begin and clear dirty pages.
Discard dirty pages from the aborted txn: clear_dirty removes them from
both the dirty set and the read cache, so subsequent reads see committed
data from disk. The freelist snapshot ensures no aborted CoW frees
corrupt future allocations. *)letrollback(Rwt:rwtxn):unitLwt.t=letto_emit=refNonein(matcht.backendwith|Mem_->(* #178: just discard the shadow — the live tree was never touched. *)t.mem_rw_shadow<-None;t.mem_savepoints<-[]|Btreest->(* Drop the per-tree cache so subsequent reads pick up the
last-committed roots from the meta-tree. Note: the meta-tree
itself may have been mutated during this txn (uncommitted puts
to it); we revert it to the last-committed root from the
header. *)Hashtbl.clearst.trees;Hashtbl.clearst.bt_append(* #356: dirty pages discarded below. *);st.meta<-Btree.createst.pager~root_page:st.current_header.root_page;(matchst.txn_freelist_snapshotwith|Somefl->Pager.set_freelistst.pagerfl;Pager.clear_dirtyst.pager;st.txn_freelist_snapshot<-None|None->());st.bt_savepoints<-[];(* #382: rollback does not change [st.current_header.txn_id], so
[active_txn_id] still reads the id this aborted txn would have used.
We capture the event here (where [st] is in scope) but emit it after
the write lock is released, for parity with [commit] — a future
observer doing real work must not stall writers while we hold the
global write lock. *)to_emit:=Some(st,Store_event.Txn_rollback{txn_id=active_txn_idst}));Rwlock.release_writet.lock;(match!to_emitwith|Some(st,ev)->emit_eventstev|None->());Lwt.return_unit;;(* ------------------------------------------------------------------ *)(* WAL checkpoint *)(* ------------------------------------------------------------------ *)(** Migrate every page currently in the WAL index to the main DB, sync
the main DB, then reset the WAL. Holds the RW mutex so no
concurrent commit can append fresh frames while we read the index.
On Mem stores or non-WAL Btree stores this is a no-op. *)letcheckpoint(t:t):unitLwt.t=matcht.backendwith|Mem_->Lwt.return_unit|Btreest->(matchst.walwith|None->Lwt.return_unit|Somewal->let*()=Rwlock.acquire_writet.lockinLwt.finalize(fun()->checkpoint_unlockedstwal)(fun()->Rwlock.release_writet.lock;Lwt.return_unit));;letwal_autocheckpoint(t:t):int=matcht.backendwith|Mem_->0|Btreest->st.wal_autocheckpoint_threshold;;letset_wal_autocheckpoint(t:t)(n:int):unit=matcht.backendwith|Mem_->()|Btreest->st.wal_autocheckpoint_threshold<-max0n;;letdurability(t:t):durability=matcht.backendwith|Mem_->Full|Btreest->(matchst.sync_modewith|`Full->Full|`Off->Off|`Batched->Batched{commits=st.batch_commits;interval_ms=st.batch_interval_ms});;(* #298/#9: durability <-> string helpers for the PRAGMA layer (Group B). *)letdurability_of_string(s:string):durabilityoption=matchString.lowercase_asciiswith|"full"->SomeFull|"off"->SomeOff|"batched"->Some(Batched{commits=default_batch_commits;interval_ms=default_batch_interval_ms})|_->None;;letstring_of_durability(d:durability):string=matchdwith|Full->"full"|Batched_->"batched"|Off->"off";;letset_durability(t:t)(d:durability):unit=matcht.backendwith|Mem_->()|Btreest->letrequested_non_full=matchdwith|Full->false|_->trueinifrequested_non_full&&st.on_committed_frames<>Nonethen((* #298: a sink mandates Full — ignore the relax request but still record
any Batched params for when the sink is later removed. The checkpoint
replica-floor gate requires every committed frame to be shipped, which
only holds under Full. *)matchdwith|Batched{commits;interval_ms}->st.batch_commits<-max0commits;st.batch_interval_ms<-max0interval_ms|_->())else(letnew_mode=matchdwith|Full->`Full|Off->`Off|Batched_->`Batchedin(* #298/#4: when the mode actually changes, reset the batched durability
counters so a long [Off] period doesn't carry a huge stale
[unsynced_commits] into [Batched] (which would immediately fsync), and
the T window restarts at mode entry. Safe because [close]/[checkpoint]
(frame-state based) remain the durability anchors; the counter is only a
trigger heuristic. *)ifst.sync_mode<>new_modethen(st.unsynced_commits<-0;st.last_sync_time<-st.clock());matchdwith|Full->st.sync_mode<-`Full|Off->st.sync_mode<-`Off|Batched{commits;interval_ms}->st.sync_mode<-`Batched;st.batch_commits<-max0commits;st.batch_interval_ms<-max0interval_ms);;(* #298: True iff a replication commit-sink is currently registered. While
active, durability is pinned to [Full]. *)letcommit_callback_active(t:t):bool=matcht.backendwith|Mem_->false|Btreest->st.on_committed_frames<>None;;(* #298/#3: force any committed-but-unsynced WAL frames to disk now. No-op in
[Full] mode, on the in-memory backend, or when nothing is pending. Group B
calls this when tightening durability so already-acked commits become durable
immediately rather than only on the next commit. *)letflush_unsynced(t:t):unitLwt.t=matcht.backendwith|Mem_->Lwt.return_unit|Btreest->letpending=st.sync_mode<>`Full&&matchst.walwith|Somew->Wal.committed_framesw>0|None->falseinifnotpendingthenLwt.return_unitelse(* #298/#5: route the fsync through the group-commit serializer rather
than calling [Pager.wal_sync] unlocked, which raced [commit_wal]'s
post-unlock fsync + cursor. A sink now forces Full, so [flush_unsynced]
only runs when NO sink is active — the previous sink-ship block here is
dead and has been removed. *)let*(_:[`Drainer|`Joiner])=group_commit_syncst.commit_queue(fun()->let*r=Pager.wal_syncst.pagerinmatchrwith|Ok()->Lwt.return_unit|Errore->Lwt.fail_with(Format.asprintf"Store.flush_unsynced: wal_sync: %a"Pager.pp_errore))inst.unsynced_commits<-0;st.last_sync_time<-st.clock();Lwt.return_unit;;letsync_batch_commits(t:t):int=matcht.backendwith|Mem_->default_batch_commits|Btreest->st.batch_commits;;letset_sync_batch_commits(t:t)(n:int):unit=matcht.backendwith|Mem_->()|Btreest->st.batch_commits<-max0n;;letsync_batch_interval_ms(t:t):int=matcht.backendwith|Mem_->default_batch_interval_ms|Btreest->st.batch_interval_ms;;letset_sync_batch_interval_ms(t:t)(n:int):unit=matcht.backendwith|Mem_->()|Btreest->st.batch_interval_ms<-max0n;;letset_clock(t:t)(c:unit->float):unit=matcht.backendwith|Mem_->()|Btreest->st.clock<-c;st.last_sync_time<-c();;(* Number of fsyncs the WAL has performed since open. Exposed for #77
group-commit testing: lets the test assert that N concurrent
autocommit fibers issue ≪ N fsyncs (proof of coalescing). *)letwal_sync_count(t:t):int=matcht.backendwith|Mem_->0|Btreest->(matchst.walwith|None->0|Somew->Granary_storage.Wal.sync_countw);;(* Diagnostic/testing accessors for #164: observe that ending a snapshot
releases its bookkeeping even when its reader closure raised. *)letactive_reader_count(t:t):int=matcht.backendwith|Mem_->0|Btreest->Hashtbl.fold(fun_cacc->acc+c)st.active_readers0;;letpinned_page_count(t:t):int=matcht.backendwith|Mem_->0|Btreest->Pager.pinned_countst.pager;;letlive_read_locks(t:t):int=Rwlock.readerst.lock(* ------------------------------------------------------------------ *)(* Savepoints (both Mem and B-tree backends) *)(* ------------------------------------------------------------------ *)(** Push a named savepoint: snapshot the current shadow state (#178). *)letsavepoint_begin(Rwt:rwtxn)name=matcht.backendwith|Memtrees->letsnap=matcht.mem_rw_shadowwith|None->Hashtbl.fold(funtidracc->(tid,!r)::acc)trees[]|Someshadow->shadowint.mem_savepoints<-(name,snap)::t.mem_savepoints;Lwt.return_unit|Btreest->lettree_roots=Hashtbl.fold(funtidbtacc->(tid,Btree.root_pagebt)::acc)st.trees[]inletsp={sp_name=name;sp_meta_root=Btree.root_pagest.meta;sp_tree_roots=tree_roots;sp_freelist=Pager.freelistst.pager;sp_n_pages=Pager.n_pagesst.pager;sp_dirty=Pager.dirty_clonest.pager;sp_txn_pool=Pager.txn_owned_pool_getst.pager}inst.bt_savepoints<-sp::st.bt_savepoints;emit_eventst(Store_event.Savepoint_begin{txn_id=active_txn_idst;name});Lwt.return_unit;;(** Release the named savepoint and all newer ones (writes are kept). *)letsavepoint_release(Rwt:rwtxn)name=matcht.backendwith|Mem_->letrecdrop=function|[]->[]|(n,_)::restwhenString.equalnname->rest|_::rest->droprestint.mem_savepoints<-dropt.mem_savepoints;Lwt.return_unit|Btreest->letrecdrop=function|[]->[]|sp::restwhenString.equalsp.sp_namename->rest|_::rest->droprestinst.bt_savepoints<-dropst.bt_savepoints;emit_eventst(Store_event.Savepoint_release{txn_id=active_txn_idst;name});Lwt.return_unit;;(** Rollback to the named savepoint: restore snapshot, drop newer savepoints,
keep the named savepoint so it can be rolled back to again. *)letsavepoint_rollback(Rwt:rwtxn)name=matcht.backendwith|Mem_->(* #178: restore the shadow to the savepoint snapshot. The live tree
was never mutated, so we just replace the shadow. *)letrecfind=function|[]->()(* savepoint not found — no-op *)|(n,snap)::restwhenString.equalnname->t.mem_rw_shadow<-Somesnap;t.mem_savepoints<-(name,snap)::rest|_::rest->findrestinfindt.mem_savepoints;Lwt.return_unit|Btreest->letrecfind=function|[]->()|sp::restwhenString.equalsp.sp_namename->(* Restore meta-tree root *)st.meta<-Btree.createst.pager~root_page:sp.sp_meta_root;(* Restore per-tree roots: drop the cache, re-populate from snapshot. *)Hashtbl.clearst.trees;List.iter(fun(tid,root)->letbt=Btree.createst.pager~root_page:rootinHashtbl.replacest.treestidbt)sp.sp_tree_roots;(* Restore freelist + n_pages + dirty set. Pages above sp.sp_n_pages
that were freshly allocated in the rolled-back range become
orphans in the file but are not in any tree, freelist, or
dirty set — harmless storage leak. *)Pager.set_freelistst.pagersp.sp_freelist;Pager.set_n_pagesst.pagersp.sp_n_pages;Pager.dirty_restorest.pagersp.sp_dirty;Pager.txn_owned_pool_setst.pagersp.sp_txn_pool;(* #356: the restored dirty pages may be an earlier version of the
cached rightmost leaf; drop the append cursor so it re-primes. *)Hashtbl.clearst.bt_append;(* Keep the named savepoint at the top so it can be re-used. *)st.bt_savepoints<-sp::rest|_::rest->findrestinfindst.bt_savepoints;emit_eventst(Store_event.Savepoint_rollback{txn_id=active_txn_idst;name});Lwt.return_unit;;(* ------------------------------------------------------------------ *)(* get / put / del *)(* ------------------------------------------------------------------ *)lettxn_store:typea.atxn->t=function|Rosnap->snap.rs_store|Rws->s;;letget:typea.atxn->tree_id->bytes->bytesoptionLwt.t=funtxtidkey->matchtxwith|Rosnap->(matchsnap.rs_store.backendwith|Mem_->(* #178: read from the snapshot captured at ro_begin so this
reader never sees uncommitted writes from a concurrent writer
that may later roll back. *)letmap=matchsnap.rs_mem_snapwith|Somesnap->mem_tree_snapsnaptid|None->Bytes_map.emptyinLwt.return(Bytes_map.find_optkeymap)|Btreest->let*r=bt_get_tree_rosnapsttidinlet*bt=unwrap_errorrinlet*g=Btree.getbtkeyin(matchgwith|Okv->Lwt.returnv|Errore->Lwt.fail_with(Format.asprintf"Store.get(ro): %a"pp_error(map_btree_erre))))|Rwt->(matcht.backendwith|Memtrees->(* #178: read from the active RW shadow, not the live tree.
The shadow contains the txn's own writes layered on top of the
pre-txn committed state; the live tree is never mutated until
commit. *)letmap=matcht.mem_rw_shadowwith|None->!(mem_treetreestid)|Someshadow->shadow_getshadowtreestidinLwt.return(Bytes_map.find_optkeymap)|Btreest->let*r=bt_get_treesttidinlet*bt=unwrap_errorrinlet*g=Btree.getbtkeyin(matchgwith|Okv->Lwt.returnv|Errore->Lwt.fail_with(Format.asprintf"Store.get(rw): %a"pp_error(map_btree_erre))));;letput(Rwt:rwtxn)tidkeyvalue:unitLwt.t=matcht.backendwith|Memtrees->(* #178: write to the shadow, not the live tree. *)(matcht.mem_rw_shadowwith|None->letr=mem_treetreestidinr:=Bytes_map.addkeyvalue!r|Someshadow->t.mem_rw_shadow<-Some(shadow_updateshadowtreestid(Bytes_map.addkeyvalue)));Lwt.return_unit|Btreest->let*r=bt_get_treesttidinlet*bt=unwrap_errorrinPager.set_write_tagst.pager(tree_tagsttid);(* #356: [put] may replace or split anywhere; invalidate the append cursor
for this tree so a subsequent append re-primes from the live rightmost. *)Hashtbl.removest.bt_appendtid;let*p=Btree.putbtkeyvaluein(matchpwith|Okbt'->Hashtbl.replacest.treestidbt';Lwt.return_unit|Errore->Lwt.fail_with(Format.asprintf"Store.put: %a"pp_error(map_btree_erre)));;letput_x(Rwt:rwtxn)tidkeyvalue:bytesoptionLwt.t=matcht.backendwith|Memtrees->letmap=matcht.mem_rw_shadowwith|None->!(mem_treetreestid)|Someshadow->shadow_getshadowtreestidin(matchBytes_map.find_optkeymapwith|Some_->Lwt.return(SomeBytes.empty)|None->(matcht.mem_rw_shadowwith|None->letr=mem_treetreestidinr:=Bytes_map.addkeyvalue!r|Someshadow->t.mem_rw_shadow<-Some(shadow_updateshadowtreestid(Bytes_map.addkeyvalue)));Lwt.returnNone)|Btreest->let*r=bt_get_treesttidinlet*bt=unwrap_errorrinPager.set_write_tagst.pager(tree_tagsttid);(* #356 append fast path: O(1) in-place append to the cached rightmost leaf
when [key] is strictly greater than the cursor's max. *)let*fast=matchHashtbl.find_optst.bt_appendtidwith|SomeacwhenBytes.comparekey(Btree.append_cursor_max_keyac)>0->let*outcome=Btree.try_inplace_appendbtac~key~valuein(matchoutcomewith|Btree.Appendedac'->Hashtbl.replacest.bt_appendtidac';Lwt.return(SomeNone)|Btree.Not_applicable->Lwt.returnNone|Btree.Append_failede->Lwt.fail_with(Format.asprintf"Store.put_x: %a"pp_error(map_btree_erre)))|_->Lwt.returnNonein(matchfastwith|Someold_opt->Lwt.returnold_opt|None->(* General path. Remember the prior cursor max to decide whether this
insert was an append worth re-priming the cursor for. *)letprev_max=Option.mapBtree.append_cursor_max_key(Hashtbl.find_optst.bt_appendtid)inlet*p=Btree.put_xbtkeyvaluein(matchpwith|Errore->Lwt.fail_with(Format.asprintf"Store.put_x: %a"pp_error(map_btree_erre))|Ok(bt',old_opt)->Hashtbl.replacest.treestidbt';(matchold_optwith|Some_->(* Conflict: nothing inserted; leave the cursor as-is. *)Lwt.returnold_opt|None->(* Inserted. If [key] extends the tree to the right (an append),
re-prime the cursor from the new rightmost leaf; otherwise it was
a middle insert and the cursor is invalidated. *)letappend_like=matchprev_maxwith|None->true(* cold start: probe whether it was an append *)|Somem->Bytes.comparekeym>0inifnotappend_likethen(Hashtbl.removest.bt_appendtid;Lwt.returnNone)elselet*rc=Btree.rightmost_append_cursorbt'in(matchrcwith|Ok(Someac)whenBytes.equal(Btree.append_cursor_max_keyac)key->Hashtbl.replacest.bt_appendtidac;Lwt.returnNone|Ok_->Hashtbl.removest.bt_appendtid;Lwt.returnNone|Errore->Lwt.fail_with(Format.asprintf"Store.put_x: %a"pp_error(map_btree_erre))))));;letdel(Rwt:rwtxn)tidkey:unitLwt.t=matcht.backendwith|Memtrees->(* #178: remove from the shadow, not the live tree. *)(matcht.mem_rw_shadowwith|None->letr=mem_treetreestidinr:=Bytes_map.removekey!r|Someshadow->t.mem_rw_shadow<-Some(shadow_updateshadowtreestid(Bytes_map.removekey)));Lwt.return_unit|Btreest->let*r=bt_get_treesttidinlet*bt=unwrap_errorrinPager.set_write_tagst.pager(tree_tagsttid);(* #356: a delete may free or restructure the rightmost leaf; invalidate. *)Hashtbl.removest.bt_appendtid;let*d=Btree.delbtkeyin(matchdwith|Okbt'->Hashtbl.replacest.treestidbt';Lwt.return_unit|Errore->Lwt.fail_with(Format.asprintf"Store.del: %a"pp_error(map_btree_erre)));;(* #174: register the page-header stamp (low 32 bits of the schema
fingerprint) for [tid]. Subsequently-written Branch/Leaf pages of that tree
carry the tag in their reserved header bytes. No-op on the in-memory
backend (no pages). *)letset_tree_tag(t:t)(tid:tree_id)(tag:int32):unit=matcht.backendwith|Mem_->()|Btreest->Hashtbl.replacest.tree_tagstidtag;;(* ------------------------------------------------------------------ *)(* Cursors *)(* ------------------------------------------------------------------ *)(* Drain a B+-tree cursor into an in-memory snapshot list. Phase 1
cursors are materialised; streaming cursors arrive later. *)letdrain_btree_cursor(c:Btree.cursor):(bytes*bytes)listLwt.t=letrecloopacc=let*r=Btree.cursor_nextcinmatchrwith|Errore->Lwt.fail_with(Format.asprintf"Store.cursor: %a"pp_error(map_btree_erre))|OkNone->Lwt.return(List.revacc)|Ok(Somekv)->loop(kv::acc)inloop[];;letcursor_open:typea.atxn->tree_id->cursorLwt.t=funtxtid->matchtxwith|Rosnap->(matchsnap.rs_store.backendwith|Mem_->(* #178: materialise from the snapshot captured at ro_begin.
Without this, a concurrent writer's uncommitted modifications
would leak into the cursor — and survive even if the writer
later rolls back. *)letmap=matchsnap.rs_mem_snapwith|Somesnap->mem_tree_snapsnaptid|None->Bytes_map.emptyinletentries=Bytes_map.bindingsmapinLwt.return{all=entries;remaining=[];ready=false}|Btreest->let*r=bt_get_tree_rosnapsttidinlet*bt=unwrap_errorrinlet*co=Btree.cursor_openbtin(matchcowith|Errore->Lwt.fail_with(Format.asprintf"Store.cursor_open(ro): %a"pp_error(map_btree_erre))|Okc->let*entries=drain_btree_cursorcinBtree.cursor_closec;Lwt.return{all=entries;remaining=[];ready=false}))|Rw_->lett=txn_storetxin(matcht.backendwith|Memtrees->(* #178: cursor materialises from the active RW shadow. *)letmap=matcht.mem_rw_shadowwith|None->!(mem_treetreestid)|Someshadow->shadow_getshadowtreestidinletentries=Bytes_map.bindingsmapinLwt.return{all=entries;remaining=[];ready=false}|Btreest->let*r=bt_get_treesttidinlet*bt=unwrap_errorrinlet*co=Btree.cursor_openbtin(matchcowith|Errore->Lwt.fail_with(Format.asprintf"Store.cursor_open: %a"pp_error(map_btree_erre))|Okc->let*entries=drain_btree_cursorcinBtree.cursor_closec;Lwt.return{all=entries;remaining=[];ready=false}));;letcursor_close_=()letcursor_firstc=c.remaining<-c.all;matchc.allwith|[]->c.ready<-false;Not_found`End|(k,_)::_->c.ready<-true;Foundk;;letcursor_seekckey=letrecfind=function|[]->c.remaining<-[];c.ready<-false;Not_found`End|(k,_)::_ascur->letcmp=Bytes.comparekkeyinifcmp>=0then(c.remaining<-cur;c.ready<-true;ifcmp=0thenFoundkelseNot_found(`Greaterk))elsefind(List.tlcur)infindc.all;;letcursor_nextc=matchc.remainingwith|[]->None|entry::rest->ifc.readythen(c.ready<-false;Someentry)else(c.remaining<-rest;matchrestwith|[]->None|next::_->Somenext);;letcursor_valuec=matchc.remainingwith|(_,v)::_whenc.ready->Somev|_->None;;(* ------------------------------------------------------------------ *)(* Native streaming seek (#228, #229) *)(* *)(* The materialised [cursor] above drains the WHOLE tree at open time so *)(* it can offer a synchronous [cursor_seek]/[cursor_next] API. For *)(* point/prefix probes (index lookups, UNIQUE pre-checks, FK checks) *)(* that O(n) drain dominates — it turns an O(log n) seek into a full *)(* table scan, and an n-row bulk insert into O(n^2). [seek_ge] instead *)(* descends the B+-tree natively in O(log n) and streams matches lazily, *)(* never materialising more than the entries the caller actually reads. *)(* Semantics match [cursor_open]+[cursor_seek]+[cursor_next]: the first *)(* [seek_next] returns the first entry with key >= [key], then ascending.*)(* ------------------------------------------------------------------ *)typeseek_impl=|SC_memof(bytes*bytes)Seq.tref|SC_btofBtree.cursor(* #235: when the backing promise is already determined — the [Mem] backend, or *)(* a B+-tree page already resident in the pager cache — [Lwt.bind] runs the *)(* caller's continuation synchronously, so a recursive [seek_next] consumer *)(* (the [gather]/[scan] loops in exec.ml) nests one OCaml frame per call *)(* instead of returning to a trampoline. Streaming a pathologically common *)(* term/value (millions of cache-resident postings) would then overflow the *)(* stack. To bound stack growth regardless of consumer shape, [seek_next] *)(* splices an [Lwt.pause] every [seek_pause_interval] calls: that defers the *)(* continuation to the scheduler, unwinding the stack. The cost is one *)(* cooperative yield per N calls — negligible. *)typeseek_cursor={mutablesc_calls:int;sc_impl:seek_impl}(* Bounds the synchronous recursion depth to <= this many frames; the yield then
amortises to one [Lwt.pause] per that many reads. 256 trades a tiny, fixed
per-scan overhead for a shallow stack ceiling. *)letseek_pause_interval=256letmk_seek_cursorsc_impl={sc_calls=0;sc_impl}letseek_ge:typea.atxn->tree_id->bytes->seek_cursorLwt.t=funtxtidkey->matchtxwith|Rosnap->(matchsnap.rs_store.backendwith|Mem_->letmap=matchsnap.rs_mem_snapwith|Somesnap->mem_tree_snapsnaptid|None->Bytes_map.emptyinLwt.return(mk_seek_cursor(SC_mem(ref(Bytes_map.to_seq_fromkeymap))))|Btreest->let*r=bt_get_tree_rosnapsttidinlet*bt=unwrap_errorrinlet*co=Btree.cursor_openbtin(matchcowith|Errore->Lwt.fail_with(Format.asprintf"Store.seek_ge(ro): %a"pp_error(map_btree_erre))|Okc->let*sr=Btree.cursor_seekckeyin(matchsrwith|Errore->Lwt.fail_with(Format.asprintf"Store.seek_ge(ro): %a"pp_error(map_btree_erre))|Ok_->Lwt.return(mk_seek_cursor(SC_btc)))))|Rw_->lett=txn_storetxin(matcht.backendwith|Memtrees->letmap=matcht.mem_rw_shadowwith|None->!(mem_treetreestid)|Someshadow->shadow_getshadowtreestidinLwt.return(mk_seek_cursor(SC_mem(ref(Bytes_map.to_seq_fromkeymap))))|Btreest->let*r=bt_get_treesttidinlet*bt=unwrap_errorrinlet*co=Btree.cursor_openbtin(matchcowith|Errore->Lwt.fail_with(Format.asprintf"Store.seek_ge: %a"pp_error(map_btree_erre))|Okc->let*sr=Btree.cursor_seekckeyin(matchsrwith|Errore->Lwt.fail_with(Format.asprintf"Store.seek_ge: %a"pp_error(map_btree_erre))|Ok_->Lwt.return(mk_seek_cursor(SC_btc)))));;(* Return the next (key, value) >= the seek key in ascending order, or [None]
when exhausted. The first call returns the positioned entry. *)letseek_next:seek_cursor->(bytes*bytes)optionLwt.t=funsc->letresult=matchsc.sc_implwith|SC_memr->(match!r()with|Seq.Nil->Lwt.return_none|Seq.Cons(kv,rest)->r:=rest;Lwt.return_somekv)|SC_btc->let*r=Btree.cursor_nextcin(matchrwith|Okkv->Lwt.returnkv|Errore->Lwt.fail_with(Format.asprintf"Store.seek_next: %a"pp_error(map_btree_erre)))in(* The cursor is advanced eagerly above; [result] already holds this call's
entry (or [None]), so splicing a pause here only defers the *return*, never
reordering or dropping a match (#235). We count calls, not matches: every
recursive consumer call nests a frame whether or not it yields a row, so the
terminal [None] call is counted too.
Yielding mid-stream is safe under the current concurrency model: a paused RO
seek streams from an immutable snapshot, and a paused RW seek holds the
single-writer lock — so no other fiber can mutate the tree under the cursor
between pause and resume. If that invariant is ever relaxed (concurrent
writers), revisit this yield point. *)sc.sc_calls<-sc.sc_calls+1;ifsc.sc_callsmodseek_pause_interval=0thenLwt.bind(Lwt.pause())(fun()->result)elseresult;;letseek_close:seek_cursor->unit=funsc->matchsc.sc_implwith|SC_mem_->()|SC_btc->Btree.cursor_closec;;letwal_modet=matcht.backendwith|Mem_->false|Btreest->st.wal<>None;;letfreelist_sizet=matcht.backendwith|Mem_->0|Btreest->Freelist.size(Pager.freelistst.pager);;letfreelist_entriest=matcht.backendwith|Mem_->[]|Btreest->Freelist.to_list(Pager.freelistst.pager);;letn_pagest=matcht.backendwith|Mem_->0L|Btreest->Pager.n_pagesst.pager;;(* Enumerate all tree_ids known to the meta tree. For VACUUM. *)letlist_tree_idst:tree_idlistLwt.t=matcht.backendwith|Memtrees->Lwt.return(Hashtbl.fold(funtid_acc->tid::acc)trees[])|Btreest->let*r=Btree.cursor_openst.metain(matchrwith|Errore->Lwt.fail_with(Format.asprintf"Store.list_tree_ids: %a"pp_error(map_btree_erre))|Okcur->letrecloopacc=let*r=Btree.cursor_nextcurinmatchrwith|Errore->Lwt.fail_with(Format.asprintf"Store.list_tree_ids: %a"pp_error(map_btree_erre))|OkNone->Lwt.return(List.revacc)|Ok(Some(k,_v))->lettid,_=Varint.decode_int64k0inloop(Int64.to_inttid::acc)inlet*result=loop[]inBtree.cursor_closecur;Lwt.returnresult);;typepage_sink=page_id:int64->page:Cstruct.t->unitLwt.t(* Shared page-image iteration for [copy_to]/[rekey_to]. Under an RO snapshot,
yields each PLAINTEXT page (page_id, buf) for page_id in [0, n), resolving the
WAL overlay (bounded to the committed_frames horizon captured at ro_begin)
before the main DB. The iteration is bounded by the snapshot-time page count
so growth during the copy does not pull pages outside the snapshot.
The buffer handed to [f] may be owned by the pager cache — a caller that
mutates it MUST copy first. *)letiter_snapshot_pagesst(Rosnap)~(f:page_id:int64->page:Cstruct.t->unitLwt.t):unitLwt.t=(* Read the page count INSIDE the snapshot (after ro_begin) so that the loop
bound n is consistent with the snapshot's WAL horizon. If a writer commits
between ro_begin and reading n_pages, the snapshot's WAL horizon already
includes those new pages; reading n_pages after ro_begin ensures we copy
them too. *)letn=Pager.n_pagesst.pagerinlethorizon=snap.rs_snap_framesinletrecloop(page_id:int64)=ifInt64.comparepage_idn>=0thenLwt.return_unitelselet*page_buf=matchst.walwith|None->(* Non-WAL path: pass ~snapshot_frames:0 so the read path skips the
dirty set entirely (pager.ml:274-276), avoiding any
concurrent-writer uncommitted data. With no WAL the
resolve_wal_page returns Ok None, falling through to
load_main_page which reads the committed on-disk state. Also pin
the page so the writer's eviction pressure doesn't drop our copy. *)let*r=Pager.read~snapshot_frames:0~pin_set:snap.rs_pinnedst.pagerpage_idin(matchrwith|Okbuf->Lwt.returnbuf|Errore->Lwt.fail_with(Format.asprintf"Store.iter_snapshot_pages(pg=%Ld): %a"page_idPager.pp_errore))|Somewal->(* WAL-mode path: snapshot overlay (WAL first, then main DB). *)(matchWal.find_page_atwalpage_id~max_frame:horizonwith|Someidx->let*r=Wal.read_framewalidxin(matchrwith|Okbuf->Lwt.returnbuf|Errore->Lwt.fail_with(Format.asprintf"Store.iter_snapshot_pages(pg=%Ld,frame=%d): %a"page_ididxWal.pp_errore))|None->let*r=Pager.read~snapshot_frames:horizon~pin_set:snap.rs_pinnedst.pagerpage_idin(matchrwith|Okbuf->Lwt.returnbuf|Errore->Lwt.fail_with(Format.asprintf"Store.iter_snapshot_pages(pg=%Ld): %a"page_idPager.pp_errore)))inlet*()=f~page_id~page:page_bufinloop(Int64.addpage_id1L)inloop0L;;(* One-shot consistent full copy via an RO snapshot + page sink (#93).
When the source is encrypted (#84), data pages (>= 2) are re-encrypted under
the source's own key before reaching the sink, so the destination is a
faithful, self-contained encrypted DB (open it with the same key) and no
user-data plaintext transits the sink. Pages 0 and 1 are plaintext headers
(carrying the enc marker + canary) and are copied verbatim.
On the Mem backend this is a no-op (there are no pages to copy). *)letcopy_to(t:t)(sink:page_sink):unitLwt.t=matcht.backendwith|Mem_->Lwt.return_unit|Btreest->with_rot(funro->iter_snapshot_pagesstro~f:(fun~page_id~page->matchst.cipherwith|SomecwhenInt64.comparepage_id2L>=0->(* Copy first — [page] may be the pager's cached buffer. *)lettmp=Cstruct.create(Cstruct.lengthpage)inCstruct.blitpage0tmp0(Cstruct.lengthpage);Crypto.encrypt_pagec~page_idtmp;sink~page_id~page:tmp|_->sink~page_id~page));;(* #215: offline key rotation. [t] must have been opened WITH THE OLD KEY so
reads decrypt to plaintext; every data page (>= 2) is re-encrypted under a
fresh cipher built from [new_key], and each header page (0, 1) has its canary
rewritten under the new key (txn parity and all other fields preserved) and
its CRC resealed. The sunk page image is a self-contained encrypted DB under
[new_key] with no WAL. Rejects a plaintext source ([Not_encrypted]) and a
wrong-length key ([Block_error]). *)letrekey_to(t:t)~(new_key:string)(sink:page_sink):(unit,error)resultLwt.t=matcht.backendwith|Mem_->Lwt.return_ok()|Btreest->(matchst.cipherwith|None->Lwt.return_errorNot_encrypted|Some_old->(matchCrypto.create~key:new_keywith|Error`Bad_key_length->Lwt.return_error(Block_error"encryption key must be 32 bytes")|Okc'->letnonce=Mirage_crypto_rng.generateCrypto.nonce_leninletcanary_tag=Crypto.make_canaryc'~nonceinlet*()=with_rot(funro->iter_snapshot_pagesstro~f:(fun~page_id~page->letlen=Cstruct.lengthpageinlettmp=Cstruct.createleninCstruct.blitpage0tmp0len;ifInt64.comparepage_id2L<0then((* header page: rewrite the canary under the new key, leaving
enc_magic + every structural field intact, then reseal CRC. *)letf=Page.read_header_fieldstmpinPage.write_header_fieldstmp{fwithPage.canary_nonce=nonce;canary_tag};Page.sealtmp;sink~page_id~page:tmp)else(Crypto.encrypt_pagec'~page_idtmp;sink~page_id~page:tmp)))inLwt.return_ok()));;(* ------------------------------------------------------------------ *)(* Replication consumer integration (#92) *)(* ------------------------------------------------------------------ *)(** Register the replication consumer's shipped position so checkpoint
truncation waits for frames to be shipped before recycling them. *)letupdate_replication_position(t:t)~shipped=matcht.backendwith|Mem_->()|Btreest->st.replication_shipped_frames<-shipped;Lwt_condition.broadcastst.reader_done_cond();;(** Bounded-yield "timeout" for the checkpoint gate's wait on the
replication floor (#207). Returns [max_int] (unbounded) by default.
[0] on the in-memory backend (no checkpoint gating). *)letreplication_gate_max_yields(t:t):int=matcht.backendwith|Mem_->0|Btreest->st.replication_gate_max_yields;;(** Set the bounded-yield budget the checkpoint gate will spend waiting for
the replication floor (a standby's acked position) to reach the
checkpoint target before proceeding anyway. See {!update_replication_position}.
Pure-Mirage has no ambient clock, so this "timeout" is a count of
cooperative [Lwt.pause] yields rather than wall-clock time. [max_int]
(the default) means wait indefinitely — a dead standby wedges the WAL,
matching the behavior before this knob existed. A finite value bounds
the wait: once spent, the checkpoint proceeds and the now-stranded
standby must re-base (#208). Negative inputs clamp to [0] (proceed
immediately if the floor is behind).
Local RO readers are never abandoned by this budget — only the
replication floor. No-op on the in-memory backend. *)letset_replication_gate_max_yields(t:t)(n:int):unit=matcht.backendwith|Mem_->()|Btreest->st.replication_gate_max_yields<-max0n;;(** Get (epoch, committed_frames) for the active WAL; [None] if no WAL. *)letreplication_state(t:t)=matcht.backendwith|Mem_->None|Btreest->(matchst.walwith|None->None|Somewal->Some(Wal.epochwal,Wal.committed_frameswal));;(* ------------------------------------------------------------------ *)(* Incremental backup (#265) *)(* ------------------------------------------------------------------ *)(** Register the backup consumer's captured position so checkpoint
truncation waits for frames to be backed up before recycling them.
Analogous to {!update_replication_position} but for the incremental
backup watermark. *)letupdate_backup_position(t:t)~shipped=matcht.backendwith|Mem_->()|Btreest->st.backup_shipped_frames<-shipped;Lwt_condition.broadcastst.reader_done_cond();;(** Get (epoch, committed_frames) for the active WAL; [None] if no WAL.
Review #8: delegates to {!replication_state} — the two functions share
the same body because both track the same WAL position. *)letbackup_state(t:t)=replication_statet(** Return the backup floor's bounded-yield budget for the checkpoint
gate. Defaults to [max_int] (unbounded) on the B+-tree backend,
[0] on the in-memory backend (no checkpoint gating). *)letbackup_gate_max_yields(t:t):int=matcht.backendwith|Mem_->0|Btreest->st.backup_gate_max_yields;;(** Set the bounded-yield budget the checkpoint gate will spend waiting
for the backup floor to reach the checkpoint target before proceeding
anyway (#265). Same semantics as {!set_replication_gate_max_yields}.
Negative inputs clamp to [0]. No-op on the in-memory backend. *)letset_backup_gate_max_yields(t:t)(n:int):unit=matcht.backendwith|Mem_->()|Btreest->st.backup_gate_max_yields<-max0n;;(** A captured WAL frame for incremental backup (#265). Contains the
full frame metadata and page payload needed to reconstruct the
database at a later point.
The {!checksum} field covers the decrypted page payload (transport
integrity for the backup frame), matching the same scheme used by
{!Granary_replication.replicated_frame}. For unencrypted WALs the
plaintext equals the on-disk page; for encrypted WALs the checksum
guards against corruption of the decrypted content during transport
or storage, not the on-disk ciphertext. *)typebackup_frame={epoch:int64;frame_idx:int;page_id:int64;is_commit:bool;page:Cstruct.t;checksum:int64;source_salt:int64;source_seed:int64}(* NOTE (review #8): backup_frame and Granary_replication.replicated_frame
are structurally identical. A future consolidation could merge them
into a shared frame type, but the two modules have no common dependency
today and the duplication is small enough to live with. *)(** Capture the committed WAL frames since a given watermark position,
returning them as a list of {!backup_frame}.
[~since_epoch] and [~since_idx] identify the watermark: frames with
indices strictly greater than [since_idx] in the current epoch are
returned. If the WAL's epoch has advanced past [since_epoch], no
frames can be captured (the caller must take a fresh base snapshot).
Returns [None] when the WAL's epoch has changed (meaning the caller's
watermark is stale and a re-base is needed). Returns [Some []] when
the watermark is current but no new frames have been committed. *)letcapture_frames_since(t:t)~since_epoch~since_idx:(backup_framelist,[>`Capture_errorofstring])resultoptionLwt.t=matcht.backendwith|Mem_->Some(Ok[])|>Lwt.return|Btreest->(matchst.walwith|None->Lwt.return(Some(Error(`Capture_error"no WAL active")))|Somewal->letcurrent_epoch=Wal.epochwalinifnot(Int64.equalcurrent_epochsince_epoch)then(* Epoch changed: the watermark is stale and the caller must
re-base (take a fresh full snapshot). *)Lwt.returnNoneelse(letcommitted=Wal.committed_frameswalin(* max_int = "no floor" sentinel; increment would overflow to min_int *)letstart=ifsince_idx=max_intthenmax_intelsesince_idx+1inifstart>=committedthenLwt.return(Some(Ok[]))else(letsalt=Wal.saltwalinletseed=Wal.seedwalinletrecloopidxacc=ifidx>=committedthenLwt.return(Some(Ok(List.revacc)))else((* A concurrent checkpoint can bump the epoch while we yield
on I/O. If the epoch changed, the watermark is stale —
signal through [None] so the caller re-bases cleanly
instead of getting an I/O error. *)letcurrent_epoch=Wal.epochwalinifnot(Int64.equalcurrent_epochsince_epoch)thenLwt.returnNoneelselet*r=Wal.read_committed_framewalidxinmatchrwith|Error_whennot(Int64.equal(Wal.epochwal)since_epoch)->Lwt.returnNone|Errore->Lwt.return(Some(Error(`Capture_error(Format.asprintf"read frame %d: %a"idxWal.pp_errore))))|Okf->ifnot(Int64.equal(Wal.epochwal)since_epoch)thenLwt.returnNoneelse(letflags=iff.is_committhen1Lelse0Linletchecksum=Wal.frame_checksum~salt~seed~page_id:f.page_id~flags~page:f.pageinletbf:backup_frame={epoch=current_epoch;frame_idx=idx;page_id=f.page_id;is_commit=f.is_commit;page=f.page;checksum;source_salt=salt;source_seed=seed}inloop(idx+1)(bf::acc)))in(* loop already returns the exact type of this branch —
None for epoch-changed, Some (Ok frames) for success,
Some (Error _) for I/O failure. Direct return. *)loopstart[])));;(** Install an asynchronous callback invoked after each WAL commit batch.
The callback receives ~epoch, ~base_idx (starting WAL frame index),
and ~count (number of committed frames). Fired via [Lwt.async] so
the commit path is never blocked by replication I/O.
When a callback is registered, the replication shipped-position
floor is initialised to the WAL's current [committed_frames] so
that checkpoint cannot recycle already-acknowledged frames before
the async sink ships its first batch. The consumer must still
call {!update_replication_position} to advance the floor as
frames are shipped.
Pass [None] to unregister (resets the floor to [max_int]). *)letset_commit_callback(t:t)(cb:(epoch:int64->base_idx:int->count:int->unitLwt.t)option)=matcht.backendwith|Mem_->Lwt.return_unit|Btreest->(matchcbwith|None->st.on_committed_frames<-None;st.replication_shipped_frames<-max_int;Lwt.return_unit|Some_->(* #336/1 + review #1: do the flush AND the pin atomically under the
write lock. [flush_unsynced] yields (group-commit drain); without the
lock a concurrent commit could interleave in the OLD mode (no fsync,
no ship — cb not yet set), and the subsequent
[sink_shipped_frames <- committed_frames] pin would then bury those
frames below the ship cursor forever (silent standby divergence). The
lock blocks new commits ([rw_begin]) for the brief flush+pin so the
cursor pins exactly the pre-registration frontier. *)let*()=Rwlock.acquire_writet.lockinLwt.finalize(fun()->(* Flush while still in the old mode — [flush_unsynced] is a no-op
once we pin [Full] below. A store opened [off]/[batched] may have
acked commits in the OS page cache; the sink ships only NEW
frames, so these historical frames would otherwise linger
crash-exposed until the next commit/checkpoint/close despite the
sink implying synchronous=full. *)let*()=flush_unsyncedtinst.on_committed_frames<-cb;(* #298: a replication commit-sink requires Full durability — the
checkpoint replica-floor gate assumes every committed frame is
shipped, which only holds when every commit fsyncs. Force Full on
registration; relaxing durability is rejected while a sink is
active (see set_durability / the PRAGMA handler). *)st.sync_mode<-`Full;(matchst.walwith|None->()|Somewal->st.replication_shipped_frames<-Wal.committed_frameswal;(* #298/#1: a sink registered mid-life ships only NEW synced
frames, not history — start the ship cursor at the current
count. *)st.sink_shipped_frames<-Wal.committed_frameswal);Lwt.return_unit)(fun()->Rwlock.release_writet.lock;Lwt.return_unit));;(* #384: map a storage-level [Pager_event.t] to a [Store_event.t], stamping the
txn id from the pager. Write-path events (alloc/write/free) always fire
inside an active RW txn, so they carry the exact id; [Page_read] may fire
outside a write txn, where [get_txn_id] returns 0 before the first txn and the
most-recent txn id between txns (best-effort).
#385: also stamps the [tree] id from [st.current_tree], set at the
[bt_get_tree]/[bt_get_tree_ro] chokepoint by the op that triggered the I/O.
Meta/system pages read while resolving a tree handle are stamped [tree = -1]
(per #174, system pages are never tagged with a user tree); a tree's own
data-page [Page_read]/[Page_alloc]/[Page_free] carry the exact tree id.
[Page_write] is best-effort: emitted at WAL-flush time, it carries whichever
tree was most recently active. *)lettranslate_pager_event(st:bt_state)(pev:Pager_event.t):Store_event.t=lettxn_id=Pager.get_txn_idst.pagerinlettree=Option.valuest.current_tree~default:(-1)inmatchpevwith|Pager_event.Page_read{page_id}->Store_event.Page_read{txn_id;tree;page=page_id}|Pager_event.Wal_read{page_id}->Store_event.Wal_read{txn_id;tree;page=page_id}|Pager_event.Page_write{page_id}->Store_event.Page_write{txn_id;tree;page=page_id}|Pager_event.Page_alloc{page_id;reused}->Store_event.Page_alloc{txn_id;tree;page=page_id;reused}|Pager_event.Page_free{page_id}->Store_event.Page_free{txn_id;tree;page=page_id};;letset_event_callback(t:t)(cb:(Store_event.t->unit)option)=matcht.backendwith|Mem_->()(* Mem backend has no bt_state; emits no events (#382). *)|Btreest->st.on_event<-cb;(matchcbwith|None->Pager.set_page_event_callbackst.pagerNone|Somef->Pager.set_page_event_callbackst.pager(Some(funpev->(* Same guarantee as [emit_event]: a faulty observer must never
break a transaction. *)tryf(translate_pager_eventstpev)with|_->())));;moduleEvent=Store_event(* ------------------------------------------------------------------ *)(* Follower mode (#172) *)(* ------------------------------------------------------------------ *)(** Enable or disable follower mode on the store. When [true],
[rw_begin] rejects write transactions so the standby's WAL does not
diverge from the master's stream. No-op on the in-memory backend. *)letset_follower(t:t)(on:bool)=matcht.backendwith|Mem_->()|Btreest->st.follower<-on;ifnotonthenst.follower_ack_position<-None;;(** True iff follower mode is active (writes are rejected). *)letis_follower(t:t)=matcht.backendwith|Mem_->false|Btreest->st.follower;;(** Record a local [Wal.committed_frames] count as the follower's last-applied
commit boundary. [ro_begin] will cap RO snapshots to this position so
readers never observe WAL frames past what has been applied on this node
(#263). The caller should supply the count from the WAL handle it applied
into, so the value lives in local committed-frame count space (no coordinate
mismatch vs. master epoch indices) without depending on WAL instance identity
between the caller and the store. No-op on the in-memory backend. *)letset_follower_ack_position(t:t)~(frames:int)=matcht.backendwith|Mem_->()|Btreest->st.follower_ack_position<-Someframes;;(** Get the recorded follower ack position (a local [Wal.committed_frames]
count), or [None] if not following or no position has been recorded yet. *)letfollower_ack_position(t:t)=matcht.backendwith|Mem_->None|Btreest->st.follower_ack_position;;letwait_for_readers_past(t:t)~target=matcht.backendwith|Mem_->Lwt.return_unit|Btreest->wait_for_readers_pastst~target~replication_max_yields:st.replication_gate_max_yields~backup_max_yields:st.backup_gate_max_yields;;[@@@ai_disclosure"ai-generated"][@@@ai_model"claude-opus-4-7"][@@@ai_provider"Anthropic"]