import { sql, type SQLWrapper } from 'drizzle-orm'
import { resolveRequiredContentCapabilities, type ContentPrincipal } from './authz.js'
import { canonicalContentMigrationHash } from './migrations/content-model.js'
import { normalizeContentTypeDefinition, type ContentTypeCapabilities, type ContentTypeDefinition } from './registry.js'
import { normalizeContentStatusDefinition, type ContentStatusDefinition } from './status.js'
import { assertActiveContentTransaction, type ContentSchemaValue, type ContentTransaction } from './schema.js'
import type { ContentOperationError } from './headless.js'

export interface ContentExportBundle { readonly version:1; readonly generatedAt:string; readonly types:readonly ContentSchemaValue[]; readonly statuses:readonly ContentSchemaValue[]; readonly taxonomies:readonly ContentSchemaValue[]; readonly terms:readonly ContentSchemaValue[]; readonly entries:readonly ContentSchemaValue[]; readonly revisions:readonly ContentSchemaValue[] }
export interface ContentImportLimits { readonly maxBytes:number;readonly maxDefinitions:number;readonly maxEntries:number;readonly maxRevisions:number;readonly batchSize:number }
export type ContentImportConflictPolicy='reject'|'preserve'|'mappedReplacement'
export type ContentImportMappingNamespace='type'|'status'|'taxonomy'|'term'|'entry'|'revision'
export type ContentImportMappingKey=`${ContentImportMappingNamespace}:${string}`
export type ContentImportMappings=Readonly<Partial<Record<ContentImportMappingKey,string>>>
export type ContentImportCapabilityBranch=keyof ContentTypeCapabilities
export interface ContentImportRequiredAuthority {readonly typeKey:string;readonly branches:readonly ContentImportCapabilityBranch[]}
export interface ContentImportScope {readonly tenantKey:string;readonly typeKeys:readonly string[]}
export interface ContentImportPlan { readonly planHash:string;readonly manifestHash:string;readonly destinationCorpusVersion:string;readonly authorizationPolicyVersion:string;readonly scope:ContentImportScope;readonly conflictPolicy:ContentImportConflictPolicy;readonly requiredAuthority:readonly ContentImportRequiredAuthority[];readonly counts:Readonly<Record<string,number>>;readonly mappings:ContentImportMappings;readonly conflicts:readonly ContentOperationError[];readonly estimatedBytes:number }
export interface ContentImportStaging {readonly hostSessionId:string;readonly visibility:'hiddenUntilPublish'}
export type ContentImportSection=ContentImportMappingNamespace
export interface ContentImportSession { readonly id:string;readonly planHash:string;readonly manifestHash:string;readonly state:'validated'|'applying'|'staged'|'published'|'failed'|'rolledBack'|'reversing'|'reversed';readonly staging:ContentImportStaging;readonly scope:ContentImportScope;readonly conflictPolicy:ContentImportConflictPolicy;readonly requiredAuthority:readonly ContentImportRequiredAuthority[];readonly authorizationPolicyVersion:string;readonly destinationCorpusVersion:string;readonly publishedCorpusVersion?:string;readonly mappings:ContentImportMappings;readonly appliedBatches:number;readonly journalVersion:string;readonly nextSection:ContentImportSection|null;readonly nextOffset:number }
export interface ContentImportAuthorizationInput {readonly scope:ContentImportScope;readonly requiredAuthority:readonly ContentImportRequiredAuthority[];readonly expectedPolicyVersion?:string}
export interface ContentImportAuthorization {assert(tx:ContentTransaction,principal:ContentPrincipal,action:'plan'|'start'|ContentImportCommandKind,input:ContentImportAuthorizationInput):Promise<{policyVersion:string}>}
export type ContentImportCommandKind='resume'|'publish'|'rollback'|'reversePublished'
export type ContentImportCommandClaim={kind:'new';session?:ContentImportSession}|{kind:'replay';session:ContentImportSession}|{kind:'conflict'}
export interface ContentImportCommand { readonly operationId:string;readonly expectedPlanHash:string;readonly expectedManifestHash:string;readonly scope:ContentImportScope;readonly requiredAuthority:readonly ContentImportRequiredAuthority[];readonly expectedAuthorizationPolicyVersion:string;readonly expectedJournalVersion:string;readonly expectedSection:ContentImportSession['nextSection'];readonly expectedOffset:number }
export interface ContentImportJournal {
  claimStart(tx:ContentTransaction,principal:ContentPrincipal,authorizationInput:ContentImportAuthorizationInput,authorization:ContentImportAuthorization,operationId:string,canonicalStartHash:string):Promise<ContentImportCommandClaim>
  create(plan:ContentImportPlan,staging:ContentImportStaging,principal:ContentPrincipal,operationId:string,tx:ContentTransaction):Promise<ContentImportSession>
  claimCommand(tx:ContentTransaction,sessionId:string,principal:ContentPrincipal,kind:ContentImportCommandKind,command:ContentImportCommand,canonicalCommandHash:string,authorization:ContentImportAuthorization):Promise<ContentImportCommandClaim>
  get(tx:ContentTransaction,sessionId:string,principal:ContentPrincipal):Promise<ContentImportSession>
  recordBatch(sessionId:string,batch:number,beforeImage:ContentSchemaValue,tx:ContentTransaction):Promise<void>
  readBatchesReverse(tx:ContentTransaction,sessionId:string,expectedPlanHash:string,expectedManifestHash:string):Promise<readonly {batch:number;beforeImage:ContentSchemaValue}[]>
  transition(sessionId:string,from:ContentImportSession['state'],to:ContentImportSession['state'],tx:ContentTransaction):Promise<void>
  completeCommand(tx:ContentTransaction,sessionId:string,kind:ContentImportCommandKind,command:ContentImportCommand,canonicalCommandHash:string,result:ContentImportSession):Promise<void>
}
export interface ContentImportLifecycleEvent {readonly id:string;readonly version:1;readonly kind:'published'|'rolledBack'|'reversed';readonly operationId:string;readonly sessionId:string;readonly planHash:string;readonly manifestHash:string;readonly principalId:string;readonly scope:ContentImportScope;readonly authorizationPolicyVersion:string;readonly affectedCounts:Readonly<Record<string,number>>}
export interface ContentImportLifecycleEffects {audit:{record(tx:ContentTransaction,event:ContentImportLifecycleEvent):Promise<void>};outbox:{enqueue(tx:ContentTransaction,event:ContentImportLifecycleEvent):Promise<void>}}
export interface ContentPublishedReversePlan {readonly reversePlanHash:string;readonly sessionId:string;readonly originalPlanHash:string;readonly manifestHash:string;readonly publishedCorpusVersion:string;readonly beforeImageHash:string;readonly affectedCounts:Readonly<Record<string,number>>;readonly scope:ContentImportScope;readonly requiredAuthority:readonly ContentImportRequiredAuthority[];readonly authorizationPolicyVersion:string}
export interface ContentPublishedReverseCommand extends ContentImportCommand {readonly expectedReversePlanHash:string;readonly expectedPublishedCorpusVersion:string}
export type ContentImportErrorCode='validation'|'authorization'|'conflict'|'stale-plan'|'stale-journal'|'integrity'|'limit'
export class ContentImportError extends Error {override readonly name='ContentImportError';constructor(readonly code:ContentImportErrorCode,readonly detail:string,readonly field?:string){super(`content import ${code}: ${detail}${field?` (${field})`:''}`)}}
export function isContentImportError(error:unknown):error is ContentImportError{return error instanceof ContentImportError}
const fail=(code:ContentImportErrorCode,detail:string,field?:string):never=>{throw new ContentImportError(code,detail,field)}

const ALL_BRANCHES=Object.freeze(['manageType','migrateAll','manageTerms','create','read','readPrivate','readProtected','editOwn','editOthers','editPrivate','editPublished','publish','deleteOwn','deleteOthers','deletePrivate','deletePublished'] satisfies readonly ContentImportCapabilityBranch[])
const UUID_RE=/^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i
const KEY_RE=/^[a-z][a-z0-9_-]*$/
const FORBIDDEN_IMPORT_KEY=/password|credential/i
function boundedString(value:unknown,field:string,max=4096):string{if(typeof value!=='string')fail('validation','must be a string',field);const normalized=(value as string).trim();if(!normalized||normalized.length>max)fail('validation',`must contain 1..${max} characters`,field);return normalized}
function plain(value:unknown,field:string):Record<string,unknown>{if(typeof value!=='object'||value===null||Array.isArray(value)||Object.getPrototypeOf(value)!==Object.prototype)fail('validation','must be a plain object',field);return value as Record<string,unknown>}
function array(value:unknown,field:string,max:number):unknown[]{if(!Array.isArray(value)||value.length>max)fail('limit',`must be an array with at most ${max} items`,field);return value as unknown[]}
function positive(value:unknown,field:string):number{if(typeof value!=='number'||!Number.isSafeInteger(value)||value<1)fail('validation','must be a positive safe integer',field);return value as number}
function nonnegative(value:unknown,field:string):number{if(typeof value!=='number'||!Number.isSafeInteger(value)||value<0)fail('validation','must be a non-negative safe integer',field);return value as number}
function iso(value:unknown,field:string):string{const raw=boundedString(value,field,128),date=new Date(raw);if(Number.isNaN(date.getTime()))fail('validation','must be an ISO timestamp',field);return date.toISOString()}
function uuid(value:unknown,field:string):string{const raw=boundedString(value,field,128);if(!UUID_RE.test(raw))fail('validation','must be a UUID',field);return raw}
function key(value:unknown,field:string):string{const raw=boundedString(value,field,128);if(!KEY_RE.test(raw))fail('validation','must be a normalized key',field);return raw}
function schemaValue(value:unknown,field='value',depth=0):ContentSchemaValue{if(depth>32)fail('limit','exceeds maximum nested depth',field);if(value===null||typeof value==='string'||typeof value==='boolean')return value;if(value instanceof Date)return value.toISOString();if(typeof value==='bigint')return value.toString();if(typeof value==='number'){if(!Number.isFinite(value))fail('validation','contains non-finite number',field);return value}if(Array.isArray(value))return Object.freeze(value.map((item,index)=>schemaValue(item,`${field}.${index}`,depth+1)));if(typeof value==='object'){if(Object.getPrototypeOf(value)!==Object.prototype)fail('validation','contains a non-plain object',field);const out:Record<string,ContentSchemaValue>={};for(const [name,child] of Object.entries(value as Record<string,unknown>)){if(FORBIDDEN_IMPORT_KEY.test(name))fail('validation','password/credential material is forbidden',`${field}.${name}`);if(child!==undefined)out[name]=schemaValue(child,`${field}.${name}`,depth+1)}return Object.freeze(out)}return fail('validation','contains unsupported value',field)}
function stableScope(scope:ContentImportScope):ContentImportScope{const tenantKey=boundedString(scope.tenantKey,'scope.tenantKey',256),typeKeys=[...new Set(scope.typeKeys.map((item,index)=>key(item,`scope.typeKeys.${index}`)))].sort();if(typeKeys.length===0||typeKeys.length>500)fail('validation','scope must contain 1..500 type keys','scope.typeKeys');return Object.freeze({tenantKey,typeKeys:Object.freeze(typeKeys)})}
function sameScope(a:ContentImportScope,b:ContentImportScope):boolean{return a.tenantKey===b.tenantKey&&a.typeKeys.length===b.typeKeys.length&&a.typeKeys.every((value,index)=>value===b.typeKeys[index])}
function stableAuthority(input:readonly ContentImportRequiredAuthority[]):readonly ContentImportRequiredAuthority[]{const map=new Map<string,Set<ContentImportCapabilityBranch>>();for(const row of input){const typeKey=key(row.typeKey,'requiredAuthority.typeKey'),set=map.get(typeKey)??new Set<ContentImportCapabilityBranch>();for(const branch of row.branches){if(!ALL_BRANCHES.includes(branch))fail('validation','unknown capability branch','requiredAuthority.branches');set.add(branch)}map.set(typeKey,set)}return Object.freeze([...map].sort(([a],[b])=>a.localeCompare(b)).map(([typeKey,branches])=>Object.freeze({typeKey,branches:Object.freeze(ALL_BRANCHES.filter((branch)=>branches.has(branch)))})))}
function sameAuthority(a:readonly ContentImportRequiredAuthority[],b:readonly ContentImportRequiredAuthority[]):boolean{return JSON.stringify(stableAuthority(a))===JSON.stringify(stableAuthority(b))}
function principalScope(principal:ContentPrincipal):string{return `${principal.tenantId??''}:${principal.id}`}
function isUniqueViolation(error:unknown):boolean{let current:unknown=error;while(current){const code=(current as {code?:unknown}).code,message=current instanceof Error?current.message:String(current);if(code==='23505'||/unique constraint|duplicate key|content_lifecycle_journal_operation_uq/i.test(message))return true;current=current instanceof Error?(current as Error&{cause?:unknown}).cause:undefined}return false}
function rowJson(value:unknown):Record<string,unknown>{const parsed=typeof value==='string'?JSON.parse(value):value;if(typeof parsed!=='object'||parsed===null||Array.isArray(parsed))fail('integrity','stored journal payload is malformed');return parsed as Record<string,unknown>}
function rowsOf<Row extends Record<string,unknown>>(result:unknown):readonly Row[]{if(Array.isArray(result))return result as Row[];const rows=(result as {rows?:unknown})?.rows;if(Array.isArray(rows))return rows as Row[];return fail('integrity','database adapter returned invalid rows')}
async function execRows<Row extends Record<string,unknown>>(tx:{execute(query:SQLWrapper):Promise<unknown>},query:SQLWrapper):Promise<readonly Row[]>{return rowsOf<Row>(await tx.execute(query))}

