compute_interceptor.h 1.7 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
// Copyright (c) 2021 PaddlePaddle Authors. All Rights Reserved.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
//     http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

#pragma once

17 18
#include <utility>

19 20 21 22 23 24 25 26 27 28 29
#include "paddle/fluid/distributed/fleet_executor/interceptor.h"

namespace paddle {
namespace distributed {

class ComputeInterceptor : public Interceptor {
 public:
  ComputeInterceptor(int64_t interceptor_id, TaskNode* node);

  void PrepareDeps();

30 31 32 33 34
  void IncreaseReady(int64_t up_id);
  void DecreaseBuff(int64_t down_id);
  bool IsInputReady();
  bool CanWriteOutput();

35
  void SendDataReadyToDownStream();
36
  void ReplyCompletedToUpStream();
37

38
  void Run();
39 40
  void Compute(const InterceptorMessage& msg);

41 42 43 44
  void HandleStop(const InterceptorMessage& msg) override;
  void ReceivedStop(int64_t up_id);
  void TryStop();

45
 private:
46 47 48 49 50 51
  // FIXME(wangxi): if use step_ and max_steps_, how to restart step_ from 0
  int64_t step_{0};
  // upstream_id-->(max_ready_size, ready_size)
  std::map<int64_t, std::pair<int64_t, int64_t>> in_readys_{};
  // downstream_id-->(max_buffer_size, used_size)
  std::map<int64_t, std::pair<int64_t, int64_t>> out_buffs_{};
52 53 54

  bool received_stop_{false};
  std::map<int64_t, bool> in_stops_{};
55 56 57 58
};

}  // namespace distributed
}  // namespace paddle