跳到内容

Agent-lightning 鸟瞰图

大量信息即将呈现

本文深入探讨了 Agent-lightning 架构。它并非入门指南或使用教程。

本文总结了 Agent-lightning (截至 v0.2) 如何将 AlgorithmRunnerLightningStore 循环连接在一起,并展示了辅助组件(TracerAdapterLLM Proxy)如何插入循环。每个部分都提供了系统不同视角的图表。

算法 ↔ 运行器 ↔ 存储 数据流

在核心层面,Agent-lightning 构建于三个主要组件之上,它们在一个协调的循环中协同工作

  • Algorithm 系统的“大脑”。它决定运行哪些任务,从结果中学习,并更新资源(如 AI 模型或提示)。
  • Runner 系统的“工作者”。它执行算法分配的任务,运行代理,并记录结果。
  • LightningStore 中央“数据库”和消息队列。它充当单一事实来源,存储任务、结果和资源,并实现算法和运行器之间的通信。

典型的训练循环数据流如下:Algorithm 将任务(称为 Rollout)排队到 LightningStore。一个 Runner 然后从队列中取出一个任务,执行它,并将结果(称为 Span)回流到存储。一旦任务完成,算法就可以从存储中查询新数据以进行学习并更新其资源。

下图展示了这种基本交互在一个简单的非并行设置中。

sequenceDiagram
    autonumber
    participant Algo as Algorithm
    participant Store as LightningStore
    participant Runner
    participant Agent

    loop Over the dataset
        Algo-->>Store: add_resources + enqueue_rollout
        Store-->>Runner: dequeue_rollout → AttemptedRollout
        Store-->>Runner: get_latest_resources
        Runner-->>Store: update_attempt("running", worker_id)
        Runner->>Agent: rollout + resources
        Agent->>Runner: reward / spans
        Runner-->>Store: add_span or add_otel_span
        Runner-->>Store: update_attempt("finished", status)
        Store-->>Algo: query_rollouts + spans
        Algo-->>Algo: Update resources (optional)
    end

实线表示直接调用,而虚线表示异步或长时间运行的操作。

关键术语

我们定义以下术语,这些术语可能有助于理解上面的图表。

  • Resources 一组要调整或训练的资产。代理对资源执行 rollout 并收集 span 数据。算法使用这些数据来更新资源。在 RL 训练中,资源是可调整的模型。在提示调整中,资源是提示模板。
  • Rollout 代理对资源执行的工作单元。一个 rollout (名词) 可能是不完整的,在这种情况下它也被称为 tasksamplejob(这些术语可以互换使用)。代理执行它自己定义的流程来处理 rollout ——这个过程也被称为“to rollout”(动词)。执行后,rollout (名词) 被认为是 complete
  • Attempt rollout 的一次执行。如果发生故障或超时,一个 rollout 可以有多次尝试。
  • Span 在 rollout 期间,代理可以生成多个 span(也称为“traces”或“events”)。记录的 span 收集在存储中,这对于理解代理行为和优化代理至关重要。
  • Reward 一个特殊的 span,定义为一个数字,用于判断 rollout 在 rollout 某个时间段内的质量。
  • Dataset 一个包含不完整的 rollout(即任务)的集合,供代理处理。双数据集(训练集、验证集)作为算法排队第一批 rollout 的初始输入。

存储

如前所述,LightningStore 是 Agent-lightning 中所有数据的中央枢纽。存储公开了一组 API,供算法和运行器与数据交互;最重要的 API 是

from agentlightning.types import AttemptedRollout, ResourcesUpdate, Span, TaskInput

class LightningStore:

    async def enqueue_rollout(self, input: TaskInput, ...) -> Rollout: ...

    async def dequeue_rollout(self) -> AttemptedRollout | None: ...

    async def add_span(self, span: Span) -> Span: ...

    async def get_latest_resources(self) -> Optional[ResourcesUpdate]: ...

    async def wait_for_rollouts(self, rollout_ids: List[str], ...): ...

    async def query_spans(self, rollout_id: str, ...): ...

    async def update_attempt(self, rollout_id: str, attempt_id: str, status: str, ...): ...

    ...

