simu.cpp 9.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291
  1. #include <vector>
  2. #include <pthread.h>
  3. #include <errno.h>
  4. #include <cstdio>
  5. #include <cstdlib>
  6. #include <iostream>
  7. #include <queue>
  8. #include "simu.h"
  9. #include "basetask.h"
  10. #include "basesched.h"
  11. using namespace std;
  12. //#define DBG_THREAD
  13. #ifdef DBG_THREAD
  14. #define _D(X) X
  15. #else
  16. #define _D(X)
  17. #endif
  18. #define forn(i,n) for(int i = 0; i < (int)(n); ++i)
  19. #define DBG(X) _D(cerr << #X << " = " << X << endl;)
  20. typedef struct cpu_ctx {
  21. int pid;
  22. int remaining;
  23. } cpu_ctx_t;
  24. /* "Globales" */
  25. // static int cur_pid;
  26. //Habria que modificar el nombre de cur_pid
  27. //lo mantuve para que sea mas facil modificar
  28. //codigo y mirar que hacia antes
  29. static vector<cpu_ctx_t> contexts;
  30. static unsigned int cur_time;
  31. static pthread_mutex_t m_sched;
  32. enum status_t {ST_EXIT, ST_IO, ST_CPU};
  33. //enum status_t {ST_CPU, ST_IO, ST_EXIT};
  34. struct task_data {
  35. pthread_t tid;
  36. int pid, running;
  37. status_t blk;
  38. int blkms;
  39. TaskBase* tsk;
  40. vector<int>* prms;
  41. pthread_mutex_t mutex;
  42. int lastcpu;
  43. };
  44. static task_data* tsks;
  45. void* task_thread(task_data* tsk) {
  46. pthread_mutex_lock(&tsk->mutex);
  47. //fprintf(stderr, "Thread %d\n", tsk->pid);
  48. _D(cerr << "tsk->tsk( " << tsk->pid << ", *tsk->prms)" << endl;)
  49. tsk->tsk(tsk->pid, *tsk->prms); // Run!
  50. //
  51. tsk->blk = ST_EXIT;
  52. _D(cerr << "TSK unlock(sched) exit" << endl;) pthread_mutex_unlock(&m_sched);
  53. return NULL;
  54. }
  55. /* Funciones llamadas por el scheduler */
  56. int current_pid(int cpu) { return contexts[cpu].pid; }
  57. int current_remaining(int cpu) { return contexts[cpu].remaining; }
  58. unsigned int current_time(void) { return cur_time; }
  59. /* Funciones llamadas por las tareas */
  60. static void uso_X(task_data* tsk, enum status_t tp, unsigned int ms) {
  61. if (ms == 0) return;
  62. tsk->blk = tp;
  63. tsk->blkms = ms;
  64. _D(cerr << "TSK unlock(sched)" << endl;) pthread_mutex_unlock(&m_sched);
  65. pthread_mutex_lock(&tsk->mutex); _D(cerr << "TSK lock(mutex["<<tsk-tsks<<"])" << endl;)
  66. }
  67. //Se agrega pid para poder ver obtener
  68. //los datos en tsks
  69. void uso_CPU(int pid, unsigned int ms) {
  70. uso_X(&(tsks[pid]), ST_CPU, ms);
  71. }
  72. void uso_IO(int pid, unsigned int ms) {
  73. uso_X(&(tsks[pid]), ST_IO, ms);
  74. }
  75. void simulate(SchedBase& sch, std::vector<ptsk>& lote, const Settings& settings) {
  76. int n = lote.size();
  77. cout.flush();
  78. if (settings.output_log != "-") {
  79. /* Oh yeah! */
  80. freopen(settings.output_log.c_str(), "wt", stdout);
  81. }
  82. tsks = (task_data*)malloc(sizeof(task_data)*n);
  83. if (!tsks) { perror("malloc(task_data)"); return; }
  84. forn(i, n) {
  85. tsks[i].pid = i;
  86. pthread_mutex_init(&(tsks[i].mutex), NULL);
  87. pthread_mutex_lock(&(tsks[i].mutex));
  88. tsks[i].tsk = lote[i].tsk;
  89. tsks[i].prms = &(lote[i].prms);
  90. tsks[i].running = 0;
  91. tsks[i].blk = ST_CPU;
  92. tsks[i].blkms = 0;
  93. tsks[i].lastcpu = -1;
  94. if (pthread_create(&(tsks[i].tid), NULL, (void*(*)(void*))task_thread, (void*)(&(tsks[i]))) < 0) {
  95. perror("Lanzando la tarea (no use muchas tareas, < 500)"); return;
  96. }
  97. }
  98. pthread_mutex_init(&m_sched, NULL);
  99. pthread_mutex_lock(&m_sched);
  100. int finished = 0;
  101. //Inicializa los cpus
  102. contexts = vector<cpu_ctx_t>(settings.num_cores);
  103. for (uint i = 0; i <settings.num_cores; i++) {
  104. contexts[i].pid = IDLE_TASK;
  105. contexts[i].remaining = 0;
  106. }
  107. cur_time = 0;
  108. priority_queue<pair<int, int> > load;
  109. forn(i, n) load.push(make_pair(-lote[i].start, -i));
  110. vector<pair<unsigned int, int> > dlote(0);
  111. forn(i, n){
  112. if(lote[i].end > 0)
  113. dlote.push_back(make_pair(lote[i].end, i));
  114. }
  115. priority_queue<pair<int, int> > unblock;
  116. int context_remain = 0; /* Remaining context_switch ticks */
  117. while (finished < n || context_remain) {
  118. if (settings.verbose) {
  119. //cerr << "--- sched, tm=" << cur_time << " pid=" << cur_pid;
  120. cerr << "--- sched, tm=" << cur_time << endl;
  121. for(uint i = 0; i < settings.num_cores; i++) {
  122. int pid=contexts[i].pid;
  123. cerr << "cpu " << i << " pid = " << pid << " rem " << contexts[i].remaining;
  124. if (pid != IDLE_TASK) { cerr << " [" << pid << " ST:"<< tsks[pid].blk << " ms:" << tsks[pid].blkms << "]"; }
  125. cerr << endl;
  126. }
  127. cerr << "--------------" << cur_time << endl;
  128. }
  129. // Load de las tareas
  130. while (!load.empty() && load.top().first >= -(int)cur_time) {
  131. int pid = -load.top().second;
  132. int deadline = lote[-load.top().second].end;
  133. load.pop();
  134. tsks[pid].running = 1;
  135. sch.load(pid,deadline);
  136. cout << "LOAD " << cur_time << " " << pid << endl;
  137. }
  138. vector<int> to_unblock;
  139. while (!unblock.empty() && unblock.top().first >= -(int)cur_time) {
  140. int pid = -unblock.top().second; unblock.pop();
  141. _D(cerr << "SCH unblock(" << pid << ")" << endl;)
  142. sch.unblock(pid); // pid
  143. to_unblock.push_back(pid);
  144. int unblocked = 0;
  145. for(uint i = 0; i < settings.num_cores && !unblocked; i++) {
  146. int it = contexts[i].pid;
  147. if (it == pid) {
  148. tsks[pid].blkms = -2;
  149. unblocked = true;
  150. }
  151. }
  152. if (!unblocked) {
  153. tsks[pid].blk = ST_CPU;
  154. tsks[pid].blkms = 0;
  155. }
  156. }
  157. //Itera por cada cpu
  158. for(uint cpu= 0; cpu < settings.num_cores; cpu++){
  159. int cpu_pid= contexts[cpu].pid;
  160. int cpu_context_remain = contexts[cpu].remaining;
  161. if (!cpu_context_remain) {
  162. int npid;
  163. if (cpu_pid == IDLE_TASK) {
  164. npid = sch.tick(cpu, TICK);
  165. _D(cerr << "SCH tick( " << cpu << " ,TICK) -> " << npid << endl;)
  166. } else {
  167. if (tsks[cpu_pid].blk == ST_CPU && !tsks[cpu_pid].blkms){
  168. _D(cerr << "SCH unlock(" << cpu_pid << ") tick" << endl;) pthread_mutex_unlock(&(tsks[cpu_pid].mutex));
  169. pthread_mutex_lock(&m_sched); _D(cerr << "SCH lock(sched)" << endl;)
  170. }
  171. switch (tsks[cpu_pid].blk) {
  172. case ST_EXIT:
  173. finished++;
  174. npid = sch.tick(cpu, EXIT);
  175. cout << "EXIT " << cur_time << " " << cpu_pid << " " << cpu << endl;
  176. _D(cerr << "SCH tick( " << cpu << " ,EXIT) -> " << npid << endl;)
  177. tsks[cpu_pid].running = 0;
  178. break;
  179. case ST_IO:
  180. if (tsks[cpu_pid].blkms >= 0) {
  181. unblock.push(make_pair(-(cur_time+tsks[cpu_pid].blkms), -cpu_pid));
  182. tsks[cpu_pid].blkms = -1;
  183. }
  184. if (tsks[cpu_pid].blkms == -2) {
  185. tsks[cpu_pid].blk = ST_CPU;
  186. tsks[cpu_pid].blkms = 0;
  187. }
  188. cout << "BLOCK " << cur_time << " " << cpu_pid << endl;
  189. npid = sch.tick(cpu, BLOCK);
  190. _D(cerr << "SCH tick( " << cpu << " ,BLOCK) -> " << npid << endl;)
  191. break;
  192. case ST_CPU:
  193. if (tsks[cpu_pid].blkms) {
  194. tsks[cpu_pid].blkms--;
  195. } else {
  196. cerr << "FATAL ERROR, this should not happend" << endl;
  197. }
  198. npid = sch.tick(cpu, TICK);
  199. _D(cerr << "SCH tick( " << cpu << " ,TICK) -> " << npid << endl;)
  200. break;
  201. }
  202. }
  203. if (npid == IDLE_TASK) {
  204. // cerr << "SCH unlock(sched)" << endl;pthread_mutex_unlock(&m_sched);
  205. } else {
  206. if (npid < 0 || npid >= n) { cerr << "Error!, scheduler sent an invalid pid="<<npid<< endl; return; }
  207. if (!tsks[npid].running) { cerr << "Error!, scheduler sent pid="<<npid << " but that process has exited." << endl; return; }
  208. // FIXME! Borrar si funciona bien. ¡¡ No hay que negar un enum, no tiene sentido eso !!
  209. // if (!tsks[npid].blk == ST_IO) { cerr << "Error!, scheduler sent pid="<<npid << " but that process is still blocked." << endl; return; }
  210. if (tsks[npid].blk == ST_EXIT) { cerr << "Error!, scheduler sent pid="<<npid << " but that process is still blocked." << endl; return; }
  211. }
  212. if (cpu_pid != npid){
  213. if (npid != IDLE_TASK) {
  214. if (settings.switch_cost > 0) {
  215. contexts[cpu].remaining += settings.switch_cost;
  216. }
  217. if ((int) cpu != tsks[npid].lastcpu && tsks[npid].lastcpu != -1) {
  218. contexts[cpu].remaining += settings.migrate_cost;
  219. }
  220. tsks[npid].lastcpu = cpu;
  221. }
  222. }
  223. cpu_pid = npid;
  224. } else {
  225. contexts[cpu].remaining--;
  226. }
  227. contexts[cpu].pid = cpu_pid;
  228. }
  229. //Hasta aca el codigo para cada CPU
  230. /* Unblock tasks at the end of the tick */
  231. for(int j=0; j<(int)to_unblock.size(); j++) cout << "UNBLOCK " << cur_time << " " << to_unblock[j] << endl;
  232. context_remain= 0;
  233. //Muestra que esta realizando cada cpu
  234. //y calcula si hay contexto total restante (ver while)
  235. for(uint i= 0; i < contexts.size(); i++){
  236. context_remain += contexts[i].remaining;
  237. if (contexts[i].remaining /*context_remain*/) {
  238. cout << "# CONTEXT CPU " << i << " " << cur_time << endl;
  239. } else{
  240. cout << "CPU "<< cur_time << " " << contexts[i].pid << " " << i <<endl;
  241. }
  242. }
  243. //Muestra si se cumple el deadline de alguna tarea en este tick
  244. forn(i, dlote.size()){
  245. if(dlote[i].first == cur_time)
  246. cout << "DEADLINE "<< cur_time << " " << dlote[i].second << endl;
  247. }
  248. cur_time++;
  249. }
  250. forn(i, n) {
  251. pthread_join( tsks[i].tid, NULL);
  252. }
  253. free(tsks);
  254. }