mirror of
https://github.com/danny-avila/LibreChat.git
synced 2026-08-03 22:32:42 +00:00
* ⚡ feat: Persist HITL checkpoints only on pause (skip clean-exit writes) With `durability: 'exit'` (set by the SDK whenever a checkpointer is active) LangGraph persists ONE checkpoint at the exit boundary on EVERY run — paused or not. So a non-paused HITL turn writes a dead checkpoint whose only fate is to be pruned by deleteAgentCheckpoint: pure write+delete churn on the common path, given HITL only ever resumes an *interrupt* checkpoint. `InterruptOnlyMongoSaver` (a MongoDBSaver subclass) persists only interrupt checkpoints and discards clean-exit ones, so a non-paused turn writes nothing. How it tells them apart (verified empirically against @langchain/langgraph, not docs): when a run interrupts, the runner calls `putWrites` with the `INTERRUPT` ("__interrupt__") channel for the checkpoint it's about to create, and that write's `config.checkpoint_id` equals the `checkpoint.id` of the `put` that immediately follows. A clean exit calls `put` with no preceding interrupt `putWrites`. So we record the checkpoint id of any interrupt `putWrites` and persist a `put` only when its `checkpoint.id` was so marked. Keying on the globally-unique checkpoint id (not thread_id) keeps this correct even when two runs race on the same conversation (the job-replacement scenario). Correctness is preserved end-to-end: interrupt checkpoints + their pending writes persist exactly as before (resume unchanged); clean checkpoints were only ever written-then-pruned, so not writing them is observationally equivalent. The eager prune stays as the backstop. Tests (mongodb-memory-server): a bare put() is discarded; an interrupt-seeded checkpoint is persisted with its __interrupt__ pending write; and an end-to-end real-graph run writes 0 checkpoints on a clean completion and a resumable one on interrupt. NOTE: a non-paused turn's deleteAgentCheckpoint now finds nothing to delete (a 0-match no-op) — a follow-up can skip that call entirely once the lingering-abandoned-pause cleanup role is reassigned to the TTL + expiry sweeper. * ⚡ feat: Drop the redundant clean-path checkpoint prune With the lazy checkpointer (InterruptOnlyMongoSaver) a non-paused turn no longer writes a clean-exit checkpoint, so the post-completion prune in chatCompletion's finally had nothing left to delete. It was also already redundant: every fresh turn runs a pre-run prune (`deleteAgentCheckpoint` before `processStream`) that clears any checkpoint orphaned by a prior abandoned pause — verified empirically that a lingering interrupt checkpoint WOULD otherwise poison a fresh turn (LangGraph continues the abandoned state + re-interrupts), and that the pre-run prune is what prevents it. The Mongo TTL remains the backstop, and the resume path still prunes after a successful finalize. Removing the clean-path prune also deletes its job-replacement race surface (round-17 F21): an older run's late finally can no longer delete a newer paused run's checkpoint, because there is no longer a clean-path prune to race. Dropped the now-dead F21 predicate test. Net per non-paused HITL turn: from {pre-run prune + checkpoint write + post-run prune} down to {pre-run prune} — no write, no post-completion delete. * 🛡️ fix: Anchor any pending-write checkpoint; stale-only eviction (Codex) Broaden the lazy saver's keep-rule from "interrupt-only" to "persist any checkpoint that carries pending writes" (renamed InterruptOnlyMongoSaver → LazyMongoSaver). This makes it robust to delta-channel graphs without changing behavior for LibreChat's graph: - K1 (P1): a delta-channel graph can write a synthetic PARENT/anchor checkpoint (no __interrupt__ mark) that the interrupt checkpoint then points at, with the delta writes stored under the parent id. The old rule discarded that parent, breaking delta-state resume. Now any checkpoint that received putWrites is persisted, so the anchor parent and its writes survive and resume can walk the chain. - K3 (P2): for the same reason, clean delta-write rows are no longer orphaned — their checkpoint is persisted alongside them. (For LibreChat's standard Annotation/messages graph a clean run makes no putWrites at all — verified empirically — so the common path still writes nothing and the optimization is unchanged.) - K2 (P2): the 1024 FIFO cap could evict a valid in-flight id whose put() was just behind Mongo I/O, mis-classifying its interrupt checkpoint as a clean exit. Replaced with time-based eviction: only ids older than 5 min (a put always follows its putWrites within ms) are swept; a recent in-flight id is never dropped, and the map grows rather than evict a valid id if nothing is stale. New integration test: a checkpoint anchored by a NON-interrupt write is persisted. Full agents/HITL suites green (108). * style(checkpointer): fix import order to satisfy sort-imports CI * fix(checkpointer): don't persist failed-turn (error-only) checkpoints LazyMongoSaver anchored on ANY pending write, so a non-paused turn that errors (LangGraph records an __error__ write then a put) was persisted and, with the clean-path prune removed, lingered until the next fresh-turn prune or the Mongo TTL. Anchor only on resumable writes — INTERRUPT or a real (non-__-prefixed) state/delta channel — so error/bookkeeping-only checkpoints are discarded at the source. Addresses Codex P3. Codex P2 (delta-stub parent orphan) is not reachable: the SDK graph uses standard Annotation/MessagesAnnotation channels (no DeltaChannel), and under durability:'exit' putWrites precedes put with a parentless boundary checkpoint — probe-confirmed against @langchain/langgraph@1.4. Documented the durability:'exit' invariant the saver depends on. Tests: error-only put discarded; e2e throwing graph persists 0 checkpoints. * fix(checkpointer): drop bookkeeping-only write batches, not just the checkpoint The prior fix stopped the failed-turn CHECKPOINT from persisting, but putWrites still forwarded the __error__ batch to MongoDBSaver.putWrites — writing a row to agent_checkpoint_writes whose parent checkpoint is then discarded. With the post-run deleteThread removed, that orphan row lingered until the Mongo TTL or the conversation's next pre-run prune. putWrites now drops a non-resumable (bookkeeping-only) batch entirely instead of forwarding it. Probed against a real MongoDBSaver (mongodb-memory-server): a throwing graph now leaves 0 checkpoints AND 0 write rows (was 0 + 1 orphan), while interrupt->resume is unaffected — the __interrupt__ write is resumable so it is still forwarded. Addresses Codex P2 (round 3). Tests: error-only put leaves no checkpoint and no write row; e2e throwing graph leaves both collections empty; new e2e interrupt->resume completes with the approval value. * fix(checkpointer): un-anchor a checkpoint whose putWrites failed; freshen comments Self-review findings on the converged PR: 1. LangGraph dispatches put() concurrently with putWrites (probe-confirmed on 1.4.5), and put() still completes when putWrites rejects — so a transient Mongo failure during the interrupt write could persist a checkpoint whose __interrupt__ row is missing (an unresumable phantom pause). putWrites now deletes the write anchor on rejection (best-effort) and rethrows, so that put() discards the checkpoint instead. The pre-recorded anchor stays where it is — recording after the await would drop slow-I/O interrupts on the success path, which the same probe showed is reachable. 2. Renamed leftovers: two comments still said InterruptOnlyMongoSaver; the class is LazyMongoSaver. 3. Documented why the pre-run prune is deliberately unconditional per HITL turn (any cheaper gate can go stale across replicas and skip the prune exactly when an orphaned interrupt exists). Test: failed putWrites → subsequent put persists nothing (14/14 green). * fix(checkpointer): bookkeeping write batches follow their checkpoint's fate The round-3 rule dropped bookkeeping-only putWrites batches (__error__/ __resume__/__no_writes__) unconditionally — batch-scoped, when the decision must be checkpoint-scoped. Probe-confirmed (langgraph 1.4.5, durability:'exit'): a Send fan-out that pauses on one sibling records the completed siblings as pure __no_writes__ batches on the RETAINED interrupt checkpoint; dropping those markers makes resume re-execute the completed siblings (side effects measured twice). Addresses Codex M2 (P2). putWrites now PARKS a bookkeeping-only batch in memory until the checkpoint's fate is known: forwarded when the checkpoint is anchored (or was just persisted — put is dispatched concurrently), dropped when put discards it. Net: an errored turn still leaves nothing durable (0 checkpoints, 0 write rows), and a retained checkpoint stores byte-for-byte what a plain MongoDBSaver would. Codex M1 (__resume__ lost on re-pause) did not reproduce: the re-pause emits [__interrupt__,__resume__] as ONE batch (anchored, forwarded whole) and a second resume on a rebuilt graph replays both answers correctly — but the fate-scoped buffering now covers a lone __resume__ batch in any ordering too. Tests: bookkeeping preserved on a retained checkpoint in either arrival order; e2e Send-sibling pause/resume with side-effect counters (was {a:2,c:2} under the drop rule, now {a:1,c:1}); error-only turn still leaves both collections empty. 16/16 green.
597 lines
22 KiB
JavaScript
597 lines
22 KiB
JavaScript
/**
|
|
* Tests for job replacement detection in ResumableAgentController
|
|
*
|
|
* Tests the following fixes from PR #11462:
|
|
* 1. Job creation timestamp tracking
|
|
* 2. Stale job detection and event skipping
|
|
* 3. Response message saving before final event emission
|
|
*/
|
|
|
|
const mockLogger = {
|
|
debug: jest.fn(),
|
|
warn: jest.fn(),
|
|
error: jest.fn(),
|
|
info: jest.fn(),
|
|
};
|
|
|
|
const mockGenerationJobManager = {
|
|
createJob: jest.fn(),
|
|
getJob: jest.fn(),
|
|
emitDone: jest.fn(),
|
|
emitChunk: jest.fn(),
|
|
completeJob: jest.fn(),
|
|
updateMetadata: jest.fn(),
|
|
setContentParts: jest.fn(),
|
|
subscribe: jest.fn(),
|
|
};
|
|
|
|
const mockSaveMessage = jest.fn();
|
|
const mockDecrementPendingRequest = jest.fn();
|
|
|
|
jest.mock('@librechat/data-schemas', () => ({
|
|
logger: mockLogger,
|
|
}));
|
|
|
|
jest.mock('@librechat/api', () => ({
|
|
isEnabled: jest.fn().mockReturnValue(false),
|
|
GenerationJobManager: mockGenerationJobManager,
|
|
getReferencedQuotes: jest.fn((quotes) => {
|
|
if (!Array.isArray(quotes)) {
|
|
return null;
|
|
}
|
|
const normalized = quotes
|
|
.filter((quote) => typeof quote === 'string' && quote.trim().length > 0)
|
|
.map((quote) => quote.trim());
|
|
return normalized.length > 0 ? normalized : null;
|
|
}),
|
|
checkAndIncrementPendingRequest: jest.fn().mockResolvedValue({ allowed: true }),
|
|
decrementPendingRequest: (...args) => mockDecrementPendingRequest(...args),
|
|
getViolationInfo: jest.fn(),
|
|
sanitizeMessageForTransmit: jest.fn((msg) => msg),
|
|
sanitizeFileForTransmit: jest.fn((file) => file),
|
|
Constants: { NO_PARENT: '00000000-0000-0000-0000-000000000000' },
|
|
}));
|
|
|
|
jest.mock('~/models', () => ({
|
|
saveMessage: (...args) => mockSaveMessage(...args),
|
|
}));
|
|
|
|
describe('Job Replacement Detection', () => {
|
|
beforeEach(() => {
|
|
jest.clearAllMocks();
|
|
});
|
|
|
|
describe('Job Creation Timestamp Tracking', () => {
|
|
it('should capture createdAt when job is created', async () => {
|
|
const streamId = 'test-stream-123';
|
|
const createdAt = Date.now();
|
|
|
|
mockGenerationJobManager.createJob.mockResolvedValue({
|
|
createdAt,
|
|
readyPromise: Promise.resolve(),
|
|
abortController: new AbortController(),
|
|
emitter: { on: jest.fn() },
|
|
});
|
|
|
|
const job = await mockGenerationJobManager.createJob(streamId, 'user-123', streamId);
|
|
|
|
expect(job.createdAt).toBe(createdAt);
|
|
});
|
|
});
|
|
|
|
describe('Job Replacement Detection Logic', () => {
|
|
/**
|
|
* Simulates the job replacement detection logic from request.js
|
|
* This is extracted for unit testing since the full controller is complex
|
|
*/
|
|
const detectJobReplacement = async (streamId, originalCreatedAt) => {
|
|
const currentJob = await mockGenerationJobManager.getJob(streamId);
|
|
return !currentJob || currentJob.createdAt !== originalCreatedAt;
|
|
};
|
|
|
|
it('should detect when job was replaced (different createdAt)', async () => {
|
|
const streamId = 'test-stream-123';
|
|
const originalCreatedAt = 1000;
|
|
const newCreatedAt = 2000;
|
|
|
|
mockGenerationJobManager.getJob.mockResolvedValue({
|
|
createdAt: newCreatedAt,
|
|
});
|
|
|
|
const wasReplaced = await detectJobReplacement(streamId, originalCreatedAt);
|
|
|
|
expect(wasReplaced).toBe(true);
|
|
});
|
|
|
|
it('should detect when job was deleted', async () => {
|
|
const streamId = 'test-stream-123';
|
|
const originalCreatedAt = 1000;
|
|
|
|
mockGenerationJobManager.getJob.mockResolvedValue(null);
|
|
|
|
const wasReplaced = await detectJobReplacement(streamId, originalCreatedAt);
|
|
|
|
expect(wasReplaced).toBe(true);
|
|
});
|
|
|
|
it('should not detect replacement when same job (same createdAt)', async () => {
|
|
const streamId = 'test-stream-123';
|
|
const originalCreatedAt = 1000;
|
|
|
|
mockGenerationJobManager.getJob.mockResolvedValue({
|
|
createdAt: originalCreatedAt,
|
|
});
|
|
|
|
const wasReplaced = await detectJobReplacement(streamId, originalCreatedAt);
|
|
|
|
expect(wasReplaced).toBe(false);
|
|
});
|
|
});
|
|
|
|
describe('Event Emission Behavior', () => {
|
|
/**
|
|
* Simulates the final event emission logic from request.js
|
|
*/
|
|
const emitFinalEventIfNotReplaced = async ({
|
|
streamId,
|
|
originalCreatedAt,
|
|
finalEvent,
|
|
userId,
|
|
}) => {
|
|
const currentJob = await mockGenerationJobManager.getJob(streamId);
|
|
const jobWasReplaced = !currentJob || currentJob.createdAt !== originalCreatedAt;
|
|
|
|
if (jobWasReplaced) {
|
|
mockLogger.debug('Skipping FINAL emit - job was replaced', {
|
|
streamId,
|
|
originalCreatedAt,
|
|
currentCreatedAt: currentJob?.createdAt,
|
|
});
|
|
await mockDecrementPendingRequest(userId);
|
|
return false;
|
|
}
|
|
|
|
mockGenerationJobManager.emitDone(streamId, finalEvent);
|
|
mockGenerationJobManager.completeJob(streamId);
|
|
await mockDecrementPendingRequest(userId);
|
|
return true;
|
|
};
|
|
|
|
it('should skip emitting when job was replaced', async () => {
|
|
const streamId = 'test-stream-123';
|
|
const originalCreatedAt = 1000;
|
|
const newCreatedAt = 2000;
|
|
const userId = 'user-123';
|
|
|
|
mockGenerationJobManager.getJob.mockResolvedValue({
|
|
createdAt: newCreatedAt,
|
|
});
|
|
|
|
const emitted = await emitFinalEventIfNotReplaced({
|
|
streamId,
|
|
originalCreatedAt,
|
|
finalEvent: { final: true },
|
|
userId,
|
|
});
|
|
|
|
expect(emitted).toBe(false);
|
|
expect(mockGenerationJobManager.emitDone).not.toHaveBeenCalled();
|
|
expect(mockGenerationJobManager.completeJob).not.toHaveBeenCalled();
|
|
expect(mockDecrementPendingRequest).toHaveBeenCalledWith(userId);
|
|
expect(mockLogger.debug).toHaveBeenCalledWith(
|
|
'Skipping FINAL emit - job was replaced',
|
|
expect.objectContaining({
|
|
streamId,
|
|
originalCreatedAt,
|
|
currentCreatedAt: newCreatedAt,
|
|
}),
|
|
);
|
|
});
|
|
|
|
it('should emit when job was not replaced', async () => {
|
|
const streamId = 'test-stream-123';
|
|
const originalCreatedAt = 1000;
|
|
const userId = 'user-123';
|
|
const finalEvent = { final: true, conversation: { conversationId: streamId } };
|
|
|
|
mockGenerationJobManager.getJob.mockResolvedValue({
|
|
createdAt: originalCreatedAt,
|
|
});
|
|
|
|
const emitted = await emitFinalEventIfNotReplaced({
|
|
streamId,
|
|
originalCreatedAt,
|
|
finalEvent,
|
|
userId,
|
|
});
|
|
|
|
expect(emitted).toBe(true);
|
|
expect(mockGenerationJobManager.emitDone).toHaveBeenCalledWith(streamId, finalEvent);
|
|
expect(mockGenerationJobManager.completeJob).toHaveBeenCalledWith(streamId);
|
|
expect(mockDecrementPendingRequest).toHaveBeenCalledWith(userId);
|
|
});
|
|
});
|
|
|
|
describe('Response Message Saving Order', () => {
|
|
/**
|
|
* Tests that response messages are saved BEFORE final events are emitted
|
|
* This prevents race conditions where clients send follow-up messages
|
|
* before the response is in the database
|
|
*/
|
|
it('should save message before emitting final event', async () => {
|
|
const callOrder = [];
|
|
|
|
mockSaveMessage.mockImplementation(async () => {
|
|
callOrder.push('saveMessage');
|
|
});
|
|
|
|
mockGenerationJobManager.emitDone.mockImplementation(() => {
|
|
callOrder.push('emitDone');
|
|
});
|
|
|
|
mockGenerationJobManager.getJob.mockResolvedValue({
|
|
createdAt: 1000,
|
|
});
|
|
|
|
// Simulate the order of operations from request.js
|
|
const streamId = 'test-stream-123';
|
|
const originalCreatedAt = 1000;
|
|
const response = { messageId: 'response-123' };
|
|
const userId = 'user-123';
|
|
|
|
// Step 1: Save message
|
|
await mockSaveMessage({}, { ...response, user: userId }, { context: 'test' });
|
|
|
|
// Step 2: Check for replacement
|
|
const currentJob = await mockGenerationJobManager.getJob(streamId);
|
|
const jobWasReplaced = !currentJob || currentJob.createdAt !== originalCreatedAt;
|
|
|
|
// Step 3: Emit if not replaced
|
|
if (!jobWasReplaced) {
|
|
mockGenerationJobManager.emitDone(streamId, { final: true });
|
|
}
|
|
|
|
expect(callOrder).toEqual(['saveMessage', 'emitDone']);
|
|
});
|
|
});
|
|
|
|
describe('Aborted Request Handling', () => {
|
|
it('should use unfinished: true instead of error: true for aborted requests', () => {
|
|
const response = { messageId: 'response-123', content: [] };
|
|
|
|
// The new format for aborted responses
|
|
const abortedResponse = { ...response, unfinished: true };
|
|
|
|
expect(abortedResponse.unfinished).toBe(true);
|
|
expect(abortedResponse.error).toBeUndefined();
|
|
});
|
|
|
|
it('should include unfinished flag in final event for aborted requests', () => {
|
|
const response = { messageId: 'response-123', content: [] };
|
|
|
|
// Old format (deprecated)
|
|
const _oldFinalEvent = {
|
|
final: true,
|
|
responseMessage: { ...response, error: true },
|
|
error: { message: 'Request was aborted' },
|
|
};
|
|
|
|
// New format (PR #11462)
|
|
const newFinalEvent = {
|
|
final: true,
|
|
responseMessage: { ...response, unfinished: true },
|
|
};
|
|
|
|
expect(newFinalEvent.responseMessage.unfinished).toBe(true);
|
|
expect(newFinalEvent.error).toBeUndefined();
|
|
expect(newFinalEvent.responseMessage.error).toBeUndefined();
|
|
});
|
|
});
|
|
});
|
|
|
|
/**
|
|
* HITL terminal-side-effect guards (PR #13942).
|
|
*
|
|
* Jobs are keyed by streamId == conversationId, so a NEW request REPLACES the running
|
|
* one on the same conversation. The replaced generation's tail (its pause attempt, its
|
|
* checkpoint prune, its resume catch-path terminal writes) must not clobber the live
|
|
* generation's state. Each guard re-reads the live job and compares createdAt against the
|
|
* generation's own captured identity before acting. These mirror the predicates in
|
|
* client.js (handleRunInterrupt / chatCompletion finally) and resume.js.
|
|
*/
|
|
describe('HITL Terminal-Side-Effect Guards', () => {
|
|
beforeEach(() => {
|
|
jest.clearAllMocks();
|
|
});
|
|
|
|
describe('F22 — pause is skipped when the generation was replaced', () => {
|
|
// Mirrors client.js handleRunInterrupt pre-check, run BEFORE approvals.pause.
|
|
const shouldPause = async ({ jobCreatedAt, streamId }) => {
|
|
if (jobCreatedAt != null) {
|
|
const liveJob = await mockGenerationJobManager.getJob(streamId);
|
|
if (!liveJob || liveJob.createdAt !== jobCreatedAt) {
|
|
return false;
|
|
}
|
|
}
|
|
return true;
|
|
};
|
|
|
|
it('does not pause when a newer job replaced this one', async () => {
|
|
mockGenerationJobManager.getJob.mockResolvedValue({ createdAt: 2000 });
|
|
expect(await shouldPause({ jobCreatedAt: 1000, streamId: 'c1' })).toBe(false);
|
|
});
|
|
|
|
it('does not pause when the job is already gone', async () => {
|
|
mockGenerationJobManager.getJob.mockResolvedValue(null);
|
|
expect(await shouldPause({ jobCreatedAt: 1000, streamId: 'c1' })).toBe(false);
|
|
});
|
|
|
|
it('pauses when this is still the live job', async () => {
|
|
mockGenerationJobManager.getJob.mockResolvedValue({ createdAt: 1000 });
|
|
expect(await shouldPause({ jobCreatedAt: 1000, streamId: 'c1' })).toBe(true);
|
|
});
|
|
|
|
it('pauses without a lookup when identity is unknown (legacy job)', async () => {
|
|
expect(await shouldPause({ jobCreatedAt: null, streamId: 'c1' })).toBe(true);
|
|
expect(mockGenerationJobManager.getJob).not.toHaveBeenCalled();
|
|
});
|
|
});
|
|
|
|
// (Removed: F21 — the chatCompletion clean-path checkpoint prune + its job-replacement
|
|
// guard no longer exist. The lazy checkpointer never writes a clean-exit checkpoint, so
|
|
// there is nothing to prune after a non-paused turn; the pre-run prune (before
|
|
// processStream) clears any orphaned interrupt checkpoint instead. See
|
|
// checkpointer.ts LazyMongoSaver and client.js chatCompletion.)
|
|
|
|
describe('F24 — resume catch-path terminal writes are skipped when replaced', () => {
|
|
// Mirrors resume.js: stillLive gate around emitError/completeJob/deleteAgentCheckpoint.
|
|
const stillLive = async ({ streamId, jobCreatedAt }) => {
|
|
let live = true;
|
|
try {
|
|
const liveJob = await mockGenerationJobManager.getJob(streamId);
|
|
live = !!liveJob && liveJob.createdAt === jobCreatedAt;
|
|
} catch {
|
|
live = true; // read failed — fail open and run the terminal writes
|
|
}
|
|
return live;
|
|
};
|
|
|
|
it('runs terminal writes when this is still the live job', async () => {
|
|
mockGenerationJobManager.getJob.mockResolvedValue({ createdAt: 1000 });
|
|
expect(await stillLive({ streamId: 'c1', jobCreatedAt: 1000 })).toBe(true);
|
|
});
|
|
|
|
it('skips terminal writes when a newer job replaced this one', async () => {
|
|
mockGenerationJobManager.getJob.mockResolvedValue({ createdAt: 2000 });
|
|
expect(await stillLive({ streamId: 'c1', jobCreatedAt: 1000 })).toBe(false);
|
|
});
|
|
|
|
it('fails open (runs terminal writes) when the liveness read throws', async () => {
|
|
mockGenerationJobManager.getJob.mockRejectedValue(new Error('store down'));
|
|
expect(await stillLive({ streamId: 'c1', jobCreatedAt: 1000 })).toBe(true);
|
|
});
|
|
});
|
|
|
|
describe('F23 — resumed turn sources files from the job, not the racy DB row', () => {
|
|
// Mirrors resume.js: prefer the body, then job.metadata.userMessage.files, then DB.
|
|
const resolveFiles = ({ bodyFiles, metaFiles, dbFiles }) => {
|
|
if (Array.isArray(bodyFiles) && bodyFiles.length > 0) {
|
|
return bodyFiles;
|
|
}
|
|
if (Array.isArray(metaFiles) && metaFiles.length > 0) {
|
|
return metaFiles;
|
|
}
|
|
return Array.isArray(dbFiles) && dbFiles.length > 0 ? dbFiles : undefined;
|
|
};
|
|
|
|
it('prefers job-metadata files over the DB row (no DB-save race)', () => {
|
|
expect(
|
|
resolveFiles({
|
|
bodyFiles: [],
|
|
metaFiles: [{ file_id: 'meta' }],
|
|
dbFiles: [{ file_id: 'db' }],
|
|
}),
|
|
).toEqual([{ file_id: 'meta' }]);
|
|
});
|
|
|
|
it('falls back to the DB row when the job has no persisted files (older job)', () => {
|
|
expect(
|
|
resolveFiles({ bodyFiles: [], metaFiles: undefined, dbFiles: [{ file_id: 'db' }] }),
|
|
).toEqual([{ file_id: 'db' }]);
|
|
});
|
|
|
|
it('keeps files already present on the resume body', () => {
|
|
expect(
|
|
resolveFiles({
|
|
bodyFiles: [{ file_id: 'body' }],
|
|
metaFiles: [{ file_id: 'meta' }],
|
|
dbFiles: [],
|
|
}),
|
|
).toEqual([{ file_id: 'body' }]);
|
|
});
|
|
});
|
|
});
|
|
|
|
/**
|
|
* Round-18 follow-ups to the guards above (Codex review 4594099963).
|
|
*/
|
|
describe('HITL Resume Fidelity Guards (round 18)', () => {
|
|
beforeEach(() => {
|
|
jest.clearAllMocks();
|
|
});
|
|
|
|
describe('G1 — resume re-checks ownership AGAIN right before terminal writes', () => {
|
|
// The start-of-finalize guard can go stale across saveMessage + title generation,
|
|
// so resume.js re-reads the live job immediately before emitDone/completeJob/prune.
|
|
// Same predicate as the catch-path (F24), applied at the success path's second point.
|
|
const stillLiveBeforeFinalize = async ({ streamId, jobCreatedAt }) => {
|
|
const liveJob = await mockGenerationJobManager.getJob(streamId);
|
|
return !!liveJob && liveJob.createdAt === jobCreatedAt;
|
|
};
|
|
|
|
it('runs terminal writes when still the live job at the second check', async () => {
|
|
mockGenerationJobManager.getJob.mockResolvedValue({ createdAt: 1000 });
|
|
expect(await stillLiveBeforeFinalize({ streamId: 'c1', jobCreatedAt: 1000 })).toBe(true);
|
|
});
|
|
|
|
it('skips terminal writes when replaced DURING finalize (after the first check passed)', async () => {
|
|
// First check passed earlier with createdAt 1000; a new request replaced it to 2000
|
|
// while saveMessage + title generation awaited. The second check must catch it.
|
|
mockGenerationJobManager.getJob.mockResolvedValue({ createdAt: 2000 });
|
|
expect(await stillLiveBeforeFinalize({ streamId: 'c1', jobCreatedAt: 1000 })).toBe(false);
|
|
});
|
|
});
|
|
|
|
describe('G2 — uploaded files are seeded into the AWAITED preliminary user message', () => {
|
|
// Mirrors getPreliminaryUserMessage: files from the request are persisted on the
|
|
// preliminary (awaited, pre-run) metadata so they land before any interrupt emits.
|
|
const buildPreliminaryUserMessage = ({ messageId, files }) => {
|
|
if (typeof messageId !== 'string' || messageId.length === 0) {
|
|
return null;
|
|
}
|
|
return {
|
|
messageId,
|
|
...(Array.isArray(files) && files.length > 0 && { files }),
|
|
};
|
|
};
|
|
|
|
it('includes files when the request carries them', () => {
|
|
const msg = buildPreliminaryUserMessage({ messageId: 'm1', files: [{ file_id: 'a' }] });
|
|
expect(msg.files).toEqual([{ file_id: 'a' }]);
|
|
});
|
|
|
|
it('omits files when none were uploaded (no empty array)', () => {
|
|
const msg = buildPreliminaryUserMessage({ messageId: 'm1', files: [] });
|
|
expect(msg).not.toHaveProperty('files');
|
|
});
|
|
});
|
|
|
|
describe('G3 — resume replays pre-pause discovered deferred tools', () => {
|
|
// Mirrors createRun's merge: discovered set is union(message-extracted, replayed),
|
|
// gated entirely on the agent actually having deferred tools.
|
|
const resolveDiscovered = ({ hasAnyDeferredTools, messageExtracted, replayed }) => {
|
|
const set = new Set();
|
|
if (hasAnyDeferredTools) {
|
|
for (const n of messageExtracted ?? []) {
|
|
set.add(n);
|
|
}
|
|
for (const n of replayed ?? []) {
|
|
set.add(n);
|
|
}
|
|
}
|
|
return set;
|
|
};
|
|
|
|
it('replays captured names on resume (messages empty) so the paused tool is present', () => {
|
|
const set = resolveDiscovered({
|
|
hasAnyDeferredTools: true,
|
|
messageExtracted: [],
|
|
replayed: ['deep_tool'],
|
|
});
|
|
expect(set.has('deep_tool')).toBe(true);
|
|
});
|
|
|
|
it('unions replayed names with message-extracted names', () => {
|
|
const set = resolveDiscovered({
|
|
hasAnyDeferredTools: true,
|
|
messageExtracted: ['from_history'],
|
|
replayed: ['deep_tool'],
|
|
});
|
|
expect([...set].sort()).toEqual(['deep_tool', 'from_history']);
|
|
});
|
|
|
|
it('is inert when the agent has no deferred tools', () => {
|
|
const set = resolveDiscovered({
|
|
hasAnyDeferredTools: false,
|
|
messageExtracted: ['x'],
|
|
replayed: ['deep_tool'],
|
|
});
|
|
expect(set.size).toBe(0);
|
|
});
|
|
});
|
|
|
|
describe("H3 — resume replays the paused turn's model parameters (ephemeral agents)", () => {
|
|
// Mirrors restoreResumeContext: spread persisted model_parameters back onto the body,
|
|
// excluding `model` (replayed via the fingerprinted RESUME_CONTEXT_KEYS path).
|
|
const replayModelParameters = (body, resumeContext) => {
|
|
const params = resumeContext?.model_parameters;
|
|
if (params && typeof params === 'object') {
|
|
const { model: _model, ...rest } = params;
|
|
Object.assign(body, rest);
|
|
}
|
|
return body;
|
|
};
|
|
|
|
it('restores non-default params (temperature, max tokens) onto the resume body', () => {
|
|
const body = { conversationId: 'c1', endpoint: 'agents' };
|
|
replayModelParameters(body, {
|
|
model_parameters: { model: 'gpt-4o', temperature: 0.2, max_tokens: 1024 },
|
|
});
|
|
expect(body).toMatchObject({ temperature: 0.2, max_tokens: 1024 });
|
|
});
|
|
|
|
it('does NOT overwrite model (kept consistent with the resume fingerprint)', () => {
|
|
const body = { model: 'pinned-model' };
|
|
replayModelParameters(body, { model_parameters: { model: 'other-model', temperature: 0.9 } });
|
|
expect(body.model).toBe('pinned-model');
|
|
});
|
|
|
|
it('overwrites a client-supplied param with the captured authoritative value', () => {
|
|
const body = { temperature: 1.0 }; // crafted/stale client value
|
|
replayModelParameters(body, { model_parameters: { temperature: 0.2 } });
|
|
expect(body.temperature).toBe(0.2);
|
|
});
|
|
|
|
it('is a no-op when nothing was captured', () => {
|
|
const body = { conversationId: 'c1' };
|
|
replayModelParameters(body, {});
|
|
expect(body).toEqual({ conversationId: 'c1' });
|
|
});
|
|
});
|
|
|
|
describe('J2 — pause unfinished-save is skipped once a fast resume took over', () => {
|
|
// Mirrors request.js: only mark the paused row unfinished while the job is STILL paused
|
|
// on THIS generation's action. A claim transitions it out of requires_action and a
|
|
// replacement bumps createdAt — either means a /resume now owns the row, so marking it
|
|
// unfinished would clobber the resumed turn's completed content. Fail open on read error.
|
|
const shouldMarkUnfinished = async ({ jobCreatedAt, streamId }) => {
|
|
let stillPaused = true;
|
|
try {
|
|
const liveJob = await mockGenerationJobManager.getJob(streamId);
|
|
stillPaused =
|
|
!!liveJob &&
|
|
liveJob.status === 'requires_action' &&
|
|
(jobCreatedAt == null || liveJob.createdAt === jobCreatedAt);
|
|
} catch {
|
|
stillPaused = true;
|
|
}
|
|
return stillPaused;
|
|
};
|
|
|
|
it('marks unfinished while still paused on this generation', async () => {
|
|
mockGenerationJobManager.getJob.mockResolvedValue({
|
|
status: 'requires_action',
|
|
createdAt: 1000,
|
|
});
|
|
expect(await shouldMarkUnfinished({ jobCreatedAt: 1000, streamId: 'c1' })).toBe(true);
|
|
});
|
|
|
|
it('skips the unfinished-save once a fast resume claimed it (no longer requires_action)', async () => {
|
|
mockGenerationJobManager.getJob.mockResolvedValue({ status: 'running', createdAt: 1000 });
|
|
expect(await shouldMarkUnfinished({ jobCreatedAt: 1000, streamId: 'c1' })).toBe(false);
|
|
});
|
|
|
|
it('skips the unfinished-save when a newer request replaced the job', async () => {
|
|
mockGenerationJobManager.getJob.mockResolvedValue({
|
|
status: 'requires_action',
|
|
createdAt: 2000,
|
|
});
|
|
expect(await shouldMarkUnfinished({ jobCreatedAt: 1000, streamId: 'c1' })).toBe(false);
|
|
});
|
|
|
|
it('fails open (marks unfinished) when the liveness read throws', async () => {
|
|
mockGenerationJobManager.getJob.mockRejectedValue(new Error('store down'));
|
|
expect(await shouldMarkUnfinished({ jobCreatedAt: 1000, streamId: 'c1' })).toBe(true);
|
|
});
|
|
});
|
|
});
|