mirror of
https://github.com/THU-MAIC/OpenMAIC.git
synced 2026-10-02 01:15:18 +08:00
refactor(media): one client-side pool commit primitive; keep refused narration instead of re-billing it (#1523)
Extracts commitToPool, the single client-side sequence for storing bytes in the asset pool, writing the allocated id back, and mirroring locally; routes the media pass, narration adoption and fresh TTS through it. A store-full refusal during TTS now retains the already-billed clip so the next load adopts it with zero provider calls. Closes #1467.
This commit is contained in:
@@ -37,9 +37,8 @@
|
||||
* and paying a provider to replace it is a decision for the author, not a
|
||||
* side effect of opening a course.
|
||||
*/
|
||||
import { putAsset } from '@/lib/media/asset-pool';
|
||||
import { commitToPool } from '@/lib/media/commit-to-pool';
|
||||
import { mayGenerateForStage } from '@/lib/classroom/generation-permission';
|
||||
import { isStorageFullFailure } from '@/lib/media/media-failure';
|
||||
import { createLogger } from '@/lib/logger';
|
||||
import { mayNameAPoolAsset } from '@/lib/media/media-placeholder';
|
||||
import { isConcreteMediaAddress } from '@/lib/media/resolve-media-ref';
|
||||
@@ -170,12 +169,14 @@ const DERIVED_KEY_ACTION_ID = /^tts_(?:request_)?s-?\d+_(.+)$/;
|
||||
*/
|
||||
const UNIQUE_ACTION_ID = /^action_[A-Za-z0-9_-]{8,}$/;
|
||||
|
||||
/** The contract code an upload failure declares, if it declares one. */
|
||||
function storageErrorCode(error: unknown): string | undefined {
|
||||
if (typeof error !== 'object' || error === null) return undefined;
|
||||
const code = (error as { code?: unknown }).code;
|
||||
return typeof code === 'string' ? code : undefined;
|
||||
}
|
||||
/**
|
||||
* What the narration write-back settled as, for the clip the loop is on.
|
||||
*
|
||||
* `departed` is not a failure of the write: it is this browser leaving the
|
||||
* course between the allocation and the rewrite, which ends the run rather than
|
||||
* skipping one clip.
|
||||
*/
|
||||
type NarrationPlacement = 'placed' | 'unplaced' | 'departed';
|
||||
|
||||
function derivedKeyIsUnique(derivedRef: string): boolean {
|
||||
const actionId = DERIVED_KEY_ACTION_ID.exec(derivedRef)?.[1];
|
||||
@@ -407,70 +408,83 @@ async function adoptCachedNarrationRun(
|
||||
// Re-checked after the read and before anything is spent.
|
||||
if (abortSignal?.aborted || !onThisCourse()) break;
|
||||
|
||||
let assetId: string;
|
||||
try {
|
||||
// Bytes first, exactly as the media path does it: a document may never
|
||||
// name narration that was not stored.
|
||||
assetId = await putAsset(
|
||||
row.blob,
|
||||
{
|
||||
contentType: row.blob.type || `audio/${row.format}`,
|
||||
...(row.duration === undefined ? {} : { durationSeconds: row.duration }),
|
||||
},
|
||||
// A write that goes through retires this course's "no room" note, at
|
||||
// the seam rather than here. Together with generated narration this is
|
||||
// the only path that can establish that for a course whose media needs
|
||||
// nothing, and it is worth naming what it costs: a few hundred bytes of
|
||||
// narration fit in headroom an image does not, so retiring the note can
|
||||
// let the next pass pay a provider for an image that is refused again.
|
||||
// Bounded at one such generation, because that pass re-marks and
|
||||
// adoption converts everything that fits in a single load, and the
|
||||
// alternative is a course whose media never generates again.
|
||||
{ stageId },
|
||||
);
|
||||
} catch (error) {
|
||||
// Bytes first, exactly as the media path does it -- through the same
|
||||
// primitive, in fact: a document may never name narration that was not
|
||||
// stored.
|
||||
//
|
||||
// No `retain` sink is handed over, and that is the whole of this caller's
|
||||
// refusal semantics: the bytes this commit would keep are the bytes it is
|
||||
// reading, already in `audioFiles` under the derived key, which is exactly
|
||||
// where the next load looks for them. A refusal here loses nothing and
|
||||
// costs no provider call.
|
||||
const outcome = await commitToPool<NarrationPlacement>({
|
||||
// A write that goes through retires this course's "no room" note, at the
|
||||
// seam rather than here. Together with generated narration this is the
|
||||
// only path that can establish that for a course whose media needs
|
||||
// nothing, and it is worth naming what it costs: a few hundred bytes of
|
||||
// narration fit in headroom an image does not, so retiring the note can
|
||||
// let the next pass pay a provider for an image that is refused again.
|
||||
// Bounded at one such generation, because that pass re-marks and adoption
|
||||
// converts everything that fits in a single load, and the alternative is
|
||||
// a course whose media never generates again.
|
||||
stageId,
|
||||
slot: action.derivedRef,
|
||||
bytes: row.blob,
|
||||
mimeType: row.blob.type || `audio/${row.format}`,
|
||||
...(row.duration === undefined ? {} : { meta: { durationSeconds: row.duration } }),
|
||||
writeBack: async (assetId) => {
|
||||
// The allocation is uncancellable, so it may finish after the course
|
||||
// was left. Its write-back is not: a document this browser no longer
|
||||
// has open would take a lock for a rewrite the live store cannot
|
||||
// mirror.
|
||||
if (!onThisCourse()) return 'departed';
|
||||
const placed = await persistNarrationReference(stageId, action.derivedRef, assetId).catch(
|
||||
(error: unknown) => {
|
||||
log.warn(`Could not write back narration ${action.derivedRef}:`, error);
|
||||
return false;
|
||||
},
|
||||
);
|
||||
return placed ? 'placed' : 'unplaced';
|
||||
},
|
||||
// Local mirror under the new id, stage-scoped so it cannot be mistaken
|
||||
// for another course's the way the derived row could be. The document
|
||||
// already points at the pool, so a failed cache write costs a
|
||||
// re-download. Nothing is mirrored for a rewrite nothing took: the new id
|
||||
// is not the one anything reads by.
|
||||
mirror: async (assetId, placement) => {
|
||||
if (placement !== 'placed') return;
|
||||
await db.audioFiles
|
||||
.put({ ...row, id: assetId, stageId, originAudioId: action.derivedRef })
|
||||
.catch((error: unknown) => {
|
||||
log.warn(`Local narration cache mirror failed for ${assetId}:`, error);
|
||||
});
|
||||
},
|
||||
});
|
||||
|
||||
if (outcome.status !== 'stored') {
|
||||
// One clip's storage failure costs that clip and nothing else. The action
|
||||
// keeps its derived id, the deck carries on, and a later load tries
|
||||
// again -- which is how a course converges the moment the ceiling moves,
|
||||
// with nothing to click and nothing to remember.
|
||||
log.warn(`Could not store cached narration ${action.derivedRef}:`, error);
|
||||
log.warn(`Could not store cached narration ${action.derivedRef}:`, outcome.error);
|
||||
unbacked += 1;
|
||||
// A refusal for room, and only that, lowers the bar for the rest of this
|
||||
// run. Any other failure -- a dropped connection, a 500 -- says nothing
|
||||
// about how much room there is, so it must not stop the next clip being
|
||||
// attempted.
|
||||
if (isStorageFullFailure(storageErrorCode(error))) {
|
||||
if (outcome.status === 'refused-retained') {
|
||||
smallestRefusedForRoom = Math.min(smallestRefusedForRoom, row.blob.size);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
// The allocation is uncancellable, so it may finish after the course was
|
||||
// left. Its write-back is not: a document this browser no longer has open
|
||||
// would take a lock for a rewrite the live store cannot mirror.
|
||||
if (!onThisCourse()) {
|
||||
if (outcome.placement === 'departed') {
|
||||
unbacked += 1;
|
||||
break;
|
||||
}
|
||||
|
||||
const placed = await persistNarrationReference(stageId, action.derivedRef, assetId).catch(
|
||||
(error: unknown) => {
|
||||
log.warn(`Could not write back narration ${action.derivedRef}:`, error);
|
||||
return false;
|
||||
},
|
||||
);
|
||||
if (!placed) {
|
||||
if (outcome.placement === 'unplaced') {
|
||||
unbacked += 1;
|
||||
continue;
|
||||
}
|
||||
|
||||
// Local mirror under the new id, stage-scoped so it cannot be mistaken for
|
||||
// another course's the way the derived row could be. The document already
|
||||
// points at the pool, so a failed cache write costs a re-download.
|
||||
await db.audioFiles
|
||||
.put({ ...row, id: assetId, stageId, originAudioId: action.derivedRef })
|
||||
.catch((error: unknown) => {
|
||||
log.warn(`Local narration cache mirror failed for ${assetId}:`, error);
|
||||
});
|
||||
adopted += 1;
|
||||
}
|
||||
|
||||
|
||||
@@ -28,7 +28,7 @@ import {
|
||||
import { resolveTTSModelForVoice } from '@/lib/audio/constants';
|
||||
import { useAgentRegistry } from '@/lib/orchestration/registry/store';
|
||||
import { generateMediaForOutlines } from '@/lib/media/media-orchestrator';
|
||||
import { putAsset } from '@/lib/media/asset-pool';
|
||||
import { commitToPool } from '@/lib/media/commit-to-pool';
|
||||
import { mayGenerateForStage } from '@/lib/classroom/generation-permission';
|
||||
import { isServerBackedMediaPersistence } from '@/lib/persistence/media-persistence';
|
||||
import { lazyBoundedMap } from '@/lib/utils/concurrency';
|
||||
@@ -473,93 +473,113 @@ export async function generateAndStoreTTS(
|
||||
// clip onto a timeline without re-decoding. null → leave undefined; the audio
|
||||
// still persists and plays.
|
||||
const duration = measureAudioDuration(bytes, data.format) ?? undefined;
|
||||
const serverBacked = isServerBackedMediaPersistence();
|
||||
// Server-backed: the bytes go to the pool and the pool allocates the
|
||||
// identity, so the id the speech action ends up holding names durable audio
|
||||
// rather than this browser's local table. Bytes land BEFORE the caller
|
||||
// stamps the action, so a document can never name narration that was not
|
||||
// stored. Browser-only keeps the historical derived key: document and audio
|
||||
// share one lifetime there, and nothing outside this browser reads either.
|
||||
let audioId: string;
|
||||
if (serverBacked) {
|
||||
const allocated = await allocatePooledAudio(blob, duration, stageId).catch((error: unknown) => {
|
||||
// Storing narration failed, not synthesizing it. A scene whose audio
|
||||
// cannot be stored keeps its text and leaves the line unvoiced and
|
||||
// retryable, exactly as an image that cannot be stored leaves its slide;
|
||||
// reporting it as a TTS failure would pause the whole deck at its first
|
||||
// slide over one clip's storage.
|
||||
log.warn('Narration storage failed; leaving the line unvoiced:', error);
|
||||
return null;
|
||||
});
|
||||
if (allocated === null) return null;
|
||||
audioId = allocated;
|
||||
} else {
|
||||
audioId = existingAudioId ?? requestId;
|
||||
}
|
||||
const cacheWrite = db.audioFiles.put({
|
||||
id: audioId,
|
||||
/** This clip's local row, under whichever id it is currently known by. */
|
||||
const cachedNarrationRow = (id: string) => ({
|
||||
id,
|
||||
stageId,
|
||||
blob,
|
||||
duration,
|
||||
format: data.format,
|
||||
format: data.format as string,
|
||||
text,
|
||||
voice: ttsVoice,
|
||||
createdAt: Date.now(),
|
||||
});
|
||||
if (serverBacked) {
|
||||
// A cache the pool already backs: a failed write costs a re-download.
|
||||
await cacheWrite.catch((error: unknown) => {
|
||||
log.warn('Local narration cache write failed for', audioId, error);
|
||||
});
|
||||
} else {
|
||||
await cacheWrite;
|
||||
const serverBacked = isServerBackedMediaPersistence();
|
||||
// Browser-only keeps the historical derived key: document and audio share one
|
||||
// lifetime there, and nothing outside this browser reads either.
|
||||
if (!serverBacked) {
|
||||
const audioId = existingAudioId ?? requestId;
|
||||
await db.audioFiles.put(cachedNarrationRow(audioId));
|
||||
return audioId;
|
||||
}
|
||||
return audioId;
|
||||
|
||||
// Server-backed: the bytes go to the pool and the pool allocates the
|
||||
// identity, so the id the speech action ends up holding names durable audio
|
||||
// rather than this browser's local table. Bytes land BEFORE the caller stamps
|
||||
// the action, so a document can never name narration that was not stored.
|
||||
const outcome = await commitToPool<void>({
|
||||
stageId,
|
||||
// The derived key, which is both what a refusal keeps the bytes under and
|
||||
// what narration adoption reads them back by on a later load.
|
||||
slot: requestId,
|
||||
bytes: blob,
|
||||
mimeType: blob.type,
|
||||
...(duration === undefined ? {} : { meta: { durationSeconds: duration } }),
|
||||
// The bytes were just bought. A full store must not be what throws them
|
||||
// away: keeping them under the derived key is what lets the next load
|
||||
// re-attempt the upload from cache instead of paying the provider again,
|
||||
// which is the same contract the media pass's retained bytes have had since
|
||||
// it learned to keep them. See the caller's handling below for the other
|
||||
// half of it -- the action has to carry this key for adoption to find them.
|
||||
//
|
||||
// The rejection is NOT swallowed, and that is the point of awaiting it: a
|
||||
// stamp is only safe once the bytes are somewhere that can be read back. A
|
||||
// local table that refuses the row leaves nothing to adopt, so the commit
|
||||
// demotes itself to `failed` and the line goes unvoiced instead of carrying
|
||||
// a derived key that resolves to nothing for the rest of the course's life.
|
||||
retain: async () => {
|
||||
await db.audioFiles.put(cachedNarrationRow(requestId));
|
||||
},
|
||||
// Nothing to write back: the action this narration belongs to is not in the
|
||||
// document yet. The caller stamps it from the id returned here, which is
|
||||
// why this path has no funnel of its own to invent one.
|
||||
writeBack: async () => undefined,
|
||||
// A cache the pool already backs: a failed write costs a re-download, and
|
||||
// the primitive holds that to be best-effort for every caller.
|
||||
mirror: async (assetId) => {
|
||||
await db.audioFiles.put(cachedNarrationRow(assetId));
|
||||
},
|
||||
});
|
||||
|
||||
if (outcome.status === 'stored') return outcome.assetId;
|
||||
if (outcome.status === 'refused-retained') {
|
||||
// The store had no room, and the bytes are kept. The action is stamped with
|
||||
// the derived key they are kept under, exactly as a refused image leaves
|
||||
// its placeholder in the slide: adoption reads that key on the next load,
|
||||
// re-attempts the upload, and writes the allocated id back with no provider
|
||||
// called. Returning null instead would leave the line unvoiced AND the
|
||||
// bytes unreachable, which is paying for the same clip on every attempt.
|
||||
log.warn(
|
||||
`Asset storage is full; keeping the narration for ${requestId} under its derived key.`,
|
||||
);
|
||||
return requestId;
|
||||
}
|
||||
// Storing narration failed for some other reason -- or the bytes could not be
|
||||
// kept -- and neither says anything about whether a later attempt would fit,
|
||||
// so nothing is left under a key a later load would take for adoptable
|
||||
// narration. A scene whose audio cannot be
|
||||
// stored keeps its text and leaves the line unvoiced and retryable, exactly
|
||||
// as an image that cannot be stored leaves its slide; reporting it as a TTS
|
||||
// failure would pause the whole deck at its first slide over one clip's
|
||||
// storage.
|
||||
log.warn('Narration storage failed; leaving the line unvoiced:', outcome.error);
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Store narration bytes in the asset pool and return the reference the
|
||||
* document should hold.
|
||||
* Why a fresh clip never replaces the bytes behind an id it is superseding.
|
||||
*
|
||||
* Regeneration always forks to a fresh id; the caller's `existingAudioId` is
|
||||
* deliberately ignored here. Replacing bytes behind a live id requires proof
|
||||
* that no other document holds it, and that proof is unavailable by
|
||||
* Regeneration always forks; the caller's `existingAudioId` is deliberately
|
||||
* ignored on the server-backed path. Replacing bytes behind a live id requires
|
||||
* proof that no other document holds it, and that proof is unavailable by
|
||||
* construction once references can leave this browser — asking the pool who
|
||||
* else holds an id would be exactly the existence oracle the asset contract
|
||||
* forbids, so `proveExclusiveAssetOwnership` fails closed under server-backed
|
||||
* persistence and every caller forks. Keeping a branch that can never be taken
|
||||
* would only describe a capability this deployment shape does not have.
|
||||
*
|
||||
* The superseded id is NOT removed here. Nothing at this point has observed
|
||||
* the new id reaching a durable document, so deleting the old bytes could
|
||||
* leave a still-referenced action pointing at nothing if the save that follows
|
||||
* fails; and the exclusivity that would make deletion safe is the same proof
|
||||
* that is unavailable. It does not have to be removed here: the save that
|
||||
* writes the new id is also the write that stops naming the old one, so the
|
||||
* server stamps the superseded entry as it lands and the collector releases it
|
||||
* after the grace period, the bytes following after their own. If that save
|
||||
* never lands, it is the NEW id that nothing committed, and it expires on
|
||||
* The superseded id is NOT removed either. Nothing at this point has observed
|
||||
* the new id reaching a durable document, so deleting the old bytes could leave
|
||||
* a still-referenced action pointing at nothing if the save that follows fails;
|
||||
* and the exclusivity that would make deletion safe is the same proof that is
|
||||
* unavailable. It does not have to be removed here: the save that writes the
|
||||
* new id is also the write that stops naming the old one, so the server stamps
|
||||
* the superseded entry as it lands and the collector releases it after the
|
||||
* grace period, the bytes following after their own. If that save never lands,
|
||||
* it is the NEW id that nothing committed, and it expires on
|
||||
* `ASSET_PENDING_TTL_MS` — either way regeneration leaves nothing permanent
|
||||
* behind.
|
||||
*/
|
||||
async function allocatePooledAudio(
|
||||
blob: Blob,
|
||||
duration: number | undefined,
|
||||
stageId: string | undefined,
|
||||
): Promise<string> {
|
||||
return putAsset(
|
||||
blob,
|
||||
{
|
||||
contentType: blob.type,
|
||||
...(duration === undefined ? {} : { durationSeconds: duration }),
|
||||
},
|
||||
// A write that goes through retires this course's "no room" note. This path
|
||||
// allocates directly rather than through the media commit, so without it a
|
||||
// course whose narration is generated rather than adopted has nothing that
|
||||
// can establish that.
|
||||
{ ...(stageId ? { stageId } : {}) },
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Drop the local copies of narration a scene has rolled back.
|
||||
@@ -599,11 +619,38 @@ export async function generateTTSForScene(
|
||||
let failedCount = 0;
|
||||
let lastError: string | undefined;
|
||||
const freshAllocations: string[] = [];
|
||||
/**
|
||||
* Actions holding retained bytes rather than a fresh allocation.
|
||||
*
|
||||
* A clip the store refused for want of room comes back under its own derived
|
||||
* key with its bytes kept in the local table. Nothing was allocated, so a
|
||||
* rollback has nothing to reclaim -- and running one anyway would delete the
|
||||
* only copy of audio that is already paid for and unstamp the key adoption
|
||||
* reads it back by, which is the double billing this whole path exists to
|
||||
* stop. A sibling line failing is not a reason to throw them away.
|
||||
*/
|
||||
const retainedRefusals = new Set<SpeechAction>();
|
||||
const serverBacked = isServerBackedMediaPersistence();
|
||||
|
||||
// Scene order keeps the provider request correlation label unique. Storage
|
||||
// identity is allocated by the pool and is never derived from this value.
|
||||
const sceneOrder = scene.order;
|
||||
|
||||
/**
|
||||
* Undo this scene's narration, keeping whatever a rollback cannot own.
|
||||
*
|
||||
* Everything in `freshAllocations` was minted for this scene and nothing else
|
||||
* holds it, so its local copy goes. A retained refusal is the exception, and
|
||||
* the only one.
|
||||
*/
|
||||
const rollBackFreshNarration = async (): Promise<void> => {
|
||||
await removeFreshTtsAllocations(freshAllocations);
|
||||
for (const action of speechActions) {
|
||||
if (retainedRefusals.has(action)) continue;
|
||||
delete action.audioId;
|
||||
}
|
||||
};
|
||||
|
||||
// Generate + store one action's audio. Failures are counted, not thrown, so
|
||||
// one bad clip never aborts the rest of the scene.
|
||||
const generateOne = async (action: SpeechAction) => {
|
||||
@@ -620,7 +667,13 @@ export async function generateTTSForScene(
|
||||
);
|
||||
if (assetId) {
|
||||
action.audioId = assetId;
|
||||
freshAllocations.push(assetId);
|
||||
// Under server-backed persistence the pool answers with an allocated
|
||||
// id, so the request key coming back means one thing only: the store
|
||||
// refused these bytes and they were kept under it. Browser-only always
|
||||
// returns the request key and always rolls back with the scene, which
|
||||
// is right there -- the bytes and the document share one lifetime.
|
||||
if (serverBacked && assetId === requestId) retainedRefusals.add(action);
|
||||
else freshAllocations.push(assetId);
|
||||
}
|
||||
} catch (error) {
|
||||
if (isAbortError(error)) throw error;
|
||||
@@ -662,14 +715,12 @@ export async function generateTTSForScene(
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
await removeFreshTtsAllocations(freshAllocations);
|
||||
for (const action of speechActions) delete action.audioId;
|
||||
await rollBackFreshNarration();
|
||||
throw error;
|
||||
}
|
||||
|
||||
if (failedCount > 0) {
|
||||
await removeFreshTtsAllocations(freshAllocations);
|
||||
for (const action of speechActions) delete action.audioId;
|
||||
await rollBackFreshNarration();
|
||||
}
|
||||
|
||||
return {
|
||||
|
||||
@@ -0,0 +1,236 @@
|
||||
'use client';
|
||||
|
||||
/**
|
||||
* The one client-side sequence that turns bytes into something a shared
|
||||
* document may name.
|
||||
*
|
||||
* Three paths reach the asset pool from this browser — the media generation
|
||||
* pass, narration adoption, and fresh TTS synthesis — and all three run the
|
||||
* same four steps in the same order: store the bytes at the pool seam, write
|
||||
* the allocated id back into the document, mirror the bytes locally under that
|
||||
* id, and decide what a refusal means. The first three steps were already
|
||||
* spelled the same way at each caller. The fourth was not, and that is what
|
||||
* this module exists to fix: the media pass kept the bytes a full store
|
||||
* refused, adoption had nothing to keep because the bytes were already local,
|
||||
* and fresh TTS threw them away — re-billing the provider on every later
|
||||
* attempt for audio this browser had already paid for. Fixing that in place
|
||||
* would have made a fourth variant, so the sequence and its refusal rule live
|
||||
* here instead.
|
||||
*
|
||||
* What this owns:
|
||||
*
|
||||
* - The pool write, through `putAsset` and therefore through its stage seam.
|
||||
* - The classification of a failure as "the store had no room for THIS write"
|
||||
* versus anything else. It was written twice before, once per caller, under
|
||||
* slightly different conditions.
|
||||
* - The retention rule: a room refusal must leave the bytes in this browser's
|
||||
* local table under the key the retry path looks for, so that a later attempt
|
||||
* re-uploads them instead of buying them again. Callers whose bytes are
|
||||
* already there (adoption reads the very row it would write) pass no `retain`
|
||||
* and say so.
|
||||
* - The order of the remaining steps, and which of them may fail the commit. A
|
||||
* reference reaches the document only after the pool answered with an id, and
|
||||
* the local mirror is written only after the write-back, so the document can
|
||||
* never name bytes that were not stored and the cache can never outlive a
|
||||
* write-back that did not happen. The mirror itself is best-effort here
|
||||
* rather than at each caller: by then the bytes are stored and the document
|
||||
* names them, so a cache write costs a re-download at worst.
|
||||
*
|
||||
* What this deliberately does NOT own:
|
||||
*
|
||||
* - The store-full marker. Clearing it lives inside `putAsset`'s stage seam and
|
||||
* nowhere else, so a successful write retires the course's "no room" note in
|
||||
* one place rather than one per caller; and *setting* it is a claim that
|
||||
* calling a provider for this course is a waste of money, which only the
|
||||
* paths that spend provider money may make. Adoption spends none, so it must
|
||||
* not write the marker — a fact several review rounds re-established. This
|
||||
* module therefore neither reads nor writes it.
|
||||
* - The write-back funnel and the local table. Generated media and narration
|
||||
* carry different references in different document shapes and mirror into
|
||||
* different tables, and each already has exactly one funnel
|
||||
* (`persistGeneratedMediaReference`, `persistNarrationReference`). Inventing a
|
||||
* third would be the opposite of this change, so the funnel and the mirror are
|
||||
* supplied by the caller and merely sequenced here.
|
||||
*/
|
||||
import type { AssetMeta } from '@openmaic/dsl';
|
||||
|
||||
import { createLogger } from '@/lib/logger';
|
||||
import { putAsset } from '@/lib/media/asset-pool';
|
||||
import { ASSET_QUOTA_EXCEEDED, isStorageFullFailure } from '@/lib/media/media-failure';
|
||||
|
||||
const log = createLogger('PoolCommit');
|
||||
|
||||
/**
|
||||
* Bytes a full store refused, handed back so nothing has to be bought twice.
|
||||
*
|
||||
* The caller gets them whether or not it supplied a `retain` sink: the media
|
||||
* pass carries them out to the failure record it writes around them, which is
|
||||
* the row that is also the only copy a pre-server-backed course has of its own
|
||||
* media.
|
||||
*/
|
||||
export interface RefusedPoolBytes {
|
||||
/** The key the retry path looks for these bytes under. */
|
||||
readonly slot: string;
|
||||
readonly bytes: Blob;
|
||||
readonly mimeType: string;
|
||||
/** The pool's own error, for the caller's record and log. */
|
||||
readonly error: unknown;
|
||||
}
|
||||
|
||||
/**
|
||||
* What one commit settled as.
|
||||
*
|
||||
* `refused-retained` is returned only once the bytes are somewhere a later
|
||||
* attempt can find them, which is what makes the name true rather than a hope.
|
||||
* Neither refusal outcome throws: a refusal is an answer, and every caller has
|
||||
* a different thing to do with it.
|
||||
*/
|
||||
export type PoolCommitOutcome<TPlacement> =
|
||||
| {
|
||||
readonly status: 'stored';
|
||||
readonly assetId: string;
|
||||
/** Whatever the caller's write-back reported. */
|
||||
readonly placement: TPlacement;
|
||||
}
|
||||
| {
|
||||
readonly status: 'refused-retained';
|
||||
readonly code: typeof ASSET_QUOTA_EXCEEDED;
|
||||
readonly error: unknown;
|
||||
readonly refused: RefusedPoolBytes;
|
||||
}
|
||||
| { readonly status: 'failed'; readonly error: unknown };
|
||||
|
||||
export interface PoolCommitPlan<TPlacement> {
|
||||
/**
|
||||
* The course these bytes belong to.
|
||||
*
|
||||
* Handed to the pool seam, which is where a write that goes through retires
|
||||
* this course's "no room" note. Optional only because one caller (fresh TTS
|
||||
* outside a course) genuinely has no course to retire it for.
|
||||
*/
|
||||
readonly stageId?: string;
|
||||
/**
|
||||
* The placeholder or derived key this element's bytes are known by locally.
|
||||
*
|
||||
* `gen_img_3` for a generation placeholder, `tts_s2_action_…` for a narration
|
||||
* key. It is not sent to the pool; it is the address a retained refusal is
|
||||
* kept under and the address a later attempt reads back.
|
||||
*/
|
||||
readonly slot: string;
|
||||
readonly bytes: Blob;
|
||||
readonly mimeType: string;
|
||||
/** Extra asset metadata, e.g. `durationSeconds` for audio. */
|
||||
readonly meta?: Readonly<Omit<AssetMeta, 'contentType'>>;
|
||||
/**
|
||||
* Keep the bytes under `slot`, so a later attempt re-uploads rather than
|
||||
* re-synthesizes. Called on a room refusal and only then.
|
||||
*
|
||||
* Omitted by a caller whose bytes are already under `slot` in its own table —
|
||||
* adoption reads exactly the row it would write — and by a caller carrying
|
||||
* them out to a record it writes itself. A caller that omits it is stating
|
||||
* that the bytes survive the refusal, not that they do not matter.
|
||||
*
|
||||
* A sink that cannot keep the bytes must REJECT rather than swallow: the
|
||||
* commit then reports `failed`, not `refused-retained`. That is the whole
|
||||
* value of the outcome's name — a caller reading `refused-retained` goes on
|
||||
* to stamp `slot` into something durable, and a stamp naming bytes no local
|
||||
* table holds is a reference that resolves to nothing for the rest of the
|
||||
* course's life.
|
||||
*/
|
||||
readonly retain?: (refused: RefusedPoolBytes) => Promise<void>;
|
||||
/**
|
||||
* Write the allocated id into the document through this reference family's
|
||||
* existing funnel, and report whatever the caller needs to know afterwards.
|
||||
*
|
||||
* Errors are the caller's: they propagate out of `commitToPool` untouched,
|
||||
* because a write-back failure carries information (whether the allocation
|
||||
* was retained) that only the caller can act on.
|
||||
*/
|
||||
readonly writeBack: (assetId: string) => Promise<TPlacement>;
|
||||
/**
|
||||
* Mirror the bytes locally under the allocated id.
|
||||
*
|
||||
* Best-effort, and enforced here rather than left to each caller's own
|
||||
* `catch`: by this point the bytes are in the pool and the document already
|
||||
* names them, so a cache write that fails costs a re-download and nothing
|
||||
* else. A commit must not be reported as failed over it.
|
||||
*/
|
||||
readonly mirror: (assetId: string, placement: TPlacement) => Promise<void>;
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether this failure is the store saying it had no room for these bytes.
|
||||
*
|
||||
* The refusal reaches this browser as the asset client's error, whose `code` is
|
||||
* the one the storage contract puts in the 507 response body. Matching on the
|
||||
* code rather than on the client's error class is deliberate: the class is not
|
||||
* always the one this bundle imported, while the code is the part of the
|
||||
* contract that crosses every boundary.
|
||||
*
|
||||
* `errorCode` — the field the generation routes' own error class carries — is
|
||||
* deliberately not read. That class is raised by the image and video API calls
|
||||
* and never by a pool write, so accepting its shape here would widen this
|
||||
* predicate past anything `putAsset` can throw, on a field whose values come
|
||||
* from a different contract.
|
||||
*/
|
||||
function poolRefusedForRoom(error: unknown): boolean {
|
||||
if (typeof error !== 'object' || error === null) return false;
|
||||
const code = (error as { code?: unknown }).code;
|
||||
return typeof code === 'string' && isStorageFullFailure(code);
|
||||
}
|
||||
|
||||
/**
|
||||
* Store bytes in the pool and point the document at them.
|
||||
*
|
||||
* Returns rather than throws for every refusal shape; a write-back that fails
|
||||
* still throws, for the reason given on `writeBack`. A mirror that fails does
|
||||
* neither: the commit already happened.
|
||||
*/
|
||||
export async function commitToPool<TPlacement>(
|
||||
plan: PoolCommitPlan<TPlacement>,
|
||||
): Promise<PoolCommitOutcome<TPlacement>> {
|
||||
let assetId: string;
|
||||
try {
|
||||
assetId = await putAsset(
|
||||
plan.bytes,
|
||||
{ contentType: plan.mimeType, ...plan.meta },
|
||||
// The stage goes to the seam, which is the one place a successful write
|
||||
// retires this course's "no room" note. A write-back that fails for its
|
||||
// own reasons afterwards therefore does not leave the course standing
|
||||
// down.
|
||||
{ ...(plan.stageId ? { stageId: plan.stageId } : {}) },
|
||||
);
|
||||
} catch (error) {
|
||||
if (!poolRefusedForRoom(error)) return { status: 'failed', error };
|
||||
const refused: RefusedPoolBytes = {
|
||||
slot: plan.slot,
|
||||
bytes: plan.bytes,
|
||||
mimeType: plan.mimeType,
|
||||
error,
|
||||
};
|
||||
// Awaited before the outcome is reported, so `refused-retained` is a
|
||||
// statement about what is on disk rather than about what was scheduled --
|
||||
// and a sink that could not keep the bytes demotes the outcome, so no
|
||||
// caller stamps a key nothing can be read back by.
|
||||
if (plan.retain) {
|
||||
try {
|
||||
await plan.retain(refused);
|
||||
} catch (retentionError) {
|
||||
log.warn(`Could not keep the bytes refused for ${plan.slot}:`, retentionError);
|
||||
return { status: 'failed', error: retentionError };
|
||||
}
|
||||
}
|
||||
return { status: 'refused-retained', code: ASSET_QUOTA_EXCEEDED, error, refused };
|
||||
}
|
||||
|
||||
const placement = await plan.writeBack(assetId);
|
||||
try {
|
||||
await plan.mirror(assetId, placement);
|
||||
} catch (error) {
|
||||
// The bytes are stored and the document names them. A cache this browser
|
||||
// could not write costs a re-download, never the media, so the commit is
|
||||
// still a commit.
|
||||
log.warn(`Local mirror failed for ${assetId}:`, error);
|
||||
}
|
||||
return { status: 'stored', assetId, placement };
|
||||
}
|
||||
+159
-105
@@ -25,7 +25,7 @@ import { mayGenerateForStage } from '@/lib/classroom/generation-permission';
|
||||
import { db, mediaFileKey, type MediaFileRecord } from '@/lib/utils/database';
|
||||
import type { SceneOutline } from '@/lib/types/generation';
|
||||
import type { MediaGenerationRequest } from '@/lib/media/types';
|
||||
import { putAsset } from '@/lib/media/asset-pool';
|
||||
import { commitToPool } from '@/lib/media/commit-to-pool';
|
||||
import {
|
||||
ASSET_QUOTA_EXCEEDED,
|
||||
isRetryableMediaFailure,
|
||||
@@ -140,18 +140,22 @@ interface RefusedMediaBytes {
|
||||
* it is reserved for refusals a retry cannot change — a provider's content
|
||||
* decision, a disabled generation setting, and a full asset store.
|
||||
*
|
||||
* The quota refusal reaches this browser as the asset client's error, whose
|
||||
* `code` is the one the storage contract puts in the response body. Matching
|
||||
* on the code rather than on the client's error class is deliberate: the class
|
||||
* is not always the one this bundle imported, while the code is the part of
|
||||
* the contract that crosses every boundary. Everything else stays retryable,
|
||||
* because everything else might work next time.
|
||||
* Exactly two error shapes can carry one, and each is named rather than probed
|
||||
* for. A generation route's refusal arrives as `MediaApiError`, whose
|
||||
* `errorCode` is the route's own; a full store arrives as the
|
||||
* `MediaStorageRefusalError` the commit raises from the pool primitive's
|
||||
* refusal outcome, whose `code` the primitive already matched against the
|
||||
* storage contract. Nothing else reaching this catch classifies a pool write:
|
||||
* the primitive owns that test now, so the structural "any object with a
|
||||
* `code`" probe this used to end with could no longer be reached by a pool
|
||||
* error that was not already wrapped, and a generalization nothing can take is
|
||||
* one more shape to keep true. Everything else stays retryable, because
|
||||
* everything else might work next time.
|
||||
*/
|
||||
function mediaFailureCode(error: unknown): string | undefined {
|
||||
if (error instanceof MediaApiError) return error.errorCode;
|
||||
if (typeof error !== 'object' || error === null) return undefined;
|
||||
const code = (error as { code?: unknown }).code;
|
||||
return code === ASSET_QUOTA_EXCEEDED ? ASSET_QUOTA_EXCEEDED : undefined;
|
||||
if (error instanceof MediaStorageRefusalError) return error.code;
|
||||
return undefined;
|
||||
}
|
||||
|
||||
function createAbortError(): Error {
|
||||
@@ -540,15 +544,66 @@ function documentSkipIndex(stageId: string): GeneratedMediaDocumentIndex | undef
|
||||
return indexGeneratedMediaReferences({ stage, scenes, generationComplete });
|
||||
}
|
||||
|
||||
/** What the generated-media write-back leaves behind for the task to read. */
|
||||
interface GeneratedMediaPlacement {
|
||||
readonly result: MediaReferenceWriteBackResult;
|
||||
readonly objectUrl: string;
|
||||
readonly posterObjectUrl?: string;
|
||||
readonly posterAssetId?: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Allocate a video's poster, and never let it cost the video.
|
||||
*
|
||||
* A poster is decorative: it is written only into a slot that has none of its
|
||||
* own. Letting its upload fail the commit would discard a stored video and send
|
||||
* the retry to submit the most expensive job in the system a second time, so a
|
||||
* poster failure costs the poster and nothing else.
|
||||
*
|
||||
* It goes through the same primitive as everything else this browser puts in
|
||||
* the pool, with both later steps empty, because a poster is a bare allocation:
|
||||
* it holds no reference of its own in the document (it rides in the slot its
|
||||
* video fills) and no row of its own in the cache (it rides in its video's
|
||||
* row). It retains nothing on a refusal for the same reason — there is no key a
|
||||
* later attempt would look for it under, and its video's row already carries
|
||||
* the bytes.
|
||||
*/
|
||||
async function allocatePoster(args: {
|
||||
elementId: string;
|
||||
stageId: string;
|
||||
posterBlob: Blob;
|
||||
posterMimeType?: string;
|
||||
}): Promise<string | undefined> {
|
||||
const outcome = await commitToPool<void>({
|
||||
stageId: args.stageId,
|
||||
slot: args.elementId,
|
||||
bytes: args.posterBlob,
|
||||
mimeType: args.posterMimeType ?? args.posterBlob.type,
|
||||
writeBack: async () => undefined,
|
||||
mirror: async () => undefined,
|
||||
});
|
||||
if (outcome.status === 'stored') return outcome.assetId;
|
||||
log.warn(`Poster allocation failed for ${args.elementId}; keeping the video:`, outcome.error);
|
||||
return undefined;
|
||||
}
|
||||
|
||||
/**
|
||||
* Commit generated bytes under server-backed persistence: pool first, then the
|
||||
* document, then the local cache, and only then the task.
|
||||
*
|
||||
* Order is the contract. A reference reaches the document only after `put`
|
||||
* returned an id, so the document can never name bytes that were not stored;
|
||||
* and a failure anywhere before the write-back leaves the placeholder in place
|
||||
* with the provider called exactly once, so the retry happens on the next
|
||||
* owner load rather than inside this run.
|
||||
* Order is the contract, and it is `commitToPool` that holds it now. A
|
||||
* reference reaches the document only after `put` returned an id, so the
|
||||
* document can never name bytes that were not stored; and a failure anywhere
|
||||
* before the write-back leaves the placeholder in place with the provider
|
||||
* called exactly once, so the retry happens on the next owner load rather than
|
||||
* inside this run.
|
||||
*
|
||||
* The refusal this path wants is the primitive's `refused-retained`, translated
|
||||
* back into the error that carries the bytes to the failure record
|
||||
* `generateSingleMedia` writes around them. That record is where this caller's
|
||||
* retention lives: the row holds the refusal's message and code alongside the
|
||||
* bytes, and only the catch that sees the failure knows those — which is why no
|
||||
* `retain` sink is handed to the primitive here.
|
||||
*/
|
||||
async function commitPooledMedia(args: {
|
||||
req: MediaGenerationRequest;
|
||||
@@ -561,102 +616,101 @@ async function commitPooledMedia(args: {
|
||||
}): Promise<void> {
|
||||
const { req, stageId, paramsJson, blob, mimeType, posterBlob, posterMimeType } = args;
|
||||
|
||||
let assetId: string;
|
||||
try {
|
||||
// The stage is handed to the seam, which is where a successful write
|
||||
// retires this course's "no room" note -- one place rather than one per
|
||||
// caller. A write-back that fails for its own reasons afterwards therefore
|
||||
// does not leave the course standing down.
|
||||
assetId = await putAsset(blob, { contentType: mimeType }, { stageId });
|
||||
} catch (error) {
|
||||
// A full store keeps its bytes. Everything else is an ordinary failure and
|
||||
// is retried from the provider, as it always was.
|
||||
if (!isStorageFullFailure(mediaFailureCode(error))) throw error;
|
||||
throw new MediaStorageRefusalError(error, ASSET_QUOTA_EXCEEDED, {
|
||||
const outcome = await commitToPool<GeneratedMediaPlacement>({
|
||||
stageId,
|
||||
slot: req.elementId,
|
||||
bytes: blob,
|
||||
mimeType,
|
||||
writeBack: async (assetId) => {
|
||||
const posterAssetId = posterBlob
|
||||
? await allocatePoster({
|
||||
elementId: req.elementId,
|
||||
stageId,
|
||||
posterBlob,
|
||||
...(posterMimeType ? { posterMimeType } : {}),
|
||||
})
|
||||
: undefined;
|
||||
|
||||
// Minted before the write-back, not after it, so the allocation the
|
||||
// funnel may park in the same turn as its decision is complete: an entry
|
||||
// drained a microtask later must carry the bytes this tab can already
|
||||
// render.
|
||||
const objectUrl = URL.createObjectURL(blob);
|
||||
const posterObjectUrl = posterBlob ? URL.createObjectURL(posterBlob) : undefined;
|
||||
const allocation: PendingMediaAllocation = {
|
||||
stageId,
|
||||
placeholderRef: req.elementId,
|
||||
assetId,
|
||||
posterAssetId,
|
||||
objectUrl,
|
||||
posterObjectUrl,
|
||||
};
|
||||
|
||||
try {
|
||||
const result = await persistGeneratedMediaReference(allocation);
|
||||
return { result, objectUrl, posterObjectUrl, posterAssetId };
|
||||
} catch (error) {
|
||||
// The funnel places or parks the allocation whenever the ids could be
|
||||
// referenced, and says so. Reclaiming is for the one case where nothing
|
||||
// can possibly hold them — otherwise a lost response would take an
|
||||
// asset the persisted document already names.
|
||||
if (error instanceof MediaReferenceWriteBackError && error.allocationRetained) throw error;
|
||||
URL.revokeObjectURL(objectUrl);
|
||||
if (posterObjectUrl) URL.revokeObjectURL(posterObjectUrl);
|
||||
// Forgetting is all a browser does here. The record outlives the parked
|
||||
// queue, so leaving it behind would let a later save stamp an id whose
|
||||
// bytes nothing references — and the placeholder it replaced would be
|
||||
// gone with it, which reads as "already generated" and stops any retry.
|
||||
//
|
||||
// The bytes themselves are NOT deleted. Asset deletion is refused to
|
||||
// every browser, because the principal it scopes to is shared and would
|
||||
// let any caller destroy another author's media. The entry does not
|
||||
// need a browser to release it: no document write ever commits this
|
||||
// allocation, so it stays pending and the collector's entry pass takes
|
||||
// it once ASSET_PENDING_TTL_MS has elapsed.
|
||||
forgetMediaAllocation(stageId, req.elementId);
|
||||
throw error;
|
||||
}
|
||||
},
|
||||
// Local cache for this tab only. The document already points at the pool,
|
||||
// so a failed cache write costs a re-download, never the media.
|
||||
mirror: async (assetId) => {
|
||||
await db.mediaFiles
|
||||
.put({
|
||||
id: mediaFileKey(stageId, assetId),
|
||||
stageId,
|
||||
type: req.type,
|
||||
blob,
|
||||
mimeType,
|
||||
size: blob.size,
|
||||
poster: posterBlob,
|
||||
placeholderRef: req.elementId,
|
||||
prompt: req.prompt,
|
||||
params: paramsJson,
|
||||
createdAt: Date.now(),
|
||||
})
|
||||
.catch((error: unknown) => {
|
||||
log.warn(`Local media cache write failed for ${assetId}:`, error);
|
||||
});
|
||||
},
|
||||
});
|
||||
|
||||
if (outcome.status === 'refused-retained') {
|
||||
// A full store keeps its bytes: they leave through the error so the failure
|
||||
// record is written around them rather than over them.
|
||||
throw new MediaStorageRefusalError(outcome.error, outcome.code, {
|
||||
blob,
|
||||
mimeType,
|
||||
...(posterBlob ? { poster: posterBlob } : {}),
|
||||
...(posterBlob && posterMimeType ? { posterMimeType } : {}),
|
||||
});
|
||||
}
|
||||
// A poster is decorative: it is written only into a slot that has none of its
|
||||
// own. Letting its upload fail the commit would discard a stored video and
|
||||
// send the retry to submit the most expensive job in the system a second
|
||||
// time, so a poster failure costs the poster and nothing else.
|
||||
let posterAssetId: string | undefined;
|
||||
if (posterBlob) {
|
||||
try {
|
||||
posterAssetId = await putAsset(
|
||||
posterBlob,
|
||||
{ contentType: posterMimeType ?? posterBlob.type },
|
||||
{ stageId },
|
||||
);
|
||||
} catch (error) {
|
||||
log.warn(`Poster allocation failed for ${req.elementId}; keeping the video:`, error);
|
||||
}
|
||||
}
|
||||
// Everything else is an ordinary failure and is retried from the provider, as
|
||||
// it always was.
|
||||
if (outcome.status === 'failed') throw outcome.error;
|
||||
|
||||
// Minted before the write-back, not after it, so the allocation the funnel
|
||||
// may park in the same turn as its decision is complete: an entry drained a
|
||||
// microtask later must carry the bytes this tab can already render.
|
||||
const objectUrl = URL.createObjectURL(blob);
|
||||
const posterObjectUrl = posterBlob ? URL.createObjectURL(posterBlob) : undefined;
|
||||
const allocation: PendingMediaAllocation = {
|
||||
stageId,
|
||||
placeholderRef: req.elementId,
|
||||
assetId,
|
||||
posterAssetId,
|
||||
objectUrl,
|
||||
posterObjectUrl,
|
||||
};
|
||||
|
||||
let outcome: MediaReferenceWriteBackResult;
|
||||
try {
|
||||
outcome = await persistGeneratedMediaReference(allocation);
|
||||
} catch (error) {
|
||||
// The funnel places or parks the allocation whenever the ids could be
|
||||
// referenced, and says so. Reclaiming is for the one case where nothing can
|
||||
// possibly hold them — otherwise a lost response would take an asset the
|
||||
// persisted document already names.
|
||||
if (error instanceof MediaReferenceWriteBackError && error.allocationRetained) throw error;
|
||||
URL.revokeObjectURL(objectUrl);
|
||||
if (posterObjectUrl) URL.revokeObjectURL(posterObjectUrl);
|
||||
// Forgetting is all a browser does here. The record outlives the parked
|
||||
// queue, so leaving it behind would let a later save stamp an id whose
|
||||
// bytes nothing references — and the placeholder it replaced would be gone
|
||||
// with it, which reads as "already generated" and stops any retry.
|
||||
//
|
||||
// The bytes themselves are NOT deleted. Asset deletion is refused to every
|
||||
// browser, because the principal it scopes to is shared and would let any
|
||||
// caller destroy another author's media. The entry does not need a browser
|
||||
// to release it: no document write ever commits this allocation, so it
|
||||
// stays pending and the collector's entry pass takes it once
|
||||
// ASSET_PENDING_TTL_MS has elapsed.
|
||||
forgetMediaAllocation(stageId, req.elementId);
|
||||
throw error;
|
||||
}
|
||||
|
||||
// Local cache for this tab only. The document already points at the pool, so
|
||||
// a failed cache write costs a re-download, never the media.
|
||||
await db.mediaFiles
|
||||
.put({
|
||||
id: mediaFileKey(stageId, assetId),
|
||||
stageId,
|
||||
type: req.type,
|
||||
blob,
|
||||
mimeType,
|
||||
size: blob.size,
|
||||
poster: posterBlob,
|
||||
placeholderRef: req.elementId,
|
||||
prompt: req.prompt,
|
||||
params: paramsJson,
|
||||
createdAt: Date.now(),
|
||||
})
|
||||
.catch((error: unknown) => {
|
||||
log.warn(`Local media cache write failed for ${assetId}:`, error);
|
||||
});
|
||||
|
||||
if (outcome === 'held') {
|
||||
const { result, objectUrl, posterObjectUrl, posterAssetId } = outcome.placement;
|
||||
if (result === 'held') {
|
||||
// The slide this media belongs to has not been built yet, which during a
|
||||
// first pass is the ordinary case rather than an edge one. The funnel holds
|
||||
// the allocation; the task stays keyed by the placeholder the document
|
||||
@@ -668,7 +722,7 @@ async function commitPooledMedia(args: {
|
||||
|
||||
useMediaGenerationStore
|
||||
.getState()
|
||||
.rekeyDone(req.elementId, assetId, objectUrl, posterObjectUrl, posterAssetId);
|
||||
.rekeyDone(req.elementId, outcome.assetId, objectUrl, posterObjectUrl, posterAssetId);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -0,0 +1,368 @@
|
||||
/**
|
||||
* Narration a full store refused is kept, so nobody pays for it twice.
|
||||
*
|
||||
* This is the defect #1467 names. The media pass has kept the bytes a full
|
||||
* store refused since it learned to: they are already paid for, the document
|
||||
* still carries the placeholder, and a later Retry re-attempts the upload
|
||||
* rather than the generation. Fresh TTS synthesis did the opposite — it dropped
|
||||
* the freshly synthesized clip before the local cache write — so every later
|
||||
* attempt called the provider again for audio this browser had already bought.
|
||||
*
|
||||
* The fix gives the TTS path the same contract through the shared commit
|
||||
* primitive, and this suite is the proof of it end to end: the refused clip is
|
||||
* kept under its derived key, the action is stamped with that key, and the next
|
||||
* load's narration adoption converts it with exactly one pool write and no
|
||||
* provider call at all.
|
||||
*
|
||||
* Both halves are loaded for real, because the claim is about what one leaves
|
||||
* for the other. The pool is doubled at the store rather than at `putAsset`, so
|
||||
* the real seam runs.
|
||||
*/
|
||||
import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest';
|
||||
|
||||
const mocks = vi.hoisted(() => ({
|
||||
getCurrentModelConfig: vi.fn(),
|
||||
settingsState: vi.fn(),
|
||||
audioGet: vi.fn(),
|
||||
audioPut: vi.fn(),
|
||||
audioDelete: vi.fn(),
|
||||
mediaPut: vi.fn(),
|
||||
mediaGet: vi.fn(),
|
||||
mediaDelete: vi.fn(),
|
||||
poolPut: vi.fn(),
|
||||
mutateDocument: vi.fn(),
|
||||
saveStageData: vi.fn(),
|
||||
saveStageDataIncremental: vi.fn(),
|
||||
isTTSProviderEnabled: vi.fn(),
|
||||
pickNarratorAgent: vi.fn(),
|
||||
resolveAgentVoiceOptions: vi.fn(),
|
||||
listAgents: vi.fn(),
|
||||
toastWarning: vi.fn(),
|
||||
serverBacked: vi.fn(),
|
||||
}));
|
||||
|
||||
vi.mock('@/lib/utils/model-config', () => ({
|
||||
getCurrentModelConfig: mocks.getCurrentModelConfig,
|
||||
}));
|
||||
vi.mock('@/lib/store/settings', () => ({
|
||||
useSettingsStore: { getState: mocks.settingsState },
|
||||
}));
|
||||
vi.mock('@/lib/document-store', () => ({ mutateDocument: mocks.mutateDocument }));
|
||||
vi.mock('@/lib/utils/stage-storage', () => ({
|
||||
saveStageData: mocks.saveStageData,
|
||||
saveStageDataIncremental: mocks.saveStageDataIncremental,
|
||||
}));
|
||||
vi.mock('@/lib/utils/database', () => ({
|
||||
mediaFileKey: (stageId: string, ref: string) => `${stageId}:${ref}`,
|
||||
db: {
|
||||
audioFiles: { get: mocks.audioGet, put: mocks.audioPut, delete: mocks.audioDelete },
|
||||
mediaFiles: {
|
||||
put: mocks.mediaPut,
|
||||
get: mocks.mediaGet,
|
||||
delete: mocks.mediaDelete,
|
||||
where: () => ({ equals: () => ({ toArray: async () => [] }) }),
|
||||
},
|
||||
},
|
||||
}));
|
||||
vi.mock('@/lib/media/asset-pool-config', async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import('@/lib/media/asset-pool-config')>();
|
||||
return {
|
||||
...actual,
|
||||
resolveConfiguredAssetPoolStore: () =>
|
||||
({ put: mocks.poolPut }) as unknown as import('@/lib/media/asset-pool-config').AssetPoolStore,
|
||||
};
|
||||
});
|
||||
vi.mock('@/lib/persistence/media-persistence', () => ({
|
||||
isServerBackedMediaPersistence: mocks.serverBacked,
|
||||
}));
|
||||
vi.mock('@/lib/audio/provider-enablement', () => ({
|
||||
isTTSProviderEnabled: mocks.isTTSProviderEnabled,
|
||||
}));
|
||||
vi.mock('@/lib/audio/agent-voice', () => ({
|
||||
pickNarratorAgent: mocks.pickNarratorAgent,
|
||||
resolveAgentVoiceOptions: mocks.resolveAgentVoiceOptions,
|
||||
}));
|
||||
vi.mock('@/lib/orchestration/registry/store', () => ({
|
||||
useAgentRegistry: { getState: () => ({ listAgents: mocks.listAgents }) },
|
||||
}));
|
||||
vi.mock('sonner', () => ({ toast: { warning: mocks.toastWarning } }));
|
||||
|
||||
const mockFetch = vi.fn() as Mock;
|
||||
vi.stubGlobal('fetch', mockFetch);
|
||||
|
||||
import { adoptCachedNarration } from '@/lib/audio/adopt-cached-narration';
|
||||
import { generateAndStoreTTS, generateTTSForScene } from '@/lib/hooks/use-scene-generator';
|
||||
import {
|
||||
noteStageGenerationOwnership,
|
||||
resetGenerationPermissionsForTests,
|
||||
} from '@/lib/classroom/generation-permission';
|
||||
import { useStageStore } from '@/lib/store/stage';
|
||||
import type { Scene } from '@/lib/types/stage';
|
||||
|
||||
const stageId = 'refused-narration-stage';
|
||||
/** What `generateTTSForScene` builds for scene order 1, action `speech-1`. */
|
||||
const derivedRef = 'tts_s1_speech-1';
|
||||
|
||||
/** What the store answers when it has no room for these bytes. */
|
||||
function quotaRefusal(): Error {
|
||||
return Object.assign(new Error('asset quota exceeded for this principal'), {
|
||||
status: 507,
|
||||
code: 'ASSET_QUOTA_EXCEEDED',
|
||||
});
|
||||
}
|
||||
|
||||
function ttsResponse() {
|
||||
return {
|
||||
ok: true,
|
||||
status: 200,
|
||||
statusText: 'OK',
|
||||
json: async () => ({ success: true, base64: btoa('narration-bytes'), format: 'wav' }),
|
||||
};
|
||||
}
|
||||
|
||||
function sceneWithOneLine(): Scene {
|
||||
return {
|
||||
id: 'scene-1',
|
||||
stageId,
|
||||
title: 'Scene',
|
||||
order: 1,
|
||||
type: 'slide',
|
||||
content: { type: 'slide', canvas: { id: 'slide-1', elements: [] } },
|
||||
actions: [{ id: 'speech-1', type: 'speech', text: 'Welcome' }],
|
||||
} as unknown as Scene;
|
||||
}
|
||||
|
||||
/** Two lines, so "this clip" and "the one next to it" are distinguishable. */
|
||||
function sceneWithTwoLines(): Scene {
|
||||
return {
|
||||
id: 'scene-1',
|
||||
stageId,
|
||||
title: 'Scene',
|
||||
order: 1,
|
||||
type: 'slide',
|
||||
content: { type: 'slide', canvas: { id: 'slide-1', elements: [] } },
|
||||
actions: [
|
||||
{ id: 'speech-1', type: 'speech', text: 'Welcome' },
|
||||
{ id: 'speech-2', type: 'speech', text: 'And then' },
|
||||
],
|
||||
} as unknown as Scene;
|
||||
}
|
||||
|
||||
function audioIdOf(scene: Scene, index = 0): string | undefined {
|
||||
return (scene as unknown as { actions: Array<{ audioId?: string }> }).actions[index]?.audioId;
|
||||
}
|
||||
|
||||
/** The local audio table, modelled so one step can read what the last wrote. */
|
||||
function modelAudioTable(): Map<string, Record<string, unknown>> {
|
||||
const rows = new Map<string, Record<string, unknown>>();
|
||||
mocks.audioPut.mockImplementation(async (row: Record<string, unknown>) => {
|
||||
rows.set(row.id as string, row);
|
||||
});
|
||||
mocks.audioGet.mockImplementation(async (id: string) => rows.get(id));
|
||||
mocks.audioDelete.mockImplementation(async (id: string) => {
|
||||
rows.delete(id);
|
||||
});
|
||||
return rows;
|
||||
}
|
||||
|
||||
/** Run the narration funnel against a document holding the same reference. */
|
||||
function serveDocument(scenes: Scene[]): Mock {
|
||||
const putScene = vi.fn().mockResolvedValue(undefined);
|
||||
mocks.mutateDocument.mockImplementation(
|
||||
async (_stageId: string, work: (document: unknown, store: unknown) => Promise<void>) => {
|
||||
await work({ scenes, stage: { id: stageId } }, { putScene, putStage: vi.fn() });
|
||||
},
|
||||
);
|
||||
return putScene;
|
||||
}
|
||||
|
||||
describe('narration refused for want of room', () => {
|
||||
beforeEach(() => {
|
||||
resetGenerationPermissionsForTests();
|
||||
mockFetch.mockReset();
|
||||
mocks.poolPut.mockReset();
|
||||
mocks.audioGet.mockReset().mockResolvedValue(undefined);
|
||||
mocks.audioPut.mockReset().mockResolvedValue(undefined);
|
||||
mocks.audioDelete.mockReset().mockResolvedValue(undefined);
|
||||
mocks.mutateDocument.mockReset();
|
||||
mocks.saveStageData.mockReset().mockResolvedValue(undefined);
|
||||
mocks.saveStageDataIncremental.mockReset().mockResolvedValue(undefined);
|
||||
mocks.mediaPut.mockReset().mockResolvedValue(undefined);
|
||||
mocks.mediaGet.mockReset().mockResolvedValue(undefined);
|
||||
mocks.mediaDelete.mockReset().mockResolvedValue(undefined);
|
||||
mocks.serverBacked.mockReset().mockReturnValue(true);
|
||||
mocks.getCurrentModelConfig.mockReturnValue({});
|
||||
mocks.settingsState.mockReturnValue({
|
||||
imageProviderId: '',
|
||||
imageProvidersConfig: {},
|
||||
imageGenerationEnabled: false,
|
||||
videoProviderId: '',
|
||||
videoProvidersConfig: {},
|
||||
videoGenerationEnabled: false,
|
||||
ttsProviderId: 'server-tts',
|
||||
ttsProvidersConfig: { 'server-tts': { apiKey: 'tts-key', modelId: 'tts-model' } },
|
||||
ttsVoice: 'narrator',
|
||||
ttsSpeed: 1,
|
||||
});
|
||||
mocks.isTTSProviderEnabled.mockReturnValue(true);
|
||||
mocks.pickNarratorAgent.mockReturnValue(undefined);
|
||||
mocks.resolveAgentVoiceOptions.mockResolvedValue({});
|
||||
mocks.listAgents.mockReturnValue([]);
|
||||
mocks.toastWarning.mockReset();
|
||||
noteStageGenerationOwnership(stageId, 'owner');
|
||||
useStageStore.setState({ stage: { id: stageId, name: 'Course' } as never, scenes: [] });
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
useStageStore.setState({ stage: null, scenes: [] });
|
||||
resetGenerationPermissionsForTests();
|
||||
});
|
||||
|
||||
// The whole point, in one run: the provider is paid once, ever.
|
||||
it('keeps the billed clip and lets the next load upload it with no provider call', async () => {
|
||||
const rows = modelAudioTable();
|
||||
mocks.poolPut.mockRejectedValueOnce(quotaRefusal());
|
||||
mockFetch.mockResolvedValueOnce(ttsResponse());
|
||||
|
||||
const scene = sceneWithOneLine();
|
||||
await expect(generateTTSForScene(scene)).resolves.toMatchObject({
|
||||
success: true,
|
||||
failedCount: 0,
|
||||
});
|
||||
|
||||
// Refused, and kept: the bytes sit under the derived key, and the action
|
||||
// carries that key, exactly as a refused image leaves its placeholder in
|
||||
// the slide.
|
||||
expect(mocks.poolPut).toHaveBeenCalledTimes(1);
|
||||
expect(mockFetch).toHaveBeenCalledTimes(1);
|
||||
expect(audioIdOf(scene)).toBe(derivedRef);
|
||||
await expect((rows.get(derivedRef)?.blob as Blob).text()).resolves.toBe('narration-bytes');
|
||||
expect(rows.get(derivedRef)).toMatchObject({ stageId, format: 'wav', text: 'Welcome' });
|
||||
|
||||
// The next load. The ceiling has moved, so the store takes it now.
|
||||
mocks.poolPut.mockResolvedValue('ast_narration_allocated');
|
||||
const documentScenes = [scene];
|
||||
const putScene = serveDocument(documentScenes);
|
||||
useStageStore.setState({ scenes: [scene] as never });
|
||||
|
||||
await expect(adoptCachedNarration(stageId)).resolves.toEqual({ adopted: 1, unbacked: 0 });
|
||||
|
||||
// One pool write for the retry, and not one more provider call.
|
||||
expect(mocks.poolPut).toHaveBeenCalledTimes(2);
|
||||
expect(mockFetch).toHaveBeenCalledTimes(1);
|
||||
const [uploaded] = mocks.poolPut.mock.calls[1] as [Blob];
|
||||
await expect(uploaded.text()).resolves.toBe('narration-bytes');
|
||||
expect(putScene).toHaveBeenCalledTimes(1);
|
||||
expect(audioIdOf(documentScenes[0])).toBe('ast_narration_allocated');
|
||||
expect(audioIdOf(useStageStore.getState().scenes[0] as Scene)).toBe('ast_narration_allocated');
|
||||
});
|
||||
|
||||
// A refusal is not a synthesis failure. Counting it as one would pause the
|
||||
// whole deck at its first slide over one clip's storage.
|
||||
it('does not fail the scene, and does not re-synthesize within the same run', async () => {
|
||||
modelAudioTable();
|
||||
mocks.poolPut.mockRejectedValue(quotaRefusal());
|
||||
mockFetch.mockResolvedValue(ttsResponse());
|
||||
|
||||
const scene = sceneWithOneLine();
|
||||
await expect(generateTTSForScene(scene)).resolves.toEqual({ success: true, failedCount: 0 });
|
||||
expect(mockFetch).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
// A sibling line failing is not a reason to throw away bytes that are
|
||||
// already paid for. The scene-level rollback reclaims what it minted for this
|
||||
// scene; a retained refusal was never minted, and unstamping it would strand
|
||||
// the only copy of that clip where nothing will ever look for it again.
|
||||
it('keeps a retained refusal when the line next to it fails', async () => {
|
||||
const rows = modelAudioTable();
|
||||
mocks.poolPut.mockRejectedValue(quotaRefusal());
|
||||
mockFetch.mockResolvedValueOnce(ttsResponse()).mockResolvedValueOnce({
|
||||
ok: false,
|
||||
status: 503,
|
||||
statusText: 'unavailable',
|
||||
json: async () => ({ error: 'provider down' }),
|
||||
});
|
||||
|
||||
const scene = sceneWithTwoLines();
|
||||
await expect(generateTTSForScene(scene)).resolves.toMatchObject({
|
||||
success: false,
|
||||
failedCount: 1,
|
||||
});
|
||||
|
||||
// The refused line keeps both halves of the contract: its bytes and the
|
||||
// key adoption reads them back by.
|
||||
expect(audioIdOf(scene, 0)).toBe(derivedRef);
|
||||
expect(rows.has(derivedRef)).toBe(true);
|
||||
expect(mocks.audioDelete).not.toHaveBeenCalledWith(derivedRef);
|
||||
// The line that failed has nothing to keep.
|
||||
expect(audioIdOf(scene, 1)).toBeUndefined();
|
||||
});
|
||||
|
||||
// A clip the pool did take is an allocation this scene minted and nothing
|
||||
// else holds, so the rollback still reclaims its local copy.
|
||||
it('still rolls back a clip the pool accepted when a sibling fails', async () => {
|
||||
modelAudioTable();
|
||||
mocks.poolPut.mockResolvedValue('ast_narration_allocated');
|
||||
mockFetch.mockResolvedValueOnce(ttsResponse()).mockResolvedValueOnce({
|
||||
ok: false,
|
||||
status: 503,
|
||||
statusText: 'unavailable',
|
||||
json: async () => ({ error: 'provider down' }),
|
||||
});
|
||||
|
||||
const scene = sceneWithTwoLines();
|
||||
await expect(generateTTSForScene(scene)).resolves.toMatchObject({
|
||||
success: false,
|
||||
failedCount: 1,
|
||||
});
|
||||
|
||||
expect(mocks.audioDelete).toHaveBeenCalledWith('ast_narration_allocated');
|
||||
expect(audioIdOf(scene, 0)).toBeUndefined();
|
||||
});
|
||||
|
||||
// `refused-retained` is a statement about what is on disk. If the local table
|
||||
// refused the row too, a stamp would name bytes nothing can read back, for
|
||||
// the rest of the course's life.
|
||||
it('leaves the line unvoiced when the refused bytes cannot be kept locally', async () => {
|
||||
mocks.audioPut.mockRejectedValue(new Error('local quota exceeded'));
|
||||
mocks.poolPut.mockRejectedValue(quotaRefusal());
|
||||
mockFetch.mockResolvedValueOnce(ttsResponse());
|
||||
|
||||
const scene = sceneWithOneLine();
|
||||
await expect(generateTTSForScene(scene)).resolves.toMatchObject({
|
||||
success: true,
|
||||
failedCount: 0,
|
||||
});
|
||||
|
||||
expect(audioIdOf(scene)).toBeUndefined();
|
||||
expect(mocks.audioPut).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
// A refusal for room says the bytes do not fit. Everything else -- a dropped
|
||||
// connection, a 500 -- says nothing about a later attempt, so nothing is left
|
||||
// under a key a later load would read as adoptable narration; the line stays
|
||||
// unvoiced and retryable instead.
|
||||
it('keeps nothing when the pool fails for a reason that is not room', async () => {
|
||||
modelAudioTable();
|
||||
mocks.poolPut.mockRejectedValue(new Error('asset store unavailable'));
|
||||
mockFetch.mockResolvedValueOnce(ttsResponse());
|
||||
|
||||
await expect(generateAndStoreTTS('tts_s2_action_1', 'Hello class')).resolves.toBeNull();
|
||||
expect(mocks.audioPut).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
// Browser-only mode has no pool to refuse anything: document and audio share
|
||||
// one lifetime, and the derived key is a complete address.
|
||||
it('leaves browser-only narration exactly as it was', async () => {
|
||||
const rows = modelAudioTable();
|
||||
mocks.serverBacked.mockReturnValue(false);
|
||||
mockFetch.mockResolvedValueOnce(ttsResponse());
|
||||
|
||||
const scene = sceneWithOneLine();
|
||||
await expect(generateTTSForScene(scene)).resolves.toEqual({ success: true, failedCount: 0 });
|
||||
|
||||
expect(mocks.poolPut).not.toHaveBeenCalled();
|
||||
expect(audioIdOf(scene)).toBe(derivedRef);
|
||||
expect(rows.get(derivedRef)).toBeDefined();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,249 @@
|
||||
/**
|
||||
* The one client-side pool commit sequence, and its one refusal rule.
|
||||
*
|
||||
* Three paths reach the pool from this browser and each used to spell the
|
||||
* sequence — and, worse, the meaning of a refusal — for itself. This suite pins
|
||||
* what the shared primitive now guarantees for all of them: the order of the
|
||||
* four steps, the difference between "the store had no room" and any other
|
||||
* failure, and the promise `refused-retained` makes about where the bytes are
|
||||
* by the time the caller reads it.
|
||||
*
|
||||
* The pool is doubled at the store rather than at `putAsset`, deliberately: the
|
||||
* stage seam that retires a course's "store is full" note lives inside
|
||||
* `putAsset`, and this module must go through it rather than around it.
|
||||
*/
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest';
|
||||
|
||||
const mocks = vi.hoisted(() => ({
|
||||
poolPut: vi.fn(),
|
||||
}));
|
||||
|
||||
vi.mock('@/lib/media/asset-pool-config', async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import('@/lib/media/asset-pool-config')>();
|
||||
return {
|
||||
...actual,
|
||||
resolveConfiguredAssetPoolStore: () =>
|
||||
({ put: mocks.poolPut }) as unknown as import('@/lib/media/asset-pool-config').AssetPoolStore,
|
||||
};
|
||||
});
|
||||
|
||||
import { commitToPool } from '@/lib/media/commit-to-pool';
|
||||
import {
|
||||
isAssetStorageFull,
|
||||
markAssetStorageFull,
|
||||
setAssetStorageFullStoreForTests,
|
||||
} from '@/lib/media/asset-storage-full';
|
||||
|
||||
const stageId = 'commit-stage';
|
||||
|
||||
function quotaRefusal(): Error {
|
||||
return Object.assign(new Error('asset quota exceeded for this principal'), {
|
||||
status: 507,
|
||||
code: 'ASSET_QUOTA_EXCEEDED',
|
||||
});
|
||||
}
|
||||
|
||||
/** The device KV the storage-full marker lives in, in memory. */
|
||||
function memoryKv() {
|
||||
const entries = new Map<string, unknown>();
|
||||
return {
|
||||
get: async <T>(key: string) => (entries.get(key) as T) ?? null,
|
||||
set: async (key: string, value: unknown) => {
|
||||
entries.set(key, value);
|
||||
},
|
||||
remove: async (key: string) => {
|
||||
entries.delete(key);
|
||||
},
|
||||
keys: async (prefix = '') => [...entries.keys()].filter((key) => key.startsWith(prefix)),
|
||||
};
|
||||
}
|
||||
|
||||
function plan(overrides: Record<string, unknown> = {}) {
|
||||
return {
|
||||
stageId,
|
||||
slot: 'gen_img_3',
|
||||
bytes: new Blob(['generated-bytes'], { type: 'image/png' }),
|
||||
mimeType: 'image/png',
|
||||
writeBack: vi.fn().mockResolvedValue('written'),
|
||||
mirror: vi.fn().mockResolvedValue(undefined),
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
describe('commitToPool', () => {
|
||||
beforeEach(() => {
|
||||
mocks.poolPut.mockReset().mockResolvedValue('ast_allocated');
|
||||
});
|
||||
|
||||
it('stores, writes back and mirrors, in that order', async () => {
|
||||
const order: string[] = [];
|
||||
mocks.poolPut.mockImplementation(async () => {
|
||||
order.push('put');
|
||||
return 'ast_allocated';
|
||||
});
|
||||
const current = plan({
|
||||
meta: { durationSeconds: 2.5 },
|
||||
writeBack: vi.fn().mockImplementation(async () => {
|
||||
order.push('writeBack');
|
||||
return 'held';
|
||||
}),
|
||||
mirror: vi.fn().mockImplementation(async () => {
|
||||
order.push('mirror');
|
||||
}),
|
||||
});
|
||||
|
||||
await expect(commitToPool(current)).resolves.toEqual({
|
||||
status: 'stored',
|
||||
assetId: 'ast_allocated',
|
||||
placement: 'held',
|
||||
});
|
||||
|
||||
expect(order).toEqual(['put', 'writeBack', 'mirror']);
|
||||
const [, meta] = mocks.poolPut.mock.calls[0] as [Blob, Record<string, unknown>];
|
||||
expect(meta).toEqual({ contentType: 'image/png', durationSeconds: 2.5 });
|
||||
expect(current.writeBack).toHaveBeenCalledWith('ast_allocated');
|
||||
expect(current.mirror).toHaveBeenCalledWith('ast_allocated', 'held');
|
||||
});
|
||||
|
||||
// A refusal is an answer, not an exception: each caller has a different thing
|
||||
// to do with it, and one of them (adoption) treats it as an ordinary step.
|
||||
it('reports a full store as refused-retained rather than throwing', async () => {
|
||||
const error = quotaRefusal();
|
||||
mocks.poolPut.mockRejectedValue(error);
|
||||
const retain = vi.fn().mockResolvedValue(undefined);
|
||||
const current = plan({ retain });
|
||||
|
||||
const outcome = await commitToPool(current);
|
||||
|
||||
expect(outcome).toMatchObject({
|
||||
status: 'refused-retained',
|
||||
code: 'ASSET_QUOTA_EXCEEDED',
|
||||
error,
|
||||
refused: { slot: 'gen_img_3', mimeType: 'image/png', error },
|
||||
});
|
||||
// Awaited before the outcome is reported, so the name is a statement about
|
||||
// what is on disk rather than about what was scheduled.
|
||||
expect(retain).toHaveBeenCalledTimes(1);
|
||||
expect(current.writeBack).not.toHaveBeenCalled();
|
||||
expect(current.mirror).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
// The bytes leave with the outcome whether or not a sink was supplied: the
|
||||
// media pass carries them out to the failure record it writes around them.
|
||||
it('hands the refused bytes back even with no retain sink', async () => {
|
||||
mocks.poolPut.mockRejectedValue(quotaRefusal());
|
||||
|
||||
const outcome = await commitToPool(plan());
|
||||
|
||||
expect(outcome.status).toBe('refused-retained');
|
||||
if (outcome.status !== 'refused-retained') return;
|
||||
await expect(outcome.refused.bytes.text()).resolves.toBe('generated-bytes');
|
||||
});
|
||||
|
||||
// `refused-retained` is what a caller reads before stamping `slot` into
|
||||
// something durable. A sink that could not keep the bytes leaves nothing to
|
||||
// read back, so the outcome must not claim otherwise.
|
||||
it('demotes a refusal to failed when the bytes could not be kept', async () => {
|
||||
mocks.poolPut.mockRejectedValue(quotaRefusal());
|
||||
const retentionError = new Error('local quota exceeded');
|
||||
const current = plan({ retain: vi.fn().mockRejectedValue(retentionError) });
|
||||
|
||||
await expect(commitToPool(current)).resolves.toEqual({
|
||||
status: 'failed',
|
||||
error: retentionError,
|
||||
});
|
||||
|
||||
expect(current.writeBack).not.toHaveBeenCalled();
|
||||
expect(current.mirror).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
// By the time the mirror runs the bytes are in the pool and the document
|
||||
// names them, so a cache this browser could not write costs a re-download and
|
||||
// nothing else. Holding that here rather than at each caller is what keeps
|
||||
// the next caller from failing a commit that already happened.
|
||||
it('reports stored even when the local mirror fails', async () => {
|
||||
const current = plan({ mirror: vi.fn().mockRejectedValue(new Error('cache unavailable')) });
|
||||
|
||||
await expect(commitToPool(current)).resolves.toEqual({
|
||||
status: 'stored',
|
||||
assetId: 'ast_allocated',
|
||||
placement: 'written',
|
||||
});
|
||||
|
||||
expect(current.mirror).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
// Anything that is not the store saying "no room" says nothing about a later
|
||||
// attempt, so nothing is kept under a key a later attempt would read.
|
||||
it('keeps nothing for a failure that is not a refusal for room', async () => {
|
||||
const error = new Error('asset store unavailable');
|
||||
mocks.poolPut.mockRejectedValue(error);
|
||||
const retain = vi.fn().mockResolvedValue(undefined);
|
||||
const current = plan({ retain });
|
||||
|
||||
await expect(commitToPool(current)).resolves.toEqual({ status: 'failed', error });
|
||||
|
||||
expect(retain).not.toHaveBeenCalled();
|
||||
expect(current.writeBack).not.toHaveBeenCalled();
|
||||
expect(current.mirror).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
// A code that is not the storage contract's is not a refusal for room either,
|
||||
// however it arrived.
|
||||
it('treats an unrelated structured code as an ordinary failure', async () => {
|
||||
mocks.poolPut.mockRejectedValue(
|
||||
Object.assign(new Error('nope'), { code: 'CONTENT_SENSITIVE' }),
|
||||
);
|
||||
|
||||
await expect(commitToPool(plan())).resolves.toMatchObject({ status: 'failed' });
|
||||
});
|
||||
|
||||
// `errorCode` is the generation routes' field, on an error class a pool write
|
||||
// cannot raise. Accepting it here would widen the predicate past anything
|
||||
// `putAsset` can throw, on values from another contract.
|
||||
it('does not read the generation routes\u2019 errorCode field', async () => {
|
||||
mocks.poolPut.mockRejectedValue(
|
||||
Object.assign(new Error('nope'), { errorCode: 'ASSET_QUOTA_EXCEEDED' }),
|
||||
);
|
||||
|
||||
await expect(commitToPool(plan())).resolves.toMatchObject({ status: 'failed' });
|
||||
});
|
||||
|
||||
// The write-back's errors are the caller's: only it knows whether the
|
||||
// allocation was retained, which decides what may be reclaimed.
|
||||
it('lets a write-back failure propagate, and mirrors nothing', async () => {
|
||||
const current = plan({ writeBack: vi.fn().mockRejectedValue(new Error('document refused')) });
|
||||
|
||||
await expect(commitToPool(current)).rejects.toThrow('document refused');
|
||||
|
||||
expect(current.mirror).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
// The marker belongs to the seam and to the paths that spend provider money.
|
||||
// The primitive reads neither and writes neither: a successful write retires
|
||||
// it because `putAsset` does that, and a refusal here sets nothing.
|
||||
it('retires the course marker through the seam, and never sets one', async () => {
|
||||
setAssetStorageFullStoreForTests(memoryKv());
|
||||
try {
|
||||
await markAssetStorageFull(stageId);
|
||||
await expect(isAssetStorageFull(stageId)).resolves.toBe(true);
|
||||
|
||||
await expect(commitToPool(plan())).resolves.toMatchObject({ status: 'stored' });
|
||||
await expect(isAssetStorageFull(stageId)).resolves.toBe(false);
|
||||
|
||||
mocks.poolPut.mockRejectedValue(quotaRefusal());
|
||||
await expect(commitToPool(plan())).resolves.toMatchObject({ status: 'refused-retained' });
|
||||
await expect(isAssetStorageFull(stageId)).resolves.toBe(false);
|
||||
} finally {
|
||||
setAssetStorageFullStoreForTests(undefined);
|
||||
}
|
||||
});
|
||||
|
||||
// A caller outside a course has no course marker to retire, so the seam is
|
||||
// handed no stage rather than a made-up one.
|
||||
it('omits the stage seam when the caller has no course', async () => {
|
||||
await expect(commitToPool(plan({ stageId: undefined }))).resolves.toMatchObject({
|
||||
status: 'stored',
|
||||
});
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user