{"version":3,"sources":["/Users/gwakko/Projects/shared-websocket/dist/chunk-JG2LC2FF.cjs","../src/utils/id.ts","../src/MessageBus.ts","../src/constants.ts","../src/TabCoordinator.ts","../src/utils/backoff.ts","../src/SharedSocket.ts","../src/WorkerSocket.ts","../src/SubscriptionManager.ts","../src/FramePipeline.ts","../src/IncomingPipeline.ts","../src/Outbox.ts","../src/SubscriptionRegistry.ts","../src/AuthManager.ts","../src/PushManager.ts","../src/SharedWebSocket.ts"],"names":[],"mappings":"AAAA;ACAO,SAAS,UAAA,CAAA,EAAqB;AACnC,EAAA,GAAA,CAAI,OAAO,OAAA,IAAW,YAAA,GAAe,MAAA,CAAO,UAAA,EAAY;AACtD,IAAA,OAAO,MAAA,CAAO,UAAA,CAAW,CAAA;AAAA,EAC3B;AACA,EAAA,OAAO,CAAA,EAAA;AACT;ADEU;AACA;AEFG;AAKX,EAAA;AAEmB,IAAA;AAEZ,IAAA;AACA,IAAA;AACH,MAAA;AACF,IAAA;AACF,EAAA;AANmB,EAAA;AANX,EAAA;AACA,iBAAA;AACA,kBAAA;AAYR,EAAA;AACQ,IAAA;AAMA,MAAA;AACF,QAAA;AACF,MAAA;AACF,IAAA;AACK,IAAA;AACL,IAAA;AACF,EAAA;AAEW,EAAA;AACJ,IAAA;AACP,EAAA;AAEA,EAAA;AACQ,IAAA;AACD,IAAA;AAEA,IAAA;AACP,EAAA;AAEM,EAAA;AACE,IAAA;AACN,IAAA;AACE,MAAA;AACE,QAAA;AACA,QAAA;AACC,MAAA;AACH,MAAA;AACA,MAAA;AACD,IAAA;AACH,EAAA;AAEc,EAAA;AACN,IAAA;AACA,MAAA;AACJ,MAAA;AACA,MAAA;AACF,IAAA;AACK,IAAA;AACL,IAAA;AACF,EAAA;AAEQ,EAAA;AAEF,IAAA;AACF,MAAA;AACA,MAAA;AACI,MAAA;AACF,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACF,IAAA;AAEM,IAAA;AACF,IAAA;AACF,MAAA;AACF,IAAA;AACF,EAAA;AAEQ,EAAA;AACD,IAAA;AACP,EAAA;AAEQ,EAAA;AACN,IAAA;AACF,EAAA;AAEQ,EAAA;AACF,IAAA;AACC,IAAA;AACH,MAAA;AACA,MAAA;AACF,IAAA;AACI,IAAA;AACN,EAAA;AAEQ,EAAA;AACD,oBAAA;AACP,EAAA;AAEQ,EAAA;AACN,IAAA;AACE,MAAA;AACA,MAAA;AACF,IAAA;AACK,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AACF;AFjBU;AACA;AG1FG;AAIA;AAAQ;AAEnB,EAAA;AAAU;AAEF,EAAA;AAAA;AAER,EAAA;AAAW;AAEX,EAAA;AAAU;AAEJ,EAAA;AAAA;AAEN,EAAA;AACF;AAGa;AAIM;AAAA;AAEjB,EAAA;AAAS;AAET,EAAA;AAAU;AAEV,EAAA;AAAkB;AAElB,EAAA;AAAgB;AAEhB,EAAA;AAAa;AAEP,EAAA;AAAA;AAEN,EAAA;AAAS;AAET,EAAA;AAAW;AAEX,EAAA;AAAa;AAEb,EAAA;AACF;AAGa;AAEA;AAIA;AACX,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACQ,EAAA;AACD,EAAA;AACD,EAAA;AACE,EAAA;AACV;AAGa;AAKA;AACX,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AAAqB;AAErB,EAAA;AACF;AAGa;AACX,EAAA;AACA,EAAA;AACA,EAAA;AAAa;AAEb,EAAA;AAAc;AAEd,EAAA;AACF;AAGa;AACX,EAAA;AAAsB;AAEtB,EAAA;AAAqB;AAErB,EAAA;AAAwB;AAExB,EAAA;AAAiB;AAEjB,EAAA;AACF;AHoEU;AACA;AI1KG;AAmCX,EAAA;AACmB,IAAA;AACA,IAAA;AAGZ,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AAMA,IAAA;AACH,MAAA;AACE,QAAA;AACE,UAAA;AACF,QAAA;AACE,UAAA;AACF,QAAA;AACD,MAAA;AACH,IAAA;AAGK,IAAA;AACH,MAAA;AACE,QAAA;AACD,MAAA;AACH,IAAA;AAGK,IAAA;AACH,MAAA;AACE,QAAA;AACE,UAAA;AACF,QAAA;AACD,MAAA;AACH,IAAA;AAKK,IAAA;AACH,MAAA;AACE,QAAA;AACE,UAAA;AACF,QAAA;AACD,MAAA;AACH,IAAA;AAMK,IAAA;AACH,MAAA;AACE,QAAA;AACE,UAAA;AACA,UAAA;AACA,UAAA;AACF,QAAA;AACD,MAAA;AACH,IAAA;AACF,EAAA;AA/DmB,EAAA;AACA,EAAA;AApCX,kBAAA;AACA,kBAAA;AACA,kBAAA;AACA,kBAAA;AACA,kBAAA;AAEA,kBAAA;AACA,kBAAA;AACA,mBAAA;AACA,mBAAA;AAA2B;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAS3B,mBAAA;AAAsC;AAEtC,mBAAA;AAAY;AAAA;AAAA;AAAA;AAAA;AAAA;AAOZ,mBAAA;AAES,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AAoEb,EAAA;AACF,IAAA;AACF,EAAA;AAEM,EAAA;AACA,IAAA;AAEJ,IAAA;AACM,MAAA;AAGJ,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACE,UAAA;AACA,UAAA;AACF,QAAA;AACA,QAAA;AAA2B,QAAA;AAE3B,QAAA;AACF,MAAA;AAIA,MAAA;AACA,MAAA;AAEA,MAAA;AAEA,MAAA;AACD,IAAA;AACH,EAAA;AAEA,EAAA;AACO,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACL,IAAA;AACF,EAAA;AAEA,EAAA;AACO,IAAA;AACL,IAAA;AACF,EAAA;AAEA,EAAA;AACO,IAAA;AACL,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAMA,EAAA;AACO,IAAA;AACL,IAAA;AACF,EAAA;AAAA;AAGA,EAAA;AACO,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAcM,EAAA;AACA,IAAA;AACC,IAAA;AACD,IAAA;AACE,MAAA;AACF,QAAA;AACE,UAAA;AACF,QAAA;AACA,QAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AACE,MAAA;AACF,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAUc,EAAA;AACR,IAAA;AAEE,IAAA;AACF,IAAA;AACA,IAAA;AAGF,MAAA;AACA,MAAA;AACF,IAAA;AAKK,IAAA;AACA,IAAA;AACC,IAAA;AACR,EAAA;AAAA;AAGQ,EAAA;AACA,IAAA;AACN,IAAA;AACM,MAAA;AACJ,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACD,MAAA;AACD,MAAA;AACA,MAAA;AACE,QAAA;AACA,QAAA;AACC,MAAA;AACJ,IAAA;AACH,EAAA;AAEQ,EAAA;AACD,IAAA;AACA,IAAA;AACA,IAAA;AACL,IAAA;AACF,EAAA;AAEQ,EAAA;AACD,IAAA;AACA,IAAA;AACH,MAAA;AACC,IAAA;AAEE,IAAA;AACP,EAAA;AAEQ,EAAA;AACF,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AACF,EAAA;AAEQ,EAAA;AACD,IAAA;AACA,IAAA;AACA,IAAA;AACC,MAAA;AACA,MAAA;AAKF,QAAA;AACA,QAAA;AACE,UAAA;AACD,QAAA;AACH,MAAA;AACC,IAAA;AACL,EAAA;AAEQ,EAAA;AACF,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AACF,EAAA;AAEQ,EAAA;AACD,IAAA;AACA,IAAA;AACD,IAAA;AACF,MAAA;AACF,IAAA;AACK,IAAA;AACA,IAAA;AACL,IAAA;AACK,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AACF;AJkHU;AACA;AKlbO;AACX,EAAA;AACG,EAAA;AACC,IAAA;AACA,IAAA;AACN,IAAA;AACF,EAAA;AACF;ALobU;AACA;AM3ZG;AA0BX,EAAA;AACU,IAAA;AAGH,IAAA;AACH,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACE,QAAA;AAEA,QAAA;AACF,MAAA;AACF,IAAA;AACF,EAAA;AAvBU,EAAA;AA1BqB,mBAAA;AACvB,mBAAA;AACA,mBAAA;AACA,mBAAA;AACA,mBAAA;AACA,mBAAA;AAAuD;AAEvD,mBAAA;AAEA,mBAAA;AACA,mBAAA;AAEA,mBAAA;AAES,EAAA;AAqCb,EAAA;AACF,IAAA;AACF,EAAA;AAEM,EAAA;AACA,IAAA;AAEC,IAAA;AAED,IAAA;AACA,IAAA;AACF,MAAA;AACF,IAAA;AAGE,MAAA;AACA,MAAA;AACF,IAAA;AACK,IAAA;AAEA,IAAA;AACH,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACF,IAAA;AAEK,IAAA;AACH,MAAA;AACI,MAAA;AACA,MAAA;AACF,QAAA;AACF,MAAA;AACE,QAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AAEK,IAAA;AACH,MAAA;AACI,MAAA;AAGF,QAAA;AACA,QAAA;AACF,MAAA;AACI,MAAA;AACF,QAAA;AACF,MAAA;AACE,QAAA;AACF,MAAA;AACF,IAAA;AAEK,IAAA;AAEL,IAAA;AACF,EAAA;AAEA,EAAA;AACO,IAAA;AACA,IAAA;AACA,IAAA;AAED,IAAA;AACF,MAAA;AACA,MAAA;AACA,MAAA;AACI,MAAA;AACF,QAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AAEK,IAAA;AACP,EAAA;AAEK,EAAA;AACC,IAAA;AACF,MAAA;AACF,IAAA;AACM,MAAA;AACF,QAAA;AACF,MAAA;AACF,IAAA;AACF,EAAA;AAEA,EAAA;AACO,IAAA;AACL,IAAA;AACF,EAAA;AAEA,EAAA;AACO,IAAA;AACL,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAQA,EAAA;AACM,IAAA;AACC,IAAA;AACA,IAAA;AAED,IAAA;AACF,MAAA;AACA,MAAA;AACA,MAAA;AACI,MAAA;AACF,QAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AAEK,IAAA;AACP,EAAA;AAEQ,EAAA;AACD,IAAA;AAED,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AAEK,IAAA;AACC,IAAA;AAEA,IAAA;AACA,MAAA;AACJ,MAAA;AACA,MAAA;AACE,QAAA;AACC,MAAA;AACL,IAAA;AAEA,IAAA;AACF,EAAA;AAEQ,EAAA;AACA,IAAA;AACN,IAAA;AACE,MAAA;AACF,IAAA;AACF,EAAA;AAEQ,EAAA;AACD,IAAA;AACA,IAAA;AAKD,MAAA;AAIA,QAAA;AACA,QAAA;AACF,MAAA;AACI,MAAA;AACF,QAAA;AACF,MAAA;AACC,IAAA;AACL,EAAA;AAEQ,EAAA;AACF,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AACF,EAAA;AAEQ,EAAA;AACF,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AACF,EAAA;AAEc,EAAA;AAER,IAAA;AACA,IAAA;AAGF,MAAA;AACI,MAAA;AAGF,QAAA;AACF,MAAA;AACF,IAAA;AACE,MAAA;AACF,IAAA;AAEK,IAAA;AAIC,IAAA;AACA,IAAA;AACN,IAAA;AAEA,IAAA;AACF,EAAA;AAEQ,EAAA;AACD,IAAA;AACL,IAAA;AACF,EAAA;AAEQ,EAAA;AACD,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AACF;AN2VU;AACA;AO7nBG;AAOX,EAAA;AACU,IAAA;AACA,IAAA;AAeP,EAAA;AAhBO,EAAA;AACA,EAAA;AARF,mBAAA;AACA,mBAAA;AAEA,mBAAA;AACA,mBAAA;AAqBJ,EAAA;AACF,IAAA;AACF,EAAA;AAEQ,EAAA;AACD,IAAA;AACL,IAAA;AACF,EAAA;AAEM,EAAA;AAEA,IAAA;AACA,IAAA;AACE,MAAA;AACF,QAAA;AACF,MAAA;AAEE,QAAA;AACA,QAAA;AACF,MAAA;AACI,MAAA;AAEF,QAAA;AACA,QAAA;AACF,MAAA;AACF,IAAA;AACE,MAAA;AACF,IAAA;AAGM,IAAA;AAGA,IAAA;AAED,IAAA;AAEA,IAAA;AACH,MAAA;AAEA,MAAA;AACE,QAAA;AACE,UAAA;AACA,UAAA;AACA,UAAA;AAEF,QAAA;AACE,UAAA;AACA,UAAA;AAEF,QAAA;AAEE,UAAA;AAEF,QAAA;AACE,UAAA;AAEF,QAAA;AACE,UAAA;AACA,UAAA;AACJ,MAAA;AACF,IAAA;AAEK,IAAA;AACH,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACD,IAAA;AACH,EAAA;AAEQ,EAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACN,IAAA;AACA,IAAA;AACF,EAAA;AAEK,EAAA;AACE,oBAAA;AACP,EAAA;AAAA;AAGA,EAAA;AACO,oBAAA;AACP,EAAA;AAEA,EAAA;AACO,oBAAA;AACL,IAAA;AACE,sBAAA;AACA,MAAA;AACI,IAAA;AACD,IAAA;AACP,EAAA;AAEA,EAAA;AACO,IAAA;AACL,IAAA;AACF,EAAA;AAEA,EAAA;AACO,IAAA;AACL,IAAA;AACF,EAAA;AAEQ,EAAA;AAGA,IAAA;AAAO;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,IAAA;AAkDP,IAAA;AACN,IAAA;AACF,EAAA;AAEQ,EAAA;AACD,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AACF;APolBU;AACA;AQ/yBG;AACH,mBAAA;AACA,mBAAA;AAEL,EAAA;AACG,IAAA;AACC,IAAA;AACH,MAAA;AACA,MAAA;AACF,IAAA;AACI,IAAA;AACJ,IAAA;AACF,EAAA;AAEK,EAAA;AACG,IAAA;AACJ,MAAA;AACA,MAAA;AACF,IAAA;AACM,IAAA;AACN,IAAA;AACF,EAAA;AAEI,EAAA;AACE,IAAA;AACF,sBAAA;AACF,IAAA;AACE,MAAA;AACF,IAAA;AACF,EAAA;AAEK,EAAA;AACE,IAAA;AACC,IAAA;AACF,IAAA;AACF,MAAA;AACF,IAAA;AACF,EAAA;AAEA,EAAA;AACE,IAAA;AACF,EAAA;AAEO,EAAA;AACC,IAAA;AACF,IAAA;AACA,IAAA;AAEE,IAAA;AACJ,MAAA;AACA,sBAAA;AACD,IAAA;AAEK,IAAA;AACJ,MAAA;AACA,sBAAA;AACF,IAAA;AACA,oBAAA;AAEI,IAAA;AACF,MAAA;AACE,QAAA;AACE,UAAA;AACF,QAAA;AACE,UAAA;AAAiC,YAAA;AAAa,UAAA;AAC9C,UAAA;AACF,QAAA;AACF,MAAA;AACF,IAAA;AACE,MAAA;AACA,sBAAA;AACF,IAAA;AACF,EAAA;AAEA,EAAA;AACO,IAAA;AACA,IAAA;AACP,EAAA;AAEQ,EAAA;AACD,IAAA;AACP,EAAA;AACF;ARwyBU;AACA;AS72BG;AAIX,EAAA;AACmB,IAAA;AACA,IAAA;AAChB,EAAA;AAFgB,EAAA;AACA,EAAA;AALX,mBAAA;AACS,mBAAA;AAA4B;AAQ7C,EAAA;AACO,IAAA;AACP,EAAA;AAEA,EAAA;AACE,IAAA;AACF,EAAA;AAAA;AAG0B,EAAA;AACnB,IAAA;AACP,EAAA;AAAA;AAGA,EAAA;AACO,IAAA;AACD,IAAA;AACA,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AACA,IAAA;AACE,MAAA;AACI,MAAA;AACF,QAAA;AACA,QAAA;AACF,MAAA;AACF,IAAA;AAEI,IAAA;AACF,MAAA;AACF,IAAA;AACE,MAAA;AACF,IAAA;AACK,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AASQ,EAAA;AACF,IAAA;AACF,MAAA;AACI,MAAA;AAEN,IAAA;AACA,IAAA;AACF,EAAA;AAAA;AAGQ,EAAA;AACF,IAAA;AACA,IAAA;AAEJ,IAAA;AACE,MAAA;AAEE,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACJ,IAAA;AAEA,IAAA;AACM,MAAA;AACH,MAAA;AACA,MAAA;AACH,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAMQ,EAAA;AACN,IAAA;AACE,MAAA;AAA2B,QAAA;AAC3B,MAAA;AACA,MAAA;AAA2B,QAAA;AAC3B,MAAA;AACA,MAAA;AAA2B,QAAA;AAC3B,MAAA;AAA2B,QAAA;AAC3B,MAAA;AAA2B,QAAA;AAC7B,IAAA;AACF,EAAA;AACF;ATy2BU;AACA;AUj+BG;AAIX,EAAA;AACmB,IAAA;AACA,IAAA;AAChB,EAAA;AAFgB,EAAA;AACA,EAAA;AALF,mBAAA;AACA,mBAAA;AAA4D;AAQnD,EAAA;AACnB,IAAA;AACP,EAAA;AAAA;AAGA,EAAA;AACO,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAKQ,EAAA;AACF,IAAA;AACJ,IAAA;AACE,MAAA;AACI,MAAA;AACF,QAAA;AACA,QAAA;AACF,MAAA;AACF,IAAA;AAEM,IAAA;AACA,IAAA;AACF,IAAA;AAGE,IAAA;AACF,IAAA;AACF,MAAA;AACF,IAAA;AAEK,IAAA;AACL,IAAA;AACF,EAAA;AACF;AV69BU;AACA;AWhgCG;AAIX,EAAA;AACmB,IAAA;AACA,IAAA;AAEA,IAAA;AACA,IAAA;AAIZ,IAAA;AACH,MAAA;AACE,QAAA;AACD,MAAA;AACH,IAAA;AAGK,IAAA;AACH,MAAA;AACE,QAAA;AACA,QAAA;AACE,UAAA;AACD,QAAA;AACF,MAAA;AACH,IAAA;AACF,EAAA;AAvBmB,EAAA;AACA,EAAA;AAEA,EAAA;AACA,EAAA;AARF,mBAAA;AACT,mBAAA;AA4BJ,EAAA;AACF,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAQM,EAAA;AACE,IAAA;AACF,IAAA;AACC,IAAA;AACP,EAAA;AAEgB,EAAA;AACV,IAAA;AACA,IAAA;AAEF,MAAA;AACI,MAAA;AACN,IAAA;AACK,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AASM,EAAA;AACE,IAAA;AACD,IAAA;AACD,IAAA;AAEA,IAAA;AACJ,IAAA;AACE,MAAA;AAGA,MAAA;AACA,MAAA;AACA,MAAA;AACF,IAAA;AACK,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAOe,EAAA;AACP,IAAA;AACN,IAAA;AACE,MAAA;AACF,IAAA;AACM,IAAA;AAEN,IAAA;AACE,MAAA;AACE,QAAA;AACC,QAAA;AACC,UAAA;AACE,YAAA;AACF,UAAA;AACF,QAAA;AACF,MAAA;AACA,MAAA;AACA,MAAA;AACE,QAAA;AACA,QAAA;AACC,MAAA;AACJ,IAAA;AACH,EAAA;AAEQ,EAAA;AACN,IAAA;AACK,IAAA;AACA,IAAA;AACP,EAAA;AACF;AXo/BU;AACA;AY7mCG;AAOX,EAAA;AACmB,IAAA;AAEA,IAAA;AACA,IAAA;AAGZ,IAAA;AACH,MAAA;AACE,QAAA;AACE,UAAA;AACA,UAAA;AACD,QAAA;AACF,MAAA;AACH,IAAA;AACF,EAAA;AAdmB,EAAA;AAEA,EAAA;AACA,EAAA;AAAA;AATF,mBAAA;AAAsC;AAEtC,mBAAA;AACT,mBAAA;AAA2B;AAoBnC,EAAA;AACO,IAAA;AACP,EAAA;AAAA;AAGA,EAAA;AACQ,IAAA;AACF,IAAA;AACC,IAAA;AACP,EAAA;AAEA,EAAA;AACO,IAAA;AACP,EAAA;AAEA,EAAA;AACO,IAAA;AACP,EAAA;AAAA;AAGA,EAAA;AACE,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAQM,EAAA;AACE,IAAA;AACD,IAAA;AAEL,IAAA;AACE,MAAA;AACF,IAAA;AACA,IAAA;AACE,MAAA;AACF,IAAA;AAEI,IAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACD,MAAA;AACH,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAOe,EAAA;AACP,IAAA;AACA,IAAA;AACA,IAAA;AAEN,IAAA;AACE,MAAA;AACE,QAAA;AACC,QAAA;AACC,UAAA;AACA,UAAA;AACF,QAAA;AACF,MAAA;AAEA,MAAA;AAEA,MAAA;AACE,QAAA;AACA,QAAA;AACC,MAAA;AACJ,IAAA;AACH,EAAA;AAEQ,EAAA;AACN,IAAA;AACK,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AACF;AZimCU;AACA;AahsCG;AAeX,EAAA;AAA6B,IAAA;AAItB,IAAA;AACH,MAAA;AACE,QAAA;AACE,UAAA;AACA,UAAA;AACF,QAAA;AACD,MAAA;AACH,IAAA;AAGK,IAAA;AACH,MAAA;AACE,QAAA;AACE,UAAA;AACA,UAAA;AACF,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACD,MAAA;AACH,IAAA;AACF,EAAA;AA5B6B,EAAA;AAdrB,mBAAA;AAAmB;AAEV,mBAAA;AAAwC;AAExC,mBAAA;AACT,mBAAA;AAAqD;AAErD,mBAAA;AAAgB;AAEhB,mBAAA;AAAa;AAEb,mBAAA;AACA,mBAAA;AAgCJ,EAAA;AACF,IAAA;AACF,EAAA;AAAA;AAGA,EAAA;AACO,IAAA;AACP,EAAA;AAEA,EAAA;AACO,IAAA;AACP,EAAA;AAEA,EAAA;AACO,IAAA;AACP,EAAA;AAEA,EAAA;AACO,IAAA;AACP,EAAA;AAEA,EAAA;AACE,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAOA,EAAA;AACO,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AAID,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AAEK,IAAA;AAEH,MAAA;AACF,IAAA;AAEK,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAOA,EAAA;AACE,IAAA;AACK,IAAA;AACL,IAAA;AACK,IAAA;AAEA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAOA,EAAA;AACO,IAAA;AACA,IAAA;AACH,MAAA;AACA,MAAA;AACF,IAAA;AACK,IAAA;AACP,EAAA;AAAA;AAGA,EAAA;AACO,IAAA;AACC,IAAA;AACF,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAQA,EAAA;AACO,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AAEA,EAAA;AACO,IAAA;AACD,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AASA,EAAA;AACM,IAAA;AACC,IAAA;AACC,IAAA;AACD,IAAA;AACD,IAAA;AACF,MAAA;AACF,IAAA;AACF,EAAA;AAEQ,EAAA;AACA,IAAA;AACN,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAOQ,EAAA;AACF,IAAA;AACC,IAAA;AAAkC,MAAA;AAA2B,IAAA;AACpE,EAAA;AAEc,EAAA;AACR,IAAA;AACA,IAAA;AACF,MAAA;AACI,MAAA;AACF,QAAA;AACA,QAAA;AACE,UAAA;AACA,UAAA;AACF,QAAA;AACE,UAAA;AACF,QAAA;AACF,MAAA;AACE,QAAA;AACF,MAAA;AACE,QAAA;AACF,MAAA;AACF,IAAA;AAEK,IAAA;AACP,EAAA;AAEQ,EAAA;AACD,IAAA;AACL,IAAA;AACK,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AACF;AbuqCU;AACA;Acx4CG;AACX,EAAA;AAA6B,IAAA;AAAwB,EAAA;AAAxB,EAAA;AAEX,EAAA;AACV,IAAA;AAGA,IAAA;AACA,IAAA;AAEF,IAAA;AACF,MAAA;AACF,IAAA;AAEA,IAAA;AACE,MAAA;AACA,MAAA;AACA,MAAA;AAGI,MAAA;AACF,QAAA;AAKA,QAAA;AACE,UAAA;AACA,UAAA;AACF,QAAA;AACF,MAAA;AAGI,MAAA;AACF,QAAA;AAMA,QAAA;AACE,UAAA;AACA,UAAA;AACA,UAAA;AAEA,UAAA;AAEA,UAAA;AACE,YAAA;AACA,YAAA;AACE,cAAA;AACA,cAAA;AAAa,YAAA;AAEjB,UAAA;AAEA,UAAA;AACF,QAAA;AACF,MAAA;AACgB,IAAA;AACpB,EAAA;AACF;Adw3CU;AACA;Ae58CJ;AACJ,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACQ,EAAA;AACR,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACF;AAEM;AACI,EAAA;AAAC,EAAA;AACF,EAAA;AAAC,EAAA;AACD,EAAA;AAAC,EAAA;AACA,EAAA;AAAC,EAAA;AACX;AAQM;AA4BO;AAwCX,EAAA;AACmB,IAAA;AACA,IAAA;AAEZ,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACH,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACD,IAAA;AAII,IAAA;AACA,IAAA;AACA,IAAA;AACH,MAAA;AACA,uBAAA;AACC,MAAA;AACD,MAAA;AACF,IAAA;AACK,IAAA;AACH,MAAA;AACC,MAAA;AACD,MAAA;AACF,IAAA;AACK,IAAA;AACH,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACD,IAAA;AACI,IAAA;AACC,MAAA;AACJ,MAAA;AACA,MAAA;AACA,MAAA;AACD,IAAA;AAKI,IAAA;AAIA,IAAA;AACH,MAAA;AACE,QAAA;AACA,wBAAA;AACD,MAAA;AACH,IAAA;AAGK,IAAA;AACH,MAAA;AAEE,QAAA;AAMA,QAAA;AACE,UAAA;AACA,UAAA;AACE,YAAA;AACA,YAAA;AACF,UAAA;AACF,QAAA;AAGA,QAAA;AACE,UAAA;AACE,YAAA;AAAM,cAAA;AAAU,YAAA;AAAW,YAAA;AAC7B,UAAA;AACF,QAAA;AACD,MAAA;AACH,IAAA;AAKK,IAAA;AACH,MAAA;AACE,QAAA;AACE,UAAA;AAKA,UAAA;AACF,QAAA;AACD,MAAA;AACH,IAAA;AAKK,IAAA;AACH,MAAA;AACE,QAAA;AACE,UAAA;AACA,UAAA;AACF,QAAA;AACD,MAAA;AACH,IAAA;AAMK,IAAA;AACH,MAAA;AACE,QAAA;AACA,QAAA;AACD,MAAA;AACH,IAAA;AAGK,IAAA;AACH,MAAA;AACA,MAAA;AACD,IAAA;AACI,IAAA;AACH,MAAA;AACA,MAAA;AACD,IAAA;AAGI,IAAA;AACH,MAAA;AACE,QAAA;AACE,UAAA;AACE,YAAA;AACA,YAAA;AACF,UAAA;AACE,YAAA;AACA,YAAA;AACF,UAAA;AACE,YAAA;AACA,YAAA;AACF,UAAA;AACE,YAAA;AACA,YAAA;AACF,UAAA;AACE,YAAA;AACA,YAAA;AACF,UAAA;AACE,YAAA;AACA,YAAA;AACF,UAAA;AACE,YAAA;AACA,YAAA;AACJ,QAAA;AACD,MAAA;AACH,IAAA;AAGI,IAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AAGE,UAAA;AAGA,UAAA;AACE,YAAA;AACF,UAAA;AACF,QAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AAMI,IAAA;AACF,MAAA;AACE,QAAA;AAIE,UAAA;AACF,QAAA;AAEE,UAAA;AACF,QAAA;AACF,MAAA;AACA,MAAA;AAGE,QAAA;AACE,UAAA;AACF,QAAA;AACF,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACF,IAAA;AACF,EAAA;AAzNmB,EAAA;AACA,EAAA;AAzCX,EAAA;AACA,EAAA;AACA,mBAAA;AACO,mBAAA;AACP,mBAAA;AACA,EAAA;AACA,mBAAA;AAA2B;AAElB,mBAAA;AACT,mBAAA;AACS,EAAA;AACA,EAAA;AAAA;AAEA,EAAA;AAAA;AAEA,EAAA;AAAA;AAEA,EAAA;AACT,mBAAA;AAA0D;AAAA;AAAA;AAAA;AAAA;AAMjD,EAAA;AAAA;AAEA,EAAA;AAAA;AAEA,EAAA;AAAA;AAET,mBAAA;AAAoD;AAAA;AAAA;AAAA;AAAA;AAAA;AAOpD,mBAAA;AA8NJ,EAAA;AACF,IAAA;AACF,EAAA;AAEI,EAAA;AACF,IAAA;AACF,EAAA;AAAA;AAGI,EAAA;AACF,IAAA;AACF,EAAA;AAAA;AAGI,EAAA;AACF,IAAA;AACF,EAAA;AAAA;AAGM,EAAA;AACE,IAAA;AACR,EAAA;AAAA;AAAA;AAKA,EAAA;AACE,IAAA;AACF,EAAA;AAAA;AAGA,EAAA;AACE,IAAA;AACF,EAAA;AAAA;AAGA,EAAA;AACE,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAYA,EAAA;AACE,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAYA,EAAA;AACM,IAAA;AACF,MAAA;AACF,IAAA;AACE,MAAA;AACF,IAAA;AACF,EAAA;AAAA;AAGA,EAAA;AACE,IAAA;AACF,EAAA;AAAA;AAGQ,EAAA;AACN,IAAA;AACF,EAAA;AAAA;AAGA,EAAA;AACE,IAAA;AACM,MAAA;AACY,IAAA;AACpB,EAAA;AAAA;AAGA,EAAA;AACE,IAAA;AACM,MAAA;AACY,IAAA;AACpB,EAAA;AAAA;AAGA,EAAA;AACE,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAiBA,EAAA;AACO,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AASA,EAAA;AACO,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAUA,EAAA;AACE,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAoBI,EAAA;AACE,IAAA;AACF,MAAA;AACF,IAAA;AACE,MAAA;AACF,IAAA;AACA,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAgBA,EAAA;AACO,IAAA;AACL,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAaA,EAAA;AACO,IAAA;AACL,IAAA;AACF,EAAA;AAAA;AAqBG,EAAA;AACD,IAAA;AACF,EAAA;AAAA;AAKK,EAAA;AACH,IAAA;AACF,EAAA;AAEI,EAAA;AACG,IAAA;AACP,EAAA;AAKO,EAAA;AACL,IAAA;AACF,EAAA;AA4BK,EAAA;AACE,IAAA;AAGC,IAAA;AACA,IAAA;AAED,IAAA;AACP,EAAA;AAEQ,EAAA;AACD,IAAA;AACD,IAAA;AACF,MAAA;AACE,QAAA;AAEF,MAAA;AACF,IAAA;AACI,IAAA;AACF,MAAA;AACE,QAAA;AAEF,MAAA;AACF,IAAA;AACF,EAAA;AAAA;AAGM,EAAA;AAIA,IAAA;AACF,MAAA;AACF,IAAA;AACA,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAQQ,EAAA;AACN,IAAA;AACE,MAAA;AACI,MAAA;AACF,QAAA;AACA,QAAA;AACF,MAAA;AACI,MAAA;AACJ,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACA,MAAA;AACE,QAAA;AACA,QAAA;AACE,UAAA;AACF,QAAA;AACD,MAAA;AACD,MAAA;AACE,QAAA;AACA,QAAA;AACF,MAAA;AACA,MAAA;AACD,IAAA;AACH,EAAA;AAAA;AAGQ,EAAA;AACD,IAAA;AACA,IAAA;AACP,EAAA;AAEW,EAAA;AACT,IAAA;AACF,EAAA;AAEU,EAAA;AACR,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAiBQ,EAAA;AAIA,IAAA;AACA,IAAA;AACF,IAAA;AAEE,IAAA;AAEI,MAAA;AACJ,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACA,MAAA;AACE,QAAA;AACA,QAAA;AACE,UAAA;AACF,QAAA;AAEE,UAAA;AACF,QAAA;AACA,QAAA;AAA2C,QAAA;AAE5C,MAAA;AACD,MAAA;AACE,QAAA;AACA,QAAA;AACF,MAAA;AACA,MAAA;AAEF,IAAA;AAGA,IAAA;AAA6B,IAAA;AAG5B,IAAA;AAGA,IAAA;AAEC,IAAA;AACA,IAAA;AACA,IAAA;AACF,IAAA;AACE,IAAA;AAEA,IAAA;AACJ,MAAA;AACA,MAAA;AACG,MAAA;AACD,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACA,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACA,MAAA;AAKE,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACA,MAAA;AACE,QAAA;AACF,MAAA;AACA,MAAA;AACE,QAAA;AACA,QAAA;AACA,wBAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACF,IAAA;AAEI,IAAA;AACF,MAAA;AACF,IAAA;AAEA,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAaA,EAAA;AACO,IAAA;AACA,IAAA;AACD,IAAA;AACF,MAAA;AACF,IAAA;AACK,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAMA,EAAA;AACO,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAwCkB,EAAA;AAChB,IAAA;AACF,EAAA;AAEA,EAAA;AACO,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAmBQ,EAAA;AACD,IAAA;AACL,IAAA;AAAe,MAAA;AAAmC,IAAA;AACpD,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAOQ,EAAA;AACF,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AACK,IAAA;AACP,EAAA;AAEQ,EAAA;AACA,IAAA;AACJ,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACF,IAAA;AAEI,IAAA;AAEF,MAAA;AACE,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AACD,MAAA;AACH,IAAA;AAGA,IAAA;AACK,MAAA;AACH,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACD,IAAA;AACH,EAAA;AAEQ,EAAA;AACD,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AAEA,IAAA;AACH,MAAA;AACI,MAAA;AACL,IAAA;AAEI,IAAA;AACH,MAAA;AACA,MAAA;AACE,QAAA;AACE,UAAA;AACA,UAAA;AACA,UAAA;AACF,QAAA;AACE,UAAA;AACA,UAAA;AACF,QAAA;AACE,UAAA;AACA,UAAA;AACF,QAAA;AACE,UAAA;AACA,UAAA;AACA,UAAA;AACJ,MAAA;AACD,IAAA;AAKI,oBAAA;AACA,IAAA;AACC,MAAA;AAAA;AAAA;AAAA;AAIH,MAAA;AACH,IAAA;AAEK,IAAA;AACP,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAmBc,EAAA;AACP,IAAA;AACC,IAAA;AACA,IAAA;AAGD,IAAA;AAEC,IAAA;AAEA,IAAA;AACR,EAAA;AAEQ,EAAA;AACD,IAAA;AACA,oBAAA;AACA,IAAA;AACD,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AACK,IAAA;AACP,EAAA;AAEQ,EAAA;AACF,IAAA;AACC,IAAA;AACA,IAAA;AACA,oBAAA;AACA,IAAA;AAEA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AAED,IAAA;AACF,MAAA;AACA,MAAA;AACF,IAAA;AAEA,IAAA;AACK,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AACF;AfwuCU;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA","file":"/Users/gwakko/Projects/shared-websocket/dist/chunk-JG2LC2FF.cjs","sourcesContent":[null,"export function generateId(): string {\n  if (typeof crypto !== 'undefined' && crypto.randomUUID) {\n    return crypto.randomUUID();\n  }\n  return `${Date.now()}-${Math.random().toString(36).slice(2, 11)}`;\n}\n","import './utils/disposable';\nimport { generateId } from './utils/id';\nimport type { BusMessage, Unsubscribe } from './types';\n\ntype Listener = (msg: BusMessage) => void;\n\nexport class MessageBus implements Disposable {\n  private channel: BroadcastChannel;\n  private listeners = new Map<string, Set<Listener>>();\n  private pendingRequests = new Map<string, { resolve: (v: unknown) => void; reject: (e: Error) => void; timer: ReturnType<typeof setTimeout> }>();\n\n  constructor(\n    channelName: string,\n    private readonly tabId: string,\n  ) {\n    this.channel = new BroadcastChannel(channelName);\n    this.channel.onmessage = (ev: MessageEvent<BusMessage>) => {\n      this.handleMessage(ev.data);\n    };\n  }\n\n  subscribe<T>(topic: string, fn: (data: T) => void): Unsubscribe {\n    const wrapper: Listener = (msg) => {\n      // `broadcast` is \"fan out to all tabs INCLUDING me\" — `broadcast()`\n      // explicitly self-delivers via handleMessage so the leader's own\n      // handlers fire for events it originated (e.g. ws:message after the\n      // leader receives a server frame). `publish` is fire-and-forget to\n      // OTHER tabs, so skip self-originated publish messages.\n      if (msg.type !== 'publish' || msg.source !== this.tabId) {\n        fn(msg.data as T);\n      }\n    };\n    this.addListener(topic, wrapper);\n    return () => this.removeListener(topic, wrapper);\n  }\n\n  publish<T>(topic: string, data: T): void {\n    this.postMessage({ topic, type: 'publish', data });\n  }\n\n  broadcast<T>(topic: string, data: T): void {\n    const msg = this.createMessage(topic, 'broadcast', data);\n    this.channel.postMessage(msg);\n    // Also deliver to self\n    this.handleMessage(msg);\n  }\n\n  async request<T, R>(topic: string, data: T, timeout = 5000): Promise<R> {\n    const msg = this.createMessage(topic, 'request', data);\n    return new Promise<R>((resolve, reject) => {\n      const timer = setTimeout(() => {\n        this.pendingRequests.delete(msg.id);\n        reject(new Error(`MessageBus.request: timeout for topic \"${topic}\"`));\n      }, timeout);\n      this.pendingRequests.set(msg.id, { resolve: resolve as (v: unknown) => void, reject, timer });\n      this.channel.postMessage(msg);\n    });\n  }\n\n  respond<T, R>(topic: string, fn: (data: T) => R | Promise<R>): Unsubscribe {\n    const wrapper: Listener = async (msg) => {\n      if (msg.type !== 'request' || msg.source === this.tabId) return;\n      const result = await fn(msg.data as T);\n      this.postMessage({ topic, type: 'response', data: { requestId: msg.id, result } });\n    };\n    this.addListener(topic, wrapper);\n    return () => this.removeListener(topic, wrapper);\n  }\n\n  private handleMessage(msg: BusMessage): void {\n    // Handle response to pending request\n    if (msg.type === 'response') {\n      const payload = msg.data as { requestId: string; result: unknown };\n      const pending = this.pendingRequests.get(payload.requestId);\n      if (pending) {\n        clearTimeout(pending.timer);\n        this.pendingRequests.delete(payload.requestId);\n        pending.resolve(payload.result);\n        return;\n      }\n    }\n\n    const listeners = this.listeners.get(msg.topic);\n    if (listeners) {\n      for (const fn of listeners) fn(msg);\n    }\n  }\n\n  private postMessage(partial: Pick<BusMessage, 'topic' | 'type' | 'data'>): void {\n    this.channel.postMessage(this.createMessage(partial.topic, partial.type, partial.data));\n  }\n\n  private createMessage(topic: string, type: BusMessage['type'], data: unknown): BusMessage {\n    return { id: generateId(), source: this.tabId, topic, type, data, timestamp: Date.now() };\n  }\n\n  private addListener(topic: string, fn: Listener): void {\n    let set = this.listeners.get(topic);\n    if (!set) {\n      set = new Set();\n      this.listeners.set(topic, set);\n    }\n    set.add(fn);\n  }\n\n  private removeListener(topic: string, fn: Listener): void {\n    this.listeners.get(topic)?.delete(fn);\n  }\n\n  [Symbol.dispose](): void {\n    for (const pending of this.pendingRequests.values()) {\n      clearTimeout(pending.timer);\n      pending.reject(new Error('MessageBus disposed'));\n    }\n    this.pendingRequests.clear();\n    this.listeners.clear();\n    this.channel.close();\n  }\n}\n","/**\n * Centralized internal identifiers and tuning defaults.\n *\n * The string values here are the cross-tab `BroadcastChannel` topic names and\n * the internal `SubscriptionManager` event keys. They are an implementation\n * detail of a single library version — every tab runs the same build — not a\n * wire/server protocol, so the literal values can change freely as long as\n * they stay consistent within a release. Keeping them in one place removes the\n * scatter of magic strings and the risk of an emit/handler typo mismatch.\n */\n\n/** BroadcastChannel name shared by every tab of one SharedWebSocket instance. */\nexport const MESSAGE_BUS_CHANNEL = 'shared-ws';\n\n// ─── Cross-tab coordination topics (TabCoordinator) ──────────────────────────\n\nexport const COORD = {\n  /** Candidate announces an election (carries its tabId). */\n  ELECTION: 'coord:election',\n  /** Current leader rejects a candidate. */\n  REJECT: 'coord:reject',\n  /** Leader liveness beat. */\n  HEARTBEAT: 'coord:heartbeat',\n  /** Leader relinquished — followers should elect. */\n  ABDICATE: 'coord:abdicate',\n  /** Health probe sent to the leader. */\n  PING: 'coord:ping',\n  /** Forced demotion demanded by an active tab taking over. */\n  STEP_DOWN: 'coord:step-down',\n} as const;\n\n/** One-shot reply topic for a health ping keyed by replyId. */\nexport const coordPongReply = (replyId: string): string => `coord:pong:${replyId}`;\n\n// ─── WebSocket fan-out / routing topics (SharedWebSocket, Outbox) ─────────────\n\nexport const BUS = {\n  /** Incoming server frame fanned out to every tab. */\n  MESSAGE: 'ws:message',\n  /** Follower → leader outbound dispatch request. */\n  DISPATCH: 'ws:dispatch',\n  /** Leader → originator ack: drop the buffered entry. */\n  DISPATCH_FLUSHED: 'ws:dispatch-flushed',\n  /** New leader asks every tab for still-pending dispatches. */\n  GATHER_PENDING: 'ws:gather-pending',\n  /** New leader asks every tab for its channels/topics. */\n  GATHER_SUBS: 'ws:gather-subs',\n  /** Cross-tab state sync (no server roundtrip). */\n  SYNC: 'ws:sync',\n  /** Request/response routed to the leader's socket. */\n  REQUEST: 'ws:request',\n  /** Follower asks the leader to reconnect. */\n  RECONNECT: 'ws:reconnect',\n  /** Follower hint to reconnect a failed socket after re-auth. */\n  AUTH_RESUME: 'ws:authenticate-resume',\n  /** Lifecycle events broadcast to all tabs. */\n  LIFECYCLE: 'ws:lifecycle',\n} as const;\n\n/** Reply topic for a pending-dispatch gather keyed by replyId. */\nexport const busPendingReply = (replyId: string): string => `ws:pending:${replyId}`;\n/** Reply topic for a subscription gather keyed by replyId. */\nexport const busSubsReply = (replyId: string): string => `ws:subs:${replyId}`;\n\n// ─── Internal lifecycle event keys (SubscriptionManager) ─────────────────────\n\nexport const LIFECYCLE = {\n  CONNECT: '$lifecycle:connect',\n  DISCONNECT: '$lifecycle:disconnect',\n  RECONNECTING: '$lifecycle:reconnecting',\n  RECONNECT_FAILED: '$lifecycle:reconnectFailed',\n  LEADER: '$lifecycle:leader',\n  ERROR: '$lifecycle:error',\n  AUTH: '$lifecycle:auth',\n  ACTIVE: '$lifecycle:active',\n} as const;\n\n/** Sync-store key holding the runtime auth token. */\nexport const AUTH_TOKEN_KEY = '$auth:token';\n\n// ─── Tuning defaults ─────────────────────────────────────────────────────────\n\n/** TabCoordinator timing defaults (ms). */\nexport const COORD_DEFAULTS = {\n  ELECTION_TIMEOUT: 200,\n  HEARTBEAT_INTERVAL: 2000,\n  LEADER_TIMEOUT: 5000,\n  LEADER_PING_TIMEOUT: 1500,\n  /** How often a follower re-checks the leader heartbeat. */\n  LEADER_CHECK_INTERVAL: 1000,\n} as const;\n\n/** SharedSocket connection defaults. */\nexport const SOCKET_DEFAULTS = {\n  RECONNECT_MAX_DELAY: 30_000,\n  HEARTBEAT_INTERVAL: 30_000,\n  SEND_BUFFER: 100,\n  /** Initial reconnect backoff before exponential growth. */\n  BACKOFF_BASE: 1000,\n  /** Default \"auth failed — stop retrying\" close code (PolicyViolation). */\n  AUTH_FAILURE_CLOSE_CODE: 1008,\n} as const;\n\n/** SharedWebSocket / Outbox defaults. */\nexport const WS_DEFAULTS = {\n  OUTBOUND_BUFFER_SIZE: 100,\n  /** Cross-tab subscription-gather window (ms). */\n  GATHER_SUBS_TIMEOUT: 150,\n  /** Cross-tab pending-dispatch-gather window (ms). */\n  GATHER_PENDING_TIMEOUT: 100,\n  /** Default request/response timeout (ms). */\n  REQUEST_TIMEOUT: 5000,\n  /** Default Channel.ready ack timeout (ms). */\n  CHANNEL_ACK_TIMEOUT: 5000,\n} as const;\n","import './utils/disposable';\nimport { MessageBus } from './MessageBus';\nimport { generateId } from './utils/id';\nimport { COORD, COORD_DEFAULTS, coordPongReply } from './constants';\nimport type { Unsubscribe } from './types';\n\ninterface CoordinatorOptions {\n  electionTimeout?: number;   // ms to wait for rejection (default 200)\n  heartbeatInterval?: number; // ms between heartbeats (default 2000)\n  leaderTimeout?: number;     // ms without heartbeat to trigger election (default 5000)\n  leaderPingTimeout?: number; // ms to wait for a leader pong on active-tab verify (default 1500)\n}\n\nexport class TabCoordinator implements Disposable {\n  private _isLeader = false;\n  private heartbeatTimer: ReturnType<typeof setInterval> | null = null;\n  private leaderCheckTimer: ReturnType<typeof setInterval> | null = null;\n  private lastHeartbeat = 0;\n  private disposed = false;\n\n  private onBecomeLeaderFns = new Set<() => void>();\n  private onLoseLeadershipFns = new Set<() => void>();\n  private onLeaderUnhealthyFns = new Set<() => void>();\n  private cleanups: Unsubscribe[] = [];\n\n  /**\n   * Optional predicate supplied by the owner that reports whether THIS tab's\n   * leader socket is actually alive. A backgrounded leader whose socket has\n   * silently died still has `_isLeader === true`, so without this check an\n   * election would keep deferring to a zombie. When set, a leader only\n   * answers health pings (and passes self-verification) while it returns true.\n   */\n  private healthCheck: (() => boolean) | null = null;\n  /** Re-entrancy guard so rapid visibility toggles don't stack verifications. */\n  private verifying = false;\n  /**\n   * Set while an election is in flight. Lets a concurrent election from a\n   * tab with a lower tabId pre-empt this one (deterministic tie-break) so two\n   * tabs electing at the same instant can't both become leader (split-brain).\n   * Calling it makes this tab yield and become a follower.\n   */\n  private electionAbort: (() => void) | null = null;\n\n  private readonly electionTimeout: number;\n  private readonly heartbeatInterval: number;\n  private readonly leaderTimeout: number;\n  private readonly leaderPingTimeout: number;\n\n  constructor(\n    private readonly bus: MessageBus,\n    private readonly tabId: string,\n    options: CoordinatorOptions = {},\n  ) {\n    this.electionTimeout = options.electionTimeout ?? COORD_DEFAULTS.ELECTION_TIMEOUT;\n    this.heartbeatInterval = options.heartbeatInterval ?? COORD_DEFAULTS.HEARTBEAT_INTERVAL;\n    this.leaderTimeout = options.leaderTimeout ?? COORD_DEFAULTS.LEADER_TIMEOUT;\n    this.leaderPingTimeout = options.leaderPingTimeout ?? COORD_DEFAULTS.LEADER_PING_TIMEOUT;\n\n    // Listen for election requests. If we're already leader, reject so the\n    // candidate backs off. If we're mid-election ourselves, the candidate\n    // with the smaller tabId wins and we yield — a deterministic tie-break\n    // that prevents two simultaneous electors from both becoming leader.\n    this.cleanups.push(\n      this.bus.subscribe<{ tabId: string }>(COORD.ELECTION, (msg) => {\n        if (this._isLeader) {\n          this.bus.publish(COORD.REJECT, { tabId: this.tabId });\n        } else if (this.electionAbort && msg.tabId < this.tabId) {\n          this.electionAbort();\n        }\n      }),\n    );\n\n    // Listen for heartbeats\n    this.cleanups.push(\n      this.bus.subscribe<{ tabId: string }>(COORD.HEARTBEAT, () => {\n        this.lastHeartbeat = Date.now();\n      }),\n    );\n\n    // Listen for abdication\n    this.cleanups.push(\n      this.bus.subscribe(COORD.ABDICATE, () => {\n        if (!this._isLeader && !this.disposed) {\n          this.elect();\n        }\n      }),\n    );\n\n    // Answer health pings — but ONLY while we are a leader whose socket is\n    // actually alive. A zombie leader (timer-throttled tab, dead socket)\n    // stays silent so the asking tab knows to take over.\n    this.cleanups.push(\n      this.bus.subscribe<{ replyId: string }>(COORD.PING, (req) => {\n        if (this._isLeader && !this.disposed && (this.healthCheck?.() ?? true)) {\n          this.bus.publish(coordPongReply(req.replyId), { tabId: this.tabId });\n        }\n      }),\n    );\n\n    // Forced step-down — another tab found us unresponsive and is taking\n    // over. Demote silently: the demanding tab runs the election itself, so\n    // we must NOT publish `coord:abdicate` (that would make every follower\n    // elect at once and risk split-brain).\n    this.cleanups.push(\n      this.bus.subscribe(COORD.STEP_DOWN, () => {\n        if (this._isLeader && !this.disposed) {\n          this._isLeader = false;\n          this.stopHeartbeat();\n          for (const fn of this.onLoseLeadershipFns) fn();\n        }\n      }),\n    );\n  }\n\n  get isLeader(): boolean {\n    return this._isLeader;\n  }\n\n  async elect(): Promise<void> {\n    if (this.disposed) return;\n\n    return new Promise<void>((resolve) => {\n      let settled = false;\n\n      // `won` true → become leader; false → yield and monitor as follower.\n      const finish = (won: boolean) => {\n        if (settled) return;\n        settled = true;\n        clearTimeout(timer);\n        unsub();\n        this.electionAbort = null;\n        if (this.disposed) {\n          resolve();\n          return;\n        }\n        if (won) this.becomeLeader();\n        else this.startLeaderCheck();\n        resolve();\n      };\n\n      // A live leader rejected us, OR a lower-tabId concurrent candidate\n      // pre-empted us — either way we step back and become a follower.\n      const unsub = this.bus.subscribe(COORD.REJECT, () => finish(false));\n      this.electionAbort = () => finish(false);\n\n      this.bus.publish(COORD.ELECTION, { tabId: this.tabId });\n\n      const timer = setTimeout(() => finish(true), this.electionTimeout);\n    });\n  }\n\n  abdicate(): void {\n    if (!this._isLeader) return;\n    this._isLeader = false;\n    this.stopHeartbeat();\n    this.bus.publish(COORD.ABDICATE, { tabId: this.tabId });\n    for (const fn of this.onLoseLeadershipFns) fn();\n  }\n\n  onBecomeLeader(fn: () => void): Unsubscribe {\n    this.onBecomeLeaderFns.add(fn);\n    return () => this.onBecomeLeaderFns.delete(fn);\n  }\n\n  onLoseLeadership(fn: () => void): Unsubscribe {\n    this.onLoseLeadershipFns.add(fn);\n    return () => this.onLoseLeadershipFns.delete(fn);\n  }\n\n  /**\n   * Fired when a self-verification finds THIS tab is leader but its socket\n   * is unhealthy (per `setHealthCheck`). The owner should reconnect.\n   */\n  onLeaderUnhealthy(fn: () => void): Unsubscribe {\n    this.onLeaderUnhealthyFns.add(fn);\n    return () => this.onLeaderUnhealthyFns.delete(fn);\n  }\n\n  /** Supply a predicate reporting whether this tab's leader socket is alive. */\n  setHealthCheck(fn: () => boolean): void {\n    this.healthCheck = fn;\n  }\n\n  /**\n   * Verify the active leader is alive and re-elect if it isn't. Call this\n   * when a tab becomes visible again after being idle: browser timer\n   * throttling can leave a backgrounded leader with a dead socket while its\n   * last heartbeat still looks recent enough that no follower would elect.\n   *\n   * - Leader tab → checks its own socket health; if unhealthy, fires\n   *   `onLeaderUnhealthy` so the owner can reconnect.\n   * - Follower tab → pings the leader and waits up to `leaderPingTimeout`.\n   *   If no healthy leader answers, forces a step-down and runs a fresh\n   *   election so this (active) tab can take over the connection.\n   */\n  async verifyLeader(): Promise<void> {\n    if (this.disposed || this.verifying) return;\n    this.verifying = true;\n    try {\n      if (this._isLeader) {\n        if (this.healthCheck && !this.healthCheck()) {\n          for (const fn of this.onLeaderUnhealthyFns) fn();\n        }\n        return;\n      }\n      await this.takeOverIfLeaderDead();\n    } finally {\n      this.verifying = false;\n    }\n  }\n\n  /**\n   * Follower-only: ping the current leader and take over the connection if no\n   * healthy leader answers. Shared by `verifyLeader()` (active-tab path) and\n   * the heartbeat-staleness check (covers the case where THIS tab stays active\n   * the whole time while a backgrounded leader's socket silently dies — there\n   * is no visibility change to trigger `verifyLeader`, but the leader-check\n   * timer keeps running and routes here).\n   */\n  private async takeOverIfLeaderDead(): Promise<void> {\n    if (this.disposed || this._isLeader) return;\n\n    const alive = await this.pingLeader();\n    if (this.disposed || this._isLeader) return;\n    if (alive) {\n      // Leader is healthy (its heartbeat may just have been throttled) —\n      // treat the pong as a fresh heartbeat so we don't immediately retry.\n      this.lastHeartbeat = Date.now();\n      return;\n    }\n\n    // No healthy leader answered — demand step-down, then take over. The\n    // step-down message is delivered before our election frame on the same\n    // ordered channel, so the old leader won't reject the election.\n    this.bus.publish(COORD.STEP_DOWN, { tabId: this.tabId });\n    this.stopLeaderCheck();\n    await this.elect();\n  }\n\n  /** Ping the current leader; resolve true if a healthy leader ponged in time. */\n  private pingLeader(): Promise<boolean> {\n    const replyId = generateId();\n    return new Promise<boolean>((resolve) => {\n      let answered = false;\n      const unsub = this.bus.subscribe(coordPongReply(replyId), () => {\n        if (answered) return;\n        answered = true;\n        unsub();\n        resolve(true);\n      });\n      this.bus.publish(COORD.PING, { replyId });\n      setTimeout(() => {\n        unsub();\n        if (!answered) resolve(false);\n      }, this.leaderPingTimeout);\n    });\n  }\n\n  private becomeLeader(): void {\n    this._isLeader = true;\n    this.stopLeaderCheck();\n    this.startHeartbeat();\n    for (const fn of this.onBecomeLeaderFns) fn();\n  }\n\n  private startHeartbeat(): void {\n    this.stopHeartbeat();\n    this.heartbeatTimer = setInterval(() => {\n      this.bus.publish(COORD.HEARTBEAT, { tabId: this.tabId });\n    }, this.heartbeatInterval);\n    // Send immediately\n    this.bus.publish(COORD.HEARTBEAT, { tabId: this.tabId });\n  }\n\n  private stopHeartbeat(): void {\n    if (this.heartbeatTimer) {\n      clearInterval(this.heartbeatTimer);\n      this.heartbeatTimer = null;\n    }\n  }\n\n  private startLeaderCheck(): void {\n    this.stopLeaderCheck();\n    this.lastHeartbeat = Date.now();\n    this.leaderCheckTimer = setInterval(() => {\n      if (this.disposed || this.verifying) return;\n      if (Date.now() - this.lastHeartbeat > this.leaderTimeout) {\n        // Heartbeat lapsed. Don't blindly elect — a zombie leader (alive tab,\n        // dead socket) would still reject the election and keep the connection\n        // stuck. Ping for real health first; take over only if nobody healthy\n        // answers. `verifying` guards against overlapping pings each tick.\n        this.verifying = true;\n        void this.takeOverIfLeaderDead().finally(() => {\n          this.verifying = false;\n        });\n      }\n    }, COORD_DEFAULTS.LEADER_CHECK_INTERVAL);\n  }\n\n  private stopLeaderCheck(): void {\n    if (this.leaderCheckTimer) {\n      clearInterval(this.leaderCheckTimer);\n      this.leaderCheckTimer = null;\n    }\n  }\n\n  [Symbol.dispose](): void {\n    this.disposed = true;\n    this.electionAbort = null;\n    if (this._isLeader) {\n      this.abdicate();\n    }\n    this.stopHeartbeat();\n    this.stopLeaderCheck();\n    for (const unsub of this.cleanups) unsub();\n    this.cleanups = [];\n    this.onBecomeLeaderFns.clear();\n    this.onLoseLeadershipFns.clear();\n    this.onLeaderUnhealthyFns.clear();\n  }\n}\n","/** Exponential backoff generator with jitter. */\nexport function* backoff(base = 1000, max = 30_000): Generator<number> {\n  let delay = base;\n  while (true) {\n    const jitter = delay * 0.25 * (Math.random() * 2 - 1);\n    yield Math.min(delay + jitter, max);\n    delay = Math.min(delay * 2, max);\n  }\n}\n","import './utils/disposable';\nimport { backoff } from './utils/backoff';\nimport { SOCKET_DEFAULTS } from './constants';\nimport type { SocketState, Unsubscribe, EventHandler } from './types';\n\ninterface SharedSocketOptions {\n  protocols?: string[];\n  reconnect?: boolean;\n  reconnectMaxDelay?: number;\n  /** Max reconnect attempts before giving up (default: Infinity). */\n  reconnectMaxRetries?: number;\n  /** Close codes that mean \"auth failed — stop reconnect.\" Default: [1008]. */\n  authFailureCloseCodes?: number[];\n  heartbeatInterval?: number;\n  /**\n   * Liveness watchdog. When `> 0`, the socket force-reconnects if NO inbound\n   * message (server data or a pong) arrives within this many ms. Detects\n   * silently-dropped connections that never fire `onclose` (sleep, network\n   * switch, captive portal). Requires the server to send periodic data or\n   * answer the heartbeat ping. Default: disabled (legacy fire-and-forget).\n   */\n  heartbeatTimeout?: number;\n  sendBuffer?: number;\n  auth?: () => string | Promise<string>;\n  authToken?: string;\n  authParam?: string;\n  /** Heartbeat payload (default: { type: \"ping\" }). */\n  pingPayload?: unknown;\n  /** Custom serializer (default: JSON.stringify). */\n  serialize?: (data: unknown) => string | ArrayBuffer | Blob;\n  /** Custom deserializer (default: JSON.parse). */\n  deserialize?: (raw: string | ArrayBuffer) => unknown;\n}\n\nexport class SharedSocket implements Disposable {\n  private ws: WebSocket | null = null;\n  private _state: SocketState = 'closed';\n  private buffer: unknown[] = [];\n  private disposed = false;\n  private heartbeatTimer: ReturnType<typeof setInterval> | null = null;\n  private reconnectTimer: ReturnType<typeof setTimeout> | null = null;\n  /** Timestamp of the last inbound message — drives the liveness watchdog. */\n  private lastInboundAt = 0;\n\n  private onMessageFns = new Set<EventHandler>();\n  private onStateChangeFns = new Set<(state: SocketState) => void>();\n\n  private reconnectAttempts = 0;\n\n  private readonly opts: Required<Omit<SharedSocketOptions, 'auth' | 'authToken' | 'authParam' | 'pingPayload' | 'serialize' | 'deserialize' | 'authFailureCloseCodes' | 'heartbeatTimeout'>> & {\n    authFailureCloseCodes: ReadonlySet<number>;\n    heartbeatTimeout: number;\n    auth?: () => string | Promise<string>;\n    authToken?: string;\n    authParam: string;\n    pingPayload: unknown;\n    serialize: (data: unknown) => string | ArrayBuffer | Blob;\n    deserialize: (raw: string | ArrayBuffer) => unknown;\n  };\n\n  constructor(\n    private url: string,\n    options: SharedSocketOptions = {},\n  ) {\n    this.opts = {\n      protocols: options.protocols ?? [],\n      reconnect: options.reconnect ?? true,\n      reconnectMaxDelay: options.reconnectMaxDelay ?? SOCKET_DEFAULTS.RECONNECT_MAX_DELAY,\n      reconnectMaxRetries: options.reconnectMaxRetries ?? Infinity,\n      authFailureCloseCodes: new Set(options.authFailureCloseCodes ?? [SOCKET_DEFAULTS.AUTH_FAILURE_CLOSE_CODE]),\n      heartbeatInterval: options.heartbeatInterval ?? SOCKET_DEFAULTS.HEARTBEAT_INTERVAL,\n      heartbeatTimeout: options.heartbeatTimeout ?? 0,\n      sendBuffer: options.sendBuffer ?? SOCKET_DEFAULTS.SEND_BUFFER,\n      auth: options.auth,\n      authToken: options.authToken,\n      authParam: options.authParam ?? 'token',\n      pingPayload: options.pingPayload ?? { type: 'ping' },\n      serialize: options.serialize ?? ((data: unknown) => JSON.stringify(data)),\n      deserialize: options.deserialize ?? ((raw: string | ArrayBuffer) => {\n        if (typeof raw === 'string') return JSON.parse(raw);\n        // ArrayBuffer → decode as UTF-8 then parse\n        return JSON.parse(new TextDecoder().decode(raw));\n      }),\n    };\n  }\n\n  get state(): SocketState {\n    return this._state;\n  }\n\n  async connect(): Promise<void> {\n    if (this.disposed) return;\n\n    this.setState('connecting');\n\n    let connectUrl: string;\n    try {\n      connectUrl = await this.buildUrl();\n    } catch {\n      // auth() threw or returned no token — pause reconnect until user\n      // provides fresh creds via ws.authenticate(token) or ws.reconnect().\n      this.setState('failed');\n      return;\n    }\n    this.ws = new WebSocket(connectUrl, this.opts.protocols);\n\n    this.ws.onopen = () => {\n      this.reconnectAttempts = 0;\n      this.lastInboundAt = Date.now();\n      this.setState('connected');\n      this.flushBuffer();\n      this.startHeartbeat();\n    };\n\n    this.ws.onmessage = (ev: MessageEvent) => {\n      this.lastInboundAt = Date.now();\n      let data: unknown;\n      try {\n        data = this.opts.deserialize(ev.data as string | ArrayBuffer);\n      } catch {\n        data = ev.data;\n      }\n      for (const fn of this.onMessageFns) fn(data);\n    };\n\n    this.ws.onclose = (ev) => {\n      this.stopHeartbeat();\n      if (this.opts.authFailureCloseCodes.has(ev.code)) {\n        // Auth-failure close code — don't burn retries with stale creds.\n        // User must call ws.authenticate(freshToken) or ws.reconnect() to resume.\n        this.setState('failed');\n        return;\n      }\n      if (!this.disposed && this.opts.reconnect) {\n        this.scheduleReconnect();\n      } else {\n        this.setState('closed');\n      }\n    };\n\n    this.ws.onerror = () => {\n      // onclose will fire after onerror\n    };\n  }\n\n  disconnect(): void {\n    this.disposed = true;\n    this.stopHeartbeat();\n    this.clearReconnect();\n\n    if (this.ws) {\n      this.ws.onclose = null;\n      this.ws.onmessage = null;\n      this.ws.onerror = null;\n      if (this.ws.readyState === WebSocket.OPEN || this.ws.readyState === WebSocket.CONNECTING) {\n        this.ws.close(1000, 'client disconnect');\n      }\n      this.ws = null;\n    }\n\n    this.setState('closed');\n  }\n\n  send(data: unknown): void {\n    if (this._state === 'connected' && this.ws?.readyState === WebSocket.OPEN) {\n      this.ws.send(this.opts.serialize(data));\n    } else if (this._state === 'reconnecting' || this._state === 'connecting') {\n      if (this.buffer.length < this.opts.sendBuffer) {\n        this.buffer.push(data);\n      }\n    }\n  }\n\n  onMessage(fn: EventHandler): Unsubscribe {\n    this.onMessageFns.add(fn);\n    return () => this.onMessageFns.delete(fn);\n  }\n\n  onStateChange(fn: (state: SocketState) => void): Unsubscribe {\n    this.onStateChangeFns.add(fn);\n    return () => this.onStateChangeFns.delete(fn);\n  }\n\n  /**\n   * Manually trigger a reconnect. Resets the retry counter and clears any\n   * scheduled backoff so the next attempt happens immediately. Use after\n   * `state === 'failed'` to let the user retry, or any time to force a\n   * fresh connection.\n   */\n  reconnect(): void {\n    if (this.disposed) return;\n    this.clearReconnect();\n    this.reconnectAttempts = 0;\n\n    if (this.ws) {\n      this.ws.onclose = null;\n      this.ws.onmessage = null;\n      this.ws.onerror = null;\n      if (this.ws.readyState === WebSocket.OPEN || this.ws.readyState === WebSocket.CONNECTING) {\n        this.ws.close(1000, 'manual reconnect');\n      }\n      this.ws = null;\n    }\n\n    void this.connect();\n  }\n\n  private scheduleReconnect(): void {\n    this.reconnectAttempts++;\n\n    if (this.reconnectAttempts > this.opts.reconnectMaxRetries) {\n      this.setState('failed');\n      return;\n    }\n\n    this.setState('reconnecting');\n    const gen = backoff(SOCKET_DEFAULTS.BACKOFF_BASE, this.opts.reconnectMaxDelay);\n\n    const attempt = () => {\n      if (this.disposed) return;\n      const delay = gen.next().value;\n      this.reconnectTimer = setTimeout(() => {\n        if (!this.disposed) this.connect();\n      }, delay);\n    };\n\n    attempt();\n  }\n\n  private flushBuffer(): void {\n    const pending = this.buffer.splice(0);\n    for (const item of pending) {\n      this.send(item);\n    }\n  }\n\n  private startHeartbeat(): void {\n    this.stopHeartbeat();\n    this.heartbeatTimer = setInterval(() => {\n      // Liveness watchdog: if enabled and we've heard nothing back within the\n      // window, the connection is silently dead (no close frame ever came).\n      // Force a reconnect rather than sending into the void forever.\n      if (\n        this.opts.heartbeatTimeout > 0 &&\n        this._state === 'connected' &&\n        Date.now() - this.lastInboundAt > this.opts.heartbeatTimeout\n      ) {\n        this.reconnect();\n        return;\n      }\n      if (this.ws?.readyState === WebSocket.OPEN) {\n        this.ws.send(this.opts.serialize(this.opts.pingPayload));\n      }\n    }, this.opts.heartbeatInterval);\n  }\n\n  private stopHeartbeat(): void {\n    if (this.heartbeatTimer) {\n      clearInterval(this.heartbeatTimer);\n      this.heartbeatTimer = null;\n    }\n  }\n\n  private clearReconnect(): void {\n    if (this.reconnectTimer) {\n      clearTimeout(this.reconnectTimer);\n      this.reconnectTimer = null;\n    }\n  }\n\n  private async buildUrl(): Promise<string> {\n    // Resolve token: callback > static > none\n    let token: string | undefined;\n    if (this.opts.auth) {\n      // If the auth callback throws, let it propagate — connect() catches and\n      // pauses reconnect until the user supplies fresh creds.\n      token = await this.opts.auth();\n      if (!token) {\n        // Configured auth callback returned no token. Treat as a fatal auth\n        // condition (don't silently connect without credentials).\n        throw new Error('SharedSocket: auth() returned no token');\n      }\n    } else if (this.opts.authToken) {\n      token = this.opts.authToken;\n    }\n\n    if (!token) return this.url;\n\n    // WebSocket URLs (ws://, wss://) are not fully supported by URL API.\n    // Convert to http(s) for parsing, then back to ws(s).\n    const httpUrl = this.url.replace(/^ws(s?):\\/\\//, 'http$1://');\n    const parsed = new URL(httpUrl);\n    parsed.searchParams.set(this.opts.authParam, token);\n\n    return parsed.toString().replace(/^http(s?):\\/\\//, 'ws$1://');\n  }\n\n  private setState(state: SocketState): void {\n    this._state = state;\n    for (const fn of this.onStateChangeFns) fn(state);\n  }\n\n  [Symbol.dispose](): void {\n    this.disconnect();\n    this.onMessageFns.clear();\n    this.onStateChangeFns.clear();\n    this.buffer = [];\n  }\n}\n","import './utils/disposable';\nimport type { SocketState, Unsubscribe, EventHandler } from './types';\n\n/**\n * WorkerSocket — WebSocket running inside a Web Worker.\n *\n * Same interface as SharedSocket, but WebSocket lives off main thread.\n * Benefits: heartbeat timers and JSON parsing don't block UI rendering.\n *\n * Use when:\n * - High message rate (50+ msgs/sec)\n * - Heavy JSON payloads\n * - UI does complex rendering that could block main thread\n *\n * Don't use when:\n * - Low message rate (simple chat, notifications)\n * - Bundle size matters (adds worker file)\n * - Debugging (Worker DevTools is less convenient)\n */\nexport class WorkerSocket implements Disposable {\n  private worker: Worker | null = null;\n  private _state: SocketState = 'closed';\n\n  private onMessageFns = new Set<EventHandler>();\n  private onStateChangeFns = new Set<(state: SocketState) => void>();\n\n  constructor(\n    private url: string,\n    private options: {\n      protocols?: string[];\n      reconnect?: boolean;\n      reconnectMaxDelay?: number;\n      reconnectMaxRetries?: number;\n      authFailureCloseCodes?: number[];\n      heartbeatInterval?: number;\n      heartbeatTimeout?: number;\n      sendBuffer?: number;\n      workerUrl?: string | URL;\n      auth?: () => string | Promise<string>;\n      authToken?: string;\n      authParam?: string;\n      pingPayload?: unknown;\n    } = {},\n  ) {}\n\n  get state(): SocketState {\n    return this._state;\n  }\n\n  private setState(s: SocketState): void {\n    this._state = s;\n    for (const fn of this.onStateChangeFns) fn(s);\n  }\n\n  async connect(): Promise<void> {\n    // Resolve auth token before sending to worker (functions can't cross worker boundary)\n    let authToken: string | undefined;\n    if (this.options.auth) {\n      try {\n        authToken = await this.options.auth();\n      } catch {\n        // auth() threw — pause reconnect until user provides fresh creds.\n        this.setState('failed');\n        return;\n      }\n      if (!authToken) {\n        // Configured auth callback returned nothing — same fail-closed behavior.\n        this.setState('failed');\n        return;\n      }\n    } else if (this.options.authToken) {\n      authToken = this.options.authToken;\n    }\n\n    // Build URL with auth token\n    const connectUrl = authToken ? this.buildUrl(authToken) : this.url;\n\n    // Create worker from inline blob if no workerUrl provided\n    const workerUrl = this.options.workerUrl ?? this.createWorkerBlob();\n\n    this.worker = new Worker(workerUrl, { type: 'module' });\n\n    this.worker.onmessage = (ev: MessageEvent) => {\n      const msg = ev.data;\n\n      switch (msg.type) {\n        case 'state':\n          this._state = msg.state;\n          for (const fn of this.onStateChangeFns) fn(msg.state);\n          break;\n\n        case 'message':\n          for (const fn of this.onMessageFns) fn(msg.data);\n          break;\n\n        case 'open':\n          // State already set via 'state' message\n          break;\n\n        case 'close':\n          break;\n\n        case 'error':\n          console.error('WorkerSocket error:', msg.message);\n          break;\n      }\n    };\n\n    this.worker.postMessage({\n      type: 'connect',\n      url: connectUrl,\n      protocols: this.options.protocols ?? [],\n      reconnect: this.options.reconnect ?? true,\n      reconnectMaxDelay: this.options.reconnectMaxDelay ?? 30_000,\n      reconnectMaxRetries: this.options.reconnectMaxRetries ?? Infinity,\n      authFailureCloseCodes: this.options.authFailureCloseCodes ?? [1008],\n      heartbeatInterval: this.options.heartbeatInterval ?? 30_000,\n      heartbeatTimeout: this.options.heartbeatTimeout ?? 0,\n      bufferSize: this.options.sendBuffer ?? 100,\n      pingPayload: this.options.pingPayload,\n    });\n  }\n\n  private buildUrl(token: string): string {\n    const param = this.options.authParam ?? 'token';\n    const httpUrl = this.url.replace(/^ws(s?):\\/\\//, 'http$1://');\n    const parsed = new URL(httpUrl);\n    parsed.searchParams.set(param, token);\n    return parsed.toString().replace(/^http(s?):\\/\\//, 'ws$1://');\n  }\n\n  send(data: unknown): void {\n    this.worker?.postMessage({ type: 'send', data });\n  }\n\n  /** Manually trigger reconnect: resets retry counter, attempts a fresh connection. */\n  reconnect(): void {\n    this.worker?.postMessage({ type: 'reconnect' });\n  }\n\n  disconnect(): void {\n    this.worker?.postMessage({ type: 'disconnect' });\n    setTimeout(() => {\n      this.worker?.terminate();\n      this.worker = null;\n    }, 100);\n    this._state = 'closed';\n  }\n\n  onMessage(fn: EventHandler): Unsubscribe {\n    this.onMessageFns.add(fn);\n    return () => this.onMessageFns.delete(fn);\n  }\n\n  onStateChange(fn: (state: SocketState) => void): Unsubscribe {\n    this.onStateChangeFns.add(fn);\n    return () => this.onStateChangeFns.delete(fn);\n  }\n\n  private createWorkerBlob(): URL {\n    // Inline the worker code as a blob URL\n    // In production, use a bundler (Vite, webpack) to handle worker imports\n    const code = `\n      let ws = null, state = 'closed', buffer = [], disposed = false;\n      let heartbeatTimer = null, reconnectTimer = null;\n      let url = '', protocols = [], shouldReconnect = true;\n      let maxDelay = 30000, maxRetries = Infinity, hbInterval = 30000, hbTimeout = 0, lastInbound = 0, maxBuf = 100;\n      let authFailCodes = new Set([1008]);\n      let delay = 1000, attempts = 0, pingPayload = '{\"type\":\"ping\"}';\n\n      function setState(s) { state = s; self.postMessage({ type: 'state', state: s }); }\n      function connect() {\n        if (disposed) return;\n        setState('connecting');\n        ws = new WebSocket(url, protocols);\n        ws.onopen = () => { attempts = 0; delay = 1000; lastInbound = Date.now(); setState('connected'); self.postMessage({ type: 'open' }); flush(); startHB(); };\n        ws.onmessage = (e) => { lastInbound = Date.now(); let d; try { d = JSON.parse(e.data); } catch { d = e.data; } self.postMessage({ type: 'message', data: d }); };\n        ws.onclose = (e) => {\n          stopHB();\n          self.postMessage({ type: 'close', code: e.code, reason: e.reason });\n          if (authFailCodes.has(e.code)) { setState('failed'); return; }\n          if (!disposed && shouldReconnect && e.code !== 1000) reconnect(); else setState('closed');\n        };\n        ws.onerror = () => { self.postMessage({ type: 'error', message: 'error' }); };\n      }\n      function send(d) { if (state === 'connected' && ws?.readyState === 1) ws.send(JSON.stringify(d)); else if (buffer.length < maxBuf) buffer.push(d); }\n      function flush() { const p = buffer.splice(0); p.forEach(send); }\n      function startHB() { stopHB(); heartbeatTimer = setInterval(() => { if (hbTimeout > 0 && state === 'connected' && Date.now() - lastInbound > hbTimeout) { manualReconnect(); return; } if (ws?.readyState === 1) ws.send(pingPayload); }, hbInterval); }\n      function stopHB() { if (heartbeatTimer) { clearInterval(heartbeatTimer); heartbeatTimer = null; } }\n      function reconnect() {\n        attempts++;\n        if (attempts > maxRetries) { setState('failed'); return; }\n        setState('reconnecting');\n        const j = delay * 0.25 * (Math.random() * 2 - 1);\n        reconnectTimer = setTimeout(() => { if (!disposed) connect(); }, Math.min(delay + j, maxDelay));\n        delay = Math.min(delay * 2, maxDelay);\n      }\n      function manualReconnect() {\n        if (disposed) return;\n        if (reconnectTimer) { clearTimeout(reconnectTimer); reconnectTimer = null; }\n        attempts = 0; delay = 1000;\n        if (ws) { ws.onclose = null; ws.onmessage = null; ws.onerror = null; if (ws.readyState < 2) ws.close(1000, 'manual reconnect'); ws = null; }\n        connect();\n      }\n      self.onmessage = (e) => {\n        const c = e.data;\n        if (c.type === 'connect') { url = c.url; protocols = c.protocols || []; shouldReconnect = c.reconnect ?? true; maxDelay = c.reconnectMaxDelay || 30000; maxRetries = c.reconnectMaxRetries ?? Infinity; if (c.authFailureCloseCodes) authFailCodes = new Set(c.authFailureCloseCodes); hbInterval = c.heartbeatInterval || 30000; hbTimeout = c.heartbeatTimeout || 0; maxBuf = c.bufferSize || 100; if (c.pingPayload) pingPayload = JSON.stringify(c.pingPayload); connect(); }\n        if (c.type === 'send') send(c.data);\n        if (c.type === 'reconnect') manualReconnect();\n        if (c.type === 'disconnect') { disposed = true; stopHB(); if (reconnectTimer) clearTimeout(reconnectTimer); if (ws) { ws.onclose = null; if (ws.readyState < 2) ws.close(1000); ws = null; } buffer = []; setState('closed'); }\n      };\n    `;\n    const blob = new Blob([code], { type: 'application/javascript' });\n    return new URL(URL.createObjectURL(blob));\n  }\n\n  [Symbol.dispose](): void {\n    this.disconnect();\n    this.onMessageFns.clear();\n    this.onStateChangeFns.clear();\n  }\n}\n","import './utils/disposable';\nimport type { EventHandler, Unsubscribe } from './types';\n\nexport class SubscriptionManager implements Disposable {\n  private handlers = new Map<string, Set<EventHandler>>();\n  private lastMessages = new Map<string, unknown>();\n\n  on(event: string, handler: EventHandler): Unsubscribe {\n    let set = this.handlers.get(event);\n    if (!set) {\n      set = new Set();\n      this.handlers.set(event, set);\n    }\n    set.add(handler);\n    return () => set!.delete(handler);\n  }\n\n  once(event: string, handler: EventHandler): Unsubscribe {\n    const wrapper: EventHandler = (data, raw) => {\n      unsub();\n      handler(data, raw);\n    };\n    const unsub = this.on(event, wrapper);\n    return unsub;\n  }\n\n  off(event: string, handler?: EventHandler): void {\n    if (handler) {\n      this.handlers.get(event)?.delete(handler);\n    } else {\n      this.handlers.delete(event);\n    }\n  }\n\n  emit(event: string, data: unknown, raw?: unknown): void {\n    this.lastMessages.set(event, data);\n    const set = this.handlers.get(event);\n    if (set) {\n      for (const fn of set) fn(data, raw);\n    }\n  }\n\n  getLastMessage(event: string): unknown | undefined {\n    return this.lastMessages.get(event);\n  }\n\n  async *stream(event: string, signal?: AbortSignal): AsyncGenerator<unknown> {\n    const queue: unknown[] = [];\n    let resolve: (() => void) | null = null;\n    let done = false;\n\n    const unsub = this.on(event, (data) => {\n      queue.push(data);\n      resolve?.();\n    });\n\n    const onAbort = () => {\n      done = true;\n      resolve?.();\n    };\n    signal?.addEventListener('abort', onAbort);\n\n    try {\n      while (!done) {\n        if (queue.length > 0) {\n          yield queue.shift()!;\n        } else {\n          await new Promise<void>((r) => { resolve = r; });\n          resolve = null;\n        }\n      }\n    } finally {\n      unsub();\n      signal?.removeEventListener('abort', onAbort);\n    }\n  }\n\n  offAll(): void {\n    this.handlers.clear();\n    this.lastMessages.clear();\n  }\n\n  [Symbol.dispose](): void {\n    this.offAll();\n  }\n}\n","import type { EventProtocol, FrameKind, FramePayload, Logger, Middleware } from './types';\n\n/** The minimum a frame sink must expose for the pipeline to write to it. */\ninterface FrameSink {\n  send(data: unknown): void;\n}\n\n/**\n * Outgoing frame pipeline — the single place a structured `(kind, payload)` is\n * turned into a wire frame and written to the socket:\n *\n *   buildFrame (custom frameBuilder → default)  →  outgoing middleware  →  send\n *\n * Extracted from SharedWebSocket so the build/middleware/send logic lives in\n * one cohesive unit. The owner keeps the socket; it's handed in via\n * `setSocket()` on each leader handover (and cleared on demotion).\n */\nexport class FramePipeline {\n  private socket: FrameSink | null = null;\n  private readonly middleware: Middleware[] = [];\n\n  constructor(\n    private readonly proto: EventProtocol,\n    private readonly log: Logger,\n  ) {}\n\n  /** Point the pipeline at the current leader socket (or `null` on demotion). */\n  setSocket(socket: FrameSink | null): void {\n    this.socket = socket;\n  }\n\n  hasSocket(): boolean {\n    return this.socket !== null;\n  }\n\n  /** Register an outgoing middleware. Return `null` from it to drop a frame. */\n  use(fn: Middleware): void {\n    this.middleware.push(fn);\n  }\n\n  /** Build, run middleware, and write to the socket. No-op without a socket. */\n  transmit(kind: FrameKind, payload: FramePayload): void {\n    if (!this.socket) return;\n    let frame: unknown = this.buildFrame(kind, payload);\n    if (frame === null) {\n      this.log.debug('[SharedWS] ✗ frameBuilder dropped frame', kind, this.frameLabel(kind, payload));\n      return;\n    }\n    for (const mw of this.middleware) {\n      frame = mw(frame);\n      if (frame === null) {\n        this.log.debug('[SharedWS] ✗ outgoing dropped by middleware', kind, this.frameLabel(kind, payload));\n        return;\n      }\n    }\n    // Auth frames carry a token — never log payload or wire frame.\n    if (kind === 'auth-login') {\n      this.log.debug('[SharedWS] → send', kind, '(token redacted)');\n    } else {\n      this.log.debug('[SharedWS] → send', kind, this.frameLabel(kind, payload), { payload, frame });\n    }\n    this.socket.send(frame);\n  }\n\n  /**\n   * Build the wire frame for a given kind. Honors custom `frameBuilder`.\n   * Return-value contract:\n   *   - any concrete value → use as the frame\n   *   - `null`             → drop the frame (intentional filter)\n   *   - `undefined`        → fall back to the default builder for this kind\n   */\n  private buildFrame(kind: FrameKind, payload: FramePayload): unknown {\n    if (this.proto.frameBuilder) {\n      const result = this.proto.frameBuilder(kind, payload);\n      if (result !== undefined) return result;\n      // undefined → fall through to default for this kind\n    }\n    return this.defaultFrameBuilder(kind, payload);\n  }\n\n  /** Legacy two-key builder — preserved as the default for back-compat. */\n  private defaultFrameBuilder(kind: FrameKind, p: FramePayload): unknown {\n    let eventName: string;\n    let dataPart: unknown;\n\n    switch (kind) {\n      case 'event':\n        // Channel-scoped events join with `:` for wire compat (Pusher convention).\n        eventName = p.channel ? `${p.channel}:${p.event ?? ''}` : (p.event ?? this.proto.defaultEvent);\n        dataPart = p.data;\n        break;\n      case 'subscribe':\n        eventName = this.proto.channelJoin;\n        dataPart = { channel: p.channel };\n        break;\n      case 'unsubscribe':\n        eventName = this.proto.channelLeave;\n        dataPart = { channel: p.channel };\n        break;\n      case 'topic-subscribe':\n        eventName = this.proto.topicSubscribe;\n        dataPart = { topic: p.topic };\n        break;\n      case 'topic-unsubscribe':\n        eventName = this.proto.topicUnsubscribe;\n        dataPart = { topic: p.topic };\n        break;\n      case 'auth-login':\n        eventName = this.proto.authLogin;\n        dataPart = { token: p.data };\n        break;\n      case 'auth-logout':\n        eventName = this.proto.authLogout;\n        dataPart = {};\n        break;\n    }\n\n    return {\n      ...(p.extras ?? {}),\n      [this.proto.eventField]: eventName,\n      [this.proto.dataField]: dataPart,\n    };\n  }\n\n  /**\n   * Human-readable headline for log lines — picks the most relevant field\n   * out of the structured payload so log scanners aren't reading objects.\n   */\n  private frameLabel(kind: FrameKind, p: FramePayload): string {\n    switch (kind) {\n      case 'event':              return p.event ?? '?';\n      case 'subscribe':\n      case 'unsubscribe':        return p.channel ?? '?';\n      case 'topic-subscribe':\n      case 'topic-unsubscribe':  return p.topic ?? '?';\n      case 'auth-login':         return '(redacted)';\n      case 'auth-logout':        return '';\n    }\n  }\n}\n","import type { EventProtocol, Logger, Middleware } from './types';\n\n/** Envelope produced from a raw incoming frame, ready to fan out to tabs. */\nexport interface IncomingEnvelope {\n  event: string;\n  data: unknown;\n  /** Full deserialized frame (post-middleware) — for handlers needing top-level fields. */\n  raw: unknown;\n}\n\n/**\n * Incoming frame pipeline — the mirror of FramePipeline. Turns a raw,\n * already-deserialized socket frame into an `{ event, data, raw }` envelope:\n *\n *   incoming middleware  →  extract event/data  →  per-event deserializer\n *\n * Runs on the leader (the tab that owns the socket); the envelope is then\n * broadcast to every tab for local fan-out. Extracted from SharedWebSocket so\n * the receive-side transform lives in one cohesive unit.\n */\nexport class IncomingPipeline {\n  private readonly middleware: Middleware[] = [];\n  private readonly deserializers = new Map<string, (data: unknown) => unknown>();\n\n  constructor(\n    private readonly proto: EventProtocol,\n    private readonly log: Logger,\n  ) {}\n\n  /** Register an incoming middleware. Return `null` from it to drop a frame. */\n  use(fn: Middleware): void {\n    this.middleware.push(fn);\n  }\n\n  /** Register a per-event deserializer, applied after global deserialize. */\n  deserializer(event: string, fn: (data: unknown) => unknown): void {\n    this.deserializers.set(event, fn);\n  }\n\n  /**\n   * Transform a raw frame into an envelope, or `null` if middleware dropped it.\n   */\n  process(raw: unknown): IncomingEnvelope | null {\n    let data: unknown = raw;\n    for (const mw of this.middleware) {\n      data = mw(data);\n      if (data === null) {\n        this.log.debug('[SharedWS] ✗ incoming dropped by middleware', { raw });\n        return null;\n      }\n    }\n\n    const msg = data as Record<string, unknown> | null | undefined;\n    const event = (msg?.[this.proto.eventField] as string) ?? this.proto.defaultEvent;\n    let payload = msg?.[this.proto.dataField] ?? data;\n\n    // Per-event deserializer transforms data after global deserialize.\n    const eventDeserializer = this.deserializers.get(event);\n    if (eventDeserializer) {\n      payload = eventDeserializer(payload);\n    }\n\n    this.log.debug('[SharedWS] ← recv', event, { data: payload, raw: data });\n    return { event, data: payload, raw: data };\n  }\n}\n","import './utils/disposable';\nimport { generateId } from './utils/id';\nimport { MessageBus } from './MessageBus';\nimport { BUS, busPendingReply, WS_DEFAULTS } from './constants';\nimport type { FrameKind, FramePayload, Logger, Unsubscribe } from './types';\n\ninterface OutboxEntry {\n  id: string;\n  kind: FrameKind;\n  payload: FramePayload;\n  enqueuedAt: number;\n}\n\n/** Subset of an OutboxEntry shared across tabs (no local-only timestamp). */\ntype ReplayEntry = Pick<OutboxEntry, 'id' | 'kind' | 'payload'>;\n\n/**\n * At-least-once outbox for follower-originated dispatches.\n *\n * When a follower calls `send()`, the frame is published over the bus for the\n * leader to write. If the leader dies between receiving the dispatch and\n * writing it to the socket, the frame would be lost — so each `event` dispatch\n * is buffered locally with a unique id. The leader broadcasts `ws:dispatch-\n * flushed` once it processes the dispatch and the originator drops the entry.\n * On leader change, the new leader gathers still-pending entries from every\n * surviving tab and replays them over the fresh socket.\n *\n * Extracted from SharedWebSocket so the buffer + gather/replay protocol is one\n * cohesive unit. (Only `event` kinds are buffered — channel/topic/auth frames\n * are re-established separately by resubscribe-on-connect.)\n */\nexport class Outbox implements Disposable {\n  private readonly pending = new Map<string, OutboxEntry>();\n  private cleanups: Unsubscribe[] = [];\n\n  constructor(\n    private readonly bus: MessageBus,\n    private readonly maxSize: number,\n    /** Writes a frame to the live socket — supplied by the owner (FramePipeline). */\n    private readonly transmit: (kind: FrameKind, payload: FramePayload) => void,\n    private readonly log: Logger,\n  ) {\n    // Originator drops its entry once the leader confirms the dispatch (or, on\n    // leader change, the new leader confirms replay).\n    this.cleanups.push(\n      this.bus.subscribe<{ id: string }>(BUS.DISPATCH_FLUSHED, (msg) => {\n        this.pending.delete(msg.id);\n      }),\n    );\n\n    // A new leader asks every tab to announce its still-pending dispatches.\n    this.cleanups.push(\n      this.bus.subscribe<{ replyId: string }>(BUS.GATHER_PENDING, (req) => {\n        if (this.pending.size === 0) return;\n        this.bus.publish(busPendingReply(req.replyId), {\n          entries: [...this.pending.values()],\n        });\n      }),\n    );\n  }\n\n  get size(): number {\n    return this.pending.size;\n  }\n\n  /**\n   * Route a follower dispatch to the leader. Buffers `event` kinds locally for\n   * replay across handover, then publishes for the leader's socket. Channel /\n   * topic / auth frames are not buffered — they're re-sent by\n   * resubscribe-on-connect, so buffering them too would double-emit.\n   */\n  route(kind: FrameKind, payload: FramePayload): void {\n    const id = generateId();\n    if (kind === 'event') this.enqueue(id, kind, payload);\n    this.bus.publish(BUS.DISPATCH, { id, kind, payload });\n  }\n\n  private enqueue(id: string, kind: FrameKind, payload: FramePayload): void {\n    if (this.maxSize <= 0) return;\n    if (this.pending.size >= this.maxSize) {\n      // Drop oldest — Map iteration order = insertion order.\n      const oldestKey = this.pending.keys().next().value;\n      if (oldestKey !== undefined) this.pending.delete(oldestKey);\n    }\n    this.pending.set(id, { id, kind, payload, enqueuedAt: Date.now() });\n  }\n\n  /**\n   * New-leader replay: gather still-pending dispatches from all tabs (including\n   * this one), transmit each over the fresh socket, then signal every\n   * originator to drop its entry. `isValid` lets the caller bail if the socket\n   * was replaced again while we were gathering (avoids replaying onto a socket\n   * that's already gone, which would drop entries that never actually sent).\n   */\n  async replay(isValid: () => boolean = () => true): Promise<void> {\n    const entries = await this.gather();\n    if (!isValid()) return;\n    if (entries.length === 0) return;\n\n    let sent = 0;\n    for (const e of entries) {\n      this.transmit(e.kind, e.payload);\n      // Remove from own pending (publish doesn't echo to self) and tell any\n      // other tab that originated the same id to drop it as well.\n      this.pending.delete(e.id);\n      this.bus.publish(BUS.DISPATCH_FLUSHED, { id: e.id });\n      sent++;\n    }\n    this.log.info('[SharedWS] replayed pending dispatches', { count: sent });\n  }\n\n  /**\n   * Cross-tab pending-dispatch gather. Broadcasts a one-shot request, collects\n   * for a short window, dedups by id (so multiple tabs holding the same id\n   * don't double-replay).\n   */\n  private gather(timeoutMs: number = WS_DEFAULTS.GATHER_PENDING_TIMEOUT): Promise<ReplayEntry[]> {\n    const seen = new Map<string, ReplayEntry>();\n    for (const e of this.pending.values()) {\n      seen.set(e.id, { id: e.id, kind: e.kind, payload: e.payload });\n    }\n    const replyId = generateId();\n\n    return new Promise((resolve) => {\n      const unsub = this.bus.subscribe<{ entries: OutboxEntry[] }>(\n        busPendingReply(replyId),\n        (msg) => {\n          for (const e of msg.entries) {\n            if (!seen.has(e.id)) seen.set(e.id, { id: e.id, kind: e.kind, payload: e.payload });\n          }\n        },\n      );\n      this.bus.publish(BUS.GATHER_PENDING, { replyId });\n      setTimeout(() => {\n        unsub();\n        resolve([...seen.values()]);\n      }, timeoutMs);\n    });\n  }\n\n  [Symbol.dispose](): void {\n    for (const unsub of this.cleanups) unsub();\n    this.cleanups = [];\n    this.pending.clear();\n  }\n}\n","import './utils/disposable';\nimport { generateId } from './utils/id';\nimport { MessageBus } from './MessageBus';\nimport { BUS, busSubsReply, WS_DEFAULTS } from './constants';\nimport type { FrameKind, FramePayload, Logger, Unsubscribe } from './types';\n\n/**\n * Tracks this tab's channel and topic subscriptions and replays them onto a\n * freshly connected leader socket.\n *\n * Two roles:\n *  - **Bookkeeping** — a refcount of channel subscriptions (so N `channel()`\n *    handles for the same name share one server-side join) and the set of\n *    subscribed topics. `channelNames()` feeds the incoming-event prefix\n *    routing in SharedWebSocket.\n *  - **Replay** — on leader handover/reconnect, gather the *union* of\n *    channels/topics across every surviving tab and re-send the join /\n *    topic-subscribe frames, so a promoted follower doesn't silently drop\n *    subscriptions any tab still cares about. It also answers other tabs'\n *    gather requests.\n *\n * Extracted from SharedWebSocket so this bookkeeping + cross-tab replay is one\n * cohesive unit. Auth-specific subscription sets (auto-leave on deauth) stay in\n * SharedWebSocket — this registry holds the full set used for routing/replay.\n */\nexport class SubscriptionRegistry implements Disposable {\n  /** Refcount of active channel subscriptions per name. */\n  private readonly channelRefs = new Map<string, number>();\n  /** All topic subscriptions (auth and non-auth). */\n  private readonly topics = new Set<string>();\n  private cleanups: Unsubscribe[] = [];\n\n  constructor(\n    private readonly bus: MessageBus,\n    /** Writes a frame to the live socket — supplied by the owner (FramePipeline). */\n    private readonly transmit: (kind: FrameKind, payload: FramePayload) => void,\n    private readonly log: Logger,\n  ) {\n    // Announce this tab's channels/topics when a new leader gathers.\n    this.cleanups.push(\n      this.bus.subscribe<{ replyId: string }>(BUS.GATHER_SUBS, (req) => {\n        this.bus.publish(busSubsReply(req.replyId), {\n          channels: [...this.channelRefs.keys()],\n          topics: [...this.topics],\n        });\n      }),\n    );\n  }\n\n  /** Track a channel subscription (refcounted across multiple handles). */\n  addChannel(name: string): void {\n    this.channelRefs.set(name, (this.channelRefs.get(name) ?? 0) + 1);\n  }\n\n  /** Drop one channel reference; forgets the channel at zero. */\n  removeChannel(name: string): void {\n    const next = (this.channelRefs.get(name) ?? 1) - 1;\n    if (next <= 0) this.channelRefs.delete(name);\n    else this.channelRefs.set(name, next);\n  }\n\n  addTopic(topic: string): void {\n    this.topics.add(topic);\n  }\n\n  removeTopic(topic: string): void {\n    this.topics.delete(topic);\n  }\n\n  /** Channel names currently held by this tab (for incoming-event routing). */\n  channelNames(): IterableIterator<string> {\n    return this.channelRefs.keys();\n  }\n\n  /**\n   * Re-establish subscriptions on a freshly connected leader socket: gather the\n   * union of channels/topics across all surviving tabs, then transmit a join /\n   * topic-subscribe for each. `isValid` aborts if the socket was replaced while\n   * we were gathering.\n   */\n  async replay(isValid: () => boolean = () => true): Promise<void> {\n    const { channels, topics } = await this.gather();\n    if (!isValid()) return;\n\n    for (const name of channels) {\n      this.transmit('subscribe', { channel: name });\n    }\n    for (const topic of topics) {\n      this.transmit('topic-subscribe', { topic });\n    }\n\n    if (channels.length || topics.length) {\n      this.log.info('[SharedWS] replayed subscriptions', {\n        channels: channels.length,\n        topics: topics.length,\n      });\n    }\n  }\n\n  /**\n   * Best-effort cross-tab gather. Broadcasts a request and collects responses\n   * for a short window. Times out gracefully — late responses are dropped. Own\n   * subs are seeded so we don't rely on BroadcastChannel echo to self.\n   */\n  private gather(timeoutMs: number = WS_DEFAULTS.GATHER_SUBS_TIMEOUT): Promise<{ channels: string[]; topics: string[] }> {\n    const channels = new Set<string>(this.channelRefs.keys());\n    const topics = new Set<string>(this.topics);\n    const replyId = generateId();\n\n    return new Promise((resolve) => {\n      const unsub = this.bus.subscribe<{ channels: string[]; topics: string[] }>(\n        busSubsReply(replyId),\n        (msg) => {\n          for (const c of msg.channels) channels.add(c);\n          for (const t of msg.topics) topics.add(t);\n        },\n      );\n\n      this.bus.publish(BUS.GATHER_SUBS, { replyId });\n\n      setTimeout(() => {\n        unsub();\n        resolve({ channels: [...channels], topics: [...topics] });\n      }, timeoutMs);\n    });\n  }\n\n  [Symbol.dispose](): void {\n    for (const unsub of this.cleanups) unsub();\n    this.cleanups = [];\n    this.channelRefs.clear();\n    this.topics.clear();\n  }\n}\n","import './utils/disposable';\nimport { MessageBus } from './MessageBus';\nimport { SubscriptionManager } from './SubscriptionManager';\nimport { BUS, LIFECYCLE, AUTH_TOKEN_KEY } from './constants';\nimport type { Channel, EventHandler, FrameKind, FramePayload, Logger, Unsubscribe } from './types';\n\n/** Everything AuthManager needs from its owner, injected to keep it decoupled. */\nexport interface AuthManagerDeps {\n  bus: MessageBus;\n  subs: SubscriptionManager;\n  /** Shared cross-tab store; the auth token lives under `AUTH_TOKEN_KEY`. */\n  syncStore: Map<string, unknown>;\n  isLeader: () => boolean;\n  /** Current leader socket state, or undefined when this tab holds no socket. */\n  socketState: () => string | undefined;\n  /** Force a reconnect of the leader socket. */\n  reconnect: () => void;\n  /** Route an outgoing frame (leader transmits; follower forwards via bus). */\n  dispatch: (kind: FrameKind, payload: FramePayload) => void;\n  /** Leave a topic subscription (SharedWebSocket.unsubscribe). */\n  unsubscribeTopic: (topic: string) => void;\n  /** Server event that signals auth was revoked (proto.authRevoked). */\n  authRevokedEvent: string;\n  /** Token provider for periodic refresh (options.refresh ?? options.auth). */\n  refresh?: () => string | Promise<string>;\n  /** Refresh interval in ms; disabled when unset or <= 0. */\n  refreshInterval?: number;\n  log: Logger;\n}\n\n/**\n * Owns runtime authentication: login/logout, cross-tab auth-state sync, the\n * leader-only token refresh timer, re-auth on reconnect, server-revocation\n * handling, and the auth-scoped channel/topic sets that auto-leave on logout.\n *\n * Extracted from SharedWebSocket so this cross-cutting concern is one unit.\n * It still leans on the owner for the things it can't own alone (the socket,\n * the dispatch pipeline, topic teardown) via the injected deps.\n */\nexport class AuthManager implements Disposable {\n  private _isAuthenticated = false;\n  /** Auth-scoped channels — auto-left on deauth/revocation. */\n  private readonly authChannels = new Map<string, Channel>();\n  /** Auth-scoped topics — auto-unsubscribed on deauth/revocation. */\n  private readonly authTopics = new Set<string>();\n  private refreshTimer: ReturnType<typeof setTimeout> | null = null;\n  /** Wall-clock time of the last token refresh — drives the catch-up check. */\n  private lastRefreshAt = 0;\n  /** Guards against overlapping refreshes (interval tick vs. catch-up). */\n  private refreshing = false;\n  /** True when the refresh loop is not running (not leader / stopped / disposed). */\n  private refreshStopped = true;\n  private cleanups: Unsubscribe[] = [];\n\n  constructor(private readonly deps: AuthManagerDeps) {\n    // Conditional resume — a follower that re-authenticates hints the leader to\n    // reconnect IFF its socket had given up (auth-failure close code), so a\n    // healthy connection isn't disrupted.\n    this.cleanups.push(\n      this.deps.bus.subscribe<void>(BUS.AUTH_RESUME, () => {\n        if (this.deps.isLeader() && this.deps.socketState() === 'failed') {\n          this.deps.log.info('[SharedWS] resume requested after auth — reconnecting failed socket');\n          this.deps.reconnect();\n        }\n      }),\n    );\n\n    // Server-initiated auth revocation — tear down auth-scoped subscriptions.\n    this.cleanups.push(\n      this.deps.subs.on(this.deps.authRevokedEvent, () => {\n        if (this.deps.isLeader()) {\n          for (const [, ch] of this.authChannels) ch.leave();\n          for (const topic of this.authTopics) this.deps.unsubscribeTopic(topic);\n        }\n        this.authChannels.clear();\n        this.authTopics.clear();\n        this._isAuthenticated = false;\n        this.deps.syncStore.delete(AUTH_TOKEN_KEY);\n        this.deps.subs.emit(LIFECYCLE.AUTH, false);\n        this.deps.log.warn('[SharedWS] auth revoked by server');\n      }),\n    );\n  }\n\n  get isAuthenticated(): boolean {\n    return this._isAuthenticated;\n  }\n\n  /** Track an auth-scoped channel so it auto-leaves on deauth/revocation. */\n  registerAuthChannel(name: string, ch: Channel): void {\n    this.authChannels.set(name, ch);\n  }\n\n  unregisterAuthChannel(name: string): void {\n    this.authChannels.delete(name);\n  }\n\n  registerAuthTopic(topic: string): void {\n    this.authTopics.add(topic);\n  }\n\n  unregisterAuthTopic(topic: string): void {\n    this.authTopics.delete(topic);\n  }\n\n  onAuthChange(fn: (authenticated: boolean) => void): Unsubscribe {\n    return this.deps.subs.on(LIFECYCLE.AUTH, fn as EventHandler);\n  }\n\n  /**\n   * Authenticate on the existing connection: sync the token to all tabs and\n   * send the auth-login frame. If the leader socket had failed (e.g. expired\n   * creds), the fresh token restarts it.\n   */\n  authenticate(token: string): void {\n    this._isAuthenticated = true;\n    this.lastRefreshAt = Date.now(); // a fresh token resets refresh staleness\n    this.deps.syncStore.set(AUTH_TOKEN_KEY, token);\n    this.deps.bus.broadcast(BUS.SYNC, { key: AUTH_TOKEN_KEY, value: token });\n    this.deps.bus.broadcast(BUS.LIFECYCLE, { type: 'auth', authenticated: true });\n    this.deps.log.info('[SharedWS] authenticated');\n\n    // If the leader's socket gave up, the new creds should restart it.\n    // reauthenticate() resends auth-login from syncStore once reconnected.\n    if (this.deps.isLeader() && this.deps.socketState() === 'failed') {\n      this.deps.reconnect();\n      return;\n    }\n\n    if (!this.deps.isLeader()) {\n      // Followers can't see leader state — hint to reconnect IFF failed.\n      this.deps.bus.publish(BUS.AUTH_RESUME, undefined);\n    }\n\n    this.deps.dispatch('auth-login', { data: token });\n  }\n\n  /**\n   * Deauthenticate: auto-leave auth channels/topics, send auth-logout, and\n   * sync the cleared state across tabs. The connection stays open for public\n   * events.\n   */\n  deauthenticate(): void {\n    for (const [, ch] of this.authChannels) ch.leave();\n    this.authChannels.clear();\n    for (const topic of this.authTopics) this.deps.unsubscribeTopic(topic);\n    this.authTopics.clear();\n\n    this._isAuthenticated = false;\n    this.deps.dispatch('auth-logout', {});\n    this.deps.syncStore.delete(AUTH_TOKEN_KEY);\n    this.deps.bus.broadcast(BUS.SYNC, { key: AUTH_TOKEN_KEY, value: undefined });\n    this.deps.bus.broadcast(BUS.LIFECYCLE, { type: 'auth', authenticated: false });\n    this.deps.log.info('[SharedWS] deauthenticated');\n  }\n\n  /**\n   * Apply an auth-state change broadcast over the bus (fired by every tab,\n   * including the originator via broadcast self-delivery): update local state,\n   * drop auth-scoped subscriptions on logout, and notify `onAuthChange`.\n   */\n  applyRemoteAuthState(authenticated: boolean | undefined): void {\n    this._isAuthenticated = !!authenticated;\n    if (!authenticated) {\n      this.authChannels.clear();\n      this.authTopics.clear();\n    }\n    this.deps.subs.emit(LIFECYCLE.AUTH, authenticated);\n  }\n\n  /** Re-send the auth-login frame from synced state after a fresh connect. */\n  reauthenticate(): void {\n    if (!this._isAuthenticated) return;\n    const token = this.deps.syncStore.get(AUTH_TOKEN_KEY) as string | undefined;\n    if (token) {\n      this.deps.dispatch('auth-login', { data: token });\n      this.deps.log.debug('[SharedWS] re-authenticated after reconnect');\n    }\n  }\n\n  /**\n   * Start the leader-only periodic token refresh. When the timer fires and the\n   * connection is authenticated, the new token flows back through\n   * `authenticate()` so subscribers stay synced and the socket re-issues\n   * auth-login. Idempotent.\n   */\n  startRefresh(): void {\n    if (!this.refreshStopped) return; // already running\n    if (!this.canRefresh() || !this.deps.isLeader()) return;\n    this.refreshStopped = false;\n    this.lastRefreshAt = Date.now();\n    this.scheduleRefresh();\n  }\n\n  stopRefresh(): void {\n    this.refreshStopped = true;\n    if (this.refreshTimer) {\n      clearTimeout(this.refreshTimer);\n      this.refreshTimer = null;\n    }\n  }\n\n  /**\n   * Catch up a refresh that a backgrounded leader missed. Browsers throttle (or\n   * freeze) timers in hidden tabs, so the periodic refresh can lapse and the\n   * token expire. Call this when the tab becomes visible again: if more than an\n   * interval has elapsed since the last refresh, refresh immediately. No-op on\n   * followers or when refresh isn't configured.\n   */\n  refreshIfStale(): void {\n    if (this.refreshStopped || this.refreshing) return;\n    if (!this.deps.isLeader() || !this._isAuthenticated) return;\n    const interval = this.deps.refreshInterval;\n    if (!this.deps.refresh || !interval || interval <= 0) return;\n    if (Date.now() - this.lastRefreshAt >= interval) {\n      void this.runRefresh();\n    }\n  }\n\n  private canRefresh(): boolean {\n    const interval = this.deps.refreshInterval;\n    return !!this.deps.refresh && !!interval && interval > 0;\n  }\n\n  /**\n   * Self-rescheduling tick (not setInterval): each run lines up the next from\n   * *now*, so a catch-up refresh on re-activation also resets the cadence\n   * instead of racing a still-pending interval.\n   */\n  private scheduleRefresh(): void {\n    if (this.refreshTimer) clearTimeout(this.refreshTimer);\n    this.refreshTimer = setTimeout(() => { void this.runRefresh(); }, this.deps.refreshInterval);\n  }\n\n  private async runRefresh(): Promise<void> {\n    if (this.refreshing) return;\n    if (this.deps.isLeader() && this._isAuthenticated && this.deps.refresh) {\n      this.refreshing = true;\n      try {\n        const token = await this.deps.refresh();\n        if (token) {\n          this.lastRefreshAt = Date.now();\n          this.authenticate(token);\n        } else {\n          this.deps.log.warn('[SharedWS] refresh() returned empty token — skipping');\n        }\n      } catch (err) {\n        this.deps.log.warn('[SharedWS] refresh() failed', err);\n      } finally {\n        this.refreshing = false;\n      }\n    }\n    // Reschedule unless we were stopped/disposed while awaiting.\n    if (!this.refreshStopped) this.scheduleRefresh();\n  }\n\n  [Symbol.dispose](): void {\n    this.stopRefresh();\n    for (const unsub of this.cleanups) unsub();\n    this.cleanups = [];\n    this.authChannels.clear();\n    this.authTopics.clear();\n  }\n}\n","import type { EventHandler, Logger, Unsubscribe } from './types';\n\n/** Configuration for `ws.push(event, config)`. */\nexport interface PushConfig<T = unknown> {\n  /** Custom render function — you decide how to display. */\n  render?: (data: T) => void;\n  /** Title for browser Notification API. */\n  title?: string | ((data: T) => string);\n  /** Body for browser Notification API. */\n  body?: string | ((data: T) => string);\n  /** Icon URL for browser Notification. */\n  icon?: string;\n  /** Tag for browser Notification deduplication. */\n  tag?: string | ((data: T) => string);\n  /**\n   * Which tab(s) show the notification:\n   * - `'active'` — only the visible/focused tab (default for render)\n   * - `'leader'` — only the leader tab (default for browser Notification)\n   * - `'all'` — every tab (critical alerts)\n   */\n  target?: 'active' | 'leader' | 'all';\n  /** Called when browser Notification is clicked. */\n  onClick?: (data: T) => void;\n}\n\n/** What PushManager needs from its owner. */\nexport interface PushManagerDeps {\n  /** Subscribe to an event (SharedWebSocket.on). */\n  on: (event: string, handler: EventHandler) => Unsubscribe;\n  isLeader: () => boolean;\n  isActive: () => boolean;\n  log: Logger;\n}\n\n/**\n * Routes incoming events to UI notifications — custom render and/or the browser\n * Notification API — with `target` deciding which tab(s) display them. Extracted\n * from SharedWebSocket so the render-vs-native dispatch lives in one unit.\n */\nexport class PushManager {\n  constructor(private readonly deps: PushManagerDeps) {}\n\n  push<T = unknown>(event: string, config: PushConfig<T>): Unsubscribe {\n    const useNativeNotification = !!config.title;\n\n    // Default target: 'active' for render, 'leader' for native.\n    const renderTarget = config.target ?? 'active';\n    const nativeTarget = config.target ?? 'leader';\n\n    if (useNativeNotification && typeof Notification !== 'undefined' && Notification.permission === 'default') {\n      Notification.requestPermission();\n    }\n\n    return this.deps.on(event, ((data: unknown) => {\n      const typed = data as T;\n      const isVisible = this.deps.isActive();\n      const isLeader = this.deps.isLeader();\n\n      // Custom render\n      if (config.render) {\n        const shouldRender =\n          renderTarget === 'all' ||\n          (renderTarget === 'active' && isVisible) ||\n          (renderTarget === 'leader' && isLeader);\n\n        if (shouldRender) {\n          config.render(typed);\n          this.deps.log.debug('[SharedWS] 🔔 render', event, `(target: ${renderTarget})`);\n        }\n      }\n\n      // Browser Notification API\n      if (useNativeNotification && typeof Notification !== 'undefined' && Notification.permission === 'granted') {\n        const shouldNotify =\n          nativeTarget === 'all' ||\n          (nativeTarget === 'leader' && isLeader) ||\n          (nativeTarget === 'active' && isVisible);\n\n        // Native notifications make sense when the tab is hidden.\n        if (shouldNotify && !isVisible) {\n          const title = typeof config.title === 'function' ? config.title(typed) : config.title!;\n          const body = typeof config.body === 'function' ? config.body(typed) : config.body;\n          const tag = typeof config.tag === 'function' ? config.tag(typed) : config.tag;\n\n          const notif = new Notification(title, { body, icon: config.icon, tag });\n\n          if (config.onClick) {\n            const handler = config.onClick;\n            notif.onclick = () => {\n              handler(typed);\n              window.focus();\n            };\n          }\n\n          this.deps.log.debug('[SharedWS] 🔔 native', title, `(target: ${nativeTarget})`);\n        }\n      }\n    }) as EventHandler);\n  }\n}\n","import './utils/disposable';\nimport { generateId } from './utils/id';\nimport { MessageBus } from './MessageBus';\nimport { TabCoordinator } from './TabCoordinator';\nimport { SharedSocket } from './SharedSocket';\nimport { WorkerSocket } from './WorkerSocket';\nimport { SubscriptionManager } from './SubscriptionManager';\nimport { FramePipeline } from './FramePipeline';\nimport { IncomingPipeline } from './IncomingPipeline';\nimport { Outbox } from './Outbox';\nimport { SubscriptionRegistry } from './SubscriptionRegistry';\nimport { AuthManager } from './AuthManager';\nimport { PushManager, type PushConfig } from './PushManager';\nimport { BUS, LIFECYCLE, WS_DEFAULTS, MESSAGE_BUS_CHANNEL } from './constants';\nimport type { SharedWebSocketOptions, TabRole, Unsubscribe, EventHandler, Channel, EventProtocol, EventMap, Logger, Middleware, FrameKind, FramePayload, ChannelAckResult } from './types';\n\nconst DEFAULT_PROTOCOL: EventProtocol = {\n  eventField: 'event',\n  dataField: 'data',\n  channelJoin: '$channel:join',\n  channelLeave: '$channel:leave',\n  ping: { type: 'ping' },\n  defaultEvent: 'message',\n  topicSubscribe: '$topic:subscribe',\n  topicUnsubscribe: '$topic:unsubscribe',\n  authLogin: '$auth:login',\n  authLogout: '$auth:logout',\n  authRevoked: '$auth:revoked',\n};\n\nconst NOOP_LOGGER: Logger = {\n  debug() {},\n  info() {},\n  warn() {},\n  error() {},\n};\n\n/**\n * Internal separator for channel-scoped subscription keys. ASCII RECORD\n * SEPARATOR (U+001E) — chosen because it cannot collide with characters\n * users put in channel or event names. Wire format keeps `:` for server\n * compatibility; this is storage-only.\n */\nconst CHANNEL_KEY_SEP = '\\u001e';\n\n/** Common interface for both SharedSocket and WorkerSocket. */\ninterface SocketAdapter {\n  readonly state: string;\n  connect(): void | Promise<void>;\n  send(data: unknown): void;\n  reconnect(): void;\n  disconnect(): void;\n  onMessage(fn: EventHandler): Unsubscribe;\n  onStateChange(fn: (state: string) => void): Unsubscribe;\n  [Symbol.dispose](): void;\n}\n\n/**\n * SharedWebSocket — shares ONE WebSocket connection across browser tabs.\n *\n * @typeParam TEvents - Event map for type-safe subscriptions.\n *\n * @example\n * // Typed events\n * type Events = {\n *   'chat.message': { text: string; userId: string };\n *   'order.created': { id: string; total: number };\n * };\n * const ws = new SharedWebSocket<Events>(url);\n * ws.on('chat.message', (msg) => msg.text); // ← msg: { text, userId }\n */\nexport class SharedWebSocket<TEvents extends EventMap = EventMap> implements Disposable {\n  private bus: MessageBus;\n  private coordinator: TabCoordinator;\n  private socket: SocketAdapter | null = null;\n  private subs = new SubscriptionManager();\n  private syncStore = new Map<string, unknown>();\n  private tabId: string;\n  private cleanups: Unsubscribe[] = [];\n  /** Removes all DOM (document/window) listeners at once on dispose. */\n  private readonly domListeners = new AbortController();\n  private disposed = false;\n  private readonly proto: EventProtocol;\n  private readonly log: Logger;\n  /** Outgoing frame building + middleware + socket write. */\n  private readonly framePipeline: FramePipeline;\n  /** Incoming frame transform: middleware + extract + per-event deserialize. */\n  private readonly incoming: IncomingPipeline;\n  /** At-least-once buffer + replay for follower-originated dispatches. */\n  private readonly outbox: Outbox;\n  private serializers = new Map<string, (data: unknown) => unknown>();\n  /**\n   * Channel/topic bookkeeping + cross-tab subscription replay. Holds the full\n   * set used for incoming-event routing and leader-handover replay; auth-scoped\n   * subscriptions are tracked separately by AuthManager for auto-leave on deauth.\n   */\n  private readonly subscriptions: SubscriptionRegistry;\n  /** Runtime auth: login/logout, token refresh, re-auth on connect, revocation. */\n  private readonly auth: AuthManager;\n  /** Routes events to render/native notifications by target tab. */\n  private readonly pushManager: PushManager;\n  /** Listeners for every raw incoming frame (post-deserialize, post-middleware). */\n  private rawFrameListeners = new Set<(raw: unknown) => void>();\n  /**\n   * Unsubscribe for the leader-only `ws:request` responder. Tracked separately\n   * from `cleanups` so it can be torn down on EACH leadership loss — otherwise\n   * a demoted tab keeps answering requests with a null socket, and every\n   * re-promotion stacks another responder (duplicate server sends).\n   */\n  private requestResponderCleanup: Unsubscribe | null = null;\n\n  constructor(\n    private readonly url: string,\n    private readonly options: SharedWebSocketOptions<TEvents> = {} as SharedWebSocketOptions<TEvents>,\n  ) {\n    this.proto = { ...DEFAULT_PROTOCOL, ...options.events };\n    this.log = options.debug ? (options.logger ?? console) : NOOP_LOGGER;\n    this.tabId = generateId();\n    this.log.debug('[SharedWS] init', { tabId: this.tabId, url });\n    this.bus = new MessageBus(MESSAGE_BUS_CHANNEL, this.tabId);\n    this.coordinator = new TabCoordinator(this.bus, this.tabId, {\n      electionTimeout: options.electionTimeout,\n      heartbeatInterval: options.leaderHeartbeat,\n      leaderTimeout: options.leaderTimeout,\n      leaderPingTimeout: options.leaderPingTimeout,\n    });\n\n    // Outgoing pipeline + outbox. The outbox replays via the pipeline's\n    // transmit, which writes to whatever socket the pipeline currently holds.\n    this.framePipeline = new FramePipeline(this.proto, this.log);\n    this.incoming = new IncomingPipeline(this.proto, this.log);\n    this.outbox = new Outbox(\n      this.bus,\n      options.outboundBufferSize ?? WS_DEFAULTS.OUTBOUND_BUFFER_SIZE,\n      (kind, payload) => this.framePipeline.transmit(kind, payload),\n      this.log,\n    );\n    this.subscriptions = new SubscriptionRegistry(\n      this.bus,\n      (kind, payload) => this.framePipeline.transmit(kind, payload),\n      this.log,\n    );\n    this.auth = new AuthManager({\n      bus: this.bus,\n      subs: this.subs,\n      syncStore: this.syncStore,\n      isLeader: () => this.coordinator.isLeader,\n      socketState: () => this.socket?.state,\n      reconnect: () => this.reconnect(),\n      dispatch: (kind, payload) => this.dispatch(kind, payload),\n      unsubscribeTopic: (topic) => this.unsubscribe(topic),\n      authRevokedEvent: this.proto.authRevoked,\n      refresh: this.options.refresh ?? this.options.auth,\n      refreshInterval: this.options.refreshTokenInterval,\n      log: this.log,\n    });\n    this.pushManager = new PushManager({\n      on: (event, handler) => this.on(event, handler),\n      isLeader: () => this.coordinator.isLeader,\n      isActive: () => this.isActive,\n      log: this.log,\n    });\n\n    // Let the coordinator decide leader health from real socket state — a\n    // backgrounded leader with a dead socket fails this check and is taken\n    // over by the next active tab instead of clinging to leadership.\n    this.coordinator.setHealthCheck(() => this.socket?.state === 'connected');\n\n    // If self-verification finds this tab is leader with a dead socket,\n    // reconnect rather than hand off — we already own the socket.\n    this.cleanups.push(\n      this.coordinator.onLeaderUnhealthy(() => {\n        this.log.warn('[SharedWS] leader socket unhealthy on activate — reconnecting');\n        this.socket?.reconnect();\n      }),\n    );\n\n    // When ANY tab receives a WS message via bus → emit to local subscribers\n    this.cleanups.push(\n      this.bus.subscribe<{ event: string; data: unknown; raw?: unknown }>(BUS.MESSAGE, (msg) => {\n        // Bare emit — fires any handler registered with the literal event name\n        this.subs.emit(msg.event, msg.data, msg.raw);\n\n        // Channel-scoped emit — for each registered channel whose name is a\n        // prefix of the incoming event (separated by ':'), also fire handlers\n        // stored under `${name}<RS>${rest}`. This lets `Channel.on('msg', h)`\n        // receive a wire event like 'chat:room:42:msg' without colon parsing.\n        for (const channelName of this.subscriptions.channelNames()) {\n          const prefix = channelName + ':';\n          if (msg.event.length > prefix.length && msg.event.startsWith(prefix)) {\n            const subEvent = msg.event.slice(prefix.length);\n            this.subs.emit(`${channelName}${CHANNEL_KEY_SEP}${subEvent}`, msg.data, msg.raw);\n          }\n        }\n\n        // Raw-frame fanout — pending Channel.ready ack matchers listen here.\n        if (this.rawFrameListeners.size > 0) {\n          for (const fn of this.rawFrameListeners) {\n            try { fn(msg.raw); } catch { /* matcher errors don't break dispatch */ }\n          }\n        }\n      }),\n    );\n\n    // Leader listens for dispatch requests from followers — re-enters\n    // transmit() so frameBuilder + outgoing middleware run on the tab that\n    // actually owns the socket.\n    this.cleanups.push(\n      this.bus.subscribe<{ kind: FrameKind; payload: FramePayload; id?: string }>(BUS.DISPATCH, (msg) => {\n        if (this.coordinator.isLeader && this.socket) {\n          this.framePipeline.transmit(msg.kind, msg.payload);\n          // Tell the originator to drop the entry from its pending buffer.\n          // Always flush — even when transmit was a no-op (middleware drop,\n          // frameBuilder returned null) — there's no point retrying a\n          // permanently-dropped frame.\n          if (msg.id) this.bus.publish(BUS.DISPATCH_FLUSHED, { id: msg.id });\n        }\n      }),\n    );\n    // The originator-side flush handler and pending-gather responder live in\n    // Outbox (it owns the buffer those operate on).\n\n    // Leader listens for reconnect requests from followers\n    this.cleanups.push(\n      this.bus.subscribe<void>(BUS.RECONNECT, () => {\n        if (this.coordinator.isLeader && this.socket) {\n          this.log.info('[SharedWS] manual reconnect requested by follower');\n          this.socket.reconnect();\n        }\n      }),\n    );\n\n    // (Auth resume + revocation live in AuthManager; channels/topics gather\n    // responder lives in SubscriptionRegistry.)\n\n    // Sync across tabs\n    this.cleanups.push(\n      this.bus.subscribe<{ key: string; value: unknown }>(BUS.SYNC, (msg) => {\n        this.syncStore.set(msg.key, msg.value);\n        this.subs.emit(`sync:${msg.key}`, msg.value);\n      }),\n    );\n\n    // Leader lifecycle\n    this.coordinator.onBecomeLeader(() => {\n      this.handleBecomeLeader();\n      this.bus.broadcast(BUS.LIFECYCLE, { type: 'leader', isLeader: true });\n    });\n    this.coordinator.onLoseLeadership(() => {\n      this.handleLoseLeadership();\n      this.bus.broadcast(BUS.LIFECYCLE, { type: 'leader', isLeader: false });\n    });\n\n    // Lifecycle events from bus (all tabs receive)\n    this.cleanups.push(\n      this.bus.subscribe<{ type: string; isLeader?: boolean; error?: unknown; authenticated?: boolean }>(BUS.LIFECYCLE, (msg) => {\n        switch (msg.type) {\n          case 'connect':\n            this.subs.emit(LIFECYCLE.CONNECT, undefined);\n            break;\n          case 'disconnect':\n            this.subs.emit(LIFECYCLE.DISCONNECT, undefined);\n            break;\n          case 'reconnecting':\n            this.subs.emit(LIFECYCLE.RECONNECTING, undefined);\n            break;\n          case 'reconnectFailed':\n            this.subs.emit(LIFECYCLE.RECONNECT_FAILED, undefined);\n            break;\n          case 'leader':\n            this.subs.emit(LIFECYCLE.LEADER, msg.isLeader);\n            break;\n          case 'error':\n            this.subs.emit(LIFECYCLE.ERROR, msg.error);\n            break;\n          case 'auth':\n            this.auth.applyRemoteAuthState(msg.authenticated);\n            break;\n        }\n      }),\n    );\n\n    // Track tab visibility\n    if (typeof document !== 'undefined') {\n      const onVisibilityChange = () => {\n        const active = !document.hidden;\n        this.subs.emit(LIFECYCLE.ACTIVE, active);\n        this.log.debug('[SharedWS]', active ? '👁 tab active' : '👁 tab hidden');\n        if (active) {\n          // Catch up a token refresh the throttled timer may have missed while\n          // backgrounded (leader-only; no-op otherwise).\n          this.auth.refreshIfStale();\n          // Make sure the leader (and its socket) survived the idle period;\n          // take over if it didn't. Opt-out via recoverOnActivate.\n          if (this.options.recoverOnActivate ?? true) {\n            void this.coordinator.verifyLeader();\n          }\n        }\n      };\n      document.addEventListener('visibilitychange', onVisibilityChange, { signal: this.domListeners.signal });\n    }\n\n    // Cleanup on tab close. Use `pagehide` rather than `beforeunload`:\n    // beforeunload doesn't fire on mobile Safari, on tab discard, or when the\n    // page enters the back/forward cache, so a leader could vanish without\n    // ever relinquishing the socket. pagehide covers all of those.\n    if (typeof window !== 'undefined') {\n      const onPageHide = (e: PageTransitionEvent) => {\n        if (e.persisted) {\n          // Entering bfcache — frozen but may be restored. Don't tear down;\n          // just relinquish leadership so a live tab takes over the socket.\n          // We re-join coordination on `pageshow` if restored.\n          this.coordinator.abdicate();\n        } else {\n          // Real unload — full teardown.\n          this[Symbol.dispose]();\n        }\n      };\n      const onPageShow = (e: PageTransitionEvent) => {\n        // Restored from bfcache — re-enter the election (becomes leader if no\n        // other tab claimed it, otherwise settles back to follower).\n        if (e.persisted && !this.disposed) {\n          void this.coordinator.elect();\n        }\n      };\n      const signal = this.domListeners.signal;\n      window.addEventListener('pagehide', onPageHide, { signal });\n      window.addEventListener('pageshow', onPageShow, { signal });\n    }\n  }\n\n  get connected(): boolean {\n    return this.socket?.state === 'connected' || !this.coordinator.isLeader;\n  }\n\n  get tabRole(): TabRole {\n    return this.coordinator.isLeader ? 'leader' : 'follower';\n  }\n\n  /** Whether the user is authenticated via runtime auth. */\n  get isAuthenticated(): boolean {\n    return this.auth.isAuthenticated;\n  }\n\n  /** Whether this tab is currently visible/focused. */\n  get isActive(): boolean {\n    return typeof document !== 'undefined' ? !document.hidden : true;\n  }\n\n  /** Start leader election and connect. */\n  async connect(): Promise<void> {\n    await this.coordinator.elect();\n  }\n\n  // ─── Lifecycle Hooks ─────────────────────────────────\n\n  /** Called when WebSocket connection opens (broadcast to all tabs). */\n  onConnect(fn: () => void): Unsubscribe {\n    return this.subs.on(LIFECYCLE.CONNECT, fn);\n  }\n\n  /** Called when WebSocket connection closes (broadcast to all tabs). */\n  onDisconnect(fn: () => void): Unsubscribe {\n    return this.subs.on(LIFECYCLE.DISCONNECT, fn);\n  }\n\n  /** Called when WebSocket starts reconnecting (broadcast to all tabs). */\n  onReconnecting(fn: () => void): Unsubscribe {\n    return this.subs.on(LIFECYCLE.RECONNECTING, fn);\n  }\n\n  /**\n   * Called when auto-reconnect gives up after exhausting `reconnectMaxRetries`.\n   * Use this to show a \"Reconnect\" UI affordance (snackbar, banner, modal)\n   * so the user can call `ws.reconnect()` to try again.\n   *\n   * @example\n   * ws.onReconnectFailed(() => {\n   *   showSnackbar('Connection lost', { action: { label: 'Reconnect', onClick: () => ws.reconnect() } });\n   * });\n   */\n  onReconnectFailed(fn: () => void): Unsubscribe {\n    return this.subs.on(LIFECYCLE.RECONNECT_FAILED, fn);\n  }\n\n  /**\n   * Manually trigger a reconnect. Resets the retry counter and attempts a\n   * fresh connection. Safe to call from any tab — the leader actually owns\n   * the socket, followers route the request via BroadcastChannel.\n   *\n   * Use after `onReconnectFailed` fires to let the user retry.\n   *\n   * @example\n   * snackbar.action('Reconnect', () => ws.reconnect());\n   */\n  reconnect(): void {\n    if (this.coordinator.isLeader && this.socket) {\n      this.socket.reconnect();\n    } else {\n      this.bus.publish(BUS.RECONNECT, undefined);\n    }\n  }\n\n  /** Called when this tab becomes leader or loses leadership. */\n  onLeaderChange(fn: (isLeader: boolean) => void): Unsubscribe {\n    return this.subs.on(LIFECYCLE.LEADER, fn as EventHandler);\n  }\n\n  /** Called on WebSocket or network error (broadcast to all tabs). */\n  onError(fn: (error: unknown) => void): Unsubscribe {\n    return this.subs.on(LIFECYCLE.ERROR, fn as EventHandler);\n  }\n\n  /** Called when this tab becomes visible/focused. */\n  onActive(fn: () => void): Unsubscribe {\n    return this.subs.on(LIFECYCLE.ACTIVE, ((isActive: unknown) => {\n      if (isActive === true) fn();\n    }) as EventHandler);\n  }\n\n  /** Called when this tab goes to background/hidden. */\n  onInactive(fn: () => void): Unsubscribe {\n    return this.subs.on(LIFECYCLE.ACTIVE, ((isActive: unknown) => {\n      if (isActive === false) fn();\n    }) as EventHandler);\n  }\n\n  /** Called on any visibility change. */\n  onVisibilityChange(fn: (isActive: boolean) => void): Unsubscribe {\n    return this.subs.on(LIFECYCLE.ACTIVE, fn as EventHandler);\n  }\n\n  // ─── Authentication ──────────────────────────────────\n\n  /**\n   * Authenticate on an existing connection. Sends auth event to server,\n   * syncs auth state across all tabs. Use for login after guest connection.\n   *\n   * @example\n   * const token = await loginApi(email, password);\n   * ws.authenticate(token);\n   *\n   * @example\n   * // React — via useSocketAuth hook\n   * const { authenticate } = useSocketAuth();\n   * authenticate(token);\n   */\n  authenticate(token: string): void {\n    this.auth.authenticate(token);\n  }\n\n  /**\n   * Deauthenticate — notifies server, auto-leaves all auth-required channels\n   * and topics, syncs state across tabs. Connection stays open for public events.\n   *\n   * @example\n   * ws.deauthenticate(); // connection stays open, auth subscriptions cleaned up\n   */\n  deauthenticate(): void {\n    this.auth.deauthenticate();\n  }\n\n  /**\n   * Called when auth state changes (authenticate, deauthenticate, or server revocation).\n   *\n   * @example\n   * ws.onAuthChange((authenticated) => {\n   *   if (!authenticated) router.push('/login');\n   * });\n   */\n  onAuthChange(fn: (authenticated: boolean) => void): Unsubscribe {\n    return this.auth.onAuthChange(fn);\n  }\n\n  // ─── Middleware ───────────────────────────────────────\n\n  /**\n   * Add middleware to transform messages before send or after receive.\n   * Return null from middleware to drop the message.\n   *\n   * @example\n   * // Add timestamp to every outgoing message\n   * ws.use('outgoing', (msg) => ({ ...msg, timestamp: Date.now() }));\n   *\n   * @example\n   * // Decrypt incoming messages\n   * ws.use('incoming', (msg) => ({ ...msg, data: decrypt(msg.data) }));\n   *\n   * @example\n   * // Drop messages from blocked users\n   * ws.use('incoming', (msg) => blockedUsers.has(msg.userId) ? null : msg);\n   */\n  use(direction: 'outgoing' | 'incoming', fn: Middleware): this {\n    if (direction === 'outgoing') {\n      this.framePipeline.use(fn);\n    } else {\n      this.incoming.use(fn);\n    }\n    return this;\n  }\n\n  // ─── Per-Event Serialization ─────────────────────────\n\n  /**\n   * Register a custom serializer for a specific event.\n   * The data is transformed before outgoing middleware and global serialize.\n   *\n   * @example\n   * // Binary for file uploads, JSON for everything else\n   * ws.serializer('file.upload', (data) => new Blob([data as ArrayBuffer]));\n   *\n   * @example\n   * // Protobuf for specific event\n   * ws.serializer('trading.order', (data) => OrderProto.encode(data).finish());\n   */\n  serializer(event: string, fn: (data: unknown) => unknown): this {\n    this.serializers.set(event, fn);\n    return this;\n  }\n\n  /**\n   * Register a custom deserializer for a specific event.\n   * The data is transformed after global deserialize and before incoming middleware.\n   *\n   * @example\n   * ws.deserializer('file.download', (data) => new Uint8Array(data as ArrayBuffer));\n   *\n   * @example\n   * // Protobuf for specific event\n   * ws.deserializer('trading.tick', (data) => TickProto.decode(data as Uint8Array));\n   */\n  deserializer(event: string, fn: (data: unknown) => unknown): this {\n    this.incoming.deserializer(event, fn);\n    return this;\n  }\n\n  // ─── Event Subscription ──────────────────────────────\n\n  /**\n   * Subscribe to server events (works in ALL tabs). Type-safe with EventMap.\n   *\n   * The handler receives `(data, raw)`:\n   * - `data` is extracted via `dataField` (default `'data'`)\n   * - `raw` is the full deserialized envelope, useful for protocols with extra\n   *   top-level fields like `id`, `kind`, `channel`, `type`, etc.\n   *\n   * @example\n   * ws.on('msg', (data, raw) => {\n   *   raw.id;    // top-level metadata\n   *   raw.kind;  // discriminator\n   * });\n   */\n  on<K extends string & keyof TEvents>(event: K, handler: EventHandler<TEvents[K]>): Unsubscribe;\n  on(event: string, handler: EventHandler<unknown>): Unsubscribe;\n  // eslint-disable-next-line @typescript-eslint/no-explicit-any\n  on(event: string, handler: (data: any, raw?: unknown) => void): Unsubscribe {\n    return this.subs.on(event, handler);\n  }\n\n  once<K extends string & keyof TEvents>(event: K, handler: EventHandler<TEvents[K]>): Unsubscribe;\n  once(event: string, handler: EventHandler<unknown>): Unsubscribe;\n  // eslint-disable-next-line @typescript-eslint/no-explicit-any\n  once(event: string, handler: (data: any, raw?: unknown) => void): Unsubscribe {\n    return this.subs.once(event, handler);\n  }\n\n  off(event: string, handler?: EventHandler): void {\n    this.subs.off(event, handler);\n  }\n\n  /** Async generator for consuming events. Type-safe with EventMap. */\n  stream<K extends string & keyof TEvents>(event: K, signal?: AbortSignal): AsyncGenerator<TEvents[K]>;\n  stream(event: string, signal?: AbortSignal): AsyncGenerator<unknown>;\n  stream(event: string, signal?: AbortSignal): AsyncGenerator<unknown> {\n    return this.subs.stream(event, signal);\n  }\n\n  /**\n   * Send message to server (auto-routed through leader). Type-safe with EventMap.\n   *\n   * The optional third argument `extras` adds top-level fields to the wire envelope.\n   * Use it for protocols that need extra envelope keys like `type`, `channel`, etc.\n   *\n   * @example\n   * // Default shape: { event, data }\n   * ws.send('chat.message', { text: 'Hello' });\n   * // → { event: 'chat.message', data: { text: 'Hello' } }\n   *\n   * @example\n   * // Pusher/Reverb-style envelope\n   * ws.send('group.member_ready',\n   *   { member_id: 'abc', ready: true },\n   *   { type: 'event', channel: 'public.group.xxx' },\n   * );\n   * // → {\n   * //     type: 'event',\n   * //     channel: 'public.group.xxx',\n   * //     event: 'group.member_ready',\n   * //     data: { member_id: 'abc', ready: true },\n   * //   }\n   */\n  send<K extends string & keyof TEvents>(event: K, data: TEvents[K], extras?: Record<string, unknown>): void;\n  send(event: string, data: unknown, extras?: Record<string, unknown>): void;\n  send(event: string, data: unknown, extras?: Record<string, unknown>): void {\n    this.assertExtrasReserved(extras);\n\n    // Per-event serializer transforms data before the frame is built\n    const eventSerializer = this.serializers.get(event);\n    const serializedData = eventSerializer ? eventSerializer(data) : data;\n\n    this.dispatch('event', { event, data: serializedData, extras });\n  }\n\n  private assertExtrasReserved(extras: Record<string, unknown> | undefined): void {\n    if (!extras) return;\n    if (this.proto.eventField in extras) {\n      throw new Error(\n        `SharedWebSocket.send: extras cannot contain reserved key \"${this.proto.eventField}\" (eventField). ` +\n          `Pass the event name as the first argument instead.`,\n      );\n    }\n    if (this.proto.dataField in extras) {\n      throw new Error(\n        `SharedWebSocket.send: extras cannot contain reserved key \"${this.proto.dataField}\" (dataField). ` +\n          `Pass the payload as the second argument instead.`,\n      );\n    }\n  }\n\n  /** Request/response through server via leader. */\n  async request<T>(event: string, data: unknown, timeout = WS_DEFAULTS.REQUEST_TIMEOUT): Promise<T> {\n    // If THIS tab owns the socket, answer locally. Routing through the bus\n    // would dead-end: `bus.respond` ignores requests from its own tabId, so a\n    // leader calling request() would never get a response and always time out.\n    if (this.coordinator.isLeader && this.socket) {\n      return this.performRequest(event, data, timeout) as Promise<T>;\n    }\n    return this.bus.request(BUS.REQUEST, { event, data, timeout }, timeout);\n  }\n\n  /**\n   * Leader-side request/response over the live socket. Sends the event, then\n   * resolves with the first matching response frame (matched by event name or\n   * a `requestId` field), or rejects on timeout. Used by both the local\n   * `request()` fast-path and the `ws:request` responder for followers.\n   */\n  private performRequest(event: string, data: unknown, timeout: number): Promise<unknown> {\n    return new Promise((resolve, reject) => {\n      const socket = this.socket;\n      if (!socket) {\n        reject(new Error('SharedWebSocket.request: no socket on leader'));\n        return;\n      }\n      let settled = false;\n      const finish = (fn: () => void) => {\n        if (settled) return;\n        settled = true;\n        clearTimeout(timer);\n        unsub();\n        fn();\n      };\n      const unsub = socket.onMessage((response: unknown) => {\n        const res = response as Record<string, unknown> | undefined;\n        if (res?.[this.proto.eventField] === event || res?.requestId) {\n          finish(() => resolve(res?.[this.proto.dataField] ?? response));\n        }\n      });\n      const timer = setTimeout(\n        () => finish(() => reject(new Error(`SharedWebSocket.request: timeout for \"${event}\"`))),\n        timeout,\n      );\n      this.framePipeline.transmit('event', { event, data });\n    });\n  }\n\n  /** Sync state across tabs (no server roundtrip). */\n  sync<T>(key: string, value: T): void {\n    this.syncStore.set(key, value);\n    this.bus.broadcast(BUS.SYNC, { key, value });\n  }\n\n  getSync<T>(key: string): T | undefined {\n    return this.syncStore.get(key) as T | undefined;\n  }\n\n  onSync<T>(key: string, fn: (value: T) => void): Unsubscribe {\n    return this.subs.on(`sync:${key}`, fn as EventHandler);\n  }\n\n  /**\n   * Subscribe to a private/scoped channel. Returns a channel handle with\n   * scoped on/send/stream methods. Sends join on subscribe, leave on unsubscribe.\n   *\n   * @example\n   * const chat = ws.channel('chat:room_123');\n   * chat.on('message', (msg) => render(msg));\n   * chat.send('message', { text: 'Hello' });\n   * chat.leave(); // sends leave + unsubscribes\n   *\n   * @example\n   * // Private notifications for tenant\n   * const notifications = ws.channel(`tenant:${tenantId}:notifications`);\n   * notifications.on('alert', (alert) => showToast(alert));\n   */\n  channel(name: string, options?: { auth?: boolean }): Channel {\n    // Set up the ack matcher BEFORE dispatching so we don't miss a fast\n    // server response. With no matcher configured, ready resolves\n    // synchronously on the next microtask after dispatch.\n    const matcher = this.proto.channelAckMatcher;\n    const ackTimeout = this.proto.channelAckTimeout ?? WS_DEFAULTS.CHANNEL_ACK_TIMEOUT;\n    let cancelReady: ((reason: Error) => void) | undefined;\n\n    const ready = matcher\n      ? new Promise<void>((resolve, reject) => {\n          let settled = false;\n          const settle = (fn: () => void) => {\n            if (settled) return;\n            settled = true;\n            clearTimeout(timer);\n            unsubAck();\n            fn();\n          };\n          const unsubAck = this.onRawFrame((frame) => {\n            let result: ChannelAckResult;\n            try {\n              result = matcher(frame, name);\n            } catch {\n              // matcher exceptions are treated as a hard reject\n              result = 'reject';\n            }\n            if (result === 'ok') settle(() => resolve());\n            else if (result === 'reject') settle(() => reject(new Error(`SharedWebSocket: subscribe rejected for channel \"${name}\"`)));\n          });\n          const timer = setTimeout(\n            () => settle(() => reject(new Error(`SharedWebSocket: subscribe ack timeout for channel \"${name}\"`))),\n            ackTimeout,\n          );\n          cancelReady = (err: Error) => settle(() => reject(err));\n        })\n      : Promise.resolve();\n\n    // Avoid noisy unhandled-rejection warnings if the user never awaits ready.\n    if (matcher) ready.catch(() => {});\n\n    // Notify server about channel subscription\n    this.dispatch('subscribe', { channel: name });\n\n    // Track this channel for incoming-event prefix routing\n    this.subscriptions.addChannel(name);\n\n    const self = this;\n    const unsubs: Unsubscribe[] = [];\n    const isAuth = options?.auth ?? false;\n    let left = false;\n    const key = (event: string) => `${name}${CHANNEL_KEY_SEP}${event}`;\n\n    const ch: Channel = {\n      name,\n      ready,\n      on(event: string, handler: EventHandler): Unsubscribe {\n        const unsub = self.subs.on(key(event), handler);\n        unsubs.push(unsub);\n        return unsub;\n      },\n      once(event: string, handler: EventHandler): Unsubscribe {\n        const unsub = self.subs.once(key(event), handler);\n        unsubs.push(unsub);\n        return unsub;\n      },\n      send(event: string, data: unknown): void {\n        // Channel name is passed structurally so a custom frameBuilder can\n        // emit it as a top-level wire field (Pusher/Reverb-style). The\n        // default builder joins as `${channel}:${event}` for back-compat.\n        // Per-event serializers are keyed on the joined name (legacy).\n        const joined = `${name}:${event}`;\n        const eventSerializer = self.serializers.get(joined) ?? self.serializers.get(event);\n        const serializedData = eventSerializer ? eventSerializer(data) : data;\n        self.dispatch('event', { event, data: serializedData, channel: name });\n      },\n      stream(event: string, signal?: AbortSignal): AsyncGenerator<unknown> {\n        return self.subs.stream(key(event), signal);\n      },\n      leave(): void {\n        if (left) return;\n        left = true;\n        cancelReady?.(new Error(`SharedWebSocket: channel \"${name}\" left before ack`));\n        self.dispatch('unsubscribe', { channel: name });\n        for (const unsub of unsubs) unsub();\n        unsubs.length = 0;\n        if (isAuth) self.auth.unregisterAuthChannel(name);\n        self.subscriptions.removeChannel(name);\n      },\n    };\n\n    if (isAuth) {\n      this.auth.registerAuthChannel(name, ch);\n    }\n\n    return ch;\n  }\n\n  // ─── Topics ──────────────────────────────────────────\n\n  /**\n   * Subscribe to a server-side topic. Server will start sending events for this topic.\n   * Sends topicSubscribe event (default: \"$topic:subscribe\").\n   *\n   * @example\n   * ws.subscribe('notifications:orders');\n   * ws.subscribe('notifications:payments');\n   * ws.subscribe(`user:${userId}:mentions`);\n   */\n  subscribe(topic: string, options?: { auth?: boolean }): void {\n    this.dispatch('topic-subscribe', { topic });\n    this.subscriptions.addTopic(topic);\n    if (options?.auth) {\n      this.auth.registerAuthTopic(topic);\n    }\n    this.log.debug('[SharedWS] subscribe topic', topic);\n  }\n\n  /**\n   * Unsubscribe from a server-side topic.\n   * Sends topicUnsubscribe event (default: \"$topic:unsubscribe\").\n   */\n  unsubscribe(topic: string): void {\n    this.dispatch('topic-unsubscribe', { topic });\n    this.subscriptions.removeTopic(topic);\n    this.auth.unregisterAuthTopic(topic);\n    this.log.debug('[SharedWS] unsubscribe topic', topic);\n  }\n\n  // ─── Push Notifications ─────────────────────────────\n\n  /**\n   * Subscribe to an event and show notifications.\n   *\n   * **target** controls which tab(s) display the notification:\n   * - `'active'` — only the currently visible tab (default for render)\n   * - `'leader'` — only the leader tab (default for browser Notification)\n   * - `'all'` — every tab (for critical alerts)\n   *\n   * @example\n   * // Custom render — sonner toast on active tab only\n   * ws.push('notification', {\n   *   render: (n) => toast(n.title),\n   *   target: 'active',  // default for render\n   * });\n   *\n   * @example\n   * // Critical alert — show in ALL tabs\n   * ws.push('payment.failed', {\n   *   render: (n) => toast.error('Payment failed!'),\n   *   target: 'all',\n   * });\n   *\n   * @example\n   * // Browser Notification — only from leader\n   * ws.push('order.created', {\n   *   title: (order) => `New Order #${order.id}`,\n   *   target: 'leader',  // default for browser Notification\n   * });\n   *\n   * @example\n   * // Both render + native with different targets\n   * ws.push('order.created', {\n   *   render: (order) => toast(`Order #${order.id}`),  // active tab\n   *   title: (order) => `New Order #${order.id}`,      // leader → native\n   * });\n   */\n  push<T = unknown>(event: string, config: PushConfig<T>): Unsubscribe {\n    return this.pushManager.push(event, config);\n  }\n\n  disconnect(): void {\n    this[Symbol.dispose]();\n  }\n\n  // ─── Frame Pipeline ─────────────────────────────────\n  //\n  // dispatch(kind, payload) is the single entry point for all outgoing\n  // frames (events, channel join/leave, topic sub/unsub, auth login/logout).\n  // - On the leader, it calls transmit() which builds the frame, runs\n  //   outgoing middleware, and writes to the socket.\n  // - On followers, it forwards { kind, payload } over BroadcastChannel;\n  //   the leader's bus subscriber re-enters transmit() so middleware\n  //   runs in exactly one place regardless of which tab originated.\n  //\n  // The actual wire shape is decided by frameBuilder (custom) or\n  // defaultFrameBuilder (legacy two-key { event, data } envelope).\n\n  /**\n   * Subscribe to every raw incoming frame (post-deserialize). Used by\n   * `Channel.ready`'s ack matcher. Internal — not part of the public API.\n   */\n  private onRawFrame(fn: (raw: unknown) => void): Unsubscribe {\n    this.rawFrameListeners.add(fn);\n    return () => { this.rawFrameListeners.delete(fn); };\n  }\n\n  /**\n   * Route a structured frame: the leader transmits directly via the frame\n   * pipeline; followers hand it to the outbox, which buffers (event kinds) and\n   * forwards over the bus for the leader to write.\n   */\n  private dispatch(kind: FrameKind, payload: FramePayload): void {\n    if (this.coordinator.isLeader && this.socket) {\n      this.framePipeline.transmit(kind, payload);\n      return;\n    }\n    this.outbox.route(kind, payload);\n  }\n\n  private createSocket(): SocketAdapter {\n    const socketOptions = {\n      protocols: this.options.protocols,\n      reconnect: this.options.reconnect,\n      reconnectMaxDelay: this.options.reconnectMaxDelay,\n      reconnectMaxRetries: this.options.reconnectMaxRetries,\n      authFailureCloseCodes: this.options.authFailureCloseCodes,\n      heartbeatInterval: this.options.heartbeatInterval,\n      heartbeatTimeout: this.options.heartbeatTimeout,\n      sendBuffer: this.options.sendBuffer,\n      pingPayload: this.proto.ping,\n    };\n\n    if (this.options.useWorker) {\n      // WebSocket runs in a Web Worker — main thread stays free\n      return new WorkerSocket(this.url, {\n        ...socketOptions,\n        workerUrl: this.options.workerUrl,\n        auth: this.options.auth,\n        authToken: this.options.authToken,\n        authParam: this.options.authParam,\n      });\n    }\n\n    // WebSocket runs in main thread (default)\n    return new SharedSocket(this.url, {\n      ...socketOptions,\n      auth: this.options.auth,\n      authToken: this.options.authToken,\n      authParam: this.options.authParam,\n      serialize: this.options.serialize,\n      deserialize: this.options.deserialize,\n    });\n  }\n\n  private handleBecomeLeader(): void {\n    this.log.info('[SharedWS] 👑 became leader');\n    this.socket = this.createSocket();\n    this.framePipeline.setSocket(this.socket);\n    this.auth.startRefresh();\n\n    this.socket.onMessage((raw: unknown) => {\n      const envelope = this.incoming.process(raw);\n      if (envelope) this.bus.broadcast(BUS.MESSAGE, envelope);\n    });\n\n    this.socket.onStateChange((state: string) => {\n      this.log.info('[SharedWS]', state === 'connected' ? '✓ connected' : state === 'reconnecting' ? '🔄 reconnecting' : state === 'failed' ? '✗ reconnect failed' : `state: ${state}`);\n      switch (state) {\n        case 'connected':\n          this.bus.broadcast(BUS.LIFECYCLE, { type: 'connect' });\n          void this.onConnected();\n          break;\n        case 'closed':\n          this.bus.broadcast(BUS.LIFECYCLE, { type: 'disconnect' });\n          break;\n        case 'reconnecting':\n          this.bus.broadcast(BUS.LIFECYCLE, { type: 'reconnecting' });\n          break;\n        case 'failed':\n          this.bus.broadcast(BUS.LIFECYCLE, { type: 'reconnectFailed' });\n          this.bus.broadcast(BUS.LIFECYCLE, { type: 'disconnect' });\n          break;\n      }\n    });\n\n    // Replace any stale responder from a previous promotion, then register a\n    // fresh one bound to the new socket. Stored outside `cleanups` so it's torn\n    // down on leadership loss (see handleLoseLeadership).\n    this.requestResponderCleanup?.();\n    this.requestResponderCleanup = this.bus.respond<{ event: string; data: unknown; timeout?: number }, unknown>(\n      BUS.REQUEST,\n      // Swallow a rejection (timeout / no socket) to undefined: the follower\n      // that asked has its own bus.request timeout and will have given up, so\n      // the late response is dropped anyway. Avoids an unhandled rejection.\n      (req) => this.performRequest(req.event, req.data, req.timeout ?? WS_DEFAULTS.REQUEST_TIMEOUT).catch(() => undefined),\n    );\n\n    void this.socket.connect();\n  }\n\n  /**\n   * Re-establish all server-side state on the freshly connected leader socket:\n   *   1. auth-login (so server accepts subsequent joins on auth channels)\n   *   2. channel-join for the union of channels held by ALL surviving tabs\n   *   3. topic-subscribe for the union of topics held by ALL surviving tabs\n   *\n   * The union covers leader handover: when a follower with handlers is\n   * promoted, no tab's subscriptions get silently dropped. Frames are sent\n   * in FIFO order over the single WebSocket, so auth precedes the joins\n   * that depend on it.\n   */\n  /**\n   * Orchestrate post-connect recovery: replay subscriptions first (so the\n   * server is ready to route events for any channels we still care about),\n   * then drain follower-pending dispatches that didn't reach the previous\n   * leader's socket.\n   */\n  private async onConnected(): Promise<void> {\n    if (!this.socket) return;\n    const socket = this.socket;\n    const stillValid = () => this.socket === socket;\n\n    // 1. Re-authenticate first so subsequent auth-channel joins succeed.\n    this.auth.reauthenticate();\n    // 2. Replay the union of channels/topics across all surviving tabs.\n    await this.subscriptions.replay(stillValid);\n    // 3. Drain follower-pending dispatches that didn't reach the old leader.\n    await this.outbox.replay(stillValid);\n  }\n\n  private handleLoseLeadership(): void {\n    this.auth.stopRefresh();\n    this.requestResponderCleanup?.();\n    this.requestResponderCleanup = null;\n    if (this.socket) {\n      this.socket[Symbol.dispose]();\n      this.socket = null;\n    }\n    this.framePipeline.setSocket(null);\n  }\n\n  [Symbol.dispose](): void {\n    if (this.disposed) return;\n    this.disposed = true;\n    this.domListeners.abort(); // removes document/window listeners\n    this.requestResponderCleanup?.();\n    this.requestResponderCleanup = null;\n\n    this.coordinator[Symbol.dispose]();\n    this.auth[Symbol.dispose]();\n    this.outbox[Symbol.dispose]();\n    this.subscriptions[Symbol.dispose]();\n    this.framePipeline.setSocket(null);\n\n    if (this.socket) {\n      this.socket[Symbol.dispose]();\n      this.socket = null;\n    }\n\n    for (const unsub of this.cleanups) unsub();\n    this.cleanups = [];\n    this.subs[Symbol.dispose]();\n    this.bus[Symbol.dispose]();\n    this.syncStore.clear();\n    this.rawFrameListeners.clear();\n  }\n}\n"]}