amoghrajesh commented on code in PR #71882:
URL: https://github.com/apache/airflow/pull/71882#discussion_r3828170879
##########
go-sdk/pkg/execution/logger.go:
##########
@@ -75,20 +87,121 @@ type socketLogHandlerShared struct {
var _ slog.Handler = (*SocketLogHandler)(nil)
+type logLevelFilter struct {
+ defaultLevel slog.Level
+ namespaceLevels map[string]slog.Level
+}
+
// NewSocketLogHandler creates a new handler. If writer is nil, messages are
// buffered until Connect() is called.
func NewSocketLogHandler(writer io.Writer, level slog.Level) *SocketLogHandler
{
+ return newSocketLogHandler(writer, newLogLevelFilter(level, nil))
+}
+
+func newSocketLogHandlerFromEnv(writer io.Writer) *SocketLogHandler {
+ defaultLevel, ok := parseLogLevel(os.Getenv(loggingLevelEnv))
+ if !ok {
+ defaultLevel = slog.LevelInfo
+ }
+ return newSocketLogHandler(
+ writer,
+ newLogLevelFilter(defaultLevel,
parseNamespaceLogLevels(os.Getenv(namespaceLevelsEnv))),
+ )
+}
+
+func newSocketLogHandler(writer io.Writer, filter *logLevelFilter)
*SocketLogHandler {
shared := &socketLogHandlerShared{}
if writer != nil {
shared.writer = writer
shared.connected = true
}
return &SocketLogHandler{
shared: shared,
- level: level,
+ filter: filter,
+ }
+}
+
+func newLogLevelFilter(
+ defaultLevel slog.Level,
+ namespaceLevels map[string]slog.Level,
+) *logLevelFilter {
+ return &logLevelFilter{
+ defaultLevel: defaultLevel,
+ namespaceLevels: namespaceLevels,
+ }
+}
+
+func parseLogLevel(value string) (slog.Level, bool) {
+ switch strings.ToUpper(strings.TrimSpace(value)) {
+ case "NOTSET":
+ return notsetLogLevel, true
+ case "DEBUG":
+ return slog.LevelDebug, true
+ case "INFO":
+ return slog.LevelInfo, true
+ case "WARN", "WARNING":
+ return slog.LevelWarn, true
+ case "ERROR":
+ return slog.LevelError, true
+ case "CRITICAL", "FATAL":
+ return criticalLogLevel, true
+ default:
+ return 0, false
+ }
+}
+
+func getAirflowLogLevelName(level slog.Level) string {
+ switch {
+ case level < slog.LevelDebug:
+ return "notset"
+ case level < slog.LevelInfo:
+ return "debug"
+ case level < slog.LevelWarn:
+ return "info"
+ case level < slog.LevelError:
+ return "warning"
+ case level < criticalLogLevel:
+ return "error"
+ default:
+ return "critical"
}
}
+func parseNamespaceLogLevels(value string) map[string]slog.Level {
+ levels := make(map[string]slog.Level)
+ entries := strings.FieldsFunc(value, func(r rune) bool {
+ return r == ',' || unicode.IsSpace(r)
+ })
+ for _, entry := range entries {
+ loggerName, levelName, ok := strings.Cut(entry, "=")
+ loggerName = strings.TrimSpace(loggerName)
+ level, validLevel := parseLogLevel(levelName)
+ if !ok || loggerName == "" || !validLevel {
+ continue
+ }
Review Comment:
Invalid entries are skipped silently. The Python equivalent
([structlog.py:211-232](https://github.com/apache/airflow/blob/main/shared/logging/src/airflow_shared/logging/structlog.py#L211-L232))
collects them and logs "Ignoring invalid namespace_levels entry: %s" for each,
so a typo like `example=DEBGU` is visible. Could we do the same?
##########
go-sdk/pkg/execution/logger.go:
##########
@@ -75,20 +87,121 @@ type socketLogHandlerShared struct {
var _ slog.Handler = (*SocketLogHandler)(nil)
+type logLevelFilter struct {
+ defaultLevel slog.Level
+ namespaceLevels map[string]slog.Level
+}
+
// NewSocketLogHandler creates a new handler. If writer is nil, messages are
// buffered until Connect() is called.
func NewSocketLogHandler(writer io.Writer, level slog.Level) *SocketLogHandler
{
+ return newSocketLogHandler(writer, newLogLevelFilter(level, nil))
+}
+
+func newSocketLogHandlerFromEnv(writer io.Writer) *SocketLogHandler {
+ defaultLevel, ok := parseLogLevel(os.Getenv(loggingLevelEnv))
+ if !ok {
+ defaultLevel = slog.LevelInfo
+ }
+ return newSocketLogHandler(
+ writer,
+ newLogLevelFilter(defaultLevel,
parseNamespaceLogLevels(os.Getenv(namespaceLevelsEnv))),
+ )
+}
+
+func newSocketLogHandler(writer io.Writer, filter *logLevelFilter)
*SocketLogHandler {
shared := &socketLogHandlerShared{}
if writer != nil {
shared.writer = writer
shared.connected = true
}
return &SocketLogHandler{
shared: shared,
- level: level,
+ filter: filter,
+ }
+}
+
+func newLogLevelFilter(
+ defaultLevel slog.Level,
+ namespaceLevels map[string]slog.Level,
+) *logLevelFilter {
+ return &logLevelFilter{
+ defaultLevel: defaultLevel,
+ namespaceLevels: namespaceLevels,
+ }
+}
+
+func parseLogLevel(value string) (slog.Level, bool) {
+ switch strings.ToUpper(strings.TrimSpace(value)) {
+ case "NOTSET":
+ return notsetLogLevel, true
+ case "DEBUG":
+ return slog.LevelDebug, true
+ case "INFO":
+ return slog.LevelInfo, true
+ case "WARN", "WARNING":
+ return slog.LevelWarn, true
+ case "ERROR":
+ return slog.LevelError, true
+ case "CRITICAL", "FATAL":
+ return criticalLogLevel, true
+ default:
+ return 0, false
+ }
+}
+
Review Comment:
The accepted level names diverge from Python log levels in both directions.
structlog has no fatal, but it does have exception (mapped to ERROR). So
namespace_levels="foo=exception" filters at ERROR for a Python task and falls
back to the global default for a Go task on the same deployment. Suggest
accepting `EXCEPTION` here, and dropping `FATAL` unless there's a reason to be
more lenient than the Python side.
##########
go-sdk/pkg/execution/logger_test.go:
##########
@@ -64,7 +65,126 @@ func TestSocketLogHandlerLevelFiltering(t *testing.T) {
var entry map[string]any
require.NoError(t, json.Unmarshal([]byte(lines[0]), &entry))
assert.Equal(t, "should appear", entry["event"])
- assert.Equal(t, "warn", entry["level"])
+ assert.Equal(t, "warning", entry["level"])
+}
+
+func TestParseLogLevel(t *testing.T) {
+ tests := map[string]slog.Level{
+ "notset": notsetLogLevel,
+ "DEBUG": slog.LevelDebug,
+ " info ": slog.LevelInfo,
+ "WARN": slog.LevelWarn,
+ "warning": slog.LevelWarn,
+ "ERROR": slog.LevelError,
+ "critical": criticalLogLevel,
+ "FATAL": criticalLogLevel,
+ }
+ for value, expected := range tests {
+ level, ok := parseLogLevel(value)
+ assert.True(t, ok, value)
+ assert.Equal(t, expected, level, value)
+ }
+
+ _, ok := parseLogLevel("verbose")
+ assert.False(t, ok)
+}
+
+func TestGetAirflowLogLevelName(t *testing.T) {
+ tests := map[slog.Level]string{
+ notsetLogLevel: "notset",
+ slog.LevelDebug - 1: "notset",
+ slog.LevelDebug: "debug",
+ slog.LevelDebug + 1: "debug",
+ slog.LevelInfo - 1: "debug",
+ slog.LevelInfo: "info",
+ slog.LevelInfo + 1: "info",
+ slog.LevelWarn - 1: "info",
+ slog.LevelWarn: "warning",
+ slog.LevelWarn + 1: "warning",
+ slog.LevelError - 1: "warning",
+ slog.LevelError: "error",
+ slog.LevelError + 1: "error",
+ criticalLogLevel - 1: "error",
+ criticalLogLevel: "critical",
+ criticalLogLevel + 1: "critical",
+ }
+ for level, expected := range tests {
+ assert.Equal(t, expected, getAirflowLogLevelName(level), level)
+ }
+}
+
+func TestParseNamespaceLogLevels(t *testing.T) {
+ assert.Equal(t, map[string]slog.Level{
+ "example": slog.LevelWarn,
+ "example.detail": slog.LevelDebug,
+ }, parseNamespaceLogLevels(
+ "example=INFO, malformed example.detail=DEBUG example=WARNING
=ERROR empty= unknown=VERBOSE",
+ ))
+}
+
+func TestLogLevelFilterUsesLongestNamespacePrefix(t *testing.T) {
+ filter := newLogLevelFilter(slog.LevelError, map[string]slog.Level{
+ "example": slog.LevelInfo,
+ "example.detail": slog.LevelDebug,
+ })
+
+ assert.Equal(t, slog.LevelDebug,
filter.getLevel("example.detail.child"))
+ assert.Equal(t, slog.LevelInfo, filter.getLevel("example.other"))
+ assert.Equal(t, slog.LevelError, filter.getLevel("exampled"))
+ assert.Equal(t, slog.LevelError, filter.getLevel(""))
+}
+
+func TestSocketLogHandlerUsesEnvironmentLevelFiltering(t *testing.T) {
+ t.Setenv(loggingLevelEnv, "ERROR")
+ t.Setenv(namespaceLevelsEnv, "example=DEBUG, example.noisy=WARNING")
+
+ var buf bytes.Buffer
+ logger := slog.New(newSocketLogHandlerFromEnv(&buf))
+ logger.Info("global filtered")
+ logger.WithGroup("example.detail").Debug("namespace debug")
+ logger.WithGroup("example.noisy.child").Info("namespace filtered")
+ logger.WithGroup("example.noisy.child").Warn("namespace warning")
+ logger.WithGroup("unrelated").Warn("unrelated filtered")
+ logger.Log(context.Background(), slog.LevelDebug+1, "unsupported level")
Review Comment:
This case doesn't cover what its message suggests. The env default in this
test is ERROR, so `slog.LevelDebug+1` is filtered by the global threshold, not
by level normalization, it would pass identically with `getAirflowLogLevelName`
deleted. To actually pin the mapping, this needs a an assertion that the
emitted level is "debug".
##########
go-sdk/pkg/execution/logger.go:
##########
@@ -134,6 +251,11 @@ func (h *SocketLogHandler) Handle(_ context.Context, r
slog.Record) error {
return true
})
+ if loggerName != "" {
+ entry["logger"] = loggerName
+ }
+ entry["level"] = getAirflowLogLevelName(r.Level)
Review Comment:
Moving this below the attribute walk means a user attribute named level (or
logger) is now overwritten by the handler rather than overwriting it. Better
behavior, but its an undeclared change with no test covering it.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]