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
4 changes: 4 additions & 0 deletions docs/MIGRATION-slogx.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
6 changes: 4 additions & 2 deletions logger/slogx/event.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
77 changes: 77 additions & 0 deletions temporalx/workflow_event.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
package temporalx

import (
"log/slog"
"regexp"
"strings"
"unicode/utf8"

sdklog "go.temporal.io/sdk/log"
"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_]*)+$`)

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
// 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,
// 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) > maxEventName || !eventNamePattern.MatchString(name) {
key = "event.invalid"
name = boundName(name, maxEventName)
}
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)
}
// 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...)
case level >= slog.LevelWarn:
l.Warn(msg, kv...)
case level >= slog.LevelInfo:
l.Info(msg, kv...)
default:
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)"
}
127 changes: 127 additions & 0 deletions temporalx/workflow_event_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,127 @@
package temporalx

import (
"log/slog"
"runtime"
"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)
}
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()
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) {
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", "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) {
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 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")
}
})
}
}
Loading