File
Blob: src/worker/do/repo/catalog/leases.ts
| 1 | import type { Logger } from "@/worker/common/logger"; |
| 2 | |
| 3 | import { asTypedStorage } from "../repoState"; |
| 4 | import type { RepoLease, RepoStateSchema } from "../repoState"; |
| 5 | import { getActivePackCatalogSnapshot } from "./state"; |
| 6 | import type { BeginReceiveResult } from "./shared"; |
| 7 | import { |
| 8 | DEFAULT_HEAD, |
| 9 | LEASE_RETRY_AFTER_SECONDS, |
| 10 | ensureRepoMetadataDefaults, |
| 11 | RECEIVE_LEASE_TTL_MS, |
| 12 | } from "./shared"; |
| 13 | |
| 14 | export async function clearExpiredLeases( |
| 15 | ctx: DurableObjectState, |
| 16 | logger?: Logger, |
| 17 | now: number = Date.now() |
| 18 | ): Promise<void> { |
| 19 | const store = asTypedStorage<RepoStateSchema>(ctx.storage); |
| 20 | const receiveLease = await store.get("receiveLease"); |
| 21 | if (receiveLease && receiveLease.expiresAt <= now) { |
| 22 | await store.delete("receiveLease"); |
| 23 | logger?.debug("lease:expired", { kind: "receive" }); |
| 24 | } |
| 25 | |
| 26 | const compactLease = await store.get("compactLease"); |
| 27 | if (compactLease && compactLease.expiresAt <= now) { |
| 28 | await store.delete("compactLease"); |
| 29 | logger?.debug("lease:expired", { kind: "compact" }); |
| 30 | } |
| 31 | } |
| 32 | |
| 33 | export async function beginReceiveLease( |
| 34 | ctx: DurableObjectState, |
| 35 | logger?: Logger |
| 36 | ): Promise<BeginReceiveResult> { |
| 37 | const store = asTypedStorage<RepoStateSchema>(ctx.storage); |
| 38 | await clearExpiredLeases(ctx, logger); |
| 39 | const existing = await store.get("receiveLease"); |
| 40 | if (existing) return { ok: false, retryAfter: LEASE_RETRY_AFTER_SECONDS }; |
| 41 | |
| 42 | const now = Date.now(); |
| 43 | const lease: RepoLease = { |
| 44 | token: crypto.randomUUID(), |
| 45 | createdAt: now, |
| 46 | expiresAt: now + RECEIVE_LEASE_TTL_MS, |
| 47 | }; |
| 48 | await store.put("receiveLease", lease); |
| 49 | |
| 50 | const activeCatalog = await getActivePackCatalogSnapshot(ctx); |
| 51 | await ensureRepoMetadataDefaults(store); |
| 52 | const refs = (await store.get("refs")) ?? []; |
| 53 | const head = (await store.get("head")) ?? DEFAULT_HEAD; |
| 54 | |
| 55 | return { |
| 56 | ok: true, |
| 57 | lease, |
| 58 | refs, |
| 59 | head, |
| 60 | refsVersion: (await store.get("refsVersion")) || 0, |
| 61 | packsetVersion: (await store.get("packsetVersion")) || 0, |
| 62 | nextPackSeq: (await store.get("nextPackSeq")) || 1, |
| 63 | activeCatalog, |
| 64 | }; |
| 65 | } |
| 66 | |
| 67 | export async function abortReceiveLease(ctx: DurableObjectState, token: string): Promise<boolean> { |
| 68 | const store = asTypedStorage<RepoStateSchema>(ctx.storage); |
| 69 | const existing = await store.get("receiveLease"); |
| 70 | if (!existing || existing.token !== token) return false; |
| 71 | await store.delete("receiveLease"); |
| 72 | return true; |
| 73 | } |
| 74 | |
| 75 | export async function abortCompactionLease( |
| 76 | ctx: DurableObjectState, |
| 77 | token: string |
| 78 | ): Promise<boolean> { |
| 79 | const store = asTypedStorage<RepoStateSchema>(ctx.storage); |
| 80 | const existing = await store.get("compactLease"); |
| 81 | if (!existing || existing.token !== token) return false; |
| 82 | await store.delete("compactLease"); |
| 83 | return true; |
| 84 | } |