diff --git a/cmd/gateway.go b/cmd/gateway.go index 3076035a..22a62ecc 100644 --- a/cmd/gateway.go +++ b/cmd/gateway.go @@ -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) diff --git a/cmd/gateway_consumer.go b/cmd/gateway_consumer.go index c6d17590..719dfd84 100644 --- a/cmd/gateway_consumer.go +++ b/cmd/gateway_consumer.go @@ -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{ diff --git a/internal/cron/service_execution.go b/internal/cron/service_execution.go index 99ba4e2a..7f01754d 100644 --- a/internal/cron/service_execution.go +++ b/internal/cron/service_execution.go @@ -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) { diff --git a/internal/sessions/key.go b/internal/sessions/key.go index 049e5635..806e5eb3 100644 --- a/internal/sessions/key.go +++ b/internal/sessions/key.go @@ -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. diff --git a/internal/store/pg/cron_scheduler.go b/internal/store/pg/cron_scheduler.go index 881abe01..75905d02 100644 --- a/internal/store/pg/cron_scheduler.go +++ b/internal/store/pg/cron_scheduler.go @@ -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, + ) + } +}