diff --git a/.github/workflows/apply-portable-transfer-domain.yml b/.github/workflows/apply-portable-transfer-domain.yml new file mode 100644 index 0000000..6e8416d --- /dev/null +++ b/.github/workflows/apply-portable-transfer-domain.yml @@ -0,0 +1,99 @@ +name: Apply portable transfer domain integration + +on: + push: + branches: + - architecture/portable-transfer + pull_request: + branches: + - master + workflow_dispatch: + +permissions: + contents: write + +concurrency: + group: apply-portable-transfer-domain + cancel-in-progress: false + +jobs: + integrate-and-verify: + if: github.actor != 'github-actions[bot]' + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + with: + ref: architecture/portable-transfer + fetch-depth: 0 + + - uses: actions/setup-node@v4 + with: + node-version: "22" + cache: npm + cache-dependency-path: frontend/package-lock.json + + - name: Validate and apply domain migration + run: | + python -m py_compile scripts/apply-portable-transfer-domain.py + python scripts/apply-portable-transfer-domain.py + git diff --check + + - name: Install frontend dependencies + working-directory: frontend + run: npm ci + + - name: Run frontend tests + id: frontend_tests + working-directory: frontend + shell: bash + run: | + set +e + npm test > ../portable-transfer-test.log 2>&1 + status=$? + echo "status=$status" >> "$GITHUB_OUTPUT" + exit 0 + + - name: Write compact failure summary + if: steps.frontend_tests.outputs.status != '0' + run: | + python - <<'PY' + from pathlib import Path + lines = Path('portable-transfer-test.log').read_text(errors='replace').splitlines() + indexes = [i for i, line in enumerate(lines) if line.startswith('not ok ')] + output = [] + for index in indexes[:3]: + output.extend(lines[index:index + 28]) + output.append('---') + Path('portable-transfer-failure-summary.txt').write_text('\n'.join(output) or 'Tests failed without a TAP not-ok block.') + PY + git config user.name "github-actions[bot]" + git config user.email "41898282+github-actions[bot]@users.noreply.github.com" + git add -f portable-transfer-failure-summary.txt + git commit -m "test: capture portable transfer assertion" + git push origin HEAD:architecture/portable-transfer + + - name: Fail when frontend tests fail + if: steps.frontend_tests.outputs.status != '0' + run: exit 1 + + - name: Audit and build frontend + working-directory: frontend + run: | + npm audit --omit=dev --audit-level=high + npm run build + + - name: Remove temporary integration files + run: | + rm scripts/apply-portable-transfer-domain.py + rm .github/workflows/apply-portable-transfer-domain.yml + rm portable-transfer-test.log + rm -f portable-transfer-failure-summary.txt + rm -rf scripts/__pycache__ + + - name: Commit verified core + run: | + git config user.name "github-actions[bot]" + git config user.email "41898282+github-actions[bot]@users.noreply.github.com" + git add -A + git commit -m "feat: add portable transfer archive and import engine" + git push origin HEAD:architecture/portable-transfer diff --git a/frontend/lib/application/browserTransferApplication.mjs b/frontend/lib/application/browserTransferApplication.mjs new file mode 100644 index 0000000..65cddfa --- /dev/null +++ b/frontend/lib/application/browserTransferApplication.mjs @@ -0,0 +1,31 @@ +import { createBrowserCampaignRepository } from "../infrastructure/adapters.mjs"; +import { + createBrowserApprovalRepository, + createBrowserAssetRepository, + createBrowserBlobStorage, + createBrowserExportRepository, + createBrowserSourceArtifactRepository, + createBrowserTransferReportRepository, +} from "../infrastructure/transferAdapters.mjs"; +import { createTransferApplication } from "../transfer/transferApplication.mjs"; + +export function createBrowserTransferApplication({ + getStorage, + campaignKey = "signalflow_recovery_library", + clock, + idService, + signer = null, +} = {}) { + return createTransferApplication({ + campaignRepository: createBrowserCampaignRepository({ getStorage, key: campaignKey }), + assetRepository: createBrowserAssetRepository({ getStorage }), + sourceArtifactRepository: createBrowserSourceArtifactRepository({ getStorage }), + approvalRepository: createBrowserApprovalRepository({ getStorage }), + exportRepository: createBrowserExportRepository({ getStorage }), + blobStorage: createBrowserBlobStorage({ getStorage }), + transferReportRepository: createBrowserTransferReportRepository({ getStorage }), + clock, + idService, + signer, + }); +} diff --git a/frontend/lib/domain/contracts.mjs b/frontend/lib/domain/contracts.mjs index 031c66d..f900094 100644 --- a/frontend/lib/domain/contracts.mjs +++ b/frontend/lib/domain/contracts.mjs @@ -17,6 +17,7 @@ export const DOMAIN_KINDS = Object.freeze({ CONNECTION: "Connection", USAGE_EVENT: "UsageEvent", AUDIT_EVENT: "AuditEvent", + TRANSFER_REPORT: "TransferReport", }); export const DOMAIN_CONTRACTS = Object.freeze({ @@ -36,6 +37,7 @@ export const DOMAIN_CONTRACTS = Object.freeze({ Connection: { idField: "connectionId", owner: "workspace", required: ["connectionId", "provider", "status"] }, UsageEvent: { idField: "usageEventId", owner: "workspace", required: ["usageEventId", "eventType"] }, AuditEvent: { idField: "auditEventId", owner: "workspace", required: ["auditEventId", "eventType"] }, + TransferReport: { idField: "transferReportId", owner: "workspace", required: ["transferReportId", "archiveId", "status"] }, }); const FORBIDDEN_FIELD = /(api[_-]?key|access[_-]?token|refresh[_-]?token|oauth[_-]?token|client[_-]?secret|password|authorization|cookie|database|dbclient|request|response)/i; diff --git a/frontend/lib/domain/ports.mjs b/frontend/lib/domain/ports.mjs index 4e88387..a73900e 100644 --- a/frontend/lib/domain/ports.mjs +++ b/frontend/lib/domain/ports.mjs @@ -1,7 +1,12 @@ export const PORT_CONTRACTS = Object.freeze({ campaignRepository: ["list", "get", "upsert", "remove"], assetRepository: ["list", "get", "upsert", "remove"], + sourceArtifactRepository: ["list", "get", "upsert", "remove"], + approvalRepository: ["list", "get", "upsert", "remove"], + exportRepository: ["list", "get", "upsert", "remove"], blobStorage: ["put", "get", "remove"], + transferReportRepository: ["list", "get", "upsert", "remove"], + archiveSigner: ["sign", "verify", "describe"], jobQueue: ["enqueue", "get", "cancel"], providerAdapter: ["test", "generate"], connectorAdapter: ["status", "publish"], diff --git a/frontend/lib/infrastructure/transferAdapters.mjs b/frontend/lib/infrastructure/transferAdapters.mjs new file mode 100644 index 0000000..1d59824 --- /dev/null +++ b/frontend/lib/infrastructure/transferAdapters.mjs @@ -0,0 +1,263 @@ +import { assertPort } from "../domain/ports.mjs"; +import { createDomainRecord, parseDomainRecord, portableClone, stableStringify } from "../domain/contracts.mjs"; + +const RECORD_SPECS = Object.freeze({ + asset: { portName: "assetRepository", kind: "Asset", idField: "assetId" }, + sourceArtifact: { portName: "sourceArtifactRepository", kind: "SourceArtifact", idField: "sourceArtifactId" }, + approval: { portName: "approvalRepository", kind: "Approval", idField: "approvalId" }, + export: { portName: "exportRepository", kind: "Export", idField: "exportId" }, + transferReport: { portName: "transferReportRepository", kind: "TransferReport", idField: "transferReportId" }, +}); + +function clone(value) { + return portableClone(value); +} + +function sortByUpdated(items) { + return [...items].sort((left, right) => String(right.updatedAt || right.createdAt || "").localeCompare(String(left.updatedAt || left.createdAt || ""))); +} + +function spec(name) { + const value = RECORD_SPECS[name]; + if (!value) throw new TypeError(`Unknown transfer record repository: ${name}.`); + return value; +} + +function normalizeRecord(kind, idField, value) { + if (value?.kind === kind) return parseDomainRecord(value, kind); + return createDomainRecord(kind, { ...value, [idField]: value?.[idField] }); +} + +function createMemoryRecordRepository({ portName, kind, idField, initial = [] }) { + const records = new Map(); + for (const item of initial) { + const normalized = normalizeRecord(kind, idField, item); + records.set(normalized[idField], normalized); + } + return assertPort(portName, { + async list() { + return sortByUpdated(Array.from(records.values())).map(clone); + }, + async get(id) { + return records.has(id) ? clone(records.get(id)) : null; + }, + async upsert(value) { + const normalized = normalizeRecord(kind, idField, value); + records.set(normalized[idField], normalized); + return clone(normalized); + }, + async remove(id) { + return records.delete(id); + }, + }); +} + +function createStoreBackedRecordRepository({ store, portName, kind, idField, prefix }) { + if (!store || ["list", "get", "set", "remove"].some((method) => typeof store[method] !== "function")) { + throw new TypeError(`${kind} store requires list/get/set/remove methods.`); + } + const keyFor = (id) => `${prefix}${id}`; + return assertPort(portName, { + async list() { + const keys = await store.list(prefix); + const values = await Promise.all(keys.map((key) => store.get(key))); + return sortByUpdated(values.filter(Boolean).map((value) => normalizeRecord(kind, idField, value))).map(clone); + }, + async get(id) { + const value = await store.get(keyFor(id)); + return value ? clone(normalizeRecord(kind, idField, value)) : null; + }, + async upsert(value) { + const normalized = normalizeRecord(kind, idField, value); + await store.set(keyFor(normalized[idField]), normalized); + return clone(normalized); + }, + async remove(id) { + return Boolean(await store.remove(keyFor(id))); + }, + }); +} + +function createBrowserRecordRepository({ getStorage, key, limit, portName, kind, idField }) { + if (typeof getStorage !== "function") throw new TypeError(`${kind} browser repository requires getStorage().`); + const storage = () => { + const target = getStorage(); + if (!target || typeof target.getItem !== "function" || typeof target.setItem !== "function") { + throw new TypeError("Browser storage is unavailable."); + } + return target; + }; + function read() { + const raw = storage().getItem(key); + if (!raw) return []; + const parsed = JSON.parse(raw); + return Array.isArray(parsed) ? parsed : []; + } + function write(items) { + const normalized = sortByUpdated(items).slice(0, limit); + storage().setItem(key, stableStringify(normalized)); + return normalized; + } + async function list() { + return sortByUpdated(read().map((value) => normalizeRecord(kind, idField, value))).map(clone); + } + async function get(id) { + const items = await list(); + return items.find((item) => item[idField] === id) || null; + } + async function upsert(value) { + const normalized = normalizeRecord(kind, idField, value); + const items = await list(); + write([normalized, ...items.filter((item) => item[idField] !== normalized[idField])]); + return clone(normalized); + } + async function remove(id) { + const items = await list(); + const next = items.filter((item) => item[idField] !== id); + write(next); + return next.length !== items.length; + } + return assertPort(portName, { list, get, upsert, remove }); +} + +function memoryRepository(name, initial = []) { + return createMemoryRecordRepository({ ...spec(name), initial }); +} + +function storeRepository(name, { store, prefix }) { + return createStoreBackedRecordRepository({ ...spec(name), store, prefix }); +} + +function browserRepository(name, { getStorage, key, limit }) { + return createBrowserRecordRepository({ ...spec(name), getStorage, key, limit }); +} + +export function createMemoryAssetRepository(initial = []) { + return memoryRepository("asset", initial); +} + +export function createStoreBackedAssetRepository({ store, prefix = "asset/" } = {}) { + return storeRepository("asset", { store, prefix }); +} + +export function createBrowserAssetRepository({ getStorage, key = "signalflow_assets_v1", limit = 250 } = {}) { + return browserRepository("asset", { getStorage, key, limit }); +} + +export function createMemorySourceArtifactRepository(initial = []) { + return memoryRepository("sourceArtifact", initial); +} + +export function createStoreBackedSourceArtifactRepository({ store, prefix = "source-artifact/" } = {}) { + return storeRepository("sourceArtifact", { store, prefix }); +} + +export function createBrowserSourceArtifactRepository({ getStorage, key = "signalflow_source_artifacts_v1", limit = 500 } = {}) { + return browserRepository("sourceArtifact", { getStorage, key, limit }); +} + +export function createMemoryApprovalRepository(initial = []) { + return memoryRepository("approval", initial); +} + +export function createStoreBackedApprovalRepository({ store, prefix = "approval/" } = {}) { + return storeRepository("approval", { store, prefix }); +} + +export function createBrowserApprovalRepository({ getStorage, key = "signalflow_approvals_v1", limit = 500 } = {}) { + return browserRepository("approval", { getStorage, key, limit }); +} + +export function createMemoryExportRepository(initial = []) { + return memoryRepository("export", initial); +} + +export function createStoreBackedExportRepository({ store, prefix = "export/" } = {}) { + return storeRepository("export", { store, prefix }); +} + +export function createBrowserExportRepository({ getStorage, key = "signalflow_exports_v1", limit = 500 } = {}) { + return browserRepository("export", { getStorage, key, limit }); +} + +export function createMemoryTransferReportRepository(initial = []) { + return memoryRepository("transferReport", initial); +} + +export function createStoreBackedTransferReportRepository({ store, prefix = "transfer-report/" } = {}) { + return storeRepository("transferReport", { store, prefix }); +} + +export function createBrowserTransferReportRepository({ getStorage, key = "signalflow_transfer_reports_v1", limit = 50 } = {}) { + return browserRepository("transferReport", { getStorage, key, limit }); +} + +function bytesToBase64(bytes) { + if (typeof Buffer !== "undefined") return Buffer.from(bytes).toString("base64"); + let binary = ""; + for (let index = 0; index < bytes.length; index += 0x8000) { + binary += String.fromCharCode(...bytes.subarray(index, index + 0x8000)); + } + return btoa(binary); +} + +function base64ToBytes(value) { + if (typeof Buffer !== "undefined") return new Uint8Array(Buffer.from(String(value), "base64")); + const binary = atob(String(value)); + const bytes = new Uint8Array(binary.length); + for (let index = 0; index < binary.length; index += 1) bytes[index] = binary.charCodeAt(index); + return bytes; +} + +function encodeBrowserBlob(value) { + if (value instanceof Uint8Array) return { type: "bytes", value: bytesToBase64(value) }; + if (value instanceof ArrayBuffer) return { type: "bytes", value: bytesToBase64(new Uint8Array(value)) }; + if (typeof value === "string") return { type: "text", value }; + return { type: "json", value: portableClone(value) }; +} + +function decodeBrowserBlob(value) { + if (!value) return null; + if (value.type === "bytes") return base64ToBytes(value.value); + if (value.type === "text") return String(value.value || ""); + if (value.type === "json") return portableClone(value.value); + throw new TypeError(`Unsupported browser blob type: ${value.type || "missing"}.`); +} + +export function createBrowserBlobStorage({ getStorage, key = "signalflow_blobs_v1" } = {}) { + if (typeof getStorage !== "function") throw new TypeError("Browser blob storage requires getStorage()."); + const storage = () => { + const target = getStorage(); + if (!target || typeof target.getItem !== "function" || typeof target.setItem !== "function") { + throw new TypeError("Browser storage is unavailable."); + } + return target; + }; + function read() { + const raw = storage().getItem(key); + if (!raw) return {}; + const parsed = JSON.parse(raw); + return parsed && typeof parsed === "object" && !Array.isArray(parsed) ? parsed : {}; + } + function write(value) { + storage().setItem(key, JSON.stringify(value)); + } + return assertPort("blobStorage", { + async put(blobId, value) { + const records = read(); + records[blobId] = encodeBrowserBlob(value); + write(records); + return { blobId }; + }, + async get(blobId) { + return decodeBrowserBlob(read()[blobId]); + }, + async remove(blobId) { + const records = read(); + const existed = Object.prototype.hasOwnProperty.call(records, blobId); + delete records[blobId]; + write(records); + return existed; + }, + }); +} diff --git a/frontend/lib/transfer/portableArchive.mjs b/frontend/lib/transfer/portableArchive.mjs new file mode 100644 index 0000000..1bdfae3 --- /dev/null +++ b/frontend/lib/transfer/portableArchive.mjs @@ -0,0 +1,338 @@ +import { portableClone, stableStringify } from "../domain/contracts.mjs"; + +export const PORTABLE_ARCHIVE_SCHEMA_VERSION = 1; +export const PORTABLE_ARCHIVE_KIND = "SignalFlowPortableArchive"; +export const DEFAULT_MAX_ARCHIVE_BYTES = 50 * 1024 * 1024; + +const SECRET_FIELD = /(api[_-]?key|access[_-]?token|refresh[_-]?token|oauth|secret|password|authorization|cookie|private[_-]?key|session[_-]?key)/i; +const PRIVATE_REFERENCE_FIELD = /(signed[_-]?url|private[_-]?(url|ref)|provider[_-]?base[_-]?url|base[_-]?url|local[_-]?(path|url)|absolute[_-]?path|filesystem[_-]?path)/i; +const PATH_FIELD = /(^|[_-])(file|folder|directory|filesystem)?path$/i; +const WINDOWS_PATH = /^[a-z]:[\\/]/i; +const POSIX_PRIVATE_PATH = /^\/(?:Users|home|root|var\/private|private|mnt|Volumes)\//i; +const FILE_URL = /^file:\/\//i; +const PRIVATE_HOST_URL = /^https?:\/\/(?:localhost|127\.0\.0\.1|0\.0\.0\.0|\[::1\]|10\.|192\.168\.|172\.(?:1[6-9]|2\d|3[01])\.)/i; + +function text(value) { + return String(value ?? "").trim(); +} + +function encoder() { + return new TextEncoder(); +} + +function cryptoSubtle() { + const subtle = globalThis.crypto?.subtle; + if (!subtle) throw new Error("Web Crypto is required for portable archive integrity."); + return subtle; +} + +function bytesToHex(bytes) { + return Array.from(bytes, (byte) => byte.toString(16).padStart(2, "0")).join(""); +} + +function bytesToBase64(bytes) { + if (typeof Buffer !== "undefined") return Buffer.from(bytes).toString("base64"); + let binary = ""; + const chunkSize = 0x8000; + for (let index = 0; index < bytes.length; index += chunkSize) { + binary += String.fromCharCode(...bytes.subarray(index, index + chunkSize)); + } + return btoa(binary); +} + +function base64ToBytes(value) { + if (typeof Buffer !== "undefined") return new Uint8Array(Buffer.from(String(value), "base64")); + const binary = atob(String(value)); + const bytes = new Uint8Array(binary.length); + for (let index = 0; index < binary.length; index += 1) bytes[index] = binary.charCodeAt(index); + return bytes; +} + +export async function sha256Hex(value) { + const bytes = value instanceof Uint8Array ? value : encoder().encode(String(value)); + const digest = await cryptoSubtle().digest("SHA-256", bytes); + return bytesToHex(new Uint8Array(digest)); +} + +function privateStringReason(value, key = "") { + const candidate = text(value); + if (!candidate) return ""; + if (SECRET_FIELD.test(key)) return "secret field"; + if (PRIVATE_REFERENCE_FIELD.test(key)) return "private deployment reference"; + if (PATH_FIELD.test(key) && (WINDOWS_PATH.test(candidate) || POSIX_PRIVATE_PATH.test(candidate) || FILE_URL.test(candidate))) { + return "local filesystem path"; + } + if (FILE_URL.test(candidate)) return "local file URL"; + if (/base[_-]?url|endpoint/i.test(key) && PRIVATE_HOST_URL.test(candidate)) return "private/local endpoint"; + return ""; +} + +export function sanitizeForPortableTransfer(value, { path = "payload", exclusions = [] } = {}) { + function visit(current, currentPath, key = "") { + if (current === null || typeof current === "boolean" || typeof current === "number") return current; + if (typeof current === "string") { + const reason = privateStringReason(current, key); + if (reason) { + exclusions.push({ path: currentPath, reason }); + return undefined; + } + return current; + } + if (current === undefined) return undefined; + if (Array.isArray(current)) { + return current + .map((item, index) => visit(item, `${currentPath}[${index}]`)) + .filter((item) => item !== undefined); + } + if (!current || typeof current !== "object" || Object.getPrototypeOf(current) !== Object.prototype) { + throw new TypeError(`${currentPath} contains a non-portable runtime object.`); + } + + const result = {}; + for (const [childKey, childValue] of Object.entries(current)) { + const childPath = `${currentPath}.${childKey}`; + if (SECRET_FIELD.test(childKey)) { + exclusions.push({ path: childPath, reason: "secret field" }); + continue; + } + if (PRIVATE_REFERENCE_FIELD.test(childKey)) { + exclusions.push({ path: childPath, reason: "private deployment reference" }); + continue; + } + const next = visit(childValue, childPath, childKey); + if (next !== undefined) result[childKey] = next; + } + return result; + } + + return visit(value, path); +} + +export function encodeBlobPayload(value) { + if (value instanceof Uint8Array) { + return { payloadFormat: "bytes", payloadBase64: bytesToBase64(value), byteLength: value.byteLength }; + } + if (value instanceof ArrayBuffer) { + const bytes = new Uint8Array(value); + return { payloadFormat: "bytes", payloadBase64: bytesToBase64(bytes), byteLength: bytes.byteLength }; + } + if (typeof value === "string") { + const bytes = encoder().encode(value); + return { payloadFormat: "text", payloadBase64: bytesToBase64(bytes), byteLength: bytes.byteLength }; + } + const serialized = stableStringify(portableClone(value)); + const bytes = encoder().encode(serialized); + return { payloadFormat: "json", payloadBase64: bytesToBase64(bytes), byteLength: bytes.byteLength }; +} + +export function decodeBlobPayload(entry) { + const bytes = base64ToBytes(entry.payloadBase64 || ""); + if (entry.payloadFormat === "bytes") return bytes; + const decoded = new TextDecoder().decode(bytes); + if (entry.payloadFormat === "text") return decoded; + if (entry.payloadFormat === "json") return JSON.parse(decoded); + throw new TypeError(`Unsupported blob payload format: ${entry.payloadFormat || "missing"}.`); +} + +export function validateArchivePath(value) { + const archivePath = text(value); + if (!archivePath || archivePath.startsWith("/") || archivePath.includes("\\") || archivePath.includes("\0")) return false; + const segments = archivePath.split("/"); + return segments[0] === "blobs" && segments.every((segment) => segment && segment !== "." && segment !== ".."); +} + +function unsignedArchive(archive) { + const { integrity, signature, ...unsigned } = archive; + void integrity; + void signature; + return unsigned; +} + +export async function createPortableArchive({ + archiveId, + createdAt, + sourceDeployment = {}, + campaigns = [], + assets = [], + sourceArtifacts = [], + approvals = [], + exports = [], + blobEntries = [], + signer = null, +} = {}) { + if (!text(archiveId)) throw new TypeError("archiveId is required."); + if (!text(createdAt)) throw new TypeError("createdAt is required."); + + const exclusions = []; + const payload = sanitizeForPortableTransfer({ + campaigns, + assets, + sourceArtifacts, + approvals, + exports, + blobEntries, + }, { path: "payload", exclusions }); + const sanitizedSource = sanitizeForPortableTransfer(sourceDeployment, { path: "sourceDeployment", exclusions }); + const estimatedAssetBytes = (payload.blobEntries || []).reduce((total, entry) => total + (Number(entry.byteLength) || 0), 0); + + const archive = { + schemaVersion: PORTABLE_ARCHIVE_SCHEMA_VERSION, + kind: PORTABLE_ARCHIVE_KIND, + archiveId: text(archiveId), + createdAt: text(createdAt), + sourceDeployment: sanitizedSource, + manifest: { + campaignCount: payload.campaigns.length, + assetCount: payload.assets.length, + sourceArtifactCount: payload.sourceArtifacts.length, + approvalCount: payload.approvals.length, + exportCount: payload.exports.length, + blobCount: payload.blobEntries.length, + estimatedAssetBytes, + exclusions, + }, + payload, + }; + + const digest = await sha256Hex(stableStringify(archive)); + archive.integrity = { algorithm: "SHA-256", digest }; + if (signer) { + const signature = await signer.sign(digest); + archive.signature = { ...signer.describe(), value: signature }; + } + return portableClone(archive); +} + +export async function validatePortableArchive(archive, { + maxBytes = DEFAULT_MAX_ARCHIVE_BYTES, + signer = null, + requireSignature = false, +} = {}) { + const errors = []; + const warnings = []; + if (!archive || typeof archive !== "object" || Array.isArray(archive)) { + return { valid: false, blocked: true, errors: [{ code: "invalid_archive", message: "Archive must be an object." }], warnings }; + } + if (archive.kind !== PORTABLE_ARCHIVE_KIND) { + errors.push({ code: "invalid_kind", message: `Expected ${PORTABLE_ARCHIVE_KIND}.` }); + } + if (!Number.isInteger(archive.schemaVersion)) { + errors.push({ code: "missing_schema", message: "Archive schema version is missing." }); + } else if (archive.schemaVersion > PORTABLE_ARCHIVE_SCHEMA_VERSION) { + errors.push({ + code: "future_schema", + message: `Archive schema ${archive.schemaVersion} is newer than supported schema ${PORTABLE_ARCHIVE_SCHEMA_VERSION}. Upgrade SignalFlow before importing.`, + }); + } else if (archive.schemaVersion < PORTABLE_ARCHIVE_SCHEMA_VERSION) { + warnings.push({ code: "legacy_archive", message: `Archive schema ${archive.schemaVersion} requires compatibility migration.` }); + } + + const serialized = stableStringify(archive); + const byteLength = encoder().encode(serialized).byteLength; + if (byteLength > maxBytes) { + errors.push({ code: "archive_too_large", message: `Archive is ${byteLength} bytes; the configured limit is ${maxBytes} bytes.` }); + } + + const blobEntries = Array.isArray(archive.payload?.blobEntries) ? archive.payload.blobEntries : []; + for (const entry of blobEntries) { + if (!validateArchivePath(entry.archivePath)) { + errors.push({ code: "archive_traversal", message: `Unsafe archive path: ${entry.archivePath || "missing"}.` }); + } + try { + const bytes = base64ToBytes(entry.payloadBase64 || ""); + if (Number(entry.byteLength) !== bytes.byteLength) { + errors.push({ code: "blob_length_mismatch", message: `Blob ${entry.blobId || entry.archivePath} length does not match its manifest.` }); + } + } catch { + errors.push({ code: "invalid_blob_encoding", message: `Blob ${entry.blobId || entry.archivePath} is not valid base64.` }); + } + } + + if (!archive.integrity?.digest || archive.integrity.algorithm !== "SHA-256") { + errors.push({ code: "missing_integrity", message: "SHA-256 integrity metadata is required." }); + } else { + const expected = await sha256Hex(stableStringify(unsignedArchive(archive))); + if (expected !== archive.integrity.digest) { + errors.push({ code: "integrity_mismatch", message: "Archive content does not match its SHA-256 digest." }); + } + } + + if (requireSignature && !archive.signature) { + errors.push({ code: "signature_required", message: "This destination requires a signed archive." }); + } + if (archive.signature) { + if (!signer) { + warnings.push({ code: "signature_unverified", message: "Archive is signed, but no matching verifier was configured." }); + if (requireSignature) errors.push({ code: "signature_unverified", message: "The required archive signature could not be verified." }); + } else { + const validSignature = await signer.verify(archive.integrity?.digest || "", archive.signature.value); + if (!validSignature) errors.push({ code: "invalid_signature", message: "Archive signature verification failed." }); + } + } + + const missingBlobAssets = (archive.payload?.assets || []).filter((asset) => { + if (!asset.blobId) return false; + return !blobEntries.some((entry) => entry.blobId === asset.blobId); + }); + if (missingBlobAssets.length) { + warnings.push({ + code: "partial_assets", + message: `${missingBlobAssets.length} asset${missingBlobAssets.length === 1 ? " is" : "s are"} missing payload data and can import as metadata only.`, + assetIds: missingBlobAssets.map((asset) => asset.assetId), + }); + } + + return { + valid: errors.length === 0, + blocked: errors.length > 0, + errors, + warnings, + byteLength, + counts: { + campaigns: archive.payload?.campaigns?.length || 0, + assets: archive.payload?.assets?.length || 0, + blobs: blobEntries.length, + sourceArtifacts: archive.payload?.sourceArtifacts?.length || 0, + approvals: archive.payload?.approvals?.length || 0, + exports: archive.payload?.exports?.length || 0, + }, + estimatedAssetBytes: archive.manifest?.estimatedAssetBytes || 0, + }; +} + +export function createHmacArchiveSigner({ secret, keyId = "signalflow-transfer" } = {}) { + const secretBytes = encoder().encode(text(secret)); + if (!secretBytes.length) throw new TypeError("HMAC signer requires a secret."); + + async function key() { + return cryptoSubtle().importKey( + "raw", + secretBytes, + { name: "HMAC", hash: "SHA-256" }, + false, + ["sign", "verify"], + ); + } + + return { + async sign(digest) { + const signature = await cryptoSubtle().sign("HMAC", await key(), encoder().encode(String(digest))); + return bytesToBase64(new Uint8Array(signature)); + }, + async verify(digest, signature) { + try { + return cryptoSubtle().verify( + "HMAC", + await key(), + base64ToBytes(signature), + encoder().encode(String(digest)), + ); + } catch { + return false; + } + }, + describe() { + return { algorithm: "HMAC-SHA-256", keyId: text(keyId) }; + }, + }; +} diff --git a/frontend/lib/transfer/transferApplication.mjs b/frontend/lib/transfer/transferApplication.mjs new file mode 100644 index 0000000..9492da3 --- /dev/null +++ b/frontend/lib/transfer/transferApplication.mjs @@ -0,0 +1,746 @@ +import { + assertPort, + createSystemClock, + createSystemIdService, +} from "../domain/ports.mjs"; +import { + createDomainRecord, + portableClone, + stableStringify, +} from "../domain/contracts.mjs"; +import { + campaignToEditorState, + createCampaignAggregate, + migrateLegacyCampaign, +} from "../domain/campaign.mjs"; +import { + createPortableArchive, + decodeBlobPayload, + encodeBlobPayload, + sha256Hex, + validatePortableArchive, +} from "./portableArchive.mjs"; + +export const TRANSFER_CONFLICT_POLICIES = Object.freeze({ + SKIP: "skip", + COPY: "copy", + REPLACE: "replace", +}); + +export const TRANSFER_STATUSES = Object.freeze({ + PREPARING: "preparing", + VALIDATING: "validating", + WARNINGS_FOUND: "warnings_found", + BLOCKED: "blocked", + SELECTING_DESTINATION: "selecting_destination", + UPLOADING: "uploading", + IMPORTING: "importing", + PARTIALLY_IMPORTED: "partially_imported", + COMPLETE: "complete", + CANCELLED: "cancelled", + FAILED: "failed", + ROLLED_BACK: "rolled_back", +}); + +const RECORD_CONFIG = Object.freeze({ + campaign: { repository: "campaignRepository", idField: "campaignId", kind: "Campaign" }, + asset: { repository: "assetRepository", idField: "assetId", kind: "Asset" }, + sourceArtifact: { repository: "sourceArtifactRepository", idField: "sourceArtifactId", kind: "SourceArtifact" }, + approval: { repository: "approvalRepository", idField: "approvalId", kind: "Approval" }, + export: { repository: "exportRepository", idField: "exportId", kind: "Export" }, +}); + +function text(value) { + return String(value ?? "").trim(); +} + +function safeArchiveSegment(value) { + const safe = text(value).replace(/[^a-zA-Z0-9._-]+/g, "-").replace(/^[-.]+|[-.]+$/g, ""); + return safe || "blob"; +} + +function unique(values) { + return Array.from(new Set(values.filter(Boolean))); +} + +function approvedDraftRecords(campaigns, clock) { + const records = []; + for (const campaign of campaigns) { + for (const [channel, draft] of Object.entries(campaign.drafts || {})) { + if (!draft?.approved) continue; + records.push(createDomainRecord("Approval", { + approvalId: `approval-${campaign.campaignId}-${channel}-${draft.current?.revisionId || "current"}`, + campaignId: campaign.campaignId, + draftId: draft.draftId, + channel, + revisionId: draft.current?.revisionId || null, + status: "approved", + historical: true, + createdAt: draft.updatedAt || campaign.updatedAt || clock.now(), + })); + } + } + return records; +} + +async function selectRecords(repository, ids = []) { + if (!ids.length) return repository.list(); + const records = await Promise.all(ids.map((id) => repository.get(id))); + return records.filter(Boolean); +} + +function recordKey(kind, id) { + return `${kind}:${id}`; +} + +function provenanceMatch(record, archive, sourceId, field) { + const provenance = record?.transferProvenance; + return Boolean( + provenance + && provenance.archiveId === archive.archiveId + && provenance[field] === sourceId, + ); +} + +function sourceFieldFor(kind) { + return `source${kind[0].toUpperCase()}${kind.slice(1)}Id`; +} + +function conflictTarget(existing, archive, source, kind, idField) { + const sourceId = source[idField]; + const provenanceField = sourceFieldFor(kind); + return existing.find((record) => provenanceMatch(record, archive, sourceId, provenanceField)) + || existing.find((record) => record[idField] === sourceId) + || null; +} + +function conflictSummary({ existing, archive, sourceRecords, kind, idField }) { + return sourceRecords.map((source) => { + const target = conflictTarget(existing, archive, source, kind, idField); + if (!target) return null; + const sameImport = provenanceMatch(target, archive, source[idField], sourceFieldFor(kind)); + return { + kind, + sourceId: source[idField], + targetId: target[idField], + type: sameImport ? "already_imported" : "id_collision", + availablePolicies: Object.values(TRANSFER_CONFLICT_POLICIES), + recommendedPolicy: "skip", + }; + }).filter(Boolean); +} + +function transferProvenance({ archive, sourceId, kind, importedAt, destinationWorkspaceId }) { + return { + archiveId: archive.archiveId, + archiveSchemaVersion: archive.schemaVersion, + sourceDeployment: portableClone(archive.sourceDeployment || {}), + [sourceFieldFor(kind)]: sourceId, + importedAt, + destinationWorkspaceId, + historical: true, + }; +} + +function reportRecord({ + transferReportId, + archive, + status, + destinationWorkspaceId = null, + conflictPolicy = "skip", + startedAt, + updatedAt, + completedAt = null, + validation = null, + summary = {}, + items = [], + journal = [], + warnings = [], + errors = [], + rollback = null, +} = {}) { + return createDomainRecord("TransferReport", { + transferReportId, + archiveId: archive.archiveId, + archiveDigest: archive.integrity?.digest || null, + archiveSchemaVersion: archive.schemaVersion, + status, + destinationWorkspaceId, + conflictPolicy, + startedAt, + updatedAt, + completedAt, + validation: validation ? portableClone(validation) : null, + summary: portableClone(summary), + items: portableClone(items), + journal: portableClone(journal), + warnings: portableClone(warnings), + errors: portableClone(errors), + rollback: rollback ? portableClone(rollback) : null, + }); +} + +function importedRecordId({ source, target, policy, idField, idService, kind }) { + if (!target) return source[idField]; + if (policy === TRANSFER_CONFLICT_POLICIES.COPY) return idService.create(kind); + return target[idField]; +} + +function updateReferences(value, idMaps) { + if (Array.isArray(value)) return value.map((item) => updateReferences(item, idMaps)); + if (!value || typeof value !== "object") return value; + const result = {}; + for (const [key, item] of Object.entries(value)) { + if (typeof item === "string") { + if (key === "campaignId" && idMaps.campaign.has(item)) result[key] = idMaps.campaign.get(item); + else if (key === "assetId" && idMaps.asset.has(item)) result[key] = idMaps.asset.get(item); + else if (key === "sourceArtifactId" && idMaps.sourceArtifact.has(item)) result[key] = idMaps.sourceArtifact.get(item); + else if (key === "draftId" && idMaps.draft.has(item)) result[key] = idMaps.draft.get(item); + else if (key === "blobId" && idMaps.blob.has(item)) result[key] = idMaps.blob.get(item); + else result[key] = item; + } else { + result[key] = updateReferences(item, idMaps); + } + } + return result; +} + +function rebuildCampaign({ source, targetId, destinationWorkspaceId, provenance, idMaps }) { + const canonical = migrateLegacyCampaign(source); + const editor = campaignToEditorState(canonical); + const campaign = createCampaignAggregate({ + ...editor, + campaignId: targetId, + workspaceId: destinationWorkspaceId, + projectId: canonical.projectId, + title: canonical.title, + status: canonical.status, + createdAt: canonical.createdAt, + updatedAt: canonical.updatedAt, + existingDrafts: canonical.drafts, + existingArchives: canonical.archives, + transferProvenance: provenance, + }); + for (const [channel, originalDraft] of Object.entries(canonical.drafts || {})) { + const importedDraft = campaign.drafts?.[channel]; + if (originalDraft?.draftId && importedDraft?.draftId) idMaps.draft.set(originalDraft.draftId, importedDraft.draftId); + } + return campaign; +} + +function buildMetadataRecord({ source, kind, idField, targetId, destinationWorkspaceId, provenance, idMaps }) { + const mapped = updateReferences(source, idMaps); + return createDomainRecord(kind, { + ...mapped, + [idField]: targetId, + workspaceId: destinationWorkspaceId, + transferProvenance: provenance, + importedHistoricalRecord: true, + }); +} + +async function restoreJournalEntry(entry, repositories, blobStorage) { + if (entry.kind === "blob") { + if (entry.previous === null) await blobStorage.remove(entry.id); + else await blobStorage.put(entry.id, decodeBlobPayload(entry.previous)); + return; + } + const config = RECORD_CONFIG[entry.kind]; + const repository = repositories[config.repository]; + if (entry.previous === null) await repository.remove(entry.id); + else await repository.upsert(entry.previous); +} + +export function createTransferApplication({ + campaignRepository, + assetRepository, + sourceArtifactRepository, + approvalRepository, + exportRepository, + blobStorage, + transferReportRepository, + signer = null, + clock = createSystemClock(), + idService = createSystemIdService("signalflow-transfer"), +} = {}) { + const repositories = { + campaignRepository: assertPort("campaignRepository", campaignRepository), + assetRepository: assertPort("assetRepository", assetRepository), + sourceArtifactRepository: assertPort("sourceArtifactRepository", sourceArtifactRepository), + approvalRepository: assertPort("approvalRepository", approvalRepository), + exportRepository: assertPort("exportRepository", exportRepository), + }; + const blobs = assertPort("blobStorage", blobStorage); + const reports = assertPort("transferReportRepository", transferReportRepository); + const applicationClock = assertPort("clock", clock); + const applicationIds = assertPort("idService", idService); + if (signer) assertPort("archiveSigner", signer); + + async function exportSelection({ + campaignIds = [], + assetIds = [], + sourceArtifactIds = [], + approvalIds = [], + exportIds = [], + sourceDeployment = {}, + } = {}) { + const createdAt = applicationClock.now(); + const campaigns = (await selectRecords(repositories.campaignRepository, campaignIds)).map(migrateLegacyCampaign); + const assets = await selectRecords(repositories.assetRepository, assetIds); + const sourceArtifacts = await selectRecords(repositories.sourceArtifactRepository, sourceArtifactIds); + const explicitApprovals = await selectRecords(repositories.approvalRepository, approvalIds); + const exports = await selectRecords(repositories.exportRepository, exportIds); + const derivedApprovals = approvedDraftRecords(campaigns, applicationClock); + const approvals = [...explicitApprovals]; + const approvalIdsSeen = new Set(approvals.map((approval) => approval.approvalId)); + for (const approval of derivedApprovals) { + if (!approvalIdsSeen.has(approval.approvalId)) approvals.push(approval); + } + + const blobEntries = []; + for (const asset of assets) { + if (!asset.blobId) continue; + const value = await blobs.get(asset.blobId); + if (value === null || value === undefined) continue; + const encoded = encodeBlobPayload(value); + blobEntries.push({ + blobId: asset.blobId, + assetId: asset.assetId, + archivePath: `blobs/${safeArchiveSegment(asset.blobId)}.bin`, + contentType: asset.contentType || "application/octet-stream", + ...encoded, + }); + } + + return createPortableArchive({ + archiveId: applicationIds.create("archive"), + createdAt, + sourceDeployment, + campaigns, + assets, + sourceArtifacts, + approvals, + exports, + blobEntries, + signer, + }); + } + + async function previewImport(archive, { + destinationWorkspaceId = "", + conflictPolicy = TRANSFER_CONFLICT_POLICIES.SKIP, + maxBytes, + requireSignature = false, + } = {}) { + const validation = await validatePortableArchive(archive, { maxBytes, signer, requireSignature }); + if (validation.blocked) { + return { + status: TRANSFER_STATUSES.BLOCKED, + validation, + destinationWorkspaceId: destinationWorkspaceId || null, + conflictPolicy, + conflicts: [], + exclusions: archive?.manifest?.exclusions || [], + }; + } + if (!text(destinationWorkspaceId)) { + return { + status: TRANSFER_STATUSES.SELECTING_DESTINATION, + validation, + destinationWorkspaceId: null, + conflictPolicy, + conflicts: [], + exclusions: archive.manifest?.exclusions || [], + }; + } + if (!Object.values(TRANSFER_CONFLICT_POLICIES).includes(conflictPolicy)) { + throw new TypeError(`Unsupported transfer conflict policy: ${conflictPolicy}.`); + } + + const existing = { + campaign: await repositories.campaignRepository.list(), + asset: await repositories.assetRepository.list(), + sourceArtifact: await repositories.sourceArtifactRepository.list(), + approval: await repositories.approvalRepository.list(), + export: await repositories.exportRepository.list(), + }; + const source = { + campaign: archive.payload?.campaigns || [], + asset: archive.payload?.assets || [], + sourceArtifact: archive.payload?.sourceArtifacts || [], + approval: archive.payload?.approvals || [], + export: archive.payload?.exports || [], + }; + const conflicts = Object.entries(RECORD_CONFIG).flatMap(([kind, config]) => conflictSummary({ + existing: existing[kind], + archive, + sourceRecords: source[kind], + kind, + idField: config.idField, + })); + const warnings = [ + ...(validation.warnings || []), + ...((archive.manifest?.exclusions || []).length + ? [{ code: "excluded_private_data", message: `${archive.manifest.exclusions.length} private or unsupported fields were excluded during export.` }] + : []), + ...(conflicts.length + ? [{ code: "conflicts_found", message: `${conflicts.length} existing record conflict${conflicts.length === 1 ? " was" : "s were"} found.` }] + : []), + ]; + + return { + status: warnings.length ? TRANSFER_STATUSES.WARNINGS_FOUND : "ready", + validation, + destinationWorkspaceId, + conflictPolicy, + conflicts, + warnings, + exclusions: archive.manifest?.exclusions || [], + estimatedAssetBytes: validation.estimatedAssetBytes, + counts: validation.counts, + }; + } + + async function persistReport(report) { + return reports.upsert(report); + } + + async function rollbackJournal(journal) { + const errors = []; + for (const entry of [...journal].reverse()) { + try { + await restoreJournalEntry(entry, repositories, blobs); + } catch (error) { + errors.push({ kind: entry.kind, id: entry.id, message: error.message }); + } + } + return { complete: errors.length === 0, errors }; + } + + async function importArchive(archive, { + destinationWorkspaceId, + conflictPolicy = TRANSFER_CONFLICT_POLICIES.SKIP, + atomic = true, + signal = null, + maxBytes, + requireSignature = false, + resumeReportId = null, + } = {}) { + const preview = await previewImport(archive, { + destinationWorkspaceId, + conflictPolicy, + maxBytes, + requireSignature, + }); + const now = applicationClock.now(); + let existingReport = resumeReportId ? await reports.get(resumeReportId) : null; + if (existingReport && existingReport.archiveDigest !== archive.integrity?.digest) { + throw new Error("The resume report belongs to a different archive payload."); + } + const transferReportId = existingReport?.transferReportId || applicationIds.create("report"); + const startedAt = existingReport?.startedAt || now; + let items = portableClone(existingReport?.items || []); + let journal = portableClone(existingReport?.journal || []); + let warnings = unique([...(existingReport?.warnings || []).map((item) => stableStringify(item)), ...(preview.warnings || []).map((item) => stableStringify(item))]) + .map((item) => JSON.parse(item)); + let errors = portableClone(existingReport?.errors || []); + + if ([TRANSFER_STATUSES.BLOCKED, TRANSFER_STATUSES.SELECTING_DESTINATION].includes(preview.status)) { + const blocked = reportRecord({ + transferReportId, + archive, + status: preview.status, + destinationWorkspaceId: destinationWorkspaceId || null, + conflictPolicy, + startedAt, + updatedAt: now, + validation: preview.validation, + summary: { preview }, + items, + journal, + warnings, + errors: preview.validation.errors || [], + }); + return persistReport(blocked); + } + + const completedKeys = new Set(items.filter((item) => ["imported", "replaced", "copied", "skipped"].includes(item.status)).map((item) => item.key)); + const idMaps = { + campaign: new Map(), + asset: new Map(), + sourceArtifact: new Map(), + approval: new Map(), + export: new Map(), + draft: new Map(), + blob: new Map(), + }; + const sourceCollections = { + campaign: archive.payload?.campaigns || [], + asset: archive.payload?.assets || [], + sourceArtifact: archive.payload?.sourceArtifacts || [], + approval: archive.payload?.approvals || [], + export: archive.payload?.exports || [], + }; + const existingCollections = { + campaign: await repositories.campaignRepository.list(), + asset: await repositories.assetRepository.list(), + sourceArtifact: await repositories.sourceArtifactRepository.list(), + approval: await repositories.approvalRepository.list(), + export: await repositories.exportRepository.list(), + }; + + const reportInProgress = () => reportRecord({ + transferReportId, + archive, + status: TRANSFER_STATUSES.IMPORTING, + destinationWorkspaceId, + conflictPolicy, + startedAt, + updatedAt: applicationClock.now(), + validation: preview.validation, + summary: { preview }, + items, + journal, + warnings, + errors, + }); + await persistReport(reportInProgress()); + + async function processRecord(kind, source) { + const config = RECORD_CONFIG[kind]; + const sourceId = source[config.idField]; + const key = recordKey(kind, sourceId); + if (completedKeys.has(key)) return; + if (signal?.aborted) throw Object.assign(new Error("Transfer cancelled by the user."), { code: "transfer_cancelled" }); + const target = conflictTarget(existingCollections[kind], archive, source, kind, config.idField); + if (target && conflictPolicy === TRANSFER_CONFLICT_POLICIES.SKIP) { + idMaps[kind].set(sourceId, target[config.idField]); + items.push({ key, kind, sourceId, targetId: target[config.idField], status: "skipped", reason: "existing_record" }); + completedKeys.add(key); + await persistReport(reportInProgress()); + return; + } + const targetId = importedRecordId({ source, target, policy: conflictPolicy, idField: config.idField, idService: applicationIds, kind }); + idMaps[kind].set(sourceId, targetId); + const previous = await repositories[config.repository].get(targetId); + const provenance = transferProvenance({ archive, sourceId, kind, importedAt: applicationClock.now(), destinationWorkspaceId }); + const imported = kind === "campaign" + ? rebuildCampaign({ source, targetId, destinationWorkspaceId, provenance, idMaps }) + : buildMetadataRecord({ source, kind: config.kind, idField: config.idField, targetId, destinationWorkspaceId, provenance, idMaps }); + await repositories[config.repository].upsert(imported); + journal.push({ kind, id: targetId, previous: previous ? portableClone(previous) : null }); + items.push({ + key, + kind, + sourceId, + targetId, + status: target ? (conflictPolicy === "copy" ? "copied" : "replaced") : "imported", + }); + completedKeys.add(key); + await persistReport(reportInProgress()); + } + + async function processAsset(source) { + const key = recordKey("asset", source.assetId); + if (completedKeys.has(key)) return; + if (signal?.aborted) throw Object.assign(new Error("Transfer cancelled by the user."), { code: "transfer_cancelled" }); + const target = conflictTarget(existingCollections.asset, archive, source, "asset", "assetId"); + if (target && conflictPolicy === TRANSFER_CONFLICT_POLICIES.SKIP) { + idMaps.asset.set(source.assetId, target.assetId); + if (source.blobId && target.blobId) idMaps.blob.set(source.blobId, target.blobId); + items.push({ key, kind: "asset", sourceId: source.assetId, targetId: target.assetId, status: "skipped", reason: "existing_record" }); + completedKeys.add(key); + await persistReport(reportInProgress()); + return; + } + + const targetAssetId = importedRecordId({ source, target, policy: conflictPolicy, idField: "assetId", idService: applicationIds, kind: "asset" }); + const targetBlobId = source.blobId + ? target && conflictPolicy === TRANSFER_CONFLICT_POLICIES.COPY + ? applicationIds.create("blob") + : source.blobId + : null; + idMaps.asset.set(source.assetId, targetAssetId); + if (source.blobId) idMaps.blob.set(source.blobId, targetBlobId); + const previousAsset = await repositories.assetRepository.get(targetAssetId); + const previousBlobValue = targetBlobId ? await blobs.get(targetBlobId) : null; + const previousBlob = previousBlobValue === null || previousBlobValue === undefined ? null : encodeBlobPayload(previousBlobValue); + const blobEntry = (archive.payload?.blobEntries || []).find((entry) => entry.blobId === source.blobId); + const localJournal = []; + try { + let availability = source.availability || "available"; + if (targetBlobId && blobEntry) { + await blobs.put(targetBlobId, decodeBlobPayload(blobEntry)); + localJournal.push({ kind: "blob", id: targetBlobId, previous: previousBlob }); + } else if (targetBlobId && !blobEntry) { + availability = "missing_payload"; + warnings.push({ code: "asset_metadata_only", assetId: source.assetId, message: `Asset ${source.assetId} imported without blob payload.` }); + } + const importedAsset = createDomainRecord("Asset", { + ...updateReferences(source, idMaps), + assetId: targetAssetId, + blobId: targetBlobId, + workspaceId: destinationWorkspaceId, + availability, + transferProvenance: transferProvenance({ + archive, + sourceId: source.assetId, + kind: "asset", + importedAt: applicationClock.now(), + destinationWorkspaceId, + }), + importedHistoricalRecord: true, + }); + await repositories.assetRepository.upsert(importedAsset); + localJournal.push({ kind: "asset", id: targetAssetId, previous: previousAsset ? portableClone(previousAsset) : null }); + journal.push(...localJournal); + items.push({ + key, + kind: "asset", + sourceId: source.assetId, + targetId: targetAssetId, + status: target ? (conflictPolicy === "copy" ? "copied" : "replaced") : "imported", + availability, + }); + completedKeys.add(key); + await persistReport(reportInProgress()); + } catch (error) { + await rollbackJournal(localJournal); + throw error; + } + } + + try { + for (const source of sourceCollections.campaign) await processRecord("campaign", source); + for (const source of sourceCollections.asset) await processAsset(source); + for (const source of sourceCollections.sourceArtifact) await processRecord("sourceArtifact", source); + for (const source of sourceCollections.approval) await processRecord("approval", source); + for (const source of sourceCollections.export) await processRecord("export", source); + + const completedAt = applicationClock.now(); + const complete = reportRecord({ + transferReportId, + archive, + status: TRANSFER_STATUSES.COMPLETE, + destinationWorkspaceId, + conflictPolicy, + startedAt, + updatedAt: completedAt, + completedAt, + validation: preview.validation, + summary: { + imported: items.filter((item) => ["imported", "replaced", "copied"].includes(item.status)).length, + skipped: items.filter((item) => item.status === "skipped").length, + warnings: warnings.length, + }, + items, + journal, + warnings, + errors, + }); + return persistReport(complete); + } catch (error) { + const cancelled = error.code === "transfer_cancelled"; + errors.push({ code: error.code || "import_failed", message: error.message }); + if (atomic && !cancelled) { + const rollback = await rollbackJournal(journal); + const failed = reportRecord({ + transferReportId, + archive, + status: TRANSFER_STATUSES.FAILED, + destinationWorkspaceId, + conflictPolicy, + startedAt, + updatedAt: applicationClock.now(), + validation: preview.validation, + summary: { importedBeforeFailure: items.length, rolledBack: rollback.complete }, + items, + journal, + warnings, + errors, + rollback, + }); + return persistReport(failed); + } + const status = cancelled + ? TRANSFER_STATUSES.CANCELLED + : items.some((item) => ["imported", "replaced", "copied"].includes(item.status)) + ? TRANSFER_STATUSES.PARTIALLY_IMPORTED + : TRANSFER_STATUSES.FAILED; + const partial = reportRecord({ + transferReportId, + archive, + status, + destinationWorkspaceId, + conflictPolicy, + startedAt, + updatedAt: applicationClock.now(), + validation: preview.validation, + summary: { + completed: items.length, + canResume: status === TRANSFER_STATUSES.PARTIALLY_IMPORTED || status === TRANSFER_STATUSES.CANCELLED, + }, + items, + journal, + warnings, + errors, + }); + return persistReport(partial); + } + } + + async function resumeImport(archive, transferReportId, options = {}) { + const report = await reports.get(transferReportId); + if (!report) throw new Error(`Transfer report ${transferReportId} was not found.`); + if (![TRANSFER_STATUSES.PARTIALLY_IMPORTED, TRANSFER_STATUSES.CANCELLED, TRANSFER_STATUSES.FAILED].includes(report.status)) { + throw new Error(`Transfer report ${transferReportId} cannot be resumed from status ${report.status}.`); + } + return importArchive(archive, { + ...options, + destinationWorkspaceId: options.destinationWorkspaceId || report.destinationWorkspaceId, + conflictPolicy: options.conflictPolicy || report.conflictPolicy, + resumeReportId: transferReportId, + atomic: options.atomic ?? false, + }); + } + + async function rollbackImport(transferReportId) { + const report = await reports.get(transferReportId); + if (!report) throw new Error(`Transfer report ${transferReportId} was not found.`); + if (report.status === TRANSFER_STATUSES.ROLLED_BACK) return report; + const rollback = await rollbackJournal(report.journal || []); + const updatedAt = applicationClock.now(); + const rolledBack = createDomainRecord("TransferReport", { + ...report, + status: rollback.complete ? TRANSFER_STATUSES.ROLLED_BACK : TRANSFER_STATUSES.FAILED, + updatedAt, + completedAt: rollback.complete ? updatedAt : report.completedAt, + rollback, + errors: [...(report.errors || []), ...rollback.errors.map((error) => ({ code: "rollback_failed", ...error }))], + }); + return persistReport(rolledBack); + } + + async function exportReport(transferReportId) { + return reports.get(transferReportId); + } + + async function listReports() { + return reports.list(); + } + + async function archiveFingerprint(archive) { + return sha256Hex(stableStringify(archive)); + } + + return { + exportSelection, + previewImport, + importArchive, + resumeImport, + rollbackImport, + exportReport, + listReports, + archiveFingerprint, + }; +} diff --git a/frontend/lib/transfer/transferState.mjs b/frontend/lib/transfer/transferState.mjs new file mode 100644 index 0000000..0c49d07 --- /dev/null +++ b/frontend/lib/transfer/transferState.mjs @@ -0,0 +1,81 @@ +import { TRANSFER_STATUSES } from "./transferApplication.mjs"; + +export function createInitialTransferState() { + return { + status: "idle", + archive: null, + preview: null, + report: null, + error: "", + destinationWorkspaceId: "local-browser", + conflictPolicy: "skip", + selectedFileName: "", + }; +} + +export function transferReducer(state, action) { + switch (action?.type) { + case "RESET": + return createInitialTransferState(); + case "SET_DESTINATION": + return { ...state, destinationWorkspaceId: String(action.value || "") }; + case "SET_CONFLICT_POLICY": + return { ...state, conflictPolicy: String(action.value || "skip") }; + case "PREPARING": + return { ...state, status: TRANSFER_STATUSES.PREPARING, error: "", report: null }; + case "VALIDATING": + return { + ...state, + status: TRANSFER_STATUSES.VALIDATING, + archive: action.archive || state.archive, + selectedFileName: action.fileName || state.selectedFileName, + error: "", + }; + case "PREVIEW_READY": + return { + ...state, + status: action.preview?.status || "ready", + preview: action.preview, + archive: action.archive || state.archive, + error: "", + }; + case "IMPORTING": + return { ...state, status: TRANSFER_STATUSES.IMPORTING, error: "" }; + case "REPORT": + return { + ...state, + status: action.report?.status || TRANSFER_STATUSES.FAILED, + report: action.report, + error: "", + }; + case "CANCELLED": + return { ...state, status: TRANSFER_STATUSES.CANCELLED, error: "" }; + case "FAILED": + return { ...state, status: TRANSFER_STATUSES.FAILED, error: String(action.error || "Transfer failed.") }; + default: + return state; + } +} + +export function selectTransferView(state) { + const status = state?.status || "idle"; + return { + status, + busy: [TRANSFER_STATUSES.PREPARING, TRANSFER_STATUSES.VALIDATING, TRANSFER_STATUSES.IMPORTING, TRANSFER_STATUSES.UPLOADING].includes(status), + canChooseFile: ![TRANSFER_STATUSES.IMPORTING, TRANSFER_STATUSES.UPLOADING].includes(status), + canImport: Boolean( + state?.archive + && state?.preview + && ![TRANSFER_STATUSES.BLOCKED, TRANSFER_STATUSES.SELECTING_DESTINATION].includes(state.preview.status) + && state.destinationWorkspaceId, + ), + canResume: [TRANSFER_STATUSES.PARTIALLY_IMPORTED, TRANSFER_STATUSES.CANCELLED].includes(state?.report?.status), + canRollback: Boolean( + state?.report?.journal?.length + && ![TRANSFER_STATUSES.ROLLED_BACK, TRANSFER_STATUSES.BLOCKED].includes(state.report.status), + ), + hasWarnings: Boolean(state?.preview?.warnings?.length || state?.report?.warnings?.length), + isComplete: status === TRANSFER_STATUSES.COMPLETE, + isBlocked: status === TRANSFER_STATUSES.BLOCKED, + }; +} diff --git a/frontend/tests/domainContracts.test.mjs b/frontend/tests/domainContracts.test.mjs index a9c0762..2f7be2c 100644 --- a/frontend/tests/domainContracts.test.mjs +++ b/frontend/tests/domainContracts.test.mjs @@ -26,6 +26,7 @@ const samples = { Connection: { connectionId: "connection-1", provider: "linkedin", status: "connected" }, UsageEvent: { usageEventId: "usage-1", eventType: "generation" }, AuditEvent: { auditEventId: "audit-1", eventType: "campaign.saved" }, + TransferReport: { transferReportId: "report-1", archiveId: "archive-1", status: "complete" }, }; test("every declared domain contract creates and round-trips a versioned record", () => { diff --git a/frontend/tests/portableArchive.test.mjs b/frontend/tests/portableArchive.test.mjs new file mode 100644 index 0000000..78ce6f6 --- /dev/null +++ b/frontend/tests/portableArchive.test.mjs @@ -0,0 +1,197 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { + createHmacArchiveSigner, + createPortableArchive, + decodeBlobPayload, + encodeBlobPayload, + PORTABLE_ARCHIVE_SCHEMA_VERSION, + sha256Hex, + validateArchivePath, + validatePortableArchive, +} from "../lib/transfer/portableArchive.mjs"; +import { stableStringify } from "../lib/domain/contracts.mjs"; + +function archiveInput(overrides = {}) { + return { + archiveId: "archive-test-1", + createdAt: "2026-07-30T18:30:00.000Z", + sourceDeployment: { profile: "local", productVersion: "0.2.0" }, + campaigns: [{ + schemaVersion: 1, + kind: "Campaign", + campaignId: "campaign-1", + title: "Portable campaign", + drafts: { + linkedin: { + schemaVersion: 1, + kind: "ChannelDraft", + draftId: "draft-1", + channel: "linkedin", + current: { content: "Edited authoritative copy" }, + generated: { content: "Generated baseline" }, + history: [], + edited: true, + approved: true, + generationRunId: "run-1", + }, + }, + brief: { notes: "Portable evidence" }, + sourceFiles: [], + documentText: [], + createdAt: "2026-07-29T10:00:00.000Z", + updatedAt: "2026-07-29T11:00:00.000Z", + }], + assets: [], + sourceArtifacts: [], + approvals: [], + exports: [], + blobEntries: [], + ...overrides, + }; +} + +function withoutIntegrity(archive) { + const { integrity, signature, ...unsigned } = archive; + void integrity; + void signature; + return unsigned; +} + +async function recalculateIntegrity(archive) { + archive.integrity = { + algorithm: "SHA-256", + digest: await sha256Hex(stableStringify(withoutIntegrity(archive))), + }; + return archive; +} + +test("portable archive is deterministic, integrity protected, and optionally HMAC signed", async () => { + const signer = createHmacArchiveSigner({ secret: "correct horse battery staple", keyId: "test-key" }); + const archive = await createPortableArchive({ ...archiveInput(), signer }); + const validated = await validatePortableArchive(archive, { signer, requireSignature: true }); + assert.equal(validated.valid, true); + assert.equal(archive.signature.algorithm, "HMAC-SHA-256"); + assert.equal(archive.signature.keyId, "test-key"); + + const same = await createPortableArchive({ ...archiveInput(), signer }); + assert.equal(stableStringify(archive), stableStringify(same)); + + const wrongSigner = createHmacArchiveSigner({ secret: "wrong secret", keyId: "test-key" }); + const invalid = await validatePortableArchive(archive, { signer: wrongSigner, requireSignature: true }); + assert.equal(invalid.valid, false); + assert.ok(invalid.errors.some((error) => error.code === "invalid_signature")); +}); + +test("secret fields, private endpoints, local paths, and signed references never enter the transferable payload", async () => { + const input = archiveInput({ + sourceDeployment: { + profile: "local", + baseUrl: "http://127.0.0.1:3000", + accessToken: "do-not-export", + }, + campaigns: [{ + ...archiveInput().campaigns[0], + brief: { + notes: "Portable evidence", + apiKey: "do-not-export", + providerBaseUrl: "http://localhost:11434", + }, + sourceFiles: [ + { name: "notes.md", path: "C:\\Users\\Ankit\\private\\notes.md" }, + { name: "home.md", filesystemPath: "/home/ankit/private/home.md" }, + ], + }], + assets: [{ + schemaVersion: 1, + kind: "Asset", + assetId: "asset-1", + assetType: "image", + signedUrl: "https://private.example/signed", + localPath: "/Users/ankit/Desktop/private.png", + }], + }); + const archive = await createPortableArchive(input); + const transferablePayload = stableStringify({ + sourceDeployment: archive.sourceDeployment, + payload: archive.payload, + }); + assert.doesNotMatch( + transferablePayload, + /do-not-export|localhost|127\.0\.0\.1|Users\\Ankit|\/home\/ankit|\/Users\/ankit|private\.example/i, + ); + assert.ok(archive.manifest.exclusions.length >= 6); + assert.ok(archive.manifest.exclusions.some((item) => item.reason === "secret field")); + assert.ok(archive.manifest.exclusions.some((item) => /path|reference|endpoint/i.test(item.reason))); +}); + +test("corrupted content and future schema versions fail safely with actionable codes", async () => { + const archive = await createPortableArchive(archiveInput()); + const corrupted = structuredClone(archive); + corrupted.payload.campaigns[0].title = "Tampered campaign"; + const invalid = await validatePortableArchive(corrupted); + assert.equal(invalid.valid, false); + assert.ok(invalid.errors.some((error) => error.code === "integrity_mismatch")); + + const future = structuredClone(archive); + future.schemaVersion = PORTABLE_ARCHIVE_SCHEMA_VERSION + 1; + await recalculateIntegrity(future); + const futureResult = await validatePortableArchive(future); + assert.equal(futureResult.valid, false); + assert.match(futureResult.errors.find((error) => error.code === "future_schema").message, /upgrade SignalFlow/i); +}); + +test("archive traversal, oversized payloads, and partial assets are reported", async () => { + const encoded = encodeBlobPayload("asset payload"); + const archive = await createPortableArchive(archiveInput({ + assets: [{ + schemaVersion: 1, + kind: "Asset", + assetId: "asset-1", + assetType: "image", + blobId: "blob-1", + }], + blobEntries: [{ + blobId: "blob-1", + assetId: "asset-1", + archivePath: "blobs/blob-1.bin", + contentType: "text/plain", + ...encoded, + }], + })); + assert.equal(validateArchivePath("blobs/blob-1.bin"), true); + assert.equal(validateArchivePath("../private.txt"), false); + assert.equal(validateArchivePath("blobs/../../private.txt"), false); + assert.equal(validateArchivePath("/blobs/private.txt"), false); + assert.equal(validateArchivePath("blobs\\private.txt"), false); + + const traversal = structuredClone(archive); + traversal.payload.blobEntries[0].archivePath = "blobs/../../private.txt"; + await recalculateIntegrity(traversal); + const traversalResult = await validatePortableArchive(traversal); + assert.ok(traversalResult.errors.some((error) => error.code === "archive_traversal")); + + const oversized = await validatePortableArchive(archive, { maxBytes: 100 }); + assert.ok(oversized.errors.some((error) => error.code === "archive_too_large")); + + const partial = await createPortableArchive(archiveInput({ + assets: [{ + schemaVersion: 1, + kind: "Asset", + assetId: "asset-missing", + assetType: "video", + blobId: "blob-missing", + }], + })); + const partialResult = await validatePortableArchive(partial); + assert.equal(partialResult.valid, true); + assert.ok(partialResult.warnings.some((warning) => warning.code === "partial_assets")); +}); + +test("blob payloads round-trip bytes, text, and JSON", () => { + const bytes = new Uint8Array([0, 1, 2, 254, 255]); + assert.deepEqual(Array.from(decodeBlobPayload(encodeBlobPayload(bytes))), Array.from(bytes)); + assert.equal(decodeBlobPayload(encodeBlobPayload("hello")), "hello"); + assert.deepEqual(decodeBlobPayload(encodeBlobPayload({ hello: "world" })), { hello: "world" }); +}); diff --git a/frontend/tests/portableTransferApplication.test.mjs b/frontend/tests/portableTransferApplication.test.mjs new file mode 100644 index 0000000..4382514 --- /dev/null +++ b/frontend/tests/portableTransferApplication.test.mjs @@ -0,0 +1,374 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { createCampaignAggregate } from "../lib/domain/campaign.mjs"; +import { createDomainRecord } from "../lib/domain/contracts.mjs"; +import { createDeterministicIdService } from "../lib/domain/ports.mjs"; +import { + createMemoryAsyncStore, + createMemoryBlobStorage, + createMemoryCampaignRepository, + createStoreBackedBlobStorage, + createStoreBackedCampaignRepository, +} from "../lib/infrastructure/adapters.mjs"; +import { + createMemoryApprovalRepository, + createMemoryAssetRepository, + createMemoryExportRepository, + createMemorySourceArtifactRepository, + createMemoryTransferReportRepository, + createStoreBackedApprovalRepository, + createStoreBackedAssetRepository, + createStoreBackedExportRepository, + createStoreBackedSourceArtifactRepository, + createStoreBackedTransferReportRepository, +} from "../lib/infrastructure/transferAdapters.mjs"; +import { + createTransferApplication, + TRANSFER_CONFLICT_POLICIES, + TRANSFER_STATUSES, +} from "../lib/transfer/transferApplication.mjs"; +import { campaignInput } from "./campaignFixtures.mjs"; + +function sequenceClock(values) { + let index = 0; + return { + now() { + const value = values[Math.min(index, values.length - 1)]; + index += 1; + return value; + }, + }; +} + +function clock() { + return sequenceClock([ + "2026-07-30T18:00:00.000Z", + "2026-07-30T18:01:00.000Z", + "2026-07-30T18:02:00.000Z", + "2026-07-30T18:03:00.000Z", + "2026-07-30T18:04:00.000Z", + "2026-07-30T18:05:00.000Z", + "2026-07-30T18:06:00.000Z", + "2026-07-30T18:07:00.000Z", + "2026-07-30T18:08:00.000Z", + "2026-07-30T18:09:00.000Z", + "2026-07-30T18:10:00.000Z", + ]); +} + +function campaignFixture(overrides = {}) { + return createCampaignAggregate(campaignInput({ + campaignId: "campaign-transfer-1", + generatedPosts: { + linkedin: "Generated LinkedIn baseline", + x: "Generated X baseline", + blog: "Generated blog baseline", + }, + channelStates: { + linkedin: { status: "generated", edited: true, approved: true, generationRunId: "run-fixture-1" }, + x: { status: "needs_review", edited: true, approved: false, generationRunId: "run-fixture-1" }, + blog: { status: "regenerated", edited: true, approved: false, generationRunId: "run-fixture-1" }, + }, + archives: [{ + archiveId: "generation-archive-1", + createdAt: "2026-07-29T12:00:00.000Z", + reason: "archive_all", + posts: { linkedin: "Previous edited LinkedIn" }, + generatedPosts: { linkedin: "Previous generated LinkedIn" }, + channelStates: { linkedin: { status: "generated", edited: true, approved: false, generationRunId: "run-old" } }, + result: { providerUsed: "gemini", generation_status: { linkedin: { status: "generated" } } }, + generationRun: { generationRunId: "run-old", provider: "gemini" }, + activeChannel: "linkedin", + revision: 4, + }], + editorState: { + revision: 9, + savedRevision: 9, + exportedRevision: 8, + lastSavedAt: "2026-07-30T10:00:00.000Z", + lastExportedAt: "2026-07-30T09:00:00.000Z", + savedSourceFingerprint: "sf1-fixture", + }, + ...overrides, + })); +} + +function metadataFixture() { + const asset = createDomainRecord("Asset", { + assetId: "asset-transfer-1", + workspaceId: "workspace-local", + assetType: "image", + blobId: "blob-transfer-1", + contentType: "text/plain", + fileName: "preview.txt", + createdAt: "2026-07-29T10:00:00.000Z", + updatedAt: "2026-07-29T10:00:00.000Z", + }); + const sourceArtifact = createDomainRecord("SourceArtifact", { + sourceArtifactId: "source-artifact-1", + campaignId: "campaign-transfer-1", + assetId: "asset-transfer-1", + artifactType: "uploaded_file", + title: "Preview source", + createdAt: "2026-07-29T10:00:00.000Z", + }); + const exportRecord = createDomainRecord("Export", { + exportId: "export-transfer-1", + campaignId: "campaign-transfer-1", + format: "markdown", + createdAt: "2026-07-29T11:00:00.000Z", + }); + return { asset, sourceArtifact, exportRecord }; +} + +function memoryApplication({ campaigns = [], assets = [], sourceArtifacts = [], approvals = [], exports = [], blobValues = {}, assetRepository = null } = {}) { + return createTransferApplication({ + campaignRepository: createMemoryCampaignRepository(campaigns), + assetRepository: assetRepository || createMemoryAssetRepository(assets), + sourceArtifactRepository: createMemorySourceArtifactRepository(sourceArtifacts), + approvalRepository: createMemoryApprovalRepository(approvals), + exportRepository: createMemoryExportRepository(exports), + blobStorage: createMemoryBlobStorage(blobValues), + transferReportRepository: createMemoryTransferReportRepository(), + clock: clock(), + idService: createDeterministicIdService("memory-transfer"), + }); +} + +function storeApplication(store = createMemoryAsyncStore()) { + return { + store, + app: createTransferApplication({ + campaignRepository: createStoreBackedCampaignRepository({ store }), + assetRepository: createStoreBackedAssetRepository({ store }), + sourceArtifactRepository: createStoreBackedSourceArtifactRepository({ store }), + approvalRepository: createStoreBackedApprovalRepository({ store }), + exportRepository: createStoreBackedExportRepository({ store }), + blobStorage: createStoreBackedBlobStorage({ store }), + transferReportRepository: createStoreBackedTransferReportRepository({ store }), + clock: clock(), + idService: createDeterministicIdService("store-transfer"), + }), + }; +} + +async function sourceArchive() { + const campaign = campaignFixture(); + const { asset, sourceArtifact, exportRecord } = metadataFixture(); + const source = memoryApplication({ + campaigns: [campaign], + assets: [asset], + sourceArtifacts: [sourceArtifact], + exports: [exportRecord], + blobValues: { "blob-transfer-1": "portable asset payload" }, + }); + const archive = await source.exportSelection({ + sourceDeployment: { profile: "local", productVersion: "0.2.0", deploymentId: "local-fixture" }, + }); + return { archive, campaign, asset, sourceArtifact, exportRecord }; +} + +test("local to store-backed hosted to fresh local round-trip preserves content history provenance and blobs", async () => { + const { archive, campaign } = await sourceArchive(); + const hosted = storeApplication(); + const preview = await hosted.app.previewImport(archive, { destinationWorkspaceId: "workspace-cloud" }); + assert.ok(["ready", TRANSFER_STATUSES.WARNINGS_FOUND].includes(preview.status)); + assert.equal(preview.counts.campaigns, 1); + assert.equal(preview.counts.assets, 1); + + const imported = await hosted.app.importArchive(archive, { destinationWorkspaceId: "workspace-cloud" }); + assert.equal(imported.status, TRANSFER_STATUSES.COMPLETE); + const hostedCampaigns = await createStoreBackedCampaignRepository({ store: hosted.store }).list(); + const hostedCampaign = hostedCampaigns[0]; + assert.equal(hostedCampaign.drafts.linkedin.current.content, campaign.drafts.linkedin.current.content); + assert.equal(hostedCampaign.drafts.linkedin.generated.content, campaign.drafts.linkedin.generated.content); + assert.equal(hostedCampaign.drafts.linkedin.approved, true); + assert.equal(hostedCampaign.archives[0].archiveId, "generation-archive-1"); + assert.equal(hostedCampaign.generationRun.createdAt, campaign.generationRun.createdAt); + assert.equal(hostedCampaign.transferProvenance.archiveId, archive.archiveId); + assert.equal(hostedCampaign.transferProvenance.historical, true); + assert.equal(hostedCampaign.transferProvenance.destinationWorkspaceId, "workspace-cloud"); + + const hostedAssets = await createStoreBackedAssetRepository({ store: hosted.store }).list(); + assert.equal(hostedAssets[0].transferProvenance.sourceAssetId, "asset-transfer-1"); + assert.equal(await createStoreBackedBlobStorage({ store: hosted.store }).get("blob-transfer-1"), "portable asset payload"); + const hostedArtifacts = await createStoreBackedSourceArtifactRepository({ store: hosted.store }).list(); + assert.equal(hostedArtifacts[0].campaignId, hostedCampaign.campaignId); + assert.equal(hostedArtifacts[0].assetId, hostedAssets[0].assetId); + assert.equal((await createStoreBackedApprovalRepository({ store: hosted.store }).list())[0].status, "approved"); + assert.equal((await createStoreBackedExportRepository({ store: hosted.store }).list())[0].campaignId, hostedCampaign.campaignId); + + const hostedArchive = await hosted.app.exportSelection({ + sourceDeployment: { profile: "hosted", productVersion: "0.2.0", deploymentId: "cloud-fixture" }, + }); + const freshLocal = memoryApplication(); + const localImport = await freshLocal.importArchive(hostedArchive, { destinationWorkspaceId: "workspace-local-restored" }); + assert.equal(localImport.status, TRANSFER_STATUSES.COMPLETE); + const localArchive = await freshLocal.exportSelection({ sourceDeployment: { profile: "local", deploymentId: "restored" } }); + const restoredCampaign = localArchive.payload.campaigns[0]; + assert.equal(restoredCampaign.drafts.linkedin.current.content, campaign.drafts.linkedin.current.content); + assert.equal(restoredCampaign.drafts.linkedin.generated.content, campaign.drafts.linkedin.generated.content); + assert.equal(restoredCampaign.drafts.linkedin.approved, true); + assert.equal(restoredCampaign.archives[0].archiveId, "generation-archive-1"); + assert.equal(restoredCampaign.transferProvenance.sourceCampaignId, hostedCampaign.campaignId); +}); + +test("re-import is idempotent with skip and creates independent IDs with copy", async () => { + const { archive } = await sourceArchive(); + const target = storeApplication(); + const first = await target.app.importArchive(archive, { destinationWorkspaceId: "workspace-cloud" }); + assert.equal(first.status, TRANSFER_STATUSES.COMPLETE); + + const secondPreview = await target.app.previewImport(archive, { destinationWorkspaceId: "workspace-cloud" }); + assert.ok(secondPreview.conflicts.some((conflict) => conflict.type === "already_imported")); + const second = await target.app.importArchive(archive, { + destinationWorkspaceId: "workspace-cloud", + conflictPolicy: TRANSFER_CONFLICT_POLICIES.SKIP, + }); + assert.equal(second.status, TRANSFER_STATUSES.COMPLETE); + assert.ok(second.items.every((item) => item.status === "skipped")); + assert.equal((await createStoreBackedCampaignRepository({ store: target.store }).list()).length, 1); + + const copied = await target.app.importArchive(archive, { + destinationWorkspaceId: "workspace-cloud", + conflictPolicy: TRANSFER_CONFLICT_POLICIES.COPY, + }); + assert.equal(copied.status, TRANSFER_STATUSES.COMPLETE); + const campaigns = await createStoreBackedCampaignRepository({ store: target.store }).list(); + assert.equal(campaigns.length, 2); + assert.notEqual(campaigns[0].campaignId, campaigns[1].campaignId); + assert.equal(campaigns[0].title, campaigns[1].title); + assert.equal((await createStoreBackedAssetRepository({ store: target.store }).list()).length, 2); +}); + +test("replace conflict restores the archive version without relabeling historical timestamps", async () => { + const { archive, campaign } = await sourceArchive(); + const altered = createCampaignAggregate({ + ...campaignInput({ + campaignId: "campaign-transfer-1", + posts: { linkedin: "Locally changed collision" }, + generatedPosts: { linkedin: "Locally changed generated" }, + channels: ["linkedin"], + createdAt: "2026-07-30T12:00:00.000Z", + updatedAt: "2026-07-30T12:00:00.000Z", + }), + }); + const target = storeApplication(); + await createStoreBackedCampaignRepository({ store: target.store }).upsert(altered); + const report = await target.app.importArchive(archive, { + destinationWorkspaceId: "workspace-cloud", + conflictPolicy: TRANSFER_CONFLICT_POLICIES.REPLACE, + }); + assert.equal(report.status, TRANSFER_STATUSES.COMPLETE); + const restored = await createStoreBackedCampaignRepository({ store: target.store }).get("campaign-transfer-1"); + assert.equal(restored.drafts.linkedin.current.content, campaign.drafts.linkedin.current.content); + assert.equal(restored.createdAt, campaign.createdAt); + assert.equal(restored.updatedAt, campaign.updatedAt); + assert.equal(restored.transferProvenance.historical, true); +}); + +test("non-atomic partial imports can resume without duplicating completed records", async () => { + const { archive } = await sourceArchive(); + const backingAssets = createMemoryAssetRepository(); + let failOnce = true; + const failingAssets = { + list: () => backingAssets.list(), + get: (id) => backingAssets.get(id), + remove: (id) => backingAssets.remove(id), + async upsert(value) { + if (failOnce) { + failOnce = false; + throw new Error("simulated asset write failure"); + } + return backingAssets.upsert(value); + }, + }; + const app = createTransferApplication({ + campaignRepository: createMemoryCampaignRepository(), + assetRepository: failingAssets, + sourceArtifactRepository: createMemorySourceArtifactRepository(), + approvalRepository: createMemoryApprovalRepository(), + exportRepository: createMemoryExportRepository(), + blobStorage: createMemoryBlobStorage(), + transferReportRepository: createMemoryTransferReportRepository(), + clock: clock(), + idService: createDeterministicIdService("resume-transfer"), + }); + + const partial = await app.importArchive(archive, { destinationWorkspaceId: "workspace-cloud", atomic: false }); + assert.equal(partial.status, TRANSFER_STATUSES.PARTIALLY_IMPORTED); + assert.ok(partial.items.some((item) => item.kind === "campaign" && item.status === "imported")); + + const resumed = await app.resumeImport(archive, partial.transferReportId, { atomic: false }); + assert.equal(resumed.status, TRANSFER_STATUSES.COMPLETE); + assert.equal(resumed.items.filter((item) => item.kind === "campaign").length, 1); + assert.equal((await backingAssets.list()).length, 1); +}); + +test("atomic import failure rolls back every prior record and manual rollback reverses a completed import", async () => { + const { archive } = await sourceArchive(); + const campaigns = createMemoryCampaignRepository(); + const assets = createMemoryAssetRepository(); + const artifacts = createMemorySourceArtifactRepository(); + const failingArtifacts = { + list: () => artifacts.list(), + get: (id) => artifacts.get(id), + remove: (id) => artifacts.remove(id), + async upsert() { throw new Error("simulated source artifact failure"); }, + }; + const reports = createMemoryTransferReportRepository(); + const atomicApp = createTransferApplication({ + campaignRepository: campaigns, + assetRepository: assets, + sourceArtifactRepository: failingArtifacts, + approvalRepository: createMemoryApprovalRepository(), + exportRepository: createMemoryExportRepository(), + blobStorage: createMemoryBlobStorage(), + transferReportRepository: reports, + clock: clock(), + idService: createDeterministicIdService("atomic-transfer"), + }); + const failed = await atomicApp.importArchive(archive, { destinationWorkspaceId: "workspace-cloud", atomic: true }); + assert.equal(failed.status, TRANSFER_STATUSES.FAILED); + assert.equal(failed.rollback.complete, true); + assert.equal((await campaigns.list()).length, 0); + assert.equal((await assets.list()).length, 0); + + const target = storeApplication(); + const complete = await target.app.importArchive(archive, { destinationWorkspaceId: "workspace-cloud" }); + assert.equal(complete.status, TRANSFER_STATUSES.COMPLETE); + const rolledBack = await target.app.rollbackImport(complete.transferReportId); + assert.equal(rolledBack.status, TRANSFER_STATUSES.ROLLED_BACK); + assert.equal((await createStoreBackedCampaignRepository({ store: target.store }).list()).length, 0); + assert.equal((await createStoreBackedAssetRepository({ store: target.store }).list()).length, 0); +}); + +test("legacy campaign records export through the canonical contract", async () => { + const legacy = { + id: "legacy-transfer-1", + title: "Legacy portable campaign", + channels: ["linkedin"], + posts: { linkedin: "Legacy edited authoritative draft" }, + result: { + providerUsed: "gemini", + posts: { linkedin: "Legacy generated baseline" }, + generation_status: { linkedin: { status: "generated" } }, + package: { project: { name: "Legacy portable campaign" } }, + }, + generationRun: campaignInput().generationRun, + brief: campaignInput().brief, + createdAt: "2026-07-28T10:00:00.000Z", + updatedAt: "2026-07-28T11:00:00.000Z", + }; + const source = memoryApplication({ campaigns: [legacy] }); + const archive = await source.exportSelection({ sourceDeployment: { profile: "local" } }); + const target = memoryApplication(); + const imported = await target.importArchive(archive, { destinationWorkspaceId: "workspace-import" }); + assert.equal(imported.status, TRANSFER_STATUSES.COMPLETE); + const exportedAgain = await target.exportSelection({ sourceDeployment: { profile: "local" } }); + const campaign = exportedAgain.payload.campaigns[0]; + assert.equal(campaign.campaignId, "legacy-transfer-1"); + assert.equal(campaign.drafts.linkedin.current.content, "Legacy edited authoritative draft"); + assert.equal(campaign.drafts.linkedin.generated.content, "Legacy generated baseline"); +}); diff --git a/frontend/tests/transferAdapters.test.mjs b/frontend/tests/transferAdapters.test.mjs new file mode 100644 index 0000000..66bd28c --- /dev/null +++ b/frontend/tests/transferAdapters.test.mjs @@ -0,0 +1,87 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { createDomainRecord } from "../lib/domain/contracts.mjs"; +import { createMemoryAsyncStore } from "../lib/infrastructure/adapters.mjs"; +import { + createBrowserApprovalRepository, + createBrowserAssetRepository, + createBrowserBlobStorage, + createBrowserExportRepository, + createBrowserSourceArtifactRepository, + createBrowserTransferReportRepository, + createMemoryApprovalRepository, + createMemoryAssetRepository, + createMemoryExportRepository, + createMemorySourceArtifactRepository, + createMemoryTransferReportRepository, + createStoreBackedApprovalRepository, + createStoreBackedAssetRepository, + createStoreBackedExportRepository, + createStoreBackedSourceArtifactRepository, + createStoreBackedTransferReportRepository, +} from "../lib/infrastructure/transferAdapters.mjs"; + +function fakeStorage() { + const values = new Map(); + return { + getItem(key) { return values.has(key) ? values.get(key) : null; }, + setItem(key, value) { values.set(key, String(value)); }, + }; +} + +async function repositoryContract(repository, record, idField) { + assert.deepEqual(await repository.list(), []); + await repository.upsert(record); + assert.equal((await repository.get(record[idField]))[idField], record[idField]); + assert.equal((await repository.list()).length, 1); + assert.equal(await repository.remove(record[idField]), true); + assert.equal(await repository.get(record[idField]), null); +} + +const fixtures = { + asset: createDomainRecord("Asset", { assetId: "asset-1", assetType: "image", createdAt: "2026-07-30T00:00:00.000Z" }), + sourceArtifact: createDomainRecord("SourceArtifact", { sourceArtifactId: "artifact-1", artifactType: "document", createdAt: "2026-07-30T00:00:00.000Z" }), + approval: createDomainRecord("Approval", { approvalId: "approval-1", status: "approved", createdAt: "2026-07-30T00:00:00.000Z" }), + export: createDomainRecord("Export", { exportId: "export-1", format: "json", createdAt: "2026-07-30T00:00:00.000Z" }), + report: createDomainRecord("TransferReport", { transferReportId: "report-1", archiveId: "archive-1", status: "complete", createdAt: "2026-07-30T00:00:00.000Z" }), +}; + +test("memory portable metadata repositories share one contract", async () => { + await repositoryContract(createMemoryAssetRepository(), fixtures.asset, "assetId"); + await repositoryContract(createMemorySourceArtifactRepository(), fixtures.sourceArtifact, "sourceArtifactId"); + await repositoryContract(createMemoryApprovalRepository(), fixtures.approval, "approvalId"); + await repositoryContract(createMemoryExportRepository(), fixtures.export, "exportId"); + await repositoryContract(createMemoryTransferReportRepository(), fixtures.report, "transferReportId"); +}); + +test("store-backed portable metadata repositories share one contract", async () => { + const store = createMemoryAsyncStore(); + await repositoryContract(createStoreBackedAssetRepository({ store }), fixtures.asset, "assetId"); + await repositoryContract(createStoreBackedSourceArtifactRepository({ store }), fixtures.sourceArtifact, "sourceArtifactId"); + await repositoryContract(createStoreBackedApprovalRepository({ store }), fixtures.approval, "approvalId"); + await repositoryContract(createStoreBackedExportRepository({ store }), fixtures.export, "exportId"); + await repositoryContract(createStoreBackedTransferReportRepository({ store }), fixtures.report, "transferReportId"); +}); + +test("browser portable metadata repositories share one contract", async () => { + const storage = fakeStorage(); + const getStorage = () => storage; + await repositoryContract(createBrowserAssetRepository({ getStorage }), fixtures.asset, "assetId"); + await repositoryContract(createBrowserSourceArtifactRepository({ getStorage }), fixtures.sourceArtifact, "sourceArtifactId"); + await repositoryContract(createBrowserApprovalRepository({ getStorage }), fixtures.approval, "approvalId"); + await repositoryContract(createBrowserExportRepository({ getStorage }), fixtures.export, "exportId"); + await repositoryContract(createBrowserTransferReportRepository({ getStorage }), fixtures.report, "transferReportId"); +}); + +test("browser blob storage preserves bytes text and JSON independently", async () => { + const storage = createBrowserBlobStorage({ getStorage: () => fakeStorage() }); + await storage.put("bytes", new Uint8Array([0, 10, 255])); + await storage.put("text", "portable text"); + await storage.put("json", { portable: true }); + assert.deepEqual(Array.from(await storage.get("bytes")), [0, 10, 255]); + assert.equal(await storage.get("text"), "portable text"); + assert.deepEqual(await storage.get("json"), { portable: true }); + assert.equal(await storage.remove("text"), true); + assert.equal(await storage.get("text"), null); +}); diff --git a/frontend/tests/transferState.test.mjs b/frontend/tests/transferState.test.mjs new file mode 100644 index 0000000..eda3295 --- /dev/null +++ b/frontend/tests/transferState.test.mjs @@ -0,0 +1,71 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { + createInitialTransferState, + selectTransferView, + transferReducer, +} from "../lib/transfer/transferState.mjs"; +import { TRANSFER_STATUSES } from "../lib/transfer/transferApplication.mjs"; + +test("transfer state covers preparation validation warnings import completion and reset", () => { + let state = createInitialTransferState(); + assert.equal(state.status, "idle"); + state = transferReducer(state, { type: "PREPARING" }); + assert.equal(selectTransferView(state).busy, true); + const archive = { archiveId: "archive-1" }; + state = transferReducer(state, { type: "VALIDATING", archive, fileName: "campaign.signalflow.json" }); + assert.equal(state.archive, archive); + assert.equal(state.selectedFileName, "campaign.signalflow.json"); + state = transferReducer(state, { + type: "PREVIEW_READY", + archive, + preview: { status: TRANSFER_STATUSES.WARNINGS_FOUND, warnings: [{ code: "excluded_private_data" }] }, + }); + assert.equal(selectTransferView(state).hasWarnings, true); + assert.equal(selectTransferView(state).canImport, true); + state = transferReducer(state, { type: "IMPORTING" }); + assert.equal(state.status, TRANSFER_STATUSES.IMPORTING); + state = transferReducer(state, { + type: "REPORT", + report: { status: TRANSFER_STATUSES.COMPLETE, journal: [{ kind: "campaign" }] }, + }); + assert.equal(selectTransferView(state).isComplete, true); + assert.equal(selectTransferView(state).canRollback, true); + assert.deepEqual(transferReducer(state, { type: "RESET" }), createInitialTransferState()); +}); + +test("blocked and destination-selection states cannot import", () => { + for (const status of [TRANSFER_STATUSES.BLOCKED, TRANSFER_STATUSES.SELECTING_DESTINATION]) { + const state = transferReducer(createInitialTransferState(), { + type: "PREVIEW_READY", + archive: { archiveId: "archive-1" }, + preview: { status }, + }); + assert.equal(selectTransferView(state).canImport, false); + assert.equal(selectTransferView(state).isBlocked, status === TRANSFER_STATUSES.BLOCKED); + } +}); + +test("partial and cancelled reports can resume while rolled-back reports cannot", () => { + for (const status of [TRANSFER_STATUSES.PARTIALLY_IMPORTED, TRANSFER_STATUSES.CANCELLED]) { + const state = transferReducer(createInitialTransferState(), { + type: "REPORT", + report: { status, journal: [{ kind: "campaign" }] }, + }); + assert.equal(selectTransferView(state).canResume, true); + assert.equal(selectTransferView(state).canRollback, true); + } + const rolledBack = transferReducer(createInitialTransferState(), { + type: "REPORT", + report: { status: TRANSFER_STATUSES.ROLLED_BACK, journal: [{ kind: "campaign" }] }, + }); + assert.equal(selectTransferView(rolledBack).canResume, false); + assert.equal(selectTransferView(rolledBack).canRollback, false); +}); + +test("failure state exposes an actionable error", () => { + const state = transferReducer(createInitialTransferState(), { type: "FAILED", error: "Archive could not be parsed." }); + assert.equal(state.status, TRANSFER_STATUSES.FAILED); + assert.equal(state.error, "Archive could not be parsed."); +}); diff --git a/portable-transfer-failure-summary.txt b/portable-transfer-failure-summary.txt new file mode 100644 index 0000000..029709a --- /dev/null +++ b/portable-transfer-failure-summary.txt @@ -0,0 +1,29 @@ +not ok 133 - browser blob storage preserves bytes text and JSON independently + --- + duration_ms: 0.890389 + type: 'test' + location: '/home/runner/work/SignalFlow-Studio/SignalFlow-Studio/frontend/tests/transferAdapters.test.mjs:77:1' + failureType: 'testCodeFailure' + error: 'object null is not iterable (cannot read property Symbol(Symbol.iterator))' + code: 'ERR_TEST_FAILURE' + name: 'TypeError' + stack: |- + Function.from () + TestContext. (file:///home/runner/work/SignalFlow-Studio/SignalFlow-Studio/frontend/tests/transferAdapters.test.mjs:82:26) + async Test.run (node:internal/test_runner/test:1054:7) + async Test.processPendingSubtests (node:internal/test_runner/test:744:7) + ... +# Subtest: transfer state covers preparation validation warnings import completion and reset +ok 134 - transfer state covers preparation validation warnings import completion and reset + --- + duration_ms: 2.33198 + type: 'test' + ... +# Subtest: blocked and destination-selection states cannot import +ok 135 - blocked and destination-selection states cannot import + --- + duration_ms: 0.346649 + type: 'test' + ... +# Subtest: partial and cancelled reports can resume while rolled-back reports cannot +--- \ No newline at end of file diff --git a/scripts/apply-portable-transfer-domain.py b/scripts/apply-portable-transfer-domain.py new file mode 100644 index 0000000..1f20200 --- /dev/null +++ b/scripts/apply-portable-transfer-domain.py @@ -0,0 +1,76 @@ +from pathlib import Path + + +def replace_once(source: str, before: str, after: str, label: str) -> str: + count = source.count(before) + if count != 1: + raise RuntimeError(f"Expected one {label}, found {count}") + return source.replace(before, after, 1) + + +campaign_path = Path("frontend/lib/domain/campaign.mjs") +campaign = campaign_path.read_text() +campaign = replace_once( + campaign, + ''' archives: cleanArchives(input.archives, input.existingArchives), + editorState, + brief: cleanBrief(input.brief || {}),''', + ''' archives: cleanArchives(input.archives, input.existingArchives), + editorState, + transferProvenance: input.transferProvenance ? portableClone(input.transferProvenance) : null, + brief: cleanBrief(input.brief || {}),''', + "campaign transfer provenance field", +) +campaign = replace_once( + campaign, + ''' archives: input?.archives || [], + editorState: input?.editorState || {''', + ''' archives: input?.archives || [], + transferProvenance: input?.transferProvenance || null, + editorState: input?.editorState || {''', + "legacy transfer provenance migration", +) +campaign = replace_once( + campaign, + ''' archives: portableClone(campaign.archives || []), + revision: campaign.editorState?.revision || 1,''', + ''' archives: portableClone(campaign.archives || []), + transferProvenance: campaign.transferProvenance ? portableClone(campaign.transferProvenance) : null, + revision: campaign.editorState?.revision || 1,''', + "editor transfer provenance projection", +) +campaign_path.write_text(campaign) + +application_path = Path("frontend/lib/transfer/transferApplication.mjs") +application = application_path.read_text() +application = replace_once( + application, + '''const RECORD_CONFIG = Object.freeze({ + campaign: { repository: "campaignRepository", idField: "campaignId", kind: "Campaign" }, + asset: { repository: "assetRepository", idField: "assetId", kind: "Asset" }, + sourceArtifact: { repository: "sourceArtifactRepository", idField: "sourceArtifactId", kind: "SourceArtifact" }, + approval: { repository: "approvalRepository", idField: "approvalId", kind: "Approval" }, + export: { repository: "exportRepository", idField: "exportId", kind: "Export" }, +});''', + '''const RECORD_CONFIG = Object.freeze({ + campaign: { repository: "campaignRepository", idField: "campaignId", kind: "Campaign", provenanceField: "sourceCampaignId" }, + asset: { repository: "assetRepository", idField: "assetId", kind: "Asset", provenanceField: "sourceAssetId" }, + sourceArtifact: { repository: "sourceArtifactRepository", idField: "sourceArtifactId", kind: "SourceArtifact", provenanceField: "sourceArtifactId" }, + approval: { repository: "approvalRepository", idField: "approvalId", kind: "Approval", provenanceField: "sourceApprovalId" }, + export: { repository: "exportRepository", idField: "exportId", kind: "Export", provenanceField: "sourceExportId" }, +});''', + "transfer provenance field map", +) +application = replace_once( + application, + '''function sourceFieldFor(kind) { + return `source${kind[0].toUpperCase()}${kind.slice(1)}Id`; +}''', + '''function sourceFieldFor(kind) { + const field = RECORD_CONFIG[kind]?.provenanceField; + if (!field) throw new TypeError(`Unknown transfer provenance kind: ${kind}.`); + return field; +}''', + "transfer provenance field selector", +) +application_path.write_text(application)