-
Notifications
You must be signed in to change notification settings - Fork 1
/
Copy pathmesos.rs
199 lines (172 loc) · 7.26 KB
/
mesos.rs
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
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
extern crate async_mesos;
#[macro_use]
extern crate failure;
extern crate futures;
extern crate hyper;
#[macro_use]
extern crate log;
extern crate mime;
extern crate protobuf;
extern crate simple_logger;
extern crate spectral;
extern crate tokio_core;
extern crate users;
#[cfg(test)]
mod integration {
use failure;
use futures::{future, Future, Stream};
use hyper::Uri;
use async_mesos::client::Client;
use async_mesos::mesos;
use async_mesos::model;
use async_mesos::scheduler;
use simple_logger;
use spectral::prelude::*;
use tokio_core::reactor::Core;
use users::{get_current_uid, get_user_by_uid};
#[test]
fn connect() {
simple_logger::init();
let mut core = Core::new().expect("Could not create Core.");
let handle = core.handle();
// Mesos message
let user = get_user_by_uid(get_current_uid()).expect("No system user found.");
let mut framework_info = mesos::FrameworkInfo::new();
framework_info.set_user(String::from(user.name()));
framework_info.set_name(String::from("Example FOO Framework"));
// Create client
let uri = "http://localhost:5050/api/v1/scheduler"
.parse::<Uri>()
.expect("Could not parse Uri.");
let client = Client::connect(&handle, uri, framework_info);
let work = client
.into_stream()
.map(|(_, events)| events)
.flatten()
.map(|event| event.get_field_type())
.take(1)
.collect();
let result = core.run(work).unwrap();
assert_that(&result).is_equal_to(vec![scheduler::Event_Type::HEARTBEAT]);
}
#[test]
fn task_launch() {
simple_logger::init();
#[derive(Debug)]
pub struct State {
pub client: Client,
pub task_id: Option<mesos::TaskID>,
}
fn build_accept_call(
state: &State,
offer_id: mesos::OfferID,
task_id: mesos::TaskID,
agent_id: mesos::AgentID
) -> Result<scheduler::Call, failure::Error> {
let cpu = model::ScalarResourceBuilder::default()
.name("cpus")
.value(0.1)
.build()?;
let mem = model::ScalarResourceBuilder::default()
.name("mem")
.value(32.0)
.build()?;
let command = model::ShellCommandBuilder::default()
.command("sleep 100000")
.build()?;
let task_info = model::TaskInfoBuilder::default()
.name("sleep_task")
.task_id(task_id)
.agent_id(agent_id)
.resources(vec![cpu, mem])
.command(command)
.build()?;
let operation = model::OfferLaunchOperationBuilder::default()
.task_info(task_info)
.build()?;
let call = state.client.accept(vec![offer_id], vec![operation]);
Ok(call)
}
let mut core = Core::new().expect("Could not create Core.");
let handle = core.handle();
// Mesos message
let user = get_user_by_uid(get_current_uid()).expect("No system user found.");
let mut framework_info = mesos::FrameworkInfo::new();
framework_info.set_user(String::from(user.name()));
framework_info.set_name(String::from("Example Rust Framework"));
framework_info.set_role(String::from("*"));
// Create client
let uri = "http://localhost:5050/api/v1/scheduler"
.parse::<Uri>()
.expect("Could not parse Uri.");
let future_client = Client::connect(&handle, uri, framework_info);
// Process events and start and stop a simple task.
let work = future_client.and_then(|(client, events)| {
let state = State {
client: client,
task_id: None,
};
events.fold(
state,
|mut state, mut event| -> Box<Future<Item = State, Error = failure::Error>> {
match event.get_field_type() {
scheduler::Event_Type::OFFERS => {
info!("Received offer.");
// We already launched a task.
if state.task_id.is_some() {
info!("Ignoring offer because task is already launching.");
return Box::new(future::result(Ok(state)));
}
// Create task for offer.
let mut offer = event.take_offers().take_offers()[0].clone();
let offer_id = offer.take_id();
let agent_id = offer.take_agent_id();
let mut task_id = mesos::TaskID::new();
task_id.set_value(String::from("my_task"));
state.task_id = Some(task_id.clone());
// Make call
if let Ok(call) = build_accept_call(&state, offer_id, task_id, agent_id) {
let s = state.client.call(&handle, call).map(|()| state);
Box::new(s)
} else {
Box::new(future::err(format_err!("Could not construct offer accept call")))
}
}
scheduler::Event_Type::UPDATE => {
info!("Received task update.");
let status = event.take_update().take_status();
let task_id = status.get_task_id().clone();
let task_state = status.get_state().clone();
debug!(
"Task {} is {:?}: {}",
task_id.get_value(),
task_state,
status.get_message()
);
let ack_call = state.client.acknowledge(status);
// Fire and forget acknowledge call.
let s = state.client.call(&handle, ack_call).map_err(|error| {
error!("Could not make acknowledge request: {}", error);
()
});
handle.spawn(s);
if task_state == mesos::TaskState::TASK_RUNNING {
// Stop framework.
let call = state.client.teardown();
let s = state.client.call(&handle, call).map(|()| state);
return Box::new(s);
} else {
Box::new(future::result(Ok(state)))
}
}
other => {
debug!("Ignore event {:?}", other);
Box::new(future::result(Ok(state)))
}
}
},
)
});
core.run(work).unwrap();
}
}