fix(cron): route jobs through scheduler for concurrency control and parallel execution

- Simplify cron session key to `agent:{agentId}:cron:{jobID}` (remove redundant `:run:{runID}`)
- Route cron jobs through scheduler's cron lane instead of calling loop.Run() directly
- Scheduler enforces per-session maxConcurrent=1, preventing same job from running concurrently
- Parallelize due job execution with goroutines + WaitGroup (PG and file store)
- Move scheduler creation before cron setup in gateway.go initialization order

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
viettranxandClaude Opus 4.6 committed 2026-02-28 16:53:15 +07:00
1 parent 73b46c3634
commit 0655849d3d
5 files changed
+135 -118

No files matched your search

+12 -11
View File
@@ -799,8 +799,18 @@ func runGateway() {
slog.Error("failed to start channels", "error", err)
}
// Start cron service with job handler
cronStore.SetOnJob(makeCronJobHandler(agentRouter, msgBus, cfg))
// Create lane-based scheduler (matching TS CommandLane pattern).
// The RunFunc resolves the agent from the RunRequest metadata.
// Must be created before cron setup so cron jobs route through the scheduler.
sched := scheduler.NewScheduler(
scheduler.DefaultLanes(),
scheduler.DefaultQueueConfig(),
makeSchedulerRunFunc(agentRouter, cfg),
)
defer sched.Stop()
// Start cron service with job handler (routes through scheduler's cron lane)
cronStore.SetOnJob(makeCronJobHandler(sched, msgBus, cfg))
if err := cronStore.Start(); err != nil {
slog.Warn("cron service failed to start", "error", err)
}
@@ -811,15 +821,6 @@ func runGateway() {
heartbeatSvc.Start()
}
// Create lane-based scheduler (matching TS CommandLane pattern).
// The RunFunc resolves the agent from the RunRequest metadata.
sched := scheduler.NewScheduler(
scheduler.DefaultLanes(),
scheduler.DefaultQueueConfig(),
makeSchedulerRunFunc(agentRouter, cfg),
)
defer sched.Stop()
// Adaptive throttle: reduce per-session concurrency when nearing the summary threshold.
// This prevents concurrent runs from racing with summarization.
// Uses calibrated token estimation (actual prompt tokens from last LLM call)
+19 -24
View File
@@ -632,36 +632,26 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents
}
}
// resolveCronAgent resolves the agent ID for a cron job, falling back to the
// config default if the requested agent doesn't exist.
func resolveCronAgent(agentID string, agents *agent.Router, cfg *config.Config) string {
if agentID == "" {
return cfg.ResolveDefaultAgentID()
}
normalized := config.NormalizeAgentID(agentID)
if _, err := agents.Get(normalized); err != nil {
slog.Warn("cron agent not found, falling back to default", "requested", agentID)
return cfg.ResolveDefaultAgentID()
}
return normalized
}
// makeCronJobHandler creates a cron job handler that sends job messages through the agent.
func makeCronJobHandler(agents *agent.Router, msgBus *bus.MessageBus, cfg *config.Config) func(job *store.CronJob) (*store.CronJobResult, error) {
// makeCronJobHandler creates a cron job handler that routes through the scheduler's cron lane.
// This ensures per-session concurrency control (same job can't run concurrently)
// and integration with /stop, /stopall commands.
func makeCronJobHandler(sched *scheduler.Scheduler, msgBus *bus.MessageBus, cfg *config.Config) func(job *store.CronJob) (*store.CronJobResult, error) {
return func(job *store.CronJob) (*store.CronJobResult, error) {
agentID := resolveCronAgent(job.AgentID, agents, cfg)
loop, err := agents.Get(agentID)
if err != nil {
return nil, fmt.Errorf("agent %s not found: %w", agentID, err)
agentID := job.AgentID
if agentID == "" {
agentID = cfg.ResolveDefaultAgentID()
} else {
agentID = config.NormalizeAgentID(agentID)
}
sessionKey := sessions.BuildCronSessionKey(agentID, job.ID, job.ID)
sessionKey := sessions.BuildCronSessionKey(agentID, job.ID)
channel := job.Payload.Channel
if channel == "" {
channel = "cron"
}
result, err := loop.Run(context.Background(), agent.RunRequest{
// Schedule through cron lane — scheduler handles agent resolution and concurrency
outCh := sched.Schedule(context.Background(), scheduler.LaneCron, agent.RunRequest{
SessionKey: sessionKey,
Message: job.Payload.Message,
Channel: channel,
@@ -672,10 +662,15 @@ func makeCronJobHandler(agents *agent.Router, msgBus *bus.MessageBus, cfg *confi
TraceName: fmt.Sprintf("Cron [%s] - %s", job.Name, agentID),
TraceTags: []string{"cron"},
})
if err != nil {
return nil, err
// Block until the scheduled run completes
outcome := <-outCh
if outcome.Err != nil {
return nil, outcome.Err
}
result := outcome.Result
// If job wants delivery to a channel, publish outbound
if job.Payload.Deliver && job.Payload.Channel != "" && job.Payload.To != "" {
msgBus.PublishOutbound(bus.OutboundMessage{
+9 -2
View File
@@ -6,6 +6,7 @@ import (
"log/slog"
"os"
"path/filepath"
"sync"
"time"
"github.com/adhocore/gronx"
@@ -175,10 +176,16 @@ func (cs *Service) checkJobs() {
cs.saveUnsafe()
cs.mu.Unlock()
// Execute jobs outside lock
// Execute jobs in parallel — scheduler enforces per-session serialization
var wg sync.WaitGroup
for _, jobID := range dueJobIDs {
cs.executeJobByID(jobID)
wg.Add(1)
go func(id string) {
defer wg.Done()
cs.executeJobByID(id)
}(jobID)
}
wg.Wait()
}
func (cs *Service) executeJobByID(jobID string) {
+8 -8
View File
@@ -10,7 +10,7 @@
// Group: {channel}:group:{groupId}
// Forum topic: {channel}:group:{groupId}:topic:{topicId}
// Subagent: subagent:{label}
// Cron: cron:{jobId}:run:{runId}
// Cron: cron:{jobId}
//
// Examples:
//
@@ -18,7 +18,7 @@
// agent:default:telegram:group:-100123456
// agent:default:telegram:group:-100123456:topic:99
// agent:default:subagent:my-task
// agent:default:cron:reminder:run:abc123
// agent:default:cron:reminder-job-id
package sessions
import (
@@ -57,18 +57,18 @@ func BuildSubagentSessionKey(agentID, label string) string {
return fmt.Sprintf("agent:%s:subagent:%s", agentID, label)
}
// BuildCronSessionKey builds the session key for a cron job run.
// BuildCronSessionKey builds the session key for a cron job.
// Each cron job gets one persistent session (all runs share the same history).
//
// agent:{agentId}:cron:{jobID}:run:{runID}
// agent:{agentId}:cron:{jobID}
//
// Guards against double-prefixing: if jobID is already a canonical session key
// (e.g. "agent:X:..."), only the rest part is used to prevent
// "agent:X:cron:agent:X:cron:..." duplication.
func BuildCronSessionKey(agentID, jobID, runID string) string {
// (e.g. "agent:X:..."), only the rest part is used.
func BuildCronSessionKey(agentID, jobID string) string {
if _, rest := ParseSessionKey(jobID); rest != "" {
jobID = rest
}
return fmt.Sprintf("agent:%s:cron:%s:run:%s", agentID, jobID, runID)
return fmt.Sprintf("agent:%s:cron:%s", agentID, jobID)
}
// BuildAgentMainSessionKey builds the shared "main" session key for an agent.
+87 -73
View File
@@ -2,6 +2,7 @@ package pg
import (
"log/slog"
"sync"
"time"
"github.com/google/uuid"
@@ -146,87 +147,100 @@ func (s *PGCronStore) checkAndRunDueJobs() {
return
}
// Clear next_run for all due jobs first to prevent duplicate fires
for _, job := range dueJobs {
// Clear next_run to prevent duplicate
if id, parseErr := uuid.Parse(job.ID); parseErr == nil {
s.db.Exec("UPDATE cron_jobs SET next_run_at = NULL WHERE id = $1", id)
}
jobCopy := job
startTime := time.Now()
// Wrap handler to fit ExecuteWithRetry's (string, error) signature
var lastResult *store.CronJobResult
resultStr, attempts, err := cron.ExecuteWithRetry(func() (string, error) {
r, e := handler(&jobCopy)
if e != nil {
return "", e
}
lastResult = r
if r != nil {
return r.Content, nil
}
return "", nil
}, s.retryCfg)
durationMS := time.Since(startTime).Milliseconds()
if attempts > 1 {
slog.Info("cron job retried", "id", job.ID, "attempts", attempts, "success", err == nil)
}
now := time.Now()
status := "ok"
var lastError *string
if err != nil {
status = "error"
errStr := err.Error()
lastError = &errStr
}
// Extract token usage from handler result
var inputTokens, outputTokens int
if lastResult != nil {
inputTokens = lastResult.InputTokens
outputTokens = lastResult.OutputTokens
}
// Log run
logID := uuid.Must(uuid.NewV7())
var summary *string
if err == nil {
s := cron.TruncateOutput(resultStr)
summary = &s
}
if id, parseErr := uuid.Parse(job.ID); parseErr == nil {
var agentUUID *uuid.UUID
if aid, aidErr := uuid.Parse(job.AgentID); aidErr == nil {
agentUUID = &aid
}
s.db.Exec(
`INSERT INTO cron_run_logs (id, job_id, agent_id, status, error, summary, duration_ms, input_tokens, output_tokens, ran_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)`,
logID, id, agentUUID, status, lastError, summary, durationMS, inputTokens, outputTokens, now,
)
}
// Recompute next run or delete
if job.DeleteAfterRun {
if id, parseErr := uuid.Parse(job.ID); parseErr == nil {
s.db.Exec("DELETE FROM cron_jobs WHERE id = $1", id)
}
} else if id, parseErr := uuid.Parse(job.ID); parseErr == nil {
schedule := job.Schedule
next := computeNextRun(&schedule, now)
s.db.Exec(
"UPDATE cron_jobs SET last_run_at = $1, last_status = $2, last_error = $3, next_run_at = $4, updated_at = $5 WHERE id = $6",
now, status, lastError, next, now, id,
)
}
}
// Execute jobs in parallel — scheduler enforces per-session serialization
var wg sync.WaitGroup
for _, job := range dueJobs {
wg.Add(1)
go func(job store.CronJob) {
defer wg.Done()
s.executeOneJob(job, handler)
}(job)
}
wg.Wait()
// Invalidate cache after job execution changed next_run_at values
s.mu.Lock()
s.cacheLoaded = false
s.mu.Unlock()
}
// executeOneJob runs a single cron job with retry, logs the result, and updates next_run_at.
func (s *PGCronStore) executeOneJob(job store.CronJob, handler func(job *store.CronJob) (*store.CronJobResult, error)) {
startTime := time.Now()
// Wrap handler to fit ExecuteWithRetry's (string, error) signature
var lastResult *store.CronJobResult
resultStr, attempts, err := cron.ExecuteWithRetry(func() (string, error) {
r, e := handler(&job)
if e != nil {
return "", e
}
lastResult = r
if r != nil {
return r.Content, nil
}
return "", nil
}, s.retryCfg)
durationMS := time.Since(startTime).Milliseconds()
if attempts > 1 {
slog.Info("cron job retried", "id", job.ID, "attempts", attempts, "success", err == nil)
}
now := time.Now()
status := "ok"
var lastError *string
if err != nil {
status = "error"
errStr := err.Error()
lastError = &errStr
}
// Extract token usage from handler result
var inputTokens, outputTokens int
if lastResult != nil {
inputTokens = lastResult.InputTokens
outputTokens = lastResult.OutputTokens
}
// Log run
logID := uuid.Must(uuid.NewV7())
var summary *string
if err == nil {
truncated := cron.TruncateOutput(resultStr)
summary = &truncated
}
if id, parseErr := uuid.Parse(job.ID); parseErr == nil {
var agentUUID *uuid.UUID
if aid, aidErr := uuid.Parse(job.AgentID); aidErr == nil {
agentUUID = &aid
}
s.db.Exec(
`INSERT INTO cron_run_logs (id, job_id, agent_id, status, error, summary, duration_ms, input_tokens, output_tokens, ran_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)`,
logID, id, agentUUID, status, lastError, summary, durationMS, inputTokens, outputTokens, now,
)
}
// Recompute next run or delete
if job.DeleteAfterRun {
if id, parseErr := uuid.Parse(job.ID); parseErr == nil {
s.db.Exec("DELETE FROM cron_jobs WHERE id = $1", id)
}
} else if id, parseErr := uuid.Parse(job.ID); parseErr == nil {
schedule := job.Schedule
next := computeNextRun(&schedule, now)
s.db.Exec(
"UPDATE cron_jobs SET last_run_at = $1, last_status = $2, last_error = $3, next_run_at = $4, updated_at = $5 WHERE id = $6",
now, status, lastError, next, now, id,
)
}
}