Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 17 additions & 7 deletions packages/core/src/wellknown.ts
Original file line number Diff line number Diff line change
Expand Up @@ -111,15 +111,25 @@ const layer = Layer.effect(
return { origin, integrationID: Integration.ID.make(origin), manifest }
})

const load = Effect.fn("WellKnown.load")(function* () {
const value = yield* kv.get(sourcesKey)
const origins = Schema.is(Sources)(value) ? value : []
const current = yield* Ref.get(cache)
const loadEntries = Effect.fn("WellKnown.loadEntries")(function* (origins: readonly string[], reuse: boolean) {
const current = Ref.getUnsafe(cache)
const entries = yield* Effect.forEach(origins, (origin) => {
const cached = current.get(origin)
if (cached) return Effect.succeed(cached)
return loadEntry(origin)
if (cached && reuse) return Effect.succeed(cached)
// An unreachable origin keeps its last known manifest, or is skipped when nothing is
// cached. One torn-down deployment must not hide every other origin's integrations.
return loadEntry(origin).pipe(
Effect.catch((error) =>
Effect.logWarning("failed to load wellknown manifest", { origin, error }).pipe(Effect.as(cached)),
),
)
})
return entries.filter((entry): entry is Entry => entry !== undefined)
})

const load = Effect.fn("WellKnown.load")(function* () {
const value = yield* kv.get(sourcesKey)
const entries = yield* loadEntries(Schema.is(Sources)(value) ? value : [], true)
yield* Ref.set(cache, new Map(entries.map((entry) => [entry.origin, entry])))
return entries
})
Expand All @@ -129,7 +139,7 @@ const layer = Layer.effect(
const value = yield* kv.get(sourcesKey)
const origins = Schema.is(Sources)(value) ? value : []
if (!origins.length) return false
const entries = yield* Effect.forEach(origins, loadEntry)
const entries = yield* loadEntries(origins, false)
const next = new Map(entries.map((entry) => [entry.origin, entry]))
const changed = !isDeepStrictEqual(Ref.getUnsafe(cache), next)
if (!changed) return false
Expand Down
6 changes: 5 additions & 1 deletion packages/core/src/wellknown/plugin.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,11 @@ export const Plugin = define({
effect: Effect.fn(function* (ctx) {
const bus = yield* Bus.Service
const wellknown = yield* WellKnown.Service
yield* wellknown.entries().pipe(Effect.orDie)
// Priming the cache is best effort. Failing here would abort the plugin before it registers
// its transform, permanently hiding every well-known integration until the next restart.
yield* wellknown
.entries()
.pipe(Effect.catch((error) => Effect.logWarning("failed to load wellknown entries", { error })))
yield* ctx.integration.transform((draft) => {
wellknown.snapshot().forEach((entry) => {
if (!entry.manifest.auth) return
Expand Down
37 changes: 37 additions & 0 deletions packages/core/test/wellknown.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -117,3 +117,40 @@ serviceIt.live("refreshes changed manifests", () =>
({ server }) => Effect.promise(() => server.stop(true)),
),
)

serviceIt.live("keeps healthy origins when another is unreachable", () =>
Effect.acquireUseRelease(
Effect.sync(() =>
Bun.serve({
port: 0,
fetch: () => Response.json({ auth: { command: ["login"], env: "TOKEN" } }),
}),
),
(server) =>
Effect.gen(function* () {
const wellknown = yield* WellKnown.Service
const kv = yield* KV.Service
// Port 1 is closed, standing in for a torn-down deployment still listed in sources.
yield* kv.set("wellknown:sources", ["http://127.0.0.1:1", server.url.origin])

expect((yield* wellknown.entries()).map((entry) => entry.origin)).toEqual([server.url.origin])
expect(wellknown.snapshot().map((entry) => entry.origin)).toEqual([server.url.origin])
}),
(server) => Effect.promise(() => server.stop(true)),
),
)

serviceIt.live("keeps the last manifest when an origin becomes unreachable", () =>
Effect.gen(function* () {
const server = Bun.serve({ port: 0, fetch: () => Response.json({ auth: { command: ["login"], env: "TOKEN" } }) })
const wellknown = yield* WellKnown.Service
const kv = yield* KV.Service
yield* kv.set("wellknown:sources", [server.url.origin])
yield* wellknown.entries()
yield* Effect.promise(() => server.stop(true))

// A temporary outage must not drop the integration and its remote config mid-session.
expect(yield* wellknown.refresh()).toBe(false)
expect(wellknown.snapshot().map((entry) => entry.origin)).toEqual([server.url.origin])
}),
)
Loading