MCPcopy Create free account
hub / github.com/UsefulSoftwareCo/executor / layer

Function layer

packages/core/integrations-registry/src/registry.ts:116–266  ·  view source on GitHub ↗
(
  config: IntegrationsRegistryConfig,
)

Source from the content-addressed store, hash-verified

114};
115
116export 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

Callers 6

defaultLayerFunction · 0.85
testing.test.tsFile · 0.85
testing.test.tsFile · 0.85
testing.test.tsFile · 0.85

Calls 9

resolveCacheDirFunction · 0.85
cacheFileForFunction · 0.85
isFetchDisabledFunction · 0.85
withWriteLockFunction · 0.85
writeCacheFunction · 0.70
parseJsonFunction · 0.70
refreshFunction · 0.70
getMethod · 0.65
executeMethod · 0.65

Tested by

no test coverage detected