diff --git a/README.md b/README.md index 5ff9fab..e9655f0 100644 --- a/README.md +++ b/README.md @@ -22,6 +22,7 @@ Whether you're building a startup or scaling an enterprise service, these librar Our philosophy is simple: **production code should be safe, testable, and easy to reason about.** Applications should: + - **Respect system signals** for shutdown/interrupts - **Be easy to start up** and log lifecycle events - **Be safe to run** in different environments without accidentally affecting production @@ -52,7 +53,7 @@ This library is the result of years of building production services that handle - [CSV Processing](packages/common-server/src/csv/README.md) - Type-safe CSV reading and writing - [Encryption](packages/common-server/src/encryption/README.md) - AES-256-GCM and RSA utilities - [Environment](packages/common-server/src/env/README.md) - Environment detection and management - - [Extensions](packages/common-server/src/extensions/README.md) - Stream and iterator utilities + - [Extensions](packages/common/src/extensions/README.md) - Stream and iterator utilities - [Hashing](packages/common-server/src/hash/README.md) - Consistent hashing for feature flags - [HTTP Client](packages/common-server/src/http/README.md) - Axios wrappers with logging and proxies - [Locking](packages/common-server/src/locking/README.md) - Distributed lock interface @@ -103,7 +104,7 @@ class Example extends ServiceBase { } } -await app(new Example()) +await app(new Example()); ``` ### Configuration @@ -117,12 +118,12 @@ For example, we can create a configuration that uses the concept of a `Provided Imagine we load `ProvidedConfig` from [our example](packages/common-server/src/config/loader.test.ts). How do we get the actual value of `/path/to/ssm`? We can resolve each value in parallel: ```typescript -const resolveConfig = async (resolver: ValueProvider, config: ProvidedConfig): Promise => { +const resolveConfig = async (resolver: ValueProvider, config: ProvidedConfig): Promise => { return autoResolve({ host: config.host, - dynamic: async () => resolver.getValue(config.dynamic) - }) -} + dynamic: async () => resolver.getValue(config.dynamic), + }); +}; ``` `autoResolve` will recursively go through the object and automagically resolve any lambda-based promises in parallel. This way you can have throughput limitation on SSM/etc and dynamically resolve your configuration with minimal friction. @@ -140,7 +141,7 @@ You may wonder how you generate a trace. To create a new one you can easily wrap ```typescript await withNewTrace(async () => { // ... -}) +}); ``` If you have a trace already provided (for example via a library like [hapi](https://hapi.dev/tutorials/logging/?lang=en_US)) you can provide a trace ID with: @@ -206,7 +207,7 @@ To see how to easily create a connection to in-memory databases (sqlite) or mysq Other than the trace logging we've mentioned above, our logger supports context-sensitive log wrapping. For example: ```typescript -const envBasedLogger = log.with({ env: currentEnvironment() }) +const envBasedLogger = log.with({ env: currentEnvironment() }); envBasedLogger.info('Booting up!'); envBasedLogger.info('Welcome'); @@ -307,7 +308,7 @@ Need to roll out features gradually or run A/B tests? Our consistent hashing uti ```typescript // Roll out to 10% of users -if (consistentChance(userId, 'new-feature', 0.10)) { +if (consistentChance(userId, 'new-feature', 0.1)) { // User is in the 10% } diff --git a/packages/common-aws/package.json b/packages/common-aws/package.json index 6ffcaf3..354b20a 100644 --- a/packages/common-aws/package.json +++ b/packages/common-aws/package.json @@ -1,6 +1,6 @@ { "name": "@paradoxical-io/common-aws", - "version": "2.5.1", + "version": "2.5.2", "description": "Common aws for paradox services", "main": "./dist/index.js", "types": "./dist/index.d.ts", diff --git a/packages/common-aws/src/dynamo/keys/test/inMemoryKvTable.ts b/packages/common-aws/src/dynamo/keys/test/inMemoryKvTable.ts index 7dafe2a..b788f8d 100644 --- a/packages/common-aws/src/dynamo/keys/test/inMemoryKvTable.ts +++ b/packages/common-aws/src/dynamo/keys/test/inMemoryKvTable.ts @@ -1,4 +1,4 @@ -import { Streams } from '@paradoxical-io/common-server'; +import { Streams } from '@paradoxical-io/common'; import { CompoundKey, notNullOrUndefined } from '@paradoxical-io/types'; import { KeyValueList, KeyValueTable } from '../keyTable'; diff --git a/packages/common-hapi/package.json b/packages/common-hapi/package.json index 47c7921..0661fb7 100644 --- a/packages/common-hapi/package.json +++ b/packages/common-hapi/package.json @@ -1,6 +1,6 @@ { "name": "@paradoxical-io/common-hapi", - "version": "2.5.1", + "version": "2.5.2", "description": "Common hapi code for paradoxical services", "files": [ "dist/**/*.js", diff --git a/packages/common-server/package.json b/packages/common-server/package.json index 224a7c6..004e59f 100644 --- a/packages/common-server/package.json +++ b/packages/common-server/package.json @@ -1,6 +1,6 @@ { "name": "@paradoxical-io/common-server", - "version": "2.5.1", + "version": "2.5.2", "description": "Common code for paradox services", "main": "./dist/index.js", "types": "./dist/index.d.ts", diff --git a/packages/common-server/src/csv/streamableCsv.ts b/packages/common-server/src/csv/streamableCsv.ts index be64ad0..07def5e 100644 --- a/packages/common-server/src/csv/streamableCsv.ts +++ b/packages/common-server/src/csv/streamableCsv.ts @@ -1,9 +1,8 @@ +import { TypedReadable, TypedTransformable } from '@paradoxical-io/common'; import { Brand } from '@paradoxical-io/types'; import { stringify } from 'csv'; import stream from 'stream'; -import { TypedReadable, TypedTransformable } from '../extensions'; - export interface StreamableOptions { dateToISO?: boolean; } diff --git a/packages/common-server/src/extensions/index.ts b/packages/common-server/src/extensions/index.ts deleted file mode 100644 index d1a032c..0000000 --- a/packages/common-server/src/extensions/index.ts +++ /dev/null @@ -1 +0,0 @@ -export * from './streams'; diff --git a/packages/common-server/src/index.ts b/packages/common-server/src/index.ts index 9730fe2..0a626ab 100644 --- a/packages/common-server/src/index.ts +++ b/packages/common-server/src/index.ts @@ -5,7 +5,6 @@ export * from './contracts'; export * from './csv'; export * from './encryption'; export * from './env'; -export * from './extensions'; export * from './hash'; export * from './http'; export * from './locking'; diff --git a/packages/common-server/src/sftp/sftp.itest.ts b/packages/common-server/src/sftp/sftp.itest.ts index 433cf3c..374e2e0 100644 --- a/packages/common-server/src/sftp/sftp.itest.ts +++ b/packages/common-server/src/sftp/sftp.itest.ts @@ -1,8 +1,8 @@ +import { Streams } from '@paradoxical-io/common'; import { safeExpect } from '@paradoxical-io/common-test'; import { Readable } from 'stream'; import { CsvStreamWriter } from '../csv'; -import { Streams } from '../extensions'; import { newSftpDocker } from './docker'; import { Sftp } from './sftp'; diff --git a/packages/common-sql/package.json b/packages/common-sql/package.json index 7716b51..1ae28a4 100644 --- a/packages/common-sql/package.json +++ b/packages/common-sql/package.json @@ -1,6 +1,6 @@ { "name": "@paradoxical-io/common-sql", - "version": "2.5.1", + "version": "2.5.2", "files": [ "dist/**/*.js", "dist/**/*.d.ts", diff --git a/packages/common-test/package.json b/packages/common-test/package.json index f0ed605..9cde306 100644 --- a/packages/common-test/package.json +++ b/packages/common-test/package.json @@ -1,6 +1,6 @@ { "name": "@paradoxical-io/common-test", - "version": "2.5.1", + "version": "2.5.2", "description": "Common test-code for jest", "files": [ "dist/**/*.js", diff --git a/packages/common/.eslintrc.js b/packages/common/.eslintrc.js index b8fbae9..1732632 100644 --- a/packages/common/.eslintrc.js +++ b/packages/common/.eslintrc.js @@ -41,7 +41,6 @@ module.exports = { 'readline', 'repl', 'smalloc', - 'stream', 'string_decoder', 'sys', 'timers', diff --git a/packages/common/package.json b/packages/common/package.json index 4094aa9..c807f88 100644 --- a/packages/common/package.json +++ b/packages/common/package.json @@ -1,6 +1,6 @@ { "name": "@paradoxical-io/common", - "version": "2.5.1", + "version": "2.5.2", "description": "Common code for all paradox projects", "files": [ "dist/**/*.js", diff --git a/packages/common-server/src/extensions/README.md b/packages/common/src/extensions/README.md similarity index 88% rename from packages/common-server/src/extensions/README.md rename to packages/common/src/extensions/README.md index 85999b5..684bb8f 100644 --- a/packages/common-server/src/extensions/README.md +++ b/packages/common/src/extensions/README.md @@ -18,7 +18,7 @@ A comprehensive utility library for working with Node.js streams, async iterator Convert Node.js streams into more manageable data structures: ```typescript -import { Streams } from '@paradoxical-io/common-server/extensions'; +import { Streams } from '@paradoxical-io/common/extensions'; import fs from 'fs'; // Convert a byte stream to a Buffer @@ -35,7 +35,7 @@ const users = await Streams.toArray(objectStream); Transform any paginated API into a streamable async iterator: ```typescript -import { Streams } from '@paradoxical-io/common-server/extensions'; +import { Streams } from '@paradoxical-io/common/extensions'; // Example: Paginate through database results interface DbResponse { @@ -45,8 +45,8 @@ interface DbResponse { const userStream = Streams.pagingAsyncIterator( undefined, // start page/token - async (token) => await db.query({ nextToken: token }), - (response) => { + async token => await db.query({ nextToken: token }), + response => { if (!response.items.length) return undefined; return [response.nextToken, response.items]; } @@ -63,7 +63,7 @@ for await (const user of userStream) { Process async iterators with functional operations: ```typescript -import { Streams } from '@paradoxical-io/common-server/extensions'; +import { Streams } from '@paradoxical-io/common/extensions'; // Take only first N items const firstTen = Streams.takeAsync(dataStream, 10); @@ -75,10 +75,7 @@ for await (const batch of batched) { } // Transform stream data -const transformed = Streams.mapAsync( - userStream, - async (user) => await enrichUserData(user) -); +const transformed = Streams.mapAsync(userStream, async user => await enrichUserData(user)); // Collect async iterator to array const allUsers = await Streams.from(userStream); @@ -89,7 +86,7 @@ const allUsers = await Streams.from(userStream); Work with sync generators using familiar functional patterns: ```typescript -import { Streams } from '@paradoxical-io/common-server/extensions'; +import { Streams } from '@paradoxical-io/common/extensions'; function* generateNumbers() { let i = 0; @@ -101,12 +98,12 @@ const firstTen = [...Streams.take(generateNumbers(), 10)]; // [0, 1, 2, 3, 4, 5, 6, 7, 8, 9] // Drop items while predicate is true -const afterFive = Streams.dropWhile(generateNumbers(), (n) => n < 5); +const afterFive = Streams.dropWhile(generateNumbers(), n => n < 5); const next5 = [...Streams.take(afterFive, 5)]; // [5, 6, 7, 8, 9] // Take items while predicate is true -const lessThanTen = Streams.takeWhile(generateNumbers(), (n) => n < 10); +const lessThanTen = Streams.takeWhile(generateNumbers(), n => n < 10); const result = [...lessThanTen]; // [0, 1, 2, 3, 4, 5, 6, 7, 8, 9] ``` @@ -118,57 +115,69 @@ const result = [...lessThanTen]; #### Static Methods **`toBuffer(readStream: Readable): Promise`** + - Reads a byte stream into a single Buffer - Handles backpressure automatically - Throws on stream errors **`toArray(readStream: TypedReadable): Promise`** + - Reads a typed object stream into an array - Useful for collecting stream results - Handles both 'end' and 'close' events **`pagingAsyncIterator(start, next, extract): AsyncGenerator`** + - Converts paginated API calls into async iterator - `start`: Initial page/token value - `next`: Function to fetch next page - `extract`: Function to extract [nextPage, results] from response **`grouped(iterator: AsyncGenerator, size: number): AsyncGenerator`** + - Groups async iterator items into fixed-size batches - Yields arrays of specified size - Last batch may be smaller **`takeAsync(iterator: AsyncGenerator, size: number): AsyncGenerator`** + - Takes first N items from async iterator - Returns new async iterator **`take(iterator: Generator, size: number): Generator`** + - Takes first N items from sync iterator - Returns new sync iterator **`takeWhile(iterator: Generator, predicate: (d: T) => boolean): Generator`** + - Takes items while predicate returns true - Stops at first false result **`dropWhile(iterator: Generator, predicate: (d: T) => boolean): Generator`** + - Drops items while predicate returns true - Yields remaining items **`map(iterator: AsyncGenerator, mapper: (data: T) => Y): AsyncGenerator`** + - Synchronously transforms async iterator values - Mapper function executes synchronously **`mapAsync(iterator: AsyncGenerator, mapper: (data: T) => Promise): AsyncGenerator`** + - Asynchronously transforms async iterator values - Awaits mapper function for each item **`from(iterator: AsyncGenerator): Promise`** + - Collects all async iterator values into array - Awaits completion of iterator ## Type Definitions ### TypedReadable + A Node.js Readable stream that pushes typed objects instead of buffers. ```typescript @@ -176,6 +185,7 @@ type TypedReadable = stream.Readable & { push(data: T): void }; ``` ### TypedTransformable + A Node.js Transform stream that accepts typed objects. ```typescript @@ -190,8 +200,8 @@ type TypedTransformable = stream.Transform & { write(data: T): void }; // Fetch all users in batches of 50 const allUsers = Streams.pagingAsyncIterator( 0, // start page - async (page) => await api.getUsers({ page, limit: 50 }), - (response) => { + async page => await api.getUsers({ page, limit: 50 }), + response => { if (response.users.length === 0) return undefined; return [page + 1, response.users]; } @@ -210,10 +220,7 @@ for await (const batch of batched) { // Complex transformation pipeline const results = await Streams.from( Streams.takeAsync( - Streams.mapAsync( - Streams.grouped(dataStream, 100), - async (batch) => await processBatch(batch) - ), + Streams.mapAsync(Streams.grouped(dataStream, 100), async batch => await processBatch(batch)), 10 // Take first 10 processed batches ) ); @@ -223,7 +230,7 @@ const results = await Streams.from( ```typescript import fs from 'fs'; -import { Streams } from '@paradoxical-io/common-server/extensions'; +import { Streams } from '@paradoxical-io/common/extensions'; // Read large file without loading into memory const stream = fs.createReadStream('huge-file.bin'); diff --git a/packages/common/src/extensions/index.ts b/packages/common/src/extensions/index.ts index 52a0a73..c4483be 100644 --- a/packages/common/src/extensions/index.ts +++ b/packages/common/src/extensions/index.ts @@ -4,4 +4,5 @@ export * from './maps'; export * from './math'; export * from './object'; export * from './sets'; +export * from './streams'; export * from './strings'; diff --git a/packages/common-server/src/extensions/streams.test.ts b/packages/common/src/extensions/streams.test.ts similarity index 100% rename from packages/common-server/src/extensions/streams.test.ts rename to packages/common/src/extensions/streams.test.ts diff --git a/packages/common-server/src/extensions/streams.ts b/packages/common/src/extensions/streams.ts similarity index 100% rename from packages/common-server/src/extensions/streams.ts rename to packages/common/src/extensions/streams.ts diff --git a/packages/types/package.json b/packages/types/package.json index a383594..dab8617 100644 --- a/packages/types/package.json +++ b/packages/types/package.json @@ -1,6 +1,6 @@ { "name": "@paradoxical-io/types", - "version": "2.5.1", + "version": "2.5.2", "description": "Paradox shared types", "files": [ "dist/**/*.js",