Skip to content
Merged
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: 2 additions & 0 deletions apps/server/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
"@pierre/diffs": "catalog:",
"effect": "catalog:",
"node-pty": "^1.1.0",
"web-push": "^3.6.7",
"yaml": "catalog:"
},
"devDependencies": {
Expand All @@ -42,6 +43,7 @@
"@t3tools/web": "workspace:*",
"@types/bun": "1.3.14",
"@types/node": "catalog:",
"@types/web-push": "^3.6.4",
"effect-acp": "workspace:*",
"effect-codex-app-server": "workspace:*",
"vite-plus": "catalog:"
Expand Down
38 changes: 36 additions & 2 deletions apps/server/src/t3x/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,14 @@ import { ServerConfig } from "../config.ts";
import { autoResumeRouteLayer } from "./autoResume/http.ts";
import { AutoResumeReactorLive } from "./autoResume/Reactor.ts";
import { AutoResumeStore, makeAutoResumeStore } from "./autoResume/state.ts";
import { resolveConfig as resolveWebPushConfig } from "./webPush/config.ts";
import { webPushRouteLayer } from "./webPush/http.ts";
import { WebPushReactorLive } from "./webPush/Reactor.ts";
import { makePushSubscriptionStore, PushSubscriptionStore } from "./webPush/state.ts";
import { makeWebPushVapid, WebPushVapid } from "./webPush/vapid.ts";

const AUTO_RESUME_STATE_FILENAME = "t3x-auto-resume.json";
const WEB_PUSH_STATE_FILENAME = "t3x-web-push-subscriptions.json";

/** Wires the durable store to the server state directory. */
const AutoResumeStoreLive = Layer.effect(
Expand All @@ -35,12 +41,37 @@ const AutoResumeStoreLive = Layer.effect(
}),
);

/** Web Push subscription store, wired to the server state directory. */
const PushSubscriptionStoreLive = Layer.effect(
PushSubscriptionStore,
Effect.gen(function* () {
const config = yield* ServerConfig;
const path = yield* Path.Path;
return yield* makePushSubscriptionStore(path.join(config.stateDir, WEB_PUSH_STATE_FILENAME));
}),
);

/** VAPID keypair (env override, else generated + persisted in the secret store). */
const WebPushVapidLive = Layer.effect(WebPushVapid, makeWebPushVapid(resolveWebPushConfig()));

/**
* Shared fork-local deps for the Web Push reactor AND routes. Defined once at module scope so
* both `T3xLayerLive` and `T3xRoutesLive` reference the SAME layer value — Effect memoises
* construction by layer identity across the shared MemoMap, so the reactor and the subscribe
* route mutate one store/keypair rather than racing two copies over a single file (the same
* property `AutoResumeStoreLive` relies on above).
*/
const WebPushDepsLive = Layer.mergeAll(PushSubscriptionStoreLive, WebPushVapidLive);

/**
* The single fork-local layer merged into the server. The auto-resume supervisor
* self-starts on construction; its store is provided here so `server.ts` merges only
* this one layer.
*/
export const T3xLayerLive = AutoResumeReactorLive.pipe(Layer.provide(AutoResumeStoreLive));
export const T3xLayerLive = Layer.mergeAll(
AutoResumeReactorLive.pipe(Layer.provide(AutoResumeStoreLive)),
WebPushReactorLive.pipe(Layer.provide(WebPushDepsLive)),
);

/**
* All fork-local HTTP routes, fanned in here for the same reason as `T3xLayerLive`:
Expand All @@ -66,4 +97,7 @@ export const T3xLayerLive = AutoResumeReactorLive.pipe(Layer.provide(AutoResumeS
* probes rather than importing these two layers, so treat it as a guard on the
* assumption, not proof of this file's graph.
*/
export const T3xRoutesLive = autoResumeRouteLayer.pipe(Layer.provide(AutoResumeStoreLive));
export const T3xRoutesLive = Layer.mergeAll(
autoResumeRouteLayer.pipe(Layer.provide(AutoResumeStoreLive)),
webPushRouteLayer.pipe(Layer.provide(WebPushDepsLive)),
);
145 changes: 145 additions & 0 deletions apps/server/src/t3x/webPush/Reactor.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,145 @@
/**
* WebPushReactor — closed-tab notification supervisor.
*
* Self-starts one scoped fiber at layer construction (no external `.start()`, so the only
* upstream seam stays the lines already in server.ts). It taps the hot orchestration
* domain-event stream, and for each thread event recomputes the awareness phase from the
* projected shell (the SAME `projectThreadAwareness` the web client uses), edge-detects, and
* on a firing transition sends a Web Push to every registered subscription. Dead
* subscriptions (404/410) are pruned as they surface.
*
* De-dup with the in-page coordinator is the service worker's job: it suppresses when a tab
* is open, so this reactor always sends and the worker only shows when no tab exists.
*
* @module t3x/webPush/Reactor
*/

