( config: IntegrationsRegistryConfig, )
| 114 | }; |
| 115 | |
| 116 | export const layer = ( |
| 117 | config: IntegrationsRegistryConfig, |
| 118 | ): Layer.Layer<IntegrationsRegistry, never, HttpClient.HttpClient | FileSystem.FileSystem> => |
| 119 | Layer.effect( |
| 120 | IntegrationsRegistry, |
| 121 | Effect.gen(function* () { |
| 122 | const fs = yield* FileSystem.FileSystem; |
| 123 | const http = yield* HttpClient.HttpClient; |
| 124 | |
| 125 | const source = |
| 126 | config.url ?? process.env.EXECUTOR_INTEGRATIONS_URL ?? DEFAULT_INTEGRATIONS_URL; |
| 127 | const cacheDir = resolveCacheDir(config.cacheDir); |
| 128 | const cacheFile = cacheFileFor(cacheDir, source); |
| 129 | const lockFile = `${cacheFile}.lock`; |
| 130 | const ttl = Duration.fromInputUnsafe(config.cacheTtl ?? Duration.hours(12)); |
| 131 | const refreshEvery = Duration.fromInputUnsafe(config.refreshInterval ?? Duration.hours(12)); |
| 132 | const disabled = config.disabled ?? isFetchDisabled(); |
| 133 | |
| 134 | const isFresh = Effect.gen(function* () { |
| 135 | const stat = yield* fs.stat(cacheFile).pipe(Effect.catch(() => Effect.succeed(undefined))); |
| 136 | if (!stat) return false; |
| 137 | const mtime = Option.getOrElse(stat.mtime, () => new Date(0)).getTime(); |
| 138 | return Date.now() - mtime < Duration.toMillis(ttl); |
| 139 | }); |
| 140 | |
| 141 | const fetchText = Effect.gen(function* () { |
| 142 | const request = HttpClientRequest.get(source).pipe( |
| 143 | HttpClientRequest.setHeader("user-agent", config.userAgent), |
| 144 | HttpClientRequest.setHeader("accept", "application/json"), |
| 145 | ); |
| 146 | const response = yield* http.execute(request); |
| 147 | return yield* response.text; |
| 148 | }).pipe(Effect.timeout(Duration.seconds(10)), Effect.withSpan("IntegrationsRegistry.fetch")); |
| 149 | |
| 150 | const writeCache = (text: string) => |
| 151 | Effect.gen(function* () { |
| 152 | yield* fs.makeDirectory(cacheDir, { recursive: true }).pipe(Effect.ignore); |
| 153 | yield* fs.writeFileString(cacheFile, text); |
| 154 | }); |
| 155 | |
| 156 | // Cross-process advisory lock so concurrent CLI invocations don't all |
| 157 | // race to refresh the same file. Best-effort: if we can't acquire, |
| 158 | // another process is already refreshing — skip and read what they wrote. |
| 159 | const withWriteLock = <A, E>(use: Effect.Effect<A, E>) => |
| 160 | Effect.acquireUseRelease( |
| 161 | Effect.gen(function* () { |
| 162 | yield* fs.makeDirectory(cacheDir, { recursive: true }).pipe(Effect.ignore); |
| 163 | return yield* fs.writeFileString(lockFile, `${process.pid}\n`, { flag: "wx" }).pipe( |
| 164 | Effect.as(true), |
| 165 | Effect.catch(() => Effect.succeed(false)), |
| 166 | ); |
| 167 | }), |
| 168 | (acquired): Effect.Effect<Option.Option<A>, E> => |
| 169 | acquired ? Effect.map(use, Option.some) : Effect.succeed(Option.none<A>()), |
| 170 | (acquired) => |
| 171 | acquired ? fs.remove(lockFile, { force: true }).pipe(Effect.ignore) : Effect.void, |
| 172 | ); |
| 173 |
no test coverage detected