Skip to content
Closed
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
2 changes: 1 addition & 1 deletion apps/server/src/auth/ServerSecretStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,7 @@ export const make = Effect.gen(function* () {
}),
),
),
Effect.withSpan("ServerSecretStore.get"),
Effect.withTracerEnabled(false),
);

const set: ServerSecretStore["Service"]["set"] = (name, value) => {
Expand Down
38 changes: 38 additions & 0 deletions apps/server/src/preview/PortScanner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import * as Layer from "effect/Layer";
import * as PlatformError from "effect/PlatformError";
import * as Scope from "effect/Scope";
import * as TestClock from "effect/testing/TestClock";
import * as Tracer from "effect/Tracer";
import { expect } from "vite-plus/test";
import { FetchHttpClient } from "effect/unstable/http";

Expand Down Expand Up @@ -665,3 +666,40 @@ effectIt.effect("does not swallow process probe interruption", () =>
}
}),
);

effectIt.effect("idle poll ticks create no spans under the ambient parent", () => {
const spanNames: Array<string> = [];
const tracer = Tracer.make({
span: (options) => {
const span = new Tracer.NativeSpan(options);
const end = span.end.bind(span);
span.end = (endTime, exit) => {
end(endTime, exit);
spanNames.push(span.name);
};
return span;
},
});
// No retain call, so every tick takes the idle early-return path.
const layer = makeProbeFailureLayer(processProbeFailure);

const idlePolls = Effect.gen(function* () {
yield* PortScanner.PortDiscovery;
yield* TestClock.adjust(Duration.seconds(30));
});

// Nesting matters: the recording tracer must already be installed when the
// ambient span is created, and the layer (which forks the poll fiber) must
// build inside that ambient span so a leaked ParentSpan would be observed.
// Assertions run after the scope closes so the ambient span has ended.
return Effect.gen(function* () {
yield* Effect.scoped(
Effect.withTracer(
Effect.withSpan("PortScannerTest.idlePollTick")(Effect.provide(idlePolls, layer)),
tracer,
),
);
expect(spanNames).toContain("PortScannerTest.idlePollTick");
expect(spanNames.filter((name) => name === "PortDiscovery.pollTick")).toHaveLength(0);
});
});
16 changes: 11 additions & 5 deletions apps/server/src/preview/PortScanner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ import * as Semaphore from "effect/Semaphore";
import { FetchHttpClient, HttpClient } from "effect/unstable/http";

import * as ProcessRunner from "../processRunner.ts";
import { forkScopedDetached } from "../serverActivation.ts";

export class PortDiscovery extends Context.Service<
PortDiscovery,
Expand Down Expand Up @@ -548,9 +549,8 @@ export const make = Effect.gen(function* PortDiscoveryMake() {
);
};