interface DefinitionVersionItem<T>{revision:number;canonicalHash:string;definition:T;origin:'code'|'db'|'import';createdAt:string}
interface TypeItem {key:string;origin:'code'|'db'|'import';version:number;revision:number;active:boolean;canonicalHash:string;definition:ContentTypeDefinition;versions:readonly DefinitionVersionItem<ContentTypeDefinition>[]}
interface StatusItem {key:string;origin:'code'|'db'|'import';version:number;revision:number;active:boolean;canonicalHash:string;definition:ContentStatusDefinition;versions:readonly DefinitionVersionItem<ContentStatusDefinition>[]}
interface TaxonomyItem {key:string;version:string}
interface TermItem {id:string;taxonomy:string;slug:string;name:string;parentId:string|null;depth:number;createdAt:string}
interface EntryItem {id:string;slug:string;type:string;title:string;body:string;status:string;visibility:'public'|'private'|'members';publishedAt:string|null;author:string;createdAt:string;updatedAt:string;parentId:string|null;menuOrder:number;templateKey:string|null;excerpt:string;featuredMedia:ContentSchemaValue|null;commentStatus:'open'|'closed';pingStatus:'open'|'closed';sticky:boolean;format:string|null;deletedAt:string|null;lastEditedBy:string;typeDefinitionRevision:number;statusDefinitionRevision:number;termIds:readonly string[]}
interface RevisionItem {id:string;entryId:string;seq:number;title:string;body:string;slug:string;type:string;termIds:readonly string[];snapshot?:ContentSchemaValue;editor:string;createdAt:string}
interface ParsedBundle {bundle:ContentExportBundle;types:readonly TypeItem[];statuses:readonly StatusItem[];taxonomies:readonly TaxonomyItem[];terms:readonly TermItem[];entries:readonly EntryItem[];revisions:readonly RevisionItem[];manifestHash:string}
function nullableString(value:unknown,field:string,max=4096):string|null{return value===null?null:boundedString(value,field,max)}
function parseDefinitionVersions<T>(value:unknown,field:string,keyValue:string,normalize:(definition:T)=>T):readonly DefinitionVersionItem<T>[] {const versions=array(value,field,1000).map((raw,index)=>{const item=plain(raw,`${field}.${index}`),definition=normalize(item.definition as T),origin=item.origin;if(origin!=='code'&&origin!=='db'&&origin!=='import')fail('validation','invalid origin',`${field}.${index}.origin`);if((definition as {key?:unknown}).key!==keyValue)fail('validation','definition version key mismatch',`${field}.${index}.definition.key`);return Object.freeze({revision:positive(item.revision,`${field}.${index}.revision`),canonicalHash:boundedString(item.canonicalHash,`${field}.${index}.canonicalHash`,128),definition,origin:origin as 'code'|'db'|'import',createdAt:iso(item.createdAt,`${field}.${index}.createdAt`)})});duplicates(versions,(item)=>String(item.revision),field);return Object.freeze(versions.sort((a,b)=>a.revision-b.revision))}
function parseType(raw:unknown,index:number):TypeItem{const item=plain(raw,`types.${index}`),definition=normalizeContentTypeDefinition(item.definition as ContentTypeDefinition),itemKey=key(item.key,`types.${index}.key`);if(definition.key!==itemKey)fail('validation','definition key mismatch',`types.${index}.definition.key`);const origin=item.origin;if(origin!=='code'&&origin!=='db'&&origin!=='import')fail('validation','invalid origin',`types.${index}.origin`);const versions=parseDefinitionVersions<ContentTypeDefinition>(item.versions,`types.${index}.versions`,itemKey,normalizeContentTypeDefinition);const active=Boolean(item.active);if(active!==definition.active)fail('validation','type active state differs from normalized definition',`types.${index}.active`);return Object.freeze({key:itemKey,origin:origin as TypeItem['origin'],version:positive(item.version,`types.${index}.version`),revision:positive(item.revision,`types.${index}.revision`),active,canonicalHash:boundedString(item.canonicalHash,`types.${index}.canonicalHash`,128),definition,versions})}
function parseStatus(raw:unknown,index:number):StatusItem{const item=plain(raw,`statuses.${index}`),definition=normalizeContentStatusDefinition(item.definition as ContentStatusDefinition),itemKey=key(item.key,`statuses.${index}.key`);if(definition.key!==itemKey)fail('validation','definition key mismatch',`statuses.${index}.definition.key`);const origin=item.origin;if(origin!=='code'&&origin!=='db'&&origin!=='import')fail('validation','invalid origin',`statuses.${index}.origin`);const versions=parseDefinitionVersions<ContentStatusDefinition>(item.versions,`statuses.${index}.versions`,itemKey,normalizeContentStatusDefinition);if(typeof item.active!=='boolean')fail('validation','status active must be boolean',`statuses.${index}.active`);return Object.freeze({key:itemKey,origin:origin as StatusItem['origin'],version:positive(item.version,`statuses.${index}.version`),revision:positive(item.revision,`statuses.${index}.revision`),active:item.active as boolean,canonicalHash:boundedString(item.canonicalHash,`statuses.${index}.canonicalHash`,128),definition,versions})}
function parseTaxonomy(raw:unknown,index:number):TaxonomyItem{const item=plain(raw,`taxonomies.${index}`);return Object.freeze({key:key(item.key,`taxonomies.${index}.key`),version:boundedString(item.version,`taxonomies.${index}.version`,128)})}
function parseTerm(raw:unknown,index:number):TermItem{const item=plain(raw,`terms.${index}`);return Object.freeze({id:uuid(item.id,`terms.${index}.id`),taxonomy:key(item.taxonomy,`terms.${index}.taxonomy`),slug:boundedString(item.slug,`terms.${index}.slug`,256),name:boundedString(item.name,`terms.${index}.name`,512),parentId:item.parentId===null?null:uuid(item.parentId,`terms.${index}.parentId`),depth:nonnegative(item.depth,`terms.${index}.depth`),createdAt:iso(item.createdAt,`terms.${index}.createdAt`)})}
function parseEntry(raw:unknown,index:number):EntryItem{const item=plain(raw,`entries.${index}`),visibility=item.visibility,commentStatus=item.commentStatus,pingStatus=item.pingStatus;if(visibility!=='public'&&visibility!=='private'&&visibility!=='members')fail('validation','invalid visibility',`entries.${index}.visibility`);if(commentStatus!=='open'&&commentStatus!=='closed')fail('validation','invalid comment status',`entries.${index}.commentStatus`);if(pingStatus!=='open'&&pingStatus!=='closed')fail('validation','invalid ping status',`entries.${index}.pingStatus`);if(typeof item.sticky!=='boolean')fail('validation','must be boolean',`entries.${index}.sticky`);const termIds=array(item.termIds,`entries.${index}.termIds`,10000).map((value,i)=>uuid(value,`entries.${index}.termIds.${i}`));return Object.freeze({id:uuid(item.id,`entries.${index}.id`),slug:boundedString(item.slug,`entries.${index}.slug`,256),type:key(item.type,`entries.${index}.type`),title:typeof item.title==='string'&&item.title.length<=300?item.title:fail('validation','invalid title',`entries.${index}.title`),body:typeof item.body==='string'&&item.body.length<=2_000_000?item.body:fail('validation','invalid body',`entries.${index}.body`),status:key(item.status,`entries.${index}.status`),visibility:visibility as EntryItem['visibility'],publishedAt:item.publishedAt===null?null:iso(item.publishedAt,`entries.${index}.publishedAt`),author:boundedString(item.author,`entries.${index}.author`,512),createdAt:iso(item.createdAt,`entries.${index}.createdAt`),updatedAt:iso(item.updatedAt,`entries.${index}.updatedAt`),parentId:item.parentId===null?null:uuid(item.parentId,`entries.${index}.parentId`),menuOrder:nonnegative(item.menuOrder,`entries.${index}.menuOrder`),templateKey:nullableString(item.templateKey,`entries.${index}.templateKey`,256),excerpt:typeof item.excerpt==='string'&&item.excerpt.length<=100_000?item.excerpt:fail('validation','invalid excerpt',`entries.${index}.excerpt`),featuredMedia:item.featuredMedia===null?null:schemaValue(item.featuredMedia,`entries.${index}.featuredMedia`),commentStatus:commentStatus as EntryItem['commentStatus'],pingStatus:pingStatus as EntryItem['pingStatus'],sticky:item.sticky as boolean,format:nullableString(item.format,`entries.${index}.format`,128),deletedAt:item.deletedAt===null?null:iso(item.deletedAt,`entries.${index}.deletedAt`),lastEditedBy:boundedString(item.lastEditedBy,`entries.${index}.lastEditedBy`,512),typeDefinitionRevision:positive(item.typeDefinitionRevision,`entries.${index}.typeDefinitionRevision`),statusDefinitionRevision:positive(item.statusDefinitionRevision,`entries.${index}.statusDefinitionRevision`),termIds:Object.freeze(termIds)})}
function parseRevision(raw:unknown,index:number):RevisionItem{const item=plain(raw,`revisions.${index}`),termIds=array(item.termIds,`revisions.${index}.termIds`,10000).map((value,i)=>uuid(value,`revisions.${index}.termIds.${i}`));return Object.freeze({id:uuid(item.id,`revisions.${index}.id`),entryId:uuid(item.entryId,`revisions.${index}.entryId`),seq:positive(item.seq,`revisions.${index}.seq`),title:typeof item.title==='string'?item.title:fail('validation','invalid title',`revisions.${index}.title`),body:typeof item.body==='string'?item.body:fail('validation','invalid body',`revisions.${index}.body`),slug:boundedString(item.slug,`revisions.${index}.slug`,256),type:key(item.type,`revisions.${index}.type`),termIds:Object.freeze(termIds),...(item.snapshot===undefined?{}:{snapshot:schemaValue(item.snapshot,`revisions.${index}.snapshot`)}),editor:boundedString(item.editor,`revisions.${index}.editor`,512),createdAt:iso(item.createdAt,`revisions.${index}.createdAt`)})}
function duplicates<T>(items:readonly T[],id:(item:T)=>string,field:string):void{const seen=new Set<string>();for(const item of items){const value=id(item);if(seen.has(value))fail('validation','contains duplicate identity',field);seen.add(value)}}
function assertTermGraph(terms:readonly TermItem[]):void{const byId=new Map(terms.map((item)=>[item.id,item]));for(const item of terms){if(item.parentId===null){if(item.depth!==0)fail('validation','root term depth must be zero',`term:${item.id}`);continue}const parent=byId.get(item.parentId);if(!parent)throw new ContentImportError('validation','term references missing parent',`term:${item.id}`);if(parent.taxonomy!==item.taxonomy)fail('validation','term parent must share taxonomy',`term:${item.id}`);if(item.depth!==parent.depth+1)fail('validation','term depth does not match parent',`term:${item.id}`);const seen=new Set([item.id]);let cursor:TermItem|undefined=parent;while(cursor){if(seen.has(cursor.id))fail('validation','term hierarchy contains a cycle',`term:${item.id}`);seen.add(cursor.id);cursor=cursor.parentId===null?undefined:byId.get(cursor.parentId)}}}
function assertEntryGraph(entries:readonly EntryItem[]):void{const byId=new Map(entries.map((item)=>[item.id,item]));for(const item of entries){if(item.parentId===null)continue;const parent=byId.get(item.parentId);if(!parent)throw new ContentImportError('validation','entry references missing parent',`entry:${item.id}`);if(parent.type!==item.type)fail('validation','entry parent must share content type',`entry:${item.id}`);const seen=new Set([item.id]);let cursor:EntryItem|undefined=parent;while(cursor){if(seen.has(cursor.id))fail('validation','entry hierarchy contains a cycle',`entry:${item.id}`);seen.add(cursor.id);cursor=cursor.parentId===null?undefined:byId.get(cursor.parentId)}}}
function assertLimits(limits:ContentImportLimits):void{for(const [name,value] of Object.entries(limits)){if(!Number.isSafeInteger(value)||value<1)fail('validation','must be a positive safe integer',`limits.${name}`)}if(limits.batchSize>1000)fail('validation','batchSize must be <= 1000','limits.batchSize')}
async function parseBundle(input:Uint8Array,limits:ContentImportLimits):Promise<ParsedBundle>{
  assertLimits(limits);if(input.byteLength>limits.maxBytes)fail('limit','input exceeds maxBytes','input')
  let decoded:unknown;try{decoded=JSON.parse(new TextDecoder('utf-8',{fatal:true}).decode(input))}catch{fail('validation','input is not valid UTF-8 JSON','input')}
  const root=plain(decoded,'bundle');schemaValue(root,'bundle');for(const name of Object.keys(root))if(!['version','generatedAt','types','statuses','taxonomies','terms','entries','revisions'].includes(name))fail('validation','unsupported bundle property',name)
  if(root.version!==1)fail('validation','unsupported bundle version','version');const generatedAt=iso(root.generatedAt,'generatedAt')
  const types=array(root.types,'types',limits.maxDefinitions).map(parseType),statuses=array(root.statuses,'statuses',limits.maxDefinitions).map(parseStatus),taxonomies=array(root.taxonomies,'taxonomies',limits.maxDefinitions).map(parseTaxonomy)
  if(types.length+statuses.length+taxonomies.length>limits.maxDefinitions)fail('limit','definitions exceed maxDefinitions','definitions')
  const terms=array(root.terms,'terms',Math.max(limits.maxEntries*10,limits.maxDefinitions)).map(parseTerm),entries=array(root.entries,'entries',limits.maxEntries).map(parseEntry),revisions=array(root.revisions,'revisions',limits.maxRevisions).map(parseRevision)
  for(const item of types){if(await canonicalContentMigrationHash(item.definition)!==item.canonicalHash)fail('validation','type canonicalHash does not match normalized definition',`type:${item.key}`);for(const version of item.versions)if(await canonicalContentMigrationHash(version.definition)!==version.canonicalHash)fail('validation','type version hash does not match normalized definition',`type:${item.key}@${version.revision}`);const current=item.versions.find((version)=>version.revision===item.revision);if(!current||current.canonicalHash!==item.canonicalHash)fail('validation','current type revision is absent or mismatched',`type:${item.key}`)}
  for(const item of statuses){if(await canonicalContentMigrationHash(item.definition)!==item.canonicalHash)fail('validation','status canonicalHash does not match normalized definition',`status:${item.key}`);for(const version of item.versions)if(await canonicalContentMigrationHash(version.definition)!==version.canonicalHash)fail('validation','status version hash does not match normalized definition',`status:${item.key}@${version.revision}`);const current=item.versions.find((version)=>version.revision===item.revision);if(!current||current.canonicalHash!==item.canonicalHash)fail('validation','current status revision is absent or mismatched',`status:${item.key}`)}
  duplicates(types,(item)=>item.key,'types');duplicates(statuses,(item)=>item.key,'statuses');duplicates(taxonomies,(item)=>item.key,'taxonomies');duplicates(terms,(item)=>item.id,'terms');duplicates(entries,(item)=>item.id,'entries');duplicates(revisions,(item)=>item.id,'revisions');assertTermGraph(terms);assertEntryGraph(entries)
  const taxonomyKeys=new Set(taxonomies.map((item)=>item.key)),termIds=new Set(terms.map((item)=>item.id)),termById=new Map(terms.map((item)=>[item.id,item])),entryIds=new Set(entries.map((item)=>item.id)),typeKeys=new Set(types.map((item)=>item.key)),statusKeys=new Set(statuses.map((item)=>item.key))
  for(const type of types){for(const definition of [type.definition,...type.versions.map((version)=>version.definition)]){for(const statusKey of definition.statusKeys??[])if(!statusKeys.has(statusKey))fail('validation','type definition references missing status',`type:${type.key}`);for(const taxonomy of definition.taxonomies)if(!taxonomyKeys.has(taxonomy))fail('validation','type definition references missing taxonomy',`type:${type.key}`)}}
  for(const term of terms){if(!taxonomyKeys.has(term.taxonomy))fail('validation','term references missing taxonomy',`term:${term.id}`);if(term.parentId!==null&&!termIds.has(term.parentId))fail('validation','term references missing parent',`term:${term.id}`)}
  for(const entry of entries){const type=types.find((item)=>item.key===entry.type),status=statuses.find((item)=>item.key===entry.status);if(!type||!typeKeys.has(entry.type))throw new ContentImportError('validation','entry references missing type',`entry:${entry.id}`);if(!status||!statusKeys.has(entry.status))throw new ContentImportError('validation','entry references missing status',`entry:${entry.id}`);const typeVersion=type.versions.find((version)=>version.revision===entry.typeDefinitionRevision),statusVersion=status.versions.find((version)=>version.revision===entry.statusDefinitionRevision);if(!typeVersion)throw new ContentImportError('validation','entry type revision is absent from bundle',`entry:${entry.id}`);if(!statusVersion)throw new ContentImportError('validation','entry status revision is absent from bundle',`entry:${entry.id}`);if(typeVersion.definition.statusKeys!==undefined&&!typeVersion.definition.statusKeys.includes(entry.status))fail('validation','entry status is not allowed by bundled type',`entry:${entry.id}`);if(entry.parentId!==null&&!entryIds.has(entry.parentId))fail('validation','entry references missing parent',`entry:${entry.id}`);for(const id of entry.termIds){const term=termById.get(id);if(!term)fail('validation','entry references missing term',`entry:${entry.id}`);if(!typeVersion.definition.taxonomies.includes(term!.taxonomy))fail('validation','entry term taxonomy is not allowed by pinned type definition',`entry:${entry.id}`)}}
  for(const revision of revisions){const entry=entries.find((item)=>item.id===revision.entryId);if(!entry||!entryIds.has(revision.entryId))throw new ContentImportError('validation','revision references missing entry',`revision:${revision.id}`);if(revision.type!==entry.type)fail('validation','revision type does not match entry',`revision:${revision.id}`);for(const id of revision.termIds)if(!termIds.has(id))fail('validation','revision references missing term',`revision:${revision.id}`)}
  const bundle:ContentExportBundle=Object.freeze({version:1,generatedAt,types:Object.freeze(types.map((item)=>schemaValue(item))),statuses:Object.freeze(statuses.map((item)=>schemaValue(item))),taxonomies:Object.freeze(taxonomies.map((item)=>schemaValue(item))),terms:Object.freeze(terms.map((item)=>schemaValue(item))),entries:Object.freeze(entries.map((item)=>schemaValue(item))),revisions:Object.freeze(revisions.map((item)=>schemaValue(item)))})
  return Object.freeze({bundle,types:Object.freeze(types),statuses:Object.freeze(statuses),taxonomies:Object.freeze(taxonomies),terms:Object.freeze(terms),entries:Object.freeze(entries),revisions:Object.freeze(revisions),manifestHash:await canonicalContentMigrationHash(bundle)})
}
export async function parseContentImport(input:Uint8Array,limits:ContentImportLimits):Promise<ContentExportBundle>{return (await parseBundle(input,limits)).bundle}

