001/*
002 * Licensed to the Apache Software Foundation (ASF) under one or more
003 * contributor license agreements. See the NOTICE file distributed with
004 * this work for additional information regarding copyright ownership.
005 * The ASF licenses this file to You under the Apache license, Version 2.0
006 * (the "License"); you may not use this file except in compliance with
007 * the License. You may obtain a copy of the License at
008 *
009 *      http://www.apache.org/licenses/LICENSE-2.0
010 *
011 * Unless required by applicable law or agreed to in writing, software
012 * distributed under the License is distributed on an "AS IS" BASIS,
013 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
014 * See the license for the specific language governing permissions and
015 * limitations under the license.
016 */
017package org.apache.logging.log4j.cassandra;
018
019import java.net.InetSocketAddress;
020import java.util.ArrayList;
021import java.util.Date;
022import java.util.List;
023
024import com.datastax.driver.core.BatchStatement;
025import com.datastax.driver.core.BoundStatement;
026import com.datastax.driver.core.Cluster;
027import com.datastax.driver.core.PreparedStatement;
028import com.datastax.driver.core.Session;
029import org.apache.logging.log4j.core.LogEvent;
030import org.apache.logging.log4j.core.appender.ManagerFactory;
031import org.apache.logging.log4j.core.appender.db.AbstractDatabaseManager;
032import org.apache.logging.log4j.core.appender.db.ColumnMapping;
033import org.apache.logging.log4j.core.config.plugins.convert.DateTypeConverter;
034import org.apache.logging.log4j.core.config.plugins.convert.TypeConverters;
035import org.apache.logging.log4j.core.net.SocketAddress;
036import org.apache.logging.log4j.spi.ThreadContextMap;
037import org.apache.logging.log4j.spi.ThreadContextStack;
038import org.apache.logging.log4j.util.ReadOnlyStringMap;
039import org.apache.logging.log4j.util.Strings;
040
041/**
042 * Manager for a Cassandra appender instance.
043 */
044public class CassandraManager extends AbstractDatabaseManager {
045
046    private static final int DEFAULT_PORT = 9042;
047
048    private final Cluster cluster;
049    private final String keyspace;
050    private final String insertQueryTemplate;
051    private final List<ColumnMapping> columnMappings;
052    private final BatchStatement batchStatement;
053    // re-usable argument binding array
054    private final Object[] values;
055
056    private Session session;
057    private PreparedStatement preparedStatement;
058
059    private CassandraManager(final String name, final int bufferSize, final Cluster cluster,
060                             final String keyspace, final String insertQueryTemplate,
061                             final List<ColumnMapping> columnMappings, final BatchStatement batchStatement) {
062        super(name, bufferSize);
063        this.cluster = cluster;
064        this.keyspace = keyspace;
065        this.insertQueryTemplate = insertQueryTemplate;
066        this.columnMappings = columnMappings;
067        this.batchStatement = batchStatement;
068        this.values = new Object[columnMappings.size()];
069    }
070
071    @Override
072    protected void startupInternal() throws Exception {
073        session = cluster.connect(keyspace);
074        preparedStatement = session.prepare(insertQueryTemplate);
075    }
076
077    @Override
078    protected boolean shutdownInternal() throws Exception {
079        session.close();
080        cluster.close();
081        return true;
082    }
083
084    @Override
085    protected void connectAndStart() {
086        // a Session automatically manages connections for us
087    }
088
089    @Override
090    protected void writeInternal(final LogEvent event) {
091        for (int i = 0; i < columnMappings.size(); i++) {
092            final ColumnMapping columnMapping = columnMappings.get(i);
093            if (ThreadContextMap.class.isAssignableFrom(columnMapping.getType())
094                || ReadOnlyStringMap.class.isAssignableFrom(columnMapping.getType())) {
095                values[i] = event.getContextData().toMap();
096            } else if (ThreadContextStack.class.isAssignableFrom(columnMapping.getType())) {
097                values[i] = event.getContextStack().asList();
098            } else if (Date.class.isAssignableFrom(columnMapping.getType())) {
099                values[i] = DateTypeConverter.fromMillis(event.getTimeMillis(), columnMapping.getType().asSubclass(Date.class));
100            } else {
101                values[i] = TypeConverters.convert(columnMapping.getLayout().toSerializable(event),
102                    columnMapping.getType(), null);
103            }
104        }
105        final BoundStatement boundStatement = preparedStatement.bind(values);
106        if (batchStatement == null) {
107            session.execute(boundStatement);
108        } else {
109            batchStatement.add(boundStatement);
110        }
111    }
112
113    @Override
114    protected boolean commitAndClose() {
115        if (batchStatement != null) {
116            session.execute(batchStatement);
117        }
118        return true;
119    }
120
121    public static CassandraManager getManager(final String name, final SocketAddress[] contactPoints,
122                                              final ColumnMapping[] columns, final boolean useTls,
123                                              final String clusterName, final String keyspace, final String table,
124                                              final String username, final String password,
125                                              final boolean useClockForTimestampGenerator, final int bufferSize,
126                                              final boolean batched, final BatchStatement.Type batchType) {
127        return getManager(name,
128            new FactoryData(contactPoints, columns, useTls, clusterName, keyspace, table, username, password,
129                useClockForTimestampGenerator, bufferSize, batched, batchType), CassandraManagerFactory.INSTANCE);
130    }
131
132    private static class CassandraManagerFactory implements ManagerFactory<CassandraManager, FactoryData> {
133
134        private static final CassandraManagerFactory INSTANCE = new CassandraManagerFactory();
135
136        @Override
137        public CassandraManager createManager(final String name, final FactoryData data) {
138            final Cluster.Builder builder = Cluster.builder()
139                .addContactPointsWithPorts(data.contactPoints)
140                .withClusterName(data.clusterName);
141            if (data.useTls) {
142                builder.withSSL();
143            }
144            if (Strings.isNotBlank(data.username)) {
145                builder.withCredentials(data.username, data.password);
146            }
147            if (data.useClockForTimestampGenerator) {
148                builder.withTimestampGenerator(new ClockTimestampGenerator());
149            }
150            final Cluster cluster = builder.build();
151
152            final StringBuilder sb = new StringBuilder("INSERT INTO ").append(data.table).append(" (");
153            for (final ColumnMapping column : data.columns) {
154                sb.append(column.getName()).append(',');
155            }
156            sb.setCharAt(sb.length() - 1, ')');
157            sb.append(" VALUES (");
158            final List<ColumnMapping> columnMappings = new ArrayList<>(data.columns.length);
159            for (final ColumnMapping column : data.columns) {
160                if (Strings.isNotEmpty(column.getLiteralValue())) {
161                    sb.append(column.getLiteralValue());
162                } else {
163                    sb.append('?');
164                    columnMappings.add(column);
165                }
166                sb.append(',');
167            }
168            sb.setCharAt(sb.length() - 1, ')');
169            final String insertQueryTemplate = sb.toString();
170            LOGGER.debug("Using CQL for appender {}: {}", name, insertQueryTemplate);
171            return new CassandraManager(name, data.getBufferSize(), cluster, data.keyspace, insertQueryTemplate,
172                columnMappings, data.batched ? new BatchStatement(data.batchType) : null);
173        }
174    }
175
176    private static class FactoryData extends AbstractFactoryData {
177        private final InetSocketAddress[] contactPoints;
178        private final ColumnMapping[] columns;
179        private final boolean useTls;
180        private final String clusterName;
181        private final String keyspace;
182        private final String table;
183        private final String username;
184        private final String password;
185        private final boolean useClockForTimestampGenerator;
186        private final boolean batched;
187        private final BatchStatement.Type batchType;
188
189        private FactoryData(final SocketAddress[] contactPoints, final ColumnMapping[] columns, final boolean useTls,
190                            final String clusterName, final String keyspace, final String table, final String username,
191                            final String password, final boolean useClockForTimestampGenerator, final int bufferSize,
192                            final boolean batched, final BatchStatement.Type batchType) {
193            super(bufferSize);
194            this.contactPoints = convertAndAddDefaultPorts(contactPoints);
195            this.columns = columns;
196            this.useTls = useTls;
197            this.clusterName = clusterName;
198            this.keyspace = keyspace;
199            this.table = table;
200            this.username = username;
201            this.password = password;
202            this.useClockForTimestampGenerator = useClockForTimestampGenerator;
203            this.batched = batched;
204            this.batchType = batchType;
205        }
206
207        private static InetSocketAddress[] convertAndAddDefaultPorts(final SocketAddress... socketAddresses) {
208            final InetSocketAddress[] inetSocketAddresses = new InetSocketAddress[socketAddresses.length];
209            for (int i = 0; i < inetSocketAddresses.length; i++) {
210                final SocketAddress socketAddress = socketAddresses[i];
211                inetSocketAddresses[i] = socketAddress.getPort() == 0
212                    ? new InetSocketAddress(socketAddress.getAddress(), DEFAULT_PORT)
213                    : socketAddress.getSocketAddress();
214            }
215            return inetSocketAddresses;
216        }
217    }
218}