Skip to content
Open
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
7 changes: 4 additions & 3 deletions pkg/frontier/edgebound/edge_dataplane.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package edgebound

import (
"github.com/singchia/frontier/pkg/frontier/misc"
"github.com/singchia/geminio"
"k8s.io/klog/v2"
)
Expand All @@ -9,7 +10,7 @@ func (em *edgeManager) acceptStream(stream geminio.Stream) {
edgeID := stream.ClientID()
streamID := stream.StreamID()
meta := stream.Meta()
klog.V(2).Infof("edge accept stream, edgeID: %d, streamID: %d, meta: %s", edgeID, streamID, meta)
klog.V(2).Infof("edge accept stream, edgeID: %d, streamID: %d, meta: %s", edgeID, streamID, misc.Redact(string(meta)))

// cache
em.streams.MSet(edgeID, streamID, stream)
Expand All @@ -23,7 +24,7 @@ func (em *edgeManager) closedStream(stream geminio.Stream) {
edgeID := stream.ClientID()
streamID := stream.StreamID()
meta := stream.Meta()
klog.V(2).Infof("edge closed stream, edgeID: %d, streamID: %d, meta: %s", edgeID, streamID, meta)
klog.V(2).Infof("edge closed stream, edgeID: %d, streamID: %d, meta: %s", edgeID, streamID, misc.Redact(string(meta)))
// cache
em.streams.MDel(edgeID, streamID)
// when the stream ends, the exchange can be noticed by functional error, so we don't update exchange
Expand All @@ -33,7 +34,7 @@ func (em *edgeManager) closedStream(stream geminio.Stream) {
func (em *edgeManager) forward(end geminio.End) {
edgeID := end.ClientID()
meta := end.Meta()
klog.V(2).Infof("edge forward raw message and rpc, edgeID: %d, meta: %s", edgeID, meta)
klog.V(2).Infof("edge forward raw message and rpc, edgeID: %d, meta: %s", edgeID, misc.Redact(string(meta)))
if em.exchange != nil {
em.exchange.ForwardToService(end)
}
Expand Down
4 changes: 4 additions & 0 deletions pkg/frontier/edgebound/edge_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"strings"
"sync"

"github.com/jumboframes/armorigo/log"
"github.com/jumboframes/armorigo/rproxy"
"github.com/jumboframes/armorigo/synchub"
"github.com/singchia/frontier/pkg/frontier/apis"
Expand Down Expand Up @@ -141,6 +142,9 @@ func (em *edgeManager) Serve() error {
func (em *edgeManager) handleConn(conn net.Conn) error {
// options for geminio End
opt := server.NewEndOptions()
// route SDK logs through klog: the SDK default logger prints raw
// connection meta to stdout on error paths, bypassing redaction.
opt.SetLog(log.NewKLog())
opt.SetTimer(em.tmr)
opt.SetDelegate(em)
// stream handler
Expand Down
13 changes: 7 additions & 6 deletions pkg/frontier/edgebound/edge_onoff.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (

"github.com/jumboframes/armorigo/synchub"
"github.com/singchia/frontier/pkg/frontier/apis"
"github.com/singchia/frontier/pkg/frontier/misc"
"github.com/singchia/frontier/pkg/frontier/repo/model"
"github.com/singchia/frontier/pkg/frontier/repo/query"
"github.com/singchia/geminio"
Expand Down Expand Up @@ -145,7 +146,7 @@ func (em *edgeManager) ConnOnline(d delegate.ConnDescriber) error {
return err
}
}
klog.V(2).Infof("edge online, edgeID: %d, meta: %s, addr: %s", edgeID, string(meta), addr)
klog.V(2).Infof("edge online, edgeID: %d, meta: %s, addr: %s", edgeID, misc.Redact(string(meta)), addr)
return nil
}

Expand All @@ -154,12 +155,12 @@ func (em *edgeManager) ConnOffline(d delegate.ConnDescriber) error {
meta := d.Meta()
addr := d.RemoteAddr()

klog.V(2).Infof("edge offline, edgeID: %d, meta: %s, addr: %s", edgeID, string(meta), addr)
klog.V(2).Infof("edge offline, edgeID: %d, meta: %s, addr: %s", edgeID, misc.Redact(string(meta)), addr)
// offline the cache
err := em.offline(edgeID, meta, addr)
if err != nil {
klog.Errorf("edge offline, cache or db offline err: %s, edgeID: %d, meta: %s, addr: %s",
err, edgeID, string(meta), addr)
err, edgeID, misc.Redact(string(meta)), addr)
return err
}
return nil
Expand All @@ -169,7 +170,7 @@ func (em *edgeManager) Heartbeat(d delegate.ConnDescriber) error {
edgeID := d.ClientID()
meta := string(d.Meta())
addr := d.RemoteAddr()
klog.V(3).Infof("edge heartbeat, edgeID: %d, meta: %s, addr: %s", edgeID, string(meta), addr)
klog.V(3).Infof("edge heartbeat, edgeID: %d, meta: %s, addr: %s", edgeID, misc.Redact(string(meta)), addr)
if em.informer != nil {
em.informer.EdgeHeartbeat(edgeID, d.Meta(), addr)
}
Expand Down Expand Up @@ -199,14 +200,14 @@ func (em *edgeManager) GetClientID(_ uint64, meta []byte) (uint64, error) {
if em.exchange != nil {
edgeID, err = em.exchange.GetEdgeID(meta)
if err == nil {
klog.V(2).Infof("edge get edgeID: %d from exchange, meta: %s", edgeID, string(meta))
klog.V(2).Infof("edge get edgeID: %d from exchange, meta: %s", edgeID, misc.Redact(string(meta)))
return edgeID, nil
}
}

if (err == apis.ErrServiceNotOnline || err == apis.ErrRPCNotOnline) && em.conf.Edgebound.EdgeIDAllocWhenNoIDServiceOn {
edgeID = em.idFactory.GetID()
klog.V(2).Infof("edge get edgeID: %d, meta: %s, after no ID acquired from exchange", edgeID, string(meta))
klog.V(2).Infof("edge get edgeID: %d, meta: %s, after no ID acquired from exchange", edgeID, misc.Redact(string(meta)))
return em.idFactory.GetID(), nil
}
return 0, err
Expand Down
16 changes: 8 additions & 8 deletions pkg/frontier/exchange/oob.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ import (
func (ex *exchange) GetEdgeID(meta []byte) (uint64, error) {
svc, err := ex.Servicebound.GetServiceByRPC(apis.RPCGetEdgeID)
if err != nil {
klog.V(2).Infof("exchange get edgeID, get service err: %s, meta: %s", err, string(meta))
klog.V(2).Infof("exchange get edgeID, get service err: %s, meta: %s", err, misc.Redact(string(meta)))
if err == apis.ErrRecordNotFound {
return 0, apis.ErrServiceNotOnline
}
Expand All @@ -32,7 +32,7 @@ func (ex *exchange) GetEdgeID(meta []byte) (uint64, error) {
opt.SetTimeout(30 * time.Second)
rsp, err := svc.Call(context.TODO(), apis.RPCGetEdgeID, req, opt)
if err != nil {
klog.V(2).Infof("exchange call service: %d, get edgeID err: %s, meta: %s", svc.ClientID(), err, meta)
klog.V(2).Infof("exchange call service: %d, get edgeID err: %s, meta: %s", svc.ClientID(), err, misc.Redact(string(meta)))
return 0, err
}
data := rsp.Data()
Expand All @@ -45,7 +45,7 @@ func (ex *exchange) GetEdgeID(meta []byte) (uint64, error) {
func (ex *exchange) EdgeOnline(edgeID uint64, meta []byte, addr net.Addr) error {
svcs, err := ex.Servicebound.GetServicesByRPC(apis.RPCEdgeOnline)
if err != nil {
klog.V(2).Infof("exchange edge online, get service err: %s, edgeID: %d, meta: %s, addr: %s", err, edgeID, string(meta), addr)
klog.V(2).Infof("exchange edge online, get service err: %s, edgeID: %d, meta: %s, addr: %s", err, edgeID, misc.Redact(string(meta)), addr)
if err == apis.ErrRecordNotFound {
return apis.ErrServiceNotOnline
}
Expand All @@ -60,7 +60,7 @@ func (ex *exchange) EdgeOnline(edgeID uint64, meta []byte, addr net.Addr) error
}
data, err := json.Marshal(event)
if err != nil {
klog.Errorf("exchange edge online, json marshal err: %s, edgeID: %d, meta: %s, addr: %s", err, edgeID, string(meta), addr)
klog.Errorf("exchange edge online, json marshal err: %s, edgeID: %d, meta: %s, addr: %s", err, edgeID, misc.Redact(string(meta)), addr)
return err
}

Expand All @@ -72,7 +72,7 @@ func (ex *exchange) EdgeOnline(edgeID uint64, meta []byte, addr net.Addr) error
opt.SetTimeout(30 * time.Second)
_, err = svc.Call(context.TODO(), apis.RPCEdgeOnline, req, opt)
if err != nil {
klog.V(2).Infof("exchange call service: %d, edge online err: %s, meta: %s, addr: %s", svc.ClientID(), err, meta, addr)
klog.V(2).Infof("exchange call service: %d, edge online err: %s, meta: %s, addr: %s", svc.ClientID(), err, misc.Redact(string(meta)), addr)
return err
}
return nil
Expand All @@ -81,7 +81,7 @@ func (ex *exchange) EdgeOnline(edgeID uint64, meta []byte, addr net.Addr) error
func (ex *exchange) EdgeOffline(edgeID uint64, meta []byte, addr net.Addr) error {
svcs, err := ex.Servicebound.GetServicesByRPC(apis.RPCEdgeOffline)
if err != nil {
klog.V(2).Infof("exchange edge offline, get service err: %s, edgeID: %d, meta: %s, addr: %s", err, edgeID, string(meta), addr)
klog.V(2).Infof("exchange edge offline, get service err: %s, edgeID: %d, meta: %s, addr: %s", err, edgeID, misc.Redact(string(meta)), addr)
if err == apis.ErrRecordNotFound {
return apis.ErrServiceNotOnline
}
Expand All @@ -98,7 +98,7 @@ func (ex *exchange) EdgeOffline(edgeID uint64, meta []byte, addr net.Addr) error
}
data, err := json.Marshal(event)
if err != nil {
klog.Errorf("exchange edge offline, json marshal err: %s, edgeID: %d, meta: %s, addr: %s", err, edgeID, string(meta), addr)
klog.Errorf("exchange edge offline, json marshal err: %s, edgeID: %d, meta: %s, addr: %s", err, edgeID, misc.Redact(string(meta)), addr)
return err
}
// call service
Expand All @@ -107,7 +107,7 @@ func (ex *exchange) EdgeOffline(edgeID uint64, meta []byte, addr net.Addr) error
opt.SetTimeout(30 * time.Second)
_, err = svc.Call(context.TODO(), apis.RPCEdgeOffline, req, opt)
if err != nil {
klog.V(2).Infof("exchange call service: %d, edge offline err: %s, meta: %s, addr: %s", svc.ClientID(), err, meta, addr)
klog.V(2).Infof("exchange call service: %d, edge offline err: %s, meta: %s, addr: %s", svc.ClientID(), err, misc.Redact(string(meta)), addr)
return err
}
return nil
Expand Down
124 changes: 124 additions & 0 deletions pkg/frontier/misc/redact.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,124 @@
package misc

import (
"bytes"
"encoding/json"
"strings"
)

// redactedValue is the uniform replacement for sensitive field values.
const redactedValue = "***"

// unparsableMetaPlaceholder is the fail-closed output for non-JSON input:
// the raw value is never echoed back, so malformed payloads cannot smuggle
// credentials into logs.
const unparsableMetaPlaceholder = "<redacted: unparsable meta>"

// nonObjectMetaPlaceholder is the fail-closed output for JSON input whose
// top-level value is not an object or an array (e.g. a bare string). A
// top-level string is a common client mistake that carries an
// already-serialized object, so it is masked as a whole instead of being
// passed through.
const nonObjectMetaPlaceholder = "<redacted: non-object meta>"

// oversizedMetaPlaceholder is the fail-closed output for oversized meta.
// Meta is client-controlled and should not flood logs anyway; returning a
// placeholder upfront also avoids parsing hostile large JSON payloads.
const oversizedMetaPlaceholder = "<redacted: meta too long>"

// maxMetaLen is the length limit for meta accepted by the log redaction path.
const maxMetaLen = 8 * 1024

// sensitiveKeyFragments holds lowercase substrings used to detect
// credential-like keys. Clients commonly carry access_key/secret_key/token
// and friends in connection meta, and those must be masked before logging.
// Substring matching covers naming variants (accessKey, secret_key_2, ...).
// It intentionally errs on the side of over-masking: innocuous keys such as
// "keyword" or "author" are masked too, which is preferred over leaking.
var sensitiveKeyFragments = []string{
"key",
"secret",
"token",
"password",
"passwd",
"pass",
"pwd",
"credential",
"auth",
"signature",
"bearer",
"session",
"cookie",
"cert",
"private",
}

// Redact masks credential-like fields in a JSON meta string and returns the
// result, for use at every call site that logs meta. Contract:
// - empty/whitespace-only input is returned as is;
// - input longer than maxMetaLen fails closed with a placeholder;
// - values of sensitive keys (case-insensitive substring match) are
// replaced with "***" at any level (objects, nested objects, objects
// inside arrays), regardless of whether the value is a string or a
// composite;
// - a top-level value that is not an object or an array fails closed with
// a placeholder;
// - non-JSON input fails closed with a placeholder, never echoing the raw
// value. If valid JSON is followed by trailing garbage, only the first
// JSON value is parsed and the trailing content is dropped.
func Redact(meta string) string {
if strings.TrimSpace(meta) == "" {
return meta
}
if len(meta) > maxMetaLen {
return oversizedMetaPlaceholder
}
var v any
dec := json.NewDecoder(bytes.NewReader([]byte(meta)))
dec.UseNumber()
if err := dec.Decode(&v); err != nil {
return unparsableMetaPlaceholder
}
switch v.(type) {
case map[string]any, []any:
default:
return nonObjectMetaPlaceholder
}
redactValue(v)
out, err := json.Marshal(v)
if err != nil {
return unparsableMetaPlaceholder
}
return string(out)
}

// redactValue recursively masks in place: map keys are matched one by one
// (a hit replaces the whole value), slices are descended element-wise, and
// other scalars are left untouched.
func redactValue(v any) {
switch val := v.(type) {
case map[string]any:
for k, item := range val {
if isSensitiveKey(k) {
val[k] = redactedValue
continue
}
redactValue(item)
}
case []any:
for _, item := range val {
redactValue(item)
}
}
}

// isSensitiveKey matches the lowercased key against the fragment table.
func isSensitiveKey(key string) bool {
lower := strings.ToLower(key)
for _, frag := range sensitiveKeyFragments {
if strings.Contains(lower, frag) {
return true
}
}
return false
}
Loading
Loading