import type { OrchestrationEvent, OrchestrationProjectShell, ThreadId } from "@t3tools/contracts";
import { projectThreadAwareness } from "@t3tools/shared/agentAwareness";
import * as Cause from "effect/Cause";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Stream from "effect/Stream";

import { ServerEnvironment } from "../../environment/ServerEnvironment.ts";
import { OrchestrationEngineService } from "../../orchestration/Services/OrchestrationEngine.ts";
import { ProjectionSnapshotQuery } from "../../orchestration/Services/ProjectionSnapshotQuery.ts";
import { buildAttentionPayload, createAttentionEdgeTracker } from "./attention.ts";
import { resolveConfig } from "./config.ts";
import { sendWebPush } from "./send.ts";
import { PushSubscriptionStore } from "./state.ts";
import { WebPushVapid } from "./vapid.ts";

const SEND_CONCURRENCY = 4;

/** Local mirror of the relay's `eventThreadId`, kept in-seam to avoid coupling to upstream. */
function eventThreadId(event: OrchestrationEvent): ThreadId | null {
const payload = event.payload as { readonly threadId?: unknown };
if (typeof payload.threadId === "string") {
return payload.threadId as ThreadId;
}
if (event.aggregateKind === "thread" && typeof event.aggregateId === "string") {
return event.aggregateId as ThreadId;
}
return null;
}

/**
* Skip events that can never change a thread's attention phase, to avoid a shell fetch +
* awareness recompute on every high-frequency domain event. Conservative denylist: anything
* not listed still flows through (correctness over perf).
*/
function isPhaseRelevantEvent(event: OrchestrationEvent): boolean {
switch (event.type) {
case "thread.meta-updated":
case "thread.runtime-mode-set":
case "thread.interaction-mode-set":
case "thread.proposed-plan-upserted":
return false;
default:
return true;
}
}

const makeSupervisor = Effect.gen(function* () {
const config = resolveConfig();
if (!config.enabled) {
yield* Effect.logInfo("t3x web-push: disabled via T3X_WEB_PUSH_ENABLED");
return;
}

const engine = yield* OrchestrationEngineService;
const snapshotQuery = yield* ProjectionSnapshotQuery;
const serverEnvironment = yield* ServerEnvironment;
const store = yield* PushSubscriptionStore;
const vapid = yield* WebPushVapid;
const tracker = createAttentionEdgeTracker();

const handleThread = (threadId: ThreadId) =>
Effect.gen(function* () {
const threadOpt = yield* snapshotQuery.getThreadShellById(threadId);
if (Option.isNone(threadOpt)) {
tracker.forget(threadId);
return;
}
const thread = threadOpt.value;
const projectOpt = yield* snapshotQuery.getProjectShellById(thread.projectId);
const environmentId = yield* serverEnvironment.getEnvironmentId;

const projectTitle = (
Option.isSome(projectOpt) ? projectOpt.value.title : "T3 Code"
) as OrchestrationProjectShell["title"];

const state = projectThreadAwareness({
environmentId,
project: { title: projectTitle },
thread,
});
// Observe unconditionally so the tracker stays primed even with no subscribers. A
// transition that occurs before the first device subscribes then still fires once a
// subscription exists, instead of being swallowed as a first-seen observation.
const kind = tracker.observe(threadId, state?.phase ?? null);
if (kind === null || state === null) {
return;
}

const subscriptions = yield* store.list;
if (subscriptions.length === 0) {
return;
}

const payload = buildAttentionPayload(state, kind);
yield* Effect.forEach(
subscriptions,
(subscription) =>
sendWebPush(vapid, subscription, payload).pipe(
Effect.flatMap((result) =>
result.expired ? store.removeByEndpoint(subscription.endpoint) : Effect.void,
),
),
{ concurrency: SEND_CONCURRENCY, discard: true },
);
});

yield* Effect.forkScoped(
Stream.runForEach(engine.streamDomainEvents, (event) => {
const threadId = eventThreadId(event);
if (threadId === null || !isPhaseRelevantEvent(event)) {
return Effect.void;
}
return handleThread(threadId).pipe(
Effect.catchCause((cause) =>
Effect.logWarning("t3x web-push: notification handler failed", {
eventType: event.type,
cause: Cause.pretty(cause),
}),
),
);
}),
);

yield* Effect.logInfo("t3x web-push: reactor started");
});

