import { BootStatus, Devtools, RejectedPushError, liveStoreVersion, MigrationsReport, SyncBackend, SyncState, UnknownError, } from '@livestore/common' import { StreamEventsOptionsFields } from '@livestore/common/leader-thread' import { EventSequenceNumber, LiveStoreEvent } from '@livestore/common/schema' import * as WebmeshWorker from '@livestore/devtools-web-common/worker' import { Schema, Transferable } from '@livestore/utils/effect' export const StorageTypeOpfs = Schema.Struct({ type: Schema.Literal('opfs'), /** * Default is `livestore-${storeId}` * * When providing this option, make sure to include the `storeId` in the path to avoid * conflicts with other LiveStore apps. */ directory: Schema.optional(Schema.String), }) export type StorageTypeOpfs = typeof StorageTypeOpfs.Type // export const StorageTypeIndexeddb = Schema.Struct({ // type: Schema.Literal('indexeddb'), // /** @default "livestore" */ // databaseName: Schema.optionalWith(Schema.String, { default: () => 'livestore' }), // /** @default "livestore-" */ // storeNamePrefix: Schema.optionalWith(Schema.String, { default: () => 'livestore-' }), // }) export const StorageType = Schema.Union( StorageTypeOpfs, // StorageTypeIndexeddb ) export type StorageType = typeof StorageType.Type export type StorageTypeEncoded = typeof StorageType.Encoded // export const SyncBackendOptions = Schema.Union(SyncBackendOptionsWebsocket) export const SyncBackendOptions = Schema.Record({ key: Schema.String, value: Schema.JsonValue }) export type SyncBackendOptions = Record export class LeaderWorkerOuterInitialMessage extends Schema.TaggedRequest()( 'InitialMessage', { payload: { port: Transferable.MessagePort, storeId: Schema.String, clientId: Schema.String }, success: Schema.Void, failure: Schema.Never, }, ) {} export class LeaderWorkerOuterRequest extends Schema.Union(LeaderWorkerOuterInitialMessage) {} // TODO unify this code with schema from node adapter export class LeaderWorkerInnerInitialMessage extends Schema.TaggedRequest()( 'InitialMessage', { payload: { storageOptions: StorageType, devtoolsEnabled: Schema.Boolean, storeId: Schema.String, clientId: Schema.String, debugInstanceId: Schema.String, syncPayloadEncoded: Schema.UndefinedOr(Schema.JsonValue), }, success: Schema.Void, failure: UnknownError, }, ) {} export class LeaderWorkerInnerBootStatusStream extends Schema.TaggedRequest()( 'BootStatusStream', { payload: {}, success: BootStatus, failure: Schema.Never, }, ) {} export class LeaderWorkerInnerPushToLeader extends Schema.TaggedRequest()( 'PushToLeader', { payload: { batch: Schema.Array(Schema.typeSchema(LiveStoreEvent.Client.Encoded)), }, success: Schema.Void as Schema.Schema, failure: RejectedPushError, }, ) {} export class LeaderWorkerInnerPullStream extends Schema.TaggedRequest()('PullStream', { payload: { cursor: Schema.typeSchema(EventSequenceNumber.Client.Composite), }, success: Schema.Struct({ payload: SyncState.PayloadUpstream, }), failure: Schema.Never, }) {} export class LeaderWorkerInnerStreamEvents extends Schema.TaggedRequest()( 'StreamEvents', { payload: StreamEventsOptionsFields, success: LiveStoreEvent.Client.Encoded, failure: Schema.Never, }, ) {} export class LeaderWorkerInnerExport extends Schema.TaggedRequest()('Export', { payload: {}, success: Transferable.Uint8Array as Schema.Schema>, failure: Schema.Never, }) {} export class LeaderWorkerInnerExportEventlog extends Schema.TaggedRequest()( 'ExportEventlog', { payload: {}, success: Transferable.Uint8Array as Schema.Schema>, failure: Schema.Never, }, ) {} export class LeaderWorkerInnerGetRecreateSnapshot extends Schema.TaggedRequest()( 'GetRecreateSnapshot', { payload: {}, success: Schema.Struct({ snapshot: Transferable.Uint8Array as Schema.Schema>, migrationsReport: MigrationsReport, }), failure: Schema.Never, }, ) {} export class LeaderWorkerInnerGetLeaderHead extends Schema.TaggedRequest()( 'GetLeaderHead', { payload: {}, success: Schema.typeSchema(EventSequenceNumber.Client.Composite), failure: Schema.Never, }, ) {} export class LeaderWorkerInnerGetLeaderSyncState extends Schema.TaggedRequest()( 'GetLeaderSyncState', { payload: {}, success: SyncState.SyncState, failure: Schema.Never, }, ) {} export class LeaderWorkerInnerSyncStateStream extends Schema.TaggedRequest()( 'SyncStateStream', { payload: {}, success: SyncState.SyncState, failure: Schema.Never, }, ) {} export class LeaderWorkerInnerGetNetworkStatus extends Schema.TaggedRequest()( 'GetNetworkStatus', { payload: {}, success: SyncBackend.NetworkStatus, failure: Schema.Never, }, ) {} export class LeaderWorkerInnerNetworkStatusStream extends Schema.TaggedRequest()( 'NetworkStatusStream', { payload: {}, success: SyncBackend.NetworkStatus, failure: Schema.Never, }, ) {} export class LeaderWorkerInnerShutdown extends Schema.TaggedRequest()('Shutdown', { payload: {}, success: Schema.Void, failure: Schema.Never, }) {} export class LeaderWorkerInnerExtraDevtoolsMessage extends Schema.TaggedRequest()( 'ExtraDevtoolsMessage', { payload: { message: Devtools.Leader.MessageToApp, }, success: Schema.Void, failure: Schema.Never, }, ) {} export const LeaderWorkerInnerRequest = Schema.Union( LeaderWorkerInnerInitialMessage, LeaderWorkerInnerBootStatusStream, LeaderWorkerInnerPushToLeader, LeaderWorkerInnerPullStream, LeaderWorkerInnerStreamEvents, LeaderWorkerInnerExport, LeaderWorkerInnerExportEventlog, LeaderWorkerInnerGetRecreateSnapshot, LeaderWorkerInnerGetLeaderHead, LeaderWorkerInnerGetLeaderSyncState, LeaderWorkerInnerSyncStateStream, LeaderWorkerInnerGetNetworkStatus, LeaderWorkerInnerNetworkStatusStream, LeaderWorkerInnerShutdown, LeaderWorkerInnerExtraDevtoolsMessage, WebmeshWorker.Schema.CreateConnection, ) export type LeaderWorkerInnerRequest = typeof LeaderWorkerInnerRequest.Type export class SharedWorkerUpdateMessagePort extends Schema.TaggedRequest()( 'UpdateMessagePort', { payload: { port: Transferable.MessagePort, // Version gate to prevent mixed LiveStore builds talking to the same SharedWorker liveStoreVersion: Schema.Literal(liveStoreVersion), /** * Initial configuration for the leader worker. This replaces the previous * two-phase SharedWorker handshake and is sent under the tab lock by the * elected leader. Subsequent calls can omit changes and will simply rebind * the port (join) without reinitializing the store. */ initial: LeaderWorkerInnerInitialMessage, }, success: Schema.Void, failure: UnknownError, }, ) {} export const SharedWorkerRequest = Schema.Union( SharedWorkerUpdateMessagePort, // Proxied requests LeaderWorkerInnerBootStatusStream, LeaderWorkerInnerPushToLeader, LeaderWorkerInnerPullStream, LeaderWorkerInnerStreamEvents, LeaderWorkerInnerExport, LeaderWorkerInnerGetRecreateSnapshot, LeaderWorkerInnerExportEventlog, LeaderWorkerInnerGetLeaderHead, LeaderWorkerInnerGetLeaderSyncState, LeaderWorkerInnerSyncStateStream, LeaderWorkerInnerGetNetworkStatus, LeaderWorkerInnerNetworkStatusStream, LeaderWorkerInnerShutdown, LeaderWorkerInnerExtraDevtoolsMessage, WebmeshWorker.Schema.CreateConnection, ) export type SharedWorkerRequest = typeof SharedWorkerRequest.Type