package agent import ( "context" "errors" "log/slog" "os" "runtime" "time" "github.com/syncova/syncova/packages/platform/logging" ) // SystemInformation beschreibt das System, auf dem der Agent läuft. type SystemInformation struct { // Hostname ist der Rechnername. Hostname string // Platform ist das Betriebssystem. Platform string // Architecture ist die Rechnerarchitektur. Architecture string } // CollectSystemInformation ermittelt die Angaben des laufenden Systems. func CollectSystemInformation() SystemInformation { hostname, hostnameError := os.Hostname() if hostnameError != nil { // Ein fehlender Rechnername verhindert den Betrieb nicht; er wird als // unbekannt gemeldet statt den Start abzubrechen. hostname = "unbekannt" } return SystemInformation{ Hostname: hostname, Platform: runtime.GOOS, Architecture: runtime.GOARCH, } } // Zeitparameter der Lebendmeldung. const ( // DefaultHeartbeatInterval ist der Abstand zweier Lebendmeldungen. // // Er liegt deutlich unter der Frist, ab der ein Agent als offline gilt: // eine einzelne ausgefallene Meldung soll noch keinen Alarm auslösen. DefaultHeartbeatInterval = 60 * time.Second // initialRetryDelay ist die Wartezeit vor dem ersten erneuten Versuch. initialRetryDelay = 5 * time.Second // maximumRetryDelay begrenzt die Wartezeit zwischen zwei Versuchen. // // Ohne Obergrenze wüchse der Abstand nach einer längeren Störung so weit, // dass der Agent nach der Behebung minutenlang stumm bliebe. maximumRetryDelay = 5 * time.Minute ) // Runner hält die Verbindung des Agents zum Control Server aufrecht. type Runner struct { // client spricht mit dem Control Server. client *Client // heartbeatInterval ist der Abstand zweier Lebendmeldungen. heartbeatInterval time.Duration // logger protokolliert den Verlauf. logger *slog.Logger // systemInformation beschreibt das laufende System. systemInformation SystemInformation // taskExecutor führt Aufträge des Servers aus. // // Er darf nil sein; dann meldet sich der Agent nur an und gibt // Lebenszeichen, ohne Aufträge zu übernehmen. Das war der Zustand vor // Phase 5 und bleibt der Rückfallweg, wenn kein Schlüsselmaterial // eingerichtet ist. taskExecutor *TaskExecutor // taskPollInterval ist der Abstand zweier Auftragsabfragen. taskPollInterval time.Duration } // DefaultTaskPollInterval ist der Abstand zweier Auftragsabfragen. // // Zehn Sekunden: Der Abstand bestimmt, wie lange ein Auftrag im Mittel liegt, // bevor der Agent ihn bemerkt — und zugleich die Grundlast, die alle Agenten // zusammen auf dem Server erzeugen. Bei hundert Agenten sind zehn Sekunden zehn // Anfragen je Sekunde; eine Sekunde wären hundert. const DefaultTaskPollInterval = 10 * time.Second // EnableTaskExecution schaltet die Auftragsausführung ein. // // Getrennt vom Erzeuger, damit ein Agent ohne Auftragsausführung derselbe // Runner bleibt: Die Betriebsschleife unterscheidet sich nicht, nur ihr Inhalt. func (runner *Runner) EnableTaskExecution(taskExecutor *TaskExecutor, pollInterval time.Duration) { if pollInterval <= 0 { pollInterval = DefaultTaskPollInterval } runner.taskExecutor = taskExecutor runner.taskPollInterval = pollInterval } // NewRunner erzeugt die Betriebsschleife des Agents. func NewRunner(client *Client, heartbeatInterval time.Duration, baseLogger *slog.Logger) *Runner { if heartbeatInterval <= 0 { heartbeatInterval = DefaultHeartbeatInterval } return &Runner{ client: client, heartbeatInterval: heartbeatInterval, logger: logging.WithComponent(baseLogger, "agent"), systemInformation: CollectSystemInformation(), } } // Run hält die Verbindung, bis der Context abgebrochen wird. // // Eine Netzunterbrechung beendet den Agent nicht: er versucht es mit wachsendem // Abstand erneut. Ein abgelehntes Token dagegen behebt sich nicht von selbst // und beendet den Lauf mit einer klaren Meldung (PROMPT.md §51). func (runner *Runner) Run(runContext context.Context) error { runner.logger.Info("agent gestartet", slog.String("hostname", runner.systemInformation.Hostname), slog.String("platform", runner.systemInformation.Platform), slog.Duration("heartbeat_intervall", runner.heartbeatInterval)) // Die erste Meldung erfolgt sofort, damit der Server den Start mitbekommt. currentRetryDelay := initialRetryDelay consecutiveFailures := 0 for { heartbeatError := runner.client.SendHeartbeat(runContext) switch { case heartbeatError == nil: if consecutiveFailures > 0 { runner.logger.Info("verbindung zum control server wiederhergestellt", slog.Int("ausgefallene_meldungen", consecutiveFailures)) } consecutiveFailures = 0 currentRetryDelay = initialRetryDelay case errors.Is(heartbeatError, ErrTokenRejected): // Der Agent wurde gesperrt oder sein Token gewechselt. Weiterversuchen // wäre sinnlos und erzeugte nur Last auf dem Server. runner.logger.Error("der control server hat das agent-token abgelehnt; der agent beendet sich", slog.String("error", heartbeatError.Error())) return heartbeatError case errors.Is(heartbeatError, ErrServerUnreachable): consecutiveFailures++ // Die Meldung wird nur beim ersten Ausfall und danach selten // wiederholt: eine längere Störung soll das Log nicht fluten. if consecutiveFailures == 1 || consecutiveFailures%10 == 0 { runner.logger.Warn("der control server ist nicht erreichbar; der agent versucht es weiter", slog.Int("ausgefallene_meldungen", consecutiveFailures), slog.Duration("naechster_versuch_in", currentRetryDelay)) } default: consecutiveFailures++ runner.logger.Error("die lebendmeldung schlug fehl", slog.String("error", heartbeatError.Error())) } // Nach einem Fehler wird früher erneut versucht als im Regelabstand, // aber mit wachsendem Abstand. waitDuration := runner.heartbeatInterval if consecutiveFailures > 0 { waitDuration = currentRetryDelay currentRetryDelay = min(currentRetryDelay*2, maximumRetryDelay) } // Solange die Verbindung steht, wird bis zur nächsten Lebendmeldung // nach Aufträgen gefragt. // // Die Auftragsabholung läuft **innerhalb** der Wartezeit und nicht in // einem eigenen Ablauf: Ein Agent, der Aufträge holt, während seine // Lebendmeldung ausgefallen ist, arbeitete gegen einen Server, der ihn // für offline hält — und dessen Auftrag längst als verwaist gilt. if consecutiveFailures == 0 && runner.taskExecutor != nil { if runner.pollAndExecuteTasks(runContext, waitDuration) { continue } runner.logger.Info("agent beendet") return nil } select { case <-time.After(waitDuration): case <-runContext.Done(): runner.logger.Info("agent beendet") return nil } } } // RegisterIfNeeded nimmt den Agent auf, sofern noch kein Token vorliegt. // // Der Aufruf ist wiederholbar: ein bereits registrierter Agent tut nichts. func (runner *Runner) RegisterIfNeeded(registerContext context.Context, enrollmentToken string) (bool, error) { if runner.client.AgentToken() != "" { return false, nil } if enrollmentToken == "" { return false, ErrNotRegistered } registrationResponse, registerError := runner.client.Register(registerContext, enrollmentToken, runner.systemInformation) if registerError != nil { return false, registerError } runner.logger.Info("agent am control server aufgenommen", slog.String("agent_id", registrationResponse.Agent.ID), slog.String("name", registrationResponse.Agent.Name)) return true, nil } // pollAndExecuteTasks fragt bis zum Ablauf der Wartezeit nach Aufträgen. // // Der Rückgabewert meldet, ob die Schleife weiterlaufen soll; false bedeutet // einen Abbruch von außen. func (runner *Runner) pollAndExecuteTasks(runContext context.Context, availableTime time.Duration) bool { deadline := time.Now().Add(availableTime) for time.Now().Before(deadline) { claimedTask, claimError := runner.client.ClaimTask(runContext) switch { case claimError == nil && claimedTask != nil: runner.executeClaimedTask(runContext, claimedTask) // Nach einem Auftrag wird sofort erneut gefragt: Liegen mehrere an, // soll der Agent sie hintereinander abarbeiten, statt zwischen // jedem eine Wartezeit einzulegen. continue case errors.Is(claimError, ErrNoContent): // Der Normalfall: nichts zu tun. case errors.Is(claimError, ErrTokenRejected): runner.logger.Error("der control server hat das agent-token abgelehnt", slog.String("error", claimError.Error())) return false case claimError != nil: // Eine gescheiterte Abfrage ist kein Grund, die Lebendmeldung // aufzugeben. Sie wird vermerkt, und der Agent versucht es beim // nächsten Durchlauf erneut. runner.logger.Warn("die auftragsabfrage schlug fehl", slog.String("error", claimError.Error())) return true } remainingTime := time.Until(deadline) if remainingTime > runner.taskPollInterval { remainingTime = runner.taskPollInterval } if remainingTime <= 0 { break } select { case <-time.After(remainingTime): case <-runContext.Done(): return false } } return true } // executeClaimedTask führt einen Auftrag aus und meldet sein Ergebnis. // // Das Ergebnis wird **in jedem Fall** gemeldet, auch bei einem Fehler: Der // Server wartet darauf, und ein Agent, der schweigt, lässt den Lauf bis zum // Ablauf der Frist hängen. func (runner *Runner) executeClaimedTask(runContext context.Context, claimedTask *ClaimedTask) { var taskResult TaskResultPayload switch claimedTask.TaskType { case "restore": taskResult = runner.taskExecutor.ExecuteRestoreTask(runContext, claimedTask) case "backup", "": taskResult = runner.taskExecutor.ExecuteBackupTask(runContext, claimedTask) default: // Eine unbekannte Auftragsart wird gemeldet, nicht geraten. // // Ein neuerer Server koennte eine Art kennen, die dieser Agent noch // nicht hat. Sie als Sicherung zu behandeln waere die bequeme und // gefaehrliche Wahl: Der Server bekaeme ein Ergebnis fuer etwas // anderes, als er beauftragt hat. taskResult = TaskResultPayload{ Status: "failed", ErrorCode: "UNKNOWN_TASK_TYPE", ErrorMessage: "Dieser Agent kennt die Auftragsart " + claimedTask.TaskType + " nicht. Aktualisieren Sie den Agenten.", FailureClass: "configuration", } } // Die Rückmeldung braucht einen eigenen Kontext: Wurde der Agent während // der Sicherung beendet, ist der Lauf-Kontext bereits abgebrochen — und das // Ergebnis ginge verloren, obwohl es feststeht. reportContext, cancelReport := context.WithTimeout(context.WithoutCancel(runContext), resultReportTimeout) defer cancelReport() if reportError := runner.client.ReportResult(reportContext, claimedTask.TaskID, taskResult); reportError != nil { runner.logger.Error("das ergebnis liess sich nicht melden; der server gibt den auftrag "+ "nach ablauf der frist als verwaist frei", slog.String("auftrag", claimedTask.TaskID), slog.String("ergebnis", taskResult.Status), slog.String("error", reportError.Error())) } } // resultReportTimeout begrenzt die Rückmeldung eines Ergebnisses. const resultReportTimeout = 30 * time.Second