-
Notifications
You must be signed in to change notification settings - Fork 6
Expand file tree
/
Copy pathIndex.ets
More file actions
147 lines (129 loc) · 3.83 KB
/
Copy pathIndex.ets
File metadata and controls
147 lines (129 loc) · 3.83 KB
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
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
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
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
import {
SynStart,
SynEnd,
wait,
notify,
SharedBoolean,
SharedString,
SharedNumber,
Syc,
isMainThread,
addFunc,
Runnable,
Thread,
} from './ThreadBridge';
export function sharedWash(runnable: Runnable) {
let archetype: Runnable;
if (runnable.className === 'Producer') {
archetype = new Producer(new CustomBlockingQueue(5));
} else if (runnable.className === 'Consumer') {
archetype = new Consumer(new CustomBlockingQueue(5));
} else {
archetype = new Thread();
}
addFunc(runnable, archetype);
runnable.run();
}
import { SynStart, SynEnd, wait, notify, Syc } from './ThreadBridge';
import { SharedNumber } from './SharedTypes';
class CustomBlockingQueue {
public syn: Syc = new Syc();
public static staticSyn: Syc = new Syc();
public className: string = 'CustomBlockingQueue';
private queue: number[];
private size = new SharedNumber(0);
private capacity = new SharedNumber();
private front = new SharedNumber(0);
private rear = new SharedNumber(-1);
private readonly synchronized_object: Syc = new Syc();
constructor(capacity: number) {
this.queue = new Array(capacity);
this.capacity.setValue(capacity);
}
public async put(item: number): Promise<void> {
SynStart(this.synchronized_object.syn);
while (this.size.getValue() === this.capacity.getValue()) {
await wait(this.synchronized_object.syn);
}
this.rear.setValue((this.rear.getValue() + 1) % this.capacity.getValue());
this.queue[this.rear.getValue()] = item;
console.log('producer: ' + item);
this.size.setValue(this.size.getValue() + 1);
notify(this.synchronized_object.syn);
SynEnd(this.synchronized_object.syn);
}
public async take(): Promise<number> {
SynStart(this.synchronized_object.syn);
while (this.size.getValue() === 0) {
await wait(this.synchronized_object.syn);
}
const item = this.queue[this.front.getValue()];
this.front.setValue((this.front.getValue() + 1) % this.capacity.getValue());
this.size.setValue(this.size.getValue() - 1);
console.log('consumer: ' + item);
notify(this.synchronized_object.syn);
SynEnd(this.synchronized_object.syn);
return item;
}
}
class Producer implements Runnable {
public syn: Syc = new Syc();
public static staticSyn: Syc = new Syc();
public className: string = 'Producer';
private queue: CustomBlockingQueue;
constructor(queue: CustomBlockingQueue) {
this.queue = queue;
}
run(): void {
for (let i = 0; i < 10; i++) {
try {
this.queue.put(i);
for (let j = 0; j < 10; j++) {}
} catch (e) {
if (e instanceof InterruptedException) {
Thread.currentThread().interrupt();
return;
}
}
}
}
}
class Consumer implements Runnable {
public syn: Syc = new Syc();
public static staticSyn: Syc = new Syc();
public className: string = 'Consumer';
private queue: CustomBlockingQueue;
constructor(queue: CustomBlockingQueue) {
this.queue = queue;
}
run(): void {
for (let i = 0; i < 10; i++) {
try {
const item = this.queue.take();
for (let j = 0; j < 5; j++) {}
} catch (e) {
if (e instanceof InterruptedException) {
Thread.currentThread().interrupt();
return;
}
}
}
}
}
class ProducerConsumerExample {
public syn: Syc = new Syc();
public static staticSyn: Syc = new Syc();
public className: string = 'ProducerConsumerExample';
static main(args: string[]): void {
const queue = new CustomBlockingQueue(5);
const producer = new Producer(queue);
const consumer = new Consumer(queue);
const producerThread = new Thread(producer);
const consumerThread = new Thread(consumer);
producerThread.start();
consumerThread.start();
}
}
if (isMainThread()) {
// You can put the entry of your code here to test.
}