mirror of
https://github.com/HeyPuter/puter.git
synced 2026-08-24 23:17:23 +00:00
Merge branch 'main' of https://github.com/HeyPuter/puter
This commit is contained in:
@@ -140,19 +140,52 @@ const XD_HTML = `<!DOCTYPE html>
|
||||
* alarm gate (server.ts) treats them as upstream failures and only
|
||||
* pages on the two we actually care about (rate-limit / auth).
|
||||
*
|
||||
* Status extraction covers the shapes we've seen in the wild:
|
||||
* - `.status` / `.statusCode` (OpenAI / Anthropic / Together / Google GenAI)
|
||||
* - `.response.status` (Replicate's `ApiError`)
|
||||
* - `.$metadata.httpStatusCode` (AWS SDK v3, e.g. Polly)
|
||||
* - status sniffed from the message string (last resort, for providers
|
||||
* that re-wrap their SDK error in `new Error(msg)` and lose the field)
|
||||
*
|
||||
* `HttpError`s thrown by drivers pass through untouched.
|
||||
*/
|
||||
const extractUpstreamStatus = (e: {
|
||||
status?: number;
|
||||
statusCode?: number;
|
||||
response?: { status?: number };
|
||||
$metadata?: { httpStatusCode?: number };
|
||||
message?: string;
|
||||
}): number | undefined => {
|
||||
const direct = e.status ?? e.statusCode;
|
||||
if (typeof direct === 'number') return direct;
|
||||
const fromResponse = e.response?.status;
|
||||
if (typeof fromResponse === 'number') return fromResponse;
|
||||
const fromAws = e.$metadata?.httpStatusCode;
|
||||
if (typeof fromAws === 'number') return fromAws;
|
||||
// Message sniff (e.g. "... failed with status 422 ...").
|
||||
// Only trust if it's adjacent to a status-indicating word to
|
||||
// avoid matching random 4xx/5xx-looking numbers in payloads.
|
||||
const msg = e.message;
|
||||
if (typeof msg === 'string') {
|
||||
const m = msg.match(/\bstatus(?:\s+code)?\s*[:=]?\s*(4\d\d|5\d\d)\b/i);
|
||||
if (m) return Number(m[1]);
|
||||
}
|
||||
return undefined;
|
||||
};
|
||||
|
||||
const translateProviderError = (err: unknown): unknown => {
|
||||
if (isHttpError(err)) return err;
|
||||
if (!err || typeof err !== 'object') return err;
|
||||
const e = err as {
|
||||
status?: number;
|
||||
statusCode?: number;
|
||||
response?: { status?: number };
|
||||
$metadata?: { httpStatusCode?: number };
|
||||
message?: string;
|
||||
error?: { code?: string; type?: string; message?: string };
|
||||
code?: string;
|
||||
};
|
||||
const status = e.status ?? e.statusCode;
|
||||
const status = extractUpstreamStatus(e);
|
||||
if (typeof status !== 'number') return err;
|
||||
|
||||
const msg = e.error?.message ?? e.message ?? 'Upstream provider error';
|
||||
|
||||
@@ -1649,24 +1649,6 @@ export class FSController extends PuterController {
|
||||
normalizedFileMetadata.multipartPartSize = multipartPartSize;
|
||||
}
|
||||
|
||||
const bucket = this.#firstDefined(
|
||||
metadataRecord.bucket,
|
||||
fallbackRecord.bucket,
|
||||
);
|
||||
if (typeof bucket === 'string' && bucket.length > 0) {
|
||||
normalizedFileMetadata.bucket = bucket;
|
||||
}
|
||||
|
||||
const bucketRegion = this.#firstDefined(
|
||||
metadataRecord.bucketRegion,
|
||||
metadataRecord.bucket_region,
|
||||
fallbackRecord.bucketRegion,
|
||||
fallbackRecord.bucket_region,
|
||||
);
|
||||
if (typeof bucketRegion === 'string' && bucketRegion.length > 0) {
|
||||
normalizedFileMetadata.bucketRegion = bucketRegion;
|
||||
}
|
||||
|
||||
const associatedAppId = this.#toNumber(
|
||||
this.#firstDefined(
|
||||
metadataRecord.associatedAppId,
|
||||
|
||||
@@ -516,7 +516,7 @@ describe('WebDAVController', () => {
|
||||
await dispatchMiddleware(
|
||||
makeReq({
|
||||
method: 'LOCK',
|
||||
path: '/lockable-file.txt',
|
||||
path: '/test/lockable-file.txt',
|
||||
actor: {
|
||||
user: {
|
||||
id: 1,
|
||||
@@ -539,8 +539,28 @@ describe('WebDAVController', () => {
|
||||
expect(captured.headers['lock-token']).toContain('urn:uuid:');
|
||||
});
|
||||
|
||||
it('rejects locking a path the user has no write access to (e.g. root)', async () => {
|
||||
const { res, captured } = makeRes();
|
||||
await dispatchMiddleware(
|
||||
makeReq({
|
||||
method: 'LOCK',
|
||||
path: '/',
|
||||
actor: {
|
||||
user: {
|
||||
id: 1,
|
||||
uuid: 'test-uuid',
|
||||
username: 'test',
|
||||
},
|
||||
},
|
||||
}),
|
||||
res,
|
||||
noop,
|
||||
);
|
||||
expect(captured.statusCode).toBe(403);
|
||||
});
|
||||
|
||||
it('rejects a second exclusive lock on the same path', async () => {
|
||||
const uniquePath = `/double-lock-${Date.now()}.txt`;
|
||||
const uniquePath = `/test/double-lock-${Date.now()}.txt`;
|
||||
|
||||
// First lock
|
||||
const { res: res1, captured: cap1 } = makeRes();
|
||||
@@ -582,7 +602,7 @@ describe('WebDAVController', () => {
|
||||
});
|
||||
|
||||
it('refreshes an existing lock when If header provides the token', async () => {
|
||||
const uniquePath = `/refresh-lock-${Date.now()}.txt`;
|
||||
const uniquePath = `/test/refresh-lock-${Date.now()}.txt`;
|
||||
const { res: res1, captured: cap1 } = makeRes();
|
||||
await dispatchMiddleware(
|
||||
makeReq({
|
||||
|
||||
@@ -128,7 +128,7 @@ export class WebDAVController extends PuterController {
|
||||
case 'MOVE':
|
||||
return this.#move(req, res, actor, davPath, redis, lockToken);
|
||||
case 'LOCK':
|
||||
return this.#lock(req, res, davPath, redis, lockToken);
|
||||
return this.#lock(req, res, actor, davPath, redis, lockToken);
|
||||
case 'UNLOCK':
|
||||
return this.#unlock(req, res, davPath, redis);
|
||||
default:
|
||||
@@ -646,12 +646,18 @@ export class WebDAVController extends PuterController {
|
||||
async #lock(
|
||||
req: Request,
|
||||
res: Response,
|
||||
actor: Actor,
|
||||
davPath: string,
|
||||
redis: unknown,
|
||||
headerToken: string | null,
|
||||
): Promise<void> {
|
||||
const r = redis as import('ioredis').Cluster;
|
||||
|
||||
// ACL must succeed before any lock state is touched — otherwise
|
||||
// an authenticated user could lock paths they don't own (e.g. `/`)
|
||||
// and block writes for everyone else.
|
||||
await this.#assertWrite(actor, davPath);
|
||||
|
||||
// Refresh existing lock
|
||||
if (headerToken) {
|
||||
const existing = await getLockIfValid(r, headerToken);
|
||||
|
||||
@@ -351,9 +351,14 @@ describe('TogetherImageProvider.generate output handling', () => {
|
||||
expect(result).toBe('data:image/png;base64,AAAA');
|
||||
});
|
||||
|
||||
it('wraps SDK errors with an "image generation error:" prefix', async () => {
|
||||
it('lets SDK errors bubble untouched so the driver boundary can classify them', async () => {
|
||||
const provider = makeProvider();
|
||||
const apiError = new Error('upstream blew up');
|
||||
// Together's SDK errors carry a `.status` field — re-wrapping
|
||||
// them in a plain Error stripped that out and caused the
|
||||
// catch-all `translateProviderError` to fall through to 500.
|
||||
const apiError = Object.assign(new Error('upstream blew up'), {
|
||||
status: 400,
|
||||
});
|
||||
generateMock.mockRejectedValueOnce(apiError);
|
||||
|
||||
await expect(
|
||||
@@ -363,7 +368,7 @@ describe('TogetherImageProvider.generate output handling', () => {
|
||||
prompt: 'hi',
|
||||
}),
|
||||
),
|
||||
).rejects.toThrow(/Together AI image generation error:.*upstream blew up/);
|
||||
).rejects.toMatchObject({ status: 400, message: 'upstream blew up' });
|
||||
|
||||
// Failure path must NOT meter usage.
|
||||
expect(incrementUsageSpy).not.toHaveBeenCalled();
|
||||
|
||||
@@ -172,43 +172,51 @@ export class TogetherImageProvider implements IImageProvider {
|
||||
model: selectedModel.id.replace('togetherai:', ''),
|
||||
}) as unknown as Together.Images.ImageGenerateParams;
|
||||
|
||||
try {
|
||||
const response = await this.#client.images.generate(request);
|
||||
if (!response?.data?.length) {
|
||||
throw new Error(
|
||||
'Together AI response did not include image data',
|
||||
);
|
||||
}
|
||||
|
||||
this.#meteringService.incrementUsage(
|
||||
actor,
|
||||
usageType,
|
||||
usageAmount,
|
||||
costInMicroCents,
|
||||
);
|
||||
|
||||
const first = response.data[0] as {
|
||||
url?: string;
|
||||
b64_json?: string;
|
||||
};
|
||||
const url =
|
||||
first.url ||
|
||||
(first.b64_json
|
||||
? `data:image/png;base64,${first.b64_json}`
|
||||
: undefined);
|
||||
|
||||
if (!url) {
|
||||
throw new Error(
|
||||
'Together AI response did not include an image URL',
|
||||
);
|
||||
}
|
||||
|
||||
return url;
|
||||
} catch (error) {
|
||||
throw new Error(
|
||||
`Together AI image generation error: ${(error as Error).message}`,
|
||||
// Let SDK errors bubble — together-ai SDK errors carry `.status`
|
||||
// which the driver-boundary `translateProviderError` maps to
|
||||
// `upstream_*` HttpErrors. Re-wrapping in `new Error(...)` would
|
||||
// strip the status field and cause these to surface as 500s.
|
||||
const response = await this.#client.images.generate(request);
|
||||
if (!response?.data?.length) {
|
||||
throw new HttpError(
|
||||
400,
|
||||
'Together AI response did not include image data',
|
||||
{
|
||||
legacyCode: 'upstream_bad_request',
|
||||
fields: { provider: 'together' },
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
this.#meteringService.incrementUsage(
|
||||
actor,
|
||||
usageType,
|
||||
usageAmount,
|
||||
costInMicroCents,
|
||||
);
|
||||
|
||||
const first = response.data[0] as {
|
||||
url?: string;
|
||||
b64_json?: string;
|
||||
};
|
||||
const url =
|
||||
first.url ||
|
||||
(first.b64_json
|
||||
? `data:image/png;base64,${first.b64_json}`
|
||||
: undefined);
|
||||
|
||||
if (!url) {
|
||||
throw new HttpError(
|
||||
400,
|
||||
'Together AI response did not include an image URL',
|
||||
{
|
||||
legacyCode: 'upstream_bad_request',
|
||||
fields: { provider: 'together' },
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
return url;
|
||||
}
|
||||
|
||||
#getModel(model?: string) {
|
||||
|
||||
@@ -317,7 +317,7 @@ describe('VoiceChangerDriver.convert success path', () => {
|
||||
// ── Error mapping ───────────────────────────────────────────────────
|
||||
|
||||
describe('VoiceChangerDriver.convert error mapping', () => {
|
||||
it('rethrows the upstream status when ElevenLabs returns an error body', async () => {
|
||||
it('maps upstream 4xx to HttpError upstream_bad_request (keeps the upstream status)', async () => {
|
||||
const { actor } = await makeUser();
|
||||
fetchSpy.mockResolvedValueOnce(
|
||||
new Response(
|
||||
@@ -336,9 +336,32 @@ describe('VoiceChangerDriver.convert error mapping', () => {
|
||||
voice_id: 'missing-voice',
|
||||
}),
|
||||
),
|
||||
).rejects.toMatchObject({ statusCode: 404 });
|
||||
).rejects.toMatchObject({
|
||||
statusCode: 404,
|
||||
legacyCode: 'upstream_bad_request',
|
||||
});
|
||||
|
||||
// No metering should be recorded on a failed call.
|
||||
expect(incrementUsageSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('maps upstream 5xx to HttpError 400 upstream_provider_unavailable (skips alert)', async () => {
|
||||
const { actor } = await makeUser();
|
||||
fetchSpy.mockResolvedValueOnce(
|
||||
new Response('boom', { status: 503 }),
|
||||
);
|
||||
|
||||
await expect(
|
||||
withActor(actor, () =>
|
||||
driver.convert({
|
||||
audio: dataUrl(Buffer.from('x'), 'audio/mpeg'),
|
||||
voice_id: 'any-voice',
|
||||
}),
|
||||
),
|
||||
).rejects.toMatchObject({
|
||||
statusCode: 400,
|
||||
legacyCode: 'upstream_provider_unavailable',
|
||||
});
|
||||
expect(incrementUsageSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -229,8 +229,32 @@ export class VoiceChangerDriver extends PuterDriver {
|
||||
detail && typeof detail === 'object' && 'detail' in detail
|
||||
? String((detail as { detail: unknown }).detail)
|
||||
: `ElevenLabs returned ${response.status}`;
|
||||
throw new HttpError(response.status, message, {
|
||||
legacyCode: 'internal_error',
|
||||
// Tag upstream status as `upstream_*` so the alarm gate
|
||||
// skips paging on ElevenLabs 5xx outages (we expose them
|
||||
// as 400 like the TTS provider does — user can't act on
|
||||
// them, but it's not our bug either).
|
||||
const legacyCode =
|
||||
response.status >= 500
|
||||
? 'upstream_provider_unavailable'
|
||||
: response.status === 401 || response.status === 403
|
||||
? 'upstream_auth_failed'
|
||||
: response.status === 429
|
||||
? 'upstream_rate_limited'
|
||||
: 'upstream_bad_request';
|
||||
const exposedStatus =
|
||||
legacyCode === 'upstream_rate_limited'
|
||||
? 429
|
||||
: legacyCode === 'upstream_auth_failed'
|
||||
? 500
|
||||
: legacyCode === 'upstream_provider_unavailable'
|
||||
? 400
|
||||
: response.status;
|
||||
throw new HttpError(exposedStatus, message, {
|
||||
legacyCode,
|
||||
fields: {
|
||||
provider: 'elevenlabs',
|
||||
upstreamStatus: response.status,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -215,32 +215,45 @@ export class XAISpeechToTextDriver extends PuterDriver {
|
||||
formData.append('file', blob, filename);
|
||||
}
|
||||
|
||||
let response: Response;
|
||||
try {
|
||||
response = await fetch(`${API_BASE}/stt`, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
Authorization: `Bearer ${this.#apiKey}`,
|
||||
},
|
||||
body: formData,
|
||||
});
|
||||
} catch (e: unknown) {
|
||||
const msg = (e as Error).message ?? String(e);
|
||||
console.error('[XAISpeechToTextDriver] API error:', msg);
|
||||
throw new HttpError(502, `xAI STT API error: ${msg}`, {
|
||||
legacyCode: 'internal_error',
|
||||
});
|
||||
}
|
||||
const response = await fetch(`${API_BASE}/stt`, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
Authorization: `Bearer ${this.#apiKey}`,
|
||||
},
|
||||
body: formData,
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errText = await response.text().catch(() => '');
|
||||
console.error(
|
||||
`[XAISpeechToTextDriver] API returned ${response.status}: ${errText}`,
|
||||
);
|
||||
// Mirrors ElevenLabs / XAITTS — map upstream status to an
|
||||
// `upstream_*` HttpError so the alarm gate skips it.
|
||||
const legacyCode =
|
||||
response.status >= 500
|
||||
? 'upstream_provider_unavailable'
|
||||
: response.status === 401 || response.status === 403
|
||||
? 'upstream_auth_failed'
|
||||
: response.status === 429
|
||||
? 'upstream_rate_limited'
|
||||
: 'upstream_bad_request';
|
||||
const exposedStatus =
|
||||
legacyCode === 'upstream_rate_limited'
|
||||
? 429
|
||||
: legacyCode === 'upstream_auth_failed'
|
||||
? 500
|
||||
: 400;
|
||||
throw new HttpError(
|
||||
502,
|
||||
`xAI STT API error (${response.status}): ${errText}`,
|
||||
{ legacyCode: 'internal_error' },
|
||||
exposedStatus,
|
||||
errText || `xAI STT request failed (status ${response.status})`,
|
||||
{
|
||||
legacyCode,
|
||||
fields: {
|
||||
provider: 'xai',
|
||||
upstreamStatus: response.status,
|
||||
},
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -476,17 +476,21 @@ describe('GeminiTTSProvider.synthesize error paths', () => {
|
||||
}
|
||||
});
|
||||
|
||||
it('wraps SDK errors as HttpError 502', async () => {
|
||||
it('lets SDK errors bubble untouched so the driver boundary can classify them', async () => {
|
||||
const provider = makeProvider();
|
||||
generateContentMock.mockRejectedValueOnce(new Error('upstream blew up'));
|
||||
// Google GenAI `ApiError`s carry `.status` — surfacing them
|
||||
// through the driver-boundary translator yields a proper
|
||||
// `upstream_*` HttpError instead of a 500/502 page.
|
||||
const apiError = Object.assign(new Error('bad voice'), { status: 400 });
|
||||
generateContentMock.mockRejectedValueOnce(apiError);
|
||||
|
||||
await expect(
|
||||
withTestActor(() => provider.synthesize({ text: 'hi' })),
|
||||
).rejects.toMatchObject({ statusCode: 502 });
|
||||
).rejects.toMatchObject({ status: 400, message: 'bad voice' });
|
||||
expect(batchIncrementUsagesSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('throws 502 when Gemini response has no inline audio data', async () => {
|
||||
it('throws 400 upstream_bad_request when Gemini response has no inline audio data', async () => {
|
||||
const provider = makeProvider();
|
||||
generateContentMock.mockResolvedValueOnce({
|
||||
candidates: [{ content: { parts: [{ text: 'no audio here' }] } }],
|
||||
@@ -495,7 +499,10 @@ describe('GeminiTTSProvider.synthesize error paths', () => {
|
||||
|
||||
await expect(
|
||||
withTestActor(() => provider.synthesize({ text: 'hi' })),
|
||||
).rejects.toMatchObject({ statusCode: 502 });
|
||||
).rejects.toMatchObject({
|
||||
statusCode: 400,
|
||||
legacyCode: 'upstream_bad_request',
|
||||
});
|
||||
expect(batchIncrementUsagesSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -229,35 +229,29 @@ export class GeminiTTSProvider extends TTSProvider {
|
||||
? `${instructions}\n\nSay the following text aloud:\n${text}`
|
||||
: `Say the following text aloud:\n${text}`;
|
||||
|
||||
// Let Google GenAI `ApiError`s bubble — they carry `.status` and
|
||||
// are mapped to `upstream_*` HttpErrors by the driver-boundary
|
||||
// translator. Catching here and wrapping as 502 hid the upstream
|
||||
// status and caused 4xx validation errors to page.
|
||||
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
||||
let response: any;
|
||||
try {
|
||||
response = await this.#client.models.generateContent({
|
||||
model,
|
||||
contents: [{ parts: [{ text: inputText }] }],
|
||||
config: {
|
||||
responseModalities: ['AUDIO'],
|
||||
speechConfig: {
|
||||
voiceConfig: {
|
||||
prebuiltVoiceConfig: { voiceName: voice },
|
||||
},
|
||||
const response: any = await this.#client.models.generateContent({
|
||||
model,
|
||||
contents: [{ parts: [{ text: inputText }] }],
|
||||
config: {
|
||||
responseModalities: ['AUDIO'],
|
||||
speechConfig: {
|
||||
voiceConfig: {
|
||||
prebuiltVoiceConfig: { voiceName: voice },
|
||||
},
|
||||
},
|
||||
});
|
||||
} catch (e: unknown) {
|
||||
const msg = (e as Error).message ?? String(e);
|
||||
console.error('[GeminiTTSProvider] API error:', msg);
|
||||
throw new HttpError(502, `Gemini TTS API error: ${msg}`, {
|
||||
legacyCode: 'internal_error',
|
||||
fields: { provider: 'gemini' },
|
||||
});
|
||||
}
|
||||
},
|
||||
});
|
||||
|
||||
// Extract audio data from response
|
||||
const part = response?.candidates?.[0]?.content?.parts?.[0];
|
||||
if (!part?.inlineData?.data) {
|
||||
throw new HttpError(502, 'Gemini TTS did not return audio data', {
|
||||
legacyCode: 'internal_error',
|
||||
throw new HttpError(400, 'Gemini TTS did not return audio data', {
|
||||
legacyCode: 'upstream_bad_request',
|
||||
fields: { provider: 'gemini' },
|
||||
});
|
||||
}
|
||||
|
||||
@@ -327,7 +327,7 @@ describe('XAITTSProvider.synthesize metering', () => {
|
||||
// ── Error paths ─────────────────────────────────────────────────────
|
||||
|
||||
describe('XAITTSProvider.synthesize error paths', () => {
|
||||
it('wraps non-OK upstream responses as HttpError 502', async () => {
|
||||
it('maps upstream 4xx to HttpError 400 upstream_bad_request', async () => {
|
||||
const provider = makeProvider();
|
||||
fetchSpy.mockResolvedValueOnce(
|
||||
new Response('bad request', { status: 400 }),
|
||||
@@ -335,17 +335,55 @@ describe('XAITTSProvider.synthesize error paths', () => {
|
||||
|
||||
await expect(
|
||||
withTestActor(() => provider.synthesize({ text: 'hi' })),
|
||||
).rejects.toMatchObject({ statusCode: 502 });
|
||||
).rejects.toMatchObject({
|
||||
statusCode: 400,
|
||||
legacyCode: 'upstream_bad_request',
|
||||
});
|
||||
expect(incrementUsageSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('wraps fetch network errors as HttpError 502 (not a thrown ENOTFOUND)', async () => {
|
||||
it('maps upstream 5xx to HttpError 400 upstream_provider_unavailable', async () => {
|
||||
const provider = makeProvider();
|
||||
fetchSpy.mockResolvedValueOnce(
|
||||
new Response('oops', { status: 503 }),
|
||||
);
|
||||
|
||||
await expect(
|
||||
withTestActor(() => provider.synthesize({ text: 'hi' })),
|
||||
).rejects.toMatchObject({
|
||||
statusCode: 400,
|
||||
legacyCode: 'upstream_provider_unavailable',
|
||||
});
|
||||
expect(incrementUsageSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('maps upstream 429 to HttpError 429 upstream_rate_limited', async () => {
|
||||
const provider = makeProvider();
|
||||
fetchSpy.mockResolvedValueOnce(
|
||||
new Response('slow down', { status: 429 }),
|
||||
);
|
||||
|
||||
await expect(
|
||||
withTestActor(() => provider.synthesize({ text: 'hi' })),
|
||||
).rejects.toMatchObject({
|
||||
statusCode: 429,
|
||||
legacyCode: 'upstream_rate_limited',
|
||||
});
|
||||
expect(incrementUsageSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('lets fetch network errors bubble so the driver boundary can decide', async () => {
|
||||
const provider = makeProvider();
|
||||
// Raw network errors don't carry an upstream status — we
|
||||
// intentionally don't wrap them here. The driver-boundary
|
||||
// catch-all will surface them as a generic 500 (which we
|
||||
// *do* want to alert on, since "we can't even reach the
|
||||
// provider" usually means something on our side is wrong).
|
||||
fetchSpy.mockRejectedValueOnce(new Error('connection reset'));
|
||||
|
||||
await expect(
|
||||
withTestActor(() => provider.synthesize({ text: 'hi' })),
|
||||
).rejects.toMatchObject({ statusCode: 502 });
|
||||
).rejects.toThrow('connection reset');
|
||||
expect(incrementUsageSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -154,34 +154,48 @@ export class XAITTSProvider extends TTSProvider {
|
||||
body.output_format = { codec };
|
||||
}
|
||||
|
||||
let response: Response;
|
||||
try {
|
||||
response = await fetch(`${API_BASE}/tts`, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
Authorization: `Bearer ${this.#apiKey}`,
|
||||
'Content-Type': 'application/json',
|
||||
},
|
||||
body: JSON.stringify(body),
|
||||
});
|
||||
} catch (e: unknown) {
|
||||
const msg = (e as Error).message ?? String(e);
|
||||
console.error('[XAITTSProvider] API error:', msg);
|
||||
throw new HttpError(502, `xAI TTS API error: ${msg}`, {
|
||||
legacyCode: 'internal_error',
|
||||
fields: { provider: 'xai' },
|
||||
});
|
||||
}
|
||||
const response = await fetch(`${API_BASE}/tts`, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
Authorization: `Bearer ${this.#apiKey}`,
|
||||
'Content-Type': 'application/json',
|
||||
},
|
||||
body: JSON.stringify(body),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errText = await response.text().catch(() => '');
|
||||
console.error(
|
||||
`[XAITTSProvider] API returned ${response.status}: ${errText}`,
|
||||
);
|
||||
// Map upstream status to an `upstream_*` HttpError so the
|
||||
// alarm gate skips it. Mirrors ElevenLabs' translator —
|
||||
// 4xx and 5xx both surface as 400 to the client (with the
|
||||
// appropriate legacyCode), 429 stays 429, auth stays 500.
|
||||
const legacyCode =
|
||||
response.status >= 500
|
||||
? 'upstream_provider_unavailable'
|
||||
: response.status === 401 || response.status === 403
|
||||
? 'upstream_auth_failed'
|
||||
: response.status === 429
|
||||
? 'upstream_rate_limited'
|
||||
: 'upstream_bad_request';
|
||||
const exposedStatus =
|
||||
legacyCode === 'upstream_rate_limited'
|
||||
? 429
|
||||
: legacyCode === 'upstream_auth_failed'
|
||||
? 500
|
||||
: 400;
|
||||
throw new HttpError(
|
||||
502,
|
||||
`xAI TTS API error (${response.status}): ${errText}`,
|
||||
{ legacyCode: 'internal_error', fields: { provider: 'xai' } },
|
||||
exposedStatus,
|
||||
errText || `xAI TTS request failed (status ${response.status})`,
|
||||
{
|
||||
legacyCode,
|
||||
fields: {
|
||||
provider: 'xai',
|
||||
upstreamStatus: response.status,
|
||||
},
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -362,7 +362,7 @@ describe('OpenAIVideoProvider.generate polling', () => {
|
||||
}
|
||||
});
|
||||
|
||||
it('throws when the polled job ends in failed state', async () => {
|
||||
it('surfaces failed jobs as HttpError 400 upstream_failed (not a 500 page)', async () => {
|
||||
const provider = makeProvider();
|
||||
videosCreateMock.mockResolvedValueOnce({
|
||||
id: 'job-fail',
|
||||
@@ -376,7 +376,11 @@ describe('OpenAIVideoProvider.generate polling', () => {
|
||||
withTestActor(() =>
|
||||
provider.generate({ prompt: 'hi', model: 'sora-2' }),
|
||||
),
|
||||
).rejects.toThrow('content policy violation');
|
||||
).rejects.toMatchObject({
|
||||
statusCode: 400,
|
||||
legacyCode: 'upstream_failed',
|
||||
message: 'content policy violation',
|
||||
});
|
||||
expect(videosDownloadContentMock).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -134,7 +134,14 @@ export class OpenAIVideoProvider extends VideoProvider {
|
||||
if (finalJob.status === 'failed') {
|
||||
const errorMessage =
|
||||
finalJob.error?.message ?? 'Video generation failed';
|
||||
throw new Error(errorMessage);
|
||||
// Same reasoning as TogetherVideoProvider — Sora's `failed`
|
||||
// status covers both user input issues (content policy) and
|
||||
// their own outages; expose as `upstream_failed` 400 so the
|
||||
// alarm gate skips it instead of paging on 500.
|
||||
throw new HttpError(400, errorMessage, {
|
||||
legacyCode: 'upstream_failed',
|
||||
fields: { provider: 'openai' },
|
||||
});
|
||||
}
|
||||
|
||||
const finalResolution =
|
||||
|
||||
@@ -409,7 +409,7 @@ describe('TogetherVideoProvider.generate polling', () => {
|
||||
}
|
||||
});
|
||||
|
||||
it('throws HttpError 500 on failed jobs with the surfaced error message', async () => {
|
||||
it('surfaces failed jobs as HttpError 400 upstream_failed (not 500) so they do not page', async () => {
|
||||
const provider = makeProvider();
|
||||
videosCreateMock.mockResolvedValueOnce({ id: 'job-3' });
|
||||
videosRetrieveMock.mockResolvedValueOnce({
|
||||
@@ -421,7 +421,8 @@ describe('TogetherVideoProvider.generate polling', () => {
|
||||
await expect(
|
||||
withTestActor(() => provider.generate({ prompt: 'hi' })),
|
||||
).rejects.toMatchObject({
|
||||
statusCode: 500,
|
||||
statusCode: 400,
|
||||
legacyCode: 'upstream_failed',
|
||||
message: 'content policy violation',
|
||||
});
|
||||
});
|
||||
|
||||
@@ -195,8 +195,14 @@ export class TogetherVideoProvider extends VideoProvider {
|
||||
finalJob?.info?.errors?.message ??
|
||||
finalJob?.info?.errors ??
|
||||
'Video generation failed';
|
||||
throw new HttpError(500, errorMessage, {
|
||||
legacyCode: 'internal_error',
|
||||
// Together returns `failed` for both user-input issues
|
||||
// (content policy / unsupported params) and their own
|
||||
// outages — we can't reliably tell from the payload, so
|
||||
// expose as 4xx with `upstream_failed`. The alarm gate
|
||||
// skips `upstream_*` legacy codes so this no longer pages.
|
||||
throw new HttpError(400, errorMessage, {
|
||||
legacyCode: 'upstream_failed',
|
||||
fields: { provider: 'together' },
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -280,9 +280,8 @@ export class FSService extends PuterService {
|
||||
return normalizedPath;
|
||||
}
|
||||
|
||||
#resolveBucket(metadata: FSEntryWriteInput): string {
|
||||
const bucket =
|
||||
metadata.bucket ?? this.config.s3_bucket ?? 'puter-local';
|
||||
#resolveBucket(): string {
|
||||
const bucket = this.config.s3_bucket ?? 'puter-local';
|
||||
if (typeof bucket !== 'string' || bucket.length === 0) {
|
||||
throw new HttpError(500, 'Missing S3 bucket configuration', {
|
||||
legacyCode: 'internal_error',
|
||||
@@ -291,12 +290,9 @@ export class FSService extends PuterService {
|
||||
return bucket;
|
||||
}
|
||||
|
||||
#resolveBucketRegion(metadata: FSEntryWriteInput): string {
|
||||
#resolveBucketRegion(): string {
|
||||
const bucketRegion =
|
||||
metadata.bucketRegion ??
|
||||
this.config.s3_region ??
|
||||
this.config.region ??
|
||||
'us-west-2';
|
||||
this.config.s3_region ?? this.config.region ?? 'us-west-2';
|
||||
|
||||
if (typeof bucketRegion !== 'string' || bucketRegion.length === 0) {
|
||||
throw new HttpError(500, 'Missing S3 region configuration', {
|
||||
@@ -345,8 +341,8 @@ export class FSService extends PuterService {
|
||||
immutable: Boolean(metadata.immutable),
|
||||
isPublic: metadata.isPublic,
|
||||
multipartPartSize: metadata.multipartPartSize,
|
||||
bucket: this.#resolveBucket(metadata),
|
||||
bucketRegion: this.#resolveBucketRegion(metadata),
|
||||
bucket: this.#resolveBucket(),
|
||||
bucketRegion: this.#resolveBucketRegion(),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -650,8 +646,14 @@ export class FSService extends PuterService {
|
||||
size: session.size,
|
||||
contentType: session.contentType,
|
||||
checksumSha256: session.checksumSha256 ?? undefined,
|
||||
bucket: session.bucket ?? undefined,
|
||||
bucketRegion: session.bucketRegion ?? undefined,
|
||||
bucket:
|
||||
session.bucket ??
|
||||
parsedMetadata.bucket ??
|
||||
this.#resolveBucket(),
|
||||
bucketRegion:
|
||||
session.bucketRegion ??
|
||||
parsedMetadata.bucketRegion ??
|
||||
this.#resolveBucketRegion(),
|
||||
overwrite: Boolean(session.overwriteTargetUid),
|
||||
};
|
||||
}
|
||||
@@ -2075,14 +2077,14 @@ export class FSService extends PuterService {
|
||||
bucket:
|
||||
session.bucket ??
|
||||
createInput.bucket ??
|
||||
this.#resolveBucket(createInput),
|
||||
this.#resolveBucket(),
|
||||
objectKey: session.objectKey,
|
||||
multipartUploadId: session.multipartUploadId,
|
||||
parts: completeParts,
|
||||
},
|
||||
session.bucketRegion ??
|
||||
createInput.bucketRegion ??
|
||||
this.#resolveBucketRegion(createInput),
|
||||
this.#resolveBucketRegion(),
|
||||
);
|
||||
}
|
||||
|
||||
@@ -2219,14 +2221,14 @@ export class FSService extends PuterService {
|
||||
bucket:
|
||||
item.session.bucket ??
|
||||
item.finalData.bucket ??
|
||||
this.#resolveBucket(item.finalData),
|
||||
this.#resolveBucket(),
|
||||
objectKey: item.session.objectKey,
|
||||
multipartUploadId: item.session.multipartUploadId,
|
||||
parts: completeParts,
|
||||
},
|
||||
item.session.bucketRegion ??
|
||||
item.finalData.bucketRegion ??
|
||||
this.#resolveBucketRegion(item.finalData),
|
||||
this.#resolveBucketRegion(),
|
||||
);
|
||||
}),
|
||||
);
|
||||
|
||||
@@ -73,13 +73,13 @@ export interface FSEntryWriteInput {
|
||||
immutable?: boolean;
|
||||
isPublic?: boolean | null;
|
||||
multipartPartSize?: number;
|
||||
bucket?: string;
|
||||
bucketRegion?: string;
|
||||
}
|
||||
|
||||
export interface FSEntryCreateInput extends FSEntryWriteInput {
|
||||
userId: number;
|
||||
uuid: string;
|
||||
bucket: string;
|
||||
bucketRegion: string;
|
||||
}
|
||||
|
||||
export interface PendingUploadSession {
|
||||
|
||||
@@ -748,8 +748,12 @@ export class FSEntryStore extends PuterStore {
|
||||
for (const requiredPath of normalizedRequiredPaths) {
|
||||
const entry = allEntries.get(requiredPath);
|
||||
if (!entry) {
|
||||
throw new Error(
|
||||
`Failed to resolve directory path: ${requiredPath}`,
|
||||
// INSERT IGNORE silently dropped the row (see sibling
|
||||
// path in #ensureDirectoryPath). Surface as 404.
|
||||
throw new HttpError(
|
||||
404,
|
||||
`Parent path does not exist: ${requiredPath}`,
|
||||
{ legacyCode: 'not_found' },
|
||||
);
|
||||
}
|
||||
if (!entry.isDir) {
|
||||
@@ -862,8 +866,15 @@ export class FSEntryStore extends PuterStore {
|
||||
}`,
|
||||
);
|
||||
}
|
||||
throw new Error(
|
||||
`Failed to resolve directory path: ${normalizedPath}`,
|
||||
// INSERT IGNORE silently dropped the row (typical cause:
|
||||
// FK violation on parent_id/user_id — e.g. the user namespace
|
||||
// doesn't exist). The path is unreachable from the client's
|
||||
// perspective; surface as 404 rather than letting a generic
|
||||
// Error become a 500 + page on-call.
|
||||
throw new HttpError(
|
||||
404,
|
||||
`Parent path does not exist: ${normalizedPath}`,
|
||||
{ legacyCode: 'not_found' },
|
||||
);
|
||||
}
|
||||
if (!resolvedEntry.isDir) {
|
||||
|
||||
@@ -60,7 +60,7 @@ const SYSTEM_NAMESPACE = `v1:${SYSTEM_ACTOR_UUID}:${GLOBAL_APP_KEY}`;
|
||||
const MAX_KEY_BYTES = 1024;
|
||||
const MAX_VALUE_BYTES = 399 * 1024;
|
||||
const BATCH_GET_CHUNK = 100;
|
||||
const PATH_CLEANER_REGEX = /[:\-+/*]/g;
|
||||
const PATH_CLEANER_REGEX = /[^A-Za-z0-9_]/g;
|
||||
|
||||
const emptyUsage = (): KVUsage => ({ read: 0, write: 0 });
|
||||
|
||||
@@ -182,7 +182,7 @@ const objectsEqual = (left: unknown, right: unknown): boolean => {
|
||||
};
|
||||
|
||||
const cleanAttrName = (chunk: string): string =>
|
||||
`#${chunk}`.replaceAll(PATH_CLEANER_REGEX, '');
|
||||
`#${chunk.replaceAll(PATH_CLEANER_REGEX, '')}`;
|
||||
|
||||
// ── SystemKVStore ────────────────────────────────────────────────────
|
||||
|
||||
|
||||
Reference in New Issue
Block a user