Skip to content

feat(Node): add Effect.ts 2nd Gen Cloud Functions sample - #1321

Open
jhuleatt wants to merge 5 commits into
mainfrom
feat/effect-functions-sample
Open

jhuleatt wants to merge 5 commits into
mainfrom
feat/effect-functions-sample

Conversation

@jhuleatt

@jhuleatt jhuleatt commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Description

This pull request adds a new code sample under Node/effect-functions demonstrating idiomatic integration of Effect 4.0 (effect@^4.0.0) with Firebase 2nd Generation Cloud Functions.

The sample deliberately leans on what the Firebase and Firestore SDKs already provide (auto-IDs, WriteBatch atomicity and create() preconditions, FieldValue.increment, Eventarc event.time, built-in gRPC retries) and uses Effect only where it adds something the SDKs don't: Schema validation at the boundaries, typed domain errors, dependency injection, and structured logging. Seven source files, readable top to bottom.

Key Architecture & Effect 4.0 Features Demonstrated

  • 2nd Gen Callable Function (createTask):
    • Parameterized configuration via firebase-functions/params (defineInt("MAX_ACTIVE_TASKS_PER_USER", { default: 100 })).
    • Request decoding and validation via Schema.decodeUnknownEffect, Schema.Trim.check(Schema.isNonEmpty()), and Schema.withDecodingDefault; schema issues surface as a typed ValidationError carrying SchemaError.message.
    • Authentication and business-rule checks yielding Schema.TaggedError domain errors directly (yield* new UnauthorizedError(...)).
    • Native Firestore auto-IDs (collection.doc()) for new tasks.
    • Zero-cast Exit & Cause inspection via Cause.findErrorOption mapping DomainError to standard HttpsError status codes.
  • 2nd Gen Firestore Trigger (onTaskWritten):
    • before/after snapshot decoding at the trigger boundary using Schema.decodeUnknownEffect(Task).
    • Allowed status transition validation wrapped in Effect.fn("validateStatusTransition"); invalid transitions are caught with Effect.catchTag, logged as warnings, and the function succeeds so Eventarc never retries a condition that cannot succeed.
    • Atomic, idempotent event processing with a single WriteBatch visible in the trigger itself: batch.create(audit_logs/{event.id}) rejects redeliveries with GrpcStatus.ALREADY_EXISTS (caught with Effect.catchIf), so the FieldValue.increment user-stats writes in the same batch apply at most once per Eventarc event ID. No reads, no transactions, no application-level retry loops.
    • Counter deltas are derived from the before/after status, so a task written directly as completed is counted correctly.
    • Audit timestamps come from the authoritative Eventarc CloudEvent event.time.
  • Serverless Production Patterns:
    • Context.Service repositories (TaskRepository, UserStatsRepository) with a single static layer each (Layer.sync), composed with Layer.mergeAll into a module-scoped ManagedRuntime (appRuntime) for warm-instance reuse.
    • Effect.fn traced repository methods and Effect.annotateLogs / Effect.withLogSpan request context.
    • Custom Logger layer built on Logger.formatStructured + Logger.map + Logger.layer, routing Effect logs into firebase-functions/logger (logger.write) with structured jsonPayload attributes and no synthetic stack traces.
  • Testing:
    • 14 integration tests using vitest and firebase-functions-test, always run against the real Firestore emulator (npm test = firebase emulators:exec --only firestore "vitest run"), covering concurrent and sequential Eventarc duplicate delivery idempotency, event.time propagation, counter deltas across every status transition, invalid-transition handling, and typed FirestoreError on corrupted documents.
    • .github/workflows/test_node.yml now installs Java 21 and firebase-tools and caches the emulator jar so test-2nd-gen exercises real Firestore in CI.

Note

The Node Unit Tests job currently fails before reaching this sample because Node/delete-unused-accounts-cron no longer compiles on main (es6-promise-pool TS2351); that is unrelated to this PR.

Add a comprehensive sample demonstrating Effect.ts integration with Firebase 2nd Gen Cloud Functions.

