feat(agent): removed 5s polling algorithm to a simple one

This commit is contained in:
Akhil Mohan
2024-03-23 21:10:42 +05:30
parent 88f7e4255e
commit 98ea2c1828
2 changed files with 83 additions and 78 deletions
+80 -77
View File
@@ -34,6 +34,9 @@ import (
const DEFAULT_INFISICAL_CLOUD_URL = "https://app.infisical.com" const DEFAULT_INFISICAL_CLOUD_URL = "https://app.infisical.com"
// duration to reduce from expiry of dynamic leases so that it gets triggered before expiry
const DYNAMIC_SECRET_PRUNE_EXPIRE_BUFFER = -15
type Config struct { type Config struct {
Infisical InfisicalConfig `yaml:"infisical"` Infisical InfisicalConfig `yaml:"infisical"`
Auth AuthConfig `yaml:"auth"` Auth AuthConfig `yaml:"auth"`
@@ -102,7 +105,7 @@ type DynamicSecretLease struct {
Slug string Slug string
ProjectSlug string ProjectSlug string
Data map[string]interface{} Data map[string]interface{}
Templates []string TemplateIDs []int
} }
type DynamicSecretLeaseManager struct { type DynamicSecretLeaseManager struct {
@@ -115,7 +118,7 @@ func (d *DynamicSecretLeaseManager) Prune() {
defer d.mutex.Unlock() defer d.mutex.Unlock()
d.leases = slices.DeleteFunc(d.leases, func(s DynamicSecretLease) bool { d.leases = slices.DeleteFunc(d.leases, func(s DynamicSecretLease) bool {
return time.Now().After(s.ExpireAt.Add(-15 * time.Second)) return time.Now().After(s.ExpireAt.Add(DYNAMIC_SECRET_PRUNE_EXPIRE_BUFFER * time.Second))
}) })
} }
@@ -131,13 +134,13 @@ func (d *DynamicSecretLeaseManager) Append(lease DynamicSecretLease) {
}) })
if index != -1 { if index != -1 {
d.leases[index].Templates = append(d.leases[index].Templates, lease.Templates...) d.leases[index].TemplateIDs = append(d.leases[index].TemplateIDs, lease.TemplateIDs...)
return return
} }
d.leases = append(d.leases, lease) d.leases = append(d.leases, lease)
} }
func (d *DynamicSecretLeaseManager) RegisterTemplate(projectSlug, environment, secretPath, slug, templateName string) { func (d *DynamicSecretLeaseManager) RegisterTemplate(projectSlug, environment, secretPath, slug string, templateId int) {
d.mutex.Lock() d.mutex.Lock()
defer d.mutex.Unlock() defer d.mutex.Unlock()
@@ -149,7 +152,7 @@ func (d *DynamicSecretLeaseManager) RegisterTemplate(projectSlug, environment, s
}) })
if index != -1 { if index != -1 {
d.leases[index].Templates = append(d.leases[index].Templates, templateName) d.leases[index].TemplateIDs = append(d.leases[index].TemplateIDs, templateId)
} }
} }
@@ -166,34 +169,31 @@ func (d *DynamicSecretLeaseManager) GetLease(projectSlug, environment, secretPat
return nil return nil
} }
// this is like etag for dynamic secret // for a given template find the first expiring lease
func (d *DynamicSecretLeaseManager) GetDynamicSecretTemplateTag(templateName string) string { // The bool indicates whether it contains valid expiry list
func (d *DynamicSecretLeaseManager) GetFirstExpiringLeaseTime(templateId int) (time.Time, bool) {
d.mutex.Lock() d.mutex.Lock()
defer d.mutex.Unlock() defer d.mutex.Unlock()
tag := "" if len(d.leases) == 0 {
for _, el := range d.leases { return time.Time{}, false
if slices.Contains(el.Templates, templateName) { }
tag += el.LeaseID
var firstExpiry time.Time
for i, el := range d.leases {
if i == 0 {
firstExpiry = el.ExpireAt
}
newLeaseTime := el.ExpireAt.Add(DYNAMIC_SECRET_PRUNE_EXPIRE_BUFFER * time.Second)
if newLeaseTime.Before(firstExpiry) {
firstExpiry = newLeaseTime
} }
} }
return tag return firstExpiry, true
} }
func NewDynamicSecretLeaseManager(sigChan chan os.Signal) *DynamicSecretLeaseManager { func NewDynamicSecretLeaseManager(sigChan chan os.Signal) *DynamicSecretLeaseManager {
manager := &DynamicSecretLeaseManager{} manager := &DynamicSecretLeaseManager{}
go func() {
for {
select {
case <-sigChan:
return
default:
time.Sleep(5 * time.Second)
manager.Prune()
}
}
}()
return manager return manager
} }
@@ -347,34 +347,43 @@ func secretTemplateFunction(accessToken string, existingEtag string, currentEtag
} }
} }
func dynamicSecretTemplateFunction(accessToken string, dynamicSecretManager *DynamicSecretLeaseManager, templateName string) func(string, string, string, string) (map[string]interface{}, error) { func dynamicSecretTemplateFunction(accessToken string, dynamicSecretManager *DynamicSecretLeaseManager, templateId int) func(...string) (map[string]interface{}, error) {
return func(projectSlug, envSlug, secretPath, slug string) (map[string]interface{}, error) { return func(args ...string) (map[string]interface{}, error) {
argLength := len(args)
if argLength != 4 && argLength != 5 {
return nil, fmt.Errorf("Invalid arguments found for dynamic-secret function. Check template %i", templateId)
}
projectSlug, envSlug, secretPath, slug, ttl := args[0], args[1], args[2], args[3], ""
if argLength == 5 {
ttl = args[4]
}
dynamicSecretData := dynamicSecretManager.GetLease(projectSlug, envSlug, secretPath, slug) dynamicSecretData := dynamicSecretManager.GetLease(projectSlug, envSlug, secretPath, slug)
if dynamicSecretData != nil { if dynamicSecretData != nil {
dynamicSecretManager.RegisterTemplate(projectSlug, envSlug, secretPath, slug, templateName) dynamicSecretManager.RegisterTemplate(projectSlug, envSlug, secretPath, slug, templateId)
return dynamicSecretData.Data, nil return dynamicSecretData.Data, nil
} }
res, err := util.CreateDynamicSecretLease(accessToken, projectSlug, envSlug, secretPath, slug) res, err := util.CreateDynamicSecretLease(accessToken, projectSlug, envSlug, secretPath, slug, ttl)
if err != nil { if err != nil {
return nil, err return nil, err
} }
dynamicSecretManager.Append(DynamicSecretLease{LeaseID: res.Lease.Id, ExpireAt: res.Lease.ExpireAt, Environment: envSlug, SecretPath: secretPath, Slug: slug, ProjectSlug: projectSlug, Data: res.Data, Templates: []string{templateName}}) dynamicSecretManager.Append(DynamicSecretLease{LeaseID: res.Lease.Id, ExpireAt: res.Lease.ExpireAt, Environment: envSlug, SecretPath: secretPath, Slug: slug, ProjectSlug: projectSlug, Data: res.Data, TemplateIDs: []int{templateId}})
return res.Data, nil return res.Data, nil
} }
} }
func ProcessTemplate(templatePath string, data interface{}, accessToken string, existingEtag string, currentEtag *string, dynamicSecretManager *DynamicSecretLeaseManager) (*bytes.Buffer, error) { func ProcessTemplate(templateId int, templatePath string, data interface{}, accessToken string, existingEtag string, currentEtag *string, dynamicSecretManager *DynamicSecretLeaseManager) (*bytes.Buffer, error) {
templateName := path.Base(templatePath)
// custom template function to fetch secrets from Infisical // custom template function to fetch secrets from Infisical
secretFunction := secretTemplateFunction(accessToken, existingEtag, currentEtag) secretFunction := secretTemplateFunction(accessToken, existingEtag, currentEtag)
dynamicSecretFunction := dynamicSecretTemplateFunction(accessToken, dynamicSecretManager, templateName) dynamicSecretFunction := dynamicSecretTemplateFunction(accessToken, dynamicSecretManager, templateId)
funcs := template.FuncMap{ funcs := template.FuncMap{
"secret": secretFunction, "secret": secretFunction,
"dynamic_secret": dynamicSecretFunction, "dynamic_secret": dynamicSecretFunction,
} }
templateName := path.Base(templatePath)
tmpl, err := template.New(templateName).Funcs(funcs).ParseFiles(templatePath) tmpl, err := template.New(templateName).Funcs(funcs).ParseFiles(templatePath)
if err != nil { if err != nil {
return nil, err return nil, err
@@ -388,7 +397,7 @@ func ProcessTemplate(templatePath string, data interface{}, accessToken string,
return &buf, nil return &buf, nil
} }
func ProcessBase64Template(encodedTemplate string, data interface{}, accessToken string, existingEtag string, currentEtag *string, dynamicSecretLeaser *DynamicSecretLeaseManager) (*bytes.Buffer, error) { func ProcessBase64Template(templateId int, encodedTemplate string, data interface{}, accessToken string, existingEtag string, currentEtag *string, dynamicSecretLeaser *DynamicSecretLeaseManager) (*bytes.Buffer, error) {
// custom template function to fetch secrets from Infisical // custom template function to fetch secrets from Infisical
decoded, err := base64.StdEncoding.DecodeString(encodedTemplate) decoded, err := base64.StdEncoding.DecodeString(encodedTemplate)
if err != nil { if err != nil {
@@ -398,7 +407,7 @@ func ProcessBase64Template(encodedTemplate string, data interface{}, accessToken
templateString := string(decoded) templateString := string(decoded)
secretFunction := secretTemplateFunction(accessToken, existingEtag, currentEtag) // TODO: Fix this secretFunction := secretTemplateFunction(accessToken, existingEtag, currentEtag) // TODO: Fix this
dynamicSecretFunction := dynamicSecretTemplateFunction(accessToken, dynamicSecretLeaser, encodedTemplate) dynamicSecretFunction := dynamicSecretTemplateFunction(accessToken, dynamicSecretLeaser, templateId)
funcs := template.FuncMap{ funcs := template.FuncMap{
"secret": secretFunction, "secret": secretFunction,
"dynamic_secret": dynamicSecretFunction, "dynamic_secret": dynamicSecretFunction,
@@ -633,7 +642,7 @@ func (tm *AgentManager) WriteTemplateToFile(bytes *bytes.Buffer, template *Templ
log.Info().Msgf("template engine: secret template at path %s has been rendered and saved to path %s", template.SourcePath, template.DestinationPath) log.Info().Msgf("template engine: secret template at path %s has been rendered and saved to path %s", template.SourcePath, template.DestinationPath)
} }
func (tm *AgentManager) MonitorSecretChanges(secretTemplate Template, sigChan chan os.Signal) { func (tm *AgentManager) MonitorSecretChanges(secretTemplate Template, templateId int, sigChan chan os.Signal) {
pollingInterval := time.Duration(5 * time.Minute) pollingInterval := time.Duration(5 * time.Minute)
@@ -652,11 +661,7 @@ func (tm *AgentManager) MonitorSecretChanges(secretTemplate Template, sigChan ch
var existingEtag string var existingEtag string
var currentEtag string var currentEtag string
var dynamicSecretTag string
var currentDynamicSecretTag string
var templateName string
var firstRun = true var firstRun = true
var pollingLapsed = pollingInterval
execTimeout := secretTemplate.Config.Execute.Timeout execTimeout := secretTemplate.Config.Execute.Timeout
execCommand := secretTemplate.Config.Execute.Command execCommand := secretTemplate.Config.Execute.Command
@@ -667,52 +672,50 @@ func (tm *AgentManager) MonitorSecretChanges(secretTemplate Template, sigChan ch
return return
default: default:
{ {
pollingLapsed = pollingLapsed - 5*time.Second tm.dynamicSecretLeases.Prune()
if !firstRun { token := tm.GetToken()
currentDynamicSecretTag = tm.dynamicSecretLeases.GetDynamicSecretTemplateTag(templateName) if token != "" {
} var processedTemplate *bytes.Buffer
shouldTrigger := currentDynamicSecretTag != dynamicSecretTag || pollingLapsed < 0 var err error
if shouldTrigger {
token := tm.GetToken()
if token != "" {
var processedTemplate *bytes.Buffer
var err error
if secretTemplate.SourcePath != "" { if secretTemplate.SourcePath != "" {
processedTemplate, err = ProcessTemplate(secretTemplate.SourcePath, nil, token, existingEtag, &currentEtag, tm.dynamicSecretLeases) processedTemplate, err = ProcessTemplate(templateId, secretTemplate.SourcePath, nil, token, existingEtag, &currentEtag, tm.dynamicSecretLeases)
templateName = path.Base(secretTemplate.SourcePath) } else {
} else { processedTemplate, err = ProcessBase64Template(templateId, secretTemplate.Base64TemplateContent, nil, token, existingEtag, &currentEtag, tm.dynamicSecretLeases)
processedTemplate, err = ProcessBase64Template(secretTemplate.Base64TemplateContent, nil, token, existingEtag, &currentEtag, tm.dynamicSecretLeases) }
templateName = path.Base(secretTemplate.Base64TemplateContent)
}
if err != nil { if err != nil {
log.Error().Msgf("unable to process template because %v", err) log.Error().Msgf("unable to process template because %v", err)
} else { } else {
if (existingEtag != currentEtag) || firstRun { if (existingEtag != currentEtag) || firstRun {
tm.WriteTemplateToFile(processedTemplate, &secretTemplate) tm.WriteTemplateToFile(processedTemplate, &secretTemplate)
existingEtag = currentEtag existingEtag = currentEtag
if !firstRun && execCommand != "" { if !firstRun && execCommand != "" {
log.Info().Msgf("executing command: %s", execCommand) log.Info().Msgf("executing command: %s", execCommand)
err := ExecuteCommandWithTimeout(execCommand, execTimeout) err := ExecuteCommandWithTimeout(execCommand, execTimeout)
if err != nil {
log.Error().Msgf("unable to execute command because %v", err)
}
if err != nil {
log.Error().Msgf("unable to execute command because %v", err)
} }
if firstRun {
firstRun = false }
} if firstRun {
firstRun = false
} }
} }
dynamicSecretTag = tm.dynamicSecretLeases.GetDynamicSecretTemplateTag(templateName)
// restart the polling because latest changes have been picked
pollingLapsed = pollingInterval
} }
time.Sleep(5 * time.Second)
// now the idea is we pick the next sleep time in which the one shorter out of
// - polling time
// - first lease that's gonna get expired in the template
firstLeaseExpiry, isValid := tm.dynamicSecretLeases.GetFirstExpiringLeaseTime(templateId)
var waitTime = pollingInterval
if isValid && firstLeaseExpiry.Sub(time.Now()) < pollingInterval {
waitTime = firstLeaseExpiry.Sub(time.Now())
}
time.Sleep(waitTime)
} else { } else {
// It fails to get the access token. So we will re-try in 3 seconds. We do this because if we don't, the user will have to wait for the next polling interval to get the first secret render. // It fails to get the access token. So we will re-try in 3 seconds. We do this because if we don't, the user will have to wait for the next polling interval to get the first secret render.
time.Sleep(3 * time.Second) time.Sleep(3 * time.Second)
@@ -807,7 +810,7 @@ var agentCmd = &cobra.Command{
for i, template := range agentConfig.Templates { for i, template := range agentConfig.Templates {
log.Info().Msgf("template engine started for template %v...", i+1) log.Info().Msgf("template engine started for template %v...", i+1)
go tm.MonitorSecretChanges(template, sigChan) go tm.MonitorSecretChanges(template, i, sigChan)
} }
for { for {
+3 -1
View File
@@ -195,7 +195,7 @@ func GetPlainTextSecretsViaMachineIdentity(accessToken string, workspaceId strin
}, nil }, nil
} }
func CreateDynamicSecretLease(accessToken string, projectSlug string, environmentName string, secretsPath string, slug string) (models.DynamicSecretLease, error) { func CreateDynamicSecretLease(accessToken string, projectSlug string, environmentName string, secretsPath string, slug string, ttl string) (models.DynamicSecretLease, error) {
httpClient := resty.New() httpClient := resty.New()
httpClient.SetAuthToken(accessToken). httpClient.SetAuthToken(accessToken).
SetHeader("Accept", "application/json") SetHeader("Accept", "application/json")
@@ -203,7 +203,9 @@ func CreateDynamicSecretLease(accessToken string, projectSlug string, environmen
dynamicSecretRequest := api.CreateDynamicSecretLeaseV1Request{ dynamicSecretRequest := api.CreateDynamicSecretLeaseV1Request{
ProjectSlug: projectSlug, ProjectSlug: projectSlug,
Environment: environmentName, Environment: environmentName,
SecretPath: secretsPath,
Slug: slug, Slug: slug,
TTL: ttl,
} }
dynamicSecret, err := api.CallCreateDynamicSecretLeaseV1(httpClient, dynamicSecretRequest) dynamicSecret, err := api.CallCreateDynamicSecretLeaseV1(httpClient, dynamicSecretRequest)