Skip to content
Open
Show file tree
Hide file tree
Changes from 8 commits
Commits
Show all changes
59 commits
Select commit Hold shift + click to select a range
dfd4100
feat: support partial (multi-request) scene uploads
LautaroPetaccio Jul 7, 2026
e8b2ea6
test: cover parallel partial-upload staging requests
LautaroPetaccio Jul 7, 2026
9057b33
fix: address partial-upload review findings
LautaroPetaccio Jul 7, 2026
4a16f90
perf: skip the permission check on resume batches of a partial upload
LautaroPetaccio Jul 8, 2026
4108a27
test: pin that losing the name permission mid-upload rejects the fina…
LautaroPetaccio Jul 8, 2026
3667eed
fix: address second-review findings on the partial-upload path
LautaroPetaccio Jul 8, 2026
6cce488
fix: guard GC pending-key projection and store all uploaded files
LautaroPetaccio Jul 9, 2026
91b2ee5
feat: order-aware pending replacement and a per-deployer staging cap
LautaroPetaccio Jul 10, 2026
597d5e3
fix: make deployment ordering and the per-deployer cap atomic
LautaroPetaccio Jul 13, 2026
4cf315e
fix: memory, GC, and cap-accounting hardening for partial uploads
LautaroPetaccio Jul 13, 2026
ab00850
feat: finalization lease so only one request finalizes a completed up…
LautaroPetaccio Jul 13, 2026
0b644ce
refactor: split partial-upload component types into per-component typ…
LautaroPetaccio Jul 13, 2026
5fcc33d
fix: address partial-deployment review findings
LautaroPetaccio Jul 14, 2026
cda077e
fix: harden partial-deployment finalize, GC, and staging hot path
LautaroPetaccio Jul 15, 2026
5a6eb94
Merge branch 'main' into feat/partial-deployments
LautaroPetaccio Jul 22, 2026
02c2f09
fix: correct persisted size and cancellation on the partial finalize …
LautaroPetaccio Jul 22, 2026
70602f7
fix: harden the partial-resume gate, GC re-check index, and duplicate…
LautaroPetaccio Jul 23, 2026
ed1bdb2
fix: harden owner reconciliation, degrade paths, and API-contract acc…
LautaroPetaccio Jul 23, 2026
eda87fc
fix: complete the whitelist casing fix and close partial-deploy edge …
LautaroPetaccio Jul 23, 2026
b84df4d
Merge branch 'main' into feat/partial-deployments
LautaroPetaccio Sep 18, 2026
b1e77de
feat: key partial uploads by entity id and let overlapping uploads co…
LautaroPetaccio Sep 23, 2026
a68f049
fix: close partial upload races and lock pool starvation
LautaroPetaccio Sep 24, 2026
436725a
fix: reject batches from another signer on a live upload
LautaroPetaccio Sep 24, 2026
bc09efa
fix: keep unadmitted batches and gc waits off upload quotas and locks
LautaroPetaccio Sep 24, 2026
d572027
fix: queue gc writers and retry a saturated lock pool
LautaroPetaccio Sep 24, 2026
d7a32b0
fix: build the gc index concurrently and bound gc writer waits
LautaroPetaccio Sep 24, 2026
818929e
fix: retry the lock pool when opening a connection times out
LautaroPetaccio Sep 24, 2026
9d0c036
fix: reuse lock connections after a failed operation
LautaroPetaccio Sep 25, 2026
70f26f8
fix: keep vanilla deploys independent of pending partial uploads
LautaroPetaccio Sep 25, 2026
e8ad942
fix: apply migrations before the server starts serving
LautaroPetaccio Sep 25, 2026
dcd5fc0
fix: tie completion receipts to the entity's latest publication
LautaroPetaccio Sep 25, 2026
d751298
fix: serialize migrations across instances with an advisory lock
LautaroPetaccio Sep 25, 2026
e99706a
feat: reject a partial query flag without the partial form field
LautaroPetaccio Sep 25, 2026
dd59b7b
fix: run migrations on the session that holds the migrations lock
LautaroPetaccio Sep 27, 2026
a6ada4d
fix: keep partial uploads within their fixed lifetime
LautaroPetaccio Sep 27, 2026
91040cd
fix: bound each source's in-flight uploads before reading the body
LautaroPetaccio Sep 28, 2026
a2d7b49
fix: start partial upload lifetimes at their first request's arrival
LautaroPetaccio Sep 28, 2026
306afac
fix: try the exclusive content lock instead of queuing on it
LautaroPetaccio Sep 28, 2026
3280dac
test: scope the deployments counter assertions to that metric
LautaroPetaccio Sep 28, 2026
1f53ca8
feat: make the trusted client ip header configurable
LautaroPetaccio Sep 29, 2026
c2f2e2e
fix: say exactly why a request timed out in 408 responses
LautaroPetaccio Sep 29, 2026
ae2f218
fix: answer full partial-upload quotas with 429 and retry-after
LautaroPetaccio Sep 29, 2026
af09cd4
feat: derive upload concurrency and file limits from the byte budget
LautaroPetaccio Sep 29, 2026
ffa10d9
fix: answer every partial batch for a published entity with 200
LautaroPetaccio Sep 29, 2026
39b38c8
feat: expire partial uploads after 1 hour and clean them every 5 minutes
LautaroPetaccio Sep 29, 2026
014177b
feat: add partial upload lifecycle, quota and cleanup metrics
LautaroPetaccio Sep 30, 2026
64c7125
fix: reject repeated form field names in multipart uploads
LautaroPetaccio Sep 30, 2026
a958049
fix: never create or charge a partial upload past its deadline
LautaroPetaccio Sep 30, 2026
e53f0bc
test: cover categories sent as repeated settings fields
LautaroPetaccio Sep 30, 2026
da53fda
fix: answer 400 for partial uploads that can never fit a quota
LautaroPetaccio Oct 4, 2026
e985749
docs: correct partial upload status codes, rollout and ordering
LautaroPetaccio Oct 4, 2026
dfcdfd7
fix: stop waiting for the content lock when a settings client disconn…
LautaroPetaccio Oct 4, 2026
daf4f8a
chore: drop the unreleased world_scenes updated_at index migration
LautaroPetaccio Oct 4, 2026
c1289f8
chore: tidy partial upload code and test responses
LautaroPetaccio Oct 4, 2026
9b93035
fix: cap post /entities form fields to what a deployment sends
LautaroPetaccio Oct 8, 2026
eb98372
feat: charge partial uploads only for the bytes they store
LautaroPetaccio Oct 8, 2026
daad8d7
fix: validate deployments before taking the content lock
LautaroPetaccio Oct 8, 2026
7131879
fix: answer multipart size and count limits with 413
LautaroPetaccio Oct 8, 2026
67e9015
feat: cap world scene size at max_scene_size
LautaroPetaccio Oct 9, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .env.default
Original file line number Diff line number Diff line change
Expand Up @@ -81,3 +81,7 @@ SHARED_SECRET_MAX_ATTEMPTS_PER_MINUTE=3
######################################
# How long to keep UNDEPLOYED scenes before permanent deletion (ms). Default: 7 days
SCENE_EVICTION_TTL_MS=604800000
# How long a partial (multi-request) deployment may stay pending before it is reclaimed (ms). Default 24h.
PENDING_DEPLOYMENT_TTL=86400000
# Max concurrent non-expired pending (partial) uploads a single deployer may have in flight. Default 10.
MAX_PENDING_DEPLOYMENTS_PER_DEPLOYER=10
15 changes: 14 additions & 1 deletion docs/database-schema.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,20 @@ The database contains the following main tables:
3. **`world_permissions`** - Stores deployment and streaming permission grants
4. **`world_permission_parcels`** - Stores parcel-level permission restrictions (normalized)
5. **`blocked`** - Stores blocked wallet addresses
6. **`migrations`** - Tracks executed database migrations (internal, managed automatically)
6. **`pending_scenes`** - Stores partial (multi-request) deployments still being uploaded (not yet live)
7. **`migrations`** - Tracks executed database migrations (internal, managed automatically)