Key features demonstrated:
- Callable function (createTask) with Schema validation and typed error handling
- Firestore trigger (onTaskWritten) with state transition validation
- Context and Layer dependency injection (Firestore, TaskRepo, AuditRepo, UserStatsRepo)
- Module-scoped ManagedRuntime for warm instance performance
- Custom Logger layer routing Effect logs to firebase-functions/logger
- Structured concurrency with Effect.all and exponential retries with Schedule
- Unit and integration tests with Vitest and Firebase Local Emulator Suite

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces a new sample project demonstrating Firebase 2nd Gen Cloud Functions integrated with Effect.ts, including a callable function, a Firestore trigger, dependency injection, and Vitest tests. The review feedback highlights a critical TypeScript compilation error in the Schema.Record definition. Additionally, it identifies state-tracking bugs in UserStatsRepository during task transitions (especially with archived states) and recommends refactoring to a single, robust updateStats method with numeric deltas. Finally, it suggests using Schedule.intersect instead of Schedule.compose for a more idiomatic retry policy.

taskId: Schema.String,
userId: Schema.String,
action: AuditAction,
details: Schema.Record({ key: Schema.String, value: Schema.Unknown }),

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

critical

In Effect Schema, Schema.Record expects two schema arguments: Schema.Record(key, value). Passing an object with key and value properties will cause a TypeScript compilation error. Update this to use the correct positional arguments.

Suggested change
details: Schema.Record({ key: Schema.String, value: Schema.Unknown }),
details: Schema.Record(Schema.String, Schema.Unknown),

Comment on lines +22 to +128
export interface UserStatsRepositoryShape {
readonly onTaskCreated: (userId: string) => Effect.Effect<void, FirestoreError>;
readonly onTaskCompleted: (userId: string) => Effect.Effect<void, FirestoreError>;
readonly onTaskReopened: (userId: string) => Effect.Effect<void, FirestoreError>;
readonly onTaskDeleted: (
userId: string,
wasCompleted: boolean
) => Effect.Effect<void, FirestoreError>;
}

export class UserStatsRepository extends Context.Tag("UserStatsRepository")<
UserStatsRepository,
UserStatsRepositoryShape
>() {}

