Lines Matching refs:conn

38 static void replicator_connection_disconnect(struct replicator_connection *conn);
41 replicator_input_line(struct replicator_connection *conn, const char *line)
53 context = hash_table_lookup(conn->requests, POINTER_CAST(id));
58 hash_table_remove(conn->requests, POINTER_CAST(id));
59 conn->callback(line[0] == '+', context);
63 static void replicator_input(struct replicator_connection *conn)
67 switch (i_stream_read(conn->input)) {
71 replicator_connection_disconnect(conn);
75 replicator_connection_disconnect(conn);
79 while ((line = i_stream_next_line(conn->input)) != NULL)
80 (void)replicator_input_line(conn, line);
84 replicator_send_buf(struct replicator_connection *conn, buffer_t *buf)
101 if (o_stream_send(conn->output, data, len) < 0) {
102 replicator_connection_disconnect(conn);
109 static int replicator_output(struct replicator_connection *conn)
113 if (o_stream_flush(conn->output) < 0) {
114 replicator_connection_disconnect(conn);
119 if (o_stream_get_buffer_used_size(conn->output) > 0) {
120 o_stream_set_flush_pending(conn->output, TRUE);
124 if (conn->queue[p]->used > 0) {
125 if (!replicator_send_buf(conn, conn->queue[p]))
136 static void replicator_connection_connect(struct replicator_connection *conn)
141 if (conn->fd != -1)
144 if (conn->port == 0) {
145 fd = net_connect_unix(conn->path);
147 i_error("net_connect_unix(%s) failed: %m", conn->path);
149 for (n = 0; n < conn->ips_count; n++) {
150 unsigned int idx = conn->ip_idx;
152 conn->ip_idx = (conn->ip_idx + 1) % conn->ips_count;
153 fd = net_connect_ip(&conn->ips[idx], conn->port, NULL);
157 net_ip2addr(&conn->ips[idx]), conn->port);
162 if (conn->to == NULL) {
163 conn->to = timeout_add(REPLICATOR_RECONNECT_MSECS,
165 conn);
170 timeout_remove(&conn->to);
171 conn->fd = fd;
172 conn->io = io_add(fd, IO_READ, replicator_input, conn);
173 conn->input = i_stream_create_fd(fd, MAX_INBUF_SIZE);
174 conn->output = o_stream_create_fd(fd, (size_t)-1);
175 o_stream_set_no_error_handling(conn->output, TRUE);
176 o_stream_nsend_str(conn->output, REPLICATOR_HANDSHAKE);
177 o_stream_set_flush_callback(conn->output, replicator_output, conn);
180 static void replicator_abort_all_requests(struct replicator_connection *conn)
185 iter = hash_table_iterate_init(conn->requests);
186 while (hash_table_iterate(iter, conn->requests, &key, &value))
187 conn->callback(FALSE, value);
189 hash_table_clear(conn->requests, TRUE);
192 static void replicator_connection_disconnect(struct replicator_connection *conn)
194 if (conn->fd == -1)
197 replicator_abort_all_requests(conn);
198 io_remove(&conn->io);
199 i_stream_destroy(&conn->input);
200 o_stream_destroy(&conn->output);
201 net_disconnect(conn->fd);
202 conn->fd = -1;
207 struct replicator_connection *conn;
210 conn = i_new(struct replicator_connection, 1);
211 conn->fd = -1;
212 hash_table_create_direct(&conn->requests, default_pool, 0);
214 conn->queue[i] = buffer_create_dynamic(default_pool, 1024);
215 return conn;
222 struct replicator_connection *conn;
224 conn = replicator_connection_create();
225 conn->callback = callback;
226 conn->path = i_strdup(path);
227 return conn;
235 struct replicator_connection *conn;
237 conn = replicator_connection_create();
238 conn->callback = callback;
239 conn->ips = i_new(struct ip_addr, ips_count);
240 memcpy(conn->ips, ips, sizeof(*ips) * ips_count);
241 conn->ips_count = ips_count;
242 conn->port = port;
243 return conn;
248 struct replicator_connection *conn = *_conn;
252 replicator_connection_disconnect(conn);
255 buffer_free(&conn->queue[i]);
257 timeout_remove(&conn->to);
258 hash_table_destroy(&conn->requests);
259 i_free(conn);
263 replicator_send(struct replicator_connection *conn,
268 if (conn->fd != -1 &&
269 o_stream_get_buffer_used_size(conn->output) == 0) {
271 o_stream_nsend(conn->output, data, data_len);
272 } else if (conn->queue[priority]->used + data_len >=
277 buffer_append(conn->queue[priority], data, data_len);
278 if (conn->output != NULL)
279 o_stream_set_flush_pending(conn->output, TRUE);
283 void replicator_connection_notify(struct replicator_connection *conn,
289 replicator_connection_connect(conn);
304 replicator_send(conn, priority, t_strdup_printf(
309 void replicator_connection_notify_sync(struct replicator_connection *conn,
314 replicator_connection_connect(conn);
316 id = ++conn->request_id_counter;
318 hash_table_insert(conn->requests, POINTER_CAST(id), context);
321 replicator_send(conn, REPLICATION_PRIORITY_SYNC, t_strdup_printf(