### Table: `pending_scenes`

Staging area for partial deployments: a scene's content can be uploaded across several `POST /entities`
requests (with `partial=true`), and the entity only becomes a live `world_scenes` row once every
referenced file is present. Intentionally has **no** foreign key to `worlds` — a half-uploaded world must
not create a `worlds` row (which would leak into listings and world-validity checks) before it goes live.
The authoritative entity bytes live in content storage under the entity id; the `entity` JSONB here is a
copy used by garbage collection. Columns: `entity_id` (PK), `world_name`, `parcels` (TEXT[]), `entity`
(JSONB), `deployer`, `created_at`/`updated_at` (TIMESTAMPTZ). At most one non-expired pending scene may
exist per world + overlapping parcels; rows expire after `PENDING_DEPLOYMENT_TTL` (default 24h) and are
removed by the daily eviction job. `created_at` anchors the deployment-TTL check for staged uploads.

---

Expand Down
8 changes: 5 additions & 3 deletions src/adapters/eviction-job.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,9 @@ const DEFAULT_TTL_MS = 7 * 24 * 60 * 60 * 1000 // 7 days
const ONE_DAY_MS = 24 * 60 * 60 * 1000

export async function createEvictionJob(
components: Pick<AppComponents, 'config' | 'logs' | 'worlds'>
components: Pick<AppComponents, 'config' | 'logs' | 'worlds' | 'pendingScenesManager'>
): Promise<IJobComponent> {
const { config, logs, worlds } = components
const { config, logs, worlds, pendingScenesManager } = components
const logger = logs.getLogger('eviction-job')
const evictionTtlMs = (await config.getNumber('SCENE_EVICTION_TTL_MS')) ?? DEFAULT_TTL_MS

Expand All @@ -16,7 +16,9 @@ export async function createEvictionJob(
async () => {
logger.info('Running eviction job...')
const evicted = await worlds.evictUndeployedWorlds(evictionTtlMs)
logger.info(`Eviction completed. Deleted ${evicted} scene(s).`)
// The pending-scenes manager owns the PENDING_DEPLOYMENT_TTL, so expiry uses its configured value.
const expiredPending = await pendingScenesManager.deleteExpired()
logger.info(`Eviction completed. Deleted ${evicted} scene(s) and ${expiredPending} expired pending upload(s).`)
},
ONE_DAY_MS,
{ repeat: true, onError: (err) => logger.error(`Eviction job failed: ${err}`) }
Expand Down
148 changes: 148 additions & 0 deletions src/adapters/pending-scenes-manager.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,148 @@
import SQL from 'sql-template-strings'
import { Entity } from '@dcl/schemas'
import { InvalidRequestError } from '@dcl/http-commons'
import { AppComponents, IPendingScenesManager, PendingScene, UpsertPendingScene } from '../types'

const DEFAULT_PENDING_DEPLOYMENT_TTL_MS = 24 * 60 * 60 * 1000 // 24 hours

type PendingSceneRow = {
entity_id: string
world_name: string
parcels: string[]
entity: Entity
deployer: string
created_at: Date
updated_at: Date
}

function toPendingScene(row: PendingSceneRow): PendingScene {
return {
entityId: row.entity_id,
worldName: row.world_name,
parcels: row.parcels,
entity: row.entity,
deployer: row.deployer,
createdAt: row.created_at,
updatedAt: row.updated_at
}
}

export async function createPendingScenesManager(
components: Pick<AppComponents, 'config' | 'database' | 'logs'>
): Promise<IPendingScenesManager> {
const { config, database, logs } = components
const logger = logs.getLogger('pending-scenes-manager')
const ttlMs = (await config.getNumber('PENDING_DEPLOYMENT_TTL')) ?? DEFAULT_PENDING_DEPLOYMENT_TTL_MS

async function getByEntityId(entityId: string): Promise<PendingScene | undefined> {
const cutoff = new Date(Date.now() - ttlMs)
const result = await database.query<PendingSceneRow>(
SQL`SELECT entity_id, world_name, parcels, entity, deployer, created_at, updated_at
FROM pending_scenes
WHERE entity_id = ${entityId} AND created_at >= ${cutoff}
LIMIT 1`
)
return result.rows.length > 0 ? toPendingScene(result.rows[0]) : undefined
}

async function upsert(input: UpsertPendingScene): Promise<PendingScene> {
const worldName = input.worldName.toLowerCase()
const expiryCutoff = new Date(Date.now() - ttlMs)

return await database.withAsyncContextTransaction(async () => {
// Serialize the "replace overlapping + insert" critical section per world (also across
// processes) so two concurrent uploads for the same world+parcels can't both insert.
await database.query(SQL`SELECT pg_advisory_xact_lock(hashtextextended(${'pending_scenes:' + worldName}, 0))`)

// The single pending slot per parcel set goes to the NEWEST scene (Decentraland deployment
// ordering: greater entity.timestamp, tie broken by greater entity id). Reject rather than
// replace when a strictly-newer overlapping upload is already in flight, so a stale/older upload
// can't evict a newer competitor's staged content (and two clients can't ping-pong evicting each
// other). A resume (same entity id) is excluded and never conflicts with itself.
const newer = await database.query(SQL`
SELECT 1 FROM pending_scenes
WHERE world_name = ${worldName}
AND parcels && ${input.parcels}::text[]
AND entity_id != ${input.entityId}
AND created_at >= ${expiryCutoff}
AND ( (entity->>'timestamp')::bigint > ${input.entity.timestamp}
OR ((entity->>'timestamp')::bigint = ${input.entity.timestamp} AND entity_id > ${input.entityId}) )
LIMIT 1
`)
if (newer.rowCount > 0) {
throw new InvalidRequestError('A newer partial upload is already in progress for one or more of these parcels.')
}

// Purge expired rows and replace any non-expired pending scene of this world whose parcels
// overlap the new one (a different entity id). Having rejected the newer-conflict above, every
// remaining overlapping row is strictly older, so replacing it is the intended "newest wins".
await database.query(SQL`
DELETE FROM pending_scenes
WHERE created_at < ${expiryCutoff}
OR (world_name = ${worldName} AND parcels && ${input.parcels}::text[] AND entity_id != ${input.entityId})
`)

const result = await database.query<PendingSceneRow>(SQL`
INSERT INTO pending_scenes (entity_id, world_name, parcels, entity, deployer, created_at, updated_at)
VALUES (${input.entityId}, ${worldName}, ${input.parcels}::text[], ${input.entity}::jsonb, ${input.deployer.toLowerCase()}, now(), now())
ON CONFLICT (entity_id) DO UPDATE SET updated_at = now()
RETURNING entity_id, world_name, parcels, entity, deployer, created_at, updated_at
`)
return toPendingScene(result.rows[0])
})
}

async function deleteByEntityId(entityId: string): Promise<void> {
await database.query(SQL`DELETE FROM pending_scenes WHERE entity_id = ${entityId}`)
}

async function countActiveByDeployer(deployer: string): Promise<number> {
const cutoff = new Date(Date.now() - ttlMs)
const result = await database.query<{ count: string }>(
SQL`SELECT COUNT(*) as count FROM pending_scenes
WHERE deployer = ${deployer.toLowerCase()} AND created_at >= ${cutoff}`
)
return parseInt(result.rows[0].count, 10)
}

async function deleteExpired(): Promise<number> {
const cutoff = new Date(Date.now() - ttlMs)
const result = await database.query(SQL`DELETE FROM pending_scenes WHERE created_at < ${cutoff}`)
const removed = result.rowCount ?? 0
if (removed > 0) {
logger.info(`Removed ${removed} expired pending scene(s)`)
}
return removed
}

async function getActivePendingKeys(): Promise<Set<string>> {
const cutoff = new Date(Date.now() - ttlMs)
// Project only the content hashes out of the entity JSONB instead of shipping every pending
// scene's full manifest: GC calls this once per delete batch, so on a large sweep the payload
// size matters more than the (tiny) row count. The jsonb_typeof guard keeps a row whose `content`
// is absent, null, or a non-array from erroring `jsonb_array_elements` ('cannot extract elements
// from a scalar') — one such row would otherwise fail the whole query and wedge GC server-wide.
const result = await database.query<{ entity_id: string; hashes: string[] | null }>(
SQL`SELECT entity_id,
ARRAY(
SELECT jsonb_array_elements(entity->'content')->>'hash'
WHERE jsonb_typeof(entity->'content') = 'array'
) AS hashes
FROM pending_scenes
WHERE created_at >= ${cutoff}`
)
const keys = new Set<string>()
for (const row of result.rows) {
// The staged entity JSON, its auth-chain blob, and every content file it references are all
// referenced by the in-flight upload even though no world_scenes row exists yet.
keys.add(row.entity_id)
keys.add(`${row.entity_id}.auth`)
for (const hash of row.hashes ?? []) {
keys.add(hash)
}
}
return keys
}

return { getByEntityId, upsert, deleteByEntityId, deleteExpired, getActivePendingKeys, countActiveByDeployer }
}
20 changes: 19 additions & 1 deletion src/components.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@ import { createWorldsIndexerComponent } from './adapters/worlds-indexer'

import { createValidator } from './logic/validations'
import { createEntityDeployer } from './adapters/entity-deployer'
import { createPendingScenesManager } from './adapters/pending-scenes-manager'
import { createPartialDeploymentsComponent } from './logic/partial-deployments'
import { createMigrationExecutor } from './adapters/migration-executor'
import { createNameDenyListChecker } from './adapters/name-deny-list-checker'
import { createDatabaseComponent } from './adapters/database-component'
Expand Down Expand Up @@ -217,6 +219,20 @@ export async function initComponents(): Promise<AppComponents> {
worldsManager
})

const pendingScenesManager = await createPendingScenesManager({ config, database, logs })

const partialDeployments = await createPartialDeploymentsComponent({
config,
coordinates,
entityDeployer,
limitsManager,
logs,
pendingScenesManager,
storage,
validator,
worldsManager
})

const migrationExecutor = createMigrationExecutor({ logs, database: database, nameOwnership, storage, worldsManager })

const notificationService = await createNotificationsClientComponent({ config, fetch, logs })
Expand Down Expand Up @@ -244,7 +260,7 @@ export async function initComponents(): Promise<AppComponents> {

const worlds = createWorldsComponent({ worldsManager, snsClient })

const evictionJob = await createEvictionJob({ config, logs, worlds })
const evictionJob = await createEvictionJob({ config, logs, worlds, pendingScenesManager })

const denyList = await createDenyListComponent({ config, fetch, logs })
const bans = await createBansComponent({ config, fetch, logs })
Expand Down Expand Up @@ -302,6 +318,8 @@ export async function initComponents(): Promise<AppComponents> {
nats,
notificationService,
participantKicker,
partialDeployments,
pendingScenesManager,
peersRegistry,
permissions,
permissionsManager,
Expand Down
82 changes: 75 additions & 7 deletions src/controllers/handlers/deploy-entity-handler.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { Entity } from '@dcl/schemas'
import { IHttpServerComponent } from '@dcl/core-commons'
import { FormDataContext, readUploadedFile, toDeploymentFile } from '../../logic/multipart'
import { bufferToStream, streamToBuffer } from '@dcl/catalyst-storage'
import { FormDataContext, toDeploymentFile } from '../../logic/multipart'
import { DeploymentFile, HandlerContextWithPath } from '../../types'
import { extractAuthChain } from '../../logic/extract-auth-chain'
import { InvalidRequestError } from '@dcl/http-commons'
Expand All @@ -19,18 +20,41 @@ function parseEntityJson(raw: string) {
}

export async function deployEntity(
ctx: FormDataContext & HandlerContextWithPath<'config' | 'entityDeployer' | 'storage' | 'validator', '/entities'>
ctx: FormDataContext &
HandlerContextWithPath<
'config' | 'entityDeployer' | 'logs' | 'partialDeployments' | 'pendingScenesManager' | 'storage' | 'validator',
'/entities'
>
): Promise<IHttpServerComponent.IResponse> {
const entityId = requireString(ctx.formData.fields.entityId?.value[0])
const authChain = extractAuthChain(ctx)

const entityFile = ctx.formData.files[entityId]
// A `partial=true` field marks a staging request of a multi-request (partial) deployment: content may
// be uploaded across several requests and the world only becomes live once all of it is present.
const isPartial = ctx.formData.fields.partial?.value[0] === 'true'

// Resolve the entity file. It must be uploaded on the first request; a later partial (resume) request
// may omit it, in which case it is read back from storage where the first request stored it.
let entityFile: DeploymentFile | undefined = ctx.formData.files[entityId]
? toDeploymentFile(ctx.formData.files[entityId])
: undefined
if (!entityFile && isPartial) {
const stored = await ctx.components.storage.retrieve(entityId)
if (stored) {
const buf = await streamToBuffer(await stored.asStream())
entityFile = { size: buf.length, getStream: () => bufferToStream(buf), asBuffer: async () => buf }
}
}
if (!entityFile) {
throw new InvalidRequestError(`Entity file "${entityId}" is missing from the request.`)
throw new InvalidRequestError(
isPartial
? `The first partial request for an entity must include the entity file "${entityId}".`
: `Entity file "${entityId}" is missing from the request.`
)
}

// The entity JSON is small, so it is safe to read fully into memory.
const entityRaw = (await readUploadedFile(entityFile)).toString()
const entityRaw = (await entityFile.asBuffer()).toString()
const entityMetadataJson = parseEntityJson(entityRaw)

const entity: Entity = {
Expand All @@ -42,6 +66,36 @@ export async function deployEntity(
for (const filesKey in ctx.formData.files) {
uploadedFiles.set(filesKey, toDeploymentFile(ctx.formData.files[filesKey]))
}
// Ensure the entity file is in the map even when it came from storage (partial resume request).
Comment thread
LautaroPetaccio marked this conversation as resolved.
Outdated
if (!uploadedFiles.has(entityId)) {
uploadedFiles.set(entityId, entityFile)
}

const baseUrl = (await ctx.components.config.getString('HTTP_BASE_URL')) || `https://${ctx.url.host}`

if (isPartial) {
const result = await ctx.components.partialDeployments.stage({
baseUrl,
entity,
entityRaw,
authChain,
files: uploadedFiles
})
if (result.complete) {
return {
status: 200,
body: {
creationTimestamp: Date.now(),
...result.result
}
}
}
return { status: 202, body: { missing: result.missing ?? [] } }
}

// Vanilla (single-request) deployment — behaves exactly as before, plus TTL anchoring on and cleanup
// of any pending upload that happens to exist for this entity.
const pending = await ctx.components.pendingScenesManager.getByEntityId(entityId)

const contentHashesInStorage = await ctx.components.storage.existMultiple(
Array.from(new Set((entity.content || []).map(($) => $.hash)))
Expand All @@ -52,15 +106,15 @@ export async function deployEntity(
entity,
files: uploadedFiles,
authChain,
contentHashesInStorage
contentHashesInStorage,
pendingCreatedAt: pending?.createdAt
})

if (!validationResult.ok()) {
throw new InvalidRequestError(`Deployment failed: ${validationResult.errors.join(', ')}`)
}

// Store the entity
const baseUrl = (await ctx.components.config.getString('HTTP_BASE_URL')) || `https://${ctx.url.host}`
const message = await ctx.components.entityDeployer.deployEntity(
baseUrl,
entity,
Expand All @@ -70,6 +124,20 @@ export async function deployEntity(
authChain
)

// If this entity had a pending (partial) upload, it is now fully deployed via the vanilla path, so
// drop its staging row. Only when one exists (the common vanilla deploy has none), and best-effort:
// the deployment already committed, so a cleanup failure must not turn a successful deploy into a
// 5xx — the stale row would otherwise expire on its own via PENDING_DEPLOYMENT_TTL.
if (pending) {
try {
await ctx.components.pendingScenesManager.deleteByEntityId(entityId)
} catch (error) {
ctx.components.logs
.getLogger('deploy-entity')
.warn(`Failed to delete pending scene after a successful deploy: ${error}`, { entityId })
}
}

return {
status: 200,
body: {
Expand Down
Loading
Loading