package application import ( "context" "io" "net/http" "git.nianxx.cn/wangxuming/NianAIGC/backend/internal/jobs" "git.nianxx.cn/wangxuming/NianAIGC/backend/internal/logging" ) // EventLogger is the application-owned write seam implemented by // *logging.Service. Keeping the seam here lets runtime composition provide a // different sink without coupling either HTTP handlers or jobs to log storage. type EventLogger interface { Append(context.Context, logging.Input) (logging.Entry, error) } // WithHTTPEventLogging records server failures while deliberately limiting // request metadata to the method and URL path. Headers, cookies, query values, // and bodies never enter the event. func WithHTTPEventLogging(next http.Handler, logger EventLogger) http.Handler { if next == nil { next = http.NotFoundHandler() } return http.HandlerFunc(func(w http.ResponseWriter, request *http.Request) { writer := &eventResponseWriter{ResponseWriter: w, status: http.StatusOK} defer func() { if recovered := recover(); recovered != nil { // Once a response has been streamed, HTTP cannot safely rewrite it. // Before the first byte, however, we can still return a stable generic // response without exposing the panic value. if !writer.wroteHeader { clearResponseHeaders(writer.Header()) http.Error(writer, "Internal Server Error", http.StatusInternalServerError) } appendEvent(request.Context(), logger, logging.Input{ Level: logging.Error, Source: "http", Message: "HTTP handler panic", Status: http.StatusInternalServerError, Method: request.Method, Path: request.URL.Path, }) if writer.wroteHeader && writer.status != http.StatusInternalServerError { // The only safe response after bytes have escaped is to abort the // connection. net/http deliberately suppresses logging for this // sentinel while preventing a truncated response from being reused. panic(http.ErrAbortHandler) } return } if writer.status >= http.StatusInternalServerError { appendEvent(request.Context(), logger, logging.Input{ Level: logging.Error, Source: "http", Message: "HTTP request failed", Status: writer.status, Method: request.Method, Path: request.URL.Path, }) } }() next.ServeHTTP(writer, request) }) } type eventResponseWriter struct { http.ResponseWriter status int wroteHeader bool } func (writer *eventResponseWriter) WriteHeader(status int) { if writer.wroteHeader { return } writer.status = status writer.wroteHeader = true writer.ResponseWriter.WriteHeader(status) } func (writer *eventResponseWriter) Write(body []byte) (int, error) { if !writer.wroteHeader { writer.WriteHeader(http.StatusOK) } return writer.ResponseWriter.Write(body) } // Unwrap lets http.ResponseController discover capabilities provided by the // original server writer without making event logging own those interfaces. func (writer *eventResponseWriter) Unwrap() http.ResponseWriter { return writer.ResponseWriter } // ReadFrom preserves io.Copy's streaming fast path for downloads while still // recording the implicit 200 response status. func (writer *eventResponseWriter) ReadFrom(source io.Reader) (int64, error) { if !writer.wroteHeader { writer.WriteHeader(http.StatusOK) } if readerFrom, ok := writer.ResponseWriter.(io.ReaderFrom); ok { return readerFrom.ReadFrom(source) } return io.Copy(writer.ResponseWriter, source) } func clearResponseHeaders(header http.Header) { for name := range header { delete(header, name) } } type eventLoggingTickRunner struct { next jobs.TickRunner logger EventLogger } // WithTickEventLogging records both whole-tick failures and per-job failures. // Error strings and worker/job inputs are intentionally omitted from events. func WithTickEventLogging(next jobs.TickRunner, logger EventLogger) eventLoggingTickRunner { return eventLoggingTickRunner{next: next, logger: logger} } func (runner eventLoggingTickRunner) Tick(ctx context.Context, workerID string) (jobs.TickResult, error) { return runner.record(ctx, func() (jobs.TickResult, error) { return runner.next.Tick(ctx, workerID) }) } // TickLimit preserves the internal Worker HTTP seam while sharing the same // best-effort event policy as the embedded loop. func (runner eventLoggingTickRunner) TickLimit(ctx context.Context, workerID string, limit int) (jobs.TickResult, error) { limited, ok := runner.next.(interface { TickLimit(context.Context, string, int) (jobs.TickResult, error) }) if !ok { return runner.Tick(ctx, workerID) } return runner.record(ctx, func() (jobs.TickResult, error) { return limited.TickLimit(ctx, workerID, limit) }) } func (runner eventLoggingTickRunner) record(ctx context.Context, tick func() (jobs.TickResult, error)) (jobs.TickResult, error) { result, err := tick() if err != nil { appendEvent(ctx, runner.logger, logging.Input{ Level: logging.Error, Source: "worker", Message: "Worker tick failed", }) return result, err } for _, item := range result.Jobs { if item.Error == "" { continue } appendEvent(ctx, runner.logger, logging.Input{ Level: logging.Error, Source: "worker", Message: "Worker job failed", }) } return result, nil } func appendEvent(ctx context.Context, logger EventLogger, input logging.Input) { if logger == nil { return } defer func() { _ = recover() }() _, _ = logger.Append(context.WithoutCancel(ctx), input) }