feat(api): 访问日志与 HTTP/运行时指标、安全响应头与面板 CSP、跨站写拦截、请求体读截止、SSE 定期重鉴权与续传、404/405/413 信封、status 按所有权裁剪、停机并发排空
This commit is contained in:
24 files changed
+1453
-55
No files matched your search
+101
-14
@@ -5,6 +5,7 @@ import (
|
||||
"context"
|
||||
"io"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||
@@ -75,6 +76,80 @@ var heartbeatInterval = 25 * time.Second
|
||||
// reassigns it.
|
||||
var writeTimeout = 30 * time.Second
|
||||
|
||||
// A stream is authorized once, when it opens, and then may run for hours. So the
|
||||
// relay asks again every streamRecheckEvery (the caller's streamGuard: is the
|
||||
// session still live, does the caller still own the server) and ends the stream
|
||||
// after streamMaxLifetime regardless. A withdrawn grant gets an "event: revoked"
|
||||
// before the close, which the panel treats as final; the lifetime cap is a plain
|
||||
// close, which EventSource answers by reconnecting through the full auth path,
|
||||
// resuming from its Last-Event-ID (see logSinceFromRequest) instead of replaying
|
||||
// the backlog. Both are vars only so a test can shrink them.
|
||||
var (
|
||||
streamRecheckEvery = time.Minute
|
||||
streamMaxLifetime = 30 * time.Minute
|
||||
)
|
||||
|
||||
// streamGuard re-checks a running stream's authorization. nil means it still
|
||||
// holds; an error ends the stream with "event: revoked". A guard that cannot
|
||||
// reach its store should return nil: an outage is no verdict on the caller, and
|
||||
// streamMaxLifetime still bounds how long the stream can outlive a revocation.
|
||||
type streamGuard func(ctx context.Context) error
|
||||
|
||||
// sseRevoked is the last event a stream whose grant was withdrawn receives.
|
||||
const sseRevoked = "event: revoked\ndata: access to this stream was withdrawn\n\n"
|
||||
|
||||
// maxResumeAge bounds how far back a Last-Event-ID may resume: past it the pod
|
||||
// has likely restarted anyway, and the tailed backlog is the better start.
|
||||
const maxResumeAge = time.Hour
|
||||
|
||||
type logSinceKey struct{}
|
||||
|
||||
// withLogSince asks the LogStreamer to start the follow at since instead of the
|
||||
// tailed backlog. It rides the context so the LogStreamer interface (and its test
|
||||
// fakes) stays one method.
|
||||
func withLogSince(ctx context.Context, since time.Time) context.Context {
|
||||
return context.WithValue(ctx, logSinceKey{}, since)
|
||||
}
|
||||
|
||||
func logSinceFromContext(ctx context.Context) (time.Time, bool) {
|
||||
t, ok := ctx.Value(logSinceKey{}).(time.Time)
|
||||
return t, ok && !t.IsZero()
|
||||
}
|
||||
|
||||
// logSinceFromRequest reads the resume point an EventSource sends on reconnect.
|
||||
// Every relayed line carries "id: <unix seconds>", so Last-Event-ID is the second
|
||||
// the client last heard from; resuming at that second may repeat a line or two
|
||||
// from it, which beats replaying 200 lines of backlog on every reconnect. Anything
|
||||
// unparsable, in the future, or older than maxResumeAge is ignored.
|
||||
func logSinceFromRequest(r *http.Request, now time.Time) (time.Time, bool) {
|
||||
raw := r.Header.Get("Last-Event-ID")
|
||||
if raw == "" {
|
||||
return time.Time{}, false
|
||||
}
|
||||
secs, err := strconv.ParseInt(raw, 10, 64)
|
||||
if err != nil {
|
||||
return time.Time{}, false
|
||||
}
|
||||
t := time.Unix(secs, 0)
|
||||
if t.After(now) || now.Sub(t) > maxResumeAge {
|
||||
return time.Time{}, false
|
||||
}
|
||||
return t, true
|
||||
}
|
||||
|
||||
// podLogOptions is the follow request both streamers send: the tailed backlog,
|
||||
// or everything since the resume point when the client is reconnecting.
|
||||
func podLogOptions(ctx context.Context, container string, tail int64) *corev1.PodLogOptions {
|
||||
opts := &corev1.PodLogOptions{Container: container, Follow: true}
|
||||
if since, ok := logSinceFromContext(ctx); ok {
|
||||
st := metav1.NewTime(since)
|
||||
opts.SinceTime = &st
|
||||
return opts
|
||||
}
|
||||
opts.TailLines = &tail
|
||||
return opts
|
||||
}
|
||||
|
||||
// relayLogStream is the shared §8 read-side relay: it copies a line-oriented log
|
||||
// source to the client as Server-Sent Events (spec §262 SSE, NOT WebSocket). It
|
||||
// is the single reusable artifact the server console (handleServerConsole) and,
|
||||
@@ -91,7 +166,11 @@ var writeTimeout = 30 * time.Second
|
||||
// r.Context(), so a client disconnect cancels it, the underlying Read errors, the
|
||||
// scan loop exits, and the deferred Close releases the upstream stream (no leaked
|
||||
// apiserver connection).
|
||||
func relayLogStream(w http.ResponseWriter, r *http.Request, src io.ReadCloser) {
|
||||
//
|
||||
// still, when non-nil, is asked every streamRecheckEvery whether the caller may
|
||||
// keep reading; the stream also ends after streamMaxLifetime (see streamGuard),
|
||||
// and when closing is closed (API.CloseStreams, at shutdown).
|
||||
func relayLogStream(w http.ResponseWriter, r *http.Request, src io.ReadCloser, still streamGuard, closing <-chan struct{}) {
|
||||
defer src.Close()
|
||||
|
||||
// SSE needs per-event flushing; without a Flusher the bytes buffer and never
|
||||
@@ -162,6 +241,10 @@ func relayLogStream(w http.ResponseWriter, r *http.Request, src io.ReadCloser) {
|
||||
|
||||
ticker := time.NewTicker(heartbeatInterval)
|
||||
defer ticker.Stop()
|
||||
recheck := time.NewTicker(streamRecheckEvery)
|
||||
defer recheck.Stop()
|
||||
expire := time.NewTimer(streamMaxLifetime)
|
||||
defer expire.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
@@ -174,9 +257,21 @@ func relayLogStream(w http.ResponseWriter, r *http.Request, src io.ReadCloser) {
|
||||
// One log line → one SSE "data:" event, written under a per-write deadline
|
||||
// so a stalled reader severs the stream (writeChunk → false) instead of
|
||||
// pinning this goroutine; the deferred Close then tears the upstream down.
|
||||
if !writeChunk(rc, w, "data: "+line+"\n\n") {
|
||||
// The id is the resume point a reconnecting EventSource sends back.
|
||||
if !writeChunk(rc, w, "id: "+strconv.FormatInt(time.Now().Unix(), 10)+"\ndata: "+line+"\n\n") {
|
||||
return
|
||||
}
|
||||
case <-recheck.C:
|
||||
if still != nil && still(ctx) != nil {
|
||||
writeChunk(rc, w, sseRevoked)
|
||||
return
|
||||
}
|
||||
case <-expire.C:
|
||||
// A plain close: the client reconnects and is authorized afresh.
|
||||
return
|
||||
case <-closing:
|
||||
// The server is shutting down; the client reconnects to the next one.
|
||||
return
|
||||
case <-ticker.C:
|
||||
// No line for a whole interval: emit a comment so the connection stays
|
||||
// warm past the proxy idle timeout — same deadline-guarded write, so a
|
||||
@@ -282,12 +377,8 @@ func (k *K8sLogStreamer) StreamLogs(ctx context.Context, name string) (io.ReadCl
|
||||
return nil, ErrNotFound
|
||||
}
|
||||
|
||||
tail := k.tailLines
|
||||
stream, err := k.clientset.CoreV1().Pods(k.namespace).GetLogs(podName, &corev1.PodLogOptions{
|
||||
Container: serverLogContainer,
|
||||
Follow: true,
|
||||
TailLines: &tail,
|
||||
}).Stream(ctx)
|
||||
stream, err := k.clientset.CoreV1().Pods(k.namespace).GetLogs(podName,
|
||||
podLogOptions(ctx, serverLogContainer, k.tailLines)).Stream(ctx)
|
||||
if err != nil {
|
||||
return nil, ErrConsoleUnavailable
|
||||
}
|
||||
@@ -354,12 +445,8 @@ func (k *K8sBuildLogStreamer) StreamLogs(ctx context.Context, buildID string) (i
|
||||
// at most one Pod, so the first match is the build's Pod.
|
||||
podName := pods.Items[0].Name
|
||||
|
||||
tail := k.tailLines
|
||||
stream, err := k.clientset.CoreV1().Pods(k.namespace).GetLogs(podName, &corev1.PodLogOptions{
|
||||
Container: build.ContainerKaniko,
|
||||
Follow: true,
|
||||
TailLines: &tail,
|
||||
}).Stream(ctx)
|
||||
stream, err := k.clientset.CoreV1().Pods(k.namespace).GetLogs(podName,
|
||||
podLogOptions(ctx, build.ContainerKaniko, k.tailLines)).Stream(ctx)
|
||||
if err != nil {
|
||||
return nil, ErrConsoleUnavailable
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user