Puedes copiar cada archivo desde aquí, o clonar el repositorio FC_CConcurrente completo:
git clone https://github.com/gilde-valeria/FC_CConcurrente.gitArchivos
Archivos
module-info.java
Programas_P5/unam.fc.concurrent.practica5/src/module-info.java
/**
*
*/
/**
*
*/
module unam.fc.concurrent.practica5 {
}
ExecuteSnapshot.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/ExecuteSnapshot.java
package unam.fc.concurrent.practica5;
// Programa 1: Programa que ejecuta el WFSnapshot, imprime el contenido de los snaps actualizados por update()
// Utiliza WFSnapshot<T>, el cual utiliza StampedSnap y StampedValue
//
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class ExecuteSnapshot {
public static void main(String[] args) {
int capacity = 4;
int init = -1;
WFSnapshot<Integer> snapshot = new WFSnapshot<Integer>(capacity,init);
ExecutorService executor = Executors.newFixedThreadPool(4);
for(int i = 0; i < 20; i++) {
final int ntask=i;
executor.execute(() -> snapshot.update(ntask));
}
executor.shutdown();
while(!executor.isTerminated()) {};
System.out.println("Snapshot Final");
StampedSnap<Integer>[] copy = (StampedSnap<Integer>[]) new StampedSnap[capacity];
copy = snapshot.collect();
for (int j = 0; j < capacity; j++) {
System.out.println(" Value: " + copy[j].value + " Owner: " +
copy[j].owner + " Stamp: " + copy[j].stamp + " Snap: ");
for (int k = 0; k < capacity; k++) {
System.out.println("Thread " + k + " has value: " + copy[j].getSnap(k));
}
}
}
}
ExecuteSnapshotQ.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/ExecuteSnapshotQ.java
package unam.fc.concurrent.practica5;
// Programa 2: Programa que ejecuta el WFSnapshotH, imprime el contenido de value del snapshot ssR
// En el ssR se guardan las respuestas y las vistas, las vistas son un scan() del ssI de invocaciones
// Utiliza WFSnapshotH<T>, el cual utiliza StampedSnapH y StampedValue
// El WFSnapshotH es diferente al WFSnapshot porque permite guardar todas las operaciones hechas en values y todos los snaps hechos en snap
// Para ello se modifica StampedSnapH (a diferencia de StampedSnap) y se le agregaron listas para values y snap
import java.util.Random;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class ExecuteSnapshotQ {
public static void main(String[] args) {
int capacity = 4;
String init = null;
ConcurrentLinkedQueue<Integer> queue = new ConcurrentLinkedQueue<>();// Implementacion linealizable de una cola no acotada
WFSnapshotH<String> snapshotI = new WFSnapshotH<String>(capacity,init);
WFSnapshotH<String> snapshotR = new WFSnapshotH<String>(capacity,init);
ExecutorService executor = Executors.newFixedThreadPool(4);
Random rand = new Random(); //Para crear num randoms
for(int i = 0; i < 10; i++) {
int numRand = rand.nextInt(2);//Si es 0 se ejecuta un enq(), si es 1 un deq()
final int ntask=i;//Los items integers que se agregan a la cola
executor.execute(new RunnableQ(ntask, queue,numRand,snapshotI,snapshotR));
}
executor.shutdown();
while(!executor.isTerminated()) {};
// Imprimir los snaps guardados en el update()
/*System.out.println("Snapshot Final in Invs");
StampedSnapH<String>[] copy = (StampedSnapH<String>[]) new StampedSnapH[capacity];
copyI = snapshotI.collect();
for (int j = 0; j < capacity; j++) {
System.out.println("\n Thread Owner: " + copy[j].owner + " Last Value: " + copy[j].value +
" Last Stamp: " + copy[j].stamp + " Number of Snaps: " + copy[j].snap.size());
for (int i = 0; i < copy[j].snap.size(); i++) {
System.out.println("The " + i + " snap has: " );
for (int k = 0; k < capacity; k++) {
System.out.println("Thread " + k + " has value: " + copy[j].getSnap(k,i)); }
}
}*/
StampedSnapH<String>[] copy = (StampedSnapH<String>[]) new StampedSnapH[capacity];
copy = snapshotR.collect();
System.out.println("Snapshot Final --- Values");
for (int j = 0; j < capacity; j++) {
System.out.println("\n Thread Owner: " + copy[j].owner +
" Last Stamp: " + copy[j].stamp + " Number of Snaps: " + copy[j].snap.size());
System.out.println("The values are: ");
for (int i = 0; i < copy[j].values.size(); i++) {
System.out.println(copy[j].values.get(i)); }
}
}
}
RunnableQ.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/RunnableQ.java
package unam.fc.concurrent.practica5;
import java.util.Arrays;
import java.util.concurrent.Callable;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.stream.Collectors;
import java.util.stream.Stream;
public class RunnableQ<T> implements Runnable{
private int item;
ConcurrentLinkedQueue<Integer> queue;
WFSnapshotH<String> snapshotI;
WFSnapshotH<String> snapshotR;
private int num;
public RunnableQ(int item, ConcurrentLinkedQueue<Integer> queue,
int num, WFSnapshotH<String> snapshotI, WFSnapshotH<String> snapshotR){
this.item = item;
this.queue = queue;
this.snapshotI = snapshotI;
this.snapshotR = snapshotR;
this.num = num;
}
@Override
public void run(){
try{
if (this.num == 0) {//Probabilidad de 50/100 de ejecutar enq
String string = String.format("enq( %s )", this.item);
snapshotI.update(string);//En el ssI se guardan las invocaciones
Boolean res = this.queue.add(this.item); // Aplico la operacion de la queue
Object[] result = (T[]) snapshotI.scan(); // Tomo la foto de las invs escritas
Stream<Object> stream = Arrays.stream(result);
String scan = stream.map( n -> n.toString() )
.collect( Collectors.joining( " + " ) );
snapshotR.update(scan + String.format(" | %s || ", res));// En el ssR se guardan las respuestas
}
if (this.num == 1) {//Probabilidad de 50/100 de ejecutar deq
snapshotI.update("deq()");
T res = (T) this.queue.poll(); // Regresa null o el valor entero
Object[] result = (T[]) snapshotI.scan();
Stream<Object> stream = Arrays.stream(result);
String scan = stream.map( n -> n.toString() )
.collect( Collectors.joining( " + " ) );
snapshotR.update(scan + String.format(" | %s || ", res));
}
}
catch(Exception e){
}
}
}
SimpleSnapshot.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/SimpleSnapshot.java
package unam.fc.concurrent.practica5;
/*
* SimpleSnapshot.java
*
* Created on January 19, 2006, 10:43 AM
*
* From "Multiprocessor Synchronization and Concurrent Data Structures",
* by Maurice Herlihy and Nir Shavit.
* Copyright 2006 Elsevier Inc. All rights reserved.
*/
import java.util.Arrays;
/**
*
* @author Maurice Herlihy
*/
public class SimpleSnapshot<T> implements Snapshot<T> {
private StampedValue<T>[] a_table; // array of atomic MRSW registers
public SimpleSnapshot(int capacity, T init) {
a_table = (StampedValue<T>[]) new StampedValue[capacity];
for (int i = 0; i < capacity; i++) {
a_table[i] = new StampedValue<T>(init);
}
}
public void update(T value) {
int me = ThreadID.get();
StampedValue<T> oldValue = a_table[me];
//oldValue.values.add(value); // Guardar el conjunto de valores hechos
StampedValue<T> newValue =
new StampedValue<T>((oldValue.stamp)+1, value);
a_table[me] = newValue;
}
private StampedValue<T>[] collect() {
StampedValue<T>[] copy = (StampedValue<T>[]) new StampedValue[a_table.length];
for (int j = 0; j < a_table.length; j++)
copy[j] = a_table[j];
return copy;
}
public T[] scan() {
StampedValue<T>[] oldCopy, newCopy;
oldCopy = collect();
collect: while (true) {
newCopy = collect();
if (! Arrays.equals(oldCopy, newCopy)) {
oldCopy = newCopy;
continue collect;
}
// clean collect
T[] result = (T[]) new Object[a_table.length];
for (int j = 0; j < a_table.length; j++)
result[j] = newCopy[j].value;
return result;
}
}
}
Snapshot.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/Snapshot.java
package unam.fc.concurrent.practica5;
/*
* Snapshot.java
*
* Created on January 19, 2006, 10:39 AM
*
* From "Multiprocessor Synchronization and Concurrent Data Structures",
* by Maurice Herlihy and Nir Shavit.
* Copyright 2006 Elsevier Inc. All rights reserved.
*/
/**
* Interface for snapshot implementations
* @author Maurice Herlihy
*/
public interface Snapshot<T> {
public void update(T v);
public T[] scan();
}
StampedSnap.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/StampedSnap.java
package unam.fc.concurrent.practica5;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Future;
/*
* StampedSnap.java
*
* Created on January 19, 2006, 9:03 PM
*
* From "Multiprocessor Synchronization and Concurrent Data Structures",
* by Maurice Herlihy and Nir Shavit.
* Copyright 2006 Elsevier Inc. All rights reserved.
*/
/**
* Labeled snapshot value.
* @param T object type
* @author Maurice Herlihy
*/
public class StampedSnap<T> extends StampedValue<T> {
public T[] snap;
public StampedSnap(T value) {
super(value);
snap = null;
}
/**
* Constructor.
* @param stamp timestamp
* @param value object value
* @param snap
*/
public StampedSnap(long stamp, T value, T[] snap) {
super(stamp, value);
this.snap = snap;
}
public T getSnap(int i) {
return this.snap[i];
}
}
StampedSnapH.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/StampedSnapH.java
package unam.fc.concurrent.practica5;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Future;
/*
* StampedSnap.java
*
* Created on January 19, 2006, 9:03 PM
*
* From "Multiprocessor Synchronization and Concurrent Data Structures",
* by Maurice Herlihy and Nir Shavit.
* Copyright 2006 Elsevier Inc. All rights reserved.
*/
/**
* Labeled snapshot value.
* @param T object type
* @author Maurice Herlihy
*/
public class StampedSnapH<T> extends StampedValue<T> {
//public T[] snap;
public List<T[]> snap = new ArrayList<T[]>();
public List<T> values = new ArrayList<T>();
public StampedSnapH(T value) {
super(value);
//snap = null;
}
/**
* Constructor.
* @param stamp timestamp
* @param value object value
* @param snap
*/
public StampedSnapH(long stamp, T value, List<T[]> s, List<T> v) {
super(stamp, value);
snap = s;
values = v;
}
public StampedSnapH(long stamp, T value) {
super(stamp, value);
}
public T getSnap(int i, int j) {
return this.snap.get(j)[i];
}
}
StampedValue.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/StampedValue.java
package unam.fc.concurrent.practica5;
/*
* StampedValue.java
*
* Created on January 19, 2006, 9:49 AM
*
* From "Multiprocessor Synchronization and Concurrent Data Structures",
* by Maurice Herlihy and Nir Shavit.
* Copyright 2006 Elsevier Inc. All rights reserved.
*/
/**
* stamps for atomic MRMW register
* @author Maurice Herlihy
*/
public class StampedValue<T> {
/**
* Counter value used for comparison.
*/
public long stamp;
/**
* Thread id of creator
*/
public int owner;
/**
* Register value.
*/
public T value;
/**
* Least value ever.
*/
public static final StampedValue MIN_VALUE = new StampedValue(null);
/**
* Constructor.
*/
public StampedValue(long stamp, T value) {
this.stamp = stamp;
this.owner = ThreadID.get();
this.value = value;
}
/**
* Constructor.
*/
public StampedValue(T value) {
this.stamp = 0;
this.value = value;
}
public static StampedValue max(StampedValue x, StampedValue y) {
if (x.stamp > y.stamp) {
return x;
} else if (x.owner > y.owner){
return x;
} else {
return y;
}
}
}
ThreadID.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/ThreadID.java
package unam.fc.concurrent.practica5;
/*
* ThreadID.java
*
* Created on January 11, 2006, 10:27 PM
*
* From "Multiprocessor Synchronization and Concurrent Data Structures",
* by Maurice Herlihy and Nir Shavit.
* Copyright 2006 Elsevier Inc. All rights reserved.
*/
/**
* Assigns unique contiguous ids to threads.
* @author Maurice Herlihyh
*/
public class ThreadID {
/**
* The next thread ID to be assigned
**/
private static volatile int nextID = 0;
/**
* My thread-local ID.
**/
private static ThreadLocalID threadID = new ThreadLocalID();
public static int get() {
return threadID.get();
}
public static void reset() {
nextID = 0;
}
private static class ThreadLocalID extends ThreadLocal<Integer> {
protected synchronized Integer initialValue() {
return nextID++;
}
}
}
WFSnapshot.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/WFSnapshot.java
package unam.fc.concurrent.practica5;
/*
* WFSnapshot.java
*
* Created on January 19, 2006, 9:20 PM
*
* From "Multiprocessor Synchronization and Concurrent Data Structures",
* by Maurice Herlihy and Nir Shavit.
* Copyright 2006 Elsevier Inc. All rights reserved.
*/
/**
* Wait-free snapshot.
* @author Maurice Herlihy
*/
public class WFSnapshot<T> implements Snapshot<T> {
private StampedSnap<T>[] a_table; // array of MRSW atomic registers
public WFSnapshot(int capacity, T init) {
a_table = (StampedSnap<T>[]) new StampedSnap[capacity];
for (int i = 0; i < a_table.length; i++) {
a_table[i] = new StampedSnap<T>(init);
}
}
public void update(T value) {
int me = ThreadID.get();
T[] snap = this.scan();
StampedSnap<T> oldValue = a_table[me];
StampedSnap<T> newValue =
new StampedSnap<T>(oldValue.stamp+1, value, snap);
a_table[me] = newValue;
}
public StampedSnap<T>[] collect() {
StampedSnap<T>[] copy = (StampedSnap<T>[]) new StampedSnap[a_table.length];
for (int j = 0; j < a_table.length; j++)
copy[j] = a_table[j];
return copy;
}
public T[] scan() {
StampedSnap<T>[] oldCopy;
StampedSnap<T>[] newCopy;
boolean[] moved = new boolean[a_table.length];
oldCopy = collect();
collect: while (true) {
newCopy = collect();
for (int j = 0; j < a_table.length; j++) {
// did any thread move?
if (oldCopy[j].stamp != newCopy[j].stamp) {
if (moved[j]) { // second move
return oldCopy[j].snap;
} else {
moved[j] = true;
oldCopy = newCopy;
continue collect;
}
}
}
// clean collect
T[] result = (T[]) new Object[a_table.length];
for (int j = 0; j < a_table.length; j++)
result[j] = newCopy[j].value;
return result;
}
}
}
WFSnapshotH.java
Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/WFSnapshotH.java
package unam.fc.concurrent.practica5;
public class WFSnapshotH<T> implements Snapshot<T> {
private StampedSnapH<T>[] a_table; // array of MRSW atomic registers
public WFSnapshotH(int capacity, T init) {
a_table = (StampedSnapH<T>[]) new StampedSnapH[capacity];
for (int i = 0; i < a_table.length; i++) {
a_table[i] = new StampedSnapH<T>(init);
}
}
public void update(T value) {
int me = ThreadID.get();
T[] snap1 = this.scan();
StampedSnapH<T> oldValue = a_table[me];
oldValue.snap.add(snap1); // Guardar el conjunto de snaps
oldValue.values.add(value); // Guardar el conjunto de valores hechos
StampedSnapH<T> newValue =
new StampedSnapH<T>(oldValue.stamp+1, value, oldValue.snap, oldValue.values);
a_table[me] = newValue;
}
public StampedSnapH<T>[] collect() {
StampedSnapH<T>[] copy = (StampedSnapH<T>[]) new StampedSnapH[a_table.length];
for (int j = 0; j < a_table.length; j++)
copy[j] = a_table[j];
return copy;
}
public T[] scan() {
StampedSnapH<T>[] oldCopy;
StampedSnapH<T>[] newCopy;
boolean[] moved = new boolean[a_table.length];
oldCopy = collect();
collect: while (true) {
newCopy = collect();
for (int j = 0; j < a_table.length; j++) {
// did any thread move?
if (oldCopy[j].stamp != newCopy[j].stamp) {
if (moved[j]) { // second move
return oldCopy[j].snap.getLast();
} else {
moved[j] = true;
oldCopy = newCopy;
continue collect;
}
}
}
// clean collect
T[] result = (T[]) new Object[a_table.length];
for (int j = 0; j < a_table.length; j++)
result[j] = (T) newCopy[j].values; // Regreso el conjunto de valores
return result;
}
}
}