Просмотр исходного кода

Read subproblem responses from the dpnodes

We do check that the response is correct, but we do not yet reconstruct
the final answer.
Ian Goldberg 14 лет назад
Родитель
Сommit
76d978789e
1 измененных файлов с 50 добавлено и 14 удалено
  1. 50 14
      controller.cc

+ 50 - 14
controller.cc

@@ -137,13 +137,13 @@ struct SubproblemProgress : Subproblem {
     unsigned int max_workers;
     unsigned int max_workers;
 
 
     // Have we found a solution?
     // Have we found a solution?
-    int solved;
+    bool solved;
     // The solution, if found.
     // The solution, if found.
     ZZ solution;
     ZZ solution;
 
 
     SubproblemProgress(unsigned short id, const ZZ &b, const ZZ &t,
     SubproblemProgress(unsigned short id, const ZZ &b, const ZZ &t,
 	    const ZZ &m, const ZZ &o, unsigned int dpf) :
 	    const ZZ &m, const ZZ &o, unsigned int dpf) :
-	    Subproblem(id, b, t, m, o, dpf), solved(0) {
+	    Subproblem(id, b, t, m, o, dpf), solved(false) {
 	// How many DPnodes should we use for a problem of this size?
 	// How many DPnodes should we use for a problem of this size?
 	desired_dpnodes = 2;
 	desired_dpnodes = 2;
 	// How many workers would we like to use?
 	// How many workers would we like to use?
@@ -156,8 +156,8 @@ struct SubproblemProgress : Subproblem {
 	}
 	}
     }
     }
 
 
-    // Stop all dpnodes and workers and reset to unstarted state
-    void reset(void) {
+    // Stop all dpnodes and workers
+    void stop(void) {
 	BESet::iterator iter;
 	BESet::iterator iter;
 	unsigned char stopcmd[1] = { 'S' };
 	unsigned char stopcmd[1] = { 'S' };
 
 
@@ -180,12 +180,16 @@ struct SubproblemProgress : Subproblem {
 
 
     // Dump for debugging purposes
     // Dump for debugging purposes
     void dump(ostream &os) const {
     void dump(ostream &os) const {
+	os << "    Subproblem " << problemid << "\n";
 	os << "    dpnodes (" << dpnodes.size() << "):\n";
 	os << "    dpnodes (" << dpnodes.size() << "):\n";
 	besetdump(dpnodes, os);
 	besetdump(dpnodes, os);
 	os << "    workers (" << workers.size() << "):\n";
 	os << "    workers (" << workers.size() << "):\n";
 	besetdump(workers, os);
 	besetdump(workers, os);
 	os << "    ipports (" << ipports.size() << "):\n";
 	os << "    ipports (" << ipports.size() << "):\n";
 	ipportsetdump(ipports, os);
 	ipportsetdump(ipports, os);
+	if (solved) {
+	    os << "    solution: " << solution << "\n\n";
+	}
     }
     }
 
 
     void worker_write(struct bufferevent *bev) {
     void worker_write(struct bufferevent *bev) {
@@ -228,14 +232,15 @@ static void find_subproblem_for_dpnode(vector<SubproblemProgress> &spv,
 		ctrlstate.dpnodes.idle.size() > 0) {
 		ctrlstate.dpnodes.idle.size() > 0) {
 	    // Get the first idle DPnode
 	    // Get the first idle DPnode
 	    BESet::iterator beviter = ctrlstate.dpnodes.idle.begin();
 	    BESet::iterator beviter = ctrlstate.dpnodes.idle.begin();
+	    struct bufferevent *firstbev = *beviter;
 
 
 	    // Allocate it to the subproblem
 	    // Allocate it to the subproblem
-	    spiter->dpnodes.insert(*beviter);
-	    ctrlstate.dpnodes.working[*beviter] = &(*spiter);
-	    ctrlstate.dpnodes.idle.erase(*beviter);
+	    spiter->dpnodes.insert(firstbev);
+	    ctrlstate.dpnodes.working[firstbev] = &(*spiter);
+	    ctrlstate.dpnodes.idle.erase(firstbev);
 
 
 	    // Tell it to start listening for DPs
 	    // Tell it to start listening for DPs
-	    spiter->bev_write(*beviter);
+	    spiter->bev_write(firstbev);
 	}
 	}
     }
     }
 }
 }
@@ -254,14 +259,15 @@ static void find_subproblem_for_worker(vector<SubproblemProgress> &spv)
 		ctrlstate.workers.idle.size() > 0) {
 		ctrlstate.workers.idle.size() > 0) {
 	    // Get the first idle worker
 	    // Get the first idle worker
 	    BESet::iterator beviter = ctrlstate.workers.idle.begin();
 	    BESet::iterator beviter = ctrlstate.workers.idle.begin();
+	    struct bufferevent *firstbev = *beviter;
 
 
 	    // Allocate it to the subproblem
 	    // Allocate it to the subproblem
-	    spiter->workers.insert(*beviter);
-	    ctrlstate.workers.working[*beviter] = &(*spiter);
-	    ctrlstate.workers.idle.erase(*beviter);
+	    spiter->workers.insert(firstbev);
+	    ctrlstate.workers.working[firstbev] = &(*spiter);
+	    ctrlstate.workers.idle.erase(firstbev);
 
 
 	    // Tell it to start working on the subproblem
 	    // Tell it to start working on the subproblem
-	    spiter->worker_write(*beviter);
+	    spiter->worker_write(firstbev);
 	}
 	}
     }
     }
 }
 }