这些接口操作于 AttemptedRolloutResourcesUpdateSpanTaskInput 实例来自 agentlightning.types

如 API 所示,存储本质上为 rollout 提供了一个队列,并为资源、span 和尝试提供存储。开发人员应仔细实现存储,以确保数据完整性和一致性,尤其是在多个运行器并行工作于多个尝试时。

存储设计为可扩展的。用户可以通过继承 LightningStore 并覆盖方法来实现自己的存储。Agent-lightning 提供了一些参考实现,例如 InMemoryLightningStore(默认)和 SqliteLightningStore(正在建设中)。当并行化时,存储可能需要特殊的包装器来确保线程/进程安全,或将计算委托给另一个进程或机器中的存储。

循环中的支持组件

虽然核心循环很简单,但 Agent-lightning 提供了几个组件来使开发更轻松、更强大。

Tracer

TracerRunner 中的一个组件,它记录代理执行期间详细的 span(事件),并将其发送到 LightningStore。代理不必手动记录每个 span,tracer 会自动检测关键方法(例如 LLM 调用)并捕获它们的输入、输出和元数据。这只需最少的精力即可提供代理行为的详细日志。

sequenceDiagram
    autonumber
    participant Store
    participant Runner
    participant Tracer
    participant Agent

    Note over Runner,Tracer: Runner manages tracer as member

    Tracer->>Agent: Apply instrumentation
    loop Until no more rollouts
        Store-->>Runner: dequeue_rollout → AttemptedRollout
        Store-->>Runner: get_latest_resources
        Runner->>Agent: training_rollout / validation_rollout
        loop For each finished span
            Agent-->>Tracer: openai.chat.completion invoked<br>agent.execute invoked<br>...
            Agent->>Tracer: emit intermediate reward
            Tracer-->>Store: add_otel_span(rollout_id, attempt_id, span)
        end
        Agent->>Runner: final reward + extra spans (if any)
        Runner-->>Store: add_span(rollout_id, attempt_id, span)
        Runner-->>Store: update_attempt(status)
    end
    Tracer->>Agent: Unapply instrumentation

上图显示了存储、追踪器和代理之间的数据流。实际上,它比这更复杂。代理不会主动发出 spans;它们是由追踪器通过钩子和工具化代理中使用的关键方法拦截的。追踪器使用回调(称为 exporter)来监控事件并记录到存储中。在 rollout 开始之前,runner 会进入一个 trace_context,然后再调用代理,将存储标识符连接到追踪器。每个 span 完成后都会通过 LightningSpanProcessor 流回存储,因此代理的工具化结果会进入 add_otel_span。如果代理的 rollout 方法返回一个数值奖励,runner 在最终化尝试之前会发出另一个 OpenTelemetry span。

钩子

Hook 实现是用户定义的的回调函数,允许你在 Runner 的生命周期中的特定点增强其行为。你可以使用钩子添加自定义日志记录,在 rollout 开始之前设置资源,或在 rollout 结束之后拆解它们。钩子可以在四个关键时刻触发:on_rollout_starton_trace_starton_trace_endon_rollout_end

用户应该特别注意 on_trace_endon_rollout_end 之间的区别。前者在追踪器退出 trace 上下文之前调用,而后者在 runner 处理完剩余的奖励和 spans,并在存储中最终化尝试之后调用。

sequenceDiagram
    autonumber
    participant Store
    participant Hooks
    participant Runner
    participant Tracer
    participant Agent

    Note over Runner,Hooks: Runner manages hooks as member

    loop Until no more rollouts
        Store-->>Runner: dequeue_rollout → AttemptedRollout
        Store-->>Runner: get_latest_resources

        Runner->>Hooks: on_rollout_start(agent, runner, rollout)
        Runner->>Agent: training_rollout / validation_rollout
        Tracer->>Agent: enter_trace_context
        activate Tracer
        Runner->>Hooks: on_trace_start(agent, runner, tracer, rollout)
        Note over Runner,Agent: Agent rollout omitted
        Runner->>Hooks: on_trace_end(agent, runner, tracer, rollout)
        Tracer->>Agent: exit_trace_context
        deactivate Tracer
        Agent->>Runner: final reward + extra spans (if any)
        Runner-->>Store: add_span(rollout_id, attempt_id, span)
        Runner->>Hooks: on_rollout_end(agent, runner, rollout, status)
    end

适配器

Adapter 是由 Algorithm 使用的组件,用于将来自 LightningStore 的原始数据转换为适合学习的格式。Runners 在执行期间将原始 spans 流式传输到存储中。 后来,算法查询这些 spans 并使用适配器将其转换为结构化数据,例如强化学习模型的训练示例。

例如,TracerTraceToTriplet 处理 OpenTelemetry spans 以创建 (prompt, response, reward) 三元组,这是许多 RL 微调算法的基本数据结构。

flowchart LR
    Runner -- (1) add_otel_span --> Store
    Store -- (2) query_spans --> Algorithm
    Algorithm -- (3) spans --> Adapter
    Adapter -- (4) transformed data --> Algorithm

LLM 代理

LLMProxy 是位于代理和算法资源之间的一个可选桥接组件。它充当所有 LLM 调用的集中端点。通常,代理 URL 会作为特殊资源添加到存储中,以便 Runner 在排队 rollout 时可以将其与其它资源一起获取。在 rollout 期间,runner 调用代理的 HTTP 端点,而不是直接调用模型后端。

这种设计提供了几个好处

  1. 工具化: 它自动捕获 LLM 交互的详细跟踪信息(提示、响应、元数据)并将其发送到存储,补充了追踪器,尤其是在代理的代码难以直接工具化时。
  2. 后端抽象: 它为各种 LLM 后端(OpenAI、Anthropic、本地模型)提供统一的接口,并可以添加诸如重试逻辑、速率限制和缓存等功能。
  3. 资源管理: 算法可以通过简单地交换代理正在使用的后端模型来动态更新代理使用的 LLM(例如,切换到新微调的模型),而无需中断代理的代码。

上述好处似乎都讨论在模型微调的上下文中。事实上,代理对于提示微调也很有用。算法可以将以下两种类型的端点之一注册到代理

  1. 由算法提供的端点: 如果算法正在内部更新 LLM 权重(例如,RL),它可以启动 LLM 推理引擎(即模型服务器)并将端点 URL 与代理注册。然后,代理将所有 LLM 调用转发到该端点。
  2. 第三方 LLM 端点: 如果算法没有更新 LLM 权重(例如,提示微调),它可以将第三方 LLM 端点注册到代理。

我们下面展示一个图表,说明代理如何适应整体数据流。

sequenceDiagram
    autonumber
    participant Algo as Algorithm
    participant LLMProxy as LLM Proxy
    participant Store
    participant Runner
    participant Agent

    Note over Algo,LLMProxy: Algorithm manages LLMProxy as member

    loop Over the Dataset
        Algo->>Algo: Launch LLM Inference Engine<br>(optional)
        Algo->>LLMProxy: Register Inference Engine<br>(optional)
        Algo-->>Store: enqueue_rollout
        LLMProxy->>Store: Proxy URL added as Resource
        Store-->>Runner: dequeue_rollout → AttemptedRollout
        Store-->>Runner: get_latest_resources
        Runner->>Agent: rollout + resources<br>(LLM Proxy URL as resource)
        loop Defined by Agent
            Agent-->>LLMProxy: LLM calls
            activate LLMProxy
            LLMProxy-->>Store: add_span or add_otel_span
            LLMProxy-->>Agent: LLM responses
            deactivate LLMProxy
            Agent-->>Runner: rewards
            Runner-->>Store: add_span or add_otel_span
        end
        Runner-->>Store: update_attempt("finished", status)
        Store-->>Algo: query_rollouts + spans
        Algo-->>Algo: Update LLM Weights<br>(optional)
    end

在这个图中,存储接收来自代理和 runner 的 spans。稍后我们将看到一个并行性问题,即代理和 runner 位于不同的机器上,并且 spans 需要从存储中获取一个特殊的计数器来确保 spans 的顺序。

训练器

Trainer 是一个高级编排器,它初始化并连接所有主要组件——AlgorithmRunnerLightningStoreTracerAdapterLLM ProxyHook。组件的生命周期可以与训练器一样长。训练器管理它们的生命周期并处理依赖注入,确保系统的每个部分都在一致且共享的环境中运行。

下面,我们演示了组件如何相互关联及其作用。我们首先阐明图中显示的作用和关系

  1. 拥有: 训练器直接构建和管理的组件(例如,runner、tracer)。
  2. 注入: 作为依赖项传递给其它组件的组件。
  3. 引用: 用于协调而不拥有关系的弱链接。
  4. 使用: 临时与之交互的组件。

例如,LightningStore 被注入到 AlgorithmRunner 中。TracerLitAgent 被注入到 runner 中。AdapterLLM Proxy 被注入到算法中。存储进一步被 runner 和算法分别注入到 tracer、adapter 和 LLM proxy 中。

flowchart TD
    %% === Left side: Algorithm domain ===
    subgraph L["Algorithm Side"]
        Algorithm["Algorithm<br>(no default)"]
        Adapter["Adapter<br>(TracerTraceToTriplet*)"]
        LLMProxy["LLM Proxy<br>(no default)"]
        Algorithm -.injects.-> Adapter
        Algorithm -.injects.-> LLMProxy
    end
    linkStyle 0,1 stroke:#896978,stroke-width:2px;

    %% === Middle: Core trainer and store ===
    subgraph M["Core"]
        Trainer["Trainer"]
        Store["LightningStore<br>(InMemory* default)"]
        Trainer --has--> Algorithm
        Trainer --has--> Store
        Trainer --has--> Adapter
        Trainer --has--> LLMProxy
    end
    linkStyle 2,3,4,5 stroke:#839791,stroke-width:2px;

    %% === Right side: Runner side ===
    subgraph R["Runner Side"]
        Runner["Runner<br>(LitAgentRunner* default)"]
        Tracer["Tracer<br>(AgentOpsTracer*)"]
        Hooks["Hooks (empty default)"]
        Agent["Agent<br>(LitAgent*)"]
        Runner -.injects.-> Tracer
        Runner -.injects.-> Store
        Runner -.injects.-> Agent
        Runner -.injects.-> Hooks
        Tracer -.injects.-> Store
        Hooks -.uses.-> Runner
        Hooks -.uses.-> Agent
        Hooks -.uses.-> Tracer
    end
    linkStyle 6,7,8,9,10 stroke:#896978,stroke-width:2px;
    linkStyle 11,12,13 stroke:#7a89c2,stroke-width:2px;

    %% === Cross-section connections ===
    Trainer --has--> Runner
    Trainer --has--> Tracer
    Trainer --has--> Hooks
    Trainer --uses--> Agent
    Algorithm -.injects.-> Store
    LLMProxy -.injects.-> Store
    Agent -.references.-> Trainer
    Runner -.references.-> Trainer
    Algorithm -.references.-> Trainer
    linkStyle 14,15,16 stroke:#839791,stroke-width:2px;
    linkStyle 17,20,21,22 stroke:#7a89c2,stroke-width:2px;
    linkStyle 18,19 stroke:#896978,stroke-width:2px;

    style L fill:none;
    style M fill:none;
    style R fill:none;

整合所有内容:强化学习示例 (VERL)

VERL 展示了算法如何使用共享基础设施。由于历史原因,代码位于 agentlightning.algorithm.verlagentlightning.verl 中。后者是遗留代码,并以令人困惑的方式重用术语,例如 Trainer。前者是一个薄包装器,符合新的算法接口。未来的版本将合并这两个版本。

强化学习旨在学习一种在状态中采取行动以最大化预期奖励的策略。对于代理,策略通常是语言模型。输入是提示(状态)。输出是生成的文本(动作)。数值评分判断质量(奖励)。(state, action, reward) 三元组 是基本的学习单元。

在 Agent-lightning 中,环境隐含在代理的工作流程中,该工作流程协调一个或多个 LLM 调用,并经常使用规则或额外的模型调用进行自我判断。在 rollout 期间,代理会发出包含 RL 训练所需的一切内容的 spans,包括 LLM 调用跟踪和数值判断/奖励信号。另一方面,“算法”具有更多的职责。

  1. 提供一个代理正在与之交互的、当前正在学习和改进的语言模型部署;
  2. 准备代理将执行的任务;
  3. 查询生成的 spans,提取三元组,并将它们转换为底层 RL 库可以使用的格式;
  4. 基于学习信号更新语言模型。

在 VERL 集成中,算法使用 vLLM 启动聊天完成端点,并使用 FSDP 进行分布式优化。它从数据集中排队任务。在 rollout 完成后,它查询 spans 并使用 TracerTraceToTriplet 将它们转换为三元组。VERL 的原生训练循环然后消耗这些三元组来更新模型权重。工作流程可以在下图总结。

sequenceDiagram
    autonumber
    participant vLLM as vLLM Chat<br>Completion Endpoint
    participant FSDP as FSDP / Megatron<br>Weights Optimizer
    participant Algo as Algorithm<br>Main Controller<br>(Main Process)
    participant Adapter as TracerTraceToTriplet
    participant LLMProxy as LLM Proxy
    participant Store as LightningStore
    participant Runner as Runner + Agent

    Note over Algo,LLMProxy: LLMProxy and Adapter are injected by Trainer as member
    Note over vLLM,Algo: Algorithm creates and owns vLLM and FSDP

    loop Over the Dataset in Batches
        Algo->>vLLM: Create Chat Completion Endpoint
        activate vLLM
        vLLM->>LLMProxy: Registered as Backend Endpoint
        LLMProxy->>Store: Proxy URL added as Resource
        par Over data samples in the batch
            Algo-->>Store: enqueue_rollout
            Store-->>Runner: Dequeue Rollout +<br>Resources (i.e., URL)
            loop One Rollout Attempt
                Runner-->>LLMProxy: LLM calls
                LLMProxy-->>vLLM: Forwarded LLM calls
                vLLM-->>LLMProxy: LLM responses
                LLMProxy-->>Store: add_span / add_otel_span
                LLMProxy-->>Runner: Forwarded LLM responses
                Runner-->>Store: add_span / add_otel_span <br> (by tracer, including rewards)
            end
            Runner-->>Store: update_attempt("finished", status)
        end
        Algo-->>Store: Poll for completed rollouts + spans
        Algo->>vLLM: Chat Completion Endpoint Sleeps
        deactivate vLLM
        Algo->>Adapter: adapt(spans)
        Adapter->>FSDP: Triplets (state, action, reward)
        activate FSDP
        FSDP-->>Algo: Updated LLM weights
        deactivate FSDP
    end

注意事项

  1. 图中存在算法注入或拥有的不同组件之间的交互,例如适配器的输出馈送到 FSDP 优化器。这只是为了说明的简单性,与实际实现略有不同,在实际实现中,算法主控制器会协调组件之间的数据流。

  2. 关于映射到 VERL。 VERL 使用经典的 RLHF 设置,其中每个动作都是单个 token,状态是直到该 token 的完整对话历史,并且奖励在最后给出。这与我们的设置非常不同,在我们的设置中,每个动作实际上是一段文本,尽管它们都被称为 RL!因此,在适配器生成三元组之后,算法将每个 (state, action, reward) 转换为一个 VERL 轨迹 (DataProto),其中包含诸如 input_idsposition_idsattention_masktoken_level_scores 之类的键。该转换发生在三元组生成之后,图中未显示。