interface StoredJournalPayload {session:ContentImportSession;batches:readonly {batch:number;beforeImage:ContentSchemaValue}[]}
interface JournalRow extends Record<string,unknown>{importId:string;hostSessionId:string;manifestHash:string;planHash:string;state:ContentImportSession['state'];version:number|string;nextSection:string|null;nextOffset:number|string|null;highWaterMark:string|null;payload:unknown}
interface Receipt {kind:'start'|ContentImportCommandKind;commandHash:string;principalScope:string;session:ContentImportSession}
function sessionFromRow(row:JournalRow):ContentImportSession{const payload=rowJson(row.payload),stored=plain(payload.session,'journal.session') as unknown as ContentImportSession,staging=plain(stored.staging,'journal.session.staging') as unknown as ContentImportStaging;return Object.freeze({...stored,id:String(row.importId),planHash:String(row.planHash),manifestHash:String(row.manifestHash),state:row.state,journalVersion:String(row.version),nextSection:row.nextSection as ContentImportSection|null,nextOffset:Number(row.nextOffset??0),...(row.highWaterMark===null?{}:{publishedCorpusVersion:String(row.highWaterMark)}),staging:Object.freeze({...staging,hostSessionId:String(row.hostSessionId)})})}
async function readJournalRow(tx:ContentTransaction,sessionId:string):Promise<JournalRow|undefined>{return (await execRows<JournalRow>(tx,sql`SELECT import_id AS "importId",host_session_id AS "hostSessionId",manifest_hash AS "manifestHash",plan_hash AS "planHash",state,version,next_section AS "nextSection",next_offset AS "nextOffset",high_water_mark AS "highWaterMark",payload FROM content_import_journal WHERE import_id=${sessionId}`))[0]}
async function readReceipt(tx:ContentTransaction,operationId:string):Promise<Receipt|undefined>{const row=(await execRows<{payload:unknown}>(tx,sql`SELECT payload FROM content_lifecycle_journal WHERE operation_id=${operationId} AND kind LIKE 'content:import:%' AND state='committed' LIMIT 1`))[0];if(!row)return;const value=rowJson(row.payload);if(typeof value.kind!=='string'||typeof value.commandHash!=='string'||typeof value.principalScope!=='string'||typeof value.session!=='object')return undefined;return value as unknown as Receipt}
async function insertPendingReceipt(tx:ContentTransaction,operationId:string,kind:'start'|ContentImportCommandKind,hash:string,scope:string):Promise<void>{const now=new Date().toISOString(),payload=schemaValue({kind,commandHash:hash,principalScope:scope});await tx.execute(sql`INSERT INTO content_lifecycle_journal (id,operation_id,kind,state,payload,created_at,updated_at) VALUES (${crypto.randomUUID()},${operationId},${`content:import:${kind}`},'pending',${JSON.stringify(payload)},${now},${now})`)}
async function commitReceipt(tx:ContentTransaction,operationId:string,kind:'start'|ContentImportCommandKind,hash:string,scope:string,session:ContentImportSession):Promise<void>{const payload=schemaValue({kind,commandHash:hash,principalScope:scope,session}),rows=await execRows<{id:string}>(tx,sql`UPDATE content_lifecycle_journal SET state='committed',payload=${JSON.stringify(payload)},updated_at=${new Date().toISOString()} WHERE operation_id=${operationId} AND state='pending' RETURNING id`);if(rows.length!==1)fail('integrity','import command claim disappeared')}

export function createDbContentImportJournal():ContentImportJournal{const journal:ContentImportJournal={
  async claimStart(tx,principal,authorizationInput,authorization,operationId,canonicalStartHash){
    assertActiveContentTransaction(tx);boundedString(operationId,'operationId',256);const auth=await authorization.assert(tx,principal,'start',authorizationInput);if(authorizationInput.expectedPolicyVersion!==undefined&&auth.policyVersion!==authorizationInput.expectedPolicyVersion)fail('authorization','import scope is unavailable')
    const receipt=await readReceipt(tx,operationId);if(receipt){if(receipt.kind!=='start'||receipt.commandHash!==canonicalStartHash||receipt.principalScope!==principalScope(principal))return {kind:'conflict'};return {kind:'replay',session:receipt.session}}
    try{await insertPendingReceipt(tx,operationId,'start',canonicalStartHash,principalScope(principal));return {kind:'new'}}catch(error){if(isUniqueViolation(error))return {kind:'conflict'};throw error}
  },
  async create(plan,staging,principal,operationId,tx){
    assertActiveContentTransaction(tx);if(staging.visibility!=='hiddenUntilPublish')fail('validation','unsupported staging visibility','staging.visibility');const id=uuid(staging.hostSessionId,'staging.hostSessionId'),session:ContentImportSession=Object.freeze({id,planHash:plan.planHash,manifestHash:plan.manifestHash,state:'validated',staging:Object.freeze(staging),scope:plan.scope,conflictPolicy:plan.conflictPolicy,requiredAuthority:plan.requiredAuthority,authorizationPolicyVersion:plan.authorizationPolicyVersion,destinationCorpusVersion:plan.destinationCorpusVersion,mappings:plan.mappings,appliedBatches:0,journalVersion:'0',nextSection:'type',nextOffset:0}),payload:StoredJournalPayload={session,batches:Object.freeze([])},now=new Date().toISOString()
    try{await tx.execute(sql`INSERT INTO content_import_journal (import_id,host_session_id,manifest_hash,plan_hash,state,version,next_section,next_offset,payload,created_at,updated_at) VALUES (${id},${staging.hostSessionId},${plan.manifestHash},${plan.planHash},'validated',0,'type',0,${JSON.stringify(schemaValue(payload))},${now},${now})`)}catch(error){if(isUniqueViolation(error))fail('conflict','staging session already exists');throw error}
    const pending=(await execRows<{payload:unknown}>(tx,sql`SELECT payload FROM content_lifecycle_journal WHERE operation_id=${operationId} AND state='pending'`))[0];if(!pending)fail('integrity','start claim disappeared');const raw=rowJson(pending!.payload),hash=boundedString(raw.commandHash,'receipt.commandHash',128),scope=boundedString(raw.principalScope,'receipt.principalScope',1024);await commitReceipt(tx,operationId,'start',hash,scope,session);return session
  },
  async claimCommand(tx,sessionId,principal,kind,command,canonicalCommandHash,authorization){
    assertActiveContentTransaction(tx);boundedString(sessionId,'sessionId',128);boundedString(command.operationId,'operationId',256)
    const authInput:ContentImportAuthorizationInput={scope:stableScope(command.scope),requiredAuthority:stableAuthority(command.requiredAuthority),expectedPolicyVersion:command.expectedAuthorizationPolicyVersion};const auth=await authorization.assert(tx,principal,kind,authInput);if(auth.policyVersion!==command.expectedAuthorizationPolicyVersion)fail('authorization','import scope is unavailable')
    const receipt=await readReceipt(tx,command.operationId);if(receipt){if(receipt.kind!==kind||receipt.commandHash!==canonicalCommandHash||receipt.principalScope!==principalScope(principal))return {kind:'conflict'};return {kind:'replay',session:receipt.session}}
    const expectedVersion=Number(command.expectedJournalVersion);if(!Number.isSafeInteger(expectedVersion)||expectedVersion<0)fail('validation','expectedJournalVersion is invalid','expectedJournalVersion');const sectionPredicate=command.expectedSection===null?sql`next_section IS NULL`:sql`next_section=${command.expectedSection}`
    const rows=await execRows<JournalRow>(tx,sql`UPDATE content_import_journal SET version=version+1,updated_at=${new Date().toISOString()} WHERE import_id=${sessionId} AND plan_hash=${command.expectedPlanHash} AND manifest_hash=${command.expectedManifestHash} AND version=${expectedVersion} AND ${sectionPredicate} AND next_offset=${command.expectedOffset} RETURNING import_id AS "importId",host_session_id AS "hostSessionId",manifest_hash AS "manifestHash",plan_hash AS "planHash",state,version,next_section AS "nextSection",next_offset AS "nextOffset",high_water_mark AS "highWaterMark",payload`)
    const row=rows[0];if(!row){if(!await readJournalRow(tx,sessionId))fail('authorization','import session is unavailable');fail('stale-journal','journal coordinate changed')}
    const session=sessionFromRow(row!);if(!sameScope(session.scope,authInput.scope)||!sameAuthority(session.requiredAuthority,authInput.requiredAuthority)||session.authorizationPolicyVersion!==command.expectedAuthorizationPolicyVersion)fail('authorization','import session is unavailable')
    try{await insertPendingReceipt(tx,command.operationId,kind,canonicalCommandHash,principalScope(principal))}catch(error){if(isUniqueViolation(error))return {kind:'conflict'};throw error}return {kind:'new',session}
  },
  async get(tx,sessionId,_principal){assertActiveContentTransaction(tx);const row=await readJournalRow(tx,sessionId);if(!row)fail('authorization','import session is unavailable');return sessionFromRow(row!)},
  async recordBatch(sessionId,batch,beforeImage,tx){
    assertActiveContentTransaction(tx);if(!Number.isSafeInteger(batch)||batch<0)fail('validation','batch must be non-negative','batch');const row=await readJournalRow(tx,sessionId);if(!row)fail('authorization','import session is unavailable');const payload=rowJson(row!.payload),batches=Array.isArray(payload.batches)?payload.batches as {batch:number;beforeImage:ContentSchemaValue}[]:[];if(batches.some((item)=>item.batch===batch))fail('integrity','duplicate import batch');const next=schemaValue({...payload,batches:[...batches,{batch,beforeImage}]});await tx.execute(sql`UPDATE content_import_journal SET payload=${JSON.stringify(next)},updated_at=${new Date().toISOString()} WHERE import_id=${sessionId}`)
  },
  async readBatchesReverse(tx,sessionId,expectedPlanHash,expectedManifestHash){
    assertActiveContentTransaction(tx);const row=await readJournalRow(tx,sessionId);if(!row||row.planHash!==expectedPlanHash||row.manifestHash!==expectedManifestHash)fail('authorization','import session is unavailable');const payload=rowJson(row!.payload),batches=Array.isArray(payload.batches)?payload.batches as {batch:number;beforeImage:ContentSchemaValue}[]:[];const sorted=[...batches].sort((a,b)=>b.batch-a.batch);for(let i=0;i<sorted.length;i++)if(sorted[i]!.batch!==sorted.length-1-i)fail('integrity','import batch ledger is not contiguous');return Object.freeze(sorted.map((item)=>Object.freeze(item)))
  },
  async transition(sessionId,from,to,tx){assertActiveContentTransaction(tx);const rows=await execRows<{importId:string}>(tx,sql`UPDATE content_import_journal SET state=${to},updated_at=${new Date().toISOString()} WHERE import_id=${sessionId} AND state=${from} RETURNING import_id AS "importId"`);if(rows.length!==1)fail('stale-journal','import state changed')},
  async completeCommand(tx,sessionId,kind,command,canonicalCommandHash,result){
    assertActiveContentTransaction(tx);const row=await readJournalRow(tx,sessionId);if(!row)fail('authorization','import session is unavailable');const payload=rowJson(row!.payload),next=schemaValue({...payload,session:result}),highWaterMark=result.publishedCorpusVersion??row!.highWaterMark,updated=await execRows<{importId:string}>(tx,sql`UPDATE content_import_journal SET state=${result.state},next_section=${result.nextSection},next_offset=${result.nextOffset},high_water_mark=${highWaterMark},payload=${JSON.stringify(next)},updated_at=${new Date().toISOString()} WHERE import_id=${sessionId} AND version=${Number(result.journalVersion)} RETURNING import_id AS "importId"`);if(updated.length!==1)fail('stale-journal','journal claim was lost')
    const pending=(await execRows<{payload:unknown}>(tx,sql`SELECT payload FROM content_lifecycle_journal WHERE operation_id=${command.operationId} AND state='pending'`))[0];if(!pending)fail('integrity','command claim disappeared');const receipt=rowJson(pending!.payload);if(receipt.commandHash!==canonicalCommandHash)fail('integrity','command claim hash changed');await commitReceipt(tx,command.operationId,kind,canonicalCommandHash,boundedString(receipt.principalScope,'receipt.principalScope',1024),result)
  },
};return Object.freeze(journal)}