@@ -309,6 +315,7 @@ typedef enum {
     CCSTATE_START,
     CCSTATE_START,
     CCSTATE_DPWAITRESP,
     CCSTATE_DPWAITRESP,
     CCSTATE_DPLISTENING,
     CCSTATE_DPLISTENING,
+    CCSTATE_DPEXPON,
     CCSTATE_END
     CCSTATE_END
 } CCState;
 } CCState;
 
 
@@ -323,6 +330,7 @@ static void controller_dpnode_reader(struct bufferevent *bev, void *ctx)
     struct evbuffer *input = bufferevent_get_input(bev);
     struct evbuffer *input = bufferevent_get_input(bev);
     ControllerConnInfo *info = (ControllerConnInfo *)ctx;
     ControllerConnInfo *info = (ControllerConnInfo *)ctx;
     unsigned char cmd[1];
     unsigned char cmd[1];
+    ZZ expon;
 
 
     while(1) {
     while(1) {
 	size_t len = evbuffer_get_length(input);
 	size_t len = evbuffer_get_length(input);
@@ -335,6 +343,9 @@ static void controller_dpnode_reader(struct bufferevent *bev, void *ctx)
 			case 'L':
 			case 'L':
 			    info->state = CCSTATE_DPLISTENING;
 			    info->state = CCSTATE_DPLISTENING;
 			    break;
 			    break;
+			case 'E':
+			    info->state = CCSTATE_DPEXPON;
+			    break;
 			default:
 			default:
 			    /* Unknown DPnode command received */
 			    /* Unknown DPnode command received */
 			    fprintf(stderr, "Unknown command in "
 			    fprintf(stderr, "Unknown command in "
@@ -366,6 +377,30 @@ static void controller_dpnode_reader(struct bufferevent *bev, void *ctx)
 		info->state = CCSTATE_DPWAITRESP;
 		info->state = CCSTATE_DPWAITRESP;
 		break;
 		break;
 
 
+	    case CCSTATE_DPEXPON:
+		// Read the subproblemid and the answer to the subproblem
+		if (len < 2 + 3*sizeof(unsigned int)) return;
+		unsigned char exponbytes[2 + 3*sizeof(unsigned int)];
+		unsigned short problemid;
+		bufferevent_read(bev, exponbytes, 2 + 3*sizeof(unsigned int));
+		memmove(&problemid, exponbytes, 2);
+		ZZFromBytes(expon, exponbytes+2, 3*sizeof(unsigned int));
+		// Find the subproblem and check the answer
+		if (ctrlstate.dpnodes.working.count(bev) > 0) {
+		    SubproblemProgress *spp = ctrlstate.dpnodes.working[bev];
+		    if (spp->problemid == problemid &&
+			    spp->target ==
+				PowerMod(spp->base, expon, spp->modulus)) {
+			// Subproblem solved!
+			spp->solution = expon;
+			spp->solved = true;
+			spp->stop();
+			schedule();
+		    }
+		}
+		info->state = CCSTATE_DPWAITRESP;
+		break;
+
 	    case CCSTATE_END:
 	    case CCSTATE_END:
 		// Shut down the connection
 		// Shut down the connection
 		delete info;
 		delete info;
@@ -388,7 +423,7 @@ static void controller_dpnode_event_cb(struct bufferevent *bev, short events,
 	    SubproblemProgress *spp = ctrlstate.dpnodes.working[bev];
 	    SubproblemProgress *spp = ctrlstate.dpnodes.working[bev];
 	    ctrlstate.dpnodes.working.erase(bev);
 	    ctrlstate.dpnodes.working.erase(bev);
 	    spp->dpnodes.erase(bev);
 	    spp->dpnodes.erase(bev);
-	    spp->reset();
+	    spp->stop();
 	} else {
 	} else {
 	    ctrlstate.dpnodes.idle.erase(bev);
 	    ctrlstate.dpnodes.idle.erase(bev);
 	}
 	}
@@ -546,7 +581,8 @@ static vector<SubproblemProgress> decomp(const ZZ_p &base, const ZZ_p &target,
 	    ZZ f = (to_ZZ(10) << 32) / SqrRoot(order);
 	    ZZ f = (to_ZZ(10) << 32) / SqrRoot(order);
 	    dpfreq = trunc_long(f, 31);
 	    dpfreq = trunc_long(f, 31);
 	}
 	}
-	ret.push_back(SubproblemProgress(curproblemid++, rep(subgroup_base),
+	++curproblemid;
+	ret.push_back(SubproblemProgress(curproblemid, rep(subgroup_base),
 				    rep(subgroup_target),
 				    rep(subgroup_target),
 				    f.factor, order, dpfreq));
 				    f.factor, order, dpfreq));
     }
     }