Merge branch 'main' of github.com:Infisical/infisical into sid/k8s-operator

This commit is contained in:
sidwebworks
2025-08-12 16:49:27 +05:30
271 changed files with 10389 additions and 6590 deletions
+116
View File
@@ -0,0 +1,116 @@
package sse
import (
"bufio"
"fmt"
"io"
"net/http"
"strings"
"sync"
"time"
)
type SSEEvent struct {
ID string
Event string
Data string
}
// SSEClient handles SSE connections
type SSEClient struct {
URL string
Client *http.Client
LastHealthCheck time.Time
mu *sync.Mutex // for safe concurrent access to LastHealthCheck
}
// NewClient creates a new SSE client
func NewClient() SSEClient {
return SSEClient{
mu: &sync.Mutex{},
Client: &http.Client{
Timeout: 0, // No timeout for streaming
},
}
}
// Connect establishes SSE connection and returns a channel of events
func (c *SSEClient) Connect(req *http.Request) (<-chan SSEEvent, <-chan error, error) {
// Set required headers for SSE
req.Header.Set("Cache-Control", "no-cache")
req.Header.Set("Connection", "keep-alive")
req.Header.Set("Content-Type", "application/json")
resp, err := c.Client.Do(req)
if err != nil {
return nil, nil, err
}
if resp.StatusCode != http.StatusOK {
resp.Body.Close()
return nil, nil, fmt.Errorf("unexpected status code: %d", resp.StatusCode)
}
eventChan := make(chan SSEEvent)
errorChan := make(chan error)
go c.stream(resp.Body, eventChan, errorChan)
return eventChan, errorChan, nil
}
func (c *SSEClient) stream(body io.ReadCloser, eventChan chan<- SSEEvent, errorChan chan<- error) {
defer body.Close()
defer close(eventChan)
defer close(errorChan)
scanner := bufio.NewScanner(body)
var event SSEEvent
for scanner.Scan() {
line := strings.TrimSpace(scanner.Text())
// End of event
if line == "" {
if event.Data != "" || event.Event != "" {
if strings.TrimSpace(event.Data) == "1" {
c.mu.Lock()
c.LastHealthCheck = time.Now()
c.mu.Unlock()
} else if event.Event != "ping" {
eventChan <- event
}
event = SSEEvent{} // Reset for next event
}
continue
}
switch {
case strings.HasPrefix(line, "data:"):
data := strings.TrimSpace(strings.TrimPrefix(line, "data:"))
if event.Data != "" {
event.Data += "\n"
}
event.Data += data
case strings.HasPrefix(line, "event:"):
event.Event = strings.TrimSpace(strings.TrimPrefix(line, "event:"))
case strings.HasPrefix(line, "id:"):
event.ID = strings.TrimSpace(strings.TrimPrefix(line, "id:"))
case strings.HasPrefix(line, "retry:"):
// Optional: parse and apply retry interval here
case strings.HasPrefix(line, ":"):
// Comment line — ignored
default:
// Unknown line format — can log/debug if needed
}
}
if err := scanner.Err(); err != nil {
errorChan <- err
}
}
+165
View File
@@ -0,0 +1,165 @@
package sse
import (
"context"
"fmt"
"net/http"
"sync"
"time"
)
type ConnectionMeta struct {
EventChan <-chan SSEEvent
ErrorChan <-chan error
LastPingAt time.Time
Cancel context.CancelFunc
}
type ConnectionRegistry struct {
Ctx context.Context
meta *ConnectionMeta
client SSEClient
mu sync.RWMutex
monitorCancel context.CancelFunc
monitorCtx context.Context
}
func NewConnectionRegistry(ctx context.Context) *ConnectionRegistry {
monitorCtx, monitorCancel := context.WithCancel(ctx)
return &ConnectionRegistry{
Ctx: ctx,
client: NewClient(),
monitorCtx: monitorCtx,
monitorCancel: monitorCancel,
}
}
// create creates a new connection
func (r *ConnectionRegistry) create(req *http.Request) (*ConnectionMeta, error) {
// Create new connection using provided request
eventChan, errorChan, err := r.client.Connect(req)
if err != nil {
return nil, fmt.Errorf("failed to connect: %w", err)
}
meta := &ConnectionMeta{
EventChan: eventChan,
ErrorChan: errorChan,
LastPingAt: time.Now(),
}
r.meta = meta
// Start cleanup monitor for this connection (NON-BLOCKING)
go r.monitor(meta)
println("Creating new connection\n")
return meta, nil
}
// Get retrieves the existing connection
func (r *ConnectionRegistry) Get() (*ConnectionMeta, bool) {
r.mu.RLock()
defer r.mu.RUnlock()
return r.meta, r.meta != nil
}
// Close closes the connection
func (r *ConnectionRegistry) Close() {
r.mu.Lock()
defer r.mu.Unlock()
if r.meta != nil {
if r.meta.Cancel != nil {
r.meta.Cancel()
}
r.meta = nil
}
// Cancel the monitor
if r.monitorCancel != nil {
r.monitorCancel()
}
}
// IsConnected returns whether there's an active connection
func (r *ConnectionRegistry) IsConnected() bool {
r.mu.RLock()
defer r.mu.RUnlock()
return r.meta != nil
}
// UpdateLastPing updates the last ping time
func (r *ConnectionRegistry) UpdateLastPing() {
r.mu.Lock()
defer r.mu.Unlock()
if r.meta != nil {
r.meta.LastPingAt = time.Now()
}
}
// monitor watches for connection closure and cleans up
func (r *ConnectionRegistry) monitor(meta *ConnectionMeta) {
ticker := time.NewTicker(30 * time.Second)
defer ticker.Stop()
for {
select {
case <-r.monitorCtx.Done():
// Context cancelled, exit monitor
return
case <-ticker.C:
r.mu.RLock()
currentMeta := r.meta
r.mu.RUnlock()
// Check if this monitor is still relevant
if currentMeta != meta {
// This connection has been replaced, exit monitor
return
}
if currentMeta != nil && time.Since(currentMeta.LastPingAt) > 2*time.Minute {
fmt.Println("Last ping was more than 2 minutes ago, closing connection")
r.mu.Lock()
if r.meta == meta { // Double-check under lock
if r.meta.Cancel != nil {
r.meta.Cancel()
}
r.meta = nil
}
r.mu.Unlock()
return // Exit monitor after cleanup
}
}
}
}
// Subscribe provides a convenient way to get events from the connection
func (r *ConnectionRegistry) Subscribe(build func() (*http.Request, error)) (<-chan SSEEvent, <-chan error, error) {
r.mu.Lock()
defer r.mu.Unlock()
// Get existing connection if available
if r.meta != nil {
return r.meta.EventChan, r.meta.ErrorChan, nil
}
req, err := build()
if err != nil {
return nil, nil, err
}
// Create new connection if none exists
meta, err := r.create(req)
if err != nil {
return nil, nil, err
}
return meta.EventChan, meta.ErrorChan, nil
}