export const WebPushReactorLive = Layer.effectDiscard(makeSupervisor);
113 changes: 113 additions & 0 deletions apps/server/src/t3x/webPush/attention.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
import type { AgentAwarenessState } from "@t3tools/shared/agentAwareness";
import { describe, expect, it } from "vite-plus/test";

import {
attentionKey,
attentionKindForEdge,
buildAttentionPayload,
createAttentionEdgeTracker,
} from "./attention.ts";

describe("attentionKindForEdge", () => {
it("announces waiting phases on entry regardless of the previous phase", () => {
expect(attentionKindForEdge("running", "waiting_for_approval")).toBe("waiting_for_approval");
expect(attentionKindForEdge("completed", "waiting_for_input")).toBe("waiting_for_input");
});

it("announces completed only out of running", () => {
expect(attentionKindForEdge("running", "completed")).toBe("completed");
expect(attentionKindForEdge("starting", "completed")).toBeNull();
expect(attentionKindForEdge(null, "completed")).toBeNull();
});

it("stays silent for non-attention phases", () => {
expect(attentionKindForEdge("waiting_for_input", "running")).toBeNull();
expect(attentionKindForEdge("running", "starting")).toBeNull();
expect(attentionKindForEdge("running", "failed")).toBeNull();
expect(attentionKindForEdge("running", "stale")).toBeNull();
});
});

describe("createAttentionEdgeTracker", () => {
it("never fires on a thread's first observation", () => {
const tracker = createAttentionEdgeTracker();
expect(tracker.observe("t1", "waiting_for_input")).toBeNull();
});

it("fires completed only on running -> completed", () => {
const tracker = createAttentionEdgeTracker();
expect(tracker.observe("t1", "running")).toBeNull(); // first-seen
expect(tracker.observe("t1", "completed")).toBe("completed");
});

it("does not fire completed when the thread was never seen running", () => {
const tracker = createAttentionEdgeTracker();
expect(tracker.observe("t1", "starting")).toBeNull(); // first-seen
expect(tracker.observe("t1", "completed")).toBeNull(); // starting -> completed
});

it("fires waiting_for_input on entry after a prior observation", () => {
const tracker = createAttentionEdgeTracker();
expect(tracker.observe("t1", "running")).toBeNull();
expect(tracker.observe("t1", "waiting_for_input")).toBe("waiting_for_input");
});

it("does not re-fire while the phase is unchanged", () => {
const tracker = createAttentionEdgeTracker();
tracker.observe("t1", "running");
expect(tracker.observe("t1", "waiting_for_approval")).toBe("waiting_for_approval");
expect(tracker.observe("t1", "waiting_for_approval")).toBeNull();
});

it("treats a null phase as an observation that never fires", () => {
const tracker = createAttentionEdgeTracker();
tracker.observe("t1", "running");
expect(tracker.observe("t1", null)).toBeNull();
// Coming back to completed from a null observation is not a running->completed edge.
expect(tracker.observe("t1", "completed")).toBeNull();
});

it("tracks threads independently", () => {
const tracker = createAttentionEdgeTracker();
tracker.observe("t1", "running");
expect(tracker.observe("t2", "completed")).toBeNull(); // t2 first-seen
expect(tracker.observe("t1", "completed")).toBe("completed");
});

it("re-arms first-seen after forget", () => {
const tracker = createAttentionEdgeTracker();
tracker.observe("t1", "running");
tracker.forget("t1");
// First observation again -> no fire even though it's a completed value.
expect(tracker.observe("t1", "completed")).toBeNull();
});
});

describe("buildAttentionPayload", () => {
it("maps an awareness state + edge into the service-worker payload", () => {
const state = {
environmentId: "env-1",
threadId: "thread-1",
projectTitle: "My Project",
threadTitle: "Fix the bug",
phase: "waiting_for_input",
headline: "Waiting for input",
modelTitle: "claude",
updatedAt: "2026-07-27T00:00:00.000Z",
deepLink: "/threads/env-1/thread-1",
} as unknown as AgentAwarenessState;

expect(buildAttentionPayload(state, "waiting_for_input")).toEqual({
title: "Waiting for input",
body: "My Project · Fix the bug",
key: "env-1::thread-1",
environmentId: "env-1",
threadId: "thread-1",
kind: "waiting_for_input",
});
});

it("keys by environment + thread", () => {
expect(attentionKey("env-1", "thread-1")).toBe("env-1::thread-1");
});
});
Loading
Loading