Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-workflow-trace-context.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"effect": patch
---

Propagate trace context through persisted cluster workflow requests.
Original file line number Diff line number Diff line change
Expand Up @@ -167,12 +167,16 @@ export const make = Effect.gen(function*() {
readonly payload: unknown
}) {
const payload = (options.rpc.payloadSchema as any).make(options.payload)
const span = yield* Effect.orDie(Effect.currentSpan)
const envelope = Envelope.makeRequest<any>({
requestId: yield* sharding.getSnowflake,
address: options.address,
tag: options.rpc._tag as any,
payload,
headers: Headers.empty
headers: Headers.empty,
traceId: span.traceId,
spanId: span.spanId,
sampled: span.sampled
})
yield* sharding.sendOutgoing(
new Message.OutgoingRequest({
Expand Down
46 changes: 45 additions & 1 deletion packages/effect/test/cluster/ClusterWorkflowEngine.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { assert, describe, expect, it } from "@effect/vitest"
import { Cause, Context, DateTime, Duration, Effect, Exit, Fiber, Layer, Option, Result, Schema } from "effect"
import { Cause, Context, DateTime, Duration, Effect, Exit, Fiber, Layer, Option, Result, Schema, Tracer } from "effect"
import { TestClock } from "effect/testing"
import {
ClusterSchema,
Expand Down Expand Up @@ -292,6 +292,50 @@ describe.concurrent("ClusterWorkflowEngine", () => {
assert.strictEqual(envelope.address.shardId.group, "workflow")
}).pipe(Effect.scoped, Effect.provide(TestWorkflowEngine)))

it.effect("propagates trace context to persisted workflow requests", () => {
let callerSpan: Tracer.NativeSpan | undefined
const tracer = Tracer.make({
span(options) {
const span = new Tracer.NativeSpan(options)
if (options.name === "WorkflowEngine.deferredDone") {
callerSpan = span
}
return span
}
})
return Effect.gen(function*() {
const driver = yield* MessageStorage.MemoryDriver
const engine = yield* WorkflowEngine
yield* engine.register(ShardedDeferredWorkflow, () => Effect.void)

const executionId = yield* ShardedDeferredWorkflow.executionId({ id: "trace-context" })
const token = DurableDeferred.tokenFromExecutionId(ShardedDeferred, {
workflow: ShardedDeferredWorkflow,
executionId
})
const journalLength = driver.journal.length
yield* DurableDeferred.done(ShardedDeferred, {
token,
exit: Exit.void
})

const envelope = driver.journal.slice(journalLength).find((envelope) =>
envelope._tag === "Request" &&
envelope.address.entityType === "Workflow/ShardedDeferredWorkflow" &&
envelope.tag === "deferred"
)
assert(envelope?._tag === "Request")
assert(callerSpan)
assert.strictEqual(envelope.traceId, callerSpan.traceId)
assert.strictEqual(envelope.spanId, callerSpan.spanId)
assert.strictEqual(envelope.sampled, callerSpan.sampled)
}).pipe(
Effect.provideService(Tracer.Tracer, tracer),
Effect.scoped,
Effect.provide(TestWorkflowEngine)
)
})

it.effect("routes activities to the workflow shard group after a partial client is cached", () =>
Effect.gen(function*() {
const driver = yield* MessageStorage.MemoryDriver
Expand Down
Loading