export const UserStatsRepositoryLive = Layer.effect(
UserStatsRepository,
Effect.gen(function* () {
const db = yield* FirestoreService;
const statsCol = db.collection("user_stats");

return {
onTaskCreated: (userId) =>
Effect.tryPromise({
try: async () => {
await statsCol.doc(userId).set(
{
userId,
activeTasks: FieldValue.increment(1),
completedTasks: FieldValue.increment(0),
lastUpdated: new Date().toISOString(),
},
{ merge: true }
);
},
catch: (cause) =>
new FirestoreError({
cause,
message: `Failed to update user stats for created task (${userId})`,
}),
}),

onTaskCompleted: (userId) =>
Effect.tryPromise({
try: async () => {
await statsCol.doc(userId).set(
{
userId,
activeTasks: FieldValue.increment(-1),
completedTasks: FieldValue.increment(1),
lastUpdated: new Date().toISOString(),
},
{ merge: true }
);
},
catch: (cause) =>
new FirestoreError({
cause,
message: `Failed to update user stats for completed task (${userId})`,
}),
}),

onTaskReopened: (userId) =>
Effect.tryPromise({
try: async () => {
await statsCol.doc(userId).set(
{
userId,
activeTasks: FieldValue.increment(1),
completedTasks: FieldValue.increment(-1),
lastUpdated: new Date().toISOString(),
},
{ merge: true }
);
},
catch: (cause) =>
new FirestoreError({
cause,
message: `Failed to update user stats for reopened task (${userId})`,
}),
}),

onTaskDeleted: (userId, wasCompleted) =>
Effect.tryPromise({
try: async () => {
await statsCol.doc(userId).set(
{
userId,
activeTasks: wasCompleted
? FieldValue.increment(0)
: FieldValue.increment(-1),
completedTasks: wasCompleted
? FieldValue.increment(-1)
: FieldValue.increment(0),
lastUpdated: new Date().toISOString(),
},
{ merge: true }
);
},
catch: (cause) =>
new FirestoreError({
cause,
message: `Failed to update user stats for deleted task (${userId})`,
}),
}),
};
})

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

The current implementation of UserStatsRepository has several correctness bugs when handling task status transitions (e.g., transitioning to/from archived or deleting archived tasks), which leads to incorrect counts for activeTasks and completedTasks.

Instead of having multiple specific methods that are prone to state-tracking bugs, we can simplify the repository to a single generic updateStats method that accepts numeric deltas. This makes the repository extremely robust and reusable.

export interface UserStatsRepositoryShape {
  readonly updateStats: (
    userId: string,
    deltas: { readonly activeDelta: number; readonly completedDelta: number }
  ) => Effect.Effect<void, FirestoreError>;
}

export class UserStatsRepository extends Context.Tag("UserStatsRepository")<
  UserStatsRepository,
  UserStatsRepositoryShape
>() {}

export const UserStatsRepositoryLive = Layer.effect(
  UserStatsRepository,
  Effect.gen(function* () {
    const db = yield* FirestoreService;
    const statsCol = db.collection("user_stats");

    return {
      updateStats: (userId, { activeDelta, completedDelta }) =>
        Effect.tryPromise({ 
          try: async () => {
            await statsCol.doc(userId).set(
              {
                userId,
                activeTasks: FieldValue.increment(activeDelta),
                completedTasks: FieldValue.increment(completedDelta),
                lastUpdated: new Date().toISOString(),
              },
              { merge: true }
            );
          },
          catch: (cause) =>
            new FirestoreError({
              cause,
              message: `Failed to update user stats for user ${userId}`,
            }),
        }),
    };
  })
);

Comment on lines +200 to +214
yield* Effect.all(
[
auditRepo.record({
id: randomUUID(),
taskId,
userId: task.userId,
action: "created",
details: { title: task.title, priority: task.priority },
timestamp: new Date().toISOString(),
}).pipe(Effect.retry(retryPolicy)),

statsRepo.onTaskCreated(task.userId).pipe(Effect.retry(retryPolicy)),
],
{ concurrency: "unbounded" }
);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

Update the task creation stats logic to use the new generic updateStats method.

Suggested change
yield* Effect.all(
[
auditRepo.record({
id: randomUUID(),
taskId,
userId: task.userId,
action: "created",
details: { title: task.title, priority: task.priority },
timestamp: new Date().toISOString(),
}).pipe(Effect.retry(retryPolicy)),
statsRepo.onTaskCreated(task.userId).pipe(Effect.retry(retryPolicy)),
],
{ concurrency: "unbounded" }
);
yield* Effect.all(
[
auditRepo.record({
id: randomUUID(),
taskId,
userId: task.userId,
action: "created",
details: { title: task.title, priority: task.priority },
timestamp: new Date().toISOString(),
}).pipe(Effect.retry(retryPolicy)),
statsRepo.updateStats(task.userId, { activeDelta: 1, completedDelta: 0 }).pipe(Effect.retry(retryPolicy)),
],
{ concurrency: "unbounded" }
);

Comment on lines +249 to +272
// Determine stats update effect based on transition
let statsEffect: Effect.Effect<void, FirestoreError> = Effect.void;
if (afterTask.status === "completed") {
statsEffect = statsRepo.onTaskCompleted(afterTask.userId);
} else if (beforeTask.status === "completed") {
statsEffect = statsRepo.onTaskReopened(afterTask.userId);
}

// Run audit logging and user statistics update concurrently with retries
yield* Effect.all(
[
auditRepo.record({
id: randomUUID(),
taskId,
userId: afterTask.userId,
action: "status_changed",
details: { from: beforeTask.status, to: afterTask.status },
timestamp: new Date().toISOString(),
}).pipe(Effect.retry(retryPolicy)),

statsEffect.pipe(Effect.retry(retryPolicy)),
],
{ concurrency: "unbounded" }
);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

Compute the precise deltas for activeTasks and completedTasks based on the status transition. This correctly handles transitions to and from the archived state, which were previously causing incorrect statistics.

        const isActive = (status: TaskStatus) => status === "todo" || status === "in_progress";
        const isCompleted = (status: TaskStatus) => status === "completed";

        const activeDelta = (isActive(afterTask.status) ? 1 : 0) - (isActive(beforeTask.status) ? 1 : 0);
        const completedDelta = (isCompleted(afterTask.status) ? 1 : 0) - (isCompleted(beforeTask.status) ? 1 : 0);

        const statsEffect = activeDelta === 0 && completedDelta === 0
          ? Effect.void
          : statsRepo.updateStats(afterTask.userId, { activeDelta, completedDelta });

        // Run audit logging and user statistics update concurrently with retries
        yield* Effect.all(
          [
            auditRepo.record({
              id: randomUUID(),
              taskId,
              userId: afterTask.userId,
              action: "status_changed",
              details: { from: beforeTask.status, to: afterTask.status },
              timestamp: new Date().toISOString(),
            }).pipe(Effect.retry(retryPolicy)),

            statsEffect.pipe(Effect.retry(retryPolicy)),
          ],
          { concurrency: "unbounded" }
        );

Comment on lines +293 to +309
yield* Effect.all(
[
auditRepo.record({
id: randomUUID(),
taskId,
userId: beforeTask.userId,
action: "deleted",
details: { title: beforeTask.title },
timestamp: new Date().toISOString(),
}).pipe(Effect.retry(retryPolicy)),

statsRepo
.onTaskDeleted(beforeTask.userId, beforeTask.status === "completed")
.pipe(Effect.retry(retryPolicy)),
],
{ concurrency: "unbounded" }
);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

Update the task deletion stats logic to use the new generic updateStats method. This ensures that deleting an archived task does not incorrectly decrement activeTasks again.

      const isActive = beforeTask.status === "todo" || beforeTask.status === "in_progress";
      const isCompleted = beforeTask.status === "completed";

      yield* Effect.all(
        [
          auditRepo.record({
            id: randomUUID(),
            taskId,
            userId: beforeTask.userId,
            action: "deleted",
            details: { title: beforeTask.title },
            timestamp: new Date().toISOString(),
          }).pipe(Effect.retry(retryPolicy)),

          statsRepo
            .updateStats(beforeTask.userId, { 
              activeDelta: isActive ? -1 : 0,
              completedDelta: isCompleted ? -1 : 0,
            })
            .pipe(Effect.retry(retryPolicy)),
        ],
        { concurrency: "unbounded" }
      );

Comment on lines +176 to +179
const retryPolicy = Schedule.exponential(Duration.millis(100)).pipe(
Schedule.jittered,
Schedule.compose(Schedule.recurs(3))
);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

In Effect, Schedule.compose is used to feed the output of one schedule as the input of another. To combine a delay schedule with a recursion limit, Schedule.intersect is the idiomatic and clearer approach.

Suggested change
const retryPolicy = Schedule.exponential(Duration.millis(100)).pipe(
Schedule.jittered,
Schedule.compose(Schedule.recurs(3))
);
const retryPolicy = Schedule.exponential(Duration.millis(100)).pipe(
Schedule.jittered,
Schedule.intersect(Schedule.recurs(3))
);

…ly essential Effect

- Replace db.runTransaction + nested Effect re-entry with a single atomic WriteBatch;
  batch.create(audit_logs/{event.id}) rejects redeliveries with ALREADY_EXISTS so
  AuditRepository.recordOnce resolves false and stats increments apply at most once
- UserStatsRepository methods become synchronous batch enqueuers with delta math
  (onTaskCreated / onStatusChanged / onTaskDeleted); drop onTaskCompleted/Reopened/Archived
- Use Eventarc event.time for audit timestamps; drop DateTime/Clock
- Schema-backed FirestoreDataConverter (withConverter(schemaConverter(...))) on all
  three collections; drop decodeUserStats
- Use Firestore auto-IDs (collection.doc()) instead of randomUUID
- Sequential decodeTask instead of Effect.all for before/after snapshots
- Tests: in-memory fallback gains batch()/withConverter, assert event.time propagation,
  typed FirestoreError on non-ALREADY_EXISTS commit errors and corrupted documents
- README + JSDoc updated; 17/17 tests pass against the Firestore emulator
…t the emulator only

Apply the complexity review: 1054 -> 628 src lines (10 -> 7 files),
1034 -> 666 test lines, with the same behaviour and one bug fixed.

- Inline the audit/stats WriteBatch into onTaskWritten; delete AuditRepository
  and the UserStatsRepository enqueuers. The Firebase idempotency recipe
  (batch.create(audit_logs/{event.id}) + FieldValue.increment in one batch,
  ALREADY_EXISTS => redelivery) now fits on one screen in the trigger
- Derive counter deltas from before/after status, so a task written directly
  with status 'completed' is no longer counted as active
- Log and succeed on InvalidTransitionError instead of throwing, so Eventarc
  never retries a condition that cannot succeed
- Delete FirestoreService, schema-converter, the self-provided static layers,
  *Live aliases, Shape interfaces, unused TaskRepository.findById/delete,
  TaskNotFoundError, and dead type exports; one way to build each service
- ValidationError carries SchemaError.message directly (drop the Standard
  Schema formatter and dead path-segment branch)
- Move the MAX_ACTIVE_TASKS_PER_USER guard into test setup (IntParam.value()
  returns 0, not the declared default, when the env var is unset)
- Delete the 194-line in-memory Firestore fake; npm test now runs
  firebase emulators:exec. CI installs Java 21 + firebase-tools and caches
  the emulator jar so test-2nd-gen exercises real Firestore
- README and JSDoc updated to match
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant