DiskReadThreadPool.hpp 7.52 KB
Newer Older
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
/******************************************************************************
 *                      DiskReadThreadPool.hpp
 *
 * This file is part of the Castor project.
 * See http://castor.web.cern.ch/castor
 *
 * Copyright (C) 2003  CERN
 * This program is free software; you can redistribute it and/or
 * modify it under the terms of the GNU General Public License
 * as published by the Free Software Foundation; either version 2
 * of the License, or (at your option) any later version.
 * This program is distributed in the hope that it will be useful,
 * but WITHOUT ANY WARRANTY; without even the implied warranty of
 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
 * GNU General Public License for more details.
 * You should have received a copy of the GNU General Public License
 * along with this program; if not, write to the Free Software
 * Foundation, Inc., 59 Temple Place - Suite 330, Boston, MA 02111-1307, USA.
 *
 * 
 *
 * @author Castor Dev team, castor-dev@cern.ch
 *****************************************************************************/

#pragma once

#include "castor/tape/tapeserver/daemon/DiskReadTask.hpp"
28
#include "castor/tape/tapeserver/daemon/TaskWatchDog.hpp"
29
30
31
#include "common/threading/BlockingQueue.hpp"
#include "common/threading/Threading.hpp"
#include "common/threading/AtomicCounter.hpp"
Victor Kotlyar's avatar
Victor Kotlyar committed
32
#include "common/log/LogContext.hpp"
Victor Kotlyar's avatar
Victor Kotlyar committed
33
#include "common/Timer.hpp"
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
#include <vector>
#include <stdint.h>

namespace castor {
namespace tape {
namespace tapeserver {
namespace daemon {
  class MigrationTaskInjector;
class DiskReadThreadPool {
public:
  /**
   * Constructor. The constructor creates the threads for the pool, but does not
   * start them.
   * @param nbThread Number of thread for reading files 
   * @param maxFilesReq maximal number of files we might require 
49
   * within a single request to the task injector
50
   * @param maxBytesReq maximal number of bytes we might require
51
52
   *  within a single request a single request to the task injector
   * @param lc log context for logging purpose
53
54
   */
  DiskReadThreadPool(int nbThread, uint64_t maxFilesReq,uint64_t maxBytesReq, 
55
          castor::tape::tapeserver::daemon::MigrationWatchDog & migrationWatchDog,
Victor Kotlyar's avatar
Victor Kotlyar committed
56
          cta::log::LogContext lc, const std::string & remoteFileProtocol,
57
          const std::string & xrootPrivateKeyPath, uint16_t xrootTimeout);
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
  
  /**
   * Destructor.
   * Simply destroys the thread, which should have been joined by the caller
   * before using waitThreads()
   */
  ~DiskReadThreadPool();
  
  /**
   * Starts the threads which were created at construction time.
   */
  void startThreads();
  
  /**
   * Waits for threads completion of all threads. Should be called before
   * destructor
   */
  void waitThreads();
  
  /**
   * Adds a DiskReadTask to the tape pool's queue of tasks.
   * @param task pointer to the new task. The thread pool takes ownership of the
   * task and will delete it after execution. Push() is not protected against races
   * with finish() as the task injection is done from a single thread (the task
   * injector)
   */
  void push(DiskReadTask *task);
  
  /**
   * Injects as many "end" tasks as there are threads in the thread pool.
   * A thread getting such an end task will exit. This method is called by the
   * task injector at the end of the tape session. It is not race protected. 
   * See push()
   */
  void finish();
  
  /**
   * Sets up the pointer to the task injector. This cannot be done at
   * construction time as both task injector and read thread pool refer to 
   * each other. This function should be called before starting the threads.
   * This is used for the feedback loop where the injector is requested to
   * fetch more work by the read thread pool when the task queue of the thread
   * pool starts to run low.
   */
  void setTaskInjector(MigrationTaskInjector* injector){
      m_injector = injector;
  }
private:
  /**
   * When the last thread finish, we log all m_pooldStat members + message
   * at the given level
   * @param level
   * @param message
   */
  void logWithStat(int level, const std::string& message);
  /** 
   * Get the next task to execute and if there is not enough tasks in queue,
   * it will ask the TaskInjector to get more jobs.
   * @return the next task to execute
   */
Victor Kotlyar's avatar
Victor Kotlyar committed
118
  DiskReadTask* popAndRequestMore(cta::log::LogContext & lc);
119
120
121
122
123
124
125
126
127
128
129
  