function bundleLimits(bundle:ContentExportBundle):ContentImportLimits{const bytes=new TextEncoder().encode(JSON.stringify(bundle)).byteLength;return Object.freeze({maxBytes:Math.max(bytes,1),maxDefinitions:Math.max(bundle.types.length+bundle.statuses.length+bundle.taxonomies.length,1),maxEntries:Math.max(bundle.entries.length,bundle.terms.length,1),maxRevisions:Math.max(bundle.revisions.length,1),batchSize:100})}
function mappingFor(parsed:ParsedBundle):Record<ContentImportMappingKey,string>{const out={} as Record<ContentImportMappingKey,string>;for(const item of parsed.types)out[`type:${item.key}`]=item.key;for(const item of parsed.statuses)out[`status:${item.key}`]=item.key;for(const item of parsed.taxonomies)out[`taxonomy:${item.key}`]=item.key;for(const item of parsed.terms)out[`term:${item.id}`]=item.id;for(const item of parsed.entries)out[`entry:${item.id}`]=item.id;for(const item of parsed.revisions)out[`revision:${item.id}`]=item.id;return out}
function definitionRevisionMappingKey(kind:'type'|'status',keyValue:string,sourceRevision:number):ContentImportMappingKey{return `${kind}:${keyValue}@${sourceRevision}` as ContentImportMappingKey}
function mappedDefinitionRevision(mappings:ContentImportMappings,kind:'type'|'status',keyValue:string,sourceRevision:number):number{const raw=mappings[definitionRevisionMappingKey(kind,keyValue,sourceRevision)],revision=Number(raw);if(!Number.isSafeInteger(revision)||revision<1)fail('stale-plan',`missing ${kind} definition revision mapping for ${keyValue}@${sourceRevision}`);return revision}
function conservativeAuthority(scope:ContentImportScope):readonly ContentImportRequiredAuthority[]{return Object.freeze(scope.typeKeys.map((typeKey)=>Object.freeze({typeKey,branches:ALL_BRANCHES})))}
function conflict(code:string,message:string,fieldPath?:string):ContentOperationError{return Object.freeze({code,category:'conflict',message,...(fieldPath?{fieldPath}:{})})}
function addBranch(branches:Map<string,Set<ContentImportCapabilityBranch>>,typeKey:string,...items:ContentImportCapabilityBranch[]):void{const set=branches.get(typeKey);if(!set)fail('authorization','import scope is unavailable');for(const item of items)set!.add(item)}
async function assertImportAuthorized(tx:ContentTransaction,authorization:ContentImportAuthorization,principal:ContentPrincipal,action:'plan'|'start'|ContentImportCommandKind,input:ContentImportAuthorizationInput):Promise<string>{try{const result=await authorization.assert(tx,principal,action,input);if(!result||typeof result.policyVersion!=='string'||!result.policyVersion)fail('authorization','import scope is unavailable');return result.policyVersion}catch(error){if(error instanceof ContentImportError)throw error;return fail('authorization','import scope is unavailable')}}
interface DefinitionCurrentRow extends Record<string,unknown>{key:string;origin:'code'|'db'|'import';version:number|string;active:boolean;currentRevision:number|string;canonicalHash:string;definition:unknown}
interface DefinitionVersionRow extends Record<string,unknown>{kind:'type'|'status';key:string;revision:number|string;canonicalHash:string}
interface TermConflictRow extends Record<string,unknown>{id:string;taxonomy:string;slug:string;parentId:string|null;depth:number|string}
interface EntryConflictRow extends Record<string,unknown>{id:string;type:string;slug:string;author:string;status:string;visibility:'public'|'private'|'members'}
interface RevisionConflictRow extends Record<string,unknown>{id:string;entryId:string}
function sqlList(values:readonly string[]):SQLWrapper{return sql.join(values.map((value)=>sql`${value}`),sql`, `)}
async function destinationCorpusVersion(tx:ContentTransaction,scope:ContentImportScope):Promise<string>{
  const mappedTypeKeys=scope.typeKeys
  const types=await execRows<Record<string,unknown>>(tx,sql`SELECT key,origin,version,active,current_revision AS "currentRevision",canonical_hash AS "canonicalHash",definition,shadowed_db_version AS "shadowedDbVersion" FROM content_type_definitions WHERE key IN (${sqlList(mappedTypeKeys)}) ORDER BY key`)
  const statuses=await execRows<Record<string,unknown>>(tx,sql`SELECT key,origin,version,active,current_revision AS "currentRevision",canonical_hash AS "canonicalHash",definition,shadowed_db_version AS "shadowedDbVersion" FROM content_status_definitions ORDER BY key`)
  const statusKeys=statuses.map((row)=>String(row.key))
  const taxonomyKeys=[...new Set(types.flatMap((row)=>{try{return normalizeContentTypeDefinition((typeof row.definition==='string'?JSON.parse(row.definition):row.definition) as ContentTypeDefinition).taxonomies}catch{return []}}))].sort()
  const definitionKeys=[...mappedTypeKeys.map((key)=>`type:${key}`),...statusKeys.map((key)=>`status:${key}`)]
  const versions=definitionKeys.length?await execRows<Record<string,unknown>>(tx,sql`SELECT definition_kind AS kind,definition_key AS key,revision,canonical_hash AS "canonicalHash",definition,origin,created_at AS "createdAt" FROM content_definition_versions WHERE (definition_kind||':'||definition_key) IN (${sqlList(definitionKeys)}) ORDER BY definition_kind,definition_key,revision`):[]
  const entries=await execRows<Record<string,unknown>>(tx,sql`SELECT id,slug,type,title,body,status,visibility,published_at AS "publishedAt",author,created_at AS "createdAt",updated_at AS "updatedAt",parent_id AS "parentId",menu_order AS "menuOrder",template_key AS "templateKey",excerpt,featured_media AS "featuredMedia",comment_status AS "commentStatus",ping_status AS "pingStatus",sticky,format,deleted_at AS "deletedAt",last_edited_by AS "lastEditedBy",type_definition_revision AS "typeDefinitionRevision",status_definition_revision AS "statusDefinitionRevision",host_session_id AS "hostSessionId" FROM content_entries WHERE type IN (${sqlList(mappedTypeKeys)}) ORDER BY id`)
  const entryIds=entries.map((row)=>String(row.id)),assignments=entryIds.length?await execRows<Record<string,unknown>>(tx,sql`SELECT entry_id AS "entryId",term_id AS "termId" FROM content_entry_terms WHERE entry_id IN (${sqlList(entryIds)}) ORDER BY entry_id,term_id`):[]
  const terms=taxonomyKeys.length?await execRows<Record<string,unknown>>(tx,sql`SELECT id,taxonomy,slug,name,parent_id AS "parentId",depth,created_at AS "createdAt",host_session_id AS "hostSessionId" FROM content_terms WHERE taxonomy IN (${sqlList(taxonomyKeys)}) ORDER BY taxonomy,depth,id`):[]
  const revisions=entryIds.length?await execRows<Record<string,unknown>>(tx,sql`SELECT id,entry_id AS "entryId",seq,title,body,slug,type,term_ids AS "termIds",snapshot,editor,created_at AS "createdAt",host_session_id AS "hostSessionId" FROM content_revisions WHERE entry_id IN (${sqlList(entryIds)}) ORDER BY entry_id,seq,id`):[]
  return canonicalContentMigrationHash({scope,mappedTypeKeys,statusKeys,taxonomyKeys,types,statuses,versions,entries,assignments,terms,revisions})
}
function mappedTypeDefinition(item:TypeItem,mappings:ContentImportMappings):ContentTypeDefinition{const keyValue=mappings[`type:${item.key}`]??item.key,statusKeys=item.definition.statusKeys?.map((key)=>mappings[`status:${key}`]??key),taxonomies=item.definition.taxonomies.map((key)=>mappings[`taxonomy:${key}`]??key);return normalizeContentTypeDefinition({...item.definition,key:keyValue,statusKeys,taxonomies})}
function mappedStatusDefinition(item:StatusItem,mappings:ContentImportMappings):ContentStatusDefinition{return normalizeContentStatusDefinition({...item.definition,key:mappings[`status:${item.key}`]??item.key})}
function authorityForImportedEntry(entry:EntryItem,status:StatusItem,principal:ContentPrincipal,mode:'create'|'replace'):ContentImportCapabilityBranch[]{const branches:ContentImportCapabilityBranch[]=[];if(mode==='create'){branches.push('create',entry.author===principal.id?'deleteOwn':'deleteOthers');if(entry.visibility==='private')branches.push('deletePrivate');if(status.definition.published)branches.push('deletePublished')}else{branches.push(entry.author===principal.id?'editOwn':'editOthers');if(entry.visibility==='private')branches.push('editPrivate');if(status.definition.published)branches.push('editPublished')}if(status.definition.published)branches.push('publish');if(entry.termIds.length)branches.push('manageTerms');return branches}
function authorityForExistingEntry(entry:EntryConflictRow,statusPublished:boolean,principal:ContentPrincipal):ContentImportCapabilityBranch[]{const branches:ContentImportCapabilityBranch[]=[entry.author===principal.id?'editOwn':'editOthers'];if(entry.visibility==='private')branches.push('editPrivate');if(statusPublished)branches.push('editPublished');return branches}
async function deterministicUuid(namespace:string,source:string,manifestHash:string):Promise<string>{const hex=await canonicalContentMigrationHash({namespace,source,manifestHash}),chars=hex.slice(0,32).split('');chars[12]='5';chars[16]=((Number.parseInt(chars[16]!,16)&3)|8).toString(16);const value=chars.join('');return `${value.slice(0,8)}-${value.slice(8,12)}-${value.slice(12,16)}-${value.slice(16,20)}-${value.slice(20)}`}

