diff --git a/packages/junior-scheduler/src/plugin.ts b/packages/junior-scheduler/src/plugin.ts index 85ec4d7d9..a22f3108c 100644 --- a/packages/junior-scheduler/src/plugin.ts +++ b/packages/junior-scheduler/src/plugin.ts @@ -81,8 +81,7 @@ function scheduledTaskDispatchSource(task: ScheduledTask): Source { return createSlackSource({ teamId: task.destination.teamId, channelId: task.destination.channelId, - - type: "priv", + type: task.conversationAccess?.visibility === "public" ? "pub" : "priv", }); } diff --git a/packages/junior/src/api/conversations/access.ts b/packages/junior/src/api/conversations/access.ts index cfd26213d..b9936004c 100644 --- a/packages/junior/src/api/conversations/access.ts +++ b/packages/junior/src/api/conversations/access.ts @@ -82,10 +82,7 @@ export async function readConversationAccessFromSql( isParticipant || (validRootConversationId !== undefined && canExposeConversationPayload({ - conversationId: validRootConversationId, - ...(visibility === "public" || visibility === "private" - ? { visibility } - : {}), + visibility: visibility ?? undefined, })); return [ row.conversationId, diff --git a/packages/junior/src/api/conversations/projection.ts b/packages/junior/src/api/conversations/projection.ts index f0e0cc700..8758b3d11 100644 --- a/packages/junior/src/api/conversations/projection.ts +++ b/packages/junior/src/api/conversations/projection.ts @@ -1,6 +1,7 @@ import { formatSlackConversationRedactedLabel, resolveSlackConversationContextFromThreadId, + type SlackConversationVisibility, } from "@/chat/slack/conversation-context"; import { parseSlackThreadId } from "@/chat/slack/context"; import { buildSlackSourceUrl } from "@/chat/slack/source-link"; @@ -115,12 +116,14 @@ function titleFromConversation(args: { canViewPrivateContent: boolean; conversation: ConversationProjectionSource; surface: ConversationSurface; + visibility?: SlackConversationVisibility; }): string { const slackThread = parseSlackThreadId(args.conversation.conversationId); const effectiveChannelName = args.conversation.channelName; const slackConversation = resolveSlackConversationContextFromThreadId({ threadId: args.conversation.conversationId, channelName: effectiveChannelName, + ...(args.visibility ? { visibility: args.visibility } : {}), }); const privateLabel = args.canViewPrivateContent ? undefined @@ -139,6 +142,7 @@ function titleFromConversation(args: { function channelNameFromConversation( conversation: ConversationProjectionSource, canViewPrivateContent: boolean, + visibility?: SlackConversationVisibility, ): string | undefined { const effectiveChannelName = conversation.channelName; const slackThread = parseSlackThreadId(conversation.conversationId); @@ -148,6 +152,7 @@ function channelNameFromConversation( const slackConversation = resolveSlackConversationContextFromThreadId({ threadId: conversation.conversationId, channelName: effectiveChannelName, + ...(visibility ? { visibility } : {}), }); if (!canViewPrivateContent) { return privateConversationLabel(slackConversation); @@ -194,6 +199,11 @@ export function conversationSummaryFromStoredConversation(args: { }): ConversationSummaryReport { const { conversation, durationMs, usage } = args; const canViewPrivateContent = args.access?.canViewPrivateContent ?? false; + const accessVisibility = args.access?.visibility; + const visibility = + accessVisibility === "public" || accessVisibility === "private" + ? accessVisibility + : undefined; const surface = surfaceFromSource( conversation.source, conversation.conversationId, @@ -207,6 +217,7 @@ export function conversationSummaryFromStoredConversation(args: { const channelName = channelNameFromConversation( conversation, canViewPrivateContent, + visibility, ); const channelNameRedacted = channelNameRedactedFromConversation( conversation, @@ -219,6 +230,7 @@ export function conversationSummaryFromStoredConversation(args: { canViewPrivateContent, conversation, surface, + visibility, }), isParticipant: args.access?.isParticipant ?? false, lastProgressAt: new Date( diff --git a/packages/junior/src/chat/agent/index.ts b/packages/junior/src/chat/agent/index.ts index d038dc4a5..ffee90649 100644 --- a/packages/junior/src/chat/agent/index.ts +++ b/packages/junior/src/chat/agent/index.ts @@ -108,6 +108,7 @@ import { toGenAiMessagesTraceAttributes, type ConversationPrivacy, } from "@/chat/conversation-privacy"; +import { resolveDestinationVisibility } from "@/chat/conversations/destination-visibility"; import { RetryableDeliveryError, assertRunRoutingConsistency, @@ -199,19 +200,19 @@ export async function executeAgentRun( if (!request.routing.destination) { throw new TypeError("Assistant reply generation requires a destination"); } - const channelId = - request.routing.destination.platform === "slack" - ? request.routing.destination.channelId - : undefined; + const destinationVisibility = await resolveDestinationVisibility({ + destination: request.routing.destination, + visibility: request.routing.destinationVisibility, + }); const conversationPrivacy = resolveConversationPrivacy({ - channelId, - conversationId: request.conversationId, - // Destination visibility is provider-neutral. Slack event context remains - // a compatibility fallback for callers that have not projected it yet. - visibility: - request.routing.destinationVisibility ?? - request.routing.slackConversation?.visibility, + visibility: destinationVisibility, }); + const resolvedRequest = destinationVisibility + ? { + ...request, + routing: { ...request.routing, destinationVisibility }, + } + : request; const credentialActor = request.routing.credentialContext?.actor; const actor = actorFromRouting(request.routing); const userActor = actor && "userId" in actor ? actor : undefined; @@ -243,9 +244,9 @@ export async function executeAgentRun( assistantUserName: botConfig.userName, }; return withLogContext(runLogContext, () => - runWithConversationPrivacy(conversationPrivacy ?? "private", () => + runWithConversationPrivacy(conversationPrivacy, () => executeAgentRunInPrivacyContext( - request, + resolvedRequest, conversationPrivacy, runLogContext, streamFn, @@ -419,6 +420,9 @@ async function executeAgentRunInPrivacyContext( resume = createResumeState({ channelName: routing.slackConversation?.name, destination: routing.destination, + ...(routing.destinationVisibility + ? { destinationVisibility: routing.destinationVisibility } + : {}), ...(routing.dispatch?.id ? { dispatchId: routing.dispatch.id } : {}), durability, getLoadedSkillNames: () => loadedSkillNamesForResume, diff --git a/packages/junior/src/chat/agent/resume.ts b/packages/junior/src/chat/agent/resume.ts index 5c83c527d..13091b821 100644 --- a/packages/junior/src/chat/agent/resume.ts +++ b/packages/junior/src/chat/agent/resume.ts @@ -41,6 +41,7 @@ import { } from "@/chat/agent/request"; import { TurnSliceLimitExceededError } from "@/chat/services/turn-limit"; import type { PluginTurnContext } from "@/chat/plugins/prompt"; +import type { ConversationPrivacy } from "@/chat/conversation-privacy"; type LoadedSessionRecordState = Awaited< ReturnType @@ -48,6 +49,7 @@ type LoadedSessionRecordState = Awaited< interface ResumeStateArgs { channelName?: string; destination: Destination; + destinationVisibility?: ConversationPrivacy; dispatchId?: string; durability: AgentRunDurability; getLoadedSkillNames: () => string[]; @@ -92,6 +94,9 @@ export function createResumeState(args: ResumeStateArgs) { channelName: args.channelName, conversationId: args.conversationId, destination: args.destination, + ...(args.destinationVisibility + ? { destinationVisibility: args.destinationVisibility } + : {}), ...(args.dispatchId ? { dispatchId: args.dispatchId } : {}), source: args.runSource, sessionId: args.turnId, diff --git a/packages/junior/src/chat/conversation-privacy.ts b/packages/junior/src/chat/conversation-privacy.ts index 445f61f96..0787ce4d0 100644 --- a/packages/junior/src/chat/conversation-privacy.ts +++ b/packages/junior/src/chat/conversation-privacy.ts @@ -7,70 +7,33 @@ import type { ThinkingContent, ToolCall, } from "@earendil-works/pi-ai"; -import { parseSlackThreadId } from "@/chat/slack/context"; +import { logWarn } from "@/chat/logging"; export type ConversationPrivacy = "public" | "private"; type TraceAttributeValue = string | number | boolean | string[]; const SAFE_METADATA_KEY_LIMIT = 20; const conversationPrivacyStorage = new AsyncLocalStorage(); -function privateNarrowingFromChannelId( - channelId: string | undefined, -): ConversationPrivacy | undefined { - const normalized = channelId?.trim(); - if (!normalized) return undefined; - // Channel-id prefixes may only narrow toward private. `C`-prefixed ids do - // not prove a conversation public: modern Slack private channels also use - // `C` prefixes, so they stay unknown without a confirmed signal. - return normalized.startsWith("D") || normalized.startsWith("G") - ? "private" - : undefined; -} - -function privateNarrowingFromConversationId( - conversationId: string | undefined, -): ConversationPrivacy | undefined { - const normalized = conversationId?.trim(); - if (!normalized) return undefined; - const slackThread = parseSlackThreadId(normalized); - if (slackThread) { - return privateNarrowingFromChannelId(slackThread.channelId); - } - if (normalized.startsWith("slack:")) { - return undefined; - } - // Non-Slack conversations (local CLI, internal runs) are private surfaces. - return "private"; -} - /** * Resolve whether a conversation may expose raw payloads. * - * Only a live source signal or persisted destination visibility can classify - * a conversation public. Identifier prefixes may only narrow classification - * toward private. Unknown stays undefined so callers fail closed to private. + * Explicit destination visibility is authoritative. Missing visibility is + * observable and fails closed without inferring from identifiers. */ export function resolveConversationPrivacy(input: { - channelId?: string; - conversationId?: string; - /** Live source or persisted visibility, when the caller has one. */ - visibility?: ConversationPrivacy; -}): ConversationPrivacy | undefined { - const narrowed = - privateNarrowingFromChannelId(input.channelId) ?? - privateNarrowingFromConversationId(input.conversationId); - if (narrowed === "private") { - return "private"; + /** Live or persisted destination visibility, when the caller has one. */ + visibility?: ConversationPrivacy | "direct" | "unknown"; +}): ConversationPrivacy { + if (input.visibility === undefined) { + logWarn("conversation.visibility.defaulted"); } - return input.visibility; + return input.visibility === "public" ? "public" : "private"; } /** Gate raw transcript/tool payload exposure to public conversations. */ export function canExposeConversationPayload(input: { - channelId?: string; - conversationId?: string; - /** Live source or persisted visibility, when the caller has one. */ - visibility?: ConversationPrivacy; + /** Live or persisted destination visibility, when the caller has one. */ + visibility?: ConversationPrivacy | "direct" | "unknown"; }): boolean { return resolveConversationPrivacy(input) === "public"; } diff --git a/packages/junior/src/chat/conversations/destination-visibility.ts b/packages/junior/src/chat/conversations/destination-visibility.ts new file mode 100644 index 000000000..a98246600 --- /dev/null +++ b/packages/junior/src/chat/conversations/destination-visibility.ts @@ -0,0 +1,21 @@ +import type { Destination } from "@sentry/junior-plugin-api"; +import { getConversationStore } from "@/chat/db"; +import type { ConversationPrivacy } from "@/chat/conversation-privacy"; + +/** Read confirmed visibility from the current signal or persisted destination. */ +export async function resolveDestinationVisibility(args: { + destination: Destination; + visibility?: ConversationPrivacy; +}): Promise { + if (args.visibility) { + return args.visibility; + } + if (args.destination.platform === "local") { + return "private"; + } + return await getConversationStore().getDestinationVisibility({ + provider: "slack", + providerDestinationId: args.destination.channelId, + providerTenantId: args.destination.teamId, + }); +} diff --git a/packages/junior/src/chat/conversations/sql/store.ts b/packages/junior/src/chat/conversations/sql/store.ts index 01ef5d9af..6a25daa11 100644 --- a/packages/junior/src/chat/conversations/sql/store.ts +++ b/packages/junior/src/chat/conversations/sql/store.ts @@ -174,7 +174,7 @@ function destinationUpsertFromDestination(args: { channelName?: string; conversationId?: string; destination: Destination | undefined; - /** Source-confirmed visibility from the current event's signal only. */ + /** Confirmed destination visibility; omit when unavailable. */ visibility?: ConversationPrivacy; }): DestinationUpsert | undefined { const { destination } = args; @@ -194,7 +194,7 @@ function destinationUpsertFromDestination(args: { providerTenantId: destination.teamId, providerDestinationId: channelId, refreshVisibility: args.visibility !== undefined, - visibility: args.visibility ?? "private", + visibility: args.visibility ?? "unknown", ...(args.channelName ? { displayName: args.channelName } : {}), metadata: { platform: "slack" }, }; @@ -719,7 +719,10 @@ export class SqlStore implements ConversationStore { if (!row) { return undefined; } - return row.visibility === "public" ? "public" : "private"; + if (row.visibility === "public" || row.visibility === "private") { + return row.visibility; + } + return undefined; } /** Serialize all durable mutations for one conversation inside a SQL transaction. */ @@ -900,9 +903,9 @@ export class SqlStore implements ConversationStore { set: { kind: sql`excluded.kind`, displayName: sql`coalesce(excluded.display_name, ${juniorDestinations.displayName})`, - // Signal-less writes insert as private but must not clobber an - // existing public/private value. Live source signals refresh this - // field so converted channels converge on the next message. + // Signal-less writes remain unknown and must not clobber an existing + // public/private value. Live source signals refresh this field so + // converted channels converge on the next message. visibility: visibilityUpdate, metadata: sql`coalesce(excluded.metadata_json, ${juniorDestinations.metadata})`, updatedAt: sql`excluded.updated_at`, diff --git a/packages/junior/src/chat/conversations/store.ts b/packages/junior/src/chat/conversations/store.ts index 2ad1fea34..eb41d72c8 100644 --- a/packages/junior/src/chat/conversations/store.ts +++ b/packages/junior/src/chat/conversations/store.ts @@ -65,7 +65,7 @@ export interface Conversation { /** Persist and read durable conversation metadata for reporting surfaces. */ export interface ConversationStore { get(args: { conversationId: string }): Promise; - /** Read persisted visibility for one destination. Missing rows fail closed. */ + /** Read confirmed public/private visibility for one destination. */ getDestinationVisibility(args: { provider: string; providerDestinationId: string; @@ -82,7 +82,7 @@ export interface ConversationStore { /** Source normalized to a stable session locator; set-once when absent. */ sessionSource?: Source; title?: string; - /** Source-confirmed visibility from the current event's signal only. */ + /** Confirmed destination visibility; omit when unavailable. */ visibility?: ConversationPrivacy; }): Promise; /** @@ -104,7 +104,7 @@ export interface ConversationStore { source?: ConversationSource; title?: string; updatedAtMs: number; - /** Source-confirmed visibility from the current event's signal only. */ + /** Confirmed destination visibility; omit when unavailable. */ visibility?: ConversationPrivacy; }): Promise; listByActivity(args?: { diff --git a/packages/junior/src/chat/pi/client.ts b/packages/junior/src/chat/pi/client.ts index eea1c5965..9f3e78009 100644 --- a/packages/junior/src/chat/pi/client.ts +++ b/packages/junior/src/chat/pi/client.ts @@ -43,7 +43,6 @@ import { import { toOptionalTrimmed } from "@/chat/optional-string"; import { getCurrentConversationPrivacy, - resolveConversationPrivacy, toCanonicalInputMessage, toCanonicalOutputMessage, toGenAiMessagesTraceAttributes, @@ -149,22 +148,7 @@ export async function completeText(params: { : toOptionalTrimmed(process.env.VERCEL_OIDC_TOKEN) ? "oidc" : "api_key"; - // Identifier metadata can only narrow toward private; the turn-scoped - // privacy context carries the source-confirmed classification. - const privacy = - resolveConversationPrivacy({ - channelId: - typeof params.metadata?.channelId === "string" - ? params.metadata.channelId - : undefined, - conversationId: - typeof params.metadata?.conversationId === "string" - ? params.metadata.conversationId - : typeof params.metadata?.threadId === "string" - ? params.metadata.threadId - : undefined, - }) ?? getCurrentConversationPrivacy(); - const effectivePrivacy = privacy ?? "private"; + const effectivePrivacy = getCurrentConversationPrivacy() ?? "private"; const messageAttributeMode = params.messageAttributeMode ?? (effectivePrivacy === "public" ? "content" : "metadata"); diff --git a/packages/junior/src/chat/runtime/reply-executor.ts b/packages/junior/src/chat/runtime/reply-executor.ts index 30581c54b..b18998550 100644 --- a/packages/junior/src/chat/runtime/reply-executor.ts +++ b/packages/junior/src/chat/runtime/reply-executor.ts @@ -130,6 +130,7 @@ import { recordAgentTurnSessionSummary, } from "@/chat/state/turn-session"; import { completeDeliveredTurn } from "@/chat/services/turn-session-record"; +import { resolveDestinationVisibility } from "@/chat/conversations/destination-visibility"; import { getConversationStore } from "@/chat/db"; import { contextProvenance, @@ -473,24 +474,25 @@ export function createReplyToThread(deps: ReplyExecutorDeps) { !options.execution && channelId ? await resolveChannelName(thread) : undefined; + const destination = requireSlackDestination( + options.destination, + "Slack reply execution", + ); const slackChannelType = resolveSlackChannelTypeFromMessage(message); const slackConversation = resolveSlackConversationContext({ channelId, channelName, channelType: slackChannelType, }); - // Source-confirmed visibility for destination persistence; undefined when - // the event carries no channel_type so existing visibility is not changed. - const destinationVisibility = - options.execution?.destinationVisibility ?? - conversationVisibilityFromSlackChannelType(slackChannelType); + const destinationVisibility = await resolveDestinationVisibility({ + destination, + visibility: + options.execution?.destinationVisibility ?? + conversationVisibilityFromSlackChannelType(slackChannelType), + }); const threadTs = getThreadTs(threadId); const assistantThreadContext = getAssistantThreadContext(message); const messageTs = getMessageTs(message); - const destination = requireSlackDestination( - options.destination, - "Slack reply execution", - ); const teamId = destination.teamId; const source = options.execution?.source ?? @@ -1357,7 +1359,7 @@ export function createReplyToThread(deps: ReplyExecutorDeps) { slackConversation, source, destination, - destinationVisibility, + ...(destinationVisibility ? { destinationVisibility } : {}), surface: options.execution?.surface ?? "slack", dispatch: options.execution?.dispatch, toolChannelId, diff --git a/packages/junior/src/chat/services/turn-session-record.ts b/packages/junior/src/chat/services/turn-session-record.ts index efe295422..7c8f82693 100644 --- a/packages/junior/src/chat/services/turn-session-record.ts +++ b/packages/junior/src/chat/services/turn-session-record.ts @@ -102,6 +102,7 @@ export async function persistRunningSessionRecord(args: { channelName?: string; conversationId: string; destination?: Destination; + destinationVisibility?: ConversationPrivacy; dispatchId?: string; source?: Source; sessionId: string; @@ -136,6 +137,7 @@ export async function persistRunningSessionRecord(args: { ...((args.destination ?? latestSessionRecord?.destination) ? { destination: args.destination ?? latestSessionRecord?.destination } : {}), + destinationVisibility: args.destinationVisibility, ...((args.dispatchId ?? latestSessionRecord?.dispatchId) ? { dispatchId: args.dispatchId ?? latestSessionRecord?.dispatchId } : {}), @@ -205,7 +207,7 @@ export async function persistCompletedSessionRecord(args: { dispatchId?: string; dispatchOutcome?: AgentDispatchOutcome; errorMessage?: string; - /** Source-confirmed destination visibility from the current event's signal. */ + /** Confirmed visibility; omit when unavailable to preserve canonical metadata. */ destinationVisibility?: ConversationPrivacy; /** Provider-owned identifier returned after visible delivery is accepted. */ resultMessageId?: string; @@ -371,6 +373,7 @@ export async function persistAuthPauseSessionRecord(args: { currentDurationMs?: number; currentUsage?: AgentTurnUsage; destination?: Destination; + destinationVisibility?: ConversationPrivacy; dispatchId?: string; source?: Source; messages: PiMessage[]; @@ -410,6 +413,7 @@ export async function persistAuthPauseSessionRecord(args: { ...((args.destination ?? latestSessionRecord?.destination) ? { destination: args.destination ?? latestSessionRecord?.destination } : {}), + destinationVisibility: args.destinationVisibility, ...((args.dispatchId ?? latestSessionRecord?.dispatchId) ? { dispatchId: args.dispatchId ?? latestSessionRecord?.dispatchId } : {}), @@ -465,6 +469,7 @@ interface ContinuationRecordInput { currentDurationMs?: number; currentUsage?: AgentTurnUsage; destination?: Destination; + destinationVisibility?: ConversationPrivacy; dispatchId?: string; source?: Source; messages: PiMessage[]; @@ -519,6 +524,7 @@ export async function persistContinuationSessionRecord( destination: args.destination ?? latestSessionRecord?.destination, } : {}), + destinationVisibility: args.destinationVisibility, ...((args.dispatchId ?? latestSessionRecord?.dispatchId) ? { dispatchId: args.dispatchId ?? latestSessionRecord?.dispatchId } : {}), @@ -565,6 +571,7 @@ export async function persistContinuationSessionRecord( ...((args.destination ?? latestSessionRecord?.destination) ? { destination: args.destination ?? latestSessionRecord?.destination } : {}), + destinationVisibility: args.destinationVisibility, ...((args.dispatchId ?? latestSessionRecord?.dispatchId) ? { dispatchId: args.dispatchId ?? latestSessionRecord?.dispatchId } : {}), diff --git a/packages/junior/src/chat/slack/conversation-context.ts b/packages/junior/src/chat/slack/conversation-context.ts index 6a5ae8717..7fad7c250 100644 --- a/packages/junior/src/chat/slack/conversation-context.ts +++ b/packages/junior/src/chat/slack/conversation-context.ts @@ -19,11 +19,7 @@ export type SlackConversationVisibility = "public" | "private"; export interface SlackConversationContext { type: SlackConversationType; name?: string; - /** - * Visibility proven by a source signal (`channel_type`) or narrowed toward - * private by a `D`/`G` id prefix. Undefined for `C`-prefixed conversations - * without a signal, which may be public or private. - */ + /** Visibility proven by source or persisted metadata. */ visibility?: SlackConversationVisibility; } @@ -86,8 +82,8 @@ function toSlackEventChannelType( /** * Map Slack's Events API channel_type to a source-confirmed visibility. * - * This is the only mapping allowed to classify a Slack conversation public; - * channel-id prefixes may only narrow classification toward private. + * Slack's event metadata is the only visibility authority; channel IDs never + * classify visibility. */ export function conversationVisibilityFromSlackChannelType( channelType: SlackEventChannelType | undefined, @@ -96,18 +92,6 @@ export function conversationVisibilityFromSlackChannelType( return channelType === "channel" ? "public" : "private"; } -// Narrows toward private only; never asserts public. -function visibilityFromChannelIdPrefix( - channelId: string | undefined, -): SlackConversationVisibility | undefined { - const normalized = normalizeSlackConversationId(channelId); - if (!normalized) return undefined; - if (normalized.startsWith("D") || normalized.startsWith("G")) { - return "private"; - } - return undefined; -} - /** Resolve Slack's raw event channel type from a Chat SDK message-like object. */ export function resolveSlackChannelTypeFromMessage( message: unknown, @@ -135,9 +119,9 @@ export function resolveSlackConversationContext(input: { if (!type) return undefined; const name = normalizeConversationName(type, input.channelName); - const visibility = - conversationVisibilityFromSlackChannelType(input.channelType) ?? - visibilityFromChannelIdPrefix(input.channelId); + const visibility = conversationVisibilityFromSlackChannelType( + input.channelType, + ); return { type, @@ -150,12 +134,24 @@ export function resolveSlackConversationContext(input: { export function resolveSlackConversationContextFromThreadId(input: { threadId?: string; channelName?: string; + visibility?: SlackConversationVisibility; }): SlackConversationContext | undefined { const slackThread = parseSlackThreadId(input.threadId); - return resolveSlackConversationContext({ + const context = resolveSlackConversationContext({ channelId: slackThread?.channelId, channelName: input.channelName, }); + if (!context || !input.visibility) { + return context; + } + return { + ...context, + type: + input.visibility === "private" && context.type === "public_channel" + ? "private_channel" + : context.type, + visibility: input.visibility, + }; } /** Render a human label for a privacy-preserving Slack conversation type. */ diff --git a/packages/junior/src/chat/state/turn-session.ts b/packages/junior/src/chat/state/turn-session.ts index 792fc9eeb..d248c3b8e 100644 --- a/packages/junior/src/chat/state/turn-session.ts +++ b/packages/junior/src/chat/state/turn-session.ts @@ -256,7 +256,7 @@ async function appendAgentTurnSessionSummary( /** Store run summary metadata in the configured conversation store. */ async function recordConversationActivityMetadata(args: { conversationStore?: ConversationStore; - /** Source-confirmed destination visibility from the current event's signal. */ + /** Confirmed destination visibility; omit when unavailable. */ destinationVisibility?: ConversationPrivacy; nowMs: number; summary: AgentTurnSessionSummary; @@ -554,7 +554,7 @@ function buildStoredRecord(args: { async function setStoredRecord(args: { conversationStore?: ConversationStore; - /** Source-confirmed destination visibility from the current event's signal. */ + /** Confirmed destination visibility; omit when unavailable. */ destinationVisibility?: ConversationPrivacy; piMessages: PiMessage[]; piMessageProvenance: ConversationMessageProvenance[]; @@ -690,7 +690,7 @@ export async function upsertAgentTurnSessionRecord(args: { dispatchId?: string; dispatchOutcome?: AgentDispatchOutcome; resultMessageId?: string; - /** Source-confirmed destination visibility from the current event's signal. */ + /** Confirmed destination visibility; omit when unavailable. */ destinationVisibility?: ConversationPrivacy; source?: Source; lastProgressAtMs?: number; @@ -881,11 +881,7 @@ export async function recordAgentTurnSessionSummary(args: { dispatchId?: string; dispatchOutcome?: AgentDispatchOutcome; resultMessageId?: string; - /** - * Source-confirmed destination visibility from the current event's signal - * (Slack `channel_type`). Leave unset when no live signal exists so an - * existing destination visibility is not overwritten. - */ + /** Confirmed destination visibility; omit when unavailable. */ destinationVisibility?: ConversationPrivacy; source?: Source; lastProgressAtMs?: number; diff --git a/packages/junior/src/chat/tools/query-conversation-events.ts b/packages/junior/src/chat/tools/query-conversation-events.ts index 81ac762bf..b7144a4b2 100644 --- a/packages/junior/src/chat/tools/query-conversation-events.ts +++ b/packages/junior/src/chat/tools/query-conversation-events.ts @@ -335,13 +335,9 @@ function assertCanQueryConversationEvents(args: { return; } - const publicPayloadAllowed = canExposeConversationPayload({ - conversationId: targetRootConversationId ?? targetConversationId, - ...(targetVisibility ? { visibility: targetVisibility } : {}), - ...(targetDestination?.platform === "slack" - ? { channelId: targetDestination.channelId } - : {}), - }); + const publicPayloadAllowed = canExposeConversationPayload( + targetVisibility ? { visibility: targetVisibility } : {}, + ); if (!publicPayloadAllowed) { throw new ToolInputError( `Conversation events are not accessible: ${targetConversationId}`, diff --git a/packages/junior/tests/component/conversation-sql-store.test.ts b/packages/junior/tests/component/conversation-sql-store.test.ts index b1053a658..df719da5e 100644 --- a/packages/junior/tests/component/conversation-sql-store.test.ts +++ b/packages/junior/tests/component/conversation-sql-store.test.ts @@ -600,15 +600,15 @@ describe("conversation SQL store", () => { } }); - it("defaults unsigned Slack destinations to private", async () => { + it("fails closed for unsigned Slack destinations without confirming visibility", async () => { const fixture = await createLocalJuniorSqlFixture(); try { const store = createSqlStore(fixture.sql); await migrateSchema(fixture.sql); - // A write without a live source signal fails closed to private even - // though the channel id is C-prefixed. + // A write without a live source signal remains unknown even though the + // channel id is C-prefixed. Conversation reads still fail closed. await store.recordActivity({ conversationId: CONVERSATION_ID, destination: inboundMessage("unsigned").destination, @@ -624,7 +624,7 @@ describe("conversation SQL store", () => { providerTenantId: "T123", providerDestinationId: "C123", }), - ).resolves.toBe("private"); + ).resolves.toBeUndefined(); } finally { await fixture.close(); } diff --git a/packages/junior/tests/component/runtime/agent-run-agent-continue.test.ts b/packages/junior/tests/component/runtime/agent-run-agent-continue.test.ts index 77bb1e8e3..0d2ddb82d 100644 --- a/packages/junior/tests/component/runtime/agent-run-agent-continue.test.ts +++ b/packages/junior/tests/component/runtime/agent-run-agent-continue.test.ts @@ -326,6 +326,7 @@ describe("agent continuation composition", () => { turnId: "turn-1", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: TEST_ACTOR, @@ -383,6 +384,7 @@ describe("agent continuation composition", () => { turnId: "turn-timeout-cap", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: TEST_ACTOR, @@ -417,6 +419,7 @@ describe("agent continuation composition", () => { turnId: "turn-short-deadline", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: TEST_ACTOR, @@ -448,6 +451,7 @@ describe("agent continuation composition", () => { omittedImageAttachmentCount: 1, }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: TEST_ACTOR, @@ -486,6 +490,7 @@ describe("agent continuation composition", () => { turnId: "turn-hung", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: TEST_ACTOR, @@ -527,6 +532,7 @@ describe("agent continuation composition", () => { turnId: "turn-retry", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: TEST_ACTOR, diff --git a/packages/junior/tests/component/runtime/agent-run-mcp-progressive-loading.test.ts b/packages/junior/tests/component/runtime/agent-run-mcp-progressive-loading.test.ts index 8b95b7e3a..9010e3ef5 100644 --- a/packages/junior/tests/component/runtime/agent-run-mcp-progressive-loading.test.ts +++ b/packages/junior/tests/component/runtime/agent-run-mcp-progressive-loading.test.ts @@ -680,6 +680,7 @@ function makeAgentRunRequest( ...(overrides.input ?? {}), }, routing: { + destinationVisibility: "private", credentialContext: { actor: { type: "user" as const, userId: "U123" }, }, diff --git a/packages/junior/tests/component/runtime/agent-run-provider-retry.test.ts b/packages/junior/tests/component/runtime/agent-run-provider-retry.test.ts index 402a18cbd..f50bf69e4 100644 --- a/packages/junior/tests/component/runtime/agent-run-provider-retry.test.ts +++ b/packages/junior/tests/component/runtime/agent-run-provider-retry.test.ts @@ -665,6 +665,7 @@ describe("agent run continuation", () => { messageText: "Make a large generated-file edit.", }, routing: { + destinationVisibility: "private", source: TEST_SOURCE, destination: TEST_DESTINATION, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -716,6 +717,7 @@ describe("agent run continuation", () => { piMessages: priorMessages, }, routing: { + destinationVisibility: "private", source: TEST_SOURCE, destination: TEST_DESTINATION, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -739,6 +741,7 @@ describe("agent run continuation", () => { turnId: "turn-1", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -817,6 +820,7 @@ describe("agent run continuation", () => { turnId: "turn-cancelled-backoff", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -853,6 +857,7 @@ describe("agent run continuation", () => { turnId: "turn-cancelled-provider", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -904,6 +909,7 @@ describe("agent run continuation", () => { turnId: "turn-steering", input: { messageText: "help me", piMessages: priorMessages }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -966,6 +972,7 @@ describe("agent run continuation", () => { turnId: "turn-steering-delivery", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -1004,6 +1011,7 @@ describe("agent run continuation", () => { turnId: "turn-terminal-delivery-yield", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -1027,6 +1035,7 @@ describe("agent run continuation", () => { turnId: "turn-delivery-steering-yield", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -1071,6 +1080,7 @@ describe("agent run continuation", () => { turnId: "turn-yield", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -1186,6 +1196,7 @@ describe("agent run continuation", () => { turnId, input: { messageText: "Delete preview-42 after I confirm." }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -1211,6 +1222,7 @@ describe("agent run continuation", () => { turnId: "turn-yield-steering", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", actor: { platform: "slack", teamId: "T123", userId: "U123" }, destination: TEST_DESTINATION, source: TEST_SOURCE, @@ -1267,6 +1279,7 @@ describe("agent run continuation", () => { turnId: "turn-yield-persist-failure", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -1300,6 +1313,7 @@ describe("agent run continuation", () => { turnId: "turn-tool-activity", input: { messageText: "run the tool" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -1339,6 +1353,7 @@ describe("agent run continuation", () => { turnId: sessionId, input: { messageText: "help me", piMessages: [checkpointedPrompt] }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, @@ -1370,6 +1385,7 @@ describe("agent run continuation", () => { turnId: "turn-steering-failure", input: { messageText: "help me" }, routing: { + destinationVisibility: "private", destination: TEST_DESTINATION, source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, diff --git a/packages/junior/tests/component/services/turn-session-record.test.ts b/packages/junior/tests/component/services/turn-session-record.test.ts index 8b7f2d914..dbf1adb23 100644 --- a/packages/junior/tests/component/services/turn-session-record.test.ts +++ b/packages/junior/tests/component/services/turn-session-record.test.ts @@ -159,15 +159,25 @@ describe("persistAuthPauseSessionRecord", () => { }); }); - it("records Slack turn activity in SQL conversation metadata", async () => { + it("records Slack turn activity without replacing confirmed visibility", async () => { vi.useFakeTimers({ now: 10_000 }); const { upsertAgentTurnSessionRecord } = await import("@/chat/state/turn-session"); const { getConversationStore } = await import("@/chat/db"); + const { resolveDestinationVisibility } = + await import("@/chat/conversations/destination-visibility"); const { appendInboundMessage } = await import("@/chat/task-execution/store"); + const conversationStore = getConversationStore(); try { + await conversationStore.recordActivity({ + conversationId: "slack:C123:turn-activity", + destination: SLACK_DESTINATION, + nowMs: 8_000, + source: "slack", + visibility: "public", + }); await appendInboundMessage({ message: { conversationId: "slack:C123:turn-activity", @@ -198,7 +208,7 @@ describe("persistAuthPauseSessionRecord", () => { }); await expect( - getConversationStore().get({ + conversationStore.get({ conversationId: "slack:C123:turn-activity", }), ).resolves.toMatchObject({ @@ -208,7 +218,11 @@ describe("persistAuthPauseSessionRecord", () => { lastActivityAtMs: 10_000, sessionSource: SLACK_SOURCE, source: "slack", + visibility: "public", }); + await expect( + resolveDestinationVisibility({ destination: SLACK_DESTINATION }), + ).resolves.toBe("public"); } finally { vi.useRealTimers(); } diff --git a/packages/junior/tests/integration/api/conversations/list.test.ts b/packages/junior/tests/integration/api/conversations/list.test.ts index d442d5172..0dff14fd1 100644 --- a/packages/junior/tests/integration/api/conversations/list.test.ts +++ b/packages/junior/tests/integration/api/conversations/list.test.ts @@ -1,6 +1,7 @@ import { describe, expect, test } from "vitest"; import { eq } from "drizzle-orm"; import { createJuniorApi } from "@/api"; +import { readConversationAccessFromSql } from "@/api/conversations/access"; import { readConversationFeedFromSql, readConversationRecordFromSql, @@ -11,6 +12,7 @@ import { createSqlStore } from "@/chat/conversations/sql/store"; import { juniorConversationEvents, juniorConversations, + juniorDestinations, juniorIdentities, juniorUsers, } from "@/db/schema"; @@ -73,7 +75,7 @@ describe("conversation list API", () => { }); await store.recordActivity({ - conversationId: "slack:D123:private-source-link", + conversationId: "slack:D123:1700000000.000200", destination: { platform: "slack", teamId: "T123", @@ -94,15 +96,57 @@ describe("conversation list API", () => { await readConversationFeedFromSql() ).conversations.find( (conversation) => - conversation.conversationId === "slack:D123:private-source-link", + conversation.conversationId === "slack:D123:1700000000.000200", ); expect(privateSummary).toBeDefined(); expect(privateSummary).not.toHaveProperty("sourceUrl"); + expect(privateSummary).toMatchObject({ + channelName: "Direct Message", + displayTitle: "Direct Message", + }); } finally { await fixture.close(); } }); + test.each(["direct", "unknown"] as const)( + "treats persisted %s visibility as private", + async (visibility) => { + const fixture = createConfiguredJuniorSqlFixture(); + const store = createSqlStore(fixture.sql); + const conversationId = `slack:C123:${visibility}-visibility`; + const channelId = `C-${visibility}`; + try { + await migrateSchema(fixture.sql); + await store.recordActivity({ + conversationId, + destination: { + platform: "slack", + teamId: "T123", + channelId, + }, + nowMs: 1_000, + source: "slack", + visibility: "private", + }); + const db = fixture.sql.db(); + await db + .update(juniorDestinations) + .set({ visibility }) + .where(eq(juniorDestinations.providerDestinationId, channelId)); + const access = await readConversationAccessFromSql(db, [ + conversationId, + ]); + expect(access.get(conversationId)).toMatchObject({ + canViewPrivateContent: false, + visibility, + }); + } finally { + await fixture.close(); + } + }, + ); + test("uses canonical and fallback actor names with provider identity fields", async () => { const fixture = createConfiguredJuniorSqlFixture(); const store = createSqlStore(fixture.sql); diff --git a/packages/junior/tests/integration/heartbeat.test.ts b/packages/junior/tests/integration/heartbeat.test.ts index 4dca16c26..b6eb19259 100644 --- a/packages/junior/tests/integration/heartbeat.test.ts +++ b/packages/junior/tests/integration/heartbeat.test.ts @@ -688,12 +688,12 @@ describe("plugin heartbeat", () => { "Post a digest. Summarize the latest state.", ); expect(dispatchRecord?.destination).toEqual(SLACK_DESTINATION); + expect(dispatchRecord?.destinationVisibility).toBe("public"); expect(dispatchRecord?.source).toEqual( createSlackSource({ teamId: "T123", channelId: "C123", - - type: "priv", + type: "pub", }), ); expect(dispatchRecord?.metadata).toMatchObject({ diff --git a/packages/junior/tests/integration/mcp-oauth-callback.test.ts b/packages/junior/tests/integration/mcp-oauth-callback.test.ts index 676218c8c..a9521d7ee 100644 --- a/packages/junior/tests/integration/mcp-oauth-callback.test.ts +++ b/packages/junior/tests/integration/mcp-oauth-callback.test.ts @@ -178,6 +178,7 @@ async function createAwaitingMcpTurnRecord(args: { sliceId: 2, state: "awaiting_resume", destination: SLACK_DESTINATION, + destinationVisibility: "public", ...(args.includeSource === false ? {} : { source: args.source ?? slackSource(args.threadTs) }), diff --git a/packages/junior/tests/integration/oauth-callback.test.ts b/packages/junior/tests/integration/oauth-callback.test.ts index 574b0c1ed..da691a0d7 100644 --- a/packages/junior/tests/integration/oauth-callback.test.ts +++ b/packages/junior/tests/integration/oauth-callback.test.ts @@ -322,6 +322,7 @@ describe("oauth callback integration", () => { sliceId: 2, state: "awaiting_resume", destination: SLACK_DESTINATION, + destinationVisibility: "public", source: storedSource, piMessages: [ { diff --git a/packages/junior/tests/unit/logging/conversation-visibility.test.ts b/packages/junior/tests/unit/logging/conversation-visibility.test.ts new file mode 100644 index 000000000..2167ce7e9 --- /dev/null +++ b/packages/junior/tests/unit/logging/conversation-visibility.test.ts @@ -0,0 +1,43 @@ +import { describe, expect, it } from "vitest"; +import { resolveConversationPrivacy } from "@/chat/conversation-privacy"; +import { registerLogRecordSink, type EmittedLogRecord } from "@/chat/logging"; + +describe("conversation visibility logging", () => { + it.each(["public", "private", "direct", "unknown"] as const)( + "does not warn for confirmed %s visibility", + (visibility) => { + const records: EmittedLogRecord[] = []; + const unregister = registerLogRecordSink((record) => { + records.push(record); + }); + + try { + resolveConversationPrivacy({ visibility }); + } finally { + unregister(); + } + + expect(records).toEqual([]); + }, + ); + + it("warns when visibility is missing", () => { + const records: EmittedLogRecord[] = []; + const unregister = registerLogRecordSink((record) => { + records.push(record); + }); + + try { + resolveConversationPrivacy({}); + } finally { + unregister(); + } + + expect(records).toEqual([ + expect.objectContaining({ + eventName: "conversation.visibility.defaulted", + level: "warn", + }), + ]); + }); +}); diff --git a/packages/junior/tests/unit/privacy/conversation-privacy.test.ts b/packages/junior/tests/unit/privacy/conversation-privacy.test.ts index daea48207..f0a60f59b 100644 --- a/packages/junior/tests/unit/privacy/conversation-privacy.test.ts +++ b/packages/junior/tests/unit/privacy/conversation-privacy.test.ts @@ -8,45 +8,24 @@ import { } from "@/chat/conversation-privacy"; describe("conversation privacy classification", () => { - it("never classifies a conversation public from a C-prefixed id alone", () => { - expect(resolveConversationPrivacy({ channelId: "C123" })).toBeUndefined(); - expect( - resolveConversationPrivacy({ conversationId: "slack:C123:1712345.0001" }), - ).toBeUndefined(); - expect(canExposeConversationPayload({ channelId: "C123" })).toBe(false); - }); + it.each(["public", "private"] as const)( + "uses confirmed %s visibility", + (visibility) => { + expect(resolveConversationPrivacy({ visibility })).toBe(visibility); + }, + ); - it("classifies public only from an explicit visibility signal", () => { - expect( - resolveConversationPrivacy({ channelId: "C123", visibility: "public" }), - ).toBe("public"); - expect( - resolveConversationPrivacy({ - conversationId: "slack:C123:details-only", - visibility: "public", - }), - ).toBe("public"); - // Slack reported channel_type group despite the C prefix. - expect( - resolveConversationPrivacy({ channelId: "C123", visibility: "private" }), - ).toBe("private"); - }); + it.each(["direct", "unknown"] as const)( + "treats confirmed %s visibility as private", + (visibility) => { + expect(resolveConversationPrivacy({ visibility })).toBe("private"); + expect(canExposeConversationPayload({ visibility })).toBe(false); + }, + ); - it("narrows toward private from D/G prefixes even against a public claim", () => { - expect(resolveConversationPrivacy({ channelId: "D123" })).toBe("private"); - expect(resolveConversationPrivacy({ channelId: "G123" })).toBe("private"); - expect( - resolveConversationPrivacy({ channelId: "D123", visibility: "public" }), - ).toBe("private"); - }); - - it("classifies non-Slack conversations private", () => { - expect( - resolveConversationPrivacy({ conversationId: "local:workspace:run-1" }), - ).toBe("private"); - expect( - resolveConversationPrivacy({ conversationId: "agent-dispatch:run-2" }), - ).toBe("private"); + it("defaults missing visibility to private", () => { + expect(resolveConversationPrivacy({})).toBe("private"); + expect(canExposeConversationPayload({})).toBe(false); }); }); diff --git a/packages/junior/tests/unit/slack/conversation-context.test.ts b/packages/junior/tests/unit/slack/conversation-context.test.ts index b640b5436..eeff9ef80 100644 --- a/packages/junior/tests/unit/slack/conversation-context.test.ts +++ b/packages/junior/tests/unit/slack/conversation-context.test.ts @@ -5,6 +5,7 @@ import { formatSlackConversationTypeLabel, resolveSlackChannelTypeFromMessage, resolveSlackConversationContext, + resolveSlackConversationContextFromThreadId, } from "@/chat/slack/conversation-context"; describe("Slack conversation prompt context", () => { @@ -67,7 +68,6 @@ describe("Slack conversation prompt context", () => { ).toEqual({ type: "private_channel_or_group_dm", name: "#maybe-private", - visibility: "private", }); }); @@ -150,4 +150,19 @@ describe("Slack conversation prompt context", () => { }), ).toBe("Public Channel"); }); + + it("uses confirmed persisted visibility when resolving stored threads", () => { + expect( + resolveSlackConversationContextFromThreadId({ + threadId: "slack:C123:1700000000.000100", + visibility: "private", + }), + ).toMatchObject({ type: "private_channel", visibility: "private" }); + expect( + resolveSlackConversationContextFromThreadId({ + threadId: "slack:D123:1700000000.000100", + visibility: "private", + }), + ).toMatchObject({ type: "direct_message", visibility: "private" }); + }); });