执行策略和并行性

读者可能已经从上图观察到,(1) 运行器和代理以及 (2) 算法之间绝对没有通信。它们唯一的重叠之处是 TrainerLightningStore。这种观察结果在训练器部分内的图表中非常清晰。这种设计允许我们灵活地独立扩展运行器和算法,这对于大规模训练至关重要。

Agent-lightning 封装了两个可执行包:一个运行器包(RunnerTracerHookLitAgent)和一个算法包(AlgorithmAdapterLLM Proxy)。两者共享 LightningStore。训练器初始化并连接这些包。

graph TD
    subgraph Runner_Side["Runner Bundle"]
        direction LR
        R[Runner] --- T[Tracer] --- H[Hooks] --- A1[Agent]
    end

    subgraph Algorithm_Side["Algorithm Bundle"]
        direction LR
        ALG[Algorithm] --- AD[Adapter] --- LLM[LLM Proxy]
    end

    S[(Store)]
    TR[Trainer]

    Runner_Side <--> S
    Algorithm_Side <--> S
    TR --> Runner_Side
    TR --> Algorithm_Side

    linkStyle 0,1,2,3,4 opacity:0;

一个 执行策略由训练器定义和拥有,它控制算法和运行器包的放置、连接、扩展和中止方式。它具有四个主要目的。

执行策略首先确定**包放置**——这两个包是在同一个线程、进程、机器上运行,还是跨多个机器运行。它们还定义**存储管理**,包装存储并指定数据如何在包之间共享。

在**可扩展性**方面,该策略可以在多个线程、进程或机器上复制运行器包,以扩展运行器侧的吞吐量。由于并行化的复杂性,算法侧仍然是单进程的。成熟的框架,如 *DeepSpeed* 和 *Megatron* 已经支持分布式模型训练,因此算法包的扩展委托给这些实现。

**中止处理**是另一个核心职责。中止可能由正常退出、任一包中的故障或用户中断触发。训练器必须包含用于包的取消接口,以便可以干净地中止包。当算法包正常退出时,该策略会向运行器包发出终止信号。如果运行器首先退出,则不会向算法发送信号,因为它可能仍在处理完成的 rollout。在发生故障或用户中断的情况下,该策略会向两个包发出中止信号;如果一个包未能响应,该策略应尝试强制终止。

Agent-lightning 当前提供两种执行策略:**共享内存**和**客户端-服务器**,如下节所述。

共享内存策略

SharedMemoryExecutionStrategy 将算法和运行器包作为单个进程中的线程运行。该策略使用 LightningStoreThreaded 包装存储,该存储使用锁来保护调用,以确保安全的并发性。

这对于轻量级调试很有用,因为组件共享一个 Python 堆,并避免序列化。它不适合重型 RL 训练或计算密集型代理。

