atproto pds in zig
Something went wrong. Try again.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707270827092710271127122713271427152716271727182719272027212722272327242725272627272728272927302731273227332734273527362737273827392740274127422743274427452746274727482749275027512752275327542755275627572758275927602761276227632764276527662767276827692770277127722773277427752776277727782779278027812782278327842785278627872788278927902791279227932794279527962797279827992800280128022803280428052806280728082809281028112812281328142815281628172818281928202821282228232824282528262827282828292830283128322833283428352836283728382839284028412842284328442845284628472848284928502851285228532854285528562857285828592860286128622863286428652866286728682869287028712872287328742875287628772878287928802881288228832884288528862887288828892890289128922893289428952896289728982899290029012902290329042905290629072908290929102911291229132914291529162917291829192920292129222923292429252926292729282929293029312932293329342935293629372938293929402941294229432944294529462947294829492950295129522953295429552956295729582959296029612962296329642965296629672968296929702971297229732974297529762977297829792980298129822983298429852986298729882989299029912992299329942995299629972998299930003001300230033004300530063007300830093010301130123013301430153016301730183019302030213022302330243025302630273028302930303031303230333034303530363037303830393040304130423043304430453046304730483049305030513052305330543055305630573058305930603061306230633064306530663067306830693070307130723073307430753076307730783079308030813082308330843085308630873088308930903091309230933094309530963097309830993100310131023103310431053106310731083109311031113112311331143115311631173118311931203121312231233124312531263127312831293130313131323133313431353136313731383139314031413142314331443145314631473148314931503151315231533154315531563157315831593160316131623163316431653166316731683169317031713172317331743175317631773178317931803181318231833184318531863187318831893190319131923193319431953196319731983199320032013202320332043205320632073208320932103211321232133214321532163217321832193220322132223223322432253226322732283229323032313232323332343235323632373238323932403241324232433244324532463247324832493250325132523253325432553256325732583259326032613262326332643265326632673268326932703271327232733274327532763277327832793280328132823283328432853286328732883289329032913292329332943295329632973298329933003301330233033304330533063307330833093310331133123313331433153316331733183319332033213322332333243325332633273328332933303331333233333334333533363337333833393340334133423343334433453346334733483349335033513352335333543355335633573358335933603361336233633364336533663367336833693370337133723373337433753376337733783379338033813382338333843385338633873388338933903391339233933394339533963397339833993400340134023403340434053406340734083409341034113412341334143415341634173418341934203421342234233424342534263427342834293430343134323433343434353436343734383439344034413442344334443445344634473448344934503451345234533454345534563457345834593460346134623463346434653466346734683469347034713472347334743475347634773478347934803481348234833484348534863487348834893490349134923493349434953496349734983499350035013502350335043505350635073508350935103511351235133514351535163517351835193520352135223523352435253526352735283529353035313532353335343535353635373538353935403541354235433544354535463547354835493550355135523553355435553556355735583559356035613562356335643565356635673568356935703571357235733574357535763577357835793580358135823583358435853586358735883589359035913592359335943595359635973598359936003601360236033604360536063607360836093610361136123613361436153616361736183619362036213622362336243625362636273628362936303631363236333634363536363637363836393640364136423643364436453646364736483649365036513652365336543655365636573658365936603661366236633664366536663667366836693670367136723673367436753676367736783679368036813682368336843685368636873688368936903691369236933694369536963697369836993700370137023703370437053706370737083709371037113712371337143715371637173718371937203721372237233724372537263727372837293730373137323733373437353736373737383739374037413742374337443745374637473748374937503751375237533754375537563757375837593760376137623763376437653766376737683769377037713772377337743775377637773778377937803781378237833784378537863787378837893790379137923793379437953796379737983799380038013802380338043805380638073808380938103811381238133814381538163817381838193820382138223823382438253826382738283829383038313832383338343835383638373838383938403841384238433844384538463847384838493850385138523853385438553856385738583859386038613862386338643865386638673868386938703871387238733874387538763877387838793880388138823883388438853886388738883889389038913892389338943895389638973898389939003901390239033904390539063907390839093910391139123913391439153916391739183919392039213922392339243925392639273928392939303931393239333934393539363937393839393940394139423943394439453946394739483949395039513952395339543955395639573958395939603961396239633964396539663967396839693970397139723973397439753976397739783979398039813982398339843985398639873988398939903991399239933994399539963997399839994000400140024003400440054006400740084009401040114012401340144015401640174018401940204021402240234024402540264027402840294030403140324033403440354036403740384039404040414042404340444045404640474048404940504051405240534054405540564057405840594060406140624063406440654066406740684069407040714072407340744075407640774078407940804081408240834084408540864087408840894090409140924093409440954096409740984099410041014102410341044105410641074108410941104111411241134114411541164117411841194120412141224123412441254126412741284129413041314132413341344135413641374138413941404141414241434144414541464147414841494150415141524153415441554156415741584159416041614162416341644165416641674168416941704171417241734174417541764177417841794180418141824183418441854186418741884189419041914192419341944195419641974198419942004201420242034204420542064207420842094210421142124213421442154216421742184219422042214222422342244225422642274228422942304231423242334234423542364237423842394240424142424243424442454246424742484249425042514252425342544255425642574258425942604261426242634264426542664267426842694270427142724273427442754276427742784279428042814282428342844285428642874288428942904291429242934294429542964297429842994300430143024303430443054306430743084309431043114312431343144315431643174318431943204321432243234324432543264327432843294330433143324333433443354336433743384339434043414342434343444345434643474348434943504351435243534354435543564357435843594360436143624363436443654366436743684369437043714372437343744375437643774378437943804381438243834384438543864387438843894390439143924393439443954396439743984399440044014402440344044405440644074408440944104411441244134414441544164417441844194420442144224423442444254426442744284429443044314432443344344435443644374438443944404441444244434444444544464447444844494450445144524453445444554456445744584459446044614462446344644465446644674468446944704471447244734474447544764477447844794480448144824483448444854486448744884489449044914492449344944495449644974498449945004501450245034504450545064507450845094510451145124513451445154516451745184519452045214522452345244525452645274528452945304531453245334534453545364537453845394540454145424543454445454546454745484549455045514552455345544555455645574558455945604561456245634564456545664567456845694570457145724573457445754576457745784579458045814582458345844585458645874588458945904591459245934594459545964597459845994600460146024603460446054606460746084609461046114612461346144615461646174618461946204621462246234624462546264627462846294630463146324633463446354636463746384639464046414642464346444645464646474648464946504651465246534654465546564657465846594660466146624663466446654666466746684669467046714672467346744675467646774678467946804681468246834684468546864687468846894690469146924693469446954696469746984699470047014702470347044705470647074708470947104711471247134714471547164717471847194720472147224723472447254726472747284729473047314732473347344735473647374738473947404741474247434744474547464747474847494750475147524753475447554756475747584759476047614762476347644765476647674768476947704771477247734774477547764777477847794780478147824783478447854786478747884789479047914792479347944795479647974798479948004801480248034804480548064807480848094810481148124813481448154816481748184819482048214822482348244825482648274828482948304831483248334834483548364837483848394840484148424843484448454846484748484849485048514852485348544855485648574858485948604861486248634864486548664867486848694870487148724873487448754876487748784879488048814882488348844885488648874888488948904891489248934894489548964897489848994900490149024903490449054906490749084909491049114912491349144915491649174918491949204921492249234924492549264927492849294930493149324933493449354936493749384939494049414942494349444945494649474948494949504951495249534954495549564957495849594960496149624963496449654966496749684969497049714972497349744975497649774978497949804981498249834984498549864987498849894990499149924993499449954996499749984999500050015002500350045005500650075008500950105011501250135014501550165017501850195020502150225023502450255026502750285029503050315032503350345035503650375038503950405041504250435044504550465047504850495050505150525053505450555056505750585059506050615062506350645065506650675068506950705071507250735074507550765077507850795080508150825083508450855086508750885089509050915092509350945095509650975098509951005101510251035104510551065107510851095110511151125113511451155116511751185119512051215122512351245125512651275128512951305131513251335134513551365137513851395140514151425143514451455146514751485149515051515152515351545155515651575158515951605161516251635164516551665167516851695170517151725173517451755176517751785179518051815182518351845185518651875188518951905191519251935194519551965197519851995200520152025203520452055206520752085209521052115212521352145215521652175218521952205221522252235224522552265227522852295230523152325233523452355236523752385239524052415242524352445245524652475248524952505251525252535254525552565257525852595260526152625263526452655266526752685269527052715272527352745275527652775278527952805281528252835284528552865287528852895290529152925293529452955296529752985299530053015302530353045305530653075308530953105311531253135314531553165317531853195320532153225323532453255326532753285329533053315332533353345335533653375338533953405341534253435344534553465347534853495350535153525353535453555356535753585359536053615362536353645365536653675368536953705371537253735374537553765377537853795380538153825383538453855386538753885389539053915392539353945395539653975398539954005401540254035404540554065407540854095410541154125413541454155416541754185419542054215422542354245425542654275428542954305431543254335434543554365437543854395440544154425443544454455446544754485449545054515452545354545455545654575458545954605461546254635464546554665467546854695470547154725473547454755476547754785479548054815482548354845485548654875488548954905491549254935494549554965497549854995500550155025503550455055506550755085509551055115512551355145515551655175518551955205521552255235524552555265527552855295530553155325533553455355536553755385539554055415542554355445545554655475548554955505551555255535554555555565557555855595560556155625563556455655566556755685569557055715572557355745575557655775578557955805581558255835584558555865587558855895590559155925593559455955596559755985599560056015602560356045605560656075608560956105611561256135614561556165617561856195620562156225623562456255626562756285629563056315632563356345635563656375638563956405641564256435644564556465647564856495650565156525653565456555656565756585659566056615662566356645665566656675668566956705671567256735674567556765677567856795680568156825683568456855686568756885689569056915692569356945695569656975698569957005701570257035704570557065707570857095710571157125713571457155716571757185719572057215722572357245725572657275728572957305731573257335734573557365737573857395740574157425743574457455746574757485749575057515752575357545755575657575758575957605761576257635764576557665767576857695770577157725773577457755776577757785779578057815782578357845785578657875788578957905791579257935794579557965797579857995800580158025803580458055806580758085809581058115812581358145815581658175818581958205821582258235824582558265827582858295830583158325833583458355836583758385839584058415842584358445845584658475848584958505851585258535854585558565857585858595860586158625863586458655866586758685869587058715872587358745875587658775878587958805881588258835884588558865887588858895890589158925893589458955896589758985899590059015902590359045905590659075908590959105911591259135914591559165917591859195920592159225923592459255926592759285929593059315932593359345935593659375938593959405941594259435944594559465947594859495950595159525953595459555956595759585959596059615962596359645965596659675968596959705971597259735974597559765977597859795980598159825983598459855986598759885989599059915992599359945995599659975998599960006001600260036004600560066007600860096010601160126013601460156016601760186019602060216022602360246025602660276028602960306031603260336034603560366037603860396040604160426043604460456046604760486049605060516052605360546055605660576058605960606061606260636064606560666067606860696070607160726073607460756076607760786079608060816082608360846085608660876088608960906091609260936094609560966097609860996100610161026103610461056106610761086109611061116112611361146115611661176118611961206121612261236124612561266127612861296130613161326133613461356136613761386139614061416142614361446145614661476148614961506151615261536154615561566157615861596160616161626163616461656166616761686169617061716172617361746175617661776178617961806181618261836184618561866187618861896190619161926193619461956196619761986199620062016202620362046205620662076208620962106211621262136214621562166217621862196220622162226223622462256226622762286229623062316232623362346235623662376238623962406241624262436244624562466247624862496250625162526253625462556256625762586259626062616262626362646265626662676268626962706271627262736274627562766277627862796280628162826283628462856286628762886289629062916292629362946295629662976298629963006301630263036304630563066307630863096310631163126313631463156316631763186319632063216322632363246325632663276328632963306331633263336334633563366337633863396340634163426343634463456346634763486349635063516352635363546355635663576358635963606361636263636364636563666367636863696370637163726373637463756376637763786379638063816382638363846385638663876388638963906391639263936394639563966397639863996400640164026403640464056406640764086409641064116412641364146415641664176418641964206421642264236424642564266427642864296430643164326433643464356436643764386439644064416442644364446445644664476448644964506451645264536454645564566457645864596460646164626463646464656466646764686469647064716472647364746475647664776478647964806481648264836484648564866487648864896490649164926493649464956496649764986499650065016502650365046505650665076508650965106511651265136514651565166517651865196520652165226523652465256526652765286529653065316532653365346535653665376538653965406541654265436544654565466547654865496550655165526553655465556556655765586559656065616562656365646565656665676568656965706571657265736574657565766577657865796580658165826583658465856586658765886589659065916592659365946595659665976598659966006601660266036604660566066607660866096610661166126613661466156616661766186619662066216622662366246625662666276628662966306631663266336634663566366637663866396640664166426643664466456646664766486649665066516652665366546655665666576658665966606661666266636664666566666667666866696670667166726673667466756676667766786679668066816682668366846685668666876688668966906691669266936694669566966697669866996700670167026703670467056706670767086709671067116712671367146715671667176718671967206721672267236724672567266727672867296730673167326733673467356736673767386739674067416742674367446745674667476748674967506751675267536754675567566757675867596760676167626763676467656766676767686769677067716772677367746775677667776778677967806781678267836784678567866787678867896790679167926793679467956796679767986799680068016802680368046805680668076808680968106811681268136814681568166817681868196820682168226823682468256826682768286829683068316832683368346835683668376838683968406841684268436844684568466847684868496850685168526853685468556856685768586859686068616862686368646865686668676868686968706871687268736874687568766877687868796880688168826883688468856886688768886889689068916892689368946895689668976898689969006901690269036904690569066907690869096910691169126913691469156916691769186919692069216922692369246925692669276928692969306931693269336934693569366937693869396940694169426943694469456946694769486949695069516952695369546955695669576958695969606961696269636964696569666967696869696970697169726973697469756976697769786979698069816982698369846985698669876988698969906991699269936994699569966997699869997000700170027003700470057006700770087009701070117012701370147015701670177018701970207021702270237024702570267027702870297030703170327033703470357036703770387039704070417042704370447045704670477048704970507051705270537054705570567057705870597060706170627063706470657066706770687069707070717072707370747075707670777078707970807081708270837084708570867087708870897090709170927093709470957096709770987099710071017102710371047105710671077108710971107111711271137114711571167117711871197120712171227123712471257126712771287129713071317132713371347135713671377138713971407141714271437144714571467147714871497150715171527153715471557156715771587159716071617162716371647165716671677168716971707171717271737174717571767177717871797180718171827183718471857186718771887189719071917192719371947195719671977198719972007201720272037204720572067207720872097210721172127213721472157216721772187219722072217222722372247225722672277228722972307231723272337234723572367237723872397240724172427243724472457246724772487249725072517252725372547255725672577258725972607261726272637264726572667267726872697270727172727273727472757276727772787279728072817282728372847285728672877288728972907291729272937294729572967297729872997300730173027303730473057306730773087309731073117312731373147315731673177318731973207321732273237324732573267327732873297330733173327333733473357336733773387339734073417342734373447345734673477348734973507351735273537354735573567357735873597360736173627363736473657366736773687369737073717372737373747375737673777378737973807381738273837384738573867387738873897390739173927393739473957396739773987399740074017402740374047405740674077408740974107411741274137414741574167417741874197420742174227423742474257426742774287429743074317432743374347435743674377438743974407441744274437444744574467447744874497450745174527453745474557456745774587459746074617462746374647465746674677468746974707471747274737474747574767477747874797480748174827483748474857486748774887489749074917492749374947495749674977498749975007501750275037504750575067507750875097510751175127513751475157516751775187519752075217522752375247525752675277528752975307531753275337534753575367537753875397540754175427543754475457546754775487549755075517552755375547555755675577558755975607561756275637564756575667567756875697570757175727573757475757576757775787579758075817582758375847585758675877588758975907591759275937594759575967597759875997600760176027603760476057606760776087609761076117612761376147615761676177618761976207621762276237624762576267627762876297630763176327633763476357636763776387639764076417642764376447645764676477648764976507651765276537654765576567657765876597660766176627663766476657666766776687669const std = @import("std");const atid = @import("../core/atid.zig");const clock = @import("../core/clock.zig");const auth = @import("../auth/tokens.zig");const blobstore = @import("blobstore.zig");const cbor_json = @import("../internal/cbor_json.zig");const eventlog = @import("eventlog.zig");const permissioned = @import("../internal/permissioned_data.zig");const sharded_locks = @import("../internal/sharded_locks.zig");const space_uri = @import("../internal/space_uri.zig");const zat = @import("zat");const zqlite = @import("zqlite");const Io = std.Io;
pub const Error = error{ InvalidCollection, InvalidCursor, InvalidRecordKey, InvalidRecordType, InvalidSwap, ValidationRequired, InvalidReservedSigningKey, MissingRecord, MissingReservedSigningKey, MissingRecordBlock, RepoNotFound, InvalidRepoPath, InvalidDagCbor, InvalidRefreshSession, InvalidAccountStatus, AccountNotFound, HandleNotAvailable, RateLimitExceeded, StoreNotInitialized,};
pub const AccountStatus = enum { active, takendown, suspended, deactivated, deleted,
pub fn asString(self: AccountStatus) []const u8 { return switch (self) { .active => "active", .takendown => "takendown", .suspended => "suspended", .deactivated => "deactivated", .deleted => "deleted", }; }
pub fn parse(raw: []const u8) AccountStatus { if (std.mem.eql(u8, raw, "takendown")) return .takendown; if (std.mem.eql(u8, raw, "suspended")) return .suspended; if (std.mem.eql(u8, raw, "deactivated")) return .deactivated; if (std.mem.eql(u8, raw, "deleted")) return .deleted; return .active; }
pub fn isActive(self: AccountStatus) bool { return self == .active; }};
pub const Record = struct { did: []const u8, collection: []const u8, rkey: []const u8, cid: []const u8, value_json: []const u8, validation_status: []const u8, rev: []const u8, seq: u64,
pub fn uri(self: Record, allocator: std.mem.Allocator) ![]const u8 { return std.fmt.allocPrint(allocator, "at://{s}/{s}/{s}", .{ self.did, self.collection, self.rkey }); }};
pub const BlobRecord = struct { mime_type: []const u8, data: []const u8,};
pub const SeqEvent = struct { seq: u64, frame: []const u8,};
pub const EmailInfo = struct { email: []const u8, email_confirmed: bool, auth_code: ?[]const u8, auth_code_expires_at: ?i64, pending_email: ?[]const u8,};
pub const CodeStatus = enum { valid, invalid, expired,};
pub const OAuthRequest = struct { request_id: []const u8, client_id: []const u8, redirect_uri: []const u8, scope: []const u8, state: []const u8, code_challenge: []const u8, code_challenge_method: []const u8, response_mode: []const u8, login_hint: ?[]const u8, dpop_jkt: ?[]const u8, expires_at: i64, sub: ?[]const u8, code: ?[]const u8, auth_method: ?[]const u8,};
pub const OAuthToken = struct { pub const Kind = enum { access, refresh, previous_refresh, used_refresh };
family_id: []const u8, did: []const u8, client_id: []const u8, scope: []const u8, access_token: []const u8, refresh_token: []const u8, access_expires_at: i64, refresh_expires_at: i64, previous_refresh_token: ?[]const u8, previous_refresh_expires_at: ?i64, revoked: bool, dpop_jkt: ?[]const u8, auth_method: ?[]const u8, kind: Kind,};
pub const SessionInfo = struct { id: []const u8, did: []const u8, handle: []const u8, auth_method: []const u8, app_password_name: ?[]const u8, controller_did: ?[]const u8, created_at: i64, updated_at: i64, last_used_at: ?i64, access_expires_at: i64, refresh_expires_at: i64, revoked_at: ?i64, active: bool,};
pub const OAuthGrantInfo = struct { did: []const u8, handle: []const u8, client_id: []const u8, scope: []const u8, created_at: i64, expires_at: i64, revoked_at: ?i64, active: bool, auth_method: ?[]const u8,};
pub const SessionTokenRow = struct { id: []const u8, did: []const u8, access_jti: []const u8, refresh_jti: []const u8, auth_method: []const u8, controller_did: ?[]const u8, created_at: i64, updated_at: i64, revoked_at: ?i64,};
pub const AppPassword = struct { name: []const u8, password_hash: []const u8, created_at: i64, privileged: bool, scopes: ?[]const u8, created_by_controller_did: ?[]const u8,};
pub const AuditEvent = struct { id: []const u8, subject_did: []const u8, actor_did: []const u8, controller_did: ?[]const u8, action: []const u8, details_json: []const u8, created_at: i64,};
pub const Passkey = struct { id: []const u8, did: []const u8, credential_id: []const u8, public_key: []const u8, sign_count: u32, friendly_name: ?[]const u8, created_at: i64, last_used: ?i64,};
pub const WebAuthnChallenge = struct { id: []const u8, did: []const u8, challenge: []const u8, kind: []const u8, state_json: []const u8, expires_at: i64,};
pub const DiscoverableWebAuthnChallenge = struct { request_id: []const u8, challenge: []const u8, expires_at: i64,};
pub const InviteCode = struct { code: []const u8, available: i64, disabled: bool, for_account: []const u8, created_by: []const u8, created_at: i64, uses: []InviteCodeUse,};
pub const ReservedSigningKey = struct { did: ?[]const u8, signing_key: []const u8, secret_key: [32]u8, expires_at: i64,};
pub const InviteCodeUse = struct { used_by: []const u8, used_at: i64,};
pub const AppPreference = struct { name: []const u8, value_json: []const u8,};
pub const Resident = struct { handle: []const u8, did: []const u8, active: bool, rev: []const u8, record_count: u64,};
pub const CollectionSummary = struct { collection: []const u8, count: u64,};
pub const ActivityBucket = struct { start_us: u64, end_us: u64, count: u64,};
pub const RecentRecord = struct { handle: []const u8, did: []const u8, collection: []const u8, rkey: []const u8, cid: []const u8, seq: u64,};
pub const ImportedRecord = struct { collection: []const u8, rkey: []const u8, cid: []const u8, blob_cids: []const []const u8,};
pub const ImportedBlock = struct { cid: []const u8, data: []const u8,};
pub const CommitInfo = struct { cid: []const u8, rev: []const u8,};
pub const WriteResult = struct { commit: CommitInfo, records: []Record,};
pub const SpaceConfig = struct { uri: []const u8, authority_did: []const u8, space_type: []const u8, skey: []const u8, managing_app: ?[]const u8, policy: []const u8, app_access_json: []const u8, is_authority: bool, deleted_at: ?i64,};
pub const SimpleSpaceMember = struct { did: []const u8, created_at: i64,};
pub const SpaceRecord = struct { space: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8, cid: []const u8, value_json: []const u8, validation_status: []const u8, repo_rev: []const u8, updated_at: i64,
pub fn uri(self: SpaceRecord, allocator: std.mem.Allocator) ![]const u8 { return space_uri.SpaceUri.formatRecord(allocator, self.space, self.repo_did, self.collection, self.rkey); }};
pub const SpaceRecordRef = struct { collection: []const u8, rkey: []const u8, cid: []const u8, value_json: ?[]const u8 = null,};
pub const SpaceRecordOplogEntry = struct { rev: []const u8, idx: i64, action: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8, cid: ?[]const u8, prev: ?[]const u8, value_json: ?[]const u8,};
pub const SpaceState = struct { set_hash: ?[]const u8, rev: ?[]const u8,};
pub const SpaceWriterState = struct { repo_did: []const u8, hash: []const u8, rev: []const u8,};
pub const CredentialRecipient = struct { service_endpoint: []const u8, repo_did: ?[]const u8, expires_at: i64,};
pub const PreparedRecord = struct { cid: []const u8, value_json: []const u8, validation_status: []const u8, blob_cids: []const []const u8,};
pub const SpaceWriteOp = union(enum) { create: struct { collection: []const u8, rkey: []const u8, prepared: PreparedRecord }, put: struct { collection: []const u8, rkey: []const u8, prepared: PreparedRecord }, update: struct { collection: []const u8, rkey: []const u8, prepared: PreparedRecord }, delete: struct { collection: []const u8, rkey: []const u8 },};
pub const SpaceWriteResult = union(enum) { create: SpaceRecord, update: SpaceRecord, delete: void,};
pub const WriteProfile = struct { validation_ns: u64 = 0, lock_wait_ns: u64 = 0, load_repo_ns: u64 = 0, stage_records_ns: u64 = 0, build_commit_ns: u64 = 0, sql_ns: u64 = 0, event_publish_ns: u64 = 0, total_ns: u64 = 0,};
pub const WriteOp = union(enum) { create: struct { collection: []const u8, rkey: ?[]const u8, value: std.json.Value, }, update: struct { collection: []const u8, rkey: []const u8, value: std.json.Value, swap: RecordSwap = .none, }, delete: struct { collection: []const u8, rkey: []const u8, swap: RecordSwap = .none, },};
pub const RecordSwap = union(enum) { none, missing, cid: []const u8,};
pub const ValidationMode = enum { known, require, skip,};
pub const WriteOptions = struct { swap_commit: ?[]const u8 = null, validate: ValidationMode = .known,};
const BlobRef = struct { cid: []const u8, uri: []const u8,};
const RepoBlockReader = struct { allocator: std.mem.Allocator, did: []const u8, locked: bool = false,
fn reader(self: *RepoBlockReader) zat.mst.BlockReader { return .{ .ctx = self, .getFn = getBlock, }; }
fn getBlock(ctx: *anyopaque, cid_raw: []const u8) anyerror!?[]const u8 { const self: *RepoBlockReader = @ptrCast(@alignCast(ctx)); const cid = try cidText(self.allocator, cid_raw);
if (self.locked) { return repoBlockDataLocked(self.allocator, self.did, cid); }
db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
return repoBlockDataLocked(self.allocator, self.did, cid); }};
fn repoBlockDataLocked(allocator: std.mem.Allocator, did: []const u8, cid: []const u8) !?[]const u8 { const row = try conn.row( \\SELECT data \\FROM repo_blocks \\WHERE did = ? AND cid = ? \\LIMIT 1 , .{ did, cid }); if (row == null) return null; defer row.?.deinit();
const data = row.?.nullableBlob(0) orelse return null; return try allocator.dupe(u8, data);}
var conn: zqlite.Conn = undefined;var initialized = false;var store_io: Io = undefined;var db_mutex: Io.Mutex = .init;var write_lanes: sharded_locks.ShardedLocks(32) = .{};var next_seq: std.atomic.Value(u64) = .init(1);
pub fn init(io: Io, path: []const u8) !void { if (initialized) return; store_io = io; if (!std.mem.eql(u8, path, ":memory:") and !std.fs.path.isAbsolute(path)) { if (std.fs.path.dirname(path)) |dir| try Io.Dir.createDirPath(.cwd(), io, dir); }
const path_z = try std.heap.page_allocator.dupeZ(u8, path); conn = try zqlite.open(path_z.ptr, zqlite.OpenFlags.Create | zqlite.OpenFlags.ReadWrite); initialized = true; errdefer close();
try conn.busyTimeout(5000); try conn.execNoArgs("PRAGMA journal_mode=WAL"); try conn.execNoArgs("PRAGMA foreign_keys=ON"); try migrate(); next_seq.store(try loadNextSeqLocked(), .release);}
pub fn close() void { if (!initialized) return; conn.close(); initialized = false;}
pub fn nowMs() i64 { return clock.nowMillis();}
pub fn randomBytes(buffer: []u8) void { store_io.random(buffer);}
pub fn randomToken(allocator: std.mem.Allocator, prefix: []const u8, comptime byte_len: usize) ![]const u8 { var bytes: [byte_len]u8 = undefined; randomBytes(&bytes); const hex = std.fmt.bytesToHex(bytes, .lower); return std.fmt.allocPrint(allocator, "{s}{s}", .{ prefix, &hex });}
pub fn currentIo() Io { return store_io;}
pub fn resolveRepo(repo: []const u8) ?auth.Account { return findAccount(std.heap.page_allocator, repo) catch null;}
pub fn findAccount(allocator: std.mem.Allocator, identifier: []const u8) !?auth.Account { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
return try findAccountLocked(allocator, identifier);}
pub fn findActiveAccount(allocator: std.mem.Allocator, identifier: []const u8) !?auth.Account { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const row = try conn.row( \\SELECT handle, did, email, password_hash \\FROM accounts \\WHERE account_status = 'active' \\ AND (lower(handle) = lower(?) OR did = ? OR lower(email) = lower(?)) \\LIMIT 1 , .{ identifier, identifier, identifier }); if (row == null) return null; defer row.?.deinit(); return try accountFromRow(allocator, row.?);}
pub fn consumeRateLimit(subject: []const u8, action: []const u8, now_ms: i64, window_ms: i64, limit: usize) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
try conn.execNoArgs("BEGIN IMMEDIATE"); errdefer conn.execNoArgs("ROLLBACK") catch {}; try conn.exec( \\DELETE FROM rate_limit_events \\WHERE action = ? AND occurred_at < ? , .{ action, now_ms - window_ms }); const row = try conn.row( \\SELECT COUNT(*) \\FROM rate_limit_events \\WHERE subject = ? AND action = ? AND occurred_at >= ? , .{ subject, action, now_ms - window_ms }); defer if (row) |value| value.deinit(); if (row != null and row.?.int(0) >= @as(i64, @intCast(limit))) { try conn.execNoArgs("ROLLBACK"); return Error.RateLimitExceeded; } try conn.exec( \\INSERT INTO rate_limit_events (subject, action, occurred_at) \\VALUES (?, ?, ?) , .{ subject, action, now_ms }); try conn.execNoArgs("COMMIT");}
pub fn searchAccounts(allocator: std.mem.Allocator, query: []const u8, limit: usize) ![]auth.Account { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const actual_limit = if (limit == 0) 25 else limit; const pattern = try std.fmt.allocPrint(allocator, "%{s}%", .{query}); var rows = if (query.len == 0) try conn.rows( \\SELECT handle, did, email, password_hash \\FROM accounts \\ORDER BY handle \\LIMIT ? , .{@as(i64, @intCast(actual_limit))}) else try conn.rows( \\SELECT handle, did, email, password_hash \\FROM accounts \\WHERE lower(handle) LIKE lower(?) OR did LIKE ? \\ORDER BY handle \\LIMIT ? , .{ pattern, pattern, @as(i64, @intCast(actual_limit)) }); defer rows.deinit();
var accounts: std.ArrayList(auth.Account) = .empty; while (rows.next()) |row| { try accounts.append(allocator, try accountFromRow(allocator, row)); } if (rows.err) |err| return err; return accounts.toOwnedSlice(allocator);}
pub fn listResidents(allocator: std.mem.Allocator, limit: usize) ![]Resident { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const actual_limit = if (limit == 0) 24 else @min(limit, 100); var rows = try conn.rows( \\SELECT a.handle, a.did, a.account_status, COALESCE(c.rev, ''), COUNT(r.uri) \\FROM accounts a \\LEFT JOIN ( \\ SELECT c1.did, c1.rev \\ FROM commits c1 \\ JOIN ( \\ SELECT did, MAX(seq) AS seq \\ FROM commits \\ GROUP BY did \\ ) latest ON latest.did = c1.did AND latest.seq = c1.seq \\) c ON c.did = a.did \\LEFT JOIN records r ON r.did = a.did \\WHERE a.handle NOT LIKE '%.test' \\GROUP BY a.did \\ORDER BY a.activated_at IS NULL, a.handle COLLATE NOCASE \\LIMIT ? , .{@as(i64, @intCast(actual_limit))}); defer rows.deinit();
var residents: std.ArrayList(Resident) = .empty; while (rows.next()) |row| { try residents.append(allocator, .{ .handle = try allocator.dupe(u8, row.text(0)), .did = try allocator.dupe(u8, row.text(1)), .active = AccountStatus.parse(row.text(2)).isActive(), .rev = try allocator.dupe(u8, row.text(3)), .record_count = @intCast(row.int(4)), }); } if (rows.err) |err| return err; return residents.toOwnedSlice(allocator);}
pub fn profileAvatarCid(allocator: std.mem.Allocator, did: []const u8) !?[]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const row = try conn.row( \\SELECT rb.data \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.did = ? AND r.collection = 'app.bsky.actor.profile' AND r.rkey = 'self' \\LIMIT 1 , .{did}); if (row == null) return null; defer row.?.deinit();
const profile_json = recordJsonFromBlock(allocator, row.?.nullableBlob(0) orelse "") catch return null; var parsed = std.json.parseFromSlice(std.json.Value, allocator, profile_json, .{}) catch return null; defer parsed.deinit(); if (parsed.value != .object) return null; const avatar = parsed.value.object.get("avatar") orelse return null; if (avatar != .object) return null; const ref = avatar.object.get("ref") orelse return null; if (ref != .object) return null; const link = ref.object.get("$link") orelse return null; if (link != .string) return null; const blob_row = try conn.row( \\SELECT 1 \\FROM blobs \\WHERE did = ? AND cid = ? \\LIMIT 1 , .{ did, link.string }); if (blob_row == null) return null; defer blob_row.?.deinit(); return try allocator.dupe(u8, link.string);}
pub fn listCollectionSummaries(allocator: std.mem.Allocator, limit: usize) ![]CollectionSummary { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const actual_limit = if (limit == 0) 12 else @min(limit, 50); var rows = try conn.rows( \\SELECT r.collection, COUNT(*) AS count \\FROM records r \\JOIN accounts a ON a.did = r.did \\WHERE a.handle NOT LIKE '%.test' \\GROUP BY r.collection \\ORDER BY count DESC, r.collection COLLATE NOCASE \\LIMIT ? , .{@as(i64, @intCast(actual_limit))}); defer rows.deinit();
var summaries: std.ArrayList(CollectionSummary) = .empty; while (rows.next()) |row| { try summaries.append(allocator, .{ .collection = try allocator.dupe(u8, row.text(0)), .count = @intCast(row.int(1)), }); } if (rows.err) |err| return err; return summaries.toOwnedSlice(allocator);}
pub fn listLandingActivity(allocator: std.mem.Allocator, bucket_count: usize) ![]ActivityBucket { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const actual_count = if (bucket_count == 0) 18 else @min(bucket_count, 48); var rows = try conn.rows( \\SELECT c.rev \\FROM commits c \\JOIN accounts a ON a.did = c.did \\WHERE a.handle NOT LIKE '%.test' \\ORDER BY c.seq ASC , .{}); defer rows.deinit();
var timestamps: std.ArrayList(u64) = .empty; defer timestamps.deinit(allocator); while (rows.next()) |row| { const timestamp = atid.timestampMicros(row.text(0)) catch continue; try timestamps.append(allocator, timestamp); } if (rows.err) |err| return err; if (timestamps.items.len == 0) return &.{};
const first = timestamps.items[0]; const last = timestamps.items[timestamps.items.len - 1]; const span = if (last > first) last - first else 1; const width = @max(@as(u64, 1), (span + actual_count - 1) / actual_count);
const buckets = try allocator.alloc(ActivityBucket, actual_count); for (buckets, 0..) |*bucket, i| { const start = first + width * i; bucket.* = .{ .start_us = start, .end_us = start + width, .count = 0, }; }
for (timestamps.items) |timestamp| { const offset = if (timestamp <= first) 0 else timestamp - first; const index = @min(actual_count - 1, @as(usize, @intCast(offset / width))); buckets[index].count += 1; } buckets[actual_count - 1].end_us = last; return buckets;}
pub fn listLandingRecentRecords(allocator: std.mem.Allocator, limit: usize) ![]RecentRecord { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const actual_limit = if (limit == 0) 8 else @min(limit, 24); var rows = try conn.rows( \\SELECT a.handle, r.did, r.collection, r.rkey, r.cid, r.seq \\FROM records r \\JOIN accounts a ON a.did = r.did \\WHERE a.handle NOT LIKE '%.test' \\ORDER BY r.seq DESC, r.collection COLLATE NOCASE \\LIMIT ? , .{@as(i64, @intCast(actual_limit))}); defer rows.deinit();
var records: std.ArrayList(RecentRecord) = .empty; while (rows.next()) |row| { try records.append(allocator, .{ .handle = try allocator.dupe(u8, row.text(0)), .did = try allocator.dupe(u8, row.text(1)), .collection = try allocator.dupe(u8, row.text(2)), .rkey = try allocator.dupe(u8, row.text(3)), .cid = try allocator.dupe(u8, row.text(4)), .seq = @intCast(row.int(5)), }); } if (rows.err) |err| return err; return records.toOwnedSlice(allocator);}
fn findAccountLocked(allocator: std.mem.Allocator, identifier: []const u8) !?auth.Account { const row = try conn.row( \\SELECT handle, did, email, password_hash \\FROM accounts \\WHERE lower(handle) = lower(?) OR did = ? OR lower(email) = lower(?) \\LIMIT 1 , .{ identifier, identifier, identifier }); if (row == null) return null; defer row.?.deinit(); return try accountFromRow(allocator, row.?);}
fn accountFromRow(allocator: std.mem.Allocator, row: zqlite.Row) !auth.Account { return .{ .handle = try allocator.dupe(u8, row.text(0)), .did = try allocator.dupe(u8, row.text(1)), .email = try allocator.dupe(u8, row.text(2)), .password = try allocator.dupe(u8, row.text(3)), };}
pub fn createAccount( allocator: std.mem.Allocator, handle: []const u8, email: []const u8, password: []const u8, did: []const u8, activated: bool,) !auth.Account { const signing_key = try generateAccountSigningKey(); return createAccountWithSigningKey(allocator, handle, email, password, did, activated, signing_key);}
pub fn createAccountWithSigningKey( allocator: std.mem.Allocator, handle: []const u8, email: []const u8, password: []const u8, did: []const u8, activated: bool, signing_key: [32]u8,) !auth.Account { return createAccountWithSigningKeyAndInvite(allocator, handle, email, password, did, activated, signing_key, null);}
pub fn createAccountWithSigningKeyAndInvite( allocator: std.mem.Allocator, handle: []const u8, email: []const u8, password: []const u8, did: []const u8, activated: bool, signing_key: [32]u8, invite_code: ?[]const u8,) !auth.Account { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
if (invite_code) |code| { try ensureInviteAvailableLocked(code); }
var salt: [16]u8 = undefined; store_io.random(&salt); const password_hash = try auth.hashPassword(allocator, password, salt); try conn.execNoArgs("BEGIN IMMEDIATE"); errdefer conn.execNoArgs("ROLLBACK") catch {}; try conn.exec( \\INSERT INTO accounts (did, handle, email, password_hash, activated_at, account_status, signing_key_type, signing_key) \\VALUES (?, ?, ?, ?, CASE WHEN ? THEN unixepoch() ELSE NULL END, CASE WHEN ? THEN 'active' ELSE 'deactivated' END, 'secp256k1', ?) , .{ did, handle, email, password_hash, activated, activated, zqlite.blob(&signing_key) }); if (invite_code) |code| { try recordInviteUseLocked(code, did); } try conn.execNoArgs("COMMIT"); return .{ .handle = try allocator.dupe(u8, handle), .did = try allocator.dupe(u8, did), .email = try allocator.dupe(u8, email), .password = try allocator.dupe(u8, password_hash), };}
pub fn generateAccountSigningKey() ![32]u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); return generateSigningKey();}
pub fn reserveSigningKey(allocator: std.mem.Allocator, did: ?[]const u8) !ReservedSigningKey { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const secret_key = try generateSigningKey(); var keypair = try zat.Keypair.fromSecretKey(.secp256k1, secret_key); const signing_key = try keypair.did(allocator); const expires_at = nowMs() + (24 * 60 * 60 * 1000); try conn.exec( \\INSERT INTO reserved_signing_keys (did, signing_key, signing_key_type, secret_key, expires_at) \\VALUES (?, ?, 'secp256k1', ?, ?) , .{ did, signing_key, zqlite.blob(&secret_key), expires_at }); return .{ .did = if (did) |value| try allocator.dupe(u8, value) else null, .signing_key = signing_key, .secret_key = secret_key, .expires_at = expires_at, };}
pub fn consumeReservedSigningKey(signing_key_or_did: []const u8, did: []const u8) ![32]u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); return consumeReservedSigningKeyLocked(signing_key_or_did, did);}
pub fn createInviteCode(allocator: std.mem.Allocator, public_url: []const u8, use_count: i64, for_account: []const u8, created_by: []const u8) ![]const u8 { if (use_count < 1) return error.InvalidUseCount; db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
var tries: usize = 0; while (tries < 16) : (tries += 1) { const code = try generateInviteCode(allocator, public_url); const existing = try conn.row("SELECT 1 FROM invite_codes WHERE code = ? LIMIT 1", .{code}); if (existing) |row| { row.deinit(); continue; } try insertInviteCodeLocked(code, use_count, for_account, created_by); return code; } return error.InviteCodeCollision;}
pub fn createInviteCodes(allocator: std.mem.Allocator, public_url: []const u8, code_count: i64, use_count: i64, for_account: []const u8, created_by: []const u8) ![][]const u8 { if (code_count < 1 or use_count < 1) return error.InvalidUseCount; const max_count: i64 = 100; if (code_count > max_count) return error.InvalidUseCount;
var codes: std.ArrayList([]const u8) = .empty; errdefer { for (codes.items) |code| allocator.free(code); codes.deinit(allocator); } var i: i64 = 0; while (i < code_count) : (i += 1) { try codes.append(allocator, try createInviteCode(allocator, public_url, use_count, for_account, created_by)); } return codes.toOwnedSlice(allocator);}
pub fn ensureBootstrapInviteCode(allocator: std.mem.Allocator, public_url: []const u8) !?[]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const accounts_row = try conn.row("SELECT COUNT(*) FROM accounts", .{}); if (accounts_row == null) return null; defer accounts_row.?.deinit(); const invites_row = try conn.row("SELECT COUNT(*) FROM invite_codes", .{}); if (invites_row == null) return null; defer invites_row.?.deinit(); const accounts = accounts_row.?.int(0); const invites = invites_row.?.int(0); if (accounts != 0 or invites != 0) return null;
const code = try generateInviteCode(allocator, public_url); try insertInviteCodeLocked(code, 1, "admin", "admin"); return code;}
pub fn getAccountInviteCodes(allocator: std.mem.Allocator, did: []const u8, include_used: bool) ![]InviteCode { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
var rows = try conn.rows( \\SELECT code, available_uses, disabled, for_account, created_by, created_at \\FROM invite_codes \\WHERE for_account = ? \\ORDER BY created_at DESC, code , .{did}); defer rows.deinit();
var codes: std.ArrayList(InviteCode) = .empty; while (rows.next()) |row| { const code_text = row.text(0); const uses = try inviteCodeUsesLocked(allocator, code_text); const available = row.int(1); const disabled = row.int(2) != 0; if (!include_used and (disabled or @as(i64, @intCast(uses.len)) >= available)) continue; try codes.append(allocator, .{ .code = try allocator.dupe(u8, code_text), .available = available, .disabled = disabled, .for_account = try allocator.dupe(u8, row.text(3)), .created_by = try allocator.dupe(u8, row.text(4)), .created_at = row.int(5), .uses = uses, }); } if (rows.err) |err| return err; return codes.toOwnedSlice(allocator);}
pub fn inviteCodeIsAvailable(code: []const u8) !bool { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); ensureInviteAvailableLocked(code) catch return false; return true;}
pub fn signingKeypair(did: []const u8) !zat.Keypair { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); return signingKeypairLocked(did);}
pub fn updateAccountHandle(did: []const u8, handle: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const existing = try conn.row( \\SELECT did \\FROM accounts \\WHERE lower(handle) = lower(?) \\LIMIT 1 , .{handle}); if (existing) |row| { defer row.deinit(); if (!std.mem.eql(u8, row.text(0), did)) return Error.HandleNotAvailable; }
const account = try conn.row( \\SELECT 1 \\FROM accounts \\WHERE did = ? , .{did}); if (account == null) return Error.AccountNotFound; account.?.deinit();
try conn.exec( \\UPDATE accounts \\SET handle = ? \\WHERE did = ? , .{ handle, did });}
pub fn setAccountActive(did: []const u8, active: bool) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); if (active) { const status = try accountStatusLocked(did); if (status != .active and status != .deactivated) return Error.InvalidAccountStatus; try conn.exec( \\UPDATE accounts \\SET account_status = 'active', \\ activated_at = COALESCE(activated_at, unixepoch()) \\WHERE did = ? , .{did}); } else { try conn.exec( \\UPDATE accounts \\SET account_status = 'deactivated', \\ activated_at = COALESCE(activated_at, unixepoch()) \\WHERE did = ? , .{did}); }}
pub fn setAccountTakendown(did: []const u8, applied: bool, status_ref: ?[]const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); if (applied) { try conn.execNoArgs("BEGIN IMMEDIATE"); errdefer conn.execNoArgs("ROLLBACK") catch {}; try conn.exec( \\UPDATE accounts \\SET account_status = 'takendown', \\ account_status_ref = ? \\WHERE did = ? , .{ status_ref, did }); try revokeAccountTokensLocked(did); try conn.execNoArgs("COMMIT"); } else { try conn.exec( \\UPDATE accounts \\SET account_status = 'active', \\ account_status_ref = NULL, \\ activated_at = COALESCE(activated_at, unixepoch()) \\WHERE did = ? AND account_status = 'takendown' , .{did}); }}
pub fn accountStatus(did: []const u8) !AccountStatus { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); return accountStatusLocked(did);}
pub fn sequenceAccountEvent(allocator: std.mem.Allocator, did: []const u8, status: AccountStatus) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const seq = try nextSeqLocked(); const frame = try accountEventFrame(allocator, seq, did, status); try insertSeqEventLocked(seq, did, "", frame); eventlog.publish(seq);}
pub fn sequenceIdentityEvent(allocator: std.mem.Allocator, did: []const u8, handle: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const seq = try nextSeqLocked(); const frame = try identityEventFrame(allocator, seq, did, handle); try insertSeqEventLocked(seq, did, "", frame); eventlog.publish(seq);}
pub fn sequenceSyncEvent(allocator: std.mem.Allocator, did: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const seq = try nextSeqLocked(); const commit = try latestCommitRawLocked(allocator, did) orelse return Error.RepoNotFound; const frame = try syncEventFrame(allocator, seq, did, commit.commit_cid_text, commit.data_cid_raw, commit.rev, commit.commit_data); try insertSeqEventLocked(seq, did, commit.commit_cid_text, frame); eventlog.publish(seq);}
pub fn isAccountActive(did: []const u8) bool { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return false; return accountActiveLocked(did) catch false;}
pub fn create( allocator: std.mem.Allocator, account: auth.Account, collection: []const u8, maybe_rkey: ?[]const u8, value: std.json.Value,) !Record { const rkey = maybe_rkey orelse try nextRkey(allocator); const result = try applyWrites(allocator, account, &.{.{ .create = .{ .collection = collection, .rkey = rkey, .value = value, } }}); if (result.records.len == 0) return Error.MissingRecord; return result.records[0];}
pub fn putOAuthRequest( request_id: []const u8, client_id: []const u8, redirect_uri: []const u8, scope: []const u8, state: []const u8, code_challenge: []const u8, code_challenge_method: []const u8, response_mode: []const u8, login_hint: ?[]const u8, dpop_jkt: ?[]const u8, expires_at: i64,) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\INSERT INTO oauth_requests \\ (request_id, client_id, redirect_uri, scope, state, code_challenge, code_challenge_method, response_mode, login_hint, dpop_jkt, expires_at) \\VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) , .{ request_id, client_id, redirect_uri, scope, state, code_challenge, code_challenge_method, response_mode, login_hint, dpop_jkt, expires_at });}
pub fn getOAuthRequest(allocator: std.mem.Allocator, request_id: []const u8) !?OAuthRequest { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const row = try conn.row( \\SELECT request_id, client_id, redirect_uri, scope, state, code_challenge, code_challenge_method, response_mode, login_hint, dpop_jkt, expires_at, sub, code, auth_method \\FROM oauth_requests \\WHERE request_id = ? , .{request_id}); if (row == null) return null; defer row.?.deinit(); return try oauthRequestFromRow(row.?, allocator);}
pub fn authorizeOAuthRequest(request_id: []const u8, did: []const u8, code: []const u8, auth_method: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE oauth_requests \\SET sub = ?, code = ?, auth_method = ? \\WHERE request_id = ? , .{ did, code, auth_method, request_id });}
pub fn getOAuthRequestByCode(allocator: std.mem.Allocator, code: []const u8) !?OAuthRequest { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const row = try conn.row( \\SELECT request_id, client_id, redirect_uri, scope, state, code_challenge, code_challenge_method, response_mode, login_hint, dpop_jkt, expires_at, sub, code, auth_method \\FROM oauth_requests \\WHERE code = ? , .{code}); if (row == null) return null; defer row.?.deinit(); return try oauthRequestFromRow(row.?, allocator);}
pub fn consumeOAuthCode(request_id: []const u8, code: []const u8) !bool { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( "UPDATE oauth_requests SET code = NULL WHERE request_id = ? AND code = ?", .{ request_id, code }, ); const row = try conn.row("SELECT changes()", .{}); if (row == null) return false; defer row.?.deinit(); return row.?.int(0) == 1;}
pub fn putOAuthToken( did: []const u8, client_id: []const u8, scope: []const u8, access_token: []const u8, refresh_token: []const u8, access_expires_at: i64, refresh_expires_at: i64, dpop_jkt: ?[]const u8, auth_method: ?[]const u8,) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\INSERT INTO oauth_tokens (family_id, access_token, refresh_token, did, client_id, scope, expires_at, access_expires_at, refresh_expires_at, dpop_jkt, auth_method) \\VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) , .{ access_token, access_token, refresh_token, did, client_id, scope, access_expires_at, access_expires_at, refresh_expires_at, dpop_jkt, auth_method });}
pub fn getOAuthToken(allocator: std.mem.Allocator, token: []const u8) !?OAuthToken { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const row = try conn.row( \\SELECT family_id, did, client_id, scope, access_token, refresh_token, access_expires_at, refresh_expires_at, \\ previous_refresh_token, previous_refresh_expires_at, revoked_at, dpop_jkt, auth_method \\FROM oauth_tokens \\WHERE access_token = ? OR refresh_token = ? OR previous_refresh_token = ? \\ORDER BY created_at DESC \\LIMIT 1 , .{ token, token, token }); if (row == null) return try getOAuthTokenByUsedRefresh(allocator, token); defer row.?.deinit(); const access_token = row.?.text(4); const refresh_token = row.?.text(5); const previous_refresh_token = row.?.nullableText(8); const kind: OAuthToken.Kind = if (std.mem.eql(u8, token, access_token)) .access else if (std.mem.eql(u8, token, refresh_token)) .refresh else if (previous_refresh_token != null and std.mem.eql(u8, token, previous_refresh_token.?)) .previous_refresh else return null; return .{ .family_id = try allocator.dupe(u8, row.?.text(0)), .did = try allocator.dupe(u8, row.?.text(1)), .client_id = try allocator.dupe(u8, row.?.text(2)), .scope = try allocator.dupe(u8, row.?.text(3)), .access_token = try allocator.dupe(u8, access_token), .refresh_token = try allocator.dupe(u8, refresh_token), .access_expires_at = row.?.int(6), .refresh_expires_at = row.?.int(7), .previous_refresh_token = if (previous_refresh_token) |text| try allocator.dupe(u8, text) else null, .previous_refresh_expires_at = row.?.nullableInt(9), .revoked = row.?.nullableInt(10) != null, .dpop_jkt = if (row.?.nullableText(11)) |text| try allocator.dupe(u8, text) else null, .auth_method = if (row.?.nullableText(12)) |text| try allocator.dupe(u8, text) else null, .kind = kind, };}
fn getOAuthTokenByUsedRefresh(allocator: std.mem.Allocator, token: []const u8) !?OAuthToken { const row = try conn.row( \\SELECT t.family_id, t.did, t.client_id, t.scope, t.access_token, t.refresh_token, t.access_expires_at, t.refresh_expires_at, \\ t.previous_refresh_token, t.previous_refresh_expires_at, t.revoked_at, t.dpop_jkt, t.auth_method \\FROM oauth_used_refresh_tokens u \\JOIN oauth_tokens t ON t.family_id = u.family_id \\WHERE u.refresh_token = ? \\LIMIT 1 , .{token}); if (row == null) return null; defer row.?.deinit(); return .{ .family_id = try allocator.dupe(u8, row.?.text(0)), .did = try allocator.dupe(u8, row.?.text(1)), .client_id = try allocator.dupe(u8, row.?.text(2)), .scope = try allocator.dupe(u8, row.?.text(3)), .access_token = try allocator.dupe(u8, row.?.text(4)), .refresh_token = try allocator.dupe(u8, row.?.text(5)), .access_expires_at = row.?.int(6), .refresh_expires_at = row.?.int(7), .previous_refresh_token = if (row.?.nullableText(8)) |text| try allocator.dupe(u8, text) else null, .previous_refresh_expires_at = row.?.nullableInt(9), .revoked = row.?.nullableInt(10) != null, .dpop_jkt = if (row.?.nullableText(11)) |text| try allocator.dupe(u8, text) else null, .auth_method = if (row.?.nullableText(12)) |text| try allocator.dupe(u8, text) else null, .kind = .used_refresh, };}
pub fn rotateOAuthToken( current_refresh_token: []const u8, new_access_token: []const u8, new_refresh_token: []const u8, access_expires_at: i64, refresh_expires_at: i64, previous_refresh_expires_at: i64,) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\INSERT OR IGNORE INTO oauth_used_refresh_tokens (refresh_token, family_id, expires_at) \\SELECT refresh_token, family_id, refresh_expires_at \\FROM oauth_tokens \\WHERE refresh_token = ? AND revoked_at IS NULL , .{current_refresh_token}); try conn.exec( \\UPDATE oauth_tokens \\SET previous_refresh_token = refresh_token, \\ previous_refresh_expires_at = ?, \\ access_token = ?, \\ refresh_token = ?, \\ expires_at = ?, \\ access_expires_at = ?, \\ refresh_expires_at = ? \\WHERE refresh_token = ? AND revoked_at IS NULL , .{ previous_refresh_expires_at, new_access_token, new_refresh_token, access_expires_at, access_expires_at, refresh_expires_at, current_refresh_token });}
pub fn recordDpopJti(jti: []const u8, expires_at: i64) !bool { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec("DELETE FROM dpop_jtis WHERE expires_at < unixepoch()", .{}); const row = try conn.row("SELECT 1 FROM dpop_jtis WHERE jti = ?", .{jti}); if (row) |found| { found.deinit(); return false; } try conn.exec( \\INSERT INTO dpop_jtis (jti, expires_at) \\VALUES (?, ?) , .{ jti, expires_at }); return true;}
/// Dev-tools only: force every OAuth access token to be expired (refresh/// tokens untouched). Access expiry is set to expired_at so callers can use/// the dev clock.pub fn expireAllOAuthAccessTokens(expired_at: i64) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE oauth_tokens SET access_expires_at = ? WHERE revoked_at IS NULL , .{expired_at});}
pub fn revokeOAuthToken(token: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE oauth_tokens \\SET revoked_at = unixepoch() \\WHERE access_token = ? OR refresh_token = ? OR previous_refresh_token = ? \\ OR family_id IN (SELECT family_id FROM oauth_used_refresh_tokens WHERE refresh_token = ?) , .{ token, token, token, token });}
pub fn revokeOAuthTokenFamily(family_id: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE oauth_tokens \\SET revoked_at = unixepoch() \\WHERE family_id = ? , .{family_id});}
pub fn listOAuthGrantsForAccount(allocator: std.mem.Allocator, did: []const u8, active_only: bool, limit: i64) ![]OAuthGrantInfo { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const sql = if (active_only) \\SELECT t.did, a.handle, t.client_id, t.scope, t.created_at, t.refresh_expires_at, t.revoked_at, t.auth_method, \\ CASE WHEN t.revoked_at IS NULL AND t.refresh_expires_at > unixepoch() THEN 1 ELSE 0 END AS active \\FROM oauth_tokens t \\JOIN accounts a ON a.did = t.did \\WHERE t.did = ? AND t.revoked_at IS NULL AND t.refresh_expires_at > unixepoch() \\ORDER BY t.created_at DESC \\LIMIT ? else \\SELECT t.did, a.handle, t.client_id, t.scope, t.created_at, t.refresh_expires_at, t.revoked_at, t.auth_method, \\ CASE WHEN t.revoked_at IS NULL AND t.refresh_expires_at > unixepoch() THEN 1 ELSE 0 END AS active \\FROM oauth_tokens t \\JOIN accounts a ON a.did = t.did \\WHERE t.did = ? \\ORDER BY t.created_at DESC \\LIMIT ? ; var rows = try conn.rows(sql, .{ did, limit }); defer rows.deinit(); return try oauthGrantsFromRows(allocator, &rows);}
pub fn listOAuthGrantsForAllAccounts(allocator: std.mem.Allocator, active_only: bool, limit: i64) ![]OAuthGrantInfo { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const sql = if (active_only) \\SELECT t.did, a.handle, t.client_id, t.scope, t.created_at, t.refresh_expires_at, t.revoked_at, t.auth_method, \\ CASE WHEN t.revoked_at IS NULL AND t.refresh_expires_at > unixepoch() THEN 1 ELSE 0 END AS active \\FROM oauth_tokens t \\JOIN accounts a ON a.did = t.did \\WHERE t.revoked_at IS NULL AND t.refresh_expires_at > unixepoch() \\ORDER BY t.created_at DESC \\LIMIT ? else \\SELECT t.did, a.handle, t.client_id, t.scope, t.created_at, t.refresh_expires_at, t.revoked_at, t.auth_method, \\ CASE WHEN t.revoked_at IS NULL AND t.refresh_expires_at > unixepoch() THEN 1 ELSE 0 END AS active \\FROM oauth_tokens t \\JOIN accounts a ON a.did = t.did \\ORDER BY t.created_at DESC \\LIMIT ? ; var rows = try conn.rows(sql, .{limit}); defer rows.deinit(); return try oauthGrantsFromRows(allocator, &rows);}
pub fn createSessionTokenRow( allocator: std.mem.Allocator, did: []const u8, access_jti: []const u8, refresh_jti: []const u8, access_expires_at: i64, refresh_expires_at: i64, auth_method: []const u8, controller_did: ?[]const u8, app_password_name: ?[]const u8,) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const id = try randomToken(allocator, "ses-", 16); errdefer allocator.free(id); try conn.exec( \\INSERT INTO session_tokens \\ (id, did, access_jti, refresh_jti, access_expires_at, refresh_expires_at, auth_method, controller_did, app_password_name) \\VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) , .{ id, did, access_jti, refresh_jti, access_expires_at, refresh_expires_at, auth_method, controller_did, app_password_name }); return id;}
pub fn listSessionsForAccount(allocator: std.mem.Allocator, did: []const u8, active_only: bool, limit: i64) ![]SessionInfo { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const sql = if (active_only) \\SELECT s.id, s.did, a.handle, s.auth_method, s.app_password_name, s.controller_did, \\ s.created_at, s.updated_at, s.last_used_at, s.access_expires_at, s.refresh_expires_at, s.revoked_at, \\ CASE WHEN s.revoked_at IS NULL AND (s.access_expires_at > unixepoch() OR s.refresh_expires_at > unixepoch()) THEN 1 ELSE 0 END AS active \\FROM session_tokens s \\JOIN accounts a ON a.did = s.did \\WHERE s.did = ? AND s.revoked_at IS NULL AND (s.access_expires_at > unixepoch() OR s.refresh_expires_at > unixepoch()) \\ORDER BY COALESCE(s.last_used_at, s.updated_at, s.created_at) DESC \\LIMIT ? else \\SELECT s.id, s.did, a.handle, s.auth_method, s.app_password_name, s.controller_did, \\ s.created_at, s.updated_at, s.last_used_at, s.access_expires_at, s.refresh_expires_at, s.revoked_at, \\ CASE WHEN s.revoked_at IS NULL AND (s.access_expires_at > unixepoch() OR s.refresh_expires_at > unixepoch()) THEN 1 ELSE 0 END AS active \\FROM session_tokens s \\JOIN accounts a ON a.did = s.did \\WHERE s.did = ? \\ORDER BY COALESCE(s.last_used_at, s.updated_at, s.created_at) DESC \\LIMIT ? ; var rows = try conn.rows(sql, .{ did, limit }); defer rows.deinit(); return try sessionsFromRows(allocator, &rows);}
pub fn listSessionsForAllAccounts(allocator: std.mem.Allocator, active_only: bool, limit: i64) ![]SessionInfo { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const sql = if (active_only) \\SELECT s.id, s.did, a.handle, s.auth_method, s.app_password_name, s.controller_did, \\ s.created_at, s.updated_at, s.last_used_at, s.access_expires_at, s.refresh_expires_at, s.revoked_at, \\ CASE WHEN s.revoked_at IS NULL AND (s.access_expires_at > unixepoch() OR s.refresh_expires_at > unixepoch()) THEN 1 ELSE 0 END AS active \\FROM session_tokens s \\JOIN accounts a ON a.did = s.did \\WHERE s.revoked_at IS NULL AND (s.access_expires_at > unixepoch() OR s.refresh_expires_at > unixepoch()) \\ORDER BY COALESCE(s.last_used_at, s.updated_at, s.created_at) DESC \\LIMIT ? else \\SELECT s.id, s.did, a.handle, s.auth_method, s.app_password_name, s.controller_did, \\ s.created_at, s.updated_at, s.last_used_at, s.access_expires_at, s.refresh_expires_at, s.revoked_at, \\ CASE WHEN s.revoked_at IS NULL AND (s.access_expires_at > unixepoch() OR s.refresh_expires_at > unixepoch()) THEN 1 ELSE 0 END AS active \\FROM session_tokens s \\JOIN accounts a ON a.did = s.did \\ORDER BY COALESCE(s.last_used_at, s.updated_at, s.created_at) DESC \\LIMIT ? ; var rows = try conn.rows(sql, .{limit}); defer rows.deinit(); return try sessionsFromRows(allocator, &rows);}
pub fn createAppPassword( allocator: std.mem.Allocator, did: []const u8, name: []const u8, password_hash: []const u8, privileged: bool, scopes: ?[]const u8, created_by_controller_did: ?[]const u8,) !AppPassword { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\INSERT INTO app_passwords \\ (did, name, password_hash, privileged, scopes, created_by_controller_did) \\VALUES (?, ?, ?, ?, ?, ?) , .{ did, name, password_hash, privileged, scopes, created_by_controller_did }); const row = (try conn.row( \\SELECT name, password_hash, created_at, privileged, scopes, created_by_controller_did \\FROM app_passwords \\WHERE did = ? AND name = ? \\LIMIT 1 , .{ did, name })) orelse return error.MissingRecord; defer row.deinit(); return appPasswordFromRow(row, allocator);}
pub fn listAppPasswords(allocator: std.mem.Allocator, did: []const u8) ![]AppPassword { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); var rows = try conn.rows( \\SELECT name, password_hash, created_at, privileged, scopes, created_by_controller_did \\FROM app_passwords \\WHERE did = ? \\ORDER BY created_at DESC, name ASC , .{did}); defer rows.deinit(); var out: std.ArrayList(AppPassword) = .empty; while (rows.next()) |row| { try out.append(allocator, try appPasswordFromRow(row, allocator)); } return out.toOwnedSlice(allocator);}
pub fn getAppPasswordByName(allocator: std.mem.Allocator, did: []const u8, name: []const u8) !?AppPassword { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const row = try conn.row( \\SELECT name, password_hash, created_at, privileged, scopes, created_by_controller_did \\FROM app_passwords \\WHERE did = ? AND name = ? \\LIMIT 1 , .{ did, name }); if (row == null) return null; defer row.?.deinit(); return try appPasswordFromRow(row.?, allocator);}
pub fn findMatchingAppPassword(allocator: std.mem.Allocator, did: []const u8, password: []const u8) !?AppPassword { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); var rows = try conn.rows( \\SELECT name, password_hash, created_at, privileged, scopes, created_by_controller_did \\FROM app_passwords \\WHERE did = ? \\ORDER BY created_at DESC , .{did}); defer rows.deinit(); while (rows.next()) |row| { const candidate = try appPasswordFromRow(row, allocator); if (auth.passwordHashMatches(candidate.password_hash, password)) return candidate; } return null;}
pub fn sessionAuthMethod(allocator: std.mem.Allocator, did: []const u8, jti: []const u8) !?[]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const row = try conn.row( \\SELECT auth_method \\FROM session_tokens \\WHERE did = ? AND (access_jti = ? OR refresh_jti = ?) AND revoked_at IS NULL \\LIMIT 1 , .{ did, jti, jti }); if (row == null) return null; defer row.?.deinit(); return try allocator.dupe(u8, row.?.text(0));}
pub fn revokeAppPassword(did: []const u8, name: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE session_tokens \\SET revoked_at = unixepoch(), updated_at = unixepoch() \\WHERE did = ? AND app_password_name = ? AND revoked_at IS NULL , .{ did, name }); try conn.exec( \\DELETE FROM app_passwords \\WHERE did = ? AND name = ? , .{ did, name });}
pub fn sessionTokenIsActive(did: []const u8, jti: []const u8, scope: []const u8) !bool { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const is_refresh = std.mem.eql(u8, scope, "com.atproto.refresh"); const row = if (is_refresh) try conn.row( \\SELECT id \\FROM session_tokens \\WHERE did = ? AND refresh_jti = ? AND refresh_expires_at > unixepoch() AND revoked_at IS NULL \\LIMIT 1 , .{ did, jti }) else try conn.row( \\SELECT id \\FROM session_tokens \\WHERE did = ? AND access_jti = ? AND access_expires_at > unixepoch() AND revoked_at IS NULL \\LIMIT 1 , .{ did, jti }); if (row) |found| { defer found.deinit(); try conn.exec("UPDATE session_tokens SET last_used_at = unixepoch() WHERE id = ?", .{found.text(0)}); return true; } return false;}
pub fn rotateSessionToken( did: []const u8, old_refresh_jti: []const u8, new_access_jti: []const u8, new_refresh_jti: []const u8, access_expires_at: i64, refresh_expires_at: i64,) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const row = try conn.row( \\SELECT id \\FROM session_tokens \\WHERE did = ? AND refresh_jti = ? AND refresh_expires_at > unixepoch() AND revoked_at IS NULL \\LIMIT 1 , .{ did, old_refresh_jti }); if (row == null) return error.InvalidRefreshSession; defer row.?.deinit(); try conn.exec( \\UPDATE session_tokens \\SET access_jti = ?, refresh_jti = ?, access_expires_at = ?, refresh_expires_at = ?, updated_at = unixepoch(), last_used_at = unixepoch() \\WHERE id = ? , .{ new_access_jti, new_refresh_jti, access_expires_at, refresh_expires_at, row.?.text(0) });}
pub fn revokeSessionToken(did: []const u8, jti: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE session_tokens \\SET revoked_at = unixepoch(), updated_at = unixepoch() \\WHERE did = ? AND (access_jti = ? OR refresh_jti = ?) , .{ did, jti, jti });}
pub fn recordAuditEvent( subject_did: []const u8, actor_did: []const u8, controller_did: ?[]const u8, action: []const u8, details_json: []const u8,) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const id = try randomToken(std.heap.page_allocator, "aud-", 16); defer std.heap.page_allocator.free(id); try conn.exec( \\INSERT INTO account_audit_log \\ (id, subject_did, actor_did, controller_did, action, details_json) \\VALUES (?, ?, ?, ?, ?, ?) , .{ id, subject_did, actor_did, controller_did, action, details_json });}
pub fn putWebAuthnChallenge(did: []const u8, kind: []const u8, challenge: []const u8, state_json: []const u8, expires_at: i64) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const id = try randomToken(std.heap.page_allocator, "wan-", 16); defer std.heap.page_allocator.free(id); try conn.exec("DELETE FROM webauthn_challenges WHERE did = ? AND kind = ?", .{ did, kind }); try conn.exec( \\INSERT INTO webauthn_challenges (id, did, challenge, kind, state_json, expires_at) \\VALUES (?, ?, ?, ?, ?, ?) , .{ id, did, challenge, kind, state_json, expires_at });}
pub fn getWebAuthnChallenge(allocator: std.mem.Allocator, did: []const u8, kind: []const u8) !?WebAuthnChallenge { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const row = try conn.row( \\SELECT id, did, challenge, kind, state_json, expires_at \\FROM webauthn_challenges \\WHERE did = ? AND kind = ? , .{ did, kind }); if (row == null) return null; defer row.?.deinit(); return .{ .id = try allocator.dupe(u8, row.?.text(0)), .did = try allocator.dupe(u8, row.?.text(1)), .challenge = try allocator.dupe(u8, row.?.text(2)), .kind = try allocator.dupe(u8, row.?.text(3)), .state_json = try allocator.dupe(u8, row.?.text(4)), .expires_at = row.?.int(5), };}
pub fn deleteWebAuthnChallenge(did: []const u8, kind: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec("DELETE FROM webauthn_challenges WHERE did = ? AND kind = ?", .{ did, kind });}
pub fn putDiscoverableWebAuthnChallenge(request_id: []const u8, challenge: []const u8, expires_at: i64) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec("DELETE FROM webauthn_discoverable_challenges WHERE request_id = ?", .{request_id}); try conn.exec( \\INSERT INTO webauthn_discoverable_challenges (request_id, challenge, expires_at) \\VALUES (?, ?, ?) , .{ request_id, challenge, expires_at });}
pub fn getDiscoverableWebAuthnChallenge(allocator: std.mem.Allocator, request_id: []const u8) !?DiscoverableWebAuthnChallenge { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const row = try conn.row( \\SELECT request_id, challenge, expires_at \\FROM webauthn_discoverable_challenges \\WHERE request_id = ? , .{request_id}); if (row == null) return null; defer row.?.deinit(); return .{ .request_id = try allocator.dupe(u8, row.?.text(0)), .challenge = try allocator.dupe(u8, row.?.text(1)), .expires_at = row.?.int(2), };}
pub fn deleteDiscoverableWebAuthnChallenge(request_id: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec("DELETE FROM webauthn_discoverable_challenges WHERE request_id = ?", .{request_id});}
pub fn savePasskey(allocator: std.mem.Allocator, did: []const u8, credential_id: []const u8, public_key: []const u8, sign_count: u32, friendly_name: ?[]const u8) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const id = try randomToken(allocator, "pk-", 16); try conn.exec( \\INSERT INTO passkeys (id, did, credential_id, public_key, sign_count, friendly_name) \\VALUES (?, ?, ?, ?, ?, ?) , .{ id, did, zqlite.blob(credential_id), zqlite.blob(public_key), @as(i64, @intCast(sign_count)), friendly_name }); return id;}
pub fn getPasskeyByCredentialId(allocator: std.mem.Allocator, credential_id: []const u8) !?Passkey { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const row = try conn.row( \\SELECT id, did, credential_id, public_key, sign_count, friendly_name, created_at, last_used \\FROM passkeys \\WHERE credential_id = ? , .{zqlite.blob(credential_id)}); if (row == null) return null; defer row.?.deinit(); return try passkeyFromRow(allocator, row.?);}
pub fn listPasskeys(allocator: std.mem.Allocator, did: []const u8) ![]Passkey { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); var rows = try conn.rows( \\SELECT id, did, credential_id, public_key, sign_count, friendly_name, created_at, last_used \\FROM passkeys \\WHERE did = ? \\ORDER BY created_at DESC , .{did}); defer rows.deinit(); var out: std.ArrayList(Passkey) = .empty; while (rows.next()) |row| try out.append(allocator, try passkeyFromRow(allocator, row)); if (rows.err) |err| return err; return out.toOwnedSlice(allocator);}
pub fn updatePasskeyUse(id: []const u8, sign_count: u32) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE passkeys \\SET sign_count = ?, last_used = unixepoch() \\WHERE id = ? , .{ @as(i64, @intCast(sign_count)), id });}
pub fn deletePasskey(did: []const u8, id: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec("DELETE FROM passkeys WHERE did = ? AND id = ?", .{ did, id });}
pub fn updatePasskeyName(did: []const u8, id: []const u8, friendly_name: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE passkeys \\SET friendly_name = ? \\WHERE did = ? AND id = ? , .{ friendly_name, did, id });}
fn passkeyFromRow(allocator: std.mem.Allocator, row: zqlite.Row) !Passkey { return .{ .id = try allocator.dupe(u8, row.text(0)), .did = try allocator.dupe(u8, row.text(1)), .credential_id = try allocator.dupe(u8, row.blob(2)), .public_key = try allocator.dupe(u8, row.blob(3)), .sign_count = @intCast(row.int(4)), .friendly_name = if (row.nullableText(5)) |text| try allocator.dupe(u8, text) else null, .created_at = row.int(6), .last_used = row.nullableInt(7), };}
pub fn put( allocator: std.mem.Allocator, account: auth.Account, collection: []const u8, rkey: []const u8, value: std.json.Value,) !Record { const result = try applyWrites(allocator, account, &.{.{ .update = .{ .collection = collection, .rkey = rkey, .value = value, } }}); if (result.records.len == 0) return Error.MissingRecord; return result.records[0];}
pub fn delete(allocator: std.mem.Allocator, account: auth.Account, collection: []const u8, rkey: []const u8) !CommitInfo { const result = try applyWrites(allocator, account, &.{.{ .delete = .{ .collection = collection, .rkey = rkey, } }}); return result.commit;}
pub fn applyWrites(allocator: std.mem.Allocator, account: auth.Account, ops: []const WriteOp) !WriteResult { return applyWritesWithOptions(allocator, account, ops, .{});}
pub fn applyWritesWithOptions(allocator: std.mem.Allocator, account: auth.Account, ops: []const WriteOp, options: WriteOptions) !WriteResult { return applyWritesMeasured(allocator, account, ops, options, null);}
pub fn applyWritesProfiled(allocator: std.mem.Allocator, account: auth.Account, ops: []const WriteOp, profile: *WriteProfile) !WriteResult { profile.* = .{}; return applyWritesMeasured(allocator, account, ops, .{}, profile);}
fn applyWritesMeasured(allocator: std.mem.Allocator, account: auth.Account, ops: []const WriteOp, options: WriteOptions, profile: ?*WriteProfile) !WriteResult { const total_start = monotonicNs(); if (ops.len == 0) return Error.MissingRecord;
const validation_start = monotonicNs(); for (ops) |op| switch (op) { .create => |create_op| { if (zat.Nsid.parse(create_op.collection) == null) return Error.InvalidCollection; if (create_op.rkey) |rkey| if (zat.Rkey.parse(rkey) == null) return Error.InvalidRecordKey; try validateRecordForWrite(create_op.collection, create_op.rkey, create_op.value, options.validate); }, .update => |update_op| { if (zat.Nsid.parse(update_op.collection) == null) return Error.InvalidCollection; if (zat.Rkey.parse(update_op.rkey) == null) return Error.InvalidRecordKey; try validateRecordForWrite(update_op.collection, update_op.rkey, update_op.value, options.validate); }, .delete => |delete_op| { if (zat.Nsid.parse(delete_op.collection) == null) return Error.InvalidCollection; if (zat.Rkey.parse(delete_op.rkey) == null) return Error.InvalidRecordKey; }, }; addProfile(profile, .validation_ns, elapsedNs(validation_start));
const lock_start = monotonicNs(); const lane = write_lanes.lock(store_io, account.did); defer lane.unlock(); addProfile(profile, .lock_wait_ns, elapsedNs(lock_start));
try requireInitialized(); const seq = try nextSeqLocked(); const load_start = monotonicNs(); const loaded = blk: { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); const current = try latestCommitRawLocked(allocator, account.did); const resolved_ops = try resolveWriteOpsLocked(allocator, ops); try requireWriteSwapsLocked(allocator, account.did, current, options, resolved_ops); break :blk .{ .current = current, .rev = try revForSeq(allocator, seq, if (current) |root| root.rev else null), .keypair = try signingKeypairLocked(account.did), .ops = resolved_ops, .prev_record_cids = try previousRecordCidsLocked(allocator, account.did, resolved_ops), }; }; const current = loaded.current; const rev = loaded.rev; const keypair = loaded.keypair; const resolved_ops = loaded.ops; const prev_record_cids = loaded.prev_record_cids;
var repo_block_reader = RepoBlockReader{ .allocator = allocator, .did = account.did, }; var tree = if (current) |root| try zat.mst.Mst.loadLazy(allocator, root.data_cid_raw, repo_block_reader.reader()) else zat.mst.Mst.init(allocator); addProfile(profile, .load_repo_ns, elapsedNs(load_start));
var records: std.ArrayList(Record) = .empty; var record_blocks: std.ArrayList(ImportedBlock) = .empty; var mst_blocks: std.ArrayList(zat.car.Block) = .empty; var blob_refs: std.ArrayList(BlobRef) = .empty;
const stage_start = monotonicNs(); for (resolved_ops) |op| switch (op) { .create => |create_op| { const rkey = create_op.rkey orelse return Error.InvalidRecordKey; const record = try stageRecordWrite(allocator, &tree, account, create_op.collection, rkey, create_op.value, options.validate, rev, seq, &record_blocks, &blob_refs); try records.append(allocator, record); }, .update => |update_op| { const record = try stageRecordWrite(allocator, &tree, account, update_op.collection, update_op.rkey, update_op.value, options.validate, rev, seq, &record_blocks, &blob_refs); try records.append(allocator, record); }, .delete => |delete_op| { const path = try repoPath(allocator, delete_op.collection, delete_op.rkey); _ = try tree.deleteReturn(path); }, }; addProfile(profile, .stage_records_ns, elapsedNs(stage_start));
const build_start = monotonicNs(); const data_cid = try tree.rootCid(); try writeMstBlocks(&tree, &mst_blocks); const commit = try signedCommitWithKeypair(allocator, account.did, rev, data_cid, if (current) |root| root.commit_cid_raw else null, &keypair); addProfile(profile, .build_commit_ns, elapsedNs(build_start));
const sql_start = monotonicNs(); db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try conn.exclusiveTransaction(); errdefer conn.rollback(); for (record_blocks.items) |block| { try conn.exec( \\INSERT INTO repo_blocks (did, cid, data, repo_rev) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(did, cid) DO UPDATE SET \\ data = excluded.data, \\ repo_rev = COALESCE(repo_blocks.repo_rev, excluded.repo_rev) , .{ account.did, block.cid, zqlite.blob(block.data), rev }); } for (mst_blocks.items) |block| { try conn.exec( \\INSERT INTO repo_blocks (did, cid, data, repo_rev) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(did, cid) DO UPDATE SET \\ data = excluded.data, \\ repo_rev = COALESCE(repo_blocks.repo_rev, excluded.repo_rev) , .{ account.did, try cidText(allocator, block.cid_raw), zqlite.blob(block.data), rev }); } try conn.exec( \\INSERT INTO repo_blocks (did, cid, data, repo_rev) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(did, cid) DO UPDATE SET \\ data = excluded.data, \\ repo_rev = COALESCE(repo_blocks.repo_rev, excluded.repo_rev) , .{ account.did, commit.cid, zqlite.blob(commit.data), rev });
for (resolved_ops) |op| switch (op) { .delete => |delete_op| { const uri = try std.fmt.allocPrint(allocator, "at://{s}/{s}/{s}", .{ account.did, delete_op.collection, delete_op.rkey }); try conn.exec("DELETE FROM expected_blobs WHERE record_uri = ?", .{uri}); try conn.exec("DELETE FROM records WHERE did = ? AND collection = ? AND rkey = ?", .{ account.did, delete_op.collection, delete_op.rkey }); }, else => {}, }; for (records.items) |record| { const uri = try record.uri(allocator); try conn.exec( \\INSERT INTO records (did, collection, rkey, uri, cid, rev, seq) \\VALUES (?, ?, ?, ?, ?, ?, ?) \\ON CONFLICT(did, collection, rkey) DO UPDATE SET \\ uri = excluded.uri, \\ cid = excluded.cid, \\ rev = excluded.rev, \\ seq = excluded.seq , .{ record.did, record.collection, record.rkey, uri, record.cid, record.rev, @as(i64, @intCast(seq)) }); try conn.exec("DELETE FROM expected_blobs WHERE record_uri = ?", .{uri}); } for (blob_refs.items) |ref| { try conn.exec( \\INSERT INTO expected_blobs (blob_cid, record_uri) \\VALUES (?, ?) \\ON CONFLICT(blob_cid, record_uri) DO UPDATE SET blob_cid = excluded.blob_cid , .{ ref.cid, ref.uri }); } try conn.exec( \\INSERT INTO commits (seq, did, cid, rev, prev) \\VALUES (?, ?, ?, ?, ?) , .{ @as(i64, @intCast(seq)), account.did, commit.cid, rev, if (current) |root| root.commit_cid_text else null }); const event_frame = try commitEventFrame( allocator, seq, account.did, commit.cid, rev, if (current) |root| root.rev else null, if (current) |root| root.data_cid_raw else null, commit.data, record_blocks.items, mst_blocks.items, resolved_ops, records.items, prev_record_cids, ); try conn.exec( \\INSERT INTO seq_events (seq, did, commit_cid, evt) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(seq) DO UPDATE SET \\ did = excluded.did, \\ commit_cid = excluded.commit_cid, \\ evt = excluded.evt , .{ @as(i64, @intCast(seq)), account.did, commit.cid, zqlite.blob(event_frame) }); try conn.commit(); addProfile(profile, .sql_ns, elapsedNs(sql_start));
const publish_start = monotonicNs(); eventlog.publish(seq); addProfile(profile, .event_publish_ns, elapsedNs(publish_start)); if (profile) |p| p.total_ns = elapsedNs(total_start);
return .{ .commit = .{ .cid = commit.cid, .rev = rev }, .records = try records.toOwnedSlice(allocator), };}const ProfileField = enum { validation_ns, lock_wait_ns, load_repo_ns, stage_records_ns, build_commit_ns, sql_ns, event_publish_ns,};
fn addProfile(profile: ?*WriteProfile, field: ProfileField, ns: u64) void { const p = profile orelse return; switch (field) { .validation_ns => p.validation_ns += ns, .lock_wait_ns => p.lock_wait_ns += ns, .load_repo_ns => p.load_repo_ns += ns, .stage_records_ns => p.stage_records_ns += ns, .build_commit_ns => p.build_commit_ns += ns, .sql_ns => p.sql_ns += ns, .event_publish_ns => p.event_publish_ns += ns, }}
fn monotonicNs() u64 { var ts: std.posix.timespec = undefined; const timestamp = switch (std.posix.errno(std.posix.system.clock_gettime(.MONOTONIC, &ts))) { .SUCCESS => ts, else => std.posix.timespec{ .sec = 0, .nsec = 0 }, }; const seconds: u64 = if (timestamp.sec < 0) 0 else @intCast(timestamp.sec); const nanos: u64 = if (timestamp.nsec < 0) 0 else @intCast(timestamp.nsec); return seconds * std.time.ns_per_s + nanos;}
fn elapsedNs(start: u64) u64 { const now = monotonicNs(); return if (now >= start) now - start else 0;}
pub fn revForSeq(allocator: std.mem.Allocator, seq: u64, current_rev: ?[]const u8) ![]const u8 { var timestamp_us = nowMicros(); if (current_rev) |rev| { const current_timestamp = atid.timestampMicros(rev) catch 0; if (timestamp_us <= current_timestamp) timestamp_us = current_timestamp + 1; } const tid = try atid.encode(timestamp_us, @intCast(seq % 1024)); return allocator.dupe(u8, &tid);}
fn nowMicros() u64 { const micros = clock.nowMicros(); return if (micros < 0) 0 else @intCast(micros);}
pub fn get(did: []const u8, collection: []const u8, rkey: []const u8) ?Record { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return null;
const row = conn.row( \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.did = ? AND r.collection = ? AND r.rkey = ? , .{ did, collection, rkey }) catch return null; if (row == null) return null; defer row.?.deinit(); return recordFromRow(row.?, std.heap.page_allocator) catch null;}
pub fn getCid(allocator: std.mem.Allocator, did: []const u8, collection: []const u8, rkey: []const u8) ?[]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return null;
const row = conn.row( \\SELECT cid \\FROM records \\WHERE did = ? AND collection = ? AND rkey = ? , .{ did, collection, rkey }) catch return null; if (row == null) return null; defer row.?.deinit(); return allocator.dupe(u8, row.?.text(0)) catch null;}
pub const RecordBlock = struct { cid: []const u8, data: []const u8,};
pub fn getRecordBlock(allocator: std.mem.Allocator, did: []const u8, collection: []const u8, rkey: []const u8) ?RecordBlock { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return null;
const row = conn.row( \\SELECT r.cid, rb.data \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.did = ? AND r.collection = ? AND r.rkey = ? , .{ did, collection, rkey }) catch return null; if (row == null) return null; defer row.?.deinit();
return .{ .cid = allocator.dupe(u8, row.?.text(0)) catch return null, .data = allocator.dupe(u8, row.?.nullableBlob(1) orelse return null) catch return null, };}
pub fn getByUri(uri: []const u8) ?Record { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return null; return getByUriLocked(uri);}
pub fn listRecentRecords(allocator: std.mem.Allocator, limit: usize) ![]Record { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
var rows = try conn.rows( \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\ORDER BY r.seq DESC \\LIMIT ? , .{@as(i64, @intCast(if (limit == 0) 100 else limit))}); defer rows.deinit();
var records: std.ArrayList(Record) = .empty; while (rows.next()) |row| { try records.append(allocator, try recordFromRow(row, allocator)); } if (rows.err) |err| return err; return records.toOwnedSlice(allocator);}
pub fn listRecords( allocator: std.mem.Allocator, did: []const u8, collection: []const u8, limit: usize,) ![]Record { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
var rows = try conn.rows( \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.did = ? AND r.collection = ? \\ORDER BY r.seq DESC \\LIMIT ? , .{ did, collection, @as(i64, @intCast(if (limit == 0) 100 else limit)) }); defer rows.deinit();
var records: std.ArrayList(Record) = .empty; while (rows.next()) |row| { try records.append(allocator, try recordFromRow(row, allocator)); } if (rows.err) |err| return err; return records.toOwnedSlice(allocator);}
pub fn listRecordsContaining( allocator: std.mem.Allocator, collection: []const u8, needle: []const u8, limit: usize,) ![]Record { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
var rows = try conn.rows( \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.collection = ? \\ORDER BY r.seq DESC \\LIMIT ? , .{ collection, @as(i64, @intCast(if (limit == 0) 100 else limit)) }); defer rows.deinit();
var records: std.ArrayList(Record) = .empty; while (rows.next()) |row| { const record = try recordFromRow(row, allocator); if (std.mem.indexOf(u8, record.value_json, needle) != null) { try records.append(allocator, record); } } if (rows.err) |err| return err; return records.toOwnedSlice(allocator);}
pub fn listRecordsByDidContaining( allocator: std.mem.Allocator, did: []const u8, collection: []const u8, needle: []const u8, limit: usize,) ![]Record { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
var rows = try conn.rows( \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.did = ? AND r.collection = ? \\ORDER BY r.seq DESC \\LIMIT ? , .{ did, collection, @as(i64, @intCast(if (limit == 0) 100 else limit)) }); defer rows.deinit();
var records: std.ArrayList(Record) = .empty; while (rows.next()) |row| { const record = try recordFromRow(row, allocator); if (std.mem.indexOf(u8, record.value_json, needle) != null) { try records.append(allocator, record); } } if (rows.err) |err| return err; return records.toOwnedSlice(allocator);}
fn getByUriLocked(uri: []const u8) ?Record { if (getByStoredUriLocked(uri)) |record| return record; const prefix = "at://"; if (!std.mem.startsWith(u8, uri, prefix)) return null; const rest = uri[prefix.len..]; const slash = std.mem.indexOfScalar(u8, rest, '/') orelse return null; const repo = rest[0..slash]; const account = (findAccountLocked(std.heap.page_allocator, repo) catch null) orelse return null; const normalized = std.fmt.allocPrint(std.heap.page_allocator, "at://{s}/{s}", .{ account.did, rest[slash + 1 ..] }) catch return null; return getByStoredUriLocked(normalized);}
fn getByStoredUriLocked(uri: []const u8) ?Record { const row = conn.row( \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.uri = ? , .{uri}) catch return null; if (row == null) return null; defer row.?.deinit(); return recordFromRow(row.?, std.heap.page_allocator) catch null;}
pub fn count(did: []const u8, collection: []const u8) usize { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return 0; const row = conn.row("SELECT count(*) FROM records WHERE did = ? AND collection = ?", .{ did, collection }) catch return 0; if (row == null) return 0; defer row.?.deinit(); return @intCast(row.?.int(0));}
pub fn countSubject(collection: []const u8, subject_did: []const u8) usize { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return 0; var rows = conn.rows( \\SELECT rb.data \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.collection = ? , .{collection}) catch return 0; defer rows.deinit();
var count_result: usize = 0; var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); while (rows.next()) |row| { _ = arena.reset(.retain_capacity); const value = recordJsonFromBlock(arena.allocator(), row.nullableBlob(0) orelse "") catch continue; if (std.mem.indexOf(u8, value, subject_did) != null) count_result += 1; } return count_result;}
pub fn listCollectionsJson(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
var rows = try conn.rows( \\SELECT DISTINCT collection \\FROM records \\WHERE did = ? \\ORDER BY collection , .{did}); defer rows.deinit();
var out: std.Io.Writer.Allocating = .init(allocator); defer out.deinit(); try out.writer.writeByte('['); var first = true; while (rows.next()) |row| { if (!first) try out.writer.writeByte(','); first = false; try out.writer.print("{f}", .{std.json.fmt(row.text(0), .{})}); } if (rows.err) |err| return err; try out.writer.writeByte(']'); return out.toOwnedSlice();}
pub fn getAppPreferences(allocator: std.mem.Allocator, did: []const u8, namespace: []const u8) ![]AppPreference { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const like_pattern = try std.fmt.allocPrint(allocator, "{s}.%", .{namespace}); var rows = try conn.rows( \\SELECT name, value_json \\FROM app_preferences \\WHERE did = ? AND (name = ? OR name LIKE ?) \\ORDER BY id ASC , .{ did, namespace, like_pattern }); defer rows.deinit();
var preferences: std.ArrayList(AppPreference) = .empty; while (rows.next()) |row| { try preferences.append(allocator, .{ .name = try allocator.dupe(u8, row.text(0)), .value_json = try allocator.dupe(u8, row.text(1)), }); } if (rows.err) |err| return err; return preferences.toOwnedSlice(allocator);}
pub fn replaceAppPreferences( allocator: std.mem.Allocator, did: []const u8, namespace: []const u8, preferences: []const AppPreference,) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const like_pattern = try std.fmt.allocPrint(allocator, "{s}.%", .{namespace}); try conn.exclusiveTransaction(); errdefer conn.rollback(); try conn.exec( \\DELETE FROM app_preferences \\WHERE did = ? AND (name = ? OR name LIKE ?) , .{ did, namespace, like_pattern }); for (preferences) |preference| { try conn.exec( \\INSERT INTO app_preferences (did, name, value_json, updated_at) \\VALUES (?, ?, ?, unixepoch()) , .{ did, preference.name, preference.value_json }); } try conn.commit();}
pub fn putBlob( allocator: std.mem.Allocator, io: Io, account: auth.Account, data: []const u8, mime_type: []const u8,) ![]const u8 { const cid = try cidForBlob(allocator, data);
db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try blobstore.put(allocator, io, account.did, cid, data); try conn.exec( \\INSERT INTO blobs (cid, did, mime_type, size) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(did, cid) DO UPDATE SET \\ mime_type = excluded.mime_type, \\ size = excluded.size , .{ cid, account.did, mime_type, @as(i64, @intCast(data.len)) }); return cid;}
pub fn getBlob( allocator: std.mem.Allocator, did: []const u8, cid: []const u8,) ?BlobRecord { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return null; const row = conn.row( \\SELECT mime_type, size \\FROM blobs \\WHERE did = ? AND cid = ? , .{ did, cid }) catch return null; const found = row orelse return null; defer found.deinit(); const size: usize = @intCast(found.int(1)); return .{ .mime_type = allocator.dupe(u8, found.text(0)) catch return null, .data = blobstore.get(allocator, did, cid, size + 1) catch return null, };}
pub fn getPublicBlob( allocator: std.mem.Allocator, did: []const u8, cid: []const u8,) ?BlobRecord { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return null; const row = conn.row( \\SELECT b.mime_type, b.size \\FROM blobs b \\WHERE b.did = ? AND b.cid = ? \\ AND EXISTS ( \\ SELECT 1 \\ FROM expected_blobs eb \\ JOIN records r ON r.uri = eb.record_uri \\ WHERE r.did = b.did AND eb.blob_cid = b.cid \\ ) , .{ did, cid }) catch return null; const found = row orelse return null; defer found.deinit(); const size: usize = @intCast(found.int(1)); return .{ .mime_type = allocator.dupe(u8, found.text(0)) catch return null, .data = blobstore.get(allocator, did, cid, size + 1) catch return null, };}
pub fn writeBlobListJson(allocator: std.mem.Allocator, did: []const u8, since: ?[]const u8, cursor: ?[]const u8, limit: usize) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const page_limit = @as(i64, @intCast(if (limit == 0) 500 else limit)); var rows = if (since) |since_rev| blk: { if (cursor) |after_cid| { break :blk try conn.rows( \\SELECT DISTINCT b.cid \\FROM blobs b \\JOIN expected_blobs eb ON eb.blob_cid = b.cid \\JOIN records r ON r.uri = eb.record_uri AND r.did = b.did \\WHERE b.did = ? AND r.rev > ? AND b.cid > ? \\ORDER BY b.cid ASC \\LIMIT ? , .{ did, since_rev, after_cid, page_limit }); } break :blk try conn.rows( \\SELECT DISTINCT b.cid \\FROM blobs b \\JOIN expected_blobs eb ON eb.blob_cid = b.cid \\JOIN records r ON r.uri = eb.record_uri AND r.did = b.did \\WHERE b.did = ? AND r.rev > ? \\ORDER BY b.cid ASC \\LIMIT ? , .{ did, since_rev, page_limit }); } else blk: { if (cursor) |after_cid| { break :blk try conn.rows( \\SELECT DISTINCT b.cid \\FROM blobs b \\JOIN expected_blobs eb ON eb.blob_cid = b.cid \\JOIN records r ON r.uri = eb.record_uri AND r.did = b.did \\WHERE b.did = ? AND b.cid > ? \\ORDER BY b.cid ASC \\LIMIT ? , .{ did, after_cid, page_limit }); } break :blk try conn.rows( \\SELECT DISTINCT b.cid \\FROM blobs b \\JOIN expected_blobs eb ON eb.blob_cid = b.cid \\JOIN records r ON r.uri = eb.record_uri AND r.did = b.did \\WHERE b.did = ? \\ORDER BY b.cid ASC \\LIMIT ? , .{ did, page_limit }); }; defer rows.deinit();
var out: std.Io.Writer.Allocating = .init(allocator); defer out.deinit(); try out.writer.writeAll("{\"cids\":["); var first = true; var last_cid: ?[]const u8 = null; while (rows.next()) |row| { if (!first) try out.writer.writeByte(','); first = false; const cid = row.text(0); last_cid = try allocator.dupe(u8, cid); try out.writer.print("{f}", .{std.json.fmt(cid, .{})}); } if (rows.err) |err| return err; try out.writer.writeByte(']'); if (last_cid) |cid| try out.writer.print(",\"cursor\":{f}", .{std.json.fmt(cid, .{})}); try out.writer.writeByte('}'); return out.toOwnedSlice();}
pub fn writeMissingBlobsJson(allocator: std.mem.Allocator, did: []const u8, cursor: ?[]const u8, limit: usize) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const page_limit = @as(i64, @intCast(if (limit == 0) 500 else limit)); var rows = if (cursor) |after_cid| try conn.rows( \\SELECT rb.blob_cid, rb.record_uri \\FROM expected_blobs rb \\JOIN records r ON r.uri = rb.record_uri \\LEFT JOIN blobs b ON b.cid = rb.blob_cid AND b.did = r.did \\WHERE r.did = ? AND b.cid IS NULL AND rb.blob_cid > ? \\ORDER BY rb.blob_cid ASC \\LIMIT ? , .{ did, after_cid, page_limit }) else try conn.rows( \\SELECT rb.blob_cid, rb.record_uri \\FROM expected_blobs rb \\JOIN records r ON r.uri = rb.record_uri \\LEFT JOIN blobs b ON b.cid = rb.blob_cid AND b.did = r.did \\WHERE r.did = ? AND b.cid IS NULL \\ORDER BY rb.blob_cid ASC \\LIMIT ? , .{ did, page_limit }); defer rows.deinit();
var out: std.Io.Writer.Allocating = .init(allocator); defer out.deinit(); try out.writer.writeAll("{\"blobs\":["); var first = true; var last_cid: ?[]const u8 = null; while (rows.next()) |row| { if (!first) try out.writer.writeByte(','); first = false; const cid = row.text(0); last_cid = try allocator.dupe(u8, cid); try out.writer.print( "{{\"cid\":{f},\"recordUri\":{f}}}", .{ std.json.fmt(cid, .{}), std.json.fmt(row.text(1), .{}) }, ); } if (rows.err) |err| return err; try out.writer.writeByte(']'); if (last_cid) |cid| try out.writer.print(",\"cursor\":{f}", .{std.json.fmt(cid, .{})}); try out.writer.writeByte('}'); return out.toOwnedSlice();}
pub fn writeAccountStatusJson(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const record_count = try scalarCountLocked( "SELECT COUNT(*) FROM records WHERE did = ?", did, ); const blob_count = try scalarCountLocked( "SELECT COUNT(*) FROM blobs WHERE did = ?", did, ); const expected_blob_count = try scalarCountLocked( \\SELECT COUNT(*) \\FROM expected_blobs rb \\JOIN records r ON r.uri = rb.record_uri \\WHERE r.did = ? , did, ); const block_count = try scalarCountLocked( "SELECT COUNT(*) FROM repo_blocks WHERE did = ?", did, ); const status = try accountStatusLocked(did); const active = status.isActive(); const root = try latestRootLocked(allocator, did);
return std.fmt.allocPrint( allocator, "{{\"activated\":{},\"validDid\":{},\"status\":{f},\"repoCommit\":{f},\"repoRev\":{f},\"repoBlocks\":{d},\"indexedRecords\":{d},\"privateStateValues\":0,\"expectedBlobs\":{d},\"importedBlobs\":{d}}}", .{ active, active, std.json.fmt(status.asString(), .{}), std.json.fmt(root.cid, .{}), std.json.fmt(root.rev, .{}), block_count, record_count, expected_blob_count, blob_count, }, );}
pub fn importRepo( allocator: std.mem.Allocator, account: auth.Account, commit_cid: []const u8, rev: []const u8, records: []const ImportedRecord, blocks: []const ImportedBlock,) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const seq = try nextSeqLocked(); try conn.exclusiveTransaction(); errdefer conn.rollback(); try conn.exec("DELETE FROM expected_blobs WHERE record_uri IN (SELECT uri FROM records WHERE did = ?)", .{account.did}); try conn.exec("DELETE FROM records WHERE did = ?", .{account.did}); try conn.exec("DELETE FROM repo_blocks WHERE did = ?", .{account.did}); try conn.exec("DELETE FROM commits WHERE did = ?", .{account.did});
for (blocks) |block| { try conn.exec( \\INSERT INTO repo_blocks (did, cid, data, repo_rev) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(did, cid) DO UPDATE SET \\ data = excluded.data, \\ repo_rev = COALESCE(repo_blocks.repo_rev, excluded.repo_rev) , .{ account.did, block.cid, zqlite.blob(block.data), rev }); } for (records) |record| { const uri = try std.fmt.allocPrint(std.heap.page_allocator, "at://{s}/{s}/{s}", .{ account.did, record.collection, record.rkey }); defer std.heap.page_allocator.free(uri); try conn.exec( \\INSERT INTO records (did, collection, rkey, uri, cid, rev, seq) \\VALUES (?, ?, ?, ?, ?, ?, ?) , .{ account.did, record.collection, record.rkey, uri, record.cid, rev, @as(i64, @intCast(seq)) }); for (record.blob_cids) |blob_cid| { try conn.exec( \\INSERT INTO expected_blobs (blob_cid, record_uri) \\VALUES (?, ?) \\ON CONFLICT(blob_cid, record_uri) DO UPDATE SET blob_cid = excluded.blob_cid , .{ blob_cid, uri }); } } try conn.exec( \\INSERT INTO commits (seq, did, cid, rev, prev) \\VALUES (?, ?, ?, ?, NULL) , .{ @as(i64, @intCast(seq)), account.did, commit_cid, rev }); const publish_event = try accountActiveLocked(account.did); if (publish_event) { const event_frame = try importedCommitEventFrame(allocator, seq, account.did, commit_cid, rev, blocks); try conn.exec( \\INSERT INTO seq_events (seq, did, commit_cid, evt) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(seq) DO UPDATE SET \\ did = excluded.did, \\ commit_cid = excluded.commit_cid, \\ evt = excluded.evt , .{ @as(i64, @intCast(seq)), account.did, commit_cid, zqlite.blob(event_frame) }); } try conn.commit(); if (publish_event) eventlog.publish(seq);}
pub fn writeRepoCar(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { return writeRepoCarSince(allocator, did, null);}
pub fn writeRepoCarSince(allocator: std.mem.Allocator, did: []const u8, since: ?[]const u8) ![]const u8 { if (since == null) return writeRepoCarFull(allocator, did);
db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const root = try latestRootLocked(allocator, did); const root_raw = try zat.multibase.base32lower.decode(allocator, root.cid[1..]); const car_root = zat.cbor.Cid{ .raw = root_raw };
const since_rev = since.?; var rows = try conn.rows( \\SELECT cid, data \\FROM repo_blocks \\WHERE did = ? AND (repo_rev IS NULL OR repo_rev > ?) \\ORDER BY repo_rev DESC, cid DESC , .{ did, since_rev }); defer rows.deinit();
var blocks: std.ArrayList(zat.car.Block) = .empty; while (rows.next()) |row| { const cid_text = row.text(0); const data = row.nullableBlob(1) orelse ""; const cid_raw = try zat.multibase.base32lower.decode(allocator, cid_text[1..]); try blocks.append(allocator, .{ .cid_raw = cid_raw, .data = try allocator.dupe(u8, data), }); } if (rows.err) |err| return err; const c: zat.car.Car = .{ .roots = &.{car_root}, .blocks = blocks.items, }; return zat.car.writeAlloc(allocator, c);}
fn writeRepoCarFull(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const root = try latestCommitRawLocked(allocator, did) orelse return Error.RepoNotFound; var blocks: std.ArrayList(zat.car.Block) = .empty; try blocks.append(allocator, .{ .cid_raw = root.commit_cid_raw, .data = root.commit_data, });
var repo_block_reader = RepoBlockReader{ .allocator = allocator, .did = did, .locked = true, }; try zat.mst.collectReachableBlocks( allocator, root.data_cid_raw, repo_block_reader.reader(), &blocks, .{ .include_records = true }, );
return zat.car.writeAlloc(allocator, .{ .roots = &.{.{ .raw = root.commit_cid_raw }}, .blocks = blocks.items, });}
pub fn writeRepoListJson(allocator: std.mem.Allocator, cursor: ?[]const u8, limit: usize) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const actual_limit = if (limit == 0) 500 else @min(limit, 1000); const parsed_cursor = try parseRepoListCursor(cursor); var rows = if (parsed_cursor) |after| try conn.rows( \\SELECT a.did, c.cid, c.rev, a.account_status, a.created_at \\FROM accounts a \\JOIN ( \\ SELECT did, MAX(seq) AS seq \\ FROM commits \\ GROUP BY did \\) latest ON latest.did = a.did \\JOIN commits c ON c.did = latest.did AND c.seq = latest.seq \\WHERE a.created_at > ? OR (a.created_at = ? AND a.did > ?) \\ORDER BY a.created_at ASC, a.did ASC \\LIMIT ? , .{ after.created_at, after.created_at, after.did, @as(i64, @intCast(actual_limit)) }) else try conn.rows( \\SELECT a.did, c.cid, c.rev, a.account_status, a.created_at \\FROM accounts a \\JOIN ( \\ SELECT did, MAX(seq) AS seq \\ FROM commits \\ GROUP BY did \\) latest ON latest.did = a.did \\JOIN commits c ON c.did = latest.did AND c.seq = latest.seq \\ORDER BY a.created_at ASC, a.did ASC \\LIMIT ? , .{@as(i64, @intCast(actual_limit))}); defer rows.deinit();
var out: std.Io.Writer.Allocating = .init(allocator); defer out.deinit(); var json: std.json.Stringify = .{ .writer = &out.writer }; try json.beginObject(); try json.objectField("repos"); try json.beginArray(); var last_created_at: ?i64 = null; var last_did: ?[]const u8 = null; while (rows.next()) |row| { const did = row.text(0); last_did = try allocator.dupe(u8, did); last_created_at = row.int(4); const status = AccountStatus.parse(row.text(3)); const active = status.isActive(); try json.beginObject(); try json.objectField("did"); try json.write(did); try json.objectField("head"); try json.write(row.text(1)); try json.objectField("rev"); try json.write(row.text(2)); try json.objectField("active"); try json.write(active); if (!active) { try json.objectField("status"); try json.write(status.asString()); } try json.endObject(); } if (rows.err) |err| return err; try json.endArray(); if (last_created_at) |created_at| { try json.objectField("cursor"); try json.write(try std.fmt.allocPrint(allocator, "{d}:{s}", .{ created_at, last_did.? })); } try json.endObject(); return out.toOwnedSlice();}
const RepoListCursor = struct { created_at: i64, did: []const u8,};
fn parseRepoListCursor(cursor: ?[]const u8) Error!?RepoListCursor { const raw = cursor orelse return null; const sep = std.mem.indexOfScalar(u8, raw, ':') orelse return Error.InvalidCursor; if (sep == 0 or sep + 1 >= raw.len) return Error.InvalidCursor; return .{ .created_at = std.fmt.parseInt(i64, raw[0..sep], 10) catch return Error.InvalidCursor, .did = raw[sep + 1 ..], };}
pub fn writeLatestCommitJson(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const root = try latestRootLocked(allocator, did); return std.fmt.allocPrint( allocator, "{{\"cid\":{f},\"rev\":{f}}}", .{ std.json.fmt(root.cid, .{}), std.json.fmt(root.rev, .{}) }, );}
pub fn listSeqEvents(allocator: std.mem.Allocator, cursor: u64, limit: usize) ![]SeqEvent { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try backfillSeqEventsLocked(allocator);
var rows = try conn.rows( \\SELECT seq, evt \\FROM seq_events \\WHERE seq > ? \\ORDER BY seq ASC \\LIMIT ? , .{ @as(i64, @intCast(cursor)), @as(i64, @intCast(@min(limit, 1000))) }); defer rows.deinit();
var events: std.ArrayList(SeqEvent) = .empty; while (rows.next()) |row| { try events.append(allocator, .{ .seq = @intCast(row.int(0)), .frame = try allocator.dupe(u8, row.nullableBlob(1) orelse ""), }); } if (rows.err) |err| return err; return events.toOwnedSlice(allocator);}
fn insertSeqEventLocked(seq: u64, did: []const u8, commit_cid: []const u8, frame: []const u8) !void { try conn.exec( \\INSERT INTO seq_events (seq, did, commit_cid, evt) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(seq) DO UPDATE SET \\ did = excluded.did, \\ commit_cid = excluded.commit_cid, \\ evt = excluded.evt , .{ @as(i64, @intCast(seq)), did, commit_cid, zqlite.blob(frame) });}
pub fn writeRepoStatusJson(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const status = try accountStatusLocked(did); const active = status.isActive(); const root = if (active) try latestRootLocked(allocator, did) else null; var out: std.Io.Writer.Allocating = .init(allocator); defer out.deinit(); var json: std.json.Stringify = .{ .writer = &out.writer }; try json.beginObject(); try json.objectField("did"); try json.write(did); try json.objectField("active"); try json.write(active); if (!active) { try json.objectField("status"); try json.write(status.asString()); } else { try json.objectField("rev"); try json.write(root.?.rev); } try json.endObject(); return out.toOwnedSlice();}
pub fn writeRecordJson(allocator: std.mem.Allocator, record: Record) ![]const u8 { var out: std.Io.Writer.Allocating = .init(allocator); defer out.deinit(); const uri = try record.uri(allocator); defer allocator.free(uri);
try out.writer.print( "{{\"uri\":{f},\"cid\":{f},\"value\":{s}}}", .{ std.json.fmt(uri, .{}), std.json.fmt(record.cid, .{}), record.value_json }, ); return out.toOwnedSlice();}
pub fn writeListJson( allocator: std.mem.Allocator, did: []const u8, collection: []const u8, cursor: ?[]const u8, reverse: bool, limit: usize,) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const page_limit = @as(i64, @intCast(if (limit == 0) 100 else limit)); var rows = if (cursor) |rkey_cursor| blk: { if (reverse) { break :blk try conn.rows( \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.did = ? AND r.collection = ? AND r.rkey > ? \\ORDER BY r.rkey ASC \\LIMIT ? , .{ did, collection, rkey_cursor, page_limit }); } break :blk try conn.rows( \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.did = ? AND r.collection = ? AND r.rkey < ? \\ORDER BY r.rkey DESC \\LIMIT ? , .{ did, collection, rkey_cursor, page_limit }); } else blk: { if (reverse) { break :blk try conn.rows( \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.did = ? AND r.collection = ? \\ORDER BY r.rkey ASC \\LIMIT ? , .{ did, collection, page_limit }); } break :blk try conn.rows( \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq \\FROM records r \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE r.did = ? AND r.collection = ? \\ORDER BY r.rkey DESC \\LIMIT ? , .{ did, collection, page_limit }); }; defer rows.deinit();
var out: std.Io.Writer.Allocating = .init(allocator); defer out.deinit(); try out.writer.writeAll("{\"records\":["); var first = true; var last_rkey: ?[]const u8 = null; while (rows.next()) |row| { const record = try recordFromRow(row, std.heap.page_allocator); if (!first) try out.writer.writeByte(','); first = false; last_rkey = try allocator.dupe(u8, row.text(2)); const uri = try record.uri(allocator); defer allocator.free(uri); try out.writer.print( "{{\"uri\":{f},\"cid\":{f},\"value\":{s}}}", .{ std.json.fmt(uri, .{}), std.json.fmt(record.cid, .{}), record.value_json }, ); } if (rows.err) |err| return err; try out.writer.writeByte(']'); if (last_rkey) |rkey| { try out.writer.print(",\"cursor\":{f}", .{std.json.fmt(rkey, .{})}); } try out.writer.writeByte('}'); return out.toOwnedSlice();}
pub fn generateRkey(allocator: std.mem.Allocator) ![]const u8 { return nextRkey(allocator);}
pub fn prepareRecordValue( allocator: std.mem.Allocator, collection: []const u8, rkey: []const u8, value: std.json.Value,) !PreparedRecord { if (zat.Nsid.parse(collection) == null) return Error.InvalidCollection; if (zat.Rkey.parse(rkey) == null) return Error.InvalidRecordKey; try validateRecordForWrite(collection, rkey, value, .known);
const cbor_value = try jsonToDagCbor(allocator, value); const record_bytes = try zat.cbor.encodeAlloc(allocator, cbor_value); const record_cid = try zat.cbor.Cid.forDagCbor(allocator, record_bytes); var blob_cids: std.ArrayList([]const u8) = .empty; try collectBlobCids(allocator, cbor_value, &blob_cids); return .{ .cid = try cidText(allocator, record_cid.raw), .value_json = try stringifyValue(allocator, value), .validation_status = validationStatusForRecord(collection, .known), .blob_cids = try blob_cids.toOwnedSlice(allocator), };}
pub const CreateSpaceInput = struct { actor_did: []const u8, authority_did: []const u8, space_type: []const u8, skey: []const u8, is_authority: bool, managing_app: ?[]const u8, policy: []const u8, app_access_json: []const u8,};
pub fn createSpace(allocator: std.mem.Allocator, input: CreateSpaceInput) !SpaceConfig { if (zat.Did.parse(input.actor_did) == null) return Error.InvalidRepoPath; if (zat.Did.parse(input.authority_did) == null) return Error.InvalidRepoPath; if (zat.Nsid.parse(input.space_type) == null) return Error.InvalidCollection; if (zat.Rkey.parse(input.skey) == null) return Error.InvalidRecordKey; if (!validSimpleSpacePolicy(input.policy)) return Error.InvalidRecordType; try validateAppAccessJson(input.app_access_json);
const uri = try space_uri.SpaceUri.format(allocator, input.authority_did, input.space_type, input.skey);
db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
if (try getSpaceLocked(allocator, input.actor_did, uri) != null) return Error.InvalidRefreshSession;
try conn.exclusiveTransaction(); errdefer conn.rollback(); try conn.exec( \\INSERT INTO permissioned_spaces ( \\ uri, authority_did, space_type, skey, managing_app, policy, app_access_json \\) VALUES (?, ?, ?, ?, ?, ?, ?) \\ON CONFLICT(uri) DO NOTHING , .{ uri, input.authority_did, input.space_type, input.skey, input.managing_app, input.policy, input.app_access_json, }); if (input.is_authority) { try conn.exec( \\UPDATE permissioned_spaces \\SET managing_app = ?, \\ policy = ?, \\ app_access_json = ? \\WHERE uri = ? , .{ input.managing_app, input.policy, input.app_access_json, uri, }); } try conn.exec( \\INSERT INTO permissioned_space_repos (space, repo_did, set_hash, rev) \\VALUES (?, ?, NULL, NULL) \\ON CONFLICT(space, repo_did) DO NOTHING , .{ uri, input.authority_did }); try conn.exec( \\INSERT INTO permissioned_space_actor_state (space, actor_did, is_authority) \\VALUES (?, ?, ?) , .{ uri, input.actor_did, @as(i64, if (input.is_authority) 1 else 0) }); if (input.is_authority) { try conn.exec( \\INSERT INTO simplespace_members (space, member_did) \\VALUES (?, ?) \\ON CONFLICT(space, member_did) DO NOTHING , .{ uri, input.authority_did }); } try conn.commit();
return (try getSpaceLocked(allocator, input.actor_did, uri)) orelse Error.RepoNotFound;}
pub fn getSpace(allocator: std.mem.Allocator, actor_did: []const u8, uri: []const u8) !?SpaceConfig { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); return getSpaceLocked(allocator, actor_did, uri);}
pub fn listSpaces( allocator: std.mem.Allocator, actor_did: []const u8, maybe_did: ?[]const u8, maybe_type: ?[]const u8, maybe_cursor: ?[]const u8, limit: usize,) ![]SpaceConfig { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const capped_limit: i64 = @intCast(@min(if (limit == 0) 50 else limit, 100)); var query: std.Io.Writer.Allocating = .init(allocator); defer query.deinit(); try query.writer.writeAll( \\SELECT s.uri, a.is_authority \\FROM permissioned_spaces s \\JOIN permissioned_space_actor_state a ON a.space = s.uri \\WHERE a.actor_did = ? AND a.deleted_at IS NULL AND s.deleted_at IS NULL ); if (maybe_did != null) try query.writer.writeAll(" AND s.authority_did = ?"); if (maybe_type != null) try query.writer.writeAll(" AND s.space_type = ?"); if (maybe_cursor != null) try query.writer.writeAll(" AND s.uri > ?"); try query.writer.writeAll(" ORDER BY s.uri ASC LIMIT ?");
var rows = if (maybe_did) |did| if (maybe_type) |space_type| if (maybe_cursor) |cursor| try conn.rows(query.written(), .{ actor_did, did, space_type, cursor, capped_limit }) else try conn.rows(query.written(), .{ actor_did, did, space_type, capped_limit }) else if (maybe_cursor) |cursor| try conn.rows(query.written(), .{ actor_did, did, cursor, capped_limit }) else try conn.rows(query.written(), .{ actor_did, did, capped_limit }) else if (maybe_type) |space_type| if (maybe_cursor) |cursor| try conn.rows(query.written(), .{ actor_did, space_type, cursor, capped_limit }) else try conn.rows(query.written(), .{ actor_did, space_type, capped_limit }) else if (maybe_cursor) |cursor| try conn.rows(query.written(), .{ actor_did, cursor, capped_limit }) else try conn.rows(query.written(), .{ actor_did, capped_limit }); defer rows.deinit();
var spaces: std.ArrayList(SpaceConfig) = .empty; while (rows.next()) |row| { var space = (try getSpaceConfigLocked(allocator, row.text(0))) orelse continue; space.is_authority = row.int(1) != 0; try spaces.append(allocator, space); } if (rows.err) |err| return err; return spaces.toOwnedSlice(allocator);}
pub fn updateSimpleSpaceConfig( space: []const u8, managing_app: ?[]const u8, clear_managing_app: bool, maybe_policy: ?[]const u8, maybe_app_access_json: ?[]const u8,) !void { if (maybe_policy) |policy| if (!validSimpleSpacePolicy(policy)) return Error.InvalidRecordType; if (maybe_app_access_json) |app_access| try validateAppAccessJson(app_access); db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); _ = (try getSpaceConfigLocked(std.heap.page_allocator, space)) orelse return Error.RepoNotFound; if (managing_app) |value| { try conn.exec("UPDATE permissioned_spaces SET managing_app = ? WHERE uri = ?", .{ value, space }); } else if (clear_managing_app) { try conn.exec("UPDATE permissioned_spaces SET managing_app = NULL WHERE uri = ?", .{space}); } if (maybe_policy) |policy| { try conn.exec("UPDATE permissioned_spaces SET policy = ? WHERE uri = ?", .{ policy, space }); } if (maybe_app_access_json) |app_access| { try conn.exec("UPDATE permissioned_spaces SET app_access_json = ? WHERE uri = ?", .{ app_access, space }); }}
pub fn addSimpleSpaceMember(space: []const u8, did: []const u8) !void { if (zat.Did.parse(did) == null) return Error.InvalidRepoPath; db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); _ = (try getSpaceConfigLocked(std.heap.page_allocator, space)) orelse return Error.RepoNotFound; try conn.exec( \\INSERT INTO simplespace_members (space, member_did) \\VALUES (?, ?) \\ON CONFLICT(space, member_did) DO NOTHING , .{ space, did });}
pub fn removeSimpleSpaceMember(space: []const u8, did: []const u8) !void { if (zat.Did.parse(did) == null) return Error.InvalidRepoPath; db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const config = (try getSpaceConfigLocked(std.heap.page_allocator, space)) orelse return Error.RepoNotFound; if (std.mem.eql(u8, config.authority_did, did)) return Error.InvalidRecordType; try conn.exec("DELETE FROM simplespace_members WHERE space = ? AND member_did = ?", .{ space, did });}
pub fn listSimpleSpaceMembers( allocator: std.mem.Allocator, space: []const u8, maybe_cursor: ?[]const u8, limit: usize,) ![]SimpleSpaceMember { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); _ = (try getSpaceConfigLocked(std.heap.page_allocator, space)) orelse return Error.RepoNotFound; const capped_limit: i64 = @intCast(@min(if (limit == 0) 50 else limit, 100)); var rows = if (maybe_cursor) |cursor| try conn.rows( \\SELECT member_did, created_at \\FROM simplespace_members \\WHERE space = ? AND member_did > ? \\ORDER BY member_did ASC \\LIMIT ? , .{ space, cursor, capped_limit }) else try conn.rows( \\SELECT member_did, created_at \\FROM simplespace_members \\WHERE space = ? \\ORDER BY member_did ASC \\LIMIT ? , .{ space, capped_limit }); defer rows.deinit(); var members: std.ArrayList(SimpleSpaceMember) = .empty; while (rows.next()) |row| { try members.append(allocator, .{ .did = try allocator.dupe(u8, row.text(0)), .created_at = row.int(1), }); } if (rows.err) |err| return err; return members.toOwnedSlice(allocator);}
pub fn simpleSpaceHasMember(space: []const u8, did: []const u8) !bool { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const row = try conn.row( "SELECT 1 FROM simplespace_members WHERE space = ? AND member_did = ?", .{ space, did }, ); if (row == null) return false; defer row.?.deinit(); return true;}
pub fn listSpaceWriters(allocator: std.mem.Allocator, space: []const u8, maybe_cursor: ?[]const u8, limit: usize) ![]SpaceWriterState { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const authority = try conn.row( \\SELECT 1 \\FROM permissioned_spaces s \\JOIN permissioned_space_actor_state a \\ ON a.space = s.uri AND a.actor_did = s.authority_did \\WHERE s.uri = ? AND s.deleted_at IS NULL AND a.deleted_at IS NULL AND a.is_authority = 1 , .{space}); if (authority == null) return Error.RepoNotFound; authority.?.deinit(); const capped_limit: i64 = @intCast(@min(if (limit == 0) 50 else limit, 100)); var rows = if (maybe_cursor) |cursor| try conn.rows( \\SELECT repo_did, hash, rev \\FROM permissioned_space_writers \\WHERE space = ? AND repo_did > ? \\ORDER BY repo_did ASC \\LIMIT ? , .{ space, cursor, capped_limit }) else try conn.rows( \\SELECT repo_did, hash, rev \\FROM permissioned_space_writers \\WHERE space = ? \\ORDER BY repo_did ASC \\LIMIT ? , .{ space, capped_limit }); defer rows.deinit(); var repos: std.ArrayList(SpaceWriterState) = .empty; while (rows.next()) |row| { try repos.append(allocator, .{ .repo_did = try allocator.dupe(u8, row.text(0)), .hash = try allocator.dupe(u8, row.blob(1)), .rev = try allocator.dupe(u8, row.text(2)), }); } if (rows.err) |err| return err; return repos.toOwnedSlice(allocator);}
pub fn recordSpaceWriter(space: []const u8, repo_did: []const u8, rev: []const u8, hash: []const u8) !void { if (hash.len != permissioned.commit_hash_bytes) return Error.InvalidRecordType; db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\INSERT INTO permissioned_space_writers (space, repo_did, rev, hash) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(space, repo_did) DO UPDATE SET \\ rev = excluded.rev, \\ hash = excluded.hash, \\ updated_at = unixepoch() , .{ space, repo_did, rev, zqlite.blob(hash) });}
pub fn markSpaceDeleted(actor_did: []const u8, space: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec("UPDATE permissioned_space_actor_state SET deleted_at = unixepoch() WHERE space = ? AND actor_did = ?", .{ space, actor_did }); try conn.exec("UPDATE permissioned_spaces SET deleted_at = unixepoch() WHERE uri = ? AND authority_did = ?", .{ space, actor_did });}
pub fn purgeAuthoritySpaceData(space: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exclusiveTransaction(); errdefer conn.rollback(); try conn.exec("DELETE FROM permissioned_space_notify_registrations WHERE space = ?", .{space}); try conn.commit();}
pub fn registerSpaceNotification(space: []const u8, repo_did: ?[]const u8, service_endpoint: []const u8, expires_at: i64) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\INSERT INTO permissioned_space_notify_registrations (space, repo_did, service_endpoint, expires_at) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(space, repo_did, service_endpoint) DO UPDATE SET \\ expires_at = excluded.expires_at , .{ space, repo_did orelse "", service_endpoint, expires_at });}
pub fn listNotificationRecipients(allocator: std.mem.Allocator, space: []const u8, repo_did: ?[]const u8, include_space_wide: bool) ![]CredentialRecipient { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.execNoArgs("DELETE FROM permissioned_space_notify_registrations WHERE expires_at <= unixepoch()"); var rows = if (repo_did) |repo| if (include_space_wide) try conn.rows( \\SELECT repo_did, service_endpoint, expires_at \\FROM permissioned_space_notify_registrations \\WHERE space = ? AND (repo_did = '' OR repo_did = ?) \\ORDER BY repo_did ASC, service_endpoint ASC , .{ space, repo }) else try conn.rows( \\SELECT repo_did, service_endpoint, expires_at \\FROM permissioned_space_notify_registrations \\WHERE space = ? AND repo_did = ? \\ORDER BY service_endpoint ASC , .{ space, repo }) else if (include_space_wide) try conn.rows( \\SELECT '', service_endpoint, max(expires_at) \\FROM permissioned_space_notify_registrations \\WHERE space = ? \\GROUP BY service_endpoint \\ORDER BY service_endpoint ASC , .{space}) else try conn.rows( \\SELECT repo_did, service_endpoint, expires_at \\FROM permissioned_space_notify_registrations \\WHERE space = ? AND repo_did = '' \\ORDER BY service_endpoint ASC , .{space}); defer rows.deinit(); var out: std.ArrayList(CredentialRecipient) = .empty; while (rows.next()) |row| { try out.append(allocator, .{ .service_endpoint = try allocator.dupe(u8, row.text(1)), .repo_did = if (row.text(0).len == 0) null else try allocator.dupe(u8, row.text(0)), .expires_at = row.int(2), }); } if (rows.err) |err| return err; return out.toOwnedSlice(allocator);}
pub fn loadSpaceRepoBlocks(allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8) ![]permissioned.RepoRecordBlock { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); var rows = try conn.rows( \\SELECT collection, rkey, cid, value_json \\FROM permissioned_space_records \\WHERE space = ? AND repo_did = ? \\ORDER BY collection ASC, rkey ASC , .{ space, repo_did }); defer rows.deinit(); var blocks: std.ArrayList(permissioned.RepoRecordBlock) = .empty; var scratch = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer scratch.deinit(); while (rows.next()) |row| { _ = scratch.reset(.retain_capacity); const scratch_allocator = scratch.allocator(); const parsed = try std.json.parseFromSlice(std.json.Value, scratch_allocator, row.blob(3), .{}); const value = try jsonToDagCbor(scratch_allocator, parsed.value); const data = try zat.cbor.encodeAlloc(allocator, value); const cid_text = row.text(2); if (cid_text.len == 0 or cid_text[0] != 'b') return Error.InvalidDagCbor; try blocks.append(allocator, .{ .path = try std.fmt.allocPrint(allocator, "{s}/{s}", .{ row.text(0), row.text(1) }), .cid = .{ .raw = try zat.multibase.base32lower.decode(allocator, cid_text[1..]) }, .data = data, }); } if (rows.err) |err| return err; return blocks.toOwnedSlice(allocator);}
pub const SpaceReplayToken = struct { jti: []const u8, expires_at: i64,};
pub fn consumeSpaceCredentialExchange(delegation: SpaceReplayToken, attestation: ?SpaceReplayToken) !bool { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exclusiveTransaction(); errdefer conn.rollback(); try conn.execNoArgs("DELETE FROM permissioned_space_used_delegations WHERE expires_at < unixepoch()"); try conn.execNoArgs("DELETE FROM permissioned_space_used_client_attestations WHERE expires_at < unixepoch()"); conn.exec( \\INSERT INTO permissioned_space_used_delegations (jti, expires_at) \\VALUES (?, ?) , .{ delegation.jti, delegation.expires_at }) catch |err| switch (err) { error.Constraint => { conn.rollback(); return false; }, else => return err, }; if (attestation) |token| { conn.exec( \\INSERT INTO permissioned_space_used_client_attestations (jti, expires_at) \\VALUES (?, ?) , .{ token.jti, token.expires_at }) catch |err| switch (err) { error.Constraint => { conn.rollback(); return false; }, else => return err, }; } try conn.commit(); return true;}
pub fn putSpaceRecord( allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8, prepared: PreparedRecord,) !SpaceRecord { const results = try applySpaceWrites(allocator, space, repo_did, &.{.{ .put = .{ .collection = collection, .rkey = rkey, .prepared = prepared, } }}); return switch (results[0]) { .create => |record| record, .update => |record| record, .delete => Error.MissingRecord, };}
pub fn createSpaceRecord( allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8, prepared: PreparedRecord,) !SpaceRecord { const results = try applySpaceWrites(allocator, space, repo_did, &.{.{ .create = .{ .collection = collection, .rkey = rkey, .prepared = prepared, } }}); return switch (results[0]) { .create => |record| record, .update => |record| record, .delete => Error.MissingRecord, };}
pub fn deleteSpaceRecord( allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8,) !void { _ = try applySpaceWrites(allocator, space, repo_did, &.{.{ .delete = .{ .collection = collection, .rkey = rkey } }});}
pub fn getSpaceRecord( allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8,) !?SpaceRecord { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); return getSpaceRecordLocked(allocator, space, repo_did, collection, rkey);}
pub fn listSpaceRecords( allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8, maybe_collection: ?[]const u8, maybe_cursor: ?[]const u8, reverse: bool, limit: usize, include_values: bool,) ![]SpaceRecordRef { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized();
const capped_limit: i64 = @intCast(@min(if (limit == 0) 50 else limit, 100)); const op = if (reverse) ">" else "<"; const direction = if (reverse) "ASC" else "DESC"; const cursor = maybe_cursor orelse ""; const parsed_cursor = parseSpaceRecordCursor(cursor); const cursor_collection = parsed_cursor.collection; const cursor_rkey = parsed_cursor.rkey;
var query = std.Io.Writer.Allocating.init(allocator); defer query.deinit(); try query.writer.print( "SELECT collection, rkey, cid, {s} FROM permissioned_space_records WHERE space = ? AND repo_did = ?", .{if (include_values) "value_json" else "NULL"}, ); if (maybe_collection != null) try query.writer.writeAll(" AND collection = ?"); if (maybe_cursor != null and cursor_collection.len > 0 and cursor_rkey.len > 0) { try query.writer.print(" AND (collection, rkey) {s} (?, ?)", .{op}); } try query.writer.print(" ORDER BY collection {s}, rkey {s} LIMIT ?", .{ direction, direction });
var rows = if (maybe_collection) |collection| if (maybe_cursor != null and cursor_collection.len > 0 and cursor_rkey.len > 0) try conn.rows(query.written(), .{ space, repo_did, collection, cursor_collection, cursor_rkey, capped_limit }) else try conn.rows(query.written(), .{ space, repo_did, collection, capped_limit }) else if (maybe_cursor != null and cursor_collection.len > 0 and cursor_rkey.len > 0) try conn.rows(query.written(), .{ space, repo_did, cursor_collection, cursor_rkey, capped_limit }) else try conn.rows(query.written(), .{ space, repo_did, capped_limit }); defer rows.deinit();
var records: std.ArrayList(SpaceRecordRef) = .empty; while (rows.next()) |row| { try records.append(allocator, .{ .collection = try allocator.dupe(u8, row.text(0)), .rkey = try allocator.dupe(u8, row.text(1)), .cid = try allocator.dupe(u8, row.text(2)), .value_json = if (row.nullableText(3)) |value| try allocator.dupe(u8, value) else null, }); } if (rows.err) |err| return err; return records.toOwnedSlice(allocator);}
pub fn applySpaceWrites( allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8, ops: []const SpaceWriteOp,) ![]SpaceWriteResult { if (zat.Did.parse(repo_did) == null) return Error.InvalidRepoPath; for (ops) |op| switch (op) { inline .create, .put, .update => |write| { if (zat.Nsid.parse(write.collection) == null) return Error.InvalidCollection; if (zat.Rkey.parse(write.rkey) == null) return Error.InvalidRecordKey; }, .delete => |write| { if (zat.Nsid.parse(write.collection) == null) return Error.InvalidCollection; if (zat.Rkey.parse(write.rkey) == null) return Error.InvalidRecordKey; }, };
db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const space_config = (try getSpaceConfigLocked(allocator, space)) orelse return Error.RepoNotFound;
try conn.exclusiveTransaction(); errdefer conn.rollback();
const state = try getRepoStateLocked(allocator, space, repo_did); var set_hash = if (state.set_hash) |bytes| try permissioned.LtHash.fromBytes(bytes) else permissioned.LtHash{}; const rev = try nextRkeyLocked(allocator);
var results: std.ArrayList(SpaceWriteResult) = .empty; var idx: i64 = 0; for (ops) |op| { switch (op) { .create => |write| { if (try getSpaceRecordLocked(allocator, space, repo_did, write.collection, write.rkey) != null) return Error.InvalidRefreshSession; try addRecordElement(allocator, &set_hash, write.collection, write.rkey, write.prepared.cid); try upsertSpaceRecordLocked(space, repo_did, write.collection, write.rkey, write.prepared, rev); try syncSpaceRecordBlobRefsLocked(space, repo_did, write.collection, write.rkey, write.prepared.blob_cids); try insertRecordOplogLocked(space, repo_did, rev, idx, "create", write.collection, write.rkey, write.prepared.cid, null); try results.append(allocator, .{ .create = (try getSpaceRecordLocked(allocator, space, repo_did, write.collection, write.rkey)) orelse return Error.MissingRecord }); }, .put => |write| { const existing = try getSpaceRecordLocked(allocator, space, repo_did, write.collection, write.rkey); if (existing) |old| try removeRecordElement(allocator, &set_hash, old.collection, old.rkey, old.cid); try addRecordElement(allocator, &set_hash, write.collection, write.rkey, write.prepared.cid); try upsertSpaceRecordLocked(space, repo_did, write.collection, write.rkey, write.prepared, rev); try syncSpaceRecordBlobRefsLocked(space, repo_did, write.collection, write.rkey, write.prepared.blob_cids); const action = if (existing == null) "create" else "update"; try insertRecordOplogLocked(space, repo_did, rev, idx, action, write.collection, write.rkey, write.prepared.cid, if (existing) |old| old.cid else null); try results.append(allocator, .{ .update = (try getSpaceRecordLocked(allocator, space, repo_did, write.collection, write.rkey)) orelse return Error.MissingRecord }); }, .update => |write| { const existing = (try getSpaceRecordLocked(allocator, space, repo_did, write.collection, write.rkey)) orelse return Error.MissingRecord; try removeRecordElement(allocator, &set_hash, existing.collection, existing.rkey, existing.cid); try addRecordElement(allocator, &set_hash, write.collection, write.rkey, write.prepared.cid); try upsertSpaceRecordLocked(space, repo_did, write.collection, write.rkey, write.prepared, rev); try syncSpaceRecordBlobRefsLocked(space, repo_did, write.collection, write.rkey, write.prepared.blob_cids); try insertRecordOplogLocked(space, repo_did, rev, idx, "update", write.collection, write.rkey, write.prepared.cid, existing.cid); try results.append(allocator, .{ .update = (try getSpaceRecordLocked(allocator, space, repo_did, write.collection, write.rkey)) orelse return Error.MissingRecord }); }, .delete => |write| { const existing = (try getSpaceRecordLocked(allocator, space, repo_did, write.collection, write.rkey)) orelse return Error.MissingRecord; try removeRecordElement(allocator, &set_hash, existing.collection, existing.rkey, existing.cid); try deleteSpaceRecordBlobRefsLocked(space, repo_did, write.collection, write.rkey); try conn.exec( \\DELETE FROM permissioned_space_records \\WHERE space = ? AND repo_did = ? AND collection = ? AND rkey = ? , .{ space, repo_did, write.collection, write.rkey }); try insertRecordOplogLocked(space, repo_did, rev, idx, "delete", write.collection, write.rkey, null, existing.cid); try results.append(allocator, .{ .delete = {} }); }, } idx += 1; }
try conn.exec( \\INSERT INTO permissioned_space_repos (space, repo_did, set_hash, rev) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(space, repo_did) DO UPDATE SET \\ set_hash = excluded.set_hash, \\ rev = excluded.rev , .{ space, repo_did, zqlite.blob(&set_hash.bytes), rev }); try conn.exec( \\INSERT INTO permissioned_space_actor_state (space, actor_did, is_authority) \\VALUES (?, ?, 0) \\ON CONFLICT(space, actor_did) DO NOTHING , .{ space, repo_did }); if (space_config.is_authority) { const digest = set_hash.digest(); try conn.exec( \\INSERT INTO permissioned_space_writers (space, repo_did, rev, hash) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(space, repo_did) DO UPDATE SET \\ rev = excluded.rev, \\ hash = excluded.hash, \\ updated_at = unixepoch() , .{ space, repo_did, rev, zqlite.blob(&digest) }); } try conn.commit(); return results.toOwnedSlice(allocator);}
pub fn getSpaceRepoState(allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8) !SpaceState { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); return getRepoStateLocked(allocator, space, repo_did);}
pub fn listSpaceRecordOplog(allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8, since: ?[]const u8, limit: usize, include_values: bool) ![]SpaceRecordOplogEntry { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const capped_limit: i64 = @intCast(@min(if (limit == 0) 100 else limit, 1000)); var rows = if (since) |rev| try conn.rows( \\SELECT o.rev, o.idx, o.action, o.repo_did, o.collection, o.rkey, o.cid, o.prev, \\ CASE WHEN ? AND r.cid = o.cid THEN r.value_json END \\FROM permissioned_space_record_oplog o \\LEFT JOIN permissioned_space_records r \\ ON r.space = o.space AND r.repo_did = o.repo_did \\ AND r.collection = o.collection AND r.rkey = o.rkey \\WHERE o.space = ? AND o.repo_did = ? AND o.rev > ? \\ORDER BY o.rev ASC, o.idx ASC \\LIMIT ? , .{ include_values, space, repo_did, rev, capped_limit }) else try conn.rows( \\SELECT o.rev, o.idx, o.action, o.repo_did, o.collection, o.rkey, o.cid, o.prev, \\ CASE WHEN ? AND r.cid = o.cid THEN r.value_json END \\FROM permissioned_space_record_oplog o \\LEFT JOIN permissioned_space_records r \\ ON r.space = o.space AND r.repo_did = o.repo_did \\ AND r.collection = o.collection AND r.rkey = o.rkey \\WHERE o.space = ? AND o.repo_did = ? \\ORDER BY o.rev ASC, o.idx ASC \\LIMIT ? , .{ include_values, space, repo_did, capped_limit }); defer rows.deinit(); var out: std.ArrayList(SpaceRecordOplogEntry) = .empty; while (rows.next()) |row| { try out.append(allocator, .{ .rev = try allocator.dupe(u8, row.text(0)), .idx = row.int(1), .action = try allocator.dupe(u8, row.text(2)), .repo_did = try allocator.dupe(u8, row.text(3)), .collection = try allocator.dupe(u8, row.text(4)), .rkey = try allocator.dupe(u8, row.text(5)), .cid = if (row.nullableText(6)) |cid| try allocator.dupe(u8, cid) else null, .prev = if (row.nullableText(7)) |prev| try allocator.dupe(u8, prev) else null, .value_json = if (row.nullableBlob(8)) |value| try allocator.dupe(u8, value) else null, }); } if (rows.err) |err| return err; return out.toOwnedSlice(allocator);}
fn getRepoStateLocked(allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8) !SpaceState { const row = try conn.row( \\SELECT set_hash, rev \\FROM permissioned_space_repos \\WHERE space = ? AND repo_did = ? , .{ space, repo_did }); if (row == null) return .{ .set_hash = null, .rev = null }; defer row.?.deinit(); return .{ .set_hash = if (row.?.nullableBlob(0)) |bytes| try allocator.dupe(u8, bytes) else null, .rev = if (row.?.nullableText(1)) |rev| try allocator.dupe(u8, rev) else null, };}
fn addRecordElement(allocator: std.mem.Allocator, hash: *permissioned.LtHash, collection: []const u8, rkey: []const u8, cid: []const u8) !void { const element = try permissioned.recordElement(allocator, collection, rkey, cid); hash.add(element);}
fn removeRecordElement(allocator: std.mem.Allocator, hash: *permissioned.LtHash, collection: []const u8, rkey: []const u8, cid: []const u8) !void { const element = try permissioned.recordElement(allocator, collection, rkey, cid); hash.remove(element);}
fn upsertSpaceRecordLocked( space: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8, prepared: PreparedRecord, rev: []const u8,) !void { try conn.exec( \\INSERT INTO permissioned_space_records (space, repo_did, collection, rkey, cid, value_json, validation_status, repo_rev, updated_at) \\VALUES (?, ?, ?, ?, ?, ?, ?, ?, unixepoch()) \\ON CONFLICT(space, repo_did, collection, rkey) DO UPDATE SET \\ cid = excluded.cid, \\ value_json = excluded.value_json, \\ validation_status = excluded.validation_status, \\ repo_rev = excluded.repo_rev, \\ updated_at = excluded.updated_at , .{ space, repo_did, collection, rkey, prepared.cid, prepared.value_json, prepared.validation_status, rev });}
fn syncSpaceRecordBlobRefsLocked( space: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8, blob_cids: []const []const u8,) !void { try deleteSpaceRecordBlobRefsLocked(space, repo_did, collection, rkey); for (blob_cids) |blob_cid| { try conn.exec( \\INSERT INTO permissioned_space_record_blobs (space, repo_did, collection, rkey, blob_cid) \\VALUES (?, ?, ?, ?, ?) \\ON CONFLICT(space, repo_did, collection, rkey, blob_cid) DO UPDATE SET \\ blob_cid = excluded.blob_cid , .{ space, repo_did, collection, rkey, blob_cid }); }}
fn deleteSpaceRecordBlobRefsLocked( space: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8,) !void { try conn.exec( \\DELETE FROM permissioned_space_record_blobs \\WHERE space = ? AND repo_did = ? AND collection = ? AND rkey = ? , .{ space, repo_did, collection, rkey });}
fn insertRecordOplogLocked( space: []const u8, repo_did: []const u8, rev: []const u8, idx: i64, action: []const u8, collection: []const u8, rkey: []const u8, cid: ?[]const u8, prev: ?[]const u8,) !void { try conn.exec( \\INSERT INTO permissioned_space_record_oplog (space, repo_did, rev, idx, action, collection, rkey, cid, prev) \\VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) , .{ space, repo_did, rev, idx, action, collection, rkey, cid, prev });}
fn getSpaceConfigLocked(allocator: std.mem.Allocator, uri: []const u8) !?SpaceConfig { const row = try conn.row( \\SELECT s.uri, s.authority_did, s.space_type, s.skey, s.managing_app, s.policy, s.app_access_json, \\ EXISTS ( \\ SELECT 1 FROM permissioned_space_actor_state a \\ WHERE a.space = s.uri AND a.actor_did = s.authority_did \\ AND a.is_authority = 1 AND a.deleted_at IS NULL \\ ) \\FROM permissioned_spaces s \\WHERE s.uri = ? AND s.deleted_at IS NULL , .{uri}); if (row == null) return null; defer row.?.deinit();
return .{ .uri = try allocator.dupe(u8, row.?.text(0)), .authority_did = try allocator.dupe(u8, row.?.text(1)), .space_type = try allocator.dupe(u8, row.?.text(2)), .skey = try allocator.dupe(u8, row.?.text(3)), .managing_app = if (row.?.nullableText(4)) |value| try allocator.dupe(u8, value) else null, .policy = try allocator.dupe(u8, row.?.text(5)), .app_access_json = try allocator.dupe(u8, row.?.text(6)), .is_authority = row.?.int(7) != 0, .deleted_at = null, };}
fn getSpaceLocked(allocator: std.mem.Allocator, actor_did: []const u8, uri: []const u8) !?SpaceConfig { const row = try conn.row( \\SELECT s.uri, s.authority_did, s.space_type, s.skey, s.managing_app, s.policy, s.app_access_json, \\ a.is_authority, a.deleted_at \\FROM permissioned_spaces s \\JOIN permissioned_space_actor_state a ON a.space = s.uri AND a.actor_did = ? \\WHERE s.uri = ? AND s.deleted_at IS NULL AND a.deleted_at IS NULL , .{ actor_did, uri }); if (row == null) return null; defer row.?.deinit();
return .{ .uri = try allocator.dupe(u8, row.?.text(0)), .authority_did = try allocator.dupe(u8, row.?.text(1)), .space_type = try allocator.dupe(u8, row.?.text(2)), .skey = try allocator.dupe(u8, row.?.text(3)), .managing_app = if (row.?.nullableText(4)) |value| try allocator.dupe(u8, value) else null, .policy = try allocator.dupe(u8, row.?.text(5)), .app_access_json = try allocator.dupe(u8, row.?.text(6)), .is_authority = row.?.int(7) != 0, .deleted_at = if (row.?.nullableInt(8)) |value| value else null, };}
fn getSpaceRecordLocked( allocator: std.mem.Allocator, space: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8,) !?SpaceRecord { const row = try conn.row( \\SELECT space, repo_did, collection, rkey, cid, value_json, validation_status, repo_rev, updated_at \\FROM permissioned_space_records \\WHERE space = ? AND repo_did = ? AND collection = ? AND rkey = ? , .{ space, repo_did, collection, rkey }); if (row == null) return null; defer row.?.deinit(); return .{ .space = try allocator.dupe(u8, row.?.text(0)), .repo_did = try allocator.dupe(u8, row.?.text(1)), .collection = try allocator.dupe(u8, row.?.text(2)), .rkey = try allocator.dupe(u8, row.?.text(3)), .cid = try allocator.dupe(u8, row.?.text(4)), .value_json = try allocator.dupe(u8, row.?.text(5)), .validation_status = try allocator.dupe(u8, row.?.text(6)), .repo_rev = try allocator.dupe(u8, row.?.text(7)), .updated_at = row.?.int(8), };}
fn parseSpaceRecordCursor(cursor: []const u8) struct { collection: []const u8, rkey: []const u8 } { const slash = std.mem.indexOfScalar(u8, cursor, '/') orelse return .{ .collection = "", .rkey = "" }; if (slash == 0 or slash + 1 >= cursor.len) return .{ .collection = "", .rkey = "" }; return .{ .collection = cursor[0..slash], .rkey = cursor[slash + 1 ..] };}
fn validSimpleSpacePolicy(policy: []const u8) bool { return std.mem.eql(u8, policy, "member-list") or std.mem.eql(u8, policy, "public") or std.mem.eql(u8, policy, "managing-app");}
fn validateAppAccessJson(raw: []const u8) !void { var parsed = try std.json.parseFromSlice(std.json.Value, std.heap.page_allocator, raw, .{}); defer parsed.deinit(); const object = switch (parsed.value) { .object => |object| object, else => return Error.InvalidRecordType, }; const kind = switch (object.get("type") orelse return Error.InvalidRecordType) { .string => |value| value, else => return Error.InvalidRecordType, }; if (std.mem.eql(u8, kind, "open")) return; if (!std.mem.eql(u8, kind, "allowList")) return Error.InvalidRecordType; const allowed = object.get("allowed") orelse return Error.InvalidRecordType; if (allowed != .array) return Error.InvalidRecordType; for (allowed.array.items) |item| { if (item != .string) return Error.InvalidRecordType; }}
fn migrate() !void { inline for (schema_statements) |sql| try conn.execNoArgs(sql); try migratePermissionedDataTables(); const had_deactivated_at = try accountsColumnExistsLocked("deactivated_at"); conn.execNoArgs("ALTER TABLE accounts ADD COLUMN email TEXT") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN email_confirmed_at INTEGER") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN auth_code TEXT") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN auth_code_expires_at INTEGER") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN pending_email TEXT") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN invites_disabled INTEGER NOT NULL DEFAULT 0") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN activated_at INTEGER") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN account_status TEXT NOT NULL DEFAULT 'active'") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN account_status_ref TEXT") catch {}; conn.execNoArgs("UPDATE accounts SET account_status = 'deactivated' WHERE activated_at IS NULL AND account_status = 'active'") catch {}; if (had_deactivated_at) { conn.execNoArgs("UPDATE accounts SET account_status = 'deactivated' WHERE deactivated_at IS NOT NULL AND account_status = 'active'") catch {}; } conn.execNoArgs("ALTER TABLE accounts ADD COLUMN signing_key_type TEXT") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN signing_key BLOB") catch {}; conn.execNoArgs("ALTER TABLE oauth_requests ADD COLUMN login_hint TEXT") catch {}; conn.execNoArgs("ALTER TABLE oauth_requests ADD COLUMN response_mode TEXT NOT NULL DEFAULT 'query'") catch {}; conn.execNoArgs("ALTER TABLE oauth_requests ADD COLUMN auth_method TEXT") catch {}; conn.execNoArgs("ALTER TABLE oauth_tokens ADD COLUMN auth_method TEXT") catch {}; conn.execNoArgs("ALTER TABLE oauth_tokens ADD COLUMN dpop_jkt TEXT") catch {}; conn.execNoArgs("ALTER TABLE oauth_tokens ADD COLUMN family_id TEXT") catch {}; conn.execNoArgs("ALTER TABLE oauth_tokens ADD COLUMN access_expires_at INTEGER") catch {}; conn.execNoArgs("ALTER TABLE oauth_tokens ADD COLUMN refresh_expires_at INTEGER") catch {}; conn.execNoArgs("ALTER TABLE oauth_tokens ADD COLUMN previous_refresh_token TEXT") catch {}; conn.execNoArgs("ALTER TABLE oauth_tokens ADD COLUMN previous_refresh_expires_at INTEGER") catch {}; conn.execNoArgs("UPDATE oauth_tokens SET family_id = access_token WHERE family_id IS NULL") catch {}; conn.execNoArgs("UPDATE oauth_tokens SET access_expires_at = expires_at WHERE access_expires_at IS NULL") catch {}; conn.execNoArgs("UPDATE oauth_tokens SET refresh_expires_at = expires_at WHERE refresh_expires_at IS NULL") catch {}; conn.execNoArgs("ALTER TABLE session_tokens ADD COLUMN app_password_name TEXT") catch {}; conn.execNoArgs("ALTER TABLE permissioned_space_records ADD COLUMN repo_rev TEXT NOT NULL DEFAULT ''") catch {}; try migrateBlobTable(); try migrateAppPreferencesTable(); conn.execNoArgs("ALTER TABLE repo_blocks ADD COLUMN repo_rev TEXT") catch {}; inline for (post_schema_statements) |sql| try conn.execNoArgs(sql); try migrateRepoBlockRevs(); try migrateRecordsTable(); try validateRecordBlocksPresent(); try migrateSeqEventsSync11(); try migrateSeqEventsSinceRev();}
fn migratePermissionedDataTables() !void { try ensureMigrationTable(); try migratePermissionedSpacesAuthorityColumns(); try migratePermissionedSpaceAuthorityFlag(); try migrateSimpleSpaceConfig(); try migratePermissionedRecordOplogRepo(); try migratePermissionedSpaceUris(); try migratePermissionedSpaceRepoHashes(); try migratePermissionedSpaceWriters(); try migratePermissionedSpaceNotifyRegistrations();}
fn migratePermissionedSpaceNotifyRegistrations() !void { const name = "permissioned-space-notify-registration-v1"; if (try migrationApplied(name)) return; // The prior table was never populated through a protocol endpoint and // cannot represent repo-scoped or expiring registrations. try conn.execNoArgs("DROP TABLE IF EXISTS permissioned_space_credential_recipients"); try conn.execNoArgs( \\CREATE TABLE IF NOT EXISTS permissioned_space_notify_registrations ( \\ space TEXT NOT NULL, \\ repo_did TEXT NOT NULL DEFAULT '', \\ service_endpoint TEXT NOT NULL, \\ expires_at INTEGER NOT NULL, \\ PRIMARY KEY (space, repo_did, service_endpoint) \\) ); try markMigrationApplied(name);}
fn migratePermissionedSpaceUris() !void { const name = "permissioned-space-at-uri"; if (try migrationApplied(name)) return;
try conn.execNoArgs("PRAGMA foreign_keys = OFF"); errdefer conn.execNoArgs("PRAGMA foreign_keys = ON") catch {}; try conn.exclusiveTransaction(); errdefer conn.rollback();
const tables = [_][]const u8{ "permissioned_space_actor_state", "simplespace_members", "permissioned_space_records", "permissioned_space_record_blobs", "permissioned_space_repos", "permissioned_space_writers", "permissioned_space_record_oplog", "permissioned_space_notify_registrations", }; inline for (tables) |table| { try conn.execNoArgs( "UPDATE " ++ table ++ " SET space = (" ++ "SELECT 'at://' || authority_did || '/space/' || space_type || '/' || skey " ++ "FROM permissioned_spaces WHERE uri = " ++ table ++ ".space" ++ ") WHERE space LIKE 'ats://%'", ); } try conn.execNoArgs( \\UPDATE permissioned_spaces \\SET uri = 'at://' || authority_did || '/space/' || space_type || '/' || skey \\WHERE uri LIKE 'ats://%' ); try markMigrationApplied(name); try conn.commit(); try conn.execNoArgs("PRAGMA foreign_keys = ON");}
fn migratePermissionedSpaceRepoHashes() !void { const name = "permissioned-space-record-element-v1"; if (try migrationApplied(name)) return;
var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); const allocator = arena.allocator(); const RepoIdentity = struct { space: []const u8, did: []const u8, has_rev: bool }; var identities: std.ArrayList(RepoIdentity) = .empty; var repos = try conn.rows("SELECT space, repo_did, rev FROM permissioned_space_repos", .{}); while (repos.next()) |repo| { try identities.append(allocator, .{ .space = try allocator.dupe(u8, repo.text(0)), .did = try allocator.dupe(u8, repo.text(1)), .has_rev = repo.nullableText(2) != null, }); } if (repos.err) |err| return err; repos.deinit();
for (identities.items) |identity| { var hash: permissioned.LtHash = .{}; var record_count: usize = 0; var records = try conn.rows( \\SELECT collection, rkey, cid \\FROM permissioned_space_records \\WHERE space = ? AND repo_did = ? , .{ identity.space, identity.did }); defer records.deinit(); while (records.next()) |record| { try addRecordElement(allocator, &hash, record.text(0), record.text(1), record.text(2)); record_count += 1; } if (records.err) |err| return err; if (!identity.has_rev and record_count == 0) continue; try conn.exec( \\UPDATE permissioned_space_repos \\SET set_hash = ? \\WHERE space = ? AND repo_did = ? , .{ zqlite.blob(&hash.bytes), identity.space, identity.did }); } try markMigrationApplied(name);}
fn migratePermissionedSpaceWriters() !void { const name = "permissioned-space-writer-set-v1"; if (try migrationApplied(name)) return;
var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); const allocator = arena.allocator(); const ExistingRepo = struct { space: []const u8, did: []const u8, rev: []const u8, state: []const u8, }; var existing: std.ArrayList(ExistingRepo) = .empty; var rows = try conn.rows( \\SELECT r.space, r.repo_did, r.rev, r.set_hash \\FROM permissioned_space_repos r \\JOIN permissioned_spaces s ON s.uri = r.space \\JOIN permissioned_space_actor_state a \\ ON a.space = s.uri AND a.actor_did = s.authority_did \\WHERE r.rev IS NOT NULL AND r.set_hash IS NOT NULL \\ AND s.deleted_at IS NULL AND a.deleted_at IS NULL AND a.is_authority = 1 , .{}); while (rows.next()) |row| { try existing.append(allocator, .{ .space = try allocator.dupe(u8, row.text(0)), .did = try allocator.dupe(u8, row.text(1)), .rev = try allocator.dupe(u8, row.text(2)), .state = try allocator.dupe(u8, row.blob(3)), }); } if (rows.err) |err| return err; rows.deinit();
for (existing.items) |repo| { const state = try permissioned.LtHash.fromBytes(repo.state); const digest = state.digest(); try conn.exec( \\INSERT INTO permissioned_space_writers (space, repo_did, rev, hash) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(space, repo_did) DO UPDATE SET \\ rev = excluded.rev, \\ hash = excluded.hash, \\ updated_at = unixepoch() , .{ repo.space, repo.did, repo.rev, zqlite.blob(&digest) }); } try markMigrationApplied(name);}
fn migrateSimpleSpaceConfig() !void { conn.execNoArgs("ALTER TABLE permissioned_spaces ADD COLUMN policy TEXT NOT NULL DEFAULT 'member-list'") catch {}; conn.execNoArgs("ALTER TABLE permissioned_spaces ADD COLUMN app_access_json TEXT NOT NULL DEFAULT '{\"type\":\"open\"}'") catch {}; conn.execNoArgs( \\UPDATE permissioned_spaces \\SET policy = CASE WHEN is_public = 1 THEN 'public' ELSE 'member-list' END \\WHERE EXISTS ( \\ SELECT 1 FROM pragma_table_info('permissioned_spaces') WHERE name = 'is_public' \\) ) catch {}; conn.execNoArgs( \\UPDATE permissioned_spaces \\SET app_access_json = CASE \\ WHEN app_access_mode = 'deny' AND app_exceptions_json != '[]' \\ THEN '{"type":"allowList","allowed":[]}' \\ ELSE '{"type":"open"}' \\END \\WHERE EXISTS ( \\ SELECT 1 FROM pragma_table_info('permissioned_spaces') WHERE name = 'app_access_mode' \\) ) catch {}; try conn.execNoArgs( \\CREATE TABLE IF NOT EXISTS simplespace_members ( \\ space TEXT NOT NULL REFERENCES permissioned_spaces(uri) ON DELETE CASCADE, \\ member_did TEXT NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ PRIMARY KEY (space, member_did) \\) ); try conn.execNoArgs("CREATE INDEX IF NOT EXISTS simplespace_members_member_idx ON simplespace_members (member_did, space)"); try conn.execNoArgs( \\INSERT OR IGNORE INTO simplespace_members (space, member_did) \\SELECT uri, authority_did \\FROM permissioned_spaces );}
fn ensureMigrationTable() !void { try conn.execNoArgs( \\CREATE TABLE IF NOT EXISTS zds_migrations ( \\ name TEXT PRIMARY KEY, \\ applied_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) );}
fn accountsColumnExistsLocked(column: []const u8) !bool { var rows = try conn.rows("PRAGMA table_info(accounts)", .{}); defer rows.deinit(); while (rows.next()) |row| { if (std.mem.eql(u8, row.text(1), column)) return true; } if (rows.err) |err| return err; return false;}
fn permissionedSpacesColumnExistsLocked(column: []const u8) !bool { var rows = try conn.rows("PRAGMA table_info(permissioned_spaces)", .{}); defer rows.deinit(); while (rows.next()) |row| { if (std.mem.eql(u8, row.text(1), column)) return true; } if (rows.err) |err| return err; return false;}
fn permissionedSpaceActorStateColumnExistsLocked(column: []const u8) !bool { var rows = try conn.rows("PRAGMA table_info(permissioned_space_actor_state)", .{}); defer rows.deinit(); while (rows.next()) |row| { if (std.mem.eql(u8, row.text(1), column)) return true; } if (rows.err) |err| return err; return false;}
fn migrationApplied(name: []const u8) !bool { const row = try conn.row("SELECT 1 FROM zds_migrations WHERE name = ?", .{name}); if (row == null) return false; defer row.?.deinit(); return true;}
fn markMigrationApplied(name: []const u8) !void { try conn.exec("INSERT OR IGNORE INTO zds_migrations (name) VALUES (?)", .{name});}
fn migratePermissionedSpacesAuthorityColumns() !void { const name = "permissioned-spaces-authority-columns"; if (try migrationApplied(name)) return; if (!try permissionedSpacesColumnExistsLocked("owner_did")) { try conn.execNoArgs("CREATE INDEX IF NOT EXISTS permissioned_spaces_authority_idx ON permissioned_spaces (authority_did, space_type, uri)"); try markMigrationApplied(name); return; } try conn.execNoArgs("DROP TABLE IF EXISTS permissioned_spaces_next"); try conn.execNoArgs( \\CREATE TABLE permissioned_spaces_next ( \\ uri TEXT PRIMARY KEY, \\ authority_did TEXT NOT NULL, \\ space_type TEXT NOT NULL, \\ skey TEXT NOT NULL, \\ managing_app TEXT, \\ is_public INTEGER NOT NULL DEFAULT 0, \\ app_access_mode TEXT NOT NULL DEFAULT 'allow', \\ app_exceptions_json TEXT NOT NULL DEFAULT '[]', \\ is_authority INTEGER NOT NULL DEFAULT 1, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ deleted_at INTEGER, \\ UNIQUE(authority_did, space_type, skey) \\) ); try conn.execNoArgs( \\INSERT INTO permissioned_spaces_next ( \\ uri, authority_did, space_type, skey, managing_app, is_public, app_access_mode, \\ app_exceptions_json, is_authority, created_at, deleted_at \\) \\SELECT uri, owner_did, space_type, skey, managing_app, is_public, app_access_mode, \\ app_exceptions_json, is_owner, created_at, deleted_at \\FROM permissioned_spaces ); try conn.execNoArgs("PRAGMA foreign_keys = OFF"); errdefer conn.execNoArgs("PRAGMA foreign_keys = ON") catch {}; try conn.execNoArgs("DROP TABLE permissioned_spaces"); try conn.execNoArgs("ALTER TABLE permissioned_spaces_next RENAME TO permissioned_spaces"); try conn.execNoArgs("PRAGMA foreign_keys = ON"); try conn.execNoArgs("CREATE INDEX IF NOT EXISTS permissioned_spaces_authority_idx ON permissioned_spaces (authority_did, space_type, uri)"); try markMigrationApplied(name);}
fn migratePermissionedSpaceAuthorityFlag() !void { const name = "permissioned-space-authority-flag"; if (try migrationApplied(name)) return; if (try permissionedSpaceActorStateColumnExistsLocked("is_owner")) { try conn.execNoArgs("DROP TABLE IF EXISTS permissioned_space_actor_state_next"); try conn.execNoArgs( \\CREATE TABLE permissioned_space_actor_state_next ( \\ space TEXT NOT NULL REFERENCES permissioned_spaces(uri) ON DELETE CASCADE, \\ actor_did TEXT NOT NULL, \\ is_authority INTEGER NOT NULL DEFAULT 0, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ deleted_at INTEGER, \\ PRIMARY KEY (space, actor_did) \\) ); try conn.execNoArgs( \\INSERT INTO permissioned_space_actor_state_next ( \\ space, actor_did, is_authority, created_at, deleted_at \\) \\SELECT space, actor_did, is_owner, created_at, deleted_at \\FROM permissioned_space_actor_state ); try conn.execNoArgs("PRAGMA foreign_keys = OFF"); errdefer conn.execNoArgs("PRAGMA foreign_keys = ON") catch {}; try conn.execNoArgs("DROP TABLE permissioned_space_actor_state"); try conn.execNoArgs("ALTER TABLE permissioned_space_actor_state_next RENAME TO permissioned_space_actor_state"); try conn.execNoArgs("PRAGMA foreign_keys = ON"); } try conn.execNoArgs( \\INSERT OR IGNORE INTO permissioned_space_actor_state (space, actor_did, is_authority, deleted_at) \\SELECT s.uri, \\ CASE WHEN s.is_authority = 1 THEN s.authority_did ELSE COALESCE(r.repo_did, s.authority_did) END, \\ s.is_authority, \\ s.deleted_at \\FROM permissioned_spaces s \\LEFT JOIN permissioned_space_repos r ON r.space = s.uri \\WHERE s.is_authority = 1 ); try markMigrationApplied(name);}
fn migratePermissionedRecordOplogRepo() !void { const name = "permissioned-record-oplog-repo"; if (try migrationApplied(name)) return; try conn.execNoArgs("DROP TABLE IF EXISTS permissioned_space_record_oplog_next"); try conn.execNoArgs( \\CREATE TABLE permissioned_space_record_oplog_next ( \\ space TEXT NOT NULL, \\ repo_did TEXT NOT NULL, \\ rev TEXT NOT NULL, \\ idx INTEGER NOT NULL, \\ action TEXT NOT NULL, \\ collection TEXT NOT NULL, \\ rkey TEXT NOT NULL, \\ cid TEXT, \\ prev TEXT, \\ PRIMARY KEY (space, repo_did, rev, idx) \\) ); try conn.execNoArgs( \\INSERT INTO permissioned_space_record_oplog_next ( \\ space, repo_did, rev, idx, action, collection, rkey, cid, prev \\) \\SELECT space, '', rev, idx, action, collection, rkey, cid, prev \\FROM permissioned_space_record_oplog ); try conn.execNoArgs("DROP TABLE permissioned_space_record_oplog"); try conn.execNoArgs("ALTER TABLE permissioned_space_record_oplog_next RENAME TO permissioned_space_record_oplog"); try markMigrationApplied(name);}
fn migrateSeqEventsSync11() !void { const name = "sync11-seq-events"; if (try migrationApplied(name)) return;
var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); const allocator = arena.allocator();
db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try conn.exclusiveTransaction(); errdefer conn.rollback(); try rebuildSeqEventsSync11Locked(allocator); try markMigrationApplied(name); try conn.commit();}
fn migrateSeqEventsSinceRev() !void { const name = "sync-since-rev-events"; if (try migrationApplied(name)) return;
var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); const allocator = arena.allocator();
db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try conn.exclusiveTransaction(); errdefer conn.rollback(); try rebuildSeqEventsSync11Locked(allocator); try markMigrationApplied(name); try conn.commit();}
fn migrateRecordsTable() !void { if (!try recordsTableHasValueJsonColumn()) return; try conn.execNoArgs("DROP TABLE IF EXISTS records_next"); try conn.execNoArgs( \\CREATE TABLE IF NOT EXISTS records_next ( \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ collection TEXT NOT NULL, \\ rkey TEXT NOT NULL, \\ uri TEXT NOT NULL UNIQUE, \\ cid TEXT NOT NULL, \\ rev TEXT NOT NULL, \\ seq INTEGER NOT NULL, \\ PRIMARY KEY (did, collection, rkey) \\) ); try conn.execNoArgs( \\INSERT INTO records_next (did, collection, rkey, uri, cid, rev, seq) \\SELECT did, collection, rkey, uri, cid, rev, seq \\FROM records ); try conn.execNoArgs("PRAGMA foreign_keys = OFF"); errdefer conn.execNoArgs("PRAGMA foreign_keys = ON") catch {}; try conn.execNoArgs("DROP TABLE records"); try conn.execNoArgs("ALTER TABLE records_next RENAME TO records"); try conn.execNoArgs("PRAGMA foreign_keys = ON"); try conn.execNoArgs("CREATE INDEX IF NOT EXISTS records_collection_idx ON records (did, collection, seq DESC)"); try conn.execNoArgs("CREATE INDEX IF NOT EXISTS records_cid_idx ON records (cid)");}
fn migrateRepoBlockRevs() !void { try conn.execNoArgs( \\UPDATE repo_blocks \\SET repo_rev = ( \\ SELECT r.rev \\ FROM records r \\ WHERE r.did = repo_blocks.did AND r.cid = repo_blocks.cid \\ LIMIT 1 \\) \\WHERE repo_rev IS NULL \\ AND EXISTS ( \\ SELECT 1 FROM records r \\ WHERE r.did = repo_blocks.did AND r.cid = repo_blocks.cid \\ ) ); try conn.execNoArgs( \\UPDATE repo_blocks \\SET repo_rev = ( \\ SELECT c.rev \\ FROM commits c \\ WHERE c.did = repo_blocks.did AND c.cid = repo_blocks.cid \\ LIMIT 1 \\) \\WHERE repo_rev IS NULL \\ AND EXISTS ( \\ SELECT 1 FROM commits c \\ WHERE c.did = repo_blocks.did AND c.cid = repo_blocks.cid \\ ) ); try conn.execNoArgs("CREATE INDEX IF NOT EXISTS repo_blocks_rev_idx ON repo_blocks (did, repo_rev DESC, cid DESC)");}
fn recordsTableHasValueJsonColumn() !bool { var rows = try conn.rows("PRAGMA table_info(records)", .{}); defer rows.deinit(); while (rows.next()) |row| { if (std.mem.eql(u8, row.text(1), "value_json")) return true; } if (rows.err) |err| return err; return false;}
fn validateRecordBlocksPresent() !void { const row = try conn.row( \\SELECT COUNT(*) \\FROM records r \\LEFT JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid \\WHERE rb.cid IS NULL , .{}); if (row == null) return; defer row.?.deinit(); if (row.?.int(0) != 0) return Error.MissingRecordBlock;}
fn migrateAppPreferencesTable() !void { if (!try appPreferencesTableHasPreferencesJsonColumn()) return; try conn.execNoArgs( \\CREATE TABLE IF NOT EXISTS app_preferences_next ( \\ id INTEGER PRIMARY KEY AUTOINCREMENT, \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ name TEXT NOT NULL, \\ value_json BLOB NOT NULL, \\ updated_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) ); var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); const allocator = arena.allocator();
{ var rows = try conn.rows( \\SELECT did, namespace, preferences_json \\FROM app_preferences , .{}); defer rows.deinit();
var preferences: std.ArrayList(AppPreference) = .empty; while (rows.next()) |row| { const did = try allocator.dupe(u8, row.text(0)); const namespace = try allocator.dupe(u8, row.text(1)); const preferences_json = try allocator.dupe(u8, row.text(2)); const parsed = std.json.parseFromSlice(std.json.Value, allocator, preferences_json, .{}) catch continue; if (parsed.value != .array) continue; for (parsed.value.array.items) |preference| { if (preference != .object) continue; const name = jsonValueString(preference, "$type") orelse continue; if (!preferenceNameInNamespace(name, namespace)) continue; const value_json = try stringifyValue(allocator, preference); try preferences.append(allocator, .{ .name = name, .value_json = value_json, }); } for (preferences.items) |preference| { try conn.exec( \\INSERT INTO app_preferences_next (did, name, value_json, updated_at) \\VALUES (?, ?, ?, unixepoch()) , .{ did, preference.name, preference.value_json }); } preferences.clearRetainingCapacity(); } if (rows.err) |err| return err; }
try conn.execNoArgs("DROP TABLE app_preferences"); try conn.execNoArgs("ALTER TABLE app_preferences_next RENAME TO app_preferences");}
fn appPreferencesTableHasPreferencesJsonColumn() !bool { var rows = try conn.rows("PRAGMA table_info(app_preferences)", .{}); defer rows.deinit(); while (rows.next()) |row| { if (std.mem.eql(u8, row.text(1), "preferences_json")) return true; } if (rows.err) |err| return err; return false;}
fn preferenceNameInNamespace(name: []const u8, namespace: []const u8) bool { return std.mem.eql(u8, name, namespace) or (std.mem.startsWith(u8, name, namespace) and name.len > namespace.len and name[namespace.len] == '.');}
fn jsonValueString(value: std.json.Value, key: []const u8) ?[]const u8 { return switch (value) { .object => |object| switch (object.get(key) orelse return null) { .string => |string| string, else => null, }, else => null, };}
fn migrateBlobTable() !void { if (!try blobTableHasDataColumn()) return; try conn.execNoArgs( \\CREATE TABLE IF NOT EXISTS blobs_next ( \\ cid TEXT NOT NULL, \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ mime_type TEXT NOT NULL, \\ size INTEGER NOT NULL, \\ storage TEXT NOT NULL DEFAULT 'disk', \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ PRIMARY KEY (did, cid) \\) ); try conn.execNoArgs( \\INSERT OR IGNORE INTO blobs_next (cid, did, mime_type, size, storage, created_at) \\SELECT cid, did, mime_type, size, 'disk', created_at \\FROM blobs ); try conn.execNoArgs("DROP TABLE blobs"); try conn.execNoArgs("ALTER TABLE blobs_next RENAME TO blobs");}
fn blobTableHasDataColumn() !bool { var rows = try conn.rows("PRAGMA table_info(blobs)", .{}); defer rows.deinit(); while (rows.next()) |row| { if (std.mem.eql(u8, row.text(1), "data")) return true; } if (rows.err) |err| return err; return false;}
pub fn getEmailInfo(allocator: std.mem.Allocator, did: []const u8) ?EmailInfo { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return null; const row = conn.row( \\SELECT email, email_confirmed_at, auth_code, auth_code_expires_at, pending_email \\FROM accounts \\WHERE did = ? , .{did}) catch return null; if (row == null) return null; defer row.?.deinit(); return .{ .email = allocator.dupe(u8, row.?.text(0)) catch return null, .email_confirmed = row.?.nullableInt(1) != null, .auth_code = if (row.?.nullableText(2)) |text| allocator.dupe(u8, text) catch null else null, .auth_code_expires_at = row.?.nullableInt(3), .pending_email = if (row.?.nullableText(4)) |text| allocator.dupe(u8, text) catch null else null, };}
pub fn setAuthCode(did: []const u8, code: []const u8, expires_at_ms: i64) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE accounts \\SET auth_code = ?, auth_code_expires_at = ? \\WHERE did = ? , .{ code, expires_at_ms, did });}
pub fn setPendingEmail(did: []const u8, pending_email: []const u8, code: []const u8, expires_at_ms: i64) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE accounts \\SET pending_email = ?, auth_code = ?, auth_code_expires_at = ? \\WHERE did = ? , .{ pending_email, code, expires_at_ms, did });}
pub fn validateAuthCode(did: []const u8, code: []const u8, now_ms: i64) CodeStatus { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return .invalid; const row = conn.row( \\SELECT auth_code, auth_code_expires_at \\FROM accounts \\WHERE did = ? , .{did}) catch return .invalid; if (row == null) return .invalid; defer row.?.deinit(); const stored = row.?.nullableText(0) orelse return .invalid; const expires_at = row.?.nullableInt(1) orelse return .invalid; if (now_ms >= expires_at) return .expired; if (!std.ascii.eqlIgnoreCase(stored, code)) return .invalid; return .valid;}
pub fn clearAuthCode(did: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE accounts \\SET auth_code = NULL, \\ auth_code_expires_at = NULL \\WHERE did = ? , .{did});}
pub fn confirmEmail(did: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE accounts \\SET email_confirmed_at = unixepoch(), \\ auth_code = NULL, \\ auth_code_expires_at = NULL \\WHERE did = ? , .{did});}
pub fn updateEmail(did: []const u8, email: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( \\UPDATE accounts \\SET email = ?, \\ email_confirmed_at = NULL, \\ pending_email = NULL, \\ auth_code = NULL, \\ auth_code_expires_at = NULL \\WHERE did = ? , .{ email, did });}
fn loadNextSeqLocked() !u64 { const row = try conn.row( \\SELECT COALESCE(MAX(seq), 0) + 1 \\FROM ( \\ SELECT seq FROM commits \\ UNION ALL \\ SELECT seq FROM seq_events \\) , .{}); if (row == null) return 1; defer row.?.deinit(); return @intCast(row.?.int(0));}
fn nextSeqLocked() !u64 { _ = try requireInitialized(); return next_seq.fetchAdd(1, .monotonic);}
fn nextRkey(allocator: std.mem.Allocator) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); return nextRkeyLocked(allocator);}
fn nextRkeyLocked(allocator: std.mem.Allocator) ![]const u8 { const seq = try nextSeqLocked(); const tid = try atid.encode(nowMicros(), @intCast(seq % 1024)); return allocator.dupe(u8, &tid);}
fn resolveWriteOpsLocked(allocator: std.mem.Allocator, ops: []const WriteOp) ![]WriteOp { var resolved = try allocator.alloc(WriteOp, ops.len); for (ops, 0..) |op, i| { resolved[i] = switch (op) { .create => |create_op| .{ .create = .{ .collection = create_op.collection, .rkey = create_op.rkey orelse try nextRkeyLocked(allocator), .value = create_op.value, } }, .update => |update_op| .{ .update = update_op }, .delete => |delete_op| .{ .delete = delete_op }, }; } return resolved;}
fn requireWriteSwapsLocked(allocator: std.mem.Allocator, did: []const u8, current: ?CurrentCommit, options: WriteOptions, ops: []const WriteOp) !void { if (options.swap_commit) |expected| { if (current == null or !std.mem.eql(u8, current.?.commit_cid_text, expected)) return Error.InvalidSwap; }
for (ops) |op| switch (op) { .create => {}, .update => |update_op| try requireRecordSwapLocked(allocator, did, update_op.collection, update_op.rkey, update_op.swap), .delete => |delete_op| try requireRecordSwapLocked(allocator, did, delete_op.collection, delete_op.rkey, delete_op.swap), };}
fn requireRecordSwapLocked(allocator: std.mem.Allocator, did: []const u8, collection: []const u8, rkey: []const u8, swap: RecordSwap) !void { switch (swap) { .none => {}, .missing => { if (try currentRecordCidLocked(allocator, did, collection, rkey) != null) return Error.InvalidSwap; }, .cid => |expected| { const actual = try currentRecordCidLocked(allocator, did, collection, rkey) orelse return Error.InvalidSwap; if (!std.mem.eql(u8, actual, expected)) return Error.InvalidSwap; }, }}
fn previousRecordCidsLocked(allocator: std.mem.Allocator, did: []const u8, ops: []const WriteOp) ![]?[]const u8 { var cids = try allocator.alloc(?[]const u8, ops.len); for (ops, 0..) |op, idx| { cids[idx] = switch (op) { .create => null, .update => |update_op| try currentRecordCidLocked(allocator, did, update_op.collection, update_op.rkey), .delete => |delete_op| try currentRecordCidLocked(allocator, did, delete_op.collection, delete_op.rkey), }; } return cids;}
fn currentRecordCidLocked(allocator: std.mem.Allocator, did: []const u8, collection: []const u8, rkey: []const u8) !?[]const u8 { const row = try conn.row( \\SELECT cid \\FROM records \\WHERE did = ? AND collection = ? AND rkey = ? \\LIMIT 1 , .{ did, collection, rkey }); if (row == null) return null; defer row.?.deinit(); return try allocator.dupe(u8, row.?.text(0));}
fn validateRecordForWrite(collection: []const u8, rkey: ?[]const u8, value: std.json.Value, mode: ValidationMode) Error!void { const object = switch (value) { .object => |object| object, else => return Error.InvalidRecordType, }; const record_type = switch (object.get("$type") orelse return Error.InvalidRecordType) { .string => |record_type| record_type, else => return Error.InvalidRecordType, }; if (!std.mem.eql(u8, record_type, collection)) return Error.InvalidRecordType; if (mode == .require and !knownRecordType(record_type)) return Error.ValidationRequired;
if (rkey) |key| try validateKnownRecordKey(record_type, key);}
fn validateKnownRecordKey(record_type: []const u8, rkey: []const u8) Error!void { if (knownRecordKeyRule(record_type)) |rule| switch (rule) { .tid => { _ = atid.decode(rkey) catch return Error.InvalidRecordKey; }, .literal_self => { if (!std.mem.eql(u8, rkey, "self")) return Error.InvalidRecordKey; }, };}
const KnownRecordKeyRule = enum { tid, literal_self };
fn validationStatusForRecord(record_type: []const u8, mode: ValidationMode) []const u8 { return if (mode != .skip and knownRecordType(record_type)) "valid" else "unknown";}
fn knownRecordType(record_type: []const u8) bool { if (knownRecordKeyRule(record_type) != null) return true;
const any_records = [_][]const u8{ "app.bsky.feed.generator", "com.atproto.lexicon.schema", }; for (any_records) |known| { if (std.mem.eql(u8, record_type, known)) return true; }
return false;}
fn knownRecordKeyRule(record_type: []const u8) ?KnownRecordKeyRule { const tid_records = [_][]const u8{ "app.bsky.feed.post", "app.bsky.feed.like", "app.bsky.feed.repost", "app.bsky.feed.threadgate", "app.bsky.feed.postgate", "app.bsky.graph.follow", "app.bsky.graph.block", "app.bsky.graph.list", "app.bsky.graph.listblock", "app.bsky.graph.listitem", "app.bsky.graph.starterpack", "app.bsky.graph.verification", }; for (tid_records) |known| { if (std.mem.eql(u8, record_type, known)) return .tid; }
const self_records = [_][]const u8{ "app.bsky.actor.profile", "app.bsky.actor.status", "app.bsky.labeler.service", "app.bsky.notification.declaration", "chat.bsky.actor.declaration", "com.germnetwork.declaration", }; for (self_records) |known| { if (std.mem.eql(u8, record_type, known)) return .literal_self; }
return null;}
fn repoPath(allocator: std.mem.Allocator, collection: []const u8, rkey: []const u8) ![]const u8 { if (zat.Nsid.parse(collection) == null) return Error.InvalidCollection; if (zat.Rkey.parse(rkey) == null) return Error.InvalidRecordKey; return std.fmt.allocPrint(allocator, "{s}/{s}", .{ collection, rkey });}
fn stageRecordWrite( allocator: std.mem.Allocator, tree: *zat.mst.Mst, account: auth.Account, collection: []const u8, rkey: []const u8, value: std.json.Value, validation_mode: ValidationMode, rev: []const u8, seq: u64, blocks: *std.ArrayList(ImportedBlock), blob_refs: *std.ArrayList(BlobRef),) !Record { const cbor_value = try jsonToDagCbor(allocator, value); const record_bytes = try zat.cbor.encodeAlloc(allocator, cbor_value); const record_cid = try zat.cbor.Cid.forDagCbor(allocator, record_bytes); const record_cid_text = try cidText(allocator, record_cid.raw); const path = try repoPath(allocator, collection, rkey); _ = try tree.putReturn(path, record_cid);
try blocks.append(allocator, .{ .cid = record_cid_text, .data = record_bytes, });
const uri = try std.fmt.allocPrint(allocator, "at://{s}/{s}/{s}", .{ account.did, collection, rkey }); var blobs: std.ArrayList([]const u8) = .empty; try collectBlobCids(allocator, cbor_value, &blobs); for (blobs.items) |blob_cid| { try blob_refs.append(allocator, .{ .cid = blob_cid, .uri = uri }); }
return .{ .did = try allocator.dupe(u8, account.did), .collection = try allocator.dupe(u8, collection), .rkey = try allocator.dupe(u8, rkey), .cid = record_cid_text, .value_json = try stringifyValue(allocator, value), .validation_status = validationStatusForRecord(collection, validation_mode), .rev = try allocator.dupe(u8, rev), .seq = seq, };}
const CurrentCommit = struct { commit_cid_text: []const u8, commit_cid_raw: []const u8, data_cid_raw: []const u8, rev: []const u8, commit_data: []const u8,};
fn latestCommitRawLocked(allocator: std.mem.Allocator, did: []const u8) !?CurrentCommit { const row = try conn.row( \\SELECT c.cid, c.rev, rb.data \\FROM commits c \\JOIN repo_blocks rb ON rb.did = c.did AND rb.cid = c.cid \\WHERE c.did = ? \\ORDER BY c.seq DESC \\LIMIT 1 , .{did}); if (row == null) return null; defer row.?.deinit();
const commit_cid_text = try allocator.dupe(u8, row.?.text(0)); const commit_cid_raw = try cidRawFromText(allocator, commit_cid_text); const rev = try allocator.dupe(u8, row.?.text(1)); const commit_data = row.?.nullableBlob(2) orelse return null; const commit_value = try zat.cbor.decodeAll(allocator, commit_data); const data = commit_value.getCid("data") orelse return error.InvalidRepoPath; return .{ .commit_cid_text = commit_cid_text, .commit_cid_raw = commit_cid_raw, .data_cid_raw = try allocator.dupe(u8, data.raw), .rev = rev, .commit_data = try allocator.dupe(u8, commit_data), };}
fn readRepoCarLocked(allocator: std.mem.Allocator, did: []const u8) !zat.car.Car { var rows = try conn.rows( \\SELECT cid, data \\FROM repo_blocks \\WHERE did = ? \\ORDER BY cid ASC , .{did}); defer rows.deinit();
var blocks: std.ArrayList(zat.car.Block) = .empty; while (rows.next()) |row| { const data = row.nullableBlob(1) orelse ""; try blocks.append(allocator, .{ .cid_raw = try cidRawFromText(allocator, row.text(0)), .data = try allocator.dupe(u8, data), }); } if (rows.err) |err| return err; return .{ .roots = &.{}, .blocks = try blocks.toOwnedSlice(allocator) };}
fn writeMstBlocks(tree: *zat.mst.Mst, out: *std.ArrayList(zat.car.Block)) !void { try tree.collectBlocks(out);}
const SignedCommit = struct { cid: []const u8, data: []const u8,};
fn signedCommit( allocator: std.mem.Allocator, did: []const u8, rev: []const u8, data_cid: zat.cbor.Cid, prev_cid_raw: ?[]const u8,) !SignedCommit { var keypair = try signingKeypairLocked(did); return signedCommitWithKeypair(allocator, did, rev, data_cid, prev_cid_raw, &keypair);}
fn signedCommitWithKeypair( allocator: std.mem.Allocator, did: []const u8, rev: []const u8, data_cid: zat.cbor.Cid, prev_cid_raw: ?[]const u8, keypair: *const zat.Keypair,) !SignedCommit { const signed = try zat.signCommit(allocator, .{ .did = did, .rev = rev, .data = data_cid, .prev = if (prev_cid_raw) |raw| .{ .raw = raw } else null, }, keypair); defer allocator.free(signed.cid.raw); return .{ .cid = try cidText(allocator, signed.cid.raw), .data = signed.bytes, };}
fn signingKeypairLocked(did: []const u8) !zat.Keypair { const row = try conn.row( \\SELECT signing_key_type, signing_key \\FROM accounts \\WHERE did = ? , .{did}); if (row == null) return Error.MissingRecord; defer row.?.deinit();
const existing = row.?.nullableBlob(1); const key_type_text = row.?.nullableText(0) orelse "secp256k1"; const key_bytes = if (existing) |bytes| blk: { if (bytes.len != 32) return Error.InvalidDagCbor; break :blk bytes[0..32].*; } else blk: { const generated = try generateSigningKey(); try conn.exec( \\UPDATE accounts \\SET signing_key_type = 'secp256k1', signing_key = ? \\WHERE did = ? , .{ zqlite.blob(&generated), did }); break :blk generated; }; const key_type: zat.multicodec.KeyType = if (std.mem.eql(u8, key_type_text, "p256")) .p256 else .secp256k1; return zat.Keypair.fromSecretKey(key_type, key_bytes);}
fn consumeReservedSigningKeyLocked(signing_key_or_did: []const u8, did: []const u8) ![32]u8 { const row = try conn.row( \\SELECT signing_key, secret_key, expires_at \\FROM reserved_signing_keys \\WHERE used_at IS NULL \\ AND expires_at > ? \\ AND (signing_key = ? OR did = ?) \\ORDER BY created_at DESC \\LIMIT 1 , .{ nowMs(), signing_key_or_did, did }); if (row == null) return Error.MissingReservedSigningKey; defer row.?.deinit();
const signing_key = row.?.text(0); const secret = row.?.blob(1); if (secret.len != 32) return Error.InvalidReservedSigningKey; const key_bytes = secret[0..32].*; try conn.exec( \\UPDATE reserved_signing_keys \\SET used_at = ? \\WHERE signing_key = ? , .{ nowMs(), signing_key }); return key_bytes;}
fn generateSigningKey() ![32]u8 { var key: [32]u8 = undefined; while (true) { store_io.random(&key); if (zat.Keypair.fromSecretKey(.secp256k1, key)) |_| return key else |_| {} }}
fn jsonToDagCbor(allocator: std.mem.Allocator, value: std.json.Value) !zat.cbor.Value { return switch (value) { .null => .null, .bool => |boolean| .{ .boolean = boolean }, .integer => |integer| if (integer >= 0) .{ .unsigned = @intCast(integer) } else .{ .negative = integer }, .float, .number_string => Error.InvalidDagCbor, .string => |string| .{ .text = string }, .array => |array| blk: { const items = try allocator.alloc(zat.cbor.Value, array.items.len); for (array.items, 0..) |item, idx| items[idx] = try jsonToDagCbor(allocator, item); break :blk .{ .array = items }; }, .object => |object| blk: { if (object.count() == 1) { if (object.get("$link")) |link_value| switch (link_value) { .string => |link| break :blk .{ .cid = .{ .raw = try cidRawFromText(allocator, link) } }, else => {}, }; if (object.get("$bytes")) |bytes_value| switch (bytes_value) { .string => |encoded| { const bytes = try allocator.alloc(u8, std.base64.standard.Decoder.calcSizeForSlice(encoded) catch return Error.InvalidDagCbor); try std.base64.standard.Decoder.decode(bytes, encoded); break :blk .{ .bytes = bytes }; }, else => {}, }; } const entries = try allocator.alloc(zat.cbor.Value.MapEntry, object.count()); var it = object.iterator(); var idx: usize = 0; while (it.next()) |entry| : (idx += 1) { entries[idx] = .{ .key = entry.key_ptr.*, .value = try jsonToDagCbor(allocator, entry.value_ptr.*), }; } break :blk .{ .map = entries }; }, };}
fn collectBlobCids(allocator: std.mem.Allocator, value: zat.cbor.Value, out: *std.ArrayList([]const u8)) !void { switch (value) { .map => |entries| { if (cborMapText(entries, "$type")) |type_name| { if (std.mem.eql(u8, type_name, "blob")) { if (cborMapValue(entries, "ref")) |ref| switch (ref) { .cid => |cid| try out.append(allocator, try cidText(allocator, cid.raw)), .map => |ref_entries| if (cborMapText(ref_entries, "$link")) |link| try out.append(allocator, try allocator.dupe(u8, link)), else => {}, }; } } for (entries) |entry| try collectBlobCids(allocator, entry.value, out); }, .array => |items| for (items) |item| try collectBlobCids(allocator, item, out), else => {}, }}
fn cborMapValue(entries: []const zat.cbor.Value.MapEntry, key: []const u8) ?zat.cbor.Value { for (entries) |entry| { if (std.mem.eql(u8, entry.key, key)) return entry.value; } return null;}
fn cborMapText(entries: []const zat.cbor.Value.MapEntry, key: []const u8) ?[]const u8 { return switch (cborMapValue(entries, key) orelse return null) { .text => |text| text, else => null, };}
fn cidRawFromText(allocator: std.mem.Allocator, cid: []const u8) ![]const u8 { if (cid.len == 0 or cid[0] != 'b') return Error.InvalidDagCbor; return zat.multibase.base32lower.decode(allocator, cid[1..]);}
fn cidText(allocator: std.mem.Allocator, raw: []const u8) ![]const u8 { return zat.multibase.base32lower.encode(allocator, raw);}
fn recordFromRow(row: zqlite.Row, allocator: std.mem.Allocator) !Record { const record_bytes = row.nullableBlob(4) orelse return Error.MissingRecordBlock; return .{ .did = try allocator.dupe(u8, row.text(0)), .collection = try allocator.dupe(u8, row.text(1)), .rkey = try allocator.dupe(u8, row.text(2)), .cid = try allocator.dupe(u8, row.text(3)), .value_json = try recordJsonFromBlock(allocator, record_bytes), .validation_status = validationStatusForRecord(row.text(1), .known), .rev = try allocator.dupe(u8, row.text(5)), .seq = @intCast(row.int(6)), };}
fn recordJsonFromBlock(allocator: std.mem.Allocator, data: []const u8) ![]const u8 { const value = zat.cbor.decodeAll(allocator, data) catch return Error.InvalidDagCbor; return cbor_json.writeAlloc(allocator, value);}
pub fn recordJsonFromDagCbor(allocator: std.mem.Allocator, data: []const u8) ![]const u8 { return recordJsonFromBlock(allocator, data);}
fn oauthRequestFromRow(row: zqlite.Row, allocator: std.mem.Allocator) !OAuthRequest { return .{ .request_id = try allocator.dupe(u8, row.text(0)), .client_id = try allocator.dupe(u8, row.text(1)), .redirect_uri = try allocator.dupe(u8, row.text(2)), .scope = try allocator.dupe(u8, row.text(3)), .state = try allocator.dupe(u8, row.text(4)), .code_challenge = try allocator.dupe(u8, row.text(5)), .code_challenge_method = try allocator.dupe(u8, row.text(6)), .response_mode = try allocator.dupe(u8, row.text(7)), .login_hint = if (row.nullableText(8)) |text| try allocator.dupe(u8, text) else null, .dpop_jkt = if (row.nullableText(9)) |text| try allocator.dupe(u8, text) else null, .expires_at = row.int(10), .sub = if (row.nullableText(11)) |text| try allocator.dupe(u8, text) else null, .code = if (row.nullableText(12)) |text| try allocator.dupe(u8, text) else null, .auth_method = if (row.nullableText(13)) |text| try allocator.dupe(u8, text) else null, };}
fn appPasswordFromRow(row: zqlite.Row, allocator: std.mem.Allocator) !AppPassword { return .{ .name = try allocator.dupe(u8, row.text(0)), .password_hash = try allocator.dupe(u8, row.text(1)), .created_at = row.int(2), .privileged = row.int(3) != 0, .scopes = if (row.nullableText(4)) |text| try allocator.dupe(u8, text) else null, .created_by_controller_did = if (row.nullableText(5)) |text| try allocator.dupe(u8, text) else null, };}
fn sessionsFromRows(allocator: std.mem.Allocator, rows: anytype) ![]SessionInfo { var out: std.ArrayList(SessionInfo) = .empty; while (rows.next()) |row| { try out.append(allocator, .{ .id = try allocator.dupe(u8, row.text(0)), .did = try allocator.dupe(u8, row.text(1)), .handle = try allocator.dupe(u8, row.text(2)), .auth_method = try allocator.dupe(u8, row.text(3)), .app_password_name = if (row.nullableText(4)) |text| try allocator.dupe(u8, text) else null, .controller_did = if (row.nullableText(5)) |text| try allocator.dupe(u8, text) else null, .created_at = row.int(6), .updated_at = row.int(7), .last_used_at = row.nullableInt(8), .access_expires_at = row.int(9), .refresh_expires_at = row.int(10), .revoked_at = row.nullableInt(11), .active = row.int(12) != 0, }); } return out.toOwnedSlice(allocator);}
fn oauthGrantsFromRows(allocator: std.mem.Allocator, rows: anytype) ![]OAuthGrantInfo { var out: std.ArrayList(OAuthGrantInfo) = .empty; while (rows.next()) |row| { try out.append(allocator, .{ .did = try allocator.dupe(u8, row.text(0)), .handle = try allocator.dupe(u8, row.text(1)), .client_id = try allocator.dupe(u8, row.text(2)), .scope = try allocator.dupe(u8, row.text(3)), .created_at = row.int(4), .expires_at = row.int(5), .revoked_at = row.nullableInt(6), .auth_method = if (row.nullableText(7)) |text| try allocator.dupe(u8, text) else null, .active = row.int(8) != 0, }); } return out.toOwnedSlice(allocator);}
fn stringifyValue(allocator: std.mem.Allocator, value: std.json.Value) ![]const u8 { var out: std.Io.Writer.Allocating = .init(allocator); defer out.deinit(); try out.writer.print("{f}", .{std.json.fmt(value, .{})}); return out.toOwnedSlice();}
const Root = struct { cid: []const u8, rev: []const u8,};
fn latestRootLocked(allocator: std.mem.Allocator, did: []const u8) !Root { const row = try conn.row( \\SELECT cid, rev \\FROM commits \\WHERE did = ? \\ORDER BY seq DESC \\LIMIT 1 , .{did}); if (row == null) return Error.RepoNotFound; defer row.?.deinit(); return .{ .cid = try allocator.dupe(u8, row.?.text(0)), .rev = try allocator.dupe(u8, row.?.text(1)), };}
fn backfillSeqEventsLocked(allocator: std.mem.Allocator) !void { var rows = try conn.rows( \\SELECT c.seq, c.did, c.cid, c.rev, c.prev, p.rev \\FROM commits c \\LEFT JOIN commits p ON p.did = c.did AND p.cid = c.prev \\LEFT JOIN seq_events e ON e.seq = c.seq \\WHERE e.seq IS NULL \\ORDER BY c.seq ASC , .{}); defer rows.deinit();
var missing: std.ArrayList(struct { seq: u64, did: []const u8, cid: []const u8, rev: []const u8, prev: ?[]const u8, since_rev: ?[]const u8, }) = .empty; while (rows.next()) |row| { try missing.append(allocator, .{ .seq = @intCast(row.int(0)), .did = try allocator.dupe(u8, row.text(1)), .cid = try allocator.dupe(u8, row.text(2)), .rev = try allocator.dupe(u8, row.text(3)), .prev = if (row.nullableText(4)) |prev| try allocator.dupe(u8, prev) else null, .since_rev = if (row.nullableText(5)) |rev| try allocator.dupe(u8, rev) else null, }); } if (rows.err) |err| return err;
for (missing.items) |commit| { { var commit_arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer commit_arena.deinit(); const commit_allocator = commit_arena.allocator(); const repo_car = writeRepoCarFromLocked(commit_allocator, commit.did, commit.cid) catch continue; const frame = commitEventFrameFromCar(commit_allocator, commit.seq, commit.did, commit.cid, commit.rev, commit.since_rev, null, repo_car, &.{}) catch continue; try conn.exec( \\INSERT INTO seq_events (seq, did, commit_cid, evt) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(seq) DO NOTHING , .{ @as(i64, @intCast(commit.seq)), commit.did, commit.cid, zqlite.blob(frame) }); } }}
fn rebuildSeqEventsSync11Locked(allocator: std.mem.Allocator) !void { var rows = try conn.rows( \\SELECT c.seq, c.did, c.cid, c.rev, c.prev, p.rev \\FROM commits c \\LEFT JOIN commits p ON p.did = c.did AND p.cid = c.prev \\ORDER BY c.seq ASC , .{}); defer rows.deinit();
var commits: std.ArrayList(struct { seq: u64, did: []const u8, cid: []const u8, rev: []const u8, prev: ?[]const u8, since_rev: ?[]const u8, }) = .empty; while (rows.next()) |row| { try commits.append(allocator, .{ .seq = @intCast(row.int(0)), .did = try allocator.dupe(u8, row.text(1)), .cid = try allocator.dupe(u8, row.text(2)), .rev = try allocator.dupe(u8, row.text(3)), .prev = if (row.nullableText(4)) |prev| try allocator.dupe(u8, prev) else null, .since_rev = if (row.nullableText(5)) |rev| try allocator.dupe(u8, rev) else null, }); } if (rows.err) |err| return err;
for (commits.items) |commit| { { var commit_arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer commit_arena.deinit(); const commit_allocator = commit_arena.allocator(); const repo_car = writeRepoCarFromLocked(commit_allocator, commit.did, commit.cid) catch continue; const prev_data_raw = if (commit.prev) |prev_commit| commitDataCidRawLocked(commit_allocator, commit.did, prev_commit) catch null else null; const frame = commitEventFrameFromCar(commit_allocator, commit.seq, commit.did, commit.cid, commit.rev, commit.since_rev, prev_data_raw, repo_car, &.{}) catch continue; try conn.exec( \\INSERT INTO seq_events (seq, did, commit_cid, evt) \\VALUES (?, ?, ?, ?) \\ON CONFLICT(seq) DO UPDATE SET \\ did = excluded.did, \\ commit_cid = excluded.commit_cid, \\ evt = excluded.evt , .{ @as(i64, @intCast(commit.seq)), commit.did, commit.cid, zqlite.blob(frame) }); } }}
fn commitDataCidRawLocked(allocator: std.mem.Allocator, did: []const u8, commit_cid: []const u8) ![]const u8 { const row = try conn.row( \\SELECT data \\FROM repo_blocks \\WHERE did = ? AND cid = ? \\LIMIT 1 , .{ did, commit_cid }); if (row == null) return Error.MissingRecordBlock; defer row.?.deinit(); const commit_data = row.?.nullableBlob(0) orelse return Error.MissingRecordBlock; const commit_value = try zat.cbor.decodeAll(allocator, commit_data); const data = commit_value.getCid("data") orelse return Error.InvalidRepoPath; return try allocator.dupe(u8, data.raw);}
fn writeRepoCarFromLocked(allocator: std.mem.Allocator, did: []const u8, root_cid: []const u8) ![]const u8 { const root_raw = try zat.multibase.base32lower.decode(allocator, root_cid[1..]); const car_root = zat.cbor.Cid{ .raw = root_raw };
var rows = try conn.rows( \\SELECT cid, data \\FROM repo_blocks \\WHERE did = ? \\ORDER BY cid ASC , .{did}); defer rows.deinit();
var blocks: std.ArrayList(zat.car.Block) = .empty; while (rows.next()) |row| { const cid_text = row.text(0); const data = row.nullableBlob(1) orelse ""; const cid_raw = try zat.multibase.base32lower.decode(allocator, cid_text[1..]); try blocks.append(allocator, .{ .cid_raw = cid_raw, .data = try allocator.dupe(u8, data), }); } if (rows.err) |err| return err; return zat.car.writeAlloc(allocator, .{ .roots = &.{car_root}, .blocks = blocks.items, });}
fn commitEventFrame( allocator: std.mem.Allocator, seq: u64, did: []const u8, commit_cid: []const u8, rev: []const u8, since_rev: ?[]const u8, prev_data_raw: ?[]const u8, commit_data: []const u8, record_blocks: []const ImportedBlock, mst_blocks: []const zat.car.Block, ops: []const WriteOp, records: []const Record, prev_record_cids: []const ?[]const u8,) ![]const u8 { const root_raw = try cidRawFromText(allocator, commit_cid); var blocks: std.ArrayList(zat.car.Block) = .empty; try blocks.append(allocator, .{ .cid_raw = root_raw, .data = commit_data }); for (record_blocks) |block| { try blocks.append(allocator, .{ .cid_raw = try cidRawFromText(allocator, block.cid), .data = block.data, }); } for (mst_blocks) |block| try blocks.append(allocator, block); const car_bytes = try zat.car.writeAlloc(allocator, .{ .roots = &.{.{ .raw = root_raw }}, .blocks = blocks.items, });
var op_values = try allocator.alloc(zat.firehose.CommitEventOp, ops.len); var record_idx: usize = 0; for (ops, 0..) |op, idx| switch (op) { .create => |create_op| { const record = records[record_idx]; record_idx += 1; op_values[idx] = try firehoseOp(allocator, .create, create_op.collection, record.rkey, record.cid, null); }, .update => |update_op| { const record = records[record_idx]; record_idx += 1; op_values[idx] = try firehoseOp(allocator, .update, update_op.collection, update_op.rkey, record.cid, prev_record_cids[idx]); }, .delete => |delete_op| { op_values[idx] = try firehoseOp(allocator, .delete, delete_op.collection, delete_op.rkey, null, prev_record_cids[idx]); }, }; return commitEventFrameFromCar(allocator, seq, did, commit_cid, rev, since_rev, prev_data_raw, car_bytes, op_values);}
fn importedCommitEventFrame( allocator: std.mem.Allocator, seq: u64, did: []const u8, commit_cid: []const u8, rev: []const u8, blocks: []const ImportedBlock,) ![]const u8 { var car_blocks: std.ArrayList(zat.car.Block) = .empty; for (blocks) |block| { try car_blocks.append(allocator, .{ .cid_raw = try cidRawFromText(allocator, block.cid), .data = block.data, }); } const root_raw = try cidRawFromText(allocator, commit_cid); const car_bytes = try zat.car.writeAlloc(allocator, .{ .roots = &.{.{ .raw = root_raw }}, .blocks = car_blocks.items, }); return commitEventFrameFromCar(allocator, seq, did, commit_cid, rev, null, null, car_bytes, &.{});}
fn accountEventFrame(allocator: std.mem.Allocator, seq: u64, did: []const u8, status: AccountStatus) ![]const u8 { const header = try eventHeader(allocator, "#account"); const body = if (status.isActive()) try zat.cbor.encodeAlloc(allocator, .{ .map = &.{ .{ .key = "seq", .value = .{ .unsigned = seq } }, .{ .key = "did", .value = .{ .text = did } }, .{ .key = "active", .value = .{ .boolean = true } }, .{ .key = "time", .value = .{ .text = try nowIso(allocator) } }, } }) else try zat.cbor.encodeAlloc(allocator, .{ .map = &.{ .{ .key = "seq", .value = .{ .unsigned = seq } }, .{ .key = "did", .value = .{ .text = did } }, .{ .key = "active", .value = .{ .boolean = false } }, .{ .key = "status", .value = .{ .text = status.asString() } }, .{ .key = "time", .value = .{ .text = try nowIso(allocator) } }, } }); return joinFrame(allocator, header, body);}
fn identityEventFrame(allocator: std.mem.Allocator, seq: u64, did: []const u8, handle: []const u8) ![]const u8 { const header = try eventHeader(allocator, "#identity"); const body = try zat.cbor.encodeAlloc(allocator, .{ .map = &.{ .{ .key = "seq", .value = .{ .unsigned = seq } }, .{ .key = "did", .value = .{ .text = did } }, .{ .key = "handle", .value = .{ .text = handle } }, .{ .key = "time", .value = .{ .text = try nowIso(allocator) } }, } }); return joinFrame(allocator, header, body);}
fn syncEventFrame( allocator: std.mem.Allocator, seq: u64, did: []const u8, commit_cid: []const u8, data_cid_raw: []const u8, rev: []const u8, commit_data: []const u8,) ![]const u8 { const root_raw = try cidRawFromText(allocator, commit_cid); const car_bytes = try zat.car.writeAlloc(allocator, .{ .roots = &.{.{ .raw = data_cid_raw }}, .blocks = &.{.{ .cid_raw = root_raw, .data = commit_data }}, }); const header = try eventHeader(allocator, "#sync"); const body = try zat.cbor.encodeAlloc(allocator, .{ .map = &.{ .{ .key = "seq", .value = .{ .unsigned = seq } }, .{ .key = "did", .value = .{ .text = did } }, .{ .key = "blocks", .value = .{ .bytes = car_bytes } }, .{ .key = "rev", .value = .{ .text = rev } }, .{ .key = "time", .value = .{ .text = try nowIso(allocator) } }, } }); return joinFrame(allocator, header, body);}
fn commitEventFrameFromCar( allocator: std.mem.Allocator, seq: u64, did: []const u8, commit_cid: []const u8, rev: []const u8, since_rev: ?[]const u8, prev_data_raw: ?[]const u8, car_bytes: []const u8, ops: []const zat.firehose.CommitEventOp,) ![]const u8 { const commit_raw = try cidRawFromText(allocator, commit_cid); return zat.firehose.encodeCommitEvent(allocator, .{ .seq = @intCast(seq), .repo_did = did, .commit_cid = .{ .raw = commit_raw }, .rev = rev, .since_rev = since_rev, .prev_data = if (prev_data_raw) |raw| .{ .raw = raw } else null, .blocks = car_bytes, .ops = ops, .time = try nowIso(allocator), });}
fn eventHeader(allocator: std.mem.Allocator, tag: []const u8) ![]const u8 { return zat.cbor.encodeAlloc(allocator, .{ .map = &.{ .{ .key = "op", .value = .{ .unsigned = 1 } }, .{ .key = "t", .value = .{ .text = tag } }, } });}
fn joinFrame(allocator: std.mem.Allocator, header: []const u8, body: []const u8) ![]const u8 { var frame = try allocator.alloc(u8, header.len + body.len); @memcpy(frame[0..header.len], header); @memcpy(frame[header.len..], body); return frame;}
fn firehoseOp( allocator: std.mem.Allocator, action: zat.CommitAction, collection: []const u8, rkey: []const u8, cid: ?[]const u8, prev: ?[]const u8,) !zat.firehose.CommitEventOp { return .{ .action = action, .collection = collection, .rkey = rkey, .cid = if (cid) |c| .{ .raw = try cidRawFromText(allocator, c) } else null, .prev = if (prev) |p| .{ .raw = try cidRawFromText(allocator, p) } else null, };}
fn nowIso(allocator: std.mem.Allocator) ![]const u8 { const seconds_i64 = clock.now(); const seconds: u64 = if (seconds_i64 < 0) 0 else @intCast(seconds_i64); const epoch_seconds = std.time.epoch.EpochSeconds{ .secs = seconds }; const year_day = epoch_seconds.getEpochDay().calculateYearDay(); const month_day = year_day.calculateMonthDay(); const day_seconds = epoch_seconds.getDaySeconds(); return std.fmt.allocPrint( allocator, "{d:0>4}-{d:0>2}-{d:0>2}T{d:0>2}:{d:0>2}:{d:0>2}.000Z", .{ year_day.year, @intFromEnum(month_day.month), month_day.day_index + 1, day_seconds.getHoursIntoDay(), day_seconds.getMinutesIntoHour(), day_seconds.getSecondsIntoMinute(), }, );}
fn scalarCountLocked(sql: [:0]const u8, did: []const u8) !u64 { const row = try conn.row(sql, .{did}); if (row == null) return 0; defer row.?.deinit(); return @intCast(row.?.int(0));}
fn accountActiveLocked(did: []const u8) !bool { return (try accountStatusLocked(did)).isActive();}
fn accountStatusLocked(did: []const u8) !AccountStatus { const row = try conn.row( \\SELECT account_status \\FROM accounts \\WHERE did = ? , .{did}); if (row == null) return Error.RepoNotFound; defer row.?.deinit(); return AccountStatus.parse(row.?.text(0));}
fn revokeAccountTokensLocked(did: []const u8) !void { try conn.exec( \\UPDATE session_tokens \\SET revoked_at = unixepoch(), updated_at = unixepoch() \\WHERE did = ? AND revoked_at IS NULL , .{did}); try conn.exec( \\UPDATE oauth_tokens \\SET revoked_at = unixepoch() \\WHERE did = ? AND revoked_at IS NULL , .{did});}
fn cidForJson(allocator: std.mem.Allocator, json: []const u8) ![]const u8 { const cid = try zat.cbor.Cid.forDagCbor(allocator, json); defer allocator.free(cid.raw); return zat.multibase.base32lower.encode(allocator, cid.raw);}
fn cidForBlob(allocator: std.mem.Allocator, data: []const u8) ![]const u8 { const cid = try zat.cbor.Cid.create(allocator, 1, 0x55, 0x12, data); defer allocator.free(cid.raw); return zat.multibase.base32lower.encode(allocator, cid.raw);}
fn requireInitialized() !void { if (!initialized) return Error.StoreNotInitialized;}
fn insertInviteCodeLocked(code: []const u8, use_count: i64, for_account: []const u8, created_by: []const u8) !void { try conn.exec( \\INSERT INTO invite_codes (code, available_uses, disabled, for_account, created_by, created_at) \\VALUES (?, ?, 0, ?, ?, unixepoch()) , .{ code, use_count, for_account, created_by });}
fn ensureInviteAvailableLocked(code: []const u8) !void { const row = try conn.row( \\SELECT available_uses, disabled, \\ (SELECT COUNT(*) FROM invite_code_uses u WHERE u.code = invite_codes.code) AS uses \\FROM invite_codes \\WHERE code = ? \\LIMIT 1 , .{code}); if (row == null) return error.InvalidInviteCode; defer row.?.deinit(); if (row.?.int(1) != 0) return error.InvalidInviteCode; if (row.?.int(2) >= row.?.int(0)) return error.InvalidInviteCode;}
fn recordInviteUseLocked(code: []const u8, used_by: []const u8) !void { try conn.exec( \\INSERT INTO invite_code_uses (code, used_by, used_at) \\VALUES (?, ?, unixepoch()) , .{ code, used_by });}
fn inviteCodeUsesLocked(allocator: std.mem.Allocator, code: []const u8) ![]InviteCodeUse { var rows = try conn.rows( \\SELECT used_by, used_at \\FROM invite_code_uses \\WHERE code = ? \\ORDER BY used_at DESC, used_by , .{code}); defer rows.deinit();
var uses: std.ArrayList(InviteCodeUse) = .empty; while (rows.next()) |row| { try uses.append(allocator, .{ .used_by = try allocator.dupe(u8, row.text(0)), .used_at = row.int(1), }); } if (rows.err) |err| return err; return uses.toOwnedSlice(allocator);}
fn generateInviteCode(allocator: std.mem.Allocator, public_url: []const u8) ![]const u8 { const host = publicHost(public_url); const normalized = try normalizedInviteHost(allocator, host); var first: [5]u8 = undefined; var second: [5]u8 = undefined; fillBase32(&first); fillBase32(&second); return std.fmt.allocPrint(allocator, "{s}-{s}-{s}", .{ normalized, &first, &second });}
fn publicHost(public_url: []const u8) []const u8 { const after_scheme = if (std.mem.indexOf(u8, public_url, "://")) |idx| public_url[idx + 3 ..] else public_url; const end = std.mem.indexOfAny(u8, after_scheme, "/?#") orelse after_scheme.len; const host_port = after_scheme[0..end]; if (std.mem.lastIndexOfScalar(u8, host_port, ':')) |idx| return host_port[0..idx]; return host_port;}
fn normalizedInviteHost(allocator: std.mem.Allocator, host: []const u8) ![]const u8 { const out = try allocator.alloc(u8, host.len); for (host, 0..) |c, i| { out[i] = switch (c) { '.' => '-', else => std.ascii.toLower(c), }; } return out;}
fn fillBase32(out: *[5]u8) void { const alphabet = "abcdefghijklmnopqrstuvwxyz234567"; var bytes: [5]u8 = undefined; randomBytes(&bytes); for (&bytes, 0..) |b, i| out[i] = alphabet[b % alphabet.len];}
const schema_statements = [_][*:0]const u8{ \\CREATE TABLE IF NOT EXISTS accounts ( \\ did TEXT PRIMARY KEY, \\ handle TEXT NOT NULL UNIQUE, \\ email TEXT NOT NULL UNIQUE, \\ email_confirmed_at INTEGER, \\ auth_code TEXT, \\ auth_code_expires_at INTEGER, \\ pending_email TEXT, \\ invites_disabled INTEGER NOT NULL DEFAULT 0, \\ activated_at INTEGER, \\ account_status TEXT NOT NULL DEFAULT 'active', \\ account_status_ref TEXT, \\ signing_key_type TEXT, \\ signing_key BLOB, \\ password_hash TEXT NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , "CREATE UNIQUE INDEX IF NOT EXISTS accounts_handle_nocase_idx ON accounts (lower(handle))", \\CREATE TABLE IF NOT EXISTS rate_limit_events ( \\ id INTEGER PRIMARY KEY AUTOINCREMENT, \\ subject TEXT NOT NULL, \\ action TEXT NOT NULL, \\ occurred_at INTEGER NOT NULL \\) , "CREATE INDEX IF NOT EXISTS rate_limit_events_lookup_idx ON rate_limit_events (subject, action, occurred_at)", \\CREATE TABLE IF NOT EXISTS records ( \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ collection TEXT NOT NULL, \\ rkey TEXT NOT NULL, \\ uri TEXT NOT NULL UNIQUE, \\ cid TEXT NOT NULL, \\ rev TEXT NOT NULL, \\ seq INTEGER NOT NULL, \\ PRIMARY KEY (did, collection, rkey) \\) , "CREATE INDEX IF NOT EXISTS records_collection_idx ON records (did, collection, seq DESC)", "CREATE INDEX IF NOT EXISTS records_cid_idx ON records (cid)", \\CREATE TABLE IF NOT EXISTS commits ( \\ seq INTEGER PRIMARY KEY, \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ cid TEXT NOT NULL, \\ rev TEXT NOT NULL, \\ prev TEXT \\) , "CREATE INDEX IF NOT EXISTS commits_did_idx ON commits (did, seq DESC)", \\CREATE TABLE IF NOT EXISTS seq_events ( \\ seq INTEGER PRIMARY KEY, \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ commit_cid TEXT NOT NULL, \\ evt BLOB NOT NULL \\) , \\CREATE TABLE IF NOT EXISTS blocks ( \\ cid TEXT PRIMARY KEY, \\ data BLOB NOT NULL \\) , \\CREATE TABLE IF NOT EXISTS blobs ( \\ cid TEXT NOT NULL, \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ mime_type TEXT NOT NULL, \\ size INTEGER NOT NULL, \\ storage TEXT NOT NULL DEFAULT 'disk', \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ PRIMARY KEY (did, cid) \\) , \\CREATE TABLE IF NOT EXISTS record_blobs ( \\ blob_cid TEXT NOT NULL, \\ record_uri TEXT NOT NULL REFERENCES records(uri) ON DELETE CASCADE, \\ PRIMARY KEY (blob_cid, record_uri) \\) , \\CREATE TABLE IF NOT EXISTS app_preferences ( \\ id INTEGER PRIMARY KEY AUTOINCREMENT, \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ name TEXT NOT NULL, \\ value_json BLOB NOT NULL, \\ updated_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , \\CREATE TABLE IF NOT EXISTS oauth_requests ( \\ request_id TEXT PRIMARY KEY, \\ client_id TEXT NOT NULL, \\ redirect_uri TEXT NOT NULL, \\ scope TEXT NOT NULL, \\ state TEXT NOT NULL, \\ code_challenge TEXT NOT NULL, \\ code_challenge_method TEXT NOT NULL, \\ response_mode TEXT NOT NULL DEFAULT 'query', \\ login_hint TEXT, \\ dpop_jkt TEXT, \\ expires_at INTEGER NOT NULL, \\ sub TEXT, \\ code TEXT UNIQUE, \\ auth_method TEXT, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , \\CREATE TABLE IF NOT EXISTS oauth_tokens ( \\ family_id TEXT NOT NULL UNIQUE, \\ access_token TEXT PRIMARY KEY, \\ refresh_token TEXT NOT NULL UNIQUE, \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ client_id TEXT NOT NULL, \\ scope TEXT NOT NULL, \\ expires_at INTEGER NOT NULL, \\ access_expires_at INTEGER NOT NULL, \\ refresh_expires_at INTEGER NOT NULL, \\ previous_refresh_token TEXT, \\ previous_refresh_expires_at INTEGER, \\ dpop_jkt TEXT, \\ auth_method TEXT, \\ revoked_at INTEGER, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , \\CREATE TABLE IF NOT EXISTS dpop_jtis ( \\ jti TEXT PRIMARY KEY, \\ expires_at INTEGER NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , \\CREATE TABLE IF NOT EXISTS session_tokens ( \\ id TEXT PRIMARY KEY, \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ access_jti TEXT NOT NULL UNIQUE, \\ refresh_jti TEXT NOT NULL UNIQUE, \\ access_expires_at INTEGER NOT NULL, \\ refresh_expires_at INTEGER NOT NULL, \\ auth_method TEXT NOT NULL, \\ controller_did TEXT, \\ app_password_name TEXT, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ updated_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ last_used_at INTEGER, \\ revoked_at INTEGER \\) , "CREATE INDEX IF NOT EXISTS session_tokens_did_idx ON session_tokens (did, created_at DESC)", "CREATE INDEX IF NOT EXISTS session_tokens_access_jti_idx ON session_tokens (access_jti)", "CREATE INDEX IF NOT EXISTS session_tokens_refresh_jti_idx ON session_tokens (refresh_jti)", "CREATE INDEX IF NOT EXISTS session_tokens_controller_idx ON session_tokens (controller_did) WHERE controller_did IS NOT NULL", "CREATE INDEX IF NOT EXISTS session_tokens_app_password_idx ON session_tokens (did, app_password_name) WHERE app_password_name IS NOT NULL", \\CREATE TABLE IF NOT EXISTS app_passwords ( \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ name TEXT NOT NULL, \\ password_hash TEXT NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ privileged INTEGER NOT NULL DEFAULT 0, \\ scopes TEXT, \\ created_by_controller_did TEXT, \\ PRIMARY KEY (did, name) \\) , "CREATE INDEX IF NOT EXISTS app_passwords_did_created_idx ON app_passwords (did, created_at DESC)", \\CREATE TABLE IF NOT EXISTS account_audit_log ( \\ id TEXT PRIMARY KEY, \\ subject_did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ actor_did TEXT NOT NULL, \\ controller_did TEXT, \\ action TEXT NOT NULL, \\ details_json TEXT NOT NULL DEFAULT '{}', \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , "CREATE INDEX IF NOT EXISTS account_audit_subject_idx ON account_audit_log (subject_did, created_at DESC)", "CREATE INDEX IF NOT EXISTS account_audit_controller_idx ON account_audit_log (controller_did, created_at DESC) WHERE controller_did IS NOT NULL", \\CREATE TABLE IF NOT EXISTS passkeys ( \\ id TEXT PRIMARY KEY, \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ credential_id BLOB NOT NULL UNIQUE, \\ public_key BLOB NOT NULL, \\ sign_count INTEGER NOT NULL DEFAULT 0, \\ friendly_name TEXT, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ last_used INTEGER \\) , "CREATE INDEX IF NOT EXISTS passkeys_did_idx ON passkeys (did)", \\CREATE TABLE IF NOT EXISTS webauthn_challenges ( \\ id TEXT PRIMARY KEY, \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ challenge TEXT NOT NULL, \\ kind TEXT NOT NULL, \\ state_json TEXT NOT NULL, \\ expires_at INTEGER NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , "CREATE INDEX IF NOT EXISTS webauthn_challenges_did_kind_idx ON webauthn_challenges (did, kind)", \\CREATE TABLE IF NOT EXISTS webauthn_discoverable_challenges ( \\ request_id TEXT PRIMARY KEY REFERENCES oauth_requests(request_id) ON DELETE CASCADE, \\ challenge TEXT NOT NULL, \\ expires_at INTEGER NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , \\CREATE TABLE IF NOT EXISTS invite_codes ( \\ code TEXT PRIMARY KEY, \\ available_uses INTEGER NOT NULL, \\ disabled INTEGER NOT NULL DEFAULT 0, \\ for_account TEXT NOT NULL, \\ created_by TEXT NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , "CREATE INDEX IF NOT EXISTS invite_codes_for_account_idx ON invite_codes (for_account)", \\CREATE TABLE IF NOT EXISTS invite_code_uses ( \\ code TEXT NOT NULL REFERENCES invite_codes(code) ON DELETE CASCADE, \\ used_by TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ used_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ PRIMARY KEY (code, used_by) \\) , \\CREATE TABLE IF NOT EXISTS reserved_signing_keys ( \\ id INTEGER PRIMARY KEY AUTOINCREMENT, \\ did TEXT, \\ signing_key TEXT NOT NULL UNIQUE, \\ signing_key_type TEXT NOT NULL, \\ secret_key BLOB NOT NULL, \\ expires_at INTEGER NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch() * 1000), \\ used_at INTEGER \\) , "CREATE INDEX IF NOT EXISTS reserved_signing_keys_did_idx ON reserved_signing_keys (did) WHERE did IS NOT NULL", "CREATE INDEX IF NOT EXISTS reserved_signing_keys_expires_idx ON reserved_signing_keys (expires_at) WHERE used_at IS NULL", \\CREATE TABLE IF NOT EXISTS permissioned_spaces ( \\ uri TEXT PRIMARY KEY, \\ authority_did TEXT NOT NULL, \\ space_type TEXT NOT NULL, \\ skey TEXT NOT NULL, \\ managing_app TEXT, \\ policy TEXT NOT NULL DEFAULT 'member-list', \\ app_access_json TEXT NOT NULL DEFAULT '{"type":"open"}', \\ is_authority INTEGER NOT NULL DEFAULT 1, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ deleted_at INTEGER, \\ UNIQUE(authority_did, space_type, skey) \\) , \\CREATE TABLE IF NOT EXISTS permissioned_space_actor_state ( \\ space TEXT NOT NULL REFERENCES permissioned_spaces(uri) ON DELETE CASCADE, \\ actor_did TEXT NOT NULL, \\ is_authority INTEGER NOT NULL DEFAULT 0, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ deleted_at INTEGER, \\ PRIMARY KEY (space, actor_did) \\) , \\CREATE TABLE IF NOT EXISTS simplespace_members ( \\ space TEXT NOT NULL REFERENCES permissioned_spaces(uri) ON DELETE CASCADE, \\ member_did TEXT NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ PRIMARY KEY (space, member_did) \\) , "CREATE INDEX IF NOT EXISTS simplespace_members_member_idx ON simplespace_members (member_did, space)", "CREATE INDEX IF NOT EXISTS permissioned_space_actor_state_actor_idx ON permissioned_space_actor_state (actor_did, space)", \\CREATE TABLE IF NOT EXISTS permissioned_space_records ( \\ space TEXT NOT NULL REFERENCES permissioned_spaces(uri) ON DELETE CASCADE, \\ repo_did TEXT NOT NULL, \\ collection TEXT NOT NULL, \\ rkey TEXT NOT NULL, \\ cid TEXT NOT NULL, \\ value_json BLOB NOT NULL, \\ validation_status TEXT NOT NULL DEFAULT 'unknown', \\ repo_rev TEXT NOT NULL DEFAULT '', \\ updated_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ PRIMARY KEY (space, repo_did, collection, rkey) \\) , "CREATE INDEX IF NOT EXISTS permissioned_space_records_list_idx ON permissioned_space_records (space, repo_did, collection, rkey)", \\CREATE TABLE IF NOT EXISTS permissioned_space_record_blobs ( \\ space TEXT NOT NULL, \\ repo_did TEXT NOT NULL, \\ collection TEXT NOT NULL, \\ rkey TEXT NOT NULL, \\ blob_cid TEXT NOT NULL, \\ PRIMARY KEY (space, repo_did, collection, rkey, blob_cid), \\ FOREIGN KEY (space, repo_did, collection, rkey) \\ REFERENCES permissioned_space_records(space, repo_did, collection, rkey) \\ ON DELETE CASCADE \\) , "CREATE INDEX IF NOT EXISTS permissioned_space_record_blobs_blob_idx ON permissioned_space_record_blobs (repo_did, blob_cid)", \\CREATE TABLE IF NOT EXISTS permissioned_space_repos ( \\ space TEXT NOT NULL REFERENCES permissioned_spaces(uri) ON DELETE CASCADE, \\ repo_did TEXT NOT NULL, \\ set_hash BLOB, \\ rev TEXT, \\ PRIMARY KEY (space, repo_did) \\) , \\CREATE TABLE IF NOT EXISTS permissioned_space_writers ( \\ space TEXT NOT NULL REFERENCES permissioned_spaces(uri) ON DELETE CASCADE, \\ repo_did TEXT NOT NULL, \\ rev TEXT NOT NULL, \\ hash BLOB NOT NULL, \\ updated_at INTEGER NOT NULL DEFAULT (unixepoch()), \\ PRIMARY KEY (space, repo_did) \\) , \\CREATE TABLE IF NOT EXISTS permissioned_space_record_oplog ( \\ space TEXT NOT NULL, \\ repo_did TEXT NOT NULL, \\ rev TEXT NOT NULL, \\ idx INTEGER NOT NULL, \\ action TEXT NOT NULL, \\ collection TEXT NOT NULL, \\ rkey TEXT NOT NULL, \\ cid TEXT, \\ prev TEXT, \\ PRIMARY KEY (space, repo_did, rev, idx) \\) , \\CREATE TABLE IF NOT EXISTS permissioned_space_notify_registrations ( \\ space TEXT NOT NULL, \\ repo_did TEXT NOT NULL DEFAULT '', \\ service_endpoint TEXT NOT NULL, \\ expires_at INTEGER NOT NULL, \\ PRIMARY KEY (space, repo_did, service_endpoint) \\) ,};
const post_schema_statements = [_][*:0]const u8{ \\CREATE TABLE IF NOT EXISTS expected_blobs ( \\ blob_cid TEXT NOT NULL, \\ record_uri TEXT NOT NULL REFERENCES records(uri) ON DELETE CASCADE, \\ PRIMARY KEY (blob_cid, record_uri) \\) , "CREATE UNIQUE INDEX IF NOT EXISTS oauth_tokens_family_id_idx ON oauth_tokens (family_id)", "CREATE UNIQUE INDEX IF NOT EXISTS oauth_tokens_previous_refresh_token_idx ON oauth_tokens (previous_refresh_token) WHERE previous_refresh_token IS NOT NULL", \\CREATE TABLE IF NOT EXISTS oauth_used_refresh_tokens ( \\ refresh_token TEXT PRIMARY KEY, \\ family_id TEXT NOT NULL REFERENCES oauth_tokens(family_id) ON DELETE CASCADE, \\ expires_at INTEGER NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , \\CREATE TABLE IF NOT EXISTS permissioned_space_used_delegations ( \\ jti TEXT PRIMARY KEY, \\ expires_at INTEGER NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , \\CREATE TABLE IF NOT EXISTS permissioned_space_used_client_attestations ( \\ jti TEXT PRIMARY KEY, \\ expires_at INTEGER NOT NULL, \\ created_at INTEGER NOT NULL DEFAULT (unixepoch()) \\) , \\CREATE TABLE IF NOT EXISTS repo_blocks ( \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, \\ cid TEXT NOT NULL, \\ data BLOB NOT NULL, \\ repo_rev TEXT, \\ PRIMARY KEY (did, cid) \\) , "CREATE INDEX IF NOT EXISTS repo_blocks_rev_idx ON repo_blocks (did, repo_rev DESC, cid DESC)",};
test "persists records in sqlite" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "alice.test", "alice@test.com", "password", "did:plc:cmadossymmii3izkabdbp5en", true, ); const parsed = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"hello\"}", .{}); defer parsed.deinit();
const rkey = "3jzfcijpj2z2a"; const record = try create(allocator, account, "app.bsky.feed.post", rkey, parsed.value); try std.testing.expectEqualStrings(rkey, record.rkey); try std.testing.expect(!try recordsTableHasValueJsonColumn());
const block_row = try conn.row( "SELECT count(*) FROM repo_blocks WHERE did = ? AND cid = ?", .{ account.did, record.cid }, ); try std.testing.expect(block_row != null); defer block_row.?.deinit(); try std.testing.expectEqual(@as(i64, 1), block_row.?.int(0));
const fetched = get(account.did, "app.bsky.feed.post", rkey).?; try std.testing.expectEqualStrings(record.cid, fetched.cid); try std.testing.expect(std.mem.indexOf(u8, fetched.value_json, "hello") != null);
const listed = try writeListJson(allocator, account.did, "app.bsky.feed.post", null, false, 10); try std.testing.expect(std.mem.indexOf(u8, listed, "hello") != null);
const repo_car = try writeRepoCar(allocator, account.did); const loaded = try zat.loadCommitFromCAR(allocator, repo_car); try std.testing.expectEqualStrings(account.did, loaded.commit.did); try std.testing.expectEqualStrings(fetched.rev, loaded.commit.rev); try std.testing.expectEqualStrings((try latestRootLocked(allocator, account.did)).cid, try cidText(allocator, loaded.commit_cid));
const tree = try zat.mst.Mst.loadFromBlocks(allocator, loaded.repo_car, loaded.commit.data_cid); const found = tree.get("app.bsky.feed.post/3jzfcijpj2z2a") orelse return error.MissingRecord; try std.testing.expectEqualStrings(record.cid, try cidText(allocator, found.raw));
try conn.exec( "DELETE FROM repo_blocks WHERE did = ? AND cid = ?", .{ account.did, record.cid }, ); try std.testing.expect(get(account.did, "app.bsky.feed.post", rkey) == null);}
test "active account resolution excludes unavailable identities" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "Alice.Example.COM", "alice-active@test.com", "password", "did:plc:activehandle", true, ); try std.testing.expect((try findActiveAccount(allocator, "alice.example.com")) != null); try setAccountActive(account.did, false); try std.testing.expect((try findAccount(allocator, "alice.example.com")) != null); try std.testing.expect((try findActiveAccount(allocator, "alice.example.com")) == null);}
test "durable rate limits are isolated by subject and window" { try init(std.Options.debug_io, ":memory:"); defer close();
try consumeRateLimit("did:plc:alice", "identity.updateHandle", 1_000, 1_000, 2); try consumeRateLimit("did:plc:alice", "identity.updateHandle", 1_500, 1_000, 2); try std.testing.expectError( error.RateLimitExceeded, consumeRateLimit("did:plc:alice", "identity.updateHandle", 1_750, 1_000, 2), ); try consumeRateLimit("did:plc:bob", "identity.updateHandle", 1_750, 1_000, 2); try consumeRateLimit("did:plc:alice", "identity.updateHandle", 2_100, 1_000, 2);}
test "repo writes enforce swap and explicit validation preconditions" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "swaps.test", "swaps@test.com", "password", "did:plc:swapscheck", true, );
const first = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"before\"}", .{}); defer first.deinit(); const rkey = "3jzfcijpj2z2b"; const created_result = try applyWritesWithOptions(allocator, account, &.{.{ .create = .{ .collection = "app.bsky.feed.post", .rkey = rkey, .value = first.value, } }}, .{}); const first_cid = created_result.records[0].cid;
const second = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"after\"}", .{}); defer second.deinit(); try std.testing.expectError(Error.InvalidSwap, applyWritesWithOptions(allocator, account, &.{.{ .update = .{ .collection = "app.bsky.feed.post", .rkey = rkey, .value = second.value, } }}, .{ .swap_commit = "bafkreiwrongcommit" })); try std.testing.expectError(Error.InvalidSwap, applyWritesWithOptions(allocator, account, &.{.{ .update = .{ .collection = "app.bsky.feed.post", .rkey = rkey, .value = second.value, .swap = .{ .cid = "bafkreiwrongrecord" }, } }}, .{}));
const updated_result = try applyWritesWithOptions(allocator, account, &.{.{ .update = .{ .collection = "app.bsky.feed.post", .rkey = rkey, .value = second.value, .swap = .{ .cid = first_cid }, } }}, .{ .swap_commit = created_result.commit.cid }); try std.testing.expectEqualStrings("after", (try std.json.parseFromSlice(std.json.Value, allocator, updated_result.records[0].value_json, .{})).value.object.get("text").?.string);
const custom = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"example.test.record\",\"text\":\"custom\"}", .{}); defer custom.deinit(); try std.testing.expectError(Error.ValidationRequired, applyWritesWithOptions(allocator, account, &.{.{ .create = .{ .collection = "example.test.record", .rkey = "custom", .value = custom.value, } }}, .{ .validate = .require }));
const custom_result = try applyWritesWithOptions(allocator, account, &.{.{ .create = .{ .collection = "example.test.record", .rkey = "custom", .value = custom.value, } }}, .{ .validate = .skip }); try std.testing.expectEqualStrings("unknown", custom_result.records[0].validation_status);}
test "repo writes lazily load existing MST blocks for mutation" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "lazy.test", "lazy@test.com", "password", "did:plc:lazywritecheck", true, );
var create_ops: [160]WriteOp = undefined; var create_parsed: [160]std.json.Parsed(std.json.Value) = undefined; var create_rkeys: [160][atid.encoded_len]u8 = undefined; for (&create_ops, &create_parsed, 0..) |*op, *parsed, i| { parsed.* = try std.json.parseFromSlice( std.json.Value, allocator, try std.fmt.allocPrint(allocator, "{{\"$type\":\"app.bsky.feed.post\",\"text\":\"record {d}\"}}", .{i}), .{}, ); create_rkeys[i] = try atid.encode(1_700_000_000_000_000 + i, 0); op.* = .{ .create = .{ .collection = "app.bsky.feed.post", .rkey = &create_rkeys[i], .value = parsed.value, } }; } defer for (&create_parsed) |*parsed| parsed.deinit();
_ = try applyWritesWithOptions(allocator, account, &create_ops, .{});
const update = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"updated\"}", .{}); defer update.deinit(); const update_rkey = create_rkeys[77]; const update_result = try applyWritesWithOptions(allocator, account, &.{.{ .update = .{ .collection = "app.bsky.feed.post", .rkey = &update_rkey, .value = update.value, } }}, .{});
const delete_rkey = create_rkeys[12]; _ = try applyWritesWithOptions(allocator, account, &.{.{ .delete = .{ .collection = "app.bsky.feed.post", .rkey = &delete_rkey, } }}, .{});
const repo_car = try writeRepoCar(allocator, account.did); const loaded = try zat.loadCommitFromCAR(allocator, repo_car); var tree = try zat.mst.Mst.loadFromBlocks(allocator, loaded.repo_car, loaded.commit.data_cid);
const updated_path = try std.fmt.allocPrint(allocator, "app.bsky.feed.post/{s}", .{update_rkey}); const updated = tree.get(updated_path) orelse return error.MissingRecord; try std.testing.expectEqualStrings(update_result.records[0].cid, try cidText(allocator, updated.raw)); const deleted_path = try std.fmt.allocPrint(allocator, "app.bsky.feed.post/{s}", .{delete_rkey}); try std.testing.expect(tree.get(deleted_path) == null); const retained_rkey = create_rkeys[159]; const retained_path = try std.fmt.allocPrint(allocator, "app.bsky.feed.post/{s}", .{retained_rkey}); const retained = tree.get(retained_path) orelse return error.MissingRecord; try std.testing.expectEqualStrings((get(account.did, "app.bsky.feed.post", &retained_rkey) orelse return error.MissingRecord).cid, try cidText(allocator, retained.raw));}
test "commit firehose since uses previous repo rev" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "firehose-since.test", "firehose-since@test.com", "password", "did:plc:firehosesince", true, );
const first = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"first\"}", .{}); defer first.deinit(); const first_rkey = "3jzfcijpj2z2c"; const first_result = try applyWritesWithOptions(allocator, account, &.{.{ .create = .{ .collection = "app.bsky.feed.post", .rkey = first_rkey, .value = first.value, } }}, .{});
const second = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"second\"}", .{}); defer second.deinit(); const second_rkey = "3jzfcijpj2z2d"; const second_result = try applyWritesWithOptions(allocator, account, &.{.{ .create = .{ .collection = "app.bsky.feed.post", .rkey = second_rkey, .value = second.value, } }}, .{});
try conn.exec("DELETE FROM seq_events WHERE did = ?", .{account.did}); try rebuildSeqEventsSync11Locked(allocator);
const events = try listSeqEvents(allocator, 0, 20); var found_second = false; for (events) |event| { const decoded = zat.firehose.decodeFrame(allocator, event.frame) catch continue; switch (decoded) { .commit => |commit| { if (!std.mem.eql(u8, commit.repo, account.did)) continue; if (!std.mem.eql(u8, commit.rev, second_result.commit.rev)) continue; found_second = true; try std.testing.expect(commit.since != null); try std.testing.expectEqualStrings(first_result.commit.rev, commit.since.?); try std.testing.expect(!std.mem.eql(u8, first_result.commit.cid, commit.since.?)); }, else => {}, } } try std.testing.expect(found_second);}
test "getRepo since filters repo blocks by revision" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "since.test", "since@test.com", "password", "did:plc:sincecheck", true, );
const first = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"first\"}", .{}); defer first.deinit(); const first_rkey = "3jzfcijpj2z2e"; const first_result = try applyWritesWithOptions(allocator, account, &.{.{ .create = .{ .collection = "app.bsky.feed.post", .rkey = first_rkey, .value = first.value, } }}, .{});
const second = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"second\"}", .{}); defer second.deinit(); const second_rkey = "3jzfcijpj2z2f"; const second_result = try applyWritesWithOptions(allocator, account, &.{.{ .create = .{ .collection = "app.bsky.feed.post", .rkey = second_rkey, .value = second.value, } }}, .{});
const full_car = try zat.car.read(allocator, try writeRepoCar(allocator, account.did)); const diff_car = try zat.car.read(allocator, try writeRepoCarSince(allocator, account.did, first_result.commit.rev)); const empty_diff_car = try zat.car.read(allocator, try writeRepoCarSince(allocator, account.did, second_result.commit.rev));
try std.testing.expect(diff_car.blocks.len < full_car.blocks.len); try std.testing.expect(diff_car.blocks.len > 0); try std.testing.expectEqual(@as(usize, 0), empty_diff_car.blocks.len); try std.testing.expectEqualStrings(second_result.commit.cid, try cidText(allocator, empty_diff_car.roots[0].raw));}
fn carContainsCid(allocator: std.mem.Allocator, c: zat.car.Car, cid: []const u8) !bool { for (c.blocks) |block| { if (std.mem.eql(u8, cid, try cidText(allocator, block.cid_raw))) return true; } return false;}
test "full getRepo exports only current reachable record blocks" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "reachable.test", "reachable@test.com", "password", "did:plc:reachablecheck", true, );
const first = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"first\"}", .{}); defer first.deinit(); const rkey = "3jzfcijpj2z2g"; const first_result = try applyWritesWithOptions(allocator, account, &.{.{ .create = .{ .collection = "app.bsky.feed.post", .rkey = rkey, .value = first.value, } }}, .{}); const first_cid = first_result.records[0].cid;
const second = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"second\"}", .{}); defer second.deinit(); const second_result = try applyWritesWithOptions(allocator, account, &.{.{ .update = .{ .collection = "app.bsky.feed.post", .rkey = rkey, .value = second.value, } }}, .{}); const second_cid = second_result.records[0].cid;
const updated_car = try zat.car.read(allocator, try writeRepoCar(allocator, account.did)); try std.testing.expect(!try carContainsCid(allocator, updated_car, first_cid)); try std.testing.expect(try carContainsCid(allocator, updated_car, second_cid));
_ = try applyWritesWithOptions(allocator, account, &.{.{ .delete = .{ .collection = "app.bsky.feed.post", .rkey = rkey, } }}, .{});
const deleted_car = try zat.car.read(allocator, try writeRepoCar(allocator, account.did)); try std.testing.expect(!try carContainsCid(allocator, deleted_car, first_cid)); try std.testing.expect(!try carContainsCid(allocator, deleted_car, second_cid));}
test "permissioned spaces store self-owned records outside public repo" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "space.test", "space@test.com", "password", "did:plc:spaceauthorityalice", true, );
const space = try createSpace(allocator, .{ .actor_did = account.did, .authority_did = account.did, .space_type = "fm.plyr.privateMedia", .skey = "self", .is_authority = true, .managing_app = "did:web:plyr.fm", .policy = "member-list", .app_access_json = "{\"type\":\"open\"}", }); try std.testing.expectEqualStrings("at://did:plc:spaceauthorityalice/space/fm.plyr.privateMedia/self", space.uri); try std.testing.expect(space.is_authority);
try std.testing.expectError(Error.InvalidRefreshSession, createSpace(allocator, .{ .actor_did = account.did, .authority_did = account.did, .space_type = "fm.plyr.privateMedia", .skey = "self", .is_authority = true, .managing_app = null, .policy = "member-list", .app_access_json = "{\"type\":\"open\"}", }));
const parsed = try std.json.parseFromSlice( std.json.Value, allocator, "{\"$type\":\"fm.plyr.track\",\"title\":\"secret track\",\"audioBlob\":{\"$type\":\"blob\",\"ref\":{\"$link\":\"bafkreihdwdcefgh4dqkjv67uzcmw7ojee6xedzdetojuzjevtenxquvyku\"},\"mimeType\":\"audio/mpeg\",\"size\":12}}", .{}, ); defer parsed.deinit();
const prepared = try prepareRecordValue(allocator, "fm.plyr.track", "track-one", parsed.value); const stored = try putSpaceRecord(allocator, space.uri, account.did, "fm.plyr.track", "track-one", prepared); try std.testing.expectEqualStrings("track-one", stored.rkey); try std.testing.expectEqualStrings("unknown", stored.validation_status); try std.testing.expect(get(account.did, "fm.plyr.track", "track-one") == null); try expectSpaceBlobRefCount(space.uri, account.did, "fm.plyr.track", "track-one", "bafkreihdwdcefgh4dqkjv67uzcmw7ojee6xedzdetojuzjevtenxquvyku", 1); try std.testing.expect(getPublicBlob(allocator, account.did, "bafkreihdwdcefgh4dqkjv67uzcmw7ojee6xedzdetojuzjevtenxquvyku") == null);
const found = (try getSpaceRecord(allocator, space.uri, account.did, "fm.plyr.track", "track-one")).?; try std.testing.expectEqualStrings(stored.cid, found.cid); try std.testing.expect(std.mem.indexOf(u8, found.value_json, "secret track") != null);
const listed = try listSpaceRecords(allocator, space.uri, account.did, "fm.plyr.track", null, false, 50, true); try std.testing.expectEqual(@as(usize, 1), listed.len); try std.testing.expectEqualStrings("fm.plyr.track", listed[0].collection); try std.testing.expectEqualStrings("track-one", listed[0].rkey); try std.testing.expect(listed[0].value_json != null); try std.testing.expect(std.mem.indexOf(u8, listed[0].value_json.?, "secret track") != null);
const listed_refs = try listSpaceRecords(allocator, space.uri, account.did, "fm.plyr.track", null, false, 50, false); try std.testing.expectEqual(@as(usize, 1), listed_refs.len); try std.testing.expect(listed_refs[0].value_json == null);
const updated_parsed = try std.json.parseFromSlice( std.json.Value, allocator, "{\"$type\":\"fm.plyr.track\",\"title\":\"secret track\",\"audioBlob\":{\"$type\":\"blob\",\"ref\":{\"$link\":\"bafkreibm6jgipb22ignhvbmf6avdkvuevprpyx6y6722fsxk7wxiofq4wu\"},\"mimeType\":\"audio/mpeg\",\"size\":12}}", .{}, ); defer updated_parsed.deinit(); const updated_prepared = try prepareRecordValue(allocator, "fm.plyr.track", "track-one", updated_parsed.value); _ = try putSpaceRecord(allocator, space.uri, account.did, "fm.plyr.track", "track-one", updated_prepared); try expectSpaceBlobRefCount(space.uri, account.did, "fm.plyr.track", "track-one", "bafkreihdwdcefgh4dqkjv67uzcmw7ojee6xedzdetojuzjevtenxquvyku", 0); try expectSpaceBlobRefCount(space.uri, account.did, "fm.plyr.track", "track-one", "bafkreibm6jgipb22ignhvbmf6avdkvuevprpyx6y6722fsxk7wxiofq4wu", 1);
try deleteSpaceRecord(allocator, space.uri, account.did, "fm.plyr.track", "track-one"); try expectSpaceBlobRefCount(space.uri, account.did, "fm.plyr.track", "track-one", "bafkreibm6jgipb22ignhvbmf6avdkvuevprpyx6y6722fsxk7wxiofq4wu", 0);}
test "permissioned space foundation migration preserves rows and rebuilds hashes" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close(); const did = "did:plc:spacemigrationalice"; const canonical = "at://did:plc:spacemigrationalice/space/fm.example.private/self"; const legacy = "ats://did:plc:spacemigrationalice/fm.example.private/self"; const space = try createSpace(allocator, .{ .actor_did = did, .authority_did = did, .space_type = "fm.example.private", .skey = "self", .is_authority = true, .managing_app = null, .policy = "member-list", .app_access_json = "{\"type\":\"open\"}", }); try std.testing.expectEqualStrings(canonical, space.uri); const parsed = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"fm.example.note\",\"text\":\"private\"}", .{}); defer parsed.deinit(); const prepared = try prepareRecordValue(allocator, "fm.example.note", "one", parsed.value); const record = try putSpaceRecord(allocator, canonical, did, "fm.example.note", "one", prepared); try conn.exec( \\INSERT INTO permissioned_space_record_blobs (space, repo_did, collection, rkey, blob_cid) \\VALUES (?, ?, ?, ?, ?) , .{ canonical, did, "fm.example.note", "one", "bafkreiaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" }); try conn.exec( \\INSERT INTO permissioned_space_notify_registrations (space, repo_did, service_endpoint, expires_at) \\VALUES (?, '', ?, 4102444800) , .{ canonical, "https://sync.example" });
try conn.execNoArgs("PRAGMA foreign_keys = OFF"); inline for (.{ "permissioned_space_actor_state", "simplespace_members", "permissioned_space_records", "permissioned_space_record_blobs", "permissioned_space_repos", "permissioned_space_writers", "permissioned_space_record_oplog", "permissioned_space_notify_registrations", }) |table| { try conn.exec("UPDATE " ++ table ++ " SET space = ? WHERE space = ?", .{ legacy, canonical }); } try conn.exec("UPDATE permissioned_spaces SET uri = ? WHERE uri = ?", .{ legacy, canonical }); try conn.exec("UPDATE permissioned_space_repos SET set_hash = zeroblob(?)", .{permissioned.lthash_state_bytes}); try conn.exec("DELETE FROM zds_migrations WHERE name IN (?, ?, ?)", .{ "permissioned-space-at-uri", "permissioned-space-record-element-v1", "permissioned-space-writer-set-v1" }); try conn.execNoArgs("PRAGMA foreign_keys = ON");
try migratePermissionedDataTables(); inline for (.{ "permissioned_spaces", "permissioned_space_actor_state", "simplespace_members", "permissioned_space_records", "permissioned_space_record_blobs", "permissioned_space_repos", "permissioned_space_writers", "permissioned_space_record_oplog", "permissioned_space_notify_registrations", }) |table| { const column = comptime if (std.mem.eql(u8, table, "permissioned_spaces")) "uri" else "space"; const row = (try conn.row("SELECT count(*) FROM " ++ table ++ " WHERE " ++ column ++ " = ?", .{canonical})).?; defer row.deinit(); try std.testing.expect(row.int(0) > 0); } const state = try getSpaceRepoState(allocator, canonical, did); var expected: permissioned.LtHash = .{}; try addRecordElement(allocator, &expected, record.collection, record.rkey, record.cid); try std.testing.expectEqualSlices(u8, &expected.bytes, state.set_hash.?); const writers = try listSpaceWriters(allocator, canonical, null, 50); try std.testing.expectEqual(@as(usize, 1), writers.len); try std.testing.expectEqualStrings(did, writers[0].repo_did); try std.testing.expectEqualSlices(u8, &expected.digest(), writers[0].hash);}
test "permissioned notification registrations are scoped and expire" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator(); try init(std.Options.debug_io, ":memory:"); defer close();
const space = "at://did:plc:authority/space/fm.example.private/self"; try registerSpaceNotification(space, null, "https://whole.example", 4102444800); try registerSpaceNotification(space, "did:plc:writer", "https://repo.example", 4102444800); try registerSpaceNotification(space, null, "https://expired.example", 1);
const whole = try listNotificationRecipients(allocator, space, null, false); try std.testing.expectEqual(@as(usize, 1), whole.len); try std.testing.expectEqualStrings("https://whole.example", whole[0].service_endpoint); const repo = try listNotificationRecipients(allocator, space, "did:plc:writer", true); try std.testing.expectEqual(@as(usize, 2), repo.len); try std.testing.expectEqualStrings("https://whole.example", repo[0].service_endpoint); try std.testing.expectEqualStrings("https://repo.example", repo[1].service_endpoint);}
test "permissioned credential exchange tokens are one use" { try init(std.Options.debug_io, ":memory:"); defer close();
const first = SpaceReplayToken{ .jti = "delegation-one", .expires_at = 4_000_000_000 }; const attestation = SpaceReplayToken{ .jti = "attestation-one", .expires_at = 4_000_000_000 }; try std.testing.expect(try consumeSpaceCredentialExchange(first, attestation)); try std.testing.expect(!try consumeSpaceCredentialExchange(first, null));
const second = SpaceReplayToken{ .jti = "delegation-two", .expires_at = 4_000_000_000 }; try std.testing.expect(!try consumeSpaceCredentialExchange(second, attestation)); try std.testing.expect(try consumeSpaceCredentialExchange(second, null));}
fn expectSpaceBlobRefCount( space: []const u8, repo_did: []const u8, collection: []const u8, rkey: []const u8, blob_cid: []const u8, expected: i64,) !void { const row = try conn.row( \\SELECT COUNT(*) \\FROM permissioned_space_record_blobs \\WHERE space = ? AND repo_did = ? AND collection = ? AND rkey = ? AND blob_cid = ? , .{ space, repo_did, collection, rkey, blob_cid }); defer row.?.deinit(); try std.testing.expectEqual(expected, row.?.int(0));}
test "permissioned spaces keep authority-local state per space URI" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const owner = try createAccount( allocator, "authority-space.test", "authority-space@test.com", "password", "did:plc:spaceauthoritylocal", true, ); const created = try createSpace(allocator, .{ .actor_did = owner.did, .authority_did = owner.did, .space_type = "fm.plyr.privateMedia", .skey = "self", .is_authority = true, .managing_app = null, .policy = "member-list", .app_access_json = "{\"type\":\"open\"}", });
const owner_view = (try getSpace(allocator, owner.did, created.uri)).?; try std.testing.expect(owner_view.is_authority); try std.testing.expect(owner_view.deleted_at == null);
const owner_spaces = try listSpaces(allocator, owner.did, owner.did, "fm.plyr.privateMedia", null, 50); try std.testing.expectEqual(@as(usize, 1), owner_spaces.len); try std.testing.expect(owner_spaces[0].is_authority);
try markSpaceDeleted(owner.did, created.uri); try std.testing.expect((try getSpace(allocator, owner.did, created.uri)) == null);}
test "permissioned space create is duplicate-checked per actor" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const owner = try createAccount(allocator, "authority-first.test", "authority-first@test.com", "password", "did:plc:ownerfirst", true); const viewer = try createAccount(allocator, "viewer-first.test", "viewer-first@test.com", "password", "did:plc:viewerfirst", true); const member = try createAccount(allocator, "member-only.test", "member-only@test.com", "password", "did:plc:memberonly", true);
const viewer_row = try createSpace(allocator, .{ .actor_did = viewer.did, .authority_did = owner.did, .space_type = "fm.plyr.privateMedia", .skey = "self", .is_authority = false, .managing_app = null, .policy = "member-list", .app_access_json = "{\"type\":\"open\"}", }); try std.testing.expect(!viewer_row.is_authority);
const owner_row = try createSpace(allocator, .{ .actor_did = owner.did, .authority_did = owner.did, .space_type = "fm.plyr.privateMedia", .skey = "self", .is_authority = true, .managing_app = "did:web:plyr.fm", .policy = "public", .app_access_json = "{\"type\":\"allowList\",\"allowed\":[\"did:web:allowed.example\"]}", }); try std.testing.expect(owner_row.is_authority); try std.testing.expectEqualStrings("public", owner_row.policy); try std.testing.expectEqualStrings("did:web:plyr.fm", owner_row.managing_app.?);
try addSimpleSpaceMember(owner_row.uri, member.did); const member_spaces = try listSpaces(allocator, member.did, owner.did, "fm.plyr.privateMedia", null, 50); try std.testing.expectEqual(@as(usize, 0), member_spaces.len);
const viewer_spaces = try listSpaces(allocator, viewer.did, owner.did, "fm.plyr.privateMedia", null, 50); try std.testing.expectEqual(@as(usize, 1), viewer_spaces.len); try std.testing.expect(!viewer_spaces[0].is_authority);
try std.testing.expectError(Error.InvalidRefreshSession, createSpace(allocator, .{ .actor_did = viewer.did, .authority_did = owner.did, .space_type = "fm.plyr.privateMedia", .skey = "self", .is_authority = false, .managing_app = null, .policy = "member-list", .app_access_json = "{\"type\":\"open\"}", }));}
test "inactive repo import does not publish firehose event" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "import-inactive.test", "import-inactive@test.com", "password", "did:plc:importinactive", true, ); const parsed = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"hello\"}", .{}); defer parsed.deinit();
const record = try create(allocator, account, "app.bsky.feed.post", "3jzfcijpj2z2h", parsed.value); const repo_car = try writeRepoCar(allocator, account.did); const loaded = try zat.loadCommitFromCAR(allocator, repo_car);
var imported_blocks: std.ArrayList(ImportedBlock) = .empty; for (loaded.repo_car.blocks) |block| { try imported_blocks.append(allocator, .{ .cid = try cidText(allocator, block.cid_raw), .data = block.data, }); } const imported_records = [_]ImportedRecord{.{ .collection = record.collection, .rkey = record.rkey, .cid = record.cid, .blob_cids = &.{}, }};
try setAccountActive(account.did, false); const before = try scalarCountLocked("SELECT COUNT(*) FROM seq_events WHERE did = ?", account.did); try importRepo( allocator, account, try cidText(allocator, loaded.commit_cid), loaded.commit.rev, &imported_records, imported_blocks.items, ); const after = try scalarCountLocked("SELECT COUNT(*) FROM seq_events WHERE did = ?", account.did); try std.testing.expectEqual(before, after); try std.testing.expectEqualStrings((try latestRootLocked(allocator, account.did)).cid, try cidText(allocator, loaded.commit_cid));}
test "put updates record index while content comes from repo blocks" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "alice.test", "alice@test.com", "password", "did:plc:cmadossymmii3izkabdbp5en", true, );
var first = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"before\"}", .{}); defer first.deinit(); const rkey = "3jzfcijpj2z2i"; _ = try create(allocator, account, "app.bsky.feed.post", rkey, first.value);
var second = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"after\"}", .{}); defer second.deinit(); const updated = try put(allocator, account, "app.bsky.feed.post", rkey, second.value);
const fetched = get(account.did, "app.bsky.feed.post", rkey).?; try std.testing.expectEqualStrings(updated.cid, fetched.cid); try std.testing.expect(std.mem.indexOf(u8, fetched.value_json, "after") != null); try std.testing.expect(std.mem.indexOf(u8, fetched.value_json, "before") == null);
const block_row = try conn.row( "SELECT count(*) FROM repo_blocks WHERE did = ? AND cid = ?", .{ account.did, updated.cid }, ); try std.testing.expect(block_row != null); defer block_row.?.deinit(); try std.testing.expectEqual(@as(i64, 1), block_row.?.int(0));}
test "stores blob metadata in sqlite and bytes in disk blobstore" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); var path_buf: [std.fs.max_path_bytes]u8 = undefined; const path_len = try tmp.dir.realPath(std.testing.io, &path_buf); const blob_root = path_buf[0..path_len]; blobstore.init(std.Options.debug_io, blob_root);
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "alice.test", "alice@test.com", "password", "did:plc:cmadossymmii3izkabdbp5en", true, ); const cid = try putBlob(allocator, std.Options.debug_io, account, "hello blob", "text/plain"); const row = try conn.row("SELECT mime_type, size FROM blobs WHERE did = ? AND cid = ?", .{ account.did, cid }); try std.testing.expect(row != null); defer row.?.deinit(); try std.testing.expectEqualStrings("text/plain", row.?.text(0)); try std.testing.expectEqual(@as(i64, 10), row.?.int(1));
const blob = getBlob(allocator, account.did, cid) orelse return error.MissingRecord; try std.testing.expectEqualStrings("text/plain", blob.mime_type); try std.testing.expectEqualStrings("hello blob", blob.data);}
test "stores larger blob metadata in sqlite and bytes in disk blobstore" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); var path_buf: [std.fs.max_path_bytes]u8 = undefined; const path_len = try tmp.dir.realPath(std.testing.io, &path_buf); const blob_root = path_buf[0..path_len]; blobstore.init(std.Options.debug_io, blob_root);
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "large-blob.test", "large-blob@test.com", "password", "did:plc:largeblob", true, );
const data = try allocator.alloc(u8, 571_880); @memset(data, 0xaa);
const cid = try putBlob(allocator, std.Options.debug_io, account, data, "image/jpeg"); const row = try conn.row("SELECT mime_type, size FROM blobs WHERE did = ? AND cid = ?", .{ account.did, cid }); try std.testing.expect(row != null); defer row.?.deinit(); try std.testing.expectEqualStrings("image/jpeg", row.?.text(0)); try std.testing.expectEqual(@as(i64, 571_880), row.?.int(1));
const blob = getBlob(allocator, account.did, cid) orelse return error.MissingRecord; try std.testing.expectEqualStrings("image/jpeg", blob.mime_type); try std.testing.expectEqualSlices(u8, data, blob.data);}
test "invite codes gate account creation and record uses" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const code = try createInviteCode(allocator, "https://pds.example.com", 1, "admin", "admin"); try std.testing.expect(std.mem.startsWith(u8, code, "pds-example-com-")); try std.testing.expect(try inviteCodeIsAvailable(code));
var signing_key: [32]u8 = .{0} ** 32; signing_key[31] = 1; const account = try createAccountWithSigningKeyAndInvite( allocator, "invited.test", "invited@test.com", "password", "did:plc:invited", true, signing_key, code, ); try std.testing.expectEqualStrings("did:plc:invited", account.did); try std.testing.expect(!try inviteCodeIsAvailable(code));
var other_key: [32]u8 = .{0} ** 32; other_key[31] = 2; try std.testing.expectError(error.InvalidInviteCode, createAccountWithSigningKeyAndInvite( allocator, "second.test", "second@test.com", "password", "did:plc:second", true, other_key, code, ));
const codes = try getAccountInviteCodes(allocator, "admin", true); try std.testing.expectEqual(@as(usize, 1), codes.len); try std.testing.expectEqual(@as(usize, 1), codes[0].uses.len); try std.testing.expectEqualStrings("did:plc:invited", codes[0].uses[0].used_by);}
test "account handle updates are unique and immediately resolvable" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const alice = try createAccount( allocator, "alice.test", "alice@test.com", "password", "did:plc:alice", true, ); _ = try createAccount( allocator, "bob.test", "bob@test.com", "password", "did:plc:bob", true, );
try updateAccountHandle(alice.did, "alice-new.test"); try std.testing.expect((try findAccount(allocator, "alice.test")) == null); const updated = (try findAccount(allocator, "ALICE-NEW.TEST")) orelse return error.AccountNotFound; try std.testing.expectEqualStrings(alice.did, updated.did); try std.testing.expectError(error.HandleNotAvailable, updateAccountHandle(alice.did, "BOB.TEST")); try std.testing.expectError(error.AccountNotFound, updateAccountHandle("did:plc:missing", "missing.test"));}
test "reserved signing keys are one-use account signing keys" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const did = "did:plc:reserved"; const reserved = try reserveSigningKey(allocator, did); try std.testing.expectEqualStrings(did, reserved.did.?); try std.testing.expect(std.mem.startsWith(u8, reserved.signing_key, "did:key:"));
const consumed = try consumeReservedSigningKey(reserved.signing_key, did); try std.testing.expectEqualSlices(u8, &reserved.secret_key, &consumed); try std.testing.expectError(error.MissingReservedSigningKey, consumeReservedSigningKey(reserved.signing_key, did));}
test "stores discoverable passkey challenges by oauth request" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
try putOAuthRequest( "req-discoverable", "https://client.example.com/oauth-client-metadata.json", "https://client.example.com/callback", "atproto", "state", "challenge", "S256", "query", null, null, 4102444800, ); try putDiscoverableWebAuthnChallenge("req-discoverable", "webauthn-challenge", 4102444800);
const stored = (try getDiscoverableWebAuthnChallenge(allocator, "req-discoverable")) orelse return error.MissingRecord; try std.testing.expectEqualStrings("req-discoverable", stored.request_id); try std.testing.expectEqualStrings("webauthn-challenge", stored.challenge); try std.testing.expectEqual(@as(i64, 4102444800), stored.expires_at);
try deleteDiscoverableWebAuthnChallenge("req-discoverable"); try std.testing.expect(try getDiscoverableWebAuthnChallenge(allocator, "req-discoverable") == null);}
test "oauth authorization codes remain readable until atomically consumed" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
try putOAuthRequest( "req-code", "https://client.example.com/oauth-client-metadata.json", "https://client.example.com/callback", "atproto", "state", "challenge", "S256", "query", null, "dpop-thumbprint", 4102444800, ); try authorizeOAuthRequest("req-code", "did:plc:code", "authorization-code", "passkey");
const before = (try getOAuthRequestByCode(allocator, "authorization-code")) orelse return error.MissingRecord; try std.testing.expectEqualStrings("req-code", before.request_id); try std.testing.expectEqualStrings("authorization-code", before.code.?);
try std.testing.expect(!try consumeOAuthCode("wrong-request", "authorization-code")); try std.testing.expect((try getOAuthRequestByCode(allocator, "authorization-code")) != null);
try std.testing.expect(try consumeOAuthCode("req-code", "authorization-code")); try std.testing.expect((try getOAuthRequestByCode(allocator, "authorization-code")) == null); try std.testing.expect(!try consumeOAuthCode("req-code", "authorization-code"));}
test "oauth app sessions stay active after access token expiry until revoked" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "oauth-session.test", "oauth-session@test.com", "password", "did:plc:oauthsession", true, );
try putOAuthToken( account.did, "https://client.example.com/oauth-client-metadata.json", "atproto", "expired-access", "current-refresh", 1, 4102444800, null, "passkey", );
const active = try listOAuthGrantsForAccount(allocator, account.did, true, 50); try std.testing.expectEqual(@as(usize, 1), active.len); try std.testing.expect(active[0].active); try std.testing.expectEqualStrings("passkey", active[0].auth_method.?);
try revokeOAuthToken("current-refresh"); try std.testing.expectEqual(@as(usize, 0), (try listOAuthGrantsForAccount(allocator, account.did, true, 50)).len);
const all = try listOAuthGrantsForAccount(allocator, account.did, false, 50); try std.testing.expectEqual(@as(usize, 1), all.len); try std.testing.expect(!all[0].active);}
test "oauth refresh rotation keeps one previous refresh token in grace" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "oauth-rotation.test", "oauth-rotation@test.com", "password", "did:plc:oauthrotation", true, );
try putOAuthToken( account.did, "https://client.example.com/oauth-client-metadata.json", "atproto", "access-one", "refresh-one", 100, 4102444800, "jkt-one", "passkey", );
const access_one = (try getOAuthToken(allocator, "access-one")) orelse return error.MissingRecord; try std.testing.expectEqual(OAuthToken.Kind.access, access_one.kind); try std.testing.expectEqualStrings("access-one", access_one.family_id); try std.testing.expectEqual(@as(i64, 100), access_one.access_expires_at); try std.testing.expectEqual(@as(i64, 4102444800), access_one.refresh_expires_at);
try rotateOAuthToken("refresh-one", "access-two", "refresh-two", 200, 4102444801, 300);
try std.testing.expect(try getOAuthToken(allocator, "access-one") == null);
const previous = (try getOAuthToken(allocator, "refresh-one")) orelse return error.MissingRecord; try std.testing.expectEqual(OAuthToken.Kind.previous_refresh, previous.kind); try std.testing.expectEqualStrings("access-two", previous.access_token); try std.testing.expectEqualStrings("refresh-two", previous.refresh_token); try std.testing.expectEqual(@as(i64, 300), previous.previous_refresh_expires_at.?);
const current = (try getOAuthToken(allocator, "refresh-two")) orelse return error.MissingRecord; try std.testing.expectEqual(OAuthToken.Kind.refresh, current.kind); try std.testing.expectEqual(@as(i64, 200), current.access_expires_at); try std.testing.expectEqual(@as(i64, 4102444801), current.refresh_expires_at);
try rotateOAuthToken("refresh-two", "access-three", "refresh-three", 400, 4102444802, 500);
const used = (try getOAuthToken(allocator, "refresh-one")) orelse return error.MissingRecord; try std.testing.expectEqual(OAuthToken.Kind.used_refresh, used.kind); try std.testing.expectEqualStrings("access-one", used.family_id); try std.testing.expectEqualStrings("refresh-three", used.refresh_token);
try revokeOAuthToken("refresh-one"); const revoked = (try getOAuthToken(allocator, "refresh-three")) orelse return error.MissingRecord; try std.testing.expect(revoked.revoked);}
test "session tokens are durable and refresh rotation invalidates old jti" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "session.test", "session@test.com", "password", "did:plc:session", true, );
const session_id = try createSessionTokenRow( allocator, account.did, "access-1", "refresh-1", 4102444800, 4102444800, "password", null, null, ); try std.testing.expect(std.mem.startsWith(u8, session_id, "ses-")); try std.testing.expect(try sessionTokenIsActive(account.did, "access-1", "com.atproto.access")); try std.testing.expect(try sessionTokenIsActive(account.did, "refresh-1", "com.atproto.refresh"));
try rotateSessionToken(account.did, "refresh-1", "access-2", "refresh-2", 4102444800, 4102444800); try std.testing.expect(!try sessionTokenIsActive(account.did, "refresh-1", "com.atproto.refresh")); try std.testing.expect(try sessionTokenIsActive(account.did, "access-2", "com.atproto.access")); try std.testing.expect(try sessionTokenIsActive(account.did, "refresh-2", "com.atproto.refresh"));
const active_sessions = try listSessionsForAccount(allocator, account.did, true, 50); try std.testing.expectEqual(@as(usize, 1), active_sessions.len); try std.testing.expectEqualStrings(session_id, active_sessions[0].id); try std.testing.expect(active_sessions[0].active);
try revokeSessionToken(account.did, "refresh-2"); try std.testing.expect(!try sessionTokenIsActive(account.did, "access-2", "com.atproto.access")); try std.testing.expect(!try sessionTokenIsActive(account.did, "refresh-2", "com.atproto.refresh")); try std.testing.expectEqual(@as(usize, 0), (try listSessionsForAccount(allocator, account.did, true, 50)).len);
const inactive_sessions = try listSessionsForAccount(allocator, account.did, false, 50); try std.testing.expectEqual(@as(usize, 1), inactive_sessions.len); try std.testing.expect(!inactive_sessions[0].active);}
test "account takedown sets precise sync status and revokes tokens" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "status.test", "status@test.com", "password", "did:plc:status", true, ); var record_json = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"hello\"}", .{}); defer record_json.deinit(); _ = try create(allocator, account, "app.bsky.feed.post", "3jzfcijpj2z2j", record_json.value);
_ = try createSessionTokenRow( allocator, account.did, "status-access", "status-refresh", 4102444800, 4102444800, "password", null, null, ); try std.testing.expect(try sessionTokenIsActive(account.did, "status-access", "com.atproto.access"));
try setAccountTakendown(account.did, true, "mod-ref-1"); try std.testing.expectEqual(AccountStatus.takendown, try accountStatus(account.did)); try std.testing.expect(!try sessionTokenIsActive(account.did, "status-access", "com.atproto.access"));
const status_json = try writeRepoStatusJson(allocator, account.did); try std.testing.expect(std.mem.indexOf(u8, status_json, "\"active\":false") != null); try std.testing.expect(std.mem.indexOf(u8, status_json, "\"status\":\"takendown\"") != null); try std.testing.expect(std.mem.indexOf(u8, status_json, "\"rev\"") == null);
const frame = try accountEventFrame(allocator, 1, account.did, .takendown); try std.testing.expect(std.mem.indexOf(u8, frame, "takendown") != null);
try setAccountTakendown(account.did, false, null); try std.testing.expectEqual(AccountStatus.active, try accountStatus(account.did));}
test "app passwords are durable matchable and revoke their sessions" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "app-pass.test", "app-pass@test.com", "password", "did:plc:apppass", true, ); const salt: [16]u8 = .{0} ** 16; const password_hash = try auth.hashPassword(allocator, "aaaa-bbbb-cccc-dddd", salt); const created = try createAppPassword( allocator, account.did, "device one", password_hash, false, "transition:generic", null, ); try std.testing.expectEqualStrings("device one", created.name); try std.testing.expectEqual(false, created.privileged);
const listed = try listAppPasswords(allocator, account.did); try std.testing.expectEqual(@as(usize, 1), listed.len); try std.testing.expectEqualStrings("device one", listed[0].name); try std.testing.expectEqualStrings("transition:generic", listed[0].scopes.?);
const matched = (try findMatchingAppPassword(allocator, account.did, "aaaa-bbbb-cccc-dddd")) orelse return error.MissingRecord; try std.testing.expectEqualStrings("device one", matched.name); try std.testing.expect(try findMatchingAppPassword(allocator, account.did, "wrong-password") == null);
_ = try createSessionTokenRow( allocator, account.did, "app-access-1", "app-refresh-1", 4102444800, 4102444800, "app_password", null, "device one", ); const method = (try sessionAuthMethod(allocator, account.did, "app-access-1")) orelse return error.MissingRecord; try std.testing.expectEqualStrings("app_password", method);
try revokeAppPassword(account.did, "device one"); try std.testing.expect(try getAppPasswordByName(allocator, account.did, "device one") == null); try std.testing.expect(try sessionAuthMethod(allocator, account.did, "app-access-1") == null); try std.testing.expect(!try sessionTokenIsActive(account.did, "app-access-1", "com.atproto.access"));}
test "account audit log stores subject actor and controller fields" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator();
try init(std.Options.debug_io, ":memory:"); defer close();
const account = try createAccount( allocator, "audit.test", "audit@test.com", "password", "did:plc:audit", true, ); try recordAuditEvent(account.did, account.did, "did:plc:controller", "repo_write", "{\"collection\":\"app.bsky.feed.post\"}");
const row = try conn.row( \\SELECT subject_did, actor_did, controller_did, action, details_json \\FROM account_audit_log \\WHERE subject_did = ? , .{account.did}); try std.testing.expect(row != null); defer row.?.deinit(); try std.testing.expectEqualStrings(account.did, row.?.text(0)); try std.testing.expectEqualStrings(account.did, row.?.text(1)); try std.testing.expectEqualStrings("did:plc:controller", row.?.text(2)); try std.testing.expectEqualStrings("repo_write", row.?.text(3)); try std.testing.expect(std.mem.indexOf(u8, row.?.text(4), "app.bsky.feed.post") != null);}
test "commit event encoding uses Zat firehose builder" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator(); const cid = "bafyreifmpxapiafzeml4ns5nedutwbsnjw5obdganycdy6mhmb2mxth5aq"; const rev = try atid.encode(1_700_000_123_456_789, 1); const since = try atid.encode(1_700_000_123_456_000, 1); const prev_data_raw = try cidRawFromText(allocator, cid);
const ops = try allocator.alloc(zat.firehose.CommitEventOp, 3); ops[0] = try firehoseOp(allocator, .create, "sh.tangled.string", "zds-create", cid, null); ops[1] = try firehoseOp(allocator, .update, "sh.tangled.string", "zds-update", cid, cid); ops[2] = try firehoseOp(allocator, .delete, "sh.tangled.string", "zds-delete", null, cid);
const frame = try commitEventFrameFromCar( allocator, 1, "did:plc:b64lsctzqnzpv6vd4ry3qktw", cid, &rev, &since, prev_data_raw, "not-a-real-car-for-this-encoding-test", ops, ); const decoded = try zat.firehose.decodeFrame(allocator, frame); try std.testing.expectEqual(@as(i64, 1), decoded.commit.seq); try std.testing.expectEqualStrings("did:plc:b64lsctzqnzpv6vd4ry3qktw", decoded.commit.repo); try std.testing.expectEqualStrings(&rev, decoded.commit.rev); try std.testing.expectEqualStrings(&since, decoded.commit.since.?); try std.testing.expectEqualSlices(u8, prev_data_raw, decoded.commit.prev_data.?.raw); try std.testing.expectEqual(@as(usize, 3), decoded.commit.ops.len); try std.testing.expectEqualStrings("sh.tangled.string/zds-create", decoded.commit.ops[0].path); try std.testing.expectEqualSlices(u8, ops[0].cid.?.raw, decoded.commit.ops[0].cid.?.raw); try std.testing.expectEqualStrings("sh.tangled.string/zds-update", decoded.commit.ops[1].path); try std.testing.expectEqualSlices(u8, ops[1].cid.?.raw, decoded.commit.ops[1].cid.?.raw); try std.testing.expectEqualSlices(u8, ops[1].prev.?.raw, decoded.commit.ops[1].prev.?.raw); try std.testing.expectEqualStrings("sh.tangled.string/zds-delete", decoded.commit.ops[2].path); try std.testing.expect(decoded.commit.ops[2].cid == null); try std.testing.expectEqualSlices(u8, ops[2].prev.?.raw, decoded.commit.ops[2].prev.?.raw);}
test "Zat firehose builder rejects commit CID since rev" { const allocator = std.testing.allocator; const cid = "bafyreifmpxapiafzeml4ns5nedutwbsnjw5obdganycdy6mhmb2mxth5aq"; const commit_raw = try cidRawFromText(allocator, cid); defer allocator.free(commit_raw); const rev = try atid.encode(1_700_000_123_456_789, 1);
try std.testing.expectError(error.InvalidSince, zat.firehose.encodeCommitEvent(allocator, .{ .seq = 1, .repo_did = "did:plc:b64lsctzqnzpv6vd4ry3qktw", .commit_cid = .{ .raw = commit_raw }, .rev = &rev, .since_rev = cid, .prev_data = null, .blocks = "not-a-real-car-for-this-encoding-test", .ops = &.{}, .time = "2026-06-22T00:00:00.000Z", }));}
test "repo rev generation does not move behind current head" { const allocator = std.testing.allocator; const future_tid = try atid.encode(nowMicros() + std.time.us_per_s, 0); const rev = try revForSeq(allocator, 42, &future_tid); defer allocator.free(rev);
const previous = try atid.timestampMicros(&future_tid); const next = try atid.timestampMicros(rev); try std.testing.expect(next > previous); try std.testing.expect(std.mem.lessThan(u8, &future_tid, rev));}