Files
WeHub Mirror 6bf8bebf51
CI / Test and Build (push) Failing after 1s
CI / Migrate Dev DB (push) Has been skipped
CI / Migrate DB (push) Has been skipped
CodeQL / Analyze actions (push) Has been cancelled
CodeQL / Analyze javascript-typescript (push) Has been cancelled
CI / Detect Version (push) Has been cancelled
CI / Detect Desktop Changes (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/cron.Dockerfile, ubuntu-latest, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build AMD64 (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/cron.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/db.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/pii.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/realtime.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-8vcpu-ubuntu-2404-arm, ./docker/app.Dockerfile, linux-arm64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Check Docs Changes (push) Has been cancelled
Publish CLI Package / publish-npm (push) Has been cancelled
Publish Python SDK / publish-pypi (push) Has been cancelled
CI / Deploy Trigger.dev (Dev) (push) Has been cancelled
Helm Chart / Lint, test, and validate chart (push) Has been cancelled
Helm Chart / Chart version bumped (push) Has been cancelled
Publish TypeScript SDK / publish-npm (push) Has been cancelled
CI / Build Dev ECR (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core) (push) Has been cancelled
CI / Promote Images (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Process Docs (push) Has been cancelled
CI / Create GitHub Release (push) Has been cancelled
CI / Check Desktop Signing Secrets (push) Has been cancelled
CI / Desktop Release (push) Has been cancelled
CI / Create Desktop Prerelease (push) Has been cancelled
CI / Desktop Prerelease Build (push) Has been cancelled
CI / Publish Desktop Prerelease (push) Has been cancelled
CI / Prune Desktop Prereleases (push) Has been cancelled
Helm Chart / Install on kind and run helm test (push) Has been cancelled
WeHub snapshot of cb28d14c6f2c081de7a0d8729a8c816c9adef67a
2026-08-10 11:17:50 +08:00

174 lines
4.2 KiB
TypeScript

import {
getLargeValueMaterializationError,
isLargeValueRef,
type LargeValueRef,
} from '@/lib/execution/payloads/large-value-ref'
const FALLBACK_TTL_MS = 15 * 60 * 1000
const MAX_IN_MEMORY_BYTES = 256 * 1024 * 1024
interface LargeValueCacheScope {
workspaceId?: string
workflowId?: string
executionId?: string
largeValueExecutionIds?: string[]
largeValueKeys?: string[]
allowLargeValueWorkflowScope?: boolean
}
const inMemoryValues = new Map<
string,
{
value: unknown
size: number
expiresAt: number
scope?: LargeValueCacheScope
recoverable: boolean
}
>()
let inMemoryBytes = 0
export function clearLargeValueCacheForTests(): void {
inMemoryValues.clear()
inMemoryBytes = 0
}
function cleanupExpiredValues(now = Date.now()): void {
for (const [id, entry] of inMemoryValues.entries()) {
if (entry.expiresAt <= now) {
inMemoryValues.delete(id)
inMemoryBytes -= entry.size
}
}
}
export function cacheLargeValue(
id: string,
value: unknown,
size: number,
scope?: LargeValueCacheScope,
options: { recoverable?: boolean } = {}
): boolean {
if (size > MAX_IN_MEMORY_BYTES) {
return false
}
cleanupExpiredValues()
const existing = inMemoryValues.get(id)
if (existing) {
inMemoryValues.delete(id)
inMemoryBytes -= existing.size
}
while (inMemoryBytes + size > MAX_IN_MEMORY_BYTES && inMemoryValues.size > 0) {
const oldestRecoverableId = Array.from(inMemoryValues.entries()).find(
([, entry]) => entry.recoverable
)?.[0]
if (!oldestRecoverableId) break
const oldest = inMemoryValues.get(oldestRecoverableId)
inMemoryValues.delete(oldestRecoverableId)
inMemoryBytes -= oldest?.size ?? 0
}
if (inMemoryBytes + size > MAX_IN_MEMORY_BYTES) {
if (existing) {
inMemoryValues.set(id, existing)
inMemoryBytes += existing.size
}
return false
}
inMemoryValues.set(id, {
value,
size,
scope,
recoverable: options.recoverable ?? false,
expiresAt: Date.now() + FALLBACK_TTL_MS,
})
inMemoryBytes += size
return true
}
function scopeMatchesRef(
ref: LargeValueRef,
cachedScope: LargeValueCacheScope | undefined,
callerScope?: LargeValueCacheScope
): boolean {
if (!cachedScope?.executionId) {
return false
}
if (ref.executionId && ref.executionId !== cachedScope.executionId) {
return false
}
if (!callerScope) {
return Boolean(ref.key) && (!ref.executionId || ref.executionId === cachedScope.executionId)
}
const allowedExecutionIds = new Set([
callerScope.executionId,
...(callerScope.largeValueExecutionIds ?? []),
])
if (ref.key && callerScope.largeValueKeys?.includes(ref.key)) {
return true
}
const workflowScopeAllowed =
callerScope.allowLargeValueWorkflowScope &&
callerScope.workspaceId === cachedScope.workspaceId &&
callerScope.workflowId === cachedScope.workflowId
return allowedExecutionIds.has(cachedScope.executionId) || Boolean(workflowScopeAllowed)
}
export function materializeLargeValueRefSync(
ref: LargeValueRef,
callerScope?: LargeValueCacheScope
): unknown {
cleanupExpiredValues()
const cached = inMemoryValues.get(ref.id)
if (!cached || !scopeMatchesRef(ref, cached.scope, callerScope)) {
return undefined
}
return cached.value
}
export function materializeLargeValueRefSyncOrThrow(
ref: LargeValueRef,
callerScope?: LargeValueCacheScope
): unknown {
const materialized = materializeLargeValueRefSync(ref, callerScope)
if (materialized === undefined) {
throw getLargeValueMaterializationError(ref)
}
return materialized
}
export function materializeLargeValueRefsSync(
value: unknown,
seen = new WeakSet<object>()
): unknown {
if (isLargeValueRef(value)) {
return materializeLargeValueRefsSync(materializeLargeValueRefSyncOrThrow(value), seen)
}
if (!value || typeof value !== 'object') {
return value
}
if (seen.has(value)) {
return value
}
seen.add(value)
if (Array.isArray(value)) {
return value.map((item) => materializeLargeValueRefsSync(item, seen))
}
return Object.fromEntries(
Object.entries(value).map(([key, entryValue]) => [
key,
materializeLargeValueRefsSync(entryValue, seen),
])
)
}