compute_interceptor.h 1.6 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
  void IncreaseReady(int64_t up_id);
  void DecreaseBuff(int64_t down_id);
  bool IsInputReady();
  bool CanWriteOutput();
34
  bool ShouldReset();
35

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

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

42 43 44
  void ReceivedStop(int64_t up_id);
  void TryStop();

45
 private:
46
  bool is_source_{false};
47
  bool is_last_{false};
48
  int64_t step_{0};
49

50 51 52 53
  // 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_{};
54 55 56

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

}  // namespace distributed
}  // namespace paddle