@dilipabc/storage-core
v0.1.0
Published
Streaming source readers, landing writers, and ingestion events for Sentinel
Maintainers
Readme
storage-core
Shared TypeScript I/O library used by poll-orchestrator. It has no server or standalone runtime.
Install
npm install @dilipabc/storage-coreThe published package exposes compiled JavaScript and TypeScript declarations through the package root. Source files and local development dependencies are not included in the npm tarball.
Public API
import {
listNewFiles,
readFromSource,
writeToLanding,
StorageCoreError,
asStorageError,
errorMessage,
ReaderService,
LandingWriterService,
KafkaEventPublisher,
} from '@dilipabc/storage-core'listNewFiles(input)selects a source driver frominput.sourceCredentials.provider.readFromSource(input)returns a readable stream for a source file.writeToLanding(stream, input)selects the landing writer fromSTORAGE_PROVIDER, uploads the stream, and publishes aningestion-eventsmetadata message.StorageCoreErrorprovides structured, machine-readable failures.asStorageError(error, code, operation)converts unknown failures into aStorageCoreError.errorMessage(error)safely extracts a message from an unknown thrown value.ReaderService,LandingWriterService, andKafkaEventPublisherprovide class-based APIs for dependency injection and explicit lifecycle management. The function APIs remain available for simple integrations.
For applications that manage lifecycle explicitly:
import {
KafkaEventPublisher,
LandingWriterService,
ReaderService,
} from '@dilipabc/storage-core'
const reader = new ReaderService()
const kafka = new KafkaEventPublisher()
const writer = new LandingWriterService((event) => kafka.publish(event))
try {
const files = await reader.listNewFiles(readInput)
// Use reader.readFromSource(...) and writer.writeToLanding(...) as needed.
console.log(files.length)
} finally {
await kafka.close()
}Complete file-transfer example
This example shows the normal flow for FTP, SFTP, S3, GCP, Azure, or MinIO sources. The source credential object is normally loaded from KMS or another secret manager.
import {
listNewFiles,
readFromSource,
writeToLanding,
type ReadInput,
type WriteInput,
} from '@dilipabc/storage-core'
const readInput: ReadInput = {
orgId: 'org-123',
fileName: '',
mimeType: '',
fileSizeBytes: 0,
sourceChannel: 'S3_INGESTION',
sourceCredentials: {
provider: 'S3',
access_key: process.env.SOURCE_ACCESS_KEY!,
secret_key: process.env.SOURCE_SECRET_KEY!,
region: 'ap-south-1',
bucket: 'source-bucket',
insuranceCompanyCode: 'ACME',
},
}
const files = await listNewFiles(readInput)
for (const file of files) {
const sourceStream = await readFromSource({
...readInput,
fileName: file.fileName,
mimeType: file.mimeType,
fileSizeBytes: file.fileSizeBytes,
filePath: file.filePath,
})
const writeInput: WriteInput = {
orgId: file.orgId,
insuranceCompanyCode: file.insuranceCompanyCode,
contextFolder: file.claimFolder,
fileName: file.fileName,
mimeType: file.mimeType,
fileSizeBytes: file.fileSizeBytes,
sourceChannel: readInput.sourceChannel,
}
const result = await writeToLanding(sourceStream, writeInput)
console.log(`Uploaded ${file.fileName} to ${result.bucketName}/${result.objectKey}`)
}listNewFiles
listNewFiles returns normalized FileDescriptor objects. The provider is selected using sourceCredentials.provider case-insensitively.
const files = await listNewFiles({
orgId: 'org-123',
fileName: '',
mimeType: '',
fileSizeBytes: 0,
sourceChannel: 'FTP_INGESTION',
sourceCredentials: {
provider: 'FTP',
host: 'ftp.example.com',
port: 21,
user: 'readonly-user',
password: process.env.FTP_PASSWORD!,
bucket: '/incoming',
source_prefix: 'claims',
insuranceCompanyCode: 'ACME',
secure: true,
},
})
for (const file of files) {
console.log({
name: file.fileName,
path: file.filePath,
size: file.fileSizeBytes,
type: file.mimeType,
})
}Provider credential shapes:
// SFTP
{ provider: 'SFTP', host, user, password, port: 22, bucket, source_prefix, insuranceCompanyCode }
// MinIO source
{ provider: 'MINIO', endpoint, access_key, secret_key, bucket, insuranceCompanyCode }
// S3 source
{ provider: 'S3', access_key, secret_key, region, bucket, endpoint?, insuranceCompanyCode }
// Google Cloud Storage source
{ provider: 'GCP', project_id, bucket_name, google_application_credentials, source_prefix?, insuranceCompanyCode }
// Azure Blob source
{ provider: 'AZURE', account_name, account_key, container, endpoint?, source_prefix?, insuranceCompanyCode }readFromSource
Use the descriptor's filePath when reading a file. The returned value is a Node.js Readable and can be piped directly to another service.
import { createWriteStream } from 'node:fs'
import { pipeline } from 'node:stream/promises'
const stream = await readFromSource({
...readInput,
fileName: 'claim.pdf',
mimeType: 'application/pdf',
fileSizeBytes: 1024,
filePath: '2026-09-10/claim-123/claim.pdf',
})
await pipeline(stream, createWriteStream('/tmp/claim.pdf'))writeToLanding
Set the landing environment variables before calling writeToLanding. The function returns a TransferResult after the storage upload.
import { Readable } from 'node:stream'
process.env.STORAGE_PROVIDER = 'S3'
process.env.AWS_REGION = 'ap-south-1'
process.env.AWS_ACCESS_KEY_ID = 'landing-access-key'
process.env.AWS_SECRET_ACCESS_KEY = 'landing-secret-key'
process.env.AWS_BUCKET = 'sentinel-landing'
process.env.KAFKA_BROKER = 'kafka.example.com:9092'
const result = await writeToLanding(
Readable.from(Buffer.from('example file content')),
{
orgId: 'org-123',
insuranceCompanyCode: 'ACME',
contextFolder: 'claim-123',
fileName: 'document.txt',
mimeType: 'text/plain',
fileSizeBytes: 20,
sourceChannel: 'API_INGESTION',
},
)
// result.success, result.message, result.objectKey, result.bucketName,
// result.storageProvider, result.fileSizeBytes
console.log(result)For MinIO, use STORAGE_PROVIDER=MINIO, MINIO_ENDPOINT, MINIO_ACCESS_KEY, MINIO_SECRET_KEY, and MINIO_BUCKET instead.
Email source example
Email content is downloaded while listing. Do not call readFromSource for an email descriptor; wrap the buffered bytes in a stream.
import { Readable } from 'node:stream'
const emailFiles = await listNewFiles({
orgId: 'org-123',
fileName: '',
mimeType: '',
fileSizeBytes: 0,
sourceChannel: 'EMAIL_INGESTION',
sourceCredentials: {
provider: 'EMAIL',
email: '[email protected]',
password: process.env.IMAP_PASSWORD!,
imap_host: 'imap.example.com',
imap_port: 993,
insuranceCompanyCode: 'ACME',
lastProcessedUid: 100,
lastUidValidity: '42',
},
})
for (const file of emailFiles) {
if (!file.emailMeta) continue
const result = await writeToLanding(
Readable.from(file.emailMeta.bufferedContent),
{
orgId: file.orgId,
insuranceCompanyCode: file.insuranceCompanyCode,
contextFolder: file.claimFolder,
fileName: file.fileName,
mimeType: file.mimeType,
fileSizeBytes: file.fileSizeBytes,
sourceChannel: 'EMAIL_INGESTION',
},
)
console.log(result.objectKey, file.emailMeta.isTranscript)
}Error handling
Use StorageCoreError.code for programmatic handling. Do not match provider error message text. Errors also expose operation, retryable, timestamp, userMessage, and toJSON() for safe API responses or structured logs.
Available error codes are INVALID_INPUT, UNSUPPORTED_PROVIDER, CONFIGURATION_ERROR, AUTHENTICATION_ERROR, NOT_FOUND, TIMEOUT, SOURCE_ERROR, STORAGE_ERROR, STREAM_ERROR, EVENT_PUBLISH_ERROR, and UNKNOWN_ERROR.
import { StorageCoreError, listNewFiles } from '@dilipabc/storage-core'
try {
await listNewFiles(input)
} catch (error) {
if (error instanceof StorageCoreError) {
console.error(error.toJSON())
switch (error.code) {
case 'CONFIGURATION_ERROR':
case 'INVALID_INPUT':
console.error('Fix configuration or request input:', error.message)
break
case 'UNSUPPORTED_PROVIDER':
console.error('Provider is not supported:', error.message)
break
default:
console.error('Storage operation failed:', error.message)
}
}
throw error
}Provider support
| Capability | Implemented | Provider values |
|---|---:|---|
| Source reader | yes | FTP, SFTP, MINIO, S3, GCP, AZURE |
| Email listing/attachment extraction | yes | EMAIL |
| Landing writer | yes | MINIO, S3 |
| Landing writer placeholder | no | GCP, AZURE |
Email attachments and generated transcripts are buffered during listing; callers should use the returned emailMeta.bufferedContent rather than call readFromSource for email files.
Generated email transcript filenames use the format email-transcript-{timestamp}-uid-{imapUid}.pdf. Attachment filenames use their sanitized original name with a timestamp, IMAP UID, and attachment index added before the extension.
Build
npm install
npm run typecheck
npm run buildThe package reads landing credentials from the caller process. MinIO uses MINIO_ENDPOINT, MINIO_ACCESS_KEY, MINIO_SECRET_KEY, and MINIO_BUCKET. S3 uses AWS_REGION, AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY, and AWS_BUCKET; AWS_ENDPOINT is optional for S3-compatible endpoints.
See docs/STORAGE_CORE_TECHNICAL_GUIDE.md for the complete API, credential contracts, provider behavior, email flow, object-key format, and Kafka event contract. See docs/NPM_PUBLISHING_GUIDE.md for the complete npmjs publishing and release process.
Publish
Authenticate with npm, verify the package contents, and publish the scoped package publicly:
npm login
npm run typecheck
npm pack --dry-run
npm publish --access publicprepublishOnly automatically runs typecheck and build before publishing. Increment version before each release, for example with npm version patch, npm version minor, or npm version major.