function planMaterial(plan:Omit<ContentImportPlan,'planHash'>):ContentSchemaValue{return schemaValue(plan)}
async function assertPlanHash(plan:ContentImportPlan):Promise<void>{const {planHash,...rest}=plan;if(await canonicalContentMigrationHash(planMaterial(rest))!==planHash)fail('integrity','import plan hash is invalid')}
export async function planContentImport(tx:ContentTransaction,authorization:ContentImportAuthorization,principal:ContentPrincipal,bundle:ContentExportBundle,input:{scope:ContentImportScope;conflictPolicy:ContentImportConflictPolicy}):Promise<ContentImportPlan>{
  assertActiveContentTransaction(tx)
  const scope=stableScope(input.scope)
  if(input.conflictPolicy!=='reject'&&input.conflictPolicy!=='preserve'&&input.conflictPolicy!=='mappedReplacement')fail('validation','invalid conflict policy','conflictPolicy')
  const conservative=conservativeAuthority(scope)
  const firstPolicy=await assertImportAuthorized(tx,authorization,principal,'plan',{scope,requiredAuthority:conservative})
  const bytes=new TextEncoder().encode(JSON.stringify(bundle))
  const parsed=await parseBundle(bytes,bundleLimits(bundle))
  if(parsed.types.some((item)=>!scope.typeKeys.includes(item.key))||parsed.entries.some((item)=>!scope.typeKeys.includes(item.type)))fail('authorization','import scope is unavailable')
  const mappings=mappingFor(parsed)
  const conflicts:ContentOperationError[]=[]
  const branches=new Map(scope.typeKeys.map((typeKey)=>[typeKey,new Set<ContentImportCapabilityBranch>()]))
  const typeKeys=parsed.types.map((item)=>item.key),statusKeys=parsed.statuses.map((item)=>item.key),taxonomyKeys=parsed.taxonomies.map((item)=>item.key)
  const currentTypes=typeKeys.length?await execRows<DefinitionCurrentRow>(tx,sql`SELECT key,origin,version,active,current_revision AS "currentRevision",canonical_hash AS "canonicalHash",definition FROM content_type_definitions WHERE key IN (${sqlList(typeKeys)})`):[]
  const currentStatuses=statusKeys.length?await execRows<DefinitionCurrentRow>(tx,sql`SELECT key,origin,version,active,current_revision AS "currentRevision",canonical_hash AS "canonicalHash",definition FROM content_status_definitions WHERE key IN (${sqlList(statusKeys)})`):[]
  const definitionKeys=[...typeKeys.map((key)=>`type:${key}`),...statusKeys.map((key)=>`status:${key}`)]
  const currentVersions=definitionKeys.length?await execRows<DefinitionVersionRow>(tx,sql`SELECT definition_kind AS kind,definition_key AS key,revision,canonical_hash AS "canonicalHash" FROM content_definition_versions WHERE (definition_kind||':'||definition_key) IN (${sqlList(definitionKeys)})`):[]
  const currentTypeByKey=new Map(currentTypes.map((row)=>[row.key,row])),currentStatusByKey=new Map(currentStatuses.map((row)=>[row.key,row]))
  for(const item of parsed.types){
    const current=currentTypeByKey.get(item.key),destVersions=currentVersions.filter((row)=>row.kind==='type'&&row.key===item.key),maxRevision=Math.max(0,...destVersions.map((row)=>Number(row.revision)));let allocated=0
    addBranch(branches,item.key,'manageType')
    for(const version of item.versions){const match=destVersions.find((row)=>row.canonicalHash===version.canonicalHash);mappings[definitionRevisionMappingKey('type',item.key,version.revision)]=String(match?Number(match.revision):maxRevision+(++allocated))}
    if(current){
      const differs=current.canonicalHash!==item.canonicalHash||boolValue(current.active)!==item.active
      if(input.conflictPolicy==='reject')conflicts.push(conflict('type-conflict',`type ${item.key} already exists`,`type:${item.key}`))
      else if(input.conflictPolicy==='preserve'&&differs)conflicts.push(conflict('type-definition-mismatch',`type ${item.key} differs from preserved destination`,`type:${item.key}`))
      else if(input.conflictPolicy==='mappedReplacement'&&current.origin==='code'&&differs&&parsed.entries.some((entry)=>entry.type===item.key))conflicts.push(conflict('code-type-owned',`type ${item.key} is code-owned and cannot be replaced for imported entries`,`type:${item.key}`))
      else if(input.conflictPolicy==='mappedReplacement'&&current.origin!=='code'&&differs)addBranch(branches,item.key,'migrateAll')
    }
  }
  for(const item of parsed.statuses){
    for(const typeKey of scope.typeKeys)addBranch(branches,typeKey,'manageType')
    const current=currentStatusByKey.get(item.key),destVersions=currentVersions.filter((row)=>row.kind==='status'&&row.key===item.key),maxRevision=Math.max(0,...destVersions.map((row)=>Number(row.revision)));let allocated=0
    for(const version of item.versions){const match=destVersions.find((row)=>row.canonicalHash===version.canonicalHash);mappings[definitionRevisionMappingKey('status',item.key,version.revision)]=String(match?Number(match.revision):maxRevision+(++allocated))}
    if(current){
      const differs=current.canonicalHash!==item.canonicalHash||boolValue(current.active)!==item.active
      if(input.conflictPolicy==='reject')conflicts.push(conflict('status-conflict',`status ${item.key} already exists`,`status:${item.key}`))
      else if(input.conflictPolicy==='preserve'&&differs)conflicts.push(conflict('status-definition-mismatch',`status ${item.key} differs from preserved destination`,`status:${item.key}`))
      else if(input.conflictPolicy==='mappedReplacement'&&current.origin==='code'&&differs&&parsed.entries.some((entry)=>entry.status===item.key))conflicts.push(conflict('code-status-owned',`status ${item.key} is code-owned and cannot be replaced for imported entries`,`status:${item.key}`))
      else if(input.conflictPolicy==='mappedReplacement'&&current.origin!=='code'&&differs)for(const typeKey of scope.typeKeys)addBranch(branches,typeKey,'migrateAll')
    }
  }
  const existingTerms=taxonomyKeys.length?await execRows<TermConflictRow>(tx,sql`SELECT id,taxonomy,slug,parent_id AS "parentId",depth FROM content_terms WHERE taxonomy IN (${sqlList(taxonomyKeys)}) ORDER BY depth,id`):[]
  const termById=new Map(existingTerms.map((row)=>[row.id,row]))
  const termByPath=new Map(existingTerms.map((row)=>[`${row.taxonomy}:${row.parentId??''}:${row.slug}`,row]))
  for(const item of [...parsed.terms].sort((a,b)=>a.depth-b.depth||a.id.localeCompare(b.id))){
    const mappedTax=mappings[`taxonomy:${item.taxonomy}`]??item.taxonomy
    const mappedParent=item.parentId===null?null:(mappings[`term:${item.parentId}`]??item.parentId)
    const byId=termById.get(item.id),byPath=termByPath.get(`${mappedTax}:${mappedParent??''}:${item.slug}`),existing=byId??byPath
    if(byId&&byPath&&byId.id!==byPath.id){conflicts.push(conflict('term-conflict-ambiguous',`term ${item.id} collides by both id and sibling slug`,`term:${item.id}`));continue}
    if(existing){if(input.conflictPolicy==='reject')conflicts.push(conflict('term-conflict',`term ${item.id} conflicts with destination term ${existing.id}`,`term:${item.id}`));else mappings[`term:${item.id}`]=existing.id}
  }
  if(parsed.terms.length){for(const type of parsed.types)if(type.definition.taxonomies.some((taxonomy)=>taxonomyKeys.includes(taxonomy)))addBranch(branches,type.key,'manageTerms')}

  const existingEntries=typeKeys.length?await execRows<EntryConflictRow>(tx,sql`SELECT id,type,slug,author,status,visibility FROM content_entries WHERE type IN (${sqlList(typeKeys)})`):[]
  const entryById=new Map(existingEntries.map((row)=>[row.id,row])),entryByPath=new Map(existingEntries.map((row)=>[`${row.type}:${row.slug}`,row]))
  const statusPublished=new Map<string,boolean>()
  for(const row of currentStatuses){try{statusPublished.set(row.key,normalizeContentStatusDefinition((typeof row.definition==='string'?JSON.parse(row.definition):row.definition) as ContentStatusDefinition).published)}catch{fail('integrity','destination status definition is malformed')}}
  const entryConflictSources=new Set<string>()
  for(const item of parsed.entries){
    const mappedType=mappings[`type:${item.type}`]??item.type,byId=entryById.get(item.id),byPath=entryByPath.get(`${mappedType}:${item.slug}`),existing=byId??byPath
    if(byId&&byPath&&byId.id!==byPath.id){conflicts.push(conflict('entry-conflict-ambiguous',`entry ${item.id} collides by both id and type/slug`,`entry:${item.id}`));continue}
    const importedStatus=parsed.statuses.find((status)=>status.key===item.status)!
    if(existing){entryConflictSources.add(item.id);if(input.conflictPolicy==='reject')conflicts.push(conflict('entry-conflict',`entry ${item.id} conflicts with destination entry ${existing.id}`,`entry:${item.id}`));else mappings[`entry:${item.id}`]=existing.id;if(input.conflictPolicy==='mappedReplacement')addBranch(branches,item.type,...authorityForExistingEntry(existing,statusPublished.get(existing.status)===true,principal),...authorityForImportedEntry(item,importedStatus,principal,'replace'))}
    else addBranch(branches,item.type,...authorityForImportedEntry(item,importedStatus,principal,'create'))
  }

  const existingRevisions=parsed.revisions.length?await execRows<RevisionConflictRow>(tx,sql`SELECT id,entry_id AS "entryId" FROM content_revisions WHERE id IN (${sqlList(parsed.revisions.map((item)=>item.id))})`):[]
  const revisionIds=new Set(existingRevisions.map((row)=>row.id))
  for(const item of parsed.revisions){
    if(revisionIds.has(item.id)){if(input.conflictPolicy==='reject')conflicts.push(conflict('revision-conflict',`revision ${item.id} already exists`,`revision:${item.id}`));else if(input.conflictPolicy==='mappedReplacement')mappings[`revision:${item.id}`]=await deterministicUuid('revision',item.id,parsed.manifestHash)}
    if(input.conflictPolicy==='preserve'&&entryConflictSources.has(item.entryId))continue
  }
  const frozenMappings=Object.freeze({...mappings}) as ContentImportMappings
  const destination=await destinationCorpusVersion(tx,scope)
  const requiredAuthority=stableAuthority([...branches].map(([typeKey,set])=>({typeKey,branches:[...set]})))
  const policy=await assertImportAuthorized(tx,authorization,principal,'plan',{scope,requiredAuthority})
  if(policy!==firstPolicy)fail('authorization','import scope is unavailable')
  const counts=Object.freeze({types:parsed.types.length,statuses:parsed.statuses.length,taxonomies:parsed.taxonomies.length,terms:parsed.terms.length,entries:parsed.entries.length,revisions:parsed.revisions.length})
  const rest:Omit<ContentImportPlan,'planHash'>={manifestHash:parsed.manifestHash,destinationCorpusVersion:destination,authorizationPolicyVersion:policy,scope,conflictPolicy:input.conflictPolicy,requiredAuthority,counts,mappings:frozenMappings,conflicts:Object.freeze(conflicts),estimatedBytes:bytes.byteLength}
  return Object.freeze({planHash:await canonicalContentMigrationHash(planMaterial(rest)),...rest})
}

function authorizationInputFromPlan(plan:ContentImportPlan):ContentImportAuthorizationInput{return Object.freeze({scope:plan.scope,requiredAuthority:plan.requiredAuthority,expectedPolicyVersion:plan.authorizationPolicyVersion})}
function authorizationInputFromCommand(command:ContentImportCommand):ContentImportAuthorizationInput{return Object.freeze({scope:stableScope(command.scope),requiredAuthority:stableAuthority(command.requiredAuthority),expectedPolicyVersion:command.expectedAuthorizationPolicyVersion})}
async function canonicalCommandHash(kind:'start'|ContentImportCommandKind,sessionId:string|null,principal:ContentPrincipal,value:unknown):Promise<string>{return canonicalContentMigrationHash(schemaValue({kind,sessionId,principalScope:principalScope(principal),value}))}
async function parsedBundleObject(bundle:ContentExportBundle):Promise<ParsedBundle>{const bytes=new TextEncoder().encode(JSON.stringify(bundle));return parseBundle(bytes,bundleLimits(bundle))}
async function revalidatePlan(tx:ContentTransaction,authorization:ContentImportAuthorization,principal:ContentPrincipal,bundle:ContentExportBundle,expected:ContentImportPlan):Promise<void>{await assertPlanHash(expected);const current=await planContentImport(tx,authorization,principal,bundle,{scope:expected.scope,conflictPolicy:expected.conflictPolicy});if(current.planHash!==expected.planHash||current.destinationCorpusVersion!==expected.destinationCorpusVersion||current.manifestHash!==expected.manifestHash)fail('stale-plan','destination corpus or plan evidence changed')}

export async function startContentImport(tx:ContentTransaction,journal:ContentImportJournal,authorization:ContentImportAuthorization,principal:ContentPrincipal,bundle:ContentExportBundle,input:{plan:ContentImportPlan;staging:ContentImportStaging;operationId:string}):Promise<ContentImportSession>{
  assertActiveContentTransaction(tx);await assertPlanHash(input.plan);const operationId=boundedString(input.operationId,'operationId',256),parsed=await parsedBundleObject(bundle);if(parsed.manifestHash!==input.plan.manifestHash)fail('stale-plan','bundle manifest differs from plan');if(input.plan.conflicts.length>0)fail('conflict','import plan contains unresolved conflicts')
  const authInput=authorizationInputFromPlan(input.plan),policy=await assertImportAuthorized(tx,authorization,principal,'start',authInput);if(policy!==input.plan.authorizationPolicyVersion)fail('authorization','import scope is unavailable')
  const currentCorpus=await destinationCorpusVersion(tx,input.plan.scope),startHash=await canonicalCommandHash('start',null,principal,{manifestHash:parsed.manifestHash,plan:input.plan,staging:input.staging,authorizationInput:authInput,currentPolicyVersion:policy,destinationCorpusVersion:currentCorpus})
  const claim=await journal.claimStart(tx,principal,authInput,authorization,operationId,startHash);if(claim.kind==='conflict')fail('conflict','operation id was already used for another import start');if(claim.kind==='replay')return claim.session
  if(currentCorpus!==input.plan.destinationCorpusVersion)fail('stale-plan','destination corpus changed after planning');await revalidatePlan(tx,authorization,principal,bundle,input.plan)
  return journal.create(input.plan,input.staging,principal,operationId,tx)
}

function sectionItems(parsed:ParsedBundle,section:ContentImportSection):readonly unknown[]{switch(section){case'type':return parsed.types;case'status':return parsed.statuses;case'taxonomy':return parsed.taxonomies;case'term':return [...parsed.terms].sort((a,b)=>a.depth-b.depth||a.id.localeCompare(b.id));case'entry':return sortEntriesParentFirst(parsed.entries);case'revision':return [...parsed.revisions].sort((a,b)=>a.entryId.localeCompare(b.entryId)||a.seq-b.seq||a.id.localeCompare(b.id))}}
function sortEntriesParentFirst(entries:readonly EntryItem[]):readonly EntryItem[]{const byId=new Map(entries.map((item)=>[item.id,item])),depthCache=new Map<string,number>();const depth=(item:EntryItem):number=>{const cached=depthCache.get(item.id);if(cached!==undefined)return cached;const result=item.parentId===null?0:depth(byId.get(item.parentId)!)+1;depthCache.set(item.id,result);return result};return Object.freeze([...entries].sort((a,b)=>depth(a)-depth(b)||a.id.localeCompare(b.id)))}
function nextCoordinate(parsed:ParsedBundle,current:ContentImportSection,offset:number,batchSize:number):{section:ContentImportSection|null;offset:number}{const order:readonly ContentImportSection[]=['type','status','taxonomy','term','entry','revision'];const items=sectionItems(parsed,current);if(offset<items.length)return {section:current,offset};let index=order.indexOf(current)+1;while(index<order.length){const section=order[index]!;if(sectionItems(parsed,section).length>0)return {section,offset:0};index++}return {section:null,offset:0}}

interface StagedBatch {section:ContentImportSection;offset:number;mutations:readonly ContentSchemaValue[]}
async function currentDefinitionSnapshot(tx:ContentTransaction,kind:'type'|'status',keyValue:string):Promise<{current:ContentSchemaValue|null;versions:readonly ContentSchemaValue[]}>{
  const currentRows=kind==='type'
    ?await execRows<Record<string,unknown>>(tx,sql`SELECT key,origin,version,active,current_revision AS "currentRevision",canonical_hash AS "canonicalHash",definition,shadowed_db_version AS "shadowedDbVersion",host_session_id AS "hostSessionId",created_at AS "createdAt",updated_at AS "updatedAt" FROM content_type_definitions WHERE key=${keyValue}`)
    :await execRows<Record<string,unknown>>(tx,sql`SELECT key,origin,version,active,current_revision AS "currentRevision",canonical_hash AS "canonicalHash",definition,shadowed_db_version AS "shadowedDbVersion",host_session_id AS "hostSessionId",created_at AS "createdAt",updated_at AS "updatedAt" FROM content_status_definitions WHERE key=${keyValue}`)
  const versions=await execRows<Record<string,unknown>>(tx,sql`SELECT definition_kind AS kind,definition_key AS key,revision,canonical_hash AS "canonicalHash",definition,origin,host_session_id AS "hostSessionId",created_at AS "createdAt" FROM content_definition_versions WHERE definition_kind=${kind} AND definition_key=${keyValue} ORDER BY revision`)
  return {current:currentRows[0]?schemaValue(currentRows[0]):null,versions:Object.freeze(versions.map((row)=>schemaValue(row)))}
}
async function mappedTypeItem(item:TypeItem,mappings:ContentImportMappings):Promise<ContentSchemaValue>{const definition=mappedTypeDefinition(item,mappings),canonicalHash=await canonicalContentMigrationHash(definition),versions=[] as ContentSchemaValue[];for(const version of item.versions){const mapped=normalizeContentTypeDefinition({...version.definition,key:mappings[`type:${item.key}`]??item.key,statusKeys:version.definition.statusKeys?.map((key)=>mappings[`status:${key}`]??key),taxonomies:version.definition.taxonomies.map((key)=>mappings[`taxonomy:${key}`]??key)}),revision=mappedDefinitionRevision(mappings,'type',item.key,version.revision);versions.push(schemaValue({revision,canonicalHash:await canonicalContentMigrationHash(mapped),definition:mapped,origin:'import',createdAt:version.createdAt}))}return schemaValue({key:mappings[`type:${item.key}`]??item.key,origin:'import',version:item.version,revision:mappedDefinitionRevision(mappings,'type',item.key,item.revision),active:item.active,canonicalHash,definition,versions})}
async function mappedStatusItem(item:StatusItem,mappings:ContentImportMappings):Promise<ContentSchemaValue>{const definition=mappedStatusDefinition(item,mappings),canonicalHash=await canonicalContentMigrationHash(definition),versions=[] as ContentSchemaValue[];for(const version of item.versions){const mapped=normalizeContentStatusDefinition({...version.definition,key:mappings[`status:${item.key}`]??item.key}),revision=mappedDefinitionRevision(mappings,'status',item.key,version.revision);versions.push(schemaValue({revision,canonicalHash:await canonicalContentMigrationHash(mapped),definition:mapped,origin:'import',createdAt:version.createdAt}))}return schemaValue({key:mappings[`status:${item.key}`]??item.key,origin:'import',version:item.version,revision:mappedDefinitionRevision(mappings,'status',item.key,item.revision),active:item.active,canonicalHash,definition,versions})}


