Skip to content
Draft
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
3 changes: 3 additions & 0 deletions apps/mobile/src/lib/threadActivity.ts
Original file line number Diff line number Diff line change
Expand Up @@ -441,6 +441,9 @@ function deriveWorkLogEntries(
if (activity.kind === "task.updated" && !isTerminalTaskUpdate(activity)) continue;
if (activity.kind === "tool.progress") continue;
if (activity.kind === "context-window.updated") continue;
// Delivery-only nudge for the post-start observation; carries no user
// content and must never render as a work-log row.
if (activity.kind === "post-start-observation") continue;
if (activity.summary === "Checkpoint captured") continue;
if (isNoContentRuntimeWarning(activity)) continue;
if (isPlanBoundaryToolActivity(activity)) continue;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ import { OrchestrationProjectionPipelineLive } from "../src/orchestration/Layers
import { OrchestrationProjectionSnapshotQueryLive } from "../src/orchestration/Layers/ProjectionSnapshotQuery.ts";
import * as ThreadBackgroundLiveness from "../src/orchestration/ThreadBackgroundLiveness.ts";
import * as ThreadPlanProgress from "../src/orchestration/ThreadPlanProgress.ts";
import * as ThreadPostStartActivity from "../src/orchestration/ThreadPostStartActivity.ts";
import { RuntimeReceiptBusTest } from "../src/orchestration/Layers/RuntimeReceiptBus.ts";
import { OrchestrationReactorLive } from "../src/orchestration/Layers/OrchestrationReactor.ts";
import { ProviderCommandReactorLive } from "../src/orchestration/Layers/ProviderCommandReactor.ts";
Expand Down Expand Up @@ -317,6 +318,7 @@ export const makeOrchestrationIntegrationHarness = (
).pipe(
Layer.provideMerge(ThreadBackgroundLiveness.layer),
Layer.provideMerge(ThreadPlanProgress.layer),
Layer.provideMerge(ThreadPostStartActivity.layer),
);
const serverSettingsLayer = ServerSettingsService.layerTest();
const runtimeIngestionLayer = ProviderRuntimeIngestionLive.pipe(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ import { OrchestrationProjectionPipelineLive } from "./ProjectionPipeline.ts";
import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQuery.ts";
import * as ThreadBackgroundLiveness from "../ThreadBackgroundLiveness.ts";
import * as ThreadPlanProgress from "../ThreadPlanProgress.ts";
import * as ThreadPostStartActivity from "../ThreadPostStartActivity.ts";
import { RuntimeReceiptBusTest } from "./RuntimeReceiptBus.ts";
import * as RuntimeReceiptBus from "../Services/RuntimeReceiptBus.ts";
import { OrchestrationEventStoreLive } from "../../persistence/Layers/OrchestrationEventStore.ts";
Expand Down Expand Up @@ -326,6 +327,7 @@ describe("CheckpointReactor", () => {
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(OrchestrationProjectionPipelineLive),
Layer.provide(OrchestrationEventStoreLive),
Layer.provide(OrchestrationCommandReceiptRepositoryLive),
Expand All @@ -335,6 +337,7 @@ describe("CheckpointReactor", () => {
const projectionSnapshotLayer = OrchestrationProjectionSnapshotQueryLive.pipe(
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(RepositoryIdentityResolver.layer),
Layer.provide(SqlitePersistenceMemory),
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ import { OrchestrationProjectionPipelineLive } from "./ProjectionPipeline.ts";
import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQuery.ts";
import * as ThreadBackgroundLiveness from "../ThreadBackgroundLiveness.ts";
import * as ThreadPlanProgress from "../ThreadPlanProgress.ts";
import * as ThreadPostStartActivity from "../ThreadPostStartActivity.ts";
import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts";
import {
OrchestrationProjectionPipeline,
Expand Down Expand Up @@ -79,6 +80,7 @@ function makeOrchestrationLayer(
).pipe(
Layer.provideMerge(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(OrchestrationEventStoreLive),
Layer.provideMerge(OrchestrationCommandReceiptRepositoryLive),
Layer.provide(
Expand Down Expand Up @@ -1514,6 +1516,8 @@ describe("OrchestrationEngine", () => {
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(OrchestrationProjectionPipelineLive),
Layer.provide(Layer.succeed(OrchestrationEventStore, flakyStore)),
Layer.provide(OrchestrationCommandReceiptRepositoryLive),
Expand Down Expand Up @@ -1622,6 +1626,8 @@ describe("OrchestrationEngine", () => {
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(Layer.succeed(OrchestrationProjectionPipeline, flakyProjectionPipeline)),
Layer.provide(OrchestrationEventStoreLive),
Layer.provide(OrchestrationCommandReceiptRepositoryLive),
Expand Down Expand Up @@ -1771,6 +1777,8 @@ describe("OrchestrationEngine", () => {
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(Layer.succeed(OrchestrationProjectionPipeline, flakyProjectionPipeline)),
Layer.provide(Layer.succeed(OrchestrationEventStore, nonTransactionalStore)),
Layer.provide(OrchestrationCommandReceiptRepositoryLive),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQu
import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts";
import * as ThreadBackgroundLiveness from "../ThreadBackgroundLiveness.ts";
import * as ThreadPlanProgress from "../ThreadPlanProgress.ts";
import * as ThreadPostStartActivity from "../ThreadPostStartActivity.ts";
import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts";
import { OrchestrationProjectionPipeline } from "../Services/ProjectionPipeline.ts";
import { ServerConfig } from "../../config.ts";
Expand Down Expand Up @@ -4358,6 +4359,7 @@ const engineLayer = it.layer(
Layer.provideMerge(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provideMerge(OrchestrationProjectionPipelineLive),
Layer.provide(OrchestrationEventStoreLive),
Layer.provide(OrchestrationCommandReceiptRepositoryLive),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import { ORCHESTRATION_PROJECTOR_NAMES } from "./ProjectionPipeline.ts";
import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQuery.ts";
import * as ThreadBackgroundLiveness from "../ThreadBackgroundLiveness.ts";
import * as ThreadPlanProgress from "../ThreadPlanProgress.ts";
import * as ThreadPostStartActivity from "../ThreadPostStartActivity.ts";
import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts";
import { encodeThreadDetailPageCursor } from "../threadDetailCursor.ts";
import { projectThreadDetailSnapshot } from "../ActivityPayloadProjection.ts";
Expand All @@ -52,6 +53,7 @@ it.effect("reads project shells without loading threads or resolving excluded pr
const layer = OrchestrationProjectionSnapshotQueryLive.pipe(
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(
Layer.succeed(RepositoryIdentityResolver.RepositoryIdentityResolver, {
resolve: (root) =>
Expand Down Expand Up @@ -103,6 +105,7 @@ const projectionSnapshotLayer = it.layer(
OrchestrationProjectionSnapshotQueryLive.pipe(
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provideMerge(RepositoryIdentityResolver.layer),
Layer.provideMerge(SqlitePersistenceMemory),
Layer.provideMerge(NodeServices.layer),
Expand Down Expand Up @@ -628,6 +631,7 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => {
hasActionableProposedPlan: false,
backgroundLiveness: null,
planProgress: null,
postStartActivity: null,
},
]);

Expand Down Expand Up @@ -2404,6 +2408,8 @@ it.effect(
const layer = OrchestrationProjectionSnapshotQueryLive.pipe(
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provideMerge(
Layer.succeed(RepositoryIdentityResolver.RepositoryIdentityResolver, {
resolve: (cwd: string) =>
Expand Down Expand Up @@ -3452,6 +3458,7 @@ it.effect("omits foreign-host PRs from legacy snapshots while preserving native
const layer = OrchestrationProjectionSnapshotQueryLive.pipe(
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(
Layer.succeed(RepositoryIdentityResolver.RepositoryIdentityResolver, {
resolve: () =>
Expand Down
17 changes: 17 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ import {
} from "@t3tools/contracts";
import { legacyLinkedPullRequestOf } from "@t3tools/shared/threadPullRequests";
import * as Arr from "effect/Array";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
Expand All @@ -55,6 +56,7 @@ import {
} from "../../persistence/Errors.ts";
import { ThreadBackgroundLivenessService } from "../ThreadBackgroundLiveness.ts";
import { ThreadPlanProgressService } from "../ThreadPlanProgress.ts";
import { ThreadPostStartActivityService } from "../ThreadPostStartActivity.ts";
import { ProjectionProject } from "../../persistence/Services/ProjectionProjects.ts";
import { ProjectionState } from "../../persistence/Services/ProjectionState.ts";
import { ProjectionThreadActivity } from "../../persistence/Services/ProjectionThreadActivities.ts";
Expand Down Expand Up @@ -493,9 +495,17 @@ function toPersistenceSqlOrDecodeError(sqlOperation: string, decodeOperation: st
const makeProjectionSnapshotQuery = Effect.gen(function* () {
const threadBackgroundLiveness = yield* ThreadBackgroundLivenessService;
const threadPlanProgress = yield* ThreadPlanProgressService;
const threadPostStartActivity = yield* ThreadPostStartActivityService;
const sql = yield* SqlClient.SqlClient;
const repositoryIdentityResolver = yield* RepositoryIdentityResolver.RepositoryIdentityResolver;
const repositoryIdentityResolutionConcurrency = 4;
// Live observation is stamped with the server clock at mapping time so the
// client can measure provider ages against the server's clock rather than
// assuming its own agrees.
const readPostStartActivity = (threadId: string, observedAt: string) => {
const state = threadPostStartActivity.getThreadPostStartActivity(threadId);
return state === null ? null : { ...state, observedAt };
};
const resolveRepositoryIdentitiesForProjects = Effect.fn(
"ProjectionSnapshotQuery.resolveRepositoryIdentitiesForProjects",
)(function* (
Expand Down Expand Up @@ -2703,6 +2713,7 @@ pending_approval_requests AS (
sessionRows.map((row) => [row.threadId, mapSessionRow(row)] as const),
);
const pullRequestsByThread = groupPullRequestRowsByThread(pullRequestRows);
const postStartObservedAt = DateTime.formatIso(yield* DateTime.now);

const snapshot = {
snapshotSequence: computeSnapshotSequence(stateRows),
Expand Down Expand Up @@ -2753,6 +2764,7 @@ pending_approval_requests AS (
row.threadId,
),
planProgress: threadPlanProgress.getThreadPlanProgress(row.threadId),
postStartActivity: readPostStartActivity(row.threadId, postStartObservedAt),
} satisfies OrchestrationThreadShell)
: Result.failVoid,
),
Expand Down Expand Up @@ -2868,6 +2880,7 @@ pending_approval_requests AS (
const sessionByThread = new Map(
sessionRows.map((row) => [row.threadId, mapSessionRow(row)] as const),
);
const postStartObservedAt = DateTime.formatIso(yield* DateTime.now);

const snapshot = {
snapshotSequence: computeSnapshotSequence(stateRows),
Expand Down Expand Up @@ -2916,6 +2929,7 @@ pending_approval_requests AS (
row.threadId,
),
planProgress: threadPlanProgress.getThreadPlanProgress(row.threadId),
postStartActivity: readPostStartActivity(row.threadId, postStartObservedAt),
})),
updatedAt: updatedAt ?? "1970-01-01T00:00:00.000Z",
};
Expand Down Expand Up @@ -3231,6 +3245,8 @@ pending_approval_requests AS (
return Option.none<OrchestrationThreadShell>();
}

const postStartObservedAt = DateTime.formatIso(yield* DateTime.now);

return Option.some({
id: threadRow.value.threadId,
projectId: threadRow.value.projectId,
Expand Down Expand Up @@ -3272,6 +3288,7 @@ pending_approval_requests AS (
threadRow.value.threadId,
),
planProgress: threadPlanProgress.getThreadPlanProgress(threadRow.value.threadId),
postStartActivity: readPostStartActivity(threadRow.value.threadId, postStartObservedAt),
} satisfies OrchestrationThreadShell);
});

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ import { OrchestrationProjectionPipelineLive } from "./ProjectionPipeline.ts";
import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQuery.ts";
import * as ThreadBackgroundLiveness from "../ThreadBackgroundLiveness.ts";
import * as ThreadPlanProgress from "../ThreadPlanProgress.ts";
import * as ThreadPostStartActivity from "../ThreadPostStartActivity.ts";
import {
providerErrorLabelFromInstanceHint,
ProviderCommandReactorLive,
Expand Down Expand Up @@ -408,6 +409,7 @@ describe("ProviderCommandReactor", () => {
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(OrchestrationProjectionPipelineLive),
Layer.provide(OrchestrationEventStoreLive),
Layer.provide(OrchestrationCommandReceiptRepositoryLive),
Expand All @@ -417,6 +419,7 @@ describe("ProviderCommandReactor", () => {
const projectionSnapshotLayer = OrchestrationProjectionSnapshotQueryLive.pipe(
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(ThreadPostStartActivity.layer),
Layer.provide(RepositoryIdentityResolver.layer),
Layer.provide(SqlitePersistenceMemory),
);
Expand Down
Loading
Loading