flowchart TB
    subgraph MainProcess
        direction TB
        subgraph AlgorithmThread [Thread 0]
            Algorithm[Algorithm bundle]
        end
        subgraph RunnerThread1 [Thread 1]
            Runner1[Runner bundle #1]
        end
        subgraph RunnerThread2 [Thread 2]
            Runner2[Runner bundle #2]
        end
        subgraph RunnerThread3 [Thread 3]
            RunnerN[Runner bundle #N]
        end
        LightningStoreFacade[LightningStoreThreaded]
        BaseStore[Underlying LightningStore]
    end
    Algorithm -- async calls --> LightningStoreFacade
    Runner1 -- async calls --> LightningStoreFacade
    Runner2 -- async calls --> LightningStoreFacade
    RunnerN -- async calls --> LightningStoreFacade
    LightningStoreFacade -->|thread-safe delegates| BaseStore

您可以配置哪个角色在主线程上运行。如果主线程运行算法,则它可以生成多个运行器线程。如果它运行一个运行器,则 n_runners 必须为 1,并且运行器位于主线程上。

客户端-服务器策略

ClientServerExecutionStrategy 在进程之间划分关注点。算法包启动一个 LightningStoreServer(HTTP API),它包装底层的存储。运行器通过 LightningStoreClient 连接,通过 REST 调用相同的接口。服务器嵌入一个客户端,以支持算法启动的子进程(例如,LLM 代理工作器),这些子进程需要通过相同的 API 与算法的进程通信。

目前,这种设计在服务器端引入了一个额外的包装器(如图所示),这有助于调试并提高了容错能力。我们可能会在未来重新审视这种设计,并强制客户端成为与存储通信的唯一方式。

flowchart TD
    subgraph Algorithm Process Group
        subgraph StoreServer[LightningStoreServer]
            StoreHttpClient[HTTP Client]
            StoreHttpServer[HTTP Server]
            StoreWrapper[LightningStore Wrapper]
            StoreHttpClient -- HTTP --> StoreHttpServer
        end
        subgraph Algorithm Bundle
            Algorithm[Algorithm Main Process]
            subgraph Another subprocess
                LLMProxy[LLM Proxy]
            end
        end
        LLMProxy -- async calls --> StoreHttpClient
        Algorithm -- async calls --> StoreWrapper
    end
    subgraph RunnerSide ["Runner Side"]
        subgraph Runner Process 1
            Runner1[Runner bundle #1]
            Runner1 -- async calls --> LightningStoreClient1
            LightningStoreClient1[LightningStoreClient]
        end
        subgraph Runner Process 2
            Runner2[Runner bundle #2]
            Runner2 -- async calls --> LightningStoreClient2
            LightningStoreClient2[LightningStoreClient]
        end
        subgraph Runner Process N
            RunnerN[Runner bundle #N]
            RunnerN -- async calls --> LightningStoreClientN
            LightningStoreClientN[LightningStoreClient]
        end
    end
    LocalStore[Underlying LightningStore]
    StoreHttpServer -->|delegates| StoreWrapper
    StoreWrapper -->|delegates| LocalStore
    LightningStoreClient1 -- HTTP --> StoreHttpServer
    LightningStoreClient2 -- HTTP --> StoreHttpServer
    LightningStoreClientN -- HTTP --> StoreHttpServer

    style RunnerSide fill:none;

在线/连续学习

连续学习在运行器报告任务和 spans 的同时保持算法循环运行。与批处理模式的主要区别

  1. 算法不会从固定数据集排队 rollout。运行器会自发报告任务/rollout 和 spans。
  2. 算法可以等待具有预期 rollout ID 集合的 rollout,但更常见的是轮询新的 rollout 和 spans,或等待到达计数。
  3. Runner 通过 step(task) 一次处理一个 rollout,而不是耗尽任务队列。当开始 rollout 时,它会通知存储,以便存储记录它。
  4. 用户或更高级别的循环控制下一个步骤使用哪些资源以及何时重试。

SpansAdapter 实现和 LLM Proxy 的工作方式相同。

sequenceDiagram
    autonumber
    actor User
    participant Runner
    participant Agent
    participant Store as LightningStore
    participant Algorithm

    Note over Algorithm: Algorithm is long-running and loops continuously

    loop Continuous Learning Loop
        activate User
        opt Decide what to do next
            User-->>Store: get_resources_by_id
            Store-->>User: Resources
            User-->>User: Prepare input for next step
        end
        User->>Runner: step(input, resources)
        activate Runner
        Runner-->>Store: Notify: start_rollout(input)
        Runner->>Agent: rollout(input, resources)
        Agent-->>Runner: add_span / reward spans
        Runner-->>Store: add_span or add_otel_span
        Runner-->>Store: update_attempt(status="finished")
        deactivate Runner
        deactivate User
        Algorithm->>Store: poll for new rollouts and spans
        opt If there is enough new data
            Store-->>Algorithm: new spans
            Algorithm->>Algorithm: adapt spans → learning signal
            Algorithm->>Store: update_resources
        end
    end