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 src/core/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -557,6 +557,8 @@ export interface DaemonSession {
cardPatchInFlight?: boolean; // true while a card PATCH is in-flight
pendingCardJson?: string; // queued card JSON — flushed when in-flight PATCH completes (latest wins)
pendingCardId?: string; // card message_id captured at schedule time — prevents stale reads when streamCardId changes between schedule and flush
pendingCardUserInitiated?: boolean; // latest queued PATCH came from an explicit card action; failures are surfaced at warn level
lastStreamingCardPatchWarnAt?: number; // in-memory warning throttle for repeated user-visible PATCH failures
frozenCards?: Map<string, FrozenCard>; // nonce → FrozenCard (historical cards' cached state for toggle)
/** Wait Mode / HTTP Sync integration: pending Promise handlers for synchronous
* webhook triggers waiting for a response in this session. Key is turnId. */
Expand Down
56 changes: 50 additions & 6 deletions src/core/worker-pool.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2895,9 +2895,9 @@ export function parkStreamCard(ds: DaemonSession): void {
* messageId is the live `streamCardId` again, and recalling it would delete
* the only card the user can see.
*/
export function recallFrozenCards(ds: DaemonSession): void {
export function recallFrozenCards(ds: DaemonSession): string[] {
if (!ds.frozenCards) ds.frozenCards = loadFrozenCards(ds.session.sessionId);
if (ds.frozenCards.size === 0) return;
if (ds.frozenCards.size === 0) return [];
const activeId = ds.streamCardId && ds.streamCardId !== CARD_POSTING_SENTINEL
? ds.streamCardId
: undefined;
Expand All @@ -2913,12 +2913,13 @@ export function recallFrozenCards(ds: DaemonSession): void {
targets.push(fc.messageId);
ds.frozenCards.delete(nonce);
}
if (targets.length === 0) return;
if (targets.length === 0) return [];
saveFrozenCards(ds.session.sessionId, ds.frozenCards);
for (const messageId of targets) {
deleteMessage(ds.larkAppId, messageId).catch(() => { /* best-effort */ });
}
logger.info(`[${tag(ds)}] Recalled ${targets.length} previous streaming card(s)`);
return targets;
}

/** A streaming-card id is only meaningful once a real Lark message id has
Expand Down Expand Up @@ -3991,6 +3992,7 @@ async function postTurnStartingStatusCard(
export async function postFreshStreamingCard(
ds: DaemonSession,
sessionReply: (rootId: string, content: string, msgType?: string, larkAppId?: string, turnId?: string) => Promise<string>,
opts?: { retireMessageId?: string },
): Promise<boolean> {
if (isDocNativeSession(ds)) return false;
if (!workerHasInitialized(ds)) return false;
Expand Down Expand Up @@ -4088,7 +4090,19 @@ export async function postFreshStreamingCard(
ds.parkedStreamCardNonce = undefined;
const predecessorIds = snapshotStreamingCardPredecessorIds(ds, messageId);
persistStreamCardState(ds);
recallFrozenCards(ds);
const recalledIds = recallFrozenCards(ds);
const retireMessageId = opts?.retireMessageId;
if (retireMessageId && retireMessageId !== messageId && !recalledIds.includes(retireMessageId)) {
if (!ds.frozenCards) ds.frozenCards = loadFrozenCards(ds.session.sessionId);
let removedCachedCard = false;
for (const [frozenNonce, frozen] of ds.frozenCards) {
if (frozen.messageId !== retireMessageId) continue;
ds.frozenCards.delete(frozenNonce);
removedCachedCard = true;
}
if (removedCachedCard) saveFrozenCards(ds.session.sessionId, ds.frozenCards);
void deleteMessage(appIdAtPost, retireMessageId).catch(() => { /* best-effort legacy-card cleanup */ });
}
flushPendingLocalCliOpenReadinessPatch(ds);
flushPendingRiffUrlPatch(ds);
flushPendingActiveRuntimePatch(ds);
Expand Down Expand Up @@ -4470,7 +4484,12 @@ export async function deliverEphemeralOrReply(
* any previously queued value — only the latest state matters). Returns
* whether the PATCH was accepted for immediate or queued delivery.
*/
export function scheduleCardPatch(ds: DaemonSession, cardJson: string, turnId?: string): boolean {
export function scheduleCardPatch(
ds: DaemonSession,
cardJson: string,
turnId?: string,
opts?: { userInitiated?: boolean },
): boolean {
// Defense-in-depth transport gate: a no-transport session (apiOnly bot or HTTP
// virtual chat) has no real Feishu card to PATCH. Callers already suppress via
// managedAuxUiSuppressed, but guarding the flush entry too means a stray direct
Expand All @@ -4487,6 +4506,10 @@ export function scheduleCardPatch(ds: DaemonSession, cardJson: string, turnId?:
// Capture the card ID now — by the time flushCardPatch runs, ds.streamCardId
// may have been overwritten by a new turn's card (CARD_POSTING_SENTINEL).
ds.pendingCardId = cardId;
// Preserve an explicit user action even if a later automatic screen render
// coalesces into the same latest-wins slot before it can be delivered.
ds.pendingCardUserInitiated =
ds.pendingCardUserInitiated === true || opts?.userInitiated === true;
if (ds.cardPatchInFlight) return true;
flushCardPatch(ds);
return true;
Expand All @@ -4495,13 +4518,16 @@ export function scheduleCardPatch(ds: DaemonSession, cardJson: string, turnId?:
function flushCardPatch(ds: DaemonSession): void {
const json = ds.pendingCardJson;
const cardId = ds.pendingCardId;
const userInitiated = ds.pendingCardUserInitiated === true;
if (!json || !cardId || cardId === CARD_POSTING_SENTINEL) {
ds.pendingCardJson = undefined;
ds.pendingCardId = undefined;
ds.pendingCardUserInitiated = undefined;
return;
}
ds.pendingCardJson = undefined;
ds.pendingCardId = undefined;
ds.pendingCardUserInitiated = undefined;
ds.cardPatchInFlight = true;
let patchSucceeded = false;
updateMessage(ds.larkAppId, cardId, json)
Expand All @@ -4527,7 +4553,24 @@ function flushCardPatch(ds: DaemonSession): void {
}
return;
}
logger.debug(`[${tag(ds)}] Failed to update streaming card: ${err}`);
const response = typeof err === 'object' && err !== null
? (err as { response?: { status?: unknown; data?: { code?: unknown; msg?: unknown; log_id?: unknown; logId?: unknown } } }).response
: undefined;
const responseData = response?.data;
const detail = [
response?.status !== undefined ? `HTTP ${String(response.status)}` : '',
typeof responseData?.code === 'number' ? `code=${responseData.code}` : '',
typeof responseData?.msg === 'string' ? responseData.msg : '',
responseData?.log_id ?? responseData?.logId
? `log_id=${String(responseData?.log_id ?? responseData?.logId)}`
: '',
].filter(Boolean).join(' ') || (err instanceof Error ? err.message : String(err));
if (userInitiated && Date.now() - (ds.lastStreamingCardPatchWarnAt ?? 0) >= 60_000) {
ds.lastStreamingCardPatchWarnAt = Date.now();
logger.warn(`[${tag(ds)}] User-triggered streaming-card PATCH failed: ${detail}`);
} else {
logger.debug(`[${tag(ds)}] Failed to update streaming card: ${detail}`);
}
})
.finally(() => {
ds.cardPatchInFlight = false;
Expand All @@ -4541,6 +4584,7 @@ function flushCardPatch(ds: DaemonSession): void {
&& ds.pendingCardJson === json) {
ds.pendingCardJson = undefined;
ds.pendingCardId = undefined;
ds.pendingCardUserInitiated = undefined;
}
if (ds.pendingCardJson) {
flushCardPatch(ds);
Expand Down
15 changes: 13 additions & 2 deletions src/im/lark/card-builder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -957,6 +957,8 @@ function pushStreamBody(
* Quick-action buttons (Esc, ^C, Tab, Space, Enter, ←↑↓→, ½屏 ↑/↓) appear
* whenever displayMode !== 'hidden'.
*/
export const STREAMING_CARD_PATCH_VERSION = '1';

export function buildStreamingCard(
sessionId: string,
rootId: string,
Expand Down Expand Up @@ -991,7 +993,13 @@ export function buildStreamingCard(
): string {
const effectiveCliId = cliId ?? 'claude-code';
const cliName = runtimeDisplayName?.trim() || getCliDisplayName(effectiveCliId);
const actionBase = { root_id: rootId, session_id: sessionId, cli_id: effectiveCliId, ...(cardNonce ? { card_nonce: cardNonce } : {}) };
const actionBase = {
root_id: rootId,
session_id: sessionId,
cli_id: effectiveCliId,
stream_card_version: STREAMING_CARD_PATCH_VERSION,
...(cardNonce ? { card_nonce: cardNonce } : {}),
};
const displayStatus = status === 'limited' && usageLimit?.retryReady ? 'retry_ready' : status;

const elements: any[] = [];
Expand Down Expand Up @@ -1202,7 +1210,10 @@ export function buildStreamingCard(
}

const card = {
config: { wide_screen_mode: true },
// Lark's ordinary message PATCH endpoint only updates cards whose original
// and replacement payloads both opt into shared updates. Streaming cards
// are group-visible mutable UI, so this is part of their wire contract.
config: { wide_screen_mode: true, update_multi: true },
header: {
title: { tag: 'plain_text', content: `🖥️ ${cliName}${serviceTierBadge ? ` ${serviceTierBadge}` : ''} · ${plainTitle(title)} — ${streamStatusLabel(status, usageLimit, locale, silentIdle)}` },
template: STREAM_STATUS_TEMPLATE_MAP[displayStatus],
Expand Down
69 changes: 66 additions & 3 deletions src/im/lark/card-handler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ import { resolveHiddenStreamingCardButtons } from './streaming-card-buttons.js';
import { canOperate, canTalk, canRunDaemonCommand } from './event-dispatcher.js';
import { isBotAdmin } from './grant-owner.js';
import { updateMessage, deleteMessage, replyMessage, sendMessage, sendUserMessage, sendEphemeralCard, getMessageDetail, isHumanOpenId, resolveUserUnionId as defaultResolveUserUnionId } from './client.js';
import { buildSessionCard, buildStreamingCard, buildTuiPromptCard, buildTuiPromptProcessingCard, buildGrantResultCard, getCliDisplayName, truncateContent, buildConfigCard, buildConfigQuotaCard, buildConfigTextCard, CONFIG_UNSET, buildRepoSelectCard, frozenIdleLabel } from './card-builder.js';
import { buildSessionCard, buildStreamingCard, buildTuiPromptCard, buildTuiPromptProcessingCard, buildGrantResultCard, getCliDisplayName, truncateContent, buildConfigCard, buildConfigQuotaCard, buildConfigTextCard, CONFIG_UNSET, buildRepoSelectCard, frozenIdleLabel, STREAMING_CARD_PATCH_VERSION } from './card-builder.js';
import { codexServiceTierBadge } from '../../services/codex-service-tier.js';
import {
findConfigField,
Expand Down Expand Up @@ -98,7 +98,7 @@ import { buildTurnContinuePrompt } from '../../services/turn-failure-notice.js';
import { loadFrozenCards, saveFrozenCards } from '../../services/frozen-card-store.js';
import { resumeStartsFresh } from '../../services/resume-fresh-policy.js';
import { cliHasNoRawPassthroughSurface } from '../../core/passthrough-commands.js';
import { forkWorker, sendWorkerInput, sendWorkerSessionInput, killWorker, closeSession as closeWorkerPoolSession, teardownAuthoritativePersistentBackingBeforeClose, scheduleCardPatch, parkStreamCard, clearUsageLimitState, cardUsageLimit, writableTerminalLinkFor, workerHasInitialized, sessionSupportsWebTerminal, readableTerminalUrlFor, resolvePrivateCardAudience, deliverWriteLinkCard, deliverEphemeralOrReply, CARD_POSTING_SENTINEL, requestSessionRestart, isSessionTransferring, getDaemonStreamingCardUsageSnapshot, withActiveSessionKeyLock, buildStreamingCardJson, canCommitStreamingCardPublication, continuePublishedStreamingCardPinChain, idleCardLabel, dshRuntimeForSession, type WorkerSessionReplyOptions } from '../../core/worker-pool.js';
import { forkWorker, sendWorkerInput, sendWorkerSessionInput, killWorker, closeSession as closeWorkerPoolSession, teardownAuthoritativePersistentBackingBeforeClose, scheduleCardPatch, parkStreamCard, clearUsageLimitState, cardUsageLimit, writableTerminalLinkFor, workerHasInitialized, sessionSupportsWebTerminal, readableTerminalUrlFor, resolvePrivateCardAudience, deliverWriteLinkCard, deliverEphemeralOrReply, CARD_POSTING_SENTINEL, requestSessionRestart, isSessionTransferring, getDaemonStreamingCardUsageSnapshot, withActiveSessionKeyLock, buildStreamingCardJson, canCommitStreamingCardPublication, continuePublishedStreamingCardPinChain, idleCardLabel, dshRuntimeForSession, postFreshStreamingCard, type WorkerSessionReplyOptions } from '../../core/worker-pool.js';
import { reconcileResumedStreamingCard } from '../../core/resume-streaming-card.js';
import { getSessionWorkingDir, buildNewTopicCliInput, getAvailableBots, persistStreamCardState, resumeSession, rememberLastCliInput, ensureSessionWhiteboard } from '../../core/session-manager.js';
import { markInitialUserTurnPending } from '../../core/initial-user-turn.js';
Expand Down Expand Up @@ -289,6 +289,18 @@ const LEGACY_SELF_HEAL_ACTIONS = new Set(['toggle_display', 'toggle_stream', 're
// In-memory (per daemon lifetime) — a restart resets it, which at worst allows
// one re-trigger on an old card; acceptable. Capped to avoid unbounded growth.
const voicedCardIds = new Set<string>();
const legacyStreamingCardMigrationIds = new Set<string>();
const LEGACY_STREAMING_CARD_MIGRATION_LIMIT = 2_000;

function claimLegacyStreamingCardMigration(messageId: string): boolean {
if (legacyStreamingCardMigrationIds.has(messageId)) return false;
if (legacyStreamingCardMigrationIds.size >= LEGACY_STREAMING_CARD_MIGRATION_LIMIT) {
const oldest = legacyStreamingCardMigrationIds.values().next().value;
if (oldest) legacyStreamingCardMigrationIds.delete(oldest);
}
legacyStreamingCardMigrationIds.add(messageId);
return true;
}

// Instruction injected into the session when the voice button is clicked. The
// model (which still has its just-sent reply in context) condenses it into
Expand Down Expand Up @@ -3745,11 +3757,62 @@ export async function handleCardAction(data: CardActionData, deps: CardHandlerDe
return { toast: { type: 'warning', content: t('card.action.session_gone', undefined, localeForBot(larkAppId)) } };
}
const clickedNonce: string | undefined = value?.card_nonce;
const needsPatchContractMigration =
value?.stream_card_version !== STREAMING_CARD_PATCH_VERSION;
const isFrozenClick = clickedNonce && ds.streamCardNonce && clickedNonce !== ds.streamCardNonce;

const nextMode = (current: DisplayMode): DisplayMode =>
current === 'hidden' ? 'screenshot' : 'hidden';

if (needsPatchContractMigration) {
const legacyMessageId = cardMessageId
?? (ds.streamCardId !== CARD_POSTING_SENTINEL ? ds.streamCardId : undefined);
if (!legacyMessageId || !claimLegacyStreamingCardMigration(legacyMessageId)) {
return { toast: { type: 'info', content: t('toast.action_received_bg', undefined, localeForBot(ds.larkAppId)) } };
}
const current: DisplayMode = ds.displayMode ?? 'hidden';
const next = nextMode(current);
ds.displayMode = next;
persistStreamCardState(ds);
if (ds.worker || isSessionTransferring(ds)) {
sendWorkerSessionInput(ds, { type: 'set_display_mode', mode: next });
}
logger.info(`[${tag(ds)}] Display mode → ${next} (legacy card migration)`);
return {
toast: {
type: 'info',
content: t('toast.action_received_bg', undefined, localeForBot(ds.larkAppId)),
},
afterAck: async () => {
let migrated = false;
try {
migrated = await postFreshStreamingCard(
ds,
deps.sessionReply,
{ retireMessageId: legacyMessageId },
);
} catch (err) {
logger.warn(`[${tag(ds)}] Legacy streaming-card migration crashed: ${err instanceof Error ? err.message : String(err)}`);
}
if (!migrated) {
// The old card did not change, so roll back the optimistic mode
// only when no newer action has superseded it. A retry will then
// request the same visible transition instead of toggling back.
if (ds.displayMode === next) {
ds.displayMode = current;
persistStreamCardState(ds);
if (ds.worker || isSessionTransferring(ds)) {
sendWorkerSessionInput(ds, { type: 'set_display_mode', mode: current });
}
}
// A failed migration must remain retryable on the next click.
legacyStreamingCardMigrationIds.delete(legacyMessageId);
logger.warn(`[${tag(ds)}] Legacy streaming-card migration failed for ${legacyMessageId.substring(0, 12)}`);
}
},
};
}

if (isFrozenClick) {
// Historical card — toggle using cached state
if (!ds.frozenCards) ds.frozenCards = loadFrozenCards(ds.session.sessionId);
Expand Down Expand Up @@ -3896,7 +3959,7 @@ export async function handleCardAction(data: CardActionData, deps: CardHandlerDe
logger.debug(`[${tag(ds)}] Failed to migrate clicked legacy card: ${err}`),
);
try { return JSON.parse(cardJson); } catch { /* fall through */ }
} else if (!scheduleCardPatch(ds, cardJson)) {
} else if (!scheduleCardPatch(ds, cardJson, undefined, { userInitiated: true })) {
// The queue can decline when live cards are disabled for this turn or
// transport/card identity is unavailable. In that case the callback
// must carry the rebuilt card so the clicked card still updates.
Expand Down
16 changes: 16 additions & 0 deletions test/card-builder.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -917,6 +917,18 @@ describe('buildStreamingCard', () => {
it('should have wide_screen_mode config', () => {
const card = parse(buildStreamingCard(SID, ROOT, URL, TITLE, CONTENT, 'working'));
expect(card.config.wide_screen_mode).toBe(true);
expect(card.config.update_multi).toBe(true);
});

it('marks every callback action with the patchable streaming-card version', () => {
const card = parse(buildStreamingCard(
SID, ROOT, URL, TITLE, CONTENT, 'working', 'claude-code', 'screenshot',
));
const callbackActions = allActions(card).filter((action: any) => action.value?.action);
expect(callbackActions.length).toBeGreaterThan(0);
for (const action of callbackActions) {
expect(action.value.stream_card_version).toBe('1');
}
});

// ── Header / status / template color ───────────────────────────────────
Expand Down Expand Up @@ -2340,6 +2352,10 @@ describe('buildPrivateSnapshotCard', () => {
.flatMap((e: any) => e.actions ?? []);
}

it('does not opt private one-shot snapshots into shared PATCH updates', () => {
expect(build().config.update_multi).toBeUndefined();
});

it('exposes open-terminal link, get_write_link, close for non-Codex/TRAE sessions, with no patch-driven controls', () => {
const card = build({ screen: 'hello' });
const btns = allButtons(card);
Expand Down
Loading
Loading