diff --git a/packages/objectql/src/index.ts b/packages/objectql/src/index.ts index 7524e2bf8d..eb730d2b30 100644 --- a/packages/objectql/src/index.ts +++ b/packages/objectql/src/index.ts @@ -80,3 +80,7 @@ export type { IntrospectedTable, IntrospectedSchema, } from './util.js'; + +// Seed loader — materializes `seed` metadata into rows (used by publishMetaItem +// and the runtime dispatcher/app plugins). +export { SeedLoaderService } from './seed-loader.js'; diff --git a/packages/objectql/src/protocol-publish-package-drafts.test.ts b/packages/objectql/src/protocol-publish-package-drafts.test.ts index 0089b74903..0c109befa9 100644 --- a/packages/objectql/src/protocol-publish-package-drafts.test.ts +++ b/packages/objectql/src/protocol-publish-package-drafts.test.ts @@ -64,4 +64,141 @@ describe('protocol.publishPackageDrafts (ADR-0033)', () => { expect(publishMetaItem).not.toHaveBeenCalled(); expect(res).toMatchObject({ success: false, publishedCount: 0, failedCount: 0 }); }); + + it('publishes seeds LAST and batch-applies their rows in ONE pass (seedApplied)', async () => { + // listDrafts order puts the seed FIRST — the partition must still publish + // the object before it (its table must exist before rows land). + const drafts = [ + { type: 'seed', name: 'project_sample' }, + { type: 'object', name: 'project' }, + { type: 'seed', name: 'task_sample' }, + ]; + const protocol = new ObjectStackProtocolImplementation({} as never); + (protocol as any).ensureOverlayIndex = async () => {}; + const seedBodyByName: Record = { + project_sample: { object: 'project', records: [{ name: 'Apollo' }] }, + task_sample: { object: 'task', records: [{ name: 'Design' }] }, + }; + (protocol as any).getOverlayRepo = () => ({ + listDrafts: async () => drafts, + get: async (ref: any, opts: any) => + opts?.state === 'draft' && seedBodyByName[ref.name] + ? { body: seedBodyByName[ref.name], hash: 'h' } + : null, + }); + const publishMetaItem = vi.spyOn(protocol, 'publishMetaItem' as never); + publishMetaItem.mockResolvedValue({ success: true, version: 'h', seq: 1 } as never); + const applySeedBodies = vi + .spyOn(protocol as any, 'applySeedBodies') + .mockResolvedValue({ success: true, inserted: 2, updated: 0 }); + + const res = await protocol.publishPackageDrafts({ packageId: 'app.pm' }); + + // Object published BEFORE the seeds, and every publish suppressed per-item apply. + expect((publishMetaItem.mock.calls[0][0] as any)).toMatchObject({ type: 'object', name: 'project' }); + for (const call of publishMetaItem.mock.calls) { + expect((call[0] as any)._skipSeedApply).toBe(true); + } + // ONE batch apply with BOTH seed bodies (cross-seed refs need a single pass). + expect(applySeedBodies).toHaveBeenCalledTimes(1); + expect(applySeedBodies.mock.calls[0][0]).toEqual([ + seedBodyByName.project_sample, + seedBodyByName.task_sample, + ]); + expect(res.seedApplied).toEqual({ success: true, inserted: 2, updated: 0 }); + }); + + it('omits seedApplied when the package has no seed drafts', async () => { + const { protocol, publishMetaItem } = makeProtocol([{ type: 'object', name: 'course' }]); + publishMetaItem.mockResolvedValue({ success: true, version: 'h', seq: 1 } as never); + const res = await protocol.publishPackageDrafts({ packageId: 'app.edu' }); + expect(res.seedApplied).toBeUndefined(); + }); +}); + +/** + * Publishing a single `seed` draft (the per-ref path: POST /meta/seed/:name/publish, + * used by the home banner) must materialize its rows too — not only the package + * route. The publish itself NEVER fails on a seed problem; it reports under + * `seedApplied`. + */ +describe('protocol.publishMetaItem — seed self-apply', () => { + function makePublishable(body: unknown) { + const protocol = new ObjectStackProtocolImplementation({} as never); + (protocol as any).ensureOverlayIndex = async () => {}; + (protocol as any).assertLockAllowsWrite = async () => null; + (protocol as any).isArtifactBacked = () => false; + (protocol as any).applyObjectRegistryMutation = () => {}; + (protocol as any).ensureObjectStorage = async () => {}; + (protocol as any).getOverlayRepo = () => ({ + promoteDraft: async () => ({ version: 'sha256:x', seq: 7, item: { body } }), + }); + const applySeedBodies = vi + .spyOn(protocol as any, 'applySeedBodies') + .mockResolvedValue({ success: true, inserted: 3, updated: 0 }); + return { protocol, applySeedBodies }; + } + + it('applies the seed body on publish and reports seedApplied', async () => { + const body = { object: 'project', records: [{ name: 'Apollo' }] }; + const { protocol, applySeedBodies } = makePublishable(body); + const res = await protocol.publishMetaItem({ type: 'seed', name: 'project_sample' }); + expect(applySeedBodies).toHaveBeenCalledWith([body], null); + expect(res.seedApplied).toEqual({ success: true, inserted: 3, updated: 0 }); + expect(res.success).toBe(true); + }); + + it('suppresses the self-apply when _skipSeedApply is set (package batch path)', async () => { + const { protocol, applySeedBodies } = makePublishable({ object: 'p', records: [] }); + const res = await protocol.publishMetaItem({ type: 'seed', name: 'p_sample', _skipSeedApply: true }); + expect(applySeedBodies).not.toHaveBeenCalled(); + expect(res.seedApplied).toBeUndefined(); + }); + + it('does not touch the loader for non-seed publishes', async () => { + const { protocol, applySeedBodies } = makePublishable({ name: 'overview' }); + const res = await protocol.publishMetaItem({ type: 'dashboard', name: 'overview' }); + expect(applySeedBodies).not.toHaveBeenCalled(); + expect(res.seedApplied).toBeUndefined(); + }); +}); + +/** + * applySeedBodies wires the real SeedLoaderService: externalId('name')-keyed + * upsert against the engine, object metadata read through the protocol's own + * getMetaItem. A smoke test with a fake engine proves rows actually land and + * the result mapping is faithful. + */ +describe('protocol.applySeedBodies — real loader smoke test', () => { + it('inserts seed records via the engine and reports counts', async () => { + const protocol = new ObjectStackProtocolImplementation({} as never); + const inserted: Array<{ object: string; record: any }> = []; + (protocol as any).engine = { + find: async () => [], + insert: async (object: string, record: any) => { + inserted.push({ object, record }); + return { id: `${object}_${inserted.length}` }; + }, + update: async () => ({}), + }; + (protocol as any).getMetaItem = async ({ name }: any) => ({ + item: { name, fields: { name: { type: 'text' } } }, + }); + + const res = await (protocol as any).applySeedBodies( + [{ object: 'project', records: [{ name: 'Apollo' }, { name: 'Gemini' }] }], + null, + ); + + expect(inserted.map((i) => i.record.name)).toEqual(['Apollo', 'Gemini']); + expect(res.success).toBe(true); + expect(res.inserted).toBe(2); + }); + + it('returns a loud failure (never throws) for an unreadable body', async () => { + const protocol = new ObjectStackProtocolImplementation({} as never); + const res = await (protocol as any).applySeedBodies([{ nope: true }], null); + expect(res.success).toBe(false); + expect(res.error).toMatch(/no readable seed bodies/); + }); }); diff --git a/packages/objectql/src/protocol.ts b/packages/objectql/src/protocol.ts index a673e5caa3..519ac8f556 100644 --- a/packages/objectql/src/protocol.ts +++ b/packages/objectql/src/protocol.ts @@ -3785,11 +3785,30 @@ export class ObjectStackProtocolImplementation implements ObjectStackProtocol { organizationId?: string; actor?: string; message?: string; + /** + * INTERNAL — `publishPackageDrafts` publishes many drafts and batch-applies + * every seed body in ONE loader pass afterwards (cross-seed references need + * multi-pass over the whole set), so it suppresses the per-item apply here. + */ + _skipSeedApply?: boolean; }): Promise<{ success: boolean; version: string; seq: number; message?: string; + /** + * Present when a `seed` draft was published: the result of materializing + * its rows. Publishing the metadata ALWAYS succeeds independently — a + * seed-load problem is surfaced here, never thrown, so callers (and UIs) + * must check `seedApplied.success` instead of assuming data went live. + */ + seedApplied?: { + success: boolean; + inserted: number; + updated: number; + error?: string; + errors?: unknown[]; + }; }> { const singularType = PLURAL_TO_SINGULAR[request.type] ?? request.type; if (!ObjectStackProtocolImplementation.isOverlayAllowed(singularType) @@ -3840,12 +3859,27 @@ export class ObjectStackProtocolImplementation implements ObjectStackProtocol { }); // Create the object's table now so it's CRUD-able without a restart. await this.ensureObjectStorage(request.type, request.name); - return { + const response: { + success: boolean; + version: string; + seq: number; + message?: string; + seedApplied?: { success: boolean; inserted: number; updated: number; error?: string; errors?: unknown[] }; + } = { success: true, version: result.version, seq: result.seq, message: `Published draft — type=${request.type}, name=${request.name} [seq=${result.seq}]`, }; + // Publishing a `seed` is what makes its rows live — materialize them + // NOW (best-effort, never fails the publish) so every publish path + // (per-ref REST publish, the home banner, package publish-drafts) + // lands data, not just metadata. The body is already in hand from + // the promote — no read-back, so no org-scope resolution pitfalls. + if (singularType === 'seed' && !request._skipSeedApply) { + response.seedApplied = await this.applySeedBodies([result.item.body], orgId); + } + return response; } catch (err: any) { if (err instanceof ConflictError) { const conflict: any = new Error( @@ -3862,6 +3896,65 @@ export class ObjectStackProtocolImplementation implements ObjectStackProtocol { } } + /** + * Materialize published `seed` bodies into data rows via the SeedLoaderService + * (externalId-keyed upsert, multi-pass for cross-seed references). Passing ALL + * of a publish's seed bodies in ONE call lets a child seed reference a parent + * seed's rows regardless of publish order. Best-effort: any failure is + * returned, never thrown — publishing metadata must not be blocked by a data + * problem, but the caller surfaces `seedApplied` so the failure is LOUD. + */ + private async applySeedBodies( + bodies: unknown[], + organizationId: string | null, + ): Promise<{ success: boolean; inserted: number; updated: number; error?: string; errors?: unknown[] }> { + try { + const seeds = bodies.filter( + (b: any) => b && typeof b.object === 'string' && Array.isArray(b.records), + ); + if (seeds.length === 0) { + return { success: false, inserted: 0, updated: 0, error: 'seed apply: no readable seed bodies' }; + } + const { SeedLoaderService } = await import('./seed-loader.js'); + const { SeedLoaderRequestSchema } = await import('@objectstack/spec/data'); + // The loader only needs `getObject` from IMetadataService (dependency + // graph + field introspection); satisfy it from the protocol's own + // metadata reads so no kernel service lookup is required. + const metadataAdapter = { + getObject: async (name: string) => { + const wrapper: any = await (this as any).getMetaItem({ + type: 'object', + name, + ...(organizationId ? { organizationId } : {}), + }); + return wrapper?.item ?? wrapper ?? null; + }, + }; + const loader = new SeedLoaderService( + this.engine as any, + metadataAdapter as any, + console as any, + ); + const request = SeedLoaderRequestSchema.parse({ + seeds, + config: { + defaultMode: 'upsert', + multiPass: true, + ...(organizationId ? { organizationId } : {}), + }, + }); + const r = await loader.load(request); + return { + success: r.success, + inserted: r.summary.totalInserted, + updated: r.summary.totalUpdated, + ...(r.errors?.length ? { errors: r.errors } : {}), + }; + } catch (e: any) { + return { success: false, inserted: 0, updated: 0, error: e?.message ?? 'seed apply failed' }; + } + } + /** * List pending DRAFT metadata (ADR-0033) for the org, optionally narrowed * by `packageId` and/or `type`. The list reads of `getMetaItems` only see @@ -3911,6 +4004,8 @@ export class ObjectStackProtocolImplementation implements ObjectStackProtocol { failedCount: number; published: Array<{ type: string; name: string; version: string }>; failed: Array<{ type: string; name: string; error: string; code?: string }>; + /** Aggregate result of materializing every published `seed` (absent when no seeds). */ + seedApplied?: { success: boolean; inserted: number; updated: number; error?: string; errors?: unknown[] }; }> { await this.ensureOverlayIndex(); const orgId = request.organizationId ?? null; @@ -3920,14 +4015,33 @@ export class ObjectStackProtocolImplementation implements ObjectStackProtocol { const published: Array<{ type: string; name: string; version: string }> = []; const failed: Array<{ type: string; name: string; error: string; code?: string }> = []; - for (const d of drafts) { + // Structure first, seeds LAST — a seed's rows can only land after its + // object's table exists (publishMetaItem creates it). Within the seeds we + // batch-apply every body in ONE loader pass below (multi-pass reference + // resolution across the whole set), so per-item apply is suppressed. + const ordered = [ + ...drafts.filter((d) => d.type !== 'seed'), + ...drafts.filter((d) => d.type === 'seed'), + ]; + const seedBodies: unknown[] = []; + + for (const d of ordered) { try { + if (d.type === 'seed') { + // Capture the body BEFORE promote (the draft row is deleted by + // the promote, and a post-publish read-back has org-scope + // resolution pitfalls — reading the draft is unambiguous). + const ref = { type: d.type, name: d.name, org: orgId ?? 'env' } as unknown as Parameters[0]; + const draft = await repo.get(ref, { state: 'draft' }); + if (draft?.body) seedBodies.push(draft.body); + } const r = await this.publishMetaItem({ type: d.type, name: d.name, ...(request.organizationId ? { organizationId: request.organizationId } : {}), ...(request.actor ? { actor: request.actor } : {}), message: `publish app package '${request.packageId}'`, + _skipSeedApply: true, }); published.push({ type: d.type, name: d.name, version: r.version }); } catch (e: any) { @@ -3946,6 +4060,9 @@ export class ObjectStackProtocolImplementation implements ObjectStackProtocol { failedCount: failed.length, published, failed, + ...(seedBodies.length > 0 + ? { seedApplied: await this.applySeedBodies(seedBodies, orgId) } + : {}), }; } diff --git a/packages/objectql/src/seed-loader.ts b/packages/objectql/src/seed-loader.ts new file mode 100644 index 0000000000..06bad9b2db --- /dev/null +++ b/packages/objectql/src/seed-loader.ts @@ -0,0 +1,848 @@ +// Copyright (c) 2025 ObjectStack. Licensed under the Apache-2.0 license. + +import type { IDataEngine, IMetadataService, ISeedLoaderService } from '@objectstack/spec/contracts'; +import type { + SeedLoaderRequest, + SeedLoaderResult, + SeedLoaderConfig, + SeedLoaderConfigInput, + ObjectDependencyGraph, + ObjectDependencyNode, + ReferenceResolution, + ReferenceResolutionError, + SeedLoadResult, + Seed, +} from '@objectstack/spec/data'; +import { SeedLoaderConfigSchema } from '@objectstack/spec/data'; +import { resolveSeedRecord } from '@objectstack/formula'; + +interface Logger { + info(message: string, meta?: Record): void; + warn(message: string, meta?: Record): void; + error(message: string, error?: Error, meta?: Record): void; + debug(message: string, meta?: Record): void; +} + +/** Default field used for externalId matching on target objects */ +const DEFAULT_EXTERNAL_ID_FIELD = 'name'; + +/** + * SeedLoaderService — Runtime implementation of ISeedLoaderService + * + * Provides metadata-driven seed data loading with: + * - Automatic lookup/master_detail reference resolution via externalId + * - Topological dependency ordering (parents before children) + * - Multi-pass loading for circular references + * - Dry-run validation mode + * - Upsert support honoring SeedSchema mode + * - Actionable error reporting + */ +export class SeedLoaderService implements ISeedLoaderService { + private engine: IDataEngine; + private metadata: IMetadataService; + private logger: Logger; + + constructor(engine: IDataEngine, metadata: IMetadataService, logger: Logger) { + this.engine = engine; + this.metadata = metadata; + this.logger = logger; + } + + // ========================================================================== + // Public API + // ========================================================================== + + async load(request: SeedLoaderRequest): Promise { + const startTime = Date.now(); + const config = request.config; + const allErrors: ReferenceResolutionError[] = []; + const allResults: SeedLoadResult[] = []; + + // 1. Filter datasets by environment + const datasets = this.filterByEnv(request.seeds, config.env); + + if (datasets.length === 0) { + return this.buildEmptyResult(config, Date.now() - startTime); + } + + // 2. Build dependency graph + const objectNames = datasets.map(d => d.object); + const graph = await this.buildDependencyGraph(objectNames); + + this.logger.info('[SeedLoader] Dependency graph built', { + objects: objectNames.length, + insertOrder: graph.insertOrder, + circularDeps: graph.circularDependencies.length, + }); + + // 3. Order datasets by topological insert order + const orderedDatasets = this.orderDatasets(datasets, graph.insertOrder); + + // 4. Build reference lookup map from metadata (field → target object) + const refMap = this.buildReferenceMap(graph); + + // 5. Pass 1: Insert/upsert records, resolving references + const insertedRecords = new Map>(); // object → externalIdValue → internalId + const deferredUpdates: DeferredUpdate[] = []; + + for (const dataset of orderedDatasets) { + const result = await this.loadDataset( + dataset, config, refMap, insertedRecords, deferredUpdates, allErrors + ); + allResults.push(result); + + if (config.haltOnError && result.errored > 0) { + this.logger.warn('[SeedLoader] Halting on first error', { object: dataset.object }); + break; + } + } + + // 6. Pass 2: Resolve deferred references (circular dependencies) + if (config.multiPass && deferredUpdates.length > 0 && !config.dryRun) { + this.logger.info('[SeedLoader] Pass 2: resolving deferred references', { + count: deferredUpdates.length, + }); + await this.resolveDeferredUpdates(deferredUpdates, insertedRecords, allResults, allErrors, config.organizationId); + } + + // 7. Build final result + const durationMs = Date.now() - startTime; + return this.buildResult(config, graph, allResults, allErrors, durationMs); + } + + async buildDependencyGraph(objectNames: string[]): Promise { + const nodes: ObjectDependencyNode[] = []; + const objectSet = new Set(objectNames); + + for (const objectName of objectNames) { + const objDef = await this.metadata.getObject(objectName) as any; + const dependsOn: string[] = []; + const references: ReferenceResolution[] = []; + + if (objDef && objDef.fields) { + const fields = objDef.fields as Record; + for (const [fieldName, fieldDef] of Object.entries(fields)) { + if ( + (fieldDef.type === 'lookup' || fieldDef.type === 'master_detail') && + fieldDef.reference + ) { + const targetObject = fieldDef.reference as string; + + // Track dependency ordering only for objects within the graph + if (objectSet.has(targetObject) && !dependsOn.includes(targetObject)) { + dependsOn.push(targetObject); + } + + // Track ALL references for resolution (target may exist in database) + references.push({ + field: fieldName, + targetObject, + targetField: DEFAULT_EXTERNAL_ID_FIELD, + fieldType: fieldDef.type as 'lookup' | 'master_detail', + }); + } + } + } + + nodes.push({ object: objectName, dependsOn, references }); + } + + // Topological sort + const { insertOrder, circularDependencies } = this.topologicalSort(nodes); + + return { nodes, insertOrder, circularDependencies }; + } + + async validate(datasets: Seed[], config?: SeedLoaderConfigInput): Promise { + const parsedConfig = SeedLoaderConfigSchema.parse({ ...config, dryRun: true }); + return this.load({ seeds: datasets, config: parsedConfig }); + } + + // ========================================================================== + // Internal: Seed Loading + // ========================================================================== + + private async loadDataset( + dataset: Seed, + config: SeedLoaderConfig, + refMap: Map, + insertedRecords: Map>, + deferredUpdates: DeferredUpdate[], + allErrors: ReferenceResolutionError[], + ): Promise { + const objectName = dataset.object; + const mode = dataset.mode || config.defaultMode; + const externalId = dataset.externalId || 'name'; + + let inserted = 0; + let updated = 0; + let skipped = 0; + let errored = 0; + let referencesResolved = 0; + let referencesDeferred = 0; + const errors: ReferenceResolutionError[] = []; + + // Ensure the object's record map exists + if (!insertedRecords.has(objectName)) { + insertedRecords.set(objectName, new Map()); + } + + // Pre-load existing records for upsert matching. When a target + // organization is set, scope the lookup so each tenant gets its + // own copy (otherwise upsert would clobber other tenants' rows + // that share the same natural key — e.g. `name: 'Acme Corp'`). + let existingRecords: Map | undefined; + if ((mode === 'upsert' || mode === 'update' || mode === 'ignore') && !config.dryRun) { + existingRecords = await this.loadExistingRecords( + objectName, + externalId, + config.organizationId, + ); + } + + // Get reference resolutions for this object + const objectRefs = refMap.get(objectName) || []; + + // Pin a single `now()` snapshot for the entire dataset so multi-pass + // loads see one logical clock — the M9 determinism guarantee for seeds. + const seedNow = new Date(); + + // Identity/context bound to seed CEL expressions. `os.user` / `os.org` + // resolve from here, so `owner_id: cel\`os.user.id\`` works. When no + // identity is supplied, `os.user` / `os.org` are simply unbound and any + // record that references them fails loudly below (rather than silently + // writing a raw Expression envelope into the column). + const seedIdentity = config.identity; + const baseEvalCtx = { + now: seedNow, + user: seedIdentity?.user, + // Fall back to the per-tenant organizationId so `os.org.id` resolves + // during per-org replay even without an explicit identity.org. + org: seedIdentity?.org ?? (config.organizationId ? { id: config.organizationId } : undefined), + env: config.env, + }; + + for (let i = 0; i < dataset.records.length; i++) { + // Resolve any embedded Expression envelopes (e.g. `cel\`daysFromNow(30)\``, + // `cel\`os.user.id\``) BEFORE reference resolution so downstream lookups + // see resolved values. + const seedResult = resolveSeedRecord( + dataset.records[i] as Record, + baseEvalCtx, + ); + if (!seedResult.ok) { + // LOUD FAILURE: a record whose dynamic values cannot be resolved is + // dropped — but never silently. Record an actionable error (so it + // surfaces in result.errors and flips success=false) instead of + // writing the unresolved Expression envelope into the database. + errored++; + const error: ReferenceResolutionError = { + sourceObject: objectName, + field: '(expression)', + targetObject: objectName, + targetField: '(expression)', + attemptedValue: dataset.records[i], + recordIndex: i, + message: + `Cannot resolve dynamic seed values for ${objectName} record #${i}: ${seedResult.error.message}. ` + + 'Records using cel`os.user.id` / cel`os.org.id` require a seed identity — ' + + 'ensure a system/admin user exists before seeding (see SeedLoaderConfig.identity).', + }; + errors.push(error); + allErrors.push(error); + this.logger.warn(`[SeedLoader] ${error.message}`); + continue; + } + const record = { ...(seedResult.value as Record) }; + + // Per-tenant tagging: when a target org is set, stamp every + // seeded row with it (unless the record itself already supplies + // an explicit organization_id — respect dataset author overrides). + // Skipped objects that don't declare `organization_id` will have + // the extra key silently ignored by the engine. + if (config.organizationId && record['organization_id'] == null) { + record['organization_id'] = config.organizationId; + } + + // Resolve references + for (const ref of objectRefs) { + const fieldValue = record[ref.field]; + if (fieldValue === undefined || fieldValue === null) continue; + + // LOUD FAILURE: a reference must be a natural-key string (or an + // internal id). An object value — e.g. the wrapper `{ externalId: 'X' }` + // — never resolves: it would otherwise fall through unresolved and reach + // the driver as a non-bindable value ("SQLite3 can only bind ..."). This + // used to be silently skipped (and only crashed on a persistent DB's + // update path), so catch it here and report the actionable fix instead. + if (typeof fieldValue === 'object') { + const wrapped = (fieldValue as Record).externalId; + const hint = + wrapped !== undefined + ? ` Pass the natural key directly: ${ref.field}: ${JSON.stringify(wrapped)}.` + : ` Pass the target's ${ref.targetField} value as a plain string.`; + const error: ReferenceResolutionError = { + sourceObject: objectName, + field: ref.field, + targetObject: ref.targetObject, + targetField: ref.targetField, + attemptedValue: fieldValue, + recordIndex: i, + message: + `Invalid reference for ${objectName}.${ref.field}: expected a ` + + `${ref.targetObject}.${ref.targetField} natural-key string but got an object.${hint}`, + }; + errors.push(error); + allErrors.push(error); + this.logger.warn(`[SeedLoader] ${error.message}`, { recordIndex: i }); + // Drop the unresolvable value so it never reaches the driver. + record[ref.field] = null; + continue; + } + + // Skip if value looks like an internal ID (not a natural key) + if (typeof fieldValue !== 'string' || this.looksLikeInternalId(fieldValue)) continue; + + // Try to resolve via already-inserted records + const targetMap = insertedRecords.get(ref.targetObject); + const resolvedId = targetMap?.get(String(fieldValue)); + + if (resolvedId) { + record[ref.field] = resolvedId; + referencesResolved++; + } else if (!config.dryRun) { + // Try to resolve from existing data in the database + const dbId = await this.resolveFromDatabase(ref.targetObject, ref.targetField, fieldValue, config.organizationId); + if (dbId) { + record[ref.field] = dbId; + referencesResolved++; + } else if (config.multiPass) { + // Defer to pass 2 + record[ref.field] = null; + deferredUpdates.push({ + objectName, + recordExternalId: String(record[externalId] ?? ''), + field: ref.field, + targetObject: ref.targetObject, + targetField: ref.targetField, + attemptedValue: fieldValue, + recordIndex: i, + }); + referencesDeferred++; + } else { + // Cannot resolve - record error + const error: ReferenceResolutionError = { + sourceObject: objectName, + field: ref.field, + targetObject: ref.targetObject, + targetField: ref.targetField, + attemptedValue: fieldValue, + recordIndex: i, + message: `Cannot resolve reference: ${objectName}.${ref.field} = '${fieldValue}' → ${ref.targetObject}.${ref.targetField} not found`, + }; + errors.push(error); + allErrors.push(error); + } + } else { + // Dry-run: attempt resolution, report error if not found + const targetMap2 = insertedRecords.get(ref.targetObject); + if (!targetMap2?.has(String(fieldValue))) { + const error: ReferenceResolutionError = { + sourceObject: objectName, + field: ref.field, + targetObject: ref.targetObject, + targetField: ref.targetField, + attemptedValue: fieldValue, + recordIndex: i, + message: `[dry-run] Reference may not resolve: ${objectName}.${ref.field} = '${fieldValue}' → ${ref.targetObject}.${ref.targetField}`, + }; + errors.push(error); + allErrors.push(error); + } + } + } + + // Insert/upsert the record + if (!config.dryRun) { + try { + const result = await this.writeRecord( + objectName, record, mode, externalId, existingRecords + ); + + if (result.action === 'inserted') inserted++; + else if (result.action === 'updated') updated++; + else if (result.action === 'skipped') skipped++; + + // Track the inserted/updated record's ID for reference resolution + const externalIdValue = String(record[externalId] ?? ''); + const internalId = result.id; + if (externalIdValue && internalId) { + insertedRecords.get(objectName)!.set(externalIdValue, String(internalId)); + } + } catch (err: any) { + // LOUD FAILURE: write errors were previously only counted + + // warn-logged, so dropped rows were invisible in result.errors and + // the boot summary. Surface them as actionable errors too, so the + // overall load is marked unsuccessful and the reason is reported. + errored++; + const error: ReferenceResolutionError = { + sourceObject: objectName, + field: '(write)', + targetObject: objectName, + targetField: externalId, + attemptedValue: record[externalId] ?? null, + recordIndex: i, + message: `Failed to write ${objectName} record #${i} (${externalId}=${String(record[externalId] ?? '')}): ${err.message}`, + }; + errors.push(error); + allErrors.push(error); + this.logger.warn(`[SeedLoader] ${error.message}`, { recordIndex: i }); + } + } else { + // Dry-run: simulate insert tracking + const externalIdValue = String(record[externalId] ?? ''); + if (externalIdValue) { + insertedRecords.get(objectName)!.set(externalIdValue, `dry-run-id-${i}`); + } + inserted++; // Count as "would be inserted" + } + } + + return { + object: objectName, + mode, + inserted, + updated, + skipped, + errored, + total: dataset.records.length, + referencesResolved, + referencesDeferred, + errors, + }; + } + + // ========================================================================== + // Internal: Reference Resolution + // ========================================================================== + + private async resolveFromDatabase( + targetObject: string, + targetField: string, + value: unknown, + organizationId?: string, + ): Promise { + try { + const where: Record = { [targetField]: value }; + // Per-tenant replay: when scoping is requested, only consider + // rows that belong to the target tenant so cross-tenant rows + // never get borrowed as a "resolved" reference (would silently + // create a cross-org FK). + if (organizationId) where.organization_id = organizationId; + const records = await this.engine.find(targetObject, { + where, + fields: ['id'], + limit: 1, + context: { isSystem: true }, + } as any); + if (records && records.length > 0) { + return String(records[0].id || records[0]._id); + } + } catch { + // Target object may not exist yet + } + return null; + } + + private async resolveDeferredUpdates( + deferredUpdates: DeferredUpdate[], + insertedRecords: Map>, + allResults: SeedLoadResult[], + allErrors: ReferenceResolutionError[], + organizationId?: string, + ): Promise { + for (const deferred of deferredUpdates) { + // Try to resolve from inserted records + const targetMap = insertedRecords.get(deferred.targetObject); + let resolvedId = targetMap?.get(String(deferred.attemptedValue)); + + // Try database fallback + if (!resolvedId) { + resolvedId = (await this.resolveFromDatabase( + deferred.targetObject, deferred.targetField, deferred.attemptedValue, organizationId + )) ?? undefined; + } + + if (resolvedId) { + // Find the record and update the reference + const objectRecordMap = insertedRecords.get(deferred.objectName); + const recordId = objectRecordMap?.get(deferred.recordExternalId); + + if (recordId) { + try { + await this.engine.update(deferred.objectName, { + id: recordId, + [deferred.field]: resolvedId, + }, { context: { isSystem: true } } as any); + + // Update result stats + const resultEntry = allResults.find(r => r.object === deferred.objectName); + if (resultEntry) { + resultEntry.referencesResolved++; + resultEntry.referencesDeferred--; + } + } catch (err: any) { + this.logger.warn('[SeedLoader] Failed to resolve deferred reference', { + object: deferred.objectName, + field: deferred.field, + error: err.message, + }); + } + } + } else { + // Still unresolved after pass 2 + const error: ReferenceResolutionError = { + sourceObject: deferred.objectName, + field: deferred.field, + targetObject: deferred.targetObject, + targetField: deferred.targetField, + attemptedValue: deferred.attemptedValue, + recordIndex: deferred.recordIndex, + message: `Deferred reference unresolved after pass 2: ${deferred.objectName}.${deferred.field} = '${deferred.attemptedValue}' → ${deferred.targetObject}.${deferred.targetField} not found`, + }; + + const resultEntry = allResults.find(r => r.object === deferred.objectName); + if (resultEntry) { + resultEntry.errors.push(error); + } + allErrors.push(error); + } + } + } + + // ========================================================================== + // Internal: Write Operations + // ========================================================================== + + /** + * Seed writes always run as a privileged system context. This bypasses + * RBAC checks (so seeds can target system tables like `sys_*`) and + * disables the SecurityPlugin's auto-injection of `organization_id` / + * `owner_id` — seeds either declare those fields explicitly per + * record, or are intentionally cross-tenant / global. + */ + private static readonly SEED_OPTIONS = { context: { isSystem: true } } as const; + + private async writeRecord( + objectName: string, + record: Record, + mode: string, + externalId: string, + existingRecords?: Map, + ): Promise<{ action: 'inserted' | 'updated' | 'skipped'; id?: string }> { + const externalIdValue = record[externalId]; + const existing = existingRecords?.get(String(externalIdValue ?? '')); + const opts = SeedLoaderService.SEED_OPTIONS as any; + + switch (mode) { + case 'insert': { + const result = await this.engine.insert(objectName, record, opts); + return { action: 'inserted', id: this.extractId(result) }; + } + + case 'update': { + if (!existing) { + return { action: 'skipped' }; + } + const id = this.extractId(existing); + await this.engine.update(objectName, { ...record, id }, opts); + return { action: 'updated', id }; + } + + case 'upsert': { + if (existing) { + const id = this.extractId(existing); + await this.engine.update(objectName, { ...record, id }, opts); + return { action: 'updated', id }; + } else { + const result = await this.engine.insert(objectName, record, opts); + return { action: 'inserted', id: this.extractId(result) }; + } + } + + case 'ignore': { + if (existing) { + return { action: 'skipped', id: this.extractId(existing) }; + } + const result = await this.engine.insert(objectName, record, opts); + return { action: 'inserted', id: this.extractId(result) }; + } + + case 'replace': { + // Replace mode: just insert (caller should have cleared the table) + const result = await this.engine.insert(objectName, record, opts); + return { action: 'inserted', id: this.extractId(result) }; + } + + default: { + const result = await this.engine.insert(objectName, record, opts); + return { action: 'inserted', id: this.extractId(result) }; + } + } + } + + // ========================================================================== + // Internal: Dependency Graph + // ========================================================================== + + /** + * Kahn's algorithm for topological sort with cycle detection. + */ + private topologicalSort( + nodes: ObjectDependencyNode[], + ): { insertOrder: string[]; circularDependencies: string[][] } { + const inDegree = new Map(); + const adjacency = new Map(); + const objectSet = new Set(nodes.map(n => n.object)); + + // Initialize + for (const node of nodes) { + inDegree.set(node.object, 0); + adjacency.set(node.object, []); + } + + // Build adjacency list and in-degree counts + for (const node of nodes) { + for (const dep of node.dependsOn) { + // Exclude self-references from ordering (e.g., employee.manager_id → employee). + // Self-referencing fields are still tracked in node.references for resolution. + if (objectSet.has(dep) && dep !== node.object) { + adjacency.get(dep)!.push(node.object); + inDegree.set(node.object, (inDegree.get(node.object) || 0) + 1); + } + } + } + + // Kahn's algorithm + const queue: string[] = []; + for (const [obj, degree] of inDegree) { + if (degree === 0) queue.push(obj); + } + + const insertOrder: string[] = []; + while (queue.length > 0) { + const current = queue.shift()!; + insertOrder.push(current); + + for (const neighbor of (adjacency.get(current) || [])) { + const newDegree = (inDegree.get(neighbor) || 0) - 1; + inDegree.set(neighbor, newDegree); + if (newDegree === 0) { + queue.push(neighbor); + } + } + } + + // Detect circular dependencies + const circularDependencies: string[][] = []; + const remaining = nodes.filter(n => !insertOrder.includes(n.object)); + + if (remaining.length > 0) { + // Find cycles using DFS + const cycles = this.findCycles(remaining); + circularDependencies.push(...cycles); + + // Add remaining objects to insertOrder (they'll need multi-pass) + for (const node of remaining) { + if (!insertOrder.includes(node.object)) { + insertOrder.push(node.object); + } + } + } + + return { insertOrder, circularDependencies }; + } + + private findCycles(nodes: ObjectDependencyNode[]): string[][] { + const cycles: string[][] = []; + const nodeMap = new Map(nodes.map(n => [n.object, n])); + const visited = new Set(); + const inStack = new Set(); + + const dfs = (current: string, path: string[]) => { + if (inStack.has(current)) { + // Found a cycle + const cycleStart = path.indexOf(current); + if (cycleStart !== -1) { + cycles.push([...path.slice(cycleStart), current]); + } + return; + } + if (visited.has(current)) return; + + visited.add(current); + inStack.add(current); + path.push(current); + + const node = nodeMap.get(current); + if (node) { + for (const dep of node.dependsOn) { + if (nodeMap.has(dep)) { + dfs(dep, [...path]); + } + } + } + + inStack.delete(current); + }; + + for (const node of nodes) { + if (!visited.has(node.object)) { + dfs(node.object, []); + } + } + + return cycles; + } + + // ========================================================================== + // Internal: Helpers + // ========================================================================== + + private filterByEnv(datasets: Seed[], env?: string): Seed[] { + if (!env) return datasets; + return datasets.filter(d => (d.env as string[]).includes(env)); + } + + private orderDatasets(datasets: Seed[], insertOrder: string[]): Seed[] { + const orderMap = new Map(insertOrder.map((name, i) => [name, i])); + return [...datasets].sort((a, b) => { + const orderA = orderMap.get(a.object) ?? Number.MAX_SAFE_INTEGER; + const orderB = orderMap.get(b.object) ?? Number.MAX_SAFE_INTEGER; + return orderA - orderB; + }); + } + + private buildReferenceMap(graph: ObjectDependencyGraph): Map { + const map = new Map(); + for (const node of graph.nodes) { + if (node.references.length > 0) { + map.set(node.object, node.references); + } + } + return map; + } + + private async loadExistingRecords( + objectName: string, + externalId: string, + organizationId?: string, + ): Promise> { + const map = new Map(); + try { + const findArgs: Record = { + fields: ['id', externalId], + context: { isSystem: true }, + }; + // Per-tenant replay: restrict to the target tenant's own rows + // so upsert key matching never returns another tenant's record + // (would silently steal/overwrite rows across orgs). + if (organizationId) findArgs.where = { organization_id: organizationId }; + const records = await this.engine.find(objectName, findArgs as any); + for (const record of records || []) { + const key = String(record[externalId] ?? ''); + if (key) { + map.set(key, record); + } + } + } catch { + // Object may not have records yet + } + return map; + } + + private looksLikeInternalId(value: string): boolean { + // UUID v4 pattern + if (/^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i.test(value)) { + return true; + } + // MongoDB ObjectId pattern (24 hex chars) + if (/^[0-9a-f]{24}$/i.test(value)) { + return true; + } + return false; + } + + private extractId(record: any): string | undefined { + if (!record) return undefined; + return String(record.id || record._id || ''); + } + + private buildEmptyResult(config: SeedLoaderConfig, durationMs: number): SeedLoaderResult { + return { + success: true, + dryRun: config.dryRun, + dependencyGraph: { nodes: [], insertOrder: [], circularDependencies: [] }, + results: [], + errors: [], + summary: { + objectsProcessed: 0, + totalRecords: 0, + totalInserted: 0, + totalUpdated: 0, + totalSkipped: 0, + totalErrored: 0, + totalReferencesResolved: 0, + totalReferencesDeferred: 0, + circularDependencyCount: 0, + durationMs, + }, + }; + } + + private buildResult( + config: SeedLoaderConfig, + graph: ObjectDependencyGraph, + results: SeedLoadResult[], + errors: ReferenceResolutionError[], + durationMs: number, + ): SeedLoaderResult { + const summary = { + objectsProcessed: results.length, + totalRecords: results.reduce((sum, r) => sum + r.total, 0), + totalInserted: results.reduce((sum, r) => sum + r.inserted, 0), + totalUpdated: results.reduce((sum, r) => sum + r.updated, 0), + totalSkipped: results.reduce((sum, r) => sum + r.skipped, 0), + totalErrored: results.reduce((sum, r) => sum + r.errored, 0), + totalReferencesResolved: results.reduce((sum, r) => sum + r.referencesResolved, 0), + totalReferencesDeferred: results.reduce((sum, r) => sum + r.referencesDeferred, 0), + circularDependencyCount: graph.circularDependencies.length, + durationMs, + }; + + const hasErrors = errors.length > 0 || summary.totalErrored > 0; + + return { + success: !hasErrors, + dryRun: config.dryRun, + dependencyGraph: graph, + results, + errors, + summary, + }; + } +} + +// ========================================================================== +// Internal Types +// ========================================================================== + +interface DeferredUpdate { + objectName: string; + recordExternalId: string; + field: string; + targetObject: string; + targetField: string; + attemptedValue: unknown; + recordIndex: number; +} diff --git a/packages/runtime/src/http-dispatcher.ts b/packages/runtime/src/http-dispatcher.ts index d389885c3d..7f8c1f29a8 100644 --- a/packages/runtime/src/http-dispatcher.ts +++ b/packages/runtime/src/http-dispatcher.ts @@ -1771,23 +1771,28 @@ export class HttpDispatcher { ...(body?.actor ? { actor: body.actor } : {}), }); // Publishing a `seed` draft is what actually loads its - // rows. Best-effort + idempotent (upsert): apply every - // just-published seed now so the data is live the moment - // the user clicks publish. A seed-load failure NEVER - // fails the publish — it is surfaced under `seedApplied`. - try { - const seedNames = ((result as any)?.published ?? []) - .filter((p: any) => p?.type === 'seed') - .map((p: any) => p.name as string); - if (seedNames.length > 0) { - (result as any).seedApplied = await this.applyPublishedSeeds( - seedNames, - organizationId, - _context, - ); + // rows. The objectql protocol now batch-applies seeds + // inside `publishPackageDrafts` itself (so EVERY publish + // path — incl. the per-ref REST publish — materializes + // data) and reports under `seedApplied`. Only fall back + // to the route-level apply for custom protocols that + // don't self-apply — never run both, or an externalId-less + // seed would double-insert. + if ((result as any)?.seedApplied === undefined) { + try { + const seedNames = ((result as any)?.published ?? []) + .filter((p: any) => p?.type === 'seed') + .map((p: any) => p.name as string); + if (seedNames.length > 0) { + (result as any).seedApplied = await this.applyPublishedSeeds( + seedNames, + organizationId, + _context, + ); + } + } catch (e: any) { + (result as any).seedApplied = { success: false, error: e?.message ?? 'seed apply failed' }; } - } catch (e: any) { - (result as any).seedApplied = { success: false, error: e?.message ?? 'seed apply failed' }; } return { handled: true, response: this.success(result) }; } catch (e: any) { diff --git a/packages/runtime/src/seed-loader.ts b/packages/runtime/src/seed-loader.ts index 06bad9b2db..1de8117681 100644 --- a/packages/runtime/src/seed-loader.ts +++ b/packages/runtime/src/seed-loader.ts @@ -1,848 +1,9 @@ // Copyright (c) 2025 ObjectStack. Licensed under the Apache-2.0 license. - -import type { IDataEngine, IMetadataService, ISeedLoaderService } from '@objectstack/spec/contracts'; -import type { - SeedLoaderRequest, - SeedLoaderResult, - SeedLoaderConfig, - SeedLoaderConfigInput, - ObjectDependencyGraph, - ObjectDependencyNode, - ReferenceResolution, - ReferenceResolutionError, - SeedLoadResult, - Seed, -} from '@objectstack/spec/data'; -import { SeedLoaderConfigSchema } from '@objectstack/spec/data'; -import { resolveSeedRecord } from '@objectstack/formula'; - -interface Logger { - info(message: string, meta?: Record): void; - warn(message: string, meta?: Record): void; - error(message: string, error?: Error, meta?: Record): void; - debug(message: string, meta?: Record): void; -} - -/** Default field used for externalId matching on target objects */ -const DEFAULT_EXTERNAL_ID_FIELD = 'name'; - -/** - * SeedLoaderService — Runtime implementation of ISeedLoaderService - * - * Provides metadata-driven seed data loading with: - * - Automatic lookup/master_detail reference resolution via externalId - * - Topological dependency ordering (parents before children) - * - Multi-pass loading for circular references - * - Dry-run validation mode - * - Upsert support honoring SeedSchema mode - * - Actionable error reporting - */ -export class SeedLoaderService implements ISeedLoaderService { - private engine: IDataEngine; - private metadata: IMetadataService; - private logger: Logger; - - constructor(engine: IDataEngine, metadata: IMetadataService, logger: Logger) { - this.engine = engine; - this.metadata = metadata; - this.logger = logger; - } - - // ========================================================================== - // Public API - // ========================================================================== - - async load(request: SeedLoaderRequest): Promise { - const startTime = Date.now(); - const config = request.config; - const allErrors: ReferenceResolutionError[] = []; - const allResults: SeedLoadResult[] = []; - - // 1. Filter datasets by environment - const datasets = this.filterByEnv(request.seeds, config.env); - - if (datasets.length === 0) { - return this.buildEmptyResult(config, Date.now() - startTime); - } - - // 2. Build dependency graph - const objectNames = datasets.map(d => d.object); - const graph = await this.buildDependencyGraph(objectNames); - - this.logger.info('[SeedLoader] Dependency graph built', { - objects: objectNames.length, - insertOrder: graph.insertOrder, - circularDeps: graph.circularDependencies.length, - }); - - // 3. Order datasets by topological insert order - const orderedDatasets = this.orderDatasets(datasets, graph.insertOrder); - - // 4. Build reference lookup map from metadata (field → target object) - const refMap = this.buildReferenceMap(graph); - - // 5. Pass 1: Insert/upsert records, resolving references - const insertedRecords = new Map>(); // object → externalIdValue → internalId - const deferredUpdates: DeferredUpdate[] = []; - - for (const dataset of orderedDatasets) { - const result = await this.loadDataset( - dataset, config, refMap, insertedRecords, deferredUpdates, allErrors - ); - allResults.push(result); - - if (config.haltOnError && result.errored > 0) { - this.logger.warn('[SeedLoader] Halting on first error', { object: dataset.object }); - break; - } - } - - // 6. Pass 2: Resolve deferred references (circular dependencies) - if (config.multiPass && deferredUpdates.length > 0 && !config.dryRun) { - this.logger.info('[SeedLoader] Pass 2: resolving deferred references', { - count: deferredUpdates.length, - }); - await this.resolveDeferredUpdates(deferredUpdates, insertedRecords, allResults, allErrors, config.organizationId); - } - - // 7. Build final result - const durationMs = Date.now() - startTime; - return this.buildResult(config, graph, allResults, allErrors, durationMs); - } - - async buildDependencyGraph(objectNames: string[]): Promise { - const nodes: ObjectDependencyNode[] = []; - const objectSet = new Set(objectNames); - - for (const objectName of objectNames) { - const objDef = await this.metadata.getObject(objectName) as any; - const dependsOn: string[] = []; - const references: ReferenceResolution[] = []; - - if (objDef && objDef.fields) { - const fields = objDef.fields as Record; - for (const [fieldName, fieldDef] of Object.entries(fields)) { - if ( - (fieldDef.type === 'lookup' || fieldDef.type === 'master_detail') && - fieldDef.reference - ) { - const targetObject = fieldDef.reference as string; - - // Track dependency ordering only for objects within the graph - if (objectSet.has(targetObject) && !dependsOn.includes(targetObject)) { - dependsOn.push(targetObject); - } - - // Track ALL references for resolution (target may exist in database) - references.push({ - field: fieldName, - targetObject, - targetField: DEFAULT_EXTERNAL_ID_FIELD, - fieldType: fieldDef.type as 'lookup' | 'master_detail', - }); - } - } - } - - nodes.push({ object: objectName, dependsOn, references }); - } - - // Topological sort - const { insertOrder, circularDependencies } = this.topologicalSort(nodes); - - return { nodes, insertOrder, circularDependencies }; - } - - async validate(datasets: Seed[], config?: SeedLoaderConfigInput): Promise { - const parsedConfig = SeedLoaderConfigSchema.parse({ ...config, dryRun: true }); - return this.load({ seeds: datasets, config: parsedConfig }); - } - - // ========================================================================== - // Internal: Seed Loading - // ========================================================================== - - private async loadDataset( - dataset: Seed, - config: SeedLoaderConfig, - refMap: Map, - insertedRecords: Map>, - deferredUpdates: DeferredUpdate[], - allErrors: ReferenceResolutionError[], - ): Promise { - const objectName = dataset.object; - const mode = dataset.mode || config.defaultMode; - const externalId = dataset.externalId || 'name'; - - let inserted = 0; - let updated = 0; - let skipped = 0; - let errored = 0; - let referencesResolved = 0; - let referencesDeferred = 0; - const errors: ReferenceResolutionError[] = []; - - // Ensure the object's record map exists - if (!insertedRecords.has(objectName)) { - insertedRecords.set(objectName, new Map()); - } - - // Pre-load existing records for upsert matching. When a target - // organization is set, scope the lookup so each tenant gets its - // own copy (otherwise upsert would clobber other tenants' rows - // that share the same natural key — e.g. `name: 'Acme Corp'`). - let existingRecords: Map | undefined; - if ((mode === 'upsert' || mode === 'update' || mode === 'ignore') && !config.dryRun) { - existingRecords = await this.loadExistingRecords( - objectName, - externalId, - config.organizationId, - ); - } - - // Get reference resolutions for this object - const objectRefs = refMap.get(objectName) || []; - - // Pin a single `now()` snapshot for the entire dataset so multi-pass - // loads see one logical clock — the M9 determinism guarantee for seeds. - const seedNow = new Date(); - - // Identity/context bound to seed CEL expressions. `os.user` / `os.org` - // resolve from here, so `owner_id: cel\`os.user.id\`` works. When no - // identity is supplied, `os.user` / `os.org` are simply unbound and any - // record that references them fails loudly below (rather than silently - // writing a raw Expression envelope into the column). - const seedIdentity = config.identity; - const baseEvalCtx = { - now: seedNow, - user: seedIdentity?.user, - // Fall back to the per-tenant organizationId so `os.org.id` resolves - // during per-org replay even without an explicit identity.org. - org: seedIdentity?.org ?? (config.organizationId ? { id: config.organizationId } : undefined), - env: config.env, - }; - - for (let i = 0; i < dataset.records.length; i++) { - // Resolve any embedded Expression envelopes (e.g. `cel\`daysFromNow(30)\``, - // `cel\`os.user.id\``) BEFORE reference resolution so downstream lookups - // see resolved values. - const seedResult = resolveSeedRecord( - dataset.records[i] as Record, - baseEvalCtx, - ); - if (!seedResult.ok) { - // LOUD FAILURE: a record whose dynamic values cannot be resolved is - // dropped — but never silently. Record an actionable error (so it - // surfaces in result.errors and flips success=false) instead of - // writing the unresolved Expression envelope into the database. - errored++; - const error: ReferenceResolutionError = { - sourceObject: objectName, - field: '(expression)', - targetObject: objectName, - targetField: '(expression)', - attemptedValue: dataset.records[i], - recordIndex: i, - message: - `Cannot resolve dynamic seed values for ${objectName} record #${i}: ${seedResult.error.message}. ` + - 'Records using cel`os.user.id` / cel`os.org.id` require a seed identity — ' + - 'ensure a system/admin user exists before seeding (see SeedLoaderConfig.identity).', - }; - errors.push(error); - allErrors.push(error); - this.logger.warn(`[SeedLoader] ${error.message}`); - continue; - } - const record = { ...(seedResult.value as Record) }; - - // Per-tenant tagging: when a target org is set, stamp every - // seeded row with it (unless the record itself already supplies - // an explicit organization_id — respect dataset author overrides). - // Skipped objects that don't declare `organization_id` will have - // the extra key silently ignored by the engine. - if (config.organizationId && record['organization_id'] == null) { - record['organization_id'] = config.organizationId; - } - - // Resolve references - for (const ref of objectRefs) { - const fieldValue = record[ref.field]; - if (fieldValue === undefined || fieldValue === null) continue; - - // LOUD FAILURE: a reference must be a natural-key string (or an - // internal id). An object value — e.g. the wrapper `{ externalId: 'X' }` - // — never resolves: it would otherwise fall through unresolved and reach - // the driver as a non-bindable value ("SQLite3 can only bind ..."). This - // used to be silently skipped (and only crashed on a persistent DB's - // update path), so catch it here and report the actionable fix instead. - if (typeof fieldValue === 'object') { - const wrapped = (fieldValue as Record).externalId; - const hint = - wrapped !== undefined - ? ` Pass the natural key directly: ${ref.field}: ${JSON.stringify(wrapped)}.` - : ` Pass the target's ${ref.targetField} value as a plain string.`; - const error: ReferenceResolutionError = { - sourceObject: objectName, - field: ref.field, - targetObject: ref.targetObject, - targetField: ref.targetField, - attemptedValue: fieldValue, - recordIndex: i, - message: - `Invalid reference for ${objectName}.${ref.field}: expected a ` + - `${ref.targetObject}.${ref.targetField} natural-key string but got an object.${hint}`, - }; - errors.push(error); - allErrors.push(error); - this.logger.warn(`[SeedLoader] ${error.message}`, { recordIndex: i }); - // Drop the unresolvable value so it never reaches the driver. - record[ref.field] = null; - continue; - } - - // Skip if value looks like an internal ID (not a natural key) - if (typeof fieldValue !== 'string' || this.looksLikeInternalId(fieldValue)) continue; - - // Try to resolve via already-inserted records - const targetMap = insertedRecords.get(ref.targetObject); - const resolvedId = targetMap?.get(String(fieldValue)); - - if (resolvedId) { - record[ref.field] = resolvedId; - referencesResolved++; - } else if (!config.dryRun) { - // Try to resolve from existing data in the database - const dbId = await this.resolveFromDatabase(ref.targetObject, ref.targetField, fieldValue, config.organizationId); - if (dbId) { - record[ref.field] = dbId; - referencesResolved++; - } else if (config.multiPass) { - // Defer to pass 2 - record[ref.field] = null; - deferredUpdates.push({ - objectName, - recordExternalId: String(record[externalId] ?? ''), - field: ref.field, - targetObject: ref.targetObject, - targetField: ref.targetField, - attemptedValue: fieldValue, - recordIndex: i, - }); - referencesDeferred++; - } else { - // Cannot resolve - record error - const error: ReferenceResolutionError = { - sourceObject: objectName, - field: ref.field, - targetObject: ref.targetObject, - targetField: ref.targetField, - attemptedValue: fieldValue, - recordIndex: i, - message: `Cannot resolve reference: ${objectName}.${ref.field} = '${fieldValue}' → ${ref.targetObject}.${ref.targetField} not found`, - }; - errors.push(error); - allErrors.push(error); - } - } else { - // Dry-run: attempt resolution, report error if not found - const targetMap2 = insertedRecords.get(ref.targetObject); - if (!targetMap2?.has(String(fieldValue))) { - const error: ReferenceResolutionError = { - sourceObject: objectName, - field: ref.field, - targetObject: ref.targetObject, - targetField: ref.targetField, - attemptedValue: fieldValue, - recordIndex: i, - message: `[dry-run] Reference may not resolve: ${objectName}.${ref.field} = '${fieldValue}' → ${ref.targetObject}.${ref.targetField}`, - }; - errors.push(error); - allErrors.push(error); - } - } - } - - // Insert/upsert the record - if (!config.dryRun) { - try { - const result = await this.writeRecord( - objectName, record, mode, externalId, existingRecords - ); - - if (result.action === 'inserted') inserted++; - else if (result.action === 'updated') updated++; - else if (result.action === 'skipped') skipped++; - - // Track the inserted/updated record's ID for reference resolution - const externalIdValue = String(record[externalId] ?? ''); - const internalId = result.id; - if (externalIdValue && internalId) { - insertedRecords.get(objectName)!.set(externalIdValue, String(internalId)); - } - } catch (err: any) { - // LOUD FAILURE: write errors were previously only counted + - // warn-logged, so dropped rows were invisible in result.errors and - // the boot summary. Surface them as actionable errors too, so the - // overall load is marked unsuccessful and the reason is reported. - errored++; - const error: ReferenceResolutionError = { - sourceObject: objectName, - field: '(write)', - targetObject: objectName, - targetField: externalId, - attemptedValue: record[externalId] ?? null, - recordIndex: i, - message: `Failed to write ${objectName} record #${i} (${externalId}=${String(record[externalId] ?? '')}): ${err.message}`, - }; - errors.push(error); - allErrors.push(error); - this.logger.warn(`[SeedLoader] ${error.message}`, { recordIndex: i }); - } - } else { - // Dry-run: simulate insert tracking - const externalIdValue = String(record[externalId] ?? ''); - if (externalIdValue) { - insertedRecords.get(objectName)!.set(externalIdValue, `dry-run-id-${i}`); - } - inserted++; // Count as "would be inserted" - } - } - - return { - object: objectName, - mode, - inserted, - updated, - skipped, - errored, - total: dataset.records.length, - referencesResolved, - referencesDeferred, - errors, - }; - } - - // ========================================================================== - // Internal: Reference Resolution - // ========================================================================== - - private async resolveFromDatabase( - targetObject: string, - targetField: string, - value: unknown, - organizationId?: string, - ): Promise { - try { - const where: Record = { [targetField]: value }; - // Per-tenant replay: when scoping is requested, only consider - // rows that belong to the target tenant so cross-tenant rows - // never get borrowed as a "resolved" reference (would silently - // create a cross-org FK). - if (organizationId) where.organization_id = organizationId; - const records = await this.engine.find(targetObject, { - where, - fields: ['id'], - limit: 1, - context: { isSystem: true }, - } as any); - if (records && records.length > 0) { - return String(records[0].id || records[0]._id); - } - } catch { - // Target object may not exist yet - } - return null; - } - - private async resolveDeferredUpdates( - deferredUpdates: DeferredUpdate[], - insertedRecords: Map>, - allResults: SeedLoadResult[], - allErrors: ReferenceResolutionError[], - organizationId?: string, - ): Promise { - for (const deferred of deferredUpdates) { - // Try to resolve from inserted records - const targetMap = insertedRecords.get(deferred.targetObject); - let resolvedId = targetMap?.get(String(deferred.attemptedValue)); - - // Try database fallback - if (!resolvedId) { - resolvedId = (await this.resolveFromDatabase( - deferred.targetObject, deferred.targetField, deferred.attemptedValue, organizationId - )) ?? undefined; - } - - if (resolvedId) { - // Find the record and update the reference - const objectRecordMap = insertedRecords.get(deferred.objectName); - const recordId = objectRecordMap?.get(deferred.recordExternalId); - - if (recordId) { - try { - await this.engine.update(deferred.objectName, { - id: recordId, - [deferred.field]: resolvedId, - }, { context: { isSystem: true } } as any); - - // Update result stats - const resultEntry = allResults.find(r => r.object === deferred.objectName); - if (resultEntry) { - resultEntry.referencesResolved++; - resultEntry.referencesDeferred--; - } - } catch (err: any) { - this.logger.warn('[SeedLoader] Failed to resolve deferred reference', { - object: deferred.objectName, - field: deferred.field, - error: err.message, - }); - } - } - } else { - // Still unresolved after pass 2 - const error: ReferenceResolutionError = { - sourceObject: deferred.objectName, - field: deferred.field, - targetObject: deferred.targetObject, - targetField: deferred.targetField, - attemptedValue: deferred.attemptedValue, - recordIndex: deferred.recordIndex, - message: `Deferred reference unresolved after pass 2: ${deferred.objectName}.${deferred.field} = '${deferred.attemptedValue}' → ${deferred.targetObject}.${deferred.targetField} not found`, - }; - - const resultEntry = allResults.find(r => r.object === deferred.objectName); - if (resultEntry) { - resultEntry.errors.push(error); - } - allErrors.push(error); - } - } - } - - // ========================================================================== - // Internal: Write Operations - // ========================================================================== - - /** - * Seed writes always run as a privileged system context. This bypasses - * RBAC checks (so seeds can target system tables like `sys_*`) and - * disables the SecurityPlugin's auto-injection of `organization_id` / - * `owner_id` — seeds either declare those fields explicitly per - * record, or are intentionally cross-tenant / global. - */ - private static readonly SEED_OPTIONS = { context: { isSystem: true } } as const; - - private async writeRecord( - objectName: string, - record: Record, - mode: string, - externalId: string, - existingRecords?: Map, - ): Promise<{ action: 'inserted' | 'updated' | 'skipped'; id?: string }> { - const externalIdValue = record[externalId]; - const existing = existingRecords?.get(String(externalIdValue ?? '')); - const opts = SeedLoaderService.SEED_OPTIONS as any; - - switch (mode) { - case 'insert': { - const result = await this.engine.insert(objectName, record, opts); - return { action: 'inserted', id: this.extractId(result) }; - } - - case 'update': { - if (!existing) { - return { action: 'skipped' }; - } - const id = this.extractId(existing); - await this.engine.update(objectName, { ...record, id }, opts); - return { action: 'updated', id }; - } - - case 'upsert': { - if (existing) { - const id = this.extractId(existing); - await this.engine.update(objectName, { ...record, id }, opts); - return { action: 'updated', id }; - } else { - const result = await this.engine.insert(objectName, record, opts); - return { action: 'inserted', id: this.extractId(result) }; - } - } - - case 'ignore': { - if (existing) { - return { action: 'skipped', id: this.extractId(existing) }; - } - const result = await this.engine.insert(objectName, record, opts); - return { action: 'inserted', id: this.extractId(result) }; - } - - case 'replace': { - // Replace mode: just insert (caller should have cleared the table) - const result = await this.engine.insert(objectName, record, opts); - return { action: 'inserted', id: this.extractId(result) }; - } - - default: { - const result = await this.engine.insert(objectName, record, opts); - return { action: 'inserted', id: this.extractId(result) }; - } - } - } - - // ========================================================================== - // Internal: Dependency Graph - // ========================================================================== - - /** - * Kahn's algorithm for topological sort with cycle detection. - */ - private topologicalSort( - nodes: ObjectDependencyNode[], - ): { insertOrder: string[]; circularDependencies: string[][] } { - const inDegree = new Map(); - const adjacency = new Map(); - const objectSet = new Set(nodes.map(n => n.object)); - - // Initialize - for (const node of nodes) { - inDegree.set(node.object, 0); - adjacency.set(node.object, []); - } - - // Build adjacency list and in-degree counts - for (const node of nodes) { - for (const dep of node.dependsOn) { - // Exclude self-references from ordering (e.g., employee.manager_id → employee). - // Self-referencing fields are still tracked in node.references for resolution. - if (objectSet.has(dep) && dep !== node.object) { - adjacency.get(dep)!.push(node.object); - inDegree.set(node.object, (inDegree.get(node.object) || 0) + 1); - } - } - } - - // Kahn's algorithm - const queue: string[] = []; - for (const [obj, degree] of inDegree) { - if (degree === 0) queue.push(obj); - } - - const insertOrder: string[] = []; - while (queue.length > 0) { - const current = queue.shift()!; - insertOrder.push(current); - - for (const neighbor of (adjacency.get(current) || [])) { - const newDegree = (inDegree.get(neighbor) || 0) - 1; - inDegree.set(neighbor, newDegree); - if (newDegree === 0) { - queue.push(neighbor); - } - } - } - - // Detect circular dependencies - const circularDependencies: string[][] = []; - const remaining = nodes.filter(n => !insertOrder.includes(n.object)); - - if (remaining.length > 0) { - // Find cycles using DFS - const cycles = this.findCycles(remaining); - circularDependencies.push(...cycles); - - // Add remaining objects to insertOrder (they'll need multi-pass) - for (const node of remaining) { - if (!insertOrder.includes(node.object)) { - insertOrder.push(node.object); - } - } - } - - return { insertOrder, circularDependencies }; - } - - private findCycles(nodes: ObjectDependencyNode[]): string[][] { - const cycles: string[][] = []; - const nodeMap = new Map(nodes.map(n => [n.object, n])); - const visited = new Set(); - const inStack = new Set(); - - const dfs = (current: string, path: string[]) => { - if (inStack.has(current)) { - // Found a cycle - const cycleStart = path.indexOf(current); - if (cycleStart !== -1) { - cycles.push([...path.slice(cycleStart), current]); - } - return; - } - if (visited.has(current)) return; - - visited.add(current); - inStack.add(current); - path.push(current); - - const node = nodeMap.get(current); - if (node) { - for (const dep of node.dependsOn) { - if (nodeMap.has(dep)) { - dfs(dep, [...path]); - } - } - } - - inStack.delete(current); - }; - - for (const node of nodes) { - if (!visited.has(node.object)) { - dfs(node.object, []); - } - } - - return cycles; - } - - // ========================================================================== - // Internal: Helpers - // ========================================================================== - - private filterByEnv(datasets: Seed[], env?: string): Seed[] { - if (!env) return datasets; - return datasets.filter(d => (d.env as string[]).includes(env)); - } - - private orderDatasets(datasets: Seed[], insertOrder: string[]): Seed[] { - const orderMap = new Map(insertOrder.map((name, i) => [name, i])); - return [...datasets].sort((a, b) => { - const orderA = orderMap.get(a.object) ?? Number.MAX_SAFE_INTEGER; - const orderB = orderMap.get(b.object) ?? Number.MAX_SAFE_INTEGER; - return orderA - orderB; - }); - } - - private buildReferenceMap(graph: ObjectDependencyGraph): Map { - const map = new Map(); - for (const node of graph.nodes) { - if (node.references.length > 0) { - map.set(node.object, node.references); - } - } - return map; - } - - private async loadExistingRecords( - objectName: string, - externalId: string, - organizationId?: string, - ): Promise> { - const map = new Map(); - try { - const findArgs: Record = { - fields: ['id', externalId], - context: { isSystem: true }, - }; - // Per-tenant replay: restrict to the target tenant's own rows - // so upsert key matching never returns another tenant's record - // (would silently steal/overwrite rows across orgs). - if (organizationId) findArgs.where = { organization_id: organizationId }; - const records = await this.engine.find(objectName, findArgs as any); - for (const record of records || []) { - const key = String(record[externalId] ?? ''); - if (key) { - map.set(key, record); - } - } - } catch { - // Object may not have records yet - } - return map; - } - - private looksLikeInternalId(value: string): boolean { - // UUID v4 pattern - if (/^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i.test(value)) { - return true; - } - // MongoDB ObjectId pattern (24 hex chars) - if (/^[0-9a-f]{24}$/i.test(value)) { - return true; - } - return false; - } - - private extractId(record: any): string | undefined { - if (!record) return undefined; - return String(record.id || record._id || ''); - } - - private buildEmptyResult(config: SeedLoaderConfig, durationMs: number): SeedLoaderResult { - return { - success: true, - dryRun: config.dryRun, - dependencyGraph: { nodes: [], insertOrder: [], circularDependencies: [] }, - results: [], - errors: [], - summary: { - objectsProcessed: 0, - totalRecords: 0, - totalInserted: 0, - totalUpdated: 0, - totalSkipped: 0, - totalErrored: 0, - totalReferencesResolved: 0, - totalReferencesDeferred: 0, - circularDependencyCount: 0, - durationMs, - }, - }; - } - - private buildResult( - config: SeedLoaderConfig, - graph: ObjectDependencyGraph, - results: SeedLoadResult[], - errors: ReferenceResolutionError[], - durationMs: number, - ): SeedLoaderResult { - const summary = { - objectsProcessed: results.length, - totalRecords: results.reduce((sum, r) => sum + r.total, 0), - totalInserted: results.reduce((sum, r) => sum + r.inserted, 0), - totalUpdated: results.reduce((sum, r) => sum + r.updated, 0), - totalSkipped: results.reduce((sum, r) => sum + r.skipped, 0), - totalErrored: results.reduce((sum, r) => sum + r.errored, 0), - totalReferencesResolved: results.reduce((sum, r) => sum + r.referencesResolved, 0), - totalReferencesDeferred: results.reduce((sum, r) => sum + r.referencesDeferred, 0), - circularDependencyCount: graph.circularDependencies.length, - durationMs, - }; - - const hasErrors = errors.length > 0 || summary.totalErrored > 0; - - return { - success: !hasErrors, - dryRun: config.dryRun, - dependencyGraph: graph, - results, - errors, - summary, - }; - } -} - -// ========================================================================== -// Internal Types -// ========================================================================== - -interface DeferredUpdate { - objectName: string; - recordExternalId: string; - field: string; - targetObject: string; - targetField: string; - attemptedValue: unknown; - recordIndex: number; -} +// +// MOVED: SeedLoaderService now lives in @objectstack/objectql so the protocol's +// `publishMetaItem` can materialize published `seed` metadata into rows on EVERY +// publish path (per-ref REST publish, package publish-drafts, dispatcher) — +// packages/rest cannot depend on runtime (runtime → rest), and objectql is the +// layer that owns both the engine and the publish primitive. This shim keeps +// the historical runtime import path working. +export { SeedLoaderService } from '@objectstack/objectql';