const pollTick = Effect.fn("PortDiscovery.pollTick")(
const pollTickActive = Effect.fn("PortDiscovery.pollTick")(
function* () {
if ((yield* Ref.get(stateRef)).retainCount <= 0) return;
const configuredUrls = [
...new Set(
[...(yield* Ref.get(stateRef)).listeners.values()].flatMap(
Expand Down Expand Up @@ -579,9 +579,15 @@ export const make = Effect.gen(function* PortDiscoveryMake() {
),
);

// Idle early-return stays outside Effect.fn so no-op ticks create no span; detach ParentSpan (#5410).
const pollTick = Effect.gen(function* () {
if ((yield* Ref.get(stateRef)).retainCount <= 0) return;
yield* pollTickActive();
});

// Single layer-scoped polling fiber. Ticks are no-ops when no client is
// currently retained, so the cost is one Ref.get every POLL_INTERVAL.
yield* Effect.forkScoped(pollTick().pipe(Effect.repeat(Schedule.spaced(POLL_INTERVAL))));
yield* forkScopedDetached(pollTick.pipe(Effect.repeat(Schedule.spaced(POLL_INTERVAL))));

const acquireRetention = Effect.fn("PortDiscovery.retain")(function* () {
const wasIdle = yield* Ref.modify(stateRef, (state) => [
Expand All @@ -591,7 +597,7 @@ export const make = Effect.gen(function* PortDiscoveryMake() {
if (wasIdle) {
// Run an immediate scan + broadcast so the new retainer doesn't have
// to wait up to POLL_INTERVAL for the first emission.
yield* pollTick();
yield* pollTick;
}
});

Expand Down Expand Up @@ -660,6 +666,6 @@ export const make = Effect.gen(function* PortDiscoveryMake() {
registerTerminalProcesses,
unregisterTerminal,
});
}).pipe(Effect.withSpan("PortDiscovery.make"));
});

export const layer = Layer.effect(PortDiscovery, make);
179 changes: 178 additions & 1 deletion apps/server/src/serverActivation.test.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,47 @@
import { expect, it } from "@effect/vitest";
import * as Context from "effect/Context";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as Option from "effect/Option";
import * as Tracer from "effect/Tracer";

import { forkParked, ServerActivation } from "./serverActivation.ts";
import { forkParked, ServerActivation, forkScopedDetached } from "./serverActivation.ts";

it.effect("forkParked detaches the activation wait and keeps the same root after release", () =>
Effect.scoped(
Effect.gen(function* () {
const ambient = Tracer.externalSpan({
traceId: "00000000000000000000000000000003",
spanId: "0000000000000003",
sampled: true,
});
const gate = yield* Deferred.make<void>();
const waiting = yield* Deferred.make<Option.Option<Tracer.AnySpan>>();
const running = yield* Deferred.make<Tracer.AnySpan>();
const activation = Effect.gen(function* () {
const parent = yield* Effect.serviceOption(Tracer.ParentSpan);
yield* Deferred.succeed(waiting, parent);
yield* Deferred.await(gate);
});

yield* forkParked(
Effect.service(Tracer.ParentSpan).pipe(
Effect.flatMap((parent) => Deferred.succeed(running, parent)),
),
).pipe(
Effect.provideService(ServerActivation, activation),
Effect.provideService(Tracer.ParentSpan, ambient),
);

const parent = Option.getOrThrow(yield* Deferred.await(waiting));
expect(parent).not.toBe(ambient);
expect(yield* Deferred.isDone(running)).toBe(false);
yield* Deferred.succeed(gate, undefined);
expect(yield* Deferred.await(running)).toBe(parent);
}),
),
);

it.effect("proves a root is parked before returning and releases it with one gate", () =>
Effect.scoped(
Expand All @@ -21,3 +60,141 @@ it.effect("proves a root is parked before returning and releases it with one gat
}),
),
);

it.effect("forkScopedDetached re-roots instead of inheriting the ambient ParentSpan", () =>
Effect.scoped(
Effect.gen(function* () {
const ambient = Tracer.externalSpan({
traceId: "00000000000000000000000000000001",
spanId: "0000000000000001",
sampled: true,
});

const detachedSpanId = yield* Effect.serviceOption(Tracer.ParentSpan).pipe(
Effect.map(Option.map((span) => span.spanId)),
forkScopedDetached,
Effect.flatMap(Fiber.join),
Effect.provideService(Tracer.ParentSpan, ambient),
);
expect(Option.isSome(detachedSpanId)).toBe(true);
if (Option.isSome(detachedSpanId)) {
expect(detachedSpanId.value).not.toBe("0000000000000001");
}

const stillHasParent = yield* Effect.serviceOption(Tracer.ParentSpan).pipe(
Effect.map(Option.map((span) => span.spanId)),
Effect.provideService(Tracer.ParentSpan, ambient),
);
expect(Option.isSome(stillHasParent)).toBe(true);
if (Option.isSome(stillHasParent)) {
expect(stillHasParent.value).toBe("0000000000000001");
}
}),
),
);

it.effect(
"forkScopedDetached keeps ParentSpan-requiring effects working under the fresh root",
() =>
Effect.scoped(
Effect.gen(function* () {
const span = yield* Effect.service(Tracer.ParentSpan).pipe(
forkScopedDetached,
Effect.flatMap(Fiber.join),
Effect.provideService(
Tracer.ParentSpan,
Tracer.externalSpan({
traceId: "00000000000000000000000000000001",
spanId: "0000000000000001",
sampled: true,
}),
),
);
expect(span.spanId).not.toBe("0000000000000001");
}),
),
);

it.effect("forkParked roots do not inherit the ambient ParentSpan", () =>
Effect.scoped(
Effect.gen(function* () {
const ambient = Tracer.externalSpan({
traceId: "00000000000000000000000000000002",
spanId: "0000000000000002",
sampled: true,
});
const observed = yield* Deferred.make<Option.Option<string>>();

yield* forkParked(
Effect.serviceOption(Tracer.ParentSpan).pipe(
Effect.map(Option.map((span) => span.spanId)),
Effect.flatMap((spanId) => Deferred.succeed(observed, spanId)),
),
).pipe(Effect.provideService(Tracer.ParentSpan, ambient));

const spanId = yield* Deferred.await(observed);
expect(Option.isSome(spanId)).toBe(true);
if (Option.isSome(spanId)) {
expect(spanId.value).not.toBe("0000000000000002");
}
}),
),
);

// Inspect actual context references, including Effect's overlay/base/cache roots.
// Service lookup alone hides replaced spans that are still retained underneath.
const retainsReference = (root: unknown, target: object): boolean => {
const seen = new Set<object>();
const pending: unknown[] = [root];
while (pending.length > 0) {
const value = pending.pop();
if (value === target) return true;
if (typeof value !== "object" || value === null || seen.has(value)) continue;
seen.add(value);
pending.push(...(value instanceof Map ? value.values() : Object.values(value)));
}
return false;
};

for (const gated of [false, true]) {
it.effect(`forkParked drops retained parent references (${gated ? "gated" : "immediate"})`, () =>
Effect.gen(function* () {
const ambient = Tracer.externalSpan({ traceId: "ambient", spanId: "ambient" });
const marker = Context.Service<{ readonly value: number }>("test/detached-marker");
const service = { value: 42 };
const observed = yield* Deferred.make<Context.Context<never>>();
const stopped = yield* Deferred.make<void>();
yield* Effect.scoped(
Effect.gen(function* () {
const observe = Effect.context<never>().pipe(
Effect.flatMap((context) => Deferred.succeed(observed, context)),
Effect.andThen(Effect.never),
Effect.ensuring(Deferred.succeed(stopped, undefined)),
);
yield* forkParked(gated ? Effect.void : observe).pipe(
Effect.provideService(ServerActivation, gated ? observe : undefined),
Effect.provideService(marker, service),
Effect.provideService(Tracer.ParentSpan, ambient),
);
const context = yield* Deferred.await(observed);
expect(Context.getUnsafe(context, marker)).toBe(service);
expect(retainsReference(context, ambient)).toBe(false);
}),
);
expect(yield* Deferred.isDone(stopped)).toBe(true);
}),
);
}

it.effect("forkScopedDetached gives the child a detached context before it starts", () =>
Effect.scoped(
Effect.gen(function* () {
const ambient = Tracer.externalSpan({ traceId: "ambient", spanId: "ambient" });
yield* Effect.gen(function* () {
const fiber = yield* forkScopedDetached(Effect.never);
expect(retainsReference(fiber.context, ambient)).toBe(false);
expect(yield* Tracer.ParentSpan).toBe(ambient);
}).pipe(Effect.provideService(Tracer.ParentSpan, ambient));
}),
),
);
29 changes: 26 additions & 3 deletions apps/server/src/serverActivation.ts
Original file line number Diff line number Diff line change
@@ -1,25 +1,48 @@
import * as Context from "effect/Context";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import type * as Fiber from "effect/Fiber";
import type * as Scope from "effect/Scope";
import * as Tracer from "effect/Tracer";
import * as NodeCrypto from "node:crypto";

export class ServerActivation extends Context.Reference<Effect.Effect<void> | undefined>(
"t3/serverActivation",
{ defaultValue: () => undefined },
) {}

// Replace and flatten the context before forking: Context.add retains the old
// parent in its overlay, while providing inside the child retains a restoration
// frame for the lifetime of the work. The external root does not collect children.
export const forkScopedDetached = <A, E, R>(
effect: Effect.Effect<A, E, R>,
): Effect.Effect<Fiber.Fiber<A, E>, never, Scope.Scope | Exclude<R, Tracer.ParentSpan>> =>
Effect.updateContext(
Effect.forkScoped(effect),
(context: Context.Context<Scope.Scope | Exclude<R, Tracer.ParentSpan>>) => {
const traceId = NodeCrypto.randomUUID().replaceAll("-", "");
const root = Tracer.externalSpan({
traceId,
spanId: traceId.slice(0, 16),
sampled: true,
});
const detached = Context.add(context, Tracer.ParentSpan, root);
return Context.makeUnsafe<Scope.Scope | R>(new Map(detached.mapUnsafe));
},
);

/** Forks a long-running root before commit and proves it is parked at the activation boundary. */
export const forkParked = <A, E, R>(
effect: Effect.Effect<A, E, R>,
): Effect.Effect<void, never, Scope.Scope | R> =>
): Effect.Effect<void, never, Scope.Scope | Exclude<R, Tracer.ParentSpan>> =>
Effect.gen(function* () {
const activation = yield* ServerActivation;
if (activation === undefined) {
yield* Effect.forkScoped(effect);
yield* forkScopedDetached(effect);
return;
}
const parked = yield* Deferred.make<void>();
yield* Effect.forkScoped(
yield* forkScopedDetached(
Deferred.succeed(parked, undefined).pipe(Effect.andThen(activation), Effect.andThen(effect)),
);
yield* Deferred.await(parked);
Expand Down
Loading