  /**
   * When a thread finishm it call this function to Add its stats to one one of the
   * Threadpool
   * @param threadStats
   */
  void addThreadStats(const DiskStats& stats);
  
  /**
   To protect addThreadStats from concurrent calls
   */
130
  cta::threading::Mutex m_statAddingProtection;
131
132
133
134
135
136
  
  /**
   * Aggregate all threads' stats 
   */
  DiskStats m_pooldStat;
  
137
138
139
  /**
   * Measure the thread pool's lifetime
   */
Victor Kotlyar's avatar
Victor Kotlyar committed
140
  cta::utils::Timer m_totalTime;
141
  
142
143
144
145
  /**
   * A disk file factory, that will create the proper type of file access class,
   * depending on the received path
   */
146
  diskFile::DiskFileFactory m_diskFileFactory;
147
148
149
150
  
  /**
   * Subclass of the thread pool's worker thread.
   */
151
  class DiskReadWorkerThread: private cta::threading::Thread {
152
153
154
  public:
    DiskReadWorkerThread(DiskReadThreadPool & parent):
    m_parent(parent),m_threadID(parent.m_nbActiveThread++),m_lc(parent.m_lc) {
Victor Kotlyar's avatar
Victor Kotlyar committed
155
156
       cta::log::LogContext::ScopedParam param(m_lc, cta::log::Param("threadID", m_threadID));
       m_lc.log(cta::log::INFO,"DisReadThread created");
157
    }
158
159
    void start() { cta::threading::Thread::start(); }
    void wait() { cta::threading::Thread::wait(); }
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
  private:
    void logWithStat(int level, const std::string& message);
    /*
     * For measuring how long  are the the different steps 
     */
    DiskStats m_threadStat;
    
    /** Pointer to the thread pool, allowing calls to popAndRequestMore,
     * and calling finish() on the task injector when the last thread
     * is finishing (thanks to the actomic counter m_parent.m_nbActiveThread) */
    DiskReadThreadPool & m_parent;
    
    /** The sequential ID of the thread, used in logs */
    const int m_threadID;
    
    /** The local copy of the log context, allowing race-free logging with context
     between threads. */
Victor Kotlyar's avatar
Victor Kotlyar committed
177
    cta::log::LogContext m_lc;
178
179
180
181
182
183
184
185
186
187
188
    
    /** The execution thread: pops and executes tasks (potentially asking for
     more) and calls task injector's finish() on exit of the last thread. */
    virtual void run();
  };
  
  /** Container for the threads */
  std::vector<DiskReadWorkerThread *> m_threads;
  
  /** The queue of pointer to tasks to be executed. We own the tasks (they are 
   * deleted by the threads after execution) */
189
  cta::threading::BlockingQueue<DiskReadTask *> m_tasks;
190
  
191
192
193
194
195
  /**
   * Reference to the watchdog, for error reporting.
   */
  castor::tape::tapeserver::daemon::MigrationWatchDog & m_watchdog;
  
196
197
198
  /** The log context. This is copied on construction to prevent interferences
   * between threads.
   */
Victor Kotlyar's avatar
Victor Kotlyar committed
199
  cta::log::LogContext m_lc;
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
  
  /** Pointer to the task injector allowing request for more work, and 
   * termination signaling */
  MigrationTaskInjector* m_injector;
  
  /** The maximum number of files we ask per request. This value is also used as
   * a threshold (half of it, indeed) to trigger the request for more work. 
   * Another request for more work is also triggered when the task FIFO gets empty.*/
  const uint64_t m_maxFilesReq;
  
  /** Same as m_maxFilesReq for size per request. */
  const uint64_t m_maxBytesReq;
  
  /** An atomic (i.e. thread safe) counter of the current number of thread (they
   are counted up at creation time and down at completion time) */
215
  cta::threading::AtomicCounter<int> m_nbActiveThread;
216
217
218
};

}}}}