async function stageMutation(tx:ContentTransaction,session:ContentImportSession,section:ContentImportSection,raw:unknown):Promise<ContentSchemaValue>{
  if(section==='type'){
    const item=raw as TypeItem,destKey=session.mappings[`type:${item.key}`]??item.key,before=await currentDefinitionSnapshot(tx,'type',destKey),current=before.current===null?null:plain(before.current,'type.before.current'),mapped=plain(await mappedTypeItem(item,session.mappings),'type.after'),same=current!==null&&String(current.canonicalHash)===String(mapped.canonicalHash)&&boolValue(current.active)===boolValue(mapped.active),action:StagedAction=current===null?'create':session.conflictPolicy==='preserve'||same?'preserve':current.origin==='code'?'shadow':'replace',destinationVersion=current===null?1:action==='preserve'?Number(current.version):Number(current.version)+1
    return schemaValue({kind:'type',action,destKey,after:{...mapped,version:destinationVersion},before})
  }
  if(section==='status'){
    const item=raw as StatusItem,destKey=session.mappings[`status:${item.key}`]??item.key,before=await currentDefinitionSnapshot(tx,'status',destKey),current=before.current===null?null:plain(before.current,'status.before.current'),mapped=plain(await mappedStatusItem(item,session.mappings),'status.after'),same=current!==null&&String(current.canonicalHash)===String(mapped.canonicalHash)&&boolValue(current.active)===boolValue(mapped.active),action:StagedAction=current===null?'create':session.conflictPolicy==='preserve'||same?'preserve':current.origin==='code'?'shadow':'replace',destinationVersion=current===null?1:action==='preserve'?Number(current.version):Number(current.version)+1
    return schemaValue({kind:'status',action,destKey,after:{...mapped,version:destinationVersion},before})
  }
  if(section==='taxonomy'){
    const item=raw as TaxonomyItem;return schemaValue({kind:'taxonomy',action:'noop',source:item,destKey:session.mappings[`taxonomy:${item.key}`]??item.key})
  }
  if(section==='term'){
    const item=raw as TermItem,destId=session.mappings[`term:${item.id}`]??item.id,destTaxonomy=session.mappings[`taxonomy:${item.taxonomy}`]??item.taxonomy,destParentId=item.parentId===null?null:(session.mappings[`term:${item.parentId}`]??item.parentId)
    const rows=await execRows<Record<string,unknown>>(tx,sql`SELECT id,taxonomy,slug,name,parent_id AS "parentId",depth,created_at AS "createdAt",host_session_id AS "hostSessionId" FROM content_terms WHERE id=${destId}`),before=rows[0]?schemaValue(rows[0]):null,action=before===null?'create':session.conflictPolicy==='preserve'?'preserve':'replace'
    return schemaValue({kind:'term',action,destId,after:{id:destId,taxonomy:destTaxonomy,slug:item.slug,name:item.name,parentId:destParentId,depth:item.depth,createdAt:item.createdAt},before})
  }
  if(section==='entry'){
    const item=raw as EntryItem,destId=session.mappings[`entry:${item.id}`]??item.id,destType=session.mappings[`type:${item.type}`]??item.type,destStatus=session.mappings[`status:${item.status}`]??item.status,destParentId=item.parentId===null?null:(session.mappings[`entry:${item.parentId}`]??item.parentId),destTermIds=item.termIds.map((id)=>session.mappings[`term:${id}`]??id)
    const rows=await execRows<Record<string,unknown>>(tx,sql`SELECT id,slug,type,title,body,status,visibility,published_at AS "publishedAt",author,created_at AS "createdAt",updated_at AS "updatedAt",parent_id AS "parentId",menu_order AS "menuOrder",template_key AS "templateKey",excerpt,featured_media AS "featuredMedia",comment_status AS "commentStatus",ping_status AS "pingStatus",sticky,format,deleted_at AS "deletedAt",last_edited_by AS "lastEditedBy",type_definition_revision AS "typeDefinitionRevision",status_definition_revision AS "statusDefinitionRevision",host_session_id AS "hostSessionId" FROM content_entries WHERE id=${destId}`),before=rows[0]?schemaValue(rows[0]):null,beforeTerms=before===null?[]:await execRows<{termId:string}>(tx,sql`SELECT term_id AS "termId" FROM content_entry_terms WHERE entry_id=${destId} ORDER BY term_id`),action=before===null?'create':session.conflictPolicy==='preserve'?'preserve':'replace'
    return schemaValue({kind:'entry',action,destId,after:{id:destId,slug:item.slug,type:destType,title:item.title,body:item.body,status:destStatus,visibility:item.visibility,publishedAt:item.publishedAt,author:item.author,createdAt:item.createdAt,updatedAt:item.updatedAt,parentId:destParentId,menuOrder:item.menuOrder,templateKey:item.templateKey,excerpt:item.excerpt,featuredMedia:item.featuredMedia,commentStatus:item.commentStatus,pingStatus:item.pingStatus,sticky:item.sticky,format:item.format,deletedAt:item.deletedAt,lastEditedBy:item.lastEditedBy,typeDefinitionRevision:mappedDefinitionRevision(session.mappings,'type',item.type,item.typeDefinitionRevision),statusDefinitionRevision:mappedDefinitionRevision(session.mappings,'status',item.status,item.statusDefinitionRevision),termIds:destTermIds},before,beforeTermIds:beforeTerms.map((row)=>row.termId)})
  }
  const item=raw as RevisionItem,destId=session.mappings[`revision:${item.id}`]??item.id,destEntryId=session.mappings[`entry:${item.entryId}`]??item.entryId,destTermIds=item.termIds.map((id)=>session.mappings[`term:${id}`]??id)
  const entryExists=(await execRows<{id:string}>(tx,sql`SELECT id FROM content_entries WHERE id=${destEntryId}`))[0]!==undefined
  const rows=await execRows<Record<string,unknown>>(tx,sql`SELECT id,entry_id AS "entryId",seq,title,body,slug,type,term_ids AS "termIds",snapshot,editor,created_at AS "createdAt",host_session_id AS "hostSessionId" FROM content_revisions WHERE id=${destId}`),before=rows[0]?schemaValue(rows[0]):null
  const action=session.conflictPolicy==='preserve'&&entryExists?'preserve':before===null?'create':session.conflictPolicy==='preserve'?'preserve':'replace'
  return schemaValue({kind:'revision',action,destId,after:{id:destId,entryId:destEntryId,seq:item.seq,title:item.title,body:item.body,slug:item.slug,type:session.mappings[`type:${item.type}`]??item.type,termIds:destTermIds,...(item.snapshot===undefined?{}:{snapshot:item.snapshot}),editor:item.editor,createdAt:item.createdAt},before})
}

function stagedBatchFromValue(value:ContentSchemaValue):StagedBatch{const objectValue=plain(value,'stagedBatch'),rawSection=objectValue.section;if(rawSection!=='type'&&rawSection!=='status'&&rawSection!=='taxonomy'&&rawSection!=='term'&&rawSection!=='entry'&&rawSection!=='revision')throw new ContentImportError('integrity','staged batch section is invalid');const section:ContentImportSection=rawSection,offset=Number(objectValue.offset),rawMutations=objectValue.mutations;if(!Number.isSafeInteger(offset)||offset<0)fail('integrity','staged batch offset is invalid');if(!Array.isArray(rawMutations))throw new ContentImportError('integrity','staged batch mutations are invalid');return Object.freeze({section,offset,mutations:Object.freeze(rawMutations.map((item:unknown,index:number)=>schemaValue(item,`stagedBatch.mutations.${index}`)))})}

export async function resumeContentImport(tx:ContentTransaction,journal:ContentImportJournal,authorization:ContentImportAuthorization,principal:ContentPrincipal,bundle:ContentExportBundle,sessionIdInput:string,command:ContentImportCommand):Promise<ContentImportSession>{
  assertActiveContentTransaction(tx);const sessionId=boundedString(sessionIdInput,'sessionId',128),parsed=await parsedBundleObject(bundle);if(parsed.manifestHash!==command.expectedManifestHash)fail('stale-plan','bundle manifest differs from command');const commandHash=await canonicalCommandHash('resume',sessionId,principal,{command,bundleManifestHash:parsed.manifestHash}),claim=await journal.claimCommand(tx,sessionId,principal,'resume',command,commandHash,authorization)
  if(claim.kind==='conflict')throw new ContentImportError('conflict','operation id was already used for another import command');if(claim.kind==='replay')return claim.session
  const session=claim.session;if(!session)throw new ContentImportError('integrity','resume claim did not return session');if(session.state!=='validated'&&session.state!=='applying')fail('stale-journal','import is not resumable');if(session.nextSection===null)throw new ContentImportError('stale-journal','import is already staged');if(parsed.manifestHash!==session.manifestHash)fail('stale-plan','bundle manifest differs from session')
  const currentPlan=await planContentImport(tx,authorization,principal,bundle,{scope:session.scope,conflictPolicy:session.conflictPolicy});if(currentPlan.planHash!==session.planHash||currentPlan.destinationCorpusVersion!==session.destinationCorpusVersion)fail('stale-plan','destination corpus or plan evidence changed')
  if(await destinationCorpusVersion(tx,session.scope)!==session.destinationCorpusVersion)fail('stale-plan','destination corpus changed after planning')
  const section=session.nextSection,offset=session.nextOffset,batchSize=bundleLimits(bundle).batchSize,items=sectionItems(parsed,section),slice=items.slice(offset,offset+batchSize),mutations:ContentSchemaValue[]=[]
  for(const item of slice)mutations.push(await stageMutation(tx,session,section,item))
  const next=nextCoordinate(parsed,section,offset+slice.length,batchSize),result:ContentImportSession=Object.freeze({...session,state:next.section===null?'staged':'applying',appliedBatches:session.appliedBatches+1,nextSection:next.section,nextOffset:next.offset})
  const batch:StagedBatch=Object.freeze({section,offset,mutations:Object.freeze(mutations)});await journal.recordBatch(sessionId,session.appliedBatches,schemaValue(batch),tx);await journal.completeCommand(tx,sessionId,'resume',command,commandHash,result);return result
}

type StagedAction='create'|'replace'|'preserve'|'shadow'|'noop'
interface ParsedMutation {kind:ContentImportSection;action:StagedAction;record:Record<string,unknown>}
function parseMutation(value:ContentSchemaValue):ParsedMutation{
  const record=plain(value,'mutation'),kind=record.kind,action=record.action
  if(kind!=='type'&&kind!=='status'&&kind!=='taxonomy'&&kind!=='term'&&kind!=='entry'&&kind!=='revision')throw new ContentImportError('integrity','staged mutation kind is invalid')
  if(action!=='create'&&action!=='replace'&&action!=='preserve'&&action!=='shadow'&&action!=='noop')throw new ContentImportError('integrity','staged mutation action is invalid')
  return {kind,action,record}
}
function mutationAfter(record:Record<string,unknown>):Record<string,unknown>{return plain(record.after,'mutation.after')}
function mutationBefore(record:Record<string,unknown>):ContentSchemaValue|null{const value=record.before;return value===null?null:schemaValue(value,'mutation.before')}
function boolValue(value:unknown):boolean{return value===true||value===1||value==='1'}
function nullableJson(value:unknown):string|null{return value===null||value===undefined?null:JSON.stringify(value)}
async function applyDefinitionMutation(tx:ContentTransaction,session:ContentImportSession,mutation:ParsedMutation):Promise<boolean>{
  if(mutation.kind!=='type'&&mutation.kind!=='status')return false
  const after=mutationAfter(mutation.record),keyValue=boundedString(after.key,'definition.key',128),revision=positive(Number(after.revision),'definition.revision'),version=positive(Number(after.version),'definition.version'),canonicalHash=boundedString(after.canonicalHash,'definition.canonicalHash',128),definition=plain(after.definition,'definition.definition'),versions=array(after.versions,'definition.versions',1000),now=new Date().toISOString()
  if(mutation.action==='preserve')return true
  for(let index=0;index<versions.length;index++){
    const row=plain(versions[index],`definition.versions.${index}`),rowRevision=positive(Number(row.revision),`definition.versions.${index}.revision`),rowHash=boundedString(row.canonicalHash,`definition.versions.${index}.canonicalHash`,128),rowDefinition=plain(row.definition,`definition.versions.${index}.definition`),createdAt=iso(row.createdAt,`definition.versions.${index}.createdAt`)
    await tx.execute(sql`INSERT INTO content_definition_versions (definition_kind,definition_key,revision,canonical_hash,definition,origin,host_session_id,created_at) VALUES (${mutation.kind},${keyValue},${rowRevision},${rowHash},${JSON.stringify(rowDefinition)},'import',${session.staging.hostSessionId},${createdAt}) ON CONFLICT (definition_kind,definition_key,revision) DO NOTHING`)
  }
  if(mutation.action==='shadow'){
    if(mutation.kind==='type')await tx.execute(sql`UPDATE content_type_definitions SET shadowed_db_version=${version},updated_at=${now} WHERE key=${keyValue} AND origin='code'`)
    else await tx.execute(sql`UPDATE content_status_definitions SET shadowed_db_version=${version},updated_at=${now} WHERE key=${keyValue} AND origin='code'`)
    return true
  }
  const active=boolValue(after.active)
  if(mutation.kind==='type')await tx.execute(sql`INSERT INTO content_type_definitions (key,origin,version,active,current_revision,canonical_hash,definition,shadowed_db_version,host_session_id,created_at,updated_at) VALUES (${keyValue},'import',${version},${active},${revision},${canonicalHash},${JSON.stringify(definition)},NULL,${session.staging.hostSessionId},${now},${now}) ON CONFLICT (key) DO UPDATE SET origin='import',version=EXCLUDED.version,active=EXCLUDED.active,current_revision=EXCLUDED.current_revision,canonical_hash=EXCLUDED.canonical_hash,definition=EXCLUDED.definition,shadowed_db_version=NULL,host_session_id=EXCLUDED.host_session_id,updated_at=EXCLUDED.updated_at`)
  else await tx.execute(sql`INSERT INTO content_status_definitions (key,origin,version,active,current_revision,canonical_hash,definition,shadowed_db_version,host_session_id,created_at,updated_at) VALUES (${keyValue},'import',${version},${active},${revision},${canonicalHash},${JSON.stringify(definition)},NULL,${session.staging.hostSessionId},${now},${now}) ON CONFLICT (key) DO UPDATE SET origin='import',version=EXCLUDED.version,active=EXCLUDED.active,current_revision=EXCLUDED.current_revision,canonical_hash=EXCLUDED.canonical_hash,definition=EXCLUDED.definition,shadowed_db_version=NULL,host_session_id=EXCLUDED.host_session_id,updated_at=EXCLUDED.updated_at`)
  return true
}

async function applyTermMutation(tx:ContentTransaction,session:ContentImportSession,mutation:ParsedMutation):Promise<boolean>{
  if(mutation.kind==='taxonomy')return true
  if(mutation.kind!=='term')return false
  if(mutation.action==='preserve')return true
  const after=mutationAfter(mutation.record),id=uuid(after.id,'term.id'),taxonomy=key(after.taxonomy,'term.taxonomy'),slug=boundedString(after.slug,'term.slug',256),name=boundedString(after.name,'term.name',512),parentId=after.parentId===null?null:uuid(after.parentId,'term.parentId'),depth=nonnegative(Number(after.depth),'term.depth'),createdAt=iso(after.createdAt,'term.createdAt')
  await tx.execute(sql`INSERT INTO content_terms (id,taxonomy,slug,name,parent_id,depth,created_at,host_session_id) VALUES (${id},${taxonomy},${slug},${name},${parentId},${depth},${createdAt},${session.staging.hostSessionId}) ON CONFLICT (id) DO UPDATE SET taxonomy=EXCLUDED.taxonomy,slug=EXCLUDED.slug,name=EXCLUDED.name,parent_id=EXCLUDED.parent_id,depth=EXCLUDED.depth,host_session_id=EXCLUDED.host_session_id`)
  return true
}
async function applyEntryMutation(tx:ContentTransaction,session:ContentImportSession,mutation:ParsedMutation):Promise<boolean>{
  if(mutation.kind!=='entry')return false
  if(mutation.action==='preserve')return true
  const after=mutationAfter(mutation.record),id=uuid(after.id,'entry.id'),slug=boundedString(after.slug,'entry.slug',256),type=key(after.type,'entry.type'),title=typeof after.title==='string'?after.title:fail('integrity','staged entry title is invalid'),body=typeof after.body==='string'?after.body:fail('integrity','staged entry body is invalid'),status=key(after.status,'entry.status'),visibility=after.visibility,author=boundedString(after.author,'entry.author',512),createdAt=iso(after.createdAt,'entry.createdAt'),updatedAt=iso(after.updatedAt,'entry.updatedAt'),parentId=after.parentId===null?null:uuid(after.parentId,'entry.parentId'),menuOrder=nonnegative(Number(after.menuOrder),'entry.menuOrder'),templateKey=after.templateKey===null?null:boundedString(after.templateKey,'entry.templateKey',256),excerpt=typeof after.excerpt==='string'?after.excerpt:fail('integrity','staged entry excerpt is invalid'),commentStatus=after.commentStatus,pingStatus=after.pingStatus,sticky=after.sticky,lastEditedBy=boundedString(after.lastEditedBy,'entry.lastEditedBy',512),typeRevision=positive(Number(after.typeDefinitionRevision),'entry.typeDefinitionRevision'),statusRevision=positive(Number(after.statusDefinitionRevision),'entry.statusDefinitionRevision'),termIds=array(after.termIds,'entry.termIds',10000).map((value,index)=>uuid(value,`entry.termIds.${index}`))
  if(visibility!=='public'&&visibility!=='private'&&visibility!=='members')fail('integrity','staged entry visibility is invalid');if(commentStatus!=='open'&&commentStatus!=='closed')fail('integrity','staged entry comment status is invalid');if(pingStatus!=='open'&&pingStatus!=='closed')fail('integrity','staged entry ping status is invalid');if(typeof sticky!=='boolean')fail('integrity','staged entry sticky value is invalid')
  const publishedAt=after.publishedAt===null?null:iso(after.publishedAt,'entry.publishedAt'),deletedAt=after.deletedAt===null?null:iso(after.deletedAt,'entry.deletedAt'),format=after.format===null?null:boundedString(after.format,'entry.format',128),featuredMedia=nullableJson(after.featuredMedia)
  await tx.execute(sql`INSERT INTO content_entries (id,slug,type,title,body,status,visibility,published_at,author,created_at,updated_at,parent_id,menu_order,template_key,excerpt,featured_media,comment_status,ping_status,sticky,format,deleted_at,last_edited_by,type_definition_revision,status_definition_revision,host_session_id) VALUES (${id},${slug},${type},${title},${body},${status},${visibility},${publishedAt},${author},${createdAt},${updatedAt},${parentId},${menuOrder},${templateKey},${excerpt},${featuredMedia},${commentStatus},${pingStatus},${sticky},${format},${deletedAt},${lastEditedBy},${typeRevision},${statusRevision},${session.staging.hostSessionId}) ON CONFLICT (id) DO UPDATE SET slug=EXCLUDED.slug,type=EXCLUDED.type,title=EXCLUDED.title,body=EXCLUDED.body,status=EXCLUDED.status,visibility=EXCLUDED.visibility,published_at=EXCLUDED.published_at,author=EXCLUDED.author,updated_at=EXCLUDED.updated_at,parent_id=EXCLUDED.parent_id,menu_order=EXCLUDED.menu_order,template_key=EXCLUDED.template_key,excerpt=EXCLUDED.excerpt,featured_media=EXCLUDED.featured_media,comment_status=EXCLUDED.comment_status,ping_status=EXCLUDED.ping_status,sticky=EXCLUDED.sticky,format=EXCLUDED.format,deleted_at=EXCLUDED.deleted_at,last_edited_by=EXCLUDED.last_edited_by,type_definition_revision=EXCLUDED.type_definition_revision,status_definition_revision=EXCLUDED.status_definition_revision,host_session_id=EXCLUDED.host_session_id`)
  await tx.execute(sql`DELETE FROM content_entry_terms WHERE entry_id=${id}`)
  for(const termId of termIds)await tx.execute(sql`INSERT INTO content_entry_terms (entry_id,term_id) VALUES (${id},${termId}) ON CONFLICT (entry_id,term_id) DO NOTHING`)
  return true
}
async function applyRevisionMutation(tx:ContentTransaction,session:ContentImportSession,mutation:ParsedMutation):Promise<boolean>{
  if(mutation.kind!=='revision')return false
  if(mutation.action==='preserve')return true
  const after=mutationAfter(mutation.record),id=uuid(after.id,'revision.id'),entryId=uuid(after.entryId,'revision.entryId'),title=typeof after.title==='string'?after.title:fail('integrity','staged revision title is invalid'),body=typeof after.body==='string'?after.body:fail('integrity','staged revision body is invalid'),slug=boundedString(after.slug,'revision.slug',256),type=key(after.type,'revision.type'),termIds=array(after.termIds,'revision.termIds',10000).map((value,index)=>uuid(value,`revision.termIds.${index}`)),snapshot=after.snapshot===undefined?null:nullableJson(after.snapshot),editor=boundedString(after.editor,'revision.editor',512),createdAt=iso(after.createdAt,'revision.createdAt')
  const exists=(await execRows<{id:string}>(tx,sql`SELECT id FROM content_revisions WHERE id=${id}`))[0]!==undefined
  if(exists)await tx.execute(sql`UPDATE content_revisions SET entry_id=${entryId},title=${title},body=${body},slug=${slug},type=${type},term_ids=${JSON.stringify(termIds)},snapshot=${snapshot},editor=${editor},created_at=${createdAt},host_session_id=${session.staging.hostSessionId} WHERE id=${id}`)
  else await tx.execute(sql`INSERT INTO content_revisions (id,entry_id,title,body,slug,type,term_ids,snapshot,editor,created_at,host_session_id) VALUES (${id},${entryId},${title},${body},${slug},${type},${JSON.stringify(termIds)},${snapshot},${editor},${createdAt},${session.staging.hostSessionId})`)
  return true
}
async function applyStagedMutation(tx:ContentTransaction,session:ContentImportSession,value:ContentSchemaValue):Promise<void>{const mutation=parseMutation(value);if(await applyDefinitionMutation(tx,session,mutation))return;if(await applyTermMutation(tx,session,mutation))return;if(await applyEntryMutation(tx,session,mutation))return;if(await applyRevisionMutation(tx,session,mutation))return;fail('integrity','unknown staged mutation')}

async function restoreDefinitionMutation(tx:ContentTransaction,session:ContentImportSession,mutation:ParsedMutation):Promise<boolean>{
  if(mutation.kind!=='type'&&mutation.kind!=='status')return false
  const after=mutationAfter(mutation.record),keyValue=boundedString(after.key,'definition.key',128),beforeBundle=plain(mutation.record.before,'definition.before'),beforeCurrent=beforeBundle.current
  await tx.execute(sql`DELETE FROM content_definition_versions WHERE definition_kind=${mutation.kind} AND definition_key=${keyValue} AND host_session_id=${session.staging.hostSessionId}`)
  if(mutation.action==='preserve')return true
  if(beforeCurrent===null){
    if(mutation.kind==='type')await tx.execute(sql`DELETE FROM content_type_definitions WHERE key=${keyValue} AND host_session_id=${session.staging.hostSessionId}`)
    else await tx.execute(sql`DELETE FROM content_status_definitions WHERE key=${keyValue} AND host_session_id=${session.staging.hostSessionId}`)
    return true
  }
  const before=plain(beforeCurrent,'definition.before.current'),origin=boundedString(before.origin,'definition.before.origin',16),version=positive(Number(before.version),'definition.before.version'),active=boolValue(before.active),revision=positive(Number(before.currentRevision),'definition.before.currentRevision'),hash=boundedString(before.canonicalHash,'definition.before.canonicalHash',128),definition=plain(before.definition,'definition.before.definition'),shadow=before.shadowedDbVersion===null||before.shadowedDbVersion===undefined?null:positive(Number(before.shadowedDbVersion),'definition.before.shadowedDbVersion'),host=before.hostSessionId===null||before.hostSessionId===undefined?null:String(before.hostSessionId),createdAt=iso(before.createdAt,'definition.before.createdAt'),updatedAt=iso(before.updatedAt,'definition.before.updatedAt')
  if(origin!=='code'&&origin!=='db'&&origin!=='import')fail('integrity','definition before-image origin is invalid')
  if(mutation.kind==='type')await tx.execute(sql`INSERT INTO content_type_definitions (key,origin,version,active,current_revision,canonical_hash,definition,shadowed_db_version,host_session_id,created_at,updated_at) VALUES (${keyValue},${origin},${version},${active},${revision},${hash},${JSON.stringify(definition)},${shadow},${host},${createdAt},${updatedAt}) ON CONFLICT (key) DO UPDATE SET origin=EXCLUDED.origin,version=EXCLUDED.version,active=EXCLUDED.active,current_revision=EXCLUDED.current_revision,canonical_hash=EXCLUDED.canonical_hash,definition=EXCLUDED.definition,shadowed_db_version=EXCLUDED.shadowed_db_version,host_session_id=EXCLUDED.host_session_id,created_at=EXCLUDED.created_at,updated_at=EXCLUDED.updated_at`)
  else await tx.execute(sql`INSERT INTO content_status_definitions (key,origin,version,active,current_revision,canonical_hash,definition,shadowed_db_version,host_session_id,created_at,updated_at) VALUES (${keyValue},${origin},${version},${active},${revision},${hash},${JSON.stringify(definition)},${shadow},${host},${createdAt},${updatedAt}) ON CONFLICT (key) DO UPDATE SET origin=EXCLUDED.origin,version=EXCLUDED.version,active=EXCLUDED.active,current_revision=EXCLUDED.current_revision,canonical_hash=EXCLUDED.canonical_hash,definition=EXCLUDED.definition,shadowed_db_version=EXCLUDED.shadowed_db_version,host_session_id=EXCLUDED.host_session_id,created_at=EXCLUDED.created_at,updated_at=EXCLUDED.updated_at`)
  return true
}
async function restoreTermMutation(tx:ContentTransaction,session:ContentImportSession,mutation:ParsedMutation):Promise<boolean>{
  if(mutation.kind==='taxonomy')return true
  if(mutation.kind!=='term')return false
  if(mutation.action==='preserve')return true
  const destId=uuid(mutation.record.destId,'term.destId'),before=mutationBefore(mutation.record)
  if(before===null){await tx.execute(sql`DELETE FROM content_terms WHERE id=${destId} AND host_session_id=${session.staging.hostSessionId}`);return true}
  const row=plain(before,'term.before'),taxonomy=key(row.taxonomy,'term.before.taxonomy'),slug=boundedString(row.slug,'term.before.slug',256),name=boundedString(row.name,'term.before.name',512),parentId=row.parentId===null?null:uuid(row.parentId,'term.before.parentId'),depth=nonnegative(Number(row.depth),'term.before.depth'),createdAt=iso(row.createdAt,'term.before.createdAt'),host=row.hostSessionId===null||row.hostSessionId===undefined?null:String(row.hostSessionId)
  await tx.execute(sql`INSERT INTO content_terms (id,taxonomy,slug,name,parent_id,depth,created_at,host_session_id) VALUES (${destId},${taxonomy},${slug},${name},${parentId},${depth},${createdAt},${host}) ON CONFLICT (id) DO UPDATE SET taxonomy=EXCLUDED.taxonomy,slug=EXCLUDED.slug,name=EXCLUDED.name,parent_id=EXCLUDED.parent_id,depth=EXCLUDED.depth,created_at=EXCLUDED.created_at,host_session_id=EXCLUDED.host_session_id`)
  return true
}

async function restoreEntryMutation(tx:ContentTransaction,session:ContentImportSession,mutation:ParsedMutation):Promise<boolean>{
  if(mutation.kind!=='entry')return false
  if(mutation.action==='preserve')return true
  const destId=uuid(mutation.record.destId,'entry.destId'),before=mutationBefore(mutation.record)
  if(before===null){await tx.execute(sql`DELETE FROM content_entries WHERE id=${destId} AND host_session_id=${session.staging.hostSessionId}`);return true}
  const row=plain(before,'entry.before'),slug=boundedString(row.slug,'entry.before.slug',256),type=key(row.type,'entry.before.type'),title=typeof row.title==='string'?row.title:fail('integrity','entry before title is invalid'),body=typeof row.body==='string'?row.body:fail('integrity','entry before body is invalid'),status=key(row.status,'entry.before.status'),visibility=row.visibility,publishedAt=row.publishedAt===null?null:iso(row.publishedAt,'entry.before.publishedAt'),author=boundedString(row.author,'entry.before.author',512),createdAt=iso(row.createdAt,'entry.before.createdAt'),updatedAt=iso(row.updatedAt,'entry.before.updatedAt'),parentId=row.parentId===null?null:uuid(row.parentId,'entry.before.parentId'),menuOrder=nonnegative(Number(row.menuOrder),'entry.before.menuOrder'),templateKey=row.templateKey===null?null:boundedString(row.templateKey,'entry.before.templateKey',256),excerpt=typeof row.excerpt==='string'?row.excerpt:fail('integrity','entry before excerpt is invalid'),commentStatus=row.commentStatus,pingStatus=row.pingStatus,sticky=boolValue(row.sticky),format=row.format===null?null:boundedString(row.format,'entry.before.format',128),deletedAt=row.deletedAt===null?null:iso(row.deletedAt,'entry.before.deletedAt'),lastEditedBy=boundedString(row.lastEditedBy,'entry.before.lastEditedBy',512),typeRevision=positive(Number(row.typeDefinitionRevision),'entry.before.typeDefinitionRevision'),statusRevision=positive(Number(row.statusDefinitionRevision),'entry.before.statusDefinitionRevision'),host=row.hostSessionId===null||row.hostSessionId===undefined?null:String(row.hostSessionId),featuredMedia=nullableJson(row.featuredMedia),beforeTerms=array(mutation.record.beforeTermIds??[],'entry.beforeTermIds',10000).map((value,index)=>uuid(value,`entry.beforeTermIds.${index}`))
  if(visibility!=='public'&&visibility!=='private'&&visibility!=='members')fail('integrity','entry before visibility is invalid');if(commentStatus!=='open'&&commentStatus!=='closed')fail('integrity','entry before comment status is invalid');if(pingStatus!=='open'&&pingStatus!=='closed')fail('integrity','entry before ping status is invalid')
  await tx.execute(sql`INSERT INTO content_entries (id,slug,type,title,body,status,visibility,published_at,author,created_at,updated_at,parent_id,menu_order,template_key,excerpt,featured_media,comment_status,ping_status,sticky,format,deleted_at,last_edited_by,type_definition_revision,status_definition_revision,host_session_id) VALUES (${destId},${slug},${type},${title},${body},${status},${visibility},${publishedAt},${author},${createdAt},${updatedAt},${parentId},${menuOrder},${templateKey},${excerpt},${featuredMedia},${commentStatus},${pingStatus},${sticky},${format},${deletedAt},${lastEditedBy},${typeRevision},${statusRevision},${host}) ON CONFLICT (id) DO UPDATE SET slug=EXCLUDED.slug,type=EXCLUDED.type,title=EXCLUDED.title,body=EXCLUDED.body,status=EXCLUDED.status,visibility=EXCLUDED.visibility,published_at=EXCLUDED.published_at,author=EXCLUDED.author,created_at=EXCLUDED.created_at,updated_at=EXCLUDED.updated_at,parent_id=EXCLUDED.parent_id,menu_order=EXCLUDED.menu_order,template_key=EXCLUDED.template_key,excerpt=EXCLUDED.excerpt,featured_media=EXCLUDED.featured_media,comment_status=EXCLUDED.comment_status,ping_status=EXCLUDED.ping_status,sticky=EXCLUDED.sticky,format=EXCLUDED.format,deleted_at=EXCLUDED.deleted_at,last_edited_by=EXCLUDED.last_edited_by,type_definition_revision=EXCLUDED.type_definition_revision,status_definition_revision=EXCLUDED.status_definition_revision,host_session_id=EXCLUDED.host_session_id`)
  await tx.execute(sql`DELETE FROM content_entry_terms WHERE entry_id=${destId}`);for(const termId of beforeTerms)await tx.execute(sql`INSERT INTO content_entry_terms (entry_id,term_id) VALUES (${destId},${termId}) ON CONFLICT (entry_id,term_id) DO NOTHING`)
  return true
}
async function restoreRevisionMutation(tx:ContentTransaction,session:ContentImportSession,mutation:ParsedMutation):Promise<boolean>{
  if(mutation.kind!=='revision')return false
  if(mutation.action==='preserve')return true
  const destId=uuid(mutation.record.destId,'revision.destId'),before=mutationBefore(mutation.record)
  if(before===null){await tx.execute(sql`DELETE FROM content_revisions WHERE id=${destId} AND host_session_id=${session.staging.hostSessionId}`);return true}
  const row=plain(before,'revision.before'),entryId=uuid(row.entryId,'revision.before.entryId'),title=typeof row.title==='string'?row.title:fail('integrity','revision before title is invalid'),body=typeof row.body==='string'?row.body:fail('integrity','revision before body is invalid'),slug=boundedString(row.slug,'revision.before.slug',256),type=key(row.type,'revision.before.type'),termIds=array(row.termIds,'revision.before.termIds',10000).map((value,index)=>uuid(value,`revision.before.termIds.${index}`)),snapshot=row.snapshot===null||row.snapshot===undefined?null:nullableJson(row.snapshot),editor=boundedString(row.editor,'revision.before.editor',512),createdAt=iso(row.createdAt,'revision.before.createdAt'),host=row.hostSessionId===null||row.hostSessionId===undefined?null:String(row.hostSessionId)
  const updated=await execRows<{id:string}>(tx,sql`UPDATE content_revisions SET entry_id=${entryId},title=${title},body=${body},slug=${slug},type=${type},term_ids=${JSON.stringify(termIds)},snapshot=${snapshot},editor=${editor},created_at=${createdAt},host_session_id=${host} WHERE id=${destId} RETURNING id`);if(updated.length!==1)fail('integrity','revision before-image target disappeared')
  return true
}
async function restoreStagedMutation(tx:ContentTransaction,session:ContentImportSession,value:ContentSchemaValue):Promise<void>{const mutation=parseMutation(value);if(await restoreRevisionMutation(tx,session,mutation))return;if(await restoreEntryMutation(tx,session,mutation))return;if(await restoreTermMutation(tx,session,mutation))return;if(await restoreDefinitionMutation(tx,session,mutation))return;fail('integrity','unknown staged mutation')}

async function readValidatedLedger(tx:ContentTransaction,journal:ContentImportJournal,session:ContentImportSession):Promise<readonly {batch:number;staged:StagedBatch}[]>{const reverse=await journal.readBatchesReverse(tx,session.id,session.planHash,session.manifestHash);if(reverse.length!==session.appliedBatches)fail('integrity','import batch ledger count differs from session');const forward=[...reverse].reverse().map((item)=>({batch:item.batch,staged:stagedBatchFromValue(item.beforeImage)}));for(let i=0;i<forward.length;i++)if(forward[i]!.batch!==i)fail('integrity','import batch ledger sequence is invalid');let lastSection=-1,lastEnd=0;for(const item of forward){const sectionIndex=(['type','status','taxonomy','term','entry','revision'] as const).indexOf(item.staged.section);if(sectionIndex<lastSection)fail('integrity','import batch section order is invalid');if(sectionIndex===lastSection&&item.staged.offset!==lastEnd)fail('integrity','import batch offset sequence is invalid');if(sectionIndex>lastSection){lastSection=sectionIndex;lastEnd=0;if(item.staged.offset!==0)fail('integrity','import section does not start at offset zero')}lastEnd=item.staged.offset+item.staged.mutations.length}return Object.freeze(forward)}
function affectedCounts(ledger:readonly {batch:number;staged:StagedBatch}[]):Readonly<Record<string,number>>{const counts:Record<string,number>={};for(const {staged} of ledger)for(const raw of staged.mutations){const mutation=parseMutation(raw);if(mutation.action==='preserve'||mutation.action==='noop')continue;counts[mutation.kind]=(counts[mutation.kind]??0)+1}return Object.freeze(counts)}
async function beforeImageHash(ledger:readonly {batch:number;staged:StagedBatch}[]):Promise<string>{return canonicalContentMigrationHash(schemaValue(ledger.map((item)=>({batch:item.batch,staged:item.staged}))))}
function lifecycleEvent(kind:ContentImportLifecycleEvent['kind'],operationId:string,session:ContentImportSession,principal:ContentPrincipal,counts:Readonly<Record<string,number>>):ContentImportLifecycleEvent{return Object.freeze({id:crypto.randomUUID(),version:1,kind,operationId,sessionId:session.id,planHash:session.planHash,manifestHash:session.manifestHash,principalId:principal.id,scope:session.scope,authorizationPolicyVersion:session.authorizationPolicyVersion,affectedCounts:counts})}
async function emitLifecycleEffects(tx:ContentTransaction,effects:ContentImportLifecycleEffects,event:ContentImportLifecycleEvent):Promise<void>{await effects.audit.record(tx,event);await effects.outbox.enqueue(tx,event)}

export async function publishContentImport(tx:ContentTransaction,journal:ContentImportJournal,authorization:ContentImportAuthorization,effects:ContentImportLifecycleEffects,principal:ContentPrincipal,sessionIdInput:string,command:ContentImportCommand):Promise<ContentImportSession>{
  assertActiveContentTransaction(tx);const sessionId=boundedString(sessionIdInput,'sessionId',128),commandHash=await canonicalCommandHash('publish',sessionId,principal,command),claim=await journal.claimCommand(tx,sessionId,principal,'publish',command,commandHash,authorization)
  if(claim.kind==='conflict')throw new ContentImportError('conflict','operation id was already used for another import command');if(claim.kind==='replay')return claim.session
  const session=claim.session;if(!session)throw new ContentImportError('integrity','publish claim did not return session');if(session.state!=='staged'||session.nextSection!==null||session.nextOffset!==0)fail('stale-journal','import is not staged for publication')
  if(await destinationCorpusVersion(tx,session.scope)!==session.destinationCorpusVersion)fail('stale-plan','destination corpus changed after staging')
  const ledger=await readValidatedLedger(tx,journal,session),counts=affectedCounts(ledger)
  for(const {staged} of ledger)for(const mutation of staged.mutations)await applyStagedMutation(tx,session,mutation)
  const publishedCorpus=await destinationCorpusVersion(tx,session.scope)
  const result:ContentImportSession=Object.freeze({...session,state:'published',publishedCorpusVersion:publishedCorpus,nextSection:null,nextOffset:0}),event=lifecycleEvent('published',command.operationId,result,principal,counts);await emitLifecycleEffects(tx,effects,event);await journal.completeCommand(tx,session.id,'publish',command,commandHash,result);return result
}

export async function rollbackContentImport(tx:ContentTransaction,journal:ContentImportJournal,authorization:ContentImportAuthorization,effects:ContentImportLifecycleEffects,principal:ContentPrincipal,sessionIdInput:string,command:ContentImportCommand):Promise<ContentImportSession>{
  assertActiveContentTransaction(tx);const sessionId=boundedString(sessionIdInput,'sessionId',128),commandHash=await canonicalCommandHash('rollback',sessionId,principal,command),claim=await journal.claimCommand(tx,sessionId,principal,'rollback',command,commandHash,authorization)
  if(claim.kind==='conflict')throw new ContentImportError('conflict','operation id was already used for another import command');if(claim.kind==='replay')return claim.session
  const session=claim.session;if(!session)throw new ContentImportError('integrity','rollback claim did not return session');if(session.state!=='validated'&&session.state!=='applying'&&session.state!=='staged')fail('stale-journal','published or terminal imports cannot use staged rollback')
  if(await destinationCorpusVersion(tx,session.scope)!==session.destinationCorpusVersion)fail('stale-plan','destination corpus changed after planning')
  const ledger=await readValidatedLedger(tx,journal,session),counts=affectedCounts(ledger);for(const item of ledger)void stagedBatchFromValue(schemaValue(item.staged))
  const result:ContentImportSession=Object.freeze({...session,state:'rolledBack',nextSection:null,nextOffset:0}),event=lifecycleEvent('rolledBack',command.operationId,result,principal,counts);await emitLifecycleEffects(tx,effects,event);await journal.completeCommand(tx,session.id,'rollback',command,commandHash,result);return result
}

function reversePlanMaterial(plan:Omit<ContentPublishedReversePlan,'reversePlanHash'>):ContentSchemaValue{return schemaValue(plan)}
async function assertReversePlanHash(plan:ContentPublishedReversePlan):Promise<void>{const {reversePlanHash,...rest}=plan;if(await canonicalContentMigrationHash(reversePlanMaterial(rest))!==reversePlanHash)fail('integrity','published reverse plan hash is invalid')}
export async function planPublishedContentImportReverse(tx:ContentTransaction,journal:ContentImportJournal,authorization:ContentImportAuthorization,principal:ContentPrincipal,sessionIdInput:string,authorizationInput:ContentImportAuthorizationInput):Promise<ContentPublishedReversePlan>{
  assertActiveContentTransaction(tx);const sessionId=boundedString(sessionIdInput,'sessionId',128),scope=stableScope(authorizationInput.scope),requiredAuthority=stableAuthority(authorizationInput.requiredAuthority),policy=await assertImportAuthorized(tx,authorization,principal,'reversePublished',{scope,requiredAuthority,...(authorizationInput.expectedPolicyVersion?{expectedPolicyVersion:authorizationInput.expectedPolicyVersion}:{})});if(authorizationInput.expectedPolicyVersion!==undefined&&policy!==authorizationInput.expectedPolicyVersion)fail('authorization','import scope is unavailable')
  const session=await journal.get(tx,sessionId,principal);if(session.state!=='published')fail('stale-journal','import is not published');if(!sameScope(scope,session.scope)||!sameAuthority(requiredAuthority,session.requiredAuthority)||policy!==session.authorizationPolicyVersion)fail('authorization','import scope is unavailable');if(!session.publishedCorpusVersion)throw new ContentImportError('integrity','published import has no corpus high-water mark')
  const publishedCorpusVersion=session.publishedCorpusVersion,currentPublishedCorpus=await destinationCorpusVersion(tx,scope);if(currentPublishedCorpus!==publishedCorpusVersion)fail('stale-plan','published corpus changed after import')
  const ledger=await readValidatedLedger(tx,journal,session),hash=await beforeImageHash(ledger),counts=affectedCounts(ledger),rest:Omit<ContentPublishedReversePlan,'reversePlanHash'>={sessionId,originalPlanHash:session.planHash,manifestHash:session.manifestHash,publishedCorpusVersion,beforeImageHash:hash,affectedCounts:counts,scope,requiredAuthority,authorizationPolicyVersion:policy}
  return Object.freeze({reversePlanHash:await canonicalContentMigrationHash(reversePlanMaterial(rest)),...rest})
}

export async function reversePublishedContentImport(tx:ContentTransaction,journal:ContentImportJournal,authorization:ContentImportAuthorization,effects:ContentImportLifecycleEffects,principal:ContentPrincipal,plan:ContentPublishedReversePlan,command:ContentPublishedReverseCommand):Promise<ContentImportSession>{
  assertActiveContentTransaction(tx);await assertReversePlanHash(plan);if(command.expectedReversePlanHash!==plan.reversePlanHash||command.expectedPublishedCorpusVersion!==plan.publishedCorpusVersion)fail('conflict','reverse command identity differs from reverse plan');const commandHash=await canonicalCommandHash('reversePublished',plan.sessionId,principal,{plan,command}),claim=await journal.claimCommand(tx,plan.sessionId,principal,'reversePublished',command,commandHash,authorization)
  if(claim.kind==='conflict')throw new ContentImportError('conflict','operation id was already used for another import command');if(claim.kind==='replay')return claim.session
  const session=claim.session;if(!session)throw new ContentImportError('integrity','reverse claim did not return session');if(session.state!=='published')fail('stale-journal','import is not published');if(session.planHash!==plan.originalPlanHash||session.manifestHash!==plan.manifestHash||!sameScope(session.scope,plan.scope)||!sameAuthority(session.requiredAuthority,plan.requiredAuthority)||session.authorizationPolicyVersion!==plan.authorizationPolicyVersion)fail('conflict','reverse plan differs from published session')
  if(await destinationCorpusVersion(tx,session.scope)!==plan.publishedCorpusVersion)fail('stale-plan','published corpus changed after reverse planning');const ledger=await readValidatedLedger(tx,journal,session);if(await beforeImageHash(ledger)!==plan.beforeImageHash)fail('stale-plan','retained before-image ledger changed')
  await journal.transition(session.id,'published','reversing',tx)
  for(const {staged} of [...ledger].reverse())for(const mutation of [...staged.mutations].reverse())await restoreStagedMutation(tx,session,mutation)
  if(await destinationCorpusVersion(tx,session.scope)!==session.destinationCorpusVersion)fail('integrity','published reverse did not restore original corpus')
  const result:ContentImportSession=Object.freeze({...session,state:'reversed',nextSection:null,nextOffset:0}),event=lifecycleEvent('reversed',command.operationId,result,principal,plan.affectedCounts);await emitLifecycleEffects(tx,effects,event);await journal.completeCommand(tx,session.id,'reversePublished',command,commandHash,result);return result
}
