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
5 changes: 5 additions & 0 deletions apps/server/src/identity/IdentityService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,11 @@ function parseMapDocument(path: string, raw: string): ReadonlyArray<IdentityMapP
return parseIdentityMapDocument(document);
}

/** Read the on-disk map without constructing the service (projection backfill). */
export function readIdentityMapPeopleFromEnv(): ReadonlyArray<IdentityMapPerson> {
return loadPeopleFromEnv();
}

function loadPeopleFromEnv(): ReadonlyArray<IdentityMapPerson> {
const configured = process.env.T3_IDENTITY_MAP_PATH?.trim();
if (configured === undefined || configured.length === 0) {
Expand Down
46 changes: 46 additions & 0 deletions apps/server/src/identity/applyThreadAttribution.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
// @effect-diagnostics nodeBuiltinImport:off
import * as NodeFS from "node:fs";
import * as NodeOS from "node:os";
import * as NodePath from "node:path";
import { afterEach, describe, expect, it } from "vite-plus/test";

import { applyStoredThreadAttribution } from "./applyThreadAttribution.ts";

const previousPath = process.env.T3_IDENTITY_MAP_PATH;

afterEach(() => {
if (previousPath === undefined) {
delete process.env.T3_IDENTITY_MAP_PATH;
} else {
process.env.T3_IDENTITY_MAP_PATH = previousPath;
}
});

describe("applyStoredThreadAttribution", () => {
it("maps a Discord snowflake on originSource to the identity-map person", () => {
const dir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-identity-"));
const mapPath = NodePath.join(dir, "identity-map.yaml");
NodeFS.writeFileSync(
mapPath,
`people:
"147977704522645504":
username: enricopolanski
name: Enrico Polanski
`,
);
process.env.T3_IDENTITY_MAP_PATH = mapPath;

const enriched = applyStoredThreadAttribution({
originSource: {
channel: "discord",
actor: { platformId: "147977704522645504", displayName: "enricopolanski" },
},
participantSummaries: [],
createdAt: "2026-09-18T00:00:00.000Z",
});

expect(enriched.originSource?.personId).toBe("enricopolanski");
expect(enriched.originSource?.username).toBe("enricopolanski");
expect(enriched.participantSummaries?.[0]?.personId).toBe("enricopolanski");
});
});
73 changes: 73 additions & 0 deletions apps/server/src/identity/applyThreadAttribution.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
import type { SourceRef, ThreadParticipantSummary } from "@t3tools/contracts";
import { IdentityUsername, PersonId } from "@t3tools/contracts";
import { enrichThreadAttribution } from "@t3tools/shared/sourceAttribution";

import { readIdentityMapPeopleFromEnv } from "./IdentityService.ts";

/**
* Re-resolve Discord/GitHub/Jira actors on stored origin stamps.
* Threads created while the identity map was empty (or before sourceHint)
* keep `actor.platformId` without personId; Mine then treats them as ours.
*/
export function applyStoredThreadAttribution(input: {
readonly originSource?: SourceRef | null | undefined;
readonly participantSummaries?: ReadonlyArray<ThreadParticipantSummary> | null | undefined;
readonly createdAt: string;
readonly messages?:
| ReadonlyArray<{
readonly role: string;
readonly createdAt: string;
readonly source?: SourceRef | undefined;
}>
| undefined;
}): {
readonly originSource?: SourceRef;
readonly participantSummaries?: ReadonlyArray<ThreadParticipantSummary>;
} {
const people = readIdentityMapPeopleFromEnv();
if (people.length === 0) {
return {
...(input.originSource !== null && input.originSource !== undefined
? { originSource: input.originSource }
: {}),
...(input.participantSummaries !== null &&
input.participantSummaries !== undefined &&
input.participantSummaries.length > 0
? { participantSummaries: input.participantSummaries }
: {}),
};
}

const enriched = enrichThreadAttribution({
originSource: input.originSource ?? null,
participantSummaries: input.participantSummaries ?? [],
createdAt: input.createdAt,
people,
messages: input.messages ?? [],
});

return {
...(enriched.originSource !== null && enriched.originSource !== undefined
? {
originSource: {
...enriched.originSource,
...(enriched.originSource.personId !== undefined
? { personId: PersonId.make(enriched.originSource.personId) }
: {}),
...(enriched.originSource.username !== undefined
? { username: IdentityUsername.make(enriched.originSource.username) }
: {}),
} as SourceRef,
}
: {}),
...(enriched.participantSummaries.length > 0
? {
participantSummaries: enriched.participantSummaries.map((entry) => ({
...entry,
personId: PersonId.make(entry.personId),
username: IdentityUsername.make(entry.username),
})) as ReadonlyArray<ThreadParticipantSummary>,
}
: {}),
};
}
46 changes: 34 additions & 12 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.ts
Original file line number Diff line number Diff line change
@@ -1,13 +1,18 @@
import {
ApprovalRequestId,
IdentityUsername,
isImportedAgentSessionMessageId,
PersonId,
type SourceRef,
UserInputAttachmentAnswerPayload,
type ChatAttachment,
type OrchestrationEvent,
type OrchestrationSessionStatus,
ThreadId,
} from "@t3tools/contracts";
import { compareDateTimeStrings } from "@t3tools/shared/dateTime";
import { withMappedPerson } from "@t3tools/shared/sourceAttribution";
import { readIdentityMapPeopleFromEnv } from "../../identity/IdentityService.ts";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
Expand Down Expand Up @@ -611,41 +616,47 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
]);

