From 13959d45f30df01aeafbac48c0bcad256e7006de Mon Sep 17 00:00:00 2001 From: Duynh ne Date: Thu, 24 Sep 2026 10:23:37 +0700 Subject: [PATCH 1/2] Emit catalog events from workflow code Most of order's catalog events are decided in workflow code: the saga's failed/compensated outcome, each compensation step, a completion or inventory commit that exhausted its retries. Only the SDK's replay-aware logger may run there, and the facade may not. temporalx.WorkflowEvent writes through workflow.GetLogger, so a replayed history writes nothing; a live run writes the event once, and the replay test proves the code ran. It applies the facade's grammar: an invalid name goes under event.invalid and a caller's own event attribute is renamed. --- docs/MIGRATION-slogx.md | 4 ++ logger/slogx/event.go | 6 +- temporalx/workflow_event.go | 52 ++++++++++++++ temporalx/workflow_event_test.go | 113 +++++++++++++++++++++++++++++++ 4 files changed, 173 insertions(+), 2 deletions(-) create mode 100644 temporalx/workflow_event.go create mode 100644 temporalx/workflow_event_test.go diff --git a/docs/MIGRATION-slogx.md b/docs/MIGRATION-slogx.md index aa3f968..e69d716 100644 --- a/docs/MIGRATION-slogx.md +++ b/docs/MIGRATION-slogx.md @@ -127,6 +127,10 @@ Four catalog names are emitted by pkg, not written by hand at a call site: `SignalWithStartWorkflow` and for `WorkflowIDConflictPolicy: USE_EXISTING`. A service that starts its workflows only through SignalWithStart (checkout's abandonment timer) therefore writes no started event. +- **Events from workflow code** — a saga outcome, a compensation step, a retry + that gave up — go through `temporalx.WorkflowEvent(ctx, level, name, msg, attrs...)`, + which writes via the SDK's replay-aware logger: a replayed history writes + nothing. Never call the facade from workflow code (temporalx v0.42.0). - `temporal.workflow.failed` — call `temporalx.WorkflowFailed(ctx, log.Slog(), type, status, slog.String("order.id", id))` where a dispatcher, reconciler or activity observes a run that ended failed, terminated or timed out. Never from workflow code: it is replayed. diff --git a/logger/slogx/event.go b/logger/slogx/event.go index d8cb203..636282b 100644 --- a/logger/slogx/event.go +++ b/logger/slogx/event.go @@ -26,8 +26,10 @@ const ( // segments, at most 64 bytes. Only Event may set the "event" attribute; a // plain Info with slog.String("event", …) is not a catalog event and the // registry lint will flag it. The one carve-out is temporalx, which owns the -// two temporal.workflow.* names and writes the key directly so the module -// does not depend on this facade; the lint must allowlist it. +// two temporal.workflow.* names and carries events written from workflow code +// (WorkflowEvent, through the SDK's replay-aware logger); it writes the key +// directly so the module does not depend on this facade, and the lint must +// allowlist it. // // A name that fails the grammar is not silently dropped and not silently // emitted as an event: the record goes out at the requested level with diff --git a/temporalx/workflow_event.go b/temporalx/workflow_event.go new file mode 100644 index 0000000..27ea121 --- /dev/null +++ b/temporalx/workflow_event.go @@ -0,0 +1,52 @@ +package temporalx + +import ( + "log/slog" + "regexp" + + "go.temporal.io/sdk/workflow" +) + +// eventNamePattern is the catalog grammar (RFC-0031 § Event catalog): dotted +// lowercase segments of [a-z][a-z0-9_]*, at least two, at most 64 bytes. +var eventNamePattern = regexp.MustCompile(`^[a-z][a-z0-9_]*(\.[a-z][a-z0-9_]*)+$`) + +// WorkflowEvent emits a catalog event from WORKFLOW code — a saga outcome, a +// compensation step, a retry that gave up — where only the SDK's replay-aware +// logger may run. It writes through workflow.GetLogger, so a replayed history +// writes nothing and a restarted worker never emits the event twice; the +// logger WithLogger installed carries it to the platform handler. Like the +// workflow-lifecycle events, it writes the "event" key directly because this +// module does not depend on the logging facade. +// +// A name outside the grammar is written under "event.invalid" instead, and a +// caller's own "event" attribute is renamed, as the facade does. Activities +// and ordinary code use the facade's Event, not this. +func WorkflowEvent(ctx workflow.Context, level slog.Level, name, msg string, attrs ...slog.Attr) { + key := "event" + if len(name) > 64 || !eventNamePattern.MatchString(name) { + key = "event.invalid" + if len(name) > 64 { + name = name[:64] + } + } + kv := make([]any, 0, len(attrs)+1) + kv = append(kv, slog.String(key, name)) + for _, a := range attrs { + if a.Key == "event" || a.Key == "event.invalid" { + a.Key += ".conflict" + } + kv = append(kv, a) + } + l := workflow.GetLogger(ctx) + switch { + case level >= slog.LevelError: + l.Error(msg, kv...) + case level >= slog.LevelWarn: + l.Warn(msg, kv...) + case level >= slog.LevelInfo: + l.Info(msg, kv...) + default: + l.Debug(msg, kv...) + } +} diff --git a/temporalx/workflow_event_test.go b/temporalx/workflow_event_test.go new file mode 100644 index 0000000..e745703 --- /dev/null +++ b/temporalx/workflow_event_test.go @@ -0,0 +1,113 @@ +package temporalx + +import ( + "log/slog" + "strings" + "sync/atomic" + "testing" + + "go.temporal.io/sdk/testsuite" + "go.temporal.io/sdk/worker" + "go.temporal.io/sdk/workflow" +) + +var eventRuns atomic.Int32 + +// eventWorkflow emits one catalog event and returns. It is registered under +// the logOnceWorkflow name for replay, whose recorded history (start, one +// decision, complete) fits any workflow that schedules nothing. +func eventWorkflow(ctx workflow.Context) error { + eventRuns.Add(1) + WorkflowEvent(ctx, slog.LevelInfo, "order.failed", "order failed", + slog.String("order.id", "8"), slog.String("outcome", "compensated")) + return nil +} + +func eventRecords(h *capture) []slog.Record { + h.mu.Lock() + defer h.mu.Unlock() + var out []slog.Record + for _, r := range h.recs { + if attrOf(r, "event") != "" || attrOf(r, "event.invalid") != "" { + out = append(out, r) + } + } + return out +} + +// A live run writes the event once with its attributes; replaying the same +// history writes it zero times — and the replay provably ran the code. +func TestWorkflowEvent_LiveOnceReplayNever(t *testing.T) { + live := &capture{} + var s testsuite.WorkflowTestSuite + s.SetLogger(sdkLogger(t, live)) + env := s.NewTestWorkflowEnvironment() + env.RegisterWorkflow(eventWorkflow) + env.ExecuteWorkflow(eventWorkflow) + if !env.IsWorkflowCompleted() || env.GetWorkflowError() != nil { + t.Fatalf("live run: completed=%v err=%v", env.IsWorkflowCompleted(), env.GetWorkflowError()) + } + ev := eventRecords(live) + if len(ev) != 1 { + t.Fatalf("live events = %d, want 1", len(ev)) + } + if r := ev[0]; attrOf(r, "event") != "order.failed" || attrOf(r, "outcome") != "compensated" || + attrOf(r, "order.id") != "8" || r.Level != slog.LevelInfo { + t.Errorf("event = %v", r) + } + + replayed := &capture{} + before := eventRuns.Load() + replayer := worker.NewWorkflowReplayer() + replayer.RegisterWorkflowWithOptions(eventWorkflow, workflow.RegisterOptions{Name: "logOnceWorkflow"}) + if err := replayer.ReplayWorkflowHistoryFromJSONFile(sdkLogger(t, replayed), "testdata/history_log_once.json"); err != nil { + t.Fatalf("replay: %v", err) + } + if eventRuns.Load() == before { + t.Fatal("replay never ran the workflow code; the zero-event check would be vacuous") + } + if n := len(eventRecords(replayed)); n != 0 { + t.Errorf("replay wrote %d events, want 0", n) + } +} + +func TestWorkflowEvent_LevelsAndGrammar(t *testing.T) { + cases := []struct { + name string + level slog.Level + ev string + key string + want slog.Level + }{ + {"error", slog.LevelError, "order.retry.exhausted", "event", slog.LevelError}, + {"warn", slog.LevelWarn, "order.compensation.completed", "event", slog.LevelWarn}, + {"debug", slog.LevelDebug, "order.debug_probe", "event", slog.LevelDebug}, + {"bad grammar", slog.LevelInfo, "Order Failed", "event.invalid", slog.LevelInfo}, + {"too long", slog.LevelInfo, "a." + strings.Repeat("b", 80), "event.invalid", slog.LevelInfo}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + h := &capture{} + var s testsuite.WorkflowTestSuite + s.SetLogger(sdkLogger(t, h)) + env := s.NewTestWorkflowEnvironment() + wf := func(ctx workflow.Context) error { + WorkflowEvent(ctx, tc.level, tc.ev, "m", slog.String("event", "spoof")) + return nil + } + env.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{Name: "wf"}) + env.ExecuteWorkflow("wf") + ev := eventRecords(h) + if len(ev) != 1 { + t.Fatalf("events = %d", len(ev)) + } + r := ev[0] + if attrOf(r, tc.key) == "" || r.Level != tc.want || len(attrOf(r, tc.key)) > 64 { + t.Errorf("record = %v %v", r.Level, r) + } + if attrOf(r, "event.conflict") != "spoof" { + t.Error("a caller's own event attribute must be renamed, not overwrite the name") + } + }) + } +} From 3140a71eae7bc9b462f8149fa43c30062a0791ca Mon Sep 17 00:00:00 2001 From: Duynh ne Date: Thu, 24 Sep 2026 10:41:02 +0700 Subject: [PATCH 2/2] Bound invalid names and name the calling workflow An invalid event name is now cut the way the facade cuts it: invalid UTF-8 repaired, the cut on a rune boundary, and a truncation marker, so event.invalid reads the same on both paths. The SDK logger skips one more frame, so the record's source is the workflow line that called WorkflowEvent rather than temporalx. The doc states when replay safety holds and how non-standard levels map; the grammar tests pin exact values, including the 64-byte edge and a cut inside a rune. --- temporalx/workflow_event.go | 41 +++++++++++++++++++++++++------- temporalx/workflow_event_test.go | 28 ++++++++++++++++------ 2 files changed, 54 insertions(+), 15 deletions(-) diff --git a/temporalx/workflow_event.go b/temporalx/workflow_event.go index 27ea121..0ab13ea 100644 --- a/temporalx/workflow_event.go +++ b/temporalx/workflow_event.go @@ -3,7 +3,10 @@ package temporalx import ( "log/slog" "regexp" + "strings" + "unicode/utf8" + sdklog "go.temporal.io/sdk/log" "go.temporal.io/sdk/workflow" ) @@ -11,6 +14,8 @@ import ( // lowercase segments of [a-z][a-z0-9_]*, at least two, at most 64 bytes. var eventNamePattern = regexp.MustCompile(`^[a-z][a-z0-9_]*(\.[a-z][a-z0-9_]*)+$`) +const maxEventName = 64 + // WorkflowEvent emits a catalog event from WORKFLOW code — a saga outcome, a // compensation step, a retry that gave up — where only the SDK's replay-aware // logger may run. It writes through workflow.GetLogger, so a replayed history @@ -19,16 +24,18 @@ var eventNamePattern = regexp.MustCompile(`^[a-z][a-z0-9_]*(\.[a-z][a-z0-9_]*)+$ // workflow-lifecycle events, it writes the "event" key directly because this // module does not depend on the logging facade. // -// A name outside the grammar is written under "event.invalid" instead, and a -// caller's own "event" attribute is renamed, as the facade does. Activities -// and ordinary code use the facade's Event, not this. +// A name outside the grammar is written under "event.invalid" instead, +// bounded as the facade bounds it, and a caller's own "event" attribute is +// renamed. The SDK logger has no level parameter, so a non-standard level is +// rounded down to Debug, Info, Warn or Error. Replay safety holds while +// worker.Options.EnableLoggingInReplay stays false (the default) and no +// workflow interceptor replaces GetLogger with a logger that is not +// replay-aware. Activities and ordinary code use the facade's Event, not this. func WorkflowEvent(ctx workflow.Context, level slog.Level, name, msg string, attrs ...slog.Attr) { key := "event" - if len(name) > 64 || !eventNamePattern.MatchString(name) { + if len(name) > maxEventName || !eventNamePattern.MatchString(name) { key = "event.invalid" - if len(name) > 64 { - name = name[:64] - } + name = boundName(name, maxEventName) } kv := make([]any, 0, len(attrs)+1) kv = append(kv, slog.String(key, name)) @@ -38,7 +45,8 @@ func WorkflowEvent(ctx workflow.Context, level slog.Level, name, msg string, att } kv = append(kv, a) } - l := workflow.GetLogger(ctx) + // Skip one frame so the record's source is WorkflowEvent's caller. + l := sdklog.Skip(workflow.GetLogger(ctx), 1) switch { case level >= slog.LevelError: l.Error(msg, kv...) @@ -50,3 +58,20 @@ func WorkflowEvent(ctx workflow.Context, level slog.Level, name, msg string, att l.Debug(msg, kv...) } } + +// boundName repairs invalid UTF-8 and cuts s to at most limit bytes on a rune +// boundary with a marker — the facade's bound, copied because this module does +// not import it, so an event.invalid value reads the same on both paths. +func boundName(s string, limit int) string { + if !utf8.ValidString(s) { + s = strings.ToValidUTF8(s, "\uFFFD") + } + if len(s) <= limit { + return s + } + cut := limit + for cut > 0 && !utf8.RuneStart(s[cut]) { + cut-- + } + return s[:cut] + "…(truncated)" +} diff --git a/temporalx/workflow_event_test.go b/temporalx/workflow_event_test.go index e745703..ef227a1 100644 --- a/temporalx/workflow_event_test.go +++ b/temporalx/workflow_event_test.go @@ -2,6 +2,7 @@ package temporalx import ( "log/slog" + "runtime" "strings" "sync/atomic" "testing" @@ -55,6 +56,9 @@ func TestWorkflowEvent_LiveOnceReplayNever(t *testing.T) { attrOf(r, "order.id") != "8" || r.Level != slog.LevelInfo { t.Errorf("event = %v", r) } + if f, _ := runtime.CallersFrames([]uintptr{ev[0].PC}).Next(); !strings.HasSuffix(f.File, "workflow_event_test.go") { + t.Errorf("source = %s:%d, want the calling workflow, not temporalx", f.File, f.Line) + } replayed := &capture{} before := eventRuns.Load() @@ -72,18 +76,25 @@ func TestWorkflowEvent_LiveOnceReplayNever(t *testing.T) { } func TestWorkflowEvent_LevelsAndGrammar(t *testing.T) { + exactly64 := "a." + strings.Repeat("b", 62) cases := []struct { name string level slog.Level ev string key string + value string want slog.Level }{ - {"error", slog.LevelError, "order.retry.exhausted", "event", slog.LevelError}, - {"warn", slog.LevelWarn, "order.compensation.completed", "event", slog.LevelWarn}, - {"debug", slog.LevelDebug, "order.debug_probe", "event", slog.LevelDebug}, - {"bad grammar", slog.LevelInfo, "Order Failed", "event.invalid", slog.LevelInfo}, - {"too long", slog.LevelInfo, "a." + strings.Repeat("b", 80), "event.invalid", slog.LevelInfo}, + {"error", slog.LevelError, "order.retry.exhausted", "event", "order.retry.exhausted", slog.LevelError}, + {"warn", slog.LevelWarn, "order.compensation.completed", "event", "order.compensation.completed", slog.LevelWarn}, + {"debug", slog.LevelDebug, "order.debug_probe", "event", "order.debug_probe", slog.LevelDebug}, + {"custom level rounds down", slog.LevelWarn + 2, "order.failed", "event", "order.failed", slog.LevelWarn}, + {"exactly 64 bytes is valid", slog.LevelInfo, exactly64, "event", exactly64, slog.LevelInfo}, + {"65 bytes is invalid", slog.LevelInfo, exactly64 + "c", "event.invalid", exactly64 + "…(truncated)", slog.LevelInfo}, + {"cut on a rune boundary", slog.LevelInfo, strings.Repeat("a", 63) + "é", "event.invalid", strings.Repeat("a", 63) + "…(truncated)", slog.LevelInfo}, + {"bad grammar", slog.LevelInfo, "Order Failed", "event.invalid", "Order Failed", slog.LevelInfo}, + {"leading digit segment", slog.LevelInfo, "order.1x", "event.invalid", "order.1x", slog.LevelInfo}, + {"trailing dot", slog.LevelInfo, "order.", "event.invalid", "order.", slog.LevelInfo}, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { @@ -102,8 +113,11 @@ func TestWorkflowEvent_LevelsAndGrammar(t *testing.T) { t.Fatalf("events = %d", len(ev)) } r := ev[0] - if attrOf(r, tc.key) == "" || r.Level != tc.want || len(attrOf(r, tc.key)) > 64 { - t.Errorf("record = %v %v", r.Level, r) + if got := attrOf(r, tc.key); got != tc.value || r.Level != tc.want { + t.Errorf("%s = %q level %v, want %q level %v", tc.key, got, r.Level, tc.value, tc.want) + } + if tc.key == "event.invalid" && attrOf(r, "event") != "" { + t.Error("an invalid name must not be written under event") } if attrOf(r, "event.conflict") != "spoof" { t.Error("a caller's own event attribute must be renamed, not overwrite the name")