// Rebuild origin + participants from projected user messages so shell stays
// consistent after resync/replay (not only first-write stamps).
// consistent after resync/replay (not only first-write stamps). Map Discord
// actor.platformId → personId when the identity map knows the snowflake —
// otherwise Mine keeps these threads as unattributed.
const identityPeople = readIdentityMapPeopleFromEnv();
type ParticipantRow = NonNullable<(typeof existingRow.value)["participantSummaries"]>[number];
const rebuiltParticipants: Array<ParticipantRow> = [];
let rebuiltOrigin = existingRow.value.originSource ?? null;
let rebuiltOrigin = existingRow.value.originSource
? withMappedPerson(existingRow.value.originSource, identityPeople)
: null;
for (const message of messages) {
if (message.role !== "user" || message.source === undefined) {
continue;
}
const messageSource = withMappedPerson(message.source, identityPeople);
if (rebuiltOrigin === null || rebuiltOrigin === undefined) {
rebuiltOrigin = message.source;
rebuiltOrigin = messageSource;
}
if (message.source.personId !== undefined && message.source.username !== undefined) {
const personId = message.source.personId;
if (messageSource.personId !== undefined && messageSource.username !== undefined) {
const personId = messageSource.personId;
const existingParticipantIndex = rebuiltParticipants.findIndex(
(entry) => entry.personId === personId,
);
if (existingParticipantIndex === -1) {
rebuiltParticipants.push({
personId,
username: message.source.username,
firstChannel: message.source.channel,
channels: [message.source.channel],
personId: PersonId.make(personId),
username: IdentityUsername.make(messageSource.username),
firstChannel: messageSource.channel,
channels: [messageSource.channel],
firstParticipatedAt: message.createdAt,
});
} else {
const existingParticipant = rebuiltParticipants[existingParticipantIndex]!;
if (!existingParticipant.channels?.includes(message.source.channel)) {
if (!existingParticipant.channels?.includes(messageSource.channel)) {
rebuiltParticipants[existingParticipantIndex] = {
...existingParticipant,
channels: [
...(existingParticipant.channels ??
(existingParticipant.firstChannel === undefined
? []
: [existingParticipant.firstChannel])),
message.source.channel,
messageSource.channel,
],
};
}
Expand All @@ -672,7 +683,18 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
pendingApprovalCount,
pendingUserInputCount,
hasActionableProposedPlan: hasActionableProposedPlan ? 1 : 0,
originSource: originSource ?? null,
originSource:
originSource === null
? null
: ({
...originSource,
...(originSource.personId !== undefined
? { personId: PersonId.make(originSource.personId) }
: {}),
...(originSource.username !== undefined
? { username: IdentityUsername.make(originSource.username) }
: {}),
} as SourceRef),
participantSummaries,
});
});
Expand Down
60 changes: 27 additions & 33 deletions apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ import {
type ThreadPullRequestLink,
} from "@t3tools/contracts";
import { legacyLinkedPullRequestOf } from "@t3tools/shared/threadPullRequests";
import { applyStoredThreadAttribution } from "../../identity/applyThreadAttribution.ts";
import * as Arr from "effect/Array";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
Expand Down Expand Up @@ -2655,12 +2656,12 @@ pending_approval_requests AS (
hasMoreActivities: false,
checkpoints: checkpointsByThread.get(row.threadId) ?? [],
session: sessionsByThread.get(row.threadId) ?? null,
...(row.originSource !== null && row.originSource !== undefined
? { originSource: row.originSource }
: {}),
...(row.participantSummaries !== null && row.participantSummaries !== undefined
? { participantSummaries: row.participantSummaries }
: {}),
...applyStoredThreadAttribution({
originSource: row.originSource ?? null,
participantSummaries: row.participantSummaries ?? null,
createdAt: row.createdAt,
messages: messagesByThread.get(row.threadId) ?? [],
}),
}));

const snapshot = {
Expand Down Expand Up @@ -3087,13 +3088,11 @@ pending_approval_requests AS (
row.projectId,
repositoryIdentities.get(row.projectId),
),
...(row.originSource !== null && row.originSource !== undefined
? { originSource: row.originSource }
: {}),
...(row.participantSummaries !== null &&
row.participantSummaries !== undefined
? { participantSummaries: row.participantSummaries }
: {}),
...applyStoredThreadAttribution({
originSource: row.originSource ?? null,
participantSummaries: row.participantSummaries ?? null,
createdAt: row.createdAt,
}),
latestTurn: latestTurnByThread.get(row.threadId) ?? null,
createdAt: row.createdAt,
updatedAt: row.updatedAt,
Expand Down Expand Up @@ -3257,12 +3256,11 @@ pending_approval_requests AS (
row.projectId,
repositoryIdentities.get(row.projectId),
),
...(row.originSource !== null && row.originSource !== undefined
? { originSource: row.originSource }
: {}),
...(row.participantSummaries !== null && row.participantSummaries !== undefined
? { participantSummaries: row.participantSummaries }
: {}),
...applyStoredThreadAttribution({
originSource: row.originSource ?? null,
participantSummaries: row.participantSummaries ?? null,
createdAt: row.createdAt,
}),
latestTurn: latestTurnByThread.get(row.threadId) ?? null,
createdAt: row.createdAt,
updatedAt: row.updatedAt,
Expand Down Expand Up @@ -3669,13 +3667,11 @@ pending_approval_requests AS (
hasPendingApprovals: threadRow.value.pendingApprovalCount > 0,
hasPendingUserInput: threadRow.value.pendingUserInputCount > 0,
hasActionableProposedPlan: threadRow.value.hasActionableProposedPlan > 0,
...(threadRow.value.originSource !== null && threadRow.value.originSource !== undefined
? { originSource: threadRow.value.originSource }
: {}),
...(threadRow.value.participantSummaries !== null &&
threadRow.value.participantSummaries !== undefined
? { participantSummaries: threadRow.value.participantSummaries }
: {}),
...applyStoredThreadAttribution({
originSource: threadRow.value.originSource ?? null,
participantSummaries: threadRow.value.participantSummaries ?? null,
createdAt: threadRow.value.createdAt,
}),
backgroundLiveness: threadBackgroundLiveness.getThreadBackgroundLiveness(
threadRow.value.threadId,
),
Expand Down Expand Up @@ -4067,13 +4063,11 @@ pending_approval_requests AS (
completedAt: row.completedAt,
})),
session: Option.isSome(sessionRow) ? mapSessionRow(sessionRow.value) : null,
...(threadRow.value.originSource !== null && threadRow.value.originSource !== undefined
? { originSource: threadRow.value.originSource }
: {}),
...(threadRow.value.participantSummaries !== null &&
threadRow.value.participantSummaries !== undefined
? { participantSummaries: threadRow.value.participantSummaries }
: {}),
...applyStoredThreadAttribution({
originSource: threadRow.value.originSource ?? null,
participantSummaries: threadRow.value.participantSummaries ?? null,
createdAt: threadRow.value.createdAt,
}),
};

return Option.some(
Expand